import { call, createChannel, each, type Operation, race, resource, sleep, spawn, type Stream, suspend, until, } from "effection"; import { createInput, type InputEvent, type InputOptions } from "clayterm"; import { useStdin } from "./stdio.js"; function nothing() { return suspend() as unknown as Operation< IteratorResult >; } export function useInput( options?: InputOptions, ): Stream { return resource(function* (provide) { let input = yield* until(createInput(options)); let stdin = yield* useStdin(); let subscription = yield* stdin; let pending = nothing(); let events = createChannel(); yield* spawn(function* () { let next = yield* subscription.next(); while (!next.done) { let result = input.scan(next.value); pending = result.pending ? rescan(result.pending.delay) : nothing(); for (let event of result.events) { yield* events.send(event); } next = yield* race([subscription.next(), pending]); } yield* events.close(); }); yield* race([provide(yield* events), drain(events)]); }); } function rescan(delay: number): ReturnType { return call(function* (): Operation> { yield* sleep(delay); return { done: false, value: new Uint8Array(), }; }); } function* drain(stream: Stream): Operation { for (let _ of yield* each(stream)) { yield* each.next(); } }