(definition: D, event: Payload<D>, commit?: PublishOptions["commit"])
| 367 | } |
| 368 | |
| 369 | function publishEvent<D extends Definition>(definition: D, event: Payload<D>, commit?: PublishOptions["commit"]) { |
| 370 | return Effect.gen(function* () { |
| 371 | if (!definition?.durable && commit) |
| 372 | return yield* Effect.die( |
| 373 | new InvalidDurableEventError({ |
| 374 | type: event.type, |
| 375 | message: "Local commit hooks require a durable event", |
| 376 | }), |
| 377 | ) |
| 378 | if (definition?.durable) { |
| 379 | const committed = yield* commitDurableEvent(definition, event as Payload, undefined, commit) |
| 380 | if (committed) { |
| 381 | event = { |
| 382 | ...event, |
| 383 | durable: { |
| 384 | aggregateID: committed.aggregateID, |
| 385 | seq: committed.seq, |
| 386 | version: definition.durable.version, |
| 387 | }, |
| 388 | } |
| 389 | yield* notify(event as Payload, true) |
| 390 | return event |
| 391 | } |
| 392 | } |
| 393 | yield* notify(event as Payload, false) |
| 394 | return event |
| 395 | }) |
| 396 | } |
| 397 | |
| 398 | const observe = (event: Payload, observer: (event: Payload) => Effect.Effect<void>) => |
| 399 | Effect.suspend(() => observer(event)).pipe( |
no test coverage detected