From ee8f4d5bf508610b046921b68e2b5b33e06158fd Mon Sep 17 00:00:00 2001 From: Matt Carey Date: Tue, 14 Jul 2026 11:56:23 +0100 Subject: [PATCH 1/7] dofs, rpc, workspace: Recover blob staging Keep blob metadata and bytes in one transaction so a receiver storage failure leaves the pull cursor pending. Cover retry behavior and repository-shaped generated output through the real wsd harness. --- packages/dofs/src/sync/blobs.test.ts | 23 +++++++++ packages/dofs/src/sync/blobs.ts | 24 ++++----- packages/rpc/src/sync-driver.test.ts | 49 +++++++++++++++++++ .../workspace/src/test-harness/shell.test.ts | 34 +++++++++++++ 4 files changed, 119 insertions(+), 11 deletions(-) diff --git a/packages/dofs/src/sync/blobs.test.ts b/packages/dofs/src/sync/blobs.test.ts index 924a9082..e9568596 100644 --- a/packages/dofs/src/sync/blobs.test.ts +++ b/packages/dofs/src/sync/blobs.test.ts @@ -3,6 +3,7 @@ import { describe, expect, it } from "vitest"; import { gc } from "../fs/gc.js"; import { withDB } from "../fs/with-db.js"; import { writeFile } from "../fs/writeFile.js"; +import { Database } from "../storage.js"; import { stageBlob } from "./blobs.js"; function sha256(bytes: Uint8Array): Uint8Array { @@ -29,6 +30,28 @@ describe("stageBlob", () => { }); }); + it("rolls back metadata and bytes when storing the bytes fails", async () => { + await withDB(async (db) => { + const bytes = new TextEncoder().encode("payload"); + const hash = sha256(bytes); + const failingDb = new Database({ + sql: { + exec: (query: string, ...bindings: unknown[]) => { + if (query.startsWith("INSERT INTO vfs_blob_bytes")) { + throw new Error("injected bytes failure"); + } + return db.sql.exec(query, ...bindings); + }, + }, + transactionSync: (closure) => db.transactionSync(closure), + }); + + expect(() => stageBlob(failingDb, hash, bytes, 1234)).toThrow("injected bytes failure"); + expect(db.scalar("SELECT COUNT(*) FROM vfs_blobs")).toBe(0); + expect(db.scalar("SELECT COUNT(*) FROM vfs_blob_bytes")).toBe(0); + }); + }); + it("is idempotent: a second call refreshes last_seen but leaves bytes alone", async () => { await withDB(async (db) => { const bytes = new TextEncoder().encode("same"); diff --git a/packages/dofs/src/sync/blobs.ts b/packages/dofs/src/sync/blobs.ts index 30c81603..0e2c812b 100644 --- a/packages/dofs/src/sync/blobs.ts +++ b/packages/dofs/src/sync/blobs.ts @@ -14,15 +14,17 @@ import type { Database } from "../storage.js"; // before calling. The function trusts the caller; a mismatched // pair would silently land under the wrong key. export function stageBlob(db: Database, hash: Uint8Array, bytes: Uint8Array, now: number): void { - db.run( - "INSERT INTO vfs_blobs (hash, size, last_seen) VALUES (?, ?, ?) ON CONFLICT(hash) DO UPDATE SET last_seen = excluded.last_seen", - hash, - bytes.byteLength, - now, - ); - db.run( - "INSERT INTO vfs_blob_bytes (hash, bytes) VALUES (?, ?) ON CONFLICT(hash) DO NOTHING", - hash, - bytes, - ); + db.transactionSync(() => { + db.run( + "INSERT INTO vfs_blobs (hash, size, last_seen) VALUES (?, ?, ?) ON CONFLICT(hash) DO UPDATE SET last_seen = excluded.last_seen", + hash, + bytes.byteLength, + now, + ); + db.run( + "INSERT INTO vfs_blob_bytes (hash, bytes) VALUES (?, ?) ON CONFLICT(hash) DO NOTHING", + hash, + bytes, + ); + }); } diff --git a/packages/rpc/src/sync-driver.test.ts b/packages/rpc/src/sync-driver.test.ts index ad78667b..11eeb260 100644 --- a/packages/rpc/src/sync-driver.test.ts +++ b/packages/rpc/src/sync-driver.test.ts @@ -991,6 +991,55 @@ describe("sync driver — streaming pullOnce", () => { }); }); +describe("sync driver — blob-stage recovery", () => { + it("keeps the pull pending and resumes after receiver storage recovers", async () => { + const upstream = makePeer(); + const receiver = makePeer(); + try { + const source = new SQLiteWorkspaceProvider(upstream.db, { now: () => 1 }); + source.mkdirSync("/repo/dist", { recursive: true }); + source.mkdirSync("/repo/node_modules/tiny", { recursive: true }); + source.writeFileSync("/repo/package.json", '{"name":"fixture"}\n'); + source.writeFileSync("/repo/dist/result.txt", "installed"); + source.writeFileSync("/repo/node_modules/tiny/index.js", "ignored"); + + let injected = false; + const failingDb = new Database({ + sql: { + exec: (query: string, ...bindings: unknown[]) => { + if (!injected && query.startsWith("INSERT INTO vfs_blob_bytes")) { + injected = true; + throw new Error("injected receiver blob storage failure"); + } + return receiver.db.sql.exec(query, ...bindings); + }, + }, + transactionSync: (closure) => receiver.db.transactionSync(closure), + }); + + await expect(pullOnce(failingDb, upstream.rpc)).rejects.toThrow( + "injected receiver blob storage failure", + ); + expect(readFetchCursor(receiver.db)).toEqual({ rev: 0, path: null }); + expect(receiver.db.scalar("SELECT COUNT(*) FROM vfs_blobs")).toBe(0); + expect(receiver.db.scalar("SELECT COUNT(*) FROM vfs_blob_bytes")).toBe(0); + + const resumed = await pullOnce(receiver.db, upstream.rpc); + const destination = new SQLiteWorkspaceProvider(receiver.db, { now: () => 1 }); + expect(resumed.applied).toBeGreaterThan(0); + expect(destination.readFileSync("/repo/dist/result.txt", "utf8")).toBe("installed"); + expect(destination.existsSync("/repo/node_modules")).toBe(false); + expect(readFetchCursor(receiver.db)).toEqual({ + rev: currentRev(upstream.db), + path: null, + }); + } finally { + upstream.close(); + receiver.close(); + } + }); +}); + describe("sync driver — push atomicity", () => { it("rolls back the entire batch when applyChanges fails mid-stream", async () => { // Construct a push with two file entries: one whose chunk bytes diff --git a/packages/workspace/src/test-harness/shell.test.ts b/packages/workspace/src/test-harness/shell.test.ts index 892cd162..14519892 100644 --- a/packages/workspace/src/test-harness/shell.test.ts +++ b/packages/workspace/src/test-harness/shell.test.ts @@ -118,10 +118,44 @@ describeIfDocker("Workspace.shell against a real wsd container", () => { it("exec result.pulled is 0 for a command that does not touch the VFS", async () => { await withWorkspace(url, async (ws) => { await ws.ready(); + await ws.pull(); const handle = await ws.shell.exec("echo cheap", { encoding: "utf8" }); const result = await handle.result(); expect(result.exitCode).toBe(0); expect(result.pulled).toBe(0); }); }); + + it("syncs repository outputs but ignores a generated dependency tree", async () => { + await withWorkspace(url, async (ws) => { + await ws.ready(); + await ws.pull(); + const root = `/workspace/repo-${crypto.randomUUID()}`; + await ws.fs.mkdir(`${root}/vendor/tiny`, { recursive: true }); + await ws.fs.writeFile( + `${root}/package.json`, + JSON.stringify({ + name: "sync-recovery-fixture", + private: true, + dependencies: { tiny: "file:vendor/tiny" }, + }), + ); + await ws.fs.writeFile(`${root}/vendor/tiny/index.js`, "installed\n"); + + const handle = await ws.shell.exec( + "mkdir -p node_modules/tiny dist && " + + "cp vendor/tiny/index.js node_modules/tiny/index.js && " + + "cat node_modules/tiny/index.js > dist/result.txt", + { cwd: root, encoding: "utf8" }, + ); + const result = await handle.result(); + + expect(result.exitCode, JSON.stringify(result)).toBe(0); + expect(result.pulled).toBeGreaterThan(0); + expect(await ws.fs.readFile(`${root}/dist/result.txt`, "utf8")).toBe("installed\n"); + await expect(ws.fs.stat(`${root}/node_modules`)).rejects.toMatchObject({ + code: "ENOENT", + }); + }); + }, 60_000); }); From 38711b157c2e128a400c02edd6e24145fdb62274 Mon Sep 17 00:00:00 2001 From: Matt Carey Date: Tue, 14 Jul 2026 12:05:27 +0100 Subject: [PATCH 2/7] dofs, rpc: Repair incomplete blobs Treat a blob as present only when its payload exists and matches the recorded size. Pull retries now fetch and restage incomplete or corrupt local objects before applying file entries and advancing the cursor. --- packages/dofs/src/sync/blobs.ts | 10 ++-- packages/dofs/src/sync/fetch.test.ts | 30 +++++++++++ packages/dofs/src/sync/fetch.ts | 13 +++-- packages/rpc/src/sync-driver.test.ts | 81 ++++++++++++++++++++++++++++ packages/rpc/src/sync-driver.ts | 9 ++-- 5 files changed, 128 insertions(+), 15 deletions(-) diff --git a/packages/dofs/src/sync/blobs.ts b/packages/dofs/src/sync/blobs.ts index 0e2c812b..d51fabe5 100644 --- a/packages/dofs/src/sync/blobs.ts +++ b/packages/dofs/src/sync/blobs.ts @@ -1,3 +1,4 @@ +import { clearBlobCache } from "../fs/blobCache.js"; import type { Database } from "../storage.js"; // Stage a chunk directly into vfs_blobs + vfs_blob_bytes without @@ -7,8 +8,8 @@ import type { Database } from "../storage.js"; // // Idempotent: a second call with the same hash refreshes // last_seen so the bytes don't get reaped by an interleaved gc. -// vfs_blob_bytes uses DO NOTHING on conflict — bytes are -// content-addressed so the stored value is always identical. +// Conflict updates also repair incomplete or size-mismatched rows +// left by an interrupted or corrupt write. // // Callers are expected to have verified that hash === sha256(bytes) // before calling. The function trusts the caller; a mismatched @@ -16,15 +17,16 @@ import type { Database } from "../storage.js"; export function stageBlob(db: Database, hash: Uint8Array, bytes: Uint8Array, now: number): void { db.transactionSync(() => { db.run( - "INSERT INTO vfs_blobs (hash, size, last_seen) VALUES (?, ?, ?) ON CONFLICT(hash) DO UPDATE SET last_seen = excluded.last_seen", + "INSERT INTO vfs_blobs (hash, size, last_seen) VALUES (?, ?, ?) ON CONFLICT(hash) DO UPDATE SET size = excluded.size, last_seen = excluded.last_seen", hash, bytes.byteLength, now, ); db.run( - "INSERT INTO vfs_blob_bytes (hash, bytes) VALUES (?, ?) ON CONFLICT(hash) DO NOTHING", + "INSERT INTO vfs_blob_bytes (hash, bytes) VALUES (?, ?) ON CONFLICT(hash) DO UPDATE SET bytes = excluded.bytes", hash, bytes, ); }); + clearBlobCache(db); } diff --git a/packages/dofs/src/sync/fetch.test.ts b/packages/dofs/src/sync/fetch.test.ts index 3aa6ab37..01b295e2 100644 --- a/packages/dofs/src/sync/fetch.test.ts +++ b/packages/dofs/src/sync/fetch.test.ts @@ -86,6 +86,36 @@ describe("hasObjects", () => { }); }); + it("treats blob metadata without bytes as missing", async () => { + await withDB(async (db) => { + await writeFile(db, "/a.txt", "alpha", {}, () => 1); + const entries = await drain(fetchChanges(db, 0)); + const known = ( + entries.find((e) => e.kind === "file") as { + chunks: { hash: Uint8Array }[]; + } + ).chunks[0].hash; + db.run("DELETE FROM vfs_blob_bytes WHERE hash = ?", known); + + expect(hasObjects(db, [known])).toEqual([]); + }); + }); + + it("treats blob bytes with the wrong length as missing", async () => { + await withDB(async (db) => { + await writeFile(db, "/a.txt", "alpha", {}, () => 1); + const entries = await drain(fetchChanges(db, 0)); + const known = ( + entries.find((e) => e.kind === "file") as { + chunks: { hash: Uint8Array }[]; + } + ).chunks[0].hash; + db.run("UPDATE vfs_blob_bytes SET bytes = ? WHERE hash = ?", new Uint8Array([1]), known); + + expect(hasObjects(db, [known])).toEqual([]); + }); + }); + it("returns an empty array when nothing matches", async () => { await withDB(async (db) => { const zero = new Uint8Array(32); diff --git a/packages/dofs/src/sync/fetch.ts b/packages/dofs/src/sync/fetch.ts index b372db42..317540fe 100644 --- a/packages/dofs/src/sync/fetch.ts +++ b/packages/dofs/src/sync/fetch.ts @@ -37,10 +37,9 @@ function toHex(bytes: Uint8Array): string { // index-backed lookups instead of one oversized statement. const PROBE_BATCH = 256; -// Subset-test the input hashes against vfs_blobs. Symmetric on both -// sides: the DO probes the container before pushObjects, and the -// container probes the DO before fetchObjects, so both sides ship -// only the bytes the receiver lacks. +// Subset-test the input hashes against complete local objects. A +// metadata row alone is not enough: the payload must exist and its +// byte length must match the advertised blob size. // // Matches the raw hash blobs through an IN (…) list so the lookup // rides the primary-key index. Present hashes are returned in input @@ -52,7 +51,11 @@ export function hasObjects(db: Database, hashes: Uint8Array[]): Uint8Array[] { const window = hashes.slice(i, i + PROBE_BATCH); const placeholders = window.map(() => "?").join(", "); const rows = db.all<{ hash: Uint8Array }>( - `SELECT hash FROM vfs_blobs WHERE hash IN (${placeholders})`, + `SELECT b.hash + FROM vfs_blobs b + JOIN vfs_blob_bytes bb ON bb.hash = b.hash + WHERE b.hash IN (${placeholders}) + AND length(bb.bytes) = b.size`, ...window, ); for (const row of rows) present.add(toHex(row.hash)); diff --git a/packages/rpc/src/sync-driver.test.ts b/packages/rpc/src/sync-driver.test.ts index 11eeb260..340884e1 100644 --- a/packages/rpc/src/sync-driver.test.ts +++ b/packages/rpc/src/sync-driver.test.ts @@ -141,6 +141,87 @@ describe("sync driver — pullOnce", () => { } }); + it("repairs blob metadata without bytes before applying the file", async () => { + const remote = makePeer(); + const local = makePeer(); + try { + const providerRemote = new SQLiteWorkspaceProvider(remote.db, { now: () => 1 }); + providerRemote.writeFileSync("/repair.txt", "recovered bytes"); + const blob = remote.db.one<{ hash: Uint8Array; size: number }>( + `SELECT c.hash, b.size + FROM vfs_chunks c + JOIN vfs_blobs b ON b.hash = c.hash + LIMIT 1`, + ); + expect(blob).toBeDefined(); + local.db.run( + "INSERT INTO vfs_blobs (hash, size, last_seen) VALUES (?, ?, ?)", + blob?.hash, + blob?.size, + 1, + ); + + let fetched: Uint8Array[] = []; + const wrapped = new Proxy(remote.rpc as object, { + get(target, prop, receiver) { + if (prop === "fetchObjects") { + return (hashes: Uint8Array[]) => { + fetched = hashes; + return Reflect.get(target, prop, receiver).call(target, hashes); + }; + } + return Reflect.get(target, prop, receiver); + }, + }) as SyncRPC; + + const applied = await pullOnce(local.db, wrapped); + const providerLocal = new SQLiteWorkspaceProvider(local.db, { now: () => 1 }); + expect(applied.applied).toBe(1); + expect(fetched).toEqual([blob?.hash]); + expect(providerLocal.readFileSync("/repair.txt", "utf8")).toBe("recovered bytes"); + expect(readFetchCursor(local.db)).toEqual({ rev: currentRev(remote.db), path: null }); + } finally { + remote.close(); + local.close(); + } + }); + + it("replaces corrupt blob bytes before applying the file", async () => { + const remote = makePeer(); + const local = makePeer(); + try { + const providerRemote = new SQLiteWorkspaceProvider(remote.db, { now: () => 1 }); + providerRemote.writeFileSync("/repair.txt", "correct bytes"); + const blob = remote.db.one<{ hash: Uint8Array; size: number }>( + `SELECT c.hash, b.size + FROM vfs_chunks c + JOIN vfs_blobs b ON b.hash = c.hash + LIMIT 1`, + ); + expect(blob).toBeDefined(); + local.db.run( + "INSERT INTO vfs_blobs (hash, size, last_seen) VALUES (?, ?, ?)", + blob?.hash, + blob?.size, + 1, + ); + local.db.run( + "INSERT INTO vfs_blob_bytes (hash, bytes) VALUES (?, ?)", + blob?.hash, + new TextEncoder().encode("bad"), + ); + + const applied = await pullOnce(local.db, remote.rpc); + const providerLocal = new SQLiteWorkspaceProvider(local.db, { now: () => 1 }); + expect(applied.applied).toBe(1); + expect(providerLocal.readFileSync("/repair.txt", "utf8")).toBe("correct bytes"); + expect(readFetchCursor(local.db)).toEqual({ rev: currentRev(remote.db), path: null }); + } finally { + remote.close(); + local.close(); + } + }); + it("is a no-op when fetchRev equals upstream currentRev", async () => { const a = makePeer(); const b = makePeer(); diff --git a/packages/rpc/src/sync-driver.ts b/packages/rpc/src/sync-driver.ts index 75028ba3..50033266 100644 --- a/packages/rpc/src/sync-driver.ts +++ b/packages/rpc/src/sync-driver.ts @@ -19,6 +19,7 @@ import { compareChangeCursors, currentRev, type Database, + hasObjects, readFetchCursor, readWatermark, type SkippedEntry, @@ -214,14 +215,10 @@ async function pullOnceImpl( const haveSubset = await remote.hasObjects(wantedHashes); const remoteHasLocally = new Set(); for (const h of haveSubset) remoteHasLocally.add(hex(h)); + const localHave = new Set(hasObjects(db, wantedHashes).map(hex)); const missing = wantedHashes.filter((h) => { const k = hex(h); - if (!remoteHasLocally.has(k)) return false; - const row = db.one<{ hash: Uint8Array }>( - "SELECT hash FROM vfs_blobs WHERE hash = ?", - h, - ); - return row === undefined; + return remoteHasLocally.has(k) && !localHave.has(k); }); if (missing.length > 0) { // Bare ReadableStream return — no envelope to dispose, From a0f1088f79e254d1c4ef422d72ecdb27fe02b9bb Mon Sep 17 00:00:00 2001 From: Matt Carey Date: Tue, 14 Jul 2026 12:16:44 +0100 Subject: [PATCH 3/7] workspace: Report pending exec sync Return command output and exit status when the post-command pull fails, while marking the structured sync outcome as pending. Bound and redact the reported error before exposing it through results and observer spans. --- docs/05_shell_interface.md | 26 ++++++-- packages/workspace/README.md | 11 +++- packages/workspace/src/index.ts | 1 + .../workspace/src/observe-integration.test.ts | 52 +++++++++++++++ packages/workspace/src/observe.test.ts | 14 ++++ packages/workspace/src/observe.ts | 27 ++++++-- packages/workspace/src/shell.test.ts | 66 ++++++++++++++++++- packages/workspace/src/shell.ts | 16 ++++- packages/workspace/src/stub.test.ts | 28 ++++++++ packages/workspace/src/stub.ts | 30 +++++++-- 10 files changed, 247 insertions(+), 24 deletions(-) diff --git a/docs/05_shell_interface.md b/docs/05_shell_interface.md index 0addbfd6..e20e78fb 100644 --- a/docs/05_shell_interface.md +++ b/docs/05_shell_interface.md @@ -101,12 +101,23 @@ type ExecEvent = | { id: string; seq: number; name: "stderr"; value: T } | { id: string; seq: number; name: "exit"; value: number }; +type ExecSyncResult = + | { status: "complete"; applied: number; skipped: SkippedEntry[] } + | { + status: "pending"; + applied: number; + skipped: SkippedEntry[]; + error: string; + }; + interface ExecResult { exitCode: number; stdout: T; stderr: T; pushed: number; // VFS changes uploaded before the command pulled: number; // VFS changes downloaded after the command + skipped: SkippedEntry[]; + sync: ExecSyncResult; } ``` @@ -139,8 +150,11 @@ const run = await workspace.shell.exec("zig build", { cwd: "/workspace", encoding: "utf8", }); -const { exitCode, stdout, stderr } = await run.result(); +const { exitCode, stdout, stderr, sync } = await run.result(); if (exitCode !== 0) throw new Error(stderr); +if (sync.status === "pending") { + console.warn(`command completed, but workspace sync is pending: ${sync.error}`); +} ``` Stream stdout as the command runs: @@ -283,9 +297,13 @@ no validation. Don't rely on the guard until it lands. - For read-write mounts, container-side writes under the mount root are mirrored back to the provider after the pull (provider first, then VFS). -- Failed pushes/pulls do not abort the command — `exec()` reports the - command's own exit code. Sync errors surface as thrown rejections - separately. +- Failed pushes/pulls do not change the command's exit code. If the + post-command pull fails, `result()` still returns the command's exit code, + stdout, and stderr. Its `sync` field has `status: "pending"`, zero applied + entries, an empty skipped list, and a bounded error string with common + credential forms redacted. A successful pull, including a clean no-op, has + `status: "complete"`. The existing `pushed`, `pulled`, and `skipped` fields + remain available for callers that only need counts. ## Wire format and backpressure diff --git a/packages/workspace/README.md b/packages/workspace/README.md index 0d0cc095..3d05d1af 100644 --- a/packages/workspace/README.md +++ b/packages/workspace/README.md @@ -246,7 +246,10 @@ export default { await ws.fs.writeFile("/notes.md", "hello"); using handle = await ws.shell.exec("ls /workspace", { encoding: "utf8" }); - const { exitCode, stdout } = await handle.result(); + const { exitCode, stdout, sync } = await handle.result(); + if (sync.status === "pending") { + console.warn("command completed before its filesystem changes synced", sync.error); + } return new Response(stdout, { status: exitCode === 0 ? 200 : 500 }); }, @@ -349,8 +352,10 @@ The span names the package emits today: `WorkspaceStub`. Contains `workspace.sync.push`, `workspace.shell.exec.spawn`, and `workspace.sync.pull` as nested children. Tagged with `workspace.shell.exit_code`, - `workspace.shell.pushed`, `workspace.shell.pulled`, and - `workspace.shell.skipped`. + `workspace.shell.pushed`, `workspace.shell.pulled`, + `workspace.shell.skipped`, and `workspace.shell.sync.status`. Pending + pulls also set `workspace.shell.sync.error` to the same bounded, + credential-redacted error returned in `ExecResult.sync`. - `workspace.fs.` — one per filesystem call routed through the stub (`readFile`, `writeFile`, `stat`, `readdir`, `find`, `ls`, `grep`, `mkdir`, `rm`). Tagged with `workspace.fs.path` and, where diff --git a/packages/workspace/src/index.ts b/packages/workspace/src/index.ts index 216cadb8..2ea9f554 100644 --- a/packages/workspace/src/index.ts +++ b/packages/workspace/src/index.ts @@ -61,6 +61,7 @@ export type { ExecHandle, ExecOptions, ExecResult, + ExecSyncResult, GetExecOptions, KillSignal, WorkspaceExecEvent, diff --git a/packages/workspace/src/observe-integration.test.ts b/packages/workspace/src/observe-integration.test.ts index 62295c38..25f6909f 100644 --- a/packages/workspace/src/observe-integration.test.ts +++ b/packages/workspace/src/observe-integration.test.ts @@ -290,6 +290,8 @@ describe("Workspace observer — shell stub", () => { expect(execSpan.attributes["workspace.shell.pushed"]).toBe(0); expect(execSpan.attributes["workspace.shell.pulled"]).toBe(0); expect(execSpan.attributes["workspace.shell.skipped"]).toBe(0); + expect(execSpan.attributes["workspace.shell.sync.status"]).toBe("complete"); + expect(execSpan.attributes["workspace.shell.sync.error"]).toBeUndefined(); // Nesting: the bracket runs push → spawn → pull inside the exec // span's callback, so all three appear as children on the recorder. @@ -298,6 +300,56 @@ describe("Workspace observer — shell stub", () => { expect(childNames).toContain("workspace.shell.exec.spawn"); expect(childNames).toContain("workspace.sync.pull"); }); + + it("reports pending sync on the exec span without exposing secrets", async () => { + const observer = makeRecorder(); + const secret = "observer-secret"; + const sync = fakeSync(); + sync.fetchChanges = async () => { + throw new Error(`WebSocket closed token=${secret}`); + }; + const shellRpc: import("@cloudflare/workspace-rpc").ShellRPC = { + async exec() { + return { + id: "exec-pending", + events: new ReadableStream({ + start(c) { + c.enqueue({ id: "exec-pending", seq: 0, name: "exit", value: 0 }); + c.close(); + }, + }), + }; + }, + async getExec() { + throw new Error("not used"); + }, + async killExec() {}, + async disposeExec() {}, + }; + const ws = new Workspace({ + storage: makeStorage(), + backends: [ + { + id: "shelled", + async connect() { + return { rpc: { sync, shell: shellRpc }, close: async () => {} }; + }, + }, + ], + observer, + }); + await ws.ready(); + using handle = await new WorkspaceShellStub(ws).exec("noop"); + const result = await handle.result(); + expect(result.sync.status).toBe("pending"); + + const execSpan = findSpan(observer.spans, "workspace.shell.exec"); + expect(execSpan.attributes["workspace.shell.sync.status"]).toBe("pending"); + expect(execSpan.attributes["workspace.shell.sync.error"]).toBe( + "WebSocket closed token=[REDACTED]", + ); + expect(JSON.stringify(execSpan.attributes)).not.toContain(secret); + }); }); function findSpan(spans: readonly RecordedSpan[], name: string): RecordedSpan { diff --git a/packages/workspace/src/observe.test.ts b/packages/workspace/src/observe.test.ts index ee98edaa..0e9e8fa0 100644 --- a/packages/workspace/src/observe.test.ts +++ b/packages/workspace/src/observe.test.ts @@ -56,6 +56,20 @@ describe("withSpan", () => { expect(recorder.spans[0].attributes).toEqual({ doubled: 14 }); }); + it("bounds and redacts error messages", async () => { + const recorder = makeRecorder(); + const secret = "trace-secret"; + await expect( + withSpan(recorder, "test.op", {}, async () => { + throw new Error(`failed token=${secret} ${"x".repeat(700)}`); + }), + ).rejects.toThrow(secret); + const recorded = recorder.spans[0].attributes["error.message"]; + expect(recorded.length).toBeLessThanOrEqual(512); + expect(recorded).toContain("failed token=[REDACTED]"); + expect(recorded).not.toContain(secret); + }); + it("records error name and message and re-throws", async () => { const recorder = makeRecorder(); await expect( diff --git a/packages/workspace/src/observe.ts b/packages/workspace/src/observe.ts index 7f0e81c0..b957ddb7 100644 --- a/packages/workspace/src/observe.ts +++ b/packages/workspace/src/observe.ts @@ -49,7 +49,9 @@ // ------ // If the wrapped work throws (or its returned promise rejects), the // span records `error.message` and `error.name` as attributes, then -// re-throws. The runtime decides what "failed span" means on its own +// re-throws. Error messages are bounded and common credential forms +// are redacted before they reach an observer. The runtime decides what +// "failed span" means on its own // side — we do not call any explicit `setStatus`-shaped API because the // Cloudflare surface does not expose one. @@ -182,7 +184,8 @@ export function withSpan( export type SpanOutcome = { ok: true; value: T } | { ok: false; error: unknown }; /** - * Records `error.message` and `error.name` on `span`. Adapters that + * Records `error.message` and `error.name` on `span`. Messages are + * bounded and common credential forms are redacted. Adapters that * want richer error reporting can subscribe to the underlying span * system directly; the workspace itself only forwards what the * Cloudflare `Span` surface accepts. @@ -190,8 +193,22 @@ export type SpanOutcome = { ok: true; value: T } | { ok: false; error: unknow function recordError(span: WorkspaceSpan, error: unknown): void { if (error instanceof Error) { span.setAttribute("error.name", error.name); - span.setAttribute("error.message", error.message); - } else { - span.setAttribute("error.message", String(error)); } + span.setAttribute("error.message", safeErrorMessage(error)); +} + +const MAX_SAFE_ERROR_LENGTH = 512; + +export function safeErrorMessage(error: unknown): string { + const message = error instanceof Error ? error.message : String(error); + return Array.from(message, (character) => + character < " " || character.charCodeAt(0) === 127 ? " " : character, + ) + .join("") + .replace( + /\b(authorization|token|api[_-]?key|password|secret|cookie)=([^\s&]+)/gi, + "$1=[REDACTED]", + ) + .replace(/\bBearer\s+[^\s,;]+/gi, "Bearer [REDACTED]") + .slice(0, MAX_SAFE_ERROR_LENGTH); } diff --git a/packages/workspace/src/shell.test.ts b/packages/workspace/src/shell.test.ts index 96a9de88..23d9d3f5 100644 --- a/packages/workspace/src/shell.test.ts +++ b/packages/workspace/src/shell.test.ts @@ -513,6 +513,14 @@ describe("WorkspaceShell.exec — push/pull bracket", () => { expect(result.pushed).toBe(5); expect(result.pulled).toBe(7); expect(result.skipped).toEqual([]); + expect(result.sync).toEqual({ status: "complete", applied: 7, skipped: [] }); + }); + + it("reports a clean no-op pull as complete", async () => { + const f = fakeRpc({ events: [exit(1, 0)] }); + const shell = new WorkspaceShell(f.rpc.shell, makeSync()); + const result = await (await shell.exec("true")).result(); + expect(result.sync).toEqual({ status: "complete", applied: 0, skipped: [] }); }); it("surfaces skipped read-only entries from the post-drain pull", async () => { @@ -588,23 +596,75 @@ describe("WorkspaceShell.exec — push/pull bracket", () => { expect(result.pulled).toBe(3); }); - it("falls back to pulled = 0 and skipped = [] when sync.pull() throws", async () => { + it("reports a pending sync after a Durable Object storage reset", async () => { + const f = fakeRpc({ + events: [stdout(1, "command output"), stderr(2, "command warning"), exit(3, 23)], + }); + const reset = "Internal error in Durable Object storage write caused object to be reset."; + const sync: Sync = { + async push() { + return 2; + }, + async pull() { + throw new Error(reset); + }, + }; + const shell = new WorkspaceShell(f.rpc.shell, sync); + const handle = await shell.exec("noop", { encoding: "utf8" }); + const result = await handle.result(); + expect(result).toMatchObject({ + exitCode: 23, + stdout: "command output", + stderr: "command warning", + pushed: 2, + pulled: 0, + skipped: [], + sync: { status: "pending", applied: 0, skipped: [], error: reset }, + }); + }); + + it("bounds and redacts pending sync errors", async () => { + const f = fakeRpc({ events: [exit(1, 0)] }); + const secret = "super-secret-value"; + const sync: Sync = { + async push() { + return 0; + }, + async pull() { + throw new Error(`transport failed token=${secret} ${"x".repeat(700)}`); + }, + }; + const shell = new WorkspaceShell(f.rpc.shell, sync); + const result = await (await shell.exec("noop")).result(); + expect(result.sync.status).toBe("pending"); + if (result.sync.status !== "pending") throw new Error("expected pending sync"); + expect(result.sync.error.length).toBeLessThanOrEqual(512); + expect(result.sync.error).toContain("transport failed token=[REDACTED]"); + expect(result.sync.error).not.toContain(secret); + }); + + it("reports a pending sync after an ordinary transport error", async () => { const f = fakeRpc({ events: [exit(1, 0)] }); const sync: Sync = { async push() { return 2; }, async pull() { - throw new Error("pull offline"); + throw new Error("WebSocket closed before pull completed"); }, }; const shell = new WorkspaceShell(f.rpc.shell, sync); const handle = await shell.exec("noop"); const result = await handle.result(); - expect(result.exitCode).toBe(0); expect(result.pushed).toBe(2); expect(result.pulled).toBe(0); expect(result.skipped).toEqual([]); + expect(result.sync).toEqual({ + status: "pending", + applied: 0, + skipped: [], + error: "WebSocket closed before pull completed", + }); }); }); diff --git a/packages/workspace/src/shell.ts b/packages/workspace/src/shell.ts index 3ae543e9..bc8e57f8 100644 --- a/packages/workspace/src/shell.ts +++ b/packages/workspace/src/shell.ts @@ -32,7 +32,7 @@ import type { ApplyResult, SkippedEntry } from "@cloudflare/dofs"; import type { ExecEvent, ShellRPC } from "@cloudflare/workspace-rpc"; -import { noopObserver, type WorkspaceObserver, withSpan } from "./observe.js"; +import { noopObserver, safeErrorMessage, type WorkspaceObserver, withSpan } from "./observe.js"; import { assertNotTemplate } from "./sh.js"; export type ExecEncoding = "utf8" | undefined; @@ -46,6 +46,10 @@ export type WorkspaceExecEvent = | { id: string; seq: number; name: "stderr"; value: Chunk } | { id: string; seq: number; name: "exit"; value: number }; +export type ExecSyncResult = + | { status: "complete"; applied: number; skipped: SkippedEntry[] } + | { status: "pending"; applied: number; skipped: SkippedEntry[]; error: string }; + export interface ExecResult { exitCode: number; stdout: Chunk; @@ -64,6 +68,9 @@ export interface ExecResult { pushed: number; pulled: number; skipped: SkippedEntry[]; + // Structured post-command sync outcome. The legacy pulled and + // skipped fields remain available for existing callers. + sync: ExecSyncResult; } export type KillSignal = "SIGTERM" | "SIGKILL" | "SIGINT" | "SIGHUP"; @@ -353,12 +360,14 @@ async function drainToResult( // values in that case. let pulled = 0; let skipped: SkippedEntry[] = []; + let syncResult: ExecSyncResult; try { const result = await sync.pull(); pulled = result.applied; skipped = result.skipped; - } catch { - // pulled / skipped stay empty + syncResult = { status: "complete", applied: pulled, skipped }; + } catch (error) { + syncResult = { status: "pending", applied: 0, skipped: [], error: safeErrorMessage(error) }; } return { exitCode, @@ -367,6 +376,7 @@ async function drainToResult( pushed, pulled, skipped, + sync: syncResult, }; } diff --git a/packages/workspace/src/stub.test.ts b/packages/workspace/src/stub.test.ts index e0af3eb7..094c83d6 100644 --- a/packages/workspace/src/stub.test.ts +++ b/packages/workspace/src/stub.test.ts @@ -337,6 +337,34 @@ describe("WorkspaceStub", () => { }); }); + it("shell result carries the structured sync outcome across Workers RPC", async () => { + using handle = new WorkspaceExecHandleStub( + Promise.resolve({ + exitCode: 0, + stdout: "done", + stderr: "", + pushed: 1, + pulled: 0, + skipped: [], + sync: { + status: "pending" as const, + applied: 0, + skipped: [], + error: "WebSocket closed before pull completed", + }, + }), + ); + await expect(handle.result()).resolves.toMatchObject({ + exitCode: 0, + sync: { + status: "pending", + applied: 0, + skipped: [], + error: "WebSocket closed before pull completed", + }, + }); + }); + it("shell.exec returns an eagerly-spawned handle", async () => { // The stub's exec() kicks off the underlying workspace.shell.exec // before returning, so the caller's first round trip already has diff --git a/packages/workspace/src/stub.ts b/packages/workspace/src/stub.ts index d79fde97..6154557c 100644 --- a/packages/workspace/src/stub.ts +++ b/packages/workspace/src/stub.ts @@ -71,7 +71,7 @@ import { encodeExecEvent } from "./exec-wire.js"; import type { GitCliInput, GitCliResult } from "./git/index.js"; import { withSpan } from "./observe.js"; import { assertNotTemplate } from "./sh.js"; -import type { ExecEncoding, ExecHandle, WorkspaceExecEvent } from "./shell.js"; +import type { ExecEncoding, ExecHandle, ExecSyncResult, WorkspaceExecEvent } from "./shell.js"; import type { Workspace } from "./workspace.js"; export interface WorkspaceExecOptions { @@ -90,6 +90,7 @@ export interface WorkspaceExecResult { exitCode: number; stdout: E extends "utf8" ? string : Uint8Array; stderr: E extends "utf8" ? string : Uint8Array; + sync: ExecSyncResult; } // Filesystem half. A direct proxy onto Workspace.fs — every @@ -260,6 +261,7 @@ interface ConsumeOutcome { pushed: number; pulled: number; skippedCount: number; + sync: ExecSyncResult; } // A deferred the exec span awaits. Resolving it (via result(), @@ -277,6 +279,16 @@ function makeConsumer(): Consumer { return { promise, resolve }; } +function emptyConsumeOutcome(exitCode = -1): ConsumeOutcome { + return { + exitCode, + pushed: 0, + pulled: 0, + skippedCount: 0, + sync: { status: "complete", applied: 0, skipped: [] }, + }; +} + // Exec handle returned from WorkspaceShellStub.exec. Holds the // underlying host ExecHandle and projects it across Workers RPC: // @@ -313,7 +325,7 @@ export class WorkspaceExecHandleStub e // stops the command rather than buffering forever. if (!this.#consumed) { this.#consumed = true; - this.#consumer.resolve({ exitCode: -1, pushed: 0, pulled: 0, skippedCount: 0 }); + this.#consumer.resolve(emptyConsumeOutcome()); this.#handle .then((handle) => handle.cancel?.()) .catch(() => { @@ -340,6 +352,7 @@ export class WorkspaceExecHandleStub e pushed: result.pushed, pulled: result.pulled, skippedCount: result.skipped.length, + sync: result.sync, }); // Let the span close before returning so its attributes are set. await this.#span.catch(() => {}); @@ -350,9 +363,10 @@ export class WorkspaceExecHandleStub e // WorkspaceExecResult shape. stdout: result.stdout as WorkspaceExecResult["stdout"], stderr: result.stderr as WorkspaceExecResult["stderr"], + sync: result.sync, }; } catch (error) { - this.#consumer.resolve({ exitCode: -1, pushed: 0, pulled: 0, skippedCount: 0 }); + this.#consumer.resolve(emptyConsumeOutcome()); throw error; } } @@ -379,18 +393,18 @@ export class WorkspaceExecHandleStub e reader.releaseLock(); reader = undefined; controller.close(); - consumer.resolve({ exitCode, pushed: 0, pulled: 0, skippedCount: 0 }); + consumer.resolve(emptyConsumeOutcome(exitCode)); return; } if (value.name === "exit") exitCode = value.value; controller.enqueue(encodeExecEvent(value as WorkspaceExecEvent)); } catch (error) { - consumer.resolve({ exitCode: -1, pushed: 0, pulled: 0, skippedCount: 0 }); + consumer.resolve(emptyConsumeOutcome()); controller.error(error); } }, cancel: async (reason) => { - consumer.resolve({ exitCode: -1, pushed: 0, pulled: 0, skippedCount: 0 }); + consumer.resolve(emptyConsumeOutcome()); await reader?.cancel(reason); reader?.releaseLock(); reader = undefined; @@ -629,6 +643,10 @@ export class WorkspaceShellStub extends RpcTarget { span.setAttribute("workspace.shell.pushed", outcome.value.pushed); span.setAttribute("workspace.shell.pulled", outcome.value.pulled); span.setAttribute("workspace.shell.skipped", outcome.value.skippedCount); + span.setAttribute("workspace.shell.sync.status", outcome.value.sync.status); + if (outcome.value.sync.status === "pending") { + span.setAttribute("workspace.shell.sync.error", outcome.value.sync.error); + } }, ); // If the spawn throws, the span rejects: forward the failure to From cd690a7becf931e37bb5cc13d679168bbc1b7323 Mon Sep 17 00:00:00 2001 From: Matt Carey Date: Tue, 14 Jul 2026 12:25:38 +0100 Subject: [PATCH 4/7] workspace: Add durable sync retries Let Durable Object hosts persist one pending pull intent per backend and wake Workspace through retryPendingSync. Failed attempts use bounded exponential backoff while exhausted intents remain available for inspection.\n\nCover scheduling, coalescing, cursor resume, convergence, exhaustion, and RPC cleanup with an in-memory scheduler. --- packages/workspace/README.md | 92 +++++++++ packages/workspace/src/index.ts | 9 +- packages/workspace/src/retry.test.ts | 284 +++++++++++++++++++++++++++ packages/workspace/src/shell.ts | 8 + packages/workspace/src/workspace.ts | 181 +++++++++++++++-- 5 files changed, 553 insertions(+), 21 deletions(-) create mode 100644 packages/workspace/src/retry.test.ts diff --git a/packages/workspace/README.md b/packages/workspace/README.md index 3d05d1af..426b92e6 100644 --- a/packages/workspace/README.md +++ b/packages/workspace/README.md @@ -234,6 +234,98 @@ A workspace with two backends that both write into [`docs/05_shell_interface.md`](../../docs/05_shell_interface.md) for the caveat. +## Durable pending-sync retries + +A command can finish after changing backend files while its post-command pull +fails. The command result reports `sync.status: "pending"`. If you configure a +`SyncRetryScheduler`, Workspace also writes one durable retry intent for that +backend. The library does not use an in-memory timer and cannot own the host +Durable Object's alarm. + +The scheduler contract contains only values that a host can persist: + +```ts +import { + type SyncRetryIntent, + type SyncRetryScheduler, + Workspace, +} from "@cloudflare/workspace"; + +const RETRY_PREFIX = "workspace:sync-retry:"; + +class DurableObjectRetryScheduler implements SyncRetryScheduler { + constructor(private readonly state: DurableObjectState) {} + + async get(backend: string): Promise { + return this.state.storage.get(`${RETRY_PREFIX}${backend}`); + } + + async schedule(intent: SyncRetryIntent): Promise { + // schedule replaces the backend's existing intent, so repeated + // failures coalesce instead of creating an alarm queue. + await this.state.storage.put(`${RETRY_PREFIX}${intent.backend}`, intent); + + const intents = await this.state.storage.list({ + prefix: RETRY_PREFIX, + }); + const next = Math.min(...[...intents.values()].map((item) => item.notBefore)); + await this.state.storage.setAlarm(next); + } + + async clear(backend: string): Promise { + await this.state.storage.delete(`${RETRY_PREFIX}${backend}`); + } +} +``` + +Pass the scheduler to `Workspace` and invoke `retryPendingSync` from the host's +alarm or another durable scheduler: + +```ts +export class WorkspaceHost extends DurableObject { + readonly scheduler = new DurableObjectRetryScheduler(this.ctx); + readonly workspace = new Workspace({ + storage: this.ctx.storage, + backends: [/* ... */], + retryScheduler: this.scheduler, + retry: { + initialDelayMs: 1_000, + maxDelayMs: 60_000, + maxAttempts: 5, + }, + }); + + async alarm(): Promise { + const intents = await this.ctx.storage.list({ + prefix: RETRY_PREFIX, + }); + const now = Date.now(); + for (const intent of intents.values()) { + if (intent.notBefore <= now) { + await this.workspace.retryPendingSync(intent.backend); + } + } + + const remaining = await this.ctx.storage.list({ + prefix: RETRY_PREFIX, + }); + if (remaining.size > 0) { + await this.ctx.storage.setAlarm( + Math.min(...[...remaining.values()].map((item) => item.notBefore)), + ); + } + } +} +``` + +`retryPendingSync(backend?)` enters the same per-backend FIFO as commands, +`push`, and `pull`. It resumes `pullOnce` from the cursor already persisted in +SQLite. Success clears the intent. Failure replaces it with the next bounded +exponential-backoff attempt. Once `maxAttempts` fails, the final intent stays +in storage and the method returns `status: "exhausted"`; the host can inspect, +alert on, or explicitly clear it. Calling the method with no pending intent +returns `status: "idle"`. + ## Worker-side consumption ```ts diff --git a/packages/workspace/src/index.ts b/packages/workspace/src/index.ts index 2ea9f554..cb373502 100644 --- a/packages/workspace/src/index.ts +++ b/packages/workspace/src/index.ts @@ -83,4 +83,11 @@ export { type WorkspaceStubHost, withWorkspace, } from "./with-workspace.js"; -export { Workspace, type WorkspaceOptions } from "./workspace.js"; +export { + type SyncRetryIntent, + type SyncRetryOptions, + type SyncRetryScheduler, + Workspace, + type WorkspaceOptions, + type WorkspaceRetryPendingSyncResult, +} from "./workspace.js"; diff --git a/packages/workspace/src/retry.test.ts b/packages/workspace/src/retry.test.ts new file mode 100644 index 00000000..a3680544 --- /dev/null +++ b/packages/workspace/src/retry.test.ts @@ -0,0 +1,284 @@ +import type { ChangeEntry } from "@cloudflare/dofs"; +import { SQLiteTestStorage } from "@cloudflare/dofs/testing"; +import { describe, expect, it } from "vitest"; + +import type { BackendHandle, WorkspaceBackend } from "./backend.js"; +import type { + SyncRetryIntent, + SyncRetryScheduler, + WorkspaceRetryPendingSyncResult, +} from "./workspace.js"; +import { Workspace } from "./workspace.js"; + +class MemoryRetryScheduler implements SyncRetryScheduler { + readonly intents = new Map(); + readonly scheduled: SyncRetryIntent[] = []; + readonly cleared: string[] = []; + + async get(backend: string): Promise { + return this.intents.get(backend); + } + + async schedule(intent: SyncRetryIntent): Promise { + this.intents.set(intent.backend, intent); + this.scheduled.push(intent); + } + + async clear(backend: string): Promise { + this.intents.delete(backend); + this.cleared.push(backend); + } +} + +function retryBackend(options: { + onExec(): void; + fetchChanges: import("@cloudflare/workspace-rpc").SyncRPC["fetchChanges"]; + close?: () => Promise; +}): WorkspaceBackend { + const sync: import("@cloudflare/workspace-rpc").SyncRPC = { + async push(input) { + return { rev: 0, appliedPushCursor: { rev: input.senderRev, path: null } }; + }, + fetchChanges: options.fetchChanges, + async readEntry() { + return null; + }, + async hasObjects(hashes) { + return hashes; + }, + fetchObjects() { + return new ReadableStream({ start: (controller) => controller.close() }); + }, + async watermarks() { + return { currentRev: 0, pushRev: 0, fetchCursor: { rev: 0, path: null } }; + }, + async pushObjects() {}, + }; + return { + id: "sandbox", + type: "fake", + async connect(): Promise { + return { + rpc: { + sync, + shell: { + async exec() { + options.onExec(); + return { + id: "command-1", + events: new ReadableStream({ + start(controller) { + controller.enqueue({ id: "command-1", seq: 1, name: "exit", value: 0 }); + controller.close(); + }, + }), + }; + }, + async getExec() { + throw new Error("not used"); + }, + async killExec() {}, + async disposeExec() {}, + }, + }, + close: options.close ?? (async () => {}), + }; + }, + }; +} + +async function runCommand(ws: Workspace): Promise { + const handle = await ws.shell.exec("build", { encoding: "utf8" }); + const result = await handle.result(); + expect(result.sync.status).toBe("pending"); + return undefined; +} + +describe("Workspace durable pending-sync retries", () => { + it("schedules the exact durable retry intent after a post-command pull failure", async () => { + const scheduler = new MemoryRetryScheduler(); + let execs = 0; + const backend = retryBackend({ + onExec: () => execs++, + async fetchChanges() { + throw new Error("backend unavailable"); + }, + }); + const ws = new Workspace({ + storage: new SQLiteTestStorage(), + backends: [backend], + retryScheduler: scheduler, + retry: { initialDelayMs: 2_000, maxDelayMs: 30_000, maxAttempts: 4 }, + now: () => 10_000, + }); + + await runCommand(ws); + + expect(execs).toBe(1); + expect(scheduler.scheduled).toEqual([{ backend: "sandbox", attempt: 1, notBefore: 12_000 }]); + }); + + it("resumes a partial batch from the persisted cursor and converges without rerunning the command", async () => { + const scheduler = new MemoryRetryScheduler(); + const after: Array<{ rev: number; path: string | null } | undefined> = []; + let fetches = 0; + let execs = 0; + const entries = Array.from( + { length: 257 }, + (_, index): ChangeEntry => ({ + kind: "delete", + rev: 1, + path: `/generated/${index.toString().padStart(3, "0")}`, + mtime: 1, + }), + ); + const backend = retryBackend({ + onExec: () => execs++, + async fetchChanges(input) { + after.push(input.after); + fetches++; + const remaining = entries.filter((entry) => { + if (!input.after || input.after.rev < entry.rev) return true; + return ( + input.after.rev === entry.rev && + input.after.path !== null && + entry.path > input.after.path + ); + }); + return { + currentCursor: { rev: 1, path: null }, + appliedPushCursor: { rev: 0, path: null }, + stream: + fetches === 1 + ? new ReadableStream({ + pull(controller) { + const entry = remaining.shift(); + if (entry !== undefined) { + controller.enqueue(entry); + return; + } + controller.error(new Error("lost after first batch")); + }, + }) + : new ReadableStream({ + start(controller) { + for (const entry of remaining) controller.enqueue(entry); + controller.close(); + }, + }), + }; + }, + }); + const ws = new Workspace({ + storage: new SQLiteTestStorage(), + backends: [backend], + retryScheduler: scheduler, + retry: { initialDelayMs: 100, maxDelayMs: 1_000, maxAttempts: 3 }, + now: () => 5_000, + }); + + await runCommand(ws); + const retried = await ws.retryPendingSync("sandbox"); + + expect(retried).toMatchObject({ status: "complete", applied: 1 }); + expect(execs).toBe(1); + expect(after).toEqual([ + { rev: 0, path: null }, + { rev: 1, path: "/generated/255" }, + ]); + expect(scheduler.intents.size).toBe(0); + expect(scheduler.cleared).toEqual(["sandbox"]); + }); + + it("coalesces repeated command failures into one pending intent per backend", async () => { + const scheduler = new MemoryRetryScheduler(); + const backend = retryBackend({ + onExec() {}, + async fetchChanges() { + throw new Error("still unavailable"); + }, + }); + const ws = new Workspace({ + storage: new SQLiteTestStorage(), + backends: [backend], + retryScheduler: scheduler, + now: () => 1_000, + }); + + await Promise.all([runCommand(ws), runCommand(ws)]); + + expect(scheduler.scheduled).toHaveLength(1); + expect(scheduler.intents.size).toBe(1); + }); + + it("reschedules with bounded exponential backoff and leaves exhaustion visible", async () => { + const scheduler = new MemoryRetryScheduler(); + const backend = retryBackend({ + onExec() {}, + async fetchChanges() { + throw new Error("still unavailable"); + }, + }); + const ws = new Workspace({ + storage: new SQLiteTestStorage(), + backends: [backend], + retryScheduler: scheduler, + retry: { initialDelayMs: 100, maxDelayMs: 150, maxAttempts: 3 }, + now: () => 1_000, + }); + + await runCommand(ws); + expect(await ws.retryPendingSync()).toMatchObject({ status: "pending", attempt: 2 }); + expect(await ws.retryPendingSync()).toMatchObject({ status: "pending", attempt: 3 }); + expect(await ws.retryPendingSync()).toMatchObject({ status: "exhausted", attempt: 3 }); + + expect(scheduler.scheduled).toEqual([ + { backend: "sandbox", attempt: 1, notBefore: 1_100 }, + { backend: "sandbox", attempt: 2, notBefore: 1_150 }, + { backend: "sandbox", attempt: 3, notBefore: 1_150 }, + ]); + expect(scheduler.intents.get("sandbox")).toEqual({ + backend: "sandbox", + attempt: 3, + notBefore: 1_150, + }); + }); + + it("disposes a failed retry envelope and closes its RPC handle", async () => { + const scheduler = new MemoryRetryScheduler(); + let closes = 0; + let disposals = 0; + const backend = retryBackend({ + onExec() {}, + async fetchChanges() { + return { + currentCursor: { rev: 1, path: null }, + appliedPushCursor: { rev: 0, path: null }, + stream: new ReadableStream({ + start(controller) { + controller.error(new Error("cancel retry stream")); + }, + }), + [Symbol.dispose]() { + disposals++; + }, + }; + }, + close: async () => { + closes++; + }, + }); + const ws = new Workspace({ + storage: new SQLiteTestStorage(), + backends: [backend], + retryScheduler: scheduler, + }); + scheduler.intents.set("sandbox", { backend: "sandbox", attempt: 1, notBefore: 0 }); + + await expect(ws.retryPendingSync()).resolves.toMatchObject({ status: "pending" }); + await ws.close(); + + expect(disposals).toBe(1); + expect(closes).toBe(1); + }); +}); diff --git a/packages/workspace/src/shell.ts b/packages/workspace/src/shell.ts index bc8e57f8..f067c789 100644 --- a/packages/workspace/src/shell.ts +++ b/packages/workspace/src/shell.ts @@ -136,6 +136,7 @@ export interface GetExecOptions { export interface Sync { push(): Promise; pull(): Promise; + onPullPending?(error: unknown): Promise; } export class WorkspaceShell { @@ -368,6 +369,13 @@ async function drainToResult( syncResult = { status: "complete", applied: pulled, skipped }; } catch (error) { syncResult = { status: "pending", applied: 0, skipped: [], error: safeErrorMessage(error) }; + try { + await sync.onPullPending?.(error); + } catch { + // The command result must remain available even when the host's + // durable scheduler is temporarily unavailable. The pending + // status keeps the missed pull visible to the caller. + } } return { exitCode, diff --git a/packages/workspace/src/workspace.ts b/packages/workspace/src/workspace.ts index 77e66eac..e1186a8d 100644 --- a/packages/workspace/src/workspace.ts +++ b/packages/workspace/src/workspace.ts @@ -31,11 +31,45 @@ import { createGitClient, type GitClient, type GitIdentity } from "./git/index.j import { MountIndex } from "./mounts/index.js"; import { buildMountRegistry, type MountValue } from "./mounts/registry.js"; import type { Mount } from "./mounts/types.js"; -import { noopObserver, type WorkspaceObserver, withSpan } from "./observe.js"; +import { noopObserver, safeErrorMessage, type WorkspaceObserver, withSpan } from "./observe.js"; import { WorkspaceShell } from "./shell.js"; import { WorkspaceStub } from "./stub.js"; import { isWorkspaceTransportFailure } from "./transport-failure.js"; +export interface SyncRetryIntent { + backend: string; + attempt: number; + notBefore: number; +} + +/** + * Durable storage boundary for pending post-command pulls. + * + * The host owns persistence and wake-up because the workspace library + * cannot own a Durable Object alarm. Each backend has at most one intent. + */ +export interface SyncRetryScheduler { + get(backend: string): Promise; + schedule(intent: SyncRetryIntent): Promise; + clear(backend: string): Promise; +} + +export interface SyncRetryOptions { + initialDelayMs?: number; + maxDelayMs?: number; + maxAttempts?: number; +} + +export type WorkspaceRetryPendingSyncResult = + | { status: "idle"; backend: string } + | { status: "complete"; backend: string; applied: number; skipped: ApplyResult["skipped"] } + | { status: "pending"; backend: string; attempt: number; notBefore: number; error: string } + | { status: "exhausted"; backend: string; attempt: number; error: string }; + +const DEFAULT_RETRY_INITIAL_DELAY_MS = 1_000; +const DEFAULT_RETRY_MAX_DELAY_MS = 60_000; +const DEFAULT_RETRY_MAX_ATTEMPTS = 5; + export interface WorkspaceOptions { // Local store backing this Workspace. In a Durable Object, pass // `ctx.storage`; in tests, pass a SQLiteTestStorage from @@ -74,6 +108,12 @@ export interface WorkspaceOptions { // adapter subpaths for the Cloudflare runtime and OpenTelemetry. observer?: WorkspaceObserver; + // Optional durable retry boundary for failed post-command pulls. + // The host persists one intent per backend and wakes the Durable + // Object at intent.notBefore to call retryPendingSync(backend). + retryScheduler?: SyncRetryScheduler; + retry?: SyncRetryOptions; + // Default identity used by commit-producing git subcommands // when neither the call site nor the relevant `GIT_AUTHOR_*` / // `GIT_COMMITTER_*` env vars supply one. Threaded through to @@ -111,6 +151,10 @@ export class Workspace { readonly #defaultBackendId: string | undefined; readonly #observer: WorkspaceObserver; readonly #now: () => number; + readonly #retryScheduler: SyncRetryScheduler | undefined; + readonly #retryInitialDelayMs: number; + readonly #retryMaxDelayMs: number; + readonly #retryMaxAttempts: number; readonly #sessionId: string; readonly #defaultGitIdentity: GitIdentity | undefined; readonly #assets: AssetsClient | undefined; @@ -142,6 +186,22 @@ export class Workspace { constructor(options: WorkspaceOptions) { this.#now = options.now ?? Date.now; + this.#retryScheduler = options.retryScheduler; + this.#retryInitialDelayMs = positiveRetryOption( + options.retry?.initialDelayMs, + DEFAULT_RETRY_INITIAL_DELAY_MS, + "initialDelayMs", + ); + this.#retryMaxDelayMs = positiveRetryOption( + options.retry?.maxDelayMs, + DEFAULT_RETRY_MAX_DELAY_MS, + "maxDelayMs", + ); + this.#retryMaxAttempts = positiveRetryOption( + options.retry?.maxAttempts, + DEFAULT_RETRY_MAX_ATTEMPTS, + "maxAttempts", + ); this.#sessionId = options.sessionId ?? ""; this.#defaultGitIdentity = options.defaultGitIdentity; this.#artifacts = options.artifacts @@ -409,26 +469,98 @@ export class Workspace { } pull(id?: string): Promise { - return this.#serialize(id, (resolvedId) => - withSpan( - this.#observer, - "workspace.sync.pull", - { "workspace.sync.backend": resolvedId }, - async () => { - if (resolvedId === undefined) return { applied: 0, skipped: [] }; - const handle = await this.#handleFor(resolvedId); - if (handle.sync === "none") return { applied: 0, skipped: [] }; - return this.#runWithInvalidation(resolvedId, handle, () => - pullOnce(this.#db, handle.rpc.sync, resolvedId), - ); - }, - (span, outcome) => { - if (!outcome.ok) return; - span.setAttribute("workspace.sync.applied", outcome.value.applied); - span.setAttribute("workspace.sync.skipped", outcome.value.skipped.length); - }, - ), + return this.#serialize(id, (resolvedId) => this.#pullResolved(resolvedId)); + } + + /** + * Run a host-scheduled pending pull from its persisted cursor. + * + * The call shares the backend's mutation FIFO with push, pull, and + * command brackets. A successful pull clears the host's durable + * intent. A failed pull advances bounded exponential backoff; the + * last failed attempt remains stored and is reported as exhausted. + */ + retryPendingSync(id?: string): Promise { + return this.#serialize(id, async (resolvedId) => { + if (resolvedId === undefined) { + throw new Error("Workspace has no backend configured for pending sync retry"); + } + const scheduler = this.#retryScheduler; + if (scheduler === undefined) { + throw new Error("Workspace has no retryScheduler configured"); + } + const intent = await scheduler.get(resolvedId); + if (intent === undefined) return { status: "idle", backend: resolvedId }; + if (intent.attempt > this.#retryMaxAttempts) { + return { + status: "exhausted", + backend: resolvedId, + attempt: intent.attempt, + error: "pending sync retry attempts exhausted", + }; + } + try { + const result = await this.#pullResolved(resolvedId); + await scheduler.clear(resolvedId); + return { + status: "complete", + backend: resolvedId, + applied: result.applied, + skipped: result.skipped, + }; + } catch (error) { + const message = safeErrorMessage(error); + if (intent.attempt >= this.#retryMaxAttempts) { + return { + status: "exhausted", + backend: resolvedId, + attempt: intent.attempt, + error: message, + }; + } + const next = this.#retryIntent(resolvedId, intent.attempt + 1); + await scheduler.schedule(next); + return { status: "pending", ...next, error: message }; + } + }); + } + + #pullResolved(resolvedId: string | undefined): Promise { + return withSpan( + this.#observer, + "workspace.sync.pull", + { "workspace.sync.backend": resolvedId }, + async () => { + if (resolvedId === undefined) return { applied: 0, skipped: [] }; + const handle = await this.#handleFor(resolvedId); + if (handle.sync === "none") return { applied: 0, skipped: [] }; + return this.#runWithInvalidation(resolvedId, handle, () => + pullOnce(this.#db, handle.rpc.sync, resolvedId), + ); + }, + (span, outcome) => { + if (!outcome.ok) return; + span.setAttribute("workspace.sync.applied", outcome.value.applied); + span.setAttribute("workspace.sync.skipped", outcome.value.skipped.length); + }, + ); + } + + async #schedulePendingSync(id: string): Promise { + const scheduler = this.#retryScheduler; + if (scheduler === undefined) return; + await this.#serialize(id, async (resolvedId) => { + if (resolvedId === undefined || (await scheduler.get(resolvedId)) !== undefined) return; + await scheduler.schedule(this.#retryIntent(resolvedId, 1)); + }); + } + + #retryIntent(backend: string, attempt: number): SyncRetryIntent { + const delay = Math.min( + this.#retryMaxDelayMs, + this.#retryInitialDelayMs * 2 ** Math.max(0, attempt - 1), ); + return { backend, attempt, notBefore: this.#now() + delay }; } // Drop a cached handle when an operation fails with a known @@ -604,6 +736,7 @@ export class Workspace { { push: () => this.push(id), pull: () => this.pull(id), + onPullPending: () => this.#schedulePendingSync(id), }, this.#observer, ); @@ -635,6 +768,14 @@ export class Workspace { } } +function positiveRetryOption(value: number | undefined, fallback: number, name: string): number { + const resolved = value ?? fallback; + if (!Number.isSafeInteger(resolved) || resolved <= 0) { + throw new Error(`Workspace retry.${name} must be a positive integer`); + } + return resolved; +} + function createDisabledArtifactsClient(): ArtifactClient { const fail = () => { throw new ArtifactError("ENOCONFIG", "Workspace Artifacts binding is not configured"); From fb58acd5da3d2ac2fa710e840b348971fb7ce6ce Mon Sep 17 00:00:00 2001 From: aron <263346377+aron-cf@users.noreply.github.com> Date: Thu, 30 Jul 2026 09:17:53 +0100 Subject: [PATCH 5/7] dofs: correct stale blob-cache immutability comment stageBlob now repairs corrupt/incomplete rows via ON CONFLICT DO UPDATE SET bytes and clears the blob cache; the normal write path (upsertChunkBlob) still uses DO NOTHING. Update the comment so it no longer claims (hash, bytes) pairs are unconditionally immutable. --- packages/dofs/src/fs/blobCache.ts | 13 ++++++++----- 1 file changed, 8 insertions(+), 5 deletions(-) diff --git a/packages/dofs/src/fs/blobCache.ts b/packages/dofs/src/fs/blobCache.ts index 16451d3c..b997ee82 100644 --- a/packages/dofs/src/fs/blobCache.ts +++ b/packages/dofs/src/fs/blobCache.ts @@ -7,11 +7,14 @@ // blob in vfs_blobs, and we then re-fetch that one blob 512 times // over the lifetime of one read pass. // -// vfs_blob_bytes is content-addressed and immutable: a stored -// (hash, bytes) pair never changes for the life of the database. -// That makes the cache trivially correct — any write that mutates a -// file produces new chunk rows with new hashes, never overwriting -// the bytes the cache holds. +// vfs_blob_bytes is content-addressed. The normal write path +// (upsertChunkBlob) uses ON CONFLICT DO NOTHING, so a correct +// (hash, bytes) pair is never overwritten and the cache stays valid +// for it. The one exception is repair: stageBlob (the sync receiver +// path) uses ON CONFLICT DO UPDATE SET bytes to replace an incomplete +// or size-mismatched payload left by an interrupted or corrupt write, +// and clears this cache afterward so a stale payload is never served +// after a repair. // // The cache is bounded (CHUNK_CACHE_MAX_ENTRIES) and per-Database so // independent test databases don't pollute each other. Eviction is From 34d7b183d2e48a9075a27d8a2af67d37c50b002b Mon Sep 17 00:00:00 2001 From: aron <263346377+aron-cf@users.noreply.github.com> Date: Thu, 30 Jul 2026 09:17:53 +0100 Subject: [PATCH 6/7] workspace: fix README alarm example looping on exhausted retries retryPendingSync leaves the final intent in storage with a past-due notBefore when a backend exhausts its attempts. The example alarm re-armed on min(notBefore), so an exhausted intent made the alarm fire immediately and repeatedly. Clear exhausted intents in the example. --- packages/workspace/README.md | 9 ++++++++- 1 file changed, 8 insertions(+), 1 deletion(-) diff --git a/packages/workspace/README.md b/packages/workspace/README.md index 426b92e6..40d56860 100644 --- a/packages/workspace/README.md +++ b/packages/workspace/README.md @@ -302,7 +302,14 @@ export class WorkspaceHost extends DurableObject { const now = Date.now(); for (const intent of intents.values()) { if (intent.notBefore <= now) { - await this.workspace.retryPendingSync(intent.backend); + const result = await this.workspace.retryPendingSync(intent.backend); + // An exhausted backend keeps its final intent in storage with a + // past-due notBefore. Clear it here so it stops driving the alarm; + // otherwise the wake-up fires immediately and forever. Inspect or + // alert before clearing if you need to surface the exhaustion. + if (result.status === "exhausted") { + await this.scheduler.clear(intent.backend); + } } } From 92c44785777761518f5a9b8d1d65ae4ededd519b Mon Sep 17 00:00:00 2001 From: aron <263346377+aron-cf@users.noreply.github.com> Date: Thu, 30 Jul 2026 09:17:53 +0100 Subject: [PATCH 7/7] workspace: fix structured sync stub test Construct the WorkspaceExecHandleStub test with a fake ExecHandle so it matches the eagerly-spawned handle implementation after the rebase. This keeps the test focused on forwarding ExecResult.sync through the Workers RPC stub. --- packages/workspace/src/stub.test.ts | 42 +++++++++++++++++++---------- 1 file changed, 28 insertions(+), 14 deletions(-) diff --git a/packages/workspace/src/stub.test.ts b/packages/workspace/src/stub.test.ts index 094c83d6..76187169 100644 --- a/packages/workspace/src/stub.test.ts +++ b/packages/workspace/src/stub.test.ts @@ -17,6 +17,7 @@ import { beforeAll, describe, expect, it } from "vitest"; import type { BackendHandle, WorkspaceBackend } from "./backend.js"; import { decodeExecEvents } from "./exec-wire.js"; +import type { ExecHandle, ExecResult } from "./shell.js"; import { WorkspaceAssetsStub, WorkspaceExecHandleStub, @@ -338,21 +339,34 @@ describe("WorkspaceStub", () => { }); it("shell result carries the structured sync outcome across Workers RPC", async () => { - using handle = new WorkspaceExecHandleStub( - Promise.resolve({ - exitCode: 0, - stdout: "done", - stderr: "", - pushed: 1, - pulled: 0, + const result = { + exitCode: 0, + stdout: "done", + stderr: "", + pushed: 1, + pulled: 0, + skipped: [], + sync: { + status: "pending" as const, + applied: 0, skipped: [], - sync: { - status: "pending" as const, - applied: 0, - skipped: [], - error: "WebSocket closed before pull completed", - }, - }), + error: "WebSocket closed before pull completed", + }, + } satisfies ExecResult<"utf8">; + const hostHandle = { result: async () => result } as ExecHandle<"utf8">; + using handle = new WorkspaceExecHandleStub<"utf8">( + Promise.resolve(hostHandle), + { + promise: Promise.resolve({ + exitCode: 0, + pushed: 0, + pulled: 0, + skippedCount: 0, + sync: { status: "complete", applied: 0, skipped: [] }, + }), + resolve() {}, + }, + Promise.resolve(), ); await expect(handle.result()).resolves.toMatchObject({ exitCode: 0,