import * as Clock from "effect/Clock";
import * as Context from "effect/Context";
import * as Crypto from "effect/Crypto";
import * as DateTime from "effect/DateTime";
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 Stream from "effect/Stream";
import * as HttpClient from "effect/http/HttpClient";
import type * as HttpServerRequest from "effect/http/HttpServerRequest";
import * as HttpServerResponse from "effect/http/HttpServerResponse";
import * as HttpApiBuilder from "effect/http-api/HttpApiBuilder";
import { EnvironmentId } from "@t3tools/contracts";
import { RelayApi, type RelayHookDeliveryProofPayload } from "@t3tools/contracts/relay";
import {
normalizeRelayIssuer,
RELAY_HOOK_DELIVERY_HEADER,
RELAY_HOOK_DELIVERY_TYP,
signRelayJwt,
} from "@t3tools/shared/relayJwt";
import * as RelayConfiguration from "../Config.ts";
import {
MANAGED_ENDPOINT_KEY_PATTERN,
managedEndpointTunnelNameForKey,
} from "../deploymentConfig.ts";
import { validateManagedEndpoint } from "../environments/EnvironmentConnector.ts";
import * as EnvironmentLinks from "../environments/EnvironmentLinks.ts";
import * as ManagedEndpointAllocations from "../environments/ManagedEndpointAllocations.ts";
import * as HookInbox from "./HookInbox.ts";
import { sendUpstream, TUNNEL_OFFLINE_STATUS } from "./upstream.ts";
export const RELAY_HOOK_PATH_PREFIX = "/v1/hooks/";
export const RELAY_HOOK_MAX_BODY_BYTES = 1_048_576;
export const RELAY_HOOK_RATE_LIMIT = { limit: 60, periodSeconds: 60 } as const;
/**
* Hook budgets are per URL, and a sender who knows an endpoint key can mint
* new URLs for free, so every endpoint also has one overall budget.
*/
export const RELAY_HOOK_ENDPOINT_RATE_LIMIT = { limit: 600, periodSeconds: 60 } as const;
/**
* Upstream statuses that mean the environment did not take the request: the
* tunnel has no origin (530) or cloudflared cannot reach the local server
* (502, 503, 504) while it restarts.
*/
export const ENVIRONMENT_UNREACHABLE_STATUSES: ReadonlySet<number> = new Set([
502,
503,
504,
TUNNEL_OFFLINE_STATUS,
]);
const DROPPED_REQUEST_HEADERS = new Set([
"host",
"connection",
"keep-alive",
"transfer-encoding",
"te",
"upgrade",
"content-length",
"cookie",
"x-real-ip",
// Only the relay may set this; a sender could otherwise collide delivery ids.
"x-t3-relay-delivery-id",
"x-t3-relay-received-at",
RELAY_HOOK_DELIVERY_HEADER,
// The environment trusts trace context only from the relay, which sets its own.
"traceparent",
"tracestate",
"b3",
]);
const DROPPED_REQUEST_HEADER_PREFIXES = ["proxy-", "cf-", "x-forwarded-", "x-b3-"];
export const isRelayHookPath = (url: string): boolean => url.startsWith(RELAY_HOOK_PATH_PREFIX);
/** Replaces the hook token (and any query) so traces and logs never record the secret. */
export const redactRelayHookUrl = (url: string): string => {
const path = url.split("?", 1)[0] ?? url;
const segments = path.split("/");
if (segments.length >= 6) {
segments[5] = "<redacted>";
}
return segments.join("/");
};
/**
* Request budget for public hook forwarding, keyed by a hash of the hook URL
* (endpoint, hook and token). Requests with a wrong token get their own
* budget, so they cannot use up a real sender's; the environment rejects them.
* Built from decoded segments, because the environment decodes them too: two
* spellings of one token (`token`, `%74oken`) must share one budget.
*/
const hookBudgetKey = (hook: {
readonly endpointKey: string;
readonly hookId: string;
readonly token: string;
}) =>
Effect.promise(() =>
crypto.subtle.digest(
"SHA-256",
// Length-prefixed, so no segment contents can make two keys collide.
new TextEncoder().encode(
[hook.endpointKey, hook.hookId, hook.token]
.map((part) => `${part.length}:${part}`)
.join(""),
),
),
).pipe(
Effect.map((digest) =>
Array.from(new Uint8Array(digest), (byte) => byte.toString(16).padStart(2, "0")).join(""),
),
);
export class HookRateLimiter extends Context.Service<
HookRateLimiter,
{
/** One hook URL's budget, keyed by `hookBudgetKey`. */
readonly allowHook: (key: string) => Effect.Effect<boolean>;
/** One endpoint's overall budget, keyed by its endpoint key. */
readonly allowEndpoint: (endpointKey: string) => Effect.Effect<boolean>;
}
>()("t3code-relay/hooks/HookForwarder/HookRateLimiter") {}
export class HookForwarder extends Context.Service<
HookForwarder,
{
readonly handle: (
request: HttpServerRequest.HttpServerRequest,
) => Effect.Effect<HttpServerResponse.HttpServerResponse>;
}
>()("t3code-relay/hooks/HookForwarder") {}
class HookBodyTooLarge extends Schema.TaggedError<HookBodyTooLarge>()("HookBodyTooLarge", {}) {}
const errorResponse = (status: number, error: string, headers?: Record<string, string>) =>
HttpServerResponse.jsonUnsafe({ error }, { status, ...(headers ? { headers } : {}) });
const hookNotFound = () => errorResponse(404, "hook_not_found");
/** Longer than the inbox holds a request, so a held delivery's proof still verifies. */
const DELIVERY_PROOF_LIFETIME_SECONDS = 25 * 60 * 60;
const signDeliveryProof = (input: {
readonly settings: RelayConfiguration.RelayConfiguration["Service"];
readonly environmentId: string;
readonly deliveryId: string;
readonly receivedAt: string;
readonly hookId: string;
readonly jti: string;
}) =>
Effect.gen(function* () {
const now = Math.floor((yield* Clock.currentTimeMillis) / 1_000);
return yield* signRelayJwt({
privateKey: Redacted.value(input.settings.cloudMintPrivateKey),
typ: RELAY_HOOK_DELIVERY_TYP,
payload: {
iss: normalizeRelayIssuer(input.settings.relayIssuer),
aud: `t3-env:${input.environmentId}`,
sub: input.environmentId,
jti: input.jti,
iat: now,
exp: now + DELIVERY_PROOF_LIFETIME_SECONDS,
environmentId: EnvironmentId.make(input.environmentId),
deliveryId: input.deliveryId,
receivedAt: input.receivedAt,
hookId: input.hookId,
} satisfies RelayHookDeliveryProofPayload,
});
}).pipe(Effect.orDie);
/** Methods a webhook can arrive with; HEAD reaches the GET route and is refused. */
const FORWARDED_METHODS = new Set(["GET", "POST", "PUT", "PATCH"]);
/**
* Whatever the environment answers is served from the relay's own origin, so
* a body must never render or run there.
*/
const SANDBOXED_RESPONSE_HEADERS = {
"x-content-type-options": "nosniff",
"content-security-policy": "sandbox; default-src 'none'",
} as const;
function parseHookPath(url: string) {
const queryIndex = url.indexOf("?");
const path = queryIndex === -1 ? url : url.slice(0, queryIndex);
const search = queryIndex === -1 ? "" : url.slice(queryIndex);
const segments = path.split("/");
// ["", "v1", "hooks", endpointKey, hookId, token]
if (segments.length !== 6) return null;
const [, , , endpointKey, rawHookId, rawToken] = segments;
if (!endpointKey || !MANAGED_ENDPOINT_KEY_PATTERN.test(endpointKey) || !rawHookId || !rawToken) {
return null;
}
try {
return {
endpointKey,
hookId: decodeURIComponent(rawHookId),
token: decodeURIComponent(rawToken),
// Forward the encoded segments byte-for-byte; the environment decodes them.
rawHookId,
rawToken,
search,
};
} catch {
return null;
}
}
function forwardedHeaders(headers: Readonly<Record<string, string>>): Record<string, string> {
const result: Record<string, string> = {};
// Headers the sender names in Connection are hop-by-hop too (RFC 9110 7.6.1).
const connectionValue =
Object.entries(headers).find(([name]) => name.toLowerCase() === "connection")?.[1] ?? "";
const nominated = new Set(
connectionValue
.split(",")
.map((name) => name.trim().toLowerCase())
.filter(Boolean),
);
for (const name in headers) {
const lower = name.toLowerCase();
if (
DROPPED_REQUEST_HEADERS.has(lower) ||
nominated.has(lower) ||
DROPPED_REQUEST_HEADER_PREFIXES.some((prefix) => lower.startsWith(prefix))
) {
continue;
}
const value = headers[name];
if (value !== undefined) result[lower] = value;
}
return result;
}
const hasNoBody = (request: HttpServerRequest.HttpServerRequest) =>
request.source instanceof Request && request.source.body === null;
const readCappedBody = (request: HttpServerRequest.HttpServerRequest) =>
Effect.suspend(() => {
if (hasNoBody(request)) {
return Effect.succeed(new Uint8Array(0));
}
const chunks: Array<Uint8Array> = [];
let total = 0;
return request.stream.pipe(
Stream.runForEach((chunk) => {
total += chunk.length;
if (total > RELAY_HOOK_MAX_BODY_BYTES) {
return Effect.fail(new HookBodyTooLarge());
}
chunks.push(chunk);
return Effect.void;
}),
Effect.map(() => {
const body = new Uint8Array(total);
let offset = 0;
for (const chunk of chunks) {
body.set(chunk, offset);
offset += chunk.length;
}
return body;
}),
);
});
/**
* The ready managed endpoint a webhook URL's endpoint key names, with whether
* its link opted in to holding webhooks while offline. The key is the tunnel
* name's hash of user and environment, so it names exactly one allocation and
* at most one active link; nobody else can link their way onto it.
*/
export const resolveHookEndpoint = Effect.fn("relay.hooks.resolve_endpoint")(function* (
endpointKey: string,
) {
const links = yield* EnvironmentLinks.EnvironmentLinks;
const allocations = yield* ManagedEndpointAllocations.ManagedEndpointAllocations;
const settings = yield* RelayConfiguration.RelayConfiguration;
if (!settings.managedEndpointNamespace) return null;
const allocation = yield* allocations.getByTunnelName(
managedEndpointTunnelNameForKey(settings.managedEndpointNamespace, endpointKey),
);
if (allocation === null) return null;
const [link] = yield* links.findActiveManagedForEnvironment({
environmentId: allocation.environmentId,
userId: allocation.userId,
});
if (!link) return null;
const result = validateManagedEndpoint({
link,
allocation,
baseDomain: settings.managedEndpointBaseDomain,
});
if (Result.isFailure(result)) return null;
return {
...result.success,
environmentId: allocation.environmentId,
holdWhileOffline: link.holdWebhooksWhileOffline,
};
});
/** The endpoint key of an environment's own managed endpoint, for authenticated callers. */
export const endpointKeyForTunnelName = (namespace: string, tunnelName: string): string | null => {
const prefix = managedEndpointTunnelNameForKey(namespace, "");
const key = tunnelName.startsWith(prefix) ? tunnelName.slice(prefix.length) : "";
return MANAGED_ENDPOINT_KEY_PATTERN.test(key) ? key : null;
};
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 rateLimiter = yield* HookRateLimiter;
const inbox = yield* HookInbox.HookInbox;
const crypto = yield* Crypto.Crypto;
const handle = Effect.fn("relay.hooks.forward")(function* (
request: HttpServerRequest.HttpServerRequest,
) {
const outcome = (value: string) => Effect.annotateCurrentSpan({ "relay.hook.outcome": value });
if (!FORWARDED_METHODS.has(request.method)) {
yield* outcome("method_not_allowed");
return errorResponse(405, "method_not_allowed", { allow: "GET, POST, PUT, PATCH" });
}
// When the sender called, not when the environment failed to answer.
const receivedAt = DateTime.formatIso(yield* DateTime.now);
const parsed = parseHookPath(request.url);
if (!parsed) {
yield* outcome("invalid_path");
return hookNotFound();
}
yield* Effect.annotateCurrentSpan({
"relay.hook.endpoint_key": parsed.endpointKey,
"relay.hook_id": parsed.hookId,
});
// A coarse budget per endpoint first, so minting new hook ids or tokens
// cannot buy unlimited lookups and forwards, or fill the inbox.
const endpointAllowed = yield* rateLimiter.allowEndpoint(parsed.endpointKey);
const hookAllowed =
endpointAllowed && (yield* rateLimiter.allowHook(yield* hookBudgetKey(parsed)));
if (!endpointAllowed || !hookAllowed) {
// Which budget ran out: the whole endpoint's, or this one hook URL's.
yield* Effect.annotateCurrentSpan({
"relay.hook.rate_limit": endpointAllowed ? "hook" : "endpoint",
});
yield* outcome("rate_limited");
return errorResponse(429, "rate_limited", {
"retry-after": String(RELAY_HOOK_RATE_LIMIT.periodSeconds),
});
}
const declaredLength = Number(request.headers["content-length"] ?? "0");
if (Number.isFinite(declaredLength) && declaredLength > RELAY_HOOK_MAX_BODY_BYTES) {
yield* outcome("payload_too_large");
return errorResponse(413, "payload_too_large");
}
const endpoint = yield* resolveHookEndpoint(parsed.endpointKey).pipe(
Effect.provideService(EnvironmentLinks.EnvironmentLinks, links),
Effect.provideService(ManagedEndpointAllocations.ManagedEndpointAllocations, allocations),
Effect.provideService(RelayConfiguration.RelayConfiguration, settings),
Effect.catch((error) =>
Effect.logWarning("Failed to resolve hook endpoint", {
endpointKey: parsed.endpointKey,
errorTag: error._tag,
}).pipe(Effect.as(null)),
),
);
if (!endpoint) {
yield* outcome("not_found");
return hookNotFound();
}
yield* Effect.annotateCurrentSpan({
"relay.environment_id": endpoint.environmentId,
"relay.hook.hold_while_offline": endpoint.holdWhileOffline,
});
const body =
request.method === "GET"
? Result.succeed(new Uint8Array(0))
: yield* readCappedBody(request).pipe(Effect.result);
if (Result.isFailure(body)) {
if (body.failure._tag === "HookBodyTooLarge") {
yield* outcome("payload_too_large");
return errorResponse(413, "payload_too_large");
}
yield* outcome("invalid_body");
return errorResponse(400, "invalid_body");
}
// One id per request, so a request that reached the environment before a
// timeout and is later delivered from the inbox runs only once.
yield* Effect.annotateCurrentSpan({ "relay.hook.body_bytes": body.success.byteLength });
const deliveryId = yield* crypto.randomUUIDv4.pipe(Effect.orDie);
// Proves to the environment that this delivery id, receive time and
// trace context came from the relay. Signed once here and stored with a
// held request, so the inbox never needs the signing key.
const proof = yield* signDeliveryProof({
settings,
environmentId: endpoint.environmentId,
deliveryId,
receivedAt,
hookId: parsed.hookId,
jti: yield* crypto.randomUUIDv4.pipe(Effect.orDie),
});
const hook = {
id: deliveryId,
receivedAt,
method: request.method,
rawHookId: parsed.rawHookId,
rawToken: parsed.rawToken,
hookKey: parsed.hookId,
query: parsed.search.replace(/^\?/, ""),
headers: { ...forwardedHeaders(request.headers), [RELAY_HOOK_DELIVERY_HEADER]: proof },
body: body.success,
};
// Held only for environments that opted in; otherwise the relay is a plain proxy.
const holdOrFail = (status: 503 | 504, error: string) =>
Effect.gen(function* () {
if (!endpoint.holdWhileOffline) {
yield* outcome(error);
return errorResponse(status, error);
}
const stored = yield* inbox
.hold({
endpointKey: parsed.endpointKey,
baseUrl: endpoint.httpBaseUrl,
hook,
})
.pipe(
Effect.catch((cause) =>
Effect.logWarning("Could not hold webhook request", {
environmentId: endpoint.environmentId,
errorTag: cause._tag,
}).pipe(Effect.as(null)),
),
);
if (stored === null) {
yield* outcome(error);
return errorResponse(status, error);
}
if (!stored) {
// The inbox span carries which cap refused it (relay.inbox.refused).
yield* outcome("inbox_full");
return errorResponse(503, "inbox_full");
}
yield* outcome("held");
return HttpServerResponse.jsonUnsafe({ queued: true }, { status: 202 });
});
const upstream = yield* sendUpstream(endpoint.httpBaseUrl, hook).pipe(
Effect.provideService(HttpClient.HttpClient, httpClient),
Effect.result,
);
if (Result.isFailure(upstream)) {
yield* Effect.annotateCurrentSpan({ "relay.hook.upstream_error": upstream.failure._tag });
return yield* holdOrFail(503, "environment_unavailable");
}
if (Option.isNone(upstream.success)) {
return yield* holdOrFail(504, "environment_timeout");
}
const response = upstream.success.value;
if (ENVIRONMENT_UNREACHABLE_STATUSES.has(response.status)) {
yield* Effect.annotateCurrentSpan({ "relay.hook.upstream_status": response.status });
return yield* holdOrFail(503, "environment_unavailable");
}
yield* Effect.annotateCurrentSpan({
"relay.hook.outcome": "forwarded",
"relay.hook.upstream_status": response.status,
...(response.outcome === undefined
? {}
: { "relay.hook.upstream_outcome": response.outcome }),
});
// Only content-type is passed through: no location (redirects are never
// followed or relayed), no cookies, no upstream infrastructure headers.
const headers = {
...SANDBOXED_RESPONSE_HEADERS,
...(response.contentType ? { "content-type": response.contentType } : {}),
};
if (response.body.length === 0) {
return HttpServerResponse.empty({ status: response.status, headers });
}
return HttpServerResponse.uint8Array(response.body, {
status: response.status,
headers,
...(response.contentType ? { contentType: response.contentType } : {}),
});
});
return HookForwarder.of({ handle });
});
export const layer = Layer.effect(HookForwarder, make);
/**
* Implements the RelayApi `hooks` group. The endpoints are raw: the forwarder
* re-reads the encoded path segments from the request so the token and hook
* id reach the environment byte for byte, and streams the body itself.
*/
export const hooksApi = HttpApiBuilder.group(
RelayApi,
"hooks",
Effect.fnUntraced(function* (handlers) {
const forwarder = yield* HookForwarder;
const forward = ({ request }: { readonly request: HttpServerRequest.HttpServerRequest }) =>
forwarder.handle(request);
return handlers
.handleRaw("forwardPost", forward)
.handleRaw("forwardPut", forward)
.handleRaw("forwardPatch", forward)
.handleRaw("forwardGet", forward);
}),
);