import {
EnvironmentHttpBadRequestError,
EnvironmentHttpConflictError,
EnvironmentHttpForbiddenError,
EnvironmentHttpInternalServerError,
EnvironmentHttpUnauthorizedError,
} from "@t3tools/contracts";
import { makeEnvironmentHttpApiClient } from "@t3tools/client-runtime/rpc";
import {
RelayCloudEnvironmentHealthProofPayload,
RelayEnvironmentHealthResponse,
RelayEnvironmentHealthResponseProofPayload,
RelayEnvironmentMintResponse,
RelayEnvironmentMintResponseProofPayload,
RelayCloudMintCredentialProofPayload,
RelayEnvironmentConnectNotAuthorizedReason,
type RelayEnvironmentConnectResponse,
type RelayEnvironmentStatusResponse,
type RelayManagedEndpoint,
} from "@t3tools/contracts/relay";
import {
normalizeRelayIssuer,
RELAY_HEALTH_REQUEST_TYP,
RELAY_HEALTH_RESPONSE_TYP,
RELAY_MINT_REQUEST_TYP,
RELAY_MINT_RESPONSE_TYP,
signRelayJwt,
verifyRelayJwt,
} from "@t3tools/shared/relayJwt";
import { stableStringify } from "@t3tools/shared/relaySigning";
import * as Context from "effect/Context";
import * as Crypto from "effect/Crypto";
import * as DateTime from "effect/DateTime";
import * as Duration from "effect/Duration";
import * as Effect from "effect/Effect";
import * as Layer from "effect/Layer";
import * as Option from "effect/Option";
import * as Redacted from "effect/Redacted";
import * as Result from "effect/Result";
import * as Schema from "effect/Schema";
import * as FetchHttpClient from "effect/http/FetchHttpClient";
import * as HttpClient from "effect/http/HttpClient";
import * as HttpClientError from "effect/http/HttpClientError";
import * as EnvironmentLinks from "./EnvironmentLinks.ts";
import * as ManagedEndpointAllocations from "./ManagedEndpointAllocations.ts";
import * as RelayConfiguration from "../Config.ts";
import { isManagedEndpointHostname } from "../deploymentConfig.ts";
function environmentConnectNotAuthorizedReasonMessage(
reason: RelayEnvironmentConnectNotAuthorizedReason,
): string {
switch (reason) {
case "client_proof_key_thumbprint_missing":
return "the client proof key thumbprint is missing";
case "environment_link_not_found":
return "no active environment link was found";
case "endpoint_provider_not_managed":
return "the linked endpoint is not relay-managed";
case "managed_endpoint_allocation_not_found":
return "no managed endpoint allocation was found";
case "managed_endpoint_base_domain_not_configured":
return "the managed endpoint base domain is not configured";
case "managed_endpoint_allocation_not_ready":
return "the managed endpoint allocation is incomplete";
case "managed_endpoint_hostname_invalid":
return "the managed endpoint hostname is invalid";
case "managed_endpoint_mismatch":
return "the linked endpoint does not match its managed allocation";
}
}
export class EnvironmentConnectNotAuthorized extends Schema.TaggedError<EnvironmentConnectNotAuthorized>()(
"EnvironmentConnectNotAuthorized",
{
environmentId: Schema.String,
operation: Schema.Literals(["connect", "status"]),
reason: RelayEnvironmentConnectNotAuthorizedReason,
},
) {
override get message(): string {
return `Environment '${this.environmentId}' is not authorized for ${this.operation}: ${environmentConnectNotAuthorizedReasonMessage(this.reason)}`;
}
}
export class EnvironmentMintRequestFailed extends Schema.TaggedError<EnvironmentMintRequestFailed>()(
"EnvironmentMintRequestFailed",
{
environmentId: Schema.String,
operation: Schema.Literals(["connect", "status"]),
cause: Schema.Defect(),
},
) {
override get message(): string {
return `Environment '${this.environmentId}' ${this.operation} request failed`;
}
}
export class EnvironmentMintRequestTimedOut extends Schema.TaggedError<EnvironmentMintRequestTimedOut>()(
"EnvironmentMintRequestTimedOut",
{
environmentId: Schema.String,
timeoutMs: Schema.Number,
},
) {
override get message(): string {
return `Environment '${this.environmentId}' mint request timed out after ${this.timeoutMs}ms`;
}
}
export class EnvironmentMintResponseInvalid extends Schema.TaggedError<EnvironmentMintResponseInvalid>()(
"EnvironmentMintResponseInvalid",
{
environmentId: Schema.String,
operation: Schema.Literals(["connect", "status"]),
},
) {
override get message(): string {
return `Environment '${this.environmentId}' returned an invalid ${this.operation} response`;
}
}
export type EnvironmentConnectorError =
| EnvironmentConnectNotAuthorized
| EnvironmentMintRequestFailed
| EnvironmentMintRequestTimedOut
| EnvironmentMintResponseInvalid
| EnvironmentLinks.EnvironmentLinkLookupPersistenceError
| ManagedEndpointAllocations.ManagedEndpointAllocationPersistenceError;
export const ENVIRONMENT_MINT_REQUEST_TIMEOUT_MS = 10_000;
const ENVIRONMENT_HEALTH_CLOCK_SKEW_MILLIS = 60 * 1_000;
export class EnvironmentConnector extends Context.Service<
EnvironmentConnector,
{
readonly connect: (input: {
readonly userId: string;
readonly environmentId: string;
readonly clientProofKeyThumbprint: string;
readonly deviceId?: string;
}) => Effect.Effect<RelayEnvironmentConnectResponse, EnvironmentConnectorError>;
readonly status: (input: {
readonly userId: string;
readonly environmentId: string;
}) => Effect.Effect<RelayEnvironmentStatusResponse, EnvironmentConnectorError>;
}
>()("t3code-relay/environments/EnvironmentConnector") {}
const decodeMintResponseProof = Schema.decodeUnknownEffect(
RelayEnvironmentMintResponseProofPayload,
);
const decodeHealthResponseProof = Schema.decodeUnknownEffect(
RelayEnvironmentHealthResponseProofPayload,
);
const isEnvironmentHealthError = Schema.is(
Schema.Union([
EnvironmentHttpBadRequestError,
EnvironmentHttpUnauthorizedError,
EnvironmentHttpForbiddenError,
EnvironmentHttpConflictError,
EnvironmentHttpInternalServerError,
]),
);
/**
* Cleanup deleted this environment's tunnel and nothing has recorded a new
* one, so the host is unreachable until it updates and requests recovery.
*/
function offlineReasonFor(
allocation: ManagedEndpointAllocations.ManagedEndpointAllocation | null,
): { readonly offlineReason?: "tunnel_released" } {
return allocation?.tunnelReleasedAt ? { offlineReason: "tunnel_released" } : {};
}
function environmentHealthRequestFailureMessage(cause: unknown): string {
return isEnvironmentHealthError(cause)
? `Managed endpoint health request failed: ${cause.message}`
: "Managed endpoint health request failed.";
}
function environmentHealthRequestFailureReason(cause: unknown): string {
if (isEnvironmentHealthError(cause)) {
return cause._tag;
}
if (HttpClientError.isHttpClientError(cause)) {
return cause.reason._tag;
}
if (Schema.isSchemaError(cause)) {
return "SchemaError";
}
return cause instanceof Error && cause.name ? cause.name : "Unknown";
}
const currentTraceId = Effect.currentSpan.pipe(
Effect.map((span) => span.traceId),
Effect.orElseSucceed(() => "unavailable"),
);
export const withoutRedirects = <A, E, R>(effect: Effect.Effect<A, E, R>) =>
effect.pipe(Effect.provideService(FetchHttpClient.RequestInit, { redirect: "manual" }));
const verifyWithEnvironmentKeys = Effect.fnUntraced(function* <A, E>(input: {
readonly token: string;
readonly typ: string;
readonly issuer: string;
readonly audience: string;
readonly nowEpochSeconds: number;
readonly environmentPublicKeys: ReadonlyArray<string>;
readonly decodePayload: (input: unknown) => Effect.Effect<A, E>;
}) {
const { decodePayload, ...rest } = input;
for (const publicKey of input.environmentPublicKeys) {
const proof = yield* verifyRelayJwt({ ...rest, publicKey }).pipe(
Effect.flatMap(decodePayload),
Effect.option,
);
if (Option.isSome(proof)) {
return proof.value;
}
// A linked environment can have rotated keys; try the remaining active keys.
}
return null;
});
function verifyEnvironmentResponse(input: {
readonly response: RelayEnvironmentMintResponse;
readonly environmentId: string;
readonly requestNonce: string;
readonly clientProofKeyThumbprint: string;
readonly environmentPublicKeys: ReadonlyArray<string>;
readonly relayIssuer: string;
readonly nowEpochSeconds: number;
}) {
return verifyWithEnvironmentKeys({
token: input.response.proof,
typ: RELAY_MINT_RESPONSE_TYP,
issuer: `t3-env:${input.environmentId}`,
audience: normalizeRelayIssuer(input.relayIssuer),
nowEpochSeconds: input.nowEpochSeconds,
environmentPublicKeys: input.environmentPublicKeys,
decodePayload: decodeMintResponseProof,
}).pipe(
Effect.map(
(proof) =>
proof !== null &&
proof.environmentId === input.environmentId &&
proof.requestNonce === input.requestNonce &&
proof.clientProofKeyThumbprint === input.clientProofKeyThumbprint &&
proof.credential === input.response.credential &&
Option.match(DateTime.make(input.response.expiresAt), {
onNone: () => false,
onSome: (expiresAt) => Math.floor(expiresAt.epochMilliseconds / 1_000) === proof.exp,
}),
),
);
}
function verifyEnvironmentHealthResponse(input: {
readonly response: RelayEnvironmentHealthResponse;
readonly environmentId: string;
readonly requestNonce: string;
readonly requestIssuedAt: DateTime.DateTime;
readonly environmentPublicKeys: ReadonlyArray<string>;
readonly relayIssuer: string;
readonly now: DateTime.DateTime;
}) {
return verifyWithEnvironmentKeys({
token: input.response.proof,
typ: RELAY_HEALTH_RESPONSE_TYP,
issuer: `t3-env:${input.environmentId}`,
audience: normalizeRelayIssuer(input.relayIssuer),
nowEpochSeconds: Math.floor(input.now.epochMilliseconds / 1_000),
environmentPublicKeys: input.environmentPublicKeys,
decodePayload: decodeHealthResponseProof,
}).pipe(
Effect.map((proof) => {
if (
proof === null ||
input.response.environmentId !== input.environmentId ||
proof.environmentId !== input.environmentId ||
proof.requestNonce !== input.requestNonce ||
proof.status !== input.response.status ||
proof.checkedAt !== input.response.checkedAt ||
stableStringify(proof.descriptor) !== stableStringify(input.response.descriptor)
) {
return false;
}
const checkedAt = DateTime.make(input.response.checkedAt);
if (Option.isNone(checkedAt)) {
return false;
}
return (
checkedAt.value.epochMilliseconds >=
input.requestIssuedAt.epochMilliseconds - ENVIRONMENT_HEALTH_CLOCK_SKEW_MILLIS &&
checkedAt.value.epochMilliseconds <=
input.now.epochMilliseconds + ENVIRONMENT_HEALTH_CLOCK_SKEW_MILLIS
);
}),
);
}
export interface ManagedEndpointValidationFailure {
readonly reason: Exclude<
RelayEnvironmentConnectNotAuthorizedReason,
"client_proof_key_thumbprint_missing" | "environment_link_not_found"
>;
/** Diagnostic span attributes; never contains secrets. */
readonly attributes?: Record<string, string | boolean>;
}
/**
* Checks that a link's relay-managed endpoint is backed by a ready allocation
* on the configured base domain and still matches what the link recorded.
*/
export function validateManagedEndpoint(input: {
readonly link: EnvironmentLinks.RelayLinkedEnvironmentRecord;
readonly allocation: ManagedEndpointAllocations.ManagedEndpointAllocation | null;
readonly baseDomain: string | undefined;
}): Result.Result<RelayManagedEndpoint, ManagedEndpointValidationFailure> {
const { link, allocation, baseDomain } = input;
if (link.endpoint.providerKind !== "cloudflare_tunnel") {
return Result.fail({
reason: "endpoint_provider_not_managed",
attributes: { "relay.authorization.endpoint_provider_kind": link.endpoint.providerKind },
});
}
if (!allocation) {
return Result.fail({ reason: "managed_endpoint_allocation_not_found" });
}
const allocationAttributes = {
"relay.authorization.allocation_hostname": allocation.hostname,
"relay.authorization.allocation_has_ready_at": allocation.readyAt !== null,
"relay.authorization.allocation_has_tunnel_id": allocation.tunnelId !== null,
"relay.authorization.allocation_has_dns_record_id": allocation.dnsRecordId !== null,
};
if (!baseDomain) {
return Result.fail({
reason: "managed_endpoint_base_domain_not_configured",
attributes: allocationAttributes,
});
}
if (
allocation.readyAt === null ||
allocation.tunnelId === null ||
allocation.dnsRecordId === null
) {
return Result.fail({
reason: "managed_endpoint_allocation_not_ready",
attributes: allocationAttributes,
});
}
if (!isManagedEndpointHostname(allocation.hostname, baseDomain)) {
return Result.fail({
reason: "managed_endpoint_hostname_invalid",
attributes: {
...allocationAttributes,
"relay.authorization.managed_endpoint_base_domain": baseDomain,
},
});
}
const endpoint = ManagedEndpointAllocations.resolveReadyManagedEndpoint({
allocation,
baseDomain,
});
if (
endpoint === null ||
endpoint.httpBaseUrl !== link.endpoint.httpBaseUrl ||
endpoint.wsBaseUrl !== link.endpoint.wsBaseUrl
) {
return Result.fail({
reason: "managed_endpoint_mismatch",
attributes: {
...allocationAttributes,
"relay.authorization.linked_http_base_url": link.endpoint.httpBaseUrl,
"relay.authorization.linked_ws_base_url": link.endpoint.wsBaseUrl,
...(endpoint
? {
"relay.authorization.resolved_http_base_url": endpoint.httpBaseUrl,
"relay.authorization.resolved_ws_base_url": endpoint.wsBaseUrl,
}
: {}),
},
});
}
return Result.succeed(endpoint);
}
const make = Effect.gen(function* () {
const links = yield* EnvironmentLinks.EnvironmentLinks;
const allocations = yield* ManagedEndpointAllocations.ManagedEndpointAllocations;
const settings = yield* RelayConfiguration.RelayConfiguration;
const httpClient = yield* HttpClient.HttpClient;
const crypto = yield* Crypto.Crypto;
const relayIssuer = normalizeRelayIssuer(settings.relayIssuer);
const makeEnvironmentClient = (httpBaseUrl: string) =>
makeEnvironmentHttpApiClient(httpBaseUrl).pipe(
Effect.provideService(HttpClient.HttpClient, httpClient),
);
const resolveManagedEndpoint = Effect.fn("relay.environment_connector.resolve_managed_endpoint")(
function* (input: {
readonly operation: "connect" | "status";
readonly link: EnvironmentLinks.RelayLinkedEnvironmentRecord;
readonly allocation: ManagedEndpointAllocations.ManagedEndpointAllocation | null;
}) {
const result = validateManagedEndpoint({
link: input.link,
allocation: input.allocation,
baseDomain: settings.managedEndpointBaseDomain,
});
if (Result.isSuccess(result)) {
return result.success;
}
if (result.failure.attributes) {
yield* Effect.annotateCurrentSpan(result.failure.attributes);
}
return yield* new EnvironmentConnectNotAuthorized({
environmentId: input.link.environmentId,
operation: input.operation,
reason: result.failure.reason,
});
},
);
return EnvironmentConnector.of({
status: Effect.fn("relay.environment_connector.status")(function* (input) {
yield* Effect.annotateCurrentSpan({
"relay.environment_id": input.environmentId,
"relay.operation": "status",
});
const { link, allocation } = yield* Effect.all(
{
link: links.getForUser(input),
allocation: allocations.get(input),
},
{ concurrency: 2 },
);
if (!link) {
return yield* new EnvironmentConnectNotAuthorized({
environmentId: input.environmentId,
operation: "status",
reason: "environment_link_not_found",
});
}
const endpoint = yield* resolveManagedEndpoint({
operation: "status",
link,
allocation,
});
const now = yield* DateTime.now;
const expiresAt = DateTime.add(now, { minutes: 2 });
const nonce = yield* crypto.randomUUIDv4.pipe(
Effect.mapError(
(cause) =>
new EnvironmentMintRequestFailed({
environmentId: input.environmentId,
operation: "status",
cause,
}),
),
);
const payload = {
iss: relayIssuer,
aud: `t3-env:${link.environmentId}`,
sub: input.userId,
jti: yield* crypto.randomUUIDv4.pipe(
Effect.mapError(
(cause) =>
new EnvironmentMintRequestFailed({
environmentId: input.environmentId,
operation: "status",
cause,
}),
),
),
iat: Math.floor(now.epochMilliseconds / 1_000),
exp: Math.floor(expiresAt.epochMilliseconds / 1_000),
environmentId: link.environmentId,
nonce,
scope: ["environment:status"],
} satisfies RelayCloudEnvironmentHealthProofPayload;
const proof = yield* signRelayJwt({
privateKey: Redacted.value(settings.cloudMintPrivateKey),
typ: RELAY_HEALTH_REQUEST_TYP,
payload,
}).pipe(
Effect.mapError(
(cause) =>
new EnvironmentMintRequestFailed({
environmentId: input.environmentId,
operation: "status",
cause,
}),
),
);
const checkedAt = DateTime.formatIso(now);
const traceId = yield* currentTraceId;
const environmentClient = yield* makeEnvironmentClient(endpoint.httpBaseUrl);
const responseOption = yield* environmentClient.connect.health({ payload: { proof } }).pipe(
withoutRedirects,
Effect.match({
onFailure: (cause) => ({ _tag: "Failure" as const, cause }),
onSuccess: (response) => ({ _tag: "Success" as const, response }),
}),
Effect.timeoutOption(Duration.millis(ENVIRONMENT_MINT_REQUEST_TIMEOUT_MS)),
);
if (Option.isNone(responseOption)) {
yield* Effect.annotateCurrentSpan({
"relay.environment_health.outcome": "timeout",
"relay.environment_health.trace_id": traceId,
});
yield* Effect.logWarning("Managed endpoint health request timed out", {
environmentId: link.environmentId,
endpoint: endpoint.httpBaseUrl,
traceId,
});
return {
environmentId: link.environmentId,
endpoint,
status: "offline" as const,
checkedAt,
error: "Managed endpoint health request timed out.",
...offlineReasonFor(allocation),
traceId,
};
}
if (responseOption.value._tag === "Failure") {
const failureReason = environmentHealthRequestFailureReason(responseOption.value.cause);
yield* Effect.annotateCurrentSpan({
"relay.environment_health.outcome": "failure",
"relay.environment_health.failure_reason": failureReason,
"relay.environment_health.trace_id": traceId,
});
yield* Effect.logWarning("Managed endpoint health request failed", {
environmentId: link.environmentId,
endpoint: endpoint.httpBaseUrl,
failureReason,
traceId,
});
return {
environmentId: link.environmentId,
endpoint,
status: "offline" as const,
checkedAt,
error: environmentHealthRequestFailureMessage(responseOption.value.cause),
...offlineReasonFor(allocation),
traceId,
};
}
const decoded = responseOption.value.response;
const verified = yield* verifyEnvironmentHealthResponse({
response: decoded,
environmentId: input.environmentId,
requestNonce: nonce,
requestIssuedAt: now,
environmentPublicKeys: [link.environmentPublicKey],
relayIssuer,
now: yield* DateTime.now,
});
if (!verified) {
return yield* new EnvironmentMintResponseInvalid({
environmentId: input.environmentId,
operation: "status",
});
}
return {
environmentId: link.environmentId,
endpoint,
status: "online" as const,
checkedAt: decoded.checkedAt,
descriptor: decoded.descriptor,
};
}),
connect: Effect.fn("relay.environment_connector.connect")(function* (input) {
yield* Effect.annotateCurrentSpan({
"relay.environment_id": input.environmentId,
"relay.operation": "connect",
"relay.connect.has_device_id": input.deviceId !== undefined,
...(input.deviceId ? { "relay.mobile.device_id": input.deviceId } : {}),
});
if (input.clientProofKeyThumbprint.trim().length === 0) {
return yield* new EnvironmentConnectNotAuthorized({
environmentId: input.environmentId,
operation: "connect",
reason: "client_proof_key_thumbprint_missing",
});
}
const { link, allocation } = yield* Effect.all(
{
link: links.getForUser(input),
allocation: allocations.get(input),
},
{ concurrency: 2 },
);
if (!link) {
return yield* new EnvironmentConnectNotAuthorized({
environmentId: input.environmentId,
operation: "connect",
reason: "environment_link_not_found",
});
}
const endpoint = yield* resolveManagedEndpoint({
operation: "connect",
link,
allocation,
});
const now = yield* DateTime.now;
const expiresAt = DateTime.add(now, { minutes: 2 });
const nonce = yield* crypto.randomUUIDv4.pipe(
Effect.mapError(
(cause) =>
new EnvironmentMintRequestFailed({
environmentId: input.environmentId,
operation: "connect",
cause,
}),
),
);
const payload = {
iss: relayIssuer,
aud: `t3-env:${link.environmentId}`,
sub: input.userId,
jti: yield* crypto.randomUUIDv4.pipe(
Effect.mapError(
(cause) =>
new EnvironmentMintRequestFailed({
environmentId: input.environmentId,
operation: "connect",
cause,
}),
),
),
iat: Math.floor(now.epochMilliseconds / 1_000),
exp: Math.floor(expiresAt.epochMilliseconds / 1_000),
environmentId: link.environmentId,
clientProofKeyThumbprint: input.clientProofKeyThumbprint,
cnf: { jkt: input.clientProofKeyThumbprint },
...(input.deviceId ? { deviceId: input.deviceId } : {}),
nonce,
scope: ["environment:connect"],
} satisfies RelayCloudMintCredentialProofPayload;
const proof = yield* signRelayJwt({
privateKey: Redacted.value(settings.cloudMintPrivateKey),
typ: RELAY_MINT_REQUEST_TYP,
payload,
}).pipe(
Effect.mapError(
(cause) =>
new EnvironmentMintRequestFailed({
environmentId: input.environmentId,
operation: "connect",
cause,
}),
),
);
const environmentClient = yield* makeEnvironmentClient(endpoint.httpBaseUrl);
const decoded = yield* environmentClient.connect
.t3MintCredential({ payload: { proof } })
.pipe(
withoutRedirects,
Effect.mapError(
(cause) =>
new EnvironmentMintRequestFailed({
environmentId: input.environmentId,
operation: "connect",
cause,
}),
),
Effect.timeoutOption(Duration.millis(ENVIRONMENT_MINT_REQUEST_TIMEOUT_MS)),
Effect.flatMap(
Option.match({
onNone: () =>
Effect.fail(
new EnvironmentMintRequestTimedOut({
environmentId: input.environmentId,
timeoutMs: ENVIRONMENT_MINT_REQUEST_TIMEOUT_MS,
}),
),
onSome: Effect.succeed,
}),
),
);
const verified = yield* verifyEnvironmentResponse({
response: decoded,
environmentId: input.environmentId,
requestNonce: nonce,
clientProofKeyThumbprint: input.clientProofKeyThumbprint,
environmentPublicKeys: [link.environmentPublicKey],
relayIssuer,
nowEpochSeconds: Math.floor(now.epochMilliseconds / 1_000),
});
if (!verified) {
return yield* new EnvironmentMintResponseInvalid({
environmentId: input.environmentId,
operation: "connect",
});
}
return {
environmentId: link.environmentId,
endpoint,
credential: decoded.credential,
expiresAt: decoded.expiresAt,
};
}),
});
});
export const layer = Layer.effect(EnvironmentConnector, make);