diff --git a/apps/server/src/orchestration-v2/Adapters/AcpAdapterV2.test.ts b/apps/server/src/orchestration-v2/Adapters/AcpAdapterV2.test.ts index 146924586763..39c9bc73d3f9 100644 --- a/apps/server/src/orchestration-v2/Adapters/AcpAdapterV2.test.ts +++ b/apps/server/src/orchestration-v2/Adapters/AcpAdapterV2.test.ts @@ -74,6 +74,7 @@ import { import type { ProviderContinuationRequest } from "../ProviderContinuationRequests.ts"; import { AcpProviderCapabilitiesV2, + acpPromptResponseTurnTokenUsage, acpProviderItemNativeId, acpScopedNativeId, acpCarryoverTerminalShouldClearContinuation, @@ -108,6 +109,95 @@ const testLayer = Layer.mergeAll(NodeServices.layer, IdAllocator.layer, serverCo const ACP_TEST_DRIVER = ProviderDriverKind.make("acp-test"); const decodeUnknownJson = Schema.decodeUnknownOption(Schema.fromJsonString(Schema.Unknown)); +describe("acpPromptResponseTurnTokenUsage", () => { + it("projects ACP prompt usage into complete main-agent turn usage", () => { + assert.deepEqual( + acpPromptResponseTurnTokenUsage( + { + inputTokens: 60, + outputTokens: 15, + totalTokens: 75, + cachedReadTokens: 40, + cachedWriteTokens: 1, + thoughtTokens: 4, + }, + { + inputTokens: 20, + outputTokens: 5, + totalTokens: 25, + cachedReadTokens: 10, + cachedWriteTokens: 0, + thoughtTokens: 1, + }, + true, + "completed", + ), + { + usageScope: "main_agent", + usageStatus: "complete", + hasSubagents: true, + inputTokens: 40, + outputTokens: 10, + cachedInputTokens: 30, + cacheCreationTokens: 1, + reasoningTokens: 3, + }, + ); + }); + + it("returns unavailable usage when ACP omits prompt usage", () => { + assert.deepEqual(acpPromptResponseTurnTokenUsage(undefined, null, false, "failed"), { + usageScope: "main_agent", + usageStatus: "unavailable", + hasSubagents: false, + }); + }); + + it("marks observed usage partial when the turn does not complete", () => { + assert.deepEqual( + acpPromptResponseTurnTokenUsage( + { inputTokens: 60, outputTokens: 15, totalTokens: 75 }, + { inputTokens: 40, outputTokens: 10, totalTokens: 50 }, + false, + "interrupted", + ), + { + usageScope: "main_agent", + usageStatus: "partial", + hasSubagents: false, + inputTokens: 20, + outputTokens: 5, + }, + ); + }); + + it("does not attribute resumed session totals or counter resets to one turn", () => { + const current = { inputTokens: 60, outputTokens: 15, totalTokens: 75 }; + const unavailable = { + usageScope: "main_agent", + usageStatus: "unavailable", + hasSubagents: false, + }; + assert.deepEqual( + acpPromptResponseTurnTokenUsage(current, null, false, "completed"), + unavailable, + ); + assert.deepEqual( + acpPromptResponseTurnTokenUsage( + current, + { + inputTokens: 90, + outputTokens: 15, + totalTokens: 105, + }, + false, + "completed", + ), + unavailable, + ); + }); +}); + describe("acpProjectedCommandExitCode", () => { const successOutput = { type: "Bash", exit_code: 0 }; const failedOutput = { type: "Bash", exit_code: 1 }; @@ -604,6 +694,338 @@ function makeTurnInput(input: { } describe("AcpAdapterV2", () => { + it.live( + "projects session usage deltas without attributing missing responses to later turns", + () => + Effect.gen(function* () { + const path = yield* Path.Path; + const instanceId = ProviderInstanceId.make("acp-test-prompt-usage"); + const threadId = ThreadId.make("thread-acp-prompt-usage"); + let promptOrdinal = 0; + const promptUsages = [ + { inputTokens: 60, outputTokens: 15, totalTokens: 75 }, + { inputTokens: 90, outputTokens: 22, totalTokens: 112 }, + undefined, + { inputTokens: 130, outputTokens: 30, totalTokens: 160 }, + { inputTokens: 150, outputTokens: 40, totalTokens: 190 }, + ]; + const adapter = makeAcpAdapterV2({ + crypto: yield* Crypto.Crypto, + instanceId, + fileSystem: yield* FileSystem.FileSystem, + idAllocator: yield* IdAllocator.IdAllocatorV2, + serverConfig: yield* ServerConfig.ServerConfig, + selfInvocation: yield* resolveSelfInvocation(), + flavor: { + driver: ACP_TEST_DRIVER, + capabilities: AcpProviderCapabilitiesV2, + makeRuntime: makeMockRuntime({ + childProcessSpawner: yield* ChildProcessSpawner.ChildProcessSpawner, + mockAgentPath: yield* path.fromFileUrl( + new URL("../../../scripts/acp-mock-agent.ts", import.meta.url), + ), + wrapRuntime: (runtime) => ({ + ...runtime, + prompt: () => { + promptOrdinal += 1; + return Effect.succeed({ + stopReason: "end_turn" as const, + usage: promptUsages[promptOrdinal - 1] ?? null, + }); + }, + }), + }), + }, + }); + const runtimePolicy = ProviderAdapterV2RuntimePolicy.make({ + runtimeMode: "full-access", + interactionMode: "default", + cwd: process.cwd(), + }); + const modelSelection = { instanceId, model: "default" } as const; + const runtime = yield* adapter.openSession({ + threadId, + providerSessionId: ProviderSessionId.make("provider-session-acp-prompt-usage"), + modelSelection, + runtimePolicy, + }); + const providerThread = yield* runtime.ensureThread({ + threadId, + modelSelection, + runtimePolicy, + }); + const now = yield* DateTime.now; + const expected = [ + { inputTokens: 60, outputTokens: 15 }, + { inputTokens: 30, outputTokens: 7 }, + null, + null, + { inputTokens: 20, outputTokens: 10 }, + ]; + for (const ordinal of [1, 2, 3, 4, 5]) { + const expectedUsage = expected[ordinal - 1]; + if (expectedUsage === undefined) return yield* Effect.die("Missing usage expectation"); + yield* runtime.startTurn( + makeTurnInput({ threadId, providerThread, instanceId, runtimePolicy, now, ordinal }), + ); + const events = Array.from( + yield* runtime.events.pipe( + Stream.takeUntil((event) => event.type === "turn.terminal"), + Stream.runCollect, + ), + ); + const completed = events.find( + (event) => + event.type === "provider_turn.updated" && event.providerTurn.status === "completed", + ); + assert.deepEqual( + completed?.type === "provider_turn.updated" + ? completed.providerTurn.turnTokenUsage + : null, + expectedUsage === null + ? { usageScope: "main_agent", usageStatus: "unavailable", hasSubagents: false } + : { + usageScope: "main_agent", + usageStatus: "complete", + hasSubagents: false, + ...expectedUsage, + }, + ); + } + }).pipe(Effect.provide(testLayer), Effect.scoped), + ); + + it.live("does not subtract pre-restart usage from a resumed runtime's counters", () => + Effect.gen(function* () { + const path = yield* Path.Path; + const instanceId = ProviderInstanceId.make("acp-test-usage-restart"); + const threadId = ThreadId.make("thread-acp-usage-restart"); + const wireSettled = yield* Deferred.make(); + const releaseCompletion = yield* Deferred.make(); + let promptOrdinal = 0; + const usages = [ + { inputTokens: 60, outputTokens: 15, totalTokens: 75 }, + { inputTokens: 90, outputTokens: 20, totalTokens: 110 }, + { inputTokens: 110, outputTokens: 25, totalTokens: 135 }, + ]; + const idAllocator = yield* IdAllocator.IdAllocatorV2; + const adapter = makeAcpAdapterV2({ + crypto: yield* Crypto.Crypto, + instanceId, + fileSystem: yield* FileSystem.FileSystem, + idAllocator, + serverConfig: yield* ServerConfig.ServerConfig, + selfInvocation: yield* resolveSelfInvocation(), + testHooks: { + afterPromptWireSettled: () => + promptOrdinal === 1 + ? Deferred.succeed(wireSettled, undefined).pipe( + Effect.andThen(Deferred.await(releaseCompletion)), + ) + : Effect.void, + }, + flavor: { + driver: ACP_TEST_DRIVER, + capabilities: AcpProviderCapabilitiesV2, + restartRuntimeAfterInterrupt: true, + terminateRuntimeProcessGroupOnInterrupt: true, + makeRuntime: makeMockRuntime({ + childProcessSpawner: yield* ChildProcessSpawner.ChildProcessSpawner, + mockAgentPath: yield* path.fromFileUrl( + new URL("../../../scripts/acp-mock-agent.ts", import.meta.url), + ), + ownDetachedProcessGroup: true, + wrapRuntime: (runtime) => ({ + ...runtime, + prompt: () => { + const usage = usages[promptOrdinal++]; + if (usage === undefined) return Effect.die("Unexpected prompt"); + return Effect.succeed({ stopReason: "end_turn" as const, usage }); + }, + }), + }), + }, + }); + const runtimePolicy = ProviderAdapterV2RuntimePolicy.make({ + runtimeMode: "full-access", + interactionMode: "default", + cwd: process.cwd(), + }); + const modelSelection = { instanceId, model: "default" } as const; + const runtime = yield* adapter.openSession({ + threadId, + providerSessionId: ProviderSessionId.make("provider-session-acp-usage-restart"), + modelSelection, + runtimePolicy, + }); + const providerThread = yield* runtime.ensureThread({ + threadId, + modelSelection, + runtimePolicy, + }); + const now = yield* DateTime.now; + yield* runtime.startTurn( + makeTurnInput({ threadId, providerThread, instanceId, runtimePolicy, now, ordinal: 1 }), + ); + yield* Deferred.await(wireSettled); + const providerTurnId = idAllocator.derive.providerTurn({ + driver: ACP_TEST_DRIVER, + nativeTurnId: acpScopedNativeId(instanceId, "mock-session-1:turn:1"), + }); + yield* runtime.interruptTurn({ providerThread, providerTurnId, requestRuntimeRestart: true }); + const interrupted = Array.from( + yield* runtime.events.pipe( + Stream.takeUntil((event) => event.type === "turn.terminal"), + Stream.runCollect, + ), + ).find( + (event) => + event.type === "provider_turn.updated" && event.providerTurn.status === "interrupted", + ); + assert.deepEqual( + interrupted?.type === "provider_turn.updated" + ? interrupted.providerTurn.turnTokenUsage + : null, + { + usageScope: "main_agent", + usageStatus: "partial", + hasSubagents: false, + inputTokens: 60, + outputTokens: 15, + }, + ); + yield* Deferred.succeed(releaseCompletion, undefined); + for (const [ordinal, expected] of [ + [2, { usageScope: "main_agent", usageStatus: "unavailable", hasSubagents: false }], + [ + 3, + { + usageScope: "main_agent", + usageStatus: "complete", + hasSubagents: false, + inputTokens: 20, + outputTokens: 5, + }, + ], + ] as const) { + yield* runtime.startTurn( + makeTurnInput({ threadId, providerThread, instanceId, runtimePolicy, now, ordinal }), + ); + const completed = Array.from( + yield* runtime.events.pipe( + Stream.takeUntil((event) => event.type === "turn.terminal"), + Stream.runCollect, + ), + ).find( + (event) => + event.type === "provider_turn.updated" && event.providerTurn.status === "completed", + ); + assert.deepEqual( + completed?.type === "provider_turn.updated" + ? completed.providerTurn.turnTokenUsage + : null, + expected, + ); + } + }).pipe(Effect.provide(testLayer), Effect.scoped), + ); + + it.live("retains returned usage when Stop finalizes before the prompt callback", () => + Effect.gen(function* () { + const path = yield* Path.Path; + const instanceId = ProviderInstanceId.make("acp-test-prompt-usage-stop"); + const threadId = ThreadId.make("thread-acp-prompt-usage-stop"); + const wireSettled = yield* Deferred.make(); + const releaseCompletion = yield* Deferred.make(); + const idAllocator = yield* IdAllocator.IdAllocatorV2; + const adapter = makeAcpAdapterV2({ + crypto: yield* Crypto.Crypto, + instanceId, + fileSystem: yield* FileSystem.FileSystem, + idAllocator, + serverConfig: yield* ServerConfig.ServerConfig, + selfInvocation: yield* resolveSelfInvocation(), + testHooks: { + afterPromptWireSettled: () => + Deferred.succeed(wireSettled, undefined).pipe( + Effect.andThen(Deferred.await(releaseCompletion)), + ), + }, + flavor: { + driver: ACP_TEST_DRIVER, + capabilities: AcpProviderCapabilitiesV2, + preserveRuntimeOnSettledInterrupt: true, + makeRuntime: makeMockRuntime({ + childProcessSpawner: yield* ChildProcessSpawner.ChildProcessSpawner, + mockAgentPath: yield* path.fromFileUrl( + new URL("../../../scripts/acp-mock-agent.ts", import.meta.url), + ), + wrapRuntime: (runtime) => ({ + ...runtime, + prompt: () => + Effect.succeed({ + stopReason: "end_turn" as const, + usage: { inputTokens: 60, outputTokens: 15, totalTokens: 75 }, + }), + }), + }), + }, + }); + const runtimePolicy = ProviderAdapterV2RuntimePolicy.make({ + runtimeMode: "full-access", + interactionMode: "default", + cwd: process.cwd(), + }); + const modelSelection = { instanceId, model: "default" } as const; + const runtime = yield* adapter.openSession({ + threadId, + providerSessionId: ProviderSessionId.make("provider-session-acp-prompt-usage-stop"), + modelSelection, + runtimePolicy, + }); + const providerThread = yield* runtime.ensureThread({ + threadId, + modelSelection, + runtimePolicy, + }); + yield* runtime.startTurn( + makeTurnInput({ + threadId, + providerThread, + instanceId, + runtimePolicy, + now: yield* DateTime.now, + }), + ); + yield* Deferred.await(wireSettled); + const providerTurnId = idAllocator.derive.providerTurn({ + driver: ACP_TEST_DRIVER, + nativeTurnId: acpScopedNativeId(instanceId, "mock-session-1:turn:1"), + }); + yield* runtime.interruptTurn({ providerThread, providerTurnId }); + const completed = Array.from( + yield* runtime.events.pipe( + Stream.takeUntil((event) => event.type === "turn.terminal"), + Stream.runCollect, + ), + ).find( + (event) => + event.type === "provider_turn.updated" && event.providerTurn.status === "interrupted", + ); + assert.deepEqual( + completed?.type === "provider_turn.updated" ? completed.providerTurn.turnTokenUsage : null, + { + usageScope: "main_agent", + usageStatus: "partial", + hasSubagents: false, + inputTokens: 60, + outputTokens: 15, + }, + ); + yield* Deferred.succeed(releaseCompletion, undefined); + }).pipe(Effect.provide(testLayer), Effect.scoped), + ); + it.live.each(["failed", "recovered", "completed", "cancelled"] as const)( "projects Mistral retry notices and their %s outcome", (outcome) => diff --git a/apps/server/src/orchestration-v2/Adapters/AcpAdapterV2.ts b/apps/server/src/orchestration-v2/Adapters/AcpAdapterV2.ts index bf01d9cb2d7d..0f7ddd44359a 100644 --- a/apps/server/src/orchestration-v2/Adapters/AcpAdapterV2.ts +++ b/apps/server/src/orchestration-v2/Adapters/AcpAdapterV2.ts @@ -30,6 +30,7 @@ import { type RuntimeRequestId, type ThreadTokenUsageSnapshot, type ThreadId, + type TurnTokenUsage, } from "@t3tools/contracts"; import { modelSelectionsEqual } from "@t3tools/shared/model"; import { type SelfInvocation, selfInvocationArgs } from "@t3tools/shared/nodeRuntime"; @@ -210,6 +211,39 @@ export interface AcpAdapterV2ExtensionContext { readonly lastProposedPlanMarkdown: Effect.Effect; } +export function acpPromptResponseTurnTokenUsage( + usage: EffectAcpSchema.Usage | null | undefined, + previous: EffectAcpSchema.Usage | null, + hasSubagents: boolean, + terminalStatus: OrchestrationV2ProviderTurn["status"], +): TurnTokenUsage { + // ACP reports session-wide counters. A resumed session, missing response, or + // counter reset leaves no trustworthy baseline for this individual turn. + if ( + usage == null || + previous === null || + usage.inputTokens < previous.inputTokens || + usage.outputTokens < previous.outputTokens + ) { + return { usageScope: "main_agent", usageStatus: "unavailable", hasSubagents }; + } + const delta = (current: number | null | undefined, before: number | null | undefined) => + current == null || before == null || current < before ? undefined : current - before; + const cachedInputTokens = delta(usage.cachedReadTokens, previous.cachedReadTokens); + const cacheCreationTokens = delta(usage.cachedWriteTokens, previous.cachedWriteTokens); + const reasoningTokens = delta(usage.thoughtTokens, previous.thoughtTokens); + return { + usageScope: "main_agent", + usageStatus: terminalStatus === "completed" ? "complete" : "partial", + hasSubagents, + inputTokens: usage.inputTokens - previous.inputTokens, + outputTokens: usage.outputTokens - previous.outputTokens, + ...(cachedInputTokens === undefined ? {} : { cachedInputTokens }), + ...(cacheCreationTokens === undefined ? {} : { cacheCreationTokens }), + ...(reasoningTokens === undefined ? {} : { reasoningTokens }), + }; +} + export interface AcpAdapterV2Flavor { /** Interprets provider-specific prompt errors before they cross into orchestration. */ readonly promptFailure?: (cause: unknown) => OrchestrationV2ProviderFailure; @@ -511,6 +545,7 @@ export interface AcpAdapterV2Options { * by exactly that on this receipt. */ readonly onDeferredFinalizeScheduled?: (debounce: Duration.Input) => Effect.Effect; + readonly afterPromptWireSettled?: () => Effect.Effect; readonly afterPromptSettledWithBackgroundWork?: () => Effect.Effect; readonly afterNativeResponseTransportClosed?: () => Effect.Effect; readonly afterHardTeardownTransportDrained?: () => Effect.Effect; @@ -1167,6 +1202,8 @@ interface ActiveAcpTurn { * complete this; settled-soft classification ORs it with `promptSettled`. */ readonly promptWireSettled: Deferred.Deferred; + promptUsage: EffectAcpSchema.Usage | undefined; + promptUsageBaseline: EffectAcpSchema.Usage | null; backgroundFinalizeGeneration: number; } @@ -1600,6 +1637,9 @@ export function makeAcpAdapterV2( const contextUsageBySessionId = yield* Ref.make( new Map(), ); + const promptUsageBySessionId = yield* Ref.make( + new Map(), + ); const nativeMetadataBySessionId = yield* Ref.make( new Map(), ); @@ -6071,6 +6111,19 @@ export function makeAcpAdapterV2( }); yield* Ref.set(activeSessionId, started.sessionId); yield* Ref.set(activeSessionSetup, started); + if (started.sessionId !== input.initialNativeThreadId) { + // A new ACP session starts at zero; a loaded one has unknown history. + yield* Ref.update(promptUsageBySessionId, (current) => + new Map(current).set(started.sessionId, { + inputTokens: 0, + outputTokens: 0, + totalTokens: 0, + cachedReadTokens: 0, + cachedWriteTokens: 0, + thoughtTokens: 0, + }), + ); + } rememberTerminalEnvironment(started.sessionId, input.threadId); const capabilities = negotiatedCapabilities(flavor.capabilities, started); const canLoadSession = started.initializeResult.agentCapabilities?.loadSession === true; @@ -6102,6 +6155,10 @@ export function makeAcpAdapterV2( driver, detail: `ACP driver cannot load or resume session ${sessionId}`, }); + // Loading a session does not establish continuity with its last prompt here. + yield* Ref.update(promptUsageBySessionId, (current) => + new Map(current).set(activated.sessionId, null), + ); rememberTerminalEnvironment(activated.sessionId, threadId); return activated; }); @@ -6344,6 +6401,12 @@ export function makeAcpAdapterV2( status, startedAt: context.startedAt, completedAt, + turnTokenUsage: acpPromptResponseTurnTokenUsage( + context.promptUsage, + context.promptUsageBaseline, + context.subagents.size > 0, + status, + ), }); const terminalizeOpenRunOwnedItems = Effect.fnUntraced(function* ( @@ -6408,6 +6471,12 @@ export function makeAcpAdapterV2( const settledStatus = context.interrupted ? "interrupted" : status; context.finalizedStatus = settledStatus; context.finalized = true; + if (context.promptUsage === undefined) { + // A cancelled/failed prompt with no response may have spent tokens. + yield* Ref.update(promptUsageBySessionId, (current) => + new Map(current).set(context.nativeThreadId, null), + ); + } const directStopQuarantine = yield* Ref.get(stoppedRunQuarantine); if (flavor.subagentsIdleOnTurnCompletion === true) { for (const subagent of context.subagents.values()) { @@ -6689,6 +6758,7 @@ export function makeAcpAdapterV2( const restartRequired = yield* Ref.get(runtimeRestartRequired); if (!restartRequired) return false; yield* restartAcpRuntime(threadId); + yield* Ref.set(promptUsageBySessionId, new Map()); yield* Ref.set(runtimeRestartRequired, false); yield* Ref.set(activeSessionId, null); yield* Ref.set(activeSessionSetup, null); @@ -6857,6 +6927,8 @@ export function makeAcpAdapterV2( promptSettled: false, promptSettledStatus: null, promptWireSettled, + promptUsage: undefined, + promptUsageBaseline: null, backgroundFinalizeGeneration: 0, }; const carryover = yield* Ref.getAndSet(carryoverSubagents, null); @@ -6999,9 +7071,22 @@ export function makeAcpAdapterV2( // Wire settlement precedes the completion callback's permit request so // settled-soft classification can observe the native return even when // the completion fiber has not yet set promptSettled under the permit. - Effect.tap(() => - Deferred.succeed(context.promptWireSettled, undefined).pipe(Effect.asVoid), + Effect.tap((result) => + Ref.modify(promptUsageBySessionId, (current) => { + if (context.finalized) return [undefined, current]; + context.promptUsage = result.usage ?? undefined; + context.promptUsageBaseline = current.get(requestedSessionId) ?? null; + return [ + undefined, + new Map(current).set(requestedSessionId, result.usage ?? null), + ]; + }).pipe( + Effect.andThen( + Deferred.succeed(context.promptWireSettled, undefined).pipe(Effect.asVoid), + ), + ), ), + Effect.tap(() => options.testHooks?.afterPromptWireSettled?.() ?? Effect.void), Effect.flatMap((result) => runRuntimeCallbackAtGeneration( promptGeneration, @@ -7669,6 +7754,9 @@ export function makeAcpAdapterV2( sessionId, acpMcpActivation(snapshotInput.providerThread.appThreadId, self), ); + yield* Ref.update(promptUsageBySessionId, (current) => + new Map(current).set(activated.sessionId, null), + ); rememberTerminalEnvironment( activated.sessionId, snapshotInput.providerThread.appThreadId, @@ -7741,6 +7829,22 @@ export function makeAcpAdapterV2( yield* Ref.set(activeSessionSetup, candidate); yield* Ref.set(activeSelection, null); yield* Ref.set(activeInteractionMode, null); + yield* Ref.set( + promptUsageBySessionId, + new Map([ + [ + candidate.sessionId, + { + inputTokens: 0, + outputTokens: 0, + totalTokens: 0, + cachedReadTokens: 0, + cachedWriteTokens: 0, + thoughtTokens: 0, + }, + ], + ]), + ); yield* Ref.set(promptInstructionStates, new Map()); yield* Ref.set(providerTurns, new Map()); yield* Ref.set(snapshot, { @@ -7832,6 +7936,9 @@ export function makeAcpAdapterV2( sourceSessionId, acpMcpActivation(forkInput.targetThreadId, self), ); + yield* Ref.update(promptUsageBySessionId, (current) => + new Map(current).set(forked.sessionId, null), + ); rememberTerminalEnvironment(forked.sessionId, forkInput.targetThreadId); yield* Ref.set(activeSessionId, forked.sessionId); yield* Ref.set(activeSessionSetup, forked);