From 938e2af1448b7efbaa7a97760c972b129ca108da Mon Sep 17 00:00:00 2001 From: Julius Marminge Date: Mon, 5 Oct 2026 15:30:00 -0700 Subject: [PATCH 1/3] fix(shared): DrainableWorker keeps running after a failed item When the processor failed or died on one item, the error escaped the worker loop and its fiber exited. Later items were never processed, and anything waiting on `drain` hung forever because the outstanding count of the queued items was never decremented. AgentAwarenessRelay's publish worker and ThreadSettlementService's sweep worker pass `process` functions that can still die, so one bad item stopped them for the rest of the server's life. The worker now catches each item's cause, logs it with `Effect.logError`, and keeps going; the outstanding count is still decremented in `Effect.ensuring`. Closing the worker's scope still stops it: in Effect 4.0.1 an interrupt of the worker fiber skips `catchCause` handlers, so only the item's own failure, defect, or self-interruption is caught. Co-Authored-By: Claude Opus 5.5 (1M context) --- packages/shared/src/DrainableWorker.test.ts | 56 +++++++++++++++++++++ packages/shared/src/DrainableWorker.ts | 16 ++++-- 2 files changed, 68 insertions(+), 4 deletions(-) diff --git a/packages/shared/src/DrainableWorker.test.ts b/packages/shared/src/DrainableWorker.test.ts index 8e4c654e2e4c..952bc507396a 100644 --- a/packages/shared/src/DrainableWorker.test.ts +++ b/packages/shared/src/DrainableWorker.test.ts @@ -1,11 +1,67 @@ import { it } from "@effect/vitest"; import { describe, expect } from "vite-plus/test"; +import * as Cause from "effect/Cause"; import * as Deferred from "effect/Deferred"; import * as Effect from "effect/Effect"; +import * as Exit from "effect/Exit"; +import * as Fiber from "effect/Fiber"; +import * as Logger from "effect/Logger"; +import * as Scope from "effect/Scope"; import { makeDrainableWorker } from "./DrainableWorker.ts"; describe("makeDrainableWorker", () => { + it.effect("logs a failed or defective item and keeps processing later items", () => { + const loggedErrors: Array = []; + const logger = Logger.make(({ logLevel, cause }) => { + if (logLevel === "Error") loggedErrors.push(Cause.squash(cause)); + }); + + return Effect.scoped( + Effect.gen(function* () { + const processed: string[] = []; + const workerFiber = yield* Deferred.make>(); + const worker = yield* makeDrainableWorker((item: string) => + Effect.gen(function* () { + yield* Deferred.succeed(workerFiber, yield* Effect.fiber); + if (item === "fail") return yield* Effect.fail("typed failure"); + if (item === "die") return yield* Effect.die("defect"); + processed.push(item); + }), + ); + + yield* worker.enqueue("fail"); + yield* worker.enqueue("die"); + yield* worker.enqueue("ok"); + + // A stopped worker never drains; its exit settles the race instead of a hang. + yield* Effect.raceFirst(worker.drain, Fiber.await(yield* Deferred.await(workerFiber))); + + expect(processed).toEqual(["ok"]); + expect(loggedErrors).toEqual(["typed failure", "defect"]); + }), + ).pipe(Effect.provide(Logger.layer([logger], { mergeWithExisting: false }))); + }); + + it.effect("stops when its scope closes during an item", () => + Effect.gen(function* () { + const scope = yield* Scope.make(); + const workerFiber = yield* Deferred.make>(); + const worker = yield* makeDrainableWorker(() => + Effect.gen(function* () { + yield* Deferred.succeed(workerFiber, yield* Effect.fiber); + return yield* Effect.never; + }), + ).pipe(Scope.provide(scope)); + + yield* worker.enqueue("item"); + const fiber = yield* Deferred.await(workerFiber); + yield* Scope.close(scope, Exit.void); + + expect(Exit.hasInterrupts(yield* Fiber.await(fiber))).toBe(true); + }), + ); + it.live("waits for work enqueued during active processing before draining", () => Effect.scoped( Effect.gen(function* () { diff --git a/packages/shared/src/DrainableWorker.ts b/packages/shared/src/DrainableWorker.ts index de40ec5e36b8..dfd21bb5660e 100644 --- a/packages/shared/src/DrainableWorker.ts +++ b/packages/shared/src/DrainableWorker.ts @@ -34,6 +34,9 @@ export interface DrainableWorker { * The worker is forked into the current scope and will be interrupted when * the scope closes. A finalizer shuts down the queue. * + * An item that fails or dies is logged and skipped; the worker keeps + * processing later items and `drain` still resolves. + * * @param process - The effect to run for each queued item. * @returns A `DrainableWorker` with `queue` and `drain`. */ @@ -45,10 +48,15 @@ export const makeDrainableWorker = ( const outstanding = yield* TxRef.make(0); yield* TxQueue.take(queue).pipe( - Effect.tap((a) => - Effect.ensuring( - process(a), - TxRef.update(outstanding, (n) => n - 1), + Effect.flatMap((a) => + // `suspend` turns a `process` that throws while building its effect + // into this item's defect instead of the loop's. + Effect.suspend(() => process(a)).pipe( + // Interrupting the worker fiber (its scope closing) skips this + // handler, so only the item's own failure, defect, or interruption + // lands here and the loop continues. + Effect.catchCause((cause) => Effect.logError("DrainableWorker item failed", cause)), + Effect.ensuring(TxRef.update(outstanding, (n) => n - 1)), ), ), Effect.forever, From 67948566391456c90e4dbd72215db32cda5a4ef4 Mon Sep 17 00:00:00 2001 From: Julius Marminge Date: Mon, 5 Oct 2026 16:41:57 -0700 Subject: [PATCH 2/3] fix(shared): DrainableWorker drains after shutdown and ignores self-cancelled items An item that only interrupted itself was logged as a failure, although every caller treats an interrupt-only cause as a cancellation. It is now skipped quietly. Closing the scope dropped queued items without uncounting them, and enqueue counted items the shut-down queue refused, so a later drain waited forever. The shutdown finalizer now clears the count and enqueue only counts accepted items. drain also returns void instead of leaking the count. Co-Authored-By: Claude Opus 5.5 (1M context) --- packages/shared/src/DrainableWorker.test.ts | 43 +++++++++++++++------ packages/shared/src/DrainableWorker.ts | 22 ++++++++--- 2 files changed, 48 insertions(+), 17 deletions(-) diff --git a/packages/shared/src/DrainableWorker.test.ts b/packages/shared/src/DrainableWorker.test.ts index 952bc507396a..bda85810fa6f 100644 --- a/packages/shared/src/DrainableWorker.test.ts +++ b/packages/shared/src/DrainableWorker.test.ts @@ -10,12 +10,17 @@ import * as Scope from "effect/Scope"; import { makeDrainableWorker } from "./DrainableWorker.ts"; +const captureErrorLogs = () => { + const loggedErrors: Array = []; + const logger = Logger.make(({ logLevel, cause }) => { + if (logLevel === "Error") loggedErrors.push(Cause.squash(cause)); + }); + return { loggedErrors, loggerLayer: Logger.layer([logger], { mergeWithExisting: false }) }; +}; + describe("makeDrainableWorker", () => { it.effect("logs a failed or defective item and keeps processing later items", () => { - const loggedErrors: Array = []; - const logger = Logger.make(({ logLevel, cause }) => { - if (logLevel === "Error") loggedErrors.push(Cause.squash(cause)); - }); + const { loggedErrors, loggerLayer } = captureErrorLogs(); return Effect.scoped( Effect.gen(function* () { @@ -26,41 +31,57 @@ describe("makeDrainableWorker", () => { yield* Deferred.succeed(workerFiber, yield* Effect.fiber); if (item === "fail") return yield* Effect.fail("typed failure"); if (item === "die") return yield* Effect.die("defect"); + if (item === "interrupt") return yield* Effect.interrupt; processed.push(item); }), ); yield* worker.enqueue("fail"); yield* worker.enqueue("die"); + yield* worker.enqueue("interrupt"); yield* worker.enqueue("ok"); // A stopped worker never drains; its exit settles the race instead of a hang. yield* Effect.raceFirst(worker.drain, Fiber.await(yield* Deferred.await(workerFiber))); expect(processed).toEqual(["ok"]); + // The item that interrupted itself was cancelled, not failed. expect(loggedErrors).toEqual(["typed failure", "defect"]); }), - ).pipe(Effect.provide(Logger.layer([logger], { mergeWithExisting: false }))); + ).pipe(Effect.provide(loggerLayer)); }); - it.effect("stops when its scope closes during an item", () => - Effect.gen(function* () { + it.effect("stops quietly when its scope closes during an item", () => { + const { loggedErrors, loggerLayer } = captureErrorLogs(); + + return Effect.gen(function* () { const scope = yield* Scope.make(); + const processed: string[] = []; const workerFiber = yield* Deferred.make>(); - const worker = yield* makeDrainableWorker(() => + const worker = yield* makeDrainableWorker((item: string) => Effect.gen(function* () { + processed.push(item); yield* Deferred.succeed(workerFiber, yield* Effect.fiber); return yield* Effect.never; }), ).pipe(Scope.provide(scope)); - yield* worker.enqueue("item"); + yield* worker.enqueue("block"); + yield* worker.enqueue("later"); const fiber = yield* Deferred.await(workerFiber); yield* Scope.close(scope, Exit.void); expect(Exit.hasInterrupts(yield* Fiber.await(fiber))).toBe(true); - }), - ); + expect(processed).toEqual(["block"]); + expect(loggedErrors).toEqual([]); + + // Queued and late items are dropped with the queue, so drain resolves at once. + yield* worker.enqueue("after close"); + const drained = yield* Effect.forkChild(worker.drain, { startImmediately: true }); + expect(drained.pollUnsafe()).toEqual(Exit.void); + expect(processed).toEqual(["block"]); + }).pipe(Effect.provide(loggerLayer)); + }); it.live("waits for work enqueued during active processing before draining", () => Effect.scoped( diff --git a/packages/shared/src/DrainableWorker.ts b/packages/shared/src/DrainableWorker.ts index dfd21bb5660e..c41ae8aefc26 100644 --- a/packages/shared/src/DrainableWorker.ts +++ b/packages/shared/src/DrainableWorker.ts @@ -8,6 +8,7 @@ * * @module DrainableWorker */ +import * as Cause from "effect/Cause"; import * as Scope from "effect/Scope"; import * as Effect from "effect/Effect"; import * as TxQueue from "effect/TxQueue"; @@ -32,7 +33,8 @@ export interface DrainableWorker { * Create a drainable worker that processes items from an unbounded queue. * * The worker is forked into the current scope and will be interrupted when - * the scope closes. A finalizer shuts down the queue. + * the scope closes. A finalizer shuts down the queue and drops queued items, + * so `drain` resolves after the scope closes instead of waiting on them. * * An item that fails or dies is logged and skipped; the worker keeps * processing later items and `drain` still resolves. @@ -44,8 +46,10 @@ export const makeDrainableWorker = ( process: (item: A) => Effect.Effect, ): Effect.Effect, never, Scope.Scope | R> => Effect.gen(function* () { - const queue = yield* Effect.acquireRelease(TxQueue.unbounded(), TxQueue.shutdown); const outstanding = yield* TxRef.make(0); + const queue = yield* Effect.acquireRelease(TxQueue.unbounded(), (queue) => + TxQueue.shutdown(queue).pipe(Effect.andThen(TxRef.set(outstanding, 0)), Effect.tx), + ); yield* TxQueue.take(queue).pipe( Effect.flatMap((a) => @@ -54,8 +58,13 @@ export const makeDrainableWorker = ( Effect.suspend(() => process(a)).pipe( // Interrupting the worker fiber (its scope closing) skips this // handler, so only the item's own failure, defect, or interruption - // lands here and the loop continues. - Effect.catchCause((cause) => Effect.logError("DrainableWorker item failed", cause)), + // lands here and the loop continues. Callers treat an item that only + // interrupted itself as cancelled, not failed. + Effect.catchCause((cause) => + Cause.hasInterruptsOnly(cause) + ? Effect.void + : Effect.logError("DrainableWorker item failed", cause), + ), Effect.ensuring(TxRef.update(outstanding, (n) => n - 1)), ), ), @@ -64,13 +73,14 @@ export const makeDrainableWorker = ( ); const drain: DrainableWorker["drain"] = TxRef.get(outstanding).pipe( - Effect.tap((n) => (n > 0 ? Effect.txRetry : Effect.void)), + Effect.flatMap((n) => (n > 0 ? Effect.txRetry : Effect.void)), Effect.tx, ); const enqueue = (element: A): Effect.Effect => TxQueue.offer(queue, element).pipe( - Effect.tap(() => TxRef.update(outstanding, (n) => n + 1)), + // A shut-down queue refuses the item, so it is never processed. + Effect.tap((offered) => (offered ? TxRef.update(outstanding, (n) => n + 1) : Effect.void)), Effect.tx, ); From 9e304cda88489eac5971503755db5e091125abbd Mon Sep 17 00:00:00 2001 From: Julius Marminge Date: Mon, 5 Oct 2026 17:59:02 -0700 Subject: [PATCH 3/3] fix(shared): DrainableWorker shutdown uncounts only the dropped items Zeroing the count on shutdown assumed the worker fiber had already exited. Under a parallel scope the in-flight item's own decrement could then run afterwards and leave the count at -1. Only the dropped items are uncounted now; the running item still uncounts itself. Co-Authored-By: Claude Opus 5.5 (1M context) --- packages/shared/src/DrainableWorker.ts | 18 ++++++++++++------ 1 file changed, 12 insertions(+), 6 deletions(-) diff --git a/packages/shared/src/DrainableWorker.ts b/packages/shared/src/DrainableWorker.ts index c41ae8aefc26..5ec2a3f9db44 100644 --- a/packages/shared/src/DrainableWorker.ts +++ b/packages/shared/src/DrainableWorker.ts @@ -40,7 +40,7 @@ export interface DrainableWorker { * processing later items and `drain` still resolves. * * @param process - The effect to run for each queued item. - * @returns A `DrainableWorker` with `queue` and `drain`. + * @returns A `DrainableWorker` with `enqueue` and `drain`. */ export const makeDrainableWorker = ( process: (item: A) => Effect.Effect, @@ -48,7 +48,13 @@ export const makeDrainableWorker = ( Effect.gen(function* () { const outstanding = yield* TxRef.make(0); const queue = yield* Effect.acquireRelease(TxQueue.unbounded(), (queue) => - TxQueue.shutdown(queue).pipe(Effect.andThen(TxRef.set(outstanding, 0)), Effect.tx), + // Uncount only the dropped items: an item still running uncounts itself, + // even when a parallel scope closes it after this finalizer. + TxQueue.clear(queue).pipe( + Effect.flatMap((dropped) => TxRef.update(outstanding, (n) => n - dropped.length)), + Effect.andThen(TxQueue.shutdown(queue)), + Effect.tx, + ), ); yield* TxQueue.take(queue).pipe( @@ -56,10 +62,10 @@ export const makeDrainableWorker = ( // `suspend` turns a `process` that throws while building its effect // into this item's defect instead of the loop's. Effect.suspend(() => process(a)).pipe( - // Interrupting the worker fiber (its scope closing) skips this - // handler, so only the item's own failure, defect, or interruption - // lands here and the loop continues. Callers treat an item that only - // interrupted itself as cancelled, not failed. + // Only the item's own failure, defect, or interruption lands here and + // the loop continues; interrupting the worker fiber still stops it. + // Callers treat an item that only interrupted itself as cancelled, + // not failed. Effect.catchCause((cause) => Cause.hasInterruptsOnly(cause) ? Effect.void