diff --git a/apps/server/src/http.test.ts b/apps/server/src/http.test.ts index 125c1d509f1..4af2ecb6457 100644 --- a/apps/server/src/http.test.ts +++ b/apps/server/src/http.test.ts @@ -1,11 +1,7 @@ import { expect, it } from "@effect/vitest"; -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 { compressHttpResponse, isLoopbackHostname, resolveDevRedirectUrl } from "./http.ts"; +import { isLoopbackHostname, resolveDevRedirectUrl } from "./http.ts"; describe("http dev routing", () => { it("treats localhost and loopback addresses as local", () => { @@ -30,49 +26,3 @@ describe("http dev routing", () => { ); }); }); - -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 response = 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"); - 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); - } - }), - ); - - 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", - ); - - 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, - ); - - expect(makeResponse("*").headers.vary).toBe("*"); - expect(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 c3070f54055..5a380be8fe2 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 "./httpCompression/HttpResponseCompression.ts"; import { traceRelayRequest } from "./cloud/traceRelayRequest.ts"; import { annotateEnvironmentRequest, @@ -67,10 +67,10 @@ function varyByAcceptEncoding(value: string | undefined): string { return values.has("*") || values.has("accept-encoding") ? value : `${value}, Accept-Encoding`; } -export function compressHttpResponse( +const compressHttpResponse = Effect.fnUntraced(function* ( response: HttpServerResponse.HttpServerResponse, acceptEncoding: string | undefined, -): HttpServerResponse.HttpServerResponse { +) { const body = response.body; if ( body._tag !== "Uint8Array" || @@ -88,30 +88,27 @@ 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 compression = yield* HttpResponseCompression.HttpResponseCompression; 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, cookies: response.cookies, contentType: body.contentType, }); -} +}); export const httpCompressionLayer = HttpRouter.middleware( (httpEffect) => - Effect.gen(function* () { - const request = yield* HttpServerRequest.HttpServerRequest; - return compressHttpResponse(yield* httpEffect, request.headers["accept-encoding"]); - }), + Effect.flatMap( + Effect.all([httpEffect, HttpServerRequest.HttpServerRequest]), + ([response, request]) => compressHttpResponse(response, request.headers["accept-encoding"]), + ), { global: true }, ); diff --git a/apps/server/src/httpCompression/HttpResponseCompression.test.ts b/apps/server/src/httpCompression/HttpResponseCompression.test.ts new file mode 100644 index 00000000000..b608720d8dd --- /dev/null +++ b/apps/server/src/httpCompression/HttpResponseCompression.test.ts @@ -0,0 +1,29 @@ +import { expect, it } from "@effect/vitest"; +import * as NodeStream from "node:stream"; +import * as Effect from "effect/Effect"; + +import * as HttpResponseCompression from "./HttpResponseCompression.ts"; + +it.effect("uses a native Node stream", () => + 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") { + expect(response.body.body).toBeInstanceOf(NodeStream.Readable); + } + }).pipe(Effect.provide(HttpResponseCompression.layerNode)), +); + +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") { + expect(response.body.body).toBeInstanceOf(ReadableStream); + } + }).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..f00c4e1a123 --- /dev/null +++ b/apps/server/src/httpCompression/HttpResponseCompression.ts @@ -0,0 +1,36 @@ +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 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, { + 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 65eee28e59b..cb68c0bba8a 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 HttpResponseCompression from "./httpCompression/HttpResponseCompression.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(HttpResponseCompression.layerNode), Layer.provide(layerConfig), ); diff --git a/apps/server/src/server.ts b/apps/server/src/server.ts index d78555cb30b..0b7c8ccd074 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, @@ -151,6 +152,9 @@ const HttpServerLive = Layer.unwrap( }), ); +const HttpResponseCompressionLive = + typeof Bun !== "undefined" ? HttpResponseCompression.layerBun : HttpResponseCompression.layerNode; + const PlatformServicesLive = Layer.unwrap( Effect.gen(function* () { if (typeof Bun !== "undefined") { @@ -516,6 +520,7 @@ export const makeServerLayer = Layer.unwrap( return serverApplicationLayer.pipe( Layer.provideMerge(RuntimeServicesLive), Layer.provideMerge(serverRelayBrokerTracingLayer), + Layer.provideMerge(HttpResponseCompressionLive), Layer.provideMerge(HttpServerLive), Layer.provide(ObservabilityLive), Layer.provideMerge(FetchHttpClient.layer),