Skip to content
Open
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
48 changes: 40 additions & 8 deletions task-buffer/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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<void> = 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.
2 changes: 1 addition & 1 deletion task-buffer/package.json
Original file line number Diff line number Diff line change
@@ -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",
Expand Down
232 changes: 231 additions & 1 deletion task-buffer/task-buffer.test.ts
Original file line number Diff line number Diff line change
@@ -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";
Comment thread
coderabbitai[bot] marked this conversation as resolved.
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<void> {
yield* until(Promise.resolve());
}

describe("TaskBuffer", () => {
it("queues up tasks when the buffer fills up", function* () {
const clock = FakeTimers.install();
Expand Down Expand Up @@ -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<string> = 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<void>();
const block = withResolvers<void>();
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<void>(), withResolvers<void>()];
const teardownStarted = [withResolvers<void>(), withResolvers<void>()];
const releaseTeardown = [withResolvers<void>(), withResolvers<void>()];
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<void>();
const secondRan = withResolvers<void>();

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<void>();
const queued = withResolvers<void>();
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);
});
});
Loading
Loading