perf: replace poll-driven queue flush and MCP child watcher with an event-driven drainable worker #489

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

Ten background loops are started unconditionally in one block (internal/server/server.go:395-410), several of them polling for work that an event could have told them about, and their tests wait on sleeps because there is no way to ask "is this worker done".

// server.go:395-410
go s.runAutoArchiveLoop(ctx)
go s.runProjectsIndexLoop(ctx)
go s.runLLMMetricsLoop(ctx)
go s.runChildSessionWatcher(ctx)
go s.runWorkflowEngine(ctx)
go s.runWorkflowTriggerEngine(ctx)
go s.runWorkflowMirror(ctx)
go s.runQueueSweep(ctx)
go s.runPromptSchedules(ctx)
go s.aaSvc().RunWatcher(ctx)

The two that are polling as a substitute for an event:

  • internal/server/queue.go:20 — queueSweepInterval = 15 * time.Second. Runs a DISTINCT query plus a per-session status check every 15 s. Its own comment acknowledges the cost.
  • internal/server/mcp_watcher.go:19 — childSessionWatchInterval = 5 * time.Second. Scans state.db for non-terminal children and calls the platform per child, every 5 s, to notice a transition the event stream already carries.

And one that polls hard for no reason: internal/platforms/opencode/models_cache.go:310-336 re-runs the full unfiltered GetSessions("", 0) every 4 s forever (sessionsRefreshInterval = 4 * time.Second), which is the heaviest query in the codebase.

Why it matters

Two costs.

Latency and load. A queued message can wait up to 15 s past the idle edge; a finished MCP child can wait up to 5 s before its parent hears. Both intervals are pure overhead when the triggering event is already available. Meanwhile the 4 s session refresher burns the expensive correlated-subquery scan whether or not anything changed.

Untestable timing. Nothing exposes "the queue is empty and the in-flight item finished", so tests either sleep or poll. That is also why the sweep exists: it is the only way to be sure a stranded row eventually drains.

Suggested fix

Add a small internal/worker package: a queue-backed serial worker with a Drain(ctx) that blocks until the queue is empty and the item currently being processed has finished.

Shape (Go, ~60 lines, no new dependency):

type Worker[T any] struct {
	mu          sync.Mutex
	cond        *sync.Cond   // broadcast on outstanding-- and on enqueue
	queue       []T
	outstanding int          // incremented atomically with the enqueue, not inferred from len(queue)
}

func (w *Worker[T]) Enqueue(v T)          // outstanding++ under the same lock as the append
func (w *Worker[T]) Drain(ctx) error      // cond.Wait until outstanding == 0

The correctness detail worth getting right: outstanding must be bumped inside the same critical section as the append, and decremented only after the handler returns, otherwise Drain can observe a false zero between take and process.

Then:

  1. Convert the queue flush to it. Enqueue on the session.idle edge; keep runQueueSweep but raise the interval and re-document it as a crash-recovery backstop, not the primary path.
  2. Convert the MCP child watcher to it. Enqueue on the child's terminal transition; keep a long-interval sweep for rows whose event was missed.
  3. Replace the sleeps in the affected tests with Drain, and add one test for the hard case: work enqueued while an item is in flight must not let Drain return early.
  4. Separately, make the 4 s session refresher on-demand or subscriber-gated (see the client-activity-lease ticket) rather than a fixed ticker.

Acceptance criteria

  • internal/worker exists with a Drain that provably waits for in-flight work; test enqueues during processing and asserts Drain has not returned.
  • Queue flush and MCP child result delivery are event-driven; their sweeps remain only as backstops with raised intervals and updated comments.
  • A queued follow-up is sent within one event round-trip of the idle edge, not up to 15 s later.
  • No time.Sleep remains in the tests for the two converted paths.
  • AGENTS.md is updated where it describes the 15 s sweep and the 5 s child poll as the mechanism.

Effort: L.

Ten background loops are started unconditionally in one block (`internal/server/server.go:395-410`), several of them polling for work that an event could have told them about, and their tests wait on sleeps because there is no way to ask "is this worker done". ```go // server.go:395-410 go s.runAutoArchiveLoop(ctx) go s.runProjectsIndexLoop(ctx) go s.runLLMMetricsLoop(ctx) go s.runChildSessionWatcher(ctx) go s.runWorkflowEngine(ctx) go s.runWorkflowTriggerEngine(ctx) go s.runWorkflowMirror(ctx) go s.runQueueSweep(ctx) go s.runPromptSchedules(ctx) go s.aaSvc().RunWatcher(ctx) ``` The two that are polling as a substitute for an event: - `internal/server/queue.go:20` — `queueSweepInterval = 15 * time.Second`. Runs a `DISTINCT` query plus a per-session status check every 15 s. Its own comment acknowledges the cost. - `internal/server/mcp_watcher.go:19` — `childSessionWatchInterval = 5 * time.Second`. Scans `state.db` for non-terminal children and calls the platform per child, every 5 s, to notice a transition the event stream already carries. And one that polls hard for no reason: `internal/platforms/opencode/models_cache.go:310-336` re-runs the full unfiltered `GetSessions("", 0)` every 4 s forever (`sessionsRefreshInterval = 4 * time.Second`), which is the heaviest query in the codebase. ## Why it matters Two costs. **Latency and load.** A queued message can wait up to 15 s past the idle edge; a finished MCP child can wait up to 5 s before its parent hears. Both intervals are pure overhead when the triggering event is already available. Meanwhile the 4 s session refresher burns the expensive correlated-subquery scan whether or not anything changed. **Untestable timing.** Nothing exposes "the queue is empty and the in-flight item finished", so tests either sleep or poll. That is also why the sweep exists: it is the only way to be sure a stranded row eventually drains. ## Suggested fix Add a small `internal/worker` package: a queue-backed serial worker with a `Drain(ctx)` that blocks until the queue is empty **and** the item currently being processed has finished. Shape (Go, ~60 lines, no new dependency): ```go type Worker[T any] struct { mu sync.Mutex cond *sync.Cond // broadcast on outstanding-- and on enqueue queue []T outstanding int // incremented atomically with the enqueue, not inferred from len(queue) } func (w *Worker[T]) Enqueue(v T) // outstanding++ under the same lock as the append func (w *Worker[T]) Drain(ctx) error // cond.Wait until outstanding == 0 ``` The correctness detail worth getting right: `outstanding` must be bumped inside the same critical section as the append, and decremented only after the handler returns, otherwise `Drain` can observe a false zero between take and process. Then: 1. Convert the queue flush to it. Enqueue on the `session.idle` edge; keep `runQueueSweep` but raise the interval and re-document it as a crash-recovery backstop, not the primary path. 2. Convert the MCP child watcher to it. Enqueue on the child's terminal transition; keep a long-interval sweep for rows whose event was missed. 3. Replace the sleeps in the affected tests with `Drain`, and add one test for the hard case: work enqueued *while* an item is in flight must not let `Drain` return early. 4. Separately, make the 4 s session refresher on-demand or subscriber-gated (see the client-activity-lease ticket) rather than a fixed ticker. ## Acceptance criteria - [ ] `internal/worker` exists with a `Drain` that provably waits for in-flight work; test enqueues during processing and asserts `Drain` has not returned. - [ ] Queue flush and MCP child result delivery are event-driven; their sweeps remain only as backstops with raised intervals and updated comments. - [ ] A queued follow-up is sent within one event round-trip of the idle edge, not up to 15 s later. - [ ] No `time.Sleep` remains in the tests for the two converted paths. - [ ] `AGENTS.md` is updated where it describes the 15 s sweep and the 5 s child poll as the mechanism. Effort: L.
dries closed this issue 2026-08-23 16:48:14 +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#489
No description provided.