diff --git a/sync/firehose/index.ts b/sync/firehose/index.ts index c4bbba8..41ba16c 100644 --- a/sync/firehose/index.ts +++ b/sync/firehose/index.ts @@ -68,6 +68,7 @@ export class Firehose { private sub: Subscription; private abortController: AbortController; private destoryDefer: Deferrable; + private matchCollection: ((col: string) => boolean) | null = null; constructor(public opts: FirehoseOptions) { this.destoryDefer = createDeferrable(); @@ -75,18 +76,37 @@ export class Firehose { if (this.opts.getCursor && this.opts.runner) { throw new Error("Must set only `getCursor` or `runner`"); } + if (opts.filterCollections) { + const exact = new Set(); + const prefixes: string[] = []; + + for (const pattern of opts.filterCollections) { + if (pattern.endsWith(".*")) { + prefixes.push(pattern.slice(0, -2)); + } else { + exact.add(pattern); + } + } + this.matchCollection = (col: string): boolean => { + if (exact.has(col)) return true; + for (const prefix of prefixes) { + if (col.startsWith(prefix)) return true; + } + return false; + }; + } this.sub = new Subscription({ ...opts, service: opts.service ?? "wss://bsky.network", method: "com.atproto.sync.subscribeRepos", signal: this.abortController.signal, - getParams: async () => { + getParams: () => { const getCursorFn = () => this.opts.runner?.getCursor() ?? this.opts.getCursor; if (!getCursorFn) { return undefined; } - const cursor = await getCursorFn(); + const cursor = getCursorFn(); return { cursor }; }, validate: (value: unknown) => { @@ -111,7 +131,7 @@ export class Firehose { const parsed = await this.parseEvt(evt); for (const write of parsed) { try { - await this.opts.handleEvent(write); + this.opts.handleEvent(write); } catch (err) { this.opts.onError(new FirehoseHandlerError(err, write)); } @@ -139,11 +159,11 @@ export class Firehose { try { if (isCommit(evt) && !this.opts.excludeCommit) { return this.opts.unauthenticatedCommits - ? await parseCommitUnauthenticated(evt, this.opts.filterCollections) + ? await parseCommitUnauthenticated(evt, this.matchCollection) : await parseCommitAuthenticated( this.opts.idResolver, evt, - this.opts.filterCollections, + this.matchCollection, ); } else if (isAccount(evt) && !this.opts.excludeAccount) { const parsed = parseAccount(evt); @@ -171,7 +191,7 @@ export class Firehose { const parsed = await this.parseEvt(evt); for (const write of parsed) { try { - await this.opts.handleEvent(write); + this.opts.handleEvent(write); } catch (err) { this.opts.onError(new FirehoseHandlerError(err, write)); } @@ -187,11 +207,11 @@ export class Firehose { export const parseCommitAuthenticated = async ( idResolver: IdResolver, evt: Commit, - filterCollections?: string[], + matchCollection?: ((col: string) => boolean) | null, forceKeyRefresh = false, ): Promise => { const did = evt.repo; - const ops = maybeFilterOps(evt.ops, filterCollections); + const ops = maybeFilterOps(evt.ops, matchCollection); if (ops.length === 0) { return []; } @@ -213,7 +233,7 @@ export const parseCommitAuthenticated = async ( }); } catch (err) { if (err instanceof RepoVerificationError && !forceKeyRefresh) { - return parseCommitAuthenticated(idResolver, evt, filterCollections, true); + return parseCommitAuthenticated(idResolver, evt, matchCollection, true); } throw err; } @@ -231,20 +251,20 @@ export const parseCommitAuthenticated = async ( export const parseCommitUnauthenticated = ( evt: Commit, - filterCollections?: string[], + matchCollection?: ((col: string) => boolean) | null, ): Promise => { - const ops = maybeFilterOps(evt.ops, filterCollections); + const ops = maybeFilterOps(evt.ops, matchCollection); return formatCommitOps(evt, ops); }; const maybeFilterOps = ( ops: RepoOp[], - filterCollections?: string[], + matchCollection?: ((col: string) => boolean) | null, ): RepoOp[] => { - if (!filterCollections) return ops; + if (!matchCollection) return ops; return ops.filter((op) => { const { collection } = parseDataKey(op.path); - return filterCollections.includes(collection); + return matchCollection(collection); }); }; diff --git a/sync/runner/memory-runner.ts b/sync/runner/memory-runner.ts index 0f343d7..784eb3b 100644 --- a/sync/runner/memory-runner.ts +++ b/sync/runner/memory-runner.ts @@ -6,6 +6,7 @@ export type MemoryRunnerOptions = { setCursor?: (cursor: number) => Promise; concurrency?: number; startCursor?: number; + setCursorInterval?: number; // milliseconds between cursor saves, 0 for immediate saves }; // A queue with arbitrarily many partitions, each processing work sequentially. @@ -15,10 +16,19 @@ export class MemoryRunner implements EventRunner { mainQueue: PQueue; partitions: Map = new Map(); cursor: number | undefined; + private lastSavedCursor: number | undefined; + private saveCursorTimer: number | undefined; + private readonly useInterval: boolean; + private readonly intervalMs: number; + private readonly setCursor: ((cursor: number) => Promise) | undefined; + private pendingSaveCursor: number | undefined; constructor(public opts: MemoryRunnerOptions = {}) { this.mainQueue = new PQueue({ concurrency: opts.concurrency ?? Infinity }); this.cursor = opts.startCursor; + this.setCursor = opts.setCursor; + this.intervalMs = opts.setCursorInterval ?? 0; + this.useInterval = this.intervalMs > 0; } getCursor(): number | undefined { @@ -50,13 +60,43 @@ export class MemoryRunner implements EventRunner { const latest = item.complete().at(-1); if (latest !== undefined) { this.cursor = latest; - if (this.opts.setCursor) { - await this.opts.setCursor(this.cursor); + if (this.setCursor) { + if (this.useInterval) { + this.scheduleIntervalSave(); + } else { + this.setCursor(this.cursor).catch(console.error); + this.lastSavedCursor = this.cursor; + } } } }); } + private scheduleIntervalSave(): void { + // Fast path: if cursor hasn't changed or timer already scheduled for this cursor + if ( + this.cursor === this.lastSavedCursor || + this.cursor === this.pendingSaveCursor + ) return; + + this.pendingSaveCursor = this.cursor; + + if (this.saveCursorTimer) { + clearTimeout(this.saveCursorTimer); + } + + this.saveCursorTimer = setTimeout(() => { + const cursorToSave = this.pendingSaveCursor!; + this.setCursor!(cursorToSave) + .then(() => { + this.lastSavedCursor = cursorToSave; + }) + .catch(console.error); + this.saveCursorTimer = undefined; + this.pendingSaveCursor = undefined; + }, this.intervalMs); + } + async processAll() { await this.mainQueue.onIdle(); } @@ -65,6 +105,26 @@ export class MemoryRunner implements EventRunner { this.mainQueue.pause(); this.mainQueue.clear(); this.partitions.forEach((p) => p.clear()); + + // Clear any pending cursor save timer and perform final save + if (this.saveCursorTimer) { + clearTimeout(this.saveCursorTimer); + this.saveCursorTimer = undefined; + } + + // Perform final cursor save if needed + if ( + this.setCursor && this.cursor !== undefined && + this.cursor !== this.lastSavedCursor + ) { + try { + await this.setCursor(this.cursor); + this.lastSavedCursor = this.cursor; + } catch (error) { + console.error("Failed to save cursor during destroy:", error); + } + } + await this.mainQueue.onIdle(); } }