diff --git a/apps/server/src/provider/Layers/ClaudeAdapter.test.ts b/apps/server/src/provider/Layers/ClaudeAdapter.test.ts index 2f0efeac5f53..8362186ebf2e 100644 --- a/apps/server/src/provider/Layers/ClaudeAdapter.test.ts +++ b/apps/server/src/provider/Layers/ClaudeAdapter.test.ts @@ -2510,7 +2510,7 @@ describe("ClaudeAdapterLive", () => { return Effect.gen(function* () { const adapter = yield* ClaudeAdapter; - const runtimeEventsFiber = yield* Stream.take(adapter.streamEvents, 6).pipe( + const runtimeEventsFiber = yield* Stream.take(adapter.streamEvents, 5).pipe( Stream.runCollect, Effect.forkChild, ); @@ -2768,12 +2768,12 @@ describe("ClaudeAdapterLive", () => { ); }); - it.effect("emits thread token usage updates from Claude task progress", () => { + it.effect("ignores task progress usage before the parent has any usage of its own", () => { const harness = makeHarness(); return Effect.gen(function* () { const adapter = yield* ClaudeAdapter; - const runtimeEventsFiber = yield* Stream.take(adapter.streamEvents, 6).pipe( + const runtimeEventsFiber = yield* Stream.take(adapter.streamEvents, 5).pipe( Stream.runCollect, Effect.forkChild, ); @@ -2801,20 +2801,207 @@ describe("ClaudeAdapterLive", () => { const runtimeEvents = Array.from(yield* Fiber.join(runtimeEventsFiber)); const usageEvent = runtimeEvents.find((event) => event.type === "thread.token-usage.updated"); const progressEvent = runtimeEvents.find((event) => event.type === "task.progress"); - assert.equal(usageEvent?.type, "thread.token-usage.updated"); - if (usageEvent?.type === "thread.token-usage.updated") { - assert.deepEqual(usageEvent.payload, { - usage: { - usedTokens: 321, - lastUsedTokens: 321, - toolUses: 2, - durationMs: 654, - }, - }); - } + // The subagent's 321 tokens are spent in its own window. With no parent + // usage recorded yet there is nothing to attach a running total to, so + // the parent meter stays silent rather than adopting the child's count. + assert.equal(usageEvent, undefined); assert.equal(progressEvent?.type, "task.progress"); - if (usageEvent && progressEvent) { - assert.notStrictEqual(usageEvent.eventId, progressEvent.eventId); + }).pipe( + Effect.provideService(Random.Random, makeDeterministicRandomService()), + Effect.provide(harness.layer), + ); + }); + + it.effect("holds the running total until it exceeds the parent's own usage", () => { + const harness = makeHarness(); + return Effect.gen(function* () { + const adapter = yield* ClaudeAdapter; + const runtimeEvents: Array = []; + const runtimeEventsFiber = yield* Stream.runForEach(adapter.streamEvents, (event) => + Effect.sync(() => runtimeEvents.push(event)), + ).pipe(Effect.forkChild); + + yield* adapter.startSession({ + threadId: THREAD_ID, + provider: ProviderDriverKind.make("claudeAgent"), + runtimeMode: "full-access", + }); + yield* adapter.sendTurn({ threadId: THREAD_ID, input: "go", attachments: [] }); + + harness.query.emit({ + type: "stream_event", + event: { + type: "message_delta", + delta: { stop_reason: null, stop_sequence: null }, + usage: { input_tokens: 3_000, output_tokens: 0 }, + }, + parent_tool_use_id: null, + session_id: "sdk-session-total-floor", + uuid: "total-floor-delta", + } as unknown as SDKMessage); + + // A cumulative figure below the parent's active usage would read as + // "3,000 of 2,000 total". The child's small tick stays off the meter. + harness.query.emit({ + type: "system", + subtype: "task_progress", + task_id: "task-total-floor", + description: "Child barely started", + usage: { total_tokens: 2_000 }, + session_id: "sdk-session-total-floor", + uuid: "total-floor-small", + } as unknown as SDKMessage); + + harness.query.emit({ + type: "system", + subtype: "task_progress", + task_id: "task-total-floor", + description: "Child past the parent", + usage: { total_tokens: 5_000 }, + session_id: "sdk-session-total-floor", + uuid: "total-floor-large", + } as unknown as SDKMessage); + + yield* Effect.yieldNow; + yield* Fiber.interrupt(runtimeEventsFiber); + + const usageEvents = runtimeEvents.filter( + (event) => event.type === "thread.token-usage.updated", + ); + assert.equal(usageEvents.length, 2); + const [afterDelta, afterLargeTick] = usageEvents; + if (afterDelta?.type === "thread.token-usage.updated") { + assert.equal(afterDelta.payload.usage.usedTokens, 3_000); + assert.equal(afterDelta.payload.usage.totalProcessedTokens, undefined); + } + if (afterLargeTick?.type === "thread.token-usage.updated") { + assert.equal(afterLargeTick.payload.usage.usedTokens, 3_000); + assert.equal(afterLargeTick.payload.usage.totalProcessedTokens, 5_000); + } + }).pipe( + Effect.provideService(Random.Random, makeDeterministicRandomService()), + Effect.provide(harness.layer), + ); + }); + + it.effect("keeps subagent tokens out of the parent context meter (#5942)", () => { + const harness = makeHarness(); + return Effect.gen(function* () { + const adapter = yield* ClaudeAdapter; + const runtimeEvents: Array = []; + const runtimeEventsFiber = yield* Stream.runForEach(adapter.streamEvents, (event) => + Effect.sync(() => runtimeEvents.push(event)), + ).pipe(Effect.forkChild); + + yield* adapter.startSession({ + threadId: THREAD_ID, + provider: ProviderDriverKind.make("claudeAgent"), + runtimeMode: "full-access", + }); + + yield* adapter.sendTurn({ + threadId: THREAD_ID, + input: "delegate this", + attachments: [], + }); + + // The parent finishes a small turn of its own. + harness.query.emit({ + type: "result", + subtype: "success", + is_error: false, + duration_ms: 10, + duration_api_ms: 8, + num_turns: 1, + result: "done", + stop_reason: "end_turn", + session_id: "sdk-session-child-usage", + usage: { input_tokens: 4_000, output_tokens: 200 }, + modelUsage: { "claude-opus-4-6": { contextWindow: 1_000_000 } }, + } as unknown as SDKMessage); + + // A background subagent then burns far more than the parent ever has. + harness.query.emit({ + type: "system", + subtype: "task_progress", + task_id: "task-child-1", + description: "Background agent doing the heavy work", + usage: { total_tokens: 900_000, tool_uses: 40, duration_ms: 1_000 }, + session_id: "sdk-session-child-usage", + uuid: "task-child-progress-1", + } as unknown as SDKMessage); + harness.query.finish(); + + yield* Effect.yieldNow; + yield* Fiber.interrupt(runtimeEventsFiber); + + const usageEvents = runtimeEvents.filter( + (event) => event.type === "thread.token-usage.updated", + ); + const latest = usageEvents.at(-1); + assert.equal(latest?.type, "thread.token-usage.updated"); + if (latest?.type === "thread.token-usage.updated") { + // The parent's own 4,200 stands; the child's 900,000 only advances + // the running total, never usedTokens. + assert.equal(latest.payload.usage.usedTokens, 4_200); + assert.equal(latest.payload.usage.totalProcessedTokens, 900_000); + } + }).pipe( + Effect.provideService(Random.Random, makeDeterministicRandomService()), + Effect.provide(harness.layer), + ); + }); + + it.effect("measures the meter against the session model's window, not a subagent's", () => { + const harness = makeHarness(); + return Effect.gen(function* () { + const adapter = yield* ClaudeAdapter; + + const runtimeEvents: Array = []; + const runtimeEventsFiber = yield* Stream.runForEach(adapter.streamEvents, (event) => + Effect.sync(() => runtimeEvents.push(event)), + ).pipe(Effect.forkChild); + + const session = yield* adapter.startSession({ + threadId: THREAD_ID, + provider: ProviderDriverKind.make("claudeAgent"), + runtimeMode: "full-access", + modelSelection: createModelSelection( + ProviderInstanceId.make("claudeAgent"), + "claude-sonnet-5", + ), + }); + + yield* adapter.sendTurn({ + threadId: session.threadId, + input: "summarize", + attachments: [], + }); + + // A 1M subagent ran alongside a 200k session model. + harness.query.emit({ + type: "result", + subtype: "success", + is_error: false, + session_id: "sdk-session-window-scope", + uuid: "result-window-scope", + usage: { input_tokens: 50_000, output_tokens: 1_000 }, + modelUsage: { + "claude-sonnet-5": { contextWindow: 200_000 }, + "claude-opus-5[1m]": { contextWindow: 1_000_000 }, + }, + } as unknown as SDKMessage); + + yield* Effect.yieldNow; + yield* Fiber.interrupt(runtimeEventsFiber); + const usageEvents = runtimeEvents.filter( + (event) => event.type === "thread.token-usage.updated", + ); + const latest = usageEvents.at(-1); + assert.equal(latest?.type, "thread.token-usage.updated"); + if (latest?.type === "thread.token-usage.updated") { + // The maximum over modelUsage would report the child's 1,000,000. + assert.equal(latest.payload.usage.maxTokens, 200_000); } }).pipe( Effect.provideService(Random.Random, makeDeterministicRandomService()), @@ -2822,6 +3009,424 @@ describe("ClaudeAdapterLive", () => { ); }); + it.effect("uses the init model's window when no model was explicitly selected", () => { + const harness = makeHarness(); + return Effect.gen(function* () { + const adapter = yield* ClaudeAdapter; + const runtimeEvents: Array = []; + const runtimeEventsFiber = yield* Stream.runForEach(adapter.streamEvents, (event) => + Effect.sync(() => runtimeEvents.push(event)), + ).pipe(Effect.forkChild); + + yield* adapter.startSession({ + threadId: THREAD_ID, + provider: ProviderDriverKind.make("claudeAgent"), + runtimeMode: "full-access", + }); + + yield* adapter.sendTurn({ + threadId: THREAD_ID, + input: "summarize", + attachments: [], + }); + + // No selection was made, so the session runs Claude Code's default and + // init is the first place its name appears. + harness.query.emit({ + type: "system", + subtype: "init", + model: "claude-sonnet-5", + session_id: "sdk-session-init-window", + uuid: "init-window", + } as unknown as SDKMessage); + + harness.query.emit({ + type: "result", + subtype: "success", + is_error: false, + duration_ms: 10, + duration_api_ms: 8, + num_turns: 1, + result: "done", + stop_reason: "end_turn", + session_id: "sdk-session-init-window", + usage: { input_tokens: 50_000, output_tokens: 1_000 }, + modelUsage: { + "claude-sonnet-5": { contextWindow: 200_000 }, + "claude-opus-5[1m]": { contextWindow: 1_000_000 }, + }, + } as unknown as SDKMessage); + harness.query.finish(); + + yield* Effect.yieldNow; + yield* Fiber.interrupt(runtimeEventsFiber); + + const usageEvents = runtimeEvents.filter( + (event) => event.type === "thread.token-usage.updated", + ); + const latest = usageEvents.at(-1); + assert.equal(latest?.type, "thread.token-usage.updated"); + if (latest?.type === "thread.token-usage.updated") { + assert.equal(latest.payload.usage.maxTokens, 200_000); + } + }).pipe( + Effect.provideService(Random.Random, makeDeterministicRandomService()), + Effect.provide(harness.layer), + ); + }); + + it.effect("carries a subagent's running total into the parent's next snapshot", () => { + const harness = makeHarness(); + return Effect.gen(function* () { + const adapter = yield* ClaudeAdapter; + const runtimeEvents: Array = []; + const runtimeEventsFiber = yield* Stream.runForEach(adapter.streamEvents, (event) => + Effect.sync(() => runtimeEvents.push(event)), + ).pipe(Effect.forkChild); + + yield* adapter.startSession({ + threadId: THREAD_ID, + provider: ProviderDriverKind.make("claudeAgent"), + runtimeMode: "full-access", + }); + + yield* adapter.sendTurn({ + threadId: THREAD_ID, + input: "delegate this", + attachments: [], + }); + + harness.query.emit({ + type: "stream_event", + event: { + type: "message_delta", + delta: { stop_reason: null, stop_sequence: null }, + usage: { input_tokens: 3_000, output_tokens: 0 }, + }, + parent_tool_use_id: null, + session_id: "sdk-session-running-total", + uuid: "parent-running-total", + } as unknown as SDKMessage); + + // The child's spend is real work the thread paid for, so it belongs in + // the running total even though it never touches the parent's own + // used count. + harness.query.emit({ + type: "system", + subtype: "task_progress", + task_id: "task-running-total", + description: "Child doing the heavy work", + usage: { total_tokens: 480_000 }, + session_id: "sdk-session-running-total", + uuid: "task-running-total-progress", + } as unknown as SDKMessage); + + yield* Effect.yieldNow; + yield* Fiber.interrupt(runtimeEventsFiber); + + const usageEvents = runtimeEvents.filter( + (event) => event.type === "thread.token-usage.updated", + ); + const latest = usageEvents.at(-1); + assert.equal(latest?.type, "thread.token-usage.updated"); + if (latest?.type === "thread.token-usage.updated") { + assert.equal(latest.payload.usage.totalProcessedTokens, 480_000); + assert.equal(latest.payload.usage.usedTokens, 3_000); + } + }).pipe( + Effect.provideService(Random.Random, makeDeterministicRandomService()), + Effect.provide(harness.layer), + ); + }); + + it.effect("keeps a subagent's running total when the parent turn completes", () => { + const harness = makeHarness(); + return Effect.gen(function* () { + const adapter = yield* ClaudeAdapter; + const runtimeEvents: Array = []; + const runtimeEventsFiber = yield* Stream.runForEach(adapter.streamEvents, (event) => + Effect.sync(() => runtimeEvents.push(event)), + ).pipe(Effect.forkChild); + + yield* adapter.startSession({ + threadId: THREAD_ID, + provider: ProviderDriverKind.make("claudeAgent"), + runtimeMode: "full-access", + }); + yield* adapter.sendTurn({ threadId: THREAD_ID, input: "go", attachments: [] }); + + harness.query.emit({ + type: "stream_event", + event: { + type: "message_delta", + delta: { stop_reason: null, stop_sequence: null }, + usage: { input_tokens: 3_000, output_tokens: 0 }, + }, + parent_tool_use_id: null, + session_id: "sdk-session-running-total-result", + uuid: "running-total-result-delta", + } as unknown as SDKMessage); + + harness.query.emit({ + type: "system", + subtype: "task_progress", + task_id: "task-running-total-result", + description: "Child doing the heavy work", + usage: { total_tokens: 900_000 }, + session_id: "sdk-session-running-total-result", + uuid: "running-total-result-task", + } as unknown as SDKMessage); + + harness.query.emit({ + type: "result", + subtype: "success", + is_error: false, + duration_ms: 10, + duration_api_ms: 8, + num_turns: 1, + result: "done", + stop_reason: "end_turn", + session_id: "sdk-session-running-total-result", + usage: { input_tokens: 4_000, output_tokens: 200 }, + modelUsage: { "claude-opus-4-6": { contextWindow: 200_000 } }, + } as unknown as SDKMessage); + harness.query.finish(); + + yield* Effect.yieldNow; + yield* Fiber.interrupt(runtimeEventsFiber); + const ev = runtimeEvents.filter((e) => e.type === "thread.token-usage.updated"); + const latest = ev.at(-1); + assert.equal(latest?.type, "thread.token-usage.updated"); + if (latest?.type === "thread.token-usage.updated") { + // The parent's own result is smaller than the running total the child + // already raised. Used tokens follow the parent; total processed is + // cumulative thread work and must not regress. + assert.equal(latest.payload.usage.usedTokens, 4_200); + assert.equal(latest.payload.usage.totalProcessedTokens, 900_000); + } + }).pipe( + Effect.provideService(Random.Random, makeDeterministicRandomService()), + Effect.provide(harness.layer), + ); + }); + + it.effect("follows a persistent refusal fallback to the model that ran", () => { + const harness = makeHarness(); + return Effect.gen(function* () { + const adapter = yield* ClaudeAdapter; + const runtimeEvents: Array = []; + const runtimeEventsFiber = yield* Stream.runForEach(adapter.streamEvents, (event) => + Effect.sync(() => runtimeEvents.push(event)), + ).pipe(Effect.forkChild); + + yield* adapter.startSession({ + threadId: THREAD_ID, + provider: ProviderDriverKind.make("claudeAgent"), + runtimeMode: "full-access", + }); + + yield* adapter.sendTurn({ + threadId: THREAD_ID, + input: "summarize", + attachments: [], + }); + + harness.query.emit({ + type: "system", + subtype: "init", + model: "claude-sonnet-5", + session_id: "sdk-session-refusal", + uuid: "refusal-init", + } as unknown as SDKMessage); + + // A refusal retry swaps the model for the rest of the session, so the + // window has to be measured against the model that actually ran. + harness.query.emit({ + type: "system", + subtype: "model_refusal_fallback", + trigger: "refusal", + direction: "retry", + original_model: "claude-sonnet-5", + fallback_model: "claude-opus-5[1m]", + request_id: null, + session_id: "sdk-session-refusal", + uuid: "refusal-swap", + } as unknown as SDKMessage); + + harness.query.emit({ + type: "result", + subtype: "success", + is_error: false, + duration_ms: 10, + duration_api_ms: 8, + num_turns: 1, + result: "done", + stop_reason: "end_turn", + session_id: "sdk-session-refusal", + usage: { input_tokens: 50_000, output_tokens: 1_000 }, + modelUsage: { + "claude-sonnet-5": { contextWindow: 200_000 }, + "claude-opus-5[1m]": { contextWindow: 1_000_000 }, + }, + } as unknown as SDKMessage); + harness.query.finish(); + + yield* Effect.yieldNow; + yield* Fiber.interrupt(runtimeEventsFiber); + + const usageEvents = runtimeEvents.filter( + (event) => event.type === "thread.token-usage.updated", + ); + const latest = usageEvents.at(-1); + assert.equal(latest?.type, "thread.token-usage.updated"); + if (latest?.type === "thread.token-usage.updated") { + assert.equal(latest.payload.usage.maxTokens, 1_000_000); + } + }).pipe( + Effect.provideService(Random.Random, makeDeterministicRandomService()), + Effect.provide(harness.layer), + ); + }); + + it.effect("measures the window against the model a mid-thread switch selected", () => { + const harness = makeHarness(); + return Effect.gen(function* () { + const adapter = yield* ClaudeAdapter; + const runtimeEvents: Array = []; + const runtimeEventsFiber = yield* Stream.runForEach(adapter.streamEvents, (event) => + Effect.sync(() => runtimeEvents.push(event)), + ).pipe(Effect.forkChild); + + yield* adapter.startSession({ + threadId: THREAD_ID, + provider: ProviderDriverKind.make("claudeAgent"), + runtimeMode: "full-access", + }); + + harness.query.emit({ + type: "system", + subtype: "init", + model: "claude-opus-4-6[1m]", + session_id: "sdk-session-switch", + uuid: "switch-init", + } as unknown as SDKMessage); + yield* Effect.yieldNow; + + // The user switches to a 200k model partway through the thread. + yield* adapter.sendTurn({ + threadId: THREAD_ID, + input: "switch", + modelSelection: { + instanceId: ProviderInstanceId.make("claudeAgent"), + model: "claude-sonnet-5", + }, + attachments: [], + }); + + harness.query.emit({ + type: "result", + subtype: "success", + is_error: false, + duration_ms: 10, + duration_api_ms: 8, + num_turns: 1, + result: "done", + stop_reason: "end_turn", + session_id: "sdk-session-switch", + usage: { input_tokens: 20_000, output_tokens: 500 }, + modelUsage: { + "claude-opus-4-6[1m]": { contextWindow: 1_000_000 }, + "claude-sonnet-5": { contextWindow: 200_000 }, + }, + } as unknown as SDKMessage); + harness.query.finish(); + + yield* Effect.yieldNow; + yield* Fiber.interrupt(runtimeEventsFiber); + + const usageEvents = runtimeEvents.filter( + (event) => event.type === "thread.token-usage.updated", + ); + const latest = usageEvents.at(-1); + assert.equal(latest?.type, "thread.token-usage.updated"); + if (latest?.type === "thread.token-usage.updated") { + assert.equal(latest.payload.usage.maxTokens, 200_000); + } + }).pipe( + Effect.provideService(Random.Random, makeDeterministicRandomService()), + Effect.provide(harness.layer), + ); + }); + + it.effect("does not re-send a refused model on the next turn with the same selection", () => { + const harness = makeHarness(); + return Effect.gen(function* () { + const adapter = yield* ClaudeAdapter; + const runtimeEventsFiber = yield* Stream.runForEach( + adapter.streamEvents, + () => Effect.void, + ).pipe(Effect.forkChild); + + yield* adapter.startSession({ + threadId: THREAD_ID, + provider: ProviderDriverKind.make("claudeAgent"), + runtimeMode: "full-access", + }); + + const selection = createModelSelection( + ProviderInstanceId.make("claudeAgent"), + "claude-opus-4-6", + ); + yield* adapter.sendTurn({ + threadId: THREAD_ID, + input: "first", + modelSelection: selection, + attachments: [], + }); + assert.deepEqual(harness.query.setModelCalls, ["claude-opus-4-6[1m]"]); + + harness.query.emit({ + type: "system", + subtype: "model_refusal_fallback", + trigger: "refusal", + direction: "retry", + original_model: "claude-opus-4-6[1m]", + fallback_model: "claude-sonnet-5", + request_id: null, + session_id: "sdk-session-refusal-resend", + uuid: "refusal-resend-swap", + } as unknown as SDKMessage); + harness.query.emit({ + type: "result", + subtype: "success", + is_error: false, + duration_ms: 10, + duration_api_ms: 8, + num_turns: 1, + result: "done", + stop_reason: "end_turn", + session_id: "sdk-session-refusal-resend", + usage: { input_tokens: 1_000, output_tokens: 10 }, + } as unknown as SDKMessage); + yield* Effect.yieldNow; + + // The selection has not changed, so the second turn must not call + // setModel at all — least of all with the model the API just refused. + yield* adapter.sendTurn({ + threadId: THREAD_ID, + input: "second", + modelSelection: selection, + attachments: [], + }); + assert.deepEqual(harness.query.setModelCalls, ["claude-opus-4-6[1m]"]); + + yield* Fiber.interrupt(runtimeEventsFiber); + }).pipe( + Effect.provideService(Random.Random, makeDeterministicRandomService()), + Effect.provide(harness.layer), + ); + }); + it.effect("emits Claude context window on result completion usage snapshots", () => { const harness = makeHarness(); return Effect.gen(function* () { @@ -2959,10 +3564,10 @@ describe("ClaudeAdapterLive", () => { return Effect.gen(function* () { const adapter = yield* ClaudeAdapter; - const runtimeEventsFiber = yield* Stream.take(adapter.streamEvents, 9).pipe( - Stream.runCollect, - Effect.forkChild, - ); + const runtimeEvents: Array = []; + const runtimeEventsFiber = yield* Stream.runForEach(adapter.streamEvents, (event) => + Effect.sync(() => runtimeEvents.push(event)), + ).pipe(Effect.forkChild); yield* adapter.startSession({ threadId: THREAD_ID, @@ -2976,6 +3581,27 @@ describe("ClaudeAdapterLive", () => { attachments: [], }); + // The parent records usage of its own first, so the task snapshot that + // follows is actually retained rather than dropped for want of a + // baseline. Without this the assertion below would hold even if + // completion silently discarded a recorded task total. + harness.query.emit({ + type: "assistant", + message: { + id: "msg-parent-baseline", + type: "message", + role: "assistant", + model: "claude-opus-4-6", + content: [{ type: "text", text: "working" }], + stop_reason: null, + stop_sequence: null, + usage: { input_tokens: 12000, output_tokens: 0 }, + }, + parent_tool_use_id: null, + session_id: "sdk-session-task-usage-clamped", + uuid: "parent-baseline-clamped", + } as unknown as SDKMessage); + harness.query.emit({ type: "system", subtype: "task_progress", @@ -3010,18 +3636,23 @@ describe("ClaudeAdapterLive", () => { } as unknown as SDKMessage); harness.query.finish(); - const runtimeEvents = Array.from(yield* Fiber.join(runtimeEventsFiber)); + yield* Effect.yieldNow; + yield* Fiber.interrupt(runtimeEventsFiber); const usageEvents = runtimeEvents.filter( (event) => event.type === "thread.token-usage.updated", ); const finalUsageEvent = usageEvents.at(-1); assert.equal(finalUsageEvent?.type, "thread.token-usage.updated"); if (finalUsageEvent?.type === "thread.token-usage.updated") { + // The task_progress 190,000 belongs to a subagent, so it survives as + // the running total only. The parent's own used count comes from its + // last assistant usage, not the cumulative result total. assert.deepEqual(finalUsageEvent.payload, { usage: { - usedTokens: 190000, - lastUsedTokens: 190000, + usedTokens: 12000, + lastUsedTokens: 12000, totalProcessedTokens: 535000, + inputTokens: 12000, maxTokens: 200000, }, }); diff --git a/apps/server/src/provider/Layers/ClaudeAdapter.ts b/apps/server/src/provider/Layers/ClaudeAdapter.ts index 6989378d8287..286bf0e6a347 100644 --- a/apps/server/src/provider/Layers/ClaudeAdapter.ts +++ b/apps/server/src/provider/Layers/ClaudeAdapter.ts @@ -279,6 +279,11 @@ interface ClaudeSessionContext { readonly startedAt: string; readonly basePermissionMode: PermissionMode | undefined; currentApiModelId: string | undefined; + /** The model the SDK reports is actually serving this session, which is not + * always the selected one: a refusal retry swaps it for the rest of the + * session. Kept apart from `currentApiModelId`, which tracks the last id + * passed to `setModel` and must keep matching the user's selection. */ + observedApiModelId: string | undefined; /** Effective effort for the session's turns; subagents without an explicit * effort override inherit this. */ currentEffort: string | undefined; @@ -455,11 +460,20 @@ function asRuntimeItemId(value: string): RuntimeItemId { return RuntimeItemId.make(value); } -function maxClaudeContextWindowFromModelUsage( +// `modelUsage` is keyed by every model that ran during the turn, subagents +// included, so the maximum could be a child's window rather than this +// session's. Fall back to it only when the session model has no entry. +function claudeContextWindowFromModelUsage( modelUsage: Record | undefined, + sessionModel: string | undefined, ): number | undefined { if (!modelUsage) return undefined; + const sessionEntry = sessionModel ? modelUsage[sessionModel] : undefined; + if (sessionEntry) { + return sessionEntry.contextWindow; + } + let maxContextWindow: number | undefined; for (const value of Object.values(modelUsage)) { const contextWindow = value.contextWindow; @@ -632,6 +646,8 @@ function compactBoundaryTokenUsageSnapshot( }); } +// A subagent's tokens are spent in its own context window, so they advance the +// thread's running total but never the parent's used count (#5942). function normalizeClaudeTaskProgressTokenUsage( value: unknown, context: ClaudeSessionContext, @@ -641,32 +657,33 @@ function normalizeClaudeTaskProgressTokenUsage( return undefined; } - const lastUsedTokens = context.lastKnownTokenUsage?.usedTokens; - const activeTokens = - lastUsedTokens !== undefined ? Math.max(totalTokens, lastUsedTokens) : totalTokens; - if (lastUsedTokens !== undefined && activeTokens === lastUsedTokens) { + const lastKnown = context.lastKnownTokenUsage; + if (!lastKnown) { return undefined; } - const usage = value as Record; - const snapshot = makeClaudeTokenUsageSnapshot({ - activeTokens, - ...(context.lastKnownContextWindow !== undefined - ? { contextWindow: context.lastKnownContextWindow } - : {}), - totalProcessedTokens: Math.max( - totalTokens, - context.lastKnownTotalProcessedTokens ?? totalTokens, - ), - }); - if (!snapshot) { + // The running total floors at the largest single contributor rather than + // summing per-task totals: the SDK does not document whether result usage + // already aggregates children, and a floor can only understate. + const totalProcessedTokens = Math.max( + totalTokens, + context.lastKnownTotalProcessedTokens ?? totalTokens, + ); + if (totalProcessedTokens === context.lastKnownTotalProcessedTokens) { + return undefined; + } + // Match makeClaudeTokenUsageSnapshot's invariant: a running total is only + // meaningful once it exceeds the parent's own used count. + if (totalProcessedTokens <= lastKnown.usedTokens) { return undefined; } + const usage = value as Record; const toolUses = finiteNonNegativeInteger(usage.tool_uses); const durationMs = finiteNonNegativeInteger(usage.duration_ms); return { - ...snapshot, + ...lastKnown, + totalProcessedTokens, ...(toolUses !== undefined ? { toolUses } : {}), ...(durationMs !== undefined ? { durationMs } : {}), }; @@ -2207,13 +2224,23 @@ export const makeClaudeAdapter = Effect.fn("makeClaudeAdapter")(function* ( errorMessage?: string, result?: SDKResultMessage, ) { - const resultContextWindow = maxClaudeContextWindowFromModelUsage(result?.modelUsage); + const resultContextWindow = claudeContextWindowFromModelUsage( + result?.modelUsage, + context.observedApiModelId ?? context.currentApiModelId, + ); if (resultContextWindow !== undefined) { context.lastKnownContextWindow = resultContextWindow; } const maxTokens = resultContextWindow ?? context.lastKnownContextWindow; - const accumulatedTotalProcessedTokens = claudeTotalProcessedTokens(result?.usage); + // The result carries only the parent's own usage, which can be smaller + // than a running total a subagent already raised; the thread total only + // ever grows. + const resultTotalProcessedTokens = claudeTotalProcessedTokens(result?.usage); + const accumulatedTotalProcessedTokens = + resultTotalProcessedTokens !== undefined + ? Math.max(resultTotalProcessedTokens, context.lastKnownTotalProcessedTokens ?? 0) + : undefined; if (accumulatedTotalProcessedTokens !== undefined) { context.lastKnownTotalProcessedTokens = accumulatedTotalProcessedTokens; } @@ -3131,7 +3158,13 @@ export const makeClaudeAdapter = Effect.fn("makeClaudeAdapter")(function* ( } switch (message.subtype) { - case "init": + case "init": { + // init is the first place the SDK names the model actually serving the + // session, which is the id `modelUsage` is keyed by. + const initModel = trimmedString(message.model); + if (initModel) { + context.observedApiModelId = initModel; + } yield* offerRuntimeEvent({ ...base, type: "session.configured", @@ -3140,6 +3173,7 @@ export const makeClaudeAdapter = Effect.fn("makeClaudeAdapter")(function* ( }, }); return; + } case "status": yield* offerRuntimeEvent({ ...base, @@ -3433,9 +3467,21 @@ export const makeClaudeAdapter = Effect.fn("makeClaudeAdapter")(function* ( yield* emitRuntimeWarning(context, message.text, message); } return; + case "model_refusal_fallback": { + // A refusal retry swaps the model for the rest of the session, so the + // window lookup has to follow it. `currentApiModelId` deliberately + // does not: it mirrors the user's selection for `setModel`, and moving + // it here would make the next turn re-send the refused model. + if (message.direction === "retry") { + const fallbackModel = trimmedString(message.fallback_model); + if (fallbackModel) { + context.observedApiModelId = fallbackModel; + } + } + return; + } // Inner protocol/UX details with no T3 surface today — consumed // deliberately so they don't masquerade as unknown-subtype warnings. - case "model_refusal_fallback": case "local_command_output": case "plugin_install": case "commands_changed": @@ -4397,6 +4443,7 @@ export const makeClaudeAdapter = Effect.fn("makeClaudeAdapter")(function* ( startedAt, basePermissionMode: permissionMode, currentApiModelId: apiModelId, + observedApiModelId: undefined, currentEffort: effectiveEffort ?? undefined, resumeSessionId: sessionId, pendingApprovals, @@ -4519,6 +4566,9 @@ export const makeClaudeAdapter = Effect.fn("makeClaudeAdapter")(function* ( catch: (cause) => toRequestError(input.threadId, "turn/setModel", cause), }); context.currentApiModelId = apiModelId; + // A deliberate switch supersedes whatever init or a refusal observed; + // leaving it stale would measure the meter against the old model. + context.observedApiModelId = apiModelId; } context.session = { ...context.session,