Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
77 changes: 77 additions & 0 deletions packages/shared/src/DrainableWorker.test.ts
Original file line number Diff line number Diff line change
@@ -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<unknown> = [];
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<Fiber.Fiber<unknown, unknown>>();
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<Fiber.Fiber<unknown, unknown>>();
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* () {
Expand Down
42 changes: 33 additions & 9 deletions packages/shared/src/DrainableWorker.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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";
Expand All @@ -32,37 +33,60 @@ export interface DrainableWorker<A> {
* 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 = <A, E, R>(
process: (item: A) => Effect.Effect<void, E, R>,
): Effect.Effect<DrainableWorker<A>, never, Scope.Scope | R> =>
Effect.gen(function* () {
const queue = yield* Effect.acquireRelease(TxQueue.unbounded<A>(), TxQueue.shutdown);
const outstanding = yield* TxRef.make(0);
const queue = yield* Effect.acquireRelease(TxQueue.unbounded<A>(), (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,
Effect.forkScoped,
);

const drain: DrainableWorker<A>["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<boolean, never, never> =>
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,
);

Expand Down
Loading