import { createChannel, Err, Ok, type Operation, type Resolve, resource, type Result, spawn, type Stream, type Task, useScope, withResolvers, } from "effection"; export interface TaskBuffer extends Operation { spawn(op: () => Operation): Operation>>; } export function useTaskBuffer(max: number): Operation { return resource(function* (provide) { let input = createChannel(); let output = createChannel, never>(); let buffer = new Set>(); let scope = yield* useScope(); let requests: SpawnRequest[] = []; yield* spawn(function* () { while (true) { if (requests.length === 0) { yield* next(input); } else if (buffer.size < max) { let request = requests.pop()!; let task = yield* scope.spawn(request.operation); buffer.add(task); yield* spawn(function* () { try { let result = Ok(yield* task); buffer.delete(task); yield* output.send(result); } catch (error) { buffer.delete(task); yield* output.send(Err(error as Error)); } }); request.resolve(task); } else { yield* next(output); } } }); 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; }, }); }); } interface SpawnRequest { operation(): Operation; resolve: Resolve>; } function* next( stream: Stream, ): Operation> { let subscription = yield* stream; return yield* subscription.next(); }