diff --git a/.changeset/fix-multipart-active-file-failure.md b/.changeset/fix-multipart-active-file-failure.md new file mode 100644 index 00000000000..bff8687e187 --- /dev/null +++ b/.changeset/fix-multipart-active-file-failure.md @@ -0,0 +1,5 @@ +--- +"effect": patch +--- + +Fix multipart file streams hanging on body read errors and preserve error causes when persisting files. diff --git a/packages/effect/src/unstable/http/Multipart.ts b/packages/effect/src/unstable/http/Multipart.ts index f60f4aceaff..9735b07fdba 100644 --- a/packages/effect/src/unstable/http/Multipart.ts +++ b/packages/effect/src/unstable/http/Multipart.ts @@ -199,6 +199,11 @@ export interface Persisted { const MultipartErrorTypeId = "~effect/http/Multipart/MultipartError" +const toMultipartError = (cause: unknown): MultipartError => + Predicate.hasProperty(cause, MultipartErrorTypeId) + ? cause as MultipartError + : MultipartError.fromReason("InternalError", cause) + /** * Error reason carried by a `MultipartError`. * @@ -480,16 +485,24 @@ export const makeChannel = (headers: Record): Channel.Channe let chunks: Array = [] let finished = false const pullChunks = Channel.fromPull( - Effect.succeed(Effect.suspend(function loop(): Pull.Pull> { - if (!Arr.isReadonlyArrayNonEmpty(chunks)) { - return finished ? Cause.done() : Effect.flatMap(pump, loop) - } - const chunk = chunks - chunks = [] - return Effect.succeed(chunk) - })) + Effect.succeed( + Effect.suspend(function loop(): Pull.Pull, IE | MultipartError> { + if (!Arr.isReadonlyArrayNonEmpty(chunks)) { + if (finished) { + return Cause.done() + } + if (Option.isSome(exit)) { + return exit.value + } + return Effect.flatMap(pump, loop) + } + const chunk = chunks + chunks = [] + return Effect.succeed(chunk) + }) + ) ) - partsBuffer.push(new FileImpl(info, pullChunks)) + partsBuffer.push(new FileImpl(info, Channel.mapError(pullChunks, toMultipartError))) return function(chunk) { if (chunk === null) { finished = true @@ -614,17 +627,14 @@ class FileImpl extends PartBase implements File { constructor( info: MP.PartInfo, - channel: Channel.Channel> + channel: Channel.Channel, MultipartError> ) { super() this.key = info.name this.name = info.filename ?? info.name this.contentType = info.contentType this.content = Stream.fromChannel(channel) - this.contentEffect = channel.pipe( - Channel.mkUint8Array, - Effect.mapError((cause) => MultipartError.fromReason("InternalError", cause)) - ) + this.contentEffect = Channel.mkUint8Array(channel) } toJSON(): unknown { @@ -641,11 +651,7 @@ class FileImpl extends PartBase implements File { const defaultWriteFile = (path: string, file: File) => Effect.flatMap( FileSystem.FileSystem, - (fs) => - Effect.mapError( - Stream.run(file.content, fs.sink(path)), - (cause) => MultipartError.fromReason("InternalError", cause) - ) + (fs) => Effect.mapError(Stream.run(file.content, fs.sink(path)), toMultipartError) ) /** diff --git a/packages/effect/test/unstable/http/Multipart.test.ts b/packages/effect/test/unstable/http/Multipart.test.ts index 8c764aa578a..b8965d12054 100644 --- a/packages/effect/test/unstable/http/Multipart.test.ts +++ b/packages/effect/test/unstable/http/Multipart.test.ts @@ -1,5 +1,5 @@ import { describe, it } from "@effect/vitest" -import { ByteSize, Effect, ErrorReporter, FileSystem, identity, Path, Schema, Stream, Unify } from "effect" +import { ByteSize, Effect, ErrorReporter, FileSystem, identity, Path, Schema, Sink, Stream, Unify } from "effect" import { HttpClientRequest, HttpIncomingMessage, @@ -240,6 +240,84 @@ describe("Multipart", () => { strictEqual(error.reason._tag, "BodyTooLarge") })) + const activeFileParts = (error: Multipart.MultipartError) => + Stream.make( + new TextEncoder().encode( + "--b\r\nContent-Disposition: form-data; name=\"file\"; filename=\"a.txt\"\r\n\r\nhello" + ) + ).pipe( + Stream.concat(Stream.fail(error)), + Stream.pipeThroughChannel(Multipart.makeChannel({ "content-type": "multipart/form-data; boundary=b" })) + ) + + it.live("propagates upstream failure while consuming an active file", () => + Effect.gen(function*() { + const upstreamError = Multipart.MultipartError.fromReason("InternalError", new Error("body-read-failed")) + let bytesRead = 0 + + const error = yield* activeFileParts(upstreamError).pipe( + Stream.runForEach((part) => + part._tag === "File" + ? Stream.runForEach(part.content, (chunk) => + Effect.sync(() => { + bytesRead += chunk.length + })).pipe(Effect.flatMap(() => Effect.die("content should have failed"))) + : Effect.die("expected file") + ), + Effect.timeout("1 second"), + Effect.flip + ) + + strictEqual(bytesRead, 5) + strictEqual(error, upstreamError) + })) + + it.live("preserves upstream failure when collecting an active file with contentEffect", () => + Effect.gen(function*() { + const upstreamError = Multipart.MultipartError.fromReason("InternalError", new Error("body-read-failed")) + + const error = yield* activeFileParts(upstreamError).pipe( + Stream.runForEach((part) => + part._tag === "File" + ? Effect.flatMap(part.contentEffect, () => Effect.die("contentEffect should have failed")) + : Effect.die("expected file") + ), + Effect.timeout("1 second"), + Effect.flip + ) + + strictEqual(error, upstreamError) + })) + + it.live("preserves the upstream error cause when persisting an active file", () => + Effect.gen(function*() { + const upstreamError = Multipart.MultipartError.fromReason("InternalError", new Error("body-read-failed")) + let bytesWritten = 0 + + const error = yield* activeFileParts(upstreamError).pipe( + Multipart.toPersisted, + Effect.provideService( + FileSystem.FileSystem, + FileSystem.makeNoop({ + makeTempDirectoryScoped: () => Effect.succeed("/tmp/multipart-test"), + sink: () => + Sink.forEach((chunk: Uint8Array) => + Effect.sync(() => { + bytesWritten += chunk.length + }) + ) + }) + ), + Effect.provide(Path.layer), + Effect.scoped, + Effect.timeout("1 second"), + Effect.flip + ) + + strictEqual(bytesWritten, 5) + strictEqual(error, upstreamError) + })) + it.effect("propagates Parse when the body ends mid-file", () => Effect.gen(function*() { const boundary = "----testboundary" diff --git a/packages/platform/bun/test/BunMultipart.test.ts b/packages/platform/bun/test/BunMultipart.test.ts new file mode 100644 index 00000000000..6519381d9c6 --- /dev/null +++ b/packages/platform/bun/test/BunMultipart.test.ts @@ -0,0 +1,48 @@ +import * as BunMultipart from "@effect/platform-bun/BunMultipart" +import { describe, it } from "@effect/vitest" +import { assertTrue, strictEqual } from "@effect/vitest/utils" +import { Effect, Stream } from "effect" + +describe("BunMultipart", () => { + it.live("propagates a request body read error while consuming an active file", () => + Effect.gen(function*() { + const cause = new Error("body-read-failed") + const body = new ReadableStream({ + start(controller) { + controller.enqueue(new TextEncoder().encode( + "--b\r\nContent-Disposition: form-data; name=\"file\"; filename=\"a.txt\"\r\n\r\nhello" + )) + }, + pull(controller) { + controller.error(cause) + } + }, { + // Delay the error until the file bytes are consumed. + highWaterMark: 0 + }) + const request = new Request("http://localhost/upload", { + method: "POST", + headers: { "content-type": "multipart/form-data; boundary=b" }, + body + }) + let bytesRead = 0 + + const error = yield* BunMultipart.stream(request).pipe( + Stream.runForEach((part) => + part._tag === "File" + ? Stream.runForEach(part.content, (chunk) => + Effect.sync(() => { + bytesRead += chunk.length + })) + : Effect.die("expected file") + ), + Effect.timeout("1 second"), + Effect.flip + ) + + strictEqual(bytesRead, 5) + assertTrue(error._tag === "MultipartError") + strictEqual(error.reason._tag, "InternalError") + strictEqual(error.reason.cause, cause) + })) +})