diff --git a/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.test.ts b/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.test.ts index d61739f72c21..c1093bff37a1 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 @@ -271,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"], { @@ -322,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). @@ -355,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)), @@ -433,6 +472,7 @@ describe("ProviderRuntimeIngestion", () => { advanceClock: (ms: number) => { clockOffsetMs += ms; }, + emitAndWaitForEnqueue, emitAndDrain, sqlCount: sqlCounter.count, setProviderSession: provider.setSession, @@ -4093,16 +4133,277 @@ 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("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; 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 +4426,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 +4489,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..cb8cbdc7486a 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"; @@ -31,6 +32,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 +129,24 @@ type TurnStartRequestedDomainEvent = Extract< { type: "thread.turn-start-requested" } >; -type ProviderDiffEvent = Extract; +type ProviderDiffEvent = Pick< + Extract, + "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 { + 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 = | { @@ -142,6 +161,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 { @@ -1029,7 +1049,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}`)), ); @@ -2640,7 +2660,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)), + ); } }; @@ -2677,18 +2699,30 @@ 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 }); + // 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({ + 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; });