| 539 | const streamAll = (): Stream.Stream<Payload> => Stream.fromPubSub(pubsub.all) |
| 540 | |
| 541 | const readAfter = (aggregateID: string, after: number) => |
| 542 | (options?.beforeAggregateRead?.(aggregateID) ?? Effect.void).pipe( |
| 543 | Effect.andThen( |
| 544 | db |
| 545 | .select() |
| 546 | .from(EventTable) |
| 547 | .where(and(eq(EventTable.aggregate_id, aggregateID), gt(EventTable.seq, after))) |
| 548 | .orderBy(asc(EventTable.seq)) |
| 549 | .all(), |
| 550 | ), |
| 551 | Effect.orDie, |
| 552 | Effect.map((rows) => |
| 553 | rows.map((event) => |
| 554 | decodeSerializedEvent({ |
| 555 | id: event.id, |
| 556 | aggregateID: event.aggregate_id, |
| 557 | seq: event.seq, |
| 558 | type: event.type, |
| 559 | data: event.data, |
| 560 | }), |
| 561 | ), |
| 562 | ), |
| 563 | ) |
| 564 | |
| 565 | const subscribeDurable = (aggregateID: string) => |
| 566 | Effect.gen(function* () { |