Skip to content
Merged
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
2 changes: 1 addition & 1 deletion apps/server/src/auth/replayMarkers.ts
Original file line number Diff line number Diff line change
Expand Up @@ -7,7 +7,7 @@ import * as Option from "effect/Option";
import * as Path from "effect/Path";
import * as Schedule from "effect/Schedule";

import { CLOUD_REPLAY_MARKER_PREFIXES } from "../cloud/CloudLink.ts";
import { CLOUD_REPLAY_MARKER_PREFIXES } from "../cloud/linkChecks.ts";
import * as ServerConfig from "../config.ts";
import { forkParked } from "../serverActivation.ts";
import { DPOP_REPLAY_MARKER_PREFIX } from "./dpop.ts";
Expand Down

Large diffs are not rendered by default.

251 changes: 18 additions & 233 deletions apps/server/src/cloud/CloudLink.ts
Original file line number Diff line number Diff line change
Expand Up @@ -7,7 +7,6 @@
import * as NodeCrypto from "node:crypto";
import {
AuthStandardClientScopes,
DESKTOP_UPDATE_RESTART_MARKER_FILE,
EnvironmentCloudEndpointUnavailableError,
type EnvironmentCloudLinkStateResult,
type EnvironmentCloudPreferencesRequest,
Expand Down Expand Up @@ -53,7 +52,6 @@ import {
verifyRelayJwt,
} from "@t3tools/shared/relayJwt";
import { isSecureRelayUrl } from "@t3tools/shared/relayUrl";
import * as Clock from "effect/Clock";
import * as Context from "effect/Context";
import * as Crypto from "effect/Crypto";
import * as DateTime from "effect/DateTime";
Expand Down Expand Up @@ -110,29 +108,24 @@ import * as ManagedEndpointRuntime from "./ManagedEndpointRuntime.ts";
import { relayUrlConfig } from "./publicConfig.ts";
import { filterRelayResponse, relayRequestError, shouldRetryCloudLink } from "./relayResponse.ts";
import {
SERVICE_STATE_FILE,
SERVICE_STOP_MARKER_FILE,
serviceStateHasPendingUpdate,
} from "./serviceProtocol.ts";

const CLOUD_MINT_NONCE_PREFIX = "cloud-mint-nonce-";
const CLOUD_MINT_JTI_PREFIX = "cloud-mint-jti-";
const CLOUD_HEALTH_NONCE_PREFIX = "cloud-health-nonce-";
const CLOUD_HEALTH_JTI_PREFIX = "cloud-health-jti-";
/** Secret store name prefixes of cloud replay markers. The server prunes expired ones. */
export const CLOUD_REPLAY_MARKER_PREFIXES = [
CLOUD_MINT_NONCE_PREFIX,
CLOUD_MINT_JTI_PREFIX,
CLOUD_HEALTH_NONCE_PREFIX,
CLOUD_HEALTH_JTI_PREFIX,
] as const;
const CLOUD_PROOF_MAX_LIFETIME_SECONDS = 5 * 60;
const CLOUD_PROOF_CLOCK_SKEW_SECONDS = 60;
// The desktop app stops its backends within seconds of writing the marker.
const DESKTOP_UPDATE_RESTART_MARKER_TTL = Duration.minutes(1);
const MANAGED_ENDPOINT_PROVISION_REQUEST_TIMEOUT = Duration.minutes(2);
const LOOPBACK_HOSTNAMES = new Set(["127.0.0.1", "::1", "localhost"]);
CLOUD_HEALTH_NONCE_PREFIX,
CLOUD_MINT_JTI_PREFIX,
CLOUD_MINT_NONCE_PREFIX,
consumeCloudReplayGuards,
hasBoundedCloudProofLifetime,
hasExactScope,
hasForwardedAuthorityHeaders,
isAllowedEndpointOrigin,
isSupportedLinkProviderKind,
linkProofScopes,
managedEndpointRuntimeConfigsMatch,
parseManagedEndpointLocalOrigin,
requestAbsoluteUrl,
} from "./linkChecks.ts";
import { desktopUpdateRestartPending, pendingUpdateHandoffExists } from "./updateHandoff.ts";

const MANAGED_ENDPOINT_PROVISION_REQUEST_TIMEOUT = Duration.minutes(2);
const failEnvironmentCloudInternalError =
(message: string) =>
(cause: unknown): Effect.Effect<never, EnvironmentHttpInternalServerError> =>
Expand Down Expand Up @@ -164,37 +157,6 @@ function stringToBytes(value: string): Uint8Array {
return new TextEncoder().encode(value);
}

export function consumeCloudReplayGuards(input: {
readonly secrets: ServerSecretStore.ServerSecretStore["Service"];
readonly names: ReadonlyArray<string>;
readonly value: Uint8Array;
}) {
return Effect.forEach(
input.names,
(name) =>
input.secrets.create(name, input.value).pipe(
Effect.as(true),
Effect.catchIf(ServerSecretStore.isSecretStoreError, (error) =>
ServerSecretStore.isSecretAlreadyExistsError(error)
? Effect.succeed(false)
: Effect.fail(error),
),
),
{ concurrency: input.names.length },
).pipe(Effect.map((created) => created.every(Boolean)));
}

function normalizePemForSignedPayload(value: string): string {
return value.trim();
}

function normalizeHostname(hostname: string): string {
return hostname
.trim()
.toLowerCase()
.replace(/^\[(.*)\]$/, "$1");
}

function validateCloudMintPublicKey(
publicKey: string,
): Effect.Effect<void, EnvironmentHttpBadRequestError> {
Expand Down Expand Up @@ -251,190 +213,13 @@ function validateRelayConfigPayload(
return Effect.void;
}

function isLoopbackHostname(hostname: string): boolean {
return LOOPBACK_HOSTNAMES.has(normalizeHostname(hostname));
}

function firstForwardedHeaderValue(value: string | undefined): string | undefined {
const first = value?.split(",")[0]?.trim();
return first && first.length > 0 ? first : undefined;
}

function requestAbsoluteUrl(request: HttpServerRequest.HttpServerRequest): string | null {
try {
return new URL(request.originalUrl).href;
} catch {
const host = firstForwardedHeaderValue(request.headers.host) ?? "127.0.0.1";
try {
return new URL(request.originalUrl, `http://${host}`).href;
} catch {
return null;
}
}
}

function hasForwardedAuthorityHeaders(request: HttpServerRequest.HttpServerRequest): boolean {
return (
firstForwardedHeaderValue(request.headers["x-forwarded-host"]) !== undefined ||
firstForwardedHeaderValue(request.headers["x-forwarded-proto"]) !== undefined
);
}

function endpointRequestPort(url: URL): number {
return Number(url.port || (url.protocol === "https:" ? 443 : 80));
}

export function parseManagedEndpointLocalOrigin(localOrigin: string) {
const url = new URL(localOrigin);
if (
localOrigin !== localOrigin.trim() ||
(url.protocol !== "http:" && url.protocol !== "https:") ||
url.username !== "" ||
url.password !== "" ||
url.pathname !== "/" ||
url.search !== "" ||
url.hash !== "" ||
localOrigin.includes("?") ||
localOrigin.includes("#")
) {
throw new Error("Invalid local origin");
}
const wsUrl = new URL(url.origin);
wsUrl.protocol = url.protocol === "https:" ? "wss:" : "ws:";
return {
httpBaseUrl: url.origin,
wsBaseUrl: wsUrl.origin,
origin: {
localHttpHost: url.hostname,
localHttpPort: endpointRequestPort(url),
} satisfies RelayManagedEndpointOrigin,
};
}

function isAllowedEndpointOrigin(input: {
readonly origin: RelayManagedEndpointOrigin;
readonly requestUrl: string;
}): boolean {
if (!isLoopbackHostname(input.origin.localHttpHost)) {
return false;
}

const url = new URL(input.requestUrl);
if (!isLoopbackHostname(url.hostname)) {
return false;
}

return input.origin.localHttpPort === endpointRequestPort(url);
}

// A managed (Cloudflare tunnel) endpoint is provisioned by the relay and must
// point at a loopback origin. A manual endpoint is reached out of band (e.g.
// Tailscale) or not advertised at all for publish-only links, so it is not
// tied to the managed-tunnel scope.
export function isSupportedLinkProviderKind(request: RelayLinkProofRequest): boolean {
return (
request.endpoint.providerKind === "cloudflare_tunnel" ||
request.endpoint.providerKind === "manual"
);
}

export function linkProofScopes(
request: RelayLinkProofRequest,
): RelayEnvironmentLinkProofPayload["scopes"] {
return request.endpoint.providerKind === "cloudflare_tunnel"
? ["agent_activity_notifications", "managed_tunnels"]
: ["agent_activity_notifications"];
}

function hasExactScope(input: {
readonly scopes: ReadonlyArray<string>;
readonly expected: string;
}): boolean {
return input.scopes.length === 1 && input.scopes[0] === input.expected;
}

function hasBoundedCloudProofLifetime(input: {
readonly iat: number;
readonly exp: number;
readonly nowSeconds: number;
}): boolean {
return (
input.exp > input.iat &&
input.exp - input.iat <= CLOUD_PROOF_MAX_LIFETIME_SECONDS &&
input.iat <= input.nowSeconds + CLOUD_PROOF_CLOCK_SKEW_SECONDS
);
}

function managedEndpointRuntimeConfigsMatch(
left: RelayManagedEndpointRuntimeConfig,
right: RelayManagedEndpointRuntimeConfig,
): boolean {
return (
left.providerKind === right.providerKind &&
left.connectorToken === right.connectorToken &&
left.tunnelId === right.tunnelId &&
left.tunnelName === right.tunnelName
);
function normalizePemForSignedPayload(value: string): string {
return value.trim();
}

const decodeCloudHealthProof = Schema.decodeUnknownEffect(RelayCloudEnvironmentHealthProofPayload);
const decodeCloudMintProof = Schema.decodeUnknownEffect(RelayCloudMintCredentialProofPayload);

// The launcher owns this durable state, so read it directly both when a trial
// decides whether it owns pre-activation cleanup and while a server tears down.
export const pendingServiceUpdateExists = Effect.gen(function* () {
const config = yield* ServerConfig.ServerConfig;
const fs = yield* FileSystem.FileSystem;
const path = yield* Path.Path;
const runtimeDir = path.join(config.baseDir, "runtime");
const stateText = yield* fs
.readFileString(path.join(runtimeDir, SERVICE_STATE_FILE))
.pipe(Effect.option);
return Option.isSome(stateText) && serviceStateHasPendingUpdate(stateText.value);
});

// A pending update alone is not proof a replacement server is coming: an
// explicit launcher stop (`t3 service uninstall`, `systemctl stop`,
// `launchctl bootout`) during
// the pending window also tears this server down. The launcher marks that case
// just before it signals the child, so pending + no marker is the handoff.
const pendingUpdateHandoffExists = Effect.gen(function* () {
if (!(yield* pendingServiceUpdateExists)) {
return false;
}
const config = yield* ServerConfig.ServerConfig;
const fs = yield* FileSystem.FileSystem;
const path = yield* Path.Path;
const runtimeDir = path.join(config.baseDir, "runtime");
const stopping = yield* fs
.exists(path.join(runtimeDir, SERVICE_STOP_MARKER_FILE))
.pipe(Effect.orElseSucceed(() => false));
return !stopping;
});

// The desktop app writes its marker right before it stops this server to
// install an update, whether a remote client or the local app started it.
// Reading consumes it, so shutdown checks it first. Only a fresh marker counts,
// so a marker the server never read (a hard kill) cannot keep the tunnel on a
// later quit.
const desktopUpdateRestartPending = Effect.gen(function* () {
const config = yield* ServerConfig.ServerConfig;
const fs = yield* FileSystem.FileSystem;
const path = yield* Path.Path;
const markerPath = path.join(config.baseDir, "runtime", DESKTOP_UPDATE_RESTART_MARKER_FILE);
const marker = yield* fs.stat(markerPath).pipe(Effect.option);
if (Option.isNone(marker)) {
return false;
}
yield* fs.remove(markerPath).pipe(Effect.ignore);
const now = yield* Clock.currentTimeMillis;
return Option.match(marker.value.mtime, {
onNone: () => false,
onSome: (writtenAt) =>
now - writtenAt.getTime() < Duration.toMillis(DESKTOP_UPDATE_RESTART_MARKER_TTL),
});
});

type ManagedTunnelRecoveryProofInput = {
readonly environmentId: RelayManagedEndpointRecoveryProofPayload["environmentId"];
readonly cloudUserId: string;
Expand Down
Loading
Loading