From d9dbb7529698eacd6acd9688ed63d13d23581350 Mon Sep 17 00:00:00 2001 From: Evelyn Osman Date: Wed, 22 Jul 2026 10:33:51 +0200 Subject: [PATCH] bugfix: persist and resume the account firehose cursor --- server/src/sync.ts | 17 ++++++++++++++--- 1 file changed, 14 insertions(+), 3 deletions(-) diff --git a/server/src/sync.ts b/server/src/sync.ts index 5357d3b..3e2cdbe 100644 --- a/server/src/sync.ts +++ b/server/src/sync.ts @@ -77,6 +77,15 @@ function attachHeartbeat(ws: WebSocket, intervalMs = 30_000) { ws.once("close", () => clearInterval(timer)); } +const ACCOUNT_STREAM_CURSOR_KEY = "account_stream_cursor"; + +export function accountStreamUrl(hostname: string, cursor: string | null): string { + return ( + `wss://${hostname}/xrpc/com.atproto.sync.subscribeRepos` + + (cursor ? `?cursor=${cursor}` : "") + ); +} + function repoStatusToAccountStatus(repo: RepoEntry): "active" | "takendown" | "deactivated" { if (repo.active !== false) return "active"; return repo.status === "deactivated" ? "deactivated" : "takendown"; @@ -253,15 +262,17 @@ export class Syncer { private connectAccountStream() { if (this.stopped) return; - const ws = new WebSocket( - `wss://${this.pds.hostname}/xrpc/com.atproto.sync.subscribeRepos`, - ); + const cursor = getSyncState(this.db, ACCOUNT_STREAM_CURSOR_KEY); + const ws = new WebSocket(accountStreamUrl(this.pds.hostname, cursor)); this.sockets.push(ws); attachHeartbeat(ws); 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)); + } const did: string | undefined = body?.did ?? body?.repo; if (!did) return; // #account carries active/status directly; #identity means handle changed; -- 2.51.2