368 lines
12 KiB
TypeScript
368 lines
12 KiB
TypeScript
import { env } from "cloudflare:workers";
|
|
import migrationSql from "../drizzle/0000_condemned_legion.sql?raw";
|
|
import claimMigrationSql from "../drizzle/0001_mushy_vance_astro.sql?raw";
|
|
import claimOwnershipMigrationSql from "../drizzle/0002_wild_ravenous.sql?raw";
|
|
import provisioningMigrationSql from "../drizzle/0003_medical_rocket_racer.sql?raw";
|
|
import provisioningFenceMigrationSql from "../drizzle/0004_loud_prodigy.sql?raw";
|
|
import feedbackMigrationSql from "../drizzle/0005_pale_corsair.sql?raw";
|
|
import serviceIncidentMigrationSql from "../drizzle/0006_clever_shocker.sql?raw";
|
|
import accountLifecycleMigrationSql from "../drizzle/0007_zippy_nomad.sql?raw";
|
|
import provisionerLivenessMigrationSql from "../drizzle/0008_far_justice.sql?raw";
|
|
import maintenanceJobMigrationSql from "../drizzle/0009_flat_robbie_robertson.sql?raw";
|
|
import betaAccessRequestMigrationSql from "../drizzle/0010_windy_toxin.sql?raw";
|
|
import betaAccessRequestTimeIndexMigrationSql from "../drizzle/0011_next_thunderball.sql?raw";
|
|
import sharedRelayControlPlaneMigrationSql from "../drizzle/0012_shared_relay_control_plane.sql?raw";
|
|
import relayMigrationFencingMigrationSql from "../drizzle/0013_relay_migration_fencing.sql?raw";
|
|
import relayTenantPurgeMigrationSql from "../drizzle/0014_relay_tenant_purge.sql?raw";
|
|
import phoneHandoffIdempotencyMigrationSql from "../drizzle/0015_phone_handoff_idempotency.sql?raw";
|
|
import phoneHandoffActivationMigrationSql from "../drizzle/0016_phone_handoff_activation.sql?raw";
|
|
import provisioningInvariantMigrationSql from "../drizzle/9000_provisioning_invariants.sql?raw";
|
|
import provisioningSlugMigrationSql from "../drizzle/9001_provisioning_slug_backfill.sql?raw";
|
|
import readyCredentialReconciliationMigrationSql from "../drizzle/9002_ready_credential_reconciliation.sql?raw";
|
|
import { LAUNCH_GATE_SEEDS } from "./launch-gates.ts";
|
|
|
|
let schemaPromise: Promise<void> | null = null;
|
|
|
|
function getDatabase(): D1Database {
|
|
if (!env.DB) {
|
|
throw new Error("Cloudflare D1 binding `DB` is unavailable.");
|
|
}
|
|
return env.DB;
|
|
}
|
|
|
|
function migrationStatements(sql: string): string[] {
|
|
return sql
|
|
.split("--> statement-breakpoint")
|
|
.map((statement) => statement.trim())
|
|
.filter(Boolean);
|
|
}
|
|
|
|
async function applyTrackedMigration(
|
|
db: D1Database,
|
|
id: string,
|
|
sql: string,
|
|
): Promise<void> {
|
|
const applied = await db
|
|
.prepare("SELECT id FROM cloud_schema_migrations WHERE id = ?")
|
|
.bind(id)
|
|
.first<{ id: string }>();
|
|
if (applied) return;
|
|
|
|
const now = new Date().toISOString();
|
|
try {
|
|
await db.batch([
|
|
...migrationStatements(sql).map((statement) => db.prepare(statement)),
|
|
db
|
|
.prepare(
|
|
`INSERT INTO cloud_schema_migrations (id, applied_at)
|
|
VALUES (?, ?)`,
|
|
)
|
|
.bind(id, now),
|
|
]);
|
|
} catch (error) {
|
|
const racedMigration = await db
|
|
.prepare("SELECT id FROM cloud_schema_migrations WHERE id = ?")
|
|
.bind(id)
|
|
.first<{ id: string }>();
|
|
if (!racedMigration) throw error;
|
|
}
|
|
}
|
|
|
|
async function migrateSchema(db: D1Database): Promise<void> {
|
|
let accountsTable = await db
|
|
.prepare(
|
|
"SELECT name FROM sqlite_schema WHERE type = 'table' AND name = 'accounts'",
|
|
)
|
|
.first<{ name: string }>();
|
|
|
|
if (!accountsTable) {
|
|
try {
|
|
await db.batch(
|
|
migrationStatements(migrationSql).map((statement) => db.prepare(statement)),
|
|
);
|
|
} catch (error) {
|
|
accountsTable = await db
|
|
.prepare(
|
|
"SELECT name FROM sqlite_schema WHERE type = 'table' AND name = 'accounts'",
|
|
)
|
|
.first<{ name: string }>();
|
|
if (!accountsTable) throw error;
|
|
}
|
|
}
|
|
|
|
const ledgerTable = await db
|
|
.prepare(
|
|
"SELECT name FROM sqlite_schema WHERE type = 'table' AND name = 'cloud_schema_migrations'",
|
|
)
|
|
.first<{ name: string }>();
|
|
const claimMigration = ledgerTable
|
|
? await db
|
|
.prepare("SELECT id FROM cloud_schema_migrations WHERE id = ?")
|
|
.bind("0001_mushy_vance_astro")
|
|
.first<{ id: string }>()
|
|
: null;
|
|
|
|
if (!claimMigration) {
|
|
const statements = migrationStatements(claimMigrationSql);
|
|
const applicableStatements = ledgerTable ? statements.slice(1) : statements;
|
|
const now = new Date().toISOString();
|
|
try {
|
|
await db.batch([
|
|
...applicableStatements.map((statement) => db.prepare(statement)),
|
|
db
|
|
.prepare(
|
|
`INSERT OR IGNORE INTO cloud_schema_migrations (id, applied_at)
|
|
VALUES ('0000_condemned_legion', ?), ('0001_mushy_vance_astro', ?)`,
|
|
)
|
|
.bind(now, now),
|
|
]);
|
|
} catch (error) {
|
|
const racedMigration = await db
|
|
.prepare("SELECT id FROM cloud_schema_migrations WHERE id = ?")
|
|
.bind("0001_mushy_vance_astro")
|
|
.first<{ id: string }>();
|
|
if (!racedMigration) throw error;
|
|
}
|
|
}
|
|
|
|
await applyTrackedMigration(
|
|
db,
|
|
"0002_wild_ravenous",
|
|
claimOwnershipMigrationSql,
|
|
);
|
|
await applyTrackedMigration(
|
|
db,
|
|
"0003_medical_rocket_racer",
|
|
provisioningMigrationSql,
|
|
);
|
|
await applyTrackedMigration(
|
|
db,
|
|
"0004_loud_prodigy",
|
|
provisioningFenceMigrationSql,
|
|
);
|
|
await applyTrackedMigration(db, "0005_pale_corsair", feedbackMigrationSql);
|
|
await applyTrackedMigration(
|
|
db,
|
|
"0006_clever_shocker",
|
|
serviceIncidentMigrationSql,
|
|
);
|
|
await applyTrackedMigration(
|
|
db,
|
|
"0007_zippy_nomad",
|
|
accountLifecycleMigrationSql,
|
|
);
|
|
await applyTrackedMigration(
|
|
db,
|
|
"0008_far_justice",
|
|
provisionerLivenessMigrationSql,
|
|
);
|
|
await applyTrackedMigration(
|
|
db,
|
|
"0009_flat_robbie_robertson",
|
|
maintenanceJobMigrationSql,
|
|
);
|
|
await applyTrackedMigration(
|
|
db,
|
|
"0010_windy_toxin",
|
|
betaAccessRequestMigrationSql,
|
|
);
|
|
await applyTrackedMigration(
|
|
db,
|
|
"0011_next_thunderball",
|
|
betaAccessRequestTimeIndexMigrationSql,
|
|
);
|
|
await applyTrackedMigration(
|
|
db,
|
|
"9000_provisioning_invariants",
|
|
provisioningInvariantMigrationSql,
|
|
);
|
|
await applyTrackedMigration(
|
|
db,
|
|
"9001_provisioning_slug_backfill",
|
|
provisioningSlugMigrationSql,
|
|
);
|
|
await applyTrackedMigration(
|
|
db,
|
|
"9002_ready_credential_reconciliation",
|
|
readyCredentialReconciliationMigrationSql,
|
|
);
|
|
await applyTrackedMigration(
|
|
db,
|
|
"0012_shared_relay_control_plane",
|
|
sharedRelayControlPlaneMigrationSql,
|
|
);
|
|
await applyTrackedMigration(
|
|
db,
|
|
"0013_relay_migration_fencing",
|
|
relayMigrationFencingMigrationSql,
|
|
);
|
|
await applyTrackedMigration(
|
|
db,
|
|
"0014_relay_tenant_purge",
|
|
relayTenantPurgeMigrationSql,
|
|
);
|
|
await applyTrackedMigration(
|
|
db,
|
|
"0015_phone_handoff_idempotency",
|
|
phoneHandoffIdempotencyMigrationSql,
|
|
);
|
|
await applyTrackedMigration(
|
|
db,
|
|
"0016_phone_handoff_activation",
|
|
phoneHandoffActivationMigrationSql,
|
|
);
|
|
await db.prepare("PRAGMA optimize").run();
|
|
}
|
|
|
|
async function seedCatalogAndGates(db: D1Database): Promise<void> {
|
|
const now = new Date().toISOString();
|
|
const gateStatements = LAUNCH_GATE_SEEDS.map(([key, priority, category, title]) =>
|
|
db
|
|
.prepare(
|
|
`INSERT OR IGNORE INTO launch_gates
|
|
(key, priority, category, title, status, notes, updated_at)
|
|
VALUES (?, ?, ?, ?, 'blocked', '', ?)`,
|
|
)
|
|
.bind(key, priority, category, title, now),
|
|
);
|
|
const gateReconciliationStatements = LAUNCH_GATE_SEEDS.map(
|
|
([key, priority, category, title]) =>
|
|
db
|
|
.prepare(
|
|
`UPDATE launch_gates
|
|
SET priority = ?, category = ?, title = ?, updated_at = ?
|
|
WHERE key = ?
|
|
AND (priority <> ? OR category <> ? OR title <> ?)`,
|
|
)
|
|
.bind(priority, category, title, now, key, priority, category, title),
|
|
);
|
|
const gateReconciliationAuditStatements = LAUNCH_GATE_SEEDS.map(
|
|
([key, priority, category, title]) =>
|
|
db
|
|
.prepare(
|
|
`INSERT OR IGNORE INTO audit_events
|
|
(id, actor_id, action, target_type, target_id, reason,
|
|
before_json, after_json, correlation_id, created_at)
|
|
SELECT 'audit_reclassify_' || key || '_free_beta_v1',
|
|
'system:bootstrap', 'launch_gate.reclassified',
|
|
'launch_gate', key,
|
|
'免费公测门禁与未来收费门禁分离',
|
|
json_object('priority', priority, 'category', category, 'title', title),
|
|
json_object('priority', ?, 'category', ?, 'title', ?),
|
|
'corr_reclassify_' || key || '_free_beta_v1', ?
|
|
FROM launch_gates
|
|
WHERE key = ?
|
|
AND (priority <> ? OR category <> ? OR title <> ?)`,
|
|
)
|
|
.bind(priority, category, title, now, key, priority, category, title),
|
|
);
|
|
|
|
await db.batch([
|
|
db
|
|
.prepare(
|
|
`INSERT OR IGNORE INTO price_versions
|
|
(id, product_code, billing_period, unit_slots, amount_minor, currency,
|
|
tax_mode, quote_ttl_seconds, status, effective_from, created_by)
|
|
VALUES (?, 'host_slot', 'month', 1, 1000, 'CNY', 'undecided', 900,
|
|
'retired', ?, 'system:seed')`,
|
|
)
|
|
.bind("price_host_month_v1", now),
|
|
db
|
|
.prepare(
|
|
`INSERT OR IGNORE INTO price_versions
|
|
(id, product_code, billing_period, unit_slots, amount_minor, currency,
|
|
tax_mode, quote_ttl_seconds, status, effective_from, created_by)
|
|
VALUES (?, 'host_slot', 'year', 1, 10000, 'CNY', 'undecided', 900,
|
|
'retired', ?, 'system:seed')`,
|
|
)
|
|
.bind("price_host_year_v1", now),
|
|
db
|
|
.prepare(
|
|
`INSERT OR IGNORE INTO beta_programs
|
|
(id, state, capacity_slots, starts_at, ends_at, grace_days, created_by)
|
|
VALUES (?, 'active', NULL, ?, NULL, 0, 'system:seed')`,
|
|
)
|
|
.bind("beta_public_v1", now),
|
|
...gateStatements,
|
|
...gateReconciliationAuditStatements,
|
|
...gateReconciliationStatements,
|
|
db
|
|
.prepare(
|
|
`INSERT OR IGNORE INTO audit_events
|
|
(id, actor_id, action, target_type, target_id, reason,
|
|
before_json, after_json, correlation_id, created_at)
|
|
SELECT 'audit_defer_paid_catalog_v1', 'system:bootstrap',
|
|
'price_catalog.retired', 'price_catalog', 'host_slot',
|
|
'免费公测阶段暂缓收费决策,保留历史价格但撤销发布状态',
|
|
'{"status":"published"}', '{"status":"retired"}',
|
|
'corr_defer_paid_catalog_v1', ?
|
|
WHERE EXISTS (
|
|
SELECT 1 FROM price_versions
|
|
WHERE product_code = 'host_slot' AND status = 'published'
|
|
)`,
|
|
)
|
|
.bind(now),
|
|
db.prepare(
|
|
`UPDATE price_versions
|
|
SET status = 'retired'
|
|
WHERE product_code = 'host_slot' AND status = 'published'`,
|
|
),
|
|
db
|
|
.prepare(
|
|
`INSERT OR IGNORE INTO audit_events
|
|
(id, actor_id, action, target_type, target_id, reason,
|
|
before_json, after_json, correlation_id, created_at)
|
|
SELECT 'audit_repair_p0_' || key, 'system:bootstrap',
|
|
'launch_gate.repaired', 'launch_gate', key,
|
|
'旧版 P0 状态缺少有效证据或使用了不适用,启动时恢复为阻止',
|
|
'{"status":"legacy_invalid"}', '{"status":"blocked"}',
|
|
'corr_repair_p0_' || key, ?
|
|
FROM launch_gates
|
|
WHERE priority = 'P0' AND (
|
|
status = 'not_applicable'
|
|
OR (status = 'passed' AND (
|
|
trim(COALESCE(owner, '')) = ''
|
|
OR trim(COALESCE(notes, '')) = ''
|
|
OR evidence_url IS NULL
|
|
OR evidence_url NOT LIKE 'https://%'
|
|
))
|
|
)`,
|
|
)
|
|
.bind(now),
|
|
db
|
|
.prepare(
|
|
`UPDATE launch_gates
|
|
SET status = 'blocked', reviewed_at = ?, updated_at = ?
|
|
WHERE priority = 'P0' AND (
|
|
status = 'not_applicable'
|
|
OR (status = 'passed' AND (
|
|
trim(COALESCE(owner, '')) = ''
|
|
OR trim(COALESCE(notes, '')) = ''
|
|
OR evidence_url IS NULL
|
|
OR evidence_url NOT LIKE 'https://%'
|
|
))
|
|
)`,
|
|
)
|
|
.bind(now, now),
|
|
]);
|
|
}
|
|
|
|
export async function ensureDatabase(): Promise<void> {
|
|
if (!schemaPromise) {
|
|
schemaPromise = (async () => {
|
|
const db = getDatabase();
|
|
await migrateSchema(db);
|
|
await seedCatalogAndGates(db);
|
|
})().catch((error) => {
|
|
schemaPromise = null;
|
|
throw error;
|
|
});
|
|
}
|
|
|
|
await schemaPromise;
|
|
}
|
|
|
|
export function getD1(): D1Database {
|
|
return getDatabase();
|
|
}
|