From 91278c6d423a6295ed3316f242712a97df49bcde Mon Sep 17 00:00:00 2001 From: Tom Scanlan Date: Wed, 10 Jun 2026 13:41:12 -0400 Subject: [PATCH] fix(ingest): bound the ingest cycle by its safety timeout ingestEvents drove the Jetstream subscription with `for await` and checked its exit conditions (caught-up and safety timeout) only inside the loop body, after an event arrived. On a low-traffic stream the async iterator blocks forever once history is replayed: the safety timeout never fires, the caller's hard timeout kills the cycle (on Workers: "Exceeded CPU Limit"), and the collected batch and cursor are never written. The per-event filter paths (unknown DID, subject filter, recordFilter) used `continue`, which also skipped the exit checks, so a flood of all-filtered events ran past the deadline too. Drive the iterator manually: check the deadline before awaiting each event, race iterator.next() against the remaining deadline so a silent stream can't outlast the safety timeout, and run the caught-up and deadline checks regardless of whether the event was filtered. Close the subscription socket best-effort on exit (not awaited, since awaiting a quiet stream's close could reintroduce the hang). Add tests covering the quiet-stream and all-filtered-flood cases. --- .changeset/ingest-quiet-stream-hang.md | 7 + .../contrail-appview/src/core/jetstream.ts | 68 ++++++-- packages/contrail/tests/ingest-hang.test.ts | 152 ++++++++++++++++++ 3 files changed, 218 insertions(+), 9 deletions(-) create mode 100644 .changeset/ingest-quiet-stream-hang.md create mode 100644 packages/contrail/tests/ingest-hang.test.ts diff --git a/.changeset/ingest-quiet-stream-hang.md b/.changeset/ingest-quiet-stream-hang.md new file mode 100644 index 0000000..d8c3aee --- /dev/null +++ b/.changeset/ingest-quiet-stream-hang.md @@ -0,0 +1,7 @@ +--- +"@atmo-dev/contrail-appview": patch +--- + +fix(ingest): stop the ingest cycle hanging on a quiet or all-filtered Jetstream + +Race `iterator.next()` against the safety timeout and check the exit conditions before each event, so a low-traffic stream no longer blocks until the caller's hard timeout kills the cycle with the batch and cursor unwritten. diff --git a/packages/contrail-appview/src/core/jetstream.ts b/packages/contrail-appview/src/core/jetstream.ts index fc48449..f96c78d 100644 --- a/packages/contrail-appview/src/core/jetstream.ts +++ b/packages/contrail-appview/src/core/jetstream.ts @@ -60,6 +60,28 @@ function getLogger(config: ContrailConfig): Logger { return config.logger ?? console; } +/** Sentinel returned by `nextWithDeadline` when the wait timed out. */ +const INGEST_TIMEOUT = Symbol("ingest-timeout"); + +/** Await the iterator's next value, but give up after `ms`. Without this a + * quiet Jetstream (the async iterator blocks forever waiting for an event that + * never arrives) holds the cycle past its safety timeout until the caller's + * hard timeout kills the isolate — so the batch and cursor are never written. */ +function nextWithDeadline( + iterator: AsyncIterator, + ms: number +): Promise | typeof INGEST_TIMEOUT> { + let timer: ReturnType; + const next = iterator.next(); + // If the timeout wins this race the next() promise stays pending; swallow a + // later rejection so it can't surface as an unhandled rejection. + next.catch(() => {}); + const timeout = new Promise((resolve) => { + timer = setTimeout(() => resolve(INGEST_TIMEOUT), ms); + }); + return Promise.race([next, timeout]).finally(() => clearTimeout(timer)); +} + export async function ingestEvents( config: ContrailConfig, cursor: number | null, @@ -111,9 +133,13 @@ export async function ingestEvents( }, }); - for await (const event of subscription) { - if (firstYieldedTimeUs === null) firstYieldedTimeUs = event.time_us; - lastYieldedTimeUs = event.time_us; + const iterator = subscription[Symbol.asyncIterator](); + type Ev = typeof subscription extends AsyncIterable ? V : never; + + // Collect or skip a single event. Filtering uses early `return` rather than + // the loop's `continue` so the loop's exit checks still run after a filtered + // event (a stream of all-filtered events must not skip the deadline). + const handleEvent = (event: Ev): void => { if (event.kind === "commit") { const { commit } = event; totalCommits++; @@ -127,7 +153,7 @@ export async function ingestEvents( if (!knownDids.has(event.did)) { filteredUnknownDid++; if (filteredDidSamples.size < 10) filteredDidSamples.add(event.did); - continue; + return; } // Subject filter: for collections with subjectField (e.g. follows // pointing at a `subject` DID), drop records whose subject isn't a @@ -139,7 +165,7 @@ export async function ingestEvents( subjectField ]; if (typeof subj === "string" && !knownDids.has(subj)) { - continue; + return; } } } @@ -152,7 +178,7 @@ export async function ingestEvents( } catch (err) { log.warn(`[ingest] recordFilter threw for ${uri}: ${err}`); } - if (!keep) continue; + if (!keep) return; } const prev = seenUris.get(uri); @@ -195,22 +221,46 @@ export async function ingestEvents( } else if (event.kind === "identity") { identityUpdates.set(event.did, event.identity.handle); } + }; - if (event.time_us >= startTimeUs) { + for (;;) { + // Run the exit checks BEFORE awaiting the next event and regardless of + // whether the previous event was filtered — otherwise a quiet stream blocks + // forever and an all-filtered flood never reaches the deadline check. + if (Date.now() >= deadline) { log.log( - `[ingest] caught up to present, stopping (last time_us=${event.time_us}, startTimeUs=${startTimeUs})` + `[ingest] safety timeout reached, stopping (deadline=${deadline}, collected=${collected.length})` ); break; } - if (Date.now() >= deadline) { + const step = await nextWithDeadline(iterator, Math.max(0, deadline - Date.now())); + if (step === INGEST_TIMEOUT) { log.log( `[ingest] safety timeout reached, stopping (deadline=${deadline}, collected=${collected.length})` ); break; } + if (step.done) break; + const event = step.value; + + if (firstYieldedTimeUs === null) firstYieldedTimeUs = event.time_us; + lastYieldedTimeUs = event.time_us; + + handleEvent(event); + + if (event.time_us >= startTimeUs) { + log.log( + `[ingest] caught up to present, stopping (last time_us=${event.time_us}, startTimeUs=${startTimeUs})` + ); + break; + } } + // Close the subscription's socket, but fire-and-forget: awaiting the + // iterator's return on a quiet stream could itself block (the hang we fix). + Promise.resolve(iterator.return?.()).catch(() => {}); + if (filteredUnknownDid > 0) { const sample = [...filteredDidSamples].join(", "); log.log( diff --git a/packages/contrail/tests/ingest-hang.test.ts b/packages/contrail/tests/ingest-hang.test.ts new file mode 100644 index 0000000..22a9499 --- /dev/null +++ b/packages/contrail/tests/ingest-hang.test.ts @@ -0,0 +1,152 @@ +import { describe, it, expect, vi, beforeEach } from "vitest"; + +// Controller the mocked Jetstream reads from. `script` is an async-generator +// factory each test sets; `abort` lets a test stop a still-running flood once +// its assertions are done (so the buggy path can't leak timers post-failure). +const jetstream = vi.hoisted(() => ({ + script: null as + | null + | ((self: { cursor: number | null }) => AsyncGenerator), + abort: false, +})); + +// Replace the real WebSocket-backed subscription with one driven by the test's +// `script`. `self.cursor` mirrors the real subscription's progress cursor, +// which ingestEvents reads back as `lastCursor`. +vi.mock("@atcute/jetstream", () => { + class MockJetstreamSubscription { + cursor: number | null = null; + constructor(opts: { cursor?: number }) { + this.cursor = typeof opts?.cursor === "number" ? opts.cursor : null; + } + async *[Symbol.asyncIterator]() { + if (!jetstream.script) throw new Error("test did not set jetstream.script"); + yield* jetstream.script(this); + } + } + return { JetstreamSubscription: MockJetstreamSubscription }; +}); + +import { ingestEvents } from "../src/core/jetstream"; +import { resolveConfig } from "../src/core/types"; +import type { ContrailConfig } from "../src/core/types"; + +const silentLogger = { log() {}, warn() {}, error() {} }; + +function commitEvent( + did: string, + collection: string, + time_us: number, + rkey: string, +) { + return { + kind: "commit" as const, + time_us, + did, + commit: { + collection, + operation: "create" as const, + rkey, + cid: "bafy" + rkey, + record: { name: "Test Event", startsAt: "2026-04-01T10:00:00Z", mode: "online" }, + }, + }; +} + +/** A discoverable-only config — events flow straight into `collected`. */ +function discoverableConfig(): ContrailConfig { + return { + ...resolveConfig({ + namespace: "com.example", + collections: { + event: { collection: "community.lexicon.calendar.event" }, + }, + }), + logger: silentLogger, + }; +} + +/** A config with a dependent collection so unknown-DID events get filtered. */ +function dependentConfig(): ContrailConfig { + return { + ...resolveConfig({ + namespace: "com.example", + collections: { + event: { collection: "community.lexicon.calendar.event" }, + follow: { collection: "app.bsky.graph.follow", discover: false }, + }, + }), + logger: silentLogger, + }; +} + +/** Reject if `p` hasn't settled within `ms` — turns a hang into a test failure + * instead of a frozen run. */ +function withTimeout(p: Promise, ms: number, label: string): Promise { + let timer: ReturnType; + const timeout = new Promise((_, reject) => { + timer = setTimeout( + () => reject(new Error(`ingestEvents did not return within ${ms}ms (${label})`)), + ms, + ); + }); + return Promise.race([p.finally(() => clearTimeout(timer)), timeout]); +} + +describe("ingestEvents — bounded by the safety timeout (om-dua7)", () => { + beforeEach(() => { + jetstream.script = null; + jetstream.abort = false; + }); + + it("returns within the safety timeout when the stream replays history then goes quiet", async () => { + jetstream.script = async function* (self: { cursor: number | null }) { + // One historical commit (time_us in the past, so never "caught up"). + self.cursor = 1_000_000; + yield commitEvent("did:plc:author", "community.lexicon.calendar.event", 1_000_000, "evt1"); + // Quiet stream: no further events ever arrive. + await new Promise(() => {}); + }; + + const result = await withTimeout( + ingestEvents(discoverableConfig(), 999_999, 150), + 2_000, + "quiet stream", + ); + + expect(result.events).toHaveLength(1); + expect(result.events[0].uri).toBe( + "at://did:plc:author/community.lexicon.calendar.event/evt1", + ); + // The replayed cursor must come back so the caller can persist it. + expect(result.lastCursor).toBe(1_000_000); + }); + + it("returns by the safety timeout even when every arriving event is filtered out", async () => { + const knownDids = new Set(); // no known DIDs -> every event filtered + jetstream.script = async function* (self) { + let t = 1_000_000; + for (let i = 0; i < 2_000 && !jetstream.abort; i++) { + t += 1_000; + self.cursor = t; + yield commitEvent("did:plc:stranger", "app.bsky.graph.follow", t, "f" + i); + await new Promise((r) => setTimeout(r, 2)); + } + }; + + try { + const result = await withTimeout( + ingestEvents(dependentConfig(), 999_999, 100, knownDids), + 2_000, + "all-filtered flood", + ); + + expect(result.events).toHaveLength(0); // all filtered, nothing collected + // ...but the cursor still advanced, so the caller persists forward progress. + expect(result.lastCursor).not.toBeNull(); + expect(result.lastCursor).toBeGreaterThan(1_000_000); + } finally { + jetstream.abort = true; + } + }); +}); -- 2.51.2