import {ActorRefFrom, ActorSystem, AnyEventObject, ObservableActorLogic} from "xstate"; import { fromEventAsyncGenerator} from "./generator"; import {batchAsync, BatchAsyncParams, mapAsync} from "@/iterator"; type MergeEvents = { (events: TIn[]): TOut } type AsyncBatchEventsInput = { merge?: MergeEvents; }& BatchAsyncParams type AsyncBatchEventsGenerator=AsyncBatchEventsInput >= { ({input, system, self, emit}:{ input: TInput; system: ActorSystem; self: ActorRefFrom>; emit: (emitted: AnyEventObject) => void; }):AsyncGenerator | Promise> } function defaultMerge( events: TIn[]) { return { type: `batch`, batch: events } } export function fromAsyncBatchEventGenerator=AsyncBatchEventsInput>(generator: ({input, system, self, emit}: { input: TInput; system: ActorSystem; self: ActorRefFrom>; emit: (emitted: AnyEventObject) => void; })=> AsyncGenerator | Promise>) { return fromEventAsyncGenerator(async function* ({input, system, self, emit}) { const merge = input.merge || defaultMerge as unknown as MergeEvents yield* mapAsync(batchAsync(await generator({input, system, self, emit}), input.split), merge) }) } export const asyncBatchEvents = fromAsyncBatchEventGenerator( async function* ({input}) { for await (const event of input.stream) { yield event; } })