diff --git a/packages/core/src/event.ts b/packages/core/src/event.ts index 03c539642b..8641c92859 100644 --- a/packages/core/src/event.ts +++ b/packages/core/src/event.ts @@ -118,6 +118,11 @@ export type SubscribePayload = D[number] extend : never export interface Subscribe { + /** + * Volatile live channel: every event published from now on, nothing before or + * across a disconnect. Consumers that need reliability combine it with `log`. + */ + (): Stream.Stream (definition: D): Stream.Stream> (definitions: D): Stream.Stream> } @@ -131,12 +136,6 @@ export interface Interface { options?: PublishOptions, ) => Effect.Effect> readonly subscribe: Subscribe - /** - * Volatile live channel: every event published from now on, nothing before, - * nothing across a disconnect. The only channel that carries non-durable - * events; consumers that need reliability combine `changes` with `log`. - */ - readonly live: () => Stream.Stream /** * Durable, ordered, gap-free per-aggregate log read. `follow: false` * completes at the end of the log; `follow: true` replays then transitions @@ -150,7 +149,7 @@ export interface Interface { }) => Stream.Stream /** Latest committed seq per aggregate. Aggregates without events are absent. */ readonly sequences: (aggregateIDs: ReadonlyArray) => Effect.Effect> - /** @deprecated Use `all()` and consume the returned stream. */ + /** @deprecated Use `subscribe()` and consume the returned stream. */ readonly listen: (listener: Subscriber) => Effect.Effect readonly project: (definition: D, projector: Subscriber) => Effect.Effect readonly replay: ( @@ -575,17 +574,16 @@ export const layerWith = (options?: LayerOptions) => ), ) - const subscribeOne = (definition: D): Stream.Stream> => - local(Stream.unwrap(getOrCreate(definition).pipe(Effect.map((pubsub) => Stream.fromPubSub(pubsub))))).pipe( - Stream.map((event) => event as Payload), - ) - + function subscribe(): Stream.Stream function subscribe(definition: D): Stream.Stream> function subscribe( definitions: D, ): Stream.Stream> - function subscribe(input: Definition | readonly Definition[]): Stream.Stream { - if (isDefinition(input)) return subscribeOne(input) + function subscribe(input?: Definition | readonly Definition[]): Stream.Stream { + if (input === undefined) return streamLive() + if (isDefinition(input)) { + return local(Stream.unwrap(getOrCreate(input).pipe(Effect.map((pubsub) => Stream.fromPubSub(pubsub))))) + } const types = new Set(input.map((definition) => definition.type)) return streamLive().pipe(Stream.filter((event) => types.has(event.type))) } @@ -737,7 +735,6 @@ export const layerWith = (options?: LayerOptions) => return Service.of({ publish, subscribe, - live: streamLive, log, sequences, listen, diff --git a/packages/core/src/plugin/host.ts b/packages/core/src/plugin/host.ts index 03c9285414..01ea0da0d3 100644 --- a/packages/core/src/plugin/host.ts +++ b/packages/core/src/plugin/host.ts @@ -159,7 +159,7 @@ export const make = Effect.fn("PluginHost.make")(function* (plugin: PluginV2.Int }), }, event: { - subscribe: () => events.live().pipe(Stream.filter(EventManifest.isServer)), + subscribe: () => events.subscribe().pipe(Stream.filter(EventManifest.isServer)), }, integration: { list: () => response(integration.list()), diff --git a/packages/core/test/event.test.ts b/packages/core/test/event.test.ts index 1f0e872063..4bff6a7abf 100644 --- a/packages/core/test/event.test.ts +++ b/packages/core/test/event.test.ts @@ -177,7 +177,7 @@ describe("EventV2", () => { Effect.gen(function* () { const events = yield* EventV2.Service const typed = yield* events.subscribe(Message).pipe(Stream.take(1), Stream.runCollect, Effect.forkScoped) - const wildcard = yield* events.live().pipe(Stream.take(1), Stream.runCollect, Effect.forkScoped) + const wildcard = yield* events.subscribe().pipe(Stream.take(1), Stream.runCollect, Effect.forkScoped) yield* Effect.yieldNow const event = yield* events.publish(Message, { text: "hello" }) @@ -258,7 +258,7 @@ describe("EventV2", () => { Effect.gen(function* () { const events = yield* EventV2.Service const received = new Array() - const fiber = yield* events.live().pipe( + const fiber = yield* events.subscribe().pipe( Stream.take(1), Stream.runForEach(() => Effect.sync(() => received.push("stream"))), Effect.forkScoped, diff --git a/packages/core/test/session-runner-tool-events.test.ts b/packages/core/test/session-runner-tool-events.test.ts index 983dca8c2c..de4700ca07 100644 --- a/packages/core/test/session-runner-tool-events.test.ts +++ b/packages/core/test/session-runner-tool-events.test.ts @@ -27,7 +27,6 @@ const capture = () => { return event }), subscribe: () => Stream.empty, - live: () => Stream.empty, log: () => Stream.empty, sequences: () => Effect.succeed(new Map()), listen: () => Effect.succeed(Effect.void),