diff --git a/apps/server/src/cli/pair.test.ts b/apps/server/src/cli/pair.test.ts index 39fe8652ce61..e641d66dadfc 100644 --- a/apps/server/src/cli/pair.test.ts +++ b/apps/server/src/cli/pair.test.ts @@ -8,10 +8,19 @@ import * as NodeServices from "@effect/platform-node/NodeServices"; import * as NetService from "@t3tools/shared/Net"; import { HostProcessEnvironment } from "@t3tools/shared/hostProcess"; import { assert, describe, expect, it } from "@effect/vitest"; +import * as Clock from "effect/Clock"; +import * as Deferred from "effect/Deferred"; +import * as Duration from "effect/Duration"; +import * as FileSystem from "effect/FileSystem"; import * as Effect from "effect/Effect"; +import * as Fiber from "effect/Fiber"; import * as Layer from "effect/Layer"; +import * as Ref from "effect/Ref"; +import * as Schedule from "effect/Schedule"; +import * as TestClock from "effect/testing/TestClock"; import * as TestConsole from "effect/testing/TestConsole"; import { Command } from "effect/cli"; +import { HttpClient, HttpClientResponse } from "effect/http"; import { cli } from "../binCli.ts"; import { @@ -25,7 +34,9 @@ import { type PersistedServerRuntimeState, } from "../serverRuntimeState.ts"; import { + awaitEnvironmentDescriptor, DevServerNotProxiableError, + discoverPairTarget, resolveDirectPairingBaseUrl, resolveTailscaleLocalTarget, } from "./pair.ts"; @@ -119,6 +130,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) => { @@ -143,6 +208,25 @@ const withDescriptorServer = (run: (origin: string) => Effect.Effect Effect.sync(() => server.close()), ); +const makeDiscoveryFixture = Effect.gen(function* () { + const fs = yield* FileSystem.FileSystem; + const baseDir = yield* fs.makeTempDirectoryScoped({ prefix: "t3-pair-probes-" }); + const candidates = yield* Effect.forEach(["userdata", "dev"] as const, (variant, index) => + Effect.gen(function* () { + const state = yield* makePersistedServerRuntimeState({ + config: { host: "127.0.0.1", devUrl: undefined }, + port: 10_000 + index, + }); + yield* persistServerRuntimeState({ + path: NodePath.join(baseDir, variant, "server-runtime.json"), + state, + }); + return { variant, state }; + }), + ); + return { baseDir, candidates }; +}); + describe("t3 pair", () => { it.effect("mints a token and prints a QR pairing URL for a live server", () => withDescriptorServer((origin) => @@ -288,4 +372,286 @@ 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.each(["Application/JSON; charset=utf-8", "application/vnd.t3+json"])( + "pairs when the descriptor arrives with content type %s", + (contentType) => + withProbeServer( + { status: 200, contentType, 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.each([ + { status: 200, contentType: "application/json" }, + { status: 200, contentType: "text/html" }, + { status: 500, contentType: "application/json" }, + ])("times out a slow-drip $status $contentType body", ({ status, contentType }) => { + // 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(status, { "content-type": contentType }); + // 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)); + }); + it.effect("overlaps probes and cancels a pending lower-priority request", () => + Effect.gen(function* () { + const { baseDir } = yield* makeDiscoveryFixture; + const lowerStarted = yield* Deferred.make(); + const interrupted = yield* Ref.make(false); + const client = HttpClient.make((request, url, signal) => + url.port === "10000" + ? Deferred.await(lowerStarted).pipe( + Effect.as(HttpClientResponse.fromWeb(request, Response.json(testDescriptor))), + ) + : Deferred.succeed(lowerStarted, signal).pipe( + Effect.andThen(Effect.never), + Effect.onInterrupt(() => Ref.set(interrupted, true)), + ), + ); + const target = yield* discoverPairTarget(baseDir).pipe( + Effect.provideService(HttpClient.HttpClient, client), + ); + expect(target.variant).toBe("userdata"); + expect(yield* Ref.get(interrupted)).toBe(true); + expect((yield* Deferred.await(lowerStarted)).aborted).toBe(true); + }).pipe(Effect.provide(NodeServices.layer)), + ); +}); + +describe("pair discovery order", () => { + it.effect("shares the probe deadline between headers and body", () => + Effect.gen(function* () { + const { baseDir } = yield* makeDiscoveryFixture; + const firstRequest = yield* Deferred.make(); + const client = HttpClient.make((request) => + Effect.gen(function* () { + yield* Deferred.succeed(firstRequest, undefined); + yield* Effect.sleep("2 seconds"); + return HttpClientResponse.fromWeb( + request, + new Response(new ReadableStream(), { + headers: { "content-type": "application/json" }, + }), + ); + }), + ); + const fiber = yield* discoverPairTarget(baseDir).pipe( + Effect.provideService(HttpClient.HttpClient, client), + Effect.flip, + Effect.forkScoped, + ); + yield* Deferred.await(firstRequest); + yield* TestClock.adjust("2 seconds"); + yield* TestClock.adjust("500 millis"); + const error = yield* Fiber.join(fiber); + expect(error._tag).toBe("NoRunningServerError"); + }).pipe(Effect.provide(NodeServices.layer)), + ); + + it.effect.each([ + { status: 200, winner: "userdata" }, + { status: 503, winner: "dev" }, + { status: 404, winner: "dev" }, + ])( + "selects $winner when the earlier probe returns $status after the later success", + ({ status, winner }) => + Effect.gen(function* () { + const { baseDir, candidates } = yield* makeDiscoveryFixture; + const lowerCompleted = yield* Deferred.make(); + const releaseFirst = yield* Deferred.make(); + const firstStarted = yield* Deferred.make(); + const client = HttpClient.make((request, url) => + Effect.gen(function* () { + if (url.port === "10000") { + yield* Deferred.succeed(firstStarted, undefined); + yield* Deferred.await(lowerCompleted); + yield* Deferred.await(releaseFirst); + return HttpClientResponse.fromWeb(request, Response.json(testDescriptor, { status })); + } + yield* Deferred.await(firstStarted); + yield* Deferred.succeed(lowerCompleted, undefined); + return HttpClientResponse.fromWeb(request, Response.json(testDescriptor)); + }), + ); + const fiber = yield* discoverPairTarget(baseDir).pipe( + Effect.provideService(HttpClient.HttpClient, client), + Effect.forkScoped, + ); + yield* Deferred.await(lowerCompleted); + yield* Deferred.succeed(releaseFirst, undefined); + const target = yield* Fiber.join(fiber); + expect(target.variant).toBe(winner); + expect(target.state).toEqual( + candidates.find((candidate) => candidate.variant === winner)?.state, + ); + expect(target.descriptor).toEqual(testDescriptor); + }).pipe(Effect.provide(NodeServices.layer)), + ); + + it.effect("interrupts all pending requests when discovery is cancelled", () => + Effect.gen(function* () { + const { baseDir } = yield* makeDiscoveryFixture; + const started = yield* Deferred.make(); + const signals = yield* Ref.make>([]); + const client = HttpClient.make((_request, _url, signal) => + Effect.gen(function* () { + const pending = yield* Ref.updateAndGet(signals, (current) => [...current, signal]); + if (pending.length === 2) { + yield* Deferred.succeed(started, undefined); + } + return yield* Effect.never; + }), + ); + const fiber = yield* discoverPairTarget(baseDir).pipe( + Effect.provideService(HttpClient.HttpClient, client), + Effect.forkScoped, + ); + yield* Deferred.await(started); + yield* Fiber.interrupt(fiber); + expect((yield* Ref.get(signals)).map((signal) => signal.aborted)).toEqual([true, true]); + }).pipe(Effect.provide(NodeServices.layer)), + ); +}); + +describe("pair tailscale probe budget", () => { + it.effect("retries all attempts but sleeps only between them", () => + Effect.gen(function* () { + const startedAt = yield* Clock.currentTimeMillis; + const attempts = yield* Ref.make>([]); + const firstAttempt = yield* Deferred.make(); + const client = HttpClient.make((request) => + Effect.gen(function* () { + const now = yield* Clock.currentTimeMillis; + yield* Ref.update(attempts, (times) => [...times, now - startedAt]); + yield* Deferred.succeed(firstAttempt, undefined); + return HttpClientResponse.fromWeb(request, new Response(null, { status: 503 })); + }), + ); + const fiber = yield* awaitEnvironmentDescriptor("http://127.0.0.1:1").pipe( + Effect.provide(Layer.succeed(HttpClient.HttpClient, client)), + Effect.forkScoped, + ); + yield* Deferred.await(firstAttempt); + yield* TestClock.adjust(Duration.millis(3_999)); + expect(yield* Ref.get(attempts)).toEqual([0, 1_000, 2_000, 3_000]); + yield* TestClock.adjust(Duration.millis(1)); + const result = yield* Fiber.join(fiber); + expect(result._tag).toBe("unreachable"); + expect(yield* Ref.get(attempts)).toEqual([0, 1_000, 2_000, 3_000, 4_000]); + expect((yield* Clock.currentTimeMillis) - startedAt).toBe(4_000); + }), + ); }); diff --git a/apps/server/src/cli/pair.ts b/apps/server/src/cli/pair.ts index 452f484e09c4..fdf33977440b 100644 --- a/apps/server/src/cli/pair.ts +++ b/apps/server/src/cli/pair.ts @@ -24,14 +24,17 @@ import { readTailscaleStatus, } from "@t3tools/tailscale"; import * as Config from "effect/Config"; +import * as Clock from "effect/Clock"; import * as Console from "effect/Console"; import * as DateTime from "effect/DateTime"; import * as Duration from "effect/Duration"; import * as Effect from "effect/Effect"; +import * as Fiber from "effect/Fiber"; 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 +58,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 +149,11 @@ export class DevServerNotProxiableError extends Schema.TaggedError + Effect.gen(function* () { + const remaining = deadline - (yield* Clock.monotonicTimeNanos); + return yield* 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(Duration.nanos(remaining > 0n ? remaining : 0n)), + ); + }); + const probeEnvironmentDescriptor = ( baseUrl: string, ): Effect.Effect => Effect.gen(function* () { const client = yield* HttpClient.HttpClient; const request = HttpClientRequest.get(new URL(WELL_KNOWN_ENVIRONMENT_PATH, baseUrl).toString()); + const deadline = (yield* Clock.monotonicTimeNanos) + Duration.toNanosUnsafe(PAIR_PROBE_TIMEOUT); const response = yield* client.execute(request).pipe( Effect.timeout(PAIR_PROBE_TIMEOUT), // Transport failure or timeout: nothing (reachable) is listening there. @@ -214,14 +254,34 @@ 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, deadline)); 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, deadline)); + 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, deadline)); + return { _tag: "not-a-t3-server" } as const; + } + const descriptor = yield* readBoundedProbeBody(response, deadline).pipe( + Effect.flatMap(decodeProbeDescriptor), Effect.mapError(() => ({ _tag: "not-a-t3-server" }) as const), ); return { _tag: "descriptor", descriptor } as const; @@ -234,7 +294,7 @@ interface DiscoveredPairTarget { readonly descriptor: ExecutionEnvironmentDescriptor; } -const discoverPairTarget = Effect.fn("pair.discoverPairTarget")(function* ( +export const discoverPairTarget = Effect.fn("pair.discoverPairTarget")(function* ( explicitBaseDir: string | undefined, ) { const bases: Array = []; @@ -253,6 +313,13 @@ const discoverPairTarget = Effect.fn("pair.discoverPairTarget")(function* ( } const checkedStatePaths: Array = []; + // Probe concurrently, but keep state-file precedence when choosing a hit: + // a later success must wait for every earlier candidate to finish. + const candidates: Array<{ + readonly baseDir: string; + readonly variant: PairStateVariant; + readonly state: PersistedServerRuntimeState; + }> = []; for (const baseDir of new Set(bases)) { for (const variant of ["userdata", "dev"] as const) { const derivedPaths = yield* ServerConfig.deriveServerPaths( @@ -272,18 +339,34 @@ const discoverPairTarget = Effect.fn("pair.discoverPairTarget")(function* ( if (!isProcessAlive(state.value.pid)) { continue; } - const probed = yield* probeEnvironmentDescriptor(state.value.origin); - if (probed._tag !== "descriptor") { - continue; - } - return { - baseDir, - variant, - state: state.value, - descriptor: probed.descriptor, - } satisfies DiscoveredPairTarget; + candidates.push({ baseDir, variant, state: state.value }); } } + const hit = yield* Effect.scoped( + Effect.gen(function* () { + const fibers = yield* Effect.forEach(candidates, (candidate) => + probeEnvironmentDescriptor(candidate.state.origin).pipe( + Effect.forkScoped, + Effect.map((fiber) => ({ ...candidate, fiber })), + ), + ); + for (const { fiber, baseDir, variant, state } of fibers) { + const result = yield* Fiber.join(fiber); + if (result._tag === "descriptor") { + return Option.some({ + baseDir, + variant, + state, + descriptor: result.descriptor, + } satisfies DiscoveredPairTarget); + } + } + return Option.none(); + }), + ); + if (Option.isSome(hit)) { + return hit.value; + } return yield* new NoRunningServerError({ checkedStatePaths }); }); @@ -344,14 +427,17 @@ const makePairServerConfig = Effect.fn(function* (input: { }); }); -const awaitEnvironmentDescriptor = Effect.fn(function* (baseUrl: string) { +export const awaitEnvironmentDescriptor = Effect.fn(function* (baseUrl: string) { let last: EnvironmentProbeResult = { _tag: "unreachable" }; for (let attempt = 0; attempt < TAILSCALE_PROBE_ATTEMPTS; attempt += 1) { last = yield* probeEnvironmentDescriptor(baseUrl); if (last._tag === "descriptor") { return last; } - yield* Effect.sleep(TAILSCALE_PROBE_RETRY_DELAY); + // No sleep after the final attempt: nothing else will use the wait. + if (attempt + 1 < TAILSCALE_PROBE_ATTEMPTS) { + yield* Effect.sleep(TAILSCALE_PROBE_RETRY_DELAY); + } } return last; });