feat(os-events): /api/os/events SSE stream and useOsEvents hook (supersedes #2309) - #2378
Conversation
- New authenticated SSE endpoint at GET /api/os/events streams typed change events (kind + id only, never payload) filtered by kinds query param - New shared client hook useOsEvents(kinds, onEvent) returns connected/stale flags, manages single per-client EventSource with exponential backoff - Reuses existing EventBus SSE plumbing and auth middleware - Backend tests cover auth gate, kind filtering, and multi-subscriber delivery - Frontend tests cover hook mount, kind filtering, dedup, stale flag, and reconnect
Supersedes PR #2309 (branch exec/tsk-jmctoa), whose lane is gone. That branch's endpoint and hook are merged in here on top of current dev, with the review findings resolved. - Subscriptions and relay tasks leaked. Both subscribe() calls and both create_task() calls ran in the handler body while teardown lived in the generator's finally. An async generator closed without ever being iterated never runs its body, so a client disconnecting between the handler returning and the stream starting leaked two bus subscriptions and two tasks, which the bus then fed forever. Setup now happens inside the generator, paired with the finally that undoes it. - The merged queue was unbounded, so a client that stopped reading grew the process without limit. It is capped at 256 and the relay drops the oldest event rather than blocking, then tells the client with an events.lagged frame carrying the drop count. - Dropped the SSE "id:" line. It is what makes a browser send Last-Event-ID, which this endpoint ignores by design, and seq restarted at 1 per connection so the ids were not stable anyway. - Hook: reopen the stream when kinds changes (the URL is fixed per connection, so a widened list never arrived), and report disconnected on an error while the stream is only CONNECTING. - The EventSource test mock had no static readyState constants, so the disconnect assertions compared undefined to undefined and passed whatever the hook did. - Documented the endpoint in docs/agent-coordination.md and moved the direct CHANGELOG.md edit into a changelog.d fragment.
|
ⓘ Qodo reviews are paused because the subscription is no longer active. Ask your workspace admin to reactivate the subscription to resume reviews. Manage billing |
|
Warning Review limit reached
Next review available in: 15 minutes You've used all free OSS reviews for now. Wait for the free limit to reset to keep reviewing this public repository. 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: Path: .coderabbit.yaml Review profile: CHILL Plan: Pro Plus Run ID: 📒 Files selected for processing (7)
📝 WalkthroughWalkthroughAdds an authenticated ChangesOS Events SSE
Estimated code review effort: 4 (Complex) | ~45 minutes Mergeability Score: 🟠 High · up to The SSE stream can still lose requested events under filtered load, hide its own lag signal from clients, and report a stalled connection as healthy indefinitely. These behaviors affect event correctness and connection availability, so the PR is not ready to merge until they are addressed. Sequence Diagram(s)sequenceDiagram
participant Desktop
participant useOsEvents
participant OSEvents
participant EventBus
Desktop->>useOsEvents: provide kinds and callback
useOsEvents->>OSEvents: GET /api/os/events?kinds=...
OSEvents->>EventBus: subscribe to user and broadcast channels
EventBus-->>OSEvents: matching event
OSEvents-->>useOsEvents: metadata-only SSE frame
useOsEvents->>Desktop: deliver validated event
Possibly related PRs
🚥 Pre-merge checks | ✅ 5✅ Passed checks (5 passed)
✨ Finishing Touches📝 Generate docstrings
🧪 Generate unit tests (beta)
Thanks for using CodeRabbit! It's free for OSS, and your support helps us grow. If you like it, consider giving us a shout-out. Comment |
|
nemotron-super review VERDICT: No blocking issues found. Automated first-pass review by the nemotron-super lane. The lead still reviews before merge. |
There was a problem hiding this comment.
Actionable comments posted: 3
🤖 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 `@desktop/src/hooks/use-os-events.ts`:
- Around line 68-80: Update the event handling around OsEvent to expose a
discriminated union for regular events and the events.lagged control event,
validating the lag payload and always forwarding valid lag events regardless of
kinds filtering. Restrict alreadySeen deduplication to regular events with
string IDs, and keep regular events subject to the existing allowlist and
validation; add coverage for lag delivery with both filtered and unfiltered
subscriptions.
- Around line 61-101: Add a named or data heartbeat to the SSE stream and handle
it in the EventSource setup alongside onmessage. Reset a client-side liveness
deadline whenever the heartbeat arrives; on expiry, close the current source,
mark the connection stale, and schedule exactly one reconnect using the existing
reconnect state and cleanup paths in connect and onerror.
In `@tinyagentos/routes/os_events.py`:
- Around line 49-66: Update _relay to accept allowed_kinds and discard events
whose kind is not requested before calling dst.put_nowait(), incrementing no lag
counter for skipped events. Pass the requested filter from the caller and remove
the later post-queue kind filtering so only allowed events consume merged queue
capacity.
🪄 Autofix
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: Path: .coderabbit.yaml
Review profile: CHILL
Plan: Pro Plus
Run ID: 300ed1f7-f683-4a1f-ba47-924fb9f375d1
📒 Files selected for processing (7)
changelog.d/tsk-ronapx-os-events-sse.mddesktop/src/hooks/use-os-events.test.tsdesktop/src/hooks/use-os-events.tsdocs/agent-coordination.mdtests/test_os_events.pytinyagentos/routes/__init__.pytinyagentos/routes/os_events.py
| es.onopen = () => { | ||
| reconnectAttemptsRef.current = 0; | ||
| setConnected(true); | ||
| setStale(false); | ||
| }; | ||
|
|
||
| es.onmessage = (msg) => { | ||
| let event: OsEvent | null; | ||
| try { | ||
| event = JSON.parse(msg.data as string) as OsEvent; | ||
| } catch { | ||
| return; | ||
| } | ||
| if (!event || typeof event !== "object") return; | ||
| if (!event.kind) return; | ||
| if (kindsRef.current.length > 0 && !kindsRef.current.includes(event.kind)) { | ||
| return; | ||
| } | ||
| if (alreadySeen(event.id)) return; | ||
| onEventRef.current(event); | ||
| }; | ||
|
|
||
| es.onerror = () => { | ||
| if (stoppedRef.current) return; | ||
| // Whatever the readyState, an error means we are not receiving events. | ||
| // Only a CLOSED stream is ours to reconnect: while CONNECTING the | ||
| // browser is already retrying, and scheduling our own reconnect on top | ||
| // of that would open a second stream. | ||
| setConnected(false); | ||
| setStale(true); | ||
| if (es.readyState === EventSource.CLOSED) { | ||
| const delay = Math.min( | ||
| RECONNECT_DELAY_MS * 2 ** reconnectAttemptsRef.current, | ||
| MAX_RECONNECT_DELAY_MS, | ||
| ); | ||
| reconnectAttemptsRef.current += 1; | ||
| reconnectTimerRef.current = setTimeout(() => { | ||
| if (!stoppedRef.current) connect(); | ||
| }, delay); | ||
| } | ||
| }; |
There was a problem hiding this comment.
🩺 Stability & Availability | 🟠 Major | 🏗️ Heavy lift
🧩 Analysis chain
🌐 Web query:
For the browser EventSource API, are SSE comment frames delivered to onmessage, and how does EventSource report a silently stalled or half-open connection?
💡 Result:
SSE comment frames, which are lines starting with a colon (e.g., : heartbeat), are not delivered to the onmessage handler or any other event listener in the browser EventSource API [1][2][3]. According to the HTML specification, these lines are ignored by the parser and serve purely as a mechanism to keep the connection alive by preventing intermediate proxies or the browser from timing out an idle stream [1][3][4]. Regarding silently stalled half-open connections, the native EventSource API does not natively detect or report them [5]. If a TCP connection is severed without a proper FIN or RST packet—common in mobile network transitions or with silent proxy timeouts—the EventSource readyState will remain 1 (OPEN) indefinitely, and no error event will be dispatched [5][6]. To handle this, developers must implement a client-side "watchdog" timer pattern [5][6]. This involves the following: 1. Heartbeat: The server must send periodic data or named events (e.g., event: ping) [3][5]. 2. Watchdog Timer: The client resets a timer whenever a message or heartbeat is received [5]. 3. Recovery: If the timer exceeds a predefined threshold (e.g., 30–60 seconds), the client assumes the connection is half-open, calls EventSource.close to clean up the socket, and re-instantiates the EventSource to trigger a reconnection [5]. Relying on the browser's built-in reconnection logic is insufficient for these "zombie" connections because the browser is unaware that the underlying TCP session is no longer functional [5][6].
Citations:
- 1: https://html.spec.whatwg.org/dev/server-sent-events.html
- 2: https://whatpr.org/html/3851/server-sent-events.html
- 3: https://developer.mozilla.org/en-US/docs/Web/API/Server-sent_events/Using_server-sent_events
- 4: https://html.spec.whatwg.org/multipage/server-sent-events.html
- 5: https://www.server-sent-events.com/backend-stream-generation-connection-management/http-keep-alive-connection-lifecycle/
- 6: https://www.server-sent-events.com/frontend-consumption-client-patterns/error-handling-and-reconnection-ux/
🏁 Script executed:
#!/bin/bash
set -eu
printf '%s\n' '--- hook ---'
sed -n '1,180p' desktop/src/hooks/use-os-events.ts
printf '%s\n' '--- related event definitions and producers ---'
rg -n -C 3 'events\.lagged|OsEvent|EventSource|heartbeat|keep-alive|event: ping|: heartbeat' desktop server src .github 2>/dev/null | head -300Repository: jaylfc/taOS
Length of output: 24448
🏁 Script executed:
#!/bin/bash
set -eu
printf '%s\n' '--- event route and SSE writers ---'
rg -n -C 4 'api/os/events|os/events|text/event-stream|EventSourceResponse|StreamingResponse|events\.lagged|last.?heartbeat|heartbeat' --glob '!desktop/src/hooks/use-os-events.ts' --glob '!desktop/src/hooks/use-os-events.test.ts' .
printf '%s\n' '--- hook tests around lifecycle behavior ---'
sed -n '100,240p' desktop/src/hooks/use-os-events.test.tsRepository: jaylfc/taOS
Length of output: 50368
🏁 Script executed:
#!/bin/bash
set -eu
printf '%s\n' '--- exact os-events route references ---'
rg -n -C 12 'api/os/events|os/events' --glob '*.py' --glob '*.ts' --glob '*.tsx' --glob '*.js' --glob '*.jsx' .
printf '%s\n' '--- SSE response implementations near exact matches ---'
rg -l 'api/os/events|os/events' --glob '*.py' --glob '*.ts' --glob '*.tsx' --glob '*.js' --glob '*.jsx' . | while IFS= read -r f; do
echo "### $f"
rg -n -C 20 'api/os/events|os/events|text/event-stream|StreamingResponse|yield|events\.lagged' "$f"
doneRepository: jaylfc/taOS
Length of output: 41255
Implement observable SSE liveness detection.
The route emits :keepalive comments every 10 seconds, but EventSource does not dispatch comments to onmessage. A silently stalled connection can remain connected: true and stale: false indefinitely.
Send a named or data heartbeat. Reset a client-side deadline when it arrives. When the deadline expires, close the source, set stale, and schedule one 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 `@desktop/src/hooks/use-os-events.ts` around lines 61 - 101, Add a named or
data heartbeat to the SSE stream and handle it in the EventSource setup
alongside onmessage. Reset a client-side liveness deadline whenever the
heartbeat arrives; on expiry, close the current source, mark the connection
stale, and schedule exactly one reconnect using the existing reconnect state and
cleanup paths in connect and onerror.
| return JSONResponse({"detail": "Service starting"}, status_code=503) | ||
|
|
||
| kinds_param = request.query_params.get("kinds", "") | ||
| allowed_kinds = ( |
There was a problem hiding this comment.
WARNING: Whitespace-only kinds parameter produces an empty allowlist that filters out every event.
The if kinds_param guard treats a whitespace-only string as truthy, so allowed_kinds becomes set() instead of None. When the client sends ?kinds= , the later event.kind not in allowed_kinds check is always True and no events are delivered. This is inconsistent with the documented behaviour ("Omitted or empty means every kind").
| allowed_kinds = ( | |
| kinds_param = request.query_params.get("kinds", "") | |
| stripped = [k.strip() for k in kinds_param.split(",") if k.strip()] | |
| allowed_kinds = set(stripped) if stripped else None |
Reply with @kilocode-bot fix it to have Kilo Code address this issue.
Code Review SummaryStatus: No Issues Found | Recommendation: Merge Files Reviewed (6 files)
Previous Review Summary (commit 27bfc74)Current summary above is authoritative. Previous snapshots are kept for context only. Previous review (commit 27bfc74)Status: 1 Issue Found | Recommendation: Address before merge Overview
Issue Details (click to expand)WARNING
Files Reviewed (7 files)
Reviewed by step-3.7-flash · Input: 57.1K · Output: 6.6K · Cached: 220.7K |
- kinds filtering moved into the relay. It ran at the yield, after an event had already taken one of the 256 slots, so unrelated traffic could evict the events a subscriber actually asked for and report a lag that, from that subscriber's view, never happened. - a kinds parameter that names no kind now means "every kind". The guard tested the raw string, and " " is truthy, so the allowlist became an empty set that matched nothing and the stream delivered silence. The set is derived first and an empty one means no filter. - events.lagged is a control frame, not a change notification. The hook dropped it whenever the caller passed a non-empty kinds list, so the one subscriber who needed to know it had missed events never heard; and its id is null, so consecutive lag frames collapsed into already-seen. It now bypasses both the kind filter and the dedupe, and OsEvent.id is typed string | null to match what the server sends. Each proven red first: the blank-kinds cases fail against the previous guard, and the lag-delivery test fails without the control-frame branch.
|
Bot round adjudicated. Both bots reviewed head Fixed.
Deferred, not declined — carded as
Verification after the fixes: 12 backend tests, 15 vitest, |
|
@coderabbitai full review Head |
|
|
|
Fresh review of head Recorded rather than glossed: kilo's check reads SUCCESS at this head with an empty description, which is the run completing, not a review artifact. The substantive re-review here is CodeRabbit's. Merging on that plus my own line review of the three folded fixes, each proven red first. |
Card:
tsk-ronapx. Supersedes #2309 (branchexec/tsk-jmctoa), which had been red and unattended since 2026-08-06 with no lane left to answer the review. STEP 0 of the card was done as instructed:git merge origin/exec/tsk-jmctoaonto current dev, which merged clean. Every item below is on top of that.Close #2309 as superseded when this merges.
1 (BLOCKING) — the subscription and task leak
Confirmed exactly as carded. Both
event_bus.subscribe()calls and bothasyncio.create_task(_relay(...))calls ran in the handler body, while teardown lived ingen()'sfinally.StreamingResponse.body_iteratoris an async generator, and an async generator closed without ever being iterated never runs its body — so thatfinallynever runs. A client disconnecting between the handler returning and the stream starting leaked two subscriptions and two never-cancelled tasks per occurrence.EventBus.subscribehands out an unboundedasyncio.Queueand_publish_to_channeldoesput_nowaitinto every registered queue, so the bus then copies every later event into queues nobody will ever drain.Fix: setup moved inside
gen(), immediately inside thetry, so thefinallycan only ever undo setup that actually happened (the handles start asNone/[]and the teardown is guarded, so a failure part-way through setup is still clean).Red first, against the code as merged — the test closes the response without consuming it and asserts the bus has no subscribers left:
2 — bounded merged queue
mergedis capped at_MERGED_MAXSIZE = 256. Chosen behaviour: drop the oldest, count the drop, and tell the client. When the queue is full the relay discards the oldest buffered event to make room and increments a counter; the generator emitsdata: {"kind": "events.lagged", "id": null, "ts": ..., "dropped": N}before the next real frame, which is the client's cue to refetch rather than assume it saw everything. The relay never blocks — a blocked relay stalls delivery for the whole connection while the bus keeps filling the two channel queues, which is the "events silently stop" outcome the card rejects._relayis now a module-level function taking(src, dst, lag)rather than a closure, so the test drives the real implementation instead of a copy of it. Red proof for the policy, produced by reverting only the relay body to the blockingawait dst.put(ev):3 — dropped the
id:lineRemoved, along with
seq. Emitting an SSEid:is what makes a browser sendLast-Event-IDon reconnect, and the module's own docstring says resume is best-effort via the bus replay buffer and notLast-Event-ID;seqalso restarted at 1 per connection, so the ids were never stable. The JSONidfield (the event's trace id) is untouched —test_stream_delivers_subscribed_kindsstill asserts on it.4 — doc-gate
docs/agent-coordination.mdgains a section for the endpoint: auth posture (session cookie, not inEXEMPT_PATHS, no registry scope reaches it), thekindsfilter, the frame shape, the no-payload rule, the keepalive, why there is noid:line, the lag behaviour, and why setup lives inside the generator. Not aDocs-Reviewedtrailer.5 — changelog
The direct
CHANGELOG.mdedit is reverted to dev's copy (git diff origin/dev -- CHANGELOG.mdis empty) and replaced withchangelog.d/tsk-ronapx-os-events-sse.md.6 —
use-os-events.ts(a) reconnect when
kindschanges. The connect effect depended only on theuseCallbackidentity, which never changed, and the URL is built once per connection — so wideningkindskept streaming the old server-side filter and the new kinds never arrived. The effect is now keyed on the serialized kinds. A companion test pins that a new array with the same contents does not churn the connection.(b)
onerrorwhile CONNECTING.connected=false/stale=trueare now set for any error, since an error means events are not arriving whatever the readyState; only a CLOSED stream schedules our own reconnect, because while CONNECTING the browser is already retrying and a second timer would open a second stream.(c) the mock had no readyState constants.
EventSource.CLOSEDwasundefinedon the mock, so the existing disconnect test setreadyState = undefinedand the hook comparedundefined === undefined— it passed no matter what the hook did. The mock now carriesCONNECTING/OPEN/CLOSED.Red proof the card asks for, with the hook's
close()removed from the cleanup:And with
onerrorreverted to the CLOSED-only guard, which is what makes (b) more than a comment:The broadcast-channel cross-user change is NOT attempted here; it stays with
tsk-3q32qc.Verification
Summary by CodeRabbit
New Features
Documentation
Tests