diff --git a/.changeset/bootstrap-sources.md b/.changeset/bootstrap-sources.md new file mode 100644 index 0000000..e7bc460 --- /dev/null +++ b/.changeset/bootstrap-sources.md @@ -0,0 +1,5 @@ +--- +"@atmo-dev/contrail": minor +--- + +Add source-neutral snapshot and ordered-change contracts plus capture-first bootstrap orchestration for fresh projection generations. Anchor the legacy PDS discovery path before its first relay request so newly created repositories remain recoverable through Jetstream replay. diff --git a/packages/contrail/src/core/backfill.ts b/packages/contrail/src/core/backfill.ts index b10fa14..85383bb 100644 --- a/packages/contrail/src/core/backfill.ts +++ b/packages/contrail/src/core/backfill.ts @@ -34,6 +34,19 @@ const BACKFILL_RETRY_BASE_MS = 15 * 60_000; const BACKFILL_RETRY_MAX_MS = 48 * 60 * 60_000; const DEFAULT_SCHEDULED_MAX_ATTEMPTS = 10; const DERIVED_PROJECTIONS_DIRTY_KEY = "backfill_derived_projections_dirty"; +/** Legacy Jetstream v1 uses wall-clock microseconds as its replay coordinate. + * Keep the same overlap Atcute uses for multi-endpoint clock skew. The future + * source adapter replaces this compatibility marker with a source-owned epoch + * and cursor. */ +const INITIAL_CAPTURE_OVERLAP_US = 10_000_000; + +async function ensureInitialReplayBoundary(db: Database): Promise { + if ((await getLastCursor(db)) !== null) return; + await saveCursor( + db, + Math.max(0, Date.now() * 1000 - INITIAL_CAPTURE_OVERLAP_US), + ); +} export interface BackfillCollectionMetrics { requests: number; @@ -683,11 +696,9 @@ async function backfillPendingWork( const metrics = emptyBackfillMetrics(); const aggregateDiagnostics: IngestDiagnosticCounts = {}; - // Anchor the jetstream cursor to now if it hasn't been set yet, so records - // emitted during backfill are replayed once jetstream starts. - if ((await getLastCursor(db)) === null) { - await saveCursor(db, Date.now() * 1000); - } + // Direct callers may begin with already-discovered work. Capture before the + // first PDS request so changes racing the sampled scan remain replayable. + await ensureInitialReplayBoundary(db); // Mark the set-based catch-up dirty before canonical writes. A crash can leave // search/count projections stale, so status stays incomplete and the next @@ -1158,6 +1169,11 @@ export async function discoverDIDs( const relays = config.relays ?? DEFAULT_RELAYS; if (relays.length === 0 || collections.length === 0) return []; + // Capture before relay discovery as well as before PDS crawling. Otherwise a + // repository created after its relay page was scanned but before the later + // PDS phase could fall before the live cursor and disappear from both paths. + await ensureInitialReplayBoundary(db); + const discovered: string[] = []; await ensureDiscoveryRows(db, collections, relays); diff --git a/packages/contrail/src/core/sources.ts b/packages/contrail/src/core/sources.ts new file mode 100644 index 0000000..1ed787f --- /dev/null +++ b/packages/contrail/src/core/sources.ts @@ -0,0 +1,311 @@ +/** Source-neutral contracts for building a fresh projection generation. */ + +export interface SourcePosition { + /** Stable logical stream identifier. */ + source: string; + /** Continuity epoch. Cursors from different epochs are never comparable. */ + epoch: string; + /** Opaque cursor interpreted only by the source adapter. */ + cursor: string; +} + +export interface SourceSemantics { + ordinaryRecords: boolean; + ordinaryDeletes: boolean; + accountLifecycle: boolean; + repositoryReplacement: boolean; + verifiedCommits: boolean; + explicitHead: boolean; +} + +export type CollectionCoverage = + | { state: "complete" } + | { state: "partial"; reason: string; unresolved?: number } + | { state: "gap"; reason: string }; + +export interface SnapshotRecord { + uri: string; + did: string; + collection: string; + rkey: string; + cid: string; + value: unknown; +} + +export interface PreparedSnapshot { + /** Provider-owned immutable snapshot identifier. */ + id: string; + provider: string; + consistency: "sampled-current-state" | "point-in-time"; + collections: Record; + semantics: SourceSemantics; + /** Upstream position represented by a point-in-time snapshot, when known. */ + through?: SourcePosition; +} + +export interface SnapshotBatch { + records: SnapshotRecord[]; + /** Opaque resume token within this exact prepared snapshot. */ + progress: string; + /** True only after every requested collection has been emitted. */ + done: boolean; +} + +export interface SnapshotSource { + readonly id: string; + /** Prepare and pin a snapshot. Record acquisition must not begin before this + * call; the bootstrap coordinator marks the change source first. */ + prepare(options: { + collections: string[]; + signal?: AbortSignal; + }): Promise; + read(options: { + snapshot: PreparedSnapshot; + progress?: string; + signal?: AbortSignal; + }): AsyncIterable; +} + +interface MutationBase { + uri: string; + did: string; + collection: string; + rkey: string; + revision?: string; + sourceTimeUs: number; + /** Per-event position when the source exposes one. */ + position?: SourcePosition; +} + +export type SourceMutation = + | (MutationBase & { + operation: "put"; + cid: string; + value: unknown; + }) + | (MutationBase & { + operation: "delete"; + }); + +export interface MutationBatch { + mutations: SourceMutation[]; + /** Everything through this position has been accounted for, including + * filtered events and an otherwise empty batch. */ + checkpoint: SourcePosition; + /** True when the requested through-position has been reached exactly. */ + caughtUp: boolean; +} + +export interface ChangeSource { + readonly id: string; + /** Return a durable replay coordinate near the current source head. */ + mark(options: { + collections: string[]; + signal?: AbortSignal; + }): Promise; + read(options: { + collections: string[]; + after: SourcePosition; + through: SourcePosition; + signal?: AbortSignal; + }): AsyncIterable; +} + +export type BootstrapPhase = "snapshot" | "catchup" | "complete"; + +/** Durable coordinator state. Snapshot progress and mutation checkpoints are + * separate because they belong to different cursor namespaces. */ +export interface BootstrapRunState { + phase: BootstrapPhase; + snapshot: PreparedSnapshot; + captureFrom: SourcePosition; + snapshotProgress: string | null; + snapshotComplete: boolean; + catchupThrough: SourcePosition | null; + changeCheckpoint: SourcePosition | null; +} + +/** Projection-owned persistence seam. Implementations commit records and the + * accompanying progress/checkpoint atomically in the destination database. */ +export interface BootstrapTarget { + load(): Promise; + begin(snapshot: PreparedSnapshot, captureFrom: SourcePosition): Promise; + applySnapshotBatch( + snapshot: PreparedSnapshot, + batch: SnapshotBatch, + ): Promise; + beginCatchup(through: SourcePosition): Promise; + applyMutationBatch(batch: MutationBatch): Promise; + complete(): Promise; +} + +export interface BootstrapResult { + snapshot: PreparedSnapshot; + captureFrom: SourcePosition; + through: SourcePosition; +} + +function positionsEqual(left: SourcePosition, right: SourcePosition): boolean { + return ( + left.source === right.source && + left.epoch === right.epoch && + left.cursor === right.cursor + ); +} + +function assertCompatiblePosition( + position: SourcePosition, + expected: SourcePosition, + label: string, +): void { + if ( + position.source !== expected.source || + position.epoch !== expected.epoch + ) { + throw new Error( + `${label} belongs to ${position.source}/${position.epoch}, expected ` + + `${expected.source}/${expected.epoch}`, + ); + } +} + +function assertCoverage( + snapshot: PreparedSnapshot, + collections: string[], + allowPartial: boolean, +): void { + for (const collection of collections) { + const coverage = snapshot.collections[collection]; + if (!coverage) { + throw new Error(`Snapshot ${snapshot.id} omitted ${collection}`); + } + if (coverage.state === "gap") { + throw new Error( + `Snapshot ${snapshot.id} has a gap for ${collection}: ${coverage.reason}`, + ); + } + if (coverage.state === "partial" && !allowPartial) { + throw new Error( + `Snapshot ${snapshot.id} is partial for ${collection}: ${coverage.reason}`, + ); + } + } +} + +/** + * Build one fresh projection using capture-first snapshot/replay semantics. + * + * The target is responsible for applying every batch through Contrail's normal + * admission/projector and persisting its progress in that same transaction. + * A prepared point-in-time snapshot may supply its own upstream boundary; + * sampled scans use the position marked before preparation begins. + */ +export async function bootstrapFreshProjection(options: { + collections: string[]; + snapshotSource: SnapshotSource; + changeSource: ChangeSource; + target: BootstrapTarget; + allowPartial?: boolean; + signal?: AbortSignal; +}): Promise { + const { + collections, + snapshotSource, + changeSource, + target, + signal, + } = options; + let state = await target.load(); + + if (!state) { + // Mark before snapshot preparation so even a provider that performs relay + // discovery while preparing cannot open a capture gap. + const marked = await changeSource.mark({ collections, signal }); + const snapshot = await snapshotSource.prepare({ collections, signal }); + assertCoverage(snapshot, collections, options.allowPartial === true); + const captureFrom = snapshot.through ?? marked; + assertCompatiblePosition(captureFrom, marked, "Snapshot boundary"); + await target.begin(snapshot, captureFrom); + state = { + phase: "snapshot", + snapshot, + captureFrom, + snapshotProgress: null, + snapshotComplete: false, + catchupThrough: null, + changeCheckpoint: null, + }; + } else { + assertCoverage(state.snapshot, collections, options.allowPartial === true); + } + + if (!state.snapshotComplete) { + let sawDone = false; + for await (const batch of snapshotSource.read({ + snapshot: state.snapshot, + ...(state.snapshotProgress === null + ? {} + : { progress: state.snapshotProgress }), + signal, + })) { + if (sawDone) { + throw new Error(`Snapshot ${state.snapshot.id} emitted data after done`); + } + await target.applySnapshotBatch(state.snapshot, batch); + state.snapshotProgress = batch.progress; + state.snapshotComplete = batch.done; + sawDone = batch.done; + } + if (!state.snapshotComplete) { + throw new Error(`Snapshot ${state.snapshot.id} ended before done`); + } + } + + if (!state.catchupThrough) { + const through = await changeSource.mark({ collections, signal }); + assertCompatiblePosition(through, state.captureFrom, "Catch-up target"); + await target.beginCatchup(through); + state.catchupThrough = through; + state.phase = "catchup"; + } + + if (state.phase !== "complete") { + const after = state.changeCheckpoint ?? state.captureFrom; + assertCompatiblePosition(after, state.catchupThrough, "Catch-up cursor"); + let caughtUp = positionsEqual(after, state.catchupThrough); + + if (!caughtUp) { + for await (const batch of changeSource.read({ + collections, + after, + through: state.catchupThrough, + signal, + })) { + assertCompatiblePosition( + batch.checkpoint, + state.catchupThrough, + "Mutation checkpoint", + ); + if (batch.caughtUp && !positionsEqual(batch.checkpoint, state.catchupThrough)) { + throw new Error("Change source reported caught up at the wrong position"); + } + await target.applyMutationBatch(batch); + state.changeCheckpoint = batch.checkpoint; + caughtUp = batch.caughtUp; + if (caughtUp) break; + } + } + + if (!caughtUp) { + throw new Error("Change source ended before the catch-up target"); + } + await target.complete(); + state.phase = "complete"; + } + + return { + snapshot: state.snapshot, + captureFrom: state.captureFrom, + through: state.catchupThrough, + }; +} diff --git a/packages/contrail/src/index.ts b/packages/contrail/src/index.ts index 6ce88b9..96058d0 100644 --- a/packages/contrail/src/index.ts +++ b/packages/contrail/src/index.ts @@ -17,6 +17,7 @@ export type { ResolvedIdentity } from "./core/client"; // Ingestion and maintenance. export * from "./core/ingest"; +export * from "./core/sources"; export * from "./core/jetstream"; export * from "./core/persistent"; export * from "./core/backfill"; diff --git a/packages/contrail/tests/backfill-status.test.ts b/packages/contrail/tests/backfill-status.test.ts index 4ca02d1..a431268 100644 --- a/packages/contrail/tests/backfill-status.test.ts +++ b/packages/contrail/tests/backfill-status.test.ts @@ -673,6 +673,40 @@ describe("scheduled backfill retries", () => { }); describe("discovery failure state", () => { + it("anchors replay before the first relay discovery request", async () => { + const db = await createTestDbWithSchema(); + const config = resolveConfig({ + namespace: "com.example", + collections: { event: { collection: EVENT } }, + relays: ["https://relay.test"], + }); + let cursorDuringRequest: number | null = null; + const fetchSpy = vi.spyOn(global, "fetch").mockImplementation(async () => { + cursorDuringRequest = + ( + await db + .prepare("SELECT time_us FROM cursor WHERE id = 1") + .first<{ time_us: number }>() + )?.time_us ?? null; + return new Response(JSON.stringify({ repos: [] }), { + status: 200, + headers: { "content-type": "application/json" }, + }); + }); + + try { + const startedAtUs = Date.now() * 1000; + await discoverDIDs(db, config, Infinity); + expect(cursorDuringRequest).not.toBeNull(); + expect(cursorDuringRequest!).toBeGreaterThanOrEqual( + startedAtUs - 10_000_000, + ); + expect(cursorDuringRequest!).toBeLessThanOrEqual(Date.now() * 1000); + } finally { + fetchSpy.mockRestore(); + } + }); + it("keeps a failed relay pending and falls back to another relay", async () => { vi.useFakeTimers(); const fetchSpy = vi.spyOn(global, "fetch").mockImplementation(async (input) => { diff --git a/packages/contrail/tests/bootstrap-sources.test.ts b/packages/contrail/tests/bootstrap-sources.test.ts new file mode 100644 index 0000000..f8f36e2 --- /dev/null +++ b/packages/contrail/tests/bootstrap-sources.test.ts @@ -0,0 +1,387 @@ +import { describe, expect, it } from "vitest"; +import { + bootstrapFreshProjection, + type BootstrapRunState, + type BootstrapTarget, + type ChangeSource, + type MutationBatch, + type PreparedSnapshot, + type SnapshotBatch, + type SnapshotRecord, + type SnapshotSource, + type SourceMutation, + type SourcePosition, +} from "../src/index"; + +const COLLECTION = "com.example.event"; +const ALICE = "did:plc:alice"; + +function position(cursor: number): SourcePosition { + return { source: "test-stream", epoch: "one", cursor: String(cursor) }; +} + +function record(rkey: string, value: number): SnapshotRecord { + return { + uri: `at://${ALICE}/${COLLECTION}/${rkey}`, + did: ALICE, + collection: COLLECTION, + rkey, + cid: `cid-${rkey}-${value}`, + value: { value }, + }; +} + +function put(rkey: string, value: number, cursor: number): SourceMutation { + return { + operation: "put", + ...record(rkey, value), + sourceTimeUs: cursor, + position: position(cursor), + }; +} + +function deletion(rkey: string, cursor: number): SourceMutation { + return { + operation: "delete", + uri: `at://${ALICE}/${COLLECTION}/${rkey}`, + did: ALICE, + collection: COLLECTION, + rkey, + sourceTimeUs: cursor, + position: position(cursor), + }; +} + +const semantics = { + ordinaryRecords: true, + ordinaryDeletes: true, + accountLifecycle: false, + repositoryReplacement: false, + verifiedCommits: false, + explicitHead: true, +}; + +class MemoryTarget implements BootstrapTarget { + state: BootstrapRunState | null = null; + records = new Map(); + snapshotBatches = 0; + mutationBatches = 0; + + async load() { + return this.state ? structuredClone(this.state) : null; + } + + async begin(snapshot: PreparedSnapshot, captureFrom: SourcePosition) { + this.state = { + phase: "snapshot", + snapshot: structuredClone(snapshot), + captureFrom: structuredClone(captureFrom), + snapshotProgress: null, + snapshotComplete: false, + catchupThrough: null, + changeCheckpoint: null, + }; + } + + async applySnapshotBatch(_snapshot: PreparedSnapshot, batch: SnapshotBatch) { + for (const item of batch.records) this.records.set(item.uri, item); + this.snapshotBatches++; + this.state!.snapshotProgress = batch.progress; + this.state!.snapshotComplete = batch.done; + } + + async beginCatchup(through: SourcePosition) { + this.state!.phase = "catchup"; + this.state!.catchupThrough = structuredClone(through); + } + + async applyMutationBatch(batch: MutationBatch) { + for (const mutation of batch.mutations) { + if (mutation.operation === "delete") { + this.records.delete(mutation.uri); + } else { + this.records.set(mutation.uri, { + uri: mutation.uri, + did: mutation.did, + collection: mutation.collection, + rkey: mutation.rkey, + cid: mutation.cid, + value: mutation.value, + }); + } + } + this.mutationBatches++; + this.state!.changeCheckpoint = structuredClone(batch.checkpoint); + } + + async complete() { + this.state!.phase = "complete"; + } +} + +function snapshotSource(options: { + snapshot: PreparedSnapshot; + batches(progress?: string): SnapshotBatch[]; + calls?: string[]; +}): SnapshotSource { + return { + id: options.snapshot.provider, + async prepare() { + options.calls?.push("prepare"); + return structuredClone(options.snapshot); + }, + async *read({ progress }) { + options.calls?.push(`snapshot:${progress ?? "start"}`); + for (const batch of options.batches(progress)) yield structuredClone(batch); + }, + }; +} + +function changeSource(options: { + marks: number[]; + mutations: SourceMutation[]; + calls?: string[]; +}): ChangeSource { + let markIndex = 0; + return { + id: "changes", + async mark() { + const cursor = options.marks[markIndex++]; + if (cursor === undefined) throw new Error("Unexpected mark"); + options.calls?.push(`mark:${cursor}`); + return position(cursor); + }, + async *read({ after, through }) { + options.calls?.push(`changes:${after.cursor}-${through.cursor}`); + const lower = Number(after.cursor); + const upper = Number(through.cursor); + const mutations = options.mutations.filter((mutation) => { + const cursor = Number(mutation.position?.cursor); + return cursor > lower && cursor <= upper; + }); + yield { + mutations, + checkpoint: structuredClone(through), + caughtUp: true, + }; + }, + }; +} + +function prepared(overrides: Partial = {}): PreparedSnapshot { + return { + id: "snapshot-one", + provider: "test-snapshot", + consistency: "sampled-current-state", + collections: { [COLLECTION]: { state: "complete" } }, + semantics, + ...overrides, + }; +} + +describe("bootstrap source orchestration", () => { + it("marks before a sampled scan and replays mutations that raced it", async () => { + const calls: string[] = []; + const snapshot = snapshotSource({ + snapshot: prepared(), + calls, + batches: () => [ + { + // Alice was sampled at different moments: A already reflects cursor + // 3, B is stale at cursor 2, and C did not exist when sampled. + records: [record("a", 3), record("b", 2)], + progress: "done", + done: true, + }, + ], + }); + const changes = changeSource({ + marks: [2, 5], + calls, + mutations: [put("a", 3, 3), deletion("b", 4), put("c", 5, 5)], + }); + const target = new MemoryTarget(); + + const result = await bootstrapFreshProjection({ + collections: [COLLECTION], + snapshotSource: snapshot, + changeSource: changes, + target, + }); + + expect(calls).toEqual([ + "mark:2", + "prepare", + "snapshot:start", + "mark:5", + "changes:2-5", + ]); + expect(result.captureFrom).toEqual(position(2)); + expect(result.through).toEqual(position(5)); + expect( + [...target.records.values()] + .map((item) => [item.rkey, (item.value as { value: number }).value]) + .sort(), + ).toEqual([ + ["a", 3], + ["c", 5], + ]); + expect(target.state?.phase).toBe("complete"); + }); + + it("uses a point-in-time snapshot boundary instead of the preliminary mark", async () => { + const calls: string[] = []; + const snapshot = snapshotSource({ + snapshot: prepared({ + consistency: "point-in-time", + through: position(5), + }), + calls, + batches: () => [ + { records: [record("a", 5)], progress: "done", done: true }, + ], + }); + const changes = changeSource({ + // The preliminary mark happens before the manifest is pinned. The + // snapshot's own boundary is the correct tail starting point. + marks: [10, 12], + calls, + mutations: [put("a", 6, 6), put("b", 11, 11)], + }); + const target = new MemoryTarget(); + + const result = await bootstrapFreshProjection({ + collections: [COLLECTION], + snapshotSource: snapshot, + changeSource: changes, + target, + }); + + expect(result.captureFrom).toEqual(position(5)); + expect(calls.at(-1)).toBe("changes:5-12"); + expect( + [...target.records.values()] + .map((item) => [item.rkey, (item.value as { value: number }).value]) + .sort(), + ).toEqual([ + ["a", 6], + ["b", 11], + ]); + }); + + it("resumes a pinned snapshot from its last committed progress", async () => { + const reads: Array = []; + let failSecondBatch = true; + const snapshot = snapshotSource({ + snapshot: prepared(), + batches(progress) { + reads.push(progress); + if (progress === "part-1") { + return [ + { records: [record("b", 2)], progress: "done", done: true }, + ]; + } + return [ + { records: [record("a", 1)], progress: "part-1", done: false }, + { records: [record("b", 2)], progress: "done", done: true }, + ]; + }, + }); + const changes = changeSource({ marks: [1, 3], mutations: [] }); + const target = new MemoryTarget(); + const apply = target.applySnapshotBatch.bind(target); + target.applySnapshotBatch = async (preparedSnapshot, batch) => { + if (batch.progress === "done" && failSecondBatch) { + failSecondBatch = false; + throw new Error("injected snapshot failure"); + } + await apply(preparedSnapshot, batch); + }; + + await expect( + bootstrapFreshProjection({ + collections: [COLLECTION], + snapshotSource: snapshot, + changeSource: changes, + target, + }), + ).rejects.toThrow("injected snapshot failure"); + + await bootstrapFreshProjection({ + collections: [COLLECTION], + snapshotSource: snapshot, + changeSource: changes, + target, + }); + + expect(reads).toEqual([undefined, "part-1"]); + expect(target.snapshotBatches).toBe(2); + expect([...target.records.values()].map((item) => item.rkey).sort()).toEqual([ + "a", + "b", + ]); + }); + + it("commits an empty change batch to prove progress through the target", async () => { + const snapshot = snapshotSource({ + snapshot: prepared(), + batches: () => [{ records: [], progress: "done", done: true }], + }); + const target = new MemoryTarget(); + + await bootstrapFreshProjection({ + collections: [COLLECTION], + snapshotSource: snapshot, + changeSource: changeSource({ marks: [1, 2], mutations: [] }), + target, + }); + + expect(target.mutationBatches).toBe(1); + expect(target.state?.changeCheckpoint).toEqual(position(2)); + expect(target.state?.phase).toBe("complete"); + }); + + it("refuses a point-in-time boundary from another source epoch", async () => { + const snapshot = snapshotSource({ + snapshot: prepared({ + consistency: "point-in-time", + through: { source: "test-stream", epoch: "old", cursor: "5" }, + }), + batches: () => [], + }); + const target = new MemoryTarget(); + + await expect( + bootstrapFreshProjection({ + collections: [COLLECTION], + snapshotSource: snapshot, + changeSource: changeSource({ marks: [10], mutations: [] }), + target, + }), + ).rejects.toThrow("Snapshot boundary belongs to test-stream/old"); + expect(target.state).toBeNull(); + }); + + it("refuses partial and gapped coverage by default", async () => { + for (const coverage of [ + { state: "partial" as const, reason: "one PDS unavailable" }, + { state: "gap" as const, reason: "source retention expired" }, + ]) { + const snapshot = snapshotSource({ + snapshot: prepared({ collections: { [COLLECTION]: coverage } }), + batches: () => [], + }); + const target = new MemoryTarget(); + + await expect( + bootstrapFreshProjection({ + collections: [COLLECTION], + snapshotSource: snapshot, + changeSource: changeSource({ marks: [1], mutations: [] }), + target, + }), + ).rejects.toThrow(coverage.reason); + expect(target.state).toBeNull(); + } + }); +});