diff --git a/CONCEPTS.md b/CONCEPTS.md index ec540cd..1b655b3 100644 --- a/CONCEPTS.md +++ b/CONCEPTS.md @@ -106,6 +106,14 @@ The curated front door at `/automations`: an admin-picked, admin-ordered set of The platform's protective shutdown of an Automation: flipping it inactive with a recorded reason and timestamp. Two triggers exist — exceeding any rate-limit window, and Dry Run expiry. Rate-limit windows (per second, minute, and hour) count live and dry-run Deliveries alike; Skips are excluded. Re-activating the Automation clears the reason and restarts the relevant rolling window (`rateLimitResetAt`); the Owner-facing page surfaces the reason as an alert until then. +One exception: a rate-limit breach on an Ingestion catch-up event throttles instead of disabling. The Delivery is still skipped, so the budget is enforced, but the Automation stays active. + +### Ingestion catch-up + +Processing of Jetstream events that are materially older than wall clock, i.e. a replay of a gap after the worker was disconnected. Classified per-event from the event's own age (`JETSTREAM_CATCH_UP_THRESHOLD_MS`), not from connection state, so the replay-to-live handover needs no cutover signal. Its only effect is on Auto-disable: a rate-limit breach during ingestion catch-up throttles the event and leaves the Automation active, because a multi-hour replay would otherwise exhaust the hourly window within seconds of reconnecting and disable healthy Automations. Throttled events are not retried, so gap recovery is best-effort. + +Distinct from **PDS catch-up** (`catchUpStalePdsRecords`), which rewrites Automation records to a user's PDS after their OAuth session is restored. Unrelated mechanisms that share a word. + ## Notifications ### Notification grant diff --git a/lib/config.test.ts b/lib/config.test.ts index 6f785bf..67917d8 100644 --- a/lib/config.test.ts +++ b/lib/config.test.ts @@ -219,3 +219,20 @@ describe("JETSTREAM_MAX_LOOKBACK_US default", () => { expect(config.jetstreamMaxLookbackUs).toBe(12 * HOUR_US); }); }); + +describe("JETSTREAM_MAX_LOOKBACK_US validation", () => { + it("rejects a non-digit value rather than parsing to NaN", async () => { + // NaN is the dangerous direction: `nowUs - stored.timeUs <= NaN` is always + // false, so every subscription would silently discard its cursor and start + // live, losing exactly the gap this knob exists to recover. + await expect(loadConfig({ JETSTREAM_MAX_LOOKBACK_US: "36h" })).rejects.toThrow( + /non-negative integer/, + ); + }); + + it("rejects a negative value", async () => { + await expect(loadConfig({ JETSTREAM_MAX_LOOKBACK_US: "-1" })).rejects.toThrow( + /non-negative integer/, + ); + }); +}); diff --git a/lib/config.ts b/lib/config.ts index 0790a3f..6d4f5e5 100644 --- a/lib/config.ts +++ b/lib/config.ts @@ -150,6 +150,20 @@ if (!Number.isInteger(jetstreamCatchUpThresholdMs) || jetstreamCatchUpThresholdM ); } +// Validated like its siblings above. Left on a bare Number() this knob fails +// in the worst possible direction: a typo (e.g. "36h") parses to NaN, the +// connect-time check `nowUs - stored.timeUs <= NaN` is always false, and every +// subscription silently discards its cursor and starts live — losing exactly +// the gap the lookback exists to recover, with no error to notice. +const jetstreamMaxLookbackUs = parseDigitOnlyInt( + env("JETSTREAM_MAX_LOOKBACK_US", String(36 * 3600 * 1_000_000)), +); +if (!Number.isInteger(jetstreamMaxLookbackUs) || jetstreamMaxLookbackUs < 0) { + throw new Error( + "JETSTREAM_MAX_LOOKBACK_US must be a non-negative integer (microseconds); 0 disables cursor resume.", + ); +} + const port = Number(env("PORT", "5176")); export const config = { @@ -195,7 +209,7 @@ export const config = { // // Also widens the cursor-row cleanup sweep, which derives from this value // (2x, in lib/db/cleanup.ts) — a row must outlive the window that can use it. - jetstreamMaxLookbackUs: Number(env("JETSTREAM_MAX_LOOKBACK_US", String(36 * 3600 * 1_000_000))), + jetstreamMaxLookbackUs, cookieSecret, secretsKey, cloudflareOriginSecret, diff --git a/lib/jetstream/handler.test.ts b/lib/jetstream/handler.test.ts index 64e8731..b4a8955 100644 --- a/lib/jetstream/handler.test.ts +++ b/lib/jetstream/handler.test.ts @@ -768,17 +768,46 @@ describe("handleMatchedEvent", () => { it("treats an event exactly at the threshold as steady-state", async () => { // The gate is strictly-greater-than, so the boundary resolves to the // conservative side: still disables. - mockCheckRateLimit.mockResolvedValueOnce({ window: "second", count: 10, limit: 10 }); + // + // The clock is frozen because the assertion is about the exact boundary: + // with a live clock, the test's Date.now() and isCatchUpEvent's Date.now() + // are separate reads, and a millisecond tick between them pushes the event + // one microsecond past the threshold and flips the outcome. + vi.useFakeTimers(); + vi.setSystemTime(new Date("2026-08-17T12:00:00.000Z")); + try { + mockCheckRateLimit.mockResolvedValueOnce({ window: "second", count: 10, limit: 10 }); - const thresholdMs = config.jetstreamCatchUpThresholdMs; - await handleMatchedEvent( - makeMatch({ - automation: { actions: [makeWebhookAction()], fetches: [] }, - event: { time_us: (Date.now() - thresholdMs) * 1000 }, - }), - ); + const thresholdMs = config.jetstreamCatchUpThresholdMs; + await handleMatchedEvent( + makeMatch({ + automation: { actions: [makeWebhookAction()], fetches: [] }, + event: { time_us: (Date.now() - thresholdMs) * 1000 }, + }), + ); - expect(mockDisableForRateLimit).toHaveBeenCalledOnce(); + expect(mockDisableForRateLimit).toHaveBeenCalledOnce(); + } finally { + vi.useRealTimers(); + } + }); + + it("treats a malformed time_us as steady-state rather than catch-up", async () => { + // The guard's direction is the safe one: an unusable timestamp must not + // silently suppress auto-disable. An inverted comparison here would open + // the exemption for every malformed event, so pin it. + mockCheckRateLimit.mockResolvedValue({ window: "hour", count: 500, limit: 500 }); + + for (const timeUs of [0, Number.NaN]) { + mockDisableForRateLimit.mockClear(); + await handleMatchedEvent( + makeMatch({ + automation: { actions: [makeWebhookAction()], fetches: [] }, + event: { time_us: timeUs }, + }), + ); + expect(mockDisableForRateLimit).toHaveBeenCalledOnce(); + } }); it("classifies nothing as catch-up when the threshold is 0", async () => { diff --git a/lib/jetstream/handler.ts b/lib/jetstream/handler.ts index 970f459..42943f5 100644 --- a/lib/jetstream/handler.ts +++ b/lib/jetstream/handler.ts @@ -616,12 +616,16 @@ async function logRateLimitRow(match: MatchedEvent, msg: string) { async function handleCatchUpThrottle(match: MatchedEvent, breach: RateLimitBreach) { const uri = match.automation.uri; if (catchUpThrottled.has(uri)) return; - catchUpThrottled.add(uri); + // Mark only after the row lands. Marking first would mean a failed insert — + // most likely during exactly the write pressure a replay burst creates — + // silently suppresses every later attempt for the rest of the episode, + // leaving the skip with no trace at all. await logRateLimitRow( match, `Rate limit exceeded while catching up: ${breach.count} actions in last ${breach.window}. Events skipped; automation left enabled.`, ); + catchUpThrottled.add(uri); } /**