From 35cb2eb8ead554efdc7b4fc435e39433d5fa35fc Mon Sep 17 00:00:00 2001 From: Taras Mankovski Date: Fri, 18 Sep 2026 21:50:59 -0400 Subject: [PATCH 1/5] =?UTF-8?q?=F0=9F=93=9D=20Correct=20the=20`task-buffer?= =?UTF-8?q?`=20README=20example=20and=20`spawn()`=20docs?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The README compressed submission, admission and completion into a single `yield* yield* buffer.spawn(...)` expression, and the `spawn()` doc comment claimed the operation "will not return until the task has actually been spawned" — which contradicts both its `Operation>>` return type and what it does. Name the intermediate values so the four steps are distinct, and describe the verified failure and cancellation behavior: a task error tears down the enclosing scope, `yield* buffer` never reports it, and only handling the error inside the spawned operation contains it. Tests cover each documented claim. No runtime change. --- task-buffer/README.md | 44 +++++++++--- task-buffer/package.json | 2 +- task-buffer/task-buffer.test.ts | 123 +++++++++++++++++++++++++++++++- task-buffer/task-buffer.ts | 11 ++- 4 files changed, 163 insertions(+), 17 deletions(-) diff --git a/task-buffer/README.md b/task-buffer/README.md index 271543b4..54f8c089 100644 --- a/task-buffer/README.md +++ b/task-buffer/README.md @@ -9,25 +9,53 @@ 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; - // Wait for all spawned tasks to complete + // 3. wait for that one task to run to completion and produce its result. + yield* task; + + // 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. 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..520c5abc 100644 --- a/task-buffer/task-buffer.test.ts +++ b/task-buffer/task-buffer.test.ts @@ -1,5 +1,14 @@ import FakeTimers from "@sinonjs/fake-timers"; -import { sleep, spawn, type Task, until } from "effection"; +import { + createScope, + run, + sleep, + spawn, + suspend, + type Task, + until, + withResolvers, +} from "effection"; import { describe, it } from "@effectionx/vitest"; import { expect } from "expect"; import { useTaskBuffer } from "./task-buffer.ts"; @@ -56,4 +65,116 @@ 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("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..b74b0a86 100644 --- a/task-buffer/task-buffer.ts +++ b/task-buffer/task-buffer.ts @@ -20,15 +20,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>>; } From f7a1d27abd6b051f1ff52654c74c2729b4db44fc Mon Sep 17 00:00:00 2001 From: Taras Mankovski Date: Fri, 18 Sep 2026 21:51:31 -0400 Subject: [PATCH 2/5] =?UTF-8?q?=F0=9F=90=9B=20Stop=20`TaskBuffer`=20from?= =?UTF-8?q?=20dropping=20a=20spawn=20request?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The dispatch loop resubscribed to its `input` channel on every iteration. Effection channels do not buffer for non-subscribers, so a `spawn()` whose `send()` landed between the emptiness check and the new subscription was lost: the request stayed in `requests`, was never spawned even with capacity to spare, and `yield* buffer` waited on it forever. A later `spawn()` did not recover it. Subscribe once during resource setup, before `provide()` can hand the buffer to a caller, so no send can arrive without a subscriber. This retires the `next()` helper. The regression test hangs without this fix. --- task-buffer/package.json | 2 +- task-buffer/task-buffer.test.ts | 20 ++++++++++++++++++++ task-buffer/task-buffer.ts | 21 +++++++++------------ 3 files changed, 30 insertions(+), 13 deletions(-) diff --git a/task-buffer/package.json b/task-buffer/package.json index 09d763dd..953415c9 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.4", + "version": "1.3.5", "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 520c5abc..e53ad09a 100644 --- a/task-buffer/task-buffer.test.ts +++ b/task-buffer/task-buffer.test.ts @@ -114,6 +114,26 @@ describe("TaskBuffer", () => { expect(queuedStarted).toEqual(false); }); + 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("propagates a task failure out of the buffer, halting its other tasks", function* () { let siblingCompleted = false; let siblingHalted = false; diff --git a/task-buffer/task-buffer.ts b/task-buffer/task-buffer.ts index b74b0a86..ea456781 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, @@ -66,10 +65,15 @@ export function useTaskBuffer(max: number): Operation { let requests: SpawnRequest[] = []; + // 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); + yield* inputs.next(); } else if (buffer.size < max) { const request = requests.pop()!; let task = yield* scope.spawn(request.operation); @@ -86,16 +90,16 @@ export function useTaskBuffer(max: number): Operation { }); request.resolve(task); } else { - yield* next(output); + yield* outputs.next(); } } }); yield* provide({ *[Symbol.iterator]() { - let outputs = yield* output; + let results = yield* output; while (buffer.size > 0 || requests.length > 0) { - yield* outputs.next(); + yield* results.next(); } }, *spawn(fn: () => Operation) { @@ -115,10 +119,3 @@ interface SpawnRequest { operation(): Operation; resolve: Resolve>; } - -function* next( - stream: Stream, -): Operation> { - let subscription = yield* stream; - return yield* subscription.next(); -} From 68b268f8670f1a0e37ca32a3356f52ed41933e80 Mon Sep 17 00:00:00 2001 From: Taras Mankovski Date: Fri, 18 Sep 2026 21:51:55 -0400 Subject: [PATCH 3/5] =?UTF-8?q?=F0=9F=90=9B=20Withdraw=20a=20`TaskBuffer`?= =?UTF-8?q?=20request=20when=20its=20admission=20is=20abandoned?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Nothing removed an entry from `requests`, so a request whose caller was halted while waiting on `yield* admission` was still admitted once room appeared, and the operation ran with nobody waiting for it. Return an operation that splices the request out when the wait for it unwinds. The cleanup is synchronous, so `finally` is safe here. Submitting without ever waiting for admission is unchanged: that request stays queued. This is an observable change in semantics rather than a pure fix, hence the minor bump. --- task-buffer/README.md | 4 ++++ task-buffer/package.json | 2 +- task-buffer/task-buffer.test.ts | 27 +++++++++++++++++++++++++++ task-buffer/task-buffer.ts | 20 +++++++++++++++++--- 4 files changed, 49 insertions(+), 4 deletions(-) diff --git a/task-buffer/README.md b/task-buffer/README.md index 54f8c089..b867dc80 100644 --- a/task-buffer/README.md +++ b/task-buffer/README.md @@ -59,3 +59,7 @@ from tearing down the caller, handle it inside the operation passed to 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 953415c9..33ea7eb5 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.5", + "version": "1.4.0", "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 e53ad09a..51182ac4 100644 --- a/task-buffer/task-buffer.test.ts +++ b/task-buffer/task-buffer.test.ts @@ -134,6 +134,33 @@ describe("TaskBuffer", () => { 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; diff --git a/task-buffer/task-buffer.ts b/task-buffer/task-buffer.ts index ea456781..7234b056 100644 --- a/task-buffer/task-buffer.ts +++ b/task-buffer/task-buffer.ts @@ -104,12 +104,26 @@ export function useTaskBuffer(max: number): Operation { }, *spawn(fn: () => Operation) { let { operation, resolve } = withResolvers>(); - requests.unshift({ + let request: SpawnRequest = { operation: fn, resolve: resolve as Resolve, - }); + }; + requests.unshift(request); yield* input.send(); - return operation; + 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); + } + } + }, + }; }, }); }); From 0726f514946236f98c4966f049c6cee34e3e8418 Mon Sep 17 00:00:00 2001 From: Taras Mankovski Date: Sun, 20 Sep 2026 09:51:36 -0400 Subject: [PATCH 4/5] Downgrade version from 1.4.0 to 1.3.4 --- task-buffer/package.json | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/task-buffer/package.json b/task-buffer/package.json index 33ea7eb5..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.4.0", + "version": "1.3.4", "keywords": ["concurrency"], "type": "module", "main": "./dist/mod.js", From d994547952aae11cf6bc14c98b3413f1424636ef Mon Sep 17 00:00:00 2001 From: Taras Mankovski Date: Sun, 20 Sep 2026 11:58:35 -0400 Subject: [PATCH 5/5] =?UTF-8?q?=F0=9F=90=9B=20Stop=20`TaskBuffer`=20from?= =?UTF-8?q?=20admitting=20queued=20work=20during=20scope=20exit?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Effection halts a task's children LIFO, so the dispatch loop — spawned first — is halted last. Every active task that settled ahead of it freed a slot, waking the loop to spawn queued requests into a scope that was already unwinding. Close the buffer in the resource body's `finally`, which runs before any child is halted, and skip admission once it is closed. --- task-buffer/task-buffer.test.ts | 62 ++++++++++++++++++++++++++++ task-buffer/task-buffer.ts | 71 +++++++++++++++++++-------------- 2 files changed, 102 insertions(+), 31 deletions(-) diff --git a/task-buffer/task-buffer.test.ts b/task-buffer/task-buffer.test.ts index 51182ac4..83a75a6b 100644 --- a/task-buffer/task-buffer.test.ts +++ b/task-buffer/task-buffer.test.ts @@ -1,9 +1,11 @@ import FakeTimers from "@sinonjs/fake-timers"; import { createScope, + ensure, run, sleep, spawn, + type Operation, suspend, type Task, until, @@ -13,6 +15,12 @@ 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(); @@ -114,6 +122,60 @@ describe("TaskBuffer", () => { 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(); diff --git a/task-buffer/task-buffer.ts b/task-buffer/task-buffer.ts index 7234b056..2e04bd3d 100644 --- a/task-buffer/task-buffer.ts +++ b/task-buffer/task-buffer.ts @@ -65,6 +65,11 @@ 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; @@ -74,7 +79,7 @@ export function useTaskBuffer(max: number): Operation { while (true) { if (requests.length === 0) { yield* inputs.next(); - } else if (buffer.size < max) { + } else if (!closed && buffer.size < max) { const request = requests.pop()!; let task = yield* scope.spawn(request.operation); buffer.add(task); @@ -95,37 +100,41 @@ export function useTaskBuffer(max: number): Operation { } }); - 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); + 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; + } }); }