diff --git a/apps/server/src/provider/Layers/ClaudeAdapter.test.ts b/apps/server/src/provider/Layers/ClaudeAdapter.test.ts index 38bdb7f1f2e7..011d9e476f26 100644 --- a/apps/server/src/provider/Layers/ClaudeAdapter.test.ts +++ b/apps/server/src/provider/Layers/ClaudeAdapter.test.ts @@ -4061,7 +4061,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, ); @@ -4914,12 +4914,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, ); @@ -4947,20 +4947,90 @@ 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 runtimeEventsFiber = yield* Stream.takeUntil( + adapter.streamEvents, + (event) => event.type === "task.completed", + ).pipe(Stream.runCollect, 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); + harness.query.emit({ + type: "system", + subtype: "task_notification", + task_id: "task-total-floor", + status: "completed", + summary: "Child finished", + usage: { total_tokens: 5_000 }, + session_id: "sdk-session-total-floor", + uuid: "total-floor-completed", + } as unknown as SDKMessage); + + const runtimeEvents = Array.from(yield* Fiber.join(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()), @@ -4968,15 +5038,96 @@ describe("ClaudeAdapterLive", () => { ); }); - it.effect("emits Claude context window on result completion usage snapshots", () => { + it.effect("adds cumulative usage deltas from separate Claude tasks", () => { const harness = makeHarness(); return Effect.gen(function* () { const adapter = yield* ClaudeAdapter; + const runtimeEventsFiber = yield* Stream.takeUntil( + adapter.streamEvents, + (event) => event.type === "task.completed", + ).pipe(Stream.runCollect, Effect.forkChild); - const runtimeEventsFiber = yield* Stream.take(adapter.streamEvents, 7).pipe( - Stream.runCollect, - Effect.forkChild, + yield* adapter.startSession({ + threadId: THREAD_ID, + provider: ProviderDriverKind.make("claudeAgent"), + runtimeMode: "full-access", + }); + yield* adapter.sendTurn({ threadId: THREAD_ID, input: "delegate", attachments: [] }); + + harness.query.emit({ + type: "result", + subtype: "success", + is_error: false, + session_id: "sdk-session-multiple-task-usage", + uuid: "multiple-task-parent-result", + usage: { input_tokens: 4_000, output_tokens: 200 }, + } as unknown as SDKMessage); + + harness.query.emit({ + type: "system", + subtype: "task_progress", + task_id: "task-usage-a", + description: "First child working", + usage: { total_tokens: 100_000, tool_uses: 10, duration_ms: 1_000 }, + session_id: "sdk-session-multiple-task-usage", + uuid: "multiple-task-a-progress-1", + } as unknown as SDKMessage); + harness.query.emit({ + type: "system", + subtype: "task_progress", + task_id: "task-usage-b", + description: "Second child working", + usage: { total_tokens: 100_000, tool_uses: 5, duration_ms: 500 }, + session_id: "sdk-session-multiple-task-usage", + uuid: "multiple-task-b-progress-1", + } as unknown as SDKMessage); + harness.query.emit({ + type: "system", + subtype: "task_progress", + task_id: "task-usage-a", + description: "First child still working", + usage: { total_tokens: 120_000, tool_uses: 12, duration_ms: 1_200 }, + session_id: "sdk-session-multiple-task-usage", + uuid: "multiple-task-a-progress-2", + } as unknown as SDKMessage); + harness.query.emit({ + type: "system", + subtype: "task_notification", + task_id: "task-usage-b", + status: "completed", + summary: "Second child finished", + usage: { total_tokens: 100_000, tool_uses: 5, duration_ms: 500 }, + session_id: "sdk-session-multiple-task-usage", + uuid: "multiple-task-b-completed", + } as unknown as SDKMessage); + + const runtimeEvents = Array.from(yield* Fiber.join(runtimeEventsFiber)); + const usageEvents = runtimeEvents.filter( + (event) => event.type === "thread.token-usage.updated", ); + assert.equal(usageEvents.length, 4); + 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.usedTokens, 4_200); + assert.equal(latest.payload.usage.totalProcessedTokens, 224_200); + assert.equal(latest.payload.usage.toolUses, 17); + assert.equal(latest.payload.usage.durationMs, 1_700); + } + }).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 runtimeEventsFiber = yield* Stream.takeUntil( + adapter.streamEvents, + (event) => event.type === "task.completed", + ).pipe(Stream.runCollect, Effect.forkChild); yield* adapter.startSession({ threadId: THREAD_ID, @@ -4986,48 +5137,57 @@ describe("ClaudeAdapterLive", () => { yield* adapter.sendTurn({ threadId: THREAD_ID, - input: "hello", + 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: 1234, - duration_api_ms: 1200, + duration_ms: 10, + duration_api_ms: 8, num_turns: 1, result: "done", stop_reason: "end_turn", - session_id: "sdk-session-result-usage", - usage: { - input_tokens: 4, - cache_creation_input_tokens: 2715, - cache_read_input_tokens: 21144, - output_tokens: 679, - }, - modelUsage: { - [SYNTHETIC_CLAUDE_CAPABLE_MODEL]: { - contextWindow: 200000, - maxOutputTokens: 64000, - }, - }, + session_id: "sdk-session-child-usage", + usage: { input_tokens: 4_000, output_tokens: 200 }, + modelUsage: { [SYNTHETIC_CLAUDE_CAPABLE_MODEL]: { 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.emit({ + type: "system", + subtype: "task_notification", + task_id: "task-child-1", + status: "completed", + summary: "Background agent finished", + usage: { total_tokens: 900_000, tool_uses: 40, duration_ms: 1_000 }, + session_id: "sdk-session-child-usage", + uuid: "task-child-completed-1", } as unknown as SDKMessage); - harness.query.finish(); const runtimeEvents = Array.from(yield* Fiber.join(runtimeEventsFiber)); - const usageEvent = runtimeEvents.find((event) => event.type === "thread.token-usage.updated"); - assert.equal(usageEvent?.type, "thread.token-usage.updated"); - if (usageEvent?.type === "thread.token-usage.updated") { - assert.deepEqual(usageEvent.payload, { - usage: { - usedTokens: 24542, - lastUsedTokens: 24542, - inputTokens: 23863, - outputTokens: 679, - maxTokens: 200000, - }, - }); + 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 advances only + // the separate running total. + assert.equal(latest.payload.usage.usedTokens, 4_200); + assert.equal(latest.payload.usage.totalProcessedTokens, 904_200); } }).pipe( Effect.provideService(Random.Random, makeDeterministicRandomService()), @@ -5035,15 +5195,71 @@ describe("ClaudeAdapterLive", () => { ); }); - it.effect("clamps oversized Claude usage to the reported context window", () => { + 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 runtimeEventsFiber = yield* Stream.take(adapter.streamEvents, 7).pipe( - Stream.runCollect, - Effect.forkChild, + const runtimeEventsFiber = yield* Stream.takeUntil( + adapter.streamEvents, + (event) => event.type === "turn.completed", + ).pipe(Stream.runCollect, Effect.forkChild); + + const session = yield* adapter.startSession({ + threadId: THREAD_ID, + provider: ProviderDriverKind.make("claudeAgent"), + runtimeMode: "full-access", + modelSelection: createModelSelection( + ProviderInstanceId.make("claudeAgent"), + SYNTHETIC_CLAUDE_STANDARD_MODEL, + [{ id: "contextWindow", value: "standard" }], + ), + }); + + 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: { + [SYNTHETIC_CLAUDE_STANDARD_MODEL]: { contextWindow: 200_000 }, + [`${SYNTHETIC_CLAUDE_CAPABLE_MODEL}[expanded]`]: { contextWindow: 1_000_000 }, + }, + } as unknown as SDKMessage); + + const runtimeEvents = Array.from(yield* Fiber.join(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()), + Effect.provide(harness.layer), + ); + }); + + 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 runtimeEventsFiber = yield* Stream.takeUntil( + adapter.streamEvents, + (event) => event.type === "turn.completed", + ).pipe(Stream.runCollect, Effect.forkChild); yield* adapter.startSession({ threadId: THREAD_ID, @@ -5053,44 +5269,45 @@ describe("ClaudeAdapterLive", () => { yield* adapter.sendTurn({ threadId: THREAD_ID, - input: "hello", + 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: SYNTHETIC_CLAUDE_STANDARD_MODEL, + 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: 1234, - duration_api_ms: 1200, + duration_ms: 10, + duration_api_ms: 8, num_turns: 1, result: "done", stop_reason: "end_turn", - session_id: "sdk-session-result-usage-clamped", - usage: { - total_tokens: 535000, - }, + session_id: "sdk-session-init-window", + usage: { input_tokens: 50_000, output_tokens: 1_000 }, modelUsage: { - [SYNTHETIC_CLAUDE_CAPABLE_MODEL]: { - contextWindow: 200000, - maxOutputTokens: 64000, - }, + [SYNTHETIC_CLAUDE_STANDARD_MODEL]: { contextWindow: 200_000 }, + [`${SYNTHETIC_CLAUDE_CAPABLE_MODEL}[expanded]`]: { contextWindow: 1_000_000 }, }, } as unknown as SDKMessage); - harness.query.finish(); const runtimeEvents = Array.from(yield* Fiber.join(runtimeEventsFiber)); - const usageEvent = runtimeEvents.find((event) => event.type === "thread.token-usage.updated"); - assert.equal(usageEvent?.type, "thread.token-usage.updated"); - if (usageEvent?.type === "thread.token-usage.updated") { - assert.deepEqual(usageEvent.payload, { - usage: { - usedTokens: 200000, - lastUsedTokens: 200000, - totalProcessedTokens: 535000, - maxTokens: 200000, - }, - }); + 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()), @@ -5098,23 +5315,630 @@ describe("ClaudeAdapterLive", () => { ); }); - it.effect( - "preserves oversized Claude result totals after task progress snapshots are recorded", - () => { - const harness = makeHarness(); - return Effect.gen(function* () { - const adapter = yield* ClaudeAdapter; + 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 runtimeEventsFiber = yield* Stream.takeUntil( + adapter.streamEvents, + (event) => event.type === "turn.completed", + ).pipe(Stream.runCollect, Effect.forkChild); - const runtimeEventsFiber = yield* Stream.take(adapter.streamEvents, 9).pipe( - Stream.runCollect, - Effect.forkChild, - ); + yield* adapter.startSession({ + threadId: THREAD_ID, + provider: ProviderDriverKind.make("claudeAgent"), + runtimeMode: "full-access", + }); - 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); + harness.query.emit({ + type: "result", + subtype: "success", + is_error: false, + session_id: "sdk-session-running-total", + uuid: "running-total-result", + } as unknown as SDKMessage); + + const runtimeEvents = Array.from(yield* Fiber.join(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 runtimeEventsFiber = yield* Stream.takeUntil( + adapter.streamEvents, + (event) => event.type === "turn.completed", + ).pipe(Stream.runCollect, 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: { [SYNTHETIC_CLAUDE_CAPABLE_MODEL]: { contextWindow: 200_000 } }, + } as unknown as SDKMessage); + + const runtimeEvents = Array.from(yield* Fiber.join(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.each([ + { scope: undefined, expectedWindow: 200_000 }, + { scope: "session" as const, expectedWindow: 200_000 }, + { scope: "local" as const, expectedWindow: 1_000_000 }, + ])("uses the parent window after a $scope refusal fallback", ({ scope, expectedWindow }) => { + const harness = makeHarness(); + return Effect.gen(function* () { + const adapter = yield* ClaudeAdapter; + const runtimeEventsFiber = yield* Stream.takeUntil( + adapter.streamEvents, + (event) => event.type === "turn.completed", + ).pipe(Stream.runCollect, 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: `${SYNTHETIC_CLAUDE_CAPABLE_MODEL}[expanded]`, + session_id: "sdk-session-refusal", + uuid: "refusal-init", + } as unknown as SDKMessage); + + // Session fallbacks change the parent model; local fallbacks do not. + harness.query.emit({ + type: "system", + subtype: "model_refusal_fallback", + trigger: "refusal", + direction: "retry", + ...(scope !== undefined ? { scope } : {}), + original_model: `${SYNTHETIC_CLAUDE_CAPABLE_MODEL}[expanded]`, + fallback_model: SYNTHETIC_CLAUDE_STANDARD_MODEL, + content: "Retrying with the fallback model.", + 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: { + [SYNTHETIC_CLAUDE_STANDARD_MODEL]: { contextWindow: 200_000 }, + [`${SYNTHETIC_CLAUDE_CAPABLE_MODEL}[expanded]`]: { contextWindow: 1_000_000 }, + }, + } as unknown as SDKMessage); + harness.query.finish(); + + const runtimeEvents = Array.from(yield* Fiber.join(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, expectedWindow); + } + }).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; + + yield* adapter.startSession({ + threadId: THREAD_ID, + provider: ProviderDriverKind.make("claudeAgent"), + runtimeMode: "full-access", + }); + + const initEventsFiber = yield* Stream.takeUntil( + adapter.streamEvents, + (event) => event.type === "session.configured", + ).pipe(Stream.runCollect, Effect.forkChild); + harness.query.emit({ + type: "system", + subtype: "init", + model: `${SYNTHETIC_CLAUDE_CAPABLE_MODEL}[expanded]`, + session_id: "sdk-session-switch", + uuid: "switch-init", + } as unknown as SDKMessage); + yield* Fiber.join(initEventsFiber); + + const runtimeEventsFiber = yield* Stream.takeUntil( + adapter.streamEvents, + (event) => event.type === "turn.completed", + ).pipe(Stream.runCollect, Effect.forkChild); + // The user switches to a 200k model partway through the thread. + yield* adapter.sendTurn({ + threadId: THREAD_ID, + input: "switch", + modelSelection: createModelSelection( + ProviderInstanceId.make("claudeAgent"), + SYNTHETIC_CLAUDE_STANDARD_MODEL, + [{ id: "contextWindow", value: "standard" }], + ), + 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: { + [`${SYNTHETIC_CLAUDE_CAPABLE_MODEL}[expanded]`]: { contextWindow: 1_000_000 }, + [SYNTHETIC_CLAUDE_STANDARD_MODEL]: { contextWindow: 200_000 }, + }, + } as unknown as SDKMessage); + + const runtimeEvents = Array.from(yield* Fiber.join(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 firstTurnEventsFiber = yield* Stream.takeUntil( + adapter.streamEvents, + (event) => event.type === "turn.completed", + ).pipe(Stream.runCollect, Effect.forkChild); + + yield* adapter.startSession({ + threadId: THREAD_ID, + provider: ProviderDriverKind.make("claudeAgent"), + runtimeMode: "full-access", + }); + + const selection = createModelSelection( + ProviderInstanceId.make("claudeAgent"), + SYNTHETIC_CLAUDE_CAPABLE_MODEL, + [{ id: "contextWindow", value: "expanded" }], + ); + yield* adapter.sendTurn({ + threadId: THREAD_ID, + input: "first", + modelSelection: selection, + attachments: [], + }); + assert.deepEqual(harness.query.setModelCalls, [ + `${SYNTHETIC_CLAUDE_CAPABLE_MODEL}[expanded]`, + ]); + + harness.query.emit({ + type: "system", + subtype: "model_refusal_fallback", + trigger: "refusal", + direction: "retry", + original_model: `${SYNTHETIC_CLAUDE_CAPABLE_MODEL}[expanded]`, + fallback_model: SYNTHETIC_CLAUDE_STANDARD_MODEL, + 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* Fiber.join(firstTurnEventsFiber); + + // 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, [ + `${SYNTHETIC_CLAUDE_CAPABLE_MODEL}[expanded]`, + ]); + }).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* () { + const adapter = yield* ClaudeAdapter; + + const runtimeEventsFiber = yield* Stream.take(adapter.streamEvents, 7).pipe( + Stream.runCollect, + Effect.forkChild, + ); + + yield* adapter.startSession({ + threadId: THREAD_ID, + provider: ProviderDriverKind.make("claudeAgent"), + runtimeMode: "full-access", + }); + + yield* adapter.sendTurn({ + threadId: THREAD_ID, + input: "hello", + attachments: [], + }); + + harness.query.emit({ + type: "result", + subtype: "success", + is_error: false, + duration_ms: 1234, + duration_api_ms: 1200, + num_turns: 1, + result: "done", + stop_reason: "end_turn", + session_id: "sdk-session-result-usage", + usage: { + input_tokens: 4, + cache_creation_input_tokens: 2715, + cache_read_input_tokens: 21144, + output_tokens: 679, + }, + modelUsage: { + [SYNTHETIC_CLAUDE_CAPABLE_MODEL]: { + contextWindow: 200000, + maxOutputTokens: 64000, + }, + }, + } as unknown as SDKMessage); + harness.query.finish(); + + const runtimeEvents = Array.from(yield* Fiber.join(runtimeEventsFiber)); + const usageEvent = runtimeEvents.find((event) => event.type === "thread.token-usage.updated"); + assert.equal(usageEvent?.type, "thread.token-usage.updated"); + if (usageEvent?.type === "thread.token-usage.updated") { + assert.deepEqual(usageEvent.payload, { + usage: { + usedTokens: 24542, + lastUsedTokens: 24542, + inputTokens: 23863, + outputTokens: 679, + maxTokens: 200000, + }, + }); + } + }).pipe( + Effect.provideService(Random.Random, makeDeterministicRandomService()), + Effect.provide(harness.layer), + ); + }); + + it.effect.each([ + { + name: "does not treat total-only Claude result usage as active context", + usage: { total_tokens: 535_000 }, + expectedUsedTokens: undefined, + }, + { + name: "does not treat a total-only result iteration as active context", + usage: { total_tokens: 535_000, iterations: [{ total_tokens: 535_000 }] }, + expectedUsedTokens: undefined, + }, + { + name: "does not use aggregate input when the selected iteration is total-only", + usage: { input_tokens: 535_000, iterations: [{ total_tokens: 535_000 }] }, + expectedUsedTokens: undefined, + }, + { + name: "uses active input and output from the selected result iteration", + usage: { + input_tokens: 535_000, + iterations: [{ input_tokens: 4_000, output_tokens: 200 }], + }, + expectedUsedTokens: 4_200, + }, + ])("$name", ({ usage, expectedUsedTokens }) => { + const harness = makeHarness(); + return Effect.gen(function* () { + const adapter = yield* ClaudeAdapter; + + const runtimeEventsFiber = yield* Stream.takeUntil( + adapter.streamEvents, + (event) => event.type === "turn.completed", + ).pipe(Stream.runCollect, Effect.forkChild); + + yield* adapter.startSession({ + threadId: THREAD_ID, + provider: ProviderDriverKind.make("claudeAgent"), + runtimeMode: "full-access", + }); + + yield* adapter.sendTurn({ + threadId: THREAD_ID, + input: "hello", + attachments: [], + }); + + harness.query.emit({ + type: "result", + subtype: "success", + is_error: false, + duration_ms: 1234, + duration_api_ms: 1200, + num_turns: 1, + result: "done", + stop_reason: "end_turn", + session_id: "sdk-session-result-total-only", + usage, + modelUsage: { + [SYNTHETIC_CLAUDE_CAPABLE_MODEL]: { + contextWindow: 200000, + maxOutputTokens: 64000, + }, + }, + } as unknown as SDKMessage); + harness.query.finish(); + + const runtimeEvents = Array.from(yield* Fiber.join(runtimeEventsFiber)); + const usageEvent = runtimeEvents.find((event) => event.type === "thread.token-usage.updated"); + if (expectedUsedTokens === undefined) { + assert.equal(usageEvent, undefined); + } else { + assert.equal(usageEvent?.type, "thread.token-usage.updated"); + if (usageEvent?.type === "thread.token-usage.updated") { + assert.equal(usageEvent.payload.usage.usedTokens, expectedUsedTokens); + } + } + }).pipe( + Effect.provideService(Random.Random, makeDeterministicRandomService()), + Effect.provide(harness.layer), + ); + }); + + it.effect("keeps subagent totals out of parent context when the result is total-only", () => { + const harness = makeHarness(); + return Effect.gen(function* () { + const adapter = yield* ClaudeAdapter; + + const runtimeEventsFiber = yield* Stream.takeUntil( + adapter.streamEvents, + (event) => event.type === "turn.completed", + ).pipe(Stream.runCollect, Effect.forkChild); + + yield* adapter.startSession({ + threadId: THREAD_ID, + provider: ProviderDriverKind.make("claudeAgent"), + runtimeMode: "full-access", + }); + + yield* adapter.sendTurn({ + threadId: THREAD_ID, + input: "hello", + attachments: [], + }); + + harness.query.emit({ + type: "system", + subtype: "task_progress", + task_id: "task-total-only", + description: "Thinking through the patch", + usage: { + total_tokens: 190000, + }, + session_id: "sdk-session-total-only-after-progress", + uuid: "task-total-only-progress", + } as unknown as SDKMessage); + + harness.query.emit({ + type: "system", + subtype: "task_notification", + task_id: "task-total-only", + status: "completed", + summary: "Task finished", + usage: { + total_tokens: 250000, + tool_uses: 100, + duration_ms: 900000, + }, + session_id: "sdk-session-total-only-after-progress", + uuid: "task-total-only-completed", + } as unknown as SDKMessage); + + harness.query.emit({ + type: "result", + subtype: "success", + is_error: false, + duration_ms: 1234, + duration_api_ms: 1200, + num_turns: 1, + result: "done", + stop_reason: "end_turn", + session_id: "sdk-session-total-only-after-progress", + usage: { + total_tokens: 535000, + }, + modelUsage: { + [SYNTHETIC_CLAUDE_CAPABLE_MODEL]: { + contextWindow: 200000, + maxOutputTokens: 64000, + }, + }, + } as unknown as SDKMessage); + harness.query.finish(); + + const runtimeEvents = Array.from(yield* Fiber.join(runtimeEventsFiber)); + const usageEvents = runtimeEvents.filter( + (event) => event.type === "thread.token-usage.updated", + ); + // The parent never reported usage of its own, so neither the children's + // totals nor the cumulative result may fabricate a meter reading. The + // children's numbers still ride on their task events. + assert.deepEqual(usageEvents, []); + const progressEvent = runtimeEvents.find((event) => event.type === "task.progress"); + const completedEvent = runtimeEvents.find((event) => event.type === "task.completed"); + assert.equal(progressEvent?.type, "task.progress"); + assert.equal(completedEvent?.type, "task.completed"); + if (progressEvent?.type === "task.progress") { + assert.equal(progressEvent.payload.typedUsage?.totalTokens, 190000); + } + if (completedEvent?.type === "task.completed") { + assert.equal(completedEvent.payload.typedUsage?.totalTokens, 250000); + } + }).pipe( + Effect.provideService(Random.Random, makeDeterministicRandomService()), + Effect.provide(harness.layer), + ); + }); + + it.effect( + "preserves oversized Claude result totals after task progress snapshots are recorded", + () => { + const harness = makeHarness(); + return Effect.gen(function* () { + const adapter = yield* ClaudeAdapter; + + const runtimeEventsFiber = yield* Stream.takeUntil( + adapter.streamEvents, + (event) => event.type === "turn.completed", + ).pipe(Stream.runCollect, Effect.forkChild); + + yield* adapter.startSession({ + threadId: THREAD_ID, + provider: ProviderDriverKind.make("claudeAgent"), + runtimeMode: "full-access", + }); yield* adapter.sendTurn({ threadId: THREAD_ID, @@ -5122,6 +5946,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: SYNTHETIC_CLAUDE_CAPABLE_MODEL, + 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", @@ -5154,7 +5999,6 @@ describe("ClaudeAdapterLive", () => { }, }, } as unknown as SDKMessage); - harness.query.finish(); const runtimeEvents = Array.from(yield* Fiber.join(runtimeEventsFiber)); const usageEvents = runtimeEvents.filter( @@ -5163,11 +6007,15 @@ describe("ClaudeAdapterLive", () => { 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 3afde2b39a7a..d8eba6aac8d0 100644 --- a/apps/server/src/provider/Layers/ClaudeAdapter.ts +++ b/apps/server/src/provider/Layers/ClaudeAdapter.ts @@ -312,6 +312,12 @@ function rememberPendingTaskModel( } } +interface ClaudeTaskCumulativeUsage { + readonly totalTokens: number; + readonly toolUses: number; + readonly durationMs: number; +} + interface ClaudeSessionContext { session: ProviderSession; startInput: Parameters[0]; @@ -322,6 +328,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; @@ -355,6 +366,8 @@ interface ClaudeSessionContext { lastKnownContextWindow: number | undefined; lastKnownTokenUsage: ThreadTokenUsageSnapshot | undefined; lastKnownTotalProcessedTokens: number | undefined; + readonly taskUsageById: Map; + taskUsageTotals: ClaudeTaskCumulativeUsage; lastAssistantUuid: string | undefined; lastThreadStartedId: string | undefined; /** Limits already announced for the running turn, keyed `window:resetsAt`. */ @@ -576,11 +589,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; @@ -824,7 +846,10 @@ function compactBoundaryTokenUsageSnapshot( return snapshotWithoutBeforeTokens; } +// 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( + taskId: string, value: unknown, context: ClaudeSessionContext, ): ThreadTokenUsageSnapshot | undefined { @@ -833,32 +858,54 @@ function normalizeClaudeTaskProgressTokenUsage( return undefined; } - const lastUsedTokens = context.lastKnownTokenUsage?.usedTokens; - const activeTokens = - lastUsedTokens !== undefined ? Math.max(totalTokens, lastUsedTokens) : totalTokens; - if (lastUsedTokens !== undefined && activeTokens === lastUsedTokens) { + const usage = value as Record; + // SDK task usage is cumulative per task. Per-task high-water marks keep + // repeated progress and completion receipts from counting the same work twice. + const previous = context.taskUsageById.get(taskId); + const next: ClaudeTaskCumulativeUsage = { + totalTokens: Math.max(previous?.totalTokens ?? 0, totalTokens), + toolUses: Math.max(previous?.toolUses ?? 0, finiteNonNegativeInteger(usage.tool_uses) ?? 0), + durationMs: Math.max( + previous?.durationMs ?? 0, + finiteNonNegativeInteger(usage.duration_ms) ?? 0, + ), + }; + const tokenDelta = next.totalTokens - (previous?.totalTokens ?? 0); + const toolUseDelta = next.toolUses - (previous?.toolUses ?? 0); + const durationDelta = next.durationMs - (previous?.durationMs ?? 0); + if (tokenDelta === 0 && toolUseDelta === 0 && durationDelta === 0) { 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) { + context.taskUsageById.set(taskId, next); + context.taskUsageTotals = { + totalTokens: context.taskUsageTotals.totalTokens + tokenDelta, + toolUses: context.taskUsageTotals.toolUses + toolUseDelta, + durationMs: context.taskUsageTotals.durationMs + durationDelta, + }; + const totalProcessedTokens = (context.lastKnownTotalProcessedTokens ?? 0) + tokenDelta; + context.lastKnownTotalProcessedTokens = totalProcessedTokens; + + const lastKnown = context.lastKnownTokenUsage; + if (!lastKnown) { + return undefined; + } + + const visibleTotal = + totalProcessedTokens > lastKnown.usedTokens ? totalProcessedTokens : undefined; + const toolUses = context.taskUsageTotals.toolUses || undefined; + const durationMs = context.taskUsageTotals.durationMs || undefined; + if ( + visibleTotal === lastKnown.totalProcessedTokens && + toolUses === lastKnown.toolUses && + durationMs === lastKnown.durationMs + ) { return undefined; } - const toolUses = finiteNonNegativeInteger(usage.tool_uses); - const durationMs = finiteNonNegativeInteger(usage.duration_ms); return { - ...snapshot, + ...lastKnown, + ...(visibleTotal !== undefined ? { totalProcessedTokens: visibleTotal } : {}), ...(toolUses !== undefined ? { toolUses } : {}), ...(durationMs !== undefined ? { durationMs } : {}), }; @@ -2485,13 +2532,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; } @@ -2501,23 +2558,27 @@ export const makeClaudeAdapter = Effect.fn("makeClaudeAdapter")(function* ( result?.usage && typeof result.usage === "object" && !Array.isArray(result.usage) ? (result.usage as Record) : undefined; - const hasResultUsageIteration = - resultUsageRecord !== undefined && lastClaudeUsageIteration(resultUsageRecord) !== undefined; + const selectedResultUsage = resultUsageRecord + ? (lastClaudeUsageIteration(resultUsageRecord) ?? resultUsageRecord) + : undefined; const resultHasActiveUsage = - resultUsageRecord !== undefined && - (hasResultUsageIteration || - claudeUsageInputTokens(resultUsageRecord) + claudeUsageOutputTokens(resultUsageRecord) > 0); + selectedResultUsage !== undefined && + claudeUsageInputTokens(selectedResultUsage) + claudeUsageOutputTokens(selectedResultUsage) > + 0; const resultTotalOnly = resultUsageRecord !== undefined && !resultHasActiveUsage && claudeTotalProcessedTokens(resultUsageRecord) !== undefined; - const resultIterationSnapshot = resultUsageRecord - ? normalizeClaudeActiveTokenUsage( - resultUsageRecord, - maxTokens, - accumulatedTotalProcessedTokens ?? context.lastKnownTotalProcessedTokens, - ) - : undefined; + // A result carrying only cumulative total_tokens is not an active-context + // reading; without this gate it would render as a full meter (#6586). + const resultIterationSnapshot = + resultUsageRecord && resultHasActiveUsage + ? normalizeClaudeActiveTokenUsage( + resultUsageRecord, + maxTokens, + accumulatedTotalProcessedTokens ?? context.lastKnownTotalProcessedTokens, + ) + : undefined; const latestAssistantSnapshot = normalizeClaudeActiveTokenUsage( context.turnState?.latestAssistantUsage, maxTokens, @@ -3433,7 +3494,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", @@ -3442,6 +3509,7 @@ export const makeClaudeAdapter = Effect.fn("makeClaudeAdapter")(function* ( }, }); return; + } case "status": yield* offerRuntimeEvent({ ...base, @@ -3598,7 +3666,7 @@ export const makeClaudeAdapter = Effect.fn("makeClaudeAdapter")(function* ( case "task_progress": { yield* emitThreadTokenUsage( context, - normalizeClaudeTaskProgressTokenUsage(message.usage, context), + normalizeClaudeTaskProgressTokenUsage(message.task_id, message.usage, context), { rawMethod: "claude/system/task_progress", rawPayload: message, @@ -3665,7 +3733,7 @@ export const makeClaudeAdapter = Effect.fn("makeClaudeAdapter")(function* ( context.liveTaskIds.delete(message.task_id); yield* emitThreadTokenUsage( context, - normalizeClaudeTaskProgressTokenUsage(message.usage, context), + normalizeClaudeTaskProgressTokenUsage(message.task_id, message.usage, context), { rawMethod: "claude/system/task_notification", rawPayload: message, @@ -3748,7 +3816,17 @@ export const makeClaudeAdapter = Effect.fn("makeClaudeAdapter")(function* ( yield* emitRuntimeWarning(context, message.text, message); } return; - case "model_refusal_fallback": + case "model_refusal_fallback": { + // A session refusal retry changes the parent model, 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" && message.scope !== "local") { + const fallbackModel = trimmedString(message.fallback_model); + if (fallbackModel) { + context.observedApiModelId = fallbackModel; + } + } // A safety fallback switched the model mid-session (e.g. Fable 5 // retried on Opus 4.8 after a flagged request). The CLI ships the // user-facing notice in `content`; surface it like high-priority @@ -3756,6 +3834,7 @@ export const makeClaudeAdapter = Effect.fn("makeClaudeAdapter")(function* ( // by a different model. yield* emitRuntimeWarning(context, message.content, message); return; + } // Inner protocol/UX details with no T3 surface today — consumed // deliberately so they don't masquerade as unknown-subtype warnings. // `background_tasks_changed` is a roster snapshot ({tasks: [...]}); the @@ -4838,6 +4917,7 @@ export const makeClaudeAdapter = Effect.fn("makeClaudeAdapter")(function* ( startedAt, basePermissionMode: permissionMode, currentApiModelId: apiModelId, + observedApiModelId: undefined, currentEffort: effectiveEffort ?? undefined, resumeSessionId: sessionId, pendingApprovals, @@ -4853,6 +4933,8 @@ export const makeClaudeAdapter = Effect.fn("makeClaudeAdapter")(function* ( lastKnownContextWindow: initialContextWindow, lastKnownTokenUsage: undefined, lastKnownTotalProcessedTokens: undefined, + taskUsageById: new Map(), + taskUsageTotals: { totalTokens: 0, toolUses: 0, durationMs: 0 }, lastAssistantUuid: resumeState?.resumeSessionAt, lastThreadStartedId: undefined, announcedUsageLimits: undefined, @@ -4968,6 +5050,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,