import { setTimeout as sleep } from 'node:timers/promises'; import { mo } from 'moroutine'; import type { Int32Atomic } from 'moroutine'; // Worker-side generator that yields [0..n) and bumps `emitted` after each // yield resumes — i.e. after the value has been accepted by the pipe. Used // to observe producer-side backpressure from the parent. export const countingGen = mo(import.meta, async function* (n: number, emitted: Int32Atomic) { for (let i = 0; i < n; i++) { yield i; emitted.add(1); } }); // Worker-side consumer that drains an AsyncIterable slowly and records the // peak observed `emitted - consumed` gap. Returns the peak. export const slowSumPeakGap = mo( import.meta, async (input: AsyncIterable, emitted: Int32Atomic, delayMs: number): Promise => { let consumed = 0; let peak = 0; for await (const _ of input) { consumed++; const gap = emitted.load() - consumed; if (gap > peak) peak = gap; await sleep(delayMs); } return peak; }, );