Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
190 changes: 190 additions & 0 deletions apps/server/src/cli/pair.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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";

Expand Down Expand Up @@ -119,6 +123,60 @@ const testDescriptor = {
capabilities: { repositoryIdentity: true },
};

const withProbeServer = <A, E, R>(
{ status, contentType, body }: { status: number; contentType: string; body: string },
run: (baseDir: string, origin: string) => Effect.Effect<A, E, R>,
) =>
Effect.scoped(
Effect.acquireUseRelease(
Effect.callback<NodeHttp.Server>((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 = <A, E, R>(run: (origin: string) => Effect.Effect<A, E, R>) =>
Effect.acquireUseRelease(
Effect.callback<NodeHttp.Server>((resume) => {
Expand Down Expand Up @@ -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<NodeHttp.ServerResponse>();
return Effect.acquireUseRelease(
Effect.callback<NodeHttp.Server>((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));
});
});
66 changes: 62 additions & 4 deletions apps/server/src/cli/pair.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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";

Expand All @@ -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;
Expand Down Expand Up @@ -143,6 +147,11 @@ export class DevServerNotProxiableError extends Schema.TaggedError<DevServerNotP

const isDevServerNotProxiableError = Schema.is(DevServerNotProxiableError);

// Compiled once: the probe decodes every candidate response through it.
const decodeProbeDescriptor = Schema.decodeUnknownEffect(
Schema.fromJsonString(ExecutionEnvironmentDescriptor),
);

/**
* The local endpoint Tailscale Serve should proxy to. Dev servers are
* single-origin, so the web dev server's port is the one to publish; the
Expand Down Expand Up @@ -200,6 +209,31 @@ type EnvironmentProbeResult =
| { readonly _tag: "unreachable" }
| { readonly _tag: "not-a-t3-server" };

const readBoundedProbeBody = (ok: HttpClientResponse.HttpClientResponse) =>
ok.stream.pipe(
Stream.runFoldEffect(
() => ({ bytes: 0, chunks: [] as Array<Uint8Array> }),
(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<EnvironmentProbeResult, never, HttpClient.HttpClient> =>
Expand All @@ -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;
Comment thread
macroscopeapp[bot] marked this conversation as resolved.
}
// 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;
Expand Down
Loading