From aa01c965b6c08db30521453f66f9c50b705f1936 Mon Sep 17 00:00:00 2001 From: ScottN-PV Date: Wed, 23 Sep 2026 16:57:08 -0400 Subject: [PATCH 1/6] fix(server): stop retaining queued Codex diff snapshots Drop unused snapshot text before the runtime notification queue and queue only checkpoint identity in ingestion. Coalesce repeated thread/turn signals through lifecycle processing, preserving completion ordering. Fixes #12883 Co-authored-by: Codex AI-Tool: OpenAI Codex AI-Harness: Codex agent runtime; specific integration/version not exposed AI-Host: T3 Code (user-confirmed) AI-Model: gpt-6-astra (GPT-6 Astra; user-confirmed) AI-Reasoning: medium (user-confirmed) AI-Contribution: Investigation, implementation, regression tests, automated verification, and drafting --- .../Layers/ProviderRuntimeIngestion.test.ts | 44 +++++++++++++++++-- .../Layers/ProviderRuntimeIngestion.ts | 39 +++++++++++++--- .../Layers/CodexSessionRuntime.test.ts | 14 ++++++ .../provider/Layers/CodexSessionRuntime.ts | 18 +++++++- .../shared/src/KeyedCoalescingWorker.test.ts | 44 +++++++++++++++++++ packages/shared/src/KeyedCoalescingWorker.ts | 14 +++++- 6 files changed, 161 insertions(+), 12 deletions(-) diff --git a/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.test.ts b/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.test.ts index d61739f72c21..6606c38d78e5 100644 --- a/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.test.ts +++ b/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.test.ts @@ -21,6 +21,7 @@ import { ProjectId, ProviderItemId, RuntimeRequestId, + RuntimeItemId, type ServerSettings, ThreadId, TurnId, @@ -60,6 +61,7 @@ import * as ThreadBackgroundLiveness from "../ThreadBackgroundLiveness.ts"; import * as ThreadPlanProgress from "../ThreadPlanProgress.ts"; import { ProviderRuntimeIngestionLive, + providerDiffSignal, splitBufferedAssistantText, } from "./ProviderRuntimeIngestion.ts"; import { DEFAULT_THREAD_TITLE } from "../threadTitles.ts"; @@ -239,6 +241,29 @@ async function waitForThread( } describe("ProviderRuntimeIngestion", () => { + it("retains only checkpoint identity when queuing a provider diff", () => { + const identity = { + type: "turn.diff.updated" as const, + eventId: asEventId("large-diff"), + threadId: ThreadId.make("thread-1"), + turnId: TurnId.make("turn-1"), + itemId: RuntimeItemId.make("item-1"), + createdAt: "2026-01-01T00:00:00.000Z", + }; + const diff = "x".repeat(4_000_000); + const signal = providerDiffSignal({ + ...identity, + provider: ProviderDriverKind.make("codex"), + payload: { unifiedDiff: diff }, + raw: { + source: "codex.app-server.notification", + method: "turn/diff/updated", + payload: { diff }, + }, + }); + expect(signal).toEqual(identity); + }); + let runtime: ManagedRuntime.ManagedRuntime< OrchestrationEngineService | ProviderRuntimeIngestionService | ProjectionSnapshotQuery, unknown @@ -4095,14 +4120,17 @@ describe("ProviderRuntimeIngestion", () => { effectIt.effect("settles the turn while repository detection for a diff is blocked", () => Effect.gen(function* () { + let detectionCount = 0; const detectionStarted = yield* Deferred.make(); const releaseDetection = yield* Deferred.make(); const harness = yield* Effect.promise(() => createHarness({ - isGitRepository: () => - Deferred.succeed(detectionStarted, undefined).pipe( + isGitRepository: () => { + detectionCount++; + return Deferred.succeed(detectionStarted, undefined).pipe( Effect.andThen(Deferred.await(releaseDetection)), - ), + ); + }, }), ); yield* Effect.addFinalizer(() => Deferred.succeed(releaseDetection, true)); @@ -4125,6 +4153,15 @@ describe("ProviderRuntimeIngestion", () => { }); yield* Deferred.await(detectionStarted); + for (let index = 0; index < 200; index++) { + harness.emit({ + ...base, + type: "turn.diff.updated", + eventId: asEventId(`evt-repeated-diff-${index}`), + payload: { unifiedDiff: "x".repeat(10_000) }, + }); + } + const settled = yield* harness.engine.streamDomainEvents.pipe( Stream.filter( (event) => @@ -4179,6 +4216,7 @@ describe("ProviderRuntimeIngestion", () => { yield* Fiber.join(nextTurnStarted); yield* Deferred.succeed(releaseDetection, true); yield* Effect.promise(harness.drain); + expect(detectionCount).toBeLessThanOrEqual(2); const released = yield* Effect.promise(harness.readModel); expect(released.threads[0]?.checkpoints).toEqual([]); expect(released.threads[0]?.latestTurn).toMatchObject({ diff --git a/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.ts b/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.ts index 0db70e491235..a619fba438c1 100644 --- a/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.ts +++ b/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.ts @@ -31,6 +31,7 @@ import * as Option from "effect/Option"; import * as Predicate from "effect/Predicate"; import * as Stream from "effect/Stream"; import { makeDrainableWorker } from "@t3tools/shared/DrainableWorker"; +import { makeKeyedCoalescingWorker } from "@t3tools/shared/KeyedCoalescingWorker"; import { formatTokens } from "@t3tools/shared/usageFormat"; import { ProviderService } from "../../provider/Services/ProviderService.ts"; @@ -127,7 +128,23 @@ type TurnStartRequestedDomainEvent = Extract< { type: "thread.turn-start-requested" } >; -type ProviderDiffEvent = Extract; +type ProviderDiffEvent = Pick< + Extract, + "type" | "eventId" | "threadId" | "turnId" | "itemId" | "createdAt" +>; + +export function providerDiffSignal( + event: Extract, +): ProviderDiffEvent { + return { + type: event.type, + eventId: event.eventId, + threadId: event.threadId, + createdAt: event.createdAt, + ...(event.turnId !== undefined ? { turnId: event.turnId } : {}), + ...(event.itemId !== undefined ? { itemId: event.itemId } : {}), + }; +} type RuntimeIngestionInput = | { @@ -1029,7 +1046,7 @@ const make = Effect.gen(function* () { const projectionThreadActivityRepository = yield* ProjectionThreadActivityRepository; const serverSettingsService = yield* ServerSettingsService; const checkpointStore = yield* CheckpointStore.CheckpointStore; - const providerCommandId = (event: ProviderRuntimeEvent, tag: string) => + const providerCommandId = (event: Pick, tag: string) => crypto.randomUUIDv4.pipe( Effect.map((uuid) => CommandId.make(`provider:${event.eventId}:${tag}:${uuid}`)), ); @@ -2678,17 +2695,27 @@ const make = Effect.gen(function* () { const workspaceCwd = checkpointContext?.worktreePath ?? checkpointContext?.workspaceRoot; if (!workspaceCwd || !(yield* checkpointStore.isGitRepository(workspaceCwd))) return; yield* worker.enqueue({ source: "diff", event }); + // Keep this key active until its lifecycle work has finished, so a slow + // lifecycle worker cannot accumulate one queued signal per snapshot either. + yield* worker.drain; + }); + const diffWorker = yield* makeKeyedCoalescingWorker({ + merge: (_current: ProviderDiffEvent, next: ProviderDiffEvent) => next, + process: (_key: string, event: ProviderDiffEvent) => + detectProviderDiffRepository(event).pipe(logIngestionFailure("diff", event)), }); - const diffWorker = yield* makeDrainableWorker((event: ProviderDiffEvent) => - detectProviderDiffRepository(event).pipe(logIngestionFailure("diff", event)), - ); const start: ProviderRuntimeIngestionShape["start"] = () => Effect.gen(function* () { yield* forkParked( Stream.runForEach(providerService.streamEvents, (event) => event.type === "turn.diff.updated" - ? diffWorker.enqueue(event) + ? event.turnId === undefined + ? Effect.void + : diffWorker.enqueue( + JSON.stringify([event.threadId, event.turnId]), + providerDiffSignal(event), + ) : worker.enqueue({ source: "runtime", event }), ), ); diff --git a/apps/server/src/provider/Layers/CodexSessionRuntime.test.ts b/apps/server/src/provider/Layers/CodexSessionRuntime.test.ts index ec113ab7c521..bec16d70a361 100644 --- a/apps/server/src/provider/Layers/CodexSessionRuntime.test.ts +++ b/apps/server/src/provider/Layers/CodexSessionRuntime.test.ts @@ -17,6 +17,7 @@ import { hasConfiguredMcpServer, isRecoverableThreadResumeError, makeMemoryConsolidationNotificationFilter, + makeCodexServerNotification, openCodexThread, readCodexThread, rollbackCodexThread, @@ -24,6 +25,19 @@ import { } from "./CodexSessionRuntime.ts"; const isCodexAppServerRequestError = Schema.is(CodexErrors.CodexAppServerRequestError); +describe("Codex notification queue payloads", () => { + it("discards large diff snapshots before buffering notifications", () => { + const params = { threadId: "thread-1", turnId: "turn-1", diff: "x".repeat(4_000_000) }; + const notification = makeCodexServerNotification("turn/diff/updated", params); + NodeAssert.ok(JSON.stringify(notification).length < 200); + NodeAssert.deepEqual(notification, { + method: "turn/diff/updated", + params: { threadId: "thread-1", turnId: "turn-1", diff: "" }, + }); + NodeAssert.equal(params.diff.length, 4_000_000); + }); +}); + describe("Codex thread history", () => { for (const numTurns of [1, 2, 3, 5]) { it.effect(`reverts ${numTurns} paginated turns at the durable boundary`, () => diff --git a/apps/server/src/provider/Layers/CodexSessionRuntime.ts b/apps/server/src/provider/Layers/CodexSessionRuntime.ts index 674d23327b65..309438b0bdd6 100644 --- a/apps/server/src/provider/Layers/CodexSessionRuntime.ts +++ b/apps/server/src/provider/Layers/CodexSessionRuntime.ts @@ -479,11 +479,25 @@ type CodexServerNotification = { }; }[CodexRpc.ServerNotificationMethod]; -function makeCodexServerNotification( +export function makeCodexServerNotification( method: M, params: CodexRpc.ServerNotificationParamsByMethod[M], ): CodexServerNotification { - return { method, params } as CodexServerNotification; + const notification = { method, params } as CodexServerNotification; + if (notification.method === "turn/diff/updated") { + // Ingestion only needs this signal to record a placeholder checkpoint. + // Drop snapshot text before the first queue, retaining the native shape + // for routing and schema validation downstream. + return { + method: notification.method, + params: { + threadId: notification.params.threadId, + turnId: notification.params.turnId, + diff: "", + }, + }; + } + return notification; } function normalizeCodexModelSlug( diff --git a/packages/shared/src/KeyedCoalescingWorker.test.ts b/packages/shared/src/KeyedCoalescingWorker.test.ts index 8bfc1a340a60..cc170c205031 100644 --- a/packages/shared/src/KeyedCoalescingWorker.test.ts +++ b/packages/shared/src/KeyedCoalescingWorker.test.ts @@ -2,10 +2,54 @@ import { it } from "@effect/vitest"; import { describe, expect } from "vite-plus/test"; import * as Deferred from "effect/Deferred"; import * as Effect from "effect/Effect"; +import * as Fiber from "effect/Fiber"; import { makeKeyedCoalescingWorker } from "./KeyedCoalescingWorker.ts"; describe("makeKeyedCoalescingWorker", () => { + it.live("coalesces a burst while draining all active and queued keys", () => + Effect.scoped( + Effect.gen(function* () { + const started = yield* Deferred.make(); + const release = yield* Deferred.make(); + const otherStarted = yield* Deferred.make(); + const releaseOther = yield* Deferred.make(); + const processed: string[] = []; + const worker = yield* makeKeyedCoalescingWorker({ + merge: (_current: number, next: number) => next, + process: (key: string, value: number) => + Effect.gen(function* () { + processed.push(`${key}:${value}`); + if (value === 0) { + yield* Deferred.succeed(started, undefined); + yield* Deferred.await(release); + } + if (key === "other-turn") { + yield* Deferred.succeed(otherStarted, undefined); + yield* Deferred.await(releaseOther); + } + }), + }); + yield* worker.enqueue("turn", 0); + yield* Deferred.await(started); + for (let index = 1; index <= 200; index++) yield* worker.enqueue("turn", index); + yield* worker.enqueue("other-turn", 1); + const drained = yield* Deferred.make(); + const waiter = yield* worker.drain.pipe( + Effect.andThen(Deferred.succeed(drained, undefined)), + Effect.forkChild({ startImmediately: true }), + ); + expect(yield* Deferred.isDone(drained)).toBe(false); + yield* Deferred.succeed(release, undefined); + yield* Deferred.await(otherStarted); + expect(yield* Deferred.isDone(drained)).toBe(false); + yield* Deferred.succeed(releaseOther, undefined); + yield* Fiber.join(waiter); + expect(processed).toEqual(["turn:0", "turn:200", "other-turn:1"]); + }), + ), + ); + it.live("waits for latest work enqueued during active processing before draining the key", () => Effect.scoped( Effect.gen(function* () { diff --git a/packages/shared/src/KeyedCoalescingWorker.ts b/packages/shared/src/KeyedCoalescingWorker.ts index f4edebfafb30..35fbf05d8e3a 100644 --- a/packages/shared/src/KeyedCoalescingWorker.ts +++ b/packages/shared/src/KeyedCoalescingWorker.ts @@ -15,6 +15,8 @@ import * as TxRef from "effect/TxRef"; export interface KeyedCoalescingWorker { readonly enqueue: (key: K, value: V) => Effect.Effect; readonly drainKey: (key: K) => Effect.Effect; + /** Wait until every key has finished, including work merged while active. */ + readonly drain: Effect.Effect; } interface KeyedCoalescingWorkerState { @@ -138,5 +140,15 @@ export const makeKeyedCoalescingWorker = (options: { Effect.tx, ); - return { enqueue, drainKey } satisfies KeyedCoalescingWorker; + const drain = TxRef.get(stateRef).pipe( + Effect.tap((state) => + state.latestByKey.size > 0 || state.queuedKeys.size > 0 || state.activeKeys.size > 0 + ? Effect.txRetry + : Effect.void, + ), + Effect.asVoid, + Effect.tx, + ); + + return { enqueue, drainKey, drain } satisfies KeyedCoalescingWorker; }); From 28649d7f41dd931707f4e186e88555fd96ba1abc Mon Sep 17 00:00:00 2001 From: ScottN-PV Date: Wed, 23 Sep 2026 19:33:55 -0400 Subject: [PATCH 2/6] fix(server): acknowledge individual queued diffs Release the active diff key when its own lifecycle work finishes, even on failure, instead of waiting for unrelated lifecycle work to drain. Add a two-thread regression with a blocked lifecycle event. Co-authored-by: Codex AI-Tool: OpenAI Codex AI-Harness: Codex agent runtime; specific integration/version not exposed AI-Host: T3 Code (user-confirmed) AI-Model: not exposed for this follow-up session AI-Reasoning: not exposed for this follow-up session AI-Contribution: Review analysis, implementation, regression test, automated verification, and drafting --- .../Layers/ProviderRuntimeIngestion.test.ts | 115 +++++++++++++++++- .../Layers/ProviderRuntimeIngestion.ts | 12 +- 2 files changed, 123 insertions(+), 4 deletions(-) diff --git a/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.test.ts b/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.test.ts index 6606c38d78e5..c3856ed89a73 100644 --- a/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.test.ts +++ b/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.test.ts @@ -296,6 +296,7 @@ describe("ProviderRuntimeIngestion", () => { threadTitle?: string; workspaceSubdirectory?: string; isGitRepository?: CheckpointStore.CheckpointStore["Service"]["isGitRepository"]; + beforeDispatch?: (command: OrchestrationCommand) => Effect.Effect; }) { const repositoryRoot = makeTempDir("t3-provider-project-"); NodeChildProcess.execFileSync("git", ["init", "--initial-branch=main"], { @@ -347,7 +348,18 @@ describe("ProviderRuntimeIngestion", () => { }; const layer = ProviderRuntimeIngestionLive.pipe( Layer.provide(Layer.succeed(Clock.Clock, shiftedClock)), - Layer.provideMerge(orchestrationLayer), + Layer.provideMerge( + Layer.effect( + OrchestrationEngineService, + Effect.map(OrchestrationEngineService, (engine) => ({ + ...engine, + dispatch: (command, dispatchOptions) => + (options?.beforeDispatch?.(command) ?? Effect.void).pipe( + Effect.andThen(engine.dispatch(command, dispatchOptions)), + ), + })), + ).pipe(Layer.provide(orchestrationLayer)), + ), Layer.provideMerge(ingestionProjectionSnapshotLayer), // Single shared liveness instance across ingestion (writer), the // engine, and the snapshot query (reader). @@ -4118,6 +4130,107 @@ describe("ProviderRuntimeIngestion", () => { }); }); + effectIt.effect("advances another diff key while unrelated lifecycle work is blocked", () => + Effect.gen(function* () { + const firstDiffStarted = yield* Deferred.make(); + const releaseFirstDiff = yield* Deferred.make(); + const lifecycleBlocked = yield* Deferred.make(); + const releaseLifecycle = yield* Deferred.make(); + const secondDetection = yield* Deferred.make(); + let detectionCount = 0; + let blockLifecycle = false; + const harness = yield* Effect.promise(() => + createHarness({ + isGitRepository: () => + Effect.gen(function* () { + if (++detectionCount === 2) yield* Deferred.succeed(secondDetection, undefined); + return true; + }), + beforeDispatch: (command) => { + if (command.type === "thread.turn.diff.complete" && command.threadId === "thread-1") { + return Deferred.succeed(firstDiffStarted, undefined).pipe( + Effect.andThen(Deferred.await(releaseFirstDiff)), + ); + } + if (blockLifecycle && command.type === "thread.session.set") { + return Deferred.succeed(lifecycleBlocked, undefined).pipe( + Effect.andThen(Deferred.await(releaseLifecycle)), + ); + } + return Effect.void; + }, + }), + ); + yield* Effect.addFinalizer(() => + Deferred.succeed(releaseFirstDiff, undefined).pipe( + Effect.andThen(Deferred.succeed(releaseLifecycle, undefined)), + ), + ); + const createdAt = "2026-01-01T00:00:00.000Z"; + yield* Effect.promise(() => + harness.dispatch({ + type: "thread.create", + commandId: CommandId.make("cmd-second-diff-thread"), + threadId: asThreadId("thread-2"), + projectId: asProjectId("project-1"), + title: "Second thread", + modelSelection: { instanceId: ProviderInstanceId.make("codex"), model: "gpt-5-codex" }, + interactionMode: DEFAULT_PROVIDER_INTERACTION_MODE, + runtimeMode: "approval-required", + branch: null, + worktreePath: null, + createdAt, + }), + ); + const first = { + provider: ProviderDriverKind.make("codex"), + threadId: asThreadId("thread-1"), + turnId: asTurnId("turn-1"), + createdAt, + }; + const second = { ...first, threadId: asThreadId("thread-2"), turnId: asTurnId("turn-2") }; + yield* Effect.promise(() => + harness.emitAndDrain([ + { ...first, type: "turn.started", eventId: asEventId("start-first") }, + { ...second, type: "turn.started", eventId: asEventId("start-second") }, + ]), + ); + harness.emit({ + ...first, + type: "turn.diff.updated", + eventId: asEventId("diff-first"), + payload: { unifiedDiff: "first" }, + }); + yield* Deferred.await(firstDiffStarted); + blockLifecycle = true; + harness.emit({ + ...first, + type: "session.state.changed", + eventId: asEventId("unrelated-state"), + payload: { state: "running" }, + }); + harness.emit({ + ...second, + type: "turn.diff.updated", + eventId: asEventId("diff-second"), + payload: { unifiedDiff: "second" }, + }); + yield* Deferred.succeed(releaseFirstDiff, undefined); + yield* Deferred.await(lifecycleBlocked); + // The second key must reach detection before the lifecycle queue becomes idle. + yield* Deferred.await(secondDetection); + yield* Deferred.succeed(releaseLifecycle, undefined); + yield* Effect.promise(harness.drain); + const snapshot = yield* Effect.promise(harness.readModel); + expect( + snapshot.threads.find((thread) => thread.id === first.threadId)?.checkpoints, + ).toHaveLength(1); + expect( + snapshot.threads.find((thread) => thread.id === second.threadId)?.checkpoints, + ).toHaveLength(1); + }), + ); + effectIt.effect("settles the turn while repository detection for a diff is blocked", () => Effect.gen(function* () { let detectionCount = 0; diff --git a/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.ts b/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.ts index a619fba438c1..8a8fdc594b68 100644 --- a/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.ts +++ b/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.ts @@ -24,6 +24,7 @@ import * as Cause from "effect/Cause"; import * as Clock from "effect/Clock"; import * as Crypto from "effect/Crypto"; import * as DateTime from "effect/DateTime"; +import * as Deferred from "effect/Deferred"; import * as Duration from "effect/Duration"; import * as Effect from "effect/Effect"; import * as Layer from "effect/Layer"; @@ -159,6 +160,7 @@ type RuntimeIngestionInput = /** A diff whose workspace the diff worker confirmed is a Git repository. */ source: "diff"; event: ProviderDiffEvent; + processed: Deferred.Deferred; }; function toTurnId(value: TurnId | string | undefined): TurnId | undefined { @@ -2657,7 +2659,9 @@ const make = Effect.gen(function* () { case "domain": return processDomainEvent(input.event); case "diff": - return recordProviderDiff(input.event); + return recordProviderDiff(input.event).pipe( + Effect.ensuring(Deferred.succeed(input.processed, undefined)), + ); } }; @@ -2694,10 +2698,12 @@ const make = Effect.gen(function* () { .pipe(Effect.map(Option.getOrUndefined)); const workspaceCwd = checkpointContext?.worktreePath ?? checkpointContext?.workspaceRoot; if (!workspaceCwd || !(yield* checkpointStore.isGitRepository(workspaceCwd))) return; - yield* worker.enqueue({ source: "diff", event }); + const processed = yield* Deferred.make(); + yield* worker.enqueue({ source: "diff", event, processed }); // Keep this key active until its lifecycle work has finished, so a slow // lifecycle worker cannot accumulate one queued signal per snapshot either. - yield* worker.drain; + // Other keys can advance without waiting for unrelated lifecycle work. + yield* Deferred.await(processed); }); const diffWorker = yield* makeKeyedCoalescingWorker({ merge: (_current: ProviderDiffEvent, next: ProviderDiffEvent) => next, From c8000459fc9d97f851ed9dcc2435b7fbad7f0f4a Mon Sep 17 00:00:00 2001 From: ScottN-PV Date: Sat, 26 Sep 2026 19:49:16 -0400 Subject: [PATCH 3/6] docs(server): describe diff queue helper contracts Document the existing behavior without changing executable code. Co-authored-by: Codex Co-authored-by: Claude AI-Tool: OpenAI Codex AI-Harness: Codex harness (integration/version not exposed) AI-Host: T3 Code AI-Model: gpt-6-astra AI-Reasoning: medium AI-Contribution: JSDoc drafting, source-equivalence checks, targeted lint, diff checks, and PR update preparation AI-Tool: Claude Code AI-Harness: Claude Code CLI 2.1.283 invoked by Codex harness AI-Host: T3 Code via PowerShell AI-Model: claude-fable-5-1 AI-Reasoning: high AI-Contribution: Independent review, docstring wording improvements, and follow-up drafting --- .../src/orchestration/Layers/ProviderRuntimeIngestion.ts | 1 + apps/server/src/provider/Layers/CodexSessionRuntime.ts | 1 + packages/shared/src/KeyedCoalescingWorker.ts | 5 +++++ 3 files changed, 7 insertions(+) diff --git a/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.ts b/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.ts index 8a8fdc594b68..45ab62a8d8ba 100644 --- a/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.ts +++ b/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.ts @@ -134,6 +134,7 @@ type ProviderDiffEvent = Pick< "type" | "eventId" | "threadId" | "turnId" | "itemId" | "createdAt" >; +/** Keeps checkpoint routing metadata without retaining the provider's diff snapshot in the queue. */ export function providerDiffSignal( event: Extract, ): ProviderDiffEvent { diff --git a/apps/server/src/provider/Layers/CodexSessionRuntime.ts b/apps/server/src/provider/Layers/CodexSessionRuntime.ts index 309438b0bdd6..e724147d5135 100644 --- a/apps/server/src/provider/Layers/CodexSessionRuntime.ts +++ b/apps/server/src/provider/Layers/CodexSessionRuntime.ts @@ -479,6 +479,7 @@ type CodexServerNotification = { }; }[CodexRpc.ServerNotificationMethod]; +/** Strips diff snapshot text before queueing while preserving the native notification shape. */ export function makeCodexServerNotification( method: M, params: CodexRpc.ServerNotificationParamsByMethod[M], diff --git a/packages/shared/src/KeyedCoalescingWorker.ts b/packages/shared/src/KeyedCoalescingWorker.ts index 35fbf05d8e3a..245788809f85 100644 --- a/packages/shared/src/KeyedCoalescingWorker.ts +++ b/packages/shared/src/KeyedCoalescingWorker.ts @@ -25,6 +25,7 @@ interface KeyedCoalescingWorkerState { readonly activeKeys: Set; } +/** Creates a scoped, serial worker that merges pending values per key and exposes per-key and whole-worker drains. */ export const makeKeyedCoalescingWorker = (options: { readonly merge: (current: V, next: V) => V; readonly process: (key: K, value: V) => Effect.Effect; @@ -37,6 +38,7 @@ export const makeKeyedCoalescingWorker = (options: { activeKeys: new Set(), }); + /** Processes a key until no merged value remains, then marks it idle. */ const processKey = (key: K, value: V): Effect.Effect => options.process(key, value).pipe( Effect.flatMap(() => @@ -58,6 +60,7 @@ export const makeKeyedCoalescingWorker = (options: { ), ); + /** Releases a failed key and requeues any value that arrived while it was active. */ const cleanupFailedKey = (key: K): Effect.Effect => TxRef.modify(stateRef, (state) => { const activeKeys = new Set(state.activeKeys); @@ -110,6 +113,7 @@ export const makeKeyedCoalescingWorker = (options: { Effect.forkScoped, ); + /** Merges a pending value atomically, queueing the key only when it is not already scheduled. */ const enqueue: KeyedCoalescingWorker["enqueue"] = (key, value) => TxRef.modify(stateRef, (state) => { const latestByKey = new Map(state.latestByKey); @@ -129,6 +133,7 @@ export const makeKeyedCoalescingWorker = (options: { Effect.asVoid, ); + /** Waits for this key's queued, active, and merged work without waiting for other keys. */ const drainKey: KeyedCoalescingWorker["drainKey"] = (key) => TxRef.get(stateRef).pipe( Effect.tap((state) => From fca74174f3c89885294d4db65f1b641a601069ea Mon Sep 17 00:00:00 2001 From: ScottN-PV Date: Mon, 28 Sep 2026 22:32:29 -0400 Subject: [PATCH 4/6] docs(server): document diff ingestion lifecycle Co-authored-by: Codex Co-authored-by: Claude AI-Tool: OpenAI Codex AI-Harness: Codex harness (integration/version not exposed) AI-Host: T3 Code AI-Model: gpt-6-astra AI-Reasoning: medium AI-Contribution: Docstring drafting, source-equivalence checks, formatting, targeted lint, and publishing preparation AI-Tool: Claude Code AI-Harness: Claude Code CLI 2.1.283 invoked by Codex harness AI-Host: T3 Code via PowerShell AI-Model: claude-fable-5-1 AI-Reasoning: high AI-Contribution: Independent implementation and docstring review, follow-up drafting --- .../Layers/ProviderRuntimeIngestion.test.ts | 1 + .../orchestration/Layers/ProviderRuntimeIngestion.ts | 10 +++++++--- 2 files changed, 8 insertions(+), 3 deletions(-) diff --git a/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.test.ts b/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.test.ts index c3856ed89a73..090b0e80cea5 100644 --- a/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.test.ts +++ b/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.test.ts @@ -291,6 +291,7 @@ describe("ProviderRuntimeIngestion", () => { } }); + /** Creates an isolated repository and ingestion runtime with controllable dispatch and Git probes. */ async function createHarness(options?: { serverSettings?: Partial; threadTitle?: string; diff --git a/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.ts b/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.ts index 45ab62a8d8ba..1a15c69bb8b9 100644 --- a/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.ts +++ b/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.ts @@ -1036,6 +1036,7 @@ export function runtimeEventToActivities( return []; } +/** Creates scoped lifecycle and diff workers, keeping repository probes off the lifecycle queue. */ const make = Effect.gen(function* () { const threadBackgroundLiveness = yield* ThreadBackgroundLivenessService; const threadPlanProgress = yield* ThreadPlanProgressService; @@ -2653,6 +2654,7 @@ const make = Effect.gen(function* () { }); }); + /** Dispatches a queued event and acknowledges each diff when its lifecycle work finishes. */ const processInput = (input: RuntimeIngestionInput) => { switch (input.source) { case "runtime": @@ -2687,9 +2689,10 @@ const make = Effect.gen(function* () { processInput(input).pipe(logIngestionFailure(input.source, input.event)), ); - // Repository detection for a diff goes through VCS subprocesses, which can - // stall behind slow or hung git. It runs on its own worker so a stuck diff - // never delays the lifecycle worker; confirmed diffs are handed back to it. + /** + * Probes the repository off the lifecycle queue, then waits only for this diff's + * lifecycle work so unrelated queued events cannot prevent other diff keys advancing. + */ const detectProviderDiffRepository = Effect.fn("detectProviderDiffRepository")(function* ( event: ProviderDiffEvent, ) { @@ -2712,6 +2715,7 @@ const make = Effect.gen(function* () { detectProviderDiffRepository(event).pipe(logIngestionFailure("diff", event)), }); + /** Subscribes to runtime and domain events, coalescing diff signals per thread and turn. */ const start: ProviderRuntimeIngestionShape["start"] = () => Effect.gen(function* () { yield* forkParked( From 8287cd1d034d6469be489ed1f674c471e3799b93 Mon Sep 17 00:00:00 2001 From: ScottN-PV Date: Mon, 28 Sep 2026 22:41:21 -0400 Subject: [PATCH 5/6] docs(server): document adjacent diff review helpers Co-authored-by: Codex Co-authored-by: Claude AI-Tool: OpenAI Codex AI-Harness: Codex harness (integration/version not exposed) AI-Host: T3 Code AI-Model: gpt-6-astra AI-Reasoning: medium AI-Contribution: Docstring drafting, source-equivalence checks, formatting, targeted lint, and publishing preparation AI-Tool: Claude Code AI-Harness: Claude Code CLI 2.1.283 invoked by Codex harness AI-Host: T3 Code via PowerShell AI-Model: claude-fable-5-1 AI-Reasoning: high AI-Contribution: Independent implementation and docstring review, follow-up drafting --- .../src/orchestration/Layers/ProviderRuntimeIngestion.test.ts | 1 + .../server/src/orchestration/Layers/ProviderRuntimeIngestion.ts | 2 ++ apps/server/src/provider/Layers/CodexSessionRuntime.ts | 1 + 3 files changed, 4 insertions(+) diff --git a/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.test.ts b/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.test.ts index 090b0e80cea5..9a584778be38 100644 --- a/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.test.ts +++ b/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.test.ts @@ -218,6 +218,7 @@ type ProviderRuntimeTestProposedPlan = ProviderRuntimeTestThread["proposedPlans" type ProviderRuntimeTestActivity = ProviderRuntimeTestThread["activities"][number]; type ProviderRuntimeTestCheckpoint = ProviderRuntimeTestThread["checkpoints"][number]; +/** Waits for a thread snapshot to satisfy the predicate, failing when the deadline expires. */ async function waitForThread( readModel: () => Promise, predicate: (thread: ProviderRuntimeTestThread) => boolean, diff --git a/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.ts b/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.ts index 1a15c69bb8b9..b6b386a71db0 100644 --- a/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.ts +++ b/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.ts @@ -164,6 +164,7 @@ type RuntimeIngestionInput = processed: Deferred.Deferred; }; +/** Normalizes a present provider turn identifier while preserving a missing identifier. */ function toTurnId(value: TurnId | string | undefined): TurnId | undefined { return value === undefined ? undefined : TurnId.make(String(value)); } @@ -486,6 +487,7 @@ function taskLinkageActivityFields(payload: Record): Record Date: Thu, 1 Oct 2026 04:13:28 -0400 Subject: [PATCH 6/6] fix(server): test a lifecycle backlog ahead of a coalesced diff Add a test that blocks lifecycle work queued ahead of the active diff and checks that repeated snapshots stay merged and every thread still gets one checkpoint. Correct the queue comment: a diff waits for its own lifecycle item, not for the whole queue. Remove doc comments on functions this fix does not change. Co-Authored-By: Claude Opus 5.5 (1M context) --- .../Layers/ProviderRuntimeIngestion.test.ts | 162 +++++++++++++++++- .../Layers/ProviderRuntimeIngestion.ts | 18 +- .../provider/Layers/CodexSessionRuntime.ts | 2 - packages/shared/src/KeyedCoalescingWorker.ts | 5 - 4 files changed, 166 insertions(+), 21 deletions(-) diff --git a/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.test.ts b/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.test.ts index 9a584778be38..c1093bff37a1 100644 --- a/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.test.ts +++ b/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.test.ts @@ -218,7 +218,6 @@ type ProviderRuntimeTestProposedPlan = ProviderRuntimeTestThread["proposedPlans" type ProviderRuntimeTestActivity = ProviderRuntimeTestThread["activities"][number]; type ProviderRuntimeTestCheckpoint = ProviderRuntimeTestThread["checkpoints"][number]; -/** Waits for a thread snapshot to satisfy the predicate, failing when the deadline expires. */ async function waitForThread( readModel: () => Promise, predicate: (thread: ProviderRuntimeTestThread) => boolean, @@ -292,7 +291,6 @@ describe("ProviderRuntimeIngestion", () => { } }); - /** Creates an isolated repository and ingestion runtime with controllable dispatch and Git probes. */ async function createHarness(options?: { serverSettings?: Partial; threadTitle?: string; @@ -394,6 +392,8 @@ describe("ProviderRuntimeIngestion", () => { const drain = () => testRuntime.runPromise(ingestion.drain); const dispatch = (command: OrchestrationCommand) => testRuntime.runPromise(engine.dispatch(command)); + const emitAndWaitForEnqueue = (events: ReadonlyArray) => + testRuntime.runPromise(provider.emitAndWaitForEnqueue(events)); const emitAndDrain = (events: ReadonlyArray) => testRuntime.runPromise( provider.emitAndWaitForEnqueue(events).pipe(Effect.andThen(ingestion.drain)), @@ -472,6 +472,7 @@ describe("ProviderRuntimeIngestion", () => { advanceClock: (ms: number) => { clockOffsetMs += ms; }, + emitAndWaitForEnqueue, emitAndDrain, sqlCount: sqlCounter.count, setProviderSession: provider.setSession, @@ -4233,6 +4234,163 @@ describe("ProviderRuntimeIngestion", () => { }), ); + effectIt.effect("keeps a diff key coalesced behind blocked lifecycle work", () => + Effect.gen(function* () { + const firstDetection = yield* Deferred.make(); + const lifecycleBlocked = yield* Deferred.make(); + const releaseLifecycle = yield* Deferred.make(); + let detectionCount = 0; + let blockLifecycle = false; + const dispatched: string[] = []; + const harness = yield* Effect.promise(() => + createHarness({ + isGitRepository: () => + Effect.gen(function* () { + if (++detectionCount === 1) yield* Deferred.succeed(firstDetection, undefined); + return true; + }), + beforeDispatch: (command) => { + if (!blockLifecycle) return Effect.void; + if ( + command.type === "thread.session.set" || + command.type === "thread.message.assistant.complete" || + command.type === "thread.turn.diff.complete" + ) { + dispatched.push(`${command.type}:${command.threadId}`); + } + return command.type === "thread.session.set" + ? Deferred.succeed(lifecycleBlocked, undefined).pipe( + Effect.andThen(Deferred.await(releaseLifecycle)), + ) + : Effect.void; + }, + }), + ); + yield* Effect.addFinalizer(() => Deferred.succeed(releaseLifecycle, undefined)); + const createdAt = "2026-01-01T00:00:00.000Z"; + yield* Effect.promise(() => + harness.dispatch({ + type: "thread.create", + commandId: CommandId.make("cmd-backlog-thread"), + threadId: asThreadId("thread-2"), + projectId: asProjectId("project-1"), + title: "Second thread", + modelSelection: { instanceId: ProviderInstanceId.make("codex"), model: "gpt-5-codex" }, + interactionMode: DEFAULT_PROVIDER_INTERACTION_MODE, + runtimeMode: "approval-required", + branch: null, + worktreePath: null, + createdAt, + }), + ); + const first = { + provider: ProviderDriverKind.make("codex"), + threadId: asThreadId("thread-1"), + turnId: asTurnId("turn-1"), + createdAt, + }; + const second = { ...first, threadId: asThreadId("thread-2"), turnId: asTurnId("turn-2") }; + yield* Effect.promise(() => + harness.emitAndDrain([ + { ...first, type: "turn.started", eventId: asEventId("start-first") }, + { ...second, type: "turn.started", eventId: asEventId("start-second") }, + ]), + ); + + blockLifecycle = true; + harness.emit({ + ...first, + type: "session.state.changed", + eventId: asEventId("blocked-state"), + payload: { state: "running" }, + }); + yield* Deferred.await(lifecycleBlocked); + yield* Effect.promise(() => + harness.emitAndWaitForEnqueue([ + { + ...second, + type: "item.completed", + eventId: asEventId("backlog-reply"), + itemId: asItemId("backlog-reply"), + payload: { + itemType: "assistant_message", + status: "completed", + detail: "Second reply.", + }, + }, + ]), + ); + // Lifecycle backlog is now queued, so this diff's lifecycle item lands behind it. + yield* Effect.promise(() => + harness.emitAndWaitForEnqueue([ + { + ...first, + type: "turn.diff.updated", + eventId: asEventId("diff-first-0"), + payload: { unifiedDiff: "first 0" }, + }, + ]), + ); + yield* Deferred.await(firstDetection); + for (let index = 1; index <= 20; index++) { + yield* Effect.promise(() => + harness.emitAndWaitForEnqueue([ + { + ...first, + type: "turn.diff.updated", + eventId: asEventId(`diff-first-${index}`), + payload: { unifiedDiff: `first ${index}` }, + }, + ]), + ); + } + yield* Effect.promise(() => + harness.emitAndWaitForEnqueue([ + { + ...second, + type: "turn.diff.updated", + eventId: asEventId("diff-second"), + payload: { unifiedDiff: "second" }, + }, + ]), + ); + // The active key waits for its own lifecycle item, so repeated signals merge + // and the second key stays queued behind it. + expect(detectionCount).toBe(1); + + yield* Deferred.succeed(releaseLifecycle, undefined); + yield* Effect.promise(harness.drain); + // One extra probe for the merged signal, then one for the second key. + expect(detectionCount).toBe(3); + expect(dispatched).toEqual([ + "thread.session.set:thread-1", + "thread.message.assistant.complete:thread-2", + "thread.turn.diff.complete:thread-1", + "thread.turn.diff.complete:thread-2", + ]); + const snapshot = yield* Effect.promise(harness.readModel); + const firstThread = snapshot.threads.find((thread) => thread.id === first.threadId); + const secondThread = snapshot.threads.find((thread) => thread.id === second.threadId); + expect(firstThread?.checkpoints).toEqual([ + expect.objectContaining({ + turnId: first.turnId, + checkpointRef: "provider-diff:diff-first-0", + checkpointTurnCount: 1, + }), + ]); + expect(secondThread?.checkpoints).toEqual([ + expect.objectContaining({ + turnId: second.turnId, + checkpointRef: "provider-diff:diff-second", + checkpointTurnCount: 1, + }), + ]); + expect(secondThread?.messages).toEqual( + expect.arrayContaining([expect.objectContaining({ text: "Second reply." })]), + ); + }), + ); + effectIt.effect("settles the turn while repository detection for a diff is blocked", () => Effect.gen(function* () { let detectionCount = 0; diff --git a/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.ts b/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.ts index b6b386a71db0..cb8cbdc7486a 100644 --- a/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.ts +++ b/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.ts @@ -164,7 +164,6 @@ type RuntimeIngestionInput = processed: Deferred.Deferred; }; -/** Normalizes a present provider turn identifier while preserving a missing identifier. */ function toTurnId(value: TurnId | string | undefined): TurnId | undefined { return value === undefined ? undefined : TurnId.make(String(value)); } @@ -487,7 +486,6 @@ function taskLinkageActivityFields(payload: Record): Record { switch (input.source) { case "runtime": @@ -2691,10 +2687,9 @@ const make = Effect.gen(function* () { processInput(input).pipe(logIngestionFailure(input.source, input.event)), ); - /** - * Probes the repository off the lifecycle queue, then waits only for this diff's - * lifecycle work so unrelated queued events cannot prevent other diff keys advancing. - */ + // Repository detection for a diff goes through VCS subprocesses, which can + // stall behind slow or hung git. It runs on its own worker so a stuck diff + // never delays the lifecycle worker; confirmed diffs are handed back to it. const detectProviderDiffRepository = Effect.fn("detectProviderDiffRepository")(function* ( event: ProviderDiffEvent, ) { @@ -2706,9 +2701,9 @@ const make = Effect.gen(function* () { if (!workspaceCwd || !(yield* checkpointStore.isGitRepository(workspaceCwd))) return; const processed = yield* Deferred.make(); yield* worker.enqueue({ source: "diff", event, processed }); - // Keep this key active until its lifecycle work has finished, so a slow - // lifecycle worker cannot accumulate one queued signal per snapshot either. - // Other keys can advance without waiting for unrelated lifecycle work. + // Hold this key until the lifecycle worker has handled this diff. Snapshots that + // arrive meanwhile merge into one pending signal for the turn instead of queueing + // a lifecycle item each. Awaiting worker.drain here would stall other diff keys. yield* Deferred.await(processed); }); const diffWorker = yield* makeKeyedCoalescingWorker({ @@ -2717,7 +2712,6 @@ const make = Effect.gen(function* () { detectProviderDiffRepository(event).pipe(logIngestionFailure("diff", event)), }); - /** Subscribes to runtime and domain events, coalescing diff signals per thread and turn. */ const start: ProviderRuntimeIngestionShape["start"] = () => Effect.gen(function* () { yield* forkParked( diff --git a/apps/server/src/provider/Layers/CodexSessionRuntime.ts b/apps/server/src/provider/Layers/CodexSessionRuntime.ts index 4448b2ea6d7c..309438b0bdd6 100644 --- a/apps/server/src/provider/Layers/CodexSessionRuntime.ts +++ b/apps/server/src/provider/Layers/CodexSessionRuntime.ts @@ -479,7 +479,6 @@ type CodexServerNotification = { }; }[CodexRpc.ServerNotificationMethod]; -/** Strips diff snapshot text before queueing while preserving the native notification shape. */ export function makeCodexServerNotification( method: M, params: CodexRpc.ServerNotificationParamsByMethod[M], @@ -501,7 +500,6 @@ export function makeCodexServerNotification { readonly activeKeys: Set; } -/** Creates a scoped, serial worker that merges pending values per key and exposes per-key and whole-worker drains. */ export const makeKeyedCoalescingWorker = (options: { readonly merge: (current: V, next: V) => V; readonly process: (key: K, value: V) => Effect.Effect; @@ -38,7 +37,6 @@ export const makeKeyedCoalescingWorker = (options: { activeKeys: new Set(), }); - /** Processes a key until no merged value remains, then marks it idle. */ const processKey = (key: K, value: V): Effect.Effect => options.process(key, value).pipe( Effect.flatMap(() => @@ -60,7 +58,6 @@ export const makeKeyedCoalescingWorker = (options: { ), ); - /** Releases a failed key and requeues any value that arrived while it was active. */ const cleanupFailedKey = (key: K): Effect.Effect => TxRef.modify(stateRef, (state) => { const activeKeys = new Set(state.activeKeys); @@ -113,7 +110,6 @@ export const makeKeyedCoalescingWorker = (options: { Effect.forkScoped, ); - /** Merges a pending value atomically, queueing the key only when it is not already scheduled. */ const enqueue: KeyedCoalescingWorker["enqueue"] = (key, value) => TxRef.modify(stateRef, (state) => { const latestByKey = new Map(state.latestByKey); @@ -133,7 +129,6 @@ export const makeKeyedCoalescingWorker = (options: { Effect.asVoid, ); - /** Waits for this key's queued, active, and merged work without waiting for other keys. */ const drainKey: KeyedCoalescingWorker["drainKey"] = (key) => TxRef.get(stateRef).pipe( Effect.tap((state) =>