diff --git a/packages/dofs/src/schema/sync.ts b/packages/dofs/src/schema/sync.ts index fdfd73e4..1c263e30 100644 --- a/packages/dofs/src/schema/sync.ts +++ b/packages/dofs/src/schema/sync.ts @@ -69,6 +69,17 @@ export const SYNC_STATEMENTS = [ path TEXT, PRIMARY KEY (k, backend) )`, + // Local revs minted by applying a pull from `backend`. A node whose + // rev is listed here still holds exactly what that backend sent, so + // a later delete from the same backend may remove it before the echo + // push. Rows at or below the push cursor are pruned. Created by the + // baseline DDL on every boot, so existing databases gain it without + // a migration. + `CREATE TABLE IF NOT EXISTS _vfs_upstream_revs ( + backend TEXT NOT NULL, + rev INTEGER NOT NULL, + PRIMARY KEY (backend, rev) + ) WITHOUT ROWID`, // Durable half of a restartable sync operation. One row per // (backend, direction): the plan's key. A restarted iterator reads // this row to recover the fixed target and generation it was working diff --git a/packages/dofs/src/sync/apply.ts b/packages/dofs/src/sync/apply.ts index ccc421a7..0f087f57 100644 --- a/packages/dofs/src/sync/apply.ts +++ b/packages/dofs/src/sync/apply.ts @@ -13,6 +13,15 @@ import { stageBlob } from "./blobs.js"; import type { ChangeEntry } from "./changes.js"; import { computeManifestHash } from "./manifests.js"; import { pathOf } from "./paths.js"; +import { + type ChangeCursor, + compareChangeCursors, + currentRev, + DEFAULT_BACKEND_ID, + readPushCursor, + readWatermark, + recordUpstreamRevs, +} from "./watermarks.js"; // One container-side change that landed under a read-only mount and // was therefore skipped rather than applied. Callers (the workspace @@ -70,6 +79,12 @@ export interface ApplyOptions { // which is fine for the container backend the package shipped // with first. backend?: string; + // Push receivers checkpoint the sender's cursor when committing a + // batch. Supplying that cursor protects local recreations from an + // already-committed delete replay; new pushes remain authoritative, + // just like incoming writes. Pulls instead protect unpushed local + // changes using this backend's local pushRev. + receivedCursor?: ChangeCursor; } const DEFAULT_MAX_BYTES = 64 * 1024 * 1024; @@ -283,94 +298,103 @@ export async function applyChanges( }; for await (const entry of entries) { - // Idempotent skip: if the entry already matches the local - // state, drop it on the floor. The check is what stops a - // pull from bumping vfs_meta.rev for entries that are - // already in place, which in turn stops the next push from - // re-shipping them. - if (options.source === "upstream" && entry.kind !== "delete") { - if (alreadyApplied(db, entry)) continue; - } - // A replayed tombstone must not delete a path that was recreated - // above the tombstone's revision. - if (options.source === "upstream" && entry.kind === "delete" && tombstoneIsStale(db, entry)) { - continue; - } - // Read-only mount guard. Entries under a registered read-only - // mount root are surfaced via the return value and not applied. - // The owning workspace's surface (Workspace.pull, exec()) folds - // these into its own return so callers see what stayed - // authoritative on the mount. - const blockingRoot = readOnlyRootFor(db, entry.path); - if (blockingRoot !== undefined) { - skipped.push({ - path: entry.path, - mountRoot: blockingRoot, - op: entry.kind === "delete" ? "delete" : "write", - reason: "read-only", - }); - continue; - } - if (entry.kind === "delete") { - try { - rm(db, entry.path, { recursive: true, force: true }); - } catch { - // Already gone is fine — idempotent apply. + const revBefore = currentRev(db); + const keepsLocalLink = rewritesUnpushedHardlink(db, entry, options); + try { + // Idempotent skip: if the entry already matches the local + // state, drop it on the floor. The check is what stops a + // pull from bumping vfs_meta.rev for entries that are + // already in place, which in turn stops the next push from + // re-shipping them. + if (options.source === "upstream" && entry.kind !== "delete") { + if (alreadyApplied(db, entry)) continue; } - applied++; - pathsInBatch++; - if (pathsInBatch >= maxPaths) flush(); - continue; - } - if (entry.kind === "dir") { - const parentResult = ensureParentDirectories(db, entry.path, entry.mtime); - if (parentResult.blockingRoot !== undefined) { + // Protect local changes without comparing independent peer rev spaces. + if ( + options.source === "upstream" && + entry.kind === "delete" && + tombstoneIsStale(db, entry, options) + ) { + continue; + } + // Read-only mount guard. Entries under a registered read-only + // mount root are surfaced via the return value and not applied. + // The owning workspace's surface (Workspace.pull, exec()) folds + // these into its own return so callers see what stayed + // authoritative on the mount. + const blockingRoot = readOnlyRootFor(db, entry.path); + if (blockingRoot !== undefined) { skipped.push({ path: entry.path, - mountRoot: parentResult.blockingRoot, - op: "write", + mountRoot: blockingRoot, + op: entry.kind === "delete" ? "delete" : "write", reason: "read-only", }); continue; } - applyDirectoryEntry(db, { ...entry, path: parentResult.path }); - applied++; - pathsInBatch++; - if (pathsInBatch >= maxPaths) flush(); - continue; - } - if (entry.kind === "symlink") { - const parentResult = ensureParentDirectories(db, entry.path, entry.mtime); - if (parentResult.blockingRoot !== undefined) { + if (entry.kind === "delete") { + try { + rm(db, entry.path, { recursive: true, force: true }); + } catch { + // Already gone is fine — idempotent apply. + } + applied++; + pathsInBatch++; + if (pathsInBatch >= maxPaths) flush(); + continue; + } + if (entry.kind === "dir") { + const parentResult = ensureParentDirectories(db, entry.path, entry.mtime); + if (parentResult.blockingRoot !== undefined) { + skipped.push({ + path: entry.path, + mountRoot: parentResult.blockingRoot, + op: "write", + reason: "read-only", + }); + continue; + } + applyDirectoryEntry(db, { ...entry, path: parentResult.path }); + applied++; + pathsInBatch++; + if (pathsInBatch >= maxPaths) flush(); + continue; + } + if (entry.kind === "symlink") { + const parentResult = ensureParentDirectories(db, entry.path, entry.mtime); + if (parentResult.blockingRoot !== undefined) { + skipped.push({ + path: entry.path, + mountRoot: parentResult.blockingRoot, + op: "write", + reason: "read-only", + }); + continue; + } + removeReplaceableFinalEntry(db, parentResult.path, "symlink"); + symlink(db, entry.target, parentResult.path, () => entry.mtime); + applied++; + pathsInBatch++; + if (pathsInBatch >= maxPaths) flush(); + continue; + } + const fileResult = applyFileEntry(db, entry, objects); + if (fileResult.blockingRoot !== undefined) { skipped.push({ path: entry.path, - mountRoot: parentResult.blockingRoot, + mountRoot: fileResult.blockingRoot, op: "write", reason: "read-only", }); continue; } - removeReplaceableFinalEntry(db, parentResult.path, "symlink"); - symlink(db, entry.target, parentResult.path, () => entry.mtime); applied++; + bytesInBatch += fileResult.total; pathsInBatch++; - if (pathsInBatch >= maxPaths) flush(); - continue; - } - const fileResult = applyFileEntry(db, entry, objects); - if (fileResult.blockingRoot !== undefined) { - skipped.push({ - path: entry.path, - mountRoot: fileResult.blockingRoot, - op: "write", - reason: "read-only", - }); - continue; + if (bytesInBatch >= maxBytes || pathsInBatch >= maxPaths) flush(); + } finally { + if (!keepsLocalLink) recordPulledRevs(db, options, revBefore); } - applied++; - bytesInBatch += fileResult.total; - pathsInBatch++; - if (bytesInBatch >= maxBytes || pathsInBatch >= maxPaths) flush(); } // Loopback suppression used to advance pushRev locally after an @@ -418,84 +442,94 @@ export function applyChangesSync( }; for (const entry of entries) { - if (options.source === "upstream" && entry.kind !== "delete") { - if (alreadyApplied(db, entry)) continue; - } - // See tombstoneIsStale: a replayed delete must not clobber a - // newer local recreation. - if (options.source === "upstream" && entry.kind === "delete" && tombstoneIsStale(db, entry)) { - continue; - } - const blockingRoot = readOnlyRootFor(db, entry.path); - if (blockingRoot !== undefined) { - skipped.push({ - path: entry.path, - mountRoot: blockingRoot, - op: entry.kind === "delete" ? "delete" : "write", - reason: "read-only", - }); - continue; - } - if (entry.kind === "delete") { - try { - rm(db, entry.path, { recursive: true, force: true }); - } catch { - // Already gone is fine — idempotent apply. + const revBefore = currentRev(db); + const keepsLocalLink = rewritesUnpushedHardlink(db, entry, options); + try { + if (options.source === "upstream" && entry.kind !== "delete") { + if (alreadyApplied(db, entry)) continue; } - applied++; - pathsInBatch++; - if (pathsInBatch >= maxPaths) flush(); - continue; - } - if (entry.kind === "dir") { - const parentResult = ensureParentDirectories(db, entry.path, entry.mtime); - if (parentResult.blockingRoot !== undefined) { + // See tombstoneIsStale: a replayed delete must not clobber a + // newer local recreation. + if ( + options.source === "upstream" && + entry.kind === "delete" && + tombstoneIsStale(db, entry, options) + ) { + continue; + } + const blockingRoot = readOnlyRootFor(db, entry.path); + if (blockingRoot !== undefined) { skipped.push({ path: entry.path, - mountRoot: parentResult.blockingRoot, - op: "write", + mountRoot: blockingRoot, + op: entry.kind === "delete" ? "delete" : "write", reason: "read-only", }); continue; } - applyDirectoryEntry(db, { ...entry, path: parentResult.path }); - applied++; - pathsInBatch++; - if (pathsInBatch >= maxPaths) flush(); - continue; - } - if (entry.kind === "symlink") { - const parentResult = ensureParentDirectories(db, entry.path, entry.mtime); - if (parentResult.blockingRoot !== undefined) { + if (entry.kind === "delete") { + try { + rm(db, entry.path, { recursive: true, force: true }); + } catch { + // Already gone is fine — idempotent apply. + } + applied++; + pathsInBatch++; + if (pathsInBatch >= maxPaths) flush(); + continue; + } + if (entry.kind === "dir") { + const parentResult = ensureParentDirectories(db, entry.path, entry.mtime); + if (parentResult.blockingRoot !== undefined) { + skipped.push({ + path: entry.path, + mountRoot: parentResult.blockingRoot, + op: "write", + reason: "read-only", + }); + continue; + } + applyDirectoryEntry(db, { ...entry, path: parentResult.path }); + applied++; + pathsInBatch++; + if (pathsInBatch >= maxPaths) flush(); + continue; + } + if (entry.kind === "symlink") { + const parentResult = ensureParentDirectories(db, entry.path, entry.mtime); + if (parentResult.blockingRoot !== undefined) { + skipped.push({ + path: entry.path, + mountRoot: parentResult.blockingRoot, + op: "write", + reason: "read-only", + }); + continue; + } + removeReplaceableFinalEntry(db, parentResult.path, "symlink"); + symlink(db, entry.target, parentResult.path, () => entry.mtime); + applied++; + pathsInBatch++; + if (pathsInBatch >= maxPaths) flush(); + continue; + } + const fileResult = applyFileEntry(db, entry, objects); + if (fileResult.blockingRoot !== undefined) { skipped.push({ path: entry.path, - mountRoot: parentResult.blockingRoot, + mountRoot: fileResult.blockingRoot, op: "write", reason: "read-only", }); continue; } - removeReplaceableFinalEntry(db, parentResult.path, "symlink"); - symlink(db, entry.target, parentResult.path, () => entry.mtime); applied++; + bytesInBatch += fileResult.total; pathsInBatch++; - if (pathsInBatch >= maxPaths) flush(); - continue; - } - const fileResult = applyFileEntry(db, entry, objects); - if (fileResult.blockingRoot !== undefined) { - skipped.push({ - path: entry.path, - mountRoot: fileResult.blockingRoot, - op: "write", - reason: "read-only", - }); - continue; + if (bytesInBatch >= maxBytes || pathsInBatch >= maxPaths) flush(); + } finally { + if (!keepsLocalLink) recordPulledRevs(db, options, revBefore); } - applied++; - bytesInBatch += fileResult.total; - pathsInBatch++; - if (bytesInBatch >= maxBytes || pathsInBatch >= maxPaths) flush(); } // See applyChanges() for why pushRev no longer advances locally @@ -567,27 +601,103 @@ function assertChunkSize(actual: number, declared: number, hash: Uint8Array, pat ); } -// Decide whether an upstream tombstone may delete the live path. -// -// A tombstone describes the path as of the revision it was stamped -// with. Because the sync cursor only advances after a block applies, -// any block interrupted before its acknowledgment is replayed — and a -// replayed tombstone whose path was recreated locally in the meantime -// would destroy content the tombstone never described. +// Peer entry.rev and local node.rev are independent counters. // -// The guard is a revision comparison: apply the delete only when the -// live path is no newer than the tombstone. A path recreated above the -// tombstone's rev is newer information than the delete, so the delete -// is stale and dropped. It is not lost work — the recreation is itself -// a change that the next push ships upstream. +// On pull, a delete yields only to local versions this backend has not +// seen: any node in the doomed subtree (a directory delete removes its +// descendants too) that sits beyond the push cursor and was not minted +// by an earlier pull from the same backend. The cursor's path matters, +// since a bounded push can ship one rev only in part. A pulled version +// awaiting its echo push still holds what the remote sent, so the +// remote's later delete wins. A reset watermark conservatively protects +// every locally authored version. // -// Local deletes are exempt. They are authored here, not replayed, so -// there is no earlier revision to compare against. -function tombstoneIsStale(db: Database, entry: ChangeEntry & { kind: "delete" }): boolean { +// On push, the receiver's committed sender cursor identifies replays +// in the sender's own rev space. New pushes are authoritative, while a +// replay must not remove a path recreated after the original commit. +function tombstoneIsStale( + db: Database, + entry: ChangeEntry & { kind: "delete" }, + options: ApplyOptions, +): boolean { const live = resolveInode(db, entry.path, { followSymlinks: false }); if (live === null) return false; - const row = db.one<{ rev: number }>("SELECT rev FROM vfs_nodes WHERE inode = ?", live.inode); - return row !== undefined && row.rev > entry.rev; + if (options.receivedCursor !== undefined) { + return compareChangeCursors({ rev: entry.rev, path: entry.path }, options.receivedCursor) <= 0; + } + const backend = options.backend ?? DEFAULT_BACKEND_ID; + const pushed = shippedCursor(db, backend); + const candidates = db.all<{ path: string; rev: number }>( + `WITH RECURSIVE tree(inode, path) AS ( + SELECT ?, ? + UNION ALL + SELECT d.child_inode, + CASE WHEN tree.path = '/' THEN '/' || d.name ELSE tree.path || '/' || d.name END + FROM vfs_dirents d JOIN tree ON d.parent_inode = tree.inode + ) + SELECT tree.path AS path, n.rev AS rev + FROM tree JOIN vfs_nodes n ON n.inode = tree.inode + WHERE n.rev >= ? + AND NOT EXISTS ( + SELECT 1 FROM _vfs_upstream_revs u WHERE u.backend = ? AND u.rev = n.rev + )`, + live.inode, + entry.path, + pushed.rev, + backend, + ); + return candidates.some((node) => compareChangeCursors(node, pushed) > 0); +} + +// The push cursor, capped by the pushRev watermark so a watermark +// reset protects local versions even if the cursor row was left behind. +function shippedCursor(db: Database, backend: string): ChangeCursor { + const watermark = readWatermark(db, "pushRev", backend); + const cursor = readPushCursor(db, backend); + return cursor.rev > watermark ? { rev: watermark, path: null } : cursor; +} + +// Remember revs minted while applying a pull, so a later delete from +// the same backend can tell them from local edits. Push receivers +// (receivedCursor set) guard replays by cursor instead and never prune, +// so they record nothing. +function recordPulledRevs(db: Database, options: ApplyOptions, revBefore: number): void { + if (options.source !== "upstream" || options.receivedCursor !== undefined) return; + recordUpstreamRevs(db, revBefore, currentRev(db), options.backend); +} + +// A pulled write to one name of a hardlinked file restamps the inode +// every name shares. When that inode holds an unpushed local version, +// such as a fresh link, the new rev must stay local: recording it as +// upstream would let a later delete of another name discard that link. +function rewritesUnpushedHardlink( + db: Database, + entry: ChangeEntry, + options: ApplyOptions, +): boolean { + if (entry.kind !== "file" || options.source !== "upstream") return false; + if (options.receivedCursor !== undefined) return false; + const live = resolveInode(db, entry.path, { followSymlinks: false }); + if (live === null || live.type !== "file") return false; + const node = db.one<{ rev: number; links: number }>( + `SELECT n.rev AS rev, + (SELECT count(*) FROM vfs_dirents d WHERE d.child_inode = n.inode) AS links + FROM vfs_nodes n WHERE n.inode = ?`, + live.inode, + ); + if (node === undefined || node.links < 2) return false; + const backend = options.backend ?? DEFAULT_BACKEND_ID; + const pushed = shippedCursor(db, backend); + // The names do not share a path, so a partially shipped rev counts + // as unpushed. + if (node.rev < pushed.rev || (node.rev === pushed.rev && pushed.path === null)) return false; + return ( + db.one( + "SELECT 1 AS hit FROM _vfs_upstream_revs WHERE backend = ? AND rev = ?", + backend, + node.rev, + ) === undefined + ); } // Compare an entry against the local node graph. Returns true when diff --git a/packages/dofs/src/sync/replay.test.ts b/packages/dofs/src/sync/replay.test.ts index d821224e..87f90f51 100644 --- a/packages/dofs/src/sync/replay.test.ts +++ b/packages/dofs/src/sync/replay.test.ts @@ -1,11 +1,15 @@ import { describe, expect, it } from "vitest"; +import { link } from "../fs/link.js"; +import { mkdir } from "../fs/mkdir.js"; import { readFile } from "../fs/readFile.js"; +import { rename } from "../fs/rename.js"; import { resolveInode } from "../fs/resolve.js"; import { withDB } from "../fs/with-db.js"; import { writeFile } from "../fs/writeFile.js"; -import { applyChanges } from "./apply.js"; +import { applyChanges, applyChangesSync } from "./apply.js"; import type { ChangeEntry } from "./changes.js"; +import { currentRev, writePushCursor, writeWatermark } from "./watermarks.js"; // The unified sync plan asserts that replaying an unacknowledged block // is safe because "revision and cursor semantics already make replay @@ -67,6 +71,7 @@ describe("block replay idempotency", () => { it("applies a tombstone twice without error when the path stays gone", async () => { await withDB(async (db) => { await writeFile(db, "/gone.txt", "bye", {}, () => 1); + writeWatermark(db, "pushRev", currentRev(db)); const entry: ChangeEntry = { kind: "delete", rev: 5, path: "/gone.txt" }; await applyChanges(db, [entry], new Map(), { source: "upstream" }); @@ -77,32 +82,25 @@ describe("block replay idempotency", () => { }); }); - // The plan's idempotency claim breaks here. The tombstone was - // produced at rev 5 and describes the file as it was then. If the - // block carrying it is interrupted before acknowledgment and the - // path is recreated locally in the meantime, replaying the - // tombstone deletes content it never described. - it("does not delete a path recreated at a newer revision than the tombstone", async () => { + it("does not delete an unpushed recreation on replay", async () => { await withDB(async (db) => { await writeFile(db, "/data.txt", "original", {}, () => 1); - // The tombstone's rev is whatever the source stamped it with. - // What matters is that the local recreation lands above it. - const tombstone: ChangeEntry = { kind: "delete", rev: 2, path: "/data.txt" }; + writeWatermark(db, "pushRev", currentRev(db)); + const tombstone: ChangeEntry = { kind: "delete", rev: 9_999, path: "/data.txt" }; // The block applied the tombstone but died before it could // acknowledge its cursor. await applyChanges(db, [tombstone], new Map(), { source: "upstream" }); expect(resolveInode(db, "/data.txt", { followSymlinks: false })).toBeNull(); - // A local write recreates the path at a revision above the - // tombstone's. + // A local write recreates the path above the local push watermark. await writeFile(db, "/data.txt", "recreated", {}, () => 2); const liveRev = db.one<{ rev: number }>( "SELECT rev FROM vfs_nodes WHERE inode = ?", resolveInode(db, "/data.txt", { followSymlinks: false })?.inode ?? 0, )?.rev ?? 0; - expect(liveRev).toBeGreaterThan(tombstone.rev); + expect(liveRev).toBeGreaterThan(1); // The interrupted block is replayed from the durable cursor. const replay = await applyChanges(db, [tombstone], new Map(), { source: "upstream" }); @@ -113,12 +111,12 @@ describe("block replay idempotency", () => { }); }); - it("still deletes a path whose live revision predates the tombstone", async () => { + it("deletes a pushed path even when the peer revision is smaller", async () => { await withDB(async (db) => { await writeFile(db, "/stale.txt", "stale", {}, () => 1); - // A tombstone from far above the live rev is a genuine delete - // the receiver has not seen yet. - const tombstone: ChangeEntry = { kind: "delete", rev: 9_999, path: "/stale.txt" }; + await writeFile(db, "/stale.txt", "newer", {}, () => 2); + writeWatermark(db, "pushRev", currentRev(db)); + const tombstone: ChangeEntry = { kind: "delete", rev: 1, path: "/stale.txt" }; const result = await applyChanges(db, [tombstone], new Map(), { source: "upstream" }); @@ -127,6 +125,250 @@ describe("block replay idempotency", () => { }); }); + for (const [name, apply] of [ + ["async", applyChanges], + ["sync", applyChangesSync], + ] as const) { + it(`${name}: uses only the selected backend's push watermark`, async () => { + await withDB(async (db) => { + await writeFile(db, "/x", "content", {}, () => 1); + writeWatermark(db, "pushRev", currentRev(db), "other"); + const entry: ChangeEntry = { kind: "delete", rev: 9_999, path: "/x" }; + expect( + (await apply(db, [entry], new Map(), { source: "upstream", backend: "linux" })).applied, + ).toBe(0); + writeWatermark(db, "pushRev", currentRev(db), "linux"); + expect( + ( + await apply(db, [{ ...entry, rev: 1 }], new Map(), { + source: "upstream", + backend: "linux", + }) + ).applied, + ).toBe(1); + }); + }); + + it(`${name}: protects unpushed edits after a watermark reset`, async () => { + await withDB(async (db) => { + await writeFile(db, "/x", "content", {}, () => 1); + writeWatermark(db, "pushRev", currentRev(db)); + writeWatermark(db, "pushRev", 0); + const result = await apply(db, [{ kind: "delete", rev: 9_999, path: "/x" }], new Map(), { + source: "upstream", + }); + expect(result.applied).toBe(0); + expect(await readFile(db, "/x", "utf8")).toBe("content"); + }); + }); + + it(`${name}: accepts new pushes but protects recreations from committed replays`, async () => { + await withDB(async (db) => { + await writeFile(db, "/x", "original", {}, () => 1); + await writeFile(db, "/x", "updated", {}, () => 2); + const entry: ChangeEntry = { kind: "delete", rev: 1, path: "/x" }; + const first = await apply(db, [entry], new Map(), { + source: "upstream", + receivedCursor: { rev: 0, path: null }, + }); + expect(first.applied).toBe(1); + await writeFile(db, "/x", "recreated", {}, () => 3); + const replay = await apply(db, [entry], new Map(), { + source: "upstream", + receivedCursor: { rev: 1, path: "/x" }, + }); + expect(replay.applied).toBe(0); + expect(await readFile(db, "/x", "utf8")).toBe("recreated"); + // A later path within the same sender rev is not a replay. + expect( + ( + await apply(db, [{ ...entry, path: "/z" }], new Map(), { + source: "upstream", + receivedCursor: { rev: 1, path: "/x" }, + }) + ).applied, + ).toBe(1); + }); + }); + } + + for (const [name, apply] of [ + ["async", applyChanges], + ["sync", applyChangesSync], + ] as const) { + it(`${name}: lets a remote delete win over an unechoed pulled file`, async () => { + await withDB(async (db) => { + writeWatermark(db, "pushRev", currentRev(db)); + const file: ChangeEntry = { + kind: "file", + rev: 1, + path: "/pulled/x", + mode: 0o644, + mtime: 1, + size: 0, + chunks: [], + }; + expect((await apply(db, [file], new Map(), { source: "upstream" })).applied).toBe(1); + // No echo push yet: the pulled versions sit above pushRev. + const result = await apply( + db, + [{ kind: "delete", rev: 2, path: "/pulled/x" }], + new Map(), + { + source: "upstream", + }, + ); + expect(result.applied).toBe(1); + expect(resolveInode(db, "/pulled/x", { followSymlinks: false })).toBeNull(); + // Pulled parents carry upstream provenance too. + expect( + ( + await apply(db, [{ kind: "delete", rev: 3, path: "/pulled" }], new Map(), { + source: "upstream", + }) + ).applied, + ).toBe(1); + }); + }); + + it(`${name}: protects a local edit made after a pull`, async () => { + await withDB(async (db) => { + writeWatermark(db, "pushRev", currentRev(db)); + const file: ChangeEntry = { + kind: "file", + rev: 1, + path: "/x", + mode: 0o644, + mtime: 1, + size: 0, + chunks: [], + }; + await apply(db, [file], new Map(), { source: "upstream" }); + await writeFile(db, "/x", "local edit", {}, () => 2); + const result = await apply(db, [{ kind: "delete", rev: 2, path: "/x" }], new Map(), { + source: "upstream", + }); + expect(result.applied).toBe(0); + expect(await readFile(db, "/x", "utf8")).toBe("local edit"); + }); + }); + + it(`${name}: protects unpushed descendants from a directory delete`, async () => { + await withDB(async (db) => { + mkdir(db, "/project", {}, () => 1); + await writeFile(db, "/project/draft", "pushed", {}, () => 1); + writeWatermark(db, "pushRev", currentRev(db)); + await writeFile(db, "/project/draft", "unpushed", {}, () => 2); + const result = await apply( + db, + [{ kind: "delete", rev: 1, path: "/project" }], + new Map(), + { + source: "upstream", + }, + ); + expect(result.applied).toBe(0); + expect(await readFile(db, "/project/draft", "utf8")).toBe("unpushed"); + }); + }); + + it(`${name}: still deletes a clean directory tree`, async () => { + await withDB(async (db) => { + mkdir(db, "/project", {}, () => 1); + await writeFile(db, "/project/done", "pushed", {}, () => 1); + writeWatermark(db, "pushRev", currentRev(db)); + const result = await apply( + db, + [{ kind: "delete", rev: 1, path: "/project" }], + new Map(), + { + source: "upstream", + }, + ); + expect(result.applied).toBe(1); + expect(resolveInode(db, "/project", { followSymlinks: false })).toBeNull(); + }); + }); + + it(`${name}: protects same-rev paths beyond a partial push cursor`, async () => { + await withDB(async (db) => { + mkdir(db, "/old", {}, () => 1); + await writeFile(db, "/old/a", "a", {}, () => 1); + await writeFile(db, "/old/b", "b", {}, () => 1); + rename(db, "/old", "/new"); + const renameRev = currentRev(db); + // A bounded push shipped /new and /new/a, but not /new/b. + writePushCursor(db, { rev: renameRev, path: "/new/a" }); + const unshipped = await apply( + db, + [{ kind: "delete", rev: 1, path: "/new/b" }], + new Map(), + { + source: "upstream", + }, + ); + expect(unshipped.applied).toBe(0); + expect(await readFile(db, "/new/b", "utf8")).toBe("b"); + const shipped = await apply(db, [{ kind: "delete", rev: 1, path: "/new/a" }], new Map(), { + source: "upstream", + }); + expect(shipped.applied).toBe(1); + expect(resolveInode(db, "/new/a", { followSymlinks: false })).toBeNull(); + }); + }); + it(`${name}: keeps an unpushed hardlink when a pull rewrites its inode`, async () => { + await withDB(async (db) => { + await writeFile(db, "/x", "content", {}, () => 1); + writeWatermark(db, "pushRev", currentRev(db)); + link(db, "/x", "/y"); + // The pull rewrites /x, which shares its inode with the + // unpushed /y, then deletes /y. + const update: ChangeEntry = { + kind: "file", + rev: 7, + path: "/x", + mode: 0o644, + mtime: 2, + size: 0, + chunks: [], + }; + const result = await apply( + db, + [update, { kind: "delete", rev: 8, path: "/y" }], + new Map(), + { source: "upstream" }, + ); + expect(result.applied).toBe(1); + expect(resolveInode(db, "/y", { followSymlinks: false })).not.toBeNull(); + }); + }); + + it(`${name}: still deletes a pushed hardlink after a pull rewrites its inode`, async () => { + await withDB(async (db) => { + await writeFile(db, "/x", "content", {}, () => 1); + link(db, "/x", "/y"); + writeWatermark(db, "pushRev", currentRev(db)); + const update: ChangeEntry = { + kind: "file", + rev: 7, + path: "/x", + mode: 0o644, + mtime: 2, + size: 0, + chunks: [], + }; + const result = await apply( + db, + [update, { kind: "delete", rev: 8, path: "/y" }], + new Map(), + { source: "upstream" }, + ); + expect(result.applied).toBe(2); + expect(resolveInode(db, "/y", { followSymlinks: false })).toBeNull(); + }); + }); + } + // A locally-authored delete is not a replay of remote state, so // the revision guard must not apply to it. it("applies a local tombstone regardless of the live revision", async () => { diff --git a/packages/dofs/src/sync/watermarks.ts b/packages/dofs/src/sync/watermarks.ts index 4fc180b4..bef45893 100644 --- a/packages/dofs/src/sync/watermarks.ts +++ b/packages/dofs/src/sync/watermarks.ts @@ -128,9 +128,26 @@ export function writePushCursor( cursor.path, ); writeWatermarkValue(db, "pushRev", cursor.rev, backend); + // Provenance below the cursor's rev can no longer matter: those + // versions have been shipped. Keep the cursor's own rev, which a + // partial push may have shipped only in part. + db.run("DELETE FROM _vfs_upstream_revs WHERE backend = ? AND rev < ?", backend, cursor.rev); }); } +// Record local revs in (after, through] as minted by applying a pull +// from `backend`. See _vfs_upstream_revs. +export function recordUpstreamRevs( + db: Database, + after: number, + through: number, + backend: string = DEFAULT_BACKEND_ID, +): void { + for (let rev = after + 1; rev <= through; rev++) { + db.run("INSERT OR IGNORE INTO _vfs_upstream_revs (backend, rev) VALUES (?, ?)", backend, rev); + } +} + export function writeFetchCursor( db: Database, cursor: ChangeCursor, diff --git a/packages/rpc/src/server.ts b/packages/rpc/src/server.ts index 4e675c3d..d6ea245e 100644 --- a/packages/rpc/src/server.ts +++ b/packages/rpc/src/server.ts @@ -142,6 +142,7 @@ class SyncRPCServer extends RpcTarget implements SyncRPC { this.db.transactionSync(() => { applyChangesSync(this.db, entries, new Map(), { source: isPeer ? "upstream" : "local", + receivedCursor: readFetchCursor(this.db), }); if (isPeer && compareChangeCursors(senderCursor, readFetchCursor(this.db)) > 0) { writeFetchCursor(this.db, senderCursor); @@ -296,6 +297,7 @@ class SyncRPCServer extends RpcTarget implements SyncRPC { } const result = applyChangesSync(this.db, decoded.entries, new Map(), { source: "upstream", + receivedCursor: readFetchCursor(this.db), }); // The footer's block cursor is the sender's checkpoint; echo it so // the sender advances only through what actually applied here. diff --git a/packages/rpc/src/sync-driver.test.ts b/packages/rpc/src/sync-driver.test.ts index 850e4af6..2a491633 100644 --- a/packages/rpc/src/sync-driver.test.ts +++ b/packages/rpc/src/sync-driver.test.ts @@ -1101,6 +1101,8 @@ describe("sync driver — streaming pullOnce", () => { await pullOnce(b.db, a.rpc); expect(providerB.readFileSync("/src/f299.txt", "utf8")).toBe("content 299"); + // A pull stamps local revs; finish the echo push before remote deletion. + await pushOnce(b.db, a.rpc); providerA.renameSync("/src", "/dst"); const renameRev = currentRev(a.db); diff --git a/packages/rpc/src/sync-engine-pack.test.ts b/packages/rpc/src/sync-engine-pack.test.ts index d987ca65..bcbb52d2 100644 --- a/packages/rpc/src/sync-engine-pack.test.ts +++ b/packages/rpc/src/sync-engine-pack.test.ts @@ -215,6 +215,8 @@ describe("pack mode pull", () => { try { seed(remote.db, 6); await drain(pullBlocks(local.db, remote.rpc, PACK_OPTIONS)); + // Complete the echo push so the peer has seen our local versions. + await drain(pushBlocks(local.db, remote.rpc, PACK_OPTIONS)); const provider = new SQLiteWorkspaceProvider(remote.db); provider.unlinkSync("/f000.txt"); @@ -420,3 +422,105 @@ describe("pack mode push", () => { } }); }); + +describe("deletes across independent peer revision spaces", () => { + for (const mode of ["entries", "pack"] as const) { + const options = { + thresholdEntries: mode === "pack" ? 1 : 1_000, + thresholdBytes: 1024 * 1024 * 1024, + }; + + it(`${mode}: pulls container deletes after pushing a higher local revision`, async () => { + const local = makePeer(); + const remote = makePeer(); + try { + const host = new SQLiteWorkspaceProvider(local.db); + const container = new SQLiteWorkspaceProvider(remote.db); + for (let i = 0; i < 20; i++) host.writeFileSync("/x", `version ${i}`); + host.writeFileSync("/y", "content"); + await drain(pushBlocks(local.db, remote.rpc, { ...options, backend: "linux" })); + container.unlinkSync("/x"); + container.unlinkSync("/y"); + const seen = await drain( + pullBlocks(local.db, remote.rpc, { ...options, backend: "linux" }), + ); + expect(seen[0].mode).toBe(mode); + expect(names(local.db)).toEqual([]); + } finally { + local.close(); + remote.close(); + } + }); + + it(`${mode}: a container delete wins over a pulled file awaiting its echo push`, async () => { + const local = makePeer(); + const remote = makePeer(); + try { + const container = new SQLiteWorkspaceProvider(remote.db); + container.writeFileSync("/x", "content"); + container.writeFileSync("/y", "content"); + await drain(pullBlocks(local.db, remote.rpc, { ...options, backend: "linux" })); + expect(names(local.db)).toEqual(["x", "y"]); + container.unlinkSync("/x"); + container.unlinkSync("/y"); + await drain(pullBlocks(local.db, remote.rpc, { ...options, backend: "linux" })); + expect(names(local.db)).toEqual([]); + await drain(pushBlocks(local.db, remote.rpc, { ...options, backend: "linux" })); + expect(names(remote.db)).toEqual([]); + } finally { + local.close(); + remote.close(); + } + }); + + it(`${mode}: a committed delete replay preserves a receiver recreation`, async () => { + const sender = makePeer(); + const receiver = makePeer(); + try { + const source = new SQLiteWorkspaceProvider(sender.db); + const target = new SQLiteWorkspaceProvider(receiver.db); + source.writeFileSync("/x", "original"); + source.writeFileSync("/y", "original"); + await drain(pushBlocks(sender.db, receiver.rpc, options)); + source.unlinkSync("/x"); + source.unlinkSync("/y"); + const failing = createSyncServer(receiver.db, { + afterApply: () => { + throw new Error("settle failed after commit"); + }, + }); + await expect(drain(pushBlocks(sender.db, failing, options))).rejects.toThrow( + "settle failed", + ); + expect(names(receiver.db)).toEqual([]); + target.writeFileSync("/x", "recreated"); + await drain(pushBlocks(sender.db, receiver.rpc, options)); + expect(target.readFileSync("/x", "utf8")).toBe("recreated"); + } finally { + sender.close(); + receiver.close(); + } + }); + + it(`${mode}: pushes deletes to a receiver whose local revision is higher`, async () => { + const sender = makePeer(); + const receiver = makePeer(); + try { + const source = new SQLiteWorkspaceProvider(sender.db); + const target = new SQLiteWorkspaceProvider(receiver.db); + for (let i = 0; i < 20; i++) target.writeFileSync("/unrelated", `version ${i}`); + source.writeFileSync("/x", "content"); + source.writeFileSync("/y", "content"); + await drain(pushBlocks(sender.db, receiver.rpc, options)); + source.unlinkSync("/x"); + source.unlinkSync("/y"); + const seen = await drain(pushBlocks(sender.db, receiver.rpc, options)); + expect(seen[0].mode).toBe(mode); + expect(names(receiver.db)).toEqual(["unrelated"]); + } finally { + sender.close(); + receiver.close(); + } + }); + } +});