fix: global SSE hub silently drops non-coalescing events on a full subscriber buffer #490

Closed
opened 2026-08-04 23:54:36 +02:00 by dries · 0 comments
Owner

The global /api/events broadcast hub silently drops events when a subscriber falls behind, and only two event types are exempt.

internal/server/broadcast.go:43-46:

var coalescingEvents = map[string]bool{
	"ocman.queue.updated":  true,
	"workflow.run.updated": true,
}

Subscribers are a chan broadcastEvent buffered at 16 (broadcast.go:52-59, :92-96). broadcast() is non-blocking (:121-135): on a full buffer a coalescing event is parked last-write-wins keyed by event + "\x00" + id (:74-86, :140-152), and everything else is dropped on the floor:

// broadcast.go:121-135 (paraphrased)
select {
case sub.ch <- ev:
default:
    if coalescingEvents[ev.name] { sub.park(ev); return }
    // dropped
}

The backstop is the frontend's 10 s notify poll (frontend/src/lib/useNotifyData.ts:21).

Why it matters

ocman.session.idle and ocman.session.changed are not in the coalescing set, so under a burst — several sessions streaming at once, a workflow fan-out, a slow tab — the events that drive queue drain, notification state, and the sidebar are the ones thrown away. The failure is invisible: no log, no metric, and the 10 s poll papers over it well enough that it looks like ordinary lag.

Buffering 16 events per subscriber is also the wrong axis. The problem is not "occasionally more than 16 events", it is "one thread emits a hundred near-identical events per second and each of them costs a downstream read".

Suggested fix

Replace the drop-or-park scheme with a windowed, aggregate-keyed batcher, and make it safe by re-reading current state instead of replaying deltas.

  1. Coalesce on a window, not on overflow. Group events over a short window (start at 50 ms, cap the batch size), keep only the newest event per aggregate key (session:<id>, run:<id>, trigger:), re-sort survivors by their original order, then emit. A burst of ten ocman.session.changed for one session becomes one; an unrelated ocman.session.idle in the same window is never stuck behind them.

  2. Make coalescing safe by construction. A coalesced event must carry only an identity, so the consumer re-reads current state. That is already true of most of them ({sessionID}), with two exceptions to fix:

    • broadcast.go:264-286 embeds a full provisional db.Session on the created path (including a hardcoded Status: "waiting").
    • internal/server/queue.go:156-178 embeds the session's entire queue in ocman.queue.updated.

    Both should shrink to an id, with the client fetching. Guard it: a deletion/removal event must still survive collapsing — if the refetch finds nothing, emit a removal rather than swallowing the event.

  3. Never silently drop. Every event type becomes coalescible under (2), so the default: drop arm disappears. If a subscriber is still hopelessly behind, close it and let the client reconnect and resync — a reconnect already triggers a refetch (frontend/src/lib/useGlobalEvents.ts:165-168).

  4. Add a counter for coalesced and for closed-behind subscribers so the behaviour is observable instead of inferred.

Acceptance criteria

  • No event type is dropped on a full subscriber buffer; the coalescingEvents allowlist is gone.
  • Events are coalesced per aggregate over a bounded window, with the batch size capped.
  • ocman.session.changed and ocman.queue.updated carry identities only; the frontend refetches. Payload size is independent of queue length and session size.
  • Test: a burst of N events for one session collapses to one delivered event while an interleaved event for a different aggregate is still delivered, in order.
  • Test: a collapsed sequence ending in a removal still delivers a removal.
  • Metrics for coalesced events and for subscribers closed for falling behind.

Effort: M.

The global `/api/events` broadcast hub silently drops events when a subscriber falls behind, and only two event types are exempt. `internal/server/broadcast.go:43-46`: ```go var coalescingEvents = map[string]bool{ "ocman.queue.updated": true, "workflow.run.updated": true, } ``` Subscribers are a `chan broadcastEvent` buffered at 16 (`broadcast.go:52-59`, `:92-96`). `broadcast()` is non-blocking (`:121-135`): on a full buffer a coalescing event is parked last-write-wins keyed by `event + "\x00" + id` (`:74-86`, `:140-152`), and **everything else is dropped on the floor**: ```go // broadcast.go:121-135 (paraphrased) select { case sub.ch <- ev: default: if coalescingEvents[ev.name] { sub.park(ev); return } // dropped } ``` The backstop is the frontend's 10 s notify poll (`frontend/src/lib/useNotifyData.ts:21`). ## Why it matters `ocman.session.idle` and `ocman.session.changed` are *not* in the coalescing set, so under a burst — several sessions streaming at once, a workflow fan-out, a slow tab — the events that drive queue drain, notification state, and the sidebar are the ones thrown away. The failure is invisible: no log, no metric, and the 10 s poll papers over it well enough that it looks like ordinary lag. Buffering 16 events per subscriber is also the wrong axis. The problem is not "occasionally more than 16 events", it is "one thread emits a hundred near-identical events per second and each of them costs a downstream read". ## Suggested fix Replace the drop-or-park scheme with a windowed, aggregate-keyed batcher, and make it safe by re-reading current state instead of replaying deltas. 1. **Coalesce on a window, not on overflow.** Group events over a short window (start at 50 ms, cap the batch size), keep only the newest event per aggregate key (`session:<id>`, `run:<id>`, `trigger:`), re-sort survivors by their original order, then emit. A burst of ten `ocman.session.changed` for one session becomes one; an unrelated `ocman.session.idle` in the same window is never stuck behind them. 2. **Make coalescing safe by construction.** A coalesced event must carry only an identity, so the consumer re-reads current state. That is already true of most of them (`{sessionID}`), with two exceptions to fix: - `broadcast.go:264-286` embeds a full provisional `db.Session` on the created path (including a hardcoded `Status: "waiting"`). - `internal/server/queue.go:156-178` embeds the session's entire queue in `ocman.queue.updated`. Both should shrink to an id, with the client fetching. Guard it: a deletion/removal event must still survive collapsing — if the refetch finds nothing, emit a removal rather than swallowing the event. 3. **Never silently drop.** Every event type becomes coalescible under (2), so the `default:` drop arm disappears. If a subscriber is still hopelessly behind, close it and let the client reconnect and resync — a reconnect already triggers a refetch (`frontend/src/lib/useGlobalEvents.ts:165-168`). 4. Add a counter for coalesced and for closed-behind subscribers so the behaviour is observable instead of inferred. ## Acceptance criteria - [ ] No event type is dropped on a full subscriber buffer; the `coalescingEvents` allowlist is gone. - [ ] Events are coalesced per aggregate over a bounded window, with the batch size capped. - [ ] `ocman.session.changed` and `ocman.queue.updated` carry identities only; the frontend refetches. Payload size is independent of queue length and session size. - [ ] Test: a burst of N events for one session collapses to one delivered event while an interleaved event for a different aggregate is still delivered, in order. - [ ] Test: a collapsed sequence ending in a removal still delivers a removal. - [ ] Metrics for coalesced events and for subscribers closed for falling behind. Effort: M.
dries closed this issue 2026-08-20 23:12:38 +02:00
Sign in to join this conversation.
No milestone
No project
No assignees
1 participant
Notifications
Due date
The due date is invalid or out of range. Please use the format "yyyy-mm-dd".

No due date set.

Dependencies

No dependencies set.

Reference
dries/ocman#490
No description provided.