From 91278c6d423a6295ed3316f242712a97df49bcde Mon Sep 17 00:00:00 2001 From: Tom Scanlan Date: Wed, 10 Jun 2026 13:41:12 -0400 Subject: [PATCH 1/2] 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 From e362ddabfa63e51da43608e30ca5e08f03d0d6a0 Mon Sep 17 00:00:00 2001 From: Tom Scanlan Date: Wed, 10 Jun 2026 14:43:33 -0400 Subject: [PATCH 2/2] test(ingest): cover caught-up break, filtered live-edge, fast-flow, identity Add regression tests locking the ingestEvents loop behavior reviewed for the quiet-stream fix: - caught-up break collects the batch + cursor when a kept event reaches the live edge (the exit path neither existing test exercised) - the caught-up check now fires on a filtered live event too (early return), deferring a following kept event to the next cycle rather than collecting it - fast-flowing events within the deadline are all collected (the next()/timeout race drops nothing) - #identity events surface as handle updates through the ingest path --- packages/contrail/tests/ingest-hang.test.ts | 102 ++++++++++++++++++++ 1 file changed, 102 insertions(+) diff --git a/packages/contrail/tests/ingest-hang.test.ts b/packages/contrail/tests/ingest-hang.test.ts index 22a9499..348bd62 100644 --- a/packages/contrail/tests/ingest-hang.test.ts +++ b/packages/contrail/tests/ingest-hang.test.ts @@ -149,4 +149,106 @@ describe("ingestEvents — bounded by the safety timeout (om-dua7)", () => { jetstream.abort = true; } }); + + it("breaks via 'caught up to present' and returns the batch + cursor when a kept event reaches the live edge", async () => { + // `live` is >= ingestEvents' startTimeUs (captured at the call below), so the + // 5s safety timeout is irrelevant — the exit must be the caught-up break. + const liveUs = Date.now() * 1000 + 5_000_000; + jetstream.script = async function* (self) { + self.cursor = 1_000_000; + yield commitEvent("did:plc:author", "community.lexicon.calendar.event", 1_000_000, "hist"); + self.cursor = liveUs; + yield commitEvent("did:plc:author", "community.lexicon.calendar.event", liveUs, "live"); + // Must NOT be reached: the caught-up break fires on `live` before this. + self.cursor = liveUs + 1_000; + yield commitEvent("did:plc:author", "community.lexicon.calendar.event", liveUs + 1_000, "after"); + await new Promise(() => {}); + }; + + const result = await withTimeout( + ingestEvents(discoverableConfig(), 999_999, 5_000), + 2_000, + "caught-up kept", + ); + + expect(result.events.map((e) => e.rkey)).toEqual(["hist", "live"]); + expect(result.lastCursor).toBe(liveUs); + }); + + it("caught-up break fires even on a filtered live event, deferring a following kept event to the next cycle", async () => { + // Behavior change vs the pre-fix loop: filtering now uses early `return`, so + // the caught-up check runs after a filtered event too. A filtered event at + // the live edge breaks the loop BEFORE the kept event that follows it. + // Pre-fix (`continue`) would have skipped the check, collected `evt`, and + // broken on it with cursor liveUs+1000. + const liveUs = Date.now() * 1000 + 5_000_000; + const knownDids = new Set(); // empty -> the follow is filtered (unknown DID) + jetstream.script = async function* (self) { + self.cursor = liveUs; + yield commitEvent("did:plc:stranger", "app.bsky.graph.follow", liveUs, "f1"); + self.cursor = liveUs + 1_000; + yield commitEvent("did:plc:author", "community.lexicon.calendar.event", liveUs + 1_000, "evt"); + await new Promise(() => {}); + }; + + const result = await withTimeout( + ingestEvents(dependentConfig(), 999_999, 5_000, knownDids), + 2_000, + "caught-up filtered", + ); + + expect(result.events).toHaveLength(0); // kept `evt` deferred, not collected this cycle + expect(result.lastCursor).toBe(liveUs); // broke at the filtered event, before `evt` + }); + + it("collects every event when they flow fast but within the safety timeout (the next()/timeout race drops nothing)", async () => { + const N = 25; + jetstream.script = async function* (self) { + let t = 1_000_000; // all historical, so caught-up never fires; deadline ends it + for (let i = 0; i < N; i++) { + t += 1_000; + self.cursor = t; + // A real await before each yield forces next() down the pending-promise + // path (not the queue fast-path), so this exercises the race directly. + await new Promise((r) => setTimeout(r, 1)); + yield commitEvent("did:plc:author", "community.lexicon.calendar.event", t, "e" + i); + } + await new Promise(() => {}); // then quiet -> safety timeout returns the batch + }; + + const result = await withTimeout( + ingestEvents(discoverableConfig(), 999_999, 300), + 2_000, + "fast flow", + ); + + expect(result.events).toHaveLength(N); + }); + + it("captures #identity events as handle updates through the ingest path", async () => { + jetstream.script = async function* (self) { + self.cursor = 1_000_000; + yield { + kind: "identity" as const, + time_us: 1_000_000, + did: "did:plc:author", + identity: { + did: "did:plc:author", + handle: "alice.test", + seq: 1, + time: "2026-04-01T10:00:00Z", + }, + }; + await new Promise(() => {}); // quiet -> safety timeout returns + }; + + const result = await withTimeout( + ingestEvents(discoverableConfig(), 999_999, 150), + 2_000, + "identity", + ); + + expect(result.identityUpdates.get("did:plc:author")).toBe("alice.test"); + expect(result.events).toHaveLength(0); // identity events are not record commits + }); }); -- 2.51.2