feat(session): serialized subscriber fan-out to prevent event reordering
- Added `#subscriberEmitGate: Promise<void>` FIFO ticket to order concurrent `#emitSessionEvent` calls and prevent event reordering. - Modified `#emitSessionEvent` to serialize subscriber fan-out: waits for previous gate before emitting, ensuring `message_start` arrives before `message_end` regardless of extension handler asymmetry. - Added `#prunedTerminalRefusal` field to retain classifier-refusal turns for post-settle readers.
This commit is contained in:
@@ -4163,30 +4163,58 @@ export class AgentSession {
|
||||
return queued;
|
||||
}
|
||||
|
||||
/**
|
||||
* Orders subscriber fan-out across concurrent `#emitSessionEvent` calls.
|
||||
* Extension emits only await when the event type has handlers, so an event
|
||||
* with no handlers could otherwise overtake an earlier event still inside
|
||||
* its extension emit — an instant refusal delivered its assistant
|
||||
* `message_end` to the TUI before its own `message_start`, skipping the
|
||||
* turn-ending error render entirely.
|
||||
*/
|
||||
#subscriberEmitGate: Promise<void> = Promise.resolve();
|
||||
|
||||
async #emitSessionEvent(event: AgentSessionEvent): Promise<void> {
|
||||
if (event.type === "message_update") {
|
||||
this.#emit(event);
|
||||
void this.#queueExtensionEvent(event);
|
||||
return;
|
||||
}
|
||||
await this.#emitExtensionEvent(event);
|
||||
// Hold the wire-level agent_end until in-flight prompts unwind. Subscribers
|
||||
// (rpc-mode, ACP, Cursor) treat agent_end as the "session is idle" signal;
|
||||
// emitting while #promptInFlightCount > 0 lets a client fire its next
|
||||
// `prompt` into a session that still reports isStreaming === true. Flush
|
||||
// happens in #endInFlight / #resetInFlight. A later agent_end (e.g. from
|
||||
// an auto-compaction turn that starts before the original prompt unwinds)
|
||||
// supersedes the pending one, which is what subscribers want — they only
|
||||
// care about the final settle.
|
||||
if (event.type === "agent_end" && this.#promptInFlightCount > 0) {
|
||||
this.#pendingAgentEndEmit = event;
|
||||
return;
|
||||
// Take a FIFO ticket before the extension emit: extension deliveries for
|
||||
// consecutive events still run concurrently, but subscriber fan-out waits
|
||||
// for every earlier event's fan-out (or deferral) to happen first.
|
||||
const previousGate = this.#subscriberEmitGate;
|
||||
const { promise: gate, resolve: releaseGate } = Promise.withResolvers<void>();
|
||||
this.#subscriberEmitGate = gate;
|
||||
try {
|
||||
await this.#emitExtensionEvent(event);
|
||||
await previousGate;
|
||||
// Hold the wire-level agent_end until in-flight prompts unwind. Subscribers
|
||||
// (rpc-mode, ACP, Cursor) treat agent_end as the "session is idle" signal;
|
||||
// emitting while #promptInFlightCount > 0 lets a client fire its next
|
||||
// `prompt` into a session that still reports isStreaming === true. Flush
|
||||
// happens in #endInFlight / #resetInFlight. A later agent_end (e.g. from
|
||||
// an auto-compaction turn that starts before the original prompt unwinds)
|
||||
// supersedes the pending one, which is what subscribers want — they only
|
||||
// care about the final settle.
|
||||
if (event.type === "agent_end" && this.#promptInFlightCount > 0) {
|
||||
this.#pendingAgentEndEmit = event;
|
||||
return;
|
||||
}
|
||||
this.#emit(event);
|
||||
} finally {
|
||||
releaseGate();
|
||||
}
|
||||
this.#emit(event);
|
||||
}
|
||||
|
||||
// Track last assistant message for auto-compaction check
|
||||
#lastAssistantMessage: AssistantMessage | undefined = undefined;
|
||||
/**
|
||||
* Classifier-refusal turn pruned from active context at settle (#3591).
|
||||
* Retained until the next run starts so post-settle readers
|
||||
* ({@link getLastAssistantMessage}: print mode, task executor) still see
|
||||
* the terminal error instead of a silently successful-looking state.
|
||||
*/
|
||||
#prunedTerminalRefusal: AssistantMessage | undefined = undefined;
|
||||
|
||||
/** Internal handler for agent events - shared by subscribe and reconnect.
|
||||
*
|
||||
|
||||
@@ -0,0 +1,127 @@
|
||||
/**
|
||||
* Subscriber event-order contract for `AgentSession`.
|
||||
*
|
||||
* Extension emits inside the session's event pipeline only await when the
|
||||
* event type has registered handlers. For a turn whose provider events all
|
||||
* land in one tick (e.g. an instant Anthropic classifier refusal: empty
|
||||
* `message_start` + terminal `message_delta` + `message_stop` in a single SSE
|
||||
* flush), an event type WITHOUT extension handlers used to overtake an earlier
|
||||
* event WITH handlers — the TUI received the assistant `message_end` before
|
||||
* its own `message_start`, so no streaming component existed yet and the
|
||||
* turn-ending error (pinned banner + inline `Error:` line) was never rendered.
|
||||
* The session must deliver events to subscribers in emission order regardless
|
||||
* of which event types extensions subscribe to.
|
||||
*/
|
||||
|
||||
import { afterAll, afterEach, beforeAll, describe, expect, it, vi } from "bun:test";
|
||||
import * as path from "node:path";
|
||||
import { Agent } from "@oh-my-pi/pi-agent-core";
|
||||
import { createMockModel } from "@oh-my-pi/pi-ai/providers/mock";
|
||||
import { getBundledModel } from "@oh-my-pi/pi-catalog/models";
|
||||
import { ModelRegistry } from "@oh-my-pi/pi-coding-agent/config/model-registry";
|
||||
import { Settings } from "@oh-my-pi/pi-coding-agent/config/settings";
|
||||
import type { ExtensionRunner } from "@oh-my-pi/pi-coding-agent/extensibility/extensions/runner";
|
||||
import { AgentSession } from "@oh-my-pi/pi-coding-agent/session/agent-session";
|
||||
import { AuthStorage } from "@oh-my-pi/pi-coding-agent/session/auth-storage";
|
||||
import { SessionManager } from "@oh-my-pi/pi-coding-agent/session/session-manager";
|
||||
import { TempDir } from "@oh-my-pi/pi-utils";
|
||||
|
||||
describe("AgentSession subscriber event order", () => {
|
||||
let tempDir: TempDir;
|
||||
let authStorage: AuthStorage;
|
||||
let modelRegistry: ModelRegistry;
|
||||
let session: AgentSession | undefined;
|
||||
|
||||
beforeAll(async () => {
|
||||
tempDir = TempDir.createSync("@pi-event-order-");
|
||||
authStorage = await AuthStorage.create(path.join(tempDir.path(), "testauth.db"));
|
||||
authStorage.setRuntimeApiKey("anthropic", "anthropic-test-key");
|
||||
modelRegistry = new ModelRegistry(authStorage, path.join(tempDir.path(), "models.yml"));
|
||||
});
|
||||
|
||||
afterAll(() => {
|
||||
authStorage.close();
|
||||
tempDir.removeSync();
|
||||
});
|
||||
|
||||
afterEach(async () => {
|
||||
if (session) {
|
||||
await session.dispose();
|
||||
session = undefined;
|
||||
}
|
||||
vi.restoreAllMocks();
|
||||
});
|
||||
|
||||
it("delivers message_start before message_end when extension handlers are asymmetric", async () => {
|
||||
const model = getBundledModel("anthropic", "claude-sonnet-4-5");
|
||||
if (!model) throw new Error("Expected bundled test model to exist");
|
||||
|
||||
// Instant terminal error turn: start and end reach the session pipeline
|
||||
// in the same tick, exactly like a real classifier refusal.
|
||||
const mock = createMockModel({
|
||||
responses: [
|
||||
{
|
||||
stopReason: "error",
|
||||
stopDetails: { type: "refusal", category: "cyber", explanation: "Declined." },
|
||||
errorMessage: "Refusal (cyber): Declined.",
|
||||
},
|
||||
],
|
||||
});
|
||||
const agent = new Agent({
|
||||
getApiKey: agentModel => `${agentModel.provider}-test-key`,
|
||||
initialState: {
|
||||
model,
|
||||
systemPrompt: ["Test"],
|
||||
tools: [],
|
||||
messages: [],
|
||||
},
|
||||
streamFn: (streamModel, context, options) => mock.stream(streamModel, context, options),
|
||||
});
|
||||
|
||||
// Handlers exist for message_start only, and each emit burns a burst of
|
||||
// microtasks — so the handler-less message_end path would skip its
|
||||
// extension await and overtake the start without subscriber-order
|
||||
// serialization.
|
||||
const extensionRunner = {
|
||||
emit: vi.fn(async () => {
|
||||
for (let hop = 0; hop < 25; hop++) await Promise.resolve();
|
||||
}),
|
||||
emitBeforeAgentStart: vi.fn().mockResolvedValue(undefined),
|
||||
hasHandlers: vi.fn((eventType: string) => eventType === "message_start"),
|
||||
emitSessionStop: vi.fn().mockResolvedValue(undefined),
|
||||
} as unknown as ExtensionRunner;
|
||||
|
||||
const settings = Settings.isolated({
|
||||
"compaction.enabled": false,
|
||||
"retry.baseDelayMs": 5,
|
||||
"retry.maxRetries": 1,
|
||||
"retry.modelFallback": false,
|
||||
});
|
||||
settings.setModelRole("default", `${model.provider}/${model.id}`);
|
||||
|
||||
session = new AgentSession({
|
||||
agent,
|
||||
sessionManager: SessionManager.inMemory(),
|
||||
settings,
|
||||
modelRegistry,
|
||||
extensionRunner,
|
||||
});
|
||||
|
||||
const order: string[] = [];
|
||||
session.subscribe(event => {
|
||||
if (event.type === "message_start" || event.type === "message_end") {
|
||||
order.push(`${event.type}:${event.message.role}`);
|
||||
}
|
||||
});
|
||||
|
||||
await session.prompt("Trigger instant refusal");
|
||||
await session.waitForIdle();
|
||||
|
||||
for (const role of ["user", "assistant"]) {
|
||||
const startIndex = order.indexOf(`message_start:${role}`);
|
||||
const endIndex = order.indexOf(`message_end:${role}`);
|
||||
expect(startIndex).toBeGreaterThanOrEqual(0);
|
||||
expect(endIndex).toBeGreaterThan(startIndex);
|
||||
}
|
||||
});
|
||||
});
|
||||
@@ -0,0 +1,103 @@
|
||||
import { describe, expect, it } from "bun:test";
|
||||
import { TranscriptContainer } from "@oh-my-pi/pi-coding-agent/modes/components/transcript-container";
|
||||
import { type Component, TUI } from "@oh-my-pi/pi-tui";
|
||||
import { StressRenderScheduler } from "../../tui/test/render-stress-scheduler";
|
||||
import { VirtualTerminal } from "../../tui/test/virtual-terminal";
|
||||
|
||||
/**
|
||||
* Finalized history block. With `tracked`, reports a post-finalize content
|
||||
* version like `AssistantMessageComponent`; otherwise it is version-untracked
|
||||
* like most tool blocks.
|
||||
*/
|
||||
class HistoryBlock implements Component {
|
||||
#lines: readonly string[];
|
||||
getTranscriptBlockVersion?: () => number;
|
||||
constructor(lines: readonly string[], tracked: boolean) {
|
||||
this.#lines = lines;
|
||||
if (tracked) this.getTranscriptBlockVersion = () => 1;
|
||||
}
|
||||
render(width: number): readonly string[] {
|
||||
return this.#lines.map(line => line.slice(0, width));
|
||||
}
|
||||
isTranscriptBlockFinalized(): boolean {
|
||||
return true;
|
||||
}
|
||||
}
|
||||
|
||||
/** Streaming live block with a settled prefix, like a streaming assistant reply. */
|
||||
class LiveBlock implements Component {
|
||||
lines: string[] = ["live-000"];
|
||||
settled = 0;
|
||||
render(width: number): readonly string[] {
|
||||
return this.lines.map(line => line.slice(0, width));
|
||||
}
|
||||
isTranscriptBlockFinalized(): boolean {
|
||||
return false;
|
||||
}
|
||||
getTranscriptBlockSettledRows(): number {
|
||||
return this.settled;
|
||||
}
|
||||
}
|
||||
|
||||
// Streams a live block behind a run of small finalized history blocks until the
|
||||
// history fully commits to native scrollback, then verifies exactly-once history
|
||||
// on the terminal tape. Regression guard: transcript-side committed-prefix
|
||||
// compaction (dropping committed rows from the local frame) shifted the frame
|
||||
// under the engine's committed-prefix ledger, the audit re-anchored, and
|
||||
// already-taped rows were recommitted below their first copy — visibly
|
||||
// duplicated blocks. The transcript now always keeps its full local frame.
|
||||
async function streamPastCommit(tracked: boolean): Promise<Map<string, number>> {
|
||||
const term = new VirtualTerminal(40, 6);
|
||||
Object.defineProperty(term, "isNativeViewportAtBottom", { configurable: true, value: () => undefined });
|
||||
const scheduler = new StressRenderScheduler();
|
||||
const tui = new TUI(term, undefined, { renderScheduler: scheduler });
|
||||
const chat = new TranscriptContainer();
|
||||
const historyRows: string[] = [];
|
||||
for (let i = 0; i < 6; i++) {
|
||||
const rows = [`box-${i}-alpha`, `box-${i}-beta`];
|
||||
historyRows.push(...rows);
|
||||
chat.addChild(new HistoryBlock(rows, tracked));
|
||||
}
|
||||
const live = new LiveBlock();
|
||||
chat.addChild(live);
|
||||
tui.addChild(chat);
|
||||
|
||||
try {
|
||||
tui.start();
|
||||
await scheduler.drain(term);
|
||||
// Grow the live block one row per frame with the settled prefix trailing
|
||||
// by one, pushing the finalized history through commit and compaction.
|
||||
for (let i = 1; i < 40; i++) {
|
||||
live.lines.push(`live-${String(i).padStart(3, "0")}`);
|
||||
live.settled = live.lines.length - 1;
|
||||
tui.requestRender();
|
||||
await scheduler.drain(term);
|
||||
}
|
||||
} finally {
|
||||
tui.stop();
|
||||
await term.flush();
|
||||
}
|
||||
|
||||
const counts = new Map<string, number>();
|
||||
for (const row of term.getScrollBuffer()) {
|
||||
const text = Bun.stripANSI(row).trimEnd();
|
||||
if (text.length === 0) continue;
|
||||
counts.set(text, (counts.get(text) ?? 0) + 1);
|
||||
}
|
||||
// Loss check alongside the duplication check: every history row must have
|
||||
// reached the tape exactly once.
|
||||
for (const row of historyRows) expect(counts.get(row) ?? 0).toBe(1);
|
||||
return counts;
|
||||
}
|
||||
|
||||
describe("transcript committed history", () => {
|
||||
it("keeps version-tracked committed history exactly once on the tape", async () => {
|
||||
const counts = await streamPastCommit(true);
|
||||
expect([...counts.entries()].filter(([, count]) => count > 1)).toEqual([]);
|
||||
});
|
||||
|
||||
it("keeps version-untracked committed history exactly once on the tape", async () => {
|
||||
const counts = await streamPastCommit(false);
|
||||
expect([...counts.entries()].filter(([, count]) => count > 1)).toEqual([]);
|
||||
});
|
||||
});
|
||||
Reference in New Issue
Block a user