diff --git a/src/agents/embedded-agent-subscribe.handlers.tools.test.ts b/src/agents/embedded-agent-subscribe.handlers.tools.test.ts index 9f709f39a339..3011d462b622 100644 --- a/src/agents/embedded-agent-subscribe.handlers.tools.test.ts +++ b/src/agents/embedded-agent-subscribe.handlers.tools.test.ts @@ -334,6 +334,50 @@ describe("handleToolExecutionStart read path checks", () => { expect(ctx.state.itemActiveIds.has("tool:tool-await-flush")).toBe(true); expect(ctx.state.itemActiveIds.has("command:tool-await-flush")).toBe(true); }); + + it("keeps processing tool start when progress callbacks throw", async () => { + const { ctx, warn, onExecutionPhase, onAgentEvent } = createTestContext(); + onExecutionPhase.mockImplementation(() => { + throw new Error("phase exploded"); + }); + onAgentEvent.mockImplementation(() => { + throw new Error("event exploded"); + }); + + const evt: ToolExecutionStartEvent = { + type: "tool_execution_start", + toolName: "exec", + toolCallId: "tool-callback-throws", + args: { command: "echo hi" }, + }; + + await handleToolExecutionStart(ctx, evt); + + expect(ctx.state.toolMetaById.has("tool-callback-throws")).toBe(true); + expect(ctx.state.itemStartedCount).toBe(2); + expect(warn).toHaveBeenCalledWith( + expect.stringContaining("tool execution phase callback failed"), + ); + expect(warn).toHaveBeenCalledWith(expect.stringContaining("tool agent event callback failed")); + }); + + it("does not leak unhandled rejections when tool start progress rejects", async () => { + const { ctx, warn, onAgentEvent } = createTestContext(); + onAgentEvent.mockRejectedValue(new Error("progress failed")); + + const evt: ToolExecutionStartEvent = { + type: "tool_execution_start", + toolName: "exec", + toolCallId: "tool-callback-rejects", + args: { command: "echo hi" }, + }; + + await handleToolExecutionStart(ctx, evt); + await Promise.resolve(); + + expect(ctx.state.toolMetaById.has("tool-callback-rejects")).toBe(true); + expect(warn).toHaveBeenCalledWith(expect.stringContaining("tool agent event callback failed")); + }); }); describe("handleToolExecutionEnd cron mutation tracking", () => { diff --git a/src/agents/embedded-agent-subscribe.handlers.tools.ts b/src/agents/embedded-agent-subscribe.handlers.tools.ts index e6f04b92582a..fab01d902318 100644 --- a/src/agents/embedded-agent-subscribe.handlers.tools.ts +++ b/src/agents/embedded-agent-subscribe.handlers.tools.ts @@ -299,12 +299,43 @@ function emitTrackedItemEvent(ctx: ToolHandlerContext, itemData: AgentItemEventD ...(ctx.params.sessionKey ? { sessionKey: ctx.params.sessionKey } : {}), data: itemData, }); - void ctx.params.onAgentEvent?.({ + emitAgentEventCallbackBestEffort(ctx, { stream: "item", data: itemData, }); } +function warnBestEffortEventFailure(ctx: ToolHandlerContext, label: string, error: unknown): void { + ctx.log.warn(`${label} callback failed: ${String(error)}`); +} + +function emitExecutionPhaseBestEffort( + ctx: ToolHandlerContext, + info: Parameters>[0], +): void { + try { + ctx.params.onExecutionPhase?.(info); + } catch (error) { + warnBestEffortEventFailure(ctx, "tool execution phase", error); + } +} + +function emitAgentEventCallbackBestEffort( + ctx: ToolHandlerContext, + event: Parameters>[0], +): void { + try { + const result = ctx.params.onAgentEvent?.(event); + if (isPromiseLike(result)) { + void Promise.resolve(result).catch((error: unknown) => { + warnBestEffortEventFailure(ctx, "tool agent event", error); + }); + } + } catch (error) { + warnBestEffortEventFailure(ctx, "tool agent event", error); + } +} + function readToolResultDetailsRecord(result: unknown): Record | undefined { return readRecordField(asOptionalObjectRecord(result)?.details); } @@ -780,7 +811,7 @@ export function handleToolExecutionStart( const args = evt.args; const runId = ctx.params.runId; ctx.state.toolExecutionSinceLastBlockReply = true; - ctx.params.onExecutionPhase?.({ + emitExecutionPhaseBestEffort(ctx, { phase: "tool_execution_started", tool: toolName, toolCallId, @@ -898,7 +929,7 @@ export function handleToolExecutionStart( }; emitTrackedItemEvent(ctx, itemData); // Best-effort typing signal; do not block tool summaries on slow emitters. - void ctx.params.onAgentEvent?.({ + emitAgentEventCallbackBestEffort(ctx, { stream: "tool", data: { phase: "start", @@ -1037,7 +1068,7 @@ export function handleToolExecutionUpdate( }; emitTrackedItemEvent(ctx, itemData); if (!toolProgress) { - void ctx.params.onAgentEvent?.({ + emitAgentEventCallbackBestEffort(ctx, { stream: "tool", data: { phase: "update", @@ -1075,7 +1106,7 @@ export function handleToolExecutionUpdate( ...(ctx.params.sessionKey ? { sessionKey: ctx.params.sessionKey } : {}), data: outputData, }); - void ctx.params.onAgentEvent?.({ + emitAgentEventCallbackBestEffort(ctx, { stream: "command_output", data: outputData, }); @@ -1322,7 +1353,7 @@ export async function handleToolExecutionEnd( : {}), }; emitTrackedItemEvent(ctx, itemData); - void ctx.params.onAgentEvent?.({ + emitAgentEventCallbackBestEffort(ctx, { stream: "tool", data: { phase: "result", @@ -1368,7 +1399,7 @@ export async function handleToolExecutionEnd( ...(ctx.params.sessionKey ? { sessionKey: ctx.params.sessionKey } : {}), data: approvalData, }); - void ctx.params.onAgentEvent?.({ + emitAgentEventCallbackBestEffort(ctx, { stream: "approval", data: approvalData, }); @@ -1435,7 +1466,7 @@ export async function handleToolExecutionEnd( ...(ctx.params.sessionKey ? { sessionKey: ctx.params.sessionKey } : {}), data: outputData, }); - void ctx.params.onAgentEvent?.({ + emitAgentEventCallbackBestEffort(ctx, { stream: "command_output", data: outputData, }); @@ -1461,7 +1492,7 @@ export async function handleToolExecutionEnd( ...(ctx.params.sessionKey ? { sessionKey: ctx.params.sessionKey } : {}), data: approvalData, }); - void ctx.params.onAgentEvent?.({ + emitAgentEventCallbackBestEffort(ctx, { stream: "approval", data: approvalData, }); @@ -1507,7 +1538,7 @@ export async function handleToolExecutionEnd( ...(ctx.params.sessionKey ? { sessionKey: ctx.params.sessionKey } : {}), data: patchData, }); - void ctx.params.onAgentEvent?.({ + emitAgentEventCallbackBestEffort(ctx, { stream: "patch", data: patchData, });