From a8de898205af7d566c4522d7c23e6c73a6bf8b35 Mon Sep 17 00:00:00 2001 From: Guillermo Casanova Date: Mon, 14 Sep 2026 23:03:36 -0300 Subject: [PATCH] fix(codex): preserve child answers and tool counts --- apps/server/src/persistence/Migrations.ts | 2 + .../Migrations/052_CodexChildUsage.ts | 28 ++++++ .../src/provider/Drivers/CodexDriver.test.ts | 2 + .../src/provider/Drivers/CodexDriver.ts | 4 +- .../src/provider/Layers/CodexAdapter.test.ts | 93 +++++++++++++++++++ .../src/provider/Layers/CodexAdapter.ts | 79 +++++++++++++--- .../provider/Layers/CodexSessionRuntime.ts | 1 + .../ProviderInstanceRegistryLive.test.ts | 3 + .../provider/Layers/ProviderRegistry.test.ts | 5 + .../provider/Layers/codexChildUsage.test.ts | 64 +++++++++++++ .../src/provider/Layers/codexChildUsage.ts | 65 +++++++++++++ 11 files changed, 332 insertions(+), 14 deletions(-) create mode 100644 apps/server/src/persistence/Migrations/052_CodexChildUsage.ts create mode 100644 apps/server/src/provider/Layers/codexChildUsage.test.ts create mode 100644 apps/server/src/provider/Layers/codexChildUsage.ts diff --git a/apps/server/src/persistence/Migrations.ts b/apps/server/src/persistence/Migrations.ts index 4adb98cc60f4..607a0b5696a7 100644 --- a/apps/server/src/persistence/Migrations.ts +++ b/apps/server/src/persistence/Migrations.ts @@ -63,6 +63,7 @@ import Migration0048 from "./Migrations/048_ProjectionThreadBranchPullRequest.ts import Migration0049 from "./Migrations/049_ProjectionThreadsActiveOrderKey.ts"; import Migration0050 from "./Migrations/050_ProjectionThreadPullRequests.ts"; import Migration0051 from "./Migrations/051_ProjectionThreadMessageContext.ts"; +import Migration0052 from "./Migrations/052_CodexChildUsage.ts"; /** * Migration loader with all migrations defined inline. @@ -126,6 +127,7 @@ const migrationEntries = [ [49, "ProjectionThreadsActiveOrderKey", Migration0049], [50, "ProjectionThreadPullRequests", Migration0050], [51, "ProjectionThreadMessageContext", Migration0051], + [52, "CodexChildUsage", Migration0052], ] as const; export const migrationManifest = migrationEntries.map(([id, name]) => [id, name] as const); diff --git a/apps/server/src/persistence/Migrations/052_CodexChildUsage.ts b/apps/server/src/persistence/Migrations/052_CodexChildUsage.ts new file mode 100644 index 000000000000..cd8a12679cd7 --- /dev/null +++ b/apps/server/src/persistence/Migrations/052_CodexChildUsage.ts @@ -0,0 +1,28 @@ +import * as Effect from "effect/Effect"; +import * as SqlClient from "effect/unstable/sql/SqlClient"; + +export default Effect.gen(function* () { + const sql = yield* SqlClient.SqlClient; + yield* sql` + CREATE TABLE codex_child_usage ( + thread_id TEXT NOT NULL, + instance_id TEXT NOT NULL, + task_id TEXT NOT NULL, + usage_json TEXT NOT NULL DEFAULT '{"totalTokens":0}', + tool_uses INTEGER NOT NULL DEFAULT 0, + PRIMARY KEY (thread_id, instance_id, task_id) + ) + `; + yield* sql` + CREATE TABLE codex_child_tool_calls ( + thread_id TEXT NOT NULL, + instance_id TEXT NOT NULL, + task_id TEXT NOT NULL, + turn_id TEXT NOT NULL, + item_id TEXT NOT NULL, + PRIMARY KEY (thread_id, instance_id, task_id, turn_id, item_id), + FOREIGN KEY (thread_id, instance_id, task_id) + REFERENCES codex_child_usage (thread_id, instance_id, task_id) ON DELETE CASCADE + ) + `; +}); diff --git a/apps/server/src/provider/Drivers/CodexDriver.test.ts b/apps/server/src/provider/Drivers/CodexDriver.test.ts index bac34db452fd..fe4083e0afd6 100644 --- a/apps/server/src/provider/Drivers/CodexDriver.test.ts +++ b/apps/server/src/provider/Drivers/CodexDriver.test.ts @@ -16,6 +16,7 @@ import * as ChildProcessSpawner from "effect/unstable/process/ChildProcessSpawne import * as BackgroundPolicy from "../../background/BackgroundPolicy.ts"; import { ServerConfig } from "../../config.ts"; +import { SqlitePersistenceMemory } from "../../persistence/Layers/Sqlite.ts"; import { ServerSettingsService } from "../../serverSettings.ts"; import { layerTest as codexResetCreditLayerTest } from "../Layers/codexResetCredit.ts"; import { NoOpProviderEventLoggers, ProviderEventLoggers } from "../Layers/ProviderEventLoggers.ts"; @@ -31,6 +32,7 @@ const testLayer = ServerConfig.layerTest(process.cwd(), { prefix: "t3-codex-driver-maintenance-", }).pipe( Layer.provideMerge(NodeServices.layer), + Layer.provideMerge(SqlitePersistenceMemory), Layer.provideMerge(ServerSettingsService.layerTest()), Layer.provideMerge(ModelManifest.layerTest), Layer.provideMerge(codexResetCreditLayerTest), diff --git a/apps/server/src/provider/Drivers/CodexDriver.ts b/apps/server/src/provider/Drivers/CodexDriver.ts index 22dd047f5c69..9835cdec9879 100644 --- a/apps/server/src/provider/Drivers/CodexDriver.ts +++ b/apps/server/src/provider/Drivers/CodexDriver.ts @@ -27,6 +27,7 @@ import * as Effect from "effect/Effect"; import * as FileSystem from "effect/FileSystem"; import * as Path from "effect/Path"; import * as Schema from "effect/Schema"; +import type * as SqlClient from "effect/unstable/sql/SqlClient"; import { HttpClient } from "effect/unstable/http"; import { ChildProcessSpawner } from "effect/unstable/process"; @@ -114,7 +115,8 @@ export type CodexDriverEnv = | Path.Path | ProviderEventLoggers | ServerConfig - | ServerSettingsService; + | ServerSettingsService + | SqlClient.SqlClient; export const CodexDriver: ProviderDriver = { driverKind: DRIVER_KIND, diff --git a/apps/server/src/provider/Layers/CodexAdapter.test.ts b/apps/server/src/provider/Layers/CodexAdapter.test.ts index f7c6036885d9..2a552ef7dc4b 100644 --- a/apps/server/src/provider/Layers/CodexAdapter.test.ts +++ b/apps/server/src/provider/Layers/CodexAdapter.test.ts @@ -36,6 +36,7 @@ import * as TestClock from "effect/testing/TestClock"; import * as CodexErrors from "effect-codex-app-server/errors"; import { ServerConfig } from "../../config.ts"; +import { SqlitePersistenceMemory } from "../../persistence/Layers/Sqlite.ts"; import { ServerSettingsService } from "../../serverSettings.ts"; import { ProviderAdapterValidationError } from "../Errors.ts"; import type { CodexAdapterShape } from "../Services/CodexAdapter.ts"; @@ -244,6 +245,7 @@ const validationLayer = it.layer( Layer.provideMerge(ServerConfig.layerTest(process.cwd(), process.cwd())), Layer.provideMerge(ServerSettingsService.layerTest()), Layer.provideMerge(providerSessionDirectoryTestLayer), + Layer.provideMerge(SqlitePersistenceMemory), Layer.provideMerge(NodeServices.layer), ), ); @@ -314,6 +316,7 @@ const sessionErrorLayer = it.layer( Layer.provideMerge(ServerConfig.layerTest(process.cwd(), process.cwd())), Layer.provideMerge(ServerSettingsService.layerTest()), Layer.provideMerge(providerSessionDirectoryTestLayer), + Layer.provideMerge(SqlitePersistenceMemory), Layer.provideMerge(NodeServices.layer), ), ); @@ -462,6 +465,7 @@ sessionErrorLayer("CodexAdapterLive session errors", (it) => { Layer.provideMerge(ServerConfig.layerTest(process.cwd(), process.cwd())), Layer.provideMerge(ServerSettingsService.layerTest()), Layer.provideMerge(providerSessionDirectoryTestLayer), + Layer.provideMerge(SqlitePersistenceMemory), Layer.provideMerge(NodeServices.layer), ); @@ -494,6 +498,7 @@ sessionErrorLayer("CodexAdapterLive session errors", (it) => { Layer.provideMerge(ServerConfig.layerTest(process.cwd(), process.cwd())), Layer.provideMerge(ServerSettingsService.layerTest()), Layer.provideMerge(providerSessionDirectoryTestLayer), + Layer.provideMerge(SqlitePersistenceMemory), Layer.provideMerge(NodeServices.layer), ); @@ -527,6 +532,7 @@ sessionErrorLayer("CodexAdapterLive session errors", (it) => { Layer.provideMerge(ServerConfig.layerTest(process.cwd(), process.cwd())), Layer.provideMerge(ServerSettingsService.layerTest()), Layer.provideMerge(providerSessionDirectoryTestLayer), + Layer.provideMerge(SqlitePersistenceMemory), Layer.provideMerge(NodeServices.layer), ); @@ -581,6 +587,7 @@ const lifecycleLayer = it.layer( Layer.provideMerge(ServerConfig.layerTest(process.cwd(), process.cwd())), Layer.provideMerge(ServerSettingsService.layerTest()), Layer.provideMerge(providerSessionDirectoryTestLayer), + Layer.provideMerge(SqlitePersistenceMemory), Layer.provideMerge(NodeServices.layer), ), ); @@ -1046,6 +1053,88 @@ lifecycleLayer("CodexAdapterLive lifecycle", (it) => { }), ); + it.effect("preserves child answers and counts each tool once", () => + Effect.gen(function* () { + const { adapter, runtime } = yield* startLifecycleRuntime(); + const eventsFiber = yield* Stream.runCollect(Stream.take(adapter.streamEvents, 5)).pipe( + Effect.forkChild, + ); + const items = [ + { type: "commandExecution", id: "tool-1", command: "cat README.md" }, + { type: "commandExecution", id: "tool-1", command: "cat README.md" }, + { type: "agentMessage", id: "answer-1", text: "" }, + { type: "agentMessage", id: "answer-1", text: "The project is T3 Code." }, + ]; + for (const [index, item] of items.entries()) { + yield* runtime.emit({ + id: asEventId(`evt-child-item-${index}`), + kind: "notification", + provider: ProviderDriverKind.make("codex"), + createdAt: "2026-01-01T00:00:00.000Z", + method: "collabAgent/item", + threadId: asThreadId("thread-1"), + payload: { agentThreadId: "child-answer", item }, + }); + } + for (const [index, extra] of [ + { method: "collabAgent/tokenUsage", tokenUsage: { total: { totalTokens: 42 } } }, + { method: "collabAgent/turnCompleted", turn: { status: "completed" } }, + ].entries()) { + yield* runtime.emit({ + id: asEventId(`evt-child-state-${index}`), + kind: "notification", + provider: ProviderDriverKind.make("codex"), + createdAt: "2026-01-01T00:00:01.000Z", + method: extra.method, + threadId: asThreadId("thread-1"), + payload: { agentThreadId: "child-answer", ...extra }, + }); + } + const events = Array.from(yield* Fiber.join(eventsFiber)); + const progress = events.filter((event) => event.type === "task.progress"); + NodeAssert.deepStrictEqual( + { summary: progress[2]?.payload.summary, usage: progress[3]?.payload.typedUsage }, + { summary: "The project is T3 Code.", usage: { totalTokens: 42, toolUses: 1 } }, + ); + NodeAssert.equal(events[4]?.type === "task.updated" && events[4].payload.status, "idle"); + yield* adapter.stopSession(asThreadId("thread-1")); + const restarted = yield* startLifecycleRuntime(); + const resumedFiber = yield* Stream.runCollect(Stream.take(adapter.streamEvents, 3)).pipe( + Effect.forkChild, + ); + for (const [index, extra] of [ + { method: "collabAgent/item", item: items[0] }, + { + method: "collabAgent/item", + item: { type: "mcpToolCall", id: "tool-2", server: "filesystem", tool: "read_file" }, + }, + { + method: "collabAgent/tokenUsage", + tokenUsage: { total: { totalTokens: 20, inputTokens: 10 } }, + }, + ].entries()) { + yield* restarted.runtime.emit({ + id: asEventId(`evt-child-resumed-${index}`), + kind: "notification", + provider: ProviderDriverKind.make("codex"), + createdAt: "2026-01-01T00:00:02.000Z", + method: extra.method, + threadId: asThreadId("thread-1"), + payload: { agentThreadId: "child-answer", ...extra }, + }); + } + const resumed = Array.from(yield* Fiber.join(resumedFiber)); + NodeAssert.deepStrictEqual( + resumed.map((event) => (event.type === "task.progress" ? event.payload.typedUsage : null)), + [ + { totalTokens: 42, toolUses: 1 }, + { totalTokens: 42, toolUses: 2 }, + { totalTokens: 42, inputTokens: 10, toolUses: 2 }, + ], + ); + }), + ); + it.effect("carries child model metadata through every task event", () => Effect.gen(function* () { const { adapter, runtime } = yield* startLifecycleRuntime(); @@ -2557,6 +2646,7 @@ const scopedLifecycleLayer = it.layer( Layer.provideMerge(ServerConfig.layerTest(process.cwd(), process.cwd())), Layer.provideMerge(ServerSettingsService.layerTest()), Layer.provideMerge(providerSessionDirectoryTestLayer), + Layer.provideMerge(SqlitePersistenceMemory), Layer.provideMerge(NodeServices.layer), ), ); @@ -2601,6 +2691,7 @@ const scopedFailureLayer = it.layer( Layer.provideMerge(ServerConfig.layerTest(process.cwd(), process.cwd())), Layer.provideMerge(ServerSettingsService.layerTest()), Layer.provideMerge(providerSessionDirectoryTestLayer), + Layer.provideMerge(SqlitePersistenceMemory), Layer.provideMerge(NodeServices.layer), ), ); @@ -2653,6 +2744,7 @@ it.effect("flushes managed native logs when the adapter layer shuts down", () => Layer.provideMerge(ServerConfig.layerTest(process.cwd(), process.cwd())), Layer.provideMerge(ServerSettingsService.layerTest()), Layer.provideMerge(providerSessionDirectoryTestLayer), + Layer.provideMerge(SqlitePersistenceMemory), Layer.provideMerge(NodeServices.layer), ); const context = yield* Layer.buildWithScope(layer, scope); @@ -2709,6 +2801,7 @@ const usageLimitLayer = it.layer( Layer.provideMerge(ServerConfig.layerTest(process.cwd(), process.cwd())), Layer.provideMerge(ServerSettingsService.layerTest()), Layer.provideMerge(providerSessionDirectoryTestLayer), + Layer.provideMerge(SqlitePersistenceMemory), Layer.provideMerge(NodeServices.layer), ), ); diff --git a/apps/server/src/provider/Layers/CodexAdapter.ts b/apps/server/src/provider/Layers/CodexAdapter.ts index b43755736ca3..49638c01e859 100644 --- a/apps/server/src/provider/Layers/CodexAdapter.ts +++ b/apps/server/src/provider/Layers/CodexAdapter.ts @@ -38,6 +38,7 @@ import * as Exit from "effect/Exit"; import * as Fiber from "effect/Fiber"; import * as FileSystem from "effect/FileSystem"; import * as Queue from "effect/Queue"; +import * as SqlClient from "effect/unstable/sql/SqlClient"; import * as Schema from "effect/Schema"; import * as Scope from "effect/Scope"; import * as Stream from "effect/Stream"; @@ -71,6 +72,7 @@ import { } from "./CodexSessionRuntime.ts"; import { type EventNdjsonLogger, makeEventNdjsonLogger } from "./EventNdjsonLogger.ts"; import { resolveCodexLaunchArgs } from "./codexLaunchArgs.ts"; +import { recordCodexChildUsage } from "./codexChildUsage.ts"; import { type CodexRateLimitSnapshot, codexRateLimitsToUpdate, @@ -768,8 +770,11 @@ function itemTitle( } } -function itemDetail(itemType: CanonicalItemType, item: CodexLifecycleItem): string | undefined { - const itemRecord = item as Record; +function itemDetail( + itemType: CanonicalItemType, + item: Record, +): string | undefined { + const itemRecord = item; const action = itemRecord.action as Record | undefined; const actionQueries = Array.isArray(action?.queries) ? action.queries : []; const candidates = [ @@ -1218,7 +1223,7 @@ function mapCollabAgentEvent( ? (tokenUsage.total as Record) : undefined; const count = (value: unknown): number | undefined => - typeof value === "number" && Number.isFinite(value) && value >= 0 ? value : undefined; + typeof value === "number" && Number.isSafeInteger(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); @@ -1259,18 +1264,13 @@ function mapCollabAgentEvent( ? (payload.item as Record) : undefined; const itemTypeRaw = typeof item?.type === "string" ? item.type : undefined; - if (!itemTypeRaw) { + if (!item || !itemTypeRaw) { return []; } - // A loose summary from the raw item: the child stream is untyped at - // this boundary (synthetic event payload), so read best-effort fields - // rather than force a schema decode. - const looseSummary = - (typeof item?.command === "string" ? item.command : undefined) ?? - (typeof item?.title === "string" ? item.title : undefined) ?? - (typeof item?.query === "string" ? item.query : undefined); const canonical = toCanonicalItemType(itemTypeRaw); - const summary = looseSummary ?? canonical.replaceAll("_", " "); + const detail = itemDetail(canonical, item); + if (canonical === "assistant_message" && !detail) return []; + const summary = detail ?? itemTitle(canonical) ?? canonical.replaceAll("_", " "); return [ { ...base, @@ -2218,6 +2218,7 @@ export const makeCodexAdapter = Effect.fn("makeCodexAdapter")(function* ( options?: CodexAdapterLiveOptions, ) { const boundInstanceId = options?.instanceId ?? ProviderInstanceId.make("codex"); + const sql = yield* SqlClient.SqlClient; const fileSystem = yield* FileSystem.FileSystem; const childProcessSpawner = yield* ChildProcessSpawner.ChildProcessSpawner; const crypto = yield* Crypto.Crypto; @@ -2427,9 +2428,61 @@ export const makeCodexAdapter = Effect.fn("makeCodexAdapter")(function* ( } return runtimeEvent; }); + for (const [index, mapped] of mappedEvents.entries()) { + if (mapped.type !== "task.progress" || !event.method.startsWith("collabAgent/")) + continue; + const payload = asUnknownRecord(event.payload); + const item = asUnknownRecord(payload?.item); + const toolItemId = + typeof item?.id === "string" && + [ + "commandExecution", + "fileChange", + "mcpToolCall", + "dynamicToolCall", + "webSearch", + "imageView", + "imageGeneration", + "collabAgentToolCall", + ].includes(String(item.type)) + ? item.id + : undefined; + if (!toolItemId && !mapped.payload.typedUsage) continue; + const usage = yield* recordCodexChildUsage(sql, { + threadId: event.threadId, + instanceId: boundInstanceId, + taskId: mapped.payload.taskId, + ...(toolItemId + ? { + toolItemId, + childTurnId: + typeof payload?.childTurnId === "string" ? payload.childTurnId : "", + } + : {}), + ...(mapped.payload.typedUsage ? { usage: mapped.payload.typedUsage } : {}), + }).pipe( + Effect.catch((cause) => + Effect.logWarning("Could not persist Codex child usage", { cause }).pipe( + Effect.as(undefined), + ), + ), + ); + // Never replace the durable usage row with a partial snapshot on failure. + const progress = { ...mapped.payload }; + delete progress.typedUsage; + mappedEvents[index] = { + ...mapped, + payload: usage ? { ...progress, typedUsage: usage } : progress, + }; + } const runtimeEvents = usageLimitError ? [usageLimitError, ...mappedEvents] - : mappedEvents; + : mappedEvents.filter( + (mapped) => + event.method !== "collabAgent/tokenUsage" || + mapped.type !== "task.progress" || + mapped.payload.typedUsage !== undefined, + ); if (runtimeEvents.length === 0) { yield* Effect.logDebug("ignoring unhandled Codex provider event", { method: event.method, diff --git a/apps/server/src/provider/Layers/CodexSessionRuntime.ts b/apps/server/src/provider/Layers/CodexSessionRuntime.ts index a04db1405912..f88e392d60c3 100644 --- a/apps/server/src/provider/Layers/CodexSessionRuntime.ts +++ b/apps/server/src/provider/Layers/CodexSessionRuntime.ts @@ -1796,6 +1796,7 @@ export const makeCodexSessionRuntime = ( payload: { ...childIdentity, item: notification.params.item, + childTurnId: notification.params.turnId, }, }); return true; diff --git a/apps/server/src/provider/Layers/ProviderInstanceRegistryLive.test.ts b/apps/server/src/provider/Layers/ProviderInstanceRegistryLive.test.ts index 25dafa5ba040..ef06688a3dc7 100644 --- a/apps/server/src/provider/Layers/ProviderInstanceRegistryLive.test.ts +++ b/apps/server/src/provider/Layers/ProviderInstanceRegistryLive.test.ts @@ -47,6 +47,7 @@ import * as BackgroundPolicy from "../../background/BackgroundPolicy.ts"; import type { BuiltInDriversEnv } from "../builtInDrivers.ts"; import { AntigravityInstallation } from "../AntigravityInstallation.ts"; import { ServerConfig } from "../../config.ts"; +import { SqlitePersistenceMemory } from "../../persistence/Layers/Sqlite.ts"; import { expandHomePath } from "../../pathExpansion.ts"; import { ServerSettingsService } from "../../serverSettings.ts"; import { ClaudeDriver } from "../Drivers/ClaudeDriver.ts"; @@ -226,6 +227,7 @@ describe("ProviderInstanceRegistryLive — multi-instance codex slice", () => { prefix: "provider-instance-registry-test", }).pipe( Layer.provideMerge(NodeServices.layer), + Layer.provideMerge(SqlitePersistenceMemory), Layer.provideMerge(BackgroundPolicyAlwaysRunLayer), Layer.provideMerge(ServerSettingsService.layerTest()), Layer.provideMerge(TestHttpClientLive), @@ -454,6 +456,7 @@ describe("ProviderInstanceRegistryLive — all drivers slice", () => { }), ), Layer.provideMerge(infraLayer), + Layer.provideMerge(SqlitePersistenceMemory), Layer.provideMerge(BackgroundPolicyAlwaysRunLayer), Layer.provideMerge(ServerSettingsService.layerTest()), Layer.provideMerge(TestHttpClientLive), diff --git a/apps/server/src/provider/Layers/ProviderRegistry.test.ts b/apps/server/src/provider/Layers/ProviderRegistry.test.ts index 988c89e1e679..210f5f58da83 100644 --- a/apps/server/src/provider/Layers/ProviderRegistry.test.ts +++ b/apps/server/src/provider/Layers/ProviderRegistry.test.ts @@ -48,6 +48,7 @@ import { ProviderRegistryLive, } from "./ProviderRegistry.ts"; import * as ServerConfig from "../../config.ts"; +import { SqlitePersistenceMemory } from "../../persistence/Layers/Sqlite.ts"; import * as ServerSettingsModule from "../../serverSettings.ts"; import { readProviderStatusCache, @@ -2290,6 +2291,7 @@ it.layer(Layer.mergeAll(NodeServices.layer, ServerSettingsModule.layerTest(), Te ), Layer.provideMerge(ModelManifest.layerTest), Layer.provideMerge(CodexResetCredit.layerTest), + Layer.provideMerge(SqlitePersistenceMemory), Layer.provideMerge(OpenCodeRuntime.OpenCodeRuntimeLive), Layer.provideMerge(BackgroundPolicyAlwaysRunLayer), // NO spawner mock — `ChildProcessSpawner` is supplied by the @@ -2391,6 +2393,7 @@ it.layer(Layer.mergeAll(NodeServices.layer, ServerSettingsModule.layerTest(), Te ), Layer.provideMerge(ModelManifest.layerTest), Layer.provideMerge(CodexResetCredit.layerTest), + Layer.provideMerge(SqlitePersistenceMemory), Layer.provideMerge(OpenCodeRuntime.OpenCodeRuntimeLive), Layer.updateService(ChildProcessSpawner.ChildProcessSpawner, (spawner) => ChildProcessSpawner.make((command) => { @@ -2507,6 +2510,7 @@ it.layer(Layer.mergeAll(NodeServices.layer, ServerSettingsModule.layerTest(), Te ), Layer.provideMerge(ModelManifest.layerTest), Layer.provideMerge(CodexResetCredit.layerTest), + Layer.provideMerge(SqlitePersistenceMemory), Layer.provideMerge(OpenCodeRuntime.OpenCodeRuntimeLive), Layer.provideMerge(NodeServices.layer), Layer.provideMerge(BackgroundPolicyAlwaysRunLayer), @@ -2569,6 +2573,7 @@ it.layer(Layer.mergeAll(NodeServices.layer, ServerSettingsModule.layerTest(), Te ), Layer.provideMerge(ModelManifest.layerTest), Layer.provideMerge(CodexResetCredit.layerTest), + Layer.provideMerge(SqlitePersistenceMemory), Layer.provideMerge(CodexResetCredit.layerTest), Layer.provideMerge(OpenCodeRuntime.OpenCodeRuntimeLive), Layer.provideMerge(BackgroundPolicyAlwaysRunLayer), diff --git a/apps/server/src/provider/Layers/codexChildUsage.test.ts b/apps/server/src/provider/Layers/codexChildUsage.test.ts new file mode 100644 index 000000000000..c1c8ffa99d53 --- /dev/null +++ b/apps/server/src/provider/Layers/codexChildUsage.test.ts @@ -0,0 +1,64 @@ +import { it, assert } from "@effect/vitest"; +import * as NodeServices from "@effect/platform-node/NodeServices"; +import { ProviderInstanceId, RuntimeTaskId, ThreadId } from "@t3tools/contracts"; +import * as Effect from "effect/Effect"; +import * as FileSystem from "effect/FileSystem"; +import * as Layer from "effect/Layer"; +import * as SqlClient from "effect/unstable/sql/SqlClient"; + +import { makeSqlitePersistenceLive } from "../../persistence/Layers/Sqlite.ts"; +import { recordCodexChildUsage } from "./codexChildUsage.ts"; + +it.layer(NodeServices.layer)("Codex child usage", (it) => { + it.effect("deduplicates replay after reopening the database and isolates child identities", () => + Effect.gen(function* () { + const fs = yield* FileSystem.FileSystem; + const directory = yield* fs.makeTempDirectoryScoped({ prefix: "codex-child-usage-" }); + const database = makeSqlitePersistenceLive(`${directory}/state.sqlite`).pipe( + Layer.provide(NodeServices.layer), + ); + const input = { + threadId: ThreadId.make("parent"), + instanceId: ProviderInstanceId.make("codex"), + taskId: RuntimeTaskId.make("child"), + childTurnId: "turn-1", + toolItemId: "tool-1", + }; + // Each call closes its SQLite connection, so no process-local state can carry the count. + const record = (patch: Partial[1]> = {}) => + Effect.gen(function* () { + const sql = yield* SqlClient.SqlClient; + return yield* recordCodexChildUsage(sql, { ...input, ...patch }); + }).pipe(Effect.provide(database)); + + assert.deepEqual( + yield* record({ usage: { totalTokens: 100, inputTokens: 80, outputTokens: 20 } }), + { + totalTokens: 100, + inputTokens: 80, + outputTokens: 20, + toolUses: 1, + }, + ); + assert.deepEqual(yield* record(), { + totalTokens: 100, + inputTokens: 80, + outputTokens: 20, + toolUses: 1, + }); + assert.deepEqual(yield* record({ childTurnId: "turn-2", usage: { totalTokens: 50 } }), { + totalTokens: 100, + inputTokens: 80, + outputTokens: 20, + toolUses: 2, + }); + for (const patch of [ + { taskId: RuntimeTaskId.make("other-child") }, + { threadId: ThreadId.make("other-parent") }, + { instanceId: ProviderInstanceId.make("other-instance") }, + ]) { + assert.deepEqual(yield* record(patch), { totalTokens: 0, toolUses: 1 }); + } + }), + ); +}); diff --git a/apps/server/src/provider/Layers/codexChildUsage.ts b/apps/server/src/provider/Layers/codexChildUsage.ts new file mode 100644 index 000000000000..aa04920c8ff6 --- /dev/null +++ b/apps/server/src/provider/Layers/codexChildUsage.ts @@ -0,0 +1,65 @@ +import { + RuntimeTaskUsage, + type ThreadId, + type ProviderInstanceId, + type RuntimeTaskId, +} from "@t3tools/contracts"; +import * as Effect from "effect/Effect"; +import * as Schema from "effect/Schema"; +import type * as SqlClient from "effect/unstable/sql/SqlClient"; + +const usageJson = Schema.fromJsonString(RuntimeTaskUsage); +const decodeUsage = Schema.decodeUnknownEffect(usageJson); +const encodeUsage = Schema.encodeEffect(usageJson); + +export const recordCodexChildUsage = Effect.fn("recordCodexChildUsage")(function* ( + sql: SqlClient.SqlClient, + input: { + threadId: ThreadId; + instanceId: ProviderInstanceId; + taskId: RuntimeTaskId; + toolItemId?: string; + childTurnId?: string; + usage?: RuntimeTaskUsage; + }, +) { + const { threadId, instanceId, taskId } = input; + return yield* sql.withTransaction( + Effect.gen(function* () { + yield* sql` + INSERT INTO codex_child_usage (thread_id, instance_id, task_id) + VALUES (${threadId}, ${instanceId}, ${taskId}) ON CONFLICT DO NOTHING + `; + const inserted = + input.toolItemId === undefined + ? [] + : yield* sql` + INSERT INTO codex_child_tool_calls (thread_id, instance_id, task_id, turn_id, item_id) + VALUES (${threadId}, ${instanceId}, ${taskId}, ${input.childTurnId ?? ""}, ${input.toolItemId}) + ON CONFLICT DO NOTHING RETURNING item_id + `; + const [previous] = yield* sql<{ usageJson: string }>` + SELECT usage_json AS "usageJson" FROM codex_child_usage + WHERE thread_id = ${threadId} AND instance_id = ${instanceId} AND task_id = ${taskId} + `; + const usage = { ...(yield* decodeUsage(previous!.usageJson)) }; + for (const key of [ + "totalTokens", + "inputTokens", + "cachedInputTokens", + "outputTokens", + "reasoningOutputTokens", + ] as const) { + const value = input.usage?.[key]; + if (value !== undefined) usage[key] = Math.max(usage[key] ?? 0, value); + } + const [updated] = yield* sql<{ toolUses: number }>` + UPDATE codex_child_usage + SET usage_json = ${yield* encodeUsage(usage)}, tool_uses = tool_uses + ${inserted.length} + WHERE thread_id = ${threadId} AND instance_id = ${instanceId} AND task_id = ${taskId} + RETURNING tool_uses AS "toolUses" + `; + return { ...usage, toolUses: updated!.toolUses }; + }), + ); +});