From c426f7a3443c8bde91cde7d80e22248714ff2b86 Mon Sep 17 00:00:00 2001 From: dietrich ayala Date: Fri, 13 Feb 2026 17:04:14 +0100 Subject: [PATCH] Add GC and tombstone tests for SyncStorage and firehose #account events --- src/replication/gc-and-tombstone.test.ts | 458 +++++++++++++++++++++++ 1 file changed, 458 insertions(+) create mode 100644 src/replication/gc-and-tombstone.test.ts diff --git a/src/replication/gc-and-tombstone.test.ts b/src/replication/gc-and-tombstone.test.ts new file mode 100644 index 0000000..add7990 --- /dev/null +++ b/src/replication/gc-and-tombstone.test.ts @@ -0,0 +1,458 @@ +/** + * Tests for replication delete/update handling: + * - SyncStorage GC methods (removeBlocks, findOrphanedCids, etc.) + * - Firehose deferred GC (needs_gc flag) + * - Tombstone handling (markTombstoned, purgeDidData) + * - FirehoseSubscription #account event parsing + */ + +import { describe, it, expect, beforeEach, afterEach } from "vitest"; +import { mkdtempSync, rmSync } from "node:fs"; +import { tmpdir } from "node:os"; +import { join } from "node:path"; +import Database from "better-sqlite3"; +import { encode as cborEncode } from "../cbor-compat.js"; + +import { SyncStorage } from "./sync-storage.js"; +import { + FirehoseSubscription, + type FirehoseAccountEvent, +} from "./firehose-subscription.js"; + +// ============================================ +// Helpers +// ============================================ + +/** Encode a firehose frame (header + body) as concatenated CBOR. */ +function encodeFrame(header: object, body: object): Buffer { + const headerBytes = cborEncode(header); + const bodyBytes = cborEncode(body); + const frame = new Uint8Array(headerBytes.length + bodyBytes.length); + frame.set(headerBytes, 0); + frame.set(bodyBytes, headerBytes.length); + return Buffer.from(frame); +} + +/** Create a mock #account frame. */ +function makeAccountFrame( + seq: number, + did: string, + active: boolean, + status?: string, +): Buffer { + const header = { op: 1, t: "#account" }; + const body = { + seq, + did, + time: new Date().toISOString(), + active, + ...(status ? { status } : {}), + }; + return encodeFrame(header, body); +} + +const DID_A = "did:plc:aaaaaa"; +const DID_B = "did:plc:bbbbbb"; + +// ============================================ +// SyncStorage GC methods +// ============================================ + +describe("SyncStorage GC methods", () => { + let tmpDir: string; + let db: InstanceType; + let storage: SyncStorage; + + beforeEach(() => { + tmpDir = mkdtempSync(join(tmpdir(), "gc-test-")); + db = new Database(join(tmpDir, "test.db")); + storage = new SyncStorage(db); + storage.initSchema(); + }); + + afterEach(() => { + db.close(); + rmSync(tmpDir, { recursive: true, force: true }); + }); + + // ---------- removeBlocks ---------- + + it("removes specific blocks for a DID", () => { + storage.upsertState({ did: DID_A, pdsEndpoint: "https://pds.example" }); + storage.trackBlocks(DID_A, ["cid1", "cid2", "cid3"]); + + expect(storage.getBlockCids(DID_A)).toHaveLength(3); + + storage.removeBlocks(DID_A, ["cid1", "cid3"]); + + const remaining = storage.getBlockCids(DID_A); + expect(remaining).toHaveLength(1); + expect(remaining).toContain("cid2"); + }); + + it("removeBlocks does nothing for empty array", () => { + storage.upsertState({ did: DID_A, pdsEndpoint: "https://pds.example" }); + storage.trackBlocks(DID_A, ["cid1"]); + storage.removeBlocks(DID_A, []); + expect(storage.getBlockCids(DID_A)).toHaveLength(1); + }); + + it("removeBlocks does not affect other DIDs", () => { + storage.upsertState({ did: DID_A, pdsEndpoint: "https://a.example" }); + storage.upsertState({ did: DID_B, pdsEndpoint: "https://b.example" }); + storage.trackBlocks(DID_A, ["cid1", "cid2"]); + storage.trackBlocks(DID_B, ["cid1", "cid3"]); + + storage.removeBlocks(DID_A, ["cid1"]); + + expect(storage.getBlockCids(DID_A)).toEqual(["cid2"]); + expect(storage.getBlockCids(DID_B)).toHaveLength(2); + expect(storage.getBlockCids(DID_B)).toContain("cid1"); + }); + + // ---------- findOrphanedCids ---------- + + it("finds orphaned CIDs with no remaining references", () => { + storage.upsertState({ did: DID_A, pdsEndpoint: "https://a.example" }); + storage.upsertState({ did: DID_B, pdsEndpoint: "https://b.example" }); + storage.trackBlocks(DID_A, ["cid1", "cid2"]); + storage.trackBlocks(DID_B, ["cid1", "cid3"]); + + // Remove DID_A's reference to cid1 and cid2 + storage.removeBlocks(DID_A, ["cid1", "cid2"]); + + // cid1 still referenced by DID_B, cid2 is orphaned + const orphaned = storage.findOrphanedCids(["cid1", "cid2"]); + expect(orphaned).toEqual(["cid2"]); + }); + + it("findOrphanedCids returns all for completely removed CIDs", () => { + storage.upsertState({ did: DID_A, pdsEndpoint: "https://a.example" }); + storage.trackBlocks(DID_A, ["cid1", "cid2"]); + storage.removeBlocks(DID_A, ["cid1", "cid2"]); + + const orphaned = storage.findOrphanedCids(["cid1", "cid2"]); + expect(orphaned).toEqual(["cid1", "cid2"]); + }); + + it("findOrphanedCids returns empty for empty input", () => { + expect(storage.findOrphanedCids([])).toEqual([]); + }); + + // ---------- removeBlobs / findOrphanedBlobCids ---------- + + it("removes specific blobs for a DID", () => { + storage.upsertState({ did: DID_A, pdsEndpoint: "https://pds.example" }); + storage.trackBlobs(DID_A, ["blob1", "blob2", "blob3"]); + + storage.removeBlobs(DID_A, ["blob1", "blob3"]); + + const remaining = storage.getBlobCids(DID_A); + expect(remaining).toHaveLength(1); + expect(remaining).toContain("blob2"); + }); + + it("finds orphaned blob CIDs", () => { + storage.upsertState({ did: DID_A, pdsEndpoint: "https://a.example" }); + storage.upsertState({ did: DID_B, pdsEndpoint: "https://b.example" }); + storage.trackBlobs(DID_A, ["blob1", "blob2"]); + storage.trackBlobs(DID_B, ["blob1"]); + + storage.removeBlobs(DID_A, ["blob1", "blob2"]); + + // blob1 still referenced by DID_B + const orphaned = storage.findOrphanedBlobCids(["blob1", "blob2"]); + expect(orphaned).toEqual(["blob2"]); + }); + + // ---------- getBlockCidSet ---------- + + it("returns block CIDs as a Set", () => { + storage.upsertState({ did: DID_A, pdsEndpoint: "https://pds.example" }); + storage.trackBlocks(DID_A, ["cid1", "cid2", "cid3"]); + + const cidSet = storage.getBlockCidSet(DID_A); + expect(cidSet).toBeInstanceOf(Set); + expect(cidSet.size).toBe(3); + expect(cidSet.has("cid1")).toBe(true); + expect(cidSet.has("cid4")).toBe(false); + }); + + // ---------- needs_gc ---------- + + it("sets and clears needs_gc flag", () => { + storage.upsertState({ did: DID_A, pdsEndpoint: "https://pds.example" }); + + expect(storage.getDidsNeedingGc()).toEqual([]); + + storage.setNeedsGc(DID_A); + expect(storage.getDidsNeedingGc()).toEqual([DID_A]); + + storage.clearNeedsGc(DID_A); + expect(storage.getDidsNeedingGc()).toEqual([]); + }); + + it("getDidsNeedingGc returns multiple flagged DIDs", () => { + storage.upsertState({ did: DID_A, pdsEndpoint: "https://a.example" }); + storage.upsertState({ did: DID_B, pdsEndpoint: "https://b.example" }); + + storage.setNeedsGc(DID_A); + storage.setNeedsGc(DID_B); + + const dids = storage.getDidsNeedingGc(); + expect(dids).toHaveLength(2); + expect(dids).toContain(DID_A); + expect(dids).toContain(DID_B); + }); + + // ---------- markTombstoned ---------- + + it("marks a DID as tombstoned", () => { + storage.upsertState({ did: DID_A, pdsEndpoint: "https://pds.example" }); + + storage.markTombstoned(DID_A); + + const state = storage.getState(DID_A); + expect(state?.status).toBe("tombstoned"); + }); + + // ---------- purgeDidData ---------- + + it("purges all data for a DID and returns removed CIDs", () => { + storage.upsertState({ did: DID_A, pdsEndpoint: "https://pds.example" }); + storage.trackBlocks(DID_A, ["cid1", "cid2"]); + storage.trackBlobs(DID_A, ["blob1"]); + storage.trackRecordPaths(DID_A, ["app.bsky.feed.post/abc"]); + + const result = storage.purgeDidData(DID_A); + + expect(result.blocksRemoved).toHaveLength(2); + expect(result.blocksRemoved).toContain("cid1"); + expect(result.blocksRemoved).toContain("cid2"); + expect(result.blobsRemoved).toEqual(["blob1"]); + + // All tracking data should be gone + expect(storage.getState(DID_A)).toBeNull(); + expect(storage.getBlockCids(DID_A)).toEqual([]); + expect(storage.getBlobCids(DID_A)).toEqual([]); + expect(storage.getRecordPaths(DID_A)).toEqual([]); + }); + + it("purgeDidData does not affect other DIDs", () => { + storage.upsertState({ did: DID_A, pdsEndpoint: "https://a.example" }); + storage.upsertState({ did: DID_B, pdsEndpoint: "https://b.example" }); + storage.trackBlocks(DID_A, ["cid1"]); + storage.trackBlocks(DID_B, ["cid1", "cid2"]); + + storage.purgeDidData(DID_A); + + expect(storage.getState(DID_B)).not.toBeNull(); + expect(storage.getBlockCids(DID_B)).toHaveLength(2); + }); + + // ---------- needsGc in SyncState ---------- + + it("SyncState includes needsGc field", () => { + storage.upsertState({ did: DID_A, pdsEndpoint: "https://pds.example" }); + + let state = storage.getState(DID_A); + expect(state?.needsGc).toBe(false); + + storage.setNeedsGc(DID_A); + state = storage.getState(DID_A); + expect(state?.needsGc).toBe(true); + }); + + // ---------- tombstoned status ---------- + + it("tombstoned status persists in SyncState", () => { + storage.upsertState({ did: DID_A, pdsEndpoint: "https://pds.example" }); + storage.markTombstoned(DID_A); + + const state = storage.getState(DID_A); + expect(state?.status).toBe("tombstoned"); + }); + + it("tombstoned DIDs appear in getAllStates", () => { + storage.upsertState({ did: DID_A, pdsEndpoint: "https://a.example" }); + storage.upsertState({ did: DID_B, pdsEndpoint: "https://b.example" }); + storage.markTombstoned(DID_A); + + const states = storage.getAllStates(); + expect(states).toHaveLength(2); + const tombstoned = states.find((s) => s.did === DID_A); + expect(tombstoned?.status).toBe("tombstoned"); + }); + + // ---------- Cross-DID safety ---------- + + it("block GC respects cross-DID sharing", () => { + storage.upsertState({ did: DID_A, pdsEndpoint: "https://a.example" }); + storage.upsertState({ did: DID_B, pdsEndpoint: "https://b.example" }); + + // Both DIDs share cid1 + storage.trackBlocks(DID_A, ["cid1", "cid2"]); + storage.trackBlocks(DID_B, ["cid1", "cid3"]); + + // Remove cid1 from DID_A only + storage.removeBlocks(DID_A, ["cid1"]); + + // cid1 is NOT orphaned (still referenced by DID_B) + expect(storage.findOrphanedCids(["cid1"])).toEqual([]); + + // Now remove cid1 from DID_B too + storage.removeBlocks(DID_B, ["cid1"]); + + // Now cid1 IS orphaned + expect(storage.findOrphanedCids(["cid1"])).toEqual(["cid1"]); + }); +}); + +// ============================================ +// FirehoseSubscription #account events +// ============================================ + +describe("FirehoseSubscription #account events", () => { + let tmpDir: string; + let server: ReturnType; + let wss: InstanceType; + let port: number; + let subscription: FirehoseSubscription; + let connectedWs: InstanceType | null; + + beforeEach(async () => { + tmpDir = mkdtempSync(join(tmpdir(), "account-test-")); + connectedWs = null; + + const { createServer: createHttpServer } = await import("node:http"); + const { WebSocketServer: WSS } = await import("ws"); + + server = createHttpServer(); + wss = new WSS({ server }); + + wss.on("connection", (ws) => { + connectedWs = ws; + }); + + await new Promise((resolve) => { + server.listen(0, "127.0.0.1", () => resolve()); + }); + const addr = server.address(); + port = typeof addr === "object" && addr ? addr.port : 0; + + subscription = new FirehoseSubscription({ + firehoseUrl: `ws://127.0.0.1:${port}/xrpc/com.atproto.sync.subscribeRepos`, + }); + }); + + afterEach(async () => { + subscription.stop(); + wss.close(); + server.close(); + rmSync(tmpDir, { recursive: true, force: true }); + // Small delay for cleanup + await new Promise((r) => setTimeout(r, 100)); + }); + + it("dispatches #account events for tracked DIDs", async () => { + const events: FirehoseAccountEvent[] = []; + subscription.onAccount((event) => { + events.push(event); + }); + + subscription.start(new Set([DID_A])); + + // Wait for WS connection + await new Promise((resolve) => { + const check = () => { + if (connectedWs) return resolve(); + setTimeout(check, 20); + }; + check(); + }); + + // Send an account deactivation event + const frame = makeAccountFrame(1, DID_A, false, "deleted"); + connectedWs!.send(frame); + + // Wait for processing + await new Promise((r) => setTimeout(r, 200)); + + expect(events).toHaveLength(1); + expect(events[0]!.did).toBe(DID_A); + expect(events[0]!.active).toBe(false); + expect(events[0]!.status).toBe("deleted"); + }); + + it("ignores #account events for untracked DIDs", async () => { + const events: FirehoseAccountEvent[] = []; + subscription.onAccount((event) => { + events.push(event); + }); + + subscription.start(new Set([DID_A])); + + await new Promise((resolve) => { + const check = () => { + if (connectedWs) return resolve(); + setTimeout(check, 20); + }; + check(); + }); + + // Send event for untracked DID + const frame = makeAccountFrame(1, DID_B, false, "deleted"); + connectedWs!.send(frame); + + await new Promise((r) => setTimeout(r, 200)); + + expect(events).toHaveLength(0); + }); + + it("updates cursor for #account events", async () => { + subscription.onAccount(() => {}); + subscription.start(new Set([DID_A])); + + await new Promise((resolve) => { + const check = () => { + if (connectedWs) return resolve(); + setTimeout(check, 20); + }; + check(); + }); + + const frame = makeAccountFrame(42, DID_A, false, "takendown"); + connectedWs!.send(frame); + + await new Promise((r) => setTimeout(r, 200)); + + expect(subscription.getCursor()).toBe(42); + }); + + it("dispatches re-activation events", async () => { + const events: FirehoseAccountEvent[] = []; + subscription.onAccount((event) => { + events.push(event); + }); + + subscription.start(new Set([DID_A])); + + await new Promise((resolve) => { + const check = () => { + if (connectedWs) return resolve(); + setTimeout(check, 20); + }; + check(); + }); + + // Deactivate then reactivate + connectedWs!.send(makeAccountFrame(1, DID_A, false, "deactivated")); + connectedWs!.send(makeAccountFrame(2, DID_A, true)); + + await new Promise((r) => setTimeout(r, 300)); + + expect(events).toHaveLength(2); + expect(events[0]!.active).toBe(false); + expect(events[1]!.active).toBe(true); + }); +}); -- 2.51.2