From 2a6d59ecd306d493380fca4f378cc36627fcc84e Mon Sep 17 00:00:00 2001 From: RoboClaw Date: Tue, 29 Sep 2026 17:53:50 -0700 Subject: [PATCH] improve(agents): avoid polling delays in async stream tests (#161434) Co-authored-by: steipete <58493+steipete@users.noreply.github.com> --- .../src/agent-loop.async-tools.test.ts | 91 +++++++++++++------ 1 file changed, 63 insertions(+), 28 deletions(-) diff --git a/packages/agent-core/src/agent-loop.async-tools.test.ts b/packages/agent-core/src/agent-loop.async-tools.test.ts index 4a49d71f7c40..23710ee9f621 100644 --- a/packages/agent-core/src/agent-loop.async-tools.test.ts +++ b/packages/agent-core/src/agent-loop.async-tools.test.ts @@ -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 = []; 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(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;