Files
nekonest-cloud/db/relay-control-plane.ts

1058 lines
40 KiB
TypeScript

import { env } from "cloudflare:workers";
import { ensureDatabase, getD1 } from "./bootstrap";
import { sha256Hex, randomDeviceToken, validateDeviceIdentity } from "./device-claim";
import { DomainError } from "./repository";
import {
type AuthorizedPhone,
RELAY_AUTHORIZATION_MAX_TTL_SECONDS,
RELAY_AUTHORIZATION_SNAPSHOT_VERSION,
RELAY_SIGNING_KEY_FOR_SNAPSHOT_SQL,
classifySnapshotPlacement,
signRelayAuthorizationSnapshot,
type AuthorizedDevice,
type RelayAuthorizationSnapshotPayload,
type SignedRelayAuthorizationSnapshot,
} from "./relay-authorization";
import {
ACTIVATE_PHONE_PRINCIPAL_SQL,
ACTIVATE_PHONE_ROUTE_SQL,
ADVANCE_AUTHORIZATION_AFTER_PHONE_ACTIVATION_SQL,
AUTHORIZE_PHONE_ROUTE_SQL,
ADVANCE_AUTHORIZATION_AFTER_PHONE_REVOKE_SQL,
CLAIM_PHONE_HANDOFF_ACTIVATION_SQL,
CONSUME_PHONE_HANDOFF_SQL,
DELETE_SUPERSEDED_PENDING_PHONE_PRINCIPALS_SQL,
DELETE_SUPERSEDED_PENDING_PHONE_ROUTES_SQL,
FINALIZE_PHONE_HANDOFF_ACTIVATION_SQL,
PHONE_FOR_NODE_REVOCATION_SQL,
REVOKE_PHONE_PRINCIPAL_SQL,
REVOKE_PHONE_ROUTES_SQL,
} from "./relay-control-sql.ts";
import {
AUTHENTICATE_RELAY_NODE_IDENTITY_SQL,
RelayMtlsIdentityError,
verifyTrustedRelayMtlsIdentity,
} from "./relay-node-identity.ts";
import {
resolveRelayPlacementRoute,
type RelayPlacementRouteRecord,
type RelayRouteResolution,
} from "./relay-routing";
export type { RelayRouteResolution } from "./relay-routing";
export type RelayNodePrincipal = {
nodeId: string;
regionId: string;
};
export type TenantConnectionState = {
connectionState: "ready" | "provisioning";
retryAfterSeconds?: number;
};
type RelaySigningKeyRecord = {
kid: string;
public_key_jwk: string;
private_key_ref: string;
};
type PlacementRecord = Omit<RelayPlacementRouteRecord, "tenant_status"> & {
tenant_status: "active" | "suspended";
};
const PHONE_HANDOFF_ACTIVATION_TTL_MS = 5 * 60_000;
function newId(prefix: string): string {
return `${prefix}_${crypto.randomUUID().replaceAll("-", "")}`;
}
function assertOrigin(value: string): string {
let url: URL;
try {
url = new URL(value);
} catch {
throw new DomainError("invalid_pwa_origin", "PWA 地址无效");
}
if (
url.origin !== value ||
!["https:", "http:"].includes(url.protocol) ||
(url.protocol === "http:" && !["localhost", "127.0.0.1", "[::1]"].includes(url.hostname))
) {
throw new DomainError("invalid_pwa_origin", "PWA 必须使用精确 HTTPS origin");
}
return url.origin;
}
function privateJwkForKid(kid: string): JsonWebKey {
const encoded = env.NEKONEST_CLOUD_RELAY_SIGNING_PRIVATE_JWKS?.trim() ?? "";
if (!encoded) {
throw new DomainError(
"relay_signing_key_unavailable",
"Relay 授权签名密钥不可用",
503,
true,
5,
);
}
let parsed: unknown;
try {
parsed = JSON.parse(encoded);
} catch {
throw new DomainError("relay_signing_key_unavailable", "Relay 授权签名密钥配置无效", 503, true, 5);
}
const candidates = Array.isArray(parsed) ? parsed : [parsed];
const match = candidates.find(
(candidate): candidate is JsonWebKey & { kid: string } =>
typeof candidate === "object" && candidate !== null &&
(candidate as { kid?: unknown }).kid === kid,
);
if (!match || match.kty !== "OKP" || match.crv !== "Ed25519" || !match.d) {
throw new DomainError("relay_signing_key_unavailable", "Relay 授权签名密钥不匹配", 503, true, 5);
}
return match;
}
async function loadSigningKey(issuedAt: string, snapshotExpiresAt: string): Promise<{
record: RelaySigningKeyRecord;
privateKey: CryptoKey;
}> {
const record = await getD1()
.prepare(RELAY_SIGNING_KEY_FOR_SNAPSHOT_SQL)
.bind(issuedAt, snapshotExpiresAt)
.first<RelaySigningKeyRecord>();
if (!record) {
throw new DomainError("relay_signing_key_unavailable", "Relay 授权签名公钥元数据不可用", 503, true, 5);
}
if (!/^env:NEKONEST_CLOUD_RELAY_SIGNING_PRIVATE_JWKS(?:#[A-Za-z0-9._:-]+)?$/u.test(record.private_key_ref)) {
throw new DomainError("relay_signing_key_unavailable", "Relay 私钥引用不受支持", 503, true, 5);
}
const privateKey = await crypto.subtle.importKey(
"jwk",
privateJwkForKid(record.kid),
{ name: "Ed25519" },
false,
["sign"],
);
return { record, privateKey };
}
export async function tenantConnectionState(accountId: string): Promise<TenantConnectionState> {
await ensureDatabase();
const row = await getD1()
.prepare(
`SELECT placements.state, nodes.status AS node_status
FROM tenant_instances AS tenants
LEFT JOIN tenant_placements AS placements ON placements.tenant_id = tenants.id
LEFT JOIN relay_nodes AS nodes ON nodes.id = placements.relay_node_id
WHERE tenants.account_id = ?`,
)
.bind(accountId)
.first<{ state: string | null; node_status: string | null }>();
if (["active", "draining"].includes(row?.state ?? "") && row?.node_status === "active") {
return { connectionState: "ready" };
}
return { connectionState: "provisioning", retryAfterSeconds: 5 };
}
export async function authenticateRelayNode(request: Request): Promise<RelayNodePrincipal> {
await ensureDatabase();
let mtlsIdentity;
try {
mtlsIdentity = await verifyTrustedRelayMtlsIdentity({
request,
assertionSecret:
env.NEKONEST_CLOUD_RELAY_INGRESS_ASSERTION_SECRET?.trim() ?? "",
});
} catch (error) {
if (error instanceof RelayMtlsIdentityError) {
throw new DomainError(error.code, error.message, error.status, error.status >= 500, 5);
}
throw error;
}
const nodeId = mtlsIdentity.nodeId;
const now = new Date().toISOString();
const db = getD1();
const result = await db
.prepare(AUTHENTICATE_RELAY_NODE_IDENTITY_SQL)
.bind(
now,
nodeId,
mtlsIdentity.spiffeId,
mtlsIdentity.certificateFingerprintSha256,
)
.run();
if (Number(result.meta.changes ?? 0) !== 1) {
throw new DomainError("relay_node_credential_invalid", "Relay 节点凭证无效", 401);
}
const node = await db
.prepare("SELECT region_id FROM relay_nodes WHERE id = ?")
.bind(nodeId)
.first<{ region_id: string }>();
if (!node) throw new DomainError("relay_node_credential_invalid", "Relay 节点凭证无效", 401);
return { nodeId, regionId: node.region_id };
}
export async function heartbeatRelayNode(input: {
principal: RelayNodePrincipal;
generation: number;
capacityTenants: number;
}): Promise<{ accepted: true; checked_at: string }> {
if (!Number.isSafeInteger(input.generation) || input.generation < 0) {
throw new DomainError("invalid_heartbeat_generation", "节点 heartbeat generation 无效");
}
if (!Number.isSafeInteger(input.capacityTenants) || input.capacityTenants < 0) {
throw new DomainError("invalid_relay_capacity", "节点容量无效");
}
const now = new Date().toISOString();
const result = await getD1()
.prepare(
`UPDATE relay_nodes
SET last_heartbeat_at = ?1, capacity_tenants = ?2,
heartbeat_generation = ?3, updated_at = ?1
WHERE id = ?4 AND region_id = ?5
AND heartbeat_generation <= ?3 AND status IN ('active', 'draining')`,
)
.bind(now, input.capacityTenants, input.generation, input.principal.nodeId, input.principal.regionId)
.run();
if (Number(result.meta.changes ?? 0) !== 1) {
throw new DomainError("stale_relay_heartbeat", "Relay 节点 heartbeat 已过期", 409);
}
return { accepted: true, checked_at: now };
}
async function placementForDevice(input: {
principal: RelayNodePrincipal;
deviceId: string;
tokenHash: string;
allowRemote?: boolean;
}): Promise<PlacementRecord> {
if (!/^host_[0-9a-f]{32}$/u.test(input.deviceId) || !/^[0-9a-f]{64}$/u.test(input.tokenHash)) {
throw new DomainError("device_credential_invalid", "设备凭证无效", 401);
}
const now = new Date().toISOString();
const row = await getD1()
.prepare(
`SELECT tenants.id AS tenant_id,
authorizations.status AS tenant_status,
regions.code AS home_region,
placements.relay_node_id,
placements.generation,
placements.state AS placement_state,
authorizations.revision AS authorization_revision,
nodes.internal_endpoint_ref,
nodes.status AS relay_node_status
FROM device_credentials AS credentials
INNER JOIN hosts ON hosts.id = credentials.host_id
INNER JOIN tenant_instances AS tenants ON tenants.account_id = hosts.account_id
INNER JOIN tenant_placements AS placements ON placements.tenant_id = tenants.id
INNER JOIN relay_regions AS regions ON regions.id = placements.home_region_id
INNER JOIN tenant_authorization_state AS authorizations ON authorizations.tenant_id = tenants.id
LEFT JOIN relay_nodes AS nodes ON nodes.id = placements.relay_node_id
WHERE hosts.id = ?1 AND credentials.token_hash = ?2
AND credentials.status = 'active' AND credentials.revoked_at IS NULL
AND (credentials.expires_at IS NULL OR credentials.expires_at > ?3)
AND hosts.lifecycle = 'active' AND hosts.slot_state = 'active'`,
)
.bind(input.deviceId, input.tokenHash, now)
.first<PlacementRecord>();
if (!row) throw new DomainError("device_credential_invalid", "设备凭证无效", 401);
if (row.tenant_status !== "active") {
throw new DomainError("access_suspended", "租户访问已暂停", 403);
}
if (!["active", "draining"].includes(row.placement_state) || !row.relay_node_id) {
throw new DomainError("service_provisioning", "租户 Relay 正在准备", 503, true, 5);
}
if (row.relay_node_status !== "active") {
throw new DomainError("region_unavailable", "目标 Relay 节点当前不可用", 503, true, 5);
}
if (!input.allowRemote && row.relay_node_id !== input.principal.nodeId) {
throw new DomainError("route_unavailable", "该节点不是租户当前写入节点", 503, true, 5);
}
return row;
}
export async function resolveDeviceRouteForRelay(input: {
principal: RelayNodePrincipal;
deviceId: string;
tokenHash: string;
}): Promise<RelayRouteResolution> {
await ensureDatabase();
return resolveRelayPlacementRoute(
await placementForDevice({ ...input, allowRemote: true }),
input.principal.nodeId,
);
}
export async function resolvePhoneRouteForRelay(input: {
principal: RelayNodePrincipal;
routeHandle: string;
phoneTokenHash: string;
}): Promise<RelayRouteResolution> {
await ensureDatabase();
const routeHash = await sha256Hex(input.routeHandle.trim());
const tokenHash = input.phoneTokenHash.trim().toLowerCase();
if (!/^[0-9a-f]{64}$/u.test(routeHash) || !/^[0-9a-f]{64}$/u.test(tokenHash)) {
throw new DomainError("phone_credential_invalid", "手机路由凭据无效", 401);
}
const activationCutoff = new Date(Date.now() - PHONE_HANDOFF_ACTIVATION_TTL_MS).toISOString();
const placement = await getD1()
.prepare(
`SELECT handles.tenant_id,
authorizations.status AS tenant_status,
regions.code AS home_region,
placements.relay_node_id,
placements.generation,
placements.state AS placement_state,
authorizations.revision AS authorization_revision,
nodes.internal_endpoint_ref,
nodes.status AS relay_node_status
FROM phone_route_handles AS handles
INNER JOIN relay_phone_principals AS phones ON phones.phone_id = handles.phone_id
INNER JOIN tenant_placements AS placements ON placements.tenant_id = handles.tenant_id
INNER JOIN relay_regions AS regions ON regions.id = placements.home_region_id
INNER JOIN tenant_authorization_state AS authorizations ON authorizations.tenant_id = handles.tenant_id
LEFT JOIN relay_nodes AS nodes ON nodes.id = placements.relay_node_id
WHERE handles.handle_hash = ?1 AND phones.token_hash = ?2
AND handles.revoked_at IS NULL AND phones.revoked_at IS NULL
AND phones.tenant_id = handles.tenant_id
AND (
(handles.status = 'active' AND phones.status = 'active')
OR (
handles.status = 'pending' AND phones.status = 'pending'
AND EXISTS (
SELECT 1 FROM phone_handoff_tickets AS tickets
WHERE tickets.tenant_id = handles.tenant_id
AND tickets.completed_phone_id = phones.phone_id
AND tickets.completed_phone_token_hash = ?2
AND tickets.completed_route_handle_hash = ?1
AND tickets.completed_at > ?3
AND tickets.activation_nonce IS NULL
)
)
)`,
)
.bind(routeHash, tokenHash, activationCutoff)
.first<PlacementRecord>();
if (!placement) throw new DomainError("phone_credential_invalid", "手机路由凭据无效", 401);
if (placement.tenant_status !== "active") {
throw new DomainError("access_suspended", "租户访问已暂停", 403);
}
return resolveRelayPlacementRoute(placement, input.principal.nodeId);
}
export async function resolveTenantRouteForRelay(input: {
principal: RelayNodePrincipal;
tenantId: string;
placementGeneration: number;
}): Promise<RelayRouteResolution> {
await ensureDatabase();
if (!/^tenant_[0-9a-f]{32}$/u.test(input.tenantId) ||
!Number.isSafeInteger(input.placementGeneration) || input.placementGeneration < 1) {
throw new DomainError("route_unavailable", "租户路由令牌无效", 404);
}
const placement = await getD1()
.prepare(
`SELECT placements.tenant_id,
authorizations.status AS tenant_status,
regions.code AS home_region,
placements.relay_node_id,
placements.generation,
placements.state AS placement_state,
authorizations.revision AS authorization_revision,
nodes.internal_endpoint_ref,
nodes.status AS relay_node_status
FROM tenant_placements AS placements
INNER JOIN relay_regions AS regions ON regions.id = placements.home_region_id
INNER JOIN tenant_authorization_state AS authorizations ON authorizations.tenant_id = placements.tenant_id
LEFT JOIN relay_nodes AS nodes ON nodes.id = placements.relay_node_id
WHERE placements.tenant_id = ?1 AND placements.generation = ?2`,
)
.bind(input.tenantId, input.placementGeneration)
.first<PlacementRecord>();
if (!placement) throw new DomainError("route_unavailable", "租户路由已过期", 404);
if (placement.tenant_status !== "active") {
throw new DomainError("access_suspended", "租户访问已暂停", 403);
}
return resolveRelayPlacementRoute(placement, input.principal.nodeId);
}
export async function resolveHandoffRouteForRelay(input: {
principal: RelayNodePrincipal;
ticket: string;
pwaOrigin: string;
}): Promise<RelayRouteResolution> {
await ensureDatabase();
const ticketHash = await sha256Hex(input.ticket.trim());
const origin = assertOrigin(input.pwaOrigin.trim());
const placement = await getD1()
.prepare(
`SELECT tickets.tenant_id,
authorizations.status AS tenant_status,
regions.code AS home_region,
placements.relay_node_id,
placements.generation,
placements.state AS placement_state,
authorizations.revision AS authorization_revision,
nodes.internal_endpoint_ref,
nodes.status AS relay_node_status
FROM phone_handoff_tickets AS tickets
INNER JOIN tenant_placements AS placements ON placements.tenant_id = tickets.tenant_id
INNER JOIN relay_regions AS regions ON regions.id = placements.home_region_id
INNER JOIN tenant_authorization_state AS authorizations ON authorizations.tenant_id = tickets.tenant_id
LEFT JOIN relay_nodes AS nodes ON nodes.id = placements.relay_node_id
WHERE tickets.ticket_hash = ?1 AND tickets.expected_origin = ?2
AND tickets.expires_at > ?3
AND (
tickets.consumed_at IS NULL
OR tickets.consumed_by_node_id = placements.relay_node_id
)`,
)
.bind(ticketHash, origin, new Date().toISOString())
.first<PlacementRecord>();
if (!placement) throw new DomainError("phone_credential_invalid", "手机交接凭证无效或已过期", 401);
if (placement.tenant_status !== "active") {
throw new DomainError("access_suspended", "租户访问已暂停", 403);
}
return resolveRelayPlacementRoute(placement, input.principal.nodeId);
}
export async function authorizeDeviceForRelay(input: {
principal: RelayNodePrincipal;
deviceId: string;
tokenHash: string;
}): Promise<{ snapshot: SignedRelayAuthorizationSnapshot; public_key_jwk: JsonWebKey }> {
await ensureDatabase();
const placement = await placementForDevice(input);
await getD1()
.prepare(
`UPDATE device_credentials
SET last_used_at = ?
WHERE host_id = ? AND token_hash = ? AND status = 'active' AND revoked_at IS NULL`,
)
.bind(new Date().toISOString(), input.deviceId, input.tokenHash)
.run();
return signedAuthorizationSnapshot(placement);
}
async function signedAuthorizationSnapshot(
placement: PlacementRecord,
): Promise<{ snapshot: SignedRelayAuthorizationSnapshot; public_key_jwk: JsonWebKey }> {
const devices = await getD1()
.prepare(
`SELECT hosts.id AS device_id, hosts.name, hosts.os,
hosts.ed25519_public, hosts.x25519_public,
credentials.token_hash AS credential_hash,
hosts.identity_fingerprint
FROM tenant_instances AS tenants
INNER JOIN hosts ON hosts.account_id = tenants.account_id
INNER JOIN device_credentials AS credentials ON credentials.host_id = hosts.id
WHERE tenants.id = ? AND hosts.lifecycle = 'active' AND hosts.slot_state = 'active'
AND credentials.status = 'active' AND credentials.revoked_at IS NULL
ORDER BY hosts.id ASC, credentials.id ASC`,
)
.bind(placement.tenant_id)
.all<AuthorizedDevice>();
const phones = await getD1()
.prepare(
`SELECT phone_id, name, token_hash AS credential_hash,
ed25519_public, x25519_public, identity_fingerprint
FROM relay_phone_principals
WHERE tenant_id = ? AND status = 'active' AND revoked_at IS NULL
ORDER BY phone_id ASC`,
)
.bind(placement.tenant_id)
.all<AuthorizedPhone>();
const nowMs = Date.now();
const issuedAt = new Date(nowMs).toISOString();
const expiresAt = new Date(
nowMs + RELAY_AUTHORIZATION_MAX_TTL_SECONDS * 1_000,
).toISOString();
const payload: RelayAuthorizationSnapshotPayload = {
snapshot_version: RELAY_AUTHORIZATION_SNAPSHOT_VERSION,
tenant_id: placement.tenant_id,
tenant_status: placement.tenant_status,
home_region: placement.home_region,
relay_node_id: placement.relay_node_id!,
placement_generation: placement.generation,
authorization_revision: placement.authorization_revision,
devices: devices.results ?? [],
phones: phones.results ?? [],
issued_at: issuedAt,
expires_at: expiresAt,
};
const signing = await loadSigningKey(issuedAt, expiresAt);
let publicKey: JsonWebKey;
try {
publicKey = JSON.parse(signing.record.public_key_jwk) as JsonWebKey;
} catch {
throw new DomainError("relay_signing_key_unavailable", "Relay 公钥元数据无效", 503, true, 5);
}
return {
snapshot: await signRelayAuthorizationSnapshot({
kid: signing.record.kid,
privateKey: signing.privateKey,
payload,
}),
public_key_jwk: publicKey,
};
}
export async function fullAuthorizationSnapshot(input: {
principal: RelayNodePrincipal;
tenantId: string;
placementGeneration: number;
}): Promise<{ snapshot: SignedRelayAuthorizationSnapshot; public_key_jwk: JsonWebKey }> {
if (!/^tenant_[0-9a-f]{32}$/u.test(input.tenantId)) {
throw new DomainError("invalid_tenant_id", "租户标识无效");
}
if (!Number.isSafeInteger(input.placementGeneration) || input.placementGeneration < 1) {
throw new DomainError("invalid_placement_generation", "placement generation 无效");
}
const placement = await getD1()
.prepare(
`SELECT placements.tenant_id, authorizations.status AS tenant_status,
regions.code AS home_region, placements.relay_node_id,
placements.generation, placements.state AS placement_state,
authorizations.revision AS authorization_revision
FROM tenant_placements AS placements
INNER JOIN relay_regions AS regions ON regions.id = placements.home_region_id
INNER JOIN tenant_authorization_state AS authorizations
ON authorizations.tenant_id = placements.tenant_id
WHERE placements.tenant_id = ?`,
)
.bind(input.tenantId)
.first<PlacementRecord>();
assertFullSnapshotPlacement({
placement,
nodeId: input.principal.nodeId,
placementGeneration: input.placementGeneration,
});
return signedAuthorizationSnapshot(placement!);
}
export function assertFullSnapshotPlacement(input: {
placement: PlacementRecord | null;
nodeId: string;
placementGeneration: number;
}): void {
const placement = input.placement;
const admission = classifySnapshotPlacement({
placement,
nodeId: input.nodeId,
expectedGeneration: input.placementGeneration,
});
if (admission === "wrong_node") {
throw new DomainError("route_unavailable", "租户不属于当前 Relay 节点", 404);
}
if (admission === "stale_generation") {
throw new DomainError("route_unavailable", "placement generation 已过期", 409, true, 1);
}
if (admission === "suspended") {
throw new DomainError("access_suspended", "租户访问已暂停", 403);
}
if (admission === "provisioning") {
throw new DomainError("service_provisioning", "租户 Relay 正在准备", 503, true, 5);
}
}
export async function authorizationRevisionDelta(input: {
principal: RelayNodePrincipal;
tenantId: string;
afterRevision: number;
}): Promise<{
tenant_id: string;
authorization_revision: number;
tenant_status: string;
changed: boolean;
checked_at: string;
}> {
if (!/^tenant_[0-9a-f]{32}$/u.test(input.tenantId)) {
throw new DomainError("invalid_tenant_id", "租户标识无效");
}
if (!Number.isSafeInteger(input.afterRevision) || input.afterRevision < 0) {
throw new DomainError("invalid_authorization_revision", "授权 revision 无效");
}
const row = await getD1()
.prepare(
`SELECT authorizations.revision, authorizations.status
FROM tenant_authorization_state AS authorizations
INNER JOIN tenant_placements AS placements ON placements.tenant_id = authorizations.tenant_id
WHERE authorizations.tenant_id = ? AND placements.relay_node_id = ?`,
)
.bind(input.tenantId, input.principal.nodeId)
.first<{ revision: number; status: string }>();
if (!row) throw new DomainError("route_unavailable", "租户不属于当前 Relay 节点", 404);
return {
tenant_id: input.tenantId,
authorization_revision: row.revision,
tenant_status: row.status,
changed: row.revision !== input.afterRevision,
checked_at: new Date().toISOString(),
};
}
export async function createPhoneHandoffTicket(input: {
accountId: string;
pwaOrigin: string;
}): Promise<{ ticket: string; pwa_url: string }> {
await ensureDatabase();
const origin = assertOrigin(input.pwaOrigin);
const tenant = await getD1()
.prepare(
`SELECT tenants.id
FROM tenant_instances AS tenants
INNER JOIN accounts ON accounts.id = tenants.account_id
WHERE tenants.account_id = ? AND accounts.status = 'active'`,
)
.bind(input.accountId)
.first<{ id: string }>();
if (!tenant) throw new DomainError("tenant_not_found", "租户不存在", 404);
const ticket = randomDeviceToken();
const ticketHash = await sha256Hex(ticket);
const nowMs = Date.now();
const now = new Date(nowMs).toISOString();
const expiresAt = new Date(nowMs + 60_000).toISOString();
await getD1()
.prepare(
`INSERT INTO phone_handoff_tickets
(id, ticket_hash, account_id, tenant_id, expected_origin, expires_at, created_at)
VALUES (?, ?, ?, ?, ?, ?, ?)`,
)
.bind(newId("handoff"), ticketHash, input.accountId, tenant.id, origin, expiresAt, now)
.run();
return {
ticket,
pwa_url: `${origin}/#handoff=${encodeURIComponent(ticket)}`,
};
}
export async function consumePhoneHandoffForRelay(input: {
principal: RelayNodePrincipal;
ticket: string;
pwaOrigin: string;
name: string;
phoneEd25519Public: string;
phoneX25519Public: string;
identityFingerprint: string;
}): Promise<{
handoff_id: string;
tenant_id: string;
name: string;
phone_ed25519_public: string;
phone_x25519_public: string;
identity_fingerprint: string;
placement_generation: number;
}> {
await ensureDatabase();
const origin = assertOrigin(input.pwaOrigin);
const ticket = input.ticket.trim().toLowerCase();
if (!/^[0-9a-f]{64}$/u.test(ticket)) {
throw new DomainError("phone_handoff_invalid", "手机交接凭证无效或已过期", 401);
}
const name = input.name.trim();
if (name.length < 1 || name.length > 48) {
throw new DomainError("invalid_phone_name", "手机名称需为 1 到 48 个字符");
}
let identity;
try {
identity = await validateDeviceIdentity({
ed25519Public: input.phoneEd25519Public,
x25519Public: input.phoneX25519Public,
identityFingerprint: input.identityFingerprint,
});
} catch {
throw new DomainError("invalid_phone_identity", "手机 E2E 身份无效");
}
const db = getD1();
const ticketHash = await sha256Hex(ticket);
const selectionSql =
`SELECT tickets.id, tickets.tenant_id, placements.generation,
tickets.consumed_at
FROM phone_handoff_tickets AS tickets
INNER JOIN tenant_placements AS placements ON placements.tenant_id = tickets.tenant_id
INNER JOIN tenant_authorization_state AS authorizations ON authorizations.tenant_id = tickets.tenant_id
WHERE tickets.ticket_hash = ? AND tickets.expected_origin = ?
AND tickets.expires_at > ?
AND placements.state IN ('active', 'draining') AND placements.relay_node_id = ?
AND authorizations.status = 'active'
AND (
tickets.consumed_at IS NULL
OR (
tickets.consumed_by_node_id = ?
AND tickets.pending_phone_name = ?
AND tickets.pending_ed25519_public = ?
AND tickets.pending_x25519_public = ?
AND tickets.pending_identity_fingerprint = ?
)
)`;
const selectRecord = () => db
.prepare(selectionSql)
.bind(
ticketHash,
origin,
new Date().toISOString(),
input.principal.nodeId,
input.principal.nodeId,
name,
identity.ed25519Public,
identity.x25519Public,
identity.fingerprint,
)
.first<{ id: string; tenant_id: string; generation: number; consumed_at: string | null }>();
let record = await selectRecord();
if (!record) throw new DomainError("phone_handoff_invalid", "手机交接凭证无效或已过期", 401);
if (record.consumed_at === null) {
const now = new Date().toISOString();
const result = await db
.prepare(CONSUME_PHONE_HANDOFF_SQL)
.bind(
now,
record.id,
ticketHash,
origin,
input.principal.nodeId,
name,
identity.ed25519Public,
identity.x25519Public,
identity.fingerprint,
)
.run();
if (Number(result.meta.changes ?? 0) !== 1) {
// An identical request may have won the one-shot consume update. Only
// resume it when every node and identity field still matches.
record = await selectRecord();
if (!record?.consumed_at) {
throw new DomainError("phone_handoff_invalid", "手机交接凭证无效或已过期", 401);
}
}
}
return {
handoff_id: record.id,
tenant_id: record.tenant_id,
name,
phone_ed25519_public: identity.ed25519Public,
phone_x25519_public: identity.x25519Public,
identity_fingerprint: identity.fingerprint,
placement_generation: record.generation,
};
}
export async function completePhoneHandoffForRelay(input: {
principal: RelayNodePrincipal;
handoffId: string;
phoneId: string;
phoneTokenHash: string;
routeHandleHash: string;
}): Promise<{ completed: true }> {
await ensureDatabase();
if (!/^handoff_[0-9a-f]{32}$/u.test(input.handoffId)) {
throw new DomainError("invalid_phone_handoff_id", "手机交接编号无效");
}
if (!/^phone_[A-Za-z0-9._:-]{1,120}$/u.test(input.phoneId)) {
throw new DomainError("invalid_phone_id", "手机身份编号无效");
}
const routeHandleHash = input.routeHandleHash.trim().toLowerCase();
const phoneTokenHash = input.phoneTokenHash.trim().toLowerCase();
if (!/^[0-9a-f]{64}$/u.test(routeHandleHash) || !/^[0-9a-f]{64}$/u.test(phoneTokenHash)) {
throw new DomainError("invalid_phone_route_hash", "手机路由摘要无效");
}
const now = new Date().toISOString();
const db = getD1();
type CompletionRecord = {
tenant_id: string;
consumed_by_node_id: string | null;
consumed_at: string | null;
completed_at: string | null;
completed_phone_id: string | null;
completed_phone_token_hash: string | null;
completed_route_handle_hash: string | null;
};
const loadCompletion = () => db
.prepare(
`SELECT tenant_id, consumed_by_node_id, consumed_at, completed_at,
completed_phone_id, completed_phone_token_hash, completed_route_handle_hash
FROM phone_handoff_tickets WHERE id = ?`,
)
.bind(input.handoffId)
.first<CompletionRecord>();
const isExactCompleted = async (record: CompletionRecord | null): Promise<boolean> => {
if (!record?.completed_at || record.consumed_by_node_id !== input.principal.nodeId ||
record.completed_phone_id !== input.phoneId ||
record.completed_phone_token_hash !== phoneTokenHash ||
record.completed_route_handle_hash !== routeHandleHash) {
return false;
}
const persisted = await db
.prepare(
`SELECT EXISTS (
SELECT 1
FROM relay_phone_principals AS phones
INNER JOIN phone_route_handles AS handles
ON handles.phone_id = phones.phone_id
AND handles.tenant_id = phones.tenant_id
WHERE phones.phone_id = ?1 AND phones.tenant_id = ?2
AND phones.token_hash = ?3 AND handles.handle_hash = ?4
AND phones.status = handles.status
AND phones.status IN ('pending', 'active')
AND phones.revoked_at IS NULL AND handles.revoked_at IS NULL
) AS completion_ok`,
)
.bind(
input.phoneId,
record.tenant_id,
phoneTokenHash,
routeHandleHash,
)
.first<{ completion_ok: number }>();
return persisted?.completion_ok === 1;
};
const initial = await loadCompletion();
if (!initial?.consumed_at || initial.consumed_by_node_id !== input.principal.nodeId) {
throw new DomainError("phone_handoff_completion_failed", "手机交接完成状态不确定", 409);
}
if (initial.completed_at) {
if (await isExactCompleted(initial)) return { completed: true };
throw new DomainError("phone_handoff_completion_conflict", "手机交接已由另一组凭据完成", 409);
}
let results;
try {
results = await db.batch([
db.prepare(DELETE_SUPERSEDED_PENDING_PHONE_ROUTES_SQL).bind(input.handoffId, now),
db.prepare(DELETE_SUPERSEDED_PENDING_PHONE_PRINCIPALS_SQL).bind(input.handoffId, now),
db
.prepare(
`INSERT INTO relay_phone_principals
(phone_id, tenant_id, token_hash, name, ed25519_public,
x25519_public, identity_fingerprint, status, created_at)
SELECT ?, tenant_id, ?, pending_phone_name, pending_ed25519_public,
pending_x25519_public, pending_identity_fingerprint, 'pending', ?
FROM phone_handoff_tickets
WHERE id = ? AND consumed_by_node_id = ? AND consumed_at IS NOT NULL
AND completed_at IS NULL AND pending_identity_fingerprint IS NOT NULL`,
)
.bind(input.phoneId, phoneTokenHash, now, input.handoffId, input.principal.nodeId),
db
.prepare(
`INSERT INTO phone_route_handles
(id, handle_hash, tenant_id, phone_id, status, created_at)
SELECT ?, ?, tenant_id, ?, 'pending', ?
FROM phone_handoff_tickets
WHERE id = ? AND consumed_by_node_id = ? AND consumed_at IS NOT NULL
AND completed_at IS NULL
AND EXISTS (
SELECT 1 FROM relay_phone_principals
WHERE phone_id = ? AND token_hash = ? AND status = 'pending'
)`,
)
.bind(
newId("route"),
routeHandleHash,
input.phoneId,
now,
input.handoffId,
input.principal.nodeId,
input.phoneId,
phoneTokenHash,
),
db
.prepare(
`UPDATE phone_handoff_tickets
SET completed_at = ?, completed_phone_id = ?,
completed_phone_token_hash = ?, completed_route_handle_hash = ?
WHERE id = ? AND consumed_by_node_id = ? AND consumed_at IS NOT NULL
AND completed_at IS NULL
AND EXISTS (SELECT 1 FROM phone_route_handles WHERE handle_hash = ? AND phone_id = ?)`,
)
.bind(
now,
input.phoneId,
phoneTokenHash,
routeHandleHash,
input.handoffId,
input.principal.nodeId,
routeHandleHash,
input.phoneId,
),
]);
} catch (error) {
if (await isExactCompleted(await loadCompletion())) return { completed: true };
throw error;
}
if (results.slice(-3).some((result) => Number(result.meta.changes ?? 0) !== 1)) {
if (await isExactCompleted(await loadCompletion())) return { completed: true };
throw new DomainError("phone_handoff_completion_failed", "手机交接完成状态不确定", 409);
}
return { completed: true };
}
export async function authorizePhoneRoute(input: {
principal: RelayNodePrincipal;
routeHandle: string;
phoneTokenHash: string;
}): Promise<{
tenant_id: string;
home_region: string;
relay_node_id: string;
placement_generation: number;
phone: {
phone_id: string;
name: string;
ed25519_public: string;
x25519_public: string;
identity_fingerprint: string;
};
}> {
await ensureDatabase();
const routeHandle = input.routeHandle.trim().toLowerCase();
const phoneTokenHash = input.phoneTokenHash.trim().toLowerCase();
if (!/^[0-9a-f]{64}$/u.test(routeHandle) || !/^[0-9a-f]{64}$/u.test(phoneTokenHash)) {
throw new DomainError("phone_route_invalid", "手机路由句柄无效", 401);
}
const hash = await sha256Hex(routeHandle);
const now = new Date().toISOString();
const db = getD1();
type AuthorizedPhoneRoute = {
tenant_id: string;
home_region: string;
relay_node_id: string;
generation: number;
phone_id: string;
name: string;
ed25519_public: string;
x25519_public: string;
identity_fingerprint: string;
};
const loadAuthorized = () => db
.prepare(AUTHORIZE_PHONE_ROUTE_SQL)
.bind(hash, phoneTokenHash)
.first<AuthorizedPhoneRoute>();
let row = await loadAuthorized();
if (!row) {
const activationNonce = newId("activation");
const activationCutoff = new Date(Date.now() - PHONE_HANDOFF_ACTIVATION_TTL_MS).toISOString();
const activationResults = await db.batch([
db.prepare(CLAIM_PHONE_HANDOFF_ACTIVATION_SQL).bind(
activationNonce,
hash,
phoneTokenHash,
activationCutoff,
input.principal.nodeId,
),
db.prepare(ACTIVATE_PHONE_PRINCIPAL_SQL).bind(activationNonce, phoneTokenHash),
db.prepare(ACTIVATE_PHONE_ROUTE_SQL).bind(activationNonce, hash, phoneTokenHash),
db.prepare(ADVANCE_AUTHORIZATION_AFTER_PHONE_ACTIVATION_SQL).bind(now, activationNonce),
db.prepare(FINALIZE_PHONE_HANDOFF_ACTIVATION_SQL).bind(now, activationNonce),
]);
row = await loadAuthorized();
if (activationResults.some((result) => Number(result.meta.changes ?? 0) !== 1) && !row) {
throw new DomainError("phone_route_invalid", "手机路由句柄无效", 401);
}
}
if (!row?.relay_node_id) throw new DomainError("phone_route_invalid", "手机路由句柄无效", 401);
if (row.relay_node_id !== input.principal.nodeId) {
throw new DomainError("route_unavailable", "该节点不是租户当前写入节点", 503, true, 5);
}
await db
.prepare("UPDATE phone_route_handles SET last_used_at = ? WHERE handle_hash = ? AND status = 'active'")
.bind(now, hash)
.run();
await db
.prepare("UPDATE relay_phone_principals SET last_used_at = ? WHERE phone_id = ? AND token_hash = ? AND status = 'active'")
.bind(now, row.phone_id, phoneTokenHash)
.run();
return {
tenant_id: row.tenant_id,
home_region: row.home_region,
relay_node_id: row.relay_node_id,
placement_generation: row.generation,
phone: {
phone_id: row.phone_id,
name: row.name,
ed25519_public: row.ed25519_public,
x25519_public: row.x25519_public,
identity_fingerprint: row.identity_fingerprint,
},
};
}
export async function revokePhoneForRelay(input: {
principal: RelayNodePrincipal;
tenantId: string;
phoneId: string;
reason: string;
}): Promise<{ revoked: true; authorization_revision: number }> {
if (!/^tenant_[0-9a-f]{32}$/u.test(input.tenantId)) {
throw new DomainError("invalid_tenant_id", "租户标识无效");
}
if (!/^phone_[A-Za-z0-9._:-]{1,120}$/u.test(input.phoneId)) {
throw new DomainError("invalid_phone_id", "手机身份编号无效");
}
const reason = input.reason.trim();
if (reason.length < 4 || reason.length > 256) {
throw new DomainError("invalid_revoke_reason", "撤销原因需为 4 到 256 个字符");
}
const db = getD1();
const placement = await db
.prepare(PHONE_FOR_NODE_REVOCATION_SQL)
.bind(input.tenantId, input.phoneId)
.first<{ relay_node_id: string | null; status: string }>();
if (!placement || placement.relay_node_id !== input.principal.nodeId) {
throw new DomainError("route_unavailable", "手机不属于当前 Relay placement", 404);
}
if (placement.status !== "active") {
throw new DomainError("phone_credential_invalid", "手机凭据已撤销", 401);
}
const now = new Date().toISOString();
const correlationId = newId("corr");
const results = await db.batch([
db
.prepare(REVOKE_PHONE_PRINCIPAL_SQL)
.bind(now, input.phoneId, input.tenantId),
db
.prepare(REVOKE_PHONE_ROUTES_SQL)
.bind(now, input.phoneId, input.tenantId),
db
.prepare(ADVANCE_AUTHORIZATION_AFTER_PHONE_REVOKE_SQL)
.bind(now, input.tenantId, input.phoneId),
db
.prepare(
`INSERT INTO audit_events
(id, actor_id, action, target_type, target_id, reason,
after_json, correlation_id, created_at)
SELECT ?, ?, 'phone.revoked', 'phone', ?, ?, ?, ?, ?
WHERE EXISTS (
SELECT 1 FROM relay_phone_principals
WHERE phone_id = ? AND tenant_id = ?
AND status = 'revoked' AND revoked_at = ?
)`,
)
.bind(
newId("audit"),
`relay:${input.principal.nodeId}`,
input.phoneId,
reason,
JSON.stringify({ tenant_id: input.tenantId }),
correlationId,
now,
input.phoneId,
input.tenantId,
now,
),
]);
if (Number(results[0]?.meta.changes ?? 0) !== 1
|| Number(results[2]?.meta.changes ?? 0) !== 1
|| Number(results[3]?.meta.changes ?? 0) !== 1) {
throw new DomainError("phone_revoke_indeterminate", "手机撤销状态不确定", 409);
}
const revision = await db
.prepare("SELECT revision FROM tenant_authorization_state WHERE tenant_id = ?")
.bind(input.tenantId)
.first<{ revision: number }>();
if (!revision) throw new DomainError("phone_revoke_indeterminate", "手机撤销状态不确定", 409);
return { revoked: true, authorization_revision: revision.revision };
}