diff --git a/.changeset/tidy-cats-stream.md b/.changeset/tidy-cats-stream.md new file mode 100644 index 00000000000..c14e87c6a17 --- /dev/null +++ b/.changeset/tidy-cats-stream.md @@ -0,0 +1,5 @@ +--- +"effect": patch +--- + +Preserve lexical ordering in streaming template interpolation. diff --git a/packages/effect/src/unstable/http/Template.ts b/packages/effect/src/unstable/http/Template.ts index bc283584232..994ecd461fd 100644 --- a/packages/effect/src/unstable/http/Template.ts +++ b/packages/effect/src/unstable/http/Template.ts @@ -215,11 +215,12 @@ export function stream>( buffer = "" } - return Stream.flatMap( - Stream.fromIterable(chunks), - (chunk) => - typeof chunk === "string" ? Stream.succeed(chunk) : Effect.isEffect(chunk) ? Stream.fromEffect(chunk) : chunk, - { concurrency: "unbounded" } + return Stream.fromIterable(chunks).pipe( + Stream.mapEffect( + (chunk) => Effect.isEffect(chunk) ? chunk : Effect.succeed(chunk), + { concurrency: "unbounded" } + ), + Stream.flatMap((chunk) => typeof chunk === "string" ? Stream.succeed(chunk) : chunk) ) } diff --git a/packages/effect/test/unstable/http/Template.test.ts b/packages/effect/test/unstable/http/Template.test.ts new file mode 100644 index 00000000000..bbe76e7bb10 --- /dev/null +++ b/packages/effect/test/unstable/http/Template.test.ts @@ -0,0 +1,38 @@ +import { assert, describe, it } from "@effect/vitest" +import { Deferred, Effect, Fiber, Stream } from "effect" +import { TestClock } from "effect/testing" +import { Template } from "effect/unstable/http" + +describe("Template", () => { + it.effect("preserves template segment order", () => + Effect.gen(function*() { + const fiber = yield* Stream.runCollect( + Template.stream`a${Effect.delay(Effect.succeed("slow"), "1 second")}b${"fast"}c` + ).pipe(Effect.forkChild) + yield* Effect.yieldNow + yield* TestClock.adjust("1 second") + const chunks = yield* Fiber.join(fiber) + assert.strictEqual(chunks.join(""), "aslowbfastc") + })) + + it.effect("evaluates effect interpolations concurrently", () => + Effect.gen(function*() { + const firstStarted = yield* Deferred.make() + const secondStarted = yield* Deferred.make() + const releaseFirst = yield* Deferred.make() + const fiber = yield* Template.stream`${ + Deferred.succeed(firstStarted, void 0).pipe( + Effect.andThen(Deferred.await(releaseFirst)), + Effect.as("first") + ) + }${Deferred.succeed(secondStarted, void 0).pipe(Effect.as("second"))}`.pipe( + Stream.runCollect, + Effect.forkChild + ) + yield* Deferred.await(firstStarted) + yield* Effect.yieldNow + assert.isTrue(yield* Deferred.isDone(secondStarted)) + yield* Deferred.succeed(releaseFirst, void 0) + assert.deepStrictEqual(yield* Fiber.join(fiber), ["first", "second"]) + })) +})