import * as Alchemy from "alchemy";
import * as Cloudflare from "alchemy/Cloudflare";
import * as Drizzle from "alchemy/Drizzle/Postgres";
import * as Config from "effect/Config";
import * as Cause from "effect/Cause";
import * as DateTime from "effect/DateTime";
import * as Crypto from "effect/Crypto";
import * as Effect from "effect/Effect";
import * as Layer from "effect/Layer";
import * as Option from "effect/Option";
import * as Redacted from "effect/Redacted";
import * as Stream from "effect/Stream";
import * as Etag from "effect/unstable/http/Etag";
import * as HttpPlatform from "effect/unstable/http/HttpPlatform";
import * as HttpRouter from "effect/unstable/http/HttpRouter";
import * as HttpApiBuilder from "effect/unstable/httpapi/HttpApiBuilder";
import * as HttpApiScalar from "effect/unstable/httpapi/HttpApiScalar";
import { RelayApi } from "@t3tools/contracts/relay";
import {
clientApi,
dpopClientApi,
healthApi,
metadataApi,
mobileApi,
RELAY_HTTP_ROUTER_CONFIG,
relayClientAuthLayer,
relayDpopClientAuthLayer,
relayCors,
relayDocsRedirectRoute,
relayEnvironmentAuthLayer,
relayNotFoundRoute,
serverApi,
traceRelayHttpRequestWith,
tokenApi,
withoutCapturedParentSpan,
} from "./http/Api.ts";
import { ManagedEndpointZone, RelayApiZone, RelayDeploymentConfig } from "./zone.ts";
import { makeRelayTraceLayer, RelayObservability } from "./observability.ts";
import * as DeliveryAttempts from "./agentActivity/DeliveryAttempts.ts";
import * as AgentActivityRows from "./agentActivity/AgentActivityRows.ts";
import * as Devices from "./agentActivity/Devices.ts";
import * as DpopProofs from "./auth/DpopProofs.ts";
import * as RelayTokens from "./auth/RelayTokens.ts";
import * as EnvironmentCredentials from "./environments/EnvironmentCredentials.ts";
import * as EnvironmentLinks from "./environments/EnvironmentLinks.ts";
import * as ManagedEndpointAllocations from "./environments/ManagedEndpointAllocations.ts";
import * as LiveActivities from "./agentActivity/LiveActivities.ts";
import * as RelayDb from "./db.ts";
import {
RelayApnsDeliveryDeadLetterQueue,
RelayApnsDeliveryQueue,
RelayFcmDeliveryQueue,
RelayFcmDeliveryDeadLetterQueue,
} from "./queues.ts";
import * as WebCrypto from "./WebCrypto.ts";
import * as FcmAssertionSigner from "./agentActivity/FcmAssertionSigner.ts";
import * as FcmClient from "./agentActivity/FcmClient.ts";
import * as FcmDeliveryQueueSender from "./agentActivity/FcmDeliveryQueueSender.ts";
import * as FcmDeliveries from "./agentActivity/FcmDeliveries.ts";
import * as FcmDeliveryQueueConsumer from "./agentActivity/FcmDeliveryQueueConsumer.ts";
import * as RelayConfiguration from "./Config.ts";
import * as AgentActivityPublisher from "./agentActivity/AgentActivityPublisher.ts";
import * as ApnsClient from "./agentActivity/ApnsClient.ts";
import * as ApnsProviderTokens from "./agentActivity/ApnsProviderTokens.ts";
import * as ApnsDeliveryQueue from "./agentActivity/ApnsDeliveryQueue.ts";
import * as ApnsDeliveries from "./agentActivity/ApnsDeliveries.ts";
import * as EnvironmentConnector from "./environments/EnvironmentConnector.ts";
import * as EnvironmentLinker from "./environments/EnvironmentLinker.ts";
import * as EnvironmentPublishSignatures from "./environments/EnvironmentPublishSignatures.ts";
import * as ManagedEndpointProvider from "./environments/ManagedEndpointProvider.ts";
import * as ManagedEndpointReaper from "./environments/ManagedEndpointReaper.ts";
import * as ManagedTunnelLimits from "./environments/ManagedTunnelLimits.ts";
import * as MobileRegistrations from "./agentActivity/MobileRegistrations.ts";
const webcryptoLayer = Layer.succeed(
Crypto.Crypto,
Crypto.make({
randomBytes: (size) => globalThis.crypto.getRandomValues(new Uint8Array(size)),
digest: (algorithm, data) =>
Effect.promise(async () => {
const input = new Uint8Array(data.length);
input.set(data);
return new Uint8Array(await globalThis.crypto.subtle.digest(algorithm, input.buffer));
}),
}),
);
const httpPlatformNotSupportedLayer = Layer.succeed(HttpPlatform.HttpPlatform, {
platform: "web",
compression: {
algorithms: new Set<HttpPlatform.CompressionAlgorithm>(),
compressResponse: (response) => Effect.succeed(response),
},
fileResponse: () => Effect.die("Relay API does not serve filesystem responses"),
fileWebResponse: () => Effect.die("Relay API does not serve file responses"),
});
const relayApiLayer = Layer.mergeAll(
healthApi,
metadataApi,
mobileApi,
clientApi,
tokenApi,
dpopClientApi,
serverApi,
);
const CloudMintKeyPair = Alchemy.KeyPair("CloudMintKeyPair");
const ApnsDeliveryJobSigningSecret = Alchemy.makeRandom("ApnsDeliveryJobSigningSecret", {
bytes: 32,
});
export class Api extends Cloudflare.Worker<Api, {}>()("Api") {}
export const ApiLive = Api.make(
RelayDeploymentConfig.pipe(
Effect.map(({ relayPublicDomain }) => ({
main: import.meta.filename,
compatibility: {
date: "2026-05-22",
flags: ["nodejs_compat"],
},
domain: relayPublicDomain,
})),
Effect.orDie,
),
Effect.gen(function* () {
//
// 1. Provision Infrastructure for the Worker to use
//
const { relayPublicOrigin, stage } = yield* RelayDeploymentConfig;
const apnsDeliveryQueue = yield* RelayApnsDeliveryQueue;
const apnsDeliveryDeadLetterQueue = yield* RelayApnsDeliveryDeadLetterQueue;
const fcmDeliveryQueue = yield* RelayFcmDeliveryQueue;
const fcmDeliveryDeadLetterQueue = yield* RelayFcmDeliveryDeadLetterQueue;
const cloudMintKeyPair = yield* CloudMintKeyPair;
const relayApiZone = yield* RelayApiZone;
const managedEndpointZone = yield* ManagedEndpointZone;
const randomApnsDeliveryJobSigningSecret = yield* ApnsDeliveryJobSigningSecret;
const observability = yield* RelayObservability;
//
// 2. Create bindings
//
const apnsEnabled = yield* Config.Boolean("APNS_ENABLED").pipe(Config.withDefault(true));
const apnsCredentials = apnsEnabled
? {
environment: yield* Config.schema(RelayConfiguration.ApnsEnvironment, "APNS_ENVIRONMENT"),
teamId: yield* Config.String("APNS_TEAM_ID"),
keyId: yield* Config.String("APNS_KEY_ID"),
bundleId: yield* Config.String("APNS_BUNDLE_ID"),
privateKey: yield* Config.Redacted("APNS_PRIVATE_KEY"),
}
: null;
const fcmServiceAccount = Option.getOrUndefined(
Option.filter(
yield* Config.option(Config.Redacted("FCM_SERVICE_ACCOUNT")),
(value) => Redacted.value(value).trim().length > 0,
),
);
const apnsDeliveryJobSigningSecret = yield* randomApnsDeliveryJobSigningSecret;
const apnsDeliveryQueueSender = yield* Cloudflare.Queues.WriteQueue(apnsDeliveryQueue);
const fcmDeliveryQueueSender = yield* Cloudflare.Queues.WriteQueue(fcmDeliveryQueue);
const axiomDatasetName = yield* observability.traces.name;
const axiomIngestToken = yield* observability.workerIngestToken.token;
const axiomTracesEndpoint = yield* observability.traces.otelTracesEndpoint;
const clerkSecretKey = yield* Config.Redacted("CLERK_SECRET_KEY");
const clerkPublishableKey = yield* Config.String("CLERK_PUBLISHABLE_KEY");
const clerkJwtAudience = yield* Config.String("CLERK_JWT_AUDIENCE");
const cloudMintPrivateKey = yield* cloudMintKeyPair.privateKey;
const cloudMintPublicKey = yield* cloudMintKeyPair.publicKey;
const hyperdrive = yield* Cloudflare.Hyperdrive.Connect(yield* RelayDb.RelayHyperdrive);
const db = yield* Drizzle.Postgres(hyperdrive.connectionString);
const managedEndpointTunnelBinding = yield* Cloudflare.Tunnel.ReadWriteTunnel();
// Keep Worker custom-domain reconciliation ordered after API zone provisioning.
yield* yield* relayApiZone.zoneId;
const managedEndpointDnsBinding = yield* Cloudflare.DNS.ReadWriteDns(managedEndpointZone);
const managedEndpointZoneName = yield* managedEndpointZone.name;
const managedEndpointCleanupMode = yield* RelayConfiguration.managedEndpointCleanupModeConfig;
//
// 3. Runtime layers and app construction
//
const alchemyRuntimeContext: Alchemy.BaseRuntimeContext = yield* Cloudflare.Worker;
const loadSettings = Effect.gen(function* () {
return RelayConfiguration.RelayConfiguration.of({
relayIssuer: relayPublicOrigin,
...(fcmServiceAccount ? { fcmServiceAccount } : {}),
apns: apnsCredentials,
apnsDeliveryJobSigningSecret: yield* apnsDeliveryJobSigningSecret,
clerkSecretKey,
clerkPublishableKey,
clerkJwtAudience,
cloudMintPrivateKey: yield* cloudMintPrivateKey,
cloudMintPublicKey: yield* cloudMintPublicKey,
managedEndpointBaseDomain: yield* managedEndpointZoneName,
managedEndpointNamespace: stage,
managedEndpointCleanupMode,
});
});
const relayTraceLayer = Layer.unwrap(
Effect.all({
tracesDatasetName: axiomDatasetName,
tracesEndpoint: axiomTracesEndpoint,
ingestToken: axiomIngestToken,
}).pipe(Effect.map(makeRelayTraceLayer)),
);
const runtimeLayer = Layer.empty.pipe(
Layer.provideMerge(MobileRegistrations.layer),
Layer.provideMerge(AgentActivityPublisher.layer),
Layer.provideMerge(EnvironmentConnector.layer),
Layer.provideMerge(EnvironmentLinker.layer),
Layer.provideMerge(
Layer.merge(EnvironmentPublishSignatures.layer, ManagedEndpointReaper.layer),
),
Layer.provideMerge(
ManagedEndpointProvider.layerCloudflareBindings(
managedEndpointTunnelBinding,
managedEndpointDnsBinding,
alchemyRuntimeContext,
),
),
Layer.provideMerge(DpopProofs.layer),
Layer.provideMerge(ApnsDeliveries.layer),
Layer.provideMerge(
FcmDeliveries.layer.pipe(
Layer.provide(
Layer.succeed(FcmDeliveryQueueSender.FcmDeliveryQueueSender, {
send: (body) =>
fcmDeliveryQueueSender
.send(body)
.pipe(Effect.provideService(Alchemy.RuntimeContext, alchemyRuntimeContext)),
}),
),
Layer.provideMerge(
FcmClient.layer.pipe(
Layer.provide(FcmAssertionSigner.layer),
Layer.provide(
Layer.succeed(WebCrypto.WebCrypto, { subtle: globalThis.crypto.subtle }),
),
),
),
),
),
Layer.provideMerge(ApnsClient.layer.pipe(Layer.provideMerge(ApnsProviderTokens.layer))),
Layer.provideMerge(
ApnsDeliveryQueue.layerCloudflareQueues(apnsDeliveryQueueSender, alchemyRuntimeContext),
),
Layer.provideMerge(Layer.mergeAll(AgentActivityRows.layer, Devices.layer)),
Layer.provideMerge(EnvironmentCredentials.layer),
Layer.provideMerge(
Layer.mergeAll(
EnvironmentLinks.layer,
ManagedEndpointAllocations.layer,
ManagedTunnelLimits.layer,
),
),
Layer.provideMerge(LiveActivities.layer),
Layer.provideMerge(DeliveryAttempts.layer),
Layer.provideMerge(RelayTokens.layer),
Layer.provideMerge(
RelayDb.RelayTransactions.layer.pipe(
Layer.provideMerge(Layer.succeed(RelayDb.RelayDb, db)),
),
),
Layer.provideMerge(Layer.effect(RelayConfiguration.RelayConfiguration, loadSettings)),
Layer.provideMerge(webcryptoLayer),
);
const appLayer = relayApiLayer.pipe(
Layer.provideMerge(relayClientAuthLayer),
Layer.provideMerge(relayDpopClientAuthLayer),
Layer.provideMerge(relayEnvironmentAuthLayer),
Layer.provide(runtimeLayer),
);
yield* Cloudflare.Queues.consumeQueueMessages<unknown>(
apnsDeliveryQueue,
{
batchSize: 10,
maxRetries: 5,
maxWaitTime: "5 seconds",
retryDelay: "30 seconds",
deadLetterQueue: apnsDeliveryDeadLetterQueue.queueName as unknown as string,
},
(stream) =>
stream.pipe(
Stream.withSpan("relay.apn_delivery_queue.process_batch"),
Stream.runForEach((message) =>
ApnsDeliveries.ApnsDeliveries.pipe(
Effect.flatMap((deliveries) => deliveries.processSignedJob(message.body)),
Effect.withSpan("relay.apn_delivery_queue.process_message"),
),
),
Effect.provide(runtimeLayer),
),
);
yield* Cloudflare.Queues.consumeQueueMessages<unknown>(
fcmDeliveryQueue,
{
batchSize: 10,
maxRetries: 5,
maxWaitTime: "1 second",
retryDelay: "30 seconds",
deadLetterQueue: fcmDeliveryDeadLetterQueue.queueName as unknown as string,
},
(stream) =>
stream.pipe(
Stream.withSpan("relay.fcm_delivery_queue.process_batch"),
Stream.runForEach(FcmDeliveryQueueConsumer.processMessage),
Effect.provide(runtimeLayer),
),
);
yield* Cloudflare.Workers.cron("*/5 * * * *", () =>
Effect.all(
[
DpopProofs.DpopProofReplay.pipe(
Effect.flatMap((dpopProofs) => dpopProofs.pruneExpired),
// Keep completed thread rows long enough to show their final state.
Effect.andThen(
Effect.all([AgentActivityRows.AgentActivityRows, DateTime.now]).pipe(
Effect.flatMap(([activityRows, now]) =>
activityRows.pruneTerminal({
updatedBefore: DateTime.formatIso(DateTime.subtract(now, { minutes: 30 })),
}),
),
),
),
Effect.catchCause((cause) =>
Cause.hasInterrupts(cause)
? Effect.interrupt
: Effect.logWarning("Failed to prune expired relay state", { cause }),
),
),
ManagedEndpointReaper.ManagedEndpointReaper.pipe(
Effect.flatMap((reaper) => reaper.sweep.pipe(Effect.timeout("2 minutes"))),
Effect.tap((result) =>
result.scanned > 0
? Effect.logInfo("Finished managed tunnel cleanup", result)
: Effect.void,
),
Effect.catchCause((cause) =>
Cause.hasInterrupts(cause)
? Effect.interrupt
: Effect.logWarning("Failed to clean up inactive managed tunnels", { cause }),
),
),
],
{ concurrency: 2, discard: true },
).pipe(
Effect.withSpan("relay.cron.prune_expired_state"),
// Export cron spans to Axiom like HTTP spans; the scope flushes them before the run ends.
Effect.provide(Layer.merge(runtimeLayer, relayTraceLayer)),
),
);
const fetch = Layer.merge(
Layer.mergeAll(
HttpApiBuilder.layer(RelayApi, { openapiPath: "/openapi.json" }).pipe(
Layer.provide(appLayer),
),
HttpApiScalar.layer(RelayApi, { path: "/docs" }),
relayDocsRedirectRoute,
).pipe(Layer.provide([Etag.layerWeak, httpPlatformNotSupportedLayer, relayCors])),
relayNotFoundRoute,
).pipe(
HttpRouter.toHttpEffect,
Effect.provideService(HttpRouter.RouterConfig, RELAY_HTTP_ROUTER_CONFIG),
withoutCapturedParentSpan,
Effect.flatMap((httpEffect) => traceRelayHttpRequestWith(httpEffect, relayTraceLayer)),
);
return { fetch };
}).pipe(
Effect.provide(
Layer.empty.pipe(
Layer.provideMerge(Cloudflare.Hyperdrive.ConnectBinding),
Layer.provideMerge(Cloudflare.Workers.CronEventSourceLive),
Layer.provideMerge(Cloudflare.Queues.WriteQueueBinding),
Layer.provideMerge(Cloudflare.Queues.EventSourceLive),
Layer.provideMerge(Cloudflare.Tunnel.ReadWriteTunnelBinding),
Layer.provideMerge(Cloudflare.DNS.ReadWriteDnsHttp),
),
),
),
);
export default ApiLive;