improve(agents): avoid polling delays in async stream tests (#161434)
Some checks are pending
ClawSweeper Dispatch / dispatch (push) Waiting to run
CodeQL / Security High (actions) (push) Waiting to run
CodeQL / Security High (channel-runtime-boundary) (push) Waiting to run
CodeQL / Security High (core-auth-secrets) (push) Waiting to run
CodeQL / Security High (mcp-process-tool-boundary) (push) Waiting to run
CodeQL / Security High (network-ssrf-boundary) (push) Waiting to run
CodeQL / Security High (plugin-trust-boundary) (push) Waiting to run
CodeQL / Security High (process-exec-boundary) (push) Waiting to run
Control UI Locale Refresh / resolve-base (push) Waiting to run
Control UI Locale Refresh / Verify generated PR App permissions (push) Blocked by required conditions
Control UI Locale Refresh / Refresh (push) Blocked by required conditions
Control UI Locale Refresh / Commit control UI locale refresh (push) Blocked by required conditions
Docs Sync Publish Repo / sync-publish-repo (push) Waiting to run
Docs / docs (push) Waiting to run
Native App Locale Refresh / Refresh native vi (push) Blocked by required conditions
Native App Locale Refresh / Refresh native zh-CN (push) Blocked by required conditions
Native App Locale Refresh / Refresh native zh-TW (push) Blocked by required conditions
Native App Locale Refresh / resolve-base (push) Waiting to run
Native App Locale Refresh / Commit native locale refresh (push) Blocked by required conditions
Native App Locale Refresh / Verify generated PR App permissions (push) Blocked by required conditions
Native App Locale Refresh / Refresh native ar (push) Blocked by required conditions
Native App Locale Refresh / Refresh native de (push) Blocked by required conditions
Native App Locale Refresh / Refresh native es (push) Blocked by required conditions
Native App Locale Refresh / Refresh native fa (push) Blocked by required conditions
Native App Locale Refresh / Refresh native fr (push) Blocked by required conditions
Native App Locale Refresh / Refresh native hi (push) Blocked by required conditions
Native App Locale Refresh / Refresh native id (push) Blocked by required conditions
Native App Locale Refresh / Refresh native it (push) Blocked by required conditions
Native App Locale Refresh / Refresh native ja-JP (push) Blocked by required conditions
Native App Locale Refresh / Refresh native ko (push) Blocked by required conditions
Native App Locale Refresh / Refresh native nl (push) Blocked by required conditions
Native App Locale Refresh / Refresh native pl (push) Blocked by required conditions
Native App Locale Refresh / Refresh native pt-BR (push) Blocked by required conditions
Native App Locale Refresh / Refresh native ru (push) Blocked by required conditions
Native App Locale Refresh / Refresh native sv (push) Blocked by required conditions
Native App Locale Refresh / Refresh native th (push) Blocked by required conditions
Native App Locale Refresh / Refresh native tr (push) Blocked by required conditions
Native App Locale Refresh / Refresh native uk (push) Blocked by required conditions
Node Runtime Conformance / scope (push) Waiting to run
Node Runtime Conformance / TypeScript contracts (push) Blocked by required conditions
Node Runtime Conformance / Rust workspace (push) Blocked by required conditions
Plugin Init Scaffold Validation / Validate provider scaffold (push) Blocked by required conditions
Plugin Init Scaffold Validation / scope (push) Waiting to run
Plugin NPM Release / preview_plugins_npm (push) Waiting to run
Plugin NPM Release / Validate release publish approval (push) Blocked by required conditions
Plugin NPM Release / preview_plugin_pack (push) Blocked by required conditions
Plugin NPM Release / Preflight plugin npm package () (push) Blocked by required conditions
Plugin NPM Release / Seal prepared plugin npm release (push) Blocked by required conditions
Plugin NPM Release / Trusted publisher OIDC exchange (push) Blocked by required conditions
Plugin NPM Release / approve_plugins_npm_release (push) Blocked by required conditions
Plugin NPM Release / Publish plugin npm package () (push) Blocked by required conditions
Plugin NPM Release / verify_plugins_npm (push) Blocked by required conditions
Vitest Cache Warm / dependencies (push) Waiting to run
Vitest Cache Warm / warm (push) Waiting to run
Workflow Sanity / actionlint (push) Waiting to run
Workflow Sanity / generated-doc-baselines (push) Waiting to run

Co-authored-by: steipete <58493+steipete@users.noreply.github.com>
This commit is contained in:
RoboClaw 2026-09-29 17:53:50 -07:00 • committed by GitHub
parent b725cc7b27
commit 2a6d59ecd3
No known key found for this signature in database
GPG key ID: B5690EEEBB952194

View file

@ -258,9 +258,13 @@ it.each([
it("keeps the live assistant visible while an async tool result starts and ends", async () => {
const response = createAssistantMessageEventStream();
const gate = createDeferred();
const lookupStarted = createDeferred();
const textUpdated = createDeferred();
const resultEnded = createDeferred();
const source = call("lookup");
const text = { type: "text" as const, text: "independent answer" };
const execute = vi.fn(async () => {
lookupStarted.resolve();
await gate.promise;
return { content: [{ type: "text" as const, text: "found" }], details: {} };
});
@ -283,11 +287,17 @@ it("keeps the live assistant visible while an async tool result starts and ends"
});
const resultStates: Array<AgentMessage | undefined> = [];
agent.subscribe((event) => {
if (event.type === "message_update" && event.assistantMessageEvent.type === "text_delta") {
textUpdated.resolve();
}
if (
(event.type === "message_start" || event.type === "message_end") &&
event.message.role === "toolResult"
) {
resultStates.push(agent.state.streamingMessage);
if (event.type === "message_end") {
resultEnded.resolve();
}
}
});
const run = agent.prompt("look up");
@ -299,19 +309,20 @@ it("keeps the live assistant visible while an async tool result starts and ends"
toolCall: source,
partial: assistant([source]),
});
await vi.waitFor(() => expect(execute).toHaveBeenCalledTimes(1));
await withTestTimeout(lookupStarted.promise, 1_000, "Lookup did not start");
expect(execute).toHaveBeenCalledTimes(1);
response.push({
type: "text_delta",
contentIndex: 1,
delta: text.text,
partial: assistant([source, text]),
});
await vi.waitFor(() =>
expect(agent.state.streamingMessage).toMatchObject({ role: "assistant", content: [text] }),
);
await withTestTimeout(textUpdated.promise, 1_000, "Assistant text not updated");
expect(agent.state.streamingMessage).toMatchObject({ role: "assistant", content: [text] });
const activeAssistant = agent.state.streamingMessage;
gate.resolve();
await vi.waitFor(() => expect(resultStates).toHaveLength(2));
await withTestTimeout(resultEnded.promise, 1_000, "Tool result did not end");
expect(resultStates).toHaveLength(2);
expect(resultStates).toEqual([activeAssistant, activeAssistant]);
expect(agent.state.streamingMessage).toBe(activeAssistant);
} finally {
@ -361,6 +372,9 @@ it.each(["stop", "length"] as const)(
it("persists async calls before admission, streams the remaining answer, and executes each call once", async () => {
const response = createAssistantMessageEventStream();
const gate = createDeferred();
const lookupStarted = createDeferred();
const textUpdated = createDeferred();
const resultPersisted = createDeferred();
const source = call("lookup");
const ordinary = call("ordinary", false);
const persisted: AgentMessage[] = [];
@ -372,6 +386,7 @@ it("persists async calls before admission, streams the remaining answer, and exe
expect(persisted).toContain(owner);
expect(owner?.turnId).toBeTruthy();
executionOrder.push("lookup");
lookupStarted.resolve();
await gate.promise;
return { content: [{ type: "text" as const, text: "lookup result" }], details: {} };
});
@ -415,6 +430,12 @@ it("persists async calls before admission, streams the remaining answer, and exe
(event) => {
events.push(event);
recordMessage(event, persisted);
if (event.type === "message_update" && event.assistantMessageEvent.type === "text_delta") {
textUpdated.resolve();
}
if (event.type === "message_end" && event.message.role === "toolResult") {
resultPersisted.resolve();
}
},
undefined,
streamFn,
@ -424,7 +445,8 @@ it("persists async calls before admission, streams the remaining answer, and exe
const prefix = assistant([source], "toolUse");
response.push({ type: "start", partial: assistant([]) });
response.push({ type: "toolcall_end", contentIndex: 0, toolCall: source, partial: prefix });
await vi.waitFor(() => expect(lookup).toHaveBeenCalledTimes(1));
await withTestTimeout(lookupStarted.promise, 1_000, "Lookup did not start");
expect(lookup).toHaveBeenCalledTimes(1);
const text = { type: "text" as const, text: "independent answer" };
const progress = assistant([source, text]);
response.push({ type: "text_delta", contentIndex: 1, delta: text.text, partial: progress });
@ -434,14 +456,15 @@ it("persists async calls before admission, streams the remaining answer, and exe
toolCall: ordinary,
partial: assistant([source, text, ordinary]),
});
await vi.waitFor(() =>
expect(
events.some(
(event) =>
event.type === "message_update" && event.assistantMessageEvent.type === "text_delta",
),
).toBe(true),
);
await withTestTimeout(textUpdated.promise, 1_000, "Assistant text not updated");
expect(
events.some(
(event) =>
event.type === "message_update" && event.assistantMessageEvent.type === "text_delta",
),
).toBe(true);
// Let the queued ordinary call reach the scheduler before checking it stayed deferred.
await setImmediate();
expect(ordinaryExecute).not.toHaveBeenCalled();
const textEvent = events.find(
(event) =>
@ -452,9 +475,8 @@ it("persists async calls before admission, streams the remaining answer, and exe
message: { content: [text] },
});
gate.resolve();
await vi.waitFor(() =>
expect(persisted.some((message) => message.role === "toolResult")).toBe(true),
);
await withTestTimeout(resultPersisted.promise, 1_000, "Tool result not persisted");
expect(persisted.some((message) => message.role === "toolResult")).toBe(true);
response.push({
type: "done",
reason: "toolUse",
@ -581,11 +603,15 @@ it.each(["error", "aborted", "output-limit"] as const)(
const response = createAssistantMessageEventStream();
const gate = createDeferred();
const persistTerminal = createDeferred();
const firstStarted = createDeferred();
const terminalRecorded = createDeferred();
const resultsPersisted = createDeferred();
const first = call("first");
const second = call("second");
const persisted: AgentMessage[] = [];
const events: AgentEvent[] = [];
const firstExecute = vi.fn<AgentTool["execute"]>(async (_id, _args, signal) => {
firstStarted.resolve();
await gate.promise;
expect(signal?.aborted).toBe(!outputLimit);
return { content: [], details: {} };
@ -607,11 +633,19 @@ it.each(["error", "aborted", "output-limit"] as const)(
async (event) => {
events.push(event);
recordMessage(event, persisted);
if (
event.type === "message_end" &&
event.message.role === "toolResult" &&
persisted.filter((message) => message.role === "toolResult").length === 2
) {
resultsPersisted.resolve();
}
if (
event.type === "message_end" &&
event.message.role === "assistant" &&
event.message.stopReason === stopReason
) {
terminalRecorded.resolve();
await persistTerminal.promise;
}
},
@ -633,7 +667,8 @@ it.each(["error", "aborted", "output-limit"] as const)(
toolCall: second,
partial: assistant([first, second]),
});
await vi.waitFor(() => expect(firstExecute).toHaveBeenCalledTimes(1));
await withTestTimeout(firstStarted.promise, 1_000, "First tool did not start");
expect(firstExecute).toHaveBeenCalledTimes(1);
const failure = {
...assistant([first, second], stopReason),
errorMessage: "stream failed",
@ -660,18 +695,18 @@ it.each(["error", "aborted", "output-limit"] as const)(
if (outputLimit) {
gate.resolve();
}
await vi.waitFor(() =>
expect(
persisted.some(
(message) => message.role === "assistant" && message.stopReason === stopReason,
),
).toBe(true),
);
await withTestTimeout(terminalRecorded.promise, 1_000, "Assistant terminal not recorded");
expect(
persisted.some(
(message) => message.role === "assistant" && message.stopReason === stopReason,
),
).toBe(true);
// Give finalization a turn to expose a missing persistence await while the gate is held.
await setImmediate();
expect(events.at(-1)?.type).not.toBe("agent_end");
gate.resolve();
await vi.waitFor(() =>
expect(persisted.filter((message) => message.role === "toolResult")).toHaveLength(2),
);
await withTestTimeout(resultsPersisted.promise, 1_000, "Tool results not persisted");
expect(persisted.filter((message) => message.role === "toolResult")).toHaveLength(2);
expect(secondExecute).toHaveBeenCalledTimes(outputLimit ? 1 : 0);
persistTerminal.resolve();
const result = await run;