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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
26 changes: 22 additions & 4 deletions docs/05_shell_interface.md
Original file line number Diff line number Diff line change
Expand Up @@ -101,12 +101,23 @@ type ExecEvent<T extends string | Uint8Array = Uint8Array> =
| { 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<T extends string | Uint8Array = Uint8Array> {
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;
}
```

Expand Down Expand Up @@ -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:
Expand Down Expand Up @@ -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

Expand Down
13 changes: 8 additions & 5 deletions packages/dofs/src/fs/blobCache.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
23 changes: 23 additions & 0 deletions packages/dofs/src/sync/blobs.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand All @@ -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: <Row extends object>(query: string, ...bindings: unknown[]) => {
if (query.startsWith("INSERT INTO vfs_blob_bytes")) {
throw new Error("injected bytes failure");
}
return db.sql.exec<Row>(query, ...bindings);
},
},
transactionSync: (closure) => db.transactionSync(closure),
});

expect(() => stageBlob(failingDb, hash, bytes, 1234)).toThrow("injected bytes failure");
expect(db.scalar<number>("SELECT COUNT(*) FROM vfs_blobs")).toBe(0);
expect(db.scalar<number>("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");
Expand Down
30 changes: 17 additions & 13 deletions packages/dofs/src/sync/blobs.ts
Original file line number Diff line number Diff line change
@@ -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
Expand All @@ -7,22 +8,25 @@ 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
// 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 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 UPDATE SET bytes = excluded.bytes",
hash,
bytes,
);
});
clearBlobCache(db);
}
30 changes: 30 additions & 0 deletions packages/dofs/src/sync/fetch.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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);
Expand Down
13 changes: 8 additions & 5 deletions packages/dofs/src/sync/fetch.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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));
Expand Down
130 changes: 130 additions & 0 deletions packages/rpc/src/sync-driver.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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();
Expand Down Expand Up @@ -991,6 +1072,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: <Row extends object>(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<Row>(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<number>("SELECT COUNT(*) FROM vfs_blobs")).toBe(0);
expect(receiver.db.scalar<number>("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
Expand Down
9 changes: 3 additions & 6 deletions packages/rpc/src/sync-driver.ts
Original file line number Diff line number Diff line change
Expand Up @@ -19,6 +19,7 @@ import {
compareChangeCursors,
currentRev,
type Database,
hasObjects,
readFetchCursor,
readWatermark,
type SkippedEntry,
Expand Down Expand Up @@ -214,14 +215,10 @@ async function pullOnceImpl(
const haveSubset = await remote.hasObjects(wantedHashes);
const remoteHasLocally = new Set<string>();
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,
Expand Down
Loading
Loading