From b736120fd5ec09844148040693bf9331ec5fc02e Mon Sep 17 00:00:00 2001 From: Julius Marminge Date: Wed, 29 Jul 2026 01:16:55 +0200 Subject: [PATCH 1/8] refactor(server): use native HTTP compression streams Co-authored-by: codex --- apps/server/src/BunHttpResponseCompression.ts | 14 ++++++ .../src/HttpResponseCompression.test.ts | 43 +++++++++++++++++++ apps/server/src/HttpResponseCompression.ts | 12 ++++++ .../server/src/NodeHttpResponseCompression.ts | 13 ++++++ apps/server/src/http.test.ts | 24 ++++++++--- apps/server/src/http.ts | 16 ++++--- apps/server/src/server.test.ts | 2 + apps/server/src/server.ts | 38 +++++++++------- 8 files changed, 135 insertions(+), 27 deletions(-) create mode 100644 apps/server/src/BunHttpResponseCompression.ts create mode 100644 apps/server/src/HttpResponseCompression.test.ts create mode 100644 apps/server/src/HttpResponseCompression.ts create mode 100644 apps/server/src/NodeHttpResponseCompression.ts diff --git a/apps/server/src/BunHttpResponseCompression.ts b/apps/server/src/BunHttpResponseCompression.ts new file mode 100644 index 00000000000..a72f46895ae --- /dev/null +++ b/apps/server/src/BunHttpResponseCompression.ts @@ -0,0 +1,14 @@ +import * as Layer from "effect/Layer"; +import * as HttpServerResponse from "effect/unstable/http/HttpServerResponse"; + +import { HttpResponseCompression } from "./HttpResponseCompression.ts"; + +export const make = HttpResponseCompression.of({ + gzip: (body, options) => + HttpServerResponse.raw( + new Response(body).body!.pipeThrough(new CompressionStream("gzip")), + options, + ), +}); + +export const layer = Layer.succeed(HttpResponseCompression, make); diff --git a/apps/server/src/HttpResponseCompression.test.ts b/apps/server/src/HttpResponseCompression.test.ts new file mode 100644 index 00000000000..0be73d34d7f --- /dev/null +++ b/apps/server/src/HttpResponseCompression.test.ts @@ -0,0 +1,43 @@ +import { expect, it } from "vite-plus/test"; +import * as NodeStream from "node:stream"; +import * as NodeZlib from "node:zlib"; + +import * as BunHttpResponseCompression from "./BunHttpResponseCompression.ts"; +import * as NodeHttpResponseCompression from "./NodeHttpResponseCompression.ts"; + +const body = new TextEncoder().encode(`{"value":"${"compressible".repeat(1_000)}"}`); + +it("creates a raw Node gzip stream", async () => { + const response = NodeHttpResponseCompression.make.gzip(body, { + contentType: "application/json", + }); + + expect(response.body._tag).toBe("Raw"); + if (response.body._tag !== "Raw") throw new Error("Expected a raw response body."); + expect(response.body.body).toBeInstanceOf(NodeStream.Readable); + if (!(response.body.body instanceof NodeStream.Readable)) { + throw new Error("Expected a Node readable."); + } + + const chunks: Array = []; + for await (const chunk of response.body.body) { + chunks.push(Buffer.from(chunk)); + } + expect(NodeZlib.gunzipSync(Buffer.concat(chunks))).toEqual(Buffer.from(body)); +}); + +it("creates a raw Bun gzip stream", async () => { + const response = BunHttpResponseCompression.make.gzip(body, { + contentType: "application/json", + }); + + expect(response.body._tag).toBe("Raw"); + if (response.body._tag !== "Raw") throw new Error("Expected a raw response body."); + expect(response.body.body).toBeInstanceOf(ReadableStream); + if (!(response.body.body instanceof ReadableStream)) { + throw new Error("Expected a Web readable stream."); + } + + const compressed = await new Response(response.body.body).arrayBuffer(); + expect(NodeZlib.gunzipSync(compressed)).toEqual(Buffer.from(body)); +}); diff --git a/apps/server/src/HttpResponseCompression.ts b/apps/server/src/HttpResponseCompression.ts new file mode 100644 index 00000000000..dee9fdc62e1 --- /dev/null +++ b/apps/server/src/HttpResponseCompression.ts @@ -0,0 +1,12 @@ +import * as Context from "effect/Context"; +import type * as HttpServerResponse from "effect/unstable/http/HttpServerResponse"; + +export class HttpResponseCompression extends Context.Service< + HttpResponseCompression, + { + readonly gzip: ( + body: Uint8Array, + options: HttpServerResponse.Options, + ) => HttpServerResponse.HttpServerResponse; + } +>()("t3/HttpResponseCompression") {} diff --git a/apps/server/src/NodeHttpResponseCompression.ts b/apps/server/src/NodeHttpResponseCompression.ts new file mode 100644 index 00000000000..b281042788e --- /dev/null +++ b/apps/server/src/NodeHttpResponseCompression.ts @@ -0,0 +1,13 @@ +import * as NodeStream from "node:stream"; +import * as NodeZlib from "node:zlib"; +import * as Layer from "effect/Layer"; +import * as HttpServerResponse from "effect/unstable/http/HttpServerResponse"; + +import { HttpResponseCompression } from "./HttpResponseCompression.ts"; + +export const make = HttpResponseCompression.of({ + gzip: (body, options) => + HttpServerResponse.raw(NodeStream.Readable.from([body]).pipe(NodeZlib.createGzip()), options), +}); + +export const layer = Layer.succeed(HttpResponseCompression, make); diff --git a/apps/server/src/http.test.ts b/apps/server/src/http.test.ts index 125c1d509f1..d953bbb8140 100644 --- a/apps/server/src/http.test.ts +++ b/apps/server/src/http.test.ts @@ -1,10 +1,11 @@ import { expect, it } from "@effect/vitest"; +import * as NodeStream from "node:stream"; import * as NodeZlib from "node:zlib"; import * as Effect from "effect/Effect"; -import * as Stream from "effect/Stream"; import { HttpServerResponse } from "effect/unstable/http"; import { describe } from "vite-plus/test"; +import * as NodeHttpResponseCompression from "./NodeHttpResponseCompression.ts"; import { compressHttpResponse, isLoopbackHostname, resolveDevRedirectUrl } from "./http.ts"; describe("http dev routing", () => { @@ -38,16 +39,27 @@ describe("http compression", () => { const response = compressHttpResponse( HttpServerResponse.text(body, { contentType: "application/json" }), "br, gzip, deflate", + NodeHttpResponseCompression.make, ); expect(response.headers["content-encoding"]).toBe("gzip"); expect(response.headers["content-length"]).toBeUndefined(); expect(response.headers.vary).toBe("Accept-Encoding"); - expect(response.body._tag).toBe("Stream"); - if (response.body._tag === "Stream") { - const chunks = yield* Stream.runCollect(Stream.orDie(response.body.stream)); - expect(NodeZlib.gunzipSync(Buffer.concat(Array.from(chunks))).toString()).toBe(body); + expect(response.body._tag).toBe("Raw"); + if (response.body._tag !== "Raw") throw new Error("Expected a raw response body."); + expect(response.body.body).toBeInstanceOf(NodeStream.Readable); + if (!(response.body.body instanceof NodeStream.Readable)) { + throw new Error("Expected a Node readable."); } + const rawBody = response.body.body; + const chunks = yield* Effect.promise(async () => { + const chunks: Array = []; + for await (const chunk of rawBody) { + chunks.push(Buffer.from(chunk)); + } + return chunks; + }); + expect(NodeZlib.gunzipSync(Buffer.concat(chunks)).toString()).toBe(body); }), ); @@ -55,6 +67,7 @@ describe("http compression", () => { const response = compressHttpResponse( HttpServerResponse.text("x".repeat(2_000), { contentType: "application/json" }), "gzip;q=0, *;q=1", + NodeHttpResponseCompression.make, ); expect(response.headers["content-encoding"]).toBeUndefined(); @@ -70,6 +83,7 @@ describe("http compression", () => { headers: { vary }, }), undefined, + NodeHttpResponseCompression.make, ); expect(makeResponse("*").headers.vary).toBe("*"); diff --git a/apps/server/src/http.ts b/apps/server/src/http.ts index c3070f54055..470dd3b954c 100644 --- a/apps/server/src/http.ts +++ b/apps/server/src/http.ts @@ -12,7 +12,6 @@ import * as FileSystem from "effect/FileSystem"; import * as Layer from "effect/Layer"; import * as Option from "effect/Option"; import * as Path from "effect/Path"; -import * as Stream from "effect/Stream"; import { cast } from "effect/Function"; import { Headers, @@ -31,6 +30,7 @@ import * as ServerConfig from "./config.ts"; import { ASSET_ROUTE_PREFIX, resolveAsset } from "./assets/AssetAccess.ts"; import * as BrowserTraceCollector from "./observability/BrowserTraceCollector.ts"; import * as EnvironmentAuth from "./auth/EnvironmentAuth.ts"; +import * as HttpResponseCompression from "./HttpResponseCompression.ts"; import { traceRelayRequest } from "./cloud/traceRelayRequest.ts"; import { annotateEnvironmentRequest, @@ -70,6 +70,7 @@ function varyByAcceptEncoding(value: string | undefined): string { export function compressHttpResponse( response: HttpServerResponse.HttpServerResponse, acceptEncoding: string | undefined, + compression: HttpResponseCompression.HttpResponseCompression["Service"], ): HttpServerResponse.HttpServerResponse { const body = response.body; if ( @@ -88,16 +89,12 @@ export function compressHttpResponse( ); if (!acceptsGzip(acceptEncoding)) return variedResponse; - const compressedBody = Stream.fromReadableStream({ - evaluate: () => new Response(body.body).body!.pipeThrough(new CompressionStream("gzip")), - onError: (cause) => cause, - }); const headers = Headers.set( Headers.remove(variedResponse.headers, "content-length"), "content-encoding", "gzip", ); - return HttpServerResponse.stream(compressedBody, { + return compression.gzip(body.body, { status: response.status, statusText: response.statusText, headers, @@ -110,7 +107,12 @@ export const httpCompressionLayer = HttpRouter.middleware( (httpEffect) => Effect.gen(function* () { const request = yield* HttpServerRequest.HttpServerRequest; - return compressHttpResponse(yield* httpEffect, request.headers["accept-encoding"]); + const compression = yield* HttpResponseCompression.HttpResponseCompression; + return compressHttpResponse( + yield* httpEffect, + request.headers["accept-encoding"], + compression, + ); }), { global: true }, ); diff --git a/apps/server/src/server.test.ts b/apps/server/src/server.test.ts index 65eee28e59b..c1ec2aee27e 100644 --- a/apps/server/src/server.test.ts +++ b/apps/server/src/server.test.ts @@ -72,6 +72,7 @@ import { vi } from "vite-plus/test"; const TEST_EPOCH = DateTime.makeUnsafe("1970-01-01T00:00:00.000Z"); import * as ServerConfig from "./config.ts"; +import * as NodeHttpResponseCompression from "./NodeHttpResponseCompression.ts"; import { makeRoutesLayer } from "./server.ts"; import { resolveAvailableEditorsForConfig } from "./ws.ts"; import * as CheckpointDiffQuery from "./checkpointing/CheckpointDiffQuery.ts"; @@ -815,6 +816,7 @@ const buildAppUnderTest = (options?: { Layer.provideMerge(ServerSecretStore.layer), Layer.provide(workspaceAndProjectServicesLayer), Layer.provideMerge(FetchHttpClient.layer), + Layer.provide(NodeHttpResponseCompression.layer), Layer.provide(layerConfig), ); diff --git a/apps/server/src/server.ts b/apps/server/src/server.ts index d78555cb30b..1ee1fe82e55 100644 --- a/apps/server/src/server.ts +++ b/apps/server/src/server.ts @@ -125,28 +125,36 @@ const RelayClientLive = Layer.unwrap( }), ); -const HttpServerLive = Layer.unwrap( +const HttpServerPlatformLive = Layer.unwrap( Effect.gen(function* () { const config = yield* ServerConfig.ServerConfig; if (typeof Bun !== "undefined") { - const BunHttpServer = yield* Effect.promise( - () => import("@effect/platform-bun/BunHttpServer"), + const [BunHttpServer, BunHttpResponseCompression] = yield* Effect.all([ + Effect.promise(() => import("@effect/platform-bun/BunHttpServer")), + Effect.promise(() => import("./BunHttpResponseCompression.ts")), + ]); + return Layer.merge( + BunHttpServer.layer({ + port: config.port, + hostname: config.host ?? "127.0.0.1", + gracefulShutdownTimeout: HTTP_PREEMPTIVE_SHUTDOWN_GRACE_MS, + }), + BunHttpResponseCompression.layer, ); - return BunHttpServer.layer({ - port: config.port, - hostname: config.host ?? "127.0.0.1", - gracefulShutdownTimeout: HTTP_PREEMPTIVE_SHUTDOWN_GRACE_MS, - }); } else { - const [NodeHttpServer, NodeHttp] = yield* Effect.all([ + const [NodeHttpServer, NodeHttp, NodeHttpResponseCompression] = yield* Effect.all([ Effect.promise(() => import("@effect/platform-node/NodeHttpServer")), Effect.promise(() => import("node:http")), + Effect.promise(() => import("./NodeHttpResponseCompression.ts")), ]); - return NodeHttpServer.layer(NodeHttp.createServer, { - host: config.host ?? "127.0.0.1", - port: config.port, - gracefulShutdownTimeout: HTTP_PREEMPTIVE_SHUTDOWN_GRACE_MS, - }); + return Layer.merge( + NodeHttpServer.layer(NodeHttp.createServer, { + host: config.host ?? "127.0.0.1", + port: config.port, + gracefulShutdownTimeout: HTTP_PREEMPTIVE_SHUTDOWN_GRACE_MS, + }), + NodeHttpResponseCompression.layer, + ); } }), ); @@ -516,7 +524,7 @@ export const makeServerLayer = Layer.unwrap( return serverApplicationLayer.pipe( Layer.provideMerge(RuntimeServicesLive), Layer.provideMerge(serverRelayBrokerTracingLayer), - Layer.provideMerge(HttpServerLive), + Layer.provideMerge(HttpServerPlatformLive), Layer.provide(ObservabilityLive), Layer.provideMerge(FetchHttpClient.layer), Layer.provideMerge(VcsProcess.layer), From 135ed07c0369a94d07b9a0948c070b9d44cbf30c Mon Sep 17 00:00:00 2001 From: Julius Marminge Date: Wed, 29 Jul 2026 01:26:02 +0200 Subject: [PATCH 2/8] refactor(server): colocate HTTP compression layers Co-authored-by: codex --- apps/server/src/BunHttpResponseCompression.ts | 14 ----- .../src/HttpResponseCompression.test.ts | 43 --------------- apps/server/src/HttpResponseCompression.ts | 12 ----- .../server/src/NodeHttpResponseCompression.ts | 13 ----- apps/server/src/http.test.ts | 17 ++++-- apps/server/src/http.ts | 2 +- .../HttpResponseCompression.test.ts | 54 +++++++++++++++++++ .../HttpResponseCompression.ts | 42 +++++++++++++++ apps/server/src/server.test.ts | 4 +- apps/server/src/server.ts | 15 +++--- 10 files changed, 118 insertions(+), 98 deletions(-) delete mode 100644 apps/server/src/BunHttpResponseCompression.ts delete mode 100644 apps/server/src/HttpResponseCompression.test.ts delete mode 100644 apps/server/src/HttpResponseCompression.ts delete mode 100644 apps/server/src/NodeHttpResponseCompression.ts create mode 100644 apps/server/src/httpCompression/HttpResponseCompression.test.ts create mode 100644 apps/server/src/httpCompression/HttpResponseCompression.ts diff --git a/apps/server/src/BunHttpResponseCompression.ts b/apps/server/src/BunHttpResponseCompression.ts deleted file mode 100644 index a72f46895ae..00000000000 --- a/apps/server/src/BunHttpResponseCompression.ts +++ /dev/null @@ -1,14 +0,0 @@ -import * as Layer from "effect/Layer"; -import * as HttpServerResponse from "effect/unstable/http/HttpServerResponse"; - -import { HttpResponseCompression } from "./HttpResponseCompression.ts"; - -export const make = HttpResponseCompression.of({ - gzip: (body, options) => - HttpServerResponse.raw( - new Response(body).body!.pipeThrough(new CompressionStream("gzip")), - options, - ), -}); - -export const layer = Layer.succeed(HttpResponseCompression, make); diff --git a/apps/server/src/HttpResponseCompression.test.ts b/apps/server/src/HttpResponseCompression.test.ts deleted file mode 100644 index 0be73d34d7f..00000000000 --- a/apps/server/src/HttpResponseCompression.test.ts +++ /dev/null @@ -1,43 +0,0 @@ -import { expect, it } from "vite-plus/test"; -import * as NodeStream from "node:stream"; -import * as NodeZlib from "node:zlib"; - -import * as BunHttpResponseCompression from "./BunHttpResponseCompression.ts"; -import * as NodeHttpResponseCompression from "./NodeHttpResponseCompression.ts"; - -const body = new TextEncoder().encode(`{"value":"${"compressible".repeat(1_000)}"}`); - -it("creates a raw Node gzip stream", async () => { - const response = NodeHttpResponseCompression.make.gzip(body, { - contentType: "application/json", - }); - - expect(response.body._tag).toBe("Raw"); - if (response.body._tag !== "Raw") throw new Error("Expected a raw response body."); - expect(response.body.body).toBeInstanceOf(NodeStream.Readable); - if (!(response.body.body instanceof NodeStream.Readable)) { - throw new Error("Expected a Node readable."); - } - - const chunks: Array = []; - for await (const chunk of response.body.body) { - chunks.push(Buffer.from(chunk)); - } - expect(NodeZlib.gunzipSync(Buffer.concat(chunks))).toEqual(Buffer.from(body)); -}); - -it("creates a raw Bun gzip stream", async () => { - const response = BunHttpResponseCompression.make.gzip(body, { - contentType: "application/json", - }); - - expect(response.body._tag).toBe("Raw"); - if (response.body._tag !== "Raw") throw new Error("Expected a raw response body."); - expect(response.body.body).toBeInstanceOf(ReadableStream); - if (!(response.body.body instanceof ReadableStream)) { - throw new Error("Expected a Web readable stream."); - } - - const compressed = await new Response(response.body.body).arrayBuffer(); - expect(NodeZlib.gunzipSync(compressed)).toEqual(Buffer.from(body)); -}); diff --git a/apps/server/src/HttpResponseCompression.ts b/apps/server/src/HttpResponseCompression.ts deleted file mode 100644 index dee9fdc62e1..00000000000 --- a/apps/server/src/HttpResponseCompression.ts +++ /dev/null @@ -1,12 +0,0 @@ -import * as Context from "effect/Context"; -import type * as HttpServerResponse from "effect/unstable/http/HttpServerResponse"; - -export class HttpResponseCompression extends Context.Service< - HttpResponseCompression, - { - readonly gzip: ( - body: Uint8Array, - options: HttpServerResponse.Options, - ) => HttpServerResponse.HttpServerResponse; - } ->()("t3/HttpResponseCompression") {} diff --git a/apps/server/src/NodeHttpResponseCompression.ts b/apps/server/src/NodeHttpResponseCompression.ts deleted file mode 100644 index b281042788e..00000000000 --- a/apps/server/src/NodeHttpResponseCompression.ts +++ /dev/null @@ -1,13 +0,0 @@ -import * as NodeStream from "node:stream"; -import * as NodeZlib from "node:zlib"; -import * as Layer from "effect/Layer"; -import * as HttpServerResponse from "effect/unstable/http/HttpServerResponse"; - -import { HttpResponseCompression } from "./HttpResponseCompression.ts"; - -export const make = HttpResponseCompression.of({ - gzip: (body, options) => - HttpServerResponse.raw(NodeStream.Readable.from([body]).pipe(NodeZlib.createGzip()), options), -}); - -export const layer = Layer.succeed(HttpResponseCompression, make); diff --git a/apps/server/src/http.test.ts b/apps/server/src/http.test.ts index d953bbb8140..c18b9b96884 100644 --- a/apps/server/src/http.test.ts +++ b/apps/server/src/http.test.ts @@ -5,9 +5,15 @@ import * as Effect from "effect/Effect"; import { HttpServerResponse } from "effect/unstable/http"; import { describe } from "vite-plus/test"; -import * as NodeHttpResponseCompression from "./NodeHttpResponseCompression.ts"; +import * as HttpResponseCompression from "./httpCompression/HttpResponseCompression.ts"; import { compressHttpResponse, isLoopbackHostname, resolveDevRedirectUrl } from "./http.ts"; +const compressionUnavailable = HttpResponseCompression.HttpResponseCompression.of({ + gzip: () => { + throw new Error("Unexpected HTTP response compression."); + }, +}); + describe("http dev routing", () => { it("treats localhost and loopback addresses as local", () => { expect(isLoopbackHostname("127.0.0.1")).toBe(true); @@ -36,10 +42,11 @@ describe("http compression", () => { it.effect("gzips large JSON responses when the client accepts it", () => Effect.gen(function* () { const body = `{"value":"${"compressible".repeat(1_000)}"}`; + const compression = yield* HttpResponseCompression.HttpResponseCompression; const response = compressHttpResponse( HttpServerResponse.text(body, { contentType: "application/json" }), "br, gzip, deflate", - NodeHttpResponseCompression.make, + compression, ); expect(response.headers["content-encoding"]).toBe("gzip"); @@ -60,14 +67,14 @@ describe("http compression", () => { return chunks; }); expect(NodeZlib.gunzipSync(Buffer.concat(chunks)).toString()).toBe(body); - }), + }).pipe(Effect.provide(HttpResponseCompression.layerNode)), ); it("keeps the original body when gzip is declined", () => { const response = compressHttpResponse( HttpServerResponse.text("x".repeat(2_000), { contentType: "application/json" }), "gzip;q=0, *;q=1", - NodeHttpResponseCompression.make, + compressionUnavailable, ); expect(response.headers["content-encoding"]).toBeUndefined(); @@ -83,7 +90,7 @@ describe("http compression", () => { headers: { vary }, }), undefined, - NodeHttpResponseCompression.make, + compressionUnavailable, ); expect(makeResponse("*").headers.vary).toBe("*"); diff --git a/apps/server/src/http.ts b/apps/server/src/http.ts index 470dd3b954c..30d94017d78 100644 --- a/apps/server/src/http.ts +++ b/apps/server/src/http.ts @@ -30,7 +30,7 @@ import * as ServerConfig from "./config.ts"; import { ASSET_ROUTE_PREFIX, resolveAsset } from "./assets/AssetAccess.ts"; import * as BrowserTraceCollector from "./observability/BrowserTraceCollector.ts"; import * as EnvironmentAuth from "./auth/EnvironmentAuth.ts"; -import * as HttpResponseCompression from "./HttpResponseCompression.ts"; +import * as HttpResponseCompression from "./httpCompression/HttpResponseCompression.ts"; import { traceRelayRequest } from "./cloud/traceRelayRequest.ts"; import { annotateEnvironmentRequest, diff --git a/apps/server/src/httpCompression/HttpResponseCompression.test.ts b/apps/server/src/httpCompression/HttpResponseCompression.test.ts new file mode 100644 index 00000000000..6051f58a947 --- /dev/null +++ b/apps/server/src/httpCompression/HttpResponseCompression.test.ts @@ -0,0 +1,54 @@ +import { expect, it } from "@effect/vitest"; +import * as NodeStream from "node:stream"; +import * as NodeZlib from "node:zlib"; +import * as Effect from "effect/Effect"; + +import * as HttpResponseCompression from "./HttpResponseCompression.ts"; + +const body = new TextEncoder().encode(`{"value":"${"compressible".repeat(1_000)}"}`); + +it.effect("creates a raw Node gzip stream", () => + Effect.gen(function* () { + const compression = yield* HttpResponseCompression.HttpResponseCompression; + const response = compression.gzip(body, { + contentType: "application/json", + }); + + expect(response.body._tag).toBe("Raw"); + if (response.body._tag !== "Raw") throw new Error("Expected a raw response body."); + expect(response.body.body).toBeInstanceOf(NodeStream.Readable); + if (!(response.body.body instanceof NodeStream.Readable)) { + throw new Error("Expected a Node readable."); + } + + const rawBody = response.body.body; + const chunks = yield* Effect.promise(async () => { + const chunks: Array = []; + for await (const chunk of rawBody) { + chunks.push(Buffer.from(chunk)); + } + return chunks; + }); + expect(NodeZlib.gunzipSync(Buffer.concat(chunks))).toEqual(Buffer.from(body)); + }).pipe(Effect.provide(HttpResponseCompression.layerNode)), +); + +it.effect("creates a raw Bun gzip stream", () => + Effect.gen(function* () { + const compression = yield* HttpResponseCompression.HttpResponseCompression; + const response = compression.gzip(body, { + contentType: "application/json", + }); + + expect(response.body._tag).toBe("Raw"); + if (response.body._tag !== "Raw") throw new Error("Expected a raw response body."); + expect(response.body.body).toBeInstanceOf(ReadableStream); + if (!(response.body.body instanceof ReadableStream)) { + throw new Error("Expected a Web readable stream."); + } + + const rawBody = response.body.body; + const compressed = yield* Effect.promise(() => new Response(rawBody).arrayBuffer()); + expect(NodeZlib.gunzipSync(compressed)).toEqual(Buffer.from(body)); + }).pipe(Effect.provide(HttpResponseCompression.layerBun)), +); diff --git a/apps/server/src/httpCompression/HttpResponseCompression.ts b/apps/server/src/httpCompression/HttpResponseCompression.ts new file mode 100644 index 00000000000..bc040e2c774 --- /dev/null +++ b/apps/server/src/httpCompression/HttpResponseCompression.ts @@ -0,0 +1,42 @@ +import * as Context from "effect/Context"; +import * as Effect from "effect/Effect"; +import * as Layer from "effect/Layer"; +import * as HttpServerResponse from "effect/unstable/http/HttpServerResponse"; + +export class HttpResponseCompression extends Context.Service< + HttpResponseCompression, + { + readonly gzip: ( + body: Uint8Array, + options: HttpServerResponse.Options, + ) => HttpServerResponse.HttpServerResponse; + } +>()("t3/httpCompression/HttpResponseCompression") {} + +export const layerNode = Layer.effect( + HttpResponseCompression, + Effect.gen(function* () { + const [NodeStream, NodeZlib] = yield* Effect.all([ + Effect.promise(() => import("node:stream")), + Effect.promise(() => import("node:zlib")), + ]); + return HttpResponseCompression.of({ + gzip: (body, options) => + HttpServerResponse.raw( + NodeStream.Readable.from([body]).pipe(NodeZlib.createGzip()), + options, + ), + }); + }), +); + +export const layerBun = Layer.succeed( + HttpResponseCompression, + HttpResponseCompression.of({ + gzip: (body, options) => + HttpServerResponse.raw( + new Response(body).body!.pipeThrough(new CompressionStream("gzip")), + options, + ), + }), +); diff --git a/apps/server/src/server.test.ts b/apps/server/src/server.test.ts index c1ec2aee27e..cb68c0bba8a 100644 --- a/apps/server/src/server.test.ts +++ b/apps/server/src/server.test.ts @@ -72,7 +72,7 @@ import { vi } from "vite-plus/test"; const TEST_EPOCH = DateTime.makeUnsafe("1970-01-01T00:00:00.000Z"); import * as ServerConfig from "./config.ts"; -import * as NodeHttpResponseCompression from "./NodeHttpResponseCompression.ts"; +import * as HttpResponseCompression from "./httpCompression/HttpResponseCompression.ts"; import { makeRoutesLayer } from "./server.ts"; import { resolveAvailableEditorsForConfig } from "./ws.ts"; import * as CheckpointDiffQuery from "./checkpointing/CheckpointDiffQuery.ts"; @@ -816,7 +816,7 @@ const buildAppUnderTest = (options?: { Layer.provideMerge(ServerSecretStore.layer), Layer.provide(workspaceAndProjectServicesLayer), Layer.provideMerge(FetchHttpClient.layer), - Layer.provide(NodeHttpResponseCompression.layer), + Layer.provide(HttpResponseCompression.layerNode), Layer.provide(layerConfig), ); diff --git a/apps/server/src/server.ts b/apps/server/src/server.ts index 1ee1fe82e55..6b8a7cc01f2 100644 --- a/apps/server/src/server.ts +++ b/apps/server/src/server.ts @@ -5,6 +5,7 @@ import { FetchHttpClient, HttpRouter, HttpServer } from "effect/unstable/http"; import * as HttpApiBuilder from "effect/unstable/httpapi/HttpApiBuilder"; import * as ServerConfig from "./config.ts"; +import * as HttpResponseCompression from "./httpCompression/HttpResponseCompression.ts"; import { otlpTracesProxyRouteLayer, assetRouteLayer, @@ -129,23 +130,21 @@ const HttpServerPlatformLive = Layer.unwrap( Effect.gen(function* () { const config = yield* ServerConfig.ServerConfig; if (typeof Bun !== "undefined") { - const [BunHttpServer, BunHttpResponseCompression] = yield* Effect.all([ - Effect.promise(() => import("@effect/platform-bun/BunHttpServer")), - Effect.promise(() => import("./BunHttpResponseCompression.ts")), - ]); + const BunHttpServer = yield* Effect.promise( + () => import("@effect/platform-bun/BunHttpServer"), + ); return Layer.merge( BunHttpServer.layer({ port: config.port, hostname: config.host ?? "127.0.0.1", gracefulShutdownTimeout: HTTP_PREEMPTIVE_SHUTDOWN_GRACE_MS, }), - BunHttpResponseCompression.layer, + HttpResponseCompression.layerBun, ); } else { - const [NodeHttpServer, NodeHttp, NodeHttpResponseCompression] = yield* Effect.all([ + const [NodeHttpServer, NodeHttp] = yield* Effect.all([ Effect.promise(() => import("@effect/platform-node/NodeHttpServer")), Effect.promise(() => import("node:http")), - Effect.promise(() => import("./NodeHttpResponseCompression.ts")), ]); return Layer.merge( NodeHttpServer.layer(NodeHttp.createServer, { @@ -153,7 +152,7 @@ const HttpServerPlatformLive = Layer.unwrap( port: config.port, gracefulShutdownTimeout: HTTP_PREEMPTIVE_SHUTDOWN_GRACE_MS, }), - NodeHttpResponseCompression.layer, + HttpResponseCompression.layerNode, ); } }), From f37416de353befa748f16b293cc42219876f4b23 Mon Sep 17 00:00:00 2001 From: Julius Marminge Date: Wed, 29 Jul 2026 01:31:20 +0200 Subject: [PATCH 3/8] refactor(server): resolve compression service from context Co-authored-by: codex --- apps/server/src/http.test.ts | 64 ++++++------- apps/server/src/http.ts | 17 ++-- .../HttpResponseCompression.test.ts | 92 ++++++++++--------- 3 files changed, 84 insertions(+), 89 deletions(-) diff --git a/apps/server/src/http.test.ts b/apps/server/src/http.test.ts index c18b9b96884..4247bfeb07d 100644 --- a/apps/server/src/http.test.ts +++ b/apps/server/src/http.test.ts @@ -8,12 +8,6 @@ import { describe } from "vite-plus/test"; import * as HttpResponseCompression from "./httpCompression/HttpResponseCompression.ts"; import { compressHttpResponse, isLoopbackHostname, resolveDevRedirectUrl } from "./http.ts"; -const compressionUnavailable = HttpResponseCompression.HttpResponseCompression.of({ - gzip: () => { - throw new Error("Unexpected HTTP response compression."); - }, -}); - describe("http dev routing", () => { it("treats localhost and loopback addresses as local", () => { expect(isLoopbackHostname("127.0.0.1")).toBe(true); @@ -38,15 +32,13 @@ describe("http dev routing", () => { }); }); -describe("http compression", () => { +it.layer(HttpResponseCompression.layerNode)("http compression", (it) => { it.effect("gzips large JSON responses when the client accepts it", () => Effect.gen(function* () { const body = `{"value":"${"compressible".repeat(1_000)}"}`; - const compression = yield* HttpResponseCompression.HttpResponseCompression; - const response = compressHttpResponse( + const response = yield* compressHttpResponse( HttpServerResponse.text(body, { contentType: "application/json" }), "br, gzip, deflate", - compression, ); expect(response.headers["content-encoding"]).toBe("gzip"); @@ -67,33 +59,37 @@ describe("http compression", () => { return chunks; }); expect(NodeZlib.gunzipSync(Buffer.concat(chunks)).toString()).toBe(body); - }).pipe(Effect.provide(HttpResponseCompression.layerNode)), + }), ); - it("keeps the original body when gzip is declined", () => { - const response = compressHttpResponse( - HttpServerResponse.text("x".repeat(2_000), { contentType: "application/json" }), - "gzip;q=0, *;q=1", - compressionUnavailable, - ); + it.effect("keeps the original body when gzip is declined", () => + Effect.gen(function* () { + const response = yield* compressHttpResponse( + HttpServerResponse.text("x".repeat(2_000), { contentType: "application/json" }), + "gzip;q=0, *;q=1", + ); - expect(response.headers["content-encoding"]).toBeUndefined(); - expect(response.headers["content-length"]).toBe("2000"); - expect(response.headers.vary).toBe("Accept-Encoding"); - }); + expect(response.headers["content-encoding"]).toBeUndefined(); + expect(response.headers["content-length"]).toBe("2000"); + expect(response.headers.vary).toBe("Accept-Encoding"); + }), + ); - it("preserves existing Vary semantics", () => { - const makeResponse = (vary: string) => - compressHttpResponse( - HttpServerResponse.text("x".repeat(2_000), { - contentType: "application/json", - headers: { vary }, - }), - undefined, - compressionUnavailable, - ); + it.effect("preserves existing Vary semantics", () => + Effect.gen(function* () { + const makeResponse = (vary: string) => + compressHttpResponse( + HttpServerResponse.text("x".repeat(2_000), { + contentType: "application/json", + headers: { vary }, + }), + undefined, + ); - expect(makeResponse("*").headers.vary).toBe("*"); - expect(makeResponse("Origin, accept-encoding").headers.vary).toBe("Origin, accept-encoding"); - }); + expect((yield* makeResponse("*")).headers.vary).toBe("*"); + expect((yield* makeResponse("Origin, accept-encoding")).headers.vary).toBe( + "Origin, accept-encoding", + ); + }), + ); }); diff --git a/apps/server/src/http.ts b/apps/server/src/http.ts index 30d94017d78..0e601374716 100644 --- a/apps/server/src/http.ts +++ b/apps/server/src/http.ts @@ -30,7 +30,7 @@ import * as ServerConfig from "./config.ts"; import { ASSET_ROUTE_PREFIX, resolveAsset } from "./assets/AssetAccess.ts"; import * as BrowserTraceCollector from "./observability/BrowserTraceCollector.ts"; import * as EnvironmentAuth from "./auth/EnvironmentAuth.ts"; -import * as HttpResponseCompression from "./httpCompression/HttpResponseCompression.ts"; +import { HttpResponseCompression } from "./httpCompression/HttpResponseCompression.ts"; import { traceRelayRequest } from "./cloud/traceRelayRequest.ts"; import { annotateEnvironmentRequest, @@ -67,11 +67,10 @@ function varyByAcceptEncoding(value: string | undefined): string { return values.has("*") || values.has("accept-encoding") ? value : `${value}, Accept-Encoding`; } -export function compressHttpResponse( +export const compressHttpResponse = Effect.fnUntraced(function* ( response: HttpServerResponse.HttpServerResponse, acceptEncoding: string | undefined, - compression: HttpResponseCompression.HttpResponseCompression["Service"], -): HttpServerResponse.HttpServerResponse { +) { const body = response.body; if ( body._tag !== "Uint8Array" || @@ -89,6 +88,7 @@ export function compressHttpResponse( ); if (!acceptsGzip(acceptEncoding)) return variedResponse; + const compression = yield* HttpResponseCompression; const headers = Headers.set( Headers.remove(variedResponse.headers, "content-length"), "content-encoding", @@ -101,18 +101,13 @@ export function compressHttpResponse( cookies: response.cookies, contentType: body.contentType, }); -} +}); export const httpCompressionLayer = HttpRouter.middleware( (httpEffect) => Effect.gen(function* () { const request = yield* HttpServerRequest.HttpServerRequest; - const compression = yield* HttpResponseCompression.HttpResponseCompression; - return compressHttpResponse( - yield* httpEffect, - request.headers["accept-encoding"], - compression, - ); + return yield* compressHttpResponse(yield* httpEffect, request.headers["accept-encoding"]); }), { global: true }, ); diff --git a/apps/server/src/httpCompression/HttpResponseCompression.test.ts b/apps/server/src/httpCompression/HttpResponseCompression.test.ts index 6051f58a947..4815ed5e755 100644 --- a/apps/server/src/httpCompression/HttpResponseCompression.test.ts +++ b/apps/server/src/httpCompression/HttpResponseCompression.test.ts @@ -7,48 +7,52 @@ import * as HttpResponseCompression from "./HttpResponseCompression.ts"; const body = new TextEncoder().encode(`{"value":"${"compressible".repeat(1_000)}"}`); -it.effect("creates a raw Node gzip stream", () => - Effect.gen(function* () { - const compression = yield* HttpResponseCompression.HttpResponseCompression; - const response = compression.gzip(body, { - contentType: "application/json", - }); - - expect(response.body._tag).toBe("Raw"); - if (response.body._tag !== "Raw") throw new Error("Expected a raw response body."); - expect(response.body.body).toBeInstanceOf(NodeStream.Readable); - if (!(response.body.body instanceof NodeStream.Readable)) { - throw new Error("Expected a Node readable."); - } - - const rawBody = response.body.body; - const chunks = yield* Effect.promise(async () => { - const chunks: Array = []; - for await (const chunk of rawBody) { - chunks.push(Buffer.from(chunk)); +it.layer(HttpResponseCompression.layerNode)("Node HTTP response compression", (it) => { + it.effect("creates a raw gzip stream", () => + Effect.gen(function* () { + const compression = yield* HttpResponseCompression.HttpResponseCompression; + const response = compression.gzip(body, { + contentType: "application/json", + }); + + expect(response.body._tag).toBe("Raw"); + if (response.body._tag !== "Raw") throw new Error("Expected a raw response body."); + expect(response.body.body).toBeInstanceOf(NodeStream.Readable); + if (!(response.body.body instanceof NodeStream.Readable)) { + throw new Error("Expected a Node readable."); } - return chunks; - }); - expect(NodeZlib.gunzipSync(Buffer.concat(chunks))).toEqual(Buffer.from(body)); - }).pipe(Effect.provide(HttpResponseCompression.layerNode)), -); - -it.effect("creates a raw Bun gzip stream", () => - Effect.gen(function* () { - const compression = yield* HttpResponseCompression.HttpResponseCompression; - const response = compression.gzip(body, { - contentType: "application/json", - }); - - expect(response.body._tag).toBe("Raw"); - if (response.body._tag !== "Raw") throw new Error("Expected a raw response body."); - expect(response.body.body).toBeInstanceOf(ReadableStream); - if (!(response.body.body instanceof ReadableStream)) { - throw new Error("Expected a Web readable stream."); - } - - const rawBody = response.body.body; - const compressed = yield* Effect.promise(() => new Response(rawBody).arrayBuffer()); - expect(NodeZlib.gunzipSync(compressed)).toEqual(Buffer.from(body)); - }).pipe(Effect.provide(HttpResponseCompression.layerBun)), -); + + const rawBody = response.body.body; + const chunks = yield* Effect.promise(async () => { + const chunks: Array = []; + for await (const chunk of rawBody) { + chunks.push(Buffer.from(chunk)); + } + return chunks; + }); + expect(NodeZlib.gunzipSync(Buffer.concat(chunks))).toEqual(Buffer.from(body)); + }), + ); +}); + +it.layer(HttpResponseCompression.layerBun)("Bun HTTP response compression", (it) => { + it.effect("creates a raw gzip stream", () => + Effect.gen(function* () { + const compression = yield* HttpResponseCompression.HttpResponseCompression; + const response = compression.gzip(body, { + contentType: "application/json", + }); + + expect(response.body._tag).toBe("Raw"); + if (response.body._tag !== "Raw") throw new Error("Expected a raw response body."); + expect(response.body.body).toBeInstanceOf(ReadableStream); + if (!(response.body.body instanceof ReadableStream)) { + throw new Error("Expected a Web readable stream."); + } + + const rawBody = response.body.body; + const compressed = yield* Effect.promise(() => new Response(rawBody).arrayBuffer()); + expect(NodeZlib.gunzipSync(compressed)).toEqual(Buffer.from(body)); + }), + ); +}); From 9b35e3e0a1e8a2a5299bd9f88b677f71765c60f7 Mon Sep 17 00:00:00 2001 From: Julius Marminge Date: Wed, 29 Jul 2026 01:34:30 +0200 Subject: [PATCH 4/8] refactor(server): separate compression platform wiring Co-authored-by: codex --- apps/server/src/http.test.ts | 17 ----------------- apps/server/src/server.ts | 34 ++++++++++++++++------------------ 2 files changed, 16 insertions(+), 35 deletions(-) diff --git a/apps/server/src/http.test.ts b/apps/server/src/http.test.ts index 4247bfeb07d..299d517ca69 100644 --- a/apps/server/src/http.test.ts +++ b/apps/server/src/http.test.ts @@ -1,6 +1,4 @@ import { expect, it } from "@effect/vitest"; -import * as NodeStream from "node:stream"; -import * as NodeZlib from "node:zlib"; import * as Effect from "effect/Effect"; import { HttpServerResponse } from "effect/unstable/http"; import { describe } from "vite-plus/test"; @@ -44,21 +42,6 @@ it.layer(HttpResponseCompression.layerNode)("http compression", (it) => { expect(response.headers["content-encoding"]).toBe("gzip"); expect(response.headers["content-length"]).toBeUndefined(); expect(response.headers.vary).toBe("Accept-Encoding"); - expect(response.body._tag).toBe("Raw"); - if (response.body._tag !== "Raw") throw new Error("Expected a raw response body."); - expect(response.body.body).toBeInstanceOf(NodeStream.Readable); - if (!(response.body.body instanceof NodeStream.Readable)) { - throw new Error("Expected a Node readable."); - } - const rawBody = response.body.body; - const chunks = yield* Effect.promise(async () => { - const chunks: Array = []; - for await (const chunk of rawBody) { - chunks.push(Buffer.from(chunk)); - } - return chunks; - }); - expect(NodeZlib.gunzipSync(Buffer.concat(chunks)).toString()).toBe(body); }), ); diff --git a/apps/server/src/server.ts b/apps/server/src/server.ts index 6b8a7cc01f2..0b7c8ccd074 100644 --- a/apps/server/src/server.ts +++ b/apps/server/src/server.ts @@ -126,38 +126,35 @@ const RelayClientLive = Layer.unwrap( }), ); -const HttpServerPlatformLive = Layer.unwrap( +const HttpServerLive = Layer.unwrap( Effect.gen(function* () { const config = yield* ServerConfig.ServerConfig; if (typeof Bun !== "undefined") { const BunHttpServer = yield* Effect.promise( () => import("@effect/platform-bun/BunHttpServer"), ); - return Layer.merge( - BunHttpServer.layer({ - port: config.port, - hostname: config.host ?? "127.0.0.1", - gracefulShutdownTimeout: HTTP_PREEMPTIVE_SHUTDOWN_GRACE_MS, - }), - HttpResponseCompression.layerBun, - ); + return BunHttpServer.layer({ + port: config.port, + hostname: config.host ?? "127.0.0.1", + gracefulShutdownTimeout: HTTP_PREEMPTIVE_SHUTDOWN_GRACE_MS, + }); } else { const [NodeHttpServer, NodeHttp] = yield* Effect.all([ Effect.promise(() => import("@effect/platform-node/NodeHttpServer")), Effect.promise(() => import("node:http")), ]); - return Layer.merge( - NodeHttpServer.layer(NodeHttp.createServer, { - host: config.host ?? "127.0.0.1", - port: config.port, - gracefulShutdownTimeout: HTTP_PREEMPTIVE_SHUTDOWN_GRACE_MS, - }), - HttpResponseCompression.layerNode, - ); + return NodeHttpServer.layer(NodeHttp.createServer, { + host: config.host ?? "127.0.0.1", + port: config.port, + gracefulShutdownTimeout: HTTP_PREEMPTIVE_SHUTDOWN_GRACE_MS, + }); } }), ); +const HttpResponseCompressionLive = + typeof Bun !== "undefined" ? HttpResponseCompression.layerBun : HttpResponseCompression.layerNode; + const PlatformServicesLive = Layer.unwrap( Effect.gen(function* () { if (typeof Bun !== "undefined") { @@ -523,7 +520,8 @@ export const makeServerLayer = Layer.unwrap( return serverApplicationLayer.pipe( Layer.provideMerge(RuntimeServicesLive), Layer.provideMerge(serverRelayBrokerTracingLayer), - Layer.provideMerge(HttpServerPlatformLive), + Layer.provideMerge(HttpResponseCompressionLive), + Layer.provideMerge(HttpServerLive), Layer.provide(ObservabilityLive), Layer.provideMerge(FetchHttpClient.layer), Layer.provideMerge(VcsProcess.layer), From 187246fabeb3a08fb09fb2b70e70afa52a2f5fb3 Mon Sep 17 00:00:00 2001 From: Julius Marminge Date: Wed, 29 Jul 2026 01:36:42 +0200 Subject: [PATCH 5/8] test(server): remove duplicate compression coverage Co-authored-by: codex --- apps/server/src/http.test.ts | 14 -------------- 1 file changed, 14 deletions(-) diff --git a/apps/server/src/http.test.ts b/apps/server/src/http.test.ts index 299d517ca69..457633734a3 100644 --- a/apps/server/src/http.test.ts +++ b/apps/server/src/http.test.ts @@ -31,20 +31,6 @@ describe("http dev routing", () => { }); it.layer(HttpResponseCompression.layerNode)("http compression", (it) => { - it.effect("gzips large JSON responses when the client accepts it", () => - Effect.gen(function* () { - const body = `{"value":"${"compressible".repeat(1_000)}"}`; - const response = yield* compressHttpResponse( - HttpServerResponse.text(body, { contentType: "application/json" }), - "br, gzip, deflate", - ); - - expect(response.headers["content-encoding"]).toBe("gzip"); - expect(response.headers["content-length"]).toBeUndefined(); - expect(response.headers.vary).toBe("Accept-Encoding"); - }), - ); - it.effect("keeps the original body when gzip is declined", () => Effect.gen(function* () { const response = yield* compressHttpResponse( From 8b4a72f7893dbac493fbee4825536d13970ff009 Mon Sep 17 00:00:00 2001 From: Julius Marminge Date: Wed, 29 Jul 2026 01:40:00 +0200 Subject: [PATCH 6/8] refactor(server): trim compression scaffolding Co-authored-by: codex --- apps/server/src/http.test.ts | 38 +---------- apps/server/src/http.ts | 2 +- .../HttpResponseCompression.test.ts | 65 +++++-------------- .../HttpResponseCompression.ts | 36 +++++----- 4 files changed, 35 insertions(+), 106 deletions(-) diff --git a/apps/server/src/http.test.ts b/apps/server/src/http.test.ts index 457633734a3..4af2ecb6457 100644 --- a/apps/server/src/http.test.ts +++ b/apps/server/src/http.test.ts @@ -1,10 +1,7 @@ import { expect, it } from "@effect/vitest"; -import * as Effect from "effect/Effect"; -import { HttpServerResponse } from "effect/unstable/http"; import { describe } from "vite-plus/test"; -import * as HttpResponseCompression from "./httpCompression/HttpResponseCompression.ts"; -import { compressHttpResponse, isLoopbackHostname, resolveDevRedirectUrl } from "./http.ts"; +import { isLoopbackHostname, resolveDevRedirectUrl } from "./http.ts"; describe("http dev routing", () => { it("treats localhost and loopback addresses as local", () => { @@ -29,36 +26,3 @@ describe("http dev routing", () => { ); }); }); - -it.layer(HttpResponseCompression.layerNode)("http compression", (it) => { - it.effect("keeps the original body when gzip is declined", () => - Effect.gen(function* () { - const response = yield* compressHttpResponse( - HttpServerResponse.text("x".repeat(2_000), { contentType: "application/json" }), - "gzip;q=0, *;q=1", - ); - - expect(response.headers["content-encoding"]).toBeUndefined(); - expect(response.headers["content-length"]).toBe("2000"); - expect(response.headers.vary).toBe("Accept-Encoding"); - }), - ); - - it.effect("preserves existing Vary semantics", () => - Effect.gen(function* () { - const makeResponse = (vary: string) => - compressHttpResponse( - HttpServerResponse.text("x".repeat(2_000), { - contentType: "application/json", - headers: { vary }, - }), - undefined, - ); - - expect((yield* makeResponse("*")).headers.vary).toBe("*"); - expect((yield* makeResponse("Origin, accept-encoding")).headers.vary).toBe( - "Origin, accept-encoding", - ); - }), - ); -}); diff --git a/apps/server/src/http.ts b/apps/server/src/http.ts index 0e601374716..bd1845a185d 100644 --- a/apps/server/src/http.ts +++ b/apps/server/src/http.ts @@ -67,7 +67,7 @@ function varyByAcceptEncoding(value: string | undefined): string { return values.has("*") || values.has("accept-encoding") ? value : `${value}, Accept-Encoding`; } -export const compressHttpResponse = Effect.fnUntraced(function* ( +const compressHttpResponse = Effect.fnUntraced(function* ( response: HttpServerResponse.HttpServerResponse, acceptEncoding: string | undefined, ) { diff --git a/apps/server/src/httpCompression/HttpResponseCompression.test.ts b/apps/server/src/httpCompression/HttpResponseCompression.test.ts index 4815ed5e755..b608720d8dd 100644 --- a/apps/server/src/httpCompression/HttpResponseCompression.test.ts +++ b/apps/server/src/httpCompression/HttpResponseCompression.test.ts @@ -1,58 +1,29 @@ import { expect, it } from "@effect/vitest"; import * as NodeStream from "node:stream"; -import * as NodeZlib from "node:zlib"; import * as Effect from "effect/Effect"; import * as HttpResponseCompression from "./HttpResponseCompression.ts"; -const body = new TextEncoder().encode(`{"value":"${"compressible".repeat(1_000)}"}`); +it.effect("uses a native Node stream", () => + Effect.gen(function* () { + const compression = yield* HttpResponseCompression.HttpResponseCompression; + const response = compression.gzip(new Uint8Array(), {}); -it.layer(HttpResponseCompression.layerNode)("Node HTTP response compression", (it) => { - it.effect("creates a raw gzip stream", () => - Effect.gen(function* () { - const compression = yield* HttpResponseCompression.HttpResponseCompression; - const response = compression.gzip(body, { - contentType: "application/json", - }); - - expect(response.body._tag).toBe("Raw"); - if (response.body._tag !== "Raw") throw new Error("Expected a raw response body."); + expect(response.body._tag).toBe("Raw"); + if (response.body._tag === "Raw") { expect(response.body.body).toBeInstanceOf(NodeStream.Readable); - if (!(response.body.body instanceof NodeStream.Readable)) { - throw new Error("Expected a Node readable."); - } - - const rawBody = response.body.body; - const chunks = yield* Effect.promise(async () => { - const chunks: Array = []; - for await (const chunk of rawBody) { - chunks.push(Buffer.from(chunk)); - } - return chunks; - }); - expect(NodeZlib.gunzipSync(Buffer.concat(chunks))).toEqual(Buffer.from(body)); - }), - ); -}); + } + }).pipe(Effect.provide(HttpResponseCompression.layerNode)), +); -it.layer(HttpResponseCompression.layerBun)("Bun HTTP response compression", (it) => { - it.effect("creates a raw gzip stream", () => - Effect.gen(function* () { - const compression = yield* HttpResponseCompression.HttpResponseCompression; - const response = compression.gzip(body, { - contentType: "application/json", - }); +it.effect("uses a native Web stream for Bun", () => + Effect.gen(function* () { + const compression = yield* HttpResponseCompression.HttpResponseCompression; + const response = compression.gzip(new Uint8Array(), {}); - expect(response.body._tag).toBe("Raw"); - if (response.body._tag !== "Raw") throw new Error("Expected a raw response body."); + expect(response.body._tag).toBe("Raw"); + if (response.body._tag === "Raw") { expect(response.body.body).toBeInstanceOf(ReadableStream); - if (!(response.body.body instanceof ReadableStream)) { - throw new Error("Expected a Web readable stream."); - } - - const rawBody = response.body.body; - const compressed = yield* Effect.promise(() => new Response(rawBody).arrayBuffer()); - expect(NodeZlib.gunzipSync(compressed)).toEqual(Buffer.from(body)); - }), - ); -}); + } + }).pipe(Effect.provide(HttpResponseCompression.layerBun)), +); diff --git a/apps/server/src/httpCompression/HttpResponseCompression.ts b/apps/server/src/httpCompression/HttpResponseCompression.ts index bc040e2c774..f00c4e1a123 100644 --- a/apps/server/src/httpCompression/HttpResponseCompression.ts +++ b/apps/server/src/httpCompression/HttpResponseCompression.ts @@ -16,27 +16,21 @@ export class HttpResponseCompression extends Context.Service< export const layerNode = Layer.effect( HttpResponseCompression, Effect.gen(function* () { - const [NodeStream, NodeZlib] = yield* Effect.all([ - Effect.promise(() => import("node:stream")), - Effect.promise(() => import("node:zlib")), - ]); - return HttpResponseCompression.of({ - gzip: (body, options) => - HttpServerResponse.raw( - NodeStream.Readable.from([body]).pipe(NodeZlib.createGzip()), - options, - ), - }); + const NodeZlib = yield* Effect.promise(() => import("node:zlib")); + return { + gzip: (body, options) => { + const stream = NodeZlib.createGzip(); + stream.end(body); + return HttpServerResponse.raw(stream, options); + }, + }; }), ); -export const layerBun = Layer.succeed( - HttpResponseCompression, - HttpResponseCompression.of({ - gzip: (body, options) => - HttpServerResponse.raw( - new Response(body).body!.pipeThrough(new CompressionStream("gzip")), - options, - ), - }), -); +export const layerBun = Layer.succeed(HttpResponseCompression, { + gzip: (body, options) => + HttpServerResponse.raw( + new Response(body).body!.pipeThrough(new CompressionStream("gzip")), + options, + ), +}); From 50d3c70dfb19aed70b5c3d8a6e7c2644a2968e80 Mon Sep 17 00:00:00 2001 From: Julius Marminge Date: Wed, 29 Jul 2026 01:54:52 +0200 Subject: [PATCH 7/8] refactor(server): follow service import convention Co-authored-by: codex --- apps/server/src/http.ts | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/apps/server/src/http.ts b/apps/server/src/http.ts index bd1845a185d..7c002036989 100644 --- a/apps/server/src/http.ts +++ b/apps/server/src/http.ts @@ -30,7 +30,7 @@ import * as ServerConfig from "./config.ts"; import { ASSET_ROUTE_PREFIX, resolveAsset } from "./assets/AssetAccess.ts"; import * as BrowserTraceCollector from "./observability/BrowserTraceCollector.ts"; import * as EnvironmentAuth from "./auth/EnvironmentAuth.ts"; -import { HttpResponseCompression } from "./httpCompression/HttpResponseCompression.ts"; +import * as HttpResponseCompression from "./httpCompression/HttpResponseCompression.ts"; import { traceRelayRequest } from "./cloud/traceRelayRequest.ts"; import { annotateEnvironmentRequest, @@ -88,7 +88,7 @@ const compressHttpResponse = Effect.fnUntraced(function* ( ); if (!acceptsGzip(acceptEncoding)) return variedResponse; - const compression = yield* HttpResponseCompression; + const compression = yield* HttpResponseCompression.HttpResponseCompression; const headers = Headers.set( Headers.remove(variedResponse.headers, "content-length"), "content-encoding", From a776bf4a6d4a72cc52bef20af4d4ab0e9aa589e5 Mon Sep 17 00:00:00 2001 From: Julius Marminge Date: Wed, 29 Jul 2026 01:59:17 +0200 Subject: [PATCH 8/8] refactor(server): flatten compression middleware Co-authored-by: codex --- apps/server/src/http.ts | 8 ++++---- 1 file changed, 4 insertions(+), 4 deletions(-) diff --git a/apps/server/src/http.ts b/apps/server/src/http.ts index 7c002036989..5a380be8fe2 100644 --- a/apps/server/src/http.ts +++ b/apps/server/src/http.ts @@ -105,10 +105,10 @@ const compressHttpResponse = Effect.fnUntraced(function* ( export const httpCompressionLayer = HttpRouter.middleware( (httpEffect) => - Effect.gen(function* () { - const request = yield* HttpServerRequest.HttpServerRequest; - return yield* compressHttpResponse(yield* httpEffect, request.headers["accept-encoding"]); - }), + Effect.flatMap( + Effect.all([httpEffect, HttpServerRequest.HttpServerRequest]), + ([response, request]) => compressHttpResponse(response, request.headers["accept-encoding"]), + ), { global: true }, );