diff --git a/task-buffer/README.md b/task-buffer/README.md index 271543b4..b867dc80 100644 --- a/task-buffer/README.md +++ b/task-buffer/README.md @@ -9,25 +9,57 @@ When this limit is reached, the `TaskBuffer` automatically queues additional spawn requests and processes them in order as capacity becomes available. This prevents resource overload while ensuring all tasks are eventually executed. +Submitting work, waiting for it to start, and waiting for it to finish are three +separate steps: + ```ts -import { run, sleep } from "effection"; +import { run, sleep, type Task } from "effection"; import { useTaskBuffer } from "@effectionx/task-buffer"; await run(function* () { - // Create a task buffer with a maximum of 2 concurrent tasks + // a buffer that keeps at most 2 tasks active at a time const buffer = yield* useTaskBuffer(2); - // These tasks will execute immediately since they're within the limit - yield* buffer.spawn(() => sleep(10)); + // 1. submit work. `spawn()` returns as soon as the request is queued; it does + // not wait for the task to start. It hands back an operation that resolves + // once the task has been admitted into the buffer. + const admission = yield* buffer.spawn(() => sleep(10)); yield* buffer.spawn(() => sleep(10)); - // This task will be queued until one of the running tasks completes + // the buffer is now full, so this request waits for capacity before it starts yield* buffer.spawn(() => sleep(10)); - // Wait for this specific task to complete - yield* yield* buffer.spawn(() => sleep(10)); + // 2. wait for admission. This resolves once there is room in the buffer and + // the task has been spawned, and returns that `Task`. + const task: Task = yield* admission; + + // 3. wait for that one task to run to completion and produce its result. + yield* task; - // Wait for all spawned tasks to complete + // 4. wait for every queued and active task to complete. yield* buffer; }); ``` + +## Failure + +A failing task is not isolated from the rest of the buffer. Its error propagates +out of the buffer and into the scope that created it, halting the buffer's other +active tasks along with everything else in that scope. The buffer is not +fail-fast in the sense of reporting the first failure to whoever is waiting on +it: `yield* buffer` waits only for queued and active tasks to drain, and is +itself halted by the unwinding scope, so it neither returns nor throws. + +`yield* task` does rethrow that task's error, but catching it there does not +contain the failure — the scope is torn down regardless. To keep one failure +from tearing down the caller, handle it inside the operation passed to +`spawn()`. + +## Cancellation + +When the scope holding the buffer exits, its active tasks are halted and any +requests still queued are never spawned. + +Abandoning the wait for admission withdraws the request. If the task holding +`yield* admission` is halted before the buffer has room for it, the operation is +never spawned. diff --git a/task-buffer/package.json b/task-buffer/package.json index cf2f1fbf..09d763dd 100644 --- a/task-buffer/package.json +++ b/task-buffer/package.json @@ -1,7 +1,7 @@ { "name": "@effectionx/task-buffer", "description": "Limit concurrent task execution with automatic queuing", - "version": "1.3.3", + "version": "1.3.4", "keywords": ["concurrency"], "type": "module", "main": "./dist/mod.js", diff --git a/task-buffer/task-buffer.test.ts b/task-buffer/task-buffer.test.ts index f019feb7..83a75a6b 100644 --- a/task-buffer/task-buffer.test.ts +++ b/task-buffer/task-buffer.test.ts @@ -1,9 +1,26 @@ import FakeTimers from "@sinonjs/fake-timers"; -import { sleep, spawn, type Task, until } from "effection"; +import { + createScope, + ensure, + run, + sleep, + spawn, + type Operation, + suspend, + type Task, + until, + withResolvers, +} from "effection"; import { describe, it } from "@effectionx/vitest"; import { expect } from "expect"; import { useTaskBuffer } from "./task-buffer.ts"; +// Park until Effection's scheduler has nothing left to run, so an assertion +// about work that did *not* happen is not just reading a queue too early. +function* settled(): Operation { + yield* until(Promise.resolve()); +} + describe("TaskBuffer", () => { it("queues up tasks when the buffer fills up", function* () { const clock = FakeTimers.install(); @@ -56,4 +73,217 @@ describe("TaskBuffer", () => { clock.uninstall(); } }); + it("admits a task, then resolves that task with its own result", function* () { + const buffer = yield* useTaskBuffer(2); + + const admission = yield* buffer.spawn(function* () { + return "result"; + }); + const task: Task = yield* admission; + + expect(yield* task).toEqual("result"); + }); + + it("halts active tasks and drops queued requests when the scope exits", function* () { + const [scope, destroy] = createScope(); + const started = withResolvers(); + const block = withResolvers(); + let activeCompleted = false; + let activeHalted = false; + let queuedStarted = false; + + yield* scope.spawn(function* () { + const buffer = yield* useTaskBuffer(1); + + yield* buffer.spawn(function* () { + started.resolve(); + try { + yield* block.operation; + activeCompleted = true; + } finally { + if (!activeCompleted) { + activeHalted = true; + } + } + }); + yield* buffer.spawn(function* () { + queuedStarted = true; + }); + + yield* block.operation; + }); + + yield* started.operation; + + yield* destroy(); + + expect(activeHalted).toEqual(true); + expect(activeCompleted).toEqual(false); + expect(queuedStarted).toEqual(false); + }); + + it("never admits a queued request once the scope begins to exit", function* () { + const [scope, destroy] = createScope(); + const started = [withResolvers(), withResolvers()]; + const teardownStarted = [withResolvers(), withResolvers()]; + const releaseTeardown = [withResolvers(), withResolvers()]; + let queuedStarted = 0; + + yield* scope.spawn(function* () { + const buffer = yield* useTaskBuffer(2); + + for (const index of [0, 1]) { + yield* buffer.spawn(function* () { + yield* ensure(function* () { + teardownStarted[index].resolve(); + yield* releaseTeardown[index].operation; + }); + started[index].resolve(); + yield* suspend(); + }); + } + + // submitted without awaiting admission, so nothing withdraws them + for (let i = 0; i < 3; i++) { + yield* buffer.spawn(function* () { + queuedStarted++; + }); + } + + yield* suspend(); + }); + + yield* started[0].operation; + yield* started[1].operation; + + const destruction = yield* spawn(destroy); + + // destruction halts the most recently admitted task first + yield* teardownStarted[1].operation; + yield* settled(); + expect(queuedStarted).toEqual(0); + + // releasing the first teardown frees a slot while the second one is still + // unwinding, which is the window the dispatch loop used to admit into + releaseTeardown[1].resolve(); + yield* teardownStarted[0].operation; + yield* settled(); + expect(queuedStarted).toEqual(0); + + releaseTeardown[0].resolve(); + yield* destruction; + + expect(queuedStarted).toEqual(0); + }); + + it("admits a request submitted after the dispatch loop has gone idle", function* () { + const buffer = yield* useTaskBuffer(5); + const started = withResolvers(); + const secondRan = withResolvers(); + + yield* buffer.spawn(function* () { + started.resolve(); + yield* suspend(); + }); + + // let the buffer go idle with capacity to spare before submitting again + yield* started.operation; + + yield* buffer.spawn(function* () { + secondRan.resolve(); + }); + + yield* secondRan.operation; + }); + + it("withdraws a queued request when its admission is abandoned", function* () { + const buffer = yield* useTaskBuffer(1); + const blocker = withResolvers(); + const queued = withResolvers(); + let queuedStarted = false; + + yield* buffer.spawn(function* () { + yield* blocker.operation; + }); + + const waiter = yield* spawn(function* () { + const admission = yield* buffer.spawn(function* () { + queuedStarted = true; + }); + queued.resolve(); + yield* admission; + }); + + yield* queued.operation; + yield* waiter.halt(); + + blocker.resolve(); + yield* buffer; + + expect(queuedStarted).toEqual(false); + }); + + it("propagates a task failure out of the buffer, halting its other tasks", function* () { + let siblingCompleted = false; + let siblingHalted = false; + let bufferOutcome = "pending"; + let error: Error | undefined; + + try { + yield* run(function* () { + const buffer = yield* useTaskBuffer(5); + + yield* buffer.spawn(function* () { + try { + yield* suspend(); + siblingCompleted = true; + } finally { + siblingHalted = !siblingCompleted; + } + }); + + yield* buffer.spawn(function* () { + throw new Error("boom"); + }); + + try { + yield* buffer; + bufferOutcome = "returned"; + } catch { + bufferOutcome = "threw"; + } + }); + } catch (e) { + error = e as Error; + } + + expect(error?.message).toEqual("boom"); + expect(siblingHalted).toEqual(true); + // the scope unwinds before `yield* buffer` can settle either way + expect(bufferOutcome).toEqual("pending"); + }); + + it("contains a failure that the spawned operation handles itself", function* () { + let siblingCompleted = false; + + yield* run(function* () { + const buffer = yield* useTaskBuffer(5); + + yield* buffer.spawn(function* () { + try { + throw new Error("boom"); + } catch { + // handled inside the operation, so it never reaches the buffer + } + }); + + yield* buffer.spawn(function* () { + siblingCompleted = true; + }); + + yield* buffer; + }); + + expect(siblingCompleted).toEqual(true); + }); }); diff --git a/task-buffer/task-buffer.ts b/task-buffer/task-buffer.ts index 591d0135..2e04bd3d 100644 --- a/task-buffer/task-buffer.ts +++ b/task-buffer/task-buffer.ts @@ -4,7 +4,6 @@ import { type Operation, type Resolve, type Result, - type Stream, type Task, createChannel, resource, @@ -20,15 +19,12 @@ import { */ export interface TaskBuffer extends Operation { /** - * Spawn `op` in the task buffer when there is room available. If - * there is room, then this operation will complete immediately. - * Otherwise, it will return once there is room in the buffer and - * the task is successfully spawned. - * `spawn()` operation will not return until the task has actually - * been spawned. + * Submit `op` to the task buffer. This operation returns as soon as the + * request has been queued; it does not wait for `op` to be spawned. * * @param op - the operation to spawn in the buffer. - * @returns the spawned task. + * @returns an operation that resolves with the spawned {@link Task} once + * there is room in the buffer and `op` has been spawned. */ spawn(op: () => Operation): Operation>>; } @@ -69,11 +65,21 @@ export function useTaskBuffer(max: number): Operation { let requests: SpawnRequest[] = []; + // Halting an active task frees a slot, and the dispatch loop is the last + // child this resource tears down, so without this it would admit queued + // work in the window between the two. + let closed = false; + + // Subscribe before the loop starts. Re-subscribing per iteration drops any + // send that lands before the new subscription is established. + let inputs = yield* input; + let outputs = yield* output; + yield* spawn(function* () { while (true) { if (requests.length === 0) { - yield* next(input); - } else if (buffer.size < max) { + yield* inputs.next(); + } else if (!closed && buffer.size < max) { const request = requests.pop()!; let task = yield* scope.spawn(request.operation); buffer.add(task); @@ -89,28 +95,46 @@ export function useTaskBuffer(max: number): Operation { }); request.resolve(task); } else { - yield* next(output); + yield* outputs.next(); } } }); - yield* provide({ - *[Symbol.iterator]() { - let outputs = yield* output; - while (buffer.size > 0 || requests.length > 0) { - yield* outputs.next(); - } - }, - *spawn(fn: () => Operation) { - let { operation, resolve } = withResolvers>(); - requests.unshift({ - operation: fn, - resolve: resolve as Resolve, - }); - yield* input.send(); - return operation; - }, - }); + try { + yield* provide({ + *[Symbol.iterator]() { + let results = yield* output; + while (buffer.size > 0 || requests.length > 0) { + yield* results.next(); + } + }, + *spawn(fn: () => Operation) { + let { operation, resolve } = withResolvers>(); + let request: SpawnRequest = { + operation: fn, + resolve: resolve as Resolve, + }; + requests.unshift(request); + yield* input.send(); + return { + *[Symbol.iterator]() { + try { + return yield* operation; + } finally { + // Abandoning the wait withdraws the request, so work nobody is + // waiting for is never admitted. + let index = requests.indexOf(request); + if (index !== -1) { + requests.splice(index, 1); + } + } + }, + }; + }, + }); + } finally { + closed = true; + } }); } @@ -118,10 +142,3 @@ interface SpawnRequest { operation(): Operation; resolve: Resolve>; } - -function* next( - stream: Stream, -): Operation> { - let subscription = yield* stream; - return yield* subscription.next(); -}