import { ensureDatabase, getD1 } from "./bootstrap"; import { type RelayNodePrincipal } from "./relay-control-plane"; import { COMPLETED_RELAY_PURGE_PROOF_SQL, relayPurgeCompletionAuditId, relayPurgeRecordCanReturn, } from "./relay-purge-proof"; import { DomainError } from "./repository"; export type RelayPurgeRecord = { id: string; deletion_request_id: string; tenant_id: string; relay_node_id: string; placement_generation: number; state: "quiescing" | "completed" | "failed"; requested_by: string; reason: string; started_at: string; updated_at: string; completed_at: string | null; evidence_sha256: string | null; last_error_code: string | null; }; export type RelayPurgeAssignment = { purge_id: string; tenant_id: string; placement_generation: number; }; function validDeletionId(value: string): boolean { return /^deletion_[0-9a-f]{32}$/u.test(value); } function validPurgeId(value: string): boolean { return /^purge_[0-9a-f]{32}$/u.test(value); } function validErrorCode(value: string): boolean { return /^[a-z][a-z0-9_]{2,63}$/u.test(value); } function changed(result: D1Result): boolean { return Number(result.meta.changes ?? 0) === 1; } async function reloadPurge(id: string): Promise { const record = await getD1() .prepare("SELECT * FROM relay_purge_jobs WHERE id = ?") .bind(id) .first(); if (!record) { throw new DomainError("relay_purge_indeterminate", "租户删除状态不确定", 503, true, 5); } return validatePurgeRecordForReturn(record); } async function completedPurgeProofIsValid(record: RelayPurgeRecord): Promise { if (record.state !== "completed" || !/^[0-9a-f]{64}$/u.test(record.evidence_sha256 ?? "")) { return false; } const completionAuditId = relayPurgeCompletionAuditId(record.id); const proof = await getD1() .prepare(COMPLETED_RELAY_PURGE_PROOF_SQL) .bind(record.id, completionAuditId, record.evidence_sha256) .first<{ proof_valid: number }>(); return proof?.proof_valid === 1; } async function validatePurgeRecordForReturn( record: RelayPurgeRecord, ): Promise { const completedProofValid = record.state === "completed" ? await completedPurgeProofIsValid(record) : false; if (!relayPurgeRecordCanReturn(record.state, completedProofValid)) { throw new DomainError( "relay_purge_indeterminate", "租户删除记录缺少完整完成证明", 503, true, 5, ); } return record; } export async function beginRelayPurge(input: { deletionRequestId: string; actorId: string; reason: string; confirmation: string; }): Promise { await ensureDatabase(); const deletionRequestId = input.deletionRequestId.trim(); const actorId = input.actorId.trim(); const reason = input.reason.trim(); if (!validDeletionId(deletionRequestId) || !actorId || reason.length < 8 || reason.length > 500) { throw new DomainError("invalid_relay_purge", "删除申请、操作者或原因无效"); } if (input.confirmation !== "DELETE TENANT DATA") { throw new DomainError("relay_purge_confirmation_required", "必须明确确认永久删除租户数据"); } const db = getD1(); const existing = await db .prepare("SELECT * FROM relay_purge_jobs WHERE deletion_request_id = ?") .bind(deletionRequestId) .first(); const now = new Date().toISOString(); if (existing) { if (existing.state !== "failed") return validatePurgeRecordForReturn(existing); const results = await db.batch([ db.prepare( `UPDATE relay_purge_jobs SET state = 'quiescing', requested_by = ?1, reason = ?2, updated_at = ?3, completed_at = NULL, last_error_code = NULL WHERE id = ?4 AND state = 'failed'`, ).bind(actorId, reason, now, existing.id), db.prepare( `UPDATE tenant_placements SET state = 'deleting', last_error_code = NULL, updated_at = ?1 WHERE tenant_id = ?2 AND relay_node_id = ?3 AND generation = ?4 AND state IN ('deleting', 'deletion_failed')`, ).bind(now, existing.tenant_id, existing.relay_node_id, existing.placement_generation), db.prepare( `INSERT INTO audit_events (id, actor_id, action, target_type, target_id, reason, created_at) SELECT ?1, ?2, 'relay.purge.retried', 'relay_purge', ?3, ?4, ?5 WHERE EXISTS (SELECT 1 FROM relay_purge_jobs WHERE id = ?3 AND state = 'quiescing')`, ).bind(`audit_${crypto.randomUUID().replaceAll("-", "")}`, actorId, existing.id, reason, now), ]); if (!changed(results[0])) { throw new DomainError("relay_purge_conflict", "租户删除重试发生冲突", 409); } return reloadPurge(existing.id); } const purgeId = `purge_${crypto.randomUUID().replaceAll("-", "")}`; try { const results = await db.batch([ db.prepare( `INSERT INTO relay_purge_jobs (id, deletion_request_id, tenant_id, relay_node_id, placement_generation, state, requested_by, reason, started_at, updated_at) SELECT ?1, deletions.id, tenants.id, placements.relay_node_id, placements.generation, 'quiescing', ?3, ?4, ?5, ?5 FROM account_deletion_requests AS deletions INNER JOIN tenant_instances AS tenants ON tenants.account_id = deletions.account_id INNER JOIN tenant_placements AS placements ON placements.tenant_id = tenants.id INNER JOIN tenant_authorization_state AS authorizations ON authorizations.tenant_id = tenants.id INNER JOIN relay_nodes AS nodes ON nodes.id = placements.relay_node_id WHERE deletions.id = ?2 AND deletions.status = 'requested' AND placements.state = 'active' AND authorizations.status = 'active' AND nodes.status IN ('active', 'draining') AND NOT EXISTS ( SELECT 1 FROM relay_migrations AS migrations WHERE migrations.tenant_id = tenants.id AND migrations.state IN ('quiescing', 'copying', 'switching', 'draining') ) AND NOT EXISTS (SELECT 1 FROM relay_purge_jobs AS purges WHERE purges.tenant_id = tenants.id)`, ).bind(purgeId, deletionRequestId, actorId, reason, now), db.prepare( `UPDATE tenant_authorization_state SET status = 'suspended', revision = revision + 1, updated_at = ?1 WHERE tenant_id = (SELECT tenant_id FROM relay_purge_jobs WHERE id = ?2) AND status = 'active'`, ).bind(now, purgeId), db.prepare( `UPDATE tenant_placements SET state = 'deleting', last_error_code = NULL, updated_at = ?1 WHERE tenant_id = (SELECT tenant_id FROM relay_purge_jobs WHERE id = ?2) AND state = 'active'`, ).bind(now, purgeId), db.prepare( `UPDATE tenant_instances SET lifecycle = 'deleting', desired_state = 'deleted', relay_ready = 0, updated_at = ?1 WHERE id = (SELECT tenant_id FROM relay_purge_jobs WHERE id = ?2)`, ).bind(now, purgeId), db.prepare( `UPDATE account_deletion_requests SET status = 'processing', updated_at = ?1 WHERE id = ?2 AND status = 'requested' AND EXISTS (SELECT 1 FROM relay_purge_jobs WHERE id = ?3)`, ).bind(now, deletionRequestId, purgeId), db.prepare( `UPDATE accounts SET status = 'deleting', updated_at = ?1 WHERE id = (SELECT account_id FROM account_deletion_requests WHERE id = ?2) AND EXISTS (SELECT 1 FROM relay_purge_jobs WHERE id = ?3)`, ).bind(now, deletionRequestId, purgeId), db.prepare( `INSERT INTO audit_events (id, actor_id, action, target_type, target_id, reason, created_at) SELECT ?1, ?2, 'relay.purge.started', 'relay_purge', ?3, ?4, ?5 WHERE EXISTS (SELECT 1 FROM relay_purge_jobs WHERE id = ?3)`, ).bind(`audit_${crypto.randomUUID().replaceAll("-", "")}`, actorId, purgeId, reason, now), ]); if (!changed(results[0])) { throw new DomainError("relay_purge_conflict", "租户不可删除、正在迁移或已经删除", 409); } } catch (error) { const raced = await db .prepare("SELECT * FROM relay_purge_jobs WHERE deletion_request_id = ?") .bind(deletionRequestId) .first(); if (raced) return validatePurgeRecordForReturn(raced); throw error; } return reloadPurge(purgeId); } export async function relayPurgeAssignments( principal: RelayNodePrincipal, ): Promise { await ensureDatabase(); const rows = await getD1() .prepare( `SELECT id AS purge_id, tenant_id, placement_generation FROM relay_purge_jobs WHERE relay_node_id = ?1 AND state = 'quiescing' ORDER BY started_at ASC, id ASC LIMIT 4`, ) .bind(principal.nodeId) .all(); return rows.results; } export async function advanceRelayPurge(input: { principal: RelayNodePrincipal; purgeId: string; action: "completed" | "failed"; evidenceSha256?: string; errorCode?: string; }): Promise { await ensureDatabase(); const purgeId = input.purgeId.trim(); if (!validPurgeId(purgeId)) { throw new DomainError("invalid_relay_purge", "租户删除任务无效"); } const db = getD1(); const record = await db .prepare("SELECT * FROM relay_purge_jobs WHERE id = ?") .bind(purgeId) .first(); if (!record) throw new DomainError("relay_purge_not_found", "租户删除任务不存在", 404); if (record.relay_node_id !== input.principal.nodeId) { throw new DomainError("relay_purge_forbidden", "节点不属于该删除任务", 403); } if (record.state === input.action) { return validatePurgeRecordForReturn(record); } if (record.state !== "quiescing") { throw new DomainError("relay_purge_fence_conflict", "租户删除阶段已经变化", 409); } const now = new Date().toISOString(); if (input.action === "failed") { const errorCode = input.errorCode?.trim().toLowerCase() ?? "relay_purge_failed"; if (!validErrorCode(errorCode)) { throw new DomainError("invalid_relay_purge_failure", "租户删除失败码无效"); } const results = await db.batch([ db.prepare( `UPDATE relay_purge_jobs SET state = 'failed', last_error_code = ?1, completed_at = ?2, updated_at = ?2 WHERE id = ?3 AND state = 'quiescing' AND relay_node_id = ?4`, ).bind(errorCode, now, purgeId, input.principal.nodeId), db.prepare( `UPDATE tenant_placements SET state = 'deletion_failed', last_error_code = ?1, updated_at = ?2 WHERE tenant_id = ?3 AND relay_node_id = ?4 AND generation = ?5 AND state = 'deleting'`, ).bind(errorCode, now, record.tenant_id, record.relay_node_id, record.placement_generation), db.prepare( `INSERT INTO audit_events (id, actor_id, action, target_type, target_id, reason, created_at) SELECT ?1, ?2, 'relay.purge.failed', 'relay_purge', ?3, ?4, ?5 WHERE EXISTS (SELECT 1 FROM relay_purge_jobs WHERE id = ?3 AND state = 'failed')`, ).bind(`audit_${crypto.randomUUID().replaceAll("-", "")}`, input.principal.nodeId, purgeId, errorCode, now), ]); if (!changed(results[0])) { throw new DomainError("relay_purge_fence_conflict", "租户删除失败栅栏冲突", 409); } return reloadPurge(purgeId); } const evidence = input.evidenceSha256?.trim().toLowerCase() ?? ""; if (!/^[0-9a-f]{64}$/u.test(evidence)) { throw new DomainError("invalid_relay_purge_evidence", "租户删除证据摘要无效"); } const completionAuditId = relayPurgeCompletionAuditId(purgeId); const results = await db.batch([ db.prepare( `DELETE FROM device_registration_replays WHERE pairing_id IN ( SELECT pairings.id FROM pairing_requests AS pairings INNER JOIN tenant_instances AS tenants ON tenants.account_id = pairings.account_id WHERE tenants.id = ?1 )`, ).bind(record.tenant_id), db.prepare( `DELETE FROM pairing_claim_attempts WHERE pairing_request_id IN ( SELECT pairings.id FROM pairing_requests AS pairings INNER JOIN tenant_instances AS tenants ON tenants.account_id = pairings.account_id WHERE tenants.id = ?1 )`, ).bind(record.tenant_id), db.prepare( `DELETE FROM device_credentials WHERE host_id IN ( SELECT hosts.id FROM hosts INNER JOIN tenant_instances AS tenants ON tenants.account_id = hosts.account_id WHERE tenants.id = ?1 )`, ).bind(record.tenant_id), db.prepare("DELETE FROM phone_route_handles WHERE tenant_id = ?").bind(record.tenant_id), db.prepare("DELETE FROM relay_phone_principals WHERE tenant_id = ?").bind(record.tenant_id), db.prepare("DELETE FROM phone_handoff_tickets WHERE tenant_id = ?").bind(record.tenant_id), db.prepare( `DELETE FROM pairing_requests WHERE account_id = (SELECT account_id FROM tenant_instances WHERE id = ?1)`, ).bind(record.tenant_id), db.prepare( `UPDATE hosts SET name = 'Deleted host', lifecycle = 'deactivated', slot_state = 'released', connection_state = 'unknown', daemon_version = NULL, recovery_fingerprint = NULL, ed25519_public = NULL, x25519_public = NULL, identity_fingerprint = NULL, claim_request_id = NULL, last_seen_at = NULL, deactivated_at = ?1 WHERE account_id = (SELECT account_id FROM tenant_instances WHERE id = ?2)`, ).bind(now, record.tenant_id), db.prepare( `UPDATE tenant_authorization_state SET status = 'deleted', revision = revision + CASE WHEN status = 'suspended' THEN 1 ELSE 0 END, updated_at = ?1 WHERE tenant_id = ?2 AND status IN ('suspended', 'deleted')`, ).bind(now, record.tenant_id), db.prepare( `UPDATE tenant_placements SET relay_node_id = NULL, generation = CASE WHEN state = 'deleting' THEN generation + 1 ELSE generation END, state = 'deleted', last_error_code = NULL, updated_at = ?1 WHERE tenant_id = ?2 AND ( (relay_node_id = ?3 AND generation = ?4 AND state = 'deleting') OR (relay_node_id IS NULL AND generation = ?4 + 1 AND state = 'deleted') )`, ).bind(now, record.tenant_id, record.relay_node_id, record.placement_generation), db.prepare( `UPDATE tenant_instances SET lifecycle = 'deleted', desired_state = 'deleted', observed_state = 'deleted', relay_ready = 0, runtime_ref = NULL, relay_origin = NULL, secret_bundle_ref = NULL, tombstoned_at = ?1, updated_at = ?1 WHERE id = ?2 AND lifecycle IN ('deleting', 'deleted') AND desired_state = 'deleted'`, ).bind(now, record.tenant_id), db.prepare( `UPDATE account_deletion_requests SET status = 'relay_purged', updated_at = ?1 WHERE id = ?2 AND status IN ('processing', 'relay_purged')`, ).bind(now, record.deletion_request_id), db.prepare( `UPDATE accounts SET status = 'relay_purged', updated_at = ?1 WHERE id = (SELECT account_id FROM account_deletion_requests WHERE id = ?2) AND status IN ('deleting', 'relay_purged')`, ).bind(now, record.deletion_request_id), db.prepare( `INSERT OR IGNORE INTO audit_events (id, actor_id, action, target_type, target_id, reason, after_json, created_at) SELECT ?1, ?2, 'relay.purge.completed', 'relay_purge', ?3, 'Relay 实时数据与备份已完成应用层逻辑删除', json_object('evidence_sha256', ?4), ?5 WHERE EXISTS ( SELECT 1 FROM relay_purge_jobs WHERE id = ?3 AND state = 'quiescing' AND relay_node_id = ?2 ) AND EXISTS ( SELECT 1 FROM tenant_authorization_state WHERE tenant_id = ?6 AND status = 'deleted' ) AND EXISTS ( SELECT 1 FROM tenant_placements WHERE tenant_id = ?6 AND relay_node_id IS NULL AND generation = ?7 + 1 AND state = 'deleted' ) AND EXISTS ( SELECT 1 FROM tenant_instances WHERE id = ?6 AND lifecycle = 'deleted' AND desired_state = 'deleted' AND observed_state = 'deleted' ) AND EXISTS ( SELECT 1 FROM account_deletion_requests WHERE id = ?8 AND status = 'relay_purged' ) AND EXISTS ( SELECT 1 FROM accounts WHERE id = (SELECT account_id FROM account_deletion_requests WHERE id = ?8) AND status = 'relay_purged' ) AND NOT EXISTS ( SELECT 1 FROM device_credentials WHERE host_id IN ( SELECT hosts.id FROM hosts INNER JOIN tenant_instances AS tenants ON tenants.account_id = hosts.account_id WHERE tenants.id = ?6 ) ) AND NOT EXISTS (SELECT 1 FROM phone_route_handles WHERE tenant_id = ?6) AND NOT EXISTS (SELECT 1 FROM relay_phone_principals WHERE tenant_id = ?6) AND NOT EXISTS (SELECT 1 FROM phone_handoff_tickets WHERE tenant_id = ?6) AND NOT EXISTS ( SELECT 1 FROM hosts WHERE account_id = (SELECT account_id FROM tenant_instances WHERE id = ?6) AND (COALESCE(lifecycle, '') != 'deactivated' OR COALESCE(slot_state, '') != 'released') )`, ).bind( completionAuditId, input.principal.nodeId, purgeId, evidence, now, record.tenant_id, record.placement_generation, record.deletion_request_id, ), db.prepare( `UPDATE relay_purge_jobs SET state = 'completed', evidence_sha256 = ?1, completed_at = ?2, updated_at = ?2, last_error_code = NULL WHERE id = ?3 AND state = 'quiescing' AND relay_node_id = ?4 AND EXISTS ( SELECT 1 FROM audit_events WHERE id = ?5 AND actor_id = ?4 AND action = 'relay.purge.completed' AND target_type = 'relay_purge' AND target_id = ?3 AND json_extract(after_json, '$.evidence_sha256') = ?1 )`, ).bind(evidence, now, purgeId, input.principal.nodeId, completionAuditId), ]); if (!changed(results.at(-1)!)) { const current = await reloadPurge(purgeId); if (current.state === "completed") return validatePurgeRecordForReturn(current); throw new DomainError("relay_purge_indeterminate", "租户删除关键状态未完整收口", 503, true, 5); } return reloadPurge(purgeId); }