Repository navigation
[PROJ-1851] Preserve Jetstream progress under firehose load - #358
AndrewNordstrom wants to merge 3 commits into
Conversation
|
Warning Review limit reachedYou’ve reached a temporary PR review limit under our Fair Usage Limits Policy. Next review available in: 18 minutes Your organization has reached its usage spending cap. Adjust your spending cap in the billing tab. How can I continue?After more reviews become available, a review can be triggered using the To avoid repeated limits, reduce automatic review volume by pausing incremental auto-reviews earlier, using label-based review opt-in, excluding WIP or generated PR titles, or requesting reviews manually when the PR is ready. If your team needs uninterrupted high-volume reviews, an organization admin can enable usage-based reviews. How do review limits work?CodeRabbit enforces per-developer PR review limits for each organization. Most developers receive the normal plan review availability. For paid Pro and Pro+ PR reviews, CodeRabbit uses adaptive limits for sustained high-volume activity. When a developer's recent PR review activity reaches the 95th percentile or higher among CodeRabbit users, additional reviews become available more gradually as earlier reviews age out of the rolling window. Please refer docs for additional details. Review details⚙️ Run configurationConfiguration used: Repository YAML (base), Organization UI (inherited) Review profile: ASSERTIVE Plan: Pro Plus Run ID: 📒 Files selected for processing (16)
WalkthroughJetstream ingestion now applies bounded WebSocket backpressure, preserves safe cursor state across reconnects, reports runtime telemetry, and evaluates freshness using cursor lag and recent event activity. Admin health, feed-health, vitals, tests, and the development journal reflect these changes. ChangesJetstream ingestion reliability
Estimated code review effort: 4 (Complex) | ~45 minutes Sequence Diagram(s)sequenceDiagram
participant JetstreamWebSocket
participant JetstreamQueue
participant processEvent
participant HealthSurface
JetstreamWebSocket->>JetstreamQueue: deliver events
JetstreamQueue->>processEvent: process bounded work
processEvent-->>JetstreamQueue: release completed slots
JetstreamQueue->>JetstreamWebSocket: pause or resume delivery
HealthSurface->>JetstreamQueue: read runtime state
HealthSurface-->>HealthSurface: evaluate cursor and event freshness
Possibly related issues
Suggested labels: 🚥 Pre-merge checks | ✅ 4 | ❌ 1❌ Failed checks (1 warning)
✅ Passed checks (4 passed)
✨ Finishing Touches🧪 Generate unit tests (beta)
✨ Simplify code
Warning Review ran into problems🔥 ProblemsThese MCP integrations need to be re-authenticated in the Integrations settings: Notion Comment |
There was a problem hiding this comment.
Actionable comments posted: 9
Caution
Some comments are outside the diff and can’t be posted inline due to platform limitations.
⚠️ Outside diff range comments (2)
src/ingestion/jetstream.ts (2)
1263-1283: 🩺 Stability & Availability | 🟡 Minor | ⚡ Quick winReset all reconnect lifecycle state in the test helper.
reset()clears the new timer and counters but leavesisShuttingDown,useFallback, reconnect/failure counters, active sockets, and metrics intervals unchanged. Reconnect/fallback tests can therefore become order-dependent or leak work into later tests. Reset or explicitly close/clear every lifecycle resource, then test reset after a fallback and pending reconnect.🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the rest with a brief reason, keep changes minimal, and validate. In `@src/ingestion/jetstream.ts` around lines 1263 - 1283, Update the test helper’s reset() method to restore all reconnect lifecycle state, including isShuttingDown, useFallback, reconnect/failure counters, active sockets, and metrics intervals, in addition to the existing fields. Explicitly close or clear any remaining lifecycle resources and reset related state so fallback and pending-reconnect tests are isolated and do not leak work into later tests.
964-975: 🗄️ Data Integrity & Integration | 🟠 Major | 🏗️ Heavy liftReject stale Jetstream messages before
processEvent. Insrc/ingestion/jetstream.ts:964, a bufferedmessagefrom the previous socket can still reachprocessJetstreamMessageDataafter reconnect, so stale commit events still hitprocessEventand its side effects; the generation checks only protect cursor/pin bookkeeping. Add a reconnect regression test that delivers an old-socket buffered message after the new connection is active.🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the rest with a brief reason, keep changes minimal, and validate. In `@src/ingestion/jetstream.ts` around lines 964 - 975, Update the JetStream message handling around the socket.on('message') callback and processJetstreamMessageData to validate sessionGeneration and socket ownership before invoking processEvent, dropping buffered messages from prior connections. Preserve processing for messages from the active connection and keep existing cursor/pin bookkeeping checks. Add a reconnect regression test that delivers a buffered old-socket message after the new connection is active and verifies it produces no processEvent side effects.Source: Path instructions
🤖 Prompt for all review comments with AI agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.
Inline comments:
In `@src/ingestion/jetstream.ts`:
- Around line 245-257: Update stopJetstream() to mark shutdown, close/detach
inbound delivery, and wait until activeEventCount reaches zero before persisting
lastCursorUs; retain draining of queued waiters as appropriate. Guard
resumeInboundIfReady() (and any releaseSlot-triggered path) with the shutdown
state so teardown cannot reopen ingestion or accept new events. Add a regression
test using a slow handler that finishes after stopJetstream() begins and verify
the persisted cursor includes its completed work.
- Around line 37-41: Increase the reserved headroom in the threshold calculation
near PAUSE_QUEUE_THRESHOLD so pausing occurs early enough to account for frames
already buffered after ws.pause(). Ensure eventQueue remains below
MAX_PENDING_EVENTS and prevents handleQueueOverload() and dropped work during
sustained input; add a soak test for continuous traffic while paused that
asserts no dropped events or overload reconnects.
In `@src/lib/health.ts`:
- Around line 81-85: Bound the startup freshness logic in the health calculation
around lastEventAgeMs, eventIsFresh, and cursorIsFresh by tracking ingestion
start or connection time. When lastEventAt or cursorLagMs is absent, use that
start time as the fallback age so a connected instance with no event or cursor
becomes unhealthy at the five-minute JETSTREAM_FRESHNESS_LIMIT_MS boundary;
retain observed event and cursor ages once available, and add tests for the
exact boundary.
In `@tests/feed-health-rescore.test.ts`:
- Around line 253-269: Update the Fastify test around registerFeedHealthRoutes
and app.inject so the request and all assertions execute inside a try block,
with await app.close() in a finally block. Ensure the Fastify instance is closed
whether the injection or assertions succeed or fail.
In `@tests/health-redaction.test.ts`:
- Around line 80-113: Extend the calculateJetstreamHealth tests to cover the
just-below five-minute boundary at 299,999 ms for both event age and cursor lag.
Add assertions using new Date(nowMs - 299_999) and cursorLagMs: 299_999,
verifying status is healthy while preserving the existing exact-boundary
unhealthy cases.
- Around line 239-251: Extend the expected telemetry assertions in
tests/health-redaction.test.ts lines 239-251 to include cursor_us,
active_events, pause_count, resume_count, and overload_reconnect_count alongside
the existing JetStream fields. Also update tests/feed-health-rescore.test.ts
lines 192-194 to assert cursorUs, activeEvents, pauseCount, resumeCount,
overloadReconnectCount, and totalDroppedEvents, using the fixture values and
preserving the existing response contract checks.
In `@tests/jetstream-backpressure.test.ts`:
- Around line 65-153: Add isolated tests in the Jetstream backpressure suite for
exceptions from the flow-control socket’s pause and resume methods. Use throwing
spies, invoke the relevant backpressure transition, and assert overload
recovery/reconnect accounting for pause failures plus detachment, runtime pause
state, and close(1011, 'backpressure_resume_failed') behavior for resume
failures. Ensure each test resets and cleans up queue state and verifies
counters against the implementation’s actual behavior.
In `@tests/jetstream-lifecycle.test.ts`:
- Around line 33-40: Update the reconnect test to invoke MockWebSocket.emitClose
with code 1006 for the abnormal peer/network closure path instead of close. In
MockWebSocket.close, validate the supplied close code and reject invalid codes
such as 1006 before calling wsCloseMock or emitClose, matching real ws behavior;
preserve valid close handling.
In `@tests/queue-saturation-metrics.test.ts`:
- Around line 72-88: Add boundary-focused cases alongside the existing
cursor-lag test for __testJetstreamQueue: verify no cursor returns null cursorUs
and cursorLagMs, a cursor equal to the fixed current time reports zero lag, and
a future cursor follows the intended non-negative/error behavior without
exposing negative lag. Use fake timers with fixed system time and reset the
queue and timers for each case.
---
Outside diff comments:
In `@src/ingestion/jetstream.ts`:
- Around line 1263-1283: Update the test helper’s reset() method to restore all
reconnect lifecycle state, including isShuttingDown, useFallback,
reconnect/failure counters, active sockets, and metrics intervals, in addition
to the existing fields. Explicitly close or clear any remaining lifecycle
resources and reset related state so fallback and pending-reconnect tests are
isolated and do not leak work into later tests.
- Around line 964-975: Update the JetStream message handling around the
socket.on('message') callback and processJetstreamMessageData to validate
sessionGeneration and socket ownership before invoking processEvent, dropping
buffered messages from prior connections. Preserve processing for messages from
the active connection and keep existing cursor/pin bookkeeping checks. Add a
reconnect regression test that delivers a buffered old-socket message after the
new connection is active and verifies it produces no processEvent side effects.
🪄 Autofix (Beta)
Fix all unresolved CodeRabbit comments on this PR:
- Push a commit to this branch (recommended)
- Create a new PR with the fixes
ℹ️ Review info
⚙️ Run configuration
Configuration used: Repository YAML (base), Organization UI (inherited)
Review profile: ASSERTIVE
Plan: Pro Plus
Run ID: 91dfbf08-5a16-4fda-a4a2-8a01a4b78603
📒 Files selected for processing (13)
docs/dev-journal.mdsrc/admin/routes/feed-health.tssrc/admin/routes/health.tssrc/admin/routes/vitals.tssrc/index.tssrc/ingestion/jetstream-health.tssrc/ingestion/jetstream.tssrc/lib/health.tstests/feed-health-rescore.test.tstests/health-redaction.test.tstests/jetstream-backpressure.test.tstests/jetstream-lifecycle.test.tstests/queue-saturation-metrics.test.ts
|
Queue control: returning this PR to draft while #360 remains the sole PI-review merge lane. Current head |
|
Exact-head preparation receipt for f248bc6:\n\n- local exact worktree: 153 files / 1,716 tests\n- npm run build, npm run verify, npm run docs:verify, web-next lint/typecheck/build: pass\n- GitHub backend/frontend/docs/report/quality/security/CodeQL gates: pass\n- production-shaped 5,000-event socket burst: 5,053.69 events/sec, 26 pause/resume cycles, zero drops/reconnects/cursor mismatch\n\nThe prior CodeRabbit findings are covered by the current code and regression tests. Automatic review skipped because this PR intentionally remains draft while #360 is the sole PI-review merge lane. After #360 deploys and smokes cleanly, this branch still needs one current-main integration, exact-head hosted review, and production recovery proof before merge. |
f248bc6 to
fc44c95
Compare
|
This PR has been inactive for 14 days. It will close in 7 days if there is no update. |
Summary
Incident evidence
Jetstream ingestion queue saturatedevents after 2026-07-12 07:57 EDTVerification
npm run verify: 153 files, 1,700 tests passed; root/CLI/SDK/legacy web/web-next builds passednpm run docs:verify: passedcd web-next && npm run lint: passedcd web-next && npx tsc --noEmit: passedgit diff --check: passedRollout boundary
Linear: PROJ-1851