From 8aae87866897c0b8bbb544bf9975177db4422f6d Mon Sep 17 00:00:00 2001 From: coyaSONG <66289470+coyaSONG@users.noreply.github.com> Date: Fri, 17 Jul 2026 07:05:36 +0900 Subject: [PATCH] Fix asyncPush bounded buffer termination --- .changeset/tame-streams-end.md | 5 +++++ packages/effect/src/internal/stream.ts | 4 ++-- packages/effect/test/Stream/async.test.ts | 15 +++++++++++++++ 3 files changed, 22 insertions(+), 2 deletions(-) create mode 100644 .changeset/tame-streams-end.md diff --git a/.changeset/tame-streams-end.md b/.changeset/tame-streams-end.md new file mode 100644 index 00000000000..6c4726838bd --- /dev/null +++ b/.changeset/tame-streams-end.md @@ -0,0 +1,5 @@ +--- +"effect": patch +--- + +Reserve queue capacity for termination signals in `Stream.asyncPush` with bounded dropping or sliding buffers. diff --git a/packages/effect/src/internal/stream.ts b/packages/effect/src/internal/stream.ts index 8c2d46d3d8b..36b0d9a3ef7 100644 --- a/packages/effect/src/internal/stream.ts +++ b/packages/effect/src/internal/stream.ts @@ -603,9 +603,9 @@ const queueFromBufferOptionsPush = ( } switch (options?.strategy) { case "sliding": - return Queue.sliding(options.bufferSize ?? 16) + return Queue.sliding((options.bufferSize ?? 16) + 1) default: - return Queue.dropping(options?.bufferSize ?? 16) + return Queue.dropping((options?.bufferSize ?? 16) + 1) } } diff --git a/packages/effect/test/Stream/async.test.ts b/packages/effect/test/Stream/async.test.ts index be4b1fdafe9..7043a74735e 100644 --- a/packages/effect/test/Stream/async.test.ts +++ b/packages/effect/test/Stream/async.test.ts @@ -438,6 +438,21 @@ describe("Stream", () => { assertTrue(Chunk.isEmpty(result)) })) + for (const strategy of ["dropping", "sliding"] as const) { + it.effect(`asyncPush - signals the end with a bounded ${strategy} buffer`, () => + Effect.gen(function*() { + const result = yield* Stream.asyncPush((emit) => { + emit.single(42) + emit.end() + return Effect.void + }, { bufferSize: 1, strategy }).pipe( + Stream.runCollect, + Effect.timeout("1 second") + ) + deepStrictEqual(Array.from(result), [42]) + })) + } + it.effect("asyncPush - handles errors", () => Effect.gen(function*() { const error = new Cause.RuntimeException("boom")