diff --git a/server/src/sync.test.ts b/server/src/sync.test.ts index 582e1b6..9ba4da9 100644 --- a/server/src/sync.test.ts +++ b/server/src/sync.test.ts @@ -48,6 +48,15 @@ test("fullSync is single-flight", async () => { assert.equal(calls, 1); }); +test("account cursor persists on flush, not per frame", () => { + const db = openDb(":memory:"); + const syncer = new Syncer(db, { hostname: "pds.test" } as unknown as PdsClient, [], stubLogger); + (syncer as any).accountCursor = 99; + assert.equal(getSyncState(db, "account_stream_cursor"), null); + (syncer as any).flushCursor(); + assert.equal(getSyncState(db, "account_stream_cursor"), "99"); +}); + test("account stream url resumes from the persisted cursor", () => { const db = openDb(":memory:"); assert.equal( diff --git a/server/src/sync.ts b/server/src/sync.ts index 23ed2d4..90fa821 100644 --- a/server/src/sync.ts +++ b/server/src/sync.ts @@ -99,6 +99,8 @@ export class Syncer { private stopped = false; private syncing = false; private timers: NodeJS.Timeout[] = []; + // latest firehose seq, persisted on the flush tick — a write per frame would hammer SQLite + private accountCursor: number | null = null; constructor( private db: Db, @@ -131,6 +133,7 @@ export class Syncer { void this.flushDirty().catch((err) => this.log.error({ err }, "dirty flush failed")); try { this.flushActivity(); + this.flushCursor(); } catch (err) { this.log.error({ err }, "activity flush failed"); } @@ -142,13 +145,19 @@ export class Syncer { for (const labeler of this.labelers) this.connectLabelStream(labeler); } - /** Shutdown seam: stop reconnects, tear down sockets, cancel timers. */ + /** Shutdown seam: stop reconnects, tear down sockets, cancel timers, flush buffers. */ stop() { this.stopped = true; for (const t of this.timers) clearInterval(t); this.timers = []; for (const ws of this.sockets) ws.terminate(); this.sockets = []; + try { + this.flushActivity(); + this.flushCursor(); + } catch (err) { + this.log.error({ err }, "final flush failed"); + } } // ---- full reconcile ------------------------------------------------------ @@ -293,9 +302,7 @@ export class Syncer { ws.on("message", (data: Buffer) => { try { const { header, body } = readFrame(data); - if (typeof body?.seq === "number") { - setSyncState(this.db, ACCOUNT_STREAM_CURSOR_KEY, String(body.seq)); - } + if (typeof body?.seq === "number") this.accountCursor = body.seq; const did: string | undefined = body?.did ?? body?.repo; if (!did) return; // #account carries active/status directly; #identity means handle changed; @@ -304,8 +311,12 @@ export class Syncer { this.dirty.add(did); } if (header?.t === "#commit") { - // UTC hour buckets; the stats API regroups them into viewer-local days - const key = `${new Date().toISOString().slice(0, 13)}|${did}`; + // UTC hour buckets; the stats API regroups them into viewer-local days. + // Bucket by the event's own time so cursor-replayed backfill after downtime + // lands in the hours it happened, not in a spike at reconnect. + const eventTime = body.time ? new Date(body.time) : new Date(); + const at = Number.isNaN(eventTime.getTime()) ? new Date() : eventTime; + const key = `${at.toISOString().slice(0, 13)}|${did}`; this.activity.set(key, (this.activity.get(key) ?? 0) + 1); } } catch { @@ -402,6 +413,12 @@ export class Syncer { this.activity.clear(); } + private flushCursor() { + if (this.accountCursor == null) return; + setSyncState(this.db, ACCOUNT_STREAM_CURSOR_KEY, String(this.accountCursor)); + this.accountCursor = null; + } + // ---- live label streams -------------------------------------------------- private connectLabelStream(labeler: WatchedLabeler) {