diff --git a/packages/coding-agent/CHANGELOG.md b/packages/coding-agent/CHANGELOG.md index 96efe0f81..732addfad 100644 --- a/packages/coding-agent/CHANGELOG.md +++ b/packages/coding-agent/CHANGELOG.md @@ -14,6 +14,10 @@ - Refined intent parameter guidance to require concise 2-6 word sentences in present participle form - Centralized per-tool timeout constants and clamping into `tool-timeouts.ts` +### Fixed + +- Fixed TTSR violations during subagent execution aborting the entire subagent run; `#waitForPostPromptRecovery()` now also awaits agent idle after TTSR/retry gates resolve, preventing `prompt()` from returning while a fire-and-forget `agent.continue()` is still streaming + ## [13.3.7] - 2026-02-27 ### Breaking Changes diff --git a/packages/coding-agent/src/session/agent-session.ts b/packages/coding-agent/src/session/agent-session.ts index d459dda4c..eec484d2f 100644 --- a/packages/coding-agent/src/session/agent-session.ts +++ b/packages/coding-agent/src/session/agent-session.ts @@ -732,18 +732,31 @@ export class AgentSession { } /** - * Wait for both retry and TTSR resume to settle. - * Loops because a TTSR continuation can trigger a retry (or vice-versa). + * Wait for retry, TTSR resume, and any background continuation to settle. + * Loops because a TTSR continuation can trigger a retry (or vice-versa), + * and fire-and-forget `agent.continue()` may still be streaming after + * the TTSR resume gate resolves. */ async #waitForPostPromptRecovery(): Promise { - while (this.#retryPromise || this.#ttsrResumePromise) { + while (true) { if (this.#retryPromise) { await this.#retryPromise; continue; } if (this.#ttsrResumePromise) { await this.#ttsrResumePromise; + continue; } + // A TTSR continuation (fire-and-forget agent.continue()) may still be + // streaming after the resume gate resolved on its first assistant + // message_end. Without this check, prompt() returns while the agent + // is still processing tool calls, causing subagent executors to hit + // AgentBusyError and dispose the session prematurely. + if (this.agent.state.isStreaming) { + await this.agent.waitForIdle(); + continue; + } + break; } } diff --git a/packages/coding-agent/test/agent-session-concurrent.test.ts b/packages/coding-agent/test/agent-session-concurrent.test.ts index c915ff4df..cb031a69b 100644 --- a/packages/coding-agent/test/agent-session-concurrent.test.ts +++ b/packages/coding-agent/test/agent-session-concurrent.test.ts @@ -6,8 +6,8 @@ import { afterEach, beforeEach, describe, expect, it } from "bun:test"; import * as fs from "node:fs"; import * as os from "node:os"; import * as path from "node:path"; -import { Agent, AgentBusyError } from "@oh-my-pi/pi-agent-core"; -import { type AssistantMessage, getBundledModel } from "@oh-my-pi/pi-ai"; +import { Agent, AgentBusyError, type AgentTool } from "@oh-my-pi/pi-agent-core"; +import { type AssistantMessage, getBundledModel, type ToolCall } from "@oh-my-pi/pi-ai"; import { AssistantMessageEventStream } from "@oh-my-pi/pi-ai/utils/event-stream"; import type { Rule } from "@oh-my-pi/pi-coding-agent/capability/rule"; import { ModelRegistry } from "@oh-my-pi/pi-coding-agent/config/model-registry"; @@ -17,6 +17,7 @@ import { AgentSession } from "@oh-my-pi/pi-coding-agent/session/agent-session"; import { AuthStorage } from "@oh-my-pi/pi-coding-agent/session/auth-storage"; import { SessionManager } from "@oh-my-pi/pi-coding-agent/session/session-manager"; import { Snowflake } from "@oh-my-pi/pi-utils"; +import { Type } from "@sinclair/typebox"; // Mock stream that mimics AssistantMessageEventStream class MockAssistantStream extends AssistantMessageEventStream {} @@ -492,4 +493,137 @@ describe("AgentSession TTSR resume gate", () => { expect(session.isStreaming).toBe(false); }); + + it("prompt() waits for TTSR continuation with tool calls to finish", async () => { + const model = getBundledModel("anthropic", "claude-sonnet-4-5")!; + let streamCallCount = 0; + let toolExecutionFinished = false; + let allTurnsCompleted = false; + + const ttsrManager = new TtsrManager({ + enabled: true, + contextMode: "discard", + interruptMode: "always", + repeatMode: "once", + repeatGap: 10, + }); + ttsrManager.addRule(testRule); + + const mockTool: AgentTool = { + name: "mock_edit", + label: "Mock Edit", + description: "A mock edit tool", + parameters: Type.Object({}), + execute: async () => { + await Bun.sleep(100); + toolExecutionFinished = true; + return { content: [{ type: "text" as const, text: "edit applied" }] }; + }, + }; + + const toolCallContent: ToolCall = { + type: "toolCall", + id: "call_test_001", + name: "mock_edit", + arguments: {}, + }; + + function makeToolCallMsg(): AssistantMessage { + return { + role: "assistant", + content: [toolCallContent], + api: "anthropic-messages", + provider: "anthropic", + model: "mock", + usage: { + input: 0, + output: 0, + cacheRead: 0, + cacheWrite: 0, + totalTokens: 0, + cost: { input: 0, output: 0, cacheRead: 0, cacheWrite: 0, total: 0 }, + }, + stopReason: "toolUse", + timestamp: Date.now(), + }; + } + + const agent = new Agent({ + getApiKey: () => "test-key", + initialState: { model, systemPrompt: "Test", tools: [mockTool] }, + streamFn: (_model, _context, options) => { + streamCallCount++; + const stream = new MockAssistantStream(); + const signal = options?.signal; + + if (streamCallCount === 1) { + // First stream: emit text that triggers TTSR, then respond to abort + queueMicrotask(() => { + const partial = makeMsg(""); + stream.push({ type: "start", partial }); + stream.push({ + type: "text_delta", + contentIndex: 0, + delta: "let val = result.unwrap(", + partial: makeMsg("let val = result.unwrap("), + }); + const checkAbort = () => { + if (signal?.aborted) { + stream.push({ + type: "error", + reason: "aborted", + error: makeMsg("let val = result.unwrap(", "aborted"), + }); + } else { + setTimeout(checkAbort, 2); + } + }; + checkAbort(); + }); + } else if (streamCallCount === 2) { + // Continuation: return assistant message with a tool call + setTimeout(() => { + const msg = makeToolCallMsg(); + stream.push({ type: "start", partial: msg }); + stream.push({ type: "done", reason: "toolUse", message: msg }); + }, 10); + } else { + // After tool execution: return final response + setTimeout(() => { + allTurnsCompleted = true; + const msg = makeMsg('Fixed: let val = result.expect("msg")'); + stream.push({ type: "start", partial: msg }); + stream.push({ type: "done", reason: "stop", message: msg }); + }, 10); + } + + return stream; + }, + }); + + const sessionManager = SessionManager.inMemory(); + const settings = Settings.isolated(); + const authStorage = await AuthStorage.create(path.join(tempDir, "testauth-tool.db")); + const modelRegistry = new ModelRegistry(authStorage, path.join(tempDir, "models.yml")); + authStorage.setRuntimeApiKey("anthropic", "test-key"); + + session = new AgentSession({ + agent, + sessionManager, + settings, + modelRegistry, + ttsrManager, + }); + + // prompt() must block until the TTSR continuation (including tool execution) completes. + // Before the fix, prompt() returned after the continuation's first assistant message_end, + // while the agent was still executing tool calls in the background. + await session.prompt("Write some Rust code"); + + // By the time prompt() returns, ALL turns must have completed + expect(toolExecutionFinished).toBe(true); + expect(allTurnsCompleted).toBe(true); + expect(streamCallCount).toBeGreaterThanOrEqual(3); + expect(session.isStreaming).toBe(false); + }); });