(input: { readonly aggregateID: string; readonly after?: number })
| 583 | }) |
| 584 | |
| 585 | const durable = (input: { readonly aggregateID: string; readonly after?: number }): Stream.Stream<Payload> => |
| 586 | Stream.unwrap( |
| 587 | Effect.gen(function* () { |
| 588 | const wakes = yield* subscribeDurable(input.aggregateID) |
| 589 | let sequence = input.after ?? -1 |
| 590 | const read = Effect.suspend(() => readAfter(input.aggregateID, sequence)).pipe( |
| 591 | Effect.tap((events) => |
| 592 | Effect.sync(() => { |
| 593 | sequence = events.at(-1)?.durable?.seq ?? sequence |
| 594 | }), |
| 595 | ), |
| 596 | ) |
| 597 | const historical = yield* read |
| 598 | const live = Stream.fromSubscription(wakes).pipe( |
| 599 | Stream.mapEffect(() => read), |
| 600 | Stream.flattenIterable, |
| 601 | ) |
| 602 | return Stream.concat(Stream.fromIterable(historical), live) |
| 603 | }), |
| 604 | ) |
| 605 | |
| 606 | const listen = (listener: Subscriber): Effect.Effect<Unsubscribe> => |
| 607 | Effect.sync(() => { |
nothing calls this directly
no test coverage detected