diff --git a/packages/agent/CHANGELOG.md b/packages/agent/CHANGELOG.md index 95962a6fb..2ef0416f3 100644 --- a/packages/agent/CHANGELOG.md +++ b/packages/agent/CHANGELOG.md @@ -8,6 +8,8 @@ ### Fixed +- Fixed an issue where legacy steering messages were prematurely consumed and dropped during in-flight tool execution polls + - Fixed an issue where skipped tool results in queued messages were incorrectly treated as completed, preventing necessary retries. - Improved branch summaries to preserve informative tool results from abandoned branches while filtering out redundant output. - Fixed interruptible tool waits to properly abort on host-provided IRC interrupts in addition to user steering. diff --git a/packages/agent/src/agent-loop.ts b/packages/agent/src/agent-loop.ts index 4a6b87d6a..bfe084b63 100644 --- a/packages/agent/src/agent-loop.ts +++ b/packages/agent/src/agent-loop.ts @@ -1671,7 +1671,6 @@ async function executeToolCalls( const { hasSteeringMessages, hasIrcInterrupts, - getSteeringMessages, interruptMode = "immediate", getToolContext, transformToolCallArguments, @@ -1724,19 +1723,18 @@ async function executeToolCalls( const checkSteering = async (): Promise => { // `signal` (external/user abort) is checked separately from the internal // abort controllers: once the run is externally aborted it is unwinding - // and the interrupt would be redundant. - if (!shouldInterruptImmediately || signal?.aborted) { + // and the interrupt would be redundant. Once an IRC/steering interrupt has + // already fired, do not poll again — especially not through a consuming + // legacy steering getter. + if (!shouldInterruptImmediately || interruptState.triggered || signal?.aborted) { return; } - // Prefer non-consuming peeks so queues still own their messages until - // the injection boundary. Fall back to consuming steering only for older - // integrations that never supplied a peek. + // Mid-batch steering detection must be non-consuming. If a direct + // integration only provides getSteeringMessages(), the queue drains at the + // injection boundary below; polling it here would strand or drop messages. let steeringQueued = false; if (hasSteeringMessages) { steeringQueued = await hasSteeringMessages(); - } else if (getSteeringMessages) { - const msgs = await getSteeringMessages(); - steeringQueued = (msgs?.length ?? 0) > 0; } if (steeringQueued) { // User steering upgrades an in-flight IRC interrupt: it aborts the @@ -1747,7 +1745,6 @@ async function executeToolCalls( } return; } - if (interruptState.triggered) return; if (hasIrcInterrupts && (await hasIrcInterrupts())) { // Peer IRC only aborts interruptible waits: a foreground bash / write // mid-execution keeps running so we never leave partial side effects. @@ -2078,13 +2075,13 @@ async function executeToolCalls( // While an interruptible tool is in flight (e.g. a `job`/`irc` wait // blocking on external work), queued steering or interrupting IRC would - // otherwise wait out the tool's own window. Poll the non-consuming queues + // otherwise wait out the tool's own window. Poll only non-consuming queues // and abort the shared tool signal so the boundary dequeue below injects // the message promptly. Gated on immediate-interrupt mode + an // interruptible tool; checkSteering is idempotent (no-op once triggered). const watchSteeringWhileRunning = shouldInterruptImmediately && - (hasSteeringMessages !== undefined || getSteeringMessages !== undefined || hasIrcInterrupts !== undefined) && + (hasSteeringMessages !== undefined || hasIrcInterrupts !== undefined) && records.some(r => r.tool?.interruptible === true); const steeringWatchTimer = watchSteeringWhileRunning ? setInterval(() => void checkSteering(), STEERING_INTERRUPT_POLL_MS) diff --git a/packages/agent/test/agent-loop.test.ts b/packages/agent/test/agent-loop.test.ts index 1070ed63d..b04d64da5 100644 --- a/packages/agent/test/agent-loop.test.ts +++ b/packages/agent/test/agent-loop.test.ts @@ -1258,6 +1258,75 @@ describe("agentLoop with AgentMessage", () => { ).toBe(true); }); + it("keeps legacy steering queued until the injection boundary when no non-consuming peek exists", async () => { + const toolSchema = type({ value: "string" }); + const executed: string[] = []; + let steerReady = false; + let steeringDrained = false; + const steeringMessage = createUserMessage("queued steering"); + + const tool: AgentTool = { + name: "echo", + label: "Echo", + description: "Echo tool", + parameters: toolSchema, + concurrency: "exclusive", + async execute(_toolCallId, params) { + executed.push(params.value); + if (params.value === "first") steerReady = true; + return { + content: [{ type: "text", text: `ok:${params.value}` }], + details: { value: params.value }, + }; + }, + }; + + const context: AgentContext = { systemPrompt: [""], messages: [], tools: [tool] }; + const mock = createMockModel({ + responses: [ + { + content: [ + { type: "toolCall", id: "tool-1", name: "echo", arguments: { value: "first" } }, + { type: "toolCall", id: "tool-2", name: "echo", arguments: { value: "second" } }, + ], + }, + { content: ["done"] }, + ], + }); + const config: AgentLoopConfig = { + model: mock.model, + convertToLlm: identityConverter, + interruptMode: "immediate", + // Legacy integrations may only have the consuming dequeue. The + // mid-batch poll must not call it; the boundary below owns the drain. + getSteeringMessages: async () => { + if (steerReady && !steeringDrained) { + steeringDrained = true; + return [steeringMessage]; + } + return []; + }, + }; + + const events: AgentEvent[] = []; + for await (const event of agentLoop([createUserMessage("start")], context, config, undefined, mock.stream)) { + events.push(event); + } + + expect(executed).toEqual(["first", "second"]); + expect(steeringDrained).toBe(true); + expect( + events.some( + e => e.type === "message_start" && e.message.role === "user" && e.message.content === "queued steering", + ), + ).toBe(true); + expect( + mock.calls[1]?.context.messages.some( + m => m.role === "user" && typeof m.content === "string" && m.content === "queued steering", + ), + ).toBe(true); + }); + it("does not abort a non-interruptible foreground tool when only IRC is queued", async () => { const toolSchema = type({}); let ircReady = false; diff --git a/packages/coding-agent/test/bash-executor.test.ts b/packages/coding-agent/test/bash-executor.test.ts index 8058e0279..c8436d7da 100644 --- a/packages/coding-agent/test/bash-executor.test.ts +++ b/packages/coding-agent/test/bash-executor.test.ts @@ -137,7 +137,6 @@ describe("executeBash", () => { expect(result.output.trim()).toBe(tempDir); }); - it("honors symlinked cwd requests in persistent shells", async () => { if (process.platform === "win32") { return;