fix(coding-agent): scope async delivery drains

This commit is contained in:
jiwangyihao
2026-05-17 15:09:04 +08:00
parent 05e1e702dd
commit 6e53286ce7
3 changed files with 86 additions and 245 deletions
@@ -383,248 +383,6 @@ describe("AgentSession concurrent prompt guard", () => {
}),
).toBe(true);
});
const model = getBundledModel("anthropic", "claude-sonnet-4-5")!;
const mock = createMockModel({ handler: () => ({ content: ["Done"] }) });
const agent = new Agent({
getApiKey: () => "test-key",
// Regression: a subscriber that fires the next prompt synchronously from the
// agent_end listener (the shape every wire transport ends up in — rpc-mode
// stdout subscriber, ACP bridge, Cursor exec) must not collide with the
// outgoing turn's still-unwinding in-flight bookkeeping. Before the wire-level
// agent_end was deferred until #promptInFlightCount drops to 0, the
// subscriber observed agent_end while Session.isStreaming was still true (the
// agent's own `isStreaming` had flipped, but #promptWithMessage's finally had
// not yet decremented the prompt-in-flight counter), and the next prompt
// threw AgentBusyError. Surfaced as `RpcCommandError: prompt: Agent is
// already processing` from omp-rpc clients (robomp triage reminder path).
it("subscriber may prompt() synchronously from agent_end without AgentBusyError", async () => {
const model = getBundledModel("anthropic", "claude-sonnet-4-5")!;
const mock = createMockModel({ handler: () => ({ content: ["Done"] }) });
const agent = new Agent({
getApiKey: () => "test-key",
initialState: { model, systemPrompt: ["Test"], tools: [] },
streamFn: mock.stream,
});
const sessionManager = SessionManager.inMemory();
const settings = Settings.isolated();
const authStorage = await AuthStorage.create(path.join(tempDir, "testauth.db"));
authStorages.push(authStorage);
const modelRegistry = new ModelRegistry(authStorage, path.join(tempDir, "models.yml"));
authStorage.setRuntimeApiKey("anthropic", "test-key");
session = new AgentSession({ agent, sessionManager, settings, modelRegistry });
const observedIsStreamingAtAgentEnd: boolean[] = [];
const reentrantPromptResults: Array<"resolved" | { error: string }> = [];
let reentrantPrompted = false;
session.subscribe(event => {
if (event.type !== "agent_end") return;
observedIsStreamingAtAgentEnd.push(session.isStreaming);
if (reentrantPrompted) return;
reentrantPrompted = true;
void session
.prompt("Second message")
.then(() => reentrantPromptResults.push("resolved"))
.catch((err: Error) => reentrantPromptResults.push({ error: err.message }));
});
await session.prompt("First message");
await waitFor(() => reentrantPromptResults.length > 0, 2000);
await session.waitForIdle();
expect(observedIsStreamingAtAgentEnd).not.toContain(true);
expect(reentrantPromptResults).toEqual(["resolved"]);
});
it("queues idle ACP client-triggered custom messages instead of starting an ownerless turn", async () => {
const model = getBundledModel("anthropic", "claude-sonnet-4-5")!;
const mock = createMockModel({ handler: () => ({ content: ["Done"] }) });
const agent = new Agent({
getApiKey: () => "test-key",
initialState: {
model,
systemPrompt: ["Test"],
tools: [],
},
convertToLlm,
streamFn: mock.stream,
});
const sessionManager = SessionManager.inMemory();
const settings = Settings.isolated();
const authStorage = await AuthStorage.create(":memory:");
authStorages.push(authStorage);
const modelRegistry = new ModelRegistry(authStorage);
authStorage.setRuntimeApiKey("anthropic", "test-key");
session = new AgentSession({
agent,
sessionManager,
settings,
modelRegistry,
});
session.setClientBridge({
capabilities: {},
deferAgentInitiatedTurns: true,
});
await session.prompt("First message");
expect(session.isStreaming).toBe(false);
const callsAfterFirstPrompt = mock.calls.length;
await session.sendCustomMessage(
{
customType: "async-result",
content: "Background result",
display: true,
attribution: "agent",
},
{ deliverAs: "followUp", triggerTurn: true },
);
expect(mock.calls).toHaveLength(callsAfterFirstPrompt);
expect(session.isStreaming).toBe(false);
await session.prompt("Next user prompt");
await session.dispose();
session = undefined as unknown as AgentSession;
expect(mock.calls).toHaveLength(callsAfterFirstPrompt + 1);
expect(
mock.calls.at(-1)?.context.messages.some(message => {
if (typeof message.content === "string") {
return message.content.includes("Background result");
}
return message.content.some(
content => content.type === "text" && content.text.includes("Background result"),
);
}),
).toBe(true);
});
streamFn: mock.stream,
});
const sessionManager = SessionManager.inMemory();
const settings = Settings.isolated();
// Regression: a subscriber that fires the next prompt synchronously from the
// agent_end listener (the shape every wire transport ends up in — rpc-mode
// stdout subscriber, ACP bridge, Cursor exec) must not collide with the
// outgoing turn's still-unwinding in-flight bookkeeping. Before the wire-level
// agent_end was deferred until #promptInFlightCount drops to 0, the
// subscriber observed agent_end while Session.isStreaming was still true (the
// agent's own `isStreaming` had flipped, but #promptWithMessage's finally had
// not yet decremented the prompt-in-flight counter), and the next prompt
// threw AgentBusyError. Surfaced as `RpcCommandError: prompt: Agent is
// already processing` from omp-rpc clients (robomp triage reminder path).
it("subscriber may prompt() synchronously from agent_end without AgentBusyError", async () => {
const model = getBundledModel("anthropic", "claude-sonnet-4-5")!;
const mock = createMockModel({ handler: () => ({ content: ["Done"] }) });
const agent = new Agent({
getApiKey: () => "test-key",
initialState: { model, systemPrompt: ["Test"], tools: [] },
streamFn: mock.stream,
});
const sessionManager = SessionManager.inMemory();
const settings = Settings.isolated();
const authStorage = await AuthStorage.create(path.join(tempDir, "testauth.db"));
authStorages.push(authStorage);
const modelRegistry = new ModelRegistry(authStorage, path.join(tempDir, "models.yml"));
authStorage.setRuntimeApiKey("anthropic", "test-key");
session = new AgentSession({ agent, sessionManager, settings, modelRegistry });
const observedIsStreamingAtAgentEnd: boolean[] = [];
const reentrantPromptResults: Array<"resolved" | { error: string }> = [];
let reentrantPrompted = false;
session.subscribe(event => {
if (event.type !== "agent_end") return;
observedIsStreamingAtAgentEnd.push(session.isStreaming);
if (reentrantPrompted) return;
reentrantPrompted = true;
void session
.prompt("Second message")
.then(() => reentrantPromptResults.push("resolved"))
.catch((err: Error) => reentrantPromptResults.push({ error: err.message }));
});
await session.prompt("First message");
await waitFor(() => reentrantPromptResults.length > 0, 2000);
await session.waitForIdle();
expect(observedIsStreamingAtAgentEnd).not.toContain(true);
expect(reentrantPromptResults).toEqual(["resolved"]);
});
it("queues idle ACP client-triggered custom messages instead of starting an ownerless turn", async () => {
const model = getBundledModel("anthropic", "claude-sonnet-4-5")!;
const mock = createMockModel({ handler: () => ({ content: ["Done"] }) });
const agent = new Agent({
getApiKey: () => "test-key",
initialState: {
model,
systemPrompt: ["Test"],
tools: [],
},
convertToLlm,
streamFn: mock.stream,
});
const sessionManager = SessionManager.inMemory();
const settings = Settings.isolated();
const authStorage = await AuthStorage.create(":memory:");
authStorages.push(authStorage);
const modelRegistry = new ModelRegistry(authStorage);
authStorage.setRuntimeApiKey("anthropic", "test-key");
session = new AgentSession({
agent,
sessionManager,
settings,
modelRegistry,
});
session.setClientBridge({
capabilities: {},
deferAgentInitiatedTurns: true,
});
await session.prompt("First message");
expect(session.isStreaming).toBe(false);
const callsAfterFirstPrompt = mock.calls.length;
await session.sendCustomMessage(
{
customType: "async-result",
content: "Background result",
display: true,
attribution: "agent",
},
{ deliverAs: "followUp", triggerTurn: true },
);
expect(mock.calls).toHaveLength(callsAfterFirstPrompt);
expect(session.isStreaming).toBe(false);
await session.prompt("Next user prompt");
await session.dispose();
session = undefined as unknown as AgentSession;
expect(mock.calls).toHaveLength(callsAfterFirstPrompt + 1);
expect(
mock.calls.at(-1)?.context.messages.some(message => {
if (typeof message.content === "string") {
return message.content.includes("Background result");
}
return message.content.some(
content => content.type === "text" && content.text.includes("Background result"),
);
}),
).toBe(true);
});
});
});
describe("AgentSession TTSR resume gate", () => {