diff --git a/apps/server/src/cli/pair.test.ts b/apps/server/src/cli/pair.test.ts index 39fe8652ce61..ba38e8aa7a83 100644 --- a/apps/server/src/cli/pair.test.ts +++ b/apps/server/src/cli/pair.test.ts @@ -9,7 +9,11 @@ import * as NetService from "@t3tools/shared/Net"; import { HostProcessEnvironment } from "@t3tools/shared/hostProcess"; import { assert, describe, expect, it } from "@effect/vitest"; import * as Effect from "effect/Effect"; +import * as Deferred from "effect/Deferred"; +import * as Fiber from "effect/Fiber"; +import * as FileSystem from "effect/FileSystem"; import * as Layer from "effect/Layer"; +import * as Schedule from "effect/Schedule"; import * as TestConsole from "effect/testing/TestConsole"; import { Command } from "effect/cli"; @@ -119,6 +123,60 @@ const testDescriptor = { capabilities: { repositoryIdentity: true }, }; +const withProbeServer = ( + { status, contentType, body }: { status: number; contentType: string; body: string }, + run: (baseDir: string, origin: string) => Effect.Effect, +) => + Effect.scoped( + Effect.acquireUseRelease( + Effect.callback((resume) => { + const server = NodeHttp.createServer((request, response) => { + if (request.url === "/.well-known/t3/environment") { + response.writeHead(status, { "content-type": contentType }); + response.end(body); + return; + } + response.writeHead(404); + response.end(); + }); + server.listen(0, "127.0.0.1", () => resume(Effect.succeed(server))); + }), + (server) => + Effect.gen(function* () { + const address = server.address(); + if (address === null || typeof address === "string") { + return yield* Effect.die(new Error("Expected a TCP address")); + } + const fs = yield* FileSystem.FileSystem; + const baseDir = yield* fs.makeTempDirectoryScoped({ prefix: "t3-pair-probe-test-" }); + yield* persistServerRuntimeState({ + path: NodePath.join(baseDir, "userdata", "server-runtime.json"), + state: yield* makePersistedServerRuntimeState({ + config: { host: "127.0.0.1", devUrl: undefined }, + port: address.port, + }), + }); + return yield* run(baseDir, `http://127.0.0.1:${String(address.port)}`); + }), + (server) => + Effect.sync(() => { + server.closeAllConnections(); + server.close(); + }), + ), + ); + +const assertPairRejected = (baseDir: string) => + Effect.gen(function* () { + const error = yield* provideCliTestLayers( + runCli(["pair", "--base-dir", baseDir]).pipe(Effect.flip), + ); + const rendered = String( + typeof error === "object" && error !== null && "cause" in error ? error.cause : error, + ); + assert.include(rendered, "No running T3 Code server found."); + }); + const withDescriptorServer = (run: (origin: string) => Effect.Effect) => Effect.acquireUseRelease( Effect.callback((resume) => { @@ -288,4 +346,136 @@ describe("t3 pair", () => { assert.include(rendered, "No running T3 Code server found."); }).pipe(Effect.provide(NodeServices.layer)), ); + + it.effect("does not decode an oversized stranger body as a server descriptor", () => + withProbeServer( + { + status: 200, + contentType: "application/json", + // Schema-valid JSON that would pair without the byte cap. + body: JSON.stringify({ ...testDescriptor, padding: "x".repeat(128 * 1024) }), + }, + assertPairRejected, + ).pipe(Effect.provide(NodeServices.layer)), + ); + + it.effect("rejects a descriptor that exceeds the cap in UTF-8 bytes only", () => + withProbeServer( + { + status: 200, + contentType: "application/json", + // Two UTF-8 bytes per character: below 64 KiB in string length, + // above it on the wire. A string-length cap would accept this. + body: JSON.stringify({ ...testDescriptor, label: "é".repeat(40 * 1024) }), + }, + assertPairRejected, + ).pipe(Effect.provide(NodeServices.layer)), + ); + + it.effect("does not pair when a valid descriptor arrives as non-JSON", () => + withProbeServer( + { + status: 200, + contentType: "text/html; charset=utf-8", + body: JSON.stringify(testDescriptor), + }, + assertPairRejected, + ).pipe(Effect.provide(NodeServices.layer)), + ); + + it.effect("does not pair when a valid descriptor arrives with an error status", () => + withProbeServer( + { status: 500, contentType: "application/json", body: JSON.stringify(testDescriptor) }, + assertPairRejected, + ).pipe(Effect.provide(NodeServices.layer)), + ); + + it.effect("pairs when the descriptor arrives with an uppercase JSON content type", () => + withProbeServer( + { + status: 200, + contentType: "Application/JSON; charset=utf-8", + body: JSON.stringify(testDescriptor), + }, + (baseDir) => + Effect.gen(function* () { + const output = yield* captureStdout(runCli(["pair", "--base-dir", baseDir])); + assert.include(output, "Pairing with pair-test ("); + assert.include(output, "/pair#token="); + }), + ).pipe(Effect.provide(NodeServices.layer)), + ); + + it.live("times out a slow-drip stranger body instead of hanging discovery", () => { + // Handoff from the raw Node request handler below: it holds the response + // open so the test fiber can drip into it. Completed synchronously from + // the callback via `doneUnsafe`, so no manual Effect runtime is created + // in the test (see `t3code(no-manual-effect-runtime-in-tests)`). + const dripTarget = Deferred.makeUnsafe(); + return Effect.acquireUseRelease( + Effect.callback((resume) => { + const server = NodeHttp.createServer((request, response) => { + if (request.url === "/.well-known/t3/environment") { + response.writeHead(200, { "content-type": "application/json" }); + // A never-ending drip that stays under the size cap: discovery + // must give up via the body timeout rather than hang on the open + // stream. + response.write(`{"environmentId":`); + Deferred.doneUnsafe(dripTarget, Effect.succeed(response)); + return; + } + response.writeHead(404); + response.end(); + }); + server.listen(0, "127.0.0.1", () => resume(Effect.succeed(server))); + }), + (server) => + Effect.scoped( + Effect.gen(function* () { + const address = server.address(); + if (address === null || typeof address === "string") { + return yield* Effect.die(new Error("Expected a TCP address")); + } + const baseDir = NodeFS.mkdtempSync( + NodePath.join(NodeOS.tmpdir(), "t3-pair-drip-test-"), + ); + const statePath = NodePath.join(baseDir, "userdata", "server-runtime.json"); + yield* persistServerRuntimeState({ + path: statePath, + state: yield* makePersistedServerRuntimeState({ + config: { host: "127.0.0.1", devUrl: undefined }, + port: address.port, + }), + }); + + const cliFiber = yield* provideCliTestLayers( + runCli(["pair", "--base-dir", baseDir]).pipe(Effect.flip), + ).pipe(Effect.forkChild); + // The probe request holds the response open; the drip loop is a + // scoped fork, so it is interrupted when the test settles. + const dripResponse = yield* Deferred.await(dripTarget); + yield* Effect.repeat( + Effect.sync(() => { + if (!dripResponse.destroyed) { + dripResponse.write(" "); + } + }), + Schedule.spaced("200 millis"), + ).pipe(Effect.forkScoped); + + const error = yield* Fiber.join(cliFiber); + + const rendered = String( + typeof error === "object" && error !== null && "cause" in error ? error.cause : error, + ); + assert.include(rendered, "No running T3 Code server found."); + }), + ), + (server) => + Effect.sync(() => { + server.closeAllConnections(); + server.close(); + }), + ).pipe(Effect.provide(NodeServices.layer)); + }); }); diff --git a/apps/server/src/cli/pair.ts b/apps/server/src/cli/pair.ts index 452f484e09c4..9921bf9e32d1 100644 --- a/apps/server/src/cli/pair.ts +++ b/apps/server/src/cli/pair.ts @@ -32,6 +32,7 @@ import * as Layer from "effect/Layer"; import * as Option from "effect/Option"; import * as References from "effect/References"; import * as Schema from "effect/Schema"; +import * as Stream from "effect/Stream"; import { Command, Flag, GlobalFlag } from "effect/cli"; import { FetchHttpClient, HttpClient, HttpClientRequest, HttpClientResponse } from "effect/http"; @@ -55,6 +56,9 @@ import { baseDirFlag, DurationFromString } from "./config.ts"; const WELL_KNOWN_ENVIRONMENT_PATH = "/.well-known/t3/environment"; const PAIR_PROBE_TIMEOUT = Duration.millis(2_500); +// The environment descriptor is a few hundred bytes; anything larger served +// from the well-known path cannot be a T3 server. +const MAX_PROBE_BODY_BYTES = 64 * 1024; // Tailscale provisions an HTTPS certificate on the first request to a fresh // serve mapping, which can take a few seconds. const TAILSCALE_PROBE_ATTEMPTS = 5; @@ -143,6 +147,11 @@ export class DevServerNotProxiableError extends Schema.TaggedError + ok.stream.pipe( + Stream.runFoldEffect( + () => ({ bytes: 0, chunks: [] as Array }), + (acc, chunk) => { + if (acc.bytes + chunk.byteLength > MAX_PROBE_BODY_BYTES) { + return Effect.fail({ _tag: "not-a-t3-server" } as const); + } + acc.bytes += chunk.byteLength; + acc.chunks.push(chunk); + return Effect.succeed(acc); + }, + ), + Effect.map(({ bytes, chunks }) => { + const merged = new Uint8Array(bytes); + let offset = 0; + for (const chunk of chunks) { + merged.set(chunk, offset); + offset += chunk.byteLength; + } + return new TextDecoder().decode(merged); + }), + Effect.timeout(PAIR_PROBE_TIMEOUT), + ); + const probeEnvironmentDescriptor = ( baseUrl: string, ): Effect.Effect => @@ -214,14 +248,38 @@ const probeEnvironmentDescriptor = ( // Bad-gateway family means a proxy (Tailscale Serve) answered for a // backend that is gone — a stale mapping, not a live occupant. Treating // it as unreachable lets `t3 pair --tailscale` repair its own mapping - // after the server's port changed. + // after the server's port changed. Drain the body so the pooled + // connection is reusable for the re-configured mapping; the drain itself + // is time-bounded and best-effort. if (response.status === 502 || response.status === 503 || response.status === 504) { + yield* Effect.ignore(readBoundedProbeBody(response)); return { _tag: "unreachable" } as const; } // Anything else that answered HTTP but not with a valid descriptor is - // some other service. - const descriptor = yield* HttpClientResponse.filterStatusOk(response).pipe( - Effect.flatMap(HttpClientResponse.schemaBodyJson(ExecutionEnvironmentDescriptor)), + // some other service. Refuse to decode a stranger's body blind: the + // descriptor is a few hundred bytes of JSON, so non-JSON content is + // classified without decoding, but the body is still drained boundedly + // so the pooled connection is reusable — a stranger that holds its body + // open must not pin pool slots across repeated probes. + const contentType = response.headers["content-type"] ?? ""; + const mediaType = contentType.split(";", 1)[0]?.trim().toLowerCase(); + if (mediaType !== "application/json" && !mediaType?.endsWith("+json")) { + yield* Effect.ignore(readBoundedProbeBody(response)); + return { _tag: "not-a-t3-server" } as const; + } + // A non-2xx answer is still a stranger, but its body must be drained + // boundedly for the same reason: an unconsumed stream pins the pooled + // connection across repeated probes. + if (response.status < 200 || response.status >= 300) { + yield* Effect.ignore(readBoundedProbeBody(response)); + return { _tag: "not-a-t3-server" } as const; + } + const descriptor = yield* readBoundedProbeBody(response).pipe( + Effect.flatMap((body) => + decodeProbeDescriptor(body).pipe( + Effect.mapError(() => ({ _tag: "not-a-t3-server" }) as const), + ), + ), Effect.mapError(() => ({ _tag: "not-a-t3-server" }) as const), ); return { _tag: "descriptor", descriptor } as const;