infra/relay/src/http/Api.ts

import { createClerkClient, verifyToken } from "@clerk/backend";
import { sql as drizzleSql } from "drizzle-orm";
import * as Crypto from "effect/Crypto";
import * as Context from "effect/Context";
import * as DateTime from "effect/DateTime";
import * as Duration from "effect/Duration";
import * as Effect from "effect/Effect";
import * as Layer from "effect/Layer";
import * as Option from "effect/Option";
import * as Record from "effect/Record";
import * as Redacted from "effect/Redacted";
import * as Schema from "effect/Schema";
import * as Tracer from "effect/Tracer";
import * as HttpEffect from "effect/unstable/http/HttpEffect";
import * as HttpMiddleware from "effect/unstable/http/HttpMiddleware";
import * as HttpRouter from "effect/unstable/http/HttpRouter";
import * as HttpServerRequest from "effect/unstable/http/HttpServerRequest";
import * as HttpServerResponse from "effect/unstable/http/HttpServerResponse";
import * as HttpTraceContext from "effect/unstable/http/HttpTraceContext";
import * as HttpApiBuilder from "effect/unstable/httpapi/HttpApiBuilder";
import * as HttpApiError from "effect/unstable/httpapi/HttpApiError";
import { encodeOAuthScope } from "@t3tools/shared/oauthScope";
import { httpHeaderRedactionLayer } from "@t3tools/shared/httpObservability";

import {
  RelayApi,
  RelayAgentActivityPublishProofExpiredError,
  RelayAgentActivityPublishProofInvalidError,
  RelayClientAuth,
  RelayClientPrincipal,
  RelayAccessTokenType,
  RelayDpopClientAuth,
  RelayEnvironmentConnectScope,
  RelayEnvironmentStatusScope,
  RelayMobileRegistrationScope,
  RelayAuthInvalidError,
  type RelayAuthInvalidReason,
  type RelayDpopFailureReason,
  RelayEnvironmentAuth,
  RelayEnvironmentConnectNotAuthorizedError,
  RelayEnvironmentEndpointTimedOutError,
  RelayEnvironmentEndpointUnavailableError,
  RelayEnvironmentLinkFailedError,
  RelayEnvironmentLinkProofExpiredError,
  RelayEnvironmentLinkProofInvalidError,
  RelayEnvironmentLinkUnavailableError,
  RelayEnvironmentLinkLimitExceededError,
  RelayEnvironmentPrincipal,
  type RelayEnvironmentConnectRequest,
  type RelayManagedEndpointOrigin,
  RelayManagedEndpointRecoveryProofPayload,
  type RelayDpopAccessTokenScope,
  RelayInternalError,
} from "@t3tools/contracts/relay";
import {
  normalizeRelayIssuer,
  RELAY_MANAGED_TUNNEL_RECOVERY_TYP,
  verifyRelayJwt,
} from "@t3tools/shared/relayJwt";

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 LiveActivities from "../agentActivity/LiveActivities.ts";
import * as RelayConfiguration from "../Config.ts";
import * as AgentActivityPublisher from "../agentActivity/AgentActivityPublisher.ts";
import * as EnvironmentConnector from "../environments/EnvironmentConnector.ts";
import * as EnvironmentLinker from "../environments/EnvironmentLinker.ts";
import * as ManagedEndpointProvider from "../environments/ManagedEndpointProvider.ts";
import * as ManagedEndpointAllocations from "../environments/ManagedEndpointAllocations.ts";
import * as EnvironmentPublishSignatures from "../environments/EnvironmentPublishSignatures.ts";
import * as MobileRegistrations from "../agentActivity/MobileRegistrations.ts";
import { withSpanAttributes } from "../observability.ts";
import * as RelayDb from "../db.ts";

// Delegated thread IDs carry escaped command provenance and exceed the router's
// default 100-character path parameter limit. Match the environment server.
export const RELAY_HTTP_ROUTER_CONFIG = {
  maxParamLength: 512,
} as const;

const relayCorsAllowedMethods = ["GET", "POST", "DELETE", "OPTIONS"] as const;
const relayCorsAllowedHeaders = [
  "authorization",
  "b3",
  "traceparent",
  "content-type",
  "dpop",
] as const;
const relayCorsExposedHeaders = ["traceparent", "www-authenticate"] as const;

const relayCorsHeaders = {
  "access-control-allow-origin": "*",
  "access-control-expose-headers": relayCorsExposedHeaders.join(","),
} as const;

const relayCorsPreflightHeaders = {
  ...relayCorsHeaders,
  "access-control-allow-methods": relayCorsAllowedMethods.join(","),
  "access-control-allow-headers": relayCorsAllowedHeaders.join(","),
  "access-control-max-age": "86400",
} as const;

const decodeManagedTunnelRecoveryProof = Schema.decodeUnknownEffect(
  RelayManagedEndpointRecoveryProofPayload,
);

const appendRelayCredentialResponseHeaders = HttpEffect.appendPreResponseHandler(
  (_request, response) =>
    Effect.succeed(
      HttpServerResponse.setHeaders(response, {
        "cache-control": "no-store",
        pragma: "no-cache",
      }),
    ),
);

const appendRelayDpopChallengeHeader = HttpEffect.appendPreResponseHandler((_request, response) =>
  Effect.succeed(
    response.status === 401
      ? HttpServerResponse.setHeader(response, "www-authenticate", "DPoP")
      : response,
  ),
);

const appendRelayTraceContextResponseHeader = Effect.gen(function* () {
  const span = yield* Effect.currentParentSpan;
  if (span._tag !== "Span") {
    return;
  }
  const traceparent = HttpTraceContext.toHeaders(span).traceparent;
  if (traceparent === undefined) {
    return;
  }
  yield* HttpEffect.appendPreResponseHandler((_request, response) =>
    Effect.succeed(HttpServerResponse.setHeader(response, "traceparent", traceparent)),
  );
}).pipe(Effect.ignore);

export const relayCors = HttpRouter.middleware(
  Effect.fnUntraced(function* <E, R>(
    httpEffect: Effect.Effect<
      HttpServerResponse.HttpServerResponse,
      E,
      HttpServerRequest.HttpServerRequest | R
    >,
  ) {
    const request = yield* HttpServerRequest.HttpServerRequest;
    if (request.method === "OPTIONS") {
      return HttpServerResponse.empty({
        status: 204,
        headers: relayCorsPreflightHeaders,
      });
    }
    const response = yield* httpEffect;
    return HttpServerResponse.setHeaders(response, relayCorsHeaders);
  }),
  { global: true },
);

export const relayNotFoundRoute = HttpRouter.add(
  "*",
  "/*",
  HttpServerResponse.empty({ status: 404 }),
);

export const relayDocsRedirectRoute = HttpRouter.add(
  "GET",
  "/",
  HttpServerResponse.redirect("/docs"),
);

// Shorter than the mobile client's 10s request timeout on purpose: when a
// request hangs (e.g. a stuck upstream query), the client would otherwise
// abort first, the invocation would die with the request span still open, and
// the batched spans would never export — leaving no server-side trace at all.
// Failing server-side first turns the hang into a completed 504 whose trace
// contains the exact child span that stalled, and the response still carries
// the traceparent back to the client.
export const RELAY_REQUEST_DEADLINE_MS = 9_000;

const relayRequestDeadline = <E, R>(
  httpEffect: Effect.Effect<
    HttpServerResponse.HttpServerResponse,
    E,
    HttpServerRequest.HttpServerRequest | R
  >,
) =>
  httpEffect.pipe(
    Effect.timeoutOption(Duration.millis(RELAY_REQUEST_DEADLINE_MS)),
    Effect.flatMap(
      Option.match({
        onNone: () =>
          Effect.gen(function* () {
            const request = yield* HttpServerRequest.HttpServerRequest;
            yield* Effect.logError("relay request exceeded deadline", {
              "http.method": request.method,
              "http.url": request.url,
              "relay.request.deadline_ms": RELAY_REQUEST_DEADLINE_MS,
            });
            yield* Effect.annotateCurrentSpan({
              "relay.request.deadline_exceeded": true,
            });
            return HttpServerResponse.jsonUnsafe(
              { error: "relay_request_deadline_exceeded" },
              { status: 504 },
            );
          }),
        onSome: Effect.succeed,
      }),
    ),
  );

export const traceRelayHttpRequest = <E, R>(
  httpEffect: Effect.Effect<
    HttpServerResponse.HttpServerResponse,
    E,
    HttpServerRequest.HttpServerRequest | R
  >,
) =>
  // HttpMiddleware finalizes its span on the dispatcher; do not close a request-scoped exporter first.
  HttpMiddleware.tracer(
    appendRelayTraceContextResponseHeader.pipe(Effect.andThen(relayRequestDeadline(httpEffect))),
  ).pipe(Effect.ensuring(Effect.yieldNow));

export const traceRelayHttpRequestWith = <E, R, LayerError, LayerRequirements>(
  httpEffect: Effect.Effect<
    HttpServerResponse.HttpServerResponse,
    E,
    HttpServerRequest.HttpServerRequest | R
  >,
  tracerLayer: Layer.Layer<never, LayerError, LayerRequirements>,
) =>
  traceRelayHttpRequest(httpEffect).pipe(
    Effect.provide(Layer.merge(tracerLayer, httpHeaderRedactionLayer)),
  );

export const withoutCapturedParentSpan = <A, E, R>(
  effect: Effect.Effect<A, E, R>,
): Effect.Effect<A, E, R> =>
  Effect.withFiber((fiber) => {
    const context = fiber.context;
    // HttpApiBuilder captures its build context for route handlers; an event parent would outlive export.
    fiber.setContext(Context.omit(Tracer.ParentSpan)(context));
    return effect.pipe(Effect.ensuring(Effect.sync(() => fiber.setContext(context))));
  });

export const relayClientAuthLayer = Layer.effect(
  RelayClientAuth,
  Effect.gen(function* () {
    const config = yield* RelayConfiguration.RelayConfiguration;
    return {
      clientBearer: Effect.fn("relay.auth.client.bearer")(function* (httpEffect, { credential }) {
        const token = readHttpAuthorizationCredential(credential);
        const verified = yield* verifyRelayClientBearerToken(config, token).pipe(
          Effect.tapError((error) =>
            Effect.annotateCurrentSpan(
              "relay.auth.clerk_verification_failure",
              clerkVerificationFailureReason(error.cause),
            ),
          ),
          Effect.catch(() => relayAuthInvalidError("invalid_bearer")),
        );
        if (!verified.sub) {
          yield* Effect.annotateCurrentSpan({
            "relay.auth.clerk_verification_failure": "missing_subject",
          });
          return yield* relayAuthInvalidError("invalid_bearer");
        }
        yield* Effect.annotateCurrentSpan({
          "relay.auth.mode": verified.mode,
          "relay.auth.subject": verified.sub,
        });

        return yield* httpEffect.pipe(
          withSpanAttributes({ "user.id": verified.sub }),
          Effect.provideService(RelayClientPrincipal, {
            userId: verified.sub,
            token,
          }),
        );
      }),
    };
  }),
);

export const relayEnvironmentAuthLayer = Layer.effect(
  RelayEnvironmentAuth,
  Effect.gen(function* () {
    const credentials = yield* EnvironmentCredentials.EnvironmentCredentials;
    return {
      environmentBearer: Effect.fn("relay.auth.environment.bearer")(function* (
        httpEffect,
        { credential },
      ) {
        const token = readHttpAuthorizationCredential(credential);
        const principal = yield* credentials.authenticate(token).pipe(
          Effect.catchTags({
            EnvironmentCredentialAuthenticatePersistenceError: () =>
              relayInternalErrorResponse("persistence_failed"),
          }),
        );
        if (principal._tag === "None") {
          return yield* relayAuthInvalidError("not_authorized");
        }
        yield* Effect.annotateCurrentSpan({
          "relay.auth.mode": "environment_credential",
        });
        return yield* httpEffect.pipe(
          withSpanAttributes({
            "relay.environment_id": principal.value.environmentId,
          }),
          Effect.provideService(RelayEnvironmentPrincipal, principal.value),
        );
      }),
    };
  }),
);

export const relayDpopClientAuthLayer = Layer.effect(
  RelayDpopClientAuth,
  Effect.gen(function* () {
    const relayTokens = yield* RelayTokens.RelayTokens;
    return {
      relayDpop: Effect.fn("relay.auth.dpop_client")(function* (httpEffect, { credential }) {
        yield* appendRelayDpopChallengeHeader;
        const request = yield* HttpServerRequest.HttpServerRequest;
        if (!isDpopAuthorizationHeader(request.headers.authorization)) {
          return yield* relayAuthInvalidError("invalid_bearer");
        }
        const token = readHttpAuthorizationCredential(credential);
        const now = yield* DateTime.now;
        const verified = yield* relayTokens.verifyDpopAccessToken({
          token,
          nowEpochSeconds: Math.floor(now.epochMilliseconds / 1_000),
        });
        if (!verified) {
          return yield* relayAuthInvalidError("invalid_bearer");
        }
        yield* Effect.annotateCurrentSpan({
          "relay.auth.mode": "dpop",
          "relay.auth.subject": verified.sub,
        });
        return yield* httpEffect.pipe(
          withSpanAttributes({ "user.id": verified.sub }),
          Effect.provideService(RelayClientPrincipal, {
            userId: verified.sub,
            token,
            proofKeyThumbprint: verified.cnf.jkt,
            dpopScopes: verified.scope,
          }),
        );
      }),
    };
  }),
);

function isDpopAuthorizationHeader(value: string | undefined): boolean {
  return /^DPoP +/iu.test(value ?? "");
}

function readHttpAuthorizationCredential(credential: Redacted.Redacted<string>): string {
  // Effect beta.73 leaves the scheme separator in decoded HTTP credentials.
  return Redacted.value(credential).trimStart();
}

export const metadataApi = HttpApiBuilder.group(
  RelayApi,
  "metadata",
  Effect.fnUntraced(function* (handlers) {
    const settings = yield* RelayConfiguration.RelayConfiguration;
    const issuer = normalizeRelayIssuer(settings.relayIssuer);
    const scopes = [
      RelayEnvironmentConnectScope,
      RelayEnvironmentStatusScope,
      RelayMobileRegistrationScope,
    ] as const;
    return handlers
      .handle("authorizationServer", () =>
        Effect.succeed({
          issuer,
          token_endpoint: `${issuer}/v1/client/dpop-token`,
          grant_types_supported: ["urn:ietf:params:oauth:grant-type:token-exchange"],
          token_endpoint_auth_methods_supported: ["none"],
          dpop_signing_alg_values_supported: ["ES256"],
          scopes_supported: scopes,
        }),
      )
      .handle("protectedResource", () =>
        Effect.succeed({
          resource: issuer,
          authorization_servers: [issuer],
          scopes_supported: scopes,
          dpop_bound_access_tokens_required: true,
          dpop_signing_alg_values_supported: ["ES256"],
        }),
      );
  }),
);

export const healthApi = HttpApiBuilder.group(
  RelayApi,
  "health",
  Effect.fnUntraced(function* (handlers) {
    const db = yield* RelayDb.RelayDb;
    return handlers.handle(
      "health",
      Effect.fn("relay.api.health")(
        function* () {
          yield* db.execute(drizzleSql`SELECT 1`);
          return { ok: true, service: "relay" as const };
        },
        Effect.catch(() => relayInternalErrorResponse("database_unavailable")),
      ),
    );
  }),
);

export const revokeEnvironmentLinkRecord = Effect.fn(
  "relay.api.client.revokeEnvironmentLinkRecord",
)(function* (input: {
  readonly userId: string;
  readonly environmentId: string;
  readonly environmentPublicKey: string;
}) {
  const transactions = yield* RelayDb.RelayTransactions;
  const links = yield* EnvironmentLinks.EnvironmentLinks;
  const credentials = yield* EnvironmentCredentials.EnvironmentCredentials;
  return yield* transactions.withTransaction(
    Effect.gen(function* () {
      const revoked = yield* links.revokeForUser({
        userId: input.userId,
        environmentId: input.environmentId,
      });
      if (revoked) {
        yield* credentials.revokeForEnvironmentPublicKey({
          environmentId: input.environmentId,
          environmentPublicKey: input.environmentPublicKey,
        });
      }
      return revoked;
    }),
  );
});

export const unlinkEnvironmentRecord = Effect.fn("relay.api.client.unlinkEnvironmentRecord")(
  function* (input: { readonly userId: string; readonly environmentId: string }) {
    const links = yield* EnvironmentLinks.EnvironmentLinks;
    const managedEndpointProvider = yield* ManagedEndpointProvider.ManagedEndpointProvider;
    const deprovisionTarget = yield* managedEndpointProvider.prepareDeprovision({
      userId: input.userId,
      environmentId: input.environmentId,
    });
    const link = yield* links.getForUser({
      userId: input.userId,
      environmentId: input.environmentId,
    });
    const unlinked =
      link === null
        ? false
        : yield* revokeEnvironmentLinkRecord({
            userId: input.userId,
            environmentId: link.environmentId,
            environmentPublicKey: link.environmentPublicKey,
          });

    // External teardown cannot share the SQL transaction. Run it only after
    // revocation commits so a database failure leaves a fully usable active
    // link. Still run teardown when the link is already revoked, allowing a
    // retry to finish cleanup after an earlier Cloudflare failure.
    const deprovisioned = yield* managedEndpointProvider.deprovision({
      userId: input.userId,
      environmentId: input.environmentId,
      target: deprovisionTarget,
    });
    if (!deprovisioned) {
      const retryTarget = yield* managedEndpointProvider.prepareDeprovision(input);
      if (retryTarget !== null && (yield* links.getForUser(input)) === null) {
        yield* managedEndpointProvider.deprovision({ ...input, target: retryTarget });
      }
    }
    return unlinked;
  },
);

type EnvironmentTunnelRecoveryProofInput = {
  readonly proof: string;
  readonly userId: string;
  readonly environmentId: string;
  readonly environmentPublicKey: string;
} & (
  | {
      readonly action: "register";
      readonly tunnelId: string;
      readonly origin: RelayManagedEndpointOrigin;
    }
  | { readonly action: "recover"; readonly origin: RelayManagedEndpointOrigin }
);

export const verifyEnvironmentTunnelRecoveryProof = Effect.fn(
  "relay.api.server.verifyEnvironmentTunnelRecoveryProof",
)(function* (input: EnvironmentTunnelRecoveryProofInput) {
  const config = yield* RelayConfiguration.RelayConfiguration;
  const now = yield* DateTime.now;
  const verified = yield* verifyRelayJwt({
    publicKey: input.environmentPublicKey,
    token: input.proof,
    typ: RELAY_MANAGED_TUNNEL_RECOVERY_TYP,
    issuer: `t3-env:${input.environmentId}`,
    audience: normalizeRelayIssuer(config.relayIssuer),
    nowEpochSeconds: Math.floor(now.epochMilliseconds / 1_000),
  }).pipe(
    Effect.flatMap(decodeManagedTunnelRecoveryProof),
    Effect.mapError(() => new HttpApiError.Unauthorized({})),
  );

  if (
    verified.environmentId !== input.environmentId ||
    verified.sub !== input.environmentId ||
    verified.cloudUserId !== input.userId ||
    verified.action !== input.action
  ) {
    return yield* new HttpApiError.Unauthorized({});
  }
  if (input.action === "register") {
    if (
      verified.action !== "register" ||
      verified.tunnelId !== input.tunnelId ||
      verified.origin.localHttpHost !== input.origin.localHttpHost ||
      verified.origin.localHttpPort !== input.origin.localHttpPort
    ) {
      return yield* new HttpApiError.Unauthorized({});
    }
    return;
  }
  if (
    verified.action !== "recover" ||
    verified.origin.localHttpHost !== input.origin.localHttpHost ||
    verified.origin.localHttpPort !== input.origin.localHttpPort
  ) {
    return yield* new HttpApiError.Unauthorized({});
  }
});

export const registerEnvironmentTunnelRecovery = Effect.fn(
  "relay.api.server.registerEnvironmentTunnelRecovery",
)(function* (input: {
  readonly userId: string;
  readonly environmentId: string;
  readonly environmentPublicKey: string;
  readonly tunnelId: string;
  readonly origin: RelayManagedEndpointOrigin;
}) {
  const links = yield* EnvironmentLinks.EnvironmentLinks;
  const allocations = yield* ManagedEndpointAllocations.ManagedEndpointAllocations;
  const managedEndpointProvider = yield* ManagedEndpointProvider.ManagedEndpointProvider;
  const link = yield* links.getForUser({
    userId: input.userId,
    environmentId: input.environmentId,
  });
  if (
    link === null ||
    link.environmentPublicKey !== input.environmentPublicKey ||
    link.endpoint.providerKind !== "cloudflare_tunnel"
  ) {
    return yield* new HttpApiError.Unauthorized({});
  }
  const status = yield* managedEndpointProvider.reconcileOrigin({
    userId: input.userId,
    environmentId: input.environmentId,
    tunnelId: input.tunnelId,
    origin: input.origin,
    endpoint: link.endpoint,
  });
  if (status === "recovery_required") {
    return { status };
  }
  if (!(yield* allocations.enableRecovery(input))) {
    return yield* new HttpApiError.Unauthorized({});
  }
  return { status };
});

export const recoverEnvironmentTunnelRecord = Effect.fn(
  "relay.api.server.recoverEnvironmentTunnelRecord",
)(function* (input: {
  readonly userId: string;
  readonly environmentId: string;
  readonly environmentPublicKey: string;
  readonly origin: RelayManagedEndpointOrigin;
}) {
  const links = yield* EnvironmentLinks.EnvironmentLinks;
  const allocations = yield* ManagedEndpointAllocations.ManagedEndpointAllocations;
  const managedEndpointProvider = yield* ManagedEndpointProvider.ManagedEndpointProvider;
  const link = yield* links.getForUser({
    userId: input.userId,
    environmentId: input.environmentId,
  });
  if (
    link === null ||
    link.environmentPublicKey !== input.environmentPublicKey ||
    link.endpoint.providerKind !== "cloudflare_tunnel"
  ) {
    return yield* new HttpApiError.Unauthorized({});
  }

  const recovered = yield* managedEndpointProvider.provision({
    userId: input.userId,
    environmentId: input.environmentId,
    origin: input.origin,
  });
  const recoveredTunnelId = recovered.runtime.tunnelId;
  if (
    recoveredTunnelId === undefined ||
    recovered.endpoint.httpBaseUrl !== link.endpoint.httpBaseUrl ||
    recovered.endpoint.wsBaseUrl !== link.endpoint.wsBaseUrl
  ) {
    if (recoveredTunnelId !== undefined) {
      yield* managedEndpointProvider
        .release({
          userId: input.userId,
          environmentId: input.environmentId,
          expectedTunnelId: recoveredTunnelId,
        })
        .pipe(
          Effect.catch((cause) =>
            Effect.logWarning("Failed to clean up a tunnel with a mismatched endpoint", {
              userId: input.userId,
              environmentId: input.environmentId,
              tunnelId: recoveredTunnelId,
              cause,
            }),
          ),
        );
    }
    return yield* new HttpApiError.Unauthorized({});
  }

  const enabled = yield* allocations.enableRecovery({
    userId: input.userId,
    environmentId: input.environmentId,
    tunnelId: recoveredTunnelId,
    environmentPublicKey: input.environmentPublicKey,
    origin: input.origin,
  });
  if (!enabled) {
    const owner = { userId: input.userId, environmentId: input.environmentId };
    const target = yield* managedEndpointProvider.prepareDeprovision(owner);
    const currentLink = target === null ? null : yield* links.getForUser(input);
    if (
      target !== null &&
      (currentLink === null || currentLink.endpoint.providerKind !== "cloudflare_tunnel")
    ) {
      yield* managedEndpointProvider.deprovision({ ...owner, target }).pipe(
        Effect.catch((cause) =>
          Effect.logWarning("Failed to clean up a tunnel after its managed link was removed", {
            userId: input.userId,
            environmentId: input.environmentId,
            cause,
          }),
        ),
      );
    }
    return yield* new HttpApiError.Unauthorized({});
  }
  return {
    endpoint: recovered.endpoint,
    endpointRuntime: recovered.runtime,
  };
});

export const mobileApi = HttpApiBuilder.group(
  RelayApi,
  "mobile",
  Effect.fnUntraced(function* (handlers) {
    const registrations = yield* MobileRegistrations.MobileRegistrations;
    const dpopProofs = yield* DpopProofs.DpopProofReplay;
    return handlers
      .handle(
        "registerDevice",
        Effect.fn("relay.api.mobile.registerDevice")(function* (args) {
          const { payload } = args;
          const { userId, token } = yield* RelayClientPrincipal;
          const proofKeyThumbprint = yield* requireDpopPrincipalScope("mobile:registration");
          yield* requireDpopThumbprint(proofKeyThumbprint, {
            expectedAccessToken: token,
          }).pipe(Effect.provideService(DpopProofs.DpopProofReplay, dpopProofs));
          return yield* registrations.registerDevice({ userId, payload });
        }, mapRelayCommonApiErrors("invalid_dpop")),
      )
      .handle(
        "registerLiveActivity",
        Effect.fn("relay.api.mobile.registerLiveActivity")(function* (args) {
          const { payload } = args;
          const { userId, token } = yield* RelayClientPrincipal;
          const proofKeyThumbprint = yield* requireDpopPrincipalScope("mobile:registration");
          yield* requireDpopThumbprint(proofKeyThumbprint, {
            expectedAccessToken: token,
          }).pipe(Effect.provideService(DpopProofs.DpopProofReplay, dpopProofs));
          return yield* registrations.registerLiveActivity({ userId, payload });
        }, mapRelayCommonApiErrors("invalid_dpop")),
      )
      .handle(
        "getAgentActivitySnapshot",
        Effect.fn("relay.api.mobile.getAgentActivitySnapshot")(function* () {
          const { userId, token } = yield* RelayClientPrincipal;
          const proofKeyThumbprint = yield* requireDpopPrincipalScope("mobile:registration");
          yield* requireDpopThumbprint(proofKeyThumbprint, {
            expectedAccessToken: token,
          }).pipe(Effect.provideService(DpopProofs.DpopProofReplay, dpopProofs));
          return yield* registrations.getAgentActivitySnapshot({ userId });
        }, mapRelayCommonApiErrors("invalid_dpop")),
      )
      .handle(
        "unregisterDevice",
        Effect.fn("relay.api.mobile.unregisterDevice")(function* (args) {
          const { params } = args;
          const { userId, token } = yield* RelayClientPrincipal;
          const proofKeyThumbprint = yield* requireDpopPrincipalScope("mobile:registration");
          yield* requireDpopThumbprint(proofKeyThumbprint, {
            expectedAccessToken: token,
          }).pipe(Effect.provideService(DpopProofs.DpopProofReplay, dpopProofs));
          return yield* registrations.unregisterDevice({ userId, deviceId: params.deviceId });
        }, mapRelayCommonApiErrors("invalid_dpop")),
      );
  }),
);

export const clientApi = HttpApiBuilder.group(
  RelayApi,
  "client",
  Effect.fnUntraced(function* (handlers) {
    const config = yield* RelayConfiguration.RelayConfiguration;
    const crypto = yield* Crypto.Crypto;
    const relayTokens = yield* RelayTokens.RelayTokens;
    const linker = yield* EnvironmentLinker.EnvironmentLinker;
    const links = yield* EnvironmentLinks.EnvironmentLinks;
    const managedEndpointProvider = yield* ManagedEndpointProvider.ManagedEndpointProvider;
    const devices = yield* Devices.Devices;
    return handlers
      .handle(
        "listEnvironments",
        Effect.fn("relay.api.client.listEnvironments")(function* () {
          const { userId } = yield* RelayClientPrincipal;
          const environments = yield* links.listForUser({ userId });
          return { environments };
        }, mapRelayCommonApiErrors("not_authorized")),
      )
      .handle(
        "listDevices",
        Effect.fn("relay.api.client.listDevices")(function* () {
          yield* appendRelayCredentialResponseHeaders;
          const { userId } = yield* RelayClientPrincipal;
          const registered = yield* devices.listForUser({ userId });
          return {
            devices: registered.flatMap((device) =>
              device.platform === "ios" && device.iosMajorVersion !== null
                ? [{ ...device, platform: "ios" as const, iosMajorVersion: device.iosMajorVersion }]
                : [],
            ),
          };
        }, mapRelayCommonApiErrors("not_authorized")),
      )
      .handle(
        "listDevicesV2",
        Effect.fn("relay.api.client.listDevicesV2")(function* () {
          yield* appendRelayCredentialResponseHeaders;
          const { userId } = yield* RelayClientPrincipal;
          return { devices: yield* devices.listForUser({ userId }) };
        }, mapRelayCommonApiErrors("not_authorized")),
      )
      .handle(
        "linkEnvironment",
        Effect.fn("relay.api.client.linkEnvironment")(
          function* (args) {
            const { payload } = args;
            yield* appendRelayCredentialResponseHeaders;
            const { userId } = yield* RelayClientPrincipal;
            const result = yield* linker.link({ userId, request: payload });
            return {
              ok: true,
              cloudUserId: userId,
              environmentId: result.environmentId,
              endpoint: result.endpoint,
              endpointRuntime: result.endpointRuntime,
              relayIssuer: config.relayIssuer,
              environmentCredential: result.environmentCredential,
              cloudMintPublicKey: config.cloudMintPublicKey,
            };
          },
          mapErrorTags({
            EnvironmentLinkProofExpired: (_error, traceId) =>
              new RelayEnvironmentLinkProofExpiredError({
                code: "environment_link_proof_expired",
                traceId,
              }),
            EnvironmentLinkProofInvalid: (linkError, traceId) =>
              new RelayEnvironmentLinkProofInvalidError({
                code: "environment_link_proof_invalid",
                reason: linkError.reason,
                traceId,
              }),
            ManagedEndpointProvisioningNotConfigured: (_error, traceId) =>
              new RelayEnvironmentLinkUnavailableError({
                code: "environment_link_unavailable",
                reason: "managed_endpoint_not_configured",
                traceId,
              }),
            ManagedEndpointProvisioningFailed: (_error, traceId) =>
              new RelayEnvironmentLinkUnavailableError({
                code: "environment_link_unavailable",
                reason: "managed_endpoint_provisioning_failed",
                traceId,
              }),
            ManagedEndpointOriginNotAllowed: (_error, traceId) =>
              new RelayEnvironmentLinkProofInvalidError({
                code: "environment_link_proof_invalid",
                reason: "origin_not_allowed",
                traceId,
              }),
            ManagedTunnelLimitExceeded: (limitError, traceId) =>
              new RelayEnvironmentLinkLimitExceededError({
                code: "environment_link_limit_exceeded",
                maxTunnels: limitError.maxTunnels,
                traceId,
              }),
            EnvironmentLinkUpsertPersistenceError: (_error, traceId) =>
              new RelayEnvironmentLinkFailedError({
                code: "environment_link_failed",
                reason: "link_persistence_failed",
                traceId,
              }),
            EnvironmentCredentialCreatePersistenceError: (_error, traceId) =>
              new RelayEnvironmentLinkFailedError({
                code: "environment_link_failed",
                reason: "credential_persistence_failed",
                traceId,
              }),
            DpopProofReplayPersistenceError: (_error, traceId) =>
              new RelayEnvironmentLinkFailedError({
                code: "environment_link_failed",
                reason: "replay_persistence_failed",
                traceId,
              }),
          }),
          mapRelayCommonApiErrors("not_authorized"),
        ),
      )
      .handle(
        "createEnvironmentLinkChallenge",
        Effect.fn("relay.api.client.createEnvironmentLinkChallenge")(function* (args) {
          yield* appendRelayCredentialResponseHeaders;
          const { userId } = yield* RelayClientPrincipal;
          const now = yield* DateTime.now;
          const expiresAt = DateTime.add(now, { minutes: 5 });
          const jti = yield* crypto.randomUUIDv4.pipe(
            Effect.catch(() => relayInternalErrorResponse("internal_error")),
          );
          const challenge = yield* relayTokens
            .issueLinkChallenge({
              userId,
              request: args.payload,
              jti,
              issuedAtEpochSeconds: Math.floor(now.epochMilliseconds / 1_000),
              expiresAtEpochSeconds: Math.floor(expiresAt.epochMilliseconds / 1_000),
            })
            .pipe(Effect.catch(() => relayInternalErrorResponse("internal_error")));
          return { challenge, expiresAt: DateTime.formatIso(expiresAt) };
        }, mapRelayCommonApiErrors("not_authorized")),
      )
      .handle(
        "unlinkEnvironment",
        Effect.fn("relay.api.client.unlinkEnvironment")(function* (args) {
          const { params } = args;
          const { userId } = yield* RelayClientPrincipal;
          const unlinked = yield* unlinkEnvironmentRecord({
            userId,
            environmentId: params.environmentId,
          }).pipe(
            Effect.catchTags({
              SqlError: () => relayInternalErrorResponse("internal_error"),
              ManagedEndpointDeprovisioningFailed: () =>
                relayInternalErrorResponse("upstream_unavailable"),
            }),
          );
          return { ok: unlinked };
        }, mapRelayCommonApiErrors("not_authorized")),
      )
      .handle(
        "releaseEnvironmentTunnel",
        Effect.fn("relay.api.client.releaseEnvironmentTunnel")(function* (args) {
          const { params } = args;
          const { userId } = yield* RelayClientPrincipal;
          // ok mirrors whether the connector token is now dead: false means a
          // concurrent provision kept the recorded tunnel alive, so the caller
          // must not discard its runtime config.
          const released = yield* managedEndpointProvider
            .release({
              userId,
              environmentId: params.environmentId,
            })
            .pipe(Effect.catch(() => relayInternalErrorResponse("upstream_unavailable")));
          return { ok: released };
        }, mapRelayCommonApiErrors("not_authorized")),
      );
  }),
);

export const tokenApi = HttpApiBuilder.group(
  RelayApi,
  "token",
  Effect.fnUntraced(function* (handlers) {
    const config = yield* RelayConfiguration.RelayConfiguration;
    const crypto = yield* Crypto.Crypto;
    const dpopProofs = yield* DpopProofs.DpopProofReplay;
    const relayTokens = yield* RelayTokens.RelayTokens;
    return handlers.handle(
      "exchangeDpopAccessToken",
      Effect.fn("relay.api.token.exchangeDpopAccessToken")(function* (args) {
        yield* appendRelayCredentialResponseHeaders;
        const issuer = normalizeRelayIssuer(config.relayIssuer);
        const requestedScopes = relayTokens.resolveDpopAccessTokenScopes({
          clientId: args.payload.client_id,
          scope: args.payload.scope,
        });
        yield* Effect.annotateCurrentSpan({
          "relay.auth.mode": "clerk_bearer_token_exchange",
          "relay.oauth.client_id": args.payload.client_id,
          "relay.oauth.scopes": args.payload.scope,
        });
        if (args.payload.resource !== issuer || requestedScopes === null) {
          return yield* new HttpApiError.Unauthorized({});
        }

        const verified = yield* verifyClerkBearerToken(config, args.payload.subject_token).pipe(
          Effect.catch(() => relayAuthInvalidError("invalid_bearer")),
        );
        if (!verified.sub || !hasExpectedClerkAudience(verified.aud, config.clerkJwtAudience)) {
          return yield* relayAuthInvalidError("invalid_bearer");
        }
        const proofKeyThumbprint = yield* requireDpopProof().pipe(
          Effect.provideService(DpopProofs.DpopProofReplay, dpopProofs),
        );
        const now = yield* DateTime.now;
        const expiresAt = DateTime.addDuration(now, RelayTokens.RELAY_DPOP_ACCESS_TOKEN_TTL);
        const jti = yield* crypto.randomUUIDv4.pipe(
          Effect.catch(() => relayInternalErrorResponse("internal_error")),
        );
        return {
          access_token: yield* relayTokens
            .issueDpopAccessToken({
              userId: verified.sub,
              proofKeyThumbprint,
              jti,
              issuedAtEpochSeconds: Math.floor(now.epochMilliseconds / 1_000),
              expiresAtEpochSeconds: Math.floor(expiresAt.epochMilliseconds / 1_000),
              clientId: args.payload.client_id,
              scopes: requestedScopes,
            })
            .pipe(Effect.catch(() => relayInternalErrorResponse("internal_error"))),
          issued_token_type: RelayAccessTokenType,
          token_type: "DPoP" as const,
          expires_in: Duration.toSeconds(RelayTokens.RELAY_DPOP_ACCESS_TOKEN_TTL),
          scope: encodeOAuthScope(requestedScopes),
        };
      }, mapRelayCommonApiErrors("invalid_dpop")),
    );
  }),
);

export const dpopClientApi = HttpApiBuilder.group(
  RelayApi,
  "dpopClient",
  Effect.fnUntraced(function* (handlers) {
    const connector = yield* EnvironmentConnector.EnvironmentConnector;
    const dpopProofs = yield* DpopProofs.DpopProofReplay;
    return handlers
      .handle(
        "connectEnvironment",
        Effect.fn("relay.api.dpopClient.connectEnvironment")(
          function* (args) {
            const { params, payload } = args;
            yield* appendRelayCredentialResponseHeaders;
            const { userId, token } = yield* RelayClientPrincipal;
            const proofKeyThumbprint = yield* requireDpopPrincipalScope("environment:connect");
            const requestedThumbprint = resolveConnectClientKeyThumbprint(payload);
            if (!requestedThumbprint || requestedThumbprint !== proofKeyThumbprint) {
              return yield* new HttpApiError.Unauthorized({});
            }
            const clientProofKeyThumbprint = yield* requireDpopThumbprint(proofKeyThumbprint, {
              expectedAccessToken: token,
            }).pipe(Effect.provideService(DpopProofs.DpopProofReplay, dpopProofs));
            return yield* connector.connect({
              userId,
              environmentId: params.environmentId,
              clientProofKeyThumbprint,
              ...(payload.deviceId ? { deviceId: payload.deviceId } : {}),
            });
          },
          mapRelayCommonApiErrors("invalid_dpop"),
          mapErrorTags({
            EnvironmentConnectNotAuthorized: (error, traceId) =>
              new RelayEnvironmentConnectNotAuthorizedError({
                code: "environment_connect_not_authorized",
                reason: error.reason,
                traceId,
              }),
            EnvironmentMintRequestFailed: (_error, traceId) =>
              new RelayEnvironmentEndpointUnavailableError({
                code: "environment_endpoint_unavailable",
                reason: "endpoint_request_failed",
                traceId,
              }),
            EnvironmentMintRequestTimedOut: (_error, traceId) =>
              new RelayEnvironmentEndpointTimedOutError({
                code: "environment_endpoint_timed_out",
                traceId,
              }),
            EnvironmentMintResponseInvalid: (_error, traceId) =>
              new RelayEnvironmentEndpointUnavailableError({
                code: "environment_endpoint_unavailable",
                reason: "endpoint_response_invalid",
                traceId,
              }),
          }),
        ),
      )
      .handle(
        "getEnvironmentStatus",
        Effect.fn("relay.api.dpopClient.getEnvironmentStatus")(
          function* (args) {
            const { params } = args;
            const { userId, token } = yield* RelayClientPrincipal;
            const proofKeyThumbprint = yield* requireDpopPrincipalScope("environment:status");
            yield* requireDpopThumbprint(proofKeyThumbprint, {
              expectedAccessToken: token,
            }).pipe(Effect.provideService(DpopProofs.DpopProofReplay, dpopProofs));
            return yield* connector.status({
              userId,
              environmentId: params.environmentId,
            });
          },
          mapRelayCommonApiErrors("invalid_dpop"),
          mapErrorTags({
            EnvironmentConnectNotAuthorized: (error, traceId) =>
              new RelayEnvironmentConnectNotAuthorizedError({
                code: "environment_connect_not_authorized",
                reason: error.reason,
                traceId,
              }),
            EnvironmentMintRequestFailed: (_error, traceId) =>
              new RelayEnvironmentEndpointUnavailableError({
                code: "environment_endpoint_unavailable",
                reason: "endpoint_request_failed",
                traceId,
              }),
            EnvironmentMintRequestTimedOut: (_error, traceId) =>
              new RelayEnvironmentEndpointTimedOutError({
                code: "environment_endpoint_timed_out",
                traceId,
              }),
            EnvironmentMintResponseInvalid: (_error, traceId) =>
              new RelayEnvironmentEndpointUnavailableError({
                code: "environment_endpoint_unavailable",
                reason: "endpoint_response_invalid",
                traceId,
              }),
          }),
        ),
      );
  }),
);

export const serverApi = HttpApiBuilder.group(
  RelayApi,
  "server",
  Effect.fnUntraced(function* (handlers) {
    const publisher = yield* AgentActivityPublisher.AgentActivityPublisher;
    const publishSignatures = yield* EnvironmentPublishSignatures.EnvironmentPublishSignatures;
    const activityHandlers = handlers.handle(
      "publishAgentActivity",
      Effect.fn("relay.api.server.publishAgentActivity")(
        function* (args) {
          const { params, payload } = args;
          const principal = yield* RelayEnvironmentPrincipal;
          if (principal.environmentId !== params.environmentId) {
            return yield* new HttpApiError.Unauthorized({});
          }
          yield* publishSignatures.verify({
            environmentId: params.environmentId,
            environmentPublicKey: principal.environmentPublicKey,
            threadId: params.threadId,
            request: payload,
          });
          return yield* publisher.publish({
            environmentId: params.environmentId,
            environmentPublicKey: principal.environmentPublicKey,
            threadId: params.threadId,
            state: payload.state,
          });
        },
        mapErrorTags({
          EnvironmentPublishPublicKeyMissing: (_error, traceId) =>
            new RelayAuthInvalidError({
              code: "auth_invalid",
              reason: "not_authorized",
              traceId,
            }),
          EnvironmentPublishSignatureExpired: (_error, traceId) =>
            new RelayAgentActivityPublishProofExpiredError({
              code: "agent_activity_publish_proof_expired",
              traceId,
            }),
          EnvironmentPublishSignatureInvalid: (_error, traceId) =>
            new RelayAgentActivityPublishProofInvalidError({
              code: "agent_activity_publish_proof_invalid",
              reason: "invalid_signature_or_payload",
              traceId,
            }),
          DpopProofReplayPersistenceError: (_error, traceId) =>
            new RelayInternalError({
              code: "internal_error",
              reason: "persistence_failed",
              traceId,
            }),
          ApnsDeliveryJobQueuePayloadInvalid: (_error, traceId) =>
            new RelayInternalError({
              code: "internal_error",
              reason: "internal_error",
              traceId,
            }),
          ApnsDeliveryJobLiveActivityAggregateMissing: (_error, traceId) =>
            new RelayInternalError({
              code: "internal_error",
              reason: "internal_error",
              traceId,
            }),
          ApnsDeliveryJobLiveActivityNotificationUnexpected: (_error, traceId) =>
            new RelayInternalError({
              code: "internal_error",
              reason: "internal_error",
              traceId,
            }),
          ApnsDeliveryJobPushNotificationMissing: (_error, traceId) =>
            new RelayInternalError({
              code: "internal_error",
              reason: "internal_error",
              traceId,
            }),
          ApnsDeliveryJobPushNotificationAggregateUnexpected: (_error, traceId) =>
            new RelayInternalError({
              code: "internal_error",
              reason: "internal_error",
              traceId,
            }),
          ApnsDeliveryJobCreatedAtInvalid: (_error, traceId) =>
            new RelayInternalError({
              code: "internal_error",
              reason: "internal_error",
              traceId,
            }),
          ApnsDeliveryJobExpiresAtInvalid: (_error, traceId) =>
            new RelayInternalError({
              code: "internal_error",
              reason: "internal_error",
              traceId,
            }),
          ApnsDeliveryJobTimeWindowInvalid: (_error, traceId) =>
            new RelayInternalError({
              code: "internal_error",
              reason: "internal_error",
              traceId,
            }),
          ApnsDeliveryJobTimeWindowTooLong: (_error, traceId) =>
            new RelayInternalError({
              code: "internal_error",
              reason: "internal_error",
              traceId,
            }),
          ApnsDeliveryJobSignatureInvalid: (_error, traceId) =>
            new RelayInternalError({
              code: "internal_error",
              reason: "internal_error",
              traceId,
            }),
          ApnsDeliveryJobExpired: (_error, traceId) =>
            new RelayInternalError({
              code: "internal_error",
              reason: "internal_error",
              traceId,
            }),
          ApnsDeliveryJobClaimInFlight: (_error, traceId) =>
            new RelayInternalError({
              code: "internal_error",
              reason: "internal_error",
              traceId,
            }),
          ApnsDeliveryQueueSendError: (_error, traceId) =>
            new RelayInternalError({
              code: "internal_error",
              reason: "upstream_unavailable",
              traceId,
            }),
          FcmDeliveryError: (_error, traceId) =>
            new RelayInternalError({
              code: "internal_error",
              reason: "upstream_unavailable",
              traceId,
            }),
        }),
        mapRelayCommonApiErrors("not_authorized"),
      ),
    );

    return activityHandlers
      .handle(
        "registerManagedEndpointRecovery",
        Effect.fn("relay.api.server.registerManagedEndpointRecovery")(
          function* ({ params, payload }) {
            const principal = yield* RelayEnvironmentPrincipal;
            if (principal.environmentId !== params.environmentId) {
              return yield* new HttpApiError.Unauthorized({});
            }
            yield* verifyEnvironmentTunnelRecoveryProof({
              action: "register",
              proof: payload.proof,
              userId: payload.cloudUserId,
              environmentId: params.environmentId,
              environmentPublicKey: principal.environmentPublicKey,
              tunnelId: payload.tunnelId,
              origin: payload.origin,
            });
            yield* appendRelayCredentialResponseHeaders;
            return yield* registerEnvironmentTunnelRecovery({
              userId: payload.cloudUserId,
              environmentId: params.environmentId,
              environmentPublicKey: principal.environmentPublicKey,
              tunnelId: payload.tunnelId,
              origin: payload.origin,
            });
          },
          Effect.catchTags({
            ManagedEndpointOriginNotAllowed: () => Effect.fail(new HttpApiError.Unauthorized({})),
            ManagedEndpointProvisioningNotConfigured: () =>
              relayInternalErrorResponse("upstream_unavailable"),
            ManagedEndpointProvisioningFailed: () =>
              relayInternalErrorResponse("upstream_unavailable"),
            ManagedTunnelLimitExceeded: () => relayInternalErrorResponse("upstream_unavailable"),
          }),
          mapRelayCommonApiErrors("not_authorized"),
        ),
      )
      .handle(
        "recoverManagedEndpoint",
        Effect.fn("relay.api.server.recoverManagedEndpoint")(
          function* ({ params, payload }) {
            const principal = yield* RelayEnvironmentPrincipal;
            if (principal.environmentId !== params.environmentId) {
              return yield* new HttpApiError.Unauthorized({});
            }
            yield* verifyEnvironmentTunnelRecoveryProof({
              action: "recover",
              proof: payload.proof,
              userId: payload.cloudUserId,
              environmentId: params.environmentId,
              environmentPublicKey: principal.environmentPublicKey,
              origin: payload.origin,
            });
            yield* appendRelayCredentialResponseHeaders;
            return yield* recoverEnvironmentTunnelRecord({
              userId: payload.cloudUserId,
              environmentId: params.environmentId,
              environmentPublicKey: principal.environmentPublicKey,
              origin: payload.origin,
            });
          },
          Effect.catchTags({
            ManagedEndpointOriginNotAllowed: () => Effect.fail(new HttpApiError.Unauthorized({})),
            ManagedEndpointProvisioningNotConfigured: () =>
              relayInternalErrorResponse("upstream_unavailable"),
            ManagedEndpointProvisioningFailed: () =>
              relayInternalErrorResponse("upstream_unavailable"),
            ManagedEndpointDeprovisioningFailed: () =>
              relayInternalErrorResponse("upstream_unavailable"),
            ManagedTunnelLimitExceeded: () => relayInternalErrorResponse("upstream_unavailable"),
          }),
          mapRelayCommonApiErrors("not_authorized"),
        ),
      );
  }),
);

class ClerkTokenVerificationFailed extends Schema.TaggedError<ClerkTokenVerificationFailed>()(
  "ClerkTokenVerificationFailed",
  {
    cause: Schema.Defect(),
  },
) {
  override get message(): string {
    return "Clerk token verification failed";
  }
}

const isHttpUnauthorized = Schema.is(HttpApiError.Unauthorized);
const isDpopProofRejected = Schema.is(DpopProofs.DpopProofRejected);

const currentTraceId = Effect.currentParentSpan.pipe(
  Effect.map((span) => span.traceId),
  Effect.orElseSucceed(() => "unavailable"),
);

const RelayCommonPersistenceError = Schema.Union([
  Devices.DeviceRegistrationPersistenceError,
  Devices.DeviceUnregistrationPersistenceError,
  Devices.DeviceListPersistenceError,
  LiveActivities.LiveActivityRegistrationPersistenceError,
  EnvironmentLinks.EnvironmentLinkUserListPersistenceError,
  EnvironmentLinks.EnvironmentLinkListPersistenceError,
  EnvironmentLinks.EnvironmentLinkLookupPersistenceError,
  EnvironmentLinks.EnvironmentLinkRevokePersistenceError,
  ManagedEndpointAllocations.ManagedEndpointAllocationPersistenceError,
  EnvironmentCredentials.EnvironmentCredentialAuthenticatePersistenceError,
  EnvironmentCredentials.EnvironmentCredentialRevokePersistenceError,
  DpopProofs.DpopProofReplayPersistenceError,
  LiveActivities.LiveActivityTargetListPersistenceError,
  AgentActivityRows.AgentActivityRowUpsertPersistenceError,
  AgentActivityRows.AgentActivityRowDeletePersistenceError,
  AgentActivityRows.AgentActivityRowListPersistenceError,
  LiveActivities.LiveActivityDeliveryMarkPersistenceError,
  DeliveryAttempts.DeliveryAttemptRecordPersistenceError,
]);
type RelayCommonPersistenceError = typeof RelayCommonPersistenceError.Type;
const isRelayCommonPersistenceError = Schema.is(RelayCommonPersistenceError);

type MapRelayCommonApiError<E> =
  | Exclude<
      E,
      HttpApiError.Unauthorized | DpopProofs.DpopProofRejected | RelayCommonPersistenceError
    >
  | (Extract<E, HttpApiError.Unauthorized> extends never ? never : RelayAuthInvalidError)
  | (Extract<E, DpopProofs.DpopProofRejected> extends never ? never : RelayAuthInvalidError)
  | (Extract<E, RelayCommonPersistenceError> extends never ? never : RelayInternalError);

export function relayDpopFailureReason(
  code: DpopProofs.DpopProofFailureCode,
): RelayDpopFailureReason {
  switch (code) {
    case "time_window":
      return "time_window";
    case "key_mismatch":
      return "key_mismatch";
    case "method_mismatch":
    case "url_mismatch":
      return "request_mismatch";
    case "access_token_hash_mismatch":
      return "token_mismatch";
    case "replayed":
      return "replay";
    case "missing_proof":
    case "malformed_proof":
    case "invalid_signature":
    case "invalid_proof":
      return "invalid_proof";
  }
}

function relayInternalErrorResponse(reason: RelayInternalError["reason"]) {
  return currentTraceId.pipe(
    Effect.flatMap((traceId) =>
      Effect.fail(new RelayInternalError({ code: "internal_error", reason, traceId })),
    ),
  );
}

function mapRelayCommonApiErrors(authReason: RelayAuthInvalidReason) {
  const mapError = Effect.fnUntraced(function* <E>(error: E) {
    const traceId = yield* currentTraceId;
    if (isDpopProofRejected(error)) {
      yield* Effect.annotateCurrentSpan({
        "relay.dpop.failure_code": error.code,
      });
      return yield* Effect.fail(
        new RelayAuthInvalidError({
          code: "auth_invalid",
          reason: authReason,
          ...(authReason === "invalid_dpop"
            ? { dpopFailureReason: relayDpopFailureReason(error.code) }
            : {}),
          traceId,
        }) as MapRelayCommonApiError<E>,
      );
    }
    if (isHttpUnauthorized(error)) {
      if (authReason === "invalid_dpop") {
        yield* Effect.annotateCurrentSpan({
          "relay.dpop.failure_code": "invalid_proof",
        });
      }
      return yield* Effect.fail(
        new RelayAuthInvalidError({
          code: "auth_invalid",
          reason: authReason,
          ...(authReason === "invalid_dpop" ? { dpopFailureReason: "invalid_proof" } : {}),
          traceId,
        }) as MapRelayCommonApiError<E>,
      );
    }
    if (isRelayCommonPersistenceError(error)) {
      return yield* Effect.fail(
        new RelayInternalError({
          code: "internal_error",
          reason: "persistence_failed",
          traceId,
        }) as MapRelayCommonApiError<E>,
      );
    }

    return yield* Effect.fail(error as MapRelayCommonApiError<E>);
  });

  return <A, E, R>(
    effect: Effect.Effect<A, E, R>,
  ): Effect.Effect<A, MapRelayCommonApiError<E>, R> => effect.pipe(Effect.catch(mapError));
}

type TaggedErrorTag<E> = Extract<E, { readonly _tag: string }>["_tag"];

type MapErrorTagCases<E> = {
  readonly [K in TaggedErrorTag<E>]+?: (
    error: Extract<E, { readonly _tag: K }>,
    traceId: string,
  ) => unknown;
};

type MappedTagError<Cases> = Cases[keyof Cases] extends (
  ...args: ReadonlyArray<never>
) => infer Error
  ? Error
  : never;

type CatchTagCases<E, Cases> = {
  readonly [K in TaggedErrorTag<E>]+?: (
    error: Extract<E, { readonly _tag: K }>,
  ) => Effect.Effect<never, MappedTagError<Cases>>;
} & (unknown extends E ? {} : { readonly [K in Exclude<keyof Cases, TaggedErrorTag<E>>]: never });

function mapErrorTags<
  E,
  Cases extends MapErrorTagCases<E> &
    (unknown extends E ? {} : { readonly [K in Exclude<keyof Cases, TaggedErrorTag<E>>]: never }),
>(cases: Cases) {
  const catchCases = Record.map(
    cases as Record.ReadonlyRecord<
      string,
      (error: never, traceId: string) => MappedTagError<Cases>
    >,
    (makeError) => (error: never) =>
      currentTraceId.pipe(Effect.flatMap((traceId) => Effect.fail(makeError(error, traceId)))),
  ) as CatchTagCases<E, Cases>;

  return <A, R>(
    self: Effect.Effect<A, E, R>,
  ): Effect.Effect<A, Exclude<E, { readonly _tag: keyof Cases }> | MappedTagError<Cases>, R> =>
    // @effect-diagnostics-next-line unsafeEffectTypeAssertion:off
    Effect.catchTags(self, catchCases) as Effect.Effect<
      A,
      Exclude<E, { readonly _tag: keyof Cases }> | MappedTagError<Cases>,
      R
    >;
}

function resolveConnectClientKeyThumbprint(payload: RelayEnvironmentConnectRequest): string | null {
  const requestedThumbprint = payload.clientKeyThumbprint ?? payload.clientProofKeyThumbprint;
  if (!requestedThumbprint) {
    return null;
  }
  if (
    payload.clientKeyThumbprint &&
    payload.clientProofKeyThumbprint &&
    payload.clientKeyThumbprint !== payload.clientProofKeyThumbprint
  ) {
    return null;
  }
  return requestedThumbprint;
}

function safeAuthFailureReason(value: string): string {
  return /^[a-z0-9._-]+$/i.test(value) ? value : "unknown";
}

function clerkVerificationFailureReason(cause: unknown): string {
  if (
    cause instanceof Error &&
    (cause.message.startsWith("Invalid JWT audience claim ") ||
      cause.message.startsWith("Invalid JWT audience claim array "))
  ) {
    return "audience_mismatch";
  }
  if (typeof cause === "object" && cause !== null && "reason" in cause) {
    const reason = (cause as { readonly reason?: unknown }).reason;
    if (typeof reason === "string" && reason.length > 0) {
      return safeAuthFailureReason(reason);
    }
  }
  if (cause instanceof Error && cause.name) {
    return safeAuthFailureReason(cause.name);
  }
  return "unknown";
}

function hasExpectedClerkAudience(audience: unknown, expectedAudience: string): boolean {
  return typeof audience === "string"
    ? audience === expectedAudience
    : Array.isArray(audience) &&
        audience.some((entry) => typeof entry === "string" && entry === expectedAudience);
}

function verifyClerkBearerToken(
  config: RelayConfiguration.RelayConfiguration["Service"],
  token: string,
) {
  return Effect.tryPromise({
    try: () =>
      verifyToken(token, {
        secretKey: Redacted.value(config.clerkSecretKey),
        audience: config.clerkJwtAudience,
      }),
    catch: (cause) => new ClerkTokenVerificationFailed({ cause }),
  }).pipe(
    Effect.withSpan("verify_clerk_bearer_token", {
      attributes: { "relay.auth.token_length": token.length },
    }),
  );
}

function verifyClerkOAuthBearerToken(
  config: RelayConfiguration.RelayConfiguration["Service"],
  token: string,
) {
  return Effect.tryPromise({
    try: async () => {
      const client = createClerkClient({
        secretKey: Redacted.value(config.clerkSecretKey),
        publishableKey: config.clerkPublishableKey,
      });
      const state = await client.authenticateRequest(
        new Request(config.relayIssuer, {
          headers: { authorization: `Bearer ${token}` },
        }),
        { acceptsToken: "oauth_token" },
      );
      const auth = state.toAuth();
      if (!state.isAuthenticated || !auth.userId) {
        throw new Error("Clerk OAuth token is not authenticated.");
      }
      return { sub: auth.userId };
    },
    catch: (cause) => new ClerkTokenVerificationFailed({ cause }),
  });
}

export function verifyRelayClientBearerToken(
  config: RelayConfiguration.RelayConfiguration["Service"],
  token: string,
) {
  return verifyClerkBearerToken(config, token).pipe(
    Effect.flatMap((verified) =>
      verified.sub && hasExpectedClerkAudience(verified.aud, config.clerkJwtAudience)
        ? Effect.succeed({ sub: verified.sub, mode: "clerk_session_bearer" as const })
        : Effect.fail(new ClerkTokenVerificationFailed({ cause: "missing_relay_audience" })),
    ),
    Effect.catch(() =>
      verifyClerkOAuthBearerToken(config, token).pipe(
        Effect.map((verified) => ({ ...verified, mode: "clerk_oauth_bearer" as const })),
      ),
    ),
  );
}

const requireDpopPrincipalScope = Effect.fn("relay.api.require_dpop_principal_scope")(function* (
  scope: RelayDpopAccessTokenScope,
) {
  yield* Effect.annotateCurrentSpan({ "relay.dpop.required_scope": scope });
  const principal = yield* RelayClientPrincipal;
  if (!principal.proofKeyThumbprint || !principal.dpopScopes?.includes(scope)) {
    return yield* new HttpApiError.Unauthorized({});
  }
  return principal.proofKeyThumbprint;
});

const requireDpopThumbprint = Effect.fn("relay.api.require_dpop_thumbprint")(function* (
  expectedThumbprint: string,
  options?: {
    readonly expectedAccessToken?: string;
  },
) {
  const request = yield* HttpServerRequest.HttpServerRequest;
  const now = yield* DateTime.now;
  const url = HttpServerRequest.toURL(request);
  if (url._tag === "None") {
    return yield* new HttpApiError.Unauthorized({});
  }
  const dpopProofs = yield* DpopProofs.DpopProofReplay;
  return yield* dpopProofs.verifyAndConsume({
    proof: request.headers.dpop,
    method: request.method,
    url: url.value.href,
    now,
    expectedThumbprint,
    ...(options?.expectedAccessToken ? { expectedAccessToken: options.expectedAccessToken } : {}),
  });
});

const requireDpopProof = Effect.fn("relay.api.require_dpop_proof")(function* (options?: {
  readonly expectedAccessToken?: string;
}) {
  const request = yield* HttpServerRequest.HttpServerRequest;
  const now = yield* DateTime.now;
  const url = HttpServerRequest.toURL(request);
  if (url._tag === "None") {
    return yield* new HttpApiError.Unauthorized({});
  }
  const dpopProofs = yield* DpopProofs.DpopProofReplay;
  return yield* dpopProofs.verifyAndConsume({
    proof: request.headers.dpop,
    method: request.method,
    url: url.value.href,
    now,
    ...(options?.expectedAccessToken ? { expectedAccessToken: options.expectedAccessToken } : {}),
  });
});

const relayAuthInvalidError = Effect.fnUntraced(function* (reason: RelayAuthInvalidReason) {
  const traceId = yield* currentTraceId;
  yield* Effect.annotateCurrentSpan({
    "relay.trace_id": traceId,
    "relay.error.outbound_tag": "RelayAuthInvalidError",
    "relay.error.outbound_reason": reason,
  });
  return yield* new RelayAuthInvalidError({ code: "auth_invalid", reason, traceId });
});