diff --git a/.changeset/native-rpc-modules.md b/.changeset/native-rpc-modules.md new file mode 100644 index 00000000..c815a6cc --- /dev/null +++ b/.changeset/native-rpc-modules.md @@ -0,0 +1,5 @@ +--- +"@cloudflare/computer": patch +--- + +`WorkerJavaScriptBackend` passes `node:fs` and host module calls between the isolate and the Durable Object as real Workers RPC values instead of JSON text with a custom byte encoding. Byte arrays now count at their real size against `maxCapabilityBytes`, so a 900-byte write fits under a 1024-byte limit where it used to be rejected. Every call still goes through one host bridge that enforces call counts, concurrency, deadlines, and byte budgets, and it now rejects functions, RPC stubs, and cycles in a request before the host acts on it. diff --git a/docs/17_isolate_javascript.md b/docs/17_isolate_javascript.md index 8abfb7da..e06ae6de 100644 --- a/docs/17_isolate_javascript.md +++ b/docs/17_isolate_javascript.md @@ -82,7 +82,7 @@ Workspace parses the graph before loading the Worker, confines every durable pat The backend admits up to twenty-four executions at a time by default. A concurrent start past that ceiling fails with `EEXEC_BUSY` instead of creating an unbounded number of Dynamic Workers. Adjust `maxConcurrentExecutions` after measuring the Durable Object and Worker Loader limits for the deployment. -Each execution also bounds combined stdout and stderr output, active event subscribers, directory entries per read, concurrent and total capability calls, and cumulative capability request and response bytes. The corresponding `maxStdioBytes`, `maxExecutionSubscribers`, `maxDirectoryEntries`, and `max*Capability*` options may be lowered for public workloads. Directory reads apply their limit in SQLite before materializing rows. Requests are checked inside the isolate before Workers RPC and again by the host. +Each execution also bounds combined stdout and stderr output, active event subscribers, directory entries per read, concurrent and total capability calls, and cumulative capability request and response bytes. The corresponding `maxStdioBytes`, `maxExecutionSubscribers`, `maxDirectoryEntries`, and `max*Capability*` options may be lowered for public workloads. Directory reads apply their limit in SQLite before materializing rows. Requests are checked inside the isolate before Workers RPC and again by the host. Every capability call goes through one host bridge that enforces these limits. Values cross as real Workers RPC values, measured as UTF-8 bytes for strings and raw bytes for byte arrays, and anything that is not plain data, such as a function, an RPC stub, or a cycle, is rejected before the host acts on it. Completed execution records remain available for replay for sixty minutes by default. The backend also keeps at most 100 completed records. Configure these bounds with `retentionMs` and `maxRetainedExecutions`. Completed records leave the in-memory active set immediately; replay reads them from SQLite. @@ -227,7 +227,7 @@ Each function receives the arguments the isolate passed, as an array of JSON-com | `access` | The backend's `"read"` or `"read-write"` access. Check it before any write. | | `resolvePath(path, { allowMissing })` | Confines a caller path to the backend root and rejects symlinks. | -The arguments come from caller code, so parse them before use. A function may return a value or a promise. The result must be JSON-compatible, and the bridge checks it at runtime: `undefined` becomes `null` and `undefined` object fields are dropped, as with `JSON.stringify`. It fits within the same capability byte limits as every other host call. A function that ignores `signal` and never settles keeps the execution in its finalizing state. +The arguments come from caller code, so parse them before use. A function may return a value or a promise. Arguments and results cross the isolate boundary as real values through Workers RPC, not as encoded text. The result must be JSON-compatible plain data, and the bridge checks it at runtime. As in JSON, an `undefined` result or array item becomes `null` and an `undefined` object field is left out, in arguments and results alike. Byte arrays work for `node:fs` calls but not for host modules. It fits within the same capability byte limits as every other host call. A function that ignores `signal` and never settles keeps the execution in its finalizing state. Specifiers and the export names of an object are checked at construction. A factory's export names are checked when the backend connects and the factory runs. A module must export at least one function, and every export name must be a JavaScript identifier name other than `default` or `then`. A reserved word such as `delete` is allowed, and caller code renames it on import: `import { delete as remove } from "ws:files"`. Importing a name the module does not export fails when the module graph links, before any code runs. diff --git a/packages/computer/src/backends/worker-javascript/module-graph.ts b/packages/computer/src/backends/worker-javascript/module-graph.ts index 2e94ca65..08de6ac8 100644 --- a/packages/computer/src/backends/worker-javascript/module-graph.ts +++ b/packages/computer/src/backends/worker-javascript/module-graph.ts @@ -342,44 +342,31 @@ function capabilitiesModule(maxCapabilityBytes: number) { } return call(namespace, method, args); } + // Arguments and results cross as real values through Workers RPC. + // The host bridge measures and limits them; this early check only + // spares an obviously oversized request the round trip. export async function call(namespace, method, args) { if (!host) throw new Error("Workspace capabilities are not installed"); - const request = JSON.stringify(args.map(encode)); - if (new TextEncoder().encode(request).byteLength > ${maxCapabilityBytes}) { + if (approximateBytes(args) > ${maxCapabilityBytes}) { throw new Error(${JSON.stringify(requestTooLargeMessage)}); } - const raw = await host.call(namespace + "." + method, request); - const payload = JSON.parse(String(raw)); + const payload = await host.call(namespace + "." + method, args); if (payload.error !== undefined) { - const detail = typeof payload.error === "string" ? { message: payload.error } : payload.error; - const error = new Error(detail.message); - if (detail.code !== undefined) error.code = detail.code; - if (detail.path !== undefined) error.path = detail.path; + const error = new Error(payload.error.message); + if (payload.error.code !== undefined) error.code = payload.error.code; + if (payload.error.path !== undefined) error.path = payload.error.path; throw error; } - return decode(payload.result); + return payload.result; } - function wrap(type, fields) { - return { __workspace_codec__: { version: 1, type, ...fields } }; - } - function encode(value) { - if (value instanceof Uint8Array) return wrap("bytes", { data: Array.from(value) }); - if (Array.isArray(value)) return wrap("array", { items: value.map(encode) }); - if (value && typeof value === "object") return wrap("object", { entries: Object.entries(value).map(([key, child]) => [key, encode(child)]) }); - return value; - } - function decode(value) { - if (!value || typeof value !== "object" || Array.isArray(value)) return value; - if (Object.keys(value).length !== 1 || !("__workspace_codec__" in value)) throw new Error("Invalid Workspace codec envelope"); - const codec = value.__workspace_codec__; - if (!codec || codec.version !== 1) throw new Error("Invalid Workspace codec envelope"); - if (codec.type === "bytes") { - if (!Array.isArray(codec.data) || !codec.data.every((byte) => Number.isInteger(byte) && byte >= 0 && byte <= 255)) throw new Error("Invalid Workspace byte value"); - return new Uint8Array(codec.data); + function approximateBytes(value) { + if (typeof value === "string") return value.length; + if (value instanceof Uint8Array) return value.byteLength; + if (Array.isArray(value)) return value.reduce((total, item) => total + approximateBytes(item), 8); + if (value && typeof value === "object") { + return Object.entries(value).reduce((total, [key, item]) => total + key.length + approximateBytes(item), 8); } - if (codec.type === "array" && Array.isArray(codec.items)) return codec.items.map(decode); - if (codec.type === "object" && Array.isArray(codec.entries)) return Object.fromEntries(codec.entries.map(([key, child]) => [key, decode(child)])); - throw new Error("Invalid Workspace codec envelope"); + return 8; } `; } diff --git a/packages/computer/src/backends/worker-javascript/worker-javascript.test.ts b/packages/computer/src/backends/worker-javascript/worker-javascript.test.ts index e75b00fb..cae6ef76 100644 --- a/packages/computer/src/backends/worker-javascript/worker-javascript.test.ts +++ b/packages/computer/src/backends/worker-javascript/worker-javascript.test.ts @@ -557,7 +557,7 @@ describe("WorkerJavaScriptBackend", () => { attachOutput(readable: ReadableStream): Promise; }, ) { - void host.call("fs.writeFile", JSON.stringify(["/workspace/output.txt", "done"])); + void host.call("fs.writeFile", ["/workspace/output.txt", "done"]); return evaluateResult(host, 1); }, }; @@ -792,9 +792,9 @@ describe("WorkerJavaScriptBackend", () => { return { async evaluate( _input: unknown, - host: { call(name: string, args: string): Promise }, + host: { call(name: string, args: unknown[]): Promise }, ) { - await host.call("host/ws:test.run", JSON.stringify([])); + await host.call("host/ws:test.run", []); }, }; }, @@ -843,9 +843,9 @@ describe("WorkerJavaScriptBackend", () => { return { evaluate( _input: unknown, - host: { call(name: string, args: string): Promise }, + host: { call(name: string, args: unknown[]): Promise }, ) { - void host.call("fs.writeFile", JSON.stringify(["/workspace/output.txt", "done"])); + void host.call("fs.writeFile", ["/workspace/output.txt", "done"]); return new Promise(() => undefined); }, }; @@ -1120,7 +1120,7 @@ describe("WorkerJavaScriptBackend", () => { initializeSchema(db, () => 0); const fs = new WorkspaceFilesystem(db); await fs.mkdir("/workspace", { recursive: true }); - let response = ""; + let response: unknown; const backend = new WorkerJavaScriptBackend({ modules: { "ws:test": { run: async () => null } }, loader: { @@ -1130,9 +1130,9 @@ describe("WorkerJavaScriptBackend", () => { return { async evaluate( _input: unknown, - host: { call(name: string, args: string): Promise }, + host: { call(name: string, args: unknown[]): Promise }, ) { - response = await host.call("host/ws:test.toString", JSON.stringify([])); + response = await host.call("host/ws:test.toString", []); }, }; }, @@ -1151,7 +1151,7 @@ describe("WorkerJavaScriptBackend", () => { for await (const _event of execution.events) { // Drain the run so the host call settles. } - expect(JSON.parse(response)).toMatchObject({ + expect(response).toMatchObject({ error: { message: expect.stringContaining("Unknown Workspace host module call") }, }); await handle.close?.(); diff --git a/packages/computer/src/runtime/bridge.test.ts b/packages/computer/src/runtime/bridge.test.ts index 4c1a83b0..423f3145 100644 --- a/packages/computer/src/runtime/bridge.test.ts +++ b/packages/computer/src/runtime/bridge.test.ts @@ -1,10 +1,10 @@ import { describe, expect, it } from "vitest"; -import { WorkspaceRuntimeBridge } from "./bridge.js"; +import { type BridgeResponse, WorkspaceRuntimeBridge } from "./bridge.js"; import type { WorkspaceRuntimeCapability } from "./capability.js"; const encoder = new TextEncoder(); -const args = JSON.stringify(["value"]); +const args = ["value"]; function bridge(limits: { maxCalls?: number; @@ -17,8 +17,9 @@ function bridge(limits: { }); } -async function message(response: Promise) { - return (JSON.parse(await response) as { error?: { message?: string } }).error?.message; +async function message(response: Promise) { + const settled = await response; + return "error" in settled ? settled.error.message : undefined; } describe("WorkspaceRuntimeBridge cumulative limits", () => { @@ -32,7 +33,8 @@ describe("WorkspaceRuntimeBridge cumulative limits", () => { }); it("accepts requests at the cumulative byte boundary and rejects the next request", async () => { - const bytes = encoder.encode(args).byteLength; + // ["value"]: 8 for the array plus 5 UTF-8 bytes for the string. + const bytes = 8 + encoder.encode("value").byteLength; const target = bridge({ maxTotalRequestBytes: bytes * 2 }); await expect(message(target.call("host/ws:test.run", args))).resolves.toBeUndefined(); await expect(message(target.call("host/ws:test.run", args))).resolves.toBeUndefined(); @@ -42,8 +44,8 @@ describe("WorkspaceRuntimeBridge cumulative limits", () => { }); it("accepts responses at the cumulative byte boundary and rejects the next response", async () => { - const sample = await bridge({}).call("host/ws:test.run", args); - const bytes = encoder.encode(sample).byteLength; + // The host function returns "ok": 2 UTF-8 bytes. + const bytes = encoder.encode("ok").byteLength; const target = bridge({ maxTotalResponseBytes: bytes * 2 }); await expect(message(target.call("host/ws:test.run", args))).resolves.toBeUndefined(); await expect(message(target.call("host/ws:test.run", args))).resolves.toBeUndefined(); @@ -51,9 +53,52 @@ describe("WorkspaceRuntimeBridge cumulative limits", () => { `responses exceed ${bytes * 2} bytes`, ); }); + + it("counts error responses against the cumulative response budget", async () => { + const target = new WorkspaceRuntimeBridge({} as WorkspaceRuntimeCapability, { + maxTotalResponseBytes: 256, + hostModules: new Map([ + [ + "ws:test", + { + fail: async () => { + throw new Error("x".repeat(100)); + }, + }, + ], + ]), + }); + await expect(message(target.call("host/ws:test.fail", []))).resolves.toBe("x".repeat(100)); + await expect(message(target.call("host/ws:test.fail", []))).resolves.toContain( + "responses exceed 256 bytes", + ); + }); }); describe("WorkspaceRuntimeBridge host modules", () => { + function echo(values: unknown[]) { + return new WorkspaceRuntimeBridge({} as WorkspaceRuntimeCapability, { + hostModules: new Map([ + [ + "ws:test", + { + run: async () => values[0], + args: async (received) => received, + }, + ], + ]), + }); + } + + it("leaves undefined fields out and turns undefined array items into null", async () => { + await expect( + echo([{ value: 1, optional: undefined, list: [1, undefined] }]).call("host/ws:test.run", []), + ).resolves.toEqual({ result: { value: 1, list: [1, null] } }); + await expect( + echo([]).call("host/ws:test.args", [undefined, { a: undefined, b: 2 }]), + ).resolves.toEqual({ result: [null, { b: 2 }] }); + }); + it("calls a host function on its module", async () => { const functions = { async name() { @@ -66,10 +111,105 @@ describe("WorkspaceRuntimeBridge host modules", () => { const target = new WorkspaceRuntimeBridge({} as WorkspaceRuntimeCapability, { hostModules: new Map([["ws:test", functions]]), }); - await expect(target.call("host/ws:test.run", args)).resolves.toBe( - JSON.stringify({ result: "ok" }), + await expect(target.call("host/ws:test.run", [])).resolves.toEqual({ result: "ok" }); + }); +}); + +describe("WorkspaceRuntimeBridge values", () => { + function echoBridge(maxPayloadBytes?: number) { + return new WorkspaceRuntimeBridge({} as WorkspaceRuntimeCapability, { + ...(maxPayloadBytes === undefined ? {} : { maxPayloadBytes }), + hostModules: new Map([["ws:test", { run: async (values) => values[0] ?? null }]]), + }); + } + + it("passes plain values through without encoding them", async () => { + await expect( + echoBridge().call("host/ws:test.run", [{ nested: [1, "two", null, { three: true }] }]), + ).resolves.toEqual({ result: { nested: [1, "two", null, { three: true }] } }); + }); + + it.each([ + ["a function", () => 1], + ["a class instance", new Date(0)], + ])("rejects %s in a request before it reaches the host", async (_label, value) => { + await expect(message(echoBridge().call("host/ws:test.run", [value]))).resolves.toContain( + "plain data", ); }); + + it("rejects a cyclic request", async () => { + const cyclic: Record = {}; + cyclic.self = cyclic; + await expect(message(echoBridge().call("host/ws:test.run", [cyclic]))).resolves.toContain( + "acyclic", + ); + }); + + it("rejects a request of many empty values by the payload limit", async () => { + await expect( + message(echoBridge(256).call("host/ws:test.run", [new Array(300).fill("")])), + ).resolves.toContain("request exceeds 256 bytes"); + await expect( + message(echoBridge(256).call("host/ws:test.run", [new Array(40).fill({})])), + ).resolves.toContain("request exceeds 256 bytes"); + }); + + it("keeps an own __proto__ field in a host module result", async () => { + const value = JSON.parse('{"__proto__": {"a": 1}, "b": 2}') as unknown; + const response = await echoBridge().call("host/ws:test.run", [value]); + expect(Object.keys((response as { result: object }).result)).toEqual(["__proto__", "b"]); + }); + + it("keeps a path that fits beside a short error message", async () => { + const path = `/${"p".repeat(600)}`; + const target = new WorkspaceRuntimeBridge({} as WorkspaceRuntimeCapability, { + maxPayloadBytes: 1024, + hostModules: new Map([ + [ + "ws:test", + { + run: async () => { + throw Object.assign(new Error("ENOENT"), { code: "ENOENT", path }); + }, + }, + ], + ]), + }); + await expect(target.call("host/ws:test.run", [])).resolves.toEqual({ + error: { message: "ENOENT", code: "ENOENT", path }, + }); + }); + + it("keeps an error with a long path within the payload limit", async () => { + const path = `/${"p".repeat(900)}`; + const target = new WorkspaceRuntimeBridge({} as WorkspaceRuntimeCapability, { + maxPayloadBytes: 1024, + hostModules: new Map([ + [ + "ws:test", + { + run: async () => { + throw Object.assign(new Error(`ENOENT: no such file, open '${path}'`), { + code: "ENOENT", + path, + }); + }, + }, + ], + ]), + }); + const response = await target.call("host/ws:test.run", []); + expect(response).toMatchObject({ error: { code: "ENOENT" } }); + expect(response).not.toHaveProperty("error.path"); + expect(encoder.encode(JSON.stringify(response)).byteLength).toBeLessThanOrEqual(1024); + }); + + it("rejects a request over the payload limit by its UTF-8 size", async () => { + await expect( + message(echoBridge(256).call("host/ws:test.run", ["é".repeat(200)])), + ).resolves.toContain("request exceeds 256 bytes"); + }); }); describe("WorkspaceRuntimeBridge assertResult", () => { diff --git a/packages/computer/src/runtime/bridge.ts b/packages/computer/src/runtime/bridge.ts index 31d4a27c..4f6468eb 100644 --- a/packages/computer/src/runtime/bridge.ts +++ b/packages/computer/src/runtime/bridge.ts @@ -1,7 +1,12 @@ import { RpcTarget } from "cloudflare:workers"; +import { utf8Prefix } from "../text-truncation.js"; import { assertRuntimeValue, type WorkspaceRuntimeCapability } from "./capability.js"; -import type { WorkspaceModuleCallContext, WorkspaceModuleFunctions } from "./types.js"; +import type { + WorkspaceModuleCallContext, + WorkspaceModuleFunctions, + WorkspaceRuntimeValue, +} from "./types.js"; export class WorkspaceRuntimeBridge extends RpcTarget { readonly #capability: WorkspaceRuntimeCapability; @@ -14,7 +19,7 @@ export class WorkspaceRuntimeBridge extends RpcTarget { readonly #maxTotalResponseBytes: number; readonly #maxResultBytes: number; readonly #onAttachOutput?: (readable: ReadableStream) => Promise; - readonly #inFlight = new Set>(); + readonly #inFlight = new Set>(); readonly #abortControllers = new Set(); #cancelled = false; #callTimedOut = false; @@ -71,13 +76,22 @@ export class WorkspaceRuntimeBridge extends RpcTarget { } } - call(name: string, argsJson: string): Promise { - const requestBytes = new TextEncoder().encode(argsJson).byteLength; + // The one entry point for isolate code. Arguments and results cross + // as real values through Workers RPC; this method is the proxy that + // keeps the limits on them, so untrusted code cannot overload the + // Durable Object with calls, concurrency, time, or bytes. + call(name: string, args: unknown[]): Promise { const reject = (message: string) => - Promise.resolve(encodeBoundedError(new Error(message), this.#maxPayloadBytes)); + Promise.resolve(boundedError(new Error(message), this.#maxPayloadBytes)); if (this.#cancelled) return reject("Workspace execution is being cancelled."); - if (requestBytes > this.#maxPayloadBytes) { - return reject(`Workspace capability request exceeds ${this.#maxPayloadBytes} bytes.`); + if (typeof name !== "string" || !Array.isArray(args)) { + return reject("Workspace capability calls take a name and an argument list."); + } + let requestBytes: number; + try { + requestBytes = measureValue(args, this.#maxPayloadBytes, "request"); + } catch (error) { + return Promise.resolve(boundedError(error, this.#maxPayloadBytes)); } if (this.#inFlight.size >= this.#maxConcurrentCalls) { return reject( @@ -97,9 +111,7 @@ export class WorkspaceRuntimeBridge extends RpcTarget { const abort = new AbortController(); this.#abortControllers.add(abort); const deadline = Date.now() + this.#maxCallDurationMs; - const operation = encodeCall(async () => { - const encodedArgs = JSON.parse(argsJson) as unknown[]; - const args = encodedArgs.map(decodeBridgeValue); + const operation = respond(async () => { if (name.startsWith("host/")) { return this.#callHostModule(name, args, { signal: abort.signal, @@ -185,10 +197,13 @@ export class WorkspaceRuntimeBridge extends RpcTarget { this.#inFlight.delete(operation); this.#abortControllers.delete(abort); }); + // Errors count against the response budget too. The budget error + // itself is short and not counted, and `maxCalls` bounds how many + // the isolate can see. return call.then((response) => { - const bytes = new TextEncoder().encode(response).byteLength; + const bytes = "result" in response ? response.bytes : errorBytes(response); if (this.#responseBytes + bytes > this.#maxTotalResponseBytes) { - return encodeBoundedError( + return boundedError( new Error( `Workspace execution capability responses exceed ${this.#maxTotalResponseBytes} bytes.`, ), @@ -196,7 +211,7 @@ export class WorkspaceRuntimeBridge extends RpcTarget { ); } this.#responseBytes += bytes; - return response; + return "result" in response ? { result: response.result } : response; }); } @@ -228,113 +243,156 @@ export class WorkspaceRuntimeBridge extends RpcTarget { if (typeof fn !== "function") { throw new Error(`Unknown Workspace host module call ${JSON.stringify(name)}.`); } - assertBridgeValues(args); // Called on its module, so a method that uses `this` still works. - const result = (await fn.call(functions, args, context)) ?? null; - assertBridgeValues([result]); - return result; + const result = await fn.call(functions, hostArguments(args), context); + return hostValue(result ?? null); } } -function assertBridgeValues( - values: unknown[], -): asserts values is import("./types.js").WorkspaceRuntimeValue[] { +function hostArguments(args: unknown[]): WorkspaceRuntimeValue[] { + return args.map((arg) => hostValue(arg ?? null)); +} + +// Host module values follow JSON: an undefined object field is left +// out, and an undefined array item becomes null, as JSON.stringify +// does. Native RPC would otherwise carry undefined through, and the +// isolate's result check rejects it. +function hostValue(value: unknown): WorkspaceRuntimeValue { const seen = new Set(); - const visit = (value: unknown): void => { + const visit = (item: unknown): WorkspaceRuntimeValue => { if ( - value === null || - typeof value === "boolean" || - typeof value === "string" || - (typeof value === "number" && Number.isFinite(value)) - ) - return; - if (typeof value !== "object") throw new Error("Host module values must be JSON-compatible."); - if (seen.has(value)) throw new Error("Host module values must be acyclic."); - seen.add(value); - if (Array.isArray(value)) for (const item of value) visit(item); - else { - const prototype = Object.getPrototypeOf(value); + item === null || + typeof item === "boolean" || + typeof item === "string" || + (typeof item === "number" && Number.isFinite(item)) + ) { + return item; + } + if (typeof item !== "object") throw new Error("Host module values must be JSON-compatible."); + if (seen.has(item)) throw new Error("Host module values must be acyclic."); + seen.add(item); + let copy: WorkspaceRuntimeValue; + if (Array.isArray(item)) { + copy = item.map((child: unknown) => (child === undefined ? null : visit(child))); + } else { + const prototype = Object.getPrototypeOf(item); if (prototype !== Object.prototype && prototype !== null) { throw new Error("Host module values must contain only plain objects."); } - // An undefined field is absent, as in JSON. encodeBridgeValue drops it. - for (const item of Object.values(value as Record)) { - if (item !== undefined) visit(item); + const fields: Record = {}; + for (const [key, child] of Object.entries(item)) { + // defineProperty keeps an own `__proto__` key a field. + if (child !== undefined) { + Object.defineProperty(fields, key, { + value: visit(child), + enumerable: true, + writable: true, + configurable: true, + }); + } } + copy = fields; } - seen.delete(value); + seen.delete(item); + return copy; }; - for (const value of values) visit(value); + return visit(value); } function decodeBytes(value: unknown): string | Uint8Array { return value instanceof Uint8Array ? value : String(value); } -function encodeBridgeValue(value: unknown): unknown { - const wrap = (type: string, fields: Record) => ({ - __workspace_codec__: { version: 1, type, ...fields }, - }); - if (value instanceof Uint8Array) return wrap("bytes", { data: Array.from(value) }); - if (Array.isArray(value)) return wrap("array", { items: value.map(encodeBridgeValue) }); - if (value && typeof value === "object") { - return wrap("object", { - entries: Object.entries(value) - .filter(([, child]) => child !== undefined) - .map(([key, child]) => [key, encodeBridgeValue(child)]), - }); - } - return value; -} +/** What the bridge sends back for one call: a value, or a bounded error. */ +export type BridgeResponse = + | { readonly result: unknown } + | { + readonly error: { readonly message: string; readonly code?: string; readonly path?: string }; + }; -function decodeBridgeValue(value: unknown): unknown { - if (!value || typeof value !== "object" || Array.isArray(value)) return value; - const record = value as Record; - if (Object.keys(record).length !== 1 || !("__workspace_codec__" in record)) { - throw new Error("Invalid Workspace codec envelope."); - } - const codec = record.__workspace_codec__ as Record | null; - if (codec?.version !== 1) throw new Error("Invalid Workspace codec envelope."); - if (codec.type === "bytes") { - if (!isByteArray(codec.data)) throw new Error("Invalid Workspace byte value."); - return new Uint8Array(codec.data); - } - if (codec.type === "array" && Array.isArray(codec.items)) { - return codec.items.map(decodeBridgeValue); - } - if (codec.type === "object" && Array.isArray(codec.entries)) { - return Object.fromEntries( - codec.entries.map((entry) => { - if (!Array.isArray(entry) || entry.length !== 2 || typeof entry[0] !== "string") { - throw new Error("Invalid Workspace object entry."); - } - return [entry[0], decodeBridgeValue(entry[1])]; - }), - ); - } - throw new Error("Invalid Workspace codec envelope."); -} +// A response before the per-execution budget check, carrying its size. +type MeasuredResponse = + | { result: unknown; bytes: number } + | Extract; -function isByteArray(value: unknown): value is number[] { - return ( - Array.isArray(value) && - value.every((byte) => Number.isInteger(byte) && byte >= 0 && byte <= 255) - ); +const encoder = new TextEncoder(); +const MAX_RESPONSE_VALUES = 4096; + +// Measure plain data the way it costs the Durable Object: UTF-8 bytes +// of strings and keys, raw bytes of byte arrays, and a small fixed cost +// per scalar. Anything that is not plain data is rejected, including +// functions and RPC stubs, which Workers RPC would otherwise carry into +// the host as live callbacks, and cycles. +function measureValue(value: unknown, maxBytes: number, kind: "request" | "response"): number { + let bytes = 0; + let values = 0; + const seen = new Set(); + const add = (count: number) => { + bytes += count; + if (bytes > maxBytes) { + throw new Error(`Workspace capability ${kind} exceeds ${maxBytes} bytes.`); + } + }; + const visit = (item: unknown): void => { + values += 1; + if (kind === "response" && values > MAX_RESPONSE_VALUES) { + throw new Error("Workspace capability response has too many values."); + } + if ( + item === null || + item === undefined || + typeof item === "boolean" || + typeof item === "number" + ) { + add(8); + return; + } + // Every value costs at least one byte, so the byte limit also bounds + // how many values the host walks, even for empty strings and objects. + if (typeof item === "string") { + add(Math.max(1, encoder.encode(item).byteLength)); + return; + } + if (item instanceof Uint8Array) { + add(Math.max(1, item.byteLength)); + return; + } + if (typeof item !== "object") { + throw new Error(`Workspace capability ${kind} values must be plain data.`); + } + if (seen.has(item)) throw new Error(`Workspace capability ${kind} values must be acyclic.`); + seen.add(item); + if (Array.isArray(item)) { + add(8); + for (const child of item) visit(child); + } else { + const prototype = Object.getPrototypeOf(item); + if (prototype !== Object.prototype && prototype !== null) { + throw new Error(`Workspace capability ${kind} values must be plain data.`); + } + add(8); + for (const [key, child] of Object.entries(item)) { + add(Math.max(1, encoder.encode(key).byteLength)); + visit(child); + } + } + seen.delete(item); + }; + visit(value); + return bytes; } function withDeadline( - call: Promise, + call: Promise, timeoutMs: number, maxPayloadBytes: number, onTimeout: () => void, -): Promise { +): Promise { let timer: ReturnType | undefined; - const timeout = new Promise((resolve) => { + const timeout = new Promise((resolve) => { timer = setTimeout(() => { onTimeout(); - resolve( - encodeBoundedError(new Error("Workspace capability call timed out."), maxPayloadBytes), - ); + resolve(boundedError(new Error("Workspace capability call timed out."), maxPayloadBytes)); }, timeoutMs); }); return Promise.race([call, timeout]).finally(() => { @@ -342,75 +400,52 @@ function withDeadline( }); } -async function encodeCall(run: () => Promise, maxPayloadBytes: number) { +async function respond( + run: () => Promise, + maxPayloadBytes: number, +): Promise { try { const result = await run(); - assertResponseWithin(result, maxPayloadBytes); - const encoded = JSON.stringify({ result: encodeBridgeValue(result) }); - if (new TextEncoder().encode(encoded).byteLength > maxPayloadBytes) { - throw new Error(`Workspace capability response exceeds ${maxPayloadBytes} bytes.`); - } - return encoded; + return { result, bytes: measureValue(result, maxPayloadBytes, "response") }; } catch (error) { - return encodeBoundedError(error, maxPayloadBytes); + return boundedError(error, maxPayloadBytes); } } -function assertResponseWithin(value: unknown, maxBytes: number) { - let bytes = 0; - let nodes = 0; - const visit = (item: unknown): void => { - nodes += 1; - if (nodes > 4096) throw new Error("Workspace capability response has too many values."); - if (typeof item === "string") bytes += item.length * 3; - else if (item instanceof Uint8Array) bytes += item.byteLength * 4; - else if (typeof item === "number" || typeof item === "boolean" || item === null) bytes += 16; - else if (Array.isArray(item)) for (const child of item) visit(child); - else if (item && typeof item === "object") { - for (const [key, child] of Object.entries(item)) { - bytes += key.length * 3; - visit(child); - } - } - if (bytes > maxBytes) { - throw new Error(`Workspace capability response exceeds ${maxBytes} bytes.`); - } - }; - visit(value); -} +// Room for the envelope's keys and punctuation around the strings. +const ERROR_OVERHEAD_BYTES = 64; -function encodeBoundedError(error: unknown, maxPayloadBytes: number) { +// An error the isolate can rebuild, cut to fit the payload limit as a +// whole. `code` and `path` carry node:fs error details. Each is kept if +// it fits beside the message, or beside half the room when the message +// is long, and the message is cut to what is left. +function boundedError(error: unknown, maxPayloadBytes: number) { const value = error as { code?: unknown; path?: unknown }; const message = error instanceof Error ? error.message : String(error); - const detailed = JSON.stringify({ - error: { - message, - ...(typeof value?.code === "string" ? { code: value.code } : {}), - ...(typeof value?.path === "string" ? { path: value.path } : {}), - }, - }); - const encoder = new TextEncoder(); - if (encoder.encode(detailed).byteLength <= maxPayloadBytes) return detailed; - - let budget = Math.max(0, maxPayloadBytes - 40); - while (budget >= 0) { - const bounded = JSON.stringify({ error: { message: truncateUtf8(message, budget) } }); - if (encoder.encode(bounded).byteLength <= maxPayloadBytes) return bounded; - budget -= 1; + let room = Math.max(0, maxPayloadBytes - ERROR_OVERHEAD_BYTES); + const messageReserve = Math.min(encoder.encode(message).byteLength, room / 2); + const details: { code?: string; path?: string } = {}; + for (const key of ["code", "path"] as const) { + const detail = value?.[key]; + if (typeof detail !== "string") continue; + const bytes = encoder.encode(detail).byteLength; + if (bytes > room - messageReserve) continue; + details[key] = detail; + room -= bytes; } - return JSON.stringify({ error: { message: "Capability call failed" } }); + return { error: { message: truncateText(message, room), ...details } }; } -function truncateUtf8(value: string, maxBytes: number) { - const bytes = new TextEncoder().encode(value); - if (bytes.byteLength <= maxBytes) return value; - let prefix = bytes.slice(0, maxBytes); - while (prefix.byteLength > 0) { - try { - return new TextDecoder("utf-8", { fatal: true }).decode(prefix); - } catch { - prefix = prefix.slice(0, -1); - } - } - return ""; +function errorBytes(response: Extract): number { + const { message, code, path } = response.error; + return ( + ERROR_OVERHEAD_BYTES + + encoder.encode(message).byteLength + + encoder.encode(code ?? "").byteLength + + encoder.encode(path ?? "").byteLength + ); +} + +function truncateText(value: string, maxBytes: number) { + return utf8Prefix(value, maxBytes).text; } diff --git a/packages/computer/src/runtime/types.ts b/packages/computer/src/runtime/types.ts index 0c0a2398..82968cee 100644 --- a/packages/computer/src/runtime/types.ts +++ b/packages/computer/src/runtime/types.ts @@ -29,10 +29,10 @@ export interface WorkspaceModuleCallContext { * * `args` holds the arguments the isolate passed, decoded from the wire. * They come from untrusted code, so parse them before use. The function - * may return a value or a promise of one. The result must be - * JSON-compatible: the bridge checks it at runtime, treats `undefined` - * as `null`, and drops `undefined` object fields, the way - * `JSON.stringify` does. + * may return a value or a promise of one. Arguments and results cross + * the isolate boundary as real values through Workers RPC, not as + * encoded text. The result must be JSON-compatible plain data: the + * bridge checks it at runtime and treats `undefined` as `null`. */ export type WorkspaceModuleFunction = ( args: readonly WorkspaceRuntimeValue[], diff --git a/packages/computer/tests/script-runner.test.ts b/packages/computer/tests/script-runner.test.ts index b017a34f..cdda1a21 100644 --- a/packages/computer/tests/script-runner.test.ts +++ b/packages/computer/tests/script-runner.test.ts @@ -100,7 +100,30 @@ describe("WorkspaceRuntime", () => { }); }); - it("round-trips bytes and marker-shaped plain objects without codec collisions", async () => { + it("moves bytes through node:fs without inflating them", async () => { + // 900 bytes fits under this fixture's 1024-byte capability limit as + // raw bytes. Encoded as JSON numbers it would be about four times + // larger and rejected. + const response = await runtime({ + source: ` + import fs from "node:fs/promises"; + export default async () => { + await fs.writeFile("/workspace/blob.bin", new Uint8Array(900).fill(255)); + const back = await fs.readFile("/workspace/blob.bin"); + return { isBytes: back instanceof Uint8Array, length: back.byteLength, last: back[899] }; + }; + `, + cwd: "/workspace", + }); + const text = await response.text(); + expect(response.status, text).toBe(200); + expect(JSON.parse(text).result, text).toMatchObject({ + status: "completed", + value: { isBytes: true, length: 900, last: 255 }, + }); + }); + + it("round-trips bytes, and plain objects shaped like the old codec, unchanged", async () => { const response = await runtime({ source: ` import fs from "node:fs/promises";