From bfae4d46c38e75f6bb78eeb2f86a75109edf8563 Mon Sep 17 00:00:00 2001 From: can1357 Date: Sun, 17 May 2026 13:15:07 +0200 Subject: [PATCH] fix(coding-agent): tracked acp tool args by session for replay - Tracked ACP tool-call inputs per session and replayed them via `toolArgsById`/`getToolArgs` plumbing. - Merged ACP tool execution end content from start and result events so command output replay preserves original args. - Scoped ACP async-job draining by session `ownerId` and `agentId` with in-flight tracking and permission-gated deferred turns. - Refactored compaction telemetry and async tests with per-test telemetry setup and asynchronous teardown resets. --- .../agent/test/compaction-telemetry.test.ts | 29 +-- packages/coding-agent/CHANGELOG.md | 8 +- .../coding-agent/src/async/job-manager.ts | 127 ++++++++----- packages/coding-agent/src/main.ts | 5 +- .../coding-agent/src/modes/acp/acp-agent.ts | 31 +++- .../src/modes/acp/acp-event-mapper.ts | 19 +- .../coding-agent/src/session/agent-session.ts | 17 +- packages/coding-agent/test/acp-agent.test.ts | 2 +- .../test/acp-event-mapper.test.ts | 173 +++++++++++++++++- .../test/agent-session-acp-permission.test.ts | 13 -- .../test/agent-session-concurrent.test.ts | 161 +++++++++++++++- .../test/async-job-manager.test.ts | 99 ++++++++-- 12 files changed, 573 insertions(+), 111 deletions(-) diff --git a/packages/agent/test/compaction-telemetry.test.ts b/packages/agent/test/compaction-telemetry.test.ts index ca87c1c50..87e3b8cad 100644 --- a/packages/agent/test/compaction-telemetry.test.ts +++ b/packages/agent/test/compaction-telemetry.test.ts @@ -2,12 +2,12 @@ * Tests for OpenTelemetry instrumentation around oneshot LLM calls: * compaction summaries, handoff document, branch summary. * - * Mirrors the InMemorySpanExporter + AsyncLocalStorageContextManager setup - * used by `otel.test.ts`. Spies on `completeSimple` to avoid real HTTP traffic + * Uses a per-test InMemorySpanExporter and explicit tracer. Spies on + * `completeSimple` to avoid real HTTP traffic * while exercising the chat-span lifecycle (`startChatSpan` → * `runInActiveSpan` → `finishChatSpan` / `failChatSpan`). */ -import { afterAll, afterEach, beforeAll, describe, expect, it, vi } from "bun:test"; +import { afterEach, beforeEach, describe, expect, it, vi } from "bun:test"; import { type CompactionPreparation, compact, @@ -27,8 +27,7 @@ import { import type { AgentMessage } from "@oh-my-pi/pi-agent-core/types"; import type { AssistantMessage, Model, Usage } from "@oh-my-pi/pi-ai"; import * as ai from "@oh-my-pi/pi-ai"; -import { context, SpanStatusCode, trace } from "@opentelemetry/api"; -import { AsyncLocalStorageContextManager } from "@opentelemetry/context-async-hooks"; +import { SpanStatusCode } from "@opentelemetry/api"; import { BasicTracerProvider, InMemorySpanExporter, @@ -49,28 +48,18 @@ const MODEL: Model = { maxTokens: 32_768, }; -const exporter = new InMemorySpanExporter(); +let exporter: InMemorySpanExporter; let provider: BasicTracerProvider; -let contextManager: AsyncLocalStorageContextManager; -beforeAll(() => { - trace.disable(); - context.disable(); - contextManager = new AsyncLocalStorageContextManager().enable(); - context.setGlobalContextManager(contextManager); +beforeEach(() => { + exporter = new InMemorySpanExporter(); provider = new BasicTracerProvider({ spanProcessors: [new SimpleSpanProcessor(exporter)] }); - trace.setGlobalTracerProvider(provider); }); -afterEach(() => { +afterEach(async () => { exporter.reset(); - vi.restoreAllMocks(); -}); - -afterAll(async () => { await provider.shutdown(); - context.disable(); - trace.disable(); + vi.restoreAllMocks(); }); function makeTelemetryConfig(): AgentTelemetryConfig { diff --git a/packages/coding-agent/CHANGELOG.md b/packages/coding-agent/CHANGELOG.md index 61770d4c3..cc9de6d42 100644 --- a/packages/coding-agent/CHANGELOG.md +++ b/packages/coding-agent/CHANGELOG.md @@ -1,12 +1,14 @@ # Changelog ## [Unreleased] -### Fixed - -- Fixed ACP ordinary file-editing calls (`edit`, `write`, `ast_edit`) incorrectly requesting `session/request_permission` before every call, while keeping permission prompts for edit operations that delete or move files; permission requests now report the gated tool call as `pending` so clients can render the approval UI instead of returning `Permission request cancelled` without a visible prompt. ### Fixed +- Fixed ACP command and custom tool-call notifications to carry the original tool arguments in replayed and final updates, so command text is preserved and raw input is no longer wrapped +- Fixed ACP async-job draining to be scoped by session owner so `getAsyncJobSnapshot` and `drainAsyncJobDeliveriesForAcp` no longer consume or expose jobs from other sessions +- Fixed async job status reporting to include in-flight completions so queued/delivering indicators remain accurate while callbacks are still running +- Fixed `deferAgentInitiatedTurns` handling during ACP async-job draining so background completion follow-up turns are delivered even when agent-initiated turns are deferred +- Fixed ACP ordinary file-editing calls (`edit`, `write`, `ast_edit`) incorrectly requesting `session/request_permission` before every call, while keeping permission prompts for edit operations that delete or move files; permission requests now report the gated tool call as `pending` so clients can render the approval UI instead of returning `Permission request cancelled` without a visible prompt. ([#1134](https://github.com/can1357/oh-my-pi/pull/1134) by [@jiwangyihao](https://github.com/jiwangyihao)) - Fixed the session tree selector to preserve a readable message column when deeply nested branch gutters would otherwise consume the viewport. ([#1144](https://github.com/can1357/oh-my-pi/issues/1144)) ## [15.1.3] - 2026-05-17 diff --git a/packages/coding-agent/src/async/job-manager.ts b/packages/coding-agent/src/async/job-manager.ts index 6e36f22b8..bde46c4b1 100644 --- a/packages/coding-agent/src/async/job-manager.ts +++ b/packages/coding-agent/src/async/job-manager.ts @@ -37,6 +37,8 @@ interface AsyncJobDelivery { attempt: number; nextAttemptAt: number; lastError?: string; + ownerId?: string; + promise?: Promise; } export interface AsyncJobDeliveryState { @@ -82,6 +84,7 @@ export class AsyncJobManager { readonly #jobs = new Map(); readonly #deliveries: AsyncJobDelivery[] = []; + readonly #inFlightDeliveries: AsyncJobDelivery[] = []; readonly #suppressedDeliveries = new Set(); readonly #watchedJobs = new Set(); readonly #evictionTimers = new Map(); @@ -221,16 +224,17 @@ export class AsyncJobManager { getDeliveryState(filter?: AsyncJobFilter): AsyncJobDeliveryState { const deliveries = this.#filterDeliveries(filter); + const inFlightDeliveries = this.#filterInFlightDeliveries(filter); const nextRetryAt = deliveries.reduce((next, delivery) => { if (next === undefined) return delivery.nextAttemptAt; return Math.min(next, delivery.nextAttemptAt); }, undefined); return { - queued: deliveries.length, - delivering: this.#deliveryLoop !== undefined && deliveries.length > 0, + queued: deliveries.length + inFlightDeliveries.length, + delivering: inFlightDeliveries.length > 0 || (this.#deliveryLoop !== undefined && deliveries.length > 0), nextRetryAt, - pendingJobIds: deliveries.map(delivery => delivery.jobId), + pendingJobIds: deliveries.concat(inFlightDeliveries).map(delivery => delivery.jobId), }; } @@ -303,6 +307,12 @@ export class AsyncJobManager { if (delivered) continue; return false; } + const inFlightDeliveries = this.#filterInFlightDeliveries(); + if (inFlightDeliveries.length > 0 && this.#filterDeliveries().length === 0) { + const delivered = await this.#waitForDeliveryPromise(inFlightDeliveries[0]?.promise, deadline); + if (delivered) continue; + return false; + } this.#ensureDeliveryLoop(); const loop = this.#deliveryLoop; @@ -338,6 +348,7 @@ export class AsyncJobManager { this.#clearEvictionTimers(); this.#jobs.clear(); this.#deliveries.length = 0; + this.#inFlightDeliveries.length = 0; this.#suppressedDeliveries.clear(); this.#watchedJobs.clear(); return drained; @@ -398,49 +409,51 @@ export class AsyncJobManager { #filterDeliveries(filter?: AsyncJobFilter): AsyncJobDelivery[] { const ownerId = filter?.ownerId; - if (!ownerId) return this.#deliveries; - return this.#deliveries.filter(delivery => this.#jobs.get(delivery.jobId)?.ownerId === ownerId); + if (!ownerId) return this.#deliveries.filter(delivery => !this.isDeliverySuppressed(delivery.jobId)); + return this.#deliveries.filter( + delivery => delivery.ownerId === ownerId && !this.isDeliverySuppressed(delivery.jobId), + ); + } + + #filterInFlightDeliveries(filter?: AsyncJobFilter): AsyncJobDelivery[] { + const ownerId = filter?.ownerId; + if (!ownerId) return this.#inFlightDeliveries.filter(delivery => !this.isDeliverySuppressed(delivery.jobId)); + return this.#inFlightDeliveries.filter( + delivery => delivery.ownerId === ownerId && !this.isDeliverySuppressed(delivery.jobId), + ); } async #deliverNextFiltered(filter: AsyncJobFilter, deadline: number): Promise { - let selected: AsyncJobDelivery | undefined; - for (const delivery of this.#deliveries) { - if (this.#jobs.get(delivery.jobId)?.ownerId !== filter.ownerId) continue; - if (this.isDeliverySuppressed(delivery.jobId)) continue; - if (!selected || delivery.nextAttemptAt < selected.nextAttemptAt) { - selected = delivery; + while (true) { + let selected: AsyncJobDelivery | undefined; + for (const delivery of this.#deliveries) { + if (delivery.ownerId !== filter.ownerId) continue; + if (this.isDeliverySuppressed(delivery.jobId)) continue; + if (!selected || delivery.nextAttemptAt < selected.nextAttemptAt) { + selected = delivery; + } } - } - if (!selected) return true; - const now = Date.now(); - if (selected.nextAttemptAt > now) { - if (selected.nextAttemptAt > deadline) return false; - await Bun.sleep(selected.nextAttemptAt - now); - } - - const index = this.#deliveries.indexOf(selected); - if (index === -1) return true; - - try { - await this.#onJobComplete(selected.jobId, selected.text, this.#jobs.get(selected.jobId)); - this.#deliveries.splice(index, 1); - } catch (error) { - selected.attempt += 1; - selected.lastError = error instanceof Error ? error.message : String(error); - selected.nextAttemptAt = Date.now() + this.#getRetryDelay(selected.attempt); - this.#deliveries.splice(index, 1); - if (!this.isDeliverySuppressed(selected.jobId)) { - this.#deliveries.push(selected); + if (!selected) { + const inFlight = this.#filterInFlightDeliveries(filter); + if (inFlight.length === 0) return true; + return this.#waitForDeliveryPromise(inFlight[0]?.promise, deadline); } - logger.warn("Async job completion delivery failed", { - jobId: selected.jobId, - attempt: selected.attempt, - nextRetryAt: selected.nextAttemptAt, - error: selected.lastError, - }); + + const now = Date.now(); + if (selected.nextAttemptAt > now) { + if (selected.nextAttemptAt > deadline) return false; + await Bun.sleep(selected.nextAttemptAt - now); + continue; + } + + const index = this.#deliveries.indexOf(selected); + if (index === -1) continue; + this.#deliveries.splice(index, 1); + if (this.isDeliverySuppressed(selected.jobId)) continue; + + return this.#waitForDeliveryPromise(this.#deliverDelivery(selected), deadline); } - return true; } isDeliverySuppressed(jobId: string): boolean { @@ -457,6 +470,7 @@ export class AsyncJobManager { text, attempt: 0, nextAttemptAt: Date.now(), + ownerId: this.#jobs.get(jobId)?.ownerId, }); this.#ensureDeliveryLoop(); } @@ -492,20 +506,25 @@ export class AsyncJobManager { if (this.#deliveries[0] !== delivery) { continue; } - // Check again after sleep if (this.isDeliverySuppressed(delivery.jobId)) { this.#deliveries.shift(); continue; } + this.#deliveries.shift(); + await this.#deliverDelivery(delivery); + } + } + + #deliverDelivery(delivery: AsyncJobDelivery): Promise { + const promise = (async () => { + this.#inFlightDeliveries.push(delivery); try { await this.#onJobComplete(delivery.jobId, delivery.text, this.#jobs.get(delivery.jobId)); - this.#deliveries.shift(); } catch (error) { delivery.attempt += 1; delivery.lastError = error instanceof Error ? error.message : String(error); delivery.nextAttemptAt = Date.now() + this.#getRetryDelay(delivery.attempt); - this.#deliveries.shift(); if (!this.isDeliverySuppressed(delivery.jobId)) { this.#deliveries.push(delivery); } @@ -515,8 +534,32 @@ export class AsyncJobManager { nextRetryAt: delivery.nextAttemptAt, error: delivery.lastError, }); + } finally { + const index = this.#inFlightDeliveries.indexOf(delivery); + if (index !== -1) this.#inFlightDeliveries.splice(index, 1); + if (this.#deliveries.length > 0) this.#ensureDeliveryLoop(); } + })(); + delivery.promise = promise; + return promise; + } + + async #waitForDeliveryPromise(promise: Promise | undefined, deadline: number): Promise { + if (!promise) return true; + if (deadline === Number.POSITIVE_INFINITY) { + await promise; + return true; } + const remainingMs = deadline - Date.now(); + if (remainingMs <= 0) return false; + let timedOut = false; + await Promise.race([ + promise, + Bun.sleep(remainingMs).then(() => { + timedOut = true; + }), + ]); + return !timedOut; } #getRetryDelay(attempt: number): number { diff --git a/packages/coding-agent/src/main.ts b/packages/coding-agent/src/main.ts index 6ace23ea1..14b4f137f 100644 --- a/packages/coding-agent/src/main.ts +++ b/packages/coding-agent/src/main.ts @@ -56,6 +56,7 @@ import { discoverAuthStorage, } from "./sdk"; import type { AgentSession } from "./session/agent-session"; +import type { AuthStorage } from "./session/auth-storage"; import { resolveResumableSession, type SessionInfo, SessionManager } from "./session/session-manager"; import { resolvePromptInput } from "./system-prompt"; import type { LspStartupServerInfo } from "./tools"; @@ -202,7 +203,7 @@ interface AcpSessionFactoryOptions { baseOptions: CreateAgentSessionOptions; settings: Settings; sessionDir?: string; - authStorage: Awaited>; + authStorage: AuthStorage; modelRegistry: ModelRegistry; parsedArgs: Pick; rawArgs: string[]; @@ -213,6 +214,7 @@ function createAcpSessionFactory(args: AcpSessionFactoryOptions): AcpSessionFact return async cwd => { const nextSettings = await args.settings.cloneForCwd(cwd); const nextSessionManager = SessionManager.create(cwd, args.sessionDir); + const agentId = `acp:${nextSessionManager.getSessionId()}`; const { session: nextSession } = await args.createSession({ ...args.baseOptions, cwd, @@ -220,6 +222,7 @@ function createAcpSessionFactory(args: AcpSessionFactoryOptions): AcpSessionFact settings: nextSettings, authStorage: args.authStorage, modelRegistry: args.modelRegistry, + agentId, hasUI: false, }); if (args.parsedArgs.apiKey && !args.baseOptions.model && nextSession.model) { diff --git a/packages/coding-agent/src/modes/acp/acp-agent.ts b/packages/coding-agent/src/modes/acp/acp-agent.ts index 91ac9d6be..7d3c8ec4c 100644 --- a/packages/coding-agent/src/modes/acp/acp-agent.ts +++ b/packages/coding-agent/src/modes/acp/acp-agent.ts @@ -140,6 +140,7 @@ type ManagedSessionRecord = { promptQueue: PromptQueueState; liveMessageId: string | undefined; liveMessageProgress: { textEmitted: boolean; thoughtEmitted: boolean } | undefined; + toolArgsById: Map; extensionsConfigured: boolean; // Installed inside `#scheduleBootstrapUpdates` (post-race-guard); released // in `#disposeSessionRecord`. Lives independent of any prompt turn. @@ -975,6 +976,7 @@ export class AcpAgent implements Agent { promptQueue: { promise: Promise.resolve(), release: undefined }, liveMessageId: undefined, liveMessageProgress: undefined, + toolArgsById: new Map(), extensionsConfigured: false, lifetimeUnsubscribe: undefined, }; @@ -1037,14 +1039,22 @@ export class AcpAgent implements Agent { return; } + if (event.type === "tool_execution_start" || event.type === "tool_execution_update") { + record.toolArgsById.set(event.toolCallId, event.args); + } + this.#prepareLiveAssistantMessage(record, event); for (const notification of mapAgentSessionEventToAcpSessionUpdates(event, record.session.sessionId, { getMessageId: message => this.#getLiveMessageId(record, message), getMessageProgress: message => this.#getLiveMessageProgress(record, message), + getToolArgs: toolCallId => record.toolArgsById.get(toolCallId), cwd: record.session.sessionManager.getCwd(), })) { await this.#connection.sessionUpdate(notification); } + if (event.type === "tool_execution_end") { + record.toolArgsById.delete(event.toolCallId); + } this.#clearLiveAssistantMessageAfterEvent(record, event); if (event.type === "agent_end") { @@ -1613,12 +1623,14 @@ export class AcpAgent implements Agent { async #replaySessionHistory(record: ManagedSessionRecord): Promise { const cwd = record.session.sessionManager.getCwd(); const replayedToolCallIds = new Set(); + const replayedToolCallArgs = new Map(); for (const message of record.session.sessionManager.buildSessionContext().messages as ReplayableMessage[]) { for (const notification of this.#messageToReplayNotifications( record.session.sessionId, message, cwd, replayedToolCallIds, + replayedToolCallArgs, )) { await this.#connection.sessionUpdate(notification); } @@ -1630,9 +1642,10 @@ export class AcpAgent implements Agent { message: ReplayableMessage, cwd: string, replayedToolCallIds: Set, + replayedToolCallArgs: Map, ): SessionNotification[] { if (message.role === "assistant") { - return this.#replayAssistantMessage(sessionId, message, cwd, replayedToolCallIds); + return this.#replayAssistantMessage(sessionId, message, cwd, replayedToolCallIds, replayedToolCallArgs); } if ( message.role === "user" || @@ -1660,7 +1673,10 @@ export class AcpAgent implements Agent { toolCallId: message.toolCallId, toolName: message.toolName, }, - { includeStart: !replayedToolCallIds.has(message.toolCallId) }, + { + includeStart: !replayedToolCallIds.has(message.toolCallId), + toolArgs: replayedToolCallArgs.get(message.toolCallId), + }, ); } if ( @@ -1683,6 +1699,7 @@ export class AcpAgent implements Agent { message: ReplayableMessage, cwd: string, replayedToolCallIds: Set, + replayedToolCallArgs: Map, ): SessionNotification[] { const notifications: SessionNotification[] = []; const messageId = crypto.randomUUID(); @@ -1734,6 +1751,7 @@ export class AcpAgent implements Agent { }); notifications.push({ sessionId, update }); replayedToolCallIds.add(toolItem.id); + replayedToolCallArgs.set(toolItem.id, args); } } } @@ -1755,7 +1773,7 @@ export class AcpAgent implements Agent { return normalizeReplayToolArguments(item.arguments).args; } if (item.type === "tool_use" && "input" in item) { - return { input: item.input }; + return item.input; } return {}; } @@ -1764,7 +1782,7 @@ export class AcpAgent implements Agent { sessionId: string, cwd: string, message: Required> & ReplayableMessage, - options: { includeStart?: boolean } = {}, + options: { includeStart?: boolean; toolArgs?: unknown } = {}, ): SessionNotification[] { const args = this.#buildReplayToolArgs(message.details); const startEvent: AgentSessionEvent = { @@ -1784,7 +1802,10 @@ export class AcpAgent implements Agent { errorMessage: message.errorMessage, }, }; - const notifications = mapAgentSessionEventToAcpSessionUpdates(endEvent, sessionId, { cwd }); + const notifications = mapAgentSessionEventToAcpSessionUpdates(endEvent, sessionId, { + cwd, + getToolArgs: toolCallId => (toolCallId === message.toolCallId ? options.toolArgs : undefined), + }); if (options.includeStart === false) { return notifications; } diff --git a/packages/coding-agent/src/modes/acp/acp-event-mapper.ts b/packages/coding-agent/src/modes/acp/acp-event-mapper.ts index 97c162531..725e86fae 100644 --- a/packages/coding-agent/src/modes/acp/acp-event-mapper.ts +++ b/packages/coding-agent/src/modes/acp/acp-event-mapper.ts @@ -18,6 +18,7 @@ interface MessageProgress { interface AcpEventMapperOptions { getMessageId?: (message: unknown) => string | undefined; getMessageProgress?: (message: unknown) => MessageProgress | undefined; + getToolArgs?: (toolCallId: string) => unknown; /** * Session cwd. Tool call locations sent to ACP clients must be absolute * (the editor host needs them to open or focus files). When provided, @@ -185,8 +186,11 @@ export function mapAgentSessionEventToAcpSessionUpdates( return [toSessionNotification(sessionId, update)]; } case "tool_execution_end": { - const diffContent = extractDiffToolCallContent(event.result); - const content = [...diffContent, ...extractToolCallContent(event.result)]; + const resultContent = [...extractDiffToolCallContent(event.result), ...extractToolCallContent(event.result)]; + const content = mergeToolUpdateContent( + buildToolStartContent(event.toolName, getToolExecutionEndArgs(event, options)), + resultContent, + ); const update: SessionUpdate = { sessionUpdate: "tool_call_update", toolCallId: event.toolCallId, @@ -380,6 +384,7 @@ function extractTodoEntries(phases: unknown[]): Array<{ content: string; status: function isTodoStatus(status: unknown): status is TodoStatus { return status === "pending" || status === "in_progress" || status === "completed" || status === "abandoned"; +} export function buildToolCallStartUpdate(input: { toolCallId: string; toolName: string; @@ -419,6 +424,16 @@ export function normalizeReplayToolArguments(value: unknown): { args: unknown } } } +function getToolExecutionEndArgs( + event: Extract, + options: AcpEventMapperOptions, +): unknown { + if ("args" in event) { + return (event as { args?: unknown }).args; + } + return options.getToolArgs?.(event.toolCallId); +} + function buildToolStartContent(toolName: string, args: unknown): ToolCallContent[] { if (!isCommandToolName(toolName)) { return []; diff --git a/packages/coding-agent/src/session/agent-session.ts b/packages/coding-agent/src/session/agent-session.ts index 4e86f9a8a..9cb07c8f5 100644 --- a/packages/coding-agent/src/session/agent-session.ts +++ b/packages/coding-agent/src/session/agent-session.ts @@ -732,6 +732,7 @@ export class AgentSession { #planReferenceSent = false; #planReferencePath = "local://PLAN.md"; #clientBridge: ClientBridge | undefined; + #allowAcpAgentInitiatedTurns = false; /** Per-session memory of allow_always / reject_always decisions for gated tools. */ #acpPermissionDecisions: Map = new Map(); @@ -2777,9 +2778,15 @@ export class AgentSession { const ownerFilter = this.#agentId ? { ownerId: this.#agentId } : undefined; const before = manager.getDeliveryState(ownerFilter); if (before.queued === 0 && !before.delivering) return false; - const drained = await manager.drainDeliveries({ timeoutMs: options?.timeoutMs, filter: ownerFilter }); - const after = manager.getDeliveryState(ownerFilter); - return drained && (before.queued !== after.queued || before.delivering !== after.delivering); + const previousAllowAcpAgentInitiatedTurns = this.#allowAcpAgentInitiatedTurns; + this.#allowAcpAgentInitiatedTurns = true; + try { + const drained = await manager.drainDeliveries({ timeoutMs: options?.timeoutMs, filter: ownerFilter }); + const after = manager.getDeliveryState(ownerFilter); + return drained && (before.queued !== after.queued || before.delivering !== after.delivering); + } finally { + this.#allowAcpAgentInitiatedTurns = previousAllowAcpAgentInitiatedTurns; + } } /** Most recent assistant message in agent state. */ @@ -4419,7 +4426,7 @@ export class AgentSession { if (options?.deliverAs === "nextTurn") { if (options?.triggerTurn) { - if (this.#clientBridge?.deferAgentInitiatedTurns) { + if (this.#clientBridge?.deferAgentInitiatedTurns && !this.#allowAcpAgentInitiatedTurns) { this.#queueHiddenNextTurnMessage(appMessage, false); return; } @@ -4438,7 +4445,7 @@ export class AgentSession { } if (options?.triggerTurn) { - if (this.#clientBridge?.deferAgentInitiatedTurns) { + if (this.#clientBridge?.deferAgentInitiatedTurns && !this.#allowAcpAgentInitiatedTurns) { this.#queueHiddenNextTurnMessage(appMessage, false); return; } diff --git a/packages/coding-agent/test/acp-agent.test.ts b/packages/coding-agent/test/acp-agent.test.ts index 1d631ea71..9c668ea5f 100644 --- a/packages/coding-agent/test/acp-agent.test.ts +++ b/packages/coding-agent/test/acp-agent.test.ts @@ -888,7 +888,7 @@ describe("ACP agent", () => { expect.objectContaining({ sessionUpdate: "tool_call", toolCallId: "toolu_custom", - rawInput: { input: "raw custom payload" }, + rawInput: "raw custom payload", }), ); diff --git a/packages/coding-agent/test/acp-event-mapper.test.ts b/packages/coding-agent/test/acp-event-mapper.test.ts index 97e69e8ec..7a501060c 100644 --- a/packages/coding-agent/test/acp-event-mapper.test.ts +++ b/packages/coding-agent/test/acp-event-mapper.test.ts @@ -1,13 +1,18 @@ import { describe, expect, it } from "bun:test"; +import * as fs from "node:fs"; +import * as os from "node:os"; import path from "node:path"; -import type { SessionNotification } from "@agentclientprotocol/sdk"; +import type { AgentSideConnection, SessionNotification } from "@agentclientprotocol/sdk"; import { zSessionNotification } from "@agentclientprotocol/sdk/dist/schema/zod.gen.js"; +import type { Model } from "@oh-my-pi/pi-ai"; +import { AcpAgent } from "../src/modes/acp/acp-agent"; import { buildToolCallStartUpdate, mapAgentSessionEventToAcpSessionUpdates, normalizeReplayToolArguments, } from "../src/modes/acp/acp-event-mapper"; -import type { AgentSessionEvent } from "../src/session/agent-session"; +import type { AgentSession, AgentSessionEvent } from "../src/session/agent-session"; +import { SessionManager } from "../src/session/session-manager"; import { expectAcpStructure, expectAcpStructureRejects } from "./helpers/acp-schema"; function makeAssistantMessage(text: string) { @@ -41,6 +46,55 @@ function expectAcpNotifications(updates: SessionNotification[]): void { } } +const TEST_MODEL: Model = { + id: "claude-sonnet-4-20250514", + name: "Claude Sonnet", + api: "anthropic-messages", + provider: "anthropic", + baseUrl: "https://example.invalid", + reasoning: true, + input: ["text", "image"], + cost: { input: 0, output: 0, cacheRead: 0, cacheWrite: 0 }, + contextWindow: 200_000, + maxTokens: 8_192, +}; + +class ReplayTestSession { + sessionManager: SessionManager; + sessionId: string; + model: Model | undefined = TEST_MODEL; + thinkingLevel: string | undefined; + customCommands: [] = []; + skills: [] = []; + extensionRunner = undefined; + settings = { get: (_key: string) => false }; + + constructor(cwd: string, sessionDir?: string) { + this.sessionManager = SessionManager.create(cwd, sessionDir); + this.sessionId = this.sessionManager.getSessionId(); + } + + getAvailableModels(): Model[] { + return [TEST_MODEL]; + } + + getAvailableThinkingLevels(): ReadonlyArray { + return []; + } + + getPlanModeState(): undefined { + return undefined; + } + + setClientBridge(_bridge: unknown): void {} + + subscribe(_listener: (event: AgentSessionEvent) => void): () => void { + return () => {}; + } + + async refreshMCPTools(_tools: unknown): Promise {} +} + describe("ACP event mapper", () => { it("attaches a stable messageId to live assistant chunks", () => { const assistantMessage = makeAssistantMessage("chunk"); @@ -331,6 +385,37 @@ describe("ACP event mapper", () => { expect(update.content).toContainEqual({ type: "terminal", terminalId: "term-1" }); }); + it("preserves command text when a command tool final update replaces content", () => { + const updates = mapAgentSessionEventToAcpSessionUpdates( + { + type: "tool_execution_end", + toolCallId: "tc-terminal-final-command", + toolName: "bash", + isError: false, + result: { + content: [{ type: "text", text: "done" }], + details: { terminalId: "term-1" }, + }, + } as AgentSessionEvent, + "session-1", + { + getToolArgs: toolCallId => + toolCallId === "tc-terminal-final-command" ? { command: "npm run check" } : undefined, + }, + ); + + expect(updates).toHaveLength(1); + expectAcpNotifications(updates); + const update = updates[0]!.update as { + sessionUpdate: string; + content?: Array<{ type: string; terminalId?: string; content?: { type: string; text?: string } }>; + }; + expect(update.sessionUpdate).toBe("tool_call_update"); + expect(update.content).toContainEqual({ type: "content", content: { type: "text", text: "$ npm run check" } }); + expect(update.content).toContainEqual({ type: "content", content: { type: "text", text: "done" } }); + expect(update.content).toContainEqual({ type: "terminal", terminalId: "term-1" }); + }); + it("keeps terminal content alongside readable error and message fields", () => { const errorUpdates = mapAgentSessionEventToAcpSessionUpdates( { @@ -494,6 +579,90 @@ describe("ACP event mapper", () => { } }); + it("replays assistant tool_use input through the ACP dispatcher without wrapping", async () => { + const root = await fs.promises.mkdtemp(path.join(os.tmpdir(), "omp-acp-replay-contract-")); + const cwd = path.join(root, "cwd"); + const sessionDir = path.join(root, "sessions"); + const initialSessionDir = path.join(root, "initial-session"); + const updates: SessionNotification[] = []; + const sessions: ReplayTestSession[] = []; + const abortController = new AbortController(); + try { + await fs.promises.mkdir(cwd, { recursive: true }); + const connection = { + sessionUpdate: async (notification: SessionNotification) => { + updates.push(notification); + }, + signal: abortController.signal, + closed: Promise.resolve(), + } as unknown as AgentSideConnection; + const agent = new AcpAgent( + connection, + async (sessionCwd: string) => { + const session = new ReplayTestSession(sessionCwd, sessionDir); + sessions.push(session); + return session as unknown as AgentSession; + }, + new ReplayTestSession(cwd, initialSessionDir) as unknown as AgentSession, + ); + const created = await agent.newSession({ cwd, mcpServers: [] }); + const session = sessions[0]!; + session.sessionManager.appendMessage({ + role: "assistant", + content: [ + { + type: "tool_use", + id: "toolu_replay_input", + name: "bash", + input: { command: "echo hi" }, + }, + ], + usage: { + input: 10, + output: 5, + cacheRead: 0, + cacheWrite: 0, + totalTokens: 15, + cost: { input: 0, output: 0, cacheRead: 0, cacheWrite: 0, total: 0 }, + }, + api: "anthropic-messages", + provider: "anthropic", + model: "claude-sonnet-4-20250514", + stopReason: "stop", + timestamp: Date.now(), + } as unknown as Parameters[0]); + session.sessionManager.appendMessage({ + role: "toolResult", + toolCallId: "toolu_replay_input", + toolName: "bash", + content: [{ type: "text", text: "done" }], + details: { terminalId: "term-replay" }, + isError: false, + timestamp: Date.now(), + }); + + updates.length = 0; + await agent.loadSession({ sessionId: created.sessionId, cwd, mcpServers: [] }); + + expectAcpNotifications(updates); + const toolCall = updates.find(update => update.update.sessionUpdate === "tool_call")?.update as + | { rawInput?: unknown; content?: unknown } + | undefined; + const finalUpdate = updates.find(update => update.update.sessionUpdate === "tool_call_update")?.update as + | { content?: unknown } + | undefined; + + expect(toolCall?.rawInput).toEqual({ command: "echo hi" }); + expect(toolCall?.rawInput).not.toEqual({ input: { command: "echo hi" } }); + expect(toolCall?.content).toEqual([{ type: "content", content: { type: "text", text: "$ echo hi" } }]); + expect(finalUpdate?.content).toContainEqual({ type: "content", content: { type: "text", text: "$ echo hi" } }); + expect(finalUpdate?.content).toContainEqual({ type: "content", content: { type: "text", text: "done" } }); + expect(finalUpdate?.content).toContainEqual({ type: "terminal", terminalId: "term-replay" }); + } finally { + abortController.abort(); + await fs.promises.rm(root, { recursive: true, force: true }); + } + }); it("builds replayed bash tool calls from JSON string arguments", () => { const replayArgs = normalizeReplayToolArguments(JSON.stringify({ command: "npm test", cwd: "/repo" })); const update = buildToolCallStartUpdate({ diff --git a/packages/coding-agent/test/agent-session-acp-permission.test.ts b/packages/coding-agent/test/agent-session-acp-permission.test.ts index f8d6042a1..7ee9f6d20 100644 --- a/packages/coding-agent/test/agent-session-acp-permission.test.ts +++ b/packages/coding-agent/test/agent-session-acp-permission.test.ts @@ -709,16 +709,3 @@ it("read tool: requestPermission is never called for non-gated tools", async () expect(permissionSpy).toHaveBeenCalledTimes(0); expect(readTool.executeCalls).toBe(1); }); - -// --------------------------------------------------------------------------- -// 5. No bridge → original tool object identity preserved (no wrapping) -// --------------------------------------------------------------------------- - -it("no bridge: original tool object is returned unchanged", async () => { - const bashTool = makeFakeTool("bash"); - session = await createSession([bashTool]); // no bridge - - await session.setActiveToolsByName(["bash"]); - const activeBash = session.agent.state.tools.find(t => t.name === "bash"); - expect(activeBash).toBe(bashTool); -}); diff --git a/packages/coding-agent/test/agent-session-concurrent.test.ts b/packages/coding-agent/test/agent-session-concurrent.test.ts index 3bea6ef5a..ed74ba4ad 100644 --- a/packages/coding-agent/test/agent-session-concurrent.test.ts +++ b/packages/coding-agent/test/agent-session-concurrent.test.ts @@ -10,6 +10,7 @@ import { Agent, AgentBusyError, type AgentTool } from "@oh-my-pi/pi-agent-core"; import { type AssistantMessage, getBundledModel, type Message, type ToolCall } from "@oh-my-pi/pi-ai"; import { createMockModel } from "@oh-my-pi/pi-ai/providers/mock"; import { AssistantMessageEventStream } from "@oh-my-pi/pi-ai/utils/event-stream"; +import { AsyncJobManager } from "@oh-my-pi/pi-coding-agent/async"; import type { Rule } from "@oh-my-pi/pi-coding-agent/capability/rule"; import { ModelRegistry } from "@oh-my-pi/pi-coding-agent/config/model-registry"; import { Settings } from "@oh-my-pi/pi-coding-agent/config/settings"; @@ -45,6 +46,7 @@ describe("AgentSession concurrent prompt guard", () => { fs.rmSync(tempDir, { recursive: true }); } vi.restoreAllMocks(); + AsyncJobManager.resetForTests(); }); async function createSession() { @@ -334,9 +336,9 @@ describe("AgentSession concurrent prompt guard", () => { const sessionManager = SessionManager.inMemory(); const settings = Settings.isolated(); - const authStorage = await AuthStorage.create(":memory:"); + const authStorage = await AuthStorage.create(path.join(tempDir, "testauth-acp-idle.db")); authStorages.push(authStorage); - const modelRegistry = new ModelRegistry(authStorage); + const modelRegistry = new ModelRegistry(authStorage, path.join(tempDir, "models-acp-idle.yml")); authStorage.setRuntimeApiKey("anthropic", "test-key"); session = new AgentSession({ @@ -383,6 +385,161 @@ describe("AgentSession concurrent prompt guard", () => { }), ).toBe(true); }); + + it("runs drained ACP async completions as owned follow-up turns despite deferred client turns", async () => { + const model = getBundledModel("anthropic", "claude-sonnet-4-5")!; + const mock = createMockModel({ handler: () => ({ content: ["Done"] }) }); + const agent = new Agent({ + getApiKey: () => "test-key", + initialState: { + model, + systemPrompt: ["Test"], + tools: [], + }, + convertToLlm, + streamFn: mock.stream, + }); + + const sessionManager = SessionManager.inMemory(); + const settings = Settings.isolated(); + const authStorage = await AuthStorage.create(path.join(tempDir, "testauth-acp-async.db")); + authStorages.push(authStorage); + const modelRegistry = new ModelRegistry(authStorage, path.join(tempDir, "models-acp-async.yml")); + authStorage.setRuntimeApiKey("anthropic", "test-key"); + + const ownerId = "acp-session-a"; + const deliveryGate = Promise.withResolvers(); + let deliveryStarted = false; + const asyncJobManager = new AsyncJobManager({ + maxRunningJobs: 2, + retentionMs: 1_000, + onJobComplete: async () => { + deliveryStarted = true; + await deliveryGate.promise; + await session.sendCustomMessage( + { + customType: "async-result", + content: "Background result", + display: true, + attribution: "agent", + }, + { deliverAs: "followUp", triggerTurn: true }, + ); + }, + }); + AsyncJobManager.setInstance(asyncJobManager); + + session = new AgentSession({ + agent, + sessionManager, + settings, + modelRegistry, + agentId: ownerId, + ownedAsyncJobManager: asyncJobManager, + }); + session.setClientBridge({ + capabilities: {}, + deferAgentInitiatedTurns: true, + }); + + await session.prompt("First message"); + expect(session.isStreaming).toBe(false); + const callsAfterFirstPrompt = mock.calls.length; + + try { + asyncJobManager.register("bash", "owned job", async () => "Background result", { + id: "owned-job", + ownerId, + }); + await waitFor(() => deliveryStarted); + + const drainedPromise = session.drainAsyncJobDeliveriesForAcp({ timeoutMs: 1_000 }); + await waitFor(() => asyncJobManager.getDeliveryState({ ownerId }).delivering); + deliveryGate.resolve(); + + await expect(drainedPromise).resolves.toBe(true); + await session.waitForIdle(); + + expect(mock.calls).toHaveLength(callsAfterFirstPrompt + 1); + expect( + mock.calls.at(-1)?.context.messages.some(message => { + if (typeof message.content === "string") { + return message.content.includes("Background result"); + } + + return message.content.some( + content => content.type === "text" && content.text.includes("Background result"), + ); + }), + ).toBe(true); + } finally { + deliveryGate.resolve(); + } + }); + + it("scopes ACP async job snapshots and drains to the owning session id", async () => { + const model = getBundledModel("anthropic", "claude-sonnet-4-5")!; + const authStorage = await AuthStorage.create(path.join(tempDir, "testauth-acp-scope.db")); + authStorages.push(authStorage); + const modelRegistry = new ModelRegistry(authStorage, path.join(tempDir, "models-acp-scope.yml")); + authStorage.setRuntimeApiKey("anthropic", "test-key"); + const settings = Settings.isolated(); + const deliveryGate = Promise.withResolvers(); + const delivered: string[] = []; + const started = new Set(); + const asyncJobManager = new AsyncJobManager({ + maxRunningJobs: 3, + retentionMs: 1_000, + onJobComplete: async jobId => { + started.add(jobId); + if (jobId === "job-a") { + await deliveryGate.promise; + } + delivered.push(jobId); + }, + }); + AsyncJobManager.setInstance(asyncJobManager); + + const agentA = new Agent({ + getApiKey: () => "test-key", + initialState: { model, systemPrompt: ["Test"], tools: [] }, + streamFn: createMockModel({ handler: () => ({ content: ["Done"] }) }).stream, + }); + const agentB = new Agent({ + getApiKey: () => "test-key", + initialState: { model, systemPrompt: ["Test"], tools: [] }, + streamFn: createMockModel({ handler: () => ({ content: ["Done"] }) }).stream, + }); + const sessionB = new AgentSession({ + agent: agentB, + sessionManager: SessionManager.inMemory(), + settings, + modelRegistry, + agentId: "acp-session-b", + }); + session = new AgentSession({ + agent: agentA, + sessionManager: SessionManager.inMemory(), + settings, + modelRegistry, + agentId: "acp-session-a", + ownedAsyncJobManager: asyncJobManager, + }); + + try { + asyncJobManager.register("bash", "A", async () => "A", { id: "job-a", ownerId: "acp-session-a" }); + await waitFor(() => started.has("job-a")); + asyncJobManager.register("bash", "B", async () => "B", { id: "job-b", ownerId: "acp-session-b" }); + await waitFor(() => asyncJobManager.getDeliveryState({ ownerId: "acp-session-b" }).queued > 0); + + expect(sessionB.getAsyncJobSnapshot()?.delivery.pendingJobIds).not.toContain("job-a"); + await expect(sessionB.drainAsyncJobDeliveriesForAcp({ timeoutMs: 1_000 })).resolves.toBe(true); + expect(delivered).toEqual(["job-b"]); + } finally { + deliveryGate.resolve(); + await sessionB.dispose(); + } + }); }); describe("AgentSession TTSR resume gate", () => { diff --git a/packages/coding-agent/test/async-job-manager.test.ts b/packages/coding-agent/test/async-job-manager.test.ts index b8a4a7c88..5caeef497 100644 --- a/packages/coding-agent/test/async-job-manager.test.ts +++ b/packages/coding-agent/test/async-job-manager.test.ts @@ -230,35 +230,104 @@ describe("AsyncJobManager", () => { }); test("scoped delivery drain returns once matching owner deliveries finish", async () => { - let firstOwnerAttempts = 0; - const secondOwnerCompletions: Array<{ jobId: string; text: string }> = []; + let mainJobId = ""; + let releaseMainDelivery = (): void => {}; + let notifyMainDeliveryStarted = (): void => {}; + const mainDeliveryStarted = new Promise(resolve => { + notifyMainDeliveryStarted = resolve; + }); + const mainDeliveryReleased = new Promise(resolve => { + releaseMainDelivery = resolve; + }); + const subagentCompletions: Array<{ jobId: string; text: string }> = []; const manager = new AsyncJobManager({ - onJobComplete: async (jobId, text, job) => { - if (job?.ownerId === "0-Main") { - firstOwnerAttempts++; - throw new Error("first owner delivery retry"); + retentionMs: 0, + onJobComplete: async (jobId, text) => { + if (jobId === mainJobId) { + notifyMainDeliveryStarted(); + await mainDeliveryReleased; + return; } - secondOwnerCompletions.push({ jobId, text }); + subagentCompletions.push({ jobId, text }); }, }); - manager.register("task", "main job", async () => "main result", { ownerId: "0-Main" }); + mainJobId = manager.register("task", "main job", async () => "main result", { ownerId: "0-Main" }); const targetJobId = manager.register("task", "subagent job", async () => "subagent result", { ownerId: "3-AuthLoader", }); await manager.waitForAll(); - const firstAttemptDeadline = Date.now() + 2_000; - while (firstOwnerAttempts === 0) { - if (Date.now() >= firstAttemptDeadline) throw new Error("Timed out waiting for first owner delivery attempt"); - await Bun.sleep(5); - } + await mainDeliveryStarted; + expect(manager.hasPendingDeliveries({ ownerId: "0-Main" })).toBe(true); const drained = await manager.drainDeliveries({ timeoutMs: 50, filter: { ownerId: "3-AuthLoader" } }); expect(drained).toBe(true); - expect(secondOwnerCompletions).toEqual([{ jobId: targetJobId, text: "subagent result" }]); + expect(subagentCompletions).toEqual([{ jobId: targetJobId, text: "subagent result" }]); expect(manager.hasPendingDeliveries({ ownerId: "3-AuthLoader" })).toBe(false); - expect(manager.hasPendingDeliveries({ ownerId: "0-Main" })).toBe(true); + + expect(manager.acknowledgeDeliveries([mainJobId])).toBe(0); + expect(manager.hasPendingDeliveries({ ownerId: "0-Main" })).toBe(false); + releaseMainDelivery(); + await Bun.sleep(0); + }); + + test("scoped delivery drain times out while a matching delivery callback is in flight", async () => { + let mainJobId = ""; + let targetJobId = ""; + let releaseMainDelivery = (): void => {}; + let notifyMainDeliveryStarted = (): void => {}; + let releaseTargetDelivery = (): void => {}; + let notifyTargetDeliveryStarted = (): void => {}; + const mainDeliveryStarted = new Promise(resolve => { + notifyMainDeliveryStarted = resolve; + }); + const mainDeliveryReleased = new Promise(resolve => { + releaseMainDelivery = resolve; + }); + const targetDeliveryStarted = new Promise(resolve => { + notifyTargetDeliveryStarted = resolve; + }); + const targetDeliveryReleased = new Promise(resolve => { + releaseTargetDelivery = resolve; + }); + const completions: string[] = []; + const manager = new AsyncJobManager({ + onJobComplete: async jobId => { + if (jobId === mainJobId) { + notifyMainDeliveryStarted(); + await mainDeliveryReleased; + return; + } + if (jobId === targetJobId) { + notifyTargetDeliveryStarted(); + await targetDeliveryReleased; + completions.push(jobId); + } + }, + }); + + mainJobId = manager.register("task", "main job", async () => "main result", { ownerId: "0-Main" }); + targetJobId = manager.register("task", "subagent job", async () => "subagent result", { + ownerId: "3-AuthLoader", + }); + await manager.waitForAll(); + await mainDeliveryStarted; + + const timedOut = await manager.drainDeliveries({ timeoutMs: 10, filter: { ownerId: "3-AuthLoader" } }); + await targetDeliveryStarted; + + expect(timedOut).toBe(false); + expect(manager.hasPendingDeliveries({ ownerId: "3-AuthLoader" })).toBe(true); + expect(completions).toEqual([]); + + releaseTargetDelivery(); + const drained = await manager.drainDeliveries({ timeoutMs: 200, filter: { ownerId: "3-AuthLoader" } }); + expect(drained).toBe(true); + expect(completions).toEqual([targetJobId]); + + releaseMainDelivery(); + expect(await manager.drainDeliveries({ timeoutMs: 200 })).toBe(true); }); test("cancelAll with ownerId only cancels matching jobs", async () => {