Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
24 changes: 23 additions & 1 deletion .github/workflows/ci.yml
Original file line number Diff line number Diff line change
Expand Up @@ -23,12 +23,34 @@ jobs:
bun-version: latest
- run: bun install
- run: bun run typecheck
- run: bun run lint
- run: bun run check
- run: bun run build
- run: bun test
- run: bun run docs:build
- run: bunx pkg-pr-new publish './packages/*'

# Bun's test runner cannot reproduce celld's request lifecycle (a finished
# or aborted request's pending work is dropped), which is where shared
# runtimes and cross-request gates hang. Boot a real celld node on the
# action server built from source.
celld-smoke:
runs-on: ubuntu-latest
env:
CELLD_VERSION: v0.6.0
ESBUILD_VERSION: 0.25.10
steps:
- uses: actions/checkout@v4
- uses: oven-sh/setup-bun@v2
with:
bun-version: latest
- run: bun install
- name: Install celld and esbuild
run: |
curl -fsSL https://celld.dev/install.sh | sh
npm install --global "esbuild@${ESBUILD_VERSION}"
echo "$HOME/.local/bin" >> "$GITHUB_PATH"
- run: bun run smoke:celld

docs:
if: github.event_name == 'push' && github.ref == 'refs/heads/main'
needs: ci
Expand Down
2 changes: 1 addition & 1 deletion bun.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

3 changes: 1 addition & 2 deletions package.json
Original file line number Diff line number Diff line change
Expand Up @@ -11,9 +11,8 @@
"test": "bun run --filter './packages/*' test",
"build": "bun run --filter './packages/*' build",
"release": "bun scripts/release.ts",
"smoke:celld": "bun scripts/celld-smoke.ts",
"typecheck": "bun run --filter './packages/*' typecheck",
"lint": "oxlint",
"fmt": "oxfmt --write .",
"docs:dev": "bun run --filter website dev",
"docs:build": "bun run --filter website build",
"docs:verify": "bun run --filter website verify",
Expand Down
33 changes: 32 additions & 1 deletion packages/oxidejs/CHANGELOG.md
Original file line number Diff line number Diff line change
@@ -1,6 +1,37 @@
# Changelog

## 0.5.4
## Unreleased

### Security

- Concurrent actions could read another request's context. The Effect RPC runtime was built once and cached, and its server fibers kept the request context of the request that started them, so under concurrency `useRequest()` / `useEnv()` / `getRequestStore()` inside an action could return a different request (its headers, cookies, env) — measured 29 of 30 concurrent requests on Bun. Every action call now runs in its own runtime on every host

### Added

- `actions.timeout` (milliseconds): an action that has not answered in time is interrupted and answers with a JSON-RPC `-32603` error. Each call of a batch has its own deadline, so one slow call no longer fails the others. A stream counts as answered at its first frame. `createActionHandler` / `createWsHooks` take it as `timeoutMs`
- The generated client rejects a call 5 s after the server deadline (`actions.timeout`), for a hung network or proxy. `createClient` takes `timeout`, `retries` (default 2) and `maxReconnects` (default 8): an HTTP POST that failed in transit (network error, 502, 503, 504) is resent only when every call in it carries an `idempotencyKey`, and a WebSocket stream stops reconnecting after `maxReconnects`
- `createActionHandler({ maxBodyBytes })` refuses a larger body with HTTP 413; the worker preset passes `bodyLimit`, which before only applied to Node and dev
- `liveQuery({ mutateWaitMs })` (default 10000) and `liveQuery(...).close()`
- `bun run smoke:celld` (and a CI job) boots a celld node on the action server built from source and checks concurrency, request context, batch timeouts and recovery

### Changed

- A batch answers each call as soon as it ends instead of after the slowest one
- `disposeActionHandler` is a no-op: nothing is cached to dispose
- The sync request-store fallback (a module-global slot) is used on WebContainer only. Workers and celld keep the store in AsyncLocalStorage across `await`; the slot could hand one request's context to a concurrent one
- A queue consumer retries only the messages whose workflow failed to start and acks the rest, instead of retrying the whole batch
- Failed producer-side workflow starts (`producerStart`) are logged instead of swallowed

### Fixed

- celld: concurrent actions hung forever and then wedged the isolate until restart. celld's `node:process` answers every unknown `versions` key with a stub function, so `process.versions.webcontainer` read as truthy and every action queued behind one cross-request gate; a request the host cancelled never released it. WebContainer detection now requires a version string
- Worker hosts (Cloudflare Workers, celld) stalled on the cached runtime: the host drops a finished request's pending work
- `withRequestEntry` no longer chains Worker requests behind each other; only WebContainer serializes entry
- Two `queue.send()` calls, a `sendBatch()`, or a `workflow.start()` in one action with an idempotency key all reused that key as their id, so every message after the first was taken for the first and its workflow never started. Ids are now `key`, `key-1`, `key-2`, … per call (stable across retries of the same RPC call), valid Workflows instance ids trimmed to 100 chars
- Queue consumers and cron handlers run inside a request store, like workflow `run()`: `useEnv()`, `queue.send()` and `workflow.start()` no longer throw `request context is unavailable` there
- A `liveQuery` mutation waits at most `mutateWaitMs` for the previous one on its topic; a mutation whose request the host dropped no longer blocks the topic for good. A mutation that a newer one overtook during that wait does not publish its older snapshot, and a mutation that finishes after `close()` fails with `LiveQueryClosedError` instead of reporting success

## 0.5.5

### Added

Expand Down
10 changes: 6 additions & 4 deletions packages/oxidejs/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -160,7 +160,7 @@ Stream actions do not support `bind` / `with`. Breaking the `for await` loop or

### Live queries

Use `liveQuery` / `publish` for snapshot streams (query = subscription, mutation = publish). Hubs are isolate-local Effect PubSub — D1 (or your DB) stays the source of truth. Prefer Effect `Stream` + `mutateEffect`:
Use `liveQuery` / `publish` for snapshot streams (query = subscription, mutation = publish). Hubs are isolate-local Effect PubSub — a publish reaches subscribers in the same isolate only, not other Worker isolates or celld nodes — and D1 (or your DB) stays the source of truth. Mutations on one topic run one at a time; one waits at most `mutateWaitMs` (default 10000) for the previous, and a mutation overtaken during that wait does not publish its older snapshot. `close()` shuts a topic's hub down; a mutation that finishes after it fails with `LiveQueryClosedError`. Prefer Effect `Stream` + `mutateEffect`:

```ts
import { Effect, Stream } from "effect";
Expand Down Expand Up @@ -258,7 +258,7 @@ await invoices.send({ orderId: "…" });
await invoices.sendBatch([{ body: { orderId: "…" } }]);
```

Oxide merges `queues.producers` / `queues.consumers` via `mergeDurableBindings` and attaches a same-worker `queue` handler that starts the workflow from each message (via `createBatch` when available, otherwise duplicate-aware `create`/`get`). Cloudflare does not return message ids from `send`, so oxide wraps bodies in an envelope with a client-chosen id (`{ idempotencyKey }` / request header / UUID) and returns `{ id }` from `send` (and `{ ids }` from `sendBatch`) for `workflow.status` polling. Queue transport options (`contentType` / `delaySeconds`) travel in the RPC payload; `signal` / `idempotencyKey` stay on `CallOptions`.
Oxide merges `queues.producers` / `queues.consumers` via `mergeDurableBindings` and attaches a same-worker `queue` handler that starts the workflow from each message (via `createBatch` when available, otherwise duplicate-aware `create`/`get`). Cloudflare does not return message ids from `send`, so oxide wraps bodies in an envelope with a client-chosen id (`{ idempotencyKey }` / request header / UUID) and returns `{ id }` from `send` (and `{ ids }` from `sendBatch`) for `workflow.status` polling. Queue transport options (`contentType` / `delaySeconds`) travel in the RPC payload; `signal` / `idempotencyKey` stay on `CallOptions`. Every message gets its own id: with a key, the first send in a request uses it and later sends (and each message of a `sendBatch`) use `key-1`, `key-2`, … (valid Workflows instance ids, trimmed to 100 chars), so a retried request reproduces the same ids. The consumer acks the messages it started and retries only the ones that failed. Queue consumers and cron handlers run inside a request store, so `useEnv()`, `queue.send()` and `workflow.start()` work there.

```ts
const { id } = await invoices.send({ orderId: "…" });
Expand Down Expand Up @@ -352,7 +352,7 @@ export class ActionRoom {

The generated worker wrapper does not create a DO — hibernation is opt-in when you own the object.

`vite dev` and `rsbuild dev` serve the endpoint via middleware. `actions: "http"` (default) serves `/__oxide/action`; `actions: "ws"` uses a WebSocket instead (`crossws` on Node/`fetch`, `WebSocketPair` on `preset: "worker"`). `actions.sameOrigin` defaults to `true` for both transports; set it to `false` only when you intentionally accept cross-origin requests. Set `actions.path` to move the endpoint. Set `actions.openrpc: true` to serve `GET /__oxide/openrpc` — an [OpenRPC](https://spec.open-rpc.org/) 1.3 document for `action()` handlers only (Effect Schema → JSON Schema; workflows/queues/schedules are omitted). OpenRPC is HTTP-only and stays off when `transport` is `"ws"`. `actionHeaders` are static headers on the shared HTTP client and are ignored for WebSocket actions.
`vite dev` and `rsbuild dev` serve the endpoint via middleware. `actions: "http"` (default) serves `/__oxide/action`; `actions: "ws"` uses a WebSocket instead (`crossws` on Node/`fetch`, `WebSocketPair` on `preset: "worker"`). `actions.sameOrigin` defaults to `true` for both transports; set it to `false` only when you intentionally accept cross-origin requests. Set `actions.path` to move the endpoint. Set `actions.openrpc: true` to serve `GET /__oxide/openrpc` — an [OpenRPC](https://spec.open-rpc.org/) 1.3 document for `action()` handlers only (Effect Schema → JSON Schema; workflows/queues/schedules are omitted). OpenRPC is HTTP-only and stays off when `transport` is `"ws"`. `actionHeaders` are static headers on the shared HTTP client and are ignored for WebSocket actions. Set `actions.timeout` (milliseconds) to bound every action: an action that has not answered in time is interrupted, and the client gets a JSON-RPC `-32603` error instead of a request that never ends. Each call of a batch has its own deadline, and a batch answers each call as soon as it ends. A stream counts as answered at its first frame. The generated client gives up 5 s after the server deadline. `bodyLimit` also bounds action request bodies on the worker preset (HTTP 413).

## Rsbuild

Expand All @@ -376,7 +376,7 @@ Same factory as Vite: client stubs, `/__oxide/action`, and `dist/server.js`.
| `workerEntry` | `src/server.ts` | Relative to project root. Default path is skipped when missing (actions-only). Explicit path must exist. |
| `outDir` | `dist` | Output root (Node / `"fetch"`) |
| `clientDir` | `client` | Must stay inside `outDir` |
| `actions` | `"http"` | `"ws"` uses WebSocket (`crossws` on Node, `WebSocketPair` with `"worker"`); object form: `{ transport, path, sameOrigin, openrpc }` (`sameOrigin: true`, `openrpc: false`) |
| `actions` | `"http"` | `"ws"` uses WebSocket (`crossws` on Node, `WebSocketPair` with `"worker"`); object form: `{ transport, path, sameOrigin, openrpc, timeout }` (`sameOrigin: true`, `openrpc: false`, no `timeout`) |
| `actionHeaders` | — | Static headers on the HTTP client |
| `middleware` | `[]` | Fetch middleware. On WS: honor Responses (auth 302/401/…); ignore only `@ilha/router/ssr` document Responses. Otherwise Response short-circuits before actions / server entry. |
| `plugins` | `[]` | Build plugins (`beforeBuild` / `afterBuild`). Pass objects or module IDs. |
Expand Down Expand Up @@ -447,6 +447,8 @@ The generated `__asset` function uses `path.join` — not `path.resolve` — so
- Body size capped at 1 MB by default (enforced on the actual body, not just `Content-Length`).
- Batch requests capped at 20 items (both HTTP and WebSocket transports).
- Effect `Defect` / `Cause` payloads are scrubbed before they leave the endpoint. Clients see plain JSON-RPC errors (`code` + `message` only). Thrown messages become `Internal error` (`-32603`). Unknown methods → `-32601`; invalid params → `-32602`.
- Every action call runs in its own Effect RPC runtime, disposed when its response ends. A shared runtime would let a concurrent action read another request's context, and on Worker hosts (Cloudflare Workers, celld) it stalls, because the host drops a finished request's pending work.
- Worker code must not share a promise, timer or runtime across requests unless it is registered with `waitUntil` or its waiters are time-bounded: on celld, work a request started stops when that request answers or its client goes away. `bun run smoke:celld` checks the action server against a real celld node.
- `actions.sameOrigin` defaults to `true`. HTTP actions require `Origin` and/or `Sec-Fetch-Site` and reject cross-site callers. WebSocket upgrades additionally allow requests that omit both headers when `Host` matches the request URL host (celld and some proxies strip those headers; `localhost` / `127.0.0.1` / `::1` count as the same host).
- StackBlitz WebContainers do not keep `AsyncLocalStorage` across `async/await`. Oxide detects `process.versions.webcontainer` and falls back to a sync request store, capturing context before Effect schedules work and serializing handler entry so concurrent requests do not stomp that store. Stream pulls re-enter the captured store. This is a demo/dev workaround, not a concurrency model for production.

Expand Down
2 changes: 1 addition & 1 deletion packages/oxidejs/package.json
Original file line number Diff line number Diff line change
@@ -1,6 +1,6 @@
{
"name": "oxidejs",
"version": "0.5.5",
"version": "0.5.6",
"description": "Vite/Rsbuild plugin. One build → dist/server.js + optional client. Server actions via *.server.ts.",
"keywords": [
"cloudflare",
Expand Down
4 changes: 4 additions & 0 deletions packages/oxidejs/src/actions.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -371,6 +371,10 @@ describe("codegen", () => {
generateClientModule("http", { authorization: "Bearer x" })
).toContain('"authorization":"Bearer x"');
expect(generateClientModule("ws")).toContain('"transport":"ws"');
// The client deadline trails the server's by a grace period.
expect(
generateClientModule("http", undefined, "/__oxide/action", 15_000)
).toContain('"timeout":20000');
});

test("client stub peels { signal } and keeps other last args", async () => {
Expand Down
26 changes: 21 additions & 5 deletions packages/oxidejs/src/actions.ts
Original file line number Diff line number Diff line change
Expand Up @@ -270,16 +270,25 @@ export const assetRelPath = function assetRelPath(

interface ClientModuleOpts {
headers?: OxidejsActionHeaders;
timeout?: number;
transport?: "ws";
url: string;
}

/** Client grace over the server deadline: the server's own timeout error
* should win; the client bound catches a hung network or proxy. */
const CLIENT_TIMEOUT_GRACE_MS = 5000;

export const generateClientModule = function generateClientModule(
transport: OxidejsActionTransport = "http",
headers?: OxidejsActionHeaders,
actionPath: string = ACTION_PATH
actionPath: string = ACTION_PATH,
serverTimeout?: number
): string {
const opts: ClientModuleOpts = { url: actionPath };
if (serverTimeout !== undefined) {
opts.timeout = serverTimeout + CLIENT_TIMEOUT_GRACE_MS;
}
if (transport === "ws") {
opts.transport = "ws";
}
Expand Down Expand Up @@ -620,6 +629,7 @@ interface WorkerWrapperOpts {
actionOpenRpc?: boolean;
actionPath?: string;
actionSameOrigin?: boolean;
actionTimeout?: number | undefined;
actions?: OxidejsActionTransport;
bodyLimit?: number;
/** Worker preset (`preset: "worker"`). */
Expand Down Expand Up @@ -847,7 +857,9 @@ const buildActionImports = function buildActionImports(
ws: boolean,
actionPath: string,
sameOrigin: boolean,
openrpc: boolean
openrpc: boolean,
timeout: number | undefined,
bodyLimit: number
): string {
if (!hasActions) {
return openrpc
Expand All @@ -861,15 +873,17 @@ const __openRpcEntries = [];
import { actionsOpenRpcEntries as __openRpcEntries } from ${JSON.stringify(VIRTUAL_ACTIONS_ID)};
`
: "";
const timeoutOpt =
timeout === undefined ? "" : `, timeoutMs: ${JSON.stringify(timeout)}`;
if (ws) {
return `${openRpcImport}import { createWsHooks } from "oxidejs/rpc";
import { actionsGroup, actionsHandlers } from ${JSON.stringify(VIRTUAL_ACTIONS_ID)};
const __ws = createWsHooks(actionsGroup, actionsHandlers, { path: ${JSON.stringify(actionPath)}, sameOrigin: ${sameOrigin} });
const __ws = createWsHooks(actionsGroup, actionsHandlers, { path: ${JSON.stringify(actionPath)}, sameOrigin: ${sameOrigin}${timeoutOpt} });
`;
}
return `${openRpcImport}import { createActionHandler } from "oxidejs/rpc";
import { actionsGroup, actionsHandlers } from ${JSON.stringify(VIRTUAL_ACTIONS_ID)};
const __rpc = createActionHandler(actionsGroup, actionsHandlers, { path: ${JSON.stringify(actionPath)}, sameOrigin: ${sameOrigin}, createContext: (req) => req[__fetch] ?? {} });
const __rpc = createActionHandler(actionsGroup, actionsHandlers, { path: ${JSON.stringify(actionPath)}, sameOrigin: ${sameOrigin}${timeoutOpt}, maxBodyBytes: ${JSON.stringify(bodyLimit)}, createContext: (req) => req[__fetch] ?? {} });
`;
};

Expand Down Expand Up @@ -1081,7 +1095,9 @@ export const generateWorkerWrapper = function generateWorkerWrapper(
ws,
actionPath,
sameOrigin,
openrpc
openrpc,
opts.actionTimeout,
bodyLimit
);
const actionMatchFn = `const __actionMatch = (p) => p === ${JSON.stringify(actionPath)} || p === ${JSON.stringify(`${actionPath}/`)};`;
const openRpcGate = buildOpenRpcGate(openrpc, actionPath);
Expand Down
4 changes: 3 additions & 1 deletion packages/oxidejs/src/context.ts
Original file line number Diff line number Diff line change
Expand Up @@ -12,6 +12,8 @@ export {
__setNeedsSyncRequestStoreForTests,
getRequestStore,
inWebcontainer,
isWebcontainerVersions,
nextRequestScopedId,
needsSyncRequestStore,
peekRequestStore,
runWithRequest,
Expand Down Expand Up @@ -104,7 +106,7 @@ export type {
ServerActionHandle,
StreamActionHandle,
} from "./action";
export { liveQuery, publish } from "./live-query";
export { LiveQueryClosedError, liveQuery, publish } from "./live-query";
export type { LiveQuery, LiveQueryOptions } from "./live-query";
export { batch } from "./rpc/batch";
export type { BatchItem, BatchResult } from "./rpc/batch";
Expand Down
11 changes: 11 additions & 0 deletions packages/oxidejs/src/core.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -85,6 +85,17 @@ describe("resolveOptions", () => {
expect(resolved.actionSameOrigin).toBe(false);
});

test("resolves actions.timeout and rejects a non-positive one", () => {
expect(
resolveOptions({ actions: { timeout: 15_000 } }, process.cwd())
.actionTimeout
).toBe(15_000);
expect(resolveOptions({}, process.cwd()).actionTimeout).toBeUndefined();
expect(() =>
resolveOptions({ actions: { timeout: 0 } }, process.cwd())
).toThrow("actions.timeout must be a positive number");
});

test("resolves actions.openrpc", () => {
expect(
resolveOptions({ actions: { openrpc: true } }, process.cwd())
Expand Down
Loading
Loading