From 4a22a8bf582ea4e0398437cc42d1e06b394796a9 Mon Sep 17 00:00:00 2001 From: Julius Marminge Date: Mon, 5 Oct 2026 15:49:40 -0700 Subject: [PATCH 1/9] refactor(server,relay): webhook capabilities live in services, not handlers Relay delivery verification is a RelayDeliveryProof service. Saving T3 Connect preferences, with its relay-first write and rollback, is a CloudPreferences service. On the relay, endpoint resolution, opting in or out of held webhooks, and waking them are HeldHooks methods; HookForwarder and the server API call it. Handlers now authorize, call one method and map errors. HeldHooks and CloudPreferences get their own tests. Co-Authored-By: Claude Opus 5.5 (1M context) --- .../server/src/cloud/CloudPreferences.test.ts | 98 +++++++++++ apps/server/src/cloud/CloudPreferences.ts | 97 +++++++++++ apps/server/src/cloud/http.ts | 62 ++----- ...DeliveryProof.ts => RelayDeliveryProof.ts} | 23 ++- .../src/scheduledTasks/webhookRoute.test.ts | 2 + .../server/src/scheduledTasks/webhookRoute.ts | 8 +- apps/server/src/server.ts | 6 +- infra/relay/src/hooks/HeldHooks.test.ts | 157 ++++++++++++++++++ infra/relay/src/hooks/HeldHooks.ts | 143 ++++++++++++++++ infra/relay/src/hooks/HookForwarder.test.ts | 2 + infra/relay/src/hooks/HookForwarder.ts | 59 +------ infra/relay/src/http/Api.test.ts | 10 +- infra/relay/src/http/Api.ts | 60 +------ infra/relay/src/worker.ts | 7 +- 14 files changed, 559 insertions(+), 175 deletions(-) create mode 100644 apps/server/src/cloud/CloudPreferences.test.ts create mode 100644 apps/server/src/cloud/CloudPreferences.ts rename apps/server/src/scheduledTasks/{relayDeliveryProof.ts => RelayDeliveryProof.ts} (81%) create mode 100644 infra/relay/src/hooks/HeldHooks.test.ts create mode 100644 infra/relay/src/hooks/HeldHooks.ts diff --git a/apps/server/src/cloud/CloudPreferences.test.ts b/apps/server/src/cloud/CloudPreferences.test.ts new file mode 100644 index 000000000000..3368b17fc52a --- /dev/null +++ b/apps/server/src/cloud/CloudPreferences.test.ts @@ -0,0 +1,98 @@ +import { EnvironmentId } from "@t3tools/contracts"; +import { assert, it } from "@effect/vitest"; +import * as Effect from "effect/Effect"; +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, + 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 }, + 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")], + ]); + 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); + return Promise.resolve(Response.json(payload)); + }, + { preconnect: () => {} }, + ); + const dependencies = Layer.mergeAll( + Layer.mock(ServerSecretStore.ServerSecretStore)({ + get: (name) => Effect.succeed(Option.fromNullishOr(stored.get(name))), + set: (name, value) => + options.failHoldWrite && name === HOLD_WEBHOOKS_WHILE_OFFLINE_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", + ); + }), + ), +); diff --git a/apps/server/src/cloud/CloudPreferences.ts b/apps/server/src/cloud/CloudPreferences.ts new file mode 100644 index 000000000000..6e26553c0d48 --- /dev/null +++ b/apps/server/src/cloud/CloudPreferences.ts @@ -0,0 +1,97 @@ +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 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, + readHoldWebhooksWhileOffline, + readRelayConnection, +} from "./config.ts"; + +const encode = (value: boolean) => new TextEncoder().encode(String(value)); + +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. Holding webhooks while + * offline is decided by the relay, so it is told first and put back if the + * local save fails. + */ + 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 pushHoldWebhooksWhileOffline = Effect.fn("CloudPreferences.pushHoldWebhooksWhileOffline")( + function* (holdWebhooksWhileOffline: boolean) { + const connection = yield* readRelayConnection(secrets); + 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 update: CloudPreferences["Service"]["update"] = Effect.fn("CloudPreferences.update")( + function* (input) { + if (input.holdWebhooksWhileOffline !== undefined) { + const next = input.holdWebhooksWhileOffline; + const previous = yield* readHoldWebhooksWhileOffline(secrets); + yield* pushHoldWebhooksWhileOffline(next); + yield* secrets + .set(HOLD_WEBHOOKS_WHILE_OFFLINE_SECRET, encode(next)) + .pipe( + Effect.tapError(() => + previous === next + ? Effect.void + : pushHoldWebhooksWhileOffline(previous).pipe(Effect.ignore), + ), + ); + } + yield* secrets.set(PUBLISH_AGENT_ACTIVITY_SECRET, encode(input.publishAgentActivity)); + yield* awarenessRelay.requestCatchUp(); + }, + Effect.catchIf( + ServerSecretStore.isSecretStoreError, + internalError("Could not persist environment cloud preferences."), + ), + ); + + return CloudPreferences.of({ update }); +}); + +export const layer = Layer.effect(CloudPreferences, make); diff --git a/apps/server/src/cloud/http.ts b/apps/server/src/cloud/http.ts index 4cb29c37b67f..a165a8aa4c4a 100644 --- a/apps/server/src/cloud/http.ts +++ b/apps/server/src/cloud/http.ts @@ -70,7 +70,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 { @@ -102,6 +101,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,57 +1358,14 @@ 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) { + function* ( + dependencies: CloudHttpDependencies, + preferences: CloudPreferences.CloudPreferences["Service"], + 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(); + yield* preferences.update(payload); return yield* readCloudLinkState(dependencies); }, Effect.catchIf( @@ -1661,12 +1618,15 @@ 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", ({ payload }) => + cloudPreferencesHandler(dependencies, preferences, payload), + ) .handle("health", ({ payload }) => cloudEnvironmentHealthHandler(dependencies, payload)) .handle("mintCredential", ({ payload }) => cloudMintCredentialHandler(dependencies, payload)) .handle("t3MintCredential", ({ payload }) => 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..f90e1b7fbef8 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. */ @@ -33,7 +33,7 @@ const json = (status: number, body: Record, outcome: string) => const handleWebhook = ( scheduledTasks: ScheduledTaskService.ScheduledTaskService["Service"], - verifyRelayDelivery: Effect.Success, + relayDeliveryProof: RelayDeliveryProof.RelayDeliveryProof["Service"], ) => ({ params, @@ -75,7 +75,7 @@ const handleWebhook = // 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 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; @@ -136,7 +136,7 @@ export const webhookHttpApiLayer = HttpApiBuilder.group( Effect.fnUntraced(function* (handlers) { const handler = handleWebhook( yield* ScheduledTaskService.ScheduledTaskService, - yield* makeRelayDeliveryVerifier, + yield* RelayDeliveryProof.RelayDeliveryProof, ); return handlers .handleRaw("webhookPost", handler) 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..6ae3b4d58e0d 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"; @@ -266,51 +261,8 @@ 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; @@ -357,10 +309,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, 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( From c470a9d22b36644a51b0a5aa5e4ae832b8ba4842 Mon Sep 17 00:00:00 2001 From: Julius Marminge Date: Mon, 5 Oct 2026 16:16:36 -0700 Subject: [PATCH 2/9] refactor(server): errors are built where they happen SecretRequests constructs each SecretRequestError at its failure point, with the underlying cause where there is one, instead of through a forwarding helper. CloudPreferences maps a store failure at each write rather than catching a whole union with a predicate. Co-Authored-By: Claude Opus 5.5 (1M context) --- apps/server/src/cloud/CloudPreferences.ts | 27 +++++++------ apps/server/src/secrets/SecretRequests.ts | 47 +++++++++++++++-------- 2 files changed, 43 insertions(+), 31 deletions(-) diff --git a/apps/server/src/cloud/CloudPreferences.ts b/apps/server/src/cloud/CloudPreferences.ts index 6e26553c0d48..048788a5c99e 100644 --- a/apps/server/src/cloud/CloudPreferences.ts +++ b/apps/server/src/cloud/CloudPreferences.ts @@ -66,29 +66,28 @@ const make = Effect.gen(function* () { }, ); + const save = (name: string, value: boolean) => + secrets + .set(name, encode(value)) + .pipe(Effect.catch(internalError("Could not persist environment cloud preferences."))); + const update: CloudPreferences["Service"]["update"] = Effect.fn("CloudPreferences.update")( function* (input) { if (input.holdWebhooksWhileOffline !== undefined) { const next = input.holdWebhooksWhileOffline; const previous = yield* readHoldWebhooksWhileOffline(secrets); yield* pushHoldWebhooksWhileOffline(next); - yield* secrets - .set(HOLD_WEBHOOKS_WHILE_OFFLINE_SECRET, encode(next)) - .pipe( - Effect.tapError(() => - previous === next - ? Effect.void - : pushHoldWebhooksWhileOffline(previous).pipe(Effect.ignore), - ), - ); + yield* save(HOLD_WEBHOOKS_WHILE_OFFLINE_SECRET, next).pipe( + Effect.tapError(() => + previous === next + ? Effect.void + : pushHoldWebhooksWhileOffline(previous).pipe(Effect.ignore), + ), + ); } - yield* secrets.set(PUBLISH_AGENT_ACTIVITY_SECRET, encode(input.publishAgentActivity)); + yield* save(PUBLISH_AGENT_ACTIVITY_SECRET, input.publishAgentActivity); yield* awarenessRelay.requestCatchUp(); }, - Effect.catchIf( - ServerSecretStore.isSecretStoreError, - internalError("Could not persist environment cloud preferences."), - ), ); return CloudPreferences.of({ update }); 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; }); From 76b8037a622844342c443acabee405c4d1f0ee5c Mon Sep 17 00:00:00 2001 From: Julius Marminge Date: Mon, 5 Oct 2026 16:24:14 -0700 Subject: [PATCH 3/9] fix(server): a failed preferences save never leaves the relay changed The agent-activity setting was saved after the relay had been told the new hold setting, so a failure there returned an error with the relay already changed. It is now saved first, before the relay hears anything. Co-Authored-By: Claude Opus 5.5 (1M context) --- .../server/src/cloud/CloudPreferences.test.ts | 22 +++++++++++++++++-- apps/server/src/cloud/CloudPreferences.ts | 4 +++- 2 files changed, 23 insertions(+), 3 deletions(-) diff --git a/apps/server/src/cloud/CloudPreferences.test.ts b/apps/server/src/cloud/CloudPreferences.test.ts index 3368b17fc52a..90933c874c25 100644 --- a/apps/server/src/cloud/CloudPreferences.test.ts +++ b/apps/server/src/cloud/CloudPreferences.test.ts @@ -11,6 +11,7 @@ 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"; @@ -19,7 +20,7 @@ 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 }, + options: { readonly failHoldWrite?: boolean; readonly failActivityWrite?: boolean }, body: (input: { readonly preferences: CloudPreferences.CloudPreferences["Service"]; readonly stored: Map; @@ -45,7 +46,8 @@ const withService = ( Layer.mock(ServerSecretStore.ServerSecretStore)({ get: (name) => Effect.succeed(Option.fromNullishOr(stored.get(name))), set: (name, value) => - options.failHoldWrite && name === HOLD_WEBHOOKS_WHILE_OFFLINE_SECRET + (options.failHoldWrite && name === HOLD_WEBHOOKS_WHILE_OFFLINE_SECRET) || + (options.failActivityWrite && name === PUBLISH_AGENT_ACTIVITY_SECRET) ? Effect.fail( new ServerSecretStore.SecretStorePersistError({ resource: name, @@ -96,3 +98,19 @@ it.effect("puts the relay back when the local save fails", () => }), ), ); + +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", + ); + }), + ), +); diff --git a/apps/server/src/cloud/CloudPreferences.ts b/apps/server/src/cloud/CloudPreferences.ts index 048788a5c99e..c2fe9d977410 100644 --- a/apps/server/src/cloud/CloudPreferences.ts +++ b/apps/server/src/cloud/CloudPreferences.ts @@ -73,6 +73,9 @@ const make = Effect.gen(function* () { const update: CloudPreferences["Service"]["update"] = Effect.fn("CloudPreferences.update")( function* (input) { + // Saved before the relay is told anything, so a failure here leaves the + // relay and the hold setting as they were. + yield* save(PUBLISH_AGENT_ACTIVITY_SECRET, input.publishAgentActivity); if (input.holdWebhooksWhileOffline !== undefined) { const next = input.holdWebhooksWhileOffline; const previous = yield* readHoldWebhooksWhileOffline(secrets); @@ -85,7 +88,6 @@ const make = Effect.gen(function* () { ), ); } - yield* save(PUBLISH_AGENT_ACTIVITY_SECRET, input.publishAgentActivity); yield* awarenessRelay.requestCatchUp(); }, ); From 8691343f4edaa208a20146e9cfd4cea16fde919b Mon Sep 17 00:00:00 2001 From: Julius Marminge Date: Mon, 5 Oct 2026 16:33:05 -0700 Subject: [PATCH 4/9] fix(server): saving T3 Connect preferences is all or nothing The activity setting is saved before the relay hears anything; if the hold change then fails at the relay or locally, the activity setting is put back, so a failed request changes neither preference. Co-Authored-By: Claude Opus 5.5 (1M context) --- .../server/src/cloud/CloudPreferences.test.ts | 26 ++++++++++++++++-- apps/server/src/cloud/CloudPreferences.ts | 27 ++++++++++++++----- 2 files changed, 44 insertions(+), 9 deletions(-) diff --git a/apps/server/src/cloud/CloudPreferences.test.ts b/apps/server/src/cloud/CloudPreferences.test.ts index 90933c874c25..78ffae10bce0 100644 --- a/apps/server/src/cloud/CloudPreferences.test.ts +++ b/apps/server/src/cloud/CloudPreferences.test.ts @@ -20,7 +20,11 @@ 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 }, + options: { + readonly failHoldWrite?: boolean; + readonly failActivityWrite?: boolean; + readonly relayFails?: boolean; + }, body: (input: { readonly preferences: CloudPreferences.CloudPreferences["Service"]; readonly stored: Map; @@ -32,13 +36,16 @@ const withService = ( [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); - return Promise.resolve(Response.json(payload)); + return options.relayFails + ? Promise.resolve(new Response("unavailable", { status: 503 })) + : Promise.resolve(Response.json(payload)); }, { preconnect: () => {} }, ); @@ -114,3 +121,18 @@ it.effect("leaves the relay untouched when the activity setting can't be saved", }), ), ); + +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", + ); + }), + ), +); diff --git a/apps/server/src/cloud/CloudPreferences.ts b/apps/server/src/cloud/CloudPreferences.ts index c2fe9d977410..36dcd42013f8 100644 --- a/apps/server/src/cloud/CloudPreferences.ts +++ b/apps/server/src/cloud/CloudPreferences.ts @@ -6,6 +6,7 @@ import { 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 ServerSecretStore from "../auth/ServerSecretStore.ts"; import * as ServerEnvironment from "../environment/ServerEnvironment.ts"; @@ -73,18 +74,30 @@ const make = Effect.gen(function* () { const update: CloudPreferences["Service"]["update"] = Effect.fn("CloudPreferences.update")( function* (input) { - // Saved before the relay is told anything, so a failure here leaves the - // relay and the hold setting as they were. + // All or nothing: the activity setting is saved first, before the relay + // is told anything, and put back if the hold change then fails. + const previousActivity = yield* secrets + .get(PUBLISH_AGENT_ACTIVITY_SECRET) + .pipe(Effect.orElseSucceed(() => Option.none())); yield* save(PUBLISH_AGENT_ACTIVITY_SECRET, input.publishAgentActivity); if (input.holdWebhooksWhileOffline !== undefined) { const next = input.holdWebhooksWhileOffline; const previous = yield* readHoldWebhooksWhileOffline(secrets); - yield* pushHoldWebhooksWhileOffline(next); - yield* save(HOLD_WEBHOOKS_WHILE_OFFLINE_SECRET, next).pipe( + yield* pushHoldWebhooksWhileOffline(next).pipe( + Effect.andThen( + save(HOLD_WEBHOOKS_WHILE_OFFLINE_SECRET, next).pipe( + Effect.tapError(() => + previous === next + ? Effect.void + : pushHoldWebhooksWhileOffline(previous).pipe(Effect.ignore), + ), + ), + ), Effect.tapError(() => - previous === next - ? Effect.void - : pushHoldWebhooksWhileOffline(previous).pipe(Effect.ignore), + Option.match(previousActivity, { + onNone: () => secrets.remove(PUBLISH_AGENT_ACTIVITY_SECRET), + onSome: (bytes) => secrets.set(PUBLISH_AGENT_ACTIVITY_SECRET, bytes), + }).pipe(Effect.ignore), ), ); } From 6187851ea358a2ade73e29ad2066cefc4a703c82 Mon Sep 17 00:00:00 2001 From: Julius Marminge Date: Mon, 5 Oct 2026 16:41:13 -0700 Subject: [PATCH 5/9] refactor(server,relay): read services from context instead of passing them in The T3 Connect config readers yield the secret store themselves rather than taking it as an argument. The webhook and preferences handlers live inside their HttpApi groups and close over the services the group acquires. The relay signs delivery proofs with the settings its forwarder already holds. Co-Authored-By: Claude Opus 5.5 (1M context) --- apps/server/src/cloud/CloudPreferences.ts | 5 +- apps/server/src/cloud/config.ts | 37 ++- apps/server/src/cloud/http.ts | 34 ++- apps/server/src/relay/HeldHooksWaker.ts | 6 +- .../server/src/scheduledTasks/webhookRoute.ts | 214 +++++++++--------- infra/relay/src/hooks/HookForwarder.ts | 56 +++-- 6 files changed, 164 insertions(+), 188 deletions(-) diff --git a/apps/server/src/cloud/CloudPreferences.ts b/apps/server/src/cloud/CloudPreferences.ts index 36dcd42013f8..1df5f0b04769 100644 --- a/apps/server/src/cloud/CloudPreferences.ts +++ b/apps/server/src/cloud/CloudPreferences.ts @@ -44,10 +44,11 @@ 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(secrets); + const connection = yield* readRelayConnection.pipe(withSecrets); if (connection === null) { return yield* new EnvironmentHttpBadRequestError({ message: "Link this environment to T3 Connect first.", @@ -82,7 +83,7 @@ const make = Effect.gen(function* () { yield* save(PUBLISH_AGENT_ACTIVITY_SECRET, input.publishAgentActivity); if (input.holdWebhooksWhileOffline !== undefined) { const next = input.holdWebhooksWhileOffline; - const previous = yield* readHoldWebhooksWhileOffline(secrets); + const previous = yield* readHoldWebhooksWhileOffline.pipe(withSecrets); yield* pushHoldWebhooksWhileOffline(next).pipe( Effect.andThen( save(HOLD_WEBHOOKS_WHILE_OFFLINE_SECRET, next).pipe( 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 a165a8aa4c4a..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, @@ -88,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, @@ -1358,22 +1355,6 @@ const cloudUnlinkHandler = Effect.fn("environment.cloud.unlink")( ), ); -const cloudPreferencesHandler = Effect.fn("environment.cloud.preferences")( - function* ( - dependencies: CloudHttpDependencies, - preferences: CloudPreferences.CloudPreferences["Service"], - payload: EnvironmentCloudPreferencesRequest, - ) { - yield* requireEnvironmentScope(AuthRelayWriteScope); - yield* preferences.update(payload); - 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 @@ -1624,8 +1605,19 @@ export const connectHttpApiLayer = HttpApiBuilder.group( .handle("relayConfig", ({ payload }) => cloudRelayConfigHandler(dependencies, payload)) .handle("linkState", () => cloudLinkStateHandler(dependencies)) .handle("unlink", () => cloudUnlinkHandler(dependencies)) - .handle("preferences", ({ payload }) => - cloudPreferencesHandler(dependencies, preferences, 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)) 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/webhookRoute.ts b/apps/server/src/scheduledTasks/webhookRoute.ts index f90e1b7fbef8..dec01a5e6b15 100644 --- a/apps/server/src/scheduledTasks/webhookRoute.ts +++ b/apps/server/src/scheduledTasks/webhookRoute.ts @@ -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"], - relayDeliveryProof: RelayDeliveryProof.RelayDeliveryProof["Service"], - ) => - ({ - 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"); - } - }); - export const webhookHttpApiLayer = HttpApiBuilder.group( EnvironmentHttpApi, "webhooks", Effect.fnUntraced(function* (handlers) { - const handler = handleWebhook( - yield* ScheduledTaskService.ScheduledTaskService, - yield* RelayDeliveryProof.RelayDeliveryProof, - ); + 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/infra/relay/src/hooks/HookForwarder.ts b/infra/relay/src/hooks/HookForwarder.ts index 6ae3b4d58e0d..8e32d9c1f185 100644 --- a/infra/relay/src/hooks/HookForwarder.ts +++ b/infra/relay/src/hooks/HookForwarder.ts @@ -138,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"]); @@ -269,6 +241,33 @@ const make = Effect.gen(function* () { 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, ) { @@ -347,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, From 9d1e4f6215a5859f791e0a4194d767ce916272ec Mon Sep 17 00:00:00 2001 From: Julius Marminge Date: Mon, 5 Oct 2026 16:42:21 -0700 Subject: [PATCH 6/9] fix(server): a preferences save that can't read the current value changes nothing The current activity setting was read with failures treated as unset, so a rollback after a failed hold change could delete the real setting. A failed read now fails the request before anything is written. Co-Authored-By: Claude Opus 5.5 (1M context) --- .../server/src/cloud/CloudPreferences.test.ts | 24 ++++++++++++++++++- apps/server/src/cloud/CloudPreferences.ts | 4 +++- 2 files changed, 26 insertions(+), 2 deletions(-) diff --git a/apps/server/src/cloud/CloudPreferences.test.ts b/apps/server/src/cloud/CloudPreferences.test.ts index 78ffae10bce0..0ceada6dd407 100644 --- a/apps/server/src/cloud/CloudPreferences.test.ts +++ b/apps/server/src/cloud/CloudPreferences.test.ts @@ -24,6 +24,7 @@ const withService = ( readonly failHoldWrite?: boolean; readonly failActivityWrite?: boolean; readonly relayFails?: boolean; + readonly failActivityRead?: boolean; }, body: (input: { readonly preferences: CloudPreferences.CloudPreferences["Service"]; @@ -51,7 +52,15 @@ const withService = ( ); const dependencies = Layer.mergeAll( Layer.mock(ServerSecretStore.ServerSecretStore)({ - get: (name) => Effect.succeed(Option.fromNullishOr(stored.get(name))), + get: (name) => + options.failActivityRead && name === PUBLISH_AGENT_ACTIVITY_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) @@ -136,3 +145,16 @@ it.effect("keeps both settings unchanged when the relay refuses the hold change" }), ), ); + +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"); + }), + ), +); diff --git a/apps/server/src/cloud/CloudPreferences.ts b/apps/server/src/cloud/CloudPreferences.ts index 1df5f0b04769..c9f79d616f42 100644 --- a/apps/server/src/cloud/CloudPreferences.ts +++ b/apps/server/src/cloud/CloudPreferences.ts @@ -77,9 +77,11 @@ const make = Effect.gen(function* () { 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. + // A failed read stops here, before anything changes: guessing "unset" + // would make a later rollback delete the real setting. const previousActivity = yield* secrets .get(PUBLISH_AGENT_ACTIVITY_SECRET) - .pipe(Effect.orElseSucceed(() => Option.none())); + .pipe(Effect.catch(internalError("Could not read environment cloud preferences."))); yield* save(PUBLISH_AGENT_ACTIVITY_SECRET, input.publishAgentActivity); if (input.holdWebhooksWhileOffline !== undefined) { const next = input.holdWebhooksWhileOffline; From a940f2c77cd068c9ec37751b2d1e6c8ae888296a Mon Sep 17 00:00:00 2001 From: Julius Marminge Date: Mon, 5 Oct 2026 16:43:09 -0700 Subject: [PATCH 7/9] docs(server): CloudPreferences describes both rollbacks, and logs a failed one Co-Authored-By: Claude Opus 5.5 (1M context) --- apps/server/src/cloud/CloudPreferences.ts | 16 +++++++++++----- 1 file changed, 11 insertions(+), 5 deletions(-) diff --git a/apps/server/src/cloud/CloudPreferences.ts b/apps/server/src/cloud/CloudPreferences.ts index c9f79d616f42..e2038642286e 100644 --- a/apps/server/src/cloud/CloudPreferences.ts +++ b/apps/server/src/cloud/CloudPreferences.ts @@ -21,6 +21,10 @@ import { 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 }))), @@ -30,9 +34,11 @@ export class CloudPreferences extends Context.Service< CloudPreferences, { /** - * Saves this environment's T3 Connect preferences. Holding webhooks while - * offline is decided by the relay, so it is told first and put back if the - * local save fails. + * 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, @@ -92,7 +98,7 @@ const make = Effect.gen(function* () { Effect.tapError(() => previous === next ? Effect.void - : pushHoldWebhooksWhileOffline(previous).pipe(Effect.ignore), + : pushHoldWebhooksWhileOffline(previous).pipe(Effect.catch(rollbackFailed)), ), ), ), @@ -100,7 +106,7 @@ const make = Effect.gen(function* () { Option.match(previousActivity, { onNone: () => secrets.remove(PUBLISH_AGENT_ACTIVITY_SECRET), onSome: (bytes) => secrets.set(PUBLISH_AGENT_ACTIVITY_SECRET, bytes), - }).pipe(Effect.ignore), + }).pipe(Effect.catch(rollbackFailed)), ), ); } From f48ff7c1daa65b3b46fcd3df37dc2310187af8f4 Mon Sep 17 00:00:00 2001 From: Julius Marminge Date: Mon, 5 Oct 2026 16:45:29 -0700 Subject: [PATCH 8/9] fix(server): overlapping T3 Connect preference saves run one at a time Two saves in flight could each leave one setting behind: the second's activity write landed while the first was still waiting on the relay. Updates now run one at a time. Co-Authored-By: Claude Opus 5.5 (1M context) --- .../server/src/cloud/CloudPreferences.test.ts | 43 +++++++++++++++++++ apps/server/src/cloud/CloudPreferences.ts | 5 +++ 2 files changed, 48 insertions(+) diff --git a/apps/server/src/cloud/CloudPreferences.test.ts b/apps/server/src/cloud/CloudPreferences.test.ts index 0ceada6dd407..a6e60f0a0918 100644 --- a/apps/server/src/cloud/CloudPreferences.test.ts +++ b/apps/server/src/cloud/CloudPreferences.test.ts @@ -1,6 +1,7 @@ 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"; @@ -25,6 +26,11 @@ const withService = ( readonly failActivityWrite?: boolean; readonly relayFails?: boolean; readonly failActivityRead?: 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"]; @@ -44,6 +50,10 @@ const withService = ( (_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)); @@ -158,3 +168,36 @@ it.effect("changes nothing when the current activity setting can't be read", () }), ), ); + +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", + ); + }), + ); +}); diff --git a/apps/server/src/cloud/CloudPreferences.ts b/apps/server/src/cloud/CloudPreferences.ts index e2038642286e..9ed60f3d4efe 100644 --- a/apps/server/src/cloud/CloudPreferences.ts +++ b/apps/server/src/cloud/CloudPreferences.ts @@ -7,6 +7,7 @@ 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"; @@ -79,6 +80,9 @@ const make = Effect.gen(function* () { .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 @@ -112,6 +116,7 @@ const make = Effect.gen(function* () { } yield* awarenessRelay.requestCatchUp(); }, + updateLock.withPermits(1), ); return CloudPreferences.of({ update }); From 874f3098e5936cf6785213b30e53273ec101b692 Mon Sep 17 00:00:00 2001 From: Julius Marminge Date: Mon, 5 Oct 2026 16:46:43 -0700 Subject: [PATCH 9/9] fix(server): a preferences save that can't read the hold setting changes nothing The current hold setting was read through a helper that treats a failed read as off, so a rollback could switch the relay off when it had been on. Both current values are now read up front, and a failed read fails the request before anything changes. Co-Authored-By: Claude Opus 5.5 (1M context) --- .../server/src/cloud/CloudPreferences.test.ts | 18 +++++++++++++++- apps/server/src/cloud/CloudPreferences.ts | 21 ++++++++++++------- 2 files changed, 30 insertions(+), 9 deletions(-) diff --git a/apps/server/src/cloud/CloudPreferences.test.ts b/apps/server/src/cloud/CloudPreferences.test.ts index a6e60f0a0918..c0e66bc49f75 100644 --- a/apps/server/src/cloud/CloudPreferences.test.ts +++ b/apps/server/src/cloud/CloudPreferences.test.ts @@ -26,6 +26,7 @@ const withService = ( 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; @@ -63,7 +64,8 @@ const withService = ( const dependencies = Layer.mergeAll( Layer.mock(ServerSecretStore.ServerSecretStore)({ get: (name) => - options.failActivityRead && name === PUBLISH_AGENT_ACTIVITY_SECRET + (options.failActivityRead && name === PUBLISH_AGENT_ACTIVITY_SECRET) || + (options.failHoldRead && name === HOLD_WEBHOOKS_WHILE_OFFLINE_SECRET) ? Effect.fail( new ServerSecretStore.SecretStoreReadError({ resource: name, @@ -201,3 +203,17 @@ it.effect("two overlapping updates are applied one after the other", () => { }), ); }); + +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 index 9ed60f3d4efe..7e765ce76959 100644 --- a/apps/server/src/cloud/CloudPreferences.ts +++ b/apps/server/src/cloud/CloudPreferences.ts @@ -16,7 +16,6 @@ import { makeRelayEnvironmentClient } from "../relay/relayEnvironmentClient.ts"; import { HOLD_WEBHOOKS_WHILE_OFFLINE_SECRET, PUBLISH_AGENT_ACTIVITY_SECRET, - readHoldWebhooksWhileOffline, readRelayConnection, } from "./config.ts"; @@ -86,16 +85,22 @@ const make = Effect.gen(function* () { 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. - // A failed read stops here, before anything changes: guessing "unset" - // would make a later rollback delete the real setting. - const previousActivity = yield* secrets - .get(PUBLISH_AGENT_ACTIVITY_SECRET) - .pipe(Effect.catch(internalError("Could not read environment cloud preferences."))); + // 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 = yield* readHoldWebhooksWhileOffline.pipe(withSecrets); + 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(