diff --git a/apps/server/src/cloud/CloudPreferences.test.ts b/apps/server/src/cloud/CloudPreferences.test.ts new file mode 100644 index 000000000000..c0e66bc49f75 --- /dev/null +++ b/apps/server/src/cloud/CloudPreferences.test.ts @@ -0,0 +1,219 @@ +import { EnvironmentId } from "@t3tools/contracts"; +import { assert, it } from "@effect/vitest"; +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 FetchHttpClient from "effect/http/FetchHttpClient"; + +import * as ServerSecretStore from "../auth/ServerSecretStore.ts"; +import * as ServerEnvironment from "../environment/ServerEnvironment.ts"; +import * as AgentAwarenessRelay from "../relay/AgentAwarenessRelay.ts"; +import * as CloudPreferences from "./CloudPreferences.ts"; +import { + HOLD_WEBHOOKS_WHILE_OFFLINE_SECRET, + PUBLISH_AGENT_ACTIVITY_SECRET, + RELAY_ENVIRONMENT_CREDENTIAL_SECRET, + RELAY_URL_SECRET, +} from "./config.ts"; + +const encode = (value: string) => new TextEncoder().encode(value); + +/** A linked environment whose secret store can refuse writes, and the relay calls it made. */ +const withService = ( + options: { + readonly failHoldWrite?: boolean; + readonly failActivityWrite?: boolean; + readonly relayFails?: boolean; + readonly failActivityRead?: boolean; + readonly failHoldRead?: boolean; + /** The first relay call reports itself, then waits for this before answering. */ + readonly holdFirstRelayCall?: { + readonly started: () => void; + readonly released: Promise; + }; + }, + body: (input: { + readonly preferences: CloudPreferences.CloudPreferences["Service"]; + readonly stored: Map; + readonly relayCalls: Array; + }) => Effect.Effect, +) => + Effect.gen(function* () { + const stored = new Map([ + [RELAY_URL_SECRET, encode("https://relay.test")], + [RELAY_ENVIRONMENT_CREDENTIAL_SECRET, encode("credential")], + [HOLD_WEBHOOKS_WHILE_OFFLINE_SECRET, encode("false")], + [PUBLISH_AGENT_ACTIVITY_SECRET, encode("false")], + ]); + const relayCalls: Array = []; + const fetch: typeof globalThis.fetch = Object.assign( + (_input: Parameters[0], init?: RequestInit) => { + const payload = JSON.parse(String(init?.body)) as { holdWebhooksWhileOffline: boolean }; + relayCalls.push(payload.holdWebhooksWhileOffline); + if (relayCalls.length === 1 && options.holdFirstRelayCall) { + options.holdFirstRelayCall.started(); + return options.holdFirstRelayCall.released.then(() => Response.json(payload)); + } + return options.relayFails + ? Promise.resolve(new Response("unavailable", { status: 503 })) + : Promise.resolve(Response.json(payload)); + }, + { preconnect: () => {} }, + ); + const dependencies = Layer.mergeAll( + Layer.mock(ServerSecretStore.ServerSecretStore)({ + get: (name) => + (options.failActivityRead && name === PUBLISH_AGENT_ACTIVITY_SECRET) || + (options.failHoldRead && name === HOLD_WEBHOOKS_WHILE_OFFLINE_SECRET) + ? Effect.fail( + new ServerSecretStore.SecretStoreReadError({ + resource: name, + cause: new Error("busy"), + }), + ) + : Effect.succeed(Option.fromNullishOr(stored.get(name))), + set: (name, value) => + (options.failHoldWrite && name === HOLD_WEBHOOKS_WHILE_OFFLINE_SECRET) || + (options.failActivityWrite && name === PUBLISH_AGENT_ACTIVITY_SECRET) + ? Effect.fail( + new ServerSecretStore.SecretStorePersistError({ + resource: name, + cause: new Error("read-only"), + }), + ) + : Effect.sync(() => void stored.set(name, value)), + }), + Layer.mock(ServerEnvironment.ServerEnvironment)({ + getEnvironmentId: Effect.succeed(EnvironmentId.make("environment-1")), + }), + Layer.mock(AgentAwarenessRelay.AgentAwarenessRelay)({ requestCatchUp: () => Effect.void }), + ); + return yield* Effect.gen(function* () { + const preferences = yield* CloudPreferences.CloudPreferences; + return yield* body({ preferences, stored, relayCalls }); + }).pipe( + Effect.provide(CloudPreferences.layer.pipe(Layer.provide(dependencies))), + Effect.provideService(FetchHttpClient.Fetch, fetch), + ); + }); + +it.effect("tells the relay before saving the hold setting locally", () => + withService({}, ({ preferences, stored, relayCalls }) => + Effect.gen(function* () { + yield* preferences.update({ publishAgentActivity: true, holdWebhooksWhileOffline: true }); + assert.deepEqual(relayCalls, [true]); + assert.equal( + new TextDecoder().decode(stored.get(HOLD_WEBHOOKS_WHILE_OFFLINE_SECRET)), + "true", + ); + }), + ), +); + +it.effect("puts the relay back when the local save fails", () => + withService({ failHoldWrite: true }, ({ preferences, stored, relayCalls }) => + Effect.gen(function* () { + const error = yield* preferences + .update({ publishAgentActivity: true, holdWebhooksWhileOffline: true }) + .pipe(Effect.flip); + assert.equal(error._tag, "EnvironmentHttpInternalServerError"); + assert.deepEqual(relayCalls, [true, false]); + assert.equal( + new TextDecoder().decode(stored.get(HOLD_WEBHOOKS_WHILE_OFFLINE_SECRET)), + "false", + ); + }), + ), +); + +it.effect("leaves the relay untouched when the activity setting can't be saved", () => + withService({ failActivityWrite: true }, ({ preferences, stored, relayCalls }) => + Effect.gen(function* () { + const error = yield* preferences + .update({ publishAgentActivity: true, holdWebhooksWhileOffline: true }) + .pipe(Effect.flip); + assert.equal(error._tag, "EnvironmentHttpInternalServerError"); + assert.deepEqual(relayCalls, []); + assert.equal( + new TextDecoder().decode(stored.get(HOLD_WEBHOOKS_WHILE_OFFLINE_SECRET)), + "false", + ); + }), + ), +); + +it.effect("keeps both settings unchanged when the relay refuses the hold change", () => + withService({ relayFails: true }, ({ preferences, stored }) => + Effect.gen(function* () { + yield* preferences + .update({ publishAgentActivity: true, holdWebhooksWhileOffline: true }) + .pipe(Effect.flip); + assert.equal(new TextDecoder().decode(stored.get(PUBLISH_AGENT_ACTIVITY_SECRET)), "false"); + assert.equal( + new TextDecoder().decode(stored.get(HOLD_WEBHOOKS_WHILE_OFFLINE_SECRET)), + "false", + ); + }), + ), +); + +it.effect("changes nothing when the current activity setting can't be read", () => + withService({ failActivityRead: true, relayFails: true }, ({ preferences, stored, relayCalls }) => + Effect.gen(function* () { + const error = yield* preferences + .update({ publishAgentActivity: true, holdWebhooksWhileOffline: true }) + .pipe(Effect.flip); + assert.equal(error._tag, "EnvironmentHttpInternalServerError"); + assert.deepEqual(relayCalls, []); + assert.equal(new TextDecoder().decode(stored.get(PUBLISH_AGENT_ACTIVITY_SECRET)), "false"); + }), + ), +); + +it.effect("two overlapping updates are applied one after the other", () => { + let release!: () => void; + let started!: () => void; + const released = new Promise((resolve) => { + release = resolve; + }); + const firstAtRelay = new Promise((resolve) => { + started = resolve; + }); + return withService({ holdFirstRelayCall: { started, released } }, ({ preferences, stored }) => + Effect.gen(function* () { + const first = yield* preferences + .update({ publishAgentActivity: true, holdWebhooksWhileOffline: true }) + .pipe(Effect.forkChild); + yield* Effect.promise(() => firstAtRelay); + const second = yield* preferences + .update({ publishAgentActivity: false, holdWebhooksWhileOffline: false }) + .pipe(Effect.forkChild); + // Give the second update every chance to run; it must wait for the first. + yield* Effect.yieldNow; + assert.equal(new TextDecoder().decode(stored.get(PUBLISH_AGENT_ACTIVITY_SECRET)), "true"); + release(); + yield* Fiber.join(first); + yield* Fiber.join(second); + assert.equal(new TextDecoder().decode(stored.get(PUBLISH_AGENT_ACTIVITY_SECRET)), "false"); + assert.equal( + new TextDecoder().decode(stored.get(HOLD_WEBHOOKS_WHILE_OFFLINE_SECRET)), + "false", + ); + }), + ); +}); + +it.effect("changes nothing when the current hold setting can't be read", () => + withService({ failHoldRead: true, failHoldWrite: true }, ({ preferences, stored, relayCalls }) => + Effect.gen(function* () { + stored.set(HOLD_WEBHOOKS_WHILE_OFFLINE_SECRET, encode("true")); + const error = yield* preferences + .update({ publishAgentActivity: true, holdWebhooksWhileOffline: false }) + .pipe(Effect.flip); + assert.equal(error._tag, "EnvironmentHttpInternalServerError"); + assert.deepEqual(relayCalls, []); + assert.equal(new TextDecoder().decode(stored.get(PUBLISH_AGENT_ACTIVITY_SECRET)), "false"); + }), + ), +); diff --git a/apps/server/src/cloud/CloudPreferences.ts b/apps/server/src/cloud/CloudPreferences.ts new file mode 100644 index 000000000000..7e765ce76959 --- /dev/null +++ b/apps/server/src/cloud/CloudPreferences.ts @@ -0,0 +1,130 @@ +import { + type EnvironmentCloudPreferencesRequest, + EnvironmentHttpBadRequestError, + EnvironmentHttpInternalServerError, +} from "@t3tools/contracts"; +import * as Context from "effect/Context"; +import * as Effect from "effect/Effect"; +import * as Layer from "effect/Layer"; +import * as Option from "effect/Option"; +import * as Semaphore from "effect/Semaphore"; + +import * as ServerSecretStore from "../auth/ServerSecretStore.ts"; +import * as ServerEnvironment from "../environment/ServerEnvironment.ts"; +import * as AgentAwarenessRelay from "../relay/AgentAwarenessRelay.ts"; +import { makeRelayEnvironmentClient } from "../relay/relayEnvironmentClient.ts"; +import { + HOLD_WEBHOOKS_WHILE_OFFLINE_SECRET, + PUBLISH_AGENT_ACTIVITY_SECRET, + readRelayConnection, +} from "./config.ts"; + +const encode = (value: boolean) => new TextEncoder().encode(String(value)); + +/** A failed rollback leaves a setting changed; it is logged, not hidden. */ +const rollbackFailed = (cause: unknown) => + Effect.logWarning("Could not roll back a T3 Connect preference", { cause }); + +const internalError = (message: string) => (cause: unknown) => + Effect.logError(message, { cause }).pipe( + Effect.andThen(Effect.fail(new EnvironmentHttpInternalServerError({ message }))), + ); + +export class CloudPreferences extends Context.Service< + CloudPreferences, + { + /** + * Saves this environment's T3 Connect preferences, all or nothing. The + * activity setting is saved first. Holding webhooks while offline is decided + * by the relay, so the relay is told before the local copy is saved. If + * either step fails, the activity setting is put back, and so is the relay + * when only the local save failed. + */ + readonly update: ( + input: EnvironmentCloudPreferencesRequest, + ) => Effect.Effect; + } +>()("t3/cloud/CloudPreferences") {} + +const make = Effect.gen(function* () { + const secrets = yield* ServerSecretStore.ServerSecretStore; + const environment = yield* ServerEnvironment.ServerEnvironment; + const awarenessRelay = yield* AgentAwarenessRelay.AgentAwarenessRelay; + const withSecrets = Effect.provideService(ServerSecretStore.ServerSecretStore, secrets); + + const pushHoldWebhooksWhileOffline = Effect.fn("CloudPreferences.pushHoldWebhooksWhileOffline")( + function* (holdWebhooksWhileOffline: boolean) { + const connection = yield* readRelayConnection.pipe(withSecrets); + if (connection === null) { + return yield* new EnvironmentHttpBadRequestError({ + message: "Link this environment to T3 Connect first.", + }); + } + const environmentId = yield* environment.getEnvironmentId; + const client = yield* makeRelayEnvironmentClient(connection); + yield* client.server + .updateLinkPreferences({ + params: { environmentId }, + payload: { holdWebhooksWhileOffline }, + }) + .pipe( + Effect.timeout("10 seconds"), + Effect.catch(internalError("Could not update T3 Connect webhook settings.")), + ); + }, + ); + + const save = (name: string, value: boolean) => + secrets + .set(name, encode(value)) + .pipe(Effect.catch(internalError("Could not persist environment cloud preferences."))); + + // One update at a time, so two requests can't each leave one setting behind. + const updateLock = yield* Semaphore.make(1); + + const update: CloudPreferences["Service"]["update"] = Effect.fn("CloudPreferences.update")( + function* (input) { + // All or nothing: the activity setting is saved first, before the relay + // is told anything, and put back if the hold change then fails. Both + // current values are read up front; a failed read stops here, before + // anything changes, because a guessed value would be the rollback target. + const readCurrent = (name: string) => + secrets + .get(name) + .pipe(Effect.catch(internalError("Could not read environment cloud preferences."))); + const previousActivity = yield* readCurrent(PUBLISH_AGENT_ACTIVITY_SECRET); + const previousHold = yield* readCurrent(HOLD_WEBHOOKS_WHILE_OFFLINE_SECRET); + yield* save(PUBLISH_AGENT_ACTIVITY_SECRET, input.publishAgentActivity); + if (input.holdWebhooksWhileOffline !== undefined) { + const next = input.holdWebhooksWhileOffline; + const previous = Option.match(previousHold, { + onNone: () => false, + onSome: (bytes) => new TextDecoder().decode(bytes) === "true", + }); + yield* pushHoldWebhooksWhileOffline(next).pipe( + Effect.andThen( + save(HOLD_WEBHOOKS_WHILE_OFFLINE_SECRET, next).pipe( + Effect.tapError(() => + previous === next + ? Effect.void + : pushHoldWebhooksWhileOffline(previous).pipe(Effect.catch(rollbackFailed)), + ), + ), + ), + Effect.tapError(() => + Option.match(previousActivity, { + onNone: () => secrets.remove(PUBLISH_AGENT_ACTIVITY_SECRET), + onSome: (bytes) => secrets.set(PUBLISH_AGENT_ACTIVITY_SECRET, bytes), + }).pipe(Effect.catch(rollbackFailed)), + ), + ); + } + yield* awarenessRelay.requestCatchUp(); + }, + updateLock.withPermits(1), + ); + + return CloudPreferences.of({ update }); +}); + +export const layer = Layer.effect(CloudPreferences, make); diff --git a/apps/server/src/cloud/config.ts b/apps/server/src/cloud/config.ts index 92bd9c856ca6..910f0820282f 100644 --- a/apps/server/src/cloud/config.ts +++ b/apps/server/src/cloud/config.ts @@ -6,7 +6,7 @@ import * as Effect from "effect/Effect"; import * as Option from "effect/Option"; import * as Schema from "effect/Schema"; -import type * as ServerSecretStore from "../auth/ServerSecretStore.ts"; +import * as ServerSecretStore from "../auth/ServerSecretStore.ts"; export const CLOUD_MINT_PUBLIC_KEY = "cloud-mint-ed25519-public-key"; export const CLOUD_ENDPOINT_RUNTIME_CONFIG = "cloud-endpoint-runtime-config"; @@ -75,11 +75,10 @@ export const readAgentActivityPublishingActive = ( ); }).pipe(Effect.orElseSucceed(() => false)); -const readSecretString = ( - secrets: ServerSecretStore.ServerSecretStore["Service"], - name: string, -): Effect.Effect => - secrets.get(name).pipe( +/** A non-empty secret as text, or null when it is missing, empty or unreadable. */ +const readSecretString = (name: string) => + ServerSecretStore.ServerSecretStore.pipe( + Effect.flatMap((secrets) => secrets.get(name)), Effect.map((bytes) => Option.isSome(bytes) && bytes.value.length > 0 ? new TextDecoder().decode(bytes.value) : null, ), @@ -87,20 +86,16 @@ const readSecretString = ( ); /** The relay URL and environment credential, or null when not linked to T3 Connect. */ -export const readRelayConnection = (secrets: ServerSecretStore.ServerSecretStore["Service"]) => - Effect.all([ - readSecretString(secrets, RELAY_URL_SECRET), - readSecretString(secrets, RELAY_ENVIRONMENT_CREDENTIAL_SECRET), - ]).pipe( - Effect.map(([url, environmentCredential]) => - url && environmentCredential ? { url, environmentCredential } : null, - ), - ); +export const readRelayConnection = Effect.all([ + readSecretString(RELAY_URL_SECRET), + readSecretString(RELAY_ENVIRONMENT_CREDENTIAL_SECRET), +]).pipe( + Effect.map(([url, environmentCredential]) => + url && environmentCredential ? { url, environmentCredential } : null, + ), +); /** Whether this environment opted in to T3 Connect holding webhooks while it is offline. */ -export const readHoldWebhooksWhileOffline = ( - secrets: ServerSecretStore.ServerSecretStore["Service"], -) => - readSecretString(secrets, HOLD_WEBHOOKS_WHILE_OFFLINE_SECRET).pipe( - Effect.map((value) => value === "true"), - ); +export const readHoldWebhooksWhileOffline = readSecretString( + HOLD_WEBHOOKS_WHILE_OFFLINE_SECRET, +).pipe(Effect.map((value) => value === "true")); diff --git a/apps/server/src/cloud/http.ts b/apps/server/src/cloud/http.ts index 4cb29c37b67f..651caa758bd4 100644 --- a/apps/server/src/cloud/http.ts +++ b/apps/server/src/cloud/http.ts @@ -8,7 +8,6 @@ import { EnvironmentCloudRelayConfigResult, EnvironmentHttpApi, EnvironmentHttpBadRequestError, - type EnvironmentCloudPreferencesRequest, EnvironmentHttpConflictError, EnvironmentHttpInternalServerError, EnvironmentHttpUnauthorizedError, @@ -70,7 +69,6 @@ import * as ServerSecretStore from "../auth/ServerSecretStore.ts"; import { requireEnvironmentScope } from "../auth/http.ts"; import * as ServerConfig from "../config.ts"; import * as ServerEnvironment from "../environment/ServerEnvironment.ts"; -import { makeRelayEnvironmentClient } from "../relay/relayEnvironmentClient.ts"; import * as AgentAwarenessRelay from "../relay/AgentAwarenessRelay.ts"; import * as ManagedEndpointRuntime from "./ManagedEndpointRuntime.ts"; import { @@ -89,8 +87,6 @@ import { encodeConfirmedOriginJson, PUBLISH_AGENT_ACTIVITY_SECRET, HOLD_WEBHOOKS_WHILE_OFFLINE_SECRET, - readHoldWebhooksWhileOffline, - readRelayConnection, RELAY_ENVIRONMENT_CREDENTIAL_SECRET, RELAY_ISSUER_SECRET, RELAY_URL_SECRET, @@ -102,6 +98,7 @@ import { setCliDesiredCloudLink, } from "./CliState.ts"; import * as CliTokenManager from "./CliTokenManager.ts"; +import * as CloudPreferences from "./CloudPreferences.ts"; import { getOrCreateEnvironmentKeyPairFromSecretStore } from "./environmentKeys.ts"; import { traceRelayRequest } from "./traceRelayRequest.ts"; import { filterRelayResponse, relayRequestError, shouldRetryCloudLink } from "./relayResponse.ts"; @@ -1358,65 +1355,6 @@ const cloudUnlinkHandler = Effect.fn("environment.cloud.unlink")( ), ); -const pushHoldWebhooksWhileOffline = Effect.fn("environment.cloud.pushHoldWebhooksWhileOffline")( - function* (dependencies: CloudHttpDependencies, holdWebhooksWhileOffline: boolean) { - const connection = yield* readRelayConnection(dependencies.secrets); - if (connection === null) { - return yield* new EnvironmentHttpBadRequestError({ - message: "Link this environment to T3 Connect first.", - }); - } - const environmentId = yield* dependencies.environment.getEnvironmentId; - const client = yield* makeRelayEnvironmentClient(connection); - yield* client.server - .updateLinkPreferences({ - params: { environmentId }, - payload: { holdWebhooksWhileOffline }, - }) - .pipe( - Effect.timeout("10 seconds"), - Effect.catch( - failEnvironmentCloudInternalError("Could not update T3 Connect webhook settings."), - ), - ); - }, -); - -const cloudPreferencesHandler = Effect.fn("environment.cloud.preferences")( - function* (dependencies: CloudHttpDependencies, payload: EnvironmentCloudPreferencesRequest) { - yield* requireEnvironmentScope(AuthRelayWriteScope); - if (payload.holdWebhooksWhileOffline !== undefined) { - // The relay decides whether to hold a request, so it is told first; the - // local copy is only saved once the relay has the same value, and the - // relay is put back if that save fails. - const previous = yield* readHoldWebhooksWhileOffline(dependencies.secrets); - yield* pushHoldWebhooksWhileOffline(dependencies, payload.holdWebhooksWhileOffline); - yield* dependencies.secrets - .set( - HOLD_WEBHOOKS_WHILE_OFFLINE_SECRET, - stringToBytes(String(payload.holdWebhooksWhileOffline)), - ) - .pipe( - Effect.tapError(() => - previous === payload.holdWebhooksWhileOffline - ? Effect.void - : pushHoldWebhooksWhileOffline(dependencies, previous).pipe(Effect.ignore), - ), - ); - } - yield* dependencies.secrets.set( - PUBLISH_AGENT_ACTIVITY_SECRET, - stringToBytes(String(payload.publishAgentActivity)), - ); - yield* dependencies.awarenessRelay.requestCatchUp(); - return yield* readCloudLinkState(dependencies); - }, - Effect.catchIf( - ServerSecretStore.isSecretStoreError, - failEnvironmentCloudInternalError("Could not persist environment cloud preferences."), - ), -); - const cloudEnvironmentHealthHandler = Effect.fn("environment.cloud.health")( function* (dependencies: CloudHttpDependencies, request: RelayCloudEnvironmentHealthRequest) { const cloudMintPublicKey = yield* dependencies.secrets @@ -1661,12 +1599,26 @@ export const connectHttpApiLayer = HttpApiBuilder.group( "connect", Effect.fnUntraced(function* (handlers) { const dependencies = yield* cloudHttpDependencies; + const preferences = yield* CloudPreferences.CloudPreferences; return handlers .handle("linkProof", ({ payload }) => cloudLinkProofHandler(dependencies, payload)) .handle("relayConfig", ({ payload }) => cloudRelayConfigHandler(dependencies, payload)) .handle("linkState", () => cloudLinkStateHandler(dependencies)) .handle("unlink", () => cloudUnlinkHandler(dependencies)) - .handle("preferences", ({ payload }) => cloudPreferencesHandler(dependencies, payload)) + .handle( + "preferences", + Effect.fn("environment.cloud.preferences")( + function* ({ payload }) { + yield* requireEnvironmentScope(AuthRelayWriteScope); + yield* preferences.update(payload); + return yield* readCloudLinkState(dependencies); + }, + Effect.catchIf( + ServerSecretStore.isSecretStoreError, + failEnvironmentCloudInternalError("Could not read environment cloud preferences."), + ), + ), + ) .handle("health", ({ payload }) => cloudEnvironmentHealthHandler(dependencies, payload)) .handle("mintCredential", ({ payload }) => cloudMintCredentialHandler(dependencies, payload)) .handle("t3MintCredential", ({ payload }) => diff --git a/apps/server/src/relay/HeldHooksWaker.ts b/apps/server/src/relay/HeldHooksWaker.ts index 9e0dae4ce624..bae1dc0705fb 100644 --- a/apps/server/src/relay/HeldHooksWaker.ts +++ b/apps/server/src/relay/HeldHooksWaker.ts @@ -3,7 +3,6 @@ import * as Layer from "effect/Layer"; import * as Schedule from "effect/Schedule"; import * as Stream from "effect/Stream"; -import * as ServerSecretStore from "../auth/ServerSecretStore.ts"; import * as CloudManagedEndpointRuntime from "../cloud/ManagedEndpointRuntime.ts"; import { readHoldWebhooksWhileOffline, readRelayConnection } from "../cloud/config.ts"; import * as ServerEnvironment from "../environment/ServerEnvironment.ts"; @@ -15,9 +14,8 @@ import { makeRelayEnvironmentClient } from "./relayEnvironmentClient.ts"; * next backoff step. Nothing happens unless the environment opted in. */ const wakeHeldHooks = Effect.fn("HeldHooksWaker.wake")(function* () { - const secrets = yield* ServerSecretStore.ServerSecretStore; - if (!(yield* readHoldWebhooksWhileOffline(secrets))) return false; - const connection = yield* readRelayConnection(secrets); + if (!(yield* readHoldWebhooksWhileOffline)) return false; + const connection = yield* readRelayConnection; if (connection === null) return false; const environmentId = yield* (yield* ServerEnvironment.ServerEnvironment).getEnvironmentId; const client = yield* makeRelayEnvironmentClient(connection); diff --git a/apps/server/src/scheduledTasks/relayDeliveryProof.ts b/apps/server/src/scheduledTasks/RelayDeliveryProof.ts similarity index 81% rename from apps/server/src/scheduledTasks/relayDeliveryProof.ts rename to apps/server/src/scheduledTasks/RelayDeliveryProof.ts index 9c277e6998d8..1012fb9a949e 100644 --- a/apps/server/src/scheduledTasks/relayDeliveryProof.ts +++ b/apps/server/src/scheduledTasks/RelayDeliveryProof.ts @@ -14,7 +14,9 @@ import { verifyRelayJwt, } from "@t3tools/shared/relayJwt"; import * as Clock from "effect/Clock"; +import * as Context from "effect/Context"; import * as Effect from "effect/Effect"; +import * as Layer from "effect/Layer"; import * as Option from "effect/Option"; import * as Schema from "effect/Schema"; @@ -29,11 +31,23 @@ const decodePayload = Schema.decodeUnknownOption(RelayHookDeliveryProofPayload); const text = (bytes: Option.Option) => Option.map(bytes, (value) => new TextDecoder().decode(value)); -/** The relay's own delivery id and receive time, when the request proves it came from the relay. */ -export const makeRelayDeliveryVerifier = Effect.gen(function* () { +export class RelayDeliveryProof extends Context.Service< + RelayDeliveryProof, + { + /** The relay's own delivery id and receive time, when the request proves it came from the relay. */ + readonly verify: (input: { + readonly headers: Readonly>; + readonly hookId: string; + }) => Effect.Effect< + Option.Option<{ readonly deliveryId: string; readonly receivedAt: string }> + >; + } +>()("t3/scheduledTasks/RelayDeliveryProof") {} + +const make = Effect.gen(function* () { const secrets = yield* ServerSecretStore.ServerSecretStore; const environment = yield* ServerEnvironment.ServerEnvironment; - return (input: { readonly headers: Readonly>; readonly hookId: string }) => + const verify: RelayDeliveryProof["Service"]["verify"] = (input) => Effect.gen(function* () { const proof = input.headers[RELAY_HOOK_DELIVERY_HEADER]; const deliveryId = input.headers["x-t3-relay-delivery-id"]; @@ -71,4 +85,7 @@ export const makeRelayDeliveryVerifier = Effect.gen(function* () { } return Option.some({ deliveryId, receivedAt }); }).pipe(Effect.orElseSucceed(Option.none), Effect.withSpan("webhook.verifyRelayDelivery")); + return RelayDeliveryProof.of({ verify }); }); + +export const layer = Layer.effect(RelayDeliveryProof, make); diff --git a/apps/server/src/scheduledTasks/webhookRoute.test.ts b/apps/server/src/scheduledTasks/webhookRoute.test.ts index 511b4f5ad875..eb7b9ddba030 100644 --- a/apps/server/src/scheduledTasks/webhookRoute.test.ts +++ b/apps/server/src/scheduledTasks/webhookRoute.test.ts @@ -26,6 +26,7 @@ import { import * as ServerSecretStore from "../auth/ServerSecretStore.ts"; import { CLOUD_MINT_PUBLIC_KEY, RELAY_ISSUER_SECRET } from "../cloud/config.ts"; import * as ServerEnvironment from "../environment/ServerEnvironment.ts"; +import * as RelayDeliveryProof from "./RelayDeliveryProof.ts"; import { WEBHOOK_MAX_BODY_BYTES, webhookHttpApiLayer } from "./webhookRoute.ts"; class WebhookTestApi extends HttpApi.make("environment").add(EnvironmentHttpApi.groups.webhooks) {} @@ -77,6 +78,7 @@ const handlerFor = ( HttpRouter.toWebHandler( HttpApiBuilder.layer(WebhookTestApi).pipe( Layer.provide(webhookHttpApiLayer), + Layer.provide(RelayDeliveryProof.layer), Layer.provide(Layer.mock(ScheduledTaskService)({ triggerWebhook: trigger })), Layer.provide( Layer.mock(ServerSecretStore.ServerSecretStore)({ diff --git a/apps/server/src/scheduledTasks/webhookRoute.ts b/apps/server/src/scheduledTasks/webhookRoute.ts index eee969f2d34c..dec01a5e6b15 100644 --- a/apps/server/src/scheduledTasks/webhookRoute.ts +++ b/apps/server/src/scheduledTasks/webhookRoute.ts @@ -11,7 +11,7 @@ import { withRelayClientTracing } from "@t3tools/shared/relayTracing"; import * as HttpApiBuilder from "effect/http-api/HttpApiBuilder"; import * as Metrics from "../observability/Metrics.ts"; -import { makeRelayDeliveryVerifier } from "./relayDeliveryProof.ts"; +import * as RelayDeliveryProof from "./RelayDeliveryProof.ts"; import * as ScheduledTaskService from "./ScheduledTaskService.ts"; /** Largest request body a webhook accepts. The relay enforces the same cap. */ @@ -23,121 +23,113 @@ const WEBHOOK_OUTCOME_HEADER = "x-t3-hook-outcome"; const json = (status: number, body: Record, outcome: string) => HttpServerResponse.jsonUnsafe(body, { status, headers: { [WEBHOOK_OUTCOME_HEADER]: outcome } }); -/** - * Handles `/api/hooks/:hookId/:token` for every accepted method. The endpoint - * is raw so the signature is checked over the exact body bytes; the service - * checks the token and signature. It is reachable directly, over the managed - * tunnel, or through the relay's stable `/v1/hooks/:endpointKey/:hookId/:token` - * URL, where the endpoint key is this environment's managed tunnel key. - */ -const handleWebhook = - ( - scheduledTasks: ScheduledTaskService.ScheduledTaskService["Service"], - verifyRelayDelivery: Effect.Success, - ) => - ({ - params, - request, - }: { - readonly params: { readonly hookId: string; readonly token: string }; - readonly request: HttpServerRequest.HttpServerRequest; - }) => - Effect.gen(function* () { - const contentLength = Number(request.headers["content-length"] ?? "0"); - if (!Number.isFinite(contentLength) || contentLength > WEBHOOK_MAX_BODY_BYTES) { - return json(413, { error: "body_too_large" }, "body_too_large"); - } - // Chunked requests carry no content-length, so the reader itself is capped. - const body = yield* request.arrayBuffer.pipe( - Effect.map((buffer) => new Uint8Array(buffer)), - Effect.provideService( - HttpIncomingMessage.MaxBodySize, - ByteSize.bytes(WEBHOOK_MAX_BODY_BYTES), - ), - Effect.option, - ); - // Refused before a task is looked up, so the service never sees them. - const tooLarge = (error: string) => - Metrics.increment(Metrics.webhookDeliveriesTotal, { - outcome: "body_too_large", - source: request.headers["x-t3-relay-delivery-id"] ? "relay" : "direct", - }).pipe(Effect.as(json(413, { error }, "body_too_large"))); - if (Option.isNone(body)) return yield* tooLarge("body_too_large_or_unreadable"); - if (body.value.byteLength > WEBHOOK_MAX_BODY_BYTES) { - return yield* tooLarge("body_too_large"); - } - - const headers: Record = {}; - for (const [name, value] of Object.entries(request.headers)) { - if (typeof value === "string") headers[name.toLowerCase()] = value; - } - const queryIndex = request.url.indexOf("?"); - // The relay's delivery id and receive time count only with its signed - // proof; the URL can also be called directly. The receive time matters - // for requests the relay held while we were offline. - const relay = yield* verifyRelayDelivery({ headers, hookId: params.hookId }); - const relayDeliveryId = Option.isSome(relay) ? relay.value.deliveryId : undefined; - const relayReceivedAt = Option.isSome(relay) ? relay.value.receivedAt : undefined; - - // A relay delivery joins the relay's trace, and goes to the T3 Connect - // tracer with it. Anyone else's traceparent is never trusted. - const relayParent = Option.isSome(relay) - ? HttpTraceContext.fromHeaders(request.headers) - : Option.none(); - const result = yield* scheduledTasks - .triggerWebhook({ - hookId: params.hookId, - token: params.token, - method: request.method, - path: `${ScheduledTaskService.WEBHOOK_ROUTE_PREFIX}/${encodeURIComponent(params.hookId)}`, - query: queryIndex === -1 ? "" : request.url.slice(queryIndex + 1), - headers, - body: body.value, - bodyText: new TextDecoder().decode(body.value), - ...(relayDeliveryId ? { relayDeliveryId } : {}), - ...(relayReceivedAt ? { receivedAt: relayReceivedAt } : {}), - }) - .pipe( - Option.isSome(relayParent) - ? (effect) => - effect.pipe(Effect.withParentSpan(relayParent.value), withRelayClientTracing) - : (effect) => effect, - // Defects too, so the sender only ever sees the fixed error body. - Effect.catchCause((cause) => - Effect.logWarning("Webhook delivery failed").pipe( - Effect.annotateLogs({ hookId: params.hookId }), - Effect.andThen(Effect.logDebug("Webhook delivery failure cause", { cause })), - Effect.as({ _tag: "error" as const }), - ), - ), - ); - - switch (result._tag) { - case "accepted": - return json(202, { deliveryId: result.deliveryId }, result.outcome); - case "not_found": - return json(404, { error: "hook_not_found" }, "not_found"); - case "rejected_signature": - return json(401, { error: "invalid_signature" }, "rejected_signature"); - case "disabled": - return json(409, { error: "hook_disabled" }, "disabled"); - case "rate_limited": - return json(429, { error: "rate_limited" }, result.outcome); - case "expired": - return json(410, { error: "delivery_too_old" }, "expired"); - case "error": - return json(500, { error: "internal_error" }, "error"); - } - }); - export const webhookHttpApiLayer = HttpApiBuilder.group( EnvironmentHttpApi, "webhooks", Effect.fnUntraced(function* (handlers) { - const handler = handleWebhook( - yield* ScheduledTaskService.ScheduledTaskService, - yield* makeRelayDeliveryVerifier, - ); + const scheduledTasks = yield* ScheduledTaskService.ScheduledTaskService; + const relayDeliveryProof = yield* RelayDeliveryProof.RelayDeliveryProof; + /** + * Handles `/api/hooks/:hookId/:token` for every accepted method. The endpoint + * is raw so the signature is checked over the exact body bytes; the service + * checks the token and signature. It is reachable directly, over the managed + * tunnel, or through the relay's stable `/v1/hooks/:endpointKey/:hookId/:token` + * URL, where the endpoint key is this environment's managed tunnel key. + */ + const handler = ({ + params, + request, + }: { + readonly params: { readonly hookId: string; readonly token: string }; + readonly request: HttpServerRequest.HttpServerRequest; + }) => + Effect.gen(function* () { + const contentLength = Number(request.headers["content-length"] ?? "0"); + if (!Number.isFinite(contentLength) || contentLength > WEBHOOK_MAX_BODY_BYTES) { + return json(413, { error: "body_too_large" }, "body_too_large"); + } + // Chunked requests carry no content-length, so the reader itself is capped. + const body = yield* request.arrayBuffer.pipe( + Effect.map((buffer) => new Uint8Array(buffer)), + Effect.provideService( + HttpIncomingMessage.MaxBodySize, + ByteSize.bytes(WEBHOOK_MAX_BODY_BYTES), + ), + Effect.option, + ); + // Refused before a task is looked up, so the service never sees them. + const tooLarge = (error: string) => + Metrics.increment(Metrics.webhookDeliveriesTotal, { + outcome: "body_too_large", + source: request.headers["x-t3-relay-delivery-id"] ? "relay" : "direct", + }).pipe(Effect.as(json(413, { error }, "body_too_large"))); + if (Option.isNone(body)) return yield* tooLarge("body_too_large_or_unreadable"); + if (body.value.byteLength > WEBHOOK_MAX_BODY_BYTES) { + return yield* tooLarge("body_too_large"); + } + + const headers: Record = {}; + for (const [name, value] of Object.entries(request.headers)) { + if (typeof value === "string") headers[name.toLowerCase()] = value; + } + const queryIndex = request.url.indexOf("?"); + // The relay's delivery id and receive time count only with its signed + // proof; the URL can also be called directly. The receive time matters + // for requests the relay held while we were offline. + const relay = yield* relayDeliveryProof.verify({ headers, hookId: params.hookId }); + const relayDeliveryId = Option.isSome(relay) ? relay.value.deliveryId : undefined; + const relayReceivedAt = Option.isSome(relay) ? relay.value.receivedAt : undefined; + + // A relay delivery joins the relay's trace, and goes to the T3 Connect + // tracer with it. Anyone else's traceparent is never trusted. + const relayParent = Option.isSome(relay) + ? HttpTraceContext.fromHeaders(request.headers) + : Option.none(); + const result = yield* scheduledTasks + .triggerWebhook({ + hookId: params.hookId, + token: params.token, + method: request.method, + path: `${ScheduledTaskService.WEBHOOK_ROUTE_PREFIX}/${encodeURIComponent(params.hookId)}`, + query: queryIndex === -1 ? "" : request.url.slice(queryIndex + 1), + headers, + body: body.value, + bodyText: new TextDecoder().decode(body.value), + ...(relayDeliveryId ? { relayDeliveryId } : {}), + ...(relayReceivedAt ? { receivedAt: relayReceivedAt } : {}), + }) + .pipe( + Option.isSome(relayParent) + ? (effect) => + effect.pipe(Effect.withParentSpan(relayParent.value), withRelayClientTracing) + : (effect) => effect, + // Defects too, so the sender only ever sees the fixed error body. + Effect.catchCause((cause) => + Effect.logWarning("Webhook delivery failed").pipe( + Effect.annotateLogs({ hookId: params.hookId }), + Effect.andThen(Effect.logDebug("Webhook delivery failure cause", { cause })), + Effect.as({ _tag: "error" as const }), + ), + ), + ); + + switch (result._tag) { + case "accepted": + return json(202, { deliveryId: result.deliveryId }, result.outcome); + case "not_found": + return json(404, { error: "hook_not_found" }, "not_found"); + case "rejected_signature": + return json(401, { error: "invalid_signature" }, "rejected_signature"); + case "disabled": + return json(409, { error: "hook_disabled" }, "disabled"); + case "rate_limited": + return json(429, { error: "rate_limited" }, result.outcome); + case "expired": + return json(410, { error: "delivery_too_old" }, "expired"); + case "error": + return json(500, { error: "internal_error" }, "error"); + } + }); return handlers .handleRaw("webhookPost", handler) .handleRaw("webhookPut", handler) diff --git a/apps/server/src/secrets/SecretRequests.ts b/apps/server/src/secrets/SecretRequests.ts index 8c2dacd4dcad..12ba9c692683 100644 --- a/apps/server/src/secrets/SecretRequests.ts +++ b/apps/server/src/secrets/SecretRequests.ts @@ -12,7 +12,6 @@ import { CommandId, SecretRef, SecretRequestError, - type SecretRequestFailureReason, type ProjectId, type SecretRequestAnswerInput, type ThreadId, @@ -62,9 +61,6 @@ const StoredSecret = Schema.fromJsonString( const encodeStored = Schema.encodeEffect(StoredSecret); const decodeStored = Schema.decodeUnknownOption(StoredSecret); -const fail = (reason: SecretRequestFailureReason, cause?: unknown) => - new SecretRequestError({ reason, ...(cause === undefined ? {} : { cause }) }); - export class SecretRequests extends Context.Service< SecretRequests, { @@ -109,19 +105,23 @@ const make = Effect.gen(function* () { turnItemTypes: ["secret_request"], messageRoles: [], }) - .pipe(Effect.mapError((cause) => fail("load_failed", cause))); + .pipe( + Effect.mapError( + (cause) => new SecretRequestError({ reason: "load_failed", cause: cause }), + ), + ); const item = records.turnItems.find((candidate) => candidate.id === input.turnItemId); if (item?.type !== "secret_request" || item.runId === null || item.nodeId === null) { - return yield* fail("not_found"); + return yield* new SecretRequestError({ reason: "not_found" }); } if (item.secretStatus !== "pending") { - return yield* fail("already_answered"); + return yield* new SecretRequestError({ reason: "already_answered" }); } // The agent is waiting inside the run that asked; once it has ended, // nobody will ever receive the ref, so a value saved now would be lost. const run = records.runs.find((candidate) => candidate.id === item.runId); if (run === undefined || ThreadManagementService.isTerminalRunStatus(run.status)) { - return yield* fail("agent_stopped"); + return yield* new SecretRequestError({ reason: "agent_stopped" }); } // Store first: the card only says saved once the value is kept. Create, // not set: a second answer racing this one must not replace the value. @@ -140,7 +140,9 @@ const make = Effect.gen(function* () { ) .pipe( Effect.catchIf(ServerSecretStore.isSecretAlreadyExistsError, () => Effect.void), - Effect.mapError((error) => fail("store_failed", error)), + Effect.mapError( + (error) => new SecretRequestError({ reason: "store_failed", cause: error }), + ), ); } const secretStatus = input.answer.type === "save" ? "saved" : "declined"; @@ -158,7 +160,9 @@ const make = Effect.gen(function* () { secretStatus, }) .pipe( - Effect.mapError((cause) => fail("record_failed", cause)), + Effect.mapError( + (cause) => new SecretRequestError({ reason: "record_failed", cause: cause }), + ), // The card still says pending, so the user can save again; the value // stored above would make that retry look already answered. Effect.tapError(() => @@ -176,11 +180,15 @@ const make = Effect.gen(function* () { turnItemTypes: ["secret_request"], messageRoles: [], }) - .pipe(Effect.mapError((cause) => fail("record_failed", cause))); + .pipe( + Effect.mapError( + (cause) => new SecretRequestError({ reason: "record_failed", cause: cause }), + ), + ); const card = recorded.turnItems.find((candidate) => candidate.id === item.id); if (card?.type === "secret_request" && card.secretStatus !== "saved") { yield* removeLogged(storeName(refFor(salt, input.threadId, item.id))); - return yield* fail("agent_stopped"); + return yield* new SecretRequestError({ reason: "agent_stopped" }); } }).pipe(Effect.withSpan("SecretRequests.answer")); @@ -217,25 +225,30 @@ const make = Effect.gen(function* () { const consumeRef = (input: { readonly ref: SecretRef; readonly projectId: ProjectId }) => Effect.gen(function* () { - if (!REF_PATTERN.test(input.ref)) return yield* fail("invalid_ref"); + if (!REF_PATTERN.test(input.ref)) + return yield* new SecretRequestError({ reason: "invalid_ref" }); const stored = yield* store .get(storeName(input.ref)) - .pipe(Effect.mapError((cause) => fail("read_failed", cause))); + .pipe( + Effect.mapError( + (cause) => new SecretRequestError({ reason: "read_failed", cause: cause }), + ), + ); const decoded = Option.flatMap(stored, (bytes) => decodeStored(new TextDecoder().decode(bytes)), ); if (Option.isNone(decoded) || decoded.value.projectId !== input.projectId) { - return yield* fail("ref_unavailable"); + return yield* new SecretRequestError({ reason: "ref_unavailable" }); } if ((yield* Clock.currentTimeMillis) - decoded.value.savedAt > SECRET_REF_TTL_MS) { // The hourly sweep retries a removal that fails here. yield* removeLogged(storeName(input.ref)); - return yield* fail("ref_expired"); + return yield* new SecretRequestError({ reason: "ref_expired" }); } // One use: the value moves into whatever consumed it. If it cannot be // deleted, it is not handed out, so a ref is never used twice. if (!(yield* removeLogged(storeName(input.ref)))) { - return yield* fail("consume_failed"); + return yield* new SecretRequestError({ reason: "consume_failed" }); } return decoded.value.value; }); diff --git a/apps/server/src/server.ts b/apps/server/src/server.ts index 3be16c3a2ffb..38d1dc6cdd98 100644 --- a/apps/server/src/server.ts +++ b/apps/server/src/server.ts @@ -119,6 +119,8 @@ import { authHttpApiLayer, environmentAuthenticatedAuthLayer } from "./auth/http import * as ReplayMarkers from "./auth/replayMarkers.ts"; import * as ServerSecretStore from "./auth/ServerSecretStore.ts"; import { webhookHttpApiLayer } from "./scheduledTasks/webhookRoute.ts"; +import * as RelayDeliveryProof from "./scheduledTasks/RelayDeliveryProof.ts"; +import * as CloudPreferences from "./cloud/CloudPreferences.ts"; import * as HeldHooksWaker from "./relay/HeldHooksWaker.ts"; import { relayHookBaseUrl, @@ -679,12 +681,12 @@ const makeRoutesLayer = Layer.mergeAll( Layer.mergeAll( HttpApiBuilder.layer(EnvironmentHttpApi).pipe( Layer.provide(authHttpApiLayer), - Layer.provide(connectHttpApiLayer), + Layer.provide(connectHttpApiLayer.pipe(Layer.provide(CloudPreferences.layer))), Layer.provide(orchestrationHttpApiLayer), Layer.provide(pullRequestHttpApiLayer), Layer.provide(projectHttpApiLayer), Layer.provide(serverEnvironmentHttpApiLayer), - Layer.provide(webhookHttpApiLayer), + Layer.provide(webhookHttpApiLayer.pipe(Layer.provide(RelayDeliveryProof.layer))), Layer.provide(environmentAuthenticatedAuthLayer), ), otlpTracesProxyRouteLayer, diff --git a/infra/relay/src/hooks/HeldHooks.test.ts b/infra/relay/src/hooks/HeldHooks.test.ts new file mode 100644 index 000000000000..25feefdf7fb1 --- /dev/null +++ b/infra/relay/src/hooks/HeldHooks.test.ts @@ -0,0 +1,157 @@ +import { describe, expect, it } from "@effect/vitest"; +import * as Effect from "effect/Effect"; +import * as Layer from "effect/Layer"; +import * as Redacted from "effect/Redacted"; + +import * as RelayConfiguration from "../Config.ts"; +import * as EnvironmentLinks from "../environments/EnvironmentLinks.ts"; +import * as ManagedEndpointAllocations from "../environments/ManagedEndpointAllocations.ts"; +import * as HeldHooks from "./HeldHooks.ts"; +import * as HookInbox from "./HookInbox.ts"; + +const settings = { + managedEndpointBaseDomain: "example.test", + managedEndpointNamespace: "dev", + cloudMintPrivateKey: Redacted.make("unused"), +} as RelayConfiguration.RelayConfiguration["Service"]; + +const environmentId = "env-hook"; +const ownKey = "0123456789abcdef"; +const otherKey = "fedcba9876543210"; + +const allocationFor = ( + userId: string, + key: string, + ready: boolean, +): ManagedEndpointAllocations.ManagedEndpointAllocation => ({ + userId, + environmentId, + hostname: `${key}.example.test`, + tunnelId: `tunnel-${key}`, + tunnelName: `t3coderelay-managedendpoint-dev-${key}`, + dnsRecordId: "dns-record-id", + readyAt: ready ? "2026-05-25T00:00:00.000Z" : null, + origin: { localHttpHost: "127.0.0.1", localHttpPort: 3773 }, + updatedAt: "2026-05-25T00:00:00.000Z", + generation: 1, +}); + +const linkFor = (userId: string, key: string, environmentPublicKey: string) => ({ + userId, + environmentId: environmentId as never, + label: "Hook env", + endpoint: { + httpBaseUrl: `https://${key}.example.test/`, + wsBaseUrl: `wss://${key}.example.test/ws`, + providerKind: "cloudflare_tunnel" as const, + }, + environmentPublicKey, + linkedAt: "2026-05-25T00:00:00.000Z", + holdWebhooksWhileOffline: true, +}); + +/** This environment's own link (user_1) and another account's link of the same environment. */ +const withService = ( + options: { readonly ownReady?: boolean; readonly pendingKeys?: ReadonlyArray }, + body: (input: { + readonly heldHooks: HeldHooks.HeldHooks["Service"]; + readonly cleared: Array; + readonly woken: Array; + }) => Effect.Effect, +) => { + const cleared: Array = []; + const woken: Array = []; + const links = [linkFor("user_1", ownKey, "own-key"), linkFor("user_2", otherKey, "other-key")]; + const allocations = [ + allocationFor("user_1", ownKey, options.ownReady ?? true), + allocationFor("user_2", otherKey, true), + ]; + const dependencies = Layer.mergeAll( + Layer.succeed(RelayConfiguration.RelayConfiguration, settings), + Layer.mock(EnvironmentLinks.EnvironmentLinks, { + findActiveManagedForEnvironment: (input) => + Effect.succeed( + links.filter( + (link) => + link.environmentId === input.environmentId && + (input.userId === undefined || link.userId === input.userId) && + (input.environmentPublicKey === undefined || + link.environmentPublicKey === input.environmentPublicKey), + ), + ), + setHoldWebhooksWhileOffline: () => Effect.void, + }), + Layer.mock(ManagedEndpointAllocations.ManagedEndpointAllocations, { + get: (input) => + Effect.succeed( + allocations.find((allocation) => allocation.userId === input.userId) ?? null, + ), + getByTunnelName: (tunnelName) => + Effect.succeed( + allocations.find((allocation) => allocation.tunnelName === tunnelName) ?? null, + ), + }), + Layer.mock(HookInbox.HookInbox, { + clear: ({ endpointKey }) => Effect.sync(() => void cleared.push(endpointKey)), + wake: ({ endpointKey }) => + Effect.sync(() => { + woken.push(endpointKey); + return (options.pendingKeys ?? []).includes(endpointKey); + }), + }), + ); + return Effect.gen(function* () { + const heldHooks = yield* HeldHooks.HeldHooks; + return yield* body({ heldHooks, cleared, woken }); + }).pipe(Effect.provide(HeldHooks.layer.pipe(Layer.provide(dependencies)))); +}; + +describe("HeldHooks", () => { + it.effect("opting out clears only the caller's own endpoints", () => + withService({}, ({ heldHooks, cleared }) => + Effect.gen(function* () { + yield* heldHooks.setHoldWhileOffline({ + environmentId, + environmentPublicKey: "own-key", + holdWebhooksWhileOffline: false, + }); + expect(cleared).toEqual([ownKey]); + }), + ), + ); + + it.effect("opting in keeps what is held", () => + withService({}, ({ heldHooks, cleared }) => + Effect.gen(function* () { + yield* heldHooks.setHoldWhileOffline({ + environmentId, + environmentPublicKey: "own-key", + holdWebhooksWhileOffline: true, + }); + expect(cleared).toEqual([]); + }), + ), + ); + + it.effect("wakes the caller's ready endpoints and reports whether anything was waiting", () => + withService({ pendingKeys: [ownKey] }, ({ heldHooks, woken }) => + Effect.gen(function* () { + expect(yield* heldHooks.wake({ environmentId, environmentPublicKey: "own-key" })).toBe( + true, + ); + expect(woken).toEqual([ownKey]); + }), + ), + ); + + it.effect("does not wake an endpoint that is not ready", () => + withService({ ownReady: false, pendingKeys: [ownKey] }, ({ heldHooks, woken }) => + Effect.gen(function* () { + expect(yield* heldHooks.wake({ environmentId, environmentPublicKey: "own-key" })).toBe( + false, + ); + expect(woken).toEqual([]); + }), + ), + ); +}); diff --git a/infra/relay/src/hooks/HeldHooks.ts b/infra/relay/src/hooks/HeldHooks.ts new file mode 100644 index 000000000000..9ed3195f3b2b --- /dev/null +++ b/infra/relay/src/hooks/HeldHooks.ts @@ -0,0 +1,143 @@ +import * as Context from "effect/Context"; +import * as Effect from "effect/Effect"; +import * as Layer from "effect/Layer"; +import * as Result from "effect/Result"; + +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"; + +/** 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; +}; + +type PersistenceError = + | EnvironmentLinks.EnvironmentLinkEnvironmentLookupPersistenceError + | ManagedEndpointAllocations.ManagedEndpointAllocationPersistenceError; + +interface HookEndpoint { + readonly httpBaseUrl: string; + readonly environmentId: string; + readonly holdWhileOffline: boolean; +} + +export class HeldHooks extends Context.Service< + HeldHooks, + { + /** + * 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. + */ + readonly resolveEndpoint: ( + endpointKey: string, + ) => Effect.Effect; + /** + * Saves an environment's opt-in. Opting out also drops what is already held, + * rather than delivering it later to an environment that said it does not + * want it. Only the links the caller's environment key proved are touched. + */ + readonly setHoldWhileOffline: (input: { + readonly environmentId: string; + readonly environmentPublicKey: string; + readonly holdWebhooksWhileOffline: boolean; + }) => Effect.Effect; + /** Delivers what the environment's endpoints hold now; true when anything was waiting. */ + readonly wake: (input: { + readonly environmentId: string; + readonly environmentPublicKey: string; + }) => Effect.Effect; + } +>()("t3code-relay/hooks/HeldHooks") {} + +const make = Effect.gen(function* () { + const links = yield* EnvironmentLinks.EnvironmentLinks; + const allocations = yield* ManagedEndpointAllocations.ManagedEndpointAllocations; + const settings = yield* RelayConfiguration.RelayConfiguration; + const inbox = yield* HookInbox.HookInbox; + + const resolveEndpoint: HeldHooks["Service"]["resolveEndpoint"] = Effect.fn( + "relay.hooks.resolve_endpoint", + )(function* (endpointKey) { + 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, + }; + }); + + /** + * Endpoint keys of the managed links the calling environment key proved. + * The environment id alone would also match other accounts' links of it. + */ + const ownEndpointKeys = Effect.fn("relay.hooks.own_endpoint_keys")(function* (input: { + readonly environmentId: string; + readonly environmentPublicKey: string; + }) { + const namespace = settings.managedEndpointNamespace; + if (!namespace) return []; + const ownLinks = yield* links.findActiveManagedForEnvironment(input); + const keys: Array = []; + for (const link of ownLinks) { + const allocation = yield* allocations.get({ + userId: link.userId, + environmentId: input.environmentId, + }); + const key = allocation ? endpointKeyForTunnelName(namespace, allocation.tunnelName) : null; + if (key !== null) keys.push(key); + } + return keys; + }); + + const setHoldWhileOffline: HeldHooks["Service"]["setHoldWhileOffline"] = Effect.fn( + "relay.hooks.set_hold_while_offline", + )(function* (input) { + yield* links.setHoldWebhooksWhileOffline(input); + if (input.holdWebhooksWhileOffline) return; + for (const endpointKey of yield* ownEndpointKeys(input)) { + yield* inbox.clear({ endpointKey }); + } + }); + + const wake: HeldHooks["Service"]["wake"] = Effect.fn("relay.hooks.wake")(function* (input) { + let pending = false; + for (const endpointKey of yield* ownEndpointKeys(input)) { + const endpoint = yield* resolveEndpoint(endpointKey); + if (endpoint === null) continue; + if (yield* inbox.wake({ endpointKey, baseUrl: endpoint.httpBaseUrl })) { + pending = true; + } + } + return pending; + }); + + return HeldHooks.of({ resolveEndpoint, setHoldWhileOffline, wake }); +}); + +export const layer = Layer.effect(HeldHooks, make); diff --git a/infra/relay/src/hooks/HookForwarder.test.ts b/infra/relay/src/hooks/HookForwarder.test.ts index aab7c991af1e..ee056a766826 100644 --- a/infra/relay/src/hooks/HookForwarder.test.ts +++ b/infra/relay/src/hooks/HookForwarder.test.ts @@ -37,6 +37,7 @@ import { } from "../http/Api.ts"; import * as HookForwarder from "./HookForwarder.ts"; import { RELAY_HOOK_DELIVERY_TYP, verifyRelayJwt } from "@t3tools/shared/relayJwt"; +import * as HeldHooks from "./HeldHooks.ts"; import * as HookInbox from "./HookInbox.ts"; import type { HeldHook } from "./HookInboxStore.ts"; import { RELAY_HOOK_UPSTREAM_TIMEOUT_MS } from "./upstream.ts"; @@ -114,6 +115,7 @@ function makeHarness(options: Harness = {}) { ), )); const forwarderLayer = HookForwarder.layer.pipe( + Layer.provideMerge(HeldHooks.layer), Layer.provide( Layer.mergeAll( Layer.succeed(RelayConfiguration.RelayConfiguration, settings), diff --git a/infra/relay/src/hooks/HookForwarder.ts b/infra/relay/src/hooks/HookForwarder.ts index f5031389dfc7..8e32d9c1f185 100644 --- a/infra/relay/src/hooks/HookForwarder.ts +++ b/infra/relay/src/hooks/HookForwarder.ts @@ -23,13 +23,8 @@ import { } 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 { MANAGED_ENDPOINT_KEY_PATTERN } from "../deploymentConfig.ts"; +import * as HeldHooks from "./HeldHooks.ts"; import * as HookInbox from "./HookInbox.ts"; import { sendUpstream, TUNNEL_OFFLINE_STATUS } from "./upstream.ts"; @@ -143,34 +138,6 @@ 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"]); @@ -266,57 +233,41 @@ const readCappedBody = (request: HttpServerRequest.HttpServerRequest) => ); }); -/** - * 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 heldHooks = yield* HeldHooks.HeldHooks; 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 signDeliveryProof = (input: { + 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(settings.cloudMintPrivateKey), + typ: RELAY_HOOK_DELIVERY_TYP, + payload: { + iss: normalizeRelayIssuer(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); + const handle = Effect.fn("relay.hooks.forward")(function* ( request: HttpServerRequest.HttpServerRequest, ) { @@ -357,10 +308,7 @@ const make = Effect.gen(function* () { 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), + const endpoint = yield* heldHooks.resolveEndpoint(parsed.endpointKey).pipe( Effect.catch((error) => Effect.logWarning("Failed to resolve hook endpoint", { endpointKey: parsed.endpointKey, @@ -398,7 +346,6 @@ const make = Effect.gen(function* () { // 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, diff --git a/infra/relay/src/http/Api.test.ts b/infra/relay/src/http/Api.test.ts index 4f1cd17e094c..4ddf5e0a3649 100644 --- a/infra/relay/src/http/Api.test.ts +++ b/infra/relay/src/http/Api.test.ts @@ -61,6 +61,7 @@ import { import * as RelayConfiguration from "../Config.ts"; import * as RelayDb from "../db.ts"; import * as EnvironmentCredentials from "../environments/EnvironmentCredentials.ts"; +import * as HeldHooks from "../hooks/HeldHooks.ts"; import * as HookInbox from "../hooks/HookInbox.ts"; import * as EnvironmentLinks from "../environments/EnvironmentLinks.ts"; import * as ManagedEndpointAllocations from "../environments/ManagedEndpointAllocations.ts"; @@ -1269,14 +1270,7 @@ describe("relay routing fallback", () => { Layer.mock(ManagedEndpointProvider.ManagedEndpointProvider, {}), ), ), - Layer.provide([ - publisher, - signatures, - Layer.mock(EnvironmentLinks.EnvironmentLinks, {}), - Layer.mock(HookInbox.HookInbox, {}), - Layer.mock(ManagedEndpointAllocations.ManagedEndpointAllocations, {}), - Layer.succeed(RelayConfiguration.RelayConfiguration, relaySettings), - ]), + Layer.provide([publisher, signatures, Layer.mock(HeldHooks.HeldHooks, {})]), ), ), Layer.provide(auth), diff --git a/infra/relay/src/http/Api.ts b/infra/relay/src/http/Api.ts index 7b308653d510..174e7e8bbbe2 100644 --- a/infra/relay/src/http/Api.ts +++ b/infra/relay/src/http/Api.ts @@ -67,6 +67,7 @@ import * as RelayTokens from "../auth/RelayTokens.ts"; import * as EnvironmentCredentials from "../environments/EnvironmentCredentials.ts"; import * as EnvironmentLinks from "../environments/EnvironmentLinks.ts"; import * as HookForwarder from "../hooks/HookForwarder.ts"; +import * as HeldHooks from "../hooks/HeldHooks.ts"; import * as HookInbox from "../hooks/HookInbox.ts"; import * as LiveActivities from "../agentActivity/LiveActivities.ts"; import * as RelayConfiguration from "../Config.ts"; @@ -554,7 +555,7 @@ export const unlinkEnvironmentRecord = Effect.fn("relay.api.client.unlinkEnviron // is already gone, and its inbox drops anything left after its TTL. const endpointKey = deprovisionTarget && input.managedEndpointNamespace - ? HookForwarder.endpointKeyForTunnelName( + ? HeldHooks.endpointKeyForTunnelName( input.managedEndpointNamespace, deprovisionTarget.tunnelName, ) @@ -1176,10 +1177,7 @@ export const serverApi = HttpApiBuilder.group( Effect.fnUntraced(function* (handlers) { const publisher = yield* AgentActivityPublisher.AgentActivityPublisher; const publishSignatures = yield* EnvironmentPublishSignatures.EnvironmentPublishSignatures; - const links = yield* EnvironmentLinks.EnvironmentLinks; - const inbox = yield* HookInbox.HookInbox; - const allocations = yield* ManagedEndpointAllocations.ManagedEndpointAllocations; - const settings = yield* RelayConfiguration.RelayConfiguration; + const heldHooks = yield* HeldHooks.HeldHooks; const requireOwnEnvironment = (environmentId: string) => Effect.gen(function* () { const principal = yield* RelayEnvironmentPrincipal; @@ -1187,29 +1185,6 @@ export const serverApi = HttpApiBuilder.group( return yield* new HttpApiError.Unauthorized({}); } }); - /** - * Endpoint keys of the managed links the calling environment key proved. - * The environment id alone would also match other accounts' links of it. - */ - const ownEndpointKeys = (environmentId: string) => - Effect.gen(function* () { - const principal = yield* RelayEnvironmentPrincipal; - const namespace = settings.managedEndpointNamespace; - if (!namespace) return []; - const ownLinks = yield* links.findActiveManagedForEnvironment({ - environmentId, - environmentPublicKey: principal.environmentPublicKey, - }); - const keys: Array = []; - for (const link of ownLinks) { - const allocation = yield* allocations.get({ userId: link.userId, environmentId }); - const key = allocation - ? HookForwarder.endpointKeyForTunnelName(namespace, allocation.tunnelName) - : null; - if (key !== null) keys.push(key); - } - return keys; - }); const activityHandlers = handlers.handle( "publishAgentActivity", Effect.fn("relay.api.server.publishAgentActivity")( @@ -1425,18 +1400,11 @@ export const serverApi = HttpApiBuilder.group( Effect.fn("relay.api.server.updateLinkPreferences")(function* ({ params, payload }) { yield* requireOwnEnvironment(params.environmentId); const principal = yield* RelayEnvironmentPrincipal; - yield* links.setHoldWebhooksWhileOffline({ + yield* heldHooks.setHoldWhileOffline({ environmentId: params.environmentId, environmentPublicKey: principal.environmentPublicKey, holdWebhooksWhileOffline: payload.holdWebhooksWhileOffline, }); - // Opting out also drops what is already held, rather than delivering - // it later to an environment that said it does not want it. - if (!payload.holdWebhooksWhileOffline) { - for (const endpointKey of yield* ownEndpointKeys(params.environmentId)) { - yield* inbox.clear({ endpointKey }); - } - } return payload; }, mapRelayCommonApiErrors("not_authorized")), ) @@ -1444,21 +1412,11 @@ export const serverApi = HttpApiBuilder.group( "wakeHeldHooks", Effect.fn("relay.api.server.wakeHeldHooks")(function* ({ params }) { yield* requireOwnEnvironment(params.environmentId); - let pending = false; - for (const endpointKey of yield* ownEndpointKeys(params.environmentId)) { - const endpoint = yield* HookForwarder.resolveHookEndpoint(endpointKey).pipe( - Effect.provideService(EnvironmentLinks.EnvironmentLinks, links), - Effect.provideService( - ManagedEndpointAllocations.ManagedEndpointAllocations, - allocations, - ), - Effect.provideService(RelayConfiguration.RelayConfiguration, settings), - ); - if (endpoint === null) continue; - if (yield* inbox.wake({ endpointKey, baseUrl: endpoint.httpBaseUrl })) { - pending = true; - } - } + const principal = yield* RelayEnvironmentPrincipal; + const pending = yield* heldHooks.wake({ + environmentId: params.environmentId, + environmentPublicKey: principal.environmentPublicKey, + }); return { pending }; }, mapRelayCommonApiErrors("not_authorized")), ); diff --git a/infra/relay/src/worker.ts b/infra/relay/src/worker.ts index 25cdd257f166..7c067589e81a 100644 --- a/infra/relay/src/worker.ts +++ b/infra/relay/src/worker.ts @@ -76,6 +76,7 @@ import * as ManagedEndpointReaper from "./environments/ManagedEndpointReaper.ts" import * as ManagedTunnelLimits from "./environments/ManagedTunnelLimits.ts"; import * as MobileRegistrations from "./agentActivity/MobileRegistrations.ts"; import * as HookForwarder from "./hooks/HookForwarder.ts"; +import * as HeldHooks from "./hooks/HeldHooks.ts"; import * as HookInbox from "./hooks/HookInbox.ts"; import { HookInboxObject, HookInboxObjectLive } from "./hooks/HookInboxObject.ts"; @@ -273,7 +274,11 @@ export const ApiLive = Api.make( Layer.provideMerge(EnvironmentConnector.layer), Layer.provideMerge(EnvironmentLinker.layer), Layer.provideMerge( - Layer.merge(EnvironmentPublishSignatures.layer, ManagedEndpointReaper.layer), + Layer.mergeAll( + EnvironmentPublishSignatures.layer, + ManagedEndpointReaper.layer, + HeldHooks.layer, + ), ), Layer.provideMerge( ManagedEndpointProvider.layerCloudflareBindings(