diff --git a/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.activity.test.ts b/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.activity.test.ts index 936041038644..87aa082cf73c 100644 --- a/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.activity.test.ts +++ b/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.activity.test.ts @@ -7,7 +7,7 @@ import { } from "@t3tools/contracts"; import { describe, expect, it } from "vite-plus/test"; -import { runtimeEventToActivities } from "./ProviderRuntimeIngestion.ts"; +import { mergeBackgroundIngestion, runtimeEventToActivities } from "./ProviderRuntimeIngestion.ts"; const base = { provider: ProviderDriverKind.make("codex"), @@ -82,3 +82,231 @@ describe("runtimeEventToActivities task progress", () => { expect(usagePayload).not.toHaveProperty("status"); }); }); + +describe("mergeBackgroundIngestion", () => { + it("preserves usage when a later task progress tick only updates activity", () => { + const taskId = RuntimeTaskId.make("agent-partial-progress"); + const usageUpdate = { + source: "runtime", + event: { + ...base, + type: "task.progress", + eventId: EventId.make("evt-task-usage"), + payload: { + taskId, + description: "Agent partial progress", + typedUsage: { totalTokens: 4_200 }, + }, + } satisfies ProviderRuntimeEvent, + } as const; + const activityUpdate = { + source: "runtime", + event: { + ...base, + type: "task.progress", + eventId: EventId.make("evt-task-activity"), + payload: { + taskId, + description: "Agent partial progress", + summary: "Running tests", + lastToolName: "exec_command", + }, + } satisfies ProviderRuntimeEvent, + } as const; + + const merged = mergeBackgroundIngestion([usageUpdate], [activityUpdate]); + + expect(merged).toHaveLength(1); + expect(merged[0]?.event.type).toBe("task.progress"); + if (merged[0]?.event.type === "task.progress") { + expect(merged[0].event.payload).toEqual({ + taskId, + description: "Agent partial progress", + typedUsage: { totalTokens: 4_200 }, + summary: "Running tests", + lastToolName: "exec_command", + }); + } + }); + + it("preserves activity and max-merges a later sparse usage update", () => { + const taskId = RuntimeTaskId.make("agent-later-usage"); + const activityUpdate = { + source: "runtime", + event: { + ...base, + type: "task.progress", + eventId: EventId.make("evt-earlier-activity"), + payload: { + taskId, + description: "Agent later usage", + summary: "Running tests", + lastToolName: "exec_command", + typedUsage: { + totalTokens: 100, + inputTokens: 80, + cachedInputTokens: 30, + outputTokens: 20, + reasoningOutputTokens: 10, + toolUses: 3, + durationMs: 1_000, + }, + }, + } satisfies ProviderRuntimeEvent, + } as const; + const usageUpdate = { + source: "runtime", + event: { + ...base, + type: "task.progress", + eventId: EventId.make("evt-later-usage"), + payload: { + taskId, + description: "Agent later usage", + typedUsage: { + totalTokens: 200, + inputTokens: 70, + outputTokens: 40, + durationMs: 900, + }, + }, + } satisfies ProviderRuntimeEvent, + } as const; + + const merged = mergeBackgroundIngestion([activityUpdate], [usageUpdate]); + + expect(merged).toHaveLength(1); + expect(merged[0]?.event.type).toBe("task.progress"); + if (merged[0]?.event.type === "task.progress") { + expect(merged[0].event.payload).toEqual({ + taskId, + description: "Agent later usage", + summary: "Running tests", + lastToolName: "exec_command", + typedUsage: { + totalTokens: 200, + inputTokens: 80, + cachedInputTokens: 30, + outputTokens: 40, + reasoningOutputTokens: 10, + toolUses: 3, + durationMs: 1_000, + }, + }); + } + }); + + it("replaces stale activity fields while preserving the last usage snapshot", () => { + const taskId = RuntimeTaskId.make("agent-replaced-progress"); + const previous = { + source: "runtime", + event: { + ...base, + type: "task.progress", + eventId: EventId.make("evt-previous-progress"), + payload: { + taskId, + description: "Agent replaced progress", + summary: "Old summary", + lastToolName: "old_tool", + error: "Old error", + phases: [{ index: 0, title: "Old phase" }], + typedUsage: { totalTokens: 4_200 }, + }, + } satisfies ProviderRuntimeEvent, + } as const; + const next = { + source: "runtime", + event: { + ...base, + type: "task.progress", + eventId: EventId.make("evt-next-progress"), + payload: { + taskId, + description: "Agent replaced progress", + phases: [{ index: 1, title: "New phase" }], + typedUsage: { totalTokens: 5_000 }, + }, + } satisfies ProviderRuntimeEvent, + } as const; + + const merged = mergeBackgroundIngestion([previous], [next]); + + expect(merged).toHaveLength(1); + expect(merged[0]?.event.type).toBe("task.progress"); + if (merged[0]?.event.type === "task.progress") { + expect(merged[0].event.payload).toEqual({ + taskId, + description: "Agent replaced progress", + phases: [{ index: 1, title: "New phase" }], + typedUsage: { totalTokens: 5_000 }, + }); + } + }); + + it("preserves explicit task reactivation order", () => { + const taskId = RuntimeTaskId.make("agent-reactivated"); + const completed = { + source: "runtime", + event: { + ...base, + type: "task.completed", + eventId: EventId.make("evt-task-completed"), + payload: { taskId, status: "completed" }, + } satisfies ProviderRuntimeEvent, + } as const; + const reactivated = { + source: "runtime", + event: { + ...base, + type: "task.updated", + eventId: EventId.make("evt-task-reactivated"), + payload: { taskId, status: "running" }, + } satisfies ProviderRuntimeEvent, + } as const; + + const merged = mergeBackgroundIngestion([completed], [reactivated]); + + expect(merged.map((input) => input.event.type)).toEqual(["task.completed", "task.updated"]); + }); + + it("preserves fields from partial task status patches", () => { + const taskId = RuntimeTaskId.make("agent-partial-update"); + const statusUpdate = { + source: "runtime", + event: { + ...base, + type: "task.updated", + eventId: EventId.make("evt-task-status"), + payload: { + taskId, + status: "completed", + }, + } satisfies ProviderRuntimeEvent, + } as const; + const completionTimeUpdate = { + source: "runtime", + event: { + ...base, + type: "task.updated", + eventId: EventId.make("evt-task-ended-at"), + payload: { + taskId, + endedAt: "2026-08-06T00:00:01.000Z", + }, + } satisfies ProviderRuntimeEvent, + } as const; + + const merged = mergeBackgroundIngestion([statusUpdate], [completionTimeUpdate]); + + expect(merged).toHaveLength(1); + expect(merged[0]?.event.type).toBe("task.updated"); + if (merged[0]?.event.type === "task.updated") { + expect(merged[0].event.payload).toEqual({ + taskId, + status: "completed", + endedAt: "2026-08-06T00:00:01.000Z", + }); + } + }); +}); diff --git a/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.test.ts b/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.test.ts index 258aa010e3e6..2300e28149e0 100644 --- a/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.test.ts +++ b/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.test.ts @@ -225,6 +225,11 @@ describe("ProviderRuntimeIngestion", () => { const workspaceRoot = makeTempDir("t3-provider-project-"); NodeFS.mkdirSync(NodePath.join(workspaceRoot, ".git")); const provider = createProviderServiceHarness(); + const backgroundLiveness = ThreadBackgroundLiveness.make(); + const backgroundLivenessLayer = Layer.succeed( + ThreadBackgroundLiveness.ThreadBackgroundLivenessService, + backgroundLiveness, + ); const orchestrationLayer = OrchestrationEngineLive.pipe( Layer.provide(OrchestrationProjectionSnapshotQueryLive), Layer.provide(OrchestrationProjectionPipelineLive), @@ -242,7 +247,7 @@ describe("ProviderRuntimeIngestion", () => { Layer.provideMerge(projectionSnapshotLayer), // Single shared liveness instance across ingestion (writer), the // engine, and the snapshot query (reader). - Layer.provideMerge(ThreadBackgroundLiveness.layer), + Layer.provideMerge(backgroundLivenessLayer), Layer.provideMerge(ThreadPlanProgress.layer), Layer.provideMerge(SqlitePersistenceMemory), Layer.provideMerge(Layer.succeed(ProviderService, provider.service)), @@ -318,6 +323,8 @@ describe("ProviderRuntimeIngestion", () => { readModel: () => Effect.runPromise(snapshotQuery.getSnapshot()), emit: provider.emit, setProviderSession: provider.setSession, + getBackgroundLiveness: (threadId: ThreadId) => + backgroundLiveness.getThreadBackgroundLiveness(threadId), drain, }; } @@ -3254,6 +3261,37 @@ describe("ProviderRuntimeIngestion", () => { ).toBe("# Plan title"); }); + it("does not revive background liveness after the provider session exits", async () => { + const harness = await createHarness(); + const threadId = asThreadId("thread-1"); + + harness.emit({ + type: "session.exited", + eventId: asEventId("evt-session-exited-before-stale-task"), + provider: ProviderDriverKind.make("claudeAgent"), + createdAt: "2026-01-01T00:00:01.000Z", + threadId, + payload: {}, + }); + await harness.drain(); + expect(harness.getBackgroundLiveness(threadId)).toBeNull(); + + harness.emit({ + type: "task.progress", + eventId: asEventId("evt-stale-task-after-session-exit"), + provider: ProviderDriverKind.make("claudeAgent"), + createdAt: "2026-01-01T00:00:02.000Z", + threadId, + payload: { + taskId: "stale-task-after-exit", + description: "This task belongs to the exited session", + }, + }); + await harness.drain(); + + expect(harness.getBackgroundLiveness(threadId)).toBeNull(); + }); + it("titles task activities with the task description, including on completion", async () => { const harness = await createHarness(); const now = "2026-01-01T00:00:00.000Z"; diff --git a/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.ts b/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.ts index 03253797242e..e85435498784 100644 --- a/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.ts +++ b/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.ts @@ -18,6 +18,8 @@ import { type OrchestrationThread, type OrchestrationThreadActivity, type ProviderRuntimeEvent, + type RuntimeTaskUsage, + type TaskProgressPayload, } from "@t3tools/contracts"; import * as Cache from "effect/Cache"; import * as Cause from "effect/Cause"; @@ -27,7 +29,7 @@ import * as Effect from "effect/Effect"; import * as Layer from "effect/Layer"; import * as Option from "effect/Option"; import * as Stream from "effect/Stream"; -import { makeDrainableWorker } from "@t3tools/shared/DrainableWorker"; +import { makePriorityCoalescingWorker } from "@t3tools/shared/PriorityCoalescingWorker"; import { ProviderService } from "../../provider/Services/ProviderService.ts"; import { ProjectionTurnRepository } from "../../persistence/Services/ProjectionTurns.ts"; @@ -114,6 +116,152 @@ type RuntimeIngestionInput = event: TurnStartRequestedDomainEvent; }; +type RuntimeIngestionBatch = ReadonlyArray; + +const TASK_LIFECYCLE_ORDER = [ + "task.started", + "task.progress", + "task.updated", + "task.completed", +] as const; + +function isTaskLifecycleInput(input: RuntimeIngestionInput): input is Extract< + RuntimeIngestionInput, + { source: "runtime" } +> & { + readonly event: Extract; +} { + return ( + input.source === "runtime" && + TASK_LIFECYCLE_ORDER.some((eventType) => eventType === input.event.type) + ); +} + +function hasTaskProgressActivityState(payload: TaskProgressPayload): boolean { + return ( + payload.typedUsage === undefined || + payload.summary !== undefined || + payload.lastToolName !== undefined || + payload.status !== undefined || + payload.error !== undefined || + payload.phases !== undefined + ); +} + +function mergeTaskUsage( + previous: RuntimeTaskUsage | undefined, + next: RuntimeTaskUsage | undefined, +): RuntimeTaskUsage | undefined { + if (!previous) return next; + if (!next) return previous; + + const max = (a: number | undefined, b: number | undefined): number | undefined => + a === undefined ? b : b === undefined ? a : Math.max(a, b); + const inputTokens = max(previous.inputTokens, next.inputTokens); + const cachedInputTokens = max(previous.cachedInputTokens, next.cachedInputTokens); + const outputTokens = max(previous.outputTokens, next.outputTokens); + const reasoningOutputTokens = max(previous.reasoningOutputTokens, next.reasoningOutputTokens); + const toolUses = max(previous.toolUses, next.toolUses); + const durationMs = max(previous.durationMs, next.durationMs); + return { + totalTokens: Math.max(previous.totalTokens, next.totalTokens), + ...(inputTokens === undefined ? {} : { inputTokens }), + ...(cachedInputTokens === undefined ? {} : { cachedInputTokens }), + ...(outputTokens === undefined ? {} : { outputTokens }), + ...(reasoningOutputTokens === undefined ? {} : { reasoningOutputTokens }), + ...(toolUses === undefined ? {} : { toolUses }), + ...(durationMs === undefined ? {} : { durationMs }), + }; +} + +function mergeTaskProgressPayload( + previous: TaskProgressPayload, + next: TaskProgressPayload, +): TaskProgressPayload { + const payload = hasTaskProgressActivityState(next) ? next : { ...previous, ...next }; + const typedUsage = mergeTaskUsage(previous.typedUsage, next.typedUsage); + return typedUsage === undefined ? payload : { ...payload, typedUsage }; +} + +export function mergeBackgroundIngestion( + current: RuntimeIngestionBatch, + next: RuntimeIngestionBatch, +): RuntimeIngestionBatch { + const combined = [...current, ...next]; + if (!combined.every(isTaskLifecycleInput)) { + return next; + } + + const latestByType = new Map<(typeof TASK_LIFECYCLE_ORDER)[number], RuntimeIngestionInput>(); + // Reinsert replacements so different lifecycle types keep their latest arrival order. + // Explicit running updates can reactivate a task after a completed attempt. + for (const input of combined) { + const previous = latestByType.get(input.event.type); + if ( + input.event.type === "task.progress" && + previous?.source === "runtime" && + previous.event.type === "task.progress" + ) { + latestByType.delete(input.event.type); + latestByType.set(input.event.type, { + ...input, + event: { + ...input.event, + payload: mergeTaskProgressPayload(previous.event.payload, input.event.payload), + }, + }); + continue; + } + if ( + input.event.type === "task.updated" && + previous?.source === "runtime" && + previous.event.type === "task.updated" + ) { + latestByType.delete(input.event.type); + latestByType.set(input.event.type, { + ...input, + event: { + ...input.event, + payload: { + ...previous.event.payload, + ...input.event.payload, + }, + }, + }); + continue; + } + latestByType.delete(input.event.type); + latestByType.set(input.event.type, input); + } + return Array.from(latestByType.values()); +} + +function backgroundIngestionKey(input: RuntimeIngestionInput): string | undefined { + if (input.source !== "runtime") { + return undefined; + } + + const event = input.event; + switch (event.type) { + // These are latest-state telemetry streams. A large agent fleet can emit + // hundreds of ticks while one database write is in flight, so keeping + // every queued intermediate value only delays user-visible replies. + case "task.started": + case "task.progress": + case "task.updated": + case "task.completed": + return `task:${event.threadId}:${event.payload.taskId}`; + case "tool.progress": + return event.payload.taskId === undefined + ? undefined + : `tool-progress:${event.threadId}:${event.payload.taskId}`; + case "thread.token-usage.updated": + return `thread-token-usage:${event.threadId}`; + default: + return undefined; + } +} + function toTurnId(value: TurnId | string | undefined): TurnId | undefined { return value === undefined ? undefined : TurnId.make(String(value)); } @@ -577,12 +725,7 @@ export function runtimeEventToActivities( event.payload.description.trim().length > 0 ? { title: truncateDetail(event.payload.description, 120) } : {}; - const hasProgressState = - event.payload.typedUsage === undefined || - event.payload.summary !== undefined || - event.payload.lastToolName !== undefined || - event.payload.status !== undefined || - event.payload.error !== undefined; + const hasProgressState = hasTaskProgressActivityState(event.payload); return [ ...(hasProgressState ? [ @@ -1965,6 +2108,11 @@ const make = Effect.gen(function* () { case "task.progress": case "task.updated": case "task.completed": { + // A realtime session exit can overtake queued background task + // telemetry. Never let those stale rows resurrect sidebar liveness. + if (thread.session?.status === "stopped") { + break; + } const payload = event.payload as { taskId: string; taskType?: string; @@ -2040,13 +2188,24 @@ const make = Effect.gen(function* () { }), ); - const worker = yield* makeDrainableWorker(processInputSafely); + const worker = yield* makePriorityCoalescingWorker({ + mergeBackground: mergeBackgroundIngestion, + process: (batch) => + Effect.forEach(batch, processInputSafely, { concurrency: 1 }).pipe(Effect.asVoid), + }); + + const enqueueInput = (input: RuntimeIngestionInput) => { + const backgroundKey = backgroundIngestionKey(input); + return backgroundKey === undefined + ? worker.enqueueRealtime([input]) + : worker.enqueueBackground(backgroundKey, [input]); + }; const start: ProviderRuntimeIngestionShape["start"] = () => Effect.gen(function* () { yield* forkParked( Stream.runForEach(providerService.streamEvents, (event) => - worker.enqueue({ source: "runtime", event }), + enqueueInput({ source: "runtime", event }), ), ); yield* forkParked( @@ -2054,7 +2213,7 @@ const make = Effect.gen(function* () { if (event.type !== "thread.turn-start-requested") { return Effect.void; } - return worker.enqueue({ source: "domain", event }); + return enqueueInput({ source: "domain", event }); }), ); }); diff --git a/packages/shared/package.json b/packages/shared/package.json index f669bd0a452c..d60d0d7817a3 100644 --- a/packages/shared/package.json +++ b/packages/shared/package.json @@ -59,6 +59,10 @@ "types": "./src/KeyedCoalescingWorker.ts", "import": "./src/KeyedCoalescingWorker.ts" }, + "./PriorityCoalescingWorker": { + "types": "./src/PriorityCoalescingWorker.ts", + "import": "./src/PriorityCoalescingWorker.ts" + }, "./schemaJson": { "types": "./src/schemaJson.ts", "import": "./src/schemaJson.ts" diff --git a/packages/shared/src/PriorityCoalescingWorker.test.ts b/packages/shared/src/PriorityCoalescingWorker.test.ts new file mode 100644 index 000000000000..428ee340d65b --- /dev/null +++ b/packages/shared/src/PriorityCoalescingWorker.test.ts @@ -0,0 +1,90 @@ +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 { makePriorityCoalescingWorker } from "./PriorityCoalescingWorker.ts"; + +describe("makePriorityCoalescingWorker", () => { + it.live("runs realtime work before queued background updates and coalesces each key", () => + Effect.scoped( + Effect.gen(function* () { + const processed: string[] = []; + const firstStarted = yield* Deferred.make(); + const releaseFirst = yield* Deferred.make(); + + const worker = yield* makePriorityCoalescingWorker({ + mergeBackground: (_current, next) => next, + process: (value) => + Effect.gen(function* () { + processed.push(value); + if (value === "background:first") { + yield* Deferred.succeed(firstStarted, undefined).pipe(Effect.orDie); + yield* Deferred.await(releaseFirst); + } + }), + }); + + yield* worker.enqueueBackground("task-1", "background:first"); + yield* Deferred.await(firstStarted); + + yield* worker.enqueueBackground("task-2", "background:stale"); + yield* worker.enqueueBackground("task-2", "background:latest"); + yield* worker.enqueueRealtime("assistant:reply"); + yield* Deferred.succeed(releaseFirst, undefined); + yield* worker.drain; + + expect(processed).toEqual(["background:first", "assistant:reply", "background:latest"]); + }), + ), + ); + + it.live("keeps processing realtime work after a processor failure", () => + Effect.scoped( + Effect.gen(function* () { + const processed = yield* Deferred.make(); + const worker = yield* makePriorityCoalescingWorker({ + mergeBackground: (_current, next) => next, + process: (value) => + value === "fail" + ? Effect.fail("expected failure") + : Deferred.succeed(processed, undefined).pipe(Effect.asVoid), + }); + + yield* worker.enqueueRealtime("fail"); + yield* worker.enqueueRealtime("succeed"); + yield* Deferred.await(processed); + yield* worker.drain; + }), + ), + ); + + it.live("processes the latest background update after a processor failure", () => + Effect.scoped( + Effect.gen(function* () { + const firstStarted = yield* Deferred.make(); + const releaseFirst = yield* Deferred.make(); + const processedLatest = yield* Deferred.make(); + const worker = yield* makePriorityCoalescingWorker({ + mergeBackground: (_current, next) => next, + process: (value) => + Effect.gen(function* () { + if (value === "first") { + yield* Deferred.succeed(firstStarted, undefined).pipe(Effect.orDie); + yield* Deferred.await(releaseFirst); + return yield* Effect.fail("expected failure"); + } + yield* Deferred.succeed(processedLatest, undefined).pipe(Effect.orDie); + }), + }); + + yield* worker.enqueueBackground("task-1", "first"); + yield* Deferred.await(firstStarted); + yield* worker.enqueueBackground("task-1", "latest"); + yield* Deferred.succeed(releaseFirst, undefined); + yield* Deferred.await(processedLatest); + yield* worker.drain; + }), + ), + ); +}); diff --git a/packages/shared/src/PriorityCoalescingWorker.ts b/packages/shared/src/PriorityCoalescingWorker.ts new file mode 100644 index 000000000000..a0510a1211c9 --- /dev/null +++ b/packages/shared/src/PriorityCoalescingWorker.ts @@ -0,0 +1,162 @@ +/** + * A single-lane worker that always drains realtime work before background + * work and keeps only the latest queued background value for each key. + * + * Realtime work cannot interrupt the item already being processed, but it is + * selected before any remaining background key. Background keys are requeued + * after one pass, so one hot key cannot starve the others. + * + * @module PriorityCoalescingWorker + */ +import * as Cause from "effect/Cause"; +import * as Effect from "effect/Effect"; +import * as HashMap from "effect/HashMap"; +import * as HashSet from "effect/HashSet"; +import * as Option from "effect/Option"; +import * as Scope from "effect/Scope"; +import * as TxQueue from "effect/TxQueue"; +import * as TxRef from "effect/TxRef"; + +export interface PriorityCoalescingWorker { + readonly enqueueRealtime: (value: V) => Effect.Effect; + readonly enqueueBackground: (key: K, value: V) => Effect.Effect; + readonly drain: Effect.Effect; +} + +interface BackgroundState { + readonly latestByKey: HashMap.HashMap; + readonly queuedKeys: HashSet.HashSet; + readonly activeKeys: HashSet.HashSet; +} + +type WorkItem = + | { readonly _tag: "Realtime"; readonly value: V } + | { readonly _tag: "Background"; readonly key: K; readonly value: V }; + +export const makePriorityCoalescingWorker = (options: { + readonly mergeBackground: (current: V, next: V) => V; + readonly process: (value: V) => Effect.Effect; +}): Effect.Effect, never, Scope.Scope | R> => + Effect.gen(function* () { + const realtimeQueue = yield* Effect.acquireRelease(TxQueue.unbounded(), TxQueue.shutdown); + const backgroundQueue = yield* Effect.acquireRelease(TxQueue.unbounded(), TxQueue.shutdown); + const backgroundState = yield* TxRef.make>({ + latestByKey: HashMap.empty(), + queuedKeys: HashSet.empty(), + activeKeys: HashSet.empty(), + }); + const outstanding = yield* TxRef.make(0); + + const takeNext = Effect.gen(function* () { + const realtime = yield* TxQueue.poll(realtimeQueue); + if (Option.isSome(realtime)) { + return { _tag: "Realtime", value: realtime.value } as const; + } + + const key = yield* TxQueue.take(backgroundQueue); + const state = yield* TxRef.get(backgroundState); + const value = HashMap.get(state.latestByKey, key); + if (Option.isNone(value)) { + return yield* Effect.txRetry; + } + + yield* TxRef.set(backgroundState, { + latestByKey: HashMap.remove(state.latestByKey, key), + queuedKeys: HashSet.remove(state.queuedKeys, key), + activeKeys: HashSet.add(state.activeKeys, key), + }); + + return { _tag: "Background", key, value: value.value } as const; + }).pipe(Effect.tx); + + const finishBackground = (key: K) => + Effect.gen(function* () { + const state = yield* TxRef.get(backgroundState); + const activeKeys = HashSet.remove(state.activeKeys, key); + + if (!HashMap.has(state.latestByKey, key)) { + yield* TxRef.set(backgroundState, { ...state, activeKeys }); + yield* TxRef.update(outstanding, (count) => count - 1); + return; + } + + yield* TxRef.set(backgroundState, { + ...state, + queuedKeys: HashSet.add(state.queuedKeys, key), + activeKeys, + }); + yield* TxQueue.offer(backgroundQueue, key); + }).pipe(Effect.tx); + + // One malformed event must not terminate the worker and strand later work. + // Scoped interruption still propagates so shutdown remains prompt. + const keepWorkerAlive = (effect: Effect.Effect) => + effect.pipe( + Effect.catchCause((cause) => + Cause.hasInterruptsOnly(cause) ? Effect.failCause(cause) : Effect.void, + ), + ); + + yield* takeNext.pipe( + Effect.flatMap((item: WorkItem) => + item._tag === "Realtime" + ? keepWorkerAlive( + options + .process(item.value) + .pipe( + Effect.ensuring(TxRef.update(outstanding, (count) => count - 1).pipe(Effect.tx)), + ), + ) + : keepWorkerAlive( + options.process(item.value).pipe(Effect.ensuring(finishBackground(item.key))), + ), + ), + Effect.forever, + Effect.forkScoped, + ); + + const enqueueRealtime: PriorityCoalescingWorker["enqueueRealtime"] = (value) => + TxQueue.offer(realtimeQueue, value).pipe( + Effect.tap(() => TxRef.update(outstanding, (count) => count + 1)), + Effect.tx, + Effect.asVoid, + ); + + const enqueueBackground: PriorityCoalescingWorker["enqueueBackground"] = (key, value) => + Effect.gen(function* () { + const state = yield* TxRef.get(backgroundState); + const existing = HashMap.get(state.latestByKey, key); + const latestByKey = HashMap.set( + state.latestByKey, + key, + Option.match(existing, { + onNone: () => value, + onSome: (current) => options.mergeBackground(current, value), + }), + ); + + if (HashSet.has(state.queuedKeys, key) || HashSet.has(state.activeKeys, key)) { + yield* TxRef.set(backgroundState, { ...state, latestByKey }); + return; + } + + yield* TxRef.set(backgroundState, { + ...state, + latestByKey, + queuedKeys: HashSet.add(state.queuedKeys, key), + }); + yield* TxRef.update(outstanding, (count) => count + 1); + yield* TxQueue.offer(backgroundQueue, key); + }).pipe(Effect.tx); + + const drain: PriorityCoalescingWorker["drain"] = TxRef.get(outstanding).pipe( + Effect.tap((count) => (count > 0 ? Effect.txRetry : Effect.void)), + Effect.tx, + ); + + return { + enqueueRealtime, + enqueueBackground, + drain, + } satisfies PriorityCoalescingWorker; + });