diff --git a/packages/shared/src/DrainableWorker.test.ts b/packages/shared/src/DrainableWorker.test.ts index 8e4c654e2e4c..bda85810fa6f 100644 --- a/packages/shared/src/DrainableWorker.test.ts +++ b/packages/shared/src/DrainableWorker.test.ts @@ -1,11 +1,88 @@ 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"; +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, loggerLayer } = captureErrorLogs(); + + 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"); + 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(loggerLayer)); + }); + + 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((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("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( Effect.gen(function* () { diff --git a/packages/shared/src/DrainableWorker.ts b/packages/shared/src/DrainableWorker.ts index de40ec5e36b8..5ec2a3f9db44 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,23 +33,45 @@ 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. * * @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, ): 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) => + // 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( - 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( + // 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 + : Effect.logError("DrainableWorker item failed", cause), + ), + Effect.ensuring(TxRef.update(outstanding, (n) => n - 1)), ), ), Effect.forever, @@ -56,13 +79,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, );