From 6108ef3d3d0ef2a598088a1bb64980580e935830 Mon Sep 17 00:00:00 2001 From: Theo Browne Date: Sat, 3 Oct 2026 01:43:08 -0700 Subject: [PATCH 01/42] fix(server): runs no longer get stuck (#15048) Co-authored-by: Claude Opus 5.5 (1M context) --- .../src/orchestration-v2/EffectOutbox.ts | 24 +- apps/server/src/orchestration-v2/EventSink.ts | 91 +++-- .../FoundationPersistence.test.ts | 159 ++++++++ .../ProviderSessionManager.test.ts | 350 +++++++++++++++++- .../ProviderSessionManager.ts | 172 +++++++-- .../ProviderTurnStartService.test.ts | 59 ++- .../ProviderTurnStartService.ts | 4 +- .../RunExecutionService.test.ts | 161 +++++++- .../orchestration-v2/RunExecutionService.ts | 54 +-- 9 files changed, 959 insertions(+), 115 deletions(-) diff --git a/apps/server/src/orchestration-v2/EffectOutbox.ts b/apps/server/src/orchestration-v2/EffectOutbox.ts index 4841a5762480..9f63d8c568b6 100644 --- a/apps/server/src/orchestration-v2/EffectOutbox.ts +++ b/apps/server/src/orchestration-v2/EffectOutbox.ts @@ -278,9 +278,16 @@ export const layer: Layer.Layer = La available, Array.from({ length: Math.min(64, Math.max(0, Math.floor(count))) }, () => undefined), ).pipe(Effect.asVoid); + // Each thread runs its effects one at a time, in enqueue (rowid) order. An earlier + // effect waiting out a retry backoff still blocks later ones, so a turn + // cannot start while a failed rollback is about to restore files. A claim + // that skips restart continuations is not blocked by them either. // Title generation is correlated metadata work, so it has its own // per-thread lane and cannot delay provider lifecycle effects. - const claimableCandidatePredicate = (availableBefore?: string) => + const claimableCandidatePredicate = ( + availableBefore?: string, + excludeRestartContinuations = false, + ) => sql` ${ availableBefore === undefined @@ -292,7 +299,18 @@ export const layer: Layer.Layer = La SELECT 1 FROM orchestration_v2_effect_outbox AS active WHERE active.thread_id = candidate.thread_id - AND active.status = 'running' + AND ( + active.status = 'running' + OR ( + active.status = 'pending' + AND active.rowid < candidate.rowid + AND ${ + excludeRestartContinuations + ? sql`active.effect_type != 'provider-runtime.continue'` + : sql`1 = 1` + } + ) + ) AND ( ( candidate.effect_type = 'thread-title.generate' @@ -492,7 +510,7 @@ export const layer: Layer.Layer = La WHERE effect_id = ( SELECT candidate.effect_id FROM orchestration_v2_effect_outbox AS candidate - WHERE ${claimableCandidatePredicate(nowIso)} + WHERE ${claimableCandidatePredicate(nowIso, excludeRestartContinuations)} AND ${excludeRestartContinuations ? sql`candidate.effect_type != 'provider-runtime.continue'` : sql`1 = 1`} ORDER BY candidate.available_at ASC, candidate.created_at ASC, candidate.effect_id ASC LIMIT 1 diff --git a/apps/server/src/orchestration-v2/EventSink.ts b/apps/server/src/orchestration-v2/EventSink.ts index c61ce6ba8b7b..71b6f0ea994f 100644 --- a/apps/server/src/orchestration-v2/EventSink.ts +++ b/apps/server/src/orchestration-v2/EventSink.ts @@ -17,6 +17,7 @@ import * as Effect from "effect/Effect"; import * as Layer from "effect/Layer"; import * as Option from "effect/Option"; import * as PubSub from "effect/PubSub"; +import * as Semaphore from "effect/Semaphore"; import * as Schema from "effect/Schema"; import * as Stream from "effect/Stream"; import * as SqlClient from "effect/unstable/sql/SqlClient"; @@ -224,6 +225,40 @@ const baseLayer: Layer.Layer< ); } }); + const publishStoredEvents = (events: ReadonlyArray) => + eventStore.publishCommitted(events).pipe(Effect.andThen(publishLiveEvents(events))); + + // Transactions commit one at a time, but each writer publishes after its + // commit. If a writer is descheduled in between, a later commit reaches + // subscribers first, and clients drop any event at or below the newest + // sequence they have applied. So a writer takes this lane as the last step + // of its transaction and holds it until it has published. Publishing never + // waits, so a writer that holds the transaction while it waits for the + // lane is not blocked for long. + const publishLane = yield* Semaphore.make(1); + const commitThenPublish = ( + transaction: Effect.Effect, + publish: (committed: A) => Effect.Effect, + ) => + Effect.suspend(() => { + let holdsLane = false; + const takeLane = publishLane.take(1).pipe( + Effect.andThen( + Effect.sync(() => { + holdsLane = true; + }), + ), + Effect.uninterruptible, + ); + return sql + .withTransaction(Effect.tap(transaction, () => takeLane)) + .pipe( + Effect.tap(publish), + Effect.ensuring( + Effect.suspend(() => (holdsLane ? publishLane.release(1) : Effect.void)), + ), + ); + }); // A user can answer after terminal normalization reads the pending request. // Recheck inside the write transaction so stale cleanup cannot erase answers. @@ -332,7 +367,7 @@ const baseLayer: Layer.Layer< "orchestration_v2.thread_id": input.events[0]?.threadId ?? null, }); - const storedEvents = yield* sql.withTransaction( + return yield* commitThenPublish( Effect.gen(function* () { const normalized = yield* normalizeEvents( input.guardPendingUserInputCancellations === true @@ -347,13 +382,14 @@ const baseLayer: Layer.Layer< yield* effectOutbox.enqueue(input.effects); return committed; }), + (storedEvents) => + Effect.gen(function* () { + if (input.effects.length > 0) { + yield* effectOutbox.notifyAvailable(input.effects.length); + } + yield* publishStoredEvents(storedEvents); + }), ); - if (input.effects.length > 0) { - yield* effectOutbox.notifyAvailable(input.effects.length); - } - yield* eventStore.publishCommitted(storedEvents); - yield* publishLiveEvents(storedEvents); - return storedEvents; }); const writeIfRunCurrentEffect = Effect.fn("orchestrationV2.EventSink.writeIfRunCurrent")( @@ -365,7 +401,7 @@ const baseLayer: Layer.Layer< "orchestration_v2.thread_id": input.threadId, }); - const result = yield* sql.withTransaction( + return yield* commitThenPublish( Effect.gen(function* () { const rows = yield* sql<{ readonly status: string; @@ -403,12 +439,8 @@ const baseLayer: Layer.Layer< yield* applyStoredEvents(storedEvents); return { committed: true as const, storedEvents }; }), + (result) => (result.committed ? publishStoredEvents(result.storedEvents) : Effect.void), ); - if (result.committed) { - yield* eventStore.publishCommitted(result.storedEvents); - yield* publishLiveEvents(result.storedEvents); - } - return result; }, ); @@ -424,7 +456,7 @@ const baseLayer: Layer.Layer< "orchestration_v2.expected_last_run_ordinal": input.expectedLastRunOrdinal, }); - const result = yield* sql.withTransaction( + return yield* commitThenPublish( Effect.gen(function* () { const rows = yield* sql<{ readonly active_attempt_id: string | null; @@ -464,12 +496,8 @@ const baseLayer: Layer.Layer< yield* applyStoredEvents(storedEvents); return { committed: true as const, storedEvents }; }), + (result) => (result.committed ? publishStoredEvents(result.storedEvents) : Effect.void), ); - if (result.committed) { - yield* eventStore.publishCommitted(result.storedEvents); - yield* publishLiveEvents(result.storedEvents); - } - return result; }); const existingCommandResult = (commandId: CommandId) => @@ -490,7 +518,7 @@ const baseLayer: Layer.Layer< const commitCommandEffect = Effect.fn("orchestrationV2.EventSink.commitCommand")(function* ( input: Parameters[0], ) { - const result = yield* sql.withTransaction( + const result = yield* commitThenPublish( Effect.gen(function* () { const reserved = yield* commandReceipts.insertIfAbsent({ commandId: input.commandId, @@ -538,15 +566,15 @@ const baseLayer: Layer.Layer< }); return { receipt, storedEvents, committed: true as const, cancelledEffectIds }; }), + (result) => + Effect.gen(function* () { + yield* effectOutbox.signalCancellations(result.cancelledEffectIds); + if (result.committed && input.effects.length > 0) { + yield* effectOutbox.notifyAvailable(input.effects.length); + } + if (result.committed) yield* publishStoredEvents(result.storedEvents); + }), ); - yield* effectOutbox.signalCancellations(result.cancelledEffectIds); - if (result.committed && input.effects.length > 0) { - yield* effectOutbox.notifyAvailable(input.effects.length); - } - if (result.committed) { - yield* eventStore.publishCommitted(result.storedEvents); - yield* publishLiveEvents(result.storedEvents); - } return { receipt: result.receipt, storedEvents: result.storedEvents, @@ -593,7 +621,7 @@ const baseLayer: Layer.Layer< const commitProjectCommandEffect = Effect.fn("orchestrationV2.EventSink.commitProjectCommand")( function* (input: Parameters[0]) { - const result = yield* sql.withTransaction( + const result = yield* commitThenPublish( Effect.gen(function* () { const reserved: CommandReceiptStore.ProjectCommandReceiptV2 = { commandId: input.commandId, @@ -613,10 +641,9 @@ const baseLayer: Layer.Layer< yield* commandReceipts.upsert(receipt); return { receipt, event }; }), + (result) => + result.event === undefined ? Effect.void : eventStore.publishCommitted([result.event]), ); - if (result.event !== undefined) { - yield* eventStore.publishCommitted([result.event]); - } return { receipt: result.receipt, committed: result.event !== undefined }; }, ); diff --git a/apps/server/src/orchestration-v2/FoundationPersistence.test.ts b/apps/server/src/orchestration-v2/FoundationPersistence.test.ts index 22b93178711c..dc7eacc59ba5 100644 --- a/apps/server/src/orchestration-v2/FoundationPersistence.test.ts +++ b/apps/server/src/orchestration-v2/FoundationPersistence.test.ts @@ -51,7 +51,9 @@ import * as EventStore from "./EventStore.ts"; import * as IdAllocator from "./IdAllocator.ts"; import * as ProjectionMaintenance from "./ProjectionMaintenance.ts"; import * as ProjectionStore from "./ProjectionStore.ts"; +import * as ProjectStore from "./ProjectStore.ts"; import * as ProviderRuntimeRecovery from "./ProviderRuntimeRecoveryService.ts"; +import * as TurnItemPositionStore from "./TurnItemPositionStore.ts"; const isLiveStreamBufferError = Schema.is(LiveStreamBufferError); const encodeJson = Schema.encodeEffect(Schema.fromJsonString(Schema.Unknown)); @@ -2303,6 +2305,76 @@ it.layer(TestLayer)("orchestration V2 foundation persistence", (it) => { }).pipe(Effect.provide(Layer.fresh(effectOutboxProvided))), ); + it.effect("keeps later thread effects behind an earlier effect waiting to retry", () => + Effect.gen(function* () { + const outbox = yield* EffectOutbox.EffectOutboxV2; + const workerId = "retry-order-worker"; + const commandId = CommandId.make("command:foundation-retry-order"); + const threadId = ThreadId.make("thread:foundation-retry-order"); + yield* outbox.enqueue([ + { + id: "effect:foundation-retry-order:z-rollback", + commandId, + threadId, + request: { + type: "provider-thread.rollback", + providerThreadId: ProviderThreadId.make("provider-thread:foundation-retry-order"), + checkpointId: CheckpointId.make("checkpoint:foundation-retry-order"), + scopeId: CheckpointScopeId.make("scope:foundation-retry-order"), + }, + }, + ]); + const rollback = yield* outbox.claimNext({ workerId, leaseDurationMs: 30_000 }); + assert.isTrue(Option.isSome(rollback)); + if (Option.isNone(rollback)) return; + yield* outbox.retry({ + effectId: rollback.value.id, + workerId, + error: "rollback failed once", + delayMs: 60_000, + }); + + // A turn the user starts during the rollback's backoff must not run first, + // even when its timestamp ties and its id sorts first. + yield* outbox.enqueue([ + { + id: "effect:foundation-retry-order:a-start", + commandId: CommandId.make("command:foundation-retry-order:start"), + threadId, + request: { type: "provider-turn.start", runId: RunId.make("run:foundation-retry-order") }, + }, + { + id: "effect:foundation-retry-order:b-title", + commandId: CommandId.make("command:foundation-retry-order:title"), + threadId, + request: { type: "thread-title.generate", kind: { type: "regenerate" } }, + }, + ]); + const title = yield* outbox.claimNext({ workerId, leaseDurationMs: 30_000 }); + assert.equal(Option.getOrUndefined(title)?.id, "effect:foundation-retry-order:b-title"); + const blocked = yield* outbox.claimNext({ workerId, leaseDurationMs: 30_000 }); + assert.isTrue(Option.isNone(blocked)); + const nextClaimable = yield* outbox.nextClaimableAt; + assert.isTrue(Option.isSome(nextClaimable)); + if (Option.isSome(nextClaimable)) { + assert.equal( + DateTime.formatIso(nextClaimable.value), + (yield* outbox.get(rollback.value.id)).pipe(Option.getOrThrow).availableAt, + ); + } + + yield* outbox.cancelUnsettled({ + threadId, + effectTypes: ["provider-thread.rollback"], + reason: "Test cleanup.", + }); + const unblocked = yield* outbox.claimNext({ workerId, leaseDurationMs: 30_000 }); + assert.equal(Option.getOrUndefined(unblocked)?.id, "effect:foundation-retry-order:a-start"); + yield* outbox.succeed({ effectId: "effect:foundation-retry-order:a-start", workerId }); + yield* outbox.succeed({ effectId: "effect:foundation-retry-order:b-title", workerId }); + }).pipe(Effect.provide(Layer.fresh(effectOutboxProvided))), + ); + it.effect("executes a retry at its durable deadline instead of the liveness interval", () => Effect.gen(function* () { const outbox = yield* EffectOutbox.EffectOutboxV2; @@ -3222,3 +3294,90 @@ it.live("keeps claiming new work after repeated idle periods", () => }).pipe(Effect.provide(workerLayer), Effect.scoped); }).pipe(Effect.provide(TestLayer)), ); + +it.effect("publishes live events in commit order across concurrent writers", () => + Effect.gen(function* () { + const firstCommitted = yield* Deferred.make(); + const releaseFirst = yield* Deferred.make(); + // The first writer's post-commit wakeup stands in for any scheduler yield + // between its commit and its publish. + const pausingOutbox = Layer.effect( + EffectOutbox.EffectOutboxV2, + Effect.gen(function* () { + const delegate = yield* EffectOutbox.EffectOutboxV2; + return EffectOutbox.EffectOutboxV2.of({ + ...delegate, + notifyAvailable: (count) => + Deferred.succeed(firstCommitted, undefined).pipe( + Effect.andThen(Deferred.await(releaseFirst)), + Effect.andThen(delegate.notifyAvailable(count)), + ), + }); + }), + ).pipe(Layer.provide(effectOutboxProvided)); + const eventSinkLayer = EventSink.layerFromStores.pipe( + Layer.provide( + Layer.mergeAll( + storesProvided, + pausingOutbox, + commandReceiptStoreProvided, + ProjectStore.layer.pipe(Layer.provide(databaseLayer)), + TurnItemPositionStore.layer.pipe(Layer.provide(databaseLayer)), + ), + ), + ); + + yield* Effect.gen(function* () { + const eventSink = yield* EventSink.EventSinkV2; + const now = yield* DateTime.now; + const first = makeThread(ThreadId.make("thread:foundation-publish-order:first"), now); + const second = makeThread(ThreadId.make("thread:foundation-publish-order:second"), now); + const published = yield* eventSink + .stream({ afterSequence: yield* eventSink.latestSequence() }) + .pipe(Stream.take(2), Stream.runCollect, Effect.forkScoped({ startImmediately: true })); + + const firstWrite = yield* eventSink + .writeWithEffects({ + events: [ + threadCreatedEvent({ id: "event:foundation-publish-order:first", thread: first, now }), + ], + effects: [ + { + id: "effect:foundation-publish-order:first", + commandId: CommandId.make("command:foundation-publish-order:first"), + threadId: first.id, + request: { type: "terminal.cleanup" }, + }, + ], + }) + .pipe(Effect.forkScoped); + yield* Deferred.await(firstCommitted); + // The second writer commits after the first. It may run as far as it can + // before the first writer resumes. + const secondWrite = yield* eventSink + .write({ + events: [ + threadCreatedEvent({ + id: "event:foundation-publish-order:second", + thread: second, + now, + }), + ], + }) + .pipe( + Effect.provideService(Scheduler.MaxOpsBeforeYield, Number.POSITIVE_INFINITY), + Effect.forkScoped, + ); + yield* Effect.yieldNow; + yield* Deferred.succeed(releaseFirst, undefined); + yield* Fiber.join(firstWrite); + yield* Fiber.join(secondWrite); + + const sequences = Array.from(yield* Fiber.join(published), (stored) => stored.sequence); + assert.deepEqual( + sequences, + [...sequences].sort((left, right) => left - right), + ); + }).pipe(Effect.provide(eventSinkLayer)); + }).pipe(Effect.provide(databaseLayer)), +); diff --git a/apps/server/src/orchestration-v2/ProviderSessionManager.test.ts b/apps/server/src/orchestration-v2/ProviderSessionManager.test.ts index a567a4ab4841..696a40d8d176 100644 --- a/apps/server/src/orchestration-v2/ProviderSessionManager.test.ts +++ b/apps/server/src/orchestration-v2/ProviderSessionManager.test.ts @@ -20,6 +20,7 @@ import * as Cause from "effect/Cause"; import * as DateTime from "effect/DateTime"; import * as Deferred from "effect/Deferred"; import * as Effect from "effect/Effect"; +import * as Exit from "effect/Exit"; import * as Fiber from "effect/Fiber"; import * as FileSystem from "effect/FileSystem"; import * as Layer from "effect/Layer"; @@ -80,6 +81,52 @@ const FailingReleaseEventSinkLayer = Layer.effect( }), ).pipe(Layer.provide(TestEventSinkLayer)); +interface FlakyReleaseWrites { + /** Which release writes fail right now. */ + readonly failing: Ref.Ref<"none" | "session" | "session-and-requests">; + /** Receives one item per failed write. */ + readonly failures: Queue.Queue; + /** Holds runtime request writes: completes `paused`, then waits for `resume`. */ + readonly pauseRequestWrites?: { + readonly paused: Deferred.Deferred; + readonly resume: Deferred.Deferred; + }; +} + +// Fails release writes with a defect, the way a failed SQL commit surfaces. +const makeFlakyReleaseEventSinkLayer = (flaky: FlakyReleaseWrites) => + Layer.effect( + EventSink.EventSinkV2, + Effect.gen(function* () { + const delegate = yield* EventSink.EventSinkV2; + return EventSink.EventSinkV2.of({ + ...delegate, + write: (input) => + Effect.gen(function* () { + const failing = yield* Ref.get(flaky.failing); + const fails = input.events.some( + (event) => + (failing !== "none" && + event.type === "provider-session.updated" && + (event.payload.status === "stopped" || event.payload.status === "error")) || + (failing === "session-and-requests" && event.type === "runtime-request.updated"), + ); + const pause = flaky.pauseRequestWrites; + if ( + pause !== undefined && + input.events.some((event) => event.type === "runtime-request.updated") + ) { + yield* Deferred.succeed(pause.paused, undefined); + yield* Deferred.await(pause.resume); + } + if (!fails) return yield* delegate.write(input); + yield* Queue.offer(flaky.failures, undefined); + return yield* Effect.die(new Error("simulated commit failure")); + }), + }); + }), + ).pipe(Layer.provide(TestEventSinkLayer)); + const CodexCapabilities: OrchestrationV2ProviderCapabilities = CodexProviderCapabilitiesV2; const ExclusiveCapabilities: OrchestrationV2ProviderCapabilities = { ...CodexCapabilities, @@ -359,15 +406,19 @@ function makeTestLayer(input: { readonly initialProviderItemIdentityVersion?: 2; }) => Effect.Effect; readonly failReleaseEventWrites?: boolean; + readonly flakyReleaseWrites?: FlakyReleaseWrites; readonly hasPendingBackgroundWork?: Effect.Effect; readonly hangSessionScopeClose?: boolean; readonly beforeUnload?: Effect.Effect; readonly serverSettingsLayer?: ReturnType; readonly projectServiceLayer?: Layer.Layer; }) { - const configuredEventSinkLayer = input.failReleaseEventWrites - ? FailingReleaseEventSinkLayer - : TestEventSinkLayer; + const configuredEventSinkLayer = + input.flakyReleaseWrites !== undefined + ? makeFlakyReleaseEventSinkLayer(input.flakyReleaseWrites) + : input.failReleaseEventWrites + ? FailingReleaseEventSinkLayer + : TestEventSinkLayer; const registryLayer = ProviderAdapterRegistry.makeSingleLayer( makeProviderAdapter(input.state, { failEventStream: input.failEventStream ?? false, @@ -2299,6 +2350,299 @@ it.effect("ProviderSessionManagerV2 marks pending runtime requests non-live on r yield* effect.pipe(Effect.provide(makeTestLayer({ state, idleTimeoutMs: 1000 }))); }), ); +it.effect("ProviderSessionManagerV2 retries release records that failed to persist", () => + Effect.gen(function* () { + const state = yield* Ref.make(emptyState); + const flaky: FlakyReleaseWrites = { + failing: yield* Ref.make<"none" | "session" | "session-and-requests">("session"), + failures: yield* Queue.unbounded(), + }; + const effect = Effect.gen(function* () { + const eventSink = yield* EventSink.EventSinkV2; + const idAllocator = yield* IdAllocator.IdAllocatorV2; + const manager = yield* ProviderSessionManager.ProviderSessionManagerV2; + const projectionStore = yield* ProjectionStore.ProjectionStoreV2; + const now = yield* DateTime.now; + const threadId = ThreadId.make("thread-provider-session-manager-release-retry"); + const providerSessionId = yield* idAllocator.allocate.providerSession({ + providerInstanceId: modelSelection.instanceId, + threadId, + }); + const providerThread = makeProviderThread({ idAllocator, threadId, providerSessionId, now }); + yield* eventSink.write({ + events: [yield* makeThreadCreatedEvent({ idAllocator, threadId, now })], + }); + const pendingRequest = yield* makePendingRuntimeRequestEvents({ + idAllocator, + threadId, + providerSessionId, + providerThread, + now, + }); + yield* eventSink.write({ events: pendingRequest.events }); + yield* manager.open({ threadId, providerSessionId, modelSelection, runtimePolicy }); + + assert.isTrue(Exit.isFailure(yield* Effect.exit(manager.close(providerSessionId)))); + yield* Queue.take(flaky.failures); + // The failed session write does not keep the approval answerable. + const afterClose = yield* projectionStore.getThreadProjection(threadId); + assert.equal(afterClose.runtimeRequests.at(-1)?.responseCapability.type, "not_resumable"); + assert.equal(afterClose.providerSessions.at(-1)?.status, "ready"); + + // The first retry fails as well, and the retries continue. + yield* TestClock.adjust("1 second"); + yield* Queue.take(flaky.failures); + yield* Ref.set(flaky.failing, "none"); + const stopped = yield* eventSink + .stream({ + threadId, + afterSequence: yield* eventSink.latestSequence({ threadId }), + eventType: "provider-session.updated", + }) + .pipe(Stream.runHead, Effect.forkScoped); + yield* TestClock.adjust("1 second"); + yield* Fiber.join(stopped); + + const afterRetry = yield* projectionStore.getThreadProjection(threadId); + assert.equal(afterRetry.providerSessions.at(-1)?.status, "stopped"); + }); + + yield* effect.pipe( + Effect.provide(makeTestLayer({ state, idleTimeoutMs: 60_000, flakyReleaseWrites: flaky })), + ); + }), +); + +it.effect("ProviderSessionManagerV2 release retries leave a replacement session alone", () => + Effect.gen(function* () { + const state = yield* Ref.make(emptyState); + const flaky: FlakyReleaseWrites = { + failing: yield* Ref.make<"none" | "session" | "session-and-requests">("none"), + failures: yield* Queue.unbounded(), + }; + const effect = Effect.gen(function* () { + const eventSink = yield* EventSink.EventSinkV2; + const idAllocator = yield* IdAllocator.IdAllocatorV2; + const manager = yield* ProviderSessionManager.ProviderSessionManagerV2; + const projectionStore = yield* ProjectionStore.ProjectionStoreV2; + const threadId = ThreadId.make("thread-provider-session-manager-release-replacement"); + const providerSessionId = yield* idAllocator.allocate.providerSession({ + providerInstanceId: modelSelection.instanceId, + threadId, + }); + const writePendingRequest = Effect.gen(function* () { + const now = yield* DateTime.now; + const request = yield* makePendingRuntimeRequestEvents({ + idAllocator, + threadId, + providerSessionId, + providerThread: makeProviderThread({ idAllocator, threadId, providerSessionId, now }), + now, + }); + yield* eventSink.write({ events: request.events }); + return request.requestId; + }); + yield* eventSink.write({ + events: [ + yield* makeThreadCreatedEvent({ idAllocator, threadId, now: yield* DateTime.now }), + ], + }); + const oldRequestId = yield* writePendingRequest; + yield* manager.open({ threadId, providerSessionId, modelSelection, runtimePolicy }); + yield* Ref.set(flaky.failing, "session-and-requests"); + assert.isTrue(Exit.isFailure(yield* Effect.exit(manager.close(providerSessionId)))); + yield* Queue.take(flaky.failures); + yield* Queue.take(flaky.failures); + yield* Ref.set(flaky.failing, "none"); + + // A replacement opens with the same id before the retry runs. + yield* TestClock.adjust("500 millis"); + yield* manager.open({ threadId, providerSessionId, modelSelection, runtimePolicy }); + const newRequestId = yield* writePendingRequest; + const replacementStatus = (yield* projectionStore.getThreadProjection( + threadId, + )).providerSessions.at(-1)?.status; + const settled = yield* eventSink + .stream({ + threadId, + afterSequence: yield* eventSink.latestSequence({ threadId }), + eventType: "runtime-request.updated", + }) + .pipe(Stream.runHead, Effect.forkScoped); + yield* TestClock.adjust("500 millis"); + yield* Fiber.join(settled); + + const projection = yield* projectionStore.getThreadProjection(threadId); + const request = (id: typeof oldRequestId) => + projection.runtimeRequests.find((candidate) => candidate.id === id); + assert.equal(request(oldRequestId)?.responseCapability.type, "not_resumable"); + assert.equal(request(newRequestId)?.responseCapability.type, "live"); + assert.equal(projection.providerSessions.at(-1)?.status, replacementStatus); + }); + + yield* effect.pipe( + Effect.provide(makeTestLayer({ state, idleTimeoutMs: 60_000, flakyReleaseWrites: flaky })), + ); + }), +); + +it.effect("ProviderSessionManagerV2 keeps each failed release's cleanup", () => + Effect.gen(function* () { + const state = yield* Ref.make(emptyState); + const flaky: FlakyReleaseWrites = { + failing: yield* Ref.make<"none" | "session" | "session-and-requests">("none"), + failures: yield* Queue.unbounded(), + }; + const effect = Effect.gen(function* () { + const eventSink = yield* EventSink.EventSinkV2; + const idAllocator = yield* IdAllocator.IdAllocatorV2; + const manager = yield* ProviderSessionManager.ProviderSessionManagerV2; + const projectionStore = yield* ProjectionStore.ProjectionStoreV2; + const now = yield* DateTime.now; + const firstThreadId = ThreadId.make("thread-provider-session-manager-release-each-a"); + const secondThreadId = ThreadId.make("thread-provider-session-manager-release-each-b"); + const providerSessionId = yield* idAllocator.allocate.providerSession({ + providerInstanceId: modelSelection.instanceId, + threadId: firstThreadId, + }); + yield* eventSink.write({ + events: [ + yield* makeThreadCreatedEvent({ idAllocator, threadId: firstThreadId, now }), + yield* makeThreadCreatedEvent({ idAllocator, threadId: secondThreadId, now }), + ], + }); + const secondThreadRequest = yield* makePendingRuntimeRequestEvents({ + idAllocator, + threadId: secondThreadId, + providerSessionId, + providerThread: makeProviderThread({ + idAllocator, + threadId: secondThreadId, + providerSessionId, + now, + }), + now, + }); + yield* eventSink.write({ events: secondThreadRequest.events }); + const failRelease = (failedWrites: number) => + Effect.gen(function* () { + yield* Ref.set(flaky.failing, "session-and-requests"); + assert.isTrue(Exit.isFailure(yield* Effect.exit(manager.close(providerSessionId)))); + yield* Effect.repeat(Queue.take(flaky.failures), { times: failedWrites - 1 }); + yield* Ref.set(flaky.failing, "none"); + }); + + // The first session serves both threads. Its replacement serves one. + yield* manager.open({ + threadId: firstThreadId, + providerSessionId, + modelSelection, + runtimePolicy, + }); + yield* manager.open({ + threadId: secondThreadId, + providerSessionId, + modelSelection, + runtimePolicy, + }); + // Both the session write and the second thread's request write fail. + yield* failRelease(2); + yield* TestClock.adjust("500 millis"); + yield* manager.open({ + threadId: firstThreadId, + providerSessionId, + modelSelection, + runtimePolicy, + }); + // Only the session write fails: the first thread has no requests. + yield* failRelease(1); + + const settled = yield* eventSink + .stream({ + threadId: secondThreadId, + afterSequence: yield* eventSink.latestSequence({ threadId: secondThreadId }), + eventType: "runtime-request.updated", + }) + .pipe(Stream.runHead, Effect.forkScoped); + yield* TestClock.adjust("500 millis"); + yield* Fiber.join(settled); + + const projection = yield* projectionStore.getThreadProjection(secondThreadId); + assert.equal(projection.runtimeRequests.at(-1)?.responseCapability.type, "not_resumable"); + }); + + yield* effect.pipe( + Effect.provide(makeTestLayer({ state, idleTimeoutMs: 60_000, flakyReleaseWrites: flaky })), + ); + }), +); + +it.effect("ProviderSessionManagerV2 settles a request the event pump persists during release", () => + Effect.gen(function* () { + const state = yield* Ref.make(emptyState); + const flaky: FlakyReleaseWrites = { + failing: yield* Ref.make<"none" | "session" | "session-and-requests">("none"), + failures: yield* Queue.unbounded(), + pauseRequestWrites: { + paused: yield* Deferred.make(), + resume: yield* Deferred.make(), + }, + }; + const pause = flaky.pauseRequestWrites!; + const effect = Effect.gen(function* () { + const eventSink = yield* EventSink.EventSinkV2; + const idAllocator = yield* IdAllocator.IdAllocatorV2; + const manager = yield* ProviderSessionManager.ProviderSessionManagerV2; + const projectionStore = yield* ProjectionStore.ProjectionStoreV2; + const now = yield* DateTime.now; + const threadId = ThreadId.make("thread-provider-session-manager-release-pump"); + const providerSessionId = yield* idAllocator.allocate.providerSession({ + providerInstanceId: modelSelection.instanceId, + threadId, + }); + yield* eventSink.write({ + events: [yield* makeThreadCreatedEvent({ idAllocator, threadId, now })], + }); + yield* manager.open({ threadId, providerSessionId, modelSelection, runtimePolicy }); + // The runtime creates the request a moment after the release starts. + const createdAt = DateTime.add(now, { seconds: 1 }); + const pendingRequest = yield* makePendingRuntimeRequestEvents({ + idAllocator, + threadId, + providerSessionId, + providerThread: makeProviderThread({ + idAllocator, + threadId, + providerSessionId, + now: createdAt, + }), + now: createdAt, + }); + const adapterEvents = (yield* Ref.get(state)).eventQueues.get(String(providerSessionId)); + assert.isDefined(adapterEvents); + yield* Queue.offerAll(adapterEvents!, pendingRequest.providerEvents); + // The event pump holds the request permit while it persists the request. + yield* Deferred.await(pause.paused); + const closed = yield* manager + .close(providerSessionId) + .pipe(Effect.forkScoped({ startImmediately: true })); + yield* TestClock.adjust("1 second"); + yield* Deferred.succeed(pause.resume, undefined); + yield* Fiber.join(closed); + + const projection = yield* projectionStore.getThreadProjection(threadId); + const request = projection.runtimeRequests.find( + (candidate) => candidate.id === pendingRequest.requestId, + ); + assert.equal(request?.responseCapability.type, "not_resumable"); + }); + + yield* effect.pipe( + Effect.provide(makeTestLayer({ state, idleTimeoutMs: 60_000, flakyReleaseWrites: flaky })), + ); + }), +); + it.effect("ProviderSessionManagerV2 terminalizes a pending input transcript item on release", () => Effect.gen(function* () { const state = yield* Ref.make(emptyState); diff --git a/apps/server/src/orchestration-v2/ProviderSessionManager.ts b/apps/server/src/orchestration-v2/ProviderSessionManager.ts index 786eb5a6bf0a..c1039c2a0207 100644 --- a/apps/server/src/orchestration-v2/ProviderSessionManager.ts +++ b/apps/server/src/orchestration-v2/ProviderSessionManager.ts @@ -17,11 +17,13 @@ import * as Duration from "effect/Duration"; import * as Effect from "effect/Effect"; import * as Exit from "effect/Exit"; import * as Fiber from "effect/Fiber"; +import * as FiberSet from "effect/FiberSet"; import * as FileSystem from "effect/FileSystem"; import * as Layer from "effect/Layer"; import * as Option from "effect/Option"; import * as Queue from "effect/Queue"; import * as Ref from "effect/Ref"; +import * as Schedule from "effect/Schedule"; import * as Schema from "effect/Schema"; import * as Scope from "effect/Scope"; import * as Semaphore from "effect/Semaphore"; @@ -361,6 +363,9 @@ export const layerWithOptions = ( ); const layerScope = yield* Effect.scope; const sessions = yield* Ref.make(new Map()); + // One retry per released entry, so a later release with the same id + // cannot drop cleanup for threads only the earlier session served. + const releaseRecordRetries = yield* FiberSet.make(); const nextSubscriberId = yield* Ref.make(0); const sessionOpen = yield* makeKeyedSerialExecutor(); // Orders a thread's attach against a detach unloading it on the same session. @@ -608,6 +613,8 @@ export const layerWithOptions = ( const writeReleasedRuntimeRequestEvents = (input: { readonly entry: LiveSessionEntry; readonly reason: ProviderSessionReleaseReason; + /** Requests created later belong to a replacement session with the same id. */ + readonly releasedAt: DateTime.Utc; }) => Effect.gen(function* () { const providerSessionId = input.entry.runtime.providerSessionId; @@ -629,7 +636,8 @@ export const layerWithOptions = ( (request) => request.status === "pending" && request.responseCapability.type === "live" && - request.responseCapability.providerSessionId === providerSessionId, + request.responseCapability.providerSessionId === providerSessionId && + DateTime.isLessThanOrEqualTo(request.createdAt, input.releasedAt), ); for (const request of releasedRequests) { @@ -708,6 +716,125 @@ export const layerWithOptions = ( } }); + // Records a released session as stopped and resolves the live runtime + // requests it left. Each write runs even if the other fails. Once a + // replacement session opens with the same id, it owns the session status, + // so only the requests are settled. + const writeReleaseRecords = (input: { + readonly entry: LiveSessionEntry; + readonly reason: ProviderSessionReleaseReason; + readonly detail?: string; + readonly releasedAt: DateTime.Utc; + readonly replaced: boolean; + }) => + Effect.all( + [ + input.replaced + ? Effect.succeed(Exit.void) + : Effect.exit(writeReleasedSessionEvents(input)), + Effect.exit( + writeReleasedRuntimeRequestEvents(input).pipe( + input.entry.requestEventPermit.withPermits(1), + ), + ), + ], + { concurrency: 1 }, + ).pipe(Effect.flatMap(Exit.asVoidAll)); + + // The session already left the live map, so a later release finds + // nothing to do. Without a retry the UI would keep a ready session and + // answerable approvals until a server restart. Each attempt holds the + // session's open lock, so it sees a replacement that opened meanwhile. + const retryReleaseRecords = ( + input: Omit[0], "replaced">, + ) => { + const providerSessionId = input.entry.runtime.providerSessionId; + const attempt = Effect.gen(function* () { + const exit = yield* Effect.exit( + sessionOpen.withLock( + providerSessionId, + Effect.gen(function* () { + const replaced = (yield* Ref.get(sessions)).has(sessionKey(providerSessionId)); + yield* writeReleaseRecords({ ...input, replaced }); + }), + ), + ); + if (Exit.isSuccess(exit) || Cause.hasInterruptsOnly(exit.cause)) return yield* exit; + yield* Effect.logWarning("orchestration-v2.provider-session-release-records-failed", { + providerSessionId, + cause: exit.cause, + }); + // A failed SQL commit is a defect, so every failure but interruption + // is retried. + return yield* Effect.fail(exit.cause); + }); + return attempt.pipe( + Effect.retry({ + schedule: Schedule.exponential("1 second").pipe( + Schedule.modifyDelay(({ duration }) => + Effect.succeed(Duration.min(duration, Duration.seconds(30))), + ), + ), + }), + Effect.delay("1 second"), + FiberSet.run(releaseRecordRetries), + ); + }; + + const logReleaseFailure = + (providerSessionId: ProviderSessionId) => + (release: Effect.Effect) => + release.pipe( + Effect.catchCause((cause) => + Effect.logWarning("orchestration-v2.provider-session-release-failed", { + providerSessionId, + cause, + }), + ), + ); + + // Removes the live entry and reads the request cleanup cutoff while + // holding the entry's request permit. A request the event pump is + // persisting for this runtime lands before the cutoff, and once the + // entry is gone the pump persists no more for it. A replacement's + // requests come after its own open. + const removeLiveEntry = (input: { + readonly providerSessionId: ProviderSessionId; + readonly onlyIfIdleGeneration?: number; + }): Effect.Effect, DateTime.Utc]> => + Effect.gen(function* () { + const key = sessionKey(input.providerSessionId); + const candidate = (yield* Ref.get(sessions)).get(key); + if (candidate === undefined) { + return [Option.none(), yield* DateTime.now] as const; + } + const removed = yield* Effect.zip( + Ref.modify(sessions, (current) => { + const existing = current.get(key); + if (existing !== candidate) { + return [existing === undefined ? "gone" : "changed", current] as const; + } + if ( + input.onlyIfIdleGeneration !== undefined && + (existing.busyCount > 0 || existing.idleGeneration !== input.onlyIfIdleGeneration) + ) { + return ["kept", current] as const; + } + const updated = new Map(current); + updated.delete(key); + return ["removed", updated] as const; + }), + DateTime.now, + ).pipe(candidate.requestEventPermit.withPermits(1)); + const [outcome, releasedAt] = removed; + // Another entry took this id while the permit was held; release it instead. + if (outcome === "changed") return yield* removeLiveEntry(input); + return [ + outcome === "removed" ? Option.some(candidate) : Option.none(), + releasedAt, + ] as const; + }); + const releaseEntry = (input: { readonly providerSessionId: ProviderSessionId; readonly reason: ProviderSessionReleaseReason; @@ -717,23 +844,8 @@ export const layerWithOptions = ( readonly gracefulSubscribers?: boolean; }) => Effect.acquireUseRelease( - Ref.modify(sessions, (current) => { - const key = sessionKey(input.providerSessionId); - const existing = current.get(key); - if (existing === undefined) { - return [Option.none(), current] as const; - } - if ( - input.onlyIfIdleGeneration !== undefined && - (existing.busyCount > 0 || existing.idleGeneration !== input.onlyIfIdleGeneration) - ) { - return [Option.none(), current] as const; - } - const updated = new Map(current); - updated.delete(key); - return [Option.some(existing), updated] as const; - }), - (entry) => + removeLiveEntry(input), + ([entry, releasedAt]) => Option.match(entry, { onNone: () => Effect.void, onSome: (entry) => @@ -794,21 +906,25 @@ export const layerWithOptions = ( Effect.forkDetach, ); } - yield* writeReleasedSessionEvents({ + const records = { entry, reason: input.reason, ...(input.detail === undefined ? {} : { detail: input.detail }), - }); - yield* writeReleasedRuntimeRequestEvents({ - entry, - reason: input.reason, - }).pipe(entry.requestEventPermit.withPermits(1)); + releasedAt, + }; + const recorded = yield* Effect.exit( + writeReleaseRecords({ ...records, replaced: false }), + ); + if (Exit.isFailure(recorded)) { + yield* retryReleaseRecords(records); + return yield* recorded; + } if (Option.isSome(closeExit) && Exit.isFailure(closeExit.value)) { return yield* Effect.failCause(closeExit.value.cause); } }), }), - (entry) => + ([entry]) => Option.match(entry, { onNone: () => Effect.void, onSome: (entry) => @@ -1495,7 +1611,7 @@ export const layerWithOptions = ( providerSessionId: entry.runtime.providerSessionId, reason: "manual_shutdown", gracefulSubscribers: true, - }).pipe(Effect.ignore); + }).pipe(logReleaseFailure(entry.runtime.providerSessionId)); return; } const cause = Exit.isFailure(exit) @@ -1516,7 +1632,7 @@ export const layerWithOptions = ( providerSessionId: entry.runtime.providerSessionId, reason: "runtime_error", detail: Cause.pretty(cause), - }).pipe(Effect.ignore); + }).pipe(logReleaseFailure(entry.runtime.providerSessionId)); }), ), Effect.forkIn(layerScope), @@ -1705,7 +1821,7 @@ export const layerWithOptions = ( providerSessionId: input.providerSessionId, reason: "runtime_error", detail: "Failed to persist the provider-session attachment.", - }).pipe(Effect.ignore), + }).pipe(logReleaseFailure(input.providerSessionId)), ), ); yield* startEventPump(entry); diff --git a/apps/server/src/orchestration-v2/ProviderTurnStartService.test.ts b/apps/server/src/orchestration-v2/ProviderTurnStartService.test.ts index 746fbe235d99..fdb45f193179 100644 --- a/apps/server/src/orchestration-v2/ProviderTurnStartService.test.ts +++ b/apps/server/src/orchestration-v2/ProviderTurnStartService.test.ts @@ -30,6 +30,7 @@ import * as ProviderAuthService from "../provider/Services/ProviderAuthService.t import * as ContextHandoffService from "./ContextHandoffService.ts"; import * as EventSink from "./EventSink.ts"; import * as IdAllocator from "./IdAllocator.ts"; +import { CodexProviderCapabilitiesV2 } from "./Adapters/CodexAdapterV2.ts"; import * as ProjectionStore from "./ProjectionStore.ts"; import { ProviderAdapterEventStreamError } from "./ProviderAdapter.ts"; import * as ProviderSessionManager from "./ProviderSessionManager.ts"; @@ -168,6 +169,8 @@ function makeLocalCommandHarness(input: { readonly interruptOpen?: boolean; readonly interruptRunBeforeOpenFailure?: boolean; readonly writeFailure?: unknown; + /** Loads the thread and starts the run, then fails every later state read. */ + readonly failReadsAfterRunning?: boolean; }) { const now = DateTime.makeUnsafe("2026-09-04T12:00:00Z"); const threadId = ThreadId.make("thread-native-account-command"); @@ -411,9 +414,40 @@ function makeLocalCommandHarness(input: { ), ), ) - : Effect.die("A local command must not open a native session."), + : input.failReadsAfterRunning === true + ? Effect.succeed({ + driver: providerThread.driver, + providerSession: { + id: providerSessionId, + driver: providerThread.driver, + providerInstanceId: newInstanceId, + status: "ready", + cwd: "/tmp/native-account-command", + model: null, + capabilities: CodexProviderCapabilitiesV2, + createdAt: now, + updatedAt: now, + lastError: null, + }, + ensureThread: () => Effect.succeed(providerThread), + } as never) + : Effect.die("A local command must not open a native session."), + ); + const startRootRun = vi.fn< + (input: RunExecutionService.RunExecutionServiceV2StartRootRunInput) => Effect.Effect + >(() => + input.failReadsAfterRunning === true + ? Effect.void + : Effect.die("A local command must not start a native turn."), + ); + const failReadIfRunning = Effect.suspend(() => + input.failReadsAfterRunning === true && + projection.runs.find((candidate) => candidate.id === runId)?.status === "running" + ? Effect.fail( + new ProjectionStore.ProjectionStoreReadError({ threadId, cause: "database unavailable" }), + ) + : Effect.void, ); - const startRootRun = vi.fn(() => Effect.die("A local command must not start a native turn.")); const tryHandlePromptCommand = vi.fn(() => input.logoutFailure === undefined ? Effect.succeed(true) @@ -471,7 +505,7 @@ function makeLocalCommandHarness(input: { ), }), getRuntimeRecoveryProjection: () => - Effect.succeed({ + Effect.as(failReadIfRunning, { ...projection, hasConversation: projection.messages.some( (m) => @@ -703,6 +737,25 @@ effectIt.effect( }), ); +effectIt.effect("does not mistake a failed state read for a superseded run", () => + Effect.gen(function* () { + const harness = makeLocalCommandHarness({ text: "Continue", failReadsAfterRunning: true }); + + yield* harness.start; + + expect(harness.projection().runs.at(-1)?.status).toBe("running"); + const controls = harness.startRootRun.mock.calls[0]?.[0]; + expect(controls).toBeDefined(); + if (controls === undefined) return; + // "false" would skip the provider turn or the terminal write and leave the + // run active. A read failure must reach the caller instead. + const startCheck = yield* Effect.flip(controls.shouldStartProviderTurn!()); + const finalizeCheck = yield* Effect.flip(controls.shouldFinalizeRun!()); + expect(startCheck._tag).toBe("ProjectionStoreReadError"); + expect(finalizeCheck._tag).toBe("ProjectionStoreReadError"); + }), +); + effectIt.effect("does not overwrite a run interrupted while its thread loads", () => Effect.gen(function* () { const harness = makeLocalCommandHarness({ diff --git a/apps/server/src/orchestration-v2/ProviderTurnStartService.ts b/apps/server/src/orchestration-v2/ProviderTurnStartService.ts index 781491ac7fd6..ce2088d911d4 100644 --- a/apps/server/src/orchestration-v2/ProviderTurnStartService.ts +++ b/apps/server/src/orchestration-v2/ProviderTurnStartService.ts @@ -124,13 +124,14 @@ export const layer: Layer.Layer< }) => { // Guards and background routing need live execution state, not a fresh // allocation of every completed message and tool output in the thread. + // `false` means the run moved on or is gone. A failed read is an error, + // so the caller fails the start or the run instead of skipping it. const isCurrentAttemptInStatus = (expectedStatus: OrchestrationV2Run["status"]) => projectionStore.getRuntimeRecoveryProjection(input.threadId).pipe( Effect.map((current) => { const run = current.runs.find((candidate) => candidate.id === input.runId); return run?.activeAttemptId === input.attemptId && run.status === expectedStatus; }), - Effect.catchCause(() => Effect.succeed(false)), ); return { isCurrentAttemptInStatus, @@ -157,7 +158,6 @@ export const layer: Layer.Layer< (run.status === "starting" || run.status === "running") ); }), - Effect.catchCause(() => Effect.succeed(false)), ), hasUnpairedRunInterruptRequest: () => projectionStore diff --git a/apps/server/src/orchestration-v2/RunExecutionService.test.ts b/apps/server/src/orchestration-v2/RunExecutionService.test.ts index 370c30058346..3d6dca6151ab 100644 --- a/apps/server/src/orchestration-v2/RunExecutionService.test.ts +++ b/apps/server/src/orchestration-v2/RunExecutionService.test.ts @@ -50,6 +50,7 @@ import { type ProviderAdapterV2Event, type ProviderAdapterV2SessionRuntime, } from "./ProviderAdapter.ts"; +import * as ProjectionStore from "./ProjectionStore.ts"; import * as ProviderEventIngestor from "./ProviderEventIngestor.ts"; import * as RunExecutionService from "./RunExecutionService.ts"; import * as RunFinalizationService from "./RunFinalizationService.ts"; @@ -590,6 +591,113 @@ it.effect("rechecks run ownership immediately before calling the provider", () = }).pipe(Effect.provide(RunExecutionTestLayer)), ); +it.effect("fails the run when its ownership check cannot be read before calling the provider", () => + Effect.gen(function* () { + const guardCalls = yield* Ref.make(0); + const providerStarts = yield* Ref.make(0); + const writes = yield* Ref.make>([]); + const threadId = ThreadId.make("thread:run-execution-start-guard-read"); + const runId = RunId.make("run:run-execution-start-guard-read"); + const attemptId = RunAttemptId.make("attempt:run-execution-start-guard-read"); + const providerInstanceId = ProviderInstanceId.make("codex"); + const testLayer = RunExecutionService.layer.pipe( + Layer.provide( + Layer.mergeAll( + Layer.mock(CheckpointService.CheckpointServiceV2)({ captureBaseline: () => Effect.void }), + Layer.mock(EventSink.EventSinkV2)({ + writeIfRunCurrent: (input) => + Ref.update(writes, (current) => [...current, ...input.events]).pipe( + Effect.as({ committed: true, storedEvents: [] }), + ), + }), + IdAllocator.layer, + Layer.mock(ProviderEventIngestor.ProviderEventIngestorV2)({ + ingestNormalized: () => Effect.succeed([]), + }), + ServerSettings.layerTest(), + ), + ), + ); + + yield* Effect.gen(function* () { + const runExecution = yield* RunExecutionService.RunExecutionServiceV2; + yield* runExecution.startRootRun({ + commandId: CommandId.make("command:run-execution-start-guard-read"), + appThread: { id: threadId } as OrchestrationV2AppThread, + providerSessionId: ProviderSessionId.make("session:run-execution-start-guard-read"), + session: { + events: Stream.never, + startTurn: () => Ref.update(providerStarts, (count) => count + 1), + } as unknown as ProviderAdapterV2SessionRuntime, + run: { id: runId, threadId, ordinal: 1, providerInstanceId } as OrchestrationV2Run, + rootNode: { + id: NodeId.make("node:run-execution-start-guard-read"), + } as OrchestrationV2ExecutionNode, + checkpointScope: { + id: CheckpointScopeId.make("checkpoint-scope:run-execution-start-guard-read"), + } as OrchestrationV2CheckpointScope, + providerThread: { + id: ProviderThreadId.make("provider-thread:run-execution-start-guard-read"), + driver, + } as OrchestrationV2ProviderThread, + attempt: { id: attemptId, providerTurnId: null } as OrchestrationV2RunAttempt, + attemptId, + providerTurnOrdinal: 1, + // The preparation check passes; the check right before the provider + // call cannot read the run. + shouldStartProviderTurn: () => + Ref.getAndUpdate(guardCalls, (calls) => calls + 1).pipe( + Effect.flatMap((calls) => + calls === 0 + ? Effect.succeed(true) + : Effect.fail( + new ProjectionStore.ProjectionStoreReadError({ + threadId, + cause: "database unavailable", + }), + ), + ), + ), + // The failure is settled by the guarded write, not by another read. + shouldFinalizeRun: () => + Effect.fail( + new ProjectionStore.ProjectionStoreReadError({ + threadId, + cause: "database unavailable", + }), + ), + message: { + messageId: MessageId.make("message:run-execution-start-guard-read"), + text: "Start while the store is down.", + attachments: [], + createdBy: "user", + creationSource: "web", + }, + modelSelection: { instanceId: providerInstanceId, model: "gpt-5.4" }, + runtimePolicy: { + runtimeMode: "full-access", + interactionMode: "default", + cwd: process.cwd(), + approvalPolicy: "never", + sandboxPolicy: { + type: "readOnly", + access: { type: "fullAccess" }, + networkAccess: false, + }, + }, + }); + }).pipe(Effect.provide(testLayer)); + + assert.equal(yield* Ref.get(guardCalls), 2); + assert.equal(yield* Ref.get(providerStarts), 0); + const runUpdate = (yield* Ref.get(writes)).find((event) => event.type === "run.updated"); + assert.equal( + runUpdate?.type === "run.updated" ? runUpdate.payload.status : undefined, + "failed", + ); + }), +); + it.effect( "dispatches only attachment-free compact commands through the native compaction path", () => @@ -3165,6 +3273,25 @@ it.effect.each(["completed", "interrupted", "cancelled", "failed"] as const)( }), ); +it.effect("records a finished run as failed when its ownership check cannot be read", () => + Effect.gen(function* () { + const { observed } = yield* captureRootRunTermination({ + key: "finalize-guard-read-failure", + shouldFinalizeRun: () => + Effect.fail( + new ProjectionStore.ProjectionStoreReadError({ + threadId: ThreadId.make("thread:finalize-guard-read-failure"), + cause: "database unavailable", + }), + ), + events: (ids) => Stream.make(rootTerminalEvent(ids, "completed")), + }); + // The fallback settles through the guarded write instead of the same + // failing read, so the run does not stay running. + assert.include(observed, "run:failed"); + }), +); + it.effect("does not refresh pull requests for auxiliary or stale provider terminals", () => Effect.gen(function* () { const { observed } = yield* captureRootRunTermination({ @@ -3254,7 +3381,7 @@ it.effect("keeps completed runs completed when pull request refresh fails", () = function captureRootRunTermination(input: { readonly key: string; - readonly shouldFinalizeRun: () => Effect.Effect; + readonly shouldFinalizeRun: () => Effect.Effect; readonly hasUnpairedRunInterruptRequest?: () => Effect.Effect; readonly seedOpenSubagent?: boolean; readonly events?: ( @@ -3278,6 +3405,17 @@ function captureRootRunTermination(input: { const ingestionDone = yield* Deferred.make(); const captureTurnItem = (payload: OrchestrationV2TurnItem) => Ref.update(writtenItems, (current) => [...current, payload]); + const captureFinalEvents = (events: ReadonlyArray) => + Effect.gen(function* () { + for (const event of events) { + if (event.type === "turn-item.updated") { + yield* captureTurnItem(event.payload); + } + if (event.type === "run.updated") { + yield* Ref.update(observed, (current) => [...current, `run:${event.payload.status}`]); + } + } + }); const testLayer = RunExecutionService.layer.pipe( Layer.provide( Layer.mergeAll( @@ -3292,22 +3430,11 @@ function captureRootRunTermination(input: { } return []; }), - writeWithEffects: (payload) => - Effect.gen(function* () { - for (const event of payload.events) { - if (event.type === "turn-item.updated") { - yield* captureTurnItem(event.payload); - } - if (event.type === "run.updated") { - yield* Ref.update(observed, (current) => [ - ...current, - `run:${event.payload.status}`, - ]); - } - } - return []; - }), - writeIfRunCurrent: () => Effect.succeed({ committed: true, storedEvents: [] }), + writeWithEffects: (payload) => captureFinalEvents(payload.events).pipe(Effect.as([])), + writeIfRunCurrent: (payload) => + captureFinalEvents(payload.events).pipe( + Effect.as({ committed: true, storedEvents: [] }), + ), }), IdAllocator.layer, Layer.mock(ProviderEventIngestor.ProviderEventIngestorV2)({ diff --git a/apps/server/src/orchestration-v2/RunExecutionService.ts b/apps/server/src/orchestration-v2/RunExecutionService.ts index 71211c126545..6aff5da24f1e 100644 --- a/apps/server/src/orchestration-v2/RunExecutionService.ts +++ b/apps/server/src/orchestration-v2/RunExecutionService.ts @@ -28,6 +28,7 @@ import * as Context from "effect/Context"; import * as Cause from "effect/Cause"; import * as DateTime from "effect/DateTime"; import * as Effect from "effect/Effect"; +import * as Exit from "effect/Exit"; import * as Fiber from "effect/Fiber"; import * as Layer from "effect/Layer"; import * as Ref from "effect/Ref"; @@ -47,6 +48,7 @@ import type { } from "./ProviderAdapter.ts"; import { ProviderAdapterTurnStartError } from "./ProviderAdapter.ts"; import * as ProviderEventIngestor from "./ProviderEventIngestor.ts"; +import type { ProjectionStoreV2Error } from "./ProjectionStore.ts"; import { makeProviderFailure, makeProviderFailureTurnItem } from "./ProviderFailure.ts"; import * as RunFinalizationService from "./RunFinalizationService.ts"; @@ -511,8 +513,8 @@ export interface RunExecutionServiceV2StartRootRunInput { >; readonly relatedThreadIds?: ReadonlyArray; readonly relatedProviderThreadIds?: ReadonlyArray; - readonly shouldStartProviderTurn?: () => Effect.Effect; - readonly shouldFinalizeRun?: () => Effect.Effect; + readonly shouldStartProviderTurn?: () => Effect.Effect; + readonly shouldFinalizeRun?: () => Effect.Effect; readonly hasUnpairedRunInterruptRequest?: () => Effect.Effect; readonly message: ProviderAdapterV2TurnMessage; readonly modelSelection: ModelSelection; @@ -557,7 +559,7 @@ export const layer: Layer.Layer< readonly checkpointScope: OrchestrationV2CheckpointScope; readonly providerThread: OrchestrationV2ProviderThread; readonly attempt: OrchestrationV2RunAttempt; - readonly shouldFinalizeRun?: () => Effect.Effect; + readonly shouldFinalizeRun?: () => Effect.Effect; readonly hasUnpairedRunInterruptRequest?: () => Effect.Effect; readonly openRunOwnedSubagents?: OpenRunOwnedSubagentProjection; readonly terminal: ProviderTerminalEvent; @@ -1286,15 +1288,12 @@ export const layer: Layer.Layer< checkpointScope: input.checkpointScope, providerThread, attempt: input.attempt, - ...(input.shouldFinalizeRun === undefined - ? {} - : { shouldFinalizeRun: input.shouldFinalizeRun }), - ...(input.hasUnpairedRunInterruptRequest === undefined - ? {} - : { - hasUnpairedRunInterruptRequest: - input.hasUnpairedRunInterruptRequest, - }), + // The failure may be the ownership + // read itself, so check in the write. + writeIfRunCurrent: { + activeAttemptId: input.attempt.id, + expectedStatus: "running", + }, openRunOwnedSubagents: openSubagents, terminal: makeFailedTerminalEvent( makeProviderFailure({ @@ -1328,10 +1327,13 @@ export const layer: Layer.Layer< Effect.forkDetach, ); - if ( - input.shouldStartProviderTurn !== undefined && - !(yield* input.shouldStartProviderTurn()) - ) { + // A failed read fails the start below, so the run is recorded as + // failed instead of staying active with no provider turn. + const shouldStart = + input.shouldStartProviderTurn === undefined + ? Exit.succeed(true) + : yield* Effect.exit(input.shouldStartProviderTurn()); + if (Exit.isSuccess(shouldStart) && !shouldStart.value) { yield* Fiber.interrupt(providerEventFiber); return; } @@ -1373,7 +1375,7 @@ export const layer: Layer.Layer< }), )) : input.session.startTurn(turnInput); - yield* startTurn.pipe( + yield* Effect.andThen(shouldStart, startTurn).pipe( Effect.catchCause((cause) => Effect.logError("orchestration V2 provider turn start failed", { runId: input.run.id, @@ -1392,20 +1394,18 @@ export const layer: Layer.Layer< checkpointScope: input.checkpointScope, providerThread, attempt: input.attempt, - ...(input.shouldFinalizeRun === undefined - ? {} - : { shouldFinalizeRun: input.shouldFinalizeRun }), - ...(input.hasUnpairedRunInterruptRequest === undefined - ? {} - : { - hasUnpairedRunInterruptRequest: - input.hasUnpairedRunInterruptRequest, - }), + // Checked in the write transaction, not by another + // read that can fail like the one before the start. + writeIfRunCurrent: { + activeAttemptId: input.attempt.id, + expectedStatus: "running", + }, openRunOwnedSubagents: openSubagents, terminal: makeFailedTerminalEvent( makeProviderFailure({ cause: Cause.squash(cause), - class: "provider_error", + // A failed ownership read is not the provider's fault. + class: Exit.isFailure(shouldStart) ? "unknown" : "provider_error", }), latestItemOrdinal + 1, ), From 56914128c1ff9dcf6585c380b7b25845d6ca5f98 Mon Sep 17 00:00:00 2001 From: Theo Browne Date: Sat, 3 Oct 2026 02:01:12 -0700 Subject: [PATCH 02/42] fix(usage): Codex Fast and Ultrafast now cost what they bill (#15101) Co-authored-by: Claude Opus 5.5 (1M context) --- apps/server/src/usage/UsageService.test.ts | 93 +++++++++++ apps/server/src/usage/UsageService.ts | 20 ++- .../src/usage/antigravityUsageReader.ts | 2 +- apps/server/src/usage/cursorUsageReader.ts | 2 +- apps/server/src/usage/opencodeUsageReader.ts | 2 +- .../server/src/usage/usageAggregation.test.ts | 5 +- apps/server/src/usage/usagePricing.test.ts | 49 +++++- apps/server/src/usage/usagePricing.ts | 158 ++++++++++++------ apps/server/src/usage/usageScanCache.test.ts | 17 +- apps/server/src/usage/usageScanCache.ts | 62 +++++-- .../src/usage/usageTranscriptReader.test.ts | 19 ++- .../server/src/usage/usageTranscriptReader.ts | 9 +- .../usage/usageTranscriptStreaming.test.ts | 2 +- .../server/src/usage/usageTranscripts.test.ts | 27 ++- apps/server/src/usage/usageTranscripts.ts | 50 ++++-- 15 files changed, 399 insertions(+), 118 deletions(-) diff --git a/apps/server/src/usage/UsageService.test.ts b/apps/server/src/usage/UsageService.test.ts index 29fcecdc234e..127679d0a449 100644 --- a/apps/server/src/usage/UsageService.test.ts +++ b/apps/server/src/usage/UsageService.test.ts @@ -34,6 +34,7 @@ import * as UsageService from "./UsageService.ts"; const encodeUnknownJson = Schema.encodeEffect(Schema.fromJsonString(Schema.Unknown)); const encodeUnknownJsonString = Schema.encodeSync(Schema.fromJsonString(Schema.Unknown)); +const decodeUnknownJsonString = Schema.decodeSync(Schema.fromJsonString(Schema.Unknown)); function claudeLine(id: number, outputTokens: number, model = "claude-fable-5"): string { return `${JSON.stringify({ @@ -788,6 +789,98 @@ describe("UsageService", () => { }).pipe(Effect.scoped), ); + it.live( + "upgrades a v4 cache: reprices live Codex tiers, keeps deleted rollouts, leaves v4 intact", + () => + Effect.gen(function* () { + const { home, settings } = yield* setup; + const sessions = NodePath.join(home, "codex", "sessions"); + const rollout = (sessionId: string, outputTokens: number) => + [ + { type: "session_meta", payload: { id: sessionId } }, + { type: "turn_context", payload: { model: "gpt-6-astra" } }, + { + type: "event_msg", + payload: { + type: "thread_settings_applied", + thread_settings: { service_tier: "ultrafast" }, + }, + }, + { + type: "event_msg", + timestamp: "2026-08-01T10:00:00Z", + payload: { + type: "token_count", + info: { last_token_usage: { input_tokens: 0, output_tokens: outputTokens } }, + }, + }, + ] + .map((line) => encodeUnknownJsonString(line)) + .join("\n") + "\n"; + const live = NodePath.join(sessions, "live.jsonl"); + const deleted = NodePath.join(sessions, "deleted.jsonl"); + yield* Effect.promise(async () => { + await NodeFSP.mkdir(sessions, { recursive: true }); + await NodeFSP.writeFile(live, rollout("live", 10)); + await NodeFSP.writeFile(deleted, rollout("deleted", 20)); + }); + + yield* Effect.gen(function* () { + const { stateDir } = yield* ServerConfig.ServerConfig; + const cachePath = NodePath.join(stateDir, "usage-scan-cache-v5.json"); + const legacyPath = NodePath.join(stateDir, "usage-scan-cache.json"); + yield* (yield* UsageService.make).readSummary(WINDOW); + + // Rewrite the cache as a v4 server left it: every Codex record at + // speed 0 (standard), and no tier in the reducer state. + const legacy = yield* Effect.promise(async () => { + const document = decodeUnknownJsonString(await NodeFSP.readFile(cachePath, "utf8")) as { + files: Record; + }; + for (const file of Object.values(document.files)) { + file.r = file.r.map((row) => [...row.slice(0, 10), 0]); + delete file.cs.speed; + } + const text = encodeUnknownJsonString({ ...document, version: 4 }); + await NodeFSP.writeFile(legacyPath, text); + await NodeFSP.rm(cachePath); + await NodeFSP.rm(deleted); + return text; + }); + + const summary = yield* (yield* UsageService.make).readSummary(WINDOW); + // The live rollout re-parses at the ultrafast rate (10 x 6); the + // deleted one keeps its saved v4 usage at the standard rate (20 x 1). + assert.strictEqual(totalOutputTokens(summary), 30); + assert.strictEqual( + summary.buckets.reduce((sum, bucket) => sum + bucket.costUsd, 0), + 80, + ); + // A v4 server sharing this state directory still finds its own cache. + assert.strictEqual( + yield* Effect.promise(() => NodeFSP.readFile(legacyPath, "utf8")), + legacy, + ); + }).pipe( + Effect.provide( + serviceLayers({ + prefix: "usage-service-v4-upgrade-test", + home, + settings, + ratesDocument: { + "gpt-6-astra": { + input_cost_per_token: 0, + output_cost_per_token: 1, + input_cost_per_token_ultrafast: 0, + output_cost_per_token_ultrafast: 6, + }, + }, + }), + ), + ); + }).pipe(Effect.scoped), + ); + it.live("preserves saved tokens, costs and sessions after transcript cleanup and restart", () => Effect.gen(function* () { const { transcript, settings, home } = yield* setup; diff --git a/apps/server/src/usage/UsageService.ts b/apps/server/src/usage/UsageService.ts index a3b07d64d71b..026641ae4d07 100644 --- a/apps/server/src/usage/UsageService.ts +++ b/apps/server/src/usage/UsageService.ts @@ -64,7 +64,9 @@ import { decodeScanCache, dedupeWithinFile, encodeScanCache, + LEGACY_SCAN_CACHE_FILE_NAME, pruneScanCache, + SCAN_CACHE_FILE_NAME, type ScanCache, } from "./usageScanCache.ts"; import type { UsageRecord } from "./usageTranscripts.ts"; @@ -168,7 +170,8 @@ export const make = Effect.gen(function* () { }; const ratesCachePath = path.join(config.stateDir, "usage-model-rates.json"); - const scanCachePath = path.join(config.stateDir, "usage-scan-cache.json"); + const scanCachePath = path.join(config.stateDir, SCAN_CACHE_FILE_NAME); + const legacyScanCachePath = path.join(config.stateDir, LEGACY_SCAN_CACHE_FILE_NAME); let rates: RateTable = new Map(); let ratesFetchedAtMs: number | null = null; let ratesStatus: UsagePricing["status"] = "unavailable"; @@ -362,10 +365,17 @@ export const make = Effect.gen(function* () { */ const ensureScanCacheLoaded = yield* Effect.cached( Effect.gen(function* () { - const document = yield* fileSystem.readFileString(scanCachePath).pipe( - Effect.flatMap((raw) => decodeScanCacheFile(raw)), - Effect.catchCause(() => Effect.succeed(null)), - ); + const readDocument = (filePath: string) => + fileSystem.readFileString(filePath).pipe( + Effect.flatMap((raw) => decodeScanCacheFile(raw)), + Effect.catchCause(() => Effect.succeed(null)), + ); + let document = yield* readDocument(scanCachePath); + if (document === null) { + document = yield* readDocument(legacyScanCachePath); + // Write the migrated cache to its own file on the next scan. + cacheDirty = document !== null; + } if (document === null) return; for (const [path, entry] of decodeScanCache(document)) fileCache.set(path, entry); const sources = decodeCachedSources(document); diff --git a/apps/server/src/usage/antigravityUsageReader.ts b/apps/server/src/usage/antigravityUsageReader.ts index 54c00b804b7a..34a88114f42e 100644 --- a/apps/server/src/usage/antigravityUsageReader.ts +++ b/apps/server/src/usage/antigravityUsageReader.ts @@ -240,7 +240,7 @@ async function readDatabase(path: string, fallbackTimestamp: number): Promise 0 ? cost : null, - fast: false, + speed: "standard", dedupeKey: id ? `opencode:${id}` : null, }; } diff --git a/apps/server/src/usage/usageAggregation.test.ts b/apps/server/src/usage/usageAggregation.test.ts index 75435de08ff7..7ae0712c12ba 100644 --- a/apps/server/src/usage/usageAggregation.test.ts +++ b/apps/server/src/usage/usageAggregation.test.ts @@ -12,7 +12,8 @@ const rates: RateTable = new Map([ outputCostPerToken: 5e-5, cacheReadCostPerToken: 1e-6, cacheCreationCostPerToken: 1.25e-5, - fastMultiplier: 1, + fast: null, + ultrafast: null, }, ], ]); @@ -32,7 +33,7 @@ function record(overrides: Partial = {}): UsageRecord { reasoningTokens: 0, }, reportedCostUsd: null, - fast: false, + speed: "standard", dedupeKey: null, ...overrides, }; diff --git a/apps/server/src/usage/usagePricing.test.ts b/apps/server/src/usage/usagePricing.test.ts index db4f68c2cb42..badd50a4901f 100644 --- a/apps/server/src/usage/usagePricing.test.ts +++ b/apps/server/src/usage/usagePricing.test.ts @@ -1,6 +1,7 @@ import { describe, expect, it } from "@effect/vitest"; import { cursorRateModel } from "./cursorUsageReader.ts"; +import type { UsageSpeed } from "./usageTranscripts.ts"; import { cacheSavingsUsd, createOverrideRateTable, @@ -23,11 +24,15 @@ describe("usage pricing", () => { outputTokens: 1_000_000, reasoningTokens: 500_000, }; - const record = (model: string, reportedCostUsd: number | null = null, fast = false) => ({ + const record = ( + model: string, + reportedCostUsd: number | null = null, + speed: UsageSpeed = "standard", + ) => ({ model, totals, reportedCostUsd, - fast, + speed, }); it("uses custom token rates ahead of public and provider-reported costs", () => { @@ -117,20 +122,46 @@ describe("usage pricing", () => { const overrides = createOverrideRateTable({ "claude-opus-5-5": { inputCostPerMillionTokens: 4, outputCostPerMillionTokens: 20 }, }); - const cost = (model: string, fast: boolean, custom?: typeof overrides) => - priceUsage(table, record(model, null, fast), custom).costUsd; + const cost = (model: string, speed: UsageSpeed, custom?: typeof overrides) => + priceUsage(table, record(model, null, speed), custom).costUsd; - expect(cost("claude-opus-5-5", true)).toBeCloseTo(2 * cost("claude-opus-5-5", false)); - expect(cacheSavingsUsd(table, record("claude-opus-5-5", null, true))).toBeCloseTo( + expect(cost("claude-opus-5-5", "fast")).toBeCloseTo(2 * cost("claude-opus-5-5", "standard")); + expect(cacheSavingsUsd(table, record("claude-opus-5-5", null, "fast"))).toBeCloseTo( 2 * cacheSavingsUsd(table, record("claude-opus-5-5")), ); // No published fast tier, and custom prices, both stay at the standard rate. - expect(cost("claude-fable-5-1", true)).toBe(cost("claude-fable-5-1", false)); - expect(cost("claude-opus-5-5", true, overrides)).toBe( - cost("claude-opus-5-5", false, overrides), + expect(cost("claude-fable-5-1", "fast")).toBe(cost("claude-fable-5-1", "standard")); + expect(cost("claude-opus-5-5", "fast", overrides)).toBe( + cost("claude-opus-5-5", "standard", overrides), ); }); + it("prices Codex priority and ultrafast requests at their published tier rates", () => { + const table = parseRateTable({ + "gpt-6-astra": { + ...rate(1e-5, 1e-6), + input_cost_per_token_priority: 2e-5, + output_cost_per_token_priority: 1e-4, + cache_read_input_token_cost_priority: 2e-6, + input_cost_per_token_ultrafast: 6e-5, + output_cost_per_token_ultrafast: 3e-4, + // No ultrafast cache rate: keeps the standard 10:1 input-to-cache ratio. + }, + "gpt-6-sol": rate(2e-6, 2e-7), + }); + const cost = (model: string, speed: UsageSpeed) => + priceUsage(table, record(model, null, speed)).costUsd; + const standard = cost("gpt-6-astra", "standard"); + + expect(cost("gpt-6-astra", "fast")).toBeCloseTo(2 * standard); + expect(cost("gpt-6-astra", "ultrafast")).toBeCloseTo(6 * standard); + expect(cacheSavingsUsd(table, record("gpt-6-astra", null, "ultrafast"))).toBeCloseTo( + 6 * cacheSavingsUsd(table, record("gpt-6-astra")), + ); + // A tier the model does not publish bills at the standard rate. + expect(cost("gpt-6-sol", "ultrafast")).toBe(cost("gpt-6-sol", "standard")); + }); + it("keeps the canonical Fable rate separate from DeepInfra in either order", () => { const canonical = ["claude-fable-5", rate(1e-5, 1e-6)] as const; const deepInfra = ["deepinfra/anthropic/claude-fable-5", rate(1e-5)] as const; diff --git a/apps/server/src/usage/usagePricing.ts b/apps/server/src/usage/usagePricing.ts index 78f6cf2c5cd9..63190eebbf37 100644 --- a/apps/server/src/usage/usagePricing.ts +++ b/apps/server/src/usage/usagePricing.ts @@ -9,33 +9,39 @@ */ import type { UsageCostSource, UsageModelPriceOverride } from "@t3tools/contracts"; -import type { UsageRecord } from "./usageTranscripts.ts"; +import type { UsageRecord, UsageSpeed } from "./usageTranscripts.ts"; -/** - * The subset of a LiteLLM entry we price against. All values are USD per token. - * - * LiteLLM also publishes tiered variants (`*_above_272k_tokens`, `*_flex`, - * `*_priority`, `*_batches`). We deliberately price at the base tier: the - * transcripts don't record which tier served a request, so anything else would - * be a guess dressed up as precision. - */ -export interface ModelRate { +/** Token rates for one billing speed. All values are USD per token. */ +export interface TokenRates { readonly inputCostPerToken: number; readonly outputCostPerToken: number; readonly cacheReadCostPerToken: number; readonly cacheCreationCostPerToken: number; +} + +/** + * The subset of a LiteLLM entry we price against: standard rates, plus rates + * for each faster speed the model publishes. A request at a speed with no + * published rates bills at the standard rates. + * + * LiteLLM also publishes `*_above_272k_tokens`, `*_flex`, and `*_batches` + * variants. Transcripts don't record those, so we don't price them. + */ +export interface ModelRate extends TokenRates { /** - * Multiple of the rates above billed for a fast-mode request, from LiteLLM's - * `provider_specific_entry.fast`. `1` when the model publishes no fast tier. + * From LiteLLM's `provider_specific_entry.fast` multiple (Claude fast mode), + * or else its `*_priority` rates (Codex `priority`). */ - readonly fastMultiplier: number; + readonly fast: TokenRates | null; + /** From LiteLLM's `*_ultrafast` rates (Codex `ultrafast`). */ + readonly ultrafast: TokenRates | null; } export type RateTable = ReadonlyMap; /** * Custom IDs keep their case, provider prefix, and variant suffix. Custom rates - * apply as entered, fast-mode requests included. + * apply as entered, at every speed. */ export function createOverrideRateTable( overrides: Readonly>, @@ -50,31 +56,68 @@ export function createOverrideRateTable( (prices.cacheReadCostPerMillionTokens ?? prices.inputCostPerMillionTokens) / 1_000_000, cacheCreationCostPerToken: (prices.cacheWriteCostPerMillionTokens ?? prices.inputCostPerMillionTokens) / 1_000_000, - fastMultiplier: 1, + fast: null, + ultrafast: null, }, ]), ); } -/** Raw shape of one LiteLLM entry, narrowed to the fields we read. */ -interface LiteLlmEntry { - readonly input_cost_per_token?: unknown; - readonly output_cost_per_token?: unknown; - readonly cache_read_input_token_cost?: unknown; - readonly cache_creation_input_token_cost?: unknown; - readonly provider_specific_entry?: unknown; -} +/** One raw LiteLLM entry. Field names carry a tier suffix, e.g. `_priority`. */ +type LiteLlmEntry = Readonly>; function finiteNumber(value: unknown): number | null { return typeof value === "number" && Number.isFinite(value) ? value : null; } +/** + * Reads one rate set, `suffix` selecting a tier such as `_priority`. Returns + * `null` without both an input and an output rate. + * + * Anthropic bills cache reads at a discount and cache writes at a premium. + * When the standard tier omits them, cached input is priced as plain input + * rather than as free. A faster tier that omits them keeps the standard tier's + * cache-to-input ratio. + */ +function readTokenRates( + entry: LiteLlmEntry, + suffix: string, + standard?: TokenRates, +): TokenRates | null { + const input = finiteNumber(entry[`input_cost_per_token${suffix}`]); + const output = finiteNumber(entry[`output_cost_per_token${suffix}`]); + if (input === null || output === null) return null; + const cacheRate = (name: string, field: "cacheReadCostPerToken" | "cacheCreationCostPerToken") => + finiteNumber(entry[`${name}${suffix}`]) ?? + (standard !== undefined && standard.inputCostPerToken > 0 + ? (standard[field] / standard.inputCostPerToken) * input + : input); + return { + inputCostPerToken: input, + outputCostPerToken: output, + cacheReadCostPerToken: cacheRate("cache_read_input_token_cost", "cacheReadCostPerToken"), + cacheCreationCostPerToken: cacheRate( + "cache_creation_input_token_cost", + "cacheCreationCostPerToken", + ), + }; +} + +function scaleTokenRates(rates: TokenRates, multiple: number): TokenRates { + return { + inputCostPerToken: rates.inputCostPerToken * multiple, + outputCostPerToken: rates.outputCostPerToken * multiple, + cacheReadCostPerToken: rates.cacheReadCostPerToken * multiple, + cacheCreationCostPerToken: rates.cacheCreationCostPerToken * multiple, + }; +} + /** Reads `provider_specific_entry.fast`, e.g. `2` for Claude Opus 5.5. */ -function fastMultiplier(entry: LiteLlmEntry): number { - const specific = entry.provider_specific_entry; - if (typeof specific !== "object" || specific === null) return 1; +function fastMultiplier(entry: LiteLlmEntry): number | null { + const specific = entry["provider_specific_entry"]; + if (typeof specific !== "object" || specific === null) return null; const fast = finiteNumber((specific as Record)["fast"]); - return fast !== null && fast > 0 ? fast : 1; + return fast !== null && fast > 0 ? fast : null; } /** @@ -94,21 +137,19 @@ export function parseRateTable(document: unknown): RateTable { for (const [name, raw] of Object.entries(document as Record)) { if (typeof raw !== "object" || raw === null) continue; const entry = raw as LiteLlmEntry; - const input = finiteNumber(entry.input_cost_per_token); - const output = finiteNumber(entry.output_cost_per_token); - if (input === null || output === null) continue; + const standard = readTokenRates(entry, ""); + if (standard === null) continue; const key = normalizeRateKey(name); if (key.length === 0) continue; + const multiple = fastMultiplier(entry); table.set(key, { - inputCostPerToken: input, - outputCostPerToken: output, - // Anthropic bills cache reads at a discount and cache writes at a - // premium. When a model omits them, cached input is priced as plain - // input rather than as free. - cacheReadCostPerToken: finiteNumber(entry.cache_read_input_token_cost) ?? input, - cacheCreationCostPerToken: finiteNumber(entry.cache_creation_input_token_cost) ?? input, - fastMultiplier: fastMultiplier(entry), + ...standard, + fast: + multiple === null + ? readTokenRates(entry, "_priority", standard) + : scaleTokenRates(standard, multiple), + ultrafast: readTokenRates(entry, "_ultrafast", standard), }); } @@ -131,16 +172,29 @@ export function parseRateTable(document: unknown): RateTable { return table; } -function sameRate(a: ModelRate, b: ModelRate): boolean { +function sameTokenRates(a: TokenRates | null, b: TokenRates | null): boolean { + if (a === null || b === null) return a === b; return ( a.inputCostPerToken === b.inputCostPerToken && a.outputCostPerToken === b.outputCostPerToken && a.cacheReadCostPerToken === b.cacheReadCostPerToken && - a.cacheCreationCostPerToken === b.cacheCreationCostPerToken && - a.fastMultiplier === b.fastMultiplier + a.cacheCreationCostPerToken === b.cacheCreationCostPerToken ); } +function sameRate(a: ModelRate, b: ModelRate): boolean { + return ( + sameTokenRates(a, b) && + sameTokenRates(a.fast, b.fast) && + sameTokenRates(a.ultrafast, b.ultrafast) + ); +} + +/** The rates a request at `speed` bills at. */ +function ratesAt(rate: ModelRate, speed: UsageSpeed): TokenRates { + return (speed === "standard" ? null : rate[speed]) ?? rate; +} + function normalizeRateKey(model: string): string { return model.trim().toLowerCase(); } @@ -186,7 +240,7 @@ export function lookupRate(table: RateTable, model: string): ModelRate | null { /** The parts of a transcript record that decide its price. */ export type PricedRecord = Pick< UsageRecord, - "model" | "rateModel" | "totals" | "fast" | "reportedCostUsd" + "model" | "rateModel" | "totals" | "speed" | "reportedCostUsd" >; export interface PricedUsage { @@ -214,14 +268,13 @@ export function priceUsage( const rate = override ?? lookupRate(table, record.rateModel ?? model); if (rate === null) return { costUsd: 0, costSource: "unpriced" }; - const standardCostUsd = - totals.uncachedInputTokens * rate.inputCostPerToken + - totals.cachedInputTokens * rate.cacheReadCostPerToken + - totals.cacheCreationTokens * rate.cacheCreationCostPerToken + - totals.outputTokens * rate.outputCostPerToken; - + const rates = ratesAt(rate, record.speed); return { - costUsd: standardCostUsd * (record.fast ? rate.fastMultiplier : 1), + costUsd: + totals.uncachedInputTokens * rates.inputCostPerToken + + totals.cachedInputTokens * rates.cacheReadCostPerToken + + totals.cacheCreationTokens * rates.cacheCreationCostPerToken + + totals.outputTokens * rates.outputCostPerToken, costSource: "modelPriced", }; } @@ -238,9 +291,6 @@ export function cacheSavingsUsd( const rate = overrides?.get(record.model.trim()) ?? lookupRate(table, record.rateModel ?? record.model); if (rate === null) return 0; - return ( - record.totals.cachedInputTokens * - (rate.inputCostPerToken - rate.cacheReadCostPerToken) * - (record.fast ? rate.fastMultiplier : 1) - ); + const rates = ratesAt(rate, record.speed); + return record.totals.cachedInputTokens * (rates.inputCostPerToken - rates.cacheReadCostPerToken); } diff --git a/apps/server/src/usage/usageScanCache.test.ts b/apps/server/src/usage/usageScanCache.test.ts index 6455b1eb7770..82ff1e36dfc1 100644 --- a/apps/server/src/usage/usageScanCache.test.ts +++ b/apps/server/src/usage/usageScanCache.test.ts @@ -24,7 +24,7 @@ function record(overrides: Partial = {}): UsageRecord { reasoningTokens: 0, }, reportedCostUsd: null, - fast: false, + speed: "standard", dedupeKey: "msg_1:", ...overrides, }; @@ -61,7 +61,7 @@ describe("scan cache round trip", () => { [ "/a.jsonl", 100, - [record(), record({ dedupeKey: "msg_2:", model: "claude-opus-5-5", fast: true })], + [record(), record({ dedupeKey: "msg_2:", model: "claude-opus-5-5", speed: "fast" })], ], ["/b.jsonl", 200, [record({ sessionId: "session-b", reportedCostUsd: 1.5 })]], ]); @@ -79,11 +79,14 @@ describe("scan cache round trip", () => { size: 80, mtimeMs: 400, provider: "codex", - records: [record({ provider: "codex", model: "gpt-5.2-codex", dedupeKey: null })], + records: [ + record({ provider: "codex", model: "gpt-6-astra", dedupeKey: null, speed: "ultrafast" }), + ], tailRecords: [], position: position({ codexState: { - model: "gpt-5.2-codex", + model: "gpt-6-astra", + speed: "ultrafast", sessionId: "session-c", lastUsageSignature: '{"input_tokens":1}', sawSessionMeta: true, @@ -128,8 +131,8 @@ describe("scan cache round trip", () => { expect(decodeScanCache(JSON.parse(JSON.stringify(poisoned))).has("/a.jsonl")).toBe(false); }); - it("drops an entry whose fast flag is not 0 or 1", () => { - const encoded = encodeScanCache(cacheWith([["/a.jsonl", 100, [record({ fast: true })]]])); + it("drops an entry whose speed is not a known index", () => { + const encoded = encodeScanCache(cacheWith([["/a.jsonl", 100, [record({ speed: "fast" })]]])); const row = encoded.files["/a.jsonl"]!.r[0]!; const poisoned = { ...encoded, @@ -139,7 +142,7 @@ describe("scan cache round trip", () => { expect(decodeScanCache(JSON.parse(JSON.stringify(poisoned))).has("/a.jsonl")).toBe(false); }); - it("rejects a document from the previous cache version", () => { + it("rejects a document from before records carried a speed", () => { const encoded = encodeScanCache(cacheWith([["/a.jsonl", 100, [record()]]])); const previous = { ...encoded, version: 3 }; diff --git a/apps/server/src/usage/usageScanCache.ts b/apps/server/src/usage/usageScanCache.ts index 79cca5cff8e2..3c6826a67986 100644 --- a/apps/server/src/usage/usageScanCache.ts +++ b/apps/server/src/usage/usageScanCache.ts @@ -17,14 +17,33 @@ import type { UsageProviderKind } from "@t3tools/contracts"; import { GUARD_LENGTH, type TranscriptParsePosition } from "./usageTranscriptReader.ts"; -import type { CodexScanState, UsageRecord } from "./usageTranscripts.ts"; +import type { CodexScanState, UsageRecord, UsageSpeed } from "./usageTranscripts.ts"; // v2: Codex fork-copy suppression changed what a file parses to, so v1 // entries would keep serving double-counted records forever. // v3: entries carry the parse position and reducer state so a grown file // re-parses only its appended bytes instead of starting over. // v4: records carry Claude fast mode, which v3 rows never captured. -const USAGE_SCAN_CACHE_VERSION = 4 as const; +// v5: Codex records carry their service tier. v4 rows store speed the same +// way, so v4 entries still load; see `decodeScanCache` for v4 Codex entries. +const USAGE_SCAN_CACHE_VERSION = 5 as const; +const SPEED_COMPATIBLE_SINCE_VERSION = 4; + +/** + * Each cache version writes its own file in the state directory. An older + * server sharing that directory cannot read a newer cache and would replace + * it, dropping saved usage for deleted transcripts. Separate files keep both. + * A v5 server reads the legacy (v4) file once, when its own file is missing. + */ +export const SCAN_CACHE_FILE_NAME = "usage-scan-cache-v5.json"; +export const LEGACY_SCAN_CACHE_FILE_NAME = "usage-scan-cache.json"; + +/** Serialised as the index into this list. */ +const SPEEDS: readonly UsageSpeed[] = ["standard", "fast", "ultrafast"]; + +function isSpeed(value: unknown): value is UsageSpeed { + return SPEEDS.some((speed) => speed === value); +} export interface CachedFile { readonly size: number; @@ -59,7 +78,7 @@ type SerializedRecord = readonly [ reasoningTokens: number, dedupeKey: string | null, reportedCostUsd: number | null, - fast: 0 | 1, + speed: number, ]; interface SerializedFile { @@ -111,7 +130,7 @@ export function encodeScanCache(cache: ScanCache): SerializedCache { record.totals.reasoningTokens, record.dedupeKey, record.reportedCostUsd, - record.fast ? 1 : 0, + SPEEDS.indexOf(record.speed), ]; const files: Record = {}; @@ -147,7 +166,14 @@ export function decodeScanCache(document: unknown): ScanCache { if (typeof document !== "object" || document === null) return cache; const root = document as Partial; - if (root.version !== USAGE_SCAN_CACHE_VERSION) return cache; + const version = root.version; + if ( + typeof version !== "number" || + version < SPEED_COMPATIBLE_SINCE_VERSION || + version > USAGE_SCAN_CACHE_VERSION + ) { + return cache; + } if (!isRecordArray(root.models) || !isRecordArray(root.sessions)) return cache; if (typeof root.files !== "object" || root.files === null) return cache; @@ -180,8 +206,9 @@ export function decodeScanCache(document: unknown): ScanCache { reasoning, dedupeKey, reportedCostUsd, - fast, + speedIndex, ] = row as SerializedRecord; + const speed = typeof speedIndex === "number" ? SPEEDS[speedIndex] : undefined; const model = typeof modelIndex === "number" ? models[modelIndex] : undefined; if ( @@ -193,7 +220,7 @@ export function decodeScanCache(document: unknown): ScanCache { !Number.isFinite(cacheCreation) || !Number.isFinite(output) || !Number.isFinite(reasoning) || - (fast !== 0 && fast !== 1) + speed === undefined ) { return null; } @@ -211,7 +238,7 @@ export function decodeScanCache(document: unknown): ScanCache { reasoningTokens: reasoning, }, reportedCostUsd: typeof reportedCostUsd === "number" ? reportedCostUsd : null, - fast: fast === 1, + speed, dedupeKey: typeof dedupeKey === "string" ? dedupeKey : null, }); } @@ -242,7 +269,11 @@ export function decodeScanCache(document: unknown): ScanCache { ) { continue; } - const codexState = decodeCodexState(entry.cs); + // v4 Codex records predate service tiers, so they all priced as standard. + // Keep them, because the rollout may be gone, but make a live rollout + // re-parse whole: no file has size -1, and a zero position cannot resume. + const legacyCodex = entry.p === "codex" && version < USAGE_SCAN_CACHE_VERSION; + const codexState = legacyCodex ? null : decodeCodexState(entry.cs); if (codexState === undefined) continue; const provider: UsageProviderKind = entry.p; @@ -251,17 +282,14 @@ export function decodeScanCache(document: unknown): ScanCache { if (records === null || tailRecords === null) continue; cache.set(path, { - size: entry.s, + size: legacyCodex ? -1 : entry.s, mtimeMs: entry.m, provider, records, tailRecords, - position: { - resumeOffset: entry.o, - guardLength: entry.gl, - guardHash: entry.gh, - codexState, - }, + position: legacyCodex + ? { resumeOffset: 0, guardLength: 0, guardHash: 0, codexState: null } + : { resumeOffset: entry.o, guardLength: entry.gl, guardHash: entry.gh, codexState }, }); } @@ -279,6 +307,7 @@ function decodeCodexState(value: unknown): CodexScanState | null | undefined { const state = value as Partial; if ( typeof state.model !== "string" || + !isSpeed(state.speed) || typeof state.sessionId !== "string" || (state.lastUsageSignature !== null && typeof state.lastUsageSignature !== "string") || typeof state.sawSessionMeta !== "boolean" || @@ -290,6 +319,7 @@ function decodeCodexState(value: unknown): CodexScanState | null | undefined { } return { model: state.model, + speed: state.speed, sessionId: state.sessionId, lastUsageSignature: state.lastUsageSignature ?? null, sawSessionMeta: state.sawSessionMeta, diff --git a/apps/server/src/usage/usageTranscriptReader.test.ts b/apps/server/src/usage/usageTranscriptReader.test.ts index 6c95a36df19b..e55c01ad4919 100644 --- a/apps/server/src/usage/usageTranscriptReader.test.ts +++ b/apps/server/src/usage/usageTranscriptReader.test.ts @@ -76,6 +76,14 @@ function codexModelLine(model: string): string { })}\n`; } +function codexTierLine(serviceTier: string): string { + return `${JSON.stringify({ + type: "event_msg", + timestamp: "2026-08-01T10:00:01Z", + payload: { type: "thread_settings_applied", thread_settings: { service_tier: serviceTier } }, + })}\n`; +} + function codexUsageLine(outputTokens: number, secondsOffset: number): string { return `${JSON.stringify({ type: "event_msg", @@ -111,19 +119,24 @@ describe("readTranscriptRecords resume", () => { it("carries the Codex reducer state across the resume boundary", async () => { const path = NodePath.join(dir, "rollout.jsonl"); - await NodeFSP.writeFile(path, codexMetaLine() + codexModelLine("gpt-5.2-codex")); + await NodeFSP.writeFile( + path, + codexMetaLine() + codexModelLine("gpt-5.2-codex") + codexTierLine("ultrafast"), + ); const first = await readTranscriptRecords(path, "codex"); assert.isNotNull(first); assert.strictEqual(first.records.length, 0); - // The appended usage event has no turn_context or session_meta of its own; - // model and session must come from the state captured before the boundary. + // The appended usage event has no turn_context, thread settings, or + // session_meta of its own; model, tier, and session must come from the + // state captured before the boundary. await NodeFSP.appendFile(path, codexUsageLine(9, 5)); const second = await readTranscriptRecords(path, "codex", first.position); assert.isNotNull(second); assert.isTrue(second.resumed); assert.strictEqual(second.records.length, 1); assert.strictEqual(second.records[0]?.model, "gpt-5.2-codex"); + assert.strictEqual(second.records[0]?.speed, "ultrafast"); assert.strictEqual(second.records[0]?.sessionId, "codex-session-1"); }); diff --git a/apps/server/src/usage/usageTranscriptReader.ts b/apps/server/src/usage/usageTranscriptReader.ts index faa686990777..4d170237568d 100644 --- a/apps/server/src/usage/usageTranscriptReader.ts +++ b/apps/server/src/usage/usageTranscriptReader.ts @@ -108,6 +108,7 @@ const USAGE_FIELDS: Record<"claude" | "codex" | "grok", SelectedFields> = { id: true, session_id: true, model: true, + thread_settings: { service_tier: true }, forked_from_id: true, source: { subagent: { thread_spawn: { parent_thread_id: true } } }, info: { last_token_usage: true }, @@ -245,9 +246,10 @@ async function guardMatches( * still match, so only appended lines are read; otherwise the whole file is * re-parsed from the start and `resumed` reports `false`. * - * Codex carries the active model on `turn_context` lines that hold no usage of - * their own, so those still have to pass through the reducer to keep model - * attribution correct. + * Codex carries the active model on `turn_context` lines and the service tier + * on `thread_settings_applied` lines. Neither holds usage of its own, but both + * still have to pass through the reducer to keep attribution and pricing + * correct. */ export async function readTranscriptRecords( filePath: string, @@ -283,6 +285,7 @@ export async function readTranscriptRecords( if ( !mightCarryUsage(line, provider) && !line.includes('"turn_context"') && + !line.includes('"thread_settings_applied"') && !line.includes('"session_meta"') ) { return; diff --git a/apps/server/src/usage/usageTranscriptStreaming.test.ts b/apps/server/src/usage/usageTranscriptStreaming.test.ts index 2187fd4aa431..fd23593e5523 100644 --- a/apps/server/src/usage/usageTranscriptStreaming.test.ts +++ b/apps/server/src/usage/usageTranscriptStreaming.test.ts @@ -122,7 +122,7 @@ describe("large usage records", () => { reasoningTokens: 0, }, reportedCostUsd: 0.25, - fast: true, + speed: "fast", dedupeKey: "m1:r-m1", }, ]); diff --git a/apps/server/src/usage/usageTranscripts.test.ts b/apps/server/src/usage/usageTranscripts.test.ts index ace3b7d18cf7..8d668a139822 100644 --- a/apps/server/src/usage/usageTranscripts.test.ts +++ b/apps/server/src/usage/usageTranscripts.test.ts @@ -53,15 +53,15 @@ describe("parseClaudeLine", () => { reasoningTokens: 0, }); expect(record?.dedupeKey).toBe("msg_1:"); - expect(record?.fast).toBe(false); + expect(record?.speed).toBe("standard"); }); it("marks fast-mode requests", () => { const line = (speed: string) => parseClaudeLine(claudeLine({ messageId: "msg_1", contentType: "text", speed })); - expect(line("fast")?.fast).toBe(true); - expect(line("standard")?.fast).toBe(false); + expect(line("fast")?.speed).toBe("fast"); + expect(line("standard")?.speed).toBe("standard"); }); it("gives every content block of one message the same dedupe key", () => { @@ -148,6 +148,27 @@ describe("parseCodexLine", () => { expect(parseCodexLine(tokenCount(100, 0, 10, 0), state)).not.toBeNull(); }); + it("carries the service tier from the latest thread settings", () => { + const settings = (thread_settings: Record) => + JSON.stringify({ + type: "event_msg", + timestamp: "2026-08-01T05:17:42.000Z", + payload: { type: "thread_settings_applied", thread_settings }, + }); + const state = initialCodexScanState(); + parseCodexLine(turnContext, state); + const speedAfter = (line: string, output: number) => { + parseCodexLine(line, state); + return parseCodexLine(tokenCount(100, 0, output, 0), state)?.speed; + }; + + expect(parseCodexLine(tokenCount(100, 0, 1, 0), state)?.speed).toBe("standard"); + expect(speedAfter(settings({ service_tier: "ultrafast" }), 2)).toBe("ultrafast"); + expect(speedAfter(settings({ service_tier: "priority" }), 3)).toBe("fast"); + // Codex omits the field when no tier was requested. + expect(speedAfter(settings({ model: "gpt-6-astra" }), 4)).toBe("standard"); + }); + // A forked/subagent rollout opens with the parent's history copied in and // every line re-stamped to the fork instant, then the ancestors' session // metas. Counting those again multiplied usage ~1.85x on real data (#5758). diff --git a/apps/server/src/usage/usageTranscripts.ts b/apps/server/src/usage/usageTranscripts.ts index c3903ef47194..eb32e664f436 100644 --- a/apps/server/src/usage/usageTranscripts.ts +++ b/apps/server/src/usage/usageTranscripts.ts @@ -8,6 +8,13 @@ */ import type { UsageProviderKind, UsageTokenTotals } from "@t3tools/contracts"; +/** + * Billing speed of a request. Faster speeds bill at a model-specific premium. + * Claude fast mode and Codex `priority` are `fast`; Codex `ultrafast` is its + * own, more expensive tier. + */ +export type UsageSpeed = "standard" | "fast" | "ultrafast"; + export interface UsageRecord { readonly provider: UsageProviderKind; readonly timestampMs: number; @@ -20,11 +27,8 @@ export interface UsageRecord { readonly sessionId: string; readonly totals: UsageTokenTotals; readonly reportedCostUsd: number | null; - /** - * Whether the request ran in fast mode, which bills at a model-specific - * multiple of the standard rate. Only Claude Code records this. - */ - readonly fast: boolean; + /** Only Claude Code and Codex record a speed; other providers are `standard`. */ + readonly speed: UsageSpeed; /** * Key for cross-file de-duplication, or `null` when the record is inherently * unique and needs no dedup. @@ -159,7 +163,7 @@ export function parseClaudeRecord(parsed: unknown): UsageRecord | null { reasoningTokens: 0, }, reportedCostUsd: typeof cost === "number" && Number.isFinite(cost) ? cost : null, - fast: usageRecord["speed"] === "fast", + speed: usageRecord["speed"] === "fast" ? "fast" : "standard", dedupeKey, }; } @@ -171,12 +175,14 @@ export function parseClaudeRecord(parsed: unknown): UsageRecord | null { /** * Rolling state for a single Codex rollout file. * - * Codex `token_count` events carry no model, so the model is carried forward - * from the most recent `turn_context`. Sessions that switch models mid-run - * attribute correctly from the switch onward. + * Codex `token_count` events carry no model or service tier, so both are + * carried forward: the model from the most recent `turn_context`, the tier from + * the most recent `thread_settings_applied`. Sessions that switch either + * mid-run attribute correctly from the switch onward. */ export interface CodexScanState { model: string; + speed: UsageSpeed; sessionId: string; lastUsageSignature: string | null; sawSessionMeta: boolean; @@ -188,6 +194,7 @@ export interface CodexScanState { export function initialCodexScanState(): CodexScanState { return { model: "", + speed: "standard", sessionId: "", lastUsageSignature: null, sawSessionMeta: false, @@ -265,6 +272,14 @@ export function parseCodexRecord(parsed: unknown, state: CodexScanState): UsageR return null; } + if (payloadType === "thread_settings_applied") { + const settings = payloadRecord["thread_settings"]; + if (typeof settings === "object" && settings !== null) { + state.speed = codexSpeed((settings as Record)["service_tier"]); + } + return null; + } + if (payloadType !== "token_count") return null; const info = payloadRecord["info"]; @@ -323,13 +338,24 @@ export function parseCodexRecord(parsed: unknown, state: CodexScanState): UsageR totals, // Codex does not report cost in the rollout. reportedCostUsd: null, - fast: false, + speed: state.speed, // Events surviving the fork-copy suppression above are unique to this // rollout, so they need no global dedup. dedupeKey: null, }; } +/** + * Maps a Codex `service_tier` to its billing speed. Codex omits the field when + * no tier was requested, which bills as standard, as do `default` and + * `standard`. `fast` is accepted as an alias of `priority`. + */ +function codexSpeed(serviceTier: unknown): UsageSpeed { + if (serviceTier === "priority" || serviceTier === "fast") return "fast"; + if (serviceTier === "ultrafast") return "ultrafast"; + return "standard"; +} + /* -------------------------------------------------------------------------- */ /* Grok Build */ /* -------------------------------------------------------------------------- */ @@ -457,7 +483,7 @@ export function parseGrokRecord(parsed: unknown): readonly UsageRecord[] { sessionId, totals: grokTotalsToUsage(topLevel), reportedCostUsd: grokCostTicksToUsd(topLevel.costUsdTicks), - fast: false, + speed: "standard", // No prompt id means we cannot tell two same-second updates apart. dedupeKey: promptId === null ? null : `${sessionId}:${promptId}:grok`, }, @@ -504,7 +530,7 @@ export function parseGrokRecord(parsed: unknown): readonly UsageRecord[] { sessionId, totals, reportedCostUsd, - fast: false, + speed: "standard", dedupeKey: promptId === null ? null : `${sessionId}:${promptId}:${entry.model}`, }); } From a5b34b25378fbbf90ce36ee101a11528ccf11c7e Mon Sep 17 00:00:00 2001 From: Theo Browne Date: Sat, 3 Oct 2026 02:07:12 -0700 Subject: [PATCH 03/42] fix(clients): a dev server left running no longer says the thread is waiting (#15114) Co-authored-by: Claude Opus 5.5 (1M context) --- apps/web/src/components/ChatView.tsx | 9 +++-- .../src/state/threadExecution.test.ts | 30 +++++++++++++++-- .../src/state/threadExecution.ts | 33 ++++++++++++++----- 3 files changed, 60 insertions(+), 12 deletions(-) diff --git a/apps/web/src/components/ChatView.tsx b/apps/web/src/components/ChatView.tsx index 8e105a49f37f..fe72ec446d34 100644 --- a/apps/web/src/components/ChatView.tsx +++ b/apps/web/src/components/ChatView.tsx @@ -6859,7 +6859,7 @@ export default function ChatView(props: ChatViewProps) { // The stack renders items[0] front-most and tucks the rest behind hover, so // ordering is priority: system banners, then the branch-mismatch notice, // and the informational parked-thread banner last — it must never cover another. - // Background work (subagent fleets, workflow runs, watch loops) can outlive + // Background work (subagent fleets, workflow runs, watch loops, dev servers) can outlive // the turn; once it settles, the composer stop button is gone, so this // banner is the only visible stop affordance. The interrupt path also // accepts a completed run while its provider still has background work. @@ -6907,9 +6907,14 @@ export default function ChatView(props: ChatViewProps) { id: `background-work:${activeThread.id}`, variant: "default", priority: "activity", + // A dev server can run for hours after the agent is done, so only work + // that will wake the agent pulses. icon: (