diff --git a/.changeset/fix-event-serial.md b/.changeset/fix-event-serial.md new file mode 100644 index 000000000000..a3f94086b71a --- /dev/null +++ b/.changeset/fix-event-serial.md @@ -0,0 +1,5 @@ +--- +"opencode": patch +--- + +Run EventV2 serial dispatch listeners sequentially instead of concurrently. diff --git a/packages/core/src/event.ts b/packages/core/src/event.ts index 35c14d23d92b..acdd595b11d2 100644 --- a/packages/core/src/event.ts +++ b/packages/core/src/event.ts @@ -683,10 +683,7 @@ export const layerWith = (options?: LayerOptions) => ): Effect.Effect> => Effect.suspend(() => { const event = placeholderEvent(definition) - return Effect.forEach(listeners, (listener) => listener(event), { - concurrency: "unbounded", - discard: false, - }) + return Effect.forEach(listeners, (listener) => listener(event)) }) const parallel = ( diff --git a/packages/core/test/event.test.ts b/packages/core/test/event.test.ts index e4329a2dde98..4396a792e893 100644 --- a/packages/core/test/event.test.ts +++ b/packages/core/test/event.test.ts @@ -268,6 +268,57 @@ describe("EventV2", () => { }), ) + it.effect("runs serial dispatch listeners one at a time", () => + Effect.gen(function* () { + const events = yield* EventV2.Service + const firstStarted = yield* Deferred.make() + const releaseFirst = yield* Deferred.make() + const secondStarted = yield* Deferred.make() + const dispatch = yield* events + .serial(Message, [ + () => + Deferred.succeed(firstStarted, undefined).pipe( + Effect.andThen(Deferred.await(releaseFirst)), + Effect.as("first"), + ), + () => Deferred.succeed(secondStarted, undefined).pipe(Effect.as("second")), + ]) + .pipe(Effect.forkChild) + + yield* Deferred.await(firstStarted) + yield* Effect.yieldNow + expect(yield* Deferred.isDone(secondStarted)).toBe(false) + + yield* Deferred.succeed(releaseFirst, undefined) + expect(yield* Fiber.join(dispatch)).toEqual(["first", "second"]) + expect(yield* Deferred.isDone(secondStarted)).toBe(true) + }), + ) + + it.effect("runs parallel dispatch listeners concurrently", () => + Effect.gen(function* () { + const events = yield* EventV2.Service + const firstStarted = yield* Deferred.make() + const releaseFirst = yield* Deferred.make() + const secondStarted = yield* Deferred.make() + const dispatch = yield* events + .parallel(Message, [ + () => + Deferred.succeed(firstStarted, undefined).pipe( + Effect.andThen(Deferred.await(releaseFirst)), + Effect.as("first"), + ), + () => Deferred.succeed(secondStarted, undefined).pipe(Effect.as("second")), + ]) + .pipe(Effect.forkChild) + + yield* Deferred.await(firstStarted) + yield* Deferred.await(secondStarted) + yield* Deferred.succeed(releaseFirst, undefined) + expect(yield* Fiber.join(dispatch)).toEqual(["first", "second"]) + }), + ) + it.effect("isolates observer defects after durable events commit", () => Effect.gen(function* () { const events = yield* EventV2.Service