From f8ff1dbe749435997450fcb2d479cd93d37b0339 Mon Sep 17 00:00:00 2001 From: T3 Code PR Stack <41898282+github-actions[bot]@users.noreply.github.com> Date: Sun, 2 Aug 2026 11:18:37 +0200 Subject: [PATCH] feat(server): isolate provider commands by thread --- .../Layers/ProviderCommandReactor.test.ts | 4 +- .../Layers/ProviderCommandReactor.ts | 9 +- docs/internals/overview.md | 3 + docs/internals/providers.md | 8 +- packages/shared/package.json | 4 + .../shared/src/KeyedDrainableWorker.test.ts | 71 ++++++++ packages/shared/src/KeyedDrainableWorker.ts | 159 ++++++++++++++++++ 7 files changed, 253 insertions(+), 5 deletions(-) create mode 100644 packages/shared/src/KeyedDrainableWorker.test.ts create mode 100644 packages/shared/src/KeyedDrainableWorker.ts diff --git a/apps/server/src/orchestration/Layers/ProviderCommandReactor.test.ts b/apps/server/src/orchestration/Layers/ProviderCommandReactor.test.ts index b36f6dbca243..0255c80eb948 100644 --- a/apps/server/src/orchestration/Layers/ProviderCommandReactor.test.ts +++ b/apps/server/src/orchestration/Layers/ProviderCommandReactor.test.ts @@ -2491,7 +2491,7 @@ describe("ProviderCommandReactor", () => { }); }); - it("bounds a hung provider interrupt so later thread starts still run", async () => { + it("isolates a hung provider interrupt so another thread starts immediately", async () => { const harness = await createHarness({ interruptTurnEffect: () => Effect.never, }); @@ -2557,7 +2557,7 @@ describe("ProviderCommandReactor", () => { const interrupted = readModel.threads.find((entry) => entry.id === ThreadId.make("thread-1")); return interrupted?.session?.status === "ready"; }); - await waitFor(() => harness.sendTurn.mock.calls.length === 1, 8_000); + await waitFor(() => harness.sendTurn.mock.calls.length === 1, 1_000); expect(harness.sendTurn.mock.calls[0]?.[0]).toMatchObject({ threadId: secondThreadId, }); diff --git a/apps/server/src/orchestration/Layers/ProviderCommandReactor.ts b/apps/server/src/orchestration/Layers/ProviderCommandReactor.ts index e2626345d121..df8df4538463 100644 --- a/apps/server/src/orchestration/Layers/ProviderCommandReactor.ts +++ b/apps/server/src/orchestration/Layers/ProviderCommandReactor.ts @@ -28,6 +28,7 @@ import * as Option from "effect/Option"; import * as Schema from "effect/Schema"; import * as Stream from "effect/Stream"; import { makeDrainableWorker } from "@t3tools/shared/DrainableWorker"; +import { makeKeyedDrainableWorker } from "@t3tools/shared/KeyedDrainableWorker"; import { resolveThreadWorkspaceCwd } from "../../checkpointing/Utils.ts"; import { @@ -1767,7 +1768,10 @@ const make = Effect.gen(function* () { }), ); - const worker = yield* makeDrainableWorker(processDomainEventSafely); + const worker = yield* makeKeyedDrainableWorker({ + key: (event: ProviderIntentEvent) => event.payload.threadId, + process: processDomainEventSafely, + }); const reconcileStartup = Effect.fn("reconcileStartup")(function* () { const bindings = yield* providerSessionDirectory.listBindings().pipe( @@ -1900,6 +1904,9 @@ const make = Effect.gen(function* () { const start: ProviderCommandReactorShape["start"] = Effect.fn("start")(function* () { const processEvent = Effect.fn("processEvent")(function* (event: OrchestrationEvent) { + if (event.type === "thread.deleted") { + return yield* worker.cancelKey(event.payload.threadId); + } if ( (event.type === "thread.meta-updated" && event.payload.regenerateTitle === true) || event.type === "thread.runtime-mode-set" || diff --git a/docs/internals/overview.md b/docs/internals/overview.md index b9454f7b58d0..2142e6f92375 100644 --- a/docs/internals/overview.md +++ b/docs/internals/overview.md @@ -100,6 +100,9 @@ Follow-up work runs asynchronously in queue-backed workers built on [`DrainableW count reaches zero, so a test can await "queue empty and current item finished" instead of sleeping. Each of the three services exposes `drain` for exactly this. +`ProviderCommandReactor` uses the keyed variant: FIFO within a thread and concurrent across threads. +Its lanes are lazy and ephemeral, so only threads with active or queued commands consume a lane. + Runtime receipts are a test-only mechanism. `RuntimeReceiptBusLive` in [`RuntimeReceiptBus.ts`][receipts] publishes nothing; only the test layer is PubSub-backed. Do not build production behavior on receipts. diff --git a/docs/internals/providers.md b/docs/internals/providers.md index a309d70f03de..05eeaebc3301 100644 --- a/docs/internals/providers.md +++ b/docs/internals/providers.md @@ -56,8 +56,7 @@ Provider output comes back as internal commands such as `thread.message.assistan ## Server-side workers Provider work flows through three queue-backed workers. All three are built with -`makeDrainableWorker` from [`DrainableWorker.ts`][worker] and expose `drain` for deterministic test -synchronization. +drainable worker primitives and expose `drain` for deterministic test synchronization. 1. [`ProviderRuntimeIngestion`][ingest] consumes provider runtime streams and emits orchestration commands. @@ -66,6 +65,11 @@ synchronization. 3. [`CheckpointReactor`][checkpoint] captures workspace checkpoints on turn start and completion, and performs reverts. +Provider commands use lazy per-thread FIFO lanes. Commands remain ordered within one thread, while +different threads run concurrently so a stalled provider cannot block unrelated conversations. A +lane is created only when work arrives and is removed as soon as it becomes idle; historical threads +do not allocate lanes. Deleting a thread interrupts and removes any active lane. + ### Buffered assistant delivery A thread in `buffered` assistant delivery mode accumulates assistant text instead of streaming each diff --git a/packages/shared/package.json b/packages/shared/package.json index 4214d57b9538..384a5df38fce 100644 --- a/packages/shared/package.json +++ b/packages/shared/package.json @@ -67,6 +67,10 @@ "types": "./src/KeyedCoalescingWorker.ts", "import": "./src/KeyedCoalescingWorker.ts" }, + "./KeyedDrainableWorker": { + "types": "./src/KeyedDrainableWorker.ts", + "import": "./src/KeyedDrainableWorker.ts" + }, "./schemaJson": { "types": "./src/schemaJson.ts", "import": "./src/schemaJson.ts" diff --git a/packages/shared/src/KeyedDrainableWorker.test.ts b/packages/shared/src/KeyedDrainableWorker.test.ts new file mode 100644 index 000000000000..cd944f8d851b --- /dev/null +++ b/packages/shared/src/KeyedDrainableWorker.test.ts @@ -0,0 +1,71 @@ +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 { makeKeyedDrainableWorker } from "./KeyedDrainableWorker.ts"; + +describe("makeKeyedDrainableWorker", () => { + it.live("keeps FIFO order per key, runs keys concurrently, and reclaims idle lanes", () => + Effect.scoped( + Effect.gen(function* () { + const blockedStarted = yield* Deferred.make(); + const releaseBlocked = yield* Deferred.make(); + const processed: string[] = []; + const worker = yield* makeKeyedDrainableWorker({ + key: (item) => item.split(":")[0] ?? item, + process: (item) => + Effect.gen(function* () { + processed.push(item); + if (item === "blocked:first") { + yield* Deferred.succeed(blockedStarted, undefined).pipe(Effect.orDie); + yield* Deferred.await(releaseBlocked); + } + }), + }); + + expect(yield* worker.activeKeyCount).toBe(0); + yield* worker.enqueue("blocked:first"); + yield* worker.enqueue("blocked:second"); + yield* Deferred.await(blockedStarted).pipe(Effect.timeout("1 second")); + yield* worker.enqueue("free:first"); + yield* worker.drainKey("free"); + + expect(processed).toEqual(["blocked:first", "free:first"]); + expect(yield* worker.activeKeyCount).toBe(1); + + yield* Deferred.succeed(releaseBlocked, undefined).pipe(Effect.orDie); + yield* worker.drain; + expect(processed).toEqual(["blocked:first", "free:first", "blocked:second"]); + expect(yield* worker.activeKeyCount).toBe(0); + }), + ), + ); + + it.live("interrupts active work and drops queued work when a key is cancelled", () => + Effect.scoped( + Effect.gen(function* () { + const started = yield* Deferred.make(); + const processed: string[] = []; + const worker = yield* makeKeyedDrainableWorker({ + key: () => "deleted-thread", + process: (item) => + Effect.gen(function* () { + processed.push(item); + yield* Deferred.succeed(started, undefined).pipe(Effect.orDie); + return yield* Effect.never; + }), + }); + + yield* worker.enqueue("active"); + yield* worker.enqueue("queued"); + yield* Deferred.await(started).pipe(Effect.timeout("1 second")); + yield* worker.cancelKey("deleted-thread"); + yield* worker.drain; + + expect(processed).toEqual(["active"]); + expect(yield* worker.activeKeyCount).toBe(0); + }), + ), + ); +}); diff --git a/packages/shared/src/KeyedDrainableWorker.ts b/packages/shared/src/KeyedDrainableWorker.ts new file mode 100644 index 000000000000..ce89a82712a7 --- /dev/null +++ b/packages/shared/src/KeyedDrainableWorker.ts @@ -0,0 +1,159 @@ +/** + * Lazy FIFO work lanes keyed by an identifier. + * + * A lane exists only while its key has active or queued work. Same-key work is + * serial, different keys run concurrently, and idle lanes are reclaimed. + */ +import * as Scope from "effect/Scope"; +import * as Effect from "effect/Effect"; +import * as Fiber from "effect/Fiber"; +import * as Semaphore from "effect/Semaphore"; +import * as TxRef from "effect/TxRef"; + +export interface KeyedDrainableWorker { + readonly enqueue: (item: A) => Effect.Effect; + readonly cancelKey: (key: K) => Effect.Effect; + readonly drain: Effect.Effect; + readonly drainKey: (key: K) => Effect.Effect; + readonly activeKeyCount: Effect.Effect; +} + +interface Lane { + readonly queue: Array; + fiber: Fiber.Fiber | undefined; +} + +interface Outstanding { + readonly total: number; + readonly byKey: Map; +} + +export const makeKeyedDrainableWorker = (options: { + readonly key: (item: A) => K; + readonly process: (item: A) => Effect.Effect; +}): Effect.Effect, never, Scope.Scope | R> => + Effect.gen(function* () { + const context = yield* Effect.context(); + const registryLock = yield* Semaphore.make(1); + const lanes = new Map>(); + const outstandingRef = yield* TxRef.make>({ + total: 0, + byKey: new Map(), + }); + + const adjustOutstanding = (key: K, delta: number) => + TxRef.update(outstandingRef, (state) => { + const byKey = new Map(state.byKey); + const next = (byKey.get(key) ?? 0) + delta; + if (next <= 0) byKey.delete(key); + else byKey.set(key, next); + return { total: state.total + delta, byKey }; + }).pipe(Effect.tx); + + const processLane = (key: K, lane: Lane): Effect.Effect => + Effect.suspend(() => + registryLock + .withPermit( + Effect.sync(() => { + const item = lane.queue.shift(); + if (item === undefined && lanes.get(key) === lane) { + lanes.delete(key); + } + return item; + }), + ) + .pipe( + Effect.flatMap((item) => + item === undefined + ? Effect.void + : options + .process(item) + .pipe( + Effect.provide(context), + Effect.ensuring(adjustOutstanding(key, -1)), + Effect.andThen(processLane(key, lane)), + ), + ), + ), + ); + + const enqueue: KeyedDrainableWorker["enqueue"] = (item) => { + const key = options.key(item); + return registryLock.withPermit( + Effect.gen(function* () { + let lane = lanes.get(key); + if (lane === undefined) { + lane = { queue: [], fiber: undefined }; + lanes.set(key, lane); + } + lane.queue.push(item); + yield* adjustOutstanding(key, 1); + if (lane.fiber === undefined) { + lane.fiber = yield* processLane(key, lane).pipe(Effect.forkDetach); + } + }), + ); + }; + + const cancelKey: KeyedDrainableWorker["cancelKey"] = (key) => + Effect.gen(function* () { + let fiber: Fiber.Fiber | undefined; + let droppedQueuedItems = 0; + yield* registryLock.withPermit( + Effect.sync(() => { + const lane = lanes.get(key); + if (lane === undefined) return; + lanes.delete(key); + droppedQueuedItems = lane.queue.length; + lane.queue.length = 0; + fiber = lane.fiber; + }), + ); + if (droppedQueuedItems > 0) { + yield* adjustOutstanding(key, -droppedQueuedItems); + } + if (fiber !== undefined) { + yield* Fiber.interrupt(fiber); + } + }); + + yield* Effect.addFinalizer(() => + registryLock + .withPermit( + Effect.sync(() => { + const fibers = Array.from(lanes.values()).flatMap((lane) => + lane.fiber === undefined ? [] : [lane.fiber], + ); + lanes.clear(); + return fibers; + }), + ) + .pipe( + Effect.flatMap((fibers) => Fiber.interruptAll(fibers)), + Effect.asVoid, + ), + ); + + const drain = TxRef.get(outstandingRef).pipe( + Effect.tap((state) => (state.total > 0 ? Effect.txRetry : Effect.void)), + Effect.asVoid, + Effect.tx, + ); + + const drainKey: KeyedDrainableWorker["drainKey"] = (key) => + TxRef.get(outstandingRef).pipe( + Effect.tap((state) => (state.byKey.has(key) ? Effect.txRetry : Effect.void)), + Effect.asVoid, + Effect.tx, + ); + + const activeKeyCount = registryLock.withPermit(Effect.sync(() => lanes.size)); + + return { + enqueue, + cancelKey, + drain, + drainKey, + activeKeyCount, + } satisfies KeyedDrainableWorker; + });