From 50382c889b645a8cafcb7ab54099065696207a75 Mon Sep 17 00:00:00 2001 From: imMxts Date: Sun, 16 Aug 2026 18:11:25 +0200 Subject: [PATCH 1/2] fix(server): exclude inherited Codex child usage (#5793) --- .../src/provider/Layers/CodexAdapter.test.ts | 235 ++++++++++++++++++ .../src/provider/Layers/CodexAdapter.ts | 111 +++++++-- 2 files changed, 330 insertions(+), 16 deletions(-) diff --git a/apps/server/src/provider/Layers/CodexAdapter.test.ts b/apps/server/src/provider/Layers/CodexAdapter.test.ts index 5358716aabe4..65e87ce3552f 100644 --- a/apps/server/src/provider/Layers/CodexAdapter.test.ts +++ b/apps/server/src/provider/Layers/CodexAdapter.test.ts @@ -15,6 +15,7 @@ import { type ProviderSession, type ProviderTurnStartResult, type ProviderUserInputAnswers, + type RuntimeTaskUsage, ThreadId, TurnId, } from "@t3tools/contracts"; @@ -1152,6 +1153,240 @@ lifecycleLayer("CodexAdapterLive lifecycle", (it) => { }), ); + it.effect("reports only child-generated usage for Codex collab agents", () => + Effect.gen(function* () { + const { adapter, runtime } = yield* startLifecycleRuntime(); + const eventsFiber = yield* adapter.streamEvents.pipe( + Stream.filter( + (event) => event.type === "thread.token-usage.updated" || event.type === "task.progress", + ), + Stream.take(7), + Stream.runCollect, + Effect.forkChild, + ); + + const parentBaseline = 5_800_000_000; + yield* runtime.emit({ + id: asEventId("evt-parent-usage-baseline"), + kind: "notification", + provider: ProviderDriverKind.make("codex"), + threadId: asThreadId("thread-1"), + turnId: asTurnId("turn-1"), + createdAt: "2026-01-01T00:00:00.000Z", + method: "thread/tokenUsage/updated", + payload: { + threadId: "thread-1", + turnId: "turn-1", + tokenUsage: { + total: { + inputTokens: parentBaseline - 100, + cachedInputTokens: 0, + outputTokens: 100, + reasoningOutputTokens: 0, + totalTokens: parentBaseline, + }, + last: { + inputTokens: 20, + cachedInputTokens: 0, + outputTokens: 5, + reasoningOutputTokens: 0, + totalTokens: 25, + }, + modelContextWindow: 258_400, + }, + }, + } satisfies ProviderEvent); + + const emitChildUsage = ( + eventId: string, + agentThreadId: string, + total: RuntimeTaskUsage, + last: RuntimeTaskUsage, + ) => + runtime.emit({ + id: asEventId(eventId), + kind: "notification", + provider: ProviderDriverKind.make("codex"), + threadId: asThreadId("thread-1"), + turnId: asTurnId("turn-1"), + createdAt: "2026-01-01T00:00:01.000Z", + method: "collabAgent/tokenUsage", + payload: { + agentThreadId, + tokenUsage: { + total, + last, + modelContextWindow: 258_400, + }, + }, + } satisfies ProviderEvent); + + const childAFirstTotal = { + totalTokens: 5_800_000_100, + inputTokens: 5_000_000_070, + cachedInputTokens: 3_000_000_040, + outputTokens: 800_000_030, + reasoningOutputTokens: 5, + } satisfies RuntimeTaskUsage; + const childAFirstLast = { + totalTokens: 100, + inputTokens: 70, + cachedInputTokens: 40, + outputTokens: 30, + reasoningOutputTokens: 5, + } satisfies RuntimeTaskUsage; + const childBZeroTotal = { + totalTokens: 5_800_050_000, + inputTokens: 5_000_040_000, + cachedInputTokens: 3_000_030_000, + outputTokens: 800_010_000, + reasoningOutputTokens: 0, + } satisfies RuntimeTaskUsage; + const childBZeroLast = { + totalTokens: 0, + inputTokens: 0, + cachedInputTokens: 0, + outputTokens: 0, + reasoningOutputTokens: 0, + } satisfies RuntimeTaskUsage; + const childBWorkTotal = { + totalTokens: 5_800_050_050, + inputTokens: 5_000_040_035, + cachedInputTokens: 3_000_030_020, + outputTokens: 800_010_015, + reasoningOutputTokens: 5, + } satisfies RuntimeTaskUsage; + const childBWorkLast = { + totalTokens: 50, + inputTokens: 35, + cachedInputTokens: 20, + outputTokens: 15, + reasoningOutputTokens: 5, + } satisfies RuntimeTaskUsage; + const childAResumeTotal = { + totalTokens: 5_800_000_130, + inputTokens: 5_000_000_090, + cachedInputTokens: 3_000_000_050, + outputTokens: 800_000_040, + reasoningOutputTokens: 10, + } satisfies RuntimeTaskUsage; + const childAResumeLast = { + totalTokens: 30, + inputTokens: 20, + cachedInputTokens: 10, + outputTokens: 10, + reasoningOutputTokens: 5, + } satisfies RuntimeTaskUsage; + const childARetryTotal = { + totalTokens: 5_800_000_150, + inputTokens: 5_000_000_100, + cachedInputTokens: 3_000_000_055, + outputTokens: 800_000_050, + reasoningOutputTokens: 15, + } satisfies RuntimeTaskUsage; + const childARetryLast = { + totalTokens: 20, + inputTokens: 10, + cachedInputTokens: 5, + outputTokens: 10, + reasoningOutputTokens: 5, + } satisfies RuntimeTaskUsage; + + yield* emitChildUsage("evt-child-a-first", "child-a", childAFirstTotal, childAFirstLast); + yield* emitChildUsage("evt-child-b-zero", "child-b", childBZeroTotal, childBZeroLast); + yield* emitChildUsage("evt-child-b-work", "child-b", childBWorkTotal, childBWorkLast); + yield* emitChildUsage("evt-child-a-resume", "child-a", childAResumeTotal, childAResumeLast); + yield* emitChildUsage("evt-child-a-retry", "child-a", childARetryTotal, childARetryLast); + yield* emitChildUsage("evt-child-a-duplicate", "child-a", childARetryTotal, childARetryLast); + + const events = Array.from(yield* Fiber.join(eventsFiber)); + const parentEvents = events.filter((event) => event.type === "thread.token-usage.updated"); + NodeAssert.equal(parentEvents.length, 1); + NodeAssert.equal(parentEvents[0]?.payload.usage.totalProcessedTokens, parentBaseline); + + const childUsage = events.flatMap((event) => { + if (event.type !== "task.progress") return []; + return [ + { + taskId: event.payload.taskId, + usage: event.payload.typedUsage, + }, + ]; + }); + NodeAssert.deepEqual(childUsage, [ + { + taskId: "child-a", + usage: { + totalTokens: 100, + inputTokens: 70, + cachedInputTokens: 40, + outputTokens: 30, + reasoningOutputTokens: 5, + }, + }, + { + taskId: "child-b", + usage: { + totalTokens: 0, + inputTokens: 0, + cachedInputTokens: 0, + outputTokens: 0, + reasoningOutputTokens: 0, + }, + }, + { + taskId: "child-b", + usage: { + totalTokens: 50, + inputTokens: 35, + cachedInputTokens: 20, + outputTokens: 15, + reasoningOutputTokens: 5, + }, + }, + { + taskId: "child-a", + usage: { + totalTokens: 130, + inputTokens: 90, + cachedInputTokens: 50, + outputTokens: 40, + reasoningOutputTokens: 10, + }, + }, + { + taskId: "child-a", + usage: { + totalTokens: 150, + inputTokens: 100, + cachedInputTokens: 55, + outputTokens: 50, + reasoningOutputTokens: 15, + }, + }, + { + taskId: "child-a", + usage: { + totalTokens: 150, + inputTokens: 100, + cachedInputTokens: 55, + outputTokens: 50, + reasoningOutputTokens: 15, + }, + }, + ]); + + const latestByChild = new Map(); + for (const child of childUsage) { + if (child.usage) latestByChild.set(child.taskId, child.usage); + } + NodeAssert.equal( + Array.from(latestByChild.values()).reduce((sum, usage) => sum + usage.totalTokens, 0), + 200, + ); + }), + ); + // Production calls startSession from a request fiber that finishes as soon as // the session exists. `Effect.forkChild` made the runtime event consumer a // child of that fiber, and Effect interrupts a fiber's children when it diff --git a/apps/server/src/provider/Layers/CodexAdapter.ts b/apps/server/src/provider/Layers/CodexAdapter.ts index 065156d36473..7acde3e0cf95 100644 --- a/apps/server/src/provider/Layers/CodexAdapter.ts +++ b/apps/server/src/provider/Layers/CodexAdapter.ts @@ -95,6 +95,11 @@ interface CodexAdapterSessionContext { stopped: boolean; } +interface CodexCollabUsageState { + readonly baseline: RuntimeTaskUsage; + readonly latest: RuntimeTaskUsage; +} + function mapCodexRuntimeError( threadId: ThreadId, method: string, @@ -507,6 +512,7 @@ function mapItemLifecycle( function mapCollabAgentEvent( event: ProviderEvent, canonicalThreadId: ThreadId, + usageByAgent: Map, ): ReadonlyArray { const payload = typeof event.payload === "object" && event.payload !== null @@ -667,8 +673,8 @@ function mapCollabAgentEvent( return []; } case "collabAgent/tokenUsage": { - // Cumulative per child thread: always the `total` breakdown, never - // `last` (which shrinks on follow-ups). Client folds max-merge. + // Calibrate the first cumulative `total` with `last` to remove copied + // history. Later snapshots stay based on `total`; `last` resets per turn. const tokenUsage = typeof payload.tokenUsage === "object" && payload.tokenUsage !== null ? (payload.tokenUsage as Record) @@ -677,29 +683,100 @@ function mapCollabAgentEvent( typeof tokenUsage?.total === "object" && tokenUsage.total !== null ? (tokenUsage.total as Record) : undefined; + const last = + typeof tokenUsage?.last === "object" && tokenUsage.last !== null + ? (tokenUsage.last as Record) + : undefined; const count = (value: unknown): number | undefined => typeof value === "number" && Number.isFinite(value) && value >= 0 ? value : undefined; // Same validation as every other field: RuntimeTaskUsage.totalTokens // is NonNegativeInt, so NaN/Infinity/negative wire values must miss. const totalTokens = count(total?.totalTokens); - if (totalTokens === undefined) { + const lastTotalTokens = count(last?.totalTokens); + if (totalTokens === undefined || lastTotalTokens === undefined) { return []; } + + const inputTokens = count(total?.inputTokens); + const cachedInputTokens = count(total?.cachedInputTokens); + const outputTokens = count(total?.outputTokens); + const reasoningOutputTokens = count(total?.reasoningOutputTokens); + const lastInputTokens = count(last?.inputTokens); + const lastCachedInputTokens = count(last?.cachedInputTokens); + const lastOutputTokens = count(last?.outputTokens); + const lastReasoningOutputTokens = count(last?.reasoningOutputTokens); + const baselineCount = ( + cumulative: number | undefined, + currentTurn: number | undefined, + ): number | undefined => + cumulative !== undefined && currentTurn !== undefined + ? Math.max(cumulative - currentTurn, 0) + : undefined; + const existing = usageByAgent.get(agentThreadId); + const baselineInputTokens = baselineCount(inputTokens, lastInputTokens); + const baselineCachedInputTokens = baselineCount(cachedInputTokens, lastCachedInputTokens); + const baselineOutputTokens = baselineCount(outputTokens, lastOutputTokens); + const baselineReasoningOutputTokens = baselineCount( + reasoningOutputTokens, + lastReasoningOutputTokens, + ); + const baseline = + existing?.baseline ?? + ({ + totalTokens: Math.max(totalTokens - lastTotalTokens, 0), + ...(baselineInputTokens !== undefined ? { inputTokens: baselineInputTokens } : {}), + ...(baselineCachedInputTokens !== undefined + ? { cachedInputTokens: baselineCachedInputTokens } + : {}), + ...(baselineOutputTokens !== undefined ? { outputTokens: baselineOutputTokens } : {}), + ...(baselineReasoningOutputTokens !== undefined + ? { reasoningOutputTokens: baselineReasoningOutputTokens } + : {}), + } satisfies RuntimeTaskUsage); + const normalizedCount = ( + cumulative: number | undefined, + offset: number | undefined, + latest: number | undefined, + ): number | undefined => + cumulative !== undefined && offset !== undefined + ? Math.max(cumulative - offset, latest ?? 0, 0) + : undefined; + const normalizedInputTokens = normalizedCount( + inputTokens, + baseline.inputTokens, + existing?.latest.inputTokens, + ); + const normalizedCachedInputTokens = normalizedCount( + cachedInputTokens, + baseline.cachedInputTokens, + existing?.latest.cachedInputTokens, + ); + const normalizedOutputTokens = normalizedCount( + outputTokens, + baseline.outputTokens, + existing?.latest.outputTokens, + ); + const normalizedReasoningOutputTokens = normalizedCount( + reasoningOutputTokens, + baseline.reasoningOutputTokens, + existing?.latest.reasoningOutputTokens, + ); const typedUsage: RuntimeTaskUsage = { - totalTokens, - ...(count(total?.inputTokens) !== undefined - ? { inputTokens: count(total?.inputTokens) } - : {}), - ...(count(total?.cachedInputTokens) !== undefined - ? { cachedInputTokens: count(total?.cachedInputTokens) } - : {}), - ...(count(total?.outputTokens) !== undefined - ? { outputTokens: count(total?.outputTokens) } + totalTokens: Math.max( + totalTokens - baseline.totalTokens, + existing?.latest.totalTokens ?? 0, + 0, + ), + ...(normalizedInputTokens !== undefined ? { inputTokens: normalizedInputTokens } : {}), + ...(normalizedCachedInputTokens !== undefined + ? { cachedInputTokens: normalizedCachedInputTokens } : {}), - ...(count(total?.reasoningOutputTokens) !== undefined - ? { reasoningOutputTokens: count(total?.reasoningOutputTokens) } + ...(normalizedOutputTokens !== undefined ? { outputTokens: normalizedOutputTokens } : {}), + ...(normalizedReasoningOutputTokens !== undefined + ? { reasoningOutputTokens: normalizedReasoningOutputTokens } : {}), }; + usageByAgent.set(agentThreadId, { baseline, latest: typedUsage }); return [ { ...base, @@ -762,9 +839,10 @@ function mapCollabAgentEvent( function mapToRuntimeEvents( event: ProviderEvent, canonicalThreadId: ThreadId, + collabUsageByAgent: Map, ): ReadonlyArray { if (event.kind === "notification" && event.method.startsWith("collabAgent/")) { - return mapCollabAgentEvent(event, canonicalThreadId); + return mapCollabAgentEvent(event, canonicalThreadId, collabUsageByAgent); } if (event.kind === "error") { if (!event.message) { @@ -1719,10 +1797,11 @@ export const makeCodexAdapter = Effect.fn("makeCodexAdapter")(function* ( // this a child of `startSession`, and Effect interrupts a fiber's // children when it completes, so the consumer died on return and every // runtime event the session emitted afterwards was dropped. + const collabUsageByAgent = new Map(); const eventFiber = yield* Stream.runForEach(runtime.events, (event) => Effect.gen(function* () { yield* writeNativeEvent(event); - const runtimeEvents = mapToRuntimeEvents(event, event.threadId); + const runtimeEvents = mapToRuntimeEvents(event, event.threadId, collabUsageByAgent); if (runtimeEvents.length === 0) { yield* Effect.logDebug("ignoring unhandled Codex provider event", { method: event.method, From f9b6ed453427146e876d994d8fda139759ec7cfd Mon Sep 17 00:00:00 2001 From: imMxts Date: Sun, 16 Aug 2026 18:45:54 +0200 Subject: [PATCH 2/2] fix(server): retain late Codex child usage fields (#5793) --- .../src/provider/Layers/CodexAdapter.test.ts | 10 ++-- .../src/provider/Layers/CodexAdapter.ts | 51 ++++++++++++------- 2 files changed, 36 insertions(+), 25 deletions(-) diff --git a/apps/server/src/provider/Layers/CodexAdapter.test.ts b/apps/server/src/provider/Layers/CodexAdapter.test.ts index 65e87ce3552f..1527b626c42e 100644 --- a/apps/server/src/provider/Layers/CodexAdapter.test.ts +++ b/apps/server/src/provider/Layers/CodexAdapter.test.ts @@ -1221,19 +1221,18 @@ lifecycleLayer("CodexAdapterLive lifecycle", (it) => { }, } satisfies ProviderEvent); + // Reasoning first appears on the resume frame. Its baseline must initialize then. const childAFirstTotal = { totalTokens: 5_800_000_100, inputTokens: 5_000_000_070, cachedInputTokens: 3_000_000_040, outputTokens: 800_000_030, - reasoningOutputTokens: 5, } satisfies RuntimeTaskUsage; const childAFirstLast = { totalTokens: 100, inputTokens: 70, cachedInputTokens: 40, outputTokens: 30, - reasoningOutputTokens: 5, } satisfies RuntimeTaskUsage; const childBZeroTotal = { totalTokens: 5_800_050_000, @@ -1321,7 +1320,6 @@ lifecycleLayer("CodexAdapterLive lifecycle", (it) => { inputTokens: 70, cachedInputTokens: 40, outputTokens: 30, - reasoningOutputTokens: 5, }, }, { @@ -1351,7 +1349,7 @@ lifecycleLayer("CodexAdapterLive lifecycle", (it) => { inputTokens: 90, cachedInputTokens: 50, outputTokens: 40, - reasoningOutputTokens: 10, + reasoningOutputTokens: 5, }, }, { @@ -1361,7 +1359,7 @@ lifecycleLayer("CodexAdapterLive lifecycle", (it) => { inputTokens: 100, cachedInputTokens: 55, outputTokens: 50, - reasoningOutputTokens: 15, + reasoningOutputTokens: 10, }, }, { @@ -1371,7 +1369,7 @@ lifecycleLayer("CodexAdapterLive lifecycle", (it) => { inputTokens: 100, cachedInputTokens: 55, outputTokens: 50, - reasoningOutputTokens: 15, + reasoningOutputTokens: 10, }, }, ]); diff --git a/apps/server/src/provider/Layers/CodexAdapter.ts b/apps/server/src/provider/Layers/CodexAdapter.ts index 7acde3e0cf95..b9a928783164 100644 --- a/apps/server/src/provider/Layers/CodexAdapter.ts +++ b/apps/server/src/provider/Layers/CodexAdapter.ts @@ -705,34 +705,47 @@ function mapCollabAgentEvent( const lastCachedInputTokens = count(last?.cachedInputTokens); const lastOutputTokens = count(last?.outputTokens); const lastReasoningOutputTokens = count(last?.reasoningOutputTokens); + const existing = usageByAgent.get(agentThreadId); const baselineCount = ( + offset: number | undefined, cumulative: number | undefined, currentTurn: number | undefined, ): number | undefined => - cumulative !== undefined && currentTurn !== undefined + offset ?? + (cumulative !== undefined && currentTurn !== undefined ? Math.max(cumulative - currentTurn, 0) - : undefined; - const existing = usageByAgent.get(agentThreadId); - const baselineInputTokens = baselineCount(inputTokens, lastInputTokens); - const baselineCachedInputTokens = baselineCount(cachedInputTokens, lastCachedInputTokens); - const baselineOutputTokens = baselineCount(outputTokens, lastOutputTokens); + : undefined); + const baselineInputTokens = baselineCount( + existing?.baseline.inputTokens, + inputTokens, + lastInputTokens, + ); + const baselineCachedInputTokens = baselineCount( + existing?.baseline.cachedInputTokens, + cachedInputTokens, + lastCachedInputTokens, + ); + const baselineOutputTokens = baselineCount( + existing?.baseline.outputTokens, + outputTokens, + lastOutputTokens, + ); const baselineReasoningOutputTokens = baselineCount( + existing?.baseline.reasoningOutputTokens, reasoningOutputTokens, lastReasoningOutputTokens, ); - const baseline = - existing?.baseline ?? - ({ - totalTokens: Math.max(totalTokens - lastTotalTokens, 0), - ...(baselineInputTokens !== undefined ? { inputTokens: baselineInputTokens } : {}), - ...(baselineCachedInputTokens !== undefined - ? { cachedInputTokens: baselineCachedInputTokens } - : {}), - ...(baselineOutputTokens !== undefined ? { outputTokens: baselineOutputTokens } : {}), - ...(baselineReasoningOutputTokens !== undefined - ? { reasoningOutputTokens: baselineReasoningOutputTokens } - : {}), - } satisfies RuntimeTaskUsage); + const baseline = { + totalTokens: existing?.baseline.totalTokens ?? Math.max(totalTokens - lastTotalTokens, 0), + ...(baselineInputTokens !== undefined ? { inputTokens: baselineInputTokens } : {}), + ...(baselineCachedInputTokens !== undefined + ? { cachedInputTokens: baselineCachedInputTokens } + : {}), + ...(baselineOutputTokens !== undefined ? { outputTokens: baselineOutputTokens } : {}), + ...(baselineReasoningOutputTokens !== undefined + ? { reasoningOutputTokens: baselineReasoningOutputTokens } + : {}), + } satisfies RuntimeTaskUsage; const normalizedCount = ( cumulative: number | undefined, offset: number | undefined,