fix(modes): serialize dispatch so a stream tail cannot overtake the coalesced flush

AgentSession.#emit fires listeners fire-and-forget, and the coalesced
message_update flush fires from its own 33ms timer — neither path awaited
the other. A rapid stream tail (message_update -> message_end ->
agent_end) could therefore run the end handlers while the flush was
suspended mid-await, agent_end removing streamingComponent before
#handleMessageEnd finalizes and records the final message (issue #7443
follow-up).

- #runSerialized chains listener dispatch and the timer flush through one
  promise chain; an in-flight run holds later events until it completes.
  Idle dispatch stays synchronous (no added microtask), preserving the
  timing the coalescing tests assert on.
- Regression test: a message_end landing while the window flush is
  suspended on init is queued behind it (initCalls 1 while suspended,
  then 2), where the pristine code ran both concurrently (2 while
  suspended). Fails without the fix.
This commit is contained in:
Slava Zavadsky
2026-08-04 16:35:56 -04:00
parent 305717b86a
commit 50ee839c51
3 changed files with 91 additions and 4 deletions
+1 -1
View File
@@ -4,7 +4,7 @@
### Fixed
- Reduced streaming CPU usage by coalescing the cumulative `message_update` deltas of a turn at the event-controller dispatch boundary: at most one streaming-state rebuild runs per ~33ms window instead of one per token, cutting the per-token handler work that dominated the CPU profile of streaming sessions (especially at high token rates) while preserving per-delta speech output. ([#7443](https://github.com/can1357/oh-my-pi/issues/7443))
- Reduced streaming CPU usage by coalescing the cumulative `message_update` deltas of a turn at the event-controller dispatch boundary: at most one streaming-state rebuild runs per ~33ms window instead of one per token, cutting the per-token handler work that dominated the CPU profile of streaming sessions (especially at high token rates) while preserving per-delta speech output. Subscriber dispatch is serialized so a rapid stream tail (`message_update` → `message_end` → `agent_end`) cannot overtake the coalesced flush. ([#7443](https://github.com/can1357/oh-my-pi/issues/7443))
## [17.2.8] - 2026-08-04
### Changed
@@ -180,6 +180,10 @@ export class EventController {
// delta before the snapshot is coalesced away.
#pendingMessageUpdate: Extract<AgentSessionEvent, { type: "message_update" }> | undefined = undefined;
#messageUpdateTimer: NodeJS.Timeout | undefined = undefined;
/** Tail of the serialized dispatch chain; see #runSerialized. */
#dispatchTail: Promise<void> = Promise.resolve();
/** Whether a chained run is currently in flight (awaiting its own awaits). */
#dispatchInFlight = false;
// Deltas already fed to speech at arrival by the coalescer. `#handleMessageUpdate`
// also vocalizes so the direct `handleEvent` path (tests, session focus replay)
// keeps working — the WeakSet makes the coalesced path speak each delta exactly
@@ -446,6 +450,17 @@ export class EventController {
}
subscribeToAgent(): void {
// Serialize non-update dispatch behind any in-flight handler run:
// AgentSession.#emit fires listeners fire-and-forget (it does not await
// listener promises), so without this a rapid stream tail
// (message_update → message_end → agent_end) could let a later callback
// overtake the coalesced flush's handler mid-await — agent_end removing
// `streamingComponent` before #handleMessageEnd finalizes and records
// the final message (issue #7443 follow-up). When the tail has settled,
// dispatch stays synchronous: the flush's streaming rebuild runs before
// the listener's first await, preserving the timing the coalescing
// tests assert on. `message_update` enqueue is itself synchronous and
// needs no serialization.
this.ctx.unsubscribe = this.ctx.session.subscribe(async (event: AgentSessionEvent) => {
// Coalesce the cumulative `message_update` deltas of a streaming turn
// into at most one handler run per window. `#handleMessageUpdate` is
@@ -460,11 +475,44 @@ export class EventController {
this.#enqueueMessageUpdate(event);
return;
}
await this.#flushPendingMessageUpdate();
await this.handleEvent(event);
await this.#runSerialized(async () => {
await this.#flushPendingMessageUpdate();
await this.handleEvent(event);
});
});
}
/**
* Run `run` in the serialized dispatch chain: when another run is already
* in flight (awaiting its own awaits), wait for it first, so a rapid
* stream tail (message_update → message_end → agent_end) cannot overtake
* the coalesced flush mid-await — agent_end removing `streamingComponent`
* before #handleMessageEnd finalizes and records the final message (issue
* #7443 follow-up). The timer-based flush uses the same chain, closing the
* same race for the ~33ms window path. When the chain is idle, `run`
* starts synchronously (no intermediate microtask), preserving the
* synchronous-flush timing the coalescing tests assert on. A rejection
* propagates to the caller (the session's fire-and-forget emit) and the
* next event starts a fresh chain link instead of being dropped.
*/
async #runSerialized(run: () => Promise<void>): Promise<void> {
if (this.#dispatchInFlight) await this.#dispatchTail.catch(() => {});
const link = run();
if (!this.#dispatchInFlight) {
this.#dispatchInFlight = true;
void link.then(
() => {
this.#dispatchInFlight = false;
},
() => {
this.#dispatchInFlight = false;
},
);
}
this.#dispatchTail = link;
await link;
}
/**
* Queue a streaming `message_update` for the next coalesced handler run.
* Speech is per-delta, so the delta is vocalized at arrival before the
@@ -482,7 +530,12 @@ export class EventController {
// Mirror AgentSession.#emit: attach a catch so a streaming rebuild
// failure surfaces as a logged warning instead of a process-level
// unhandled rejection (the timer path has no listener to attach one).
void this.#flushPendingMessageUpdate().catch(err => {
// Runs inside the serialized dispatch chain so a message_end /
// agent_end landing mid-window cannot overtake this flush (issue
// #7443 follow-up).
void this.#runSerialized(async () => {
await this.#flushPendingMessageUpdate();
}).catch(err => {
logger.warn("Message update flush rejected", {
error: err instanceof Error ? err.message : String(err),
});
@@ -138,4 +138,38 @@ describe("EventController message_update coalescing", () => {
expect(pushDelta).toHaveBeenNthCalledWith(2, "one two ");
expect(pushDelta).toHaveBeenNthCalledWith(3, "one two three ");
});
it("serializes a tail event behind an in-flight window flush", async () => {
// The coalesced flush fires from a 33ms timer, NOT from the listener
// path, so AgentSession's fire-and-forget dispatch cannot serialize it:
// a message_end landing mid-flush used to run its handler concurrently,
// both calling init while the flush was suspended. The dispatch chain
// must hold the tail event until the window flush completed.
const { ctx, emit } = createStreamingFixture();
ctx.isInitialized = false;
const initGate = Promise.withResolvers<void>();
let initCalls = 0;
ctx.init = vi.fn(async () => {
initCalls += 1;
if (initCalls === 1) await initGate.promise;
});
emit(messageUpdate("tok1 tok2"));
await Bun.sleep(45); // window fires; flush suspends on init (call 1)
emit({ type: "message_end", message: assistantMessage("tok1 tok2") } as Extract<
AgentSessionEvent,
{ type: "message_end" }
>);
await Bun.sleep(0);
// The end handler must be queued behind the suspended flush, not
// running alongside it (which would double-init).
expect(initCalls).toBe(1);
initGate.resolve();
await Bun.sleep(0);
// Flush completed, then the end handler ran to completion.
expect(initCalls).toBe(2);
});
});