From 5bed2bae9fec0e682ad26fcdf00afcbc4c07e814 Mon Sep 17 00:00:00 2001 From: ryan Date: Fri, 19 Jun 2026 21:31:21 +0800 Subject: [PATCH] migrate database --- config.example.yaml | 4 +- docs/changelog/index.md | 6 + .../db/migrator/202606050001_bridge_legacy.go | 204 +++++++ .../202606200006_migrate_legacy_data.go | 538 ++++++++++++++++++ internal/db/migrator/bridge_test.go | 270 +++++++++ internal/db/migrator/clickhouse.go | 26 +- internal/db/migrator/migrator.go | 24 +- 7 files changed, 1062 insertions(+), 10 deletions(-) create mode 100644 internal/db/migrator/202606050001_bridge_legacy.go create mode 100644 internal/db/migrator/202606200006_migrate_legacy_data.go create mode 100644 internal/db/migrator/bridge_test.go diff --git a/config.example.yaml b/config.example.yaml index 033f3d90..4d995c13 100644 --- a/config.example.yaml +++ b/config.example.yaml @@ -24,8 +24,8 @@ database: sqlite_path: "openflare.db" # PostgreSQL 禁用时使用此 SQLite 文件路径 host: "127.0.0.1" port: 5432 - username: "postgres" - password: "postgres" + username: "openflare" + password: "replace-with-strong-password" database: "openflare" max_idle_conn: 16 max_open_conn: 128 diff --git a/docs/changelog/index.md b/docs/changelog/index.md index d47a46fb..1843f5b0 100644 --- a/docs/changelog/index.md +++ b/docs/changelog/index.md @@ -24,6 +24,9 @@ sidebar: false ### 新增 +- 新增 `goose` 数据库平滑升级桥接机制:在全新命名空间中通过数据迁移(而非改表名)迁移合并 legacy 旧版 SQLite/PostgreSQL 生产数据(Schema `202606040004` 及以下),并在完成后清理 `legacy_` 临时表。 +- 新增 SQLite 迁移底层 `goose_db_version` 表的主键 `AUTOINCREMENT` 自动补全修复,避免由于旧版 Goose schema 限制导致的新升级写入冲突。 +- 新增 `goose.NewProvider` 及 `goose.WithDisableGlobalRegistry` 用于隔离 ClickHouse 与 SQLite/PostgreSQL 间的 Go 代码全局迁移污染。 - 新增 `internal/db/batchwriter` 通用批量写入框架,支持各业务域独立队列实例、按条数/时间 flush、非阻塞入队与优雅停机。 - 业务层接入批量写入:`risk_control` 审计日志迁移至 `batchwriter`;OpenFlare 可观测时序与节点访问日志通过 `internal/apps/openflare/chwriter` 异步 flush,移除写前 `SELECT count()` 去重。 @@ -32,6 +35,9 @@ sidebar: false ### 修复 +- 修复 PostgreSQL 下 legacy 桥接迁移 `202606050001` / `202606200006` 使用 `?` 占位符导致 `syntax error at end of input` 的启动失败问题。 +- 修复 PostgreSQL legacy 数据迁移时旧表可空字段写入新表 `NOT NULL` 列触发约束错误的问题(`INSERT ... SELECT` 不会自动套用列默认值)。 +- 修复 PostgreSQL legacy `dns_accounts` 迁移未转义保留字列名 `authorization` 导致语法错误的问题。 - 修复总览看板「24 小时请求趋势」摘要误展示 24 小时累计值的问题:趋势图摘要改为「当前小时」桶数据,顶部 24h 统计改为按小时趋势聚合。 - 修复访问日志 ClickHouse 聚合查询因 `trim(x) AS x` 别名与表列同名导致总览看板地域分布及 IP 统计失败的问题。 - 修复访问日志页 `count()` 扫描类型不匹配(ClickHouse `UInt64` 写入 `int64`)导致列表计数失败的问题。 diff --git a/internal/db/migrator/202606050001_bridge_legacy.go b/internal/db/migrator/202606050001_bridge_legacy.go new file mode 100644 index 00000000..26ba4e0e --- /dev/null +++ b/internal/db/migrator/202606050001_bridge_legacy.go @@ -0,0 +1,204 @@ +// Copyright 2026 Arctel.net +// SPDX-License-Identifier: Apache-2.0 + +package migrator + +import ( + "context" + "database/sql" + "fmt" + "strings" + + "github.com/pressly/goose/v3" +) + +func init() { + goose.AddMigrationContext(up202606050001, down202606050001) +} + +func up202606050001(ctx context.Context, tx *sql.Tx) error { + dialect := gooseDialect() + if dialect == dialectSqlite { + if err := fixSqliteGooseDbVersion(ctx, tx); err != nil { + return err + } + } + tableExistsQuery := tableExistsSQL(dialect) + + checkTable := func(name string) (bool, error) { + var count int + err := tx.QueryRowContext(ctx, tableExistsQuery, name).Scan(&count) + return count > 0, err + } + + // Check if this is a legacy database (by checking if the "options" table exists). + // On a clean install, "options" won't exist, so this will be a quick no-op. + hasOptions, err := checkTable("options") + if err != nil { + return fmt.Errorf("check table options failed: %w", err) + } + + if !hasOptions { + // Clean install or already migrated, no action needed. + return nil + } + + // Helper to drop table if it exists + dropTable := func(tableName string) error { + dropSQL := fmt.Sprintf("DROP TABLE IF EXISTS %s", tableName) + if dialect == dialectPostgres { + dropSQL += cascadeSuffix + } + if _, err := tx.ExecContext(ctx, dropSQL); err != nil { + return fmt.Errorf("drop table %s failed: %w", tableName, err) + } + return nil + } + + // 1. Drop legacy observability and log tables to free space and avoid conflict + obsPrefixes := []string{ + "node_system_profiles", + "node_health_events", + "node_access_logs_", + "node_observation_openresties_", + "node_observation_frpcs_", + "node_observation_frps_", + "node_metric_snapshots_", + "node_request_reports_", + } + + for _, prefix := range obsPrefixes { + if err := dropTablesWithPrefix(ctx, tx, tablesWithPrefixSQL(dialect), prefix, dropTable); err != nil { + return err + } + } + + // 2. Rename custom and platform legacy tables to legacy_ prefix + tablesToRename := []string{ + "users", + "auth_sources", + "external_accounts", + "options", + "origins", + "apply_logs", + "proxy_routes", + "nodes", + "waf_rule_groups", + "waf_rule_group_bindings", + "waf_ip_groups", + "tls_certificates", + "managed_domains", + "dns_accounts", + "acme_accounts", + "config_versions", + "pages_projects", + "pages_deployments", + "pages_deployment_files", + } + + for _, oldName := range tablesToRename { + if err := renameLegacyTable(ctx, tx, oldName, checkTable, dropTable); err != nil { + return err + } + } + + return nil +} + +func dropTablesWithPrefix(ctx context.Context, tx *sql.Tx, query, prefix string, dropTable func(string) error) error { + rows, err := tx.QueryContext(ctx, query, prefix+"%") + if err != nil { + return fmt.Errorf("query legacy tables for prefix %s failed: %w", prefix, err) + } + defer func() { + _ = rows.Close() + }() + + var tables []string + for rows.Next() { + var name string + if err := rows.Scan(&name); err != nil { + return err + } + tables = append(tables, name) + } + + for _, table := range tables { + if err := dropTable(table); err != nil { + return err + } + } + return nil +} + +func renameLegacyTable(ctx context.Context, tx *sql.Tx, oldName string, checkTable func(string) (bool, error), dropTable func(string) error) error { + exists, err := checkTable(oldName) + if err != nil { + return err + } + if !exists { + return nil + } + + newName := "legacy_" + oldName + newExists, err := checkTable(newName) + if err != nil { + return err + } + if newExists { + if err := dropTable(newName); err != nil { + return err + } + } + + renameSQL := fmt.Sprintf("ALTER TABLE %s RENAME TO %s", oldName, newName) + if _, err := tx.ExecContext(ctx, renameSQL); err != nil { + return fmt.Errorf("rename table %s to %s failed: %w", oldName, newName, err) + } + return nil +} + +func down202606050001(_ context.Context, _ *sql.Tx) error { + // Down migration is a no-op as the legacy DB is rolled forward during migration. + return nil +} + +func fixSqliteGooseDbVersion(ctx context.Context, tx *sql.Tx) error { + var createSQL string + err := tx.QueryRowContext(ctx, "SELECT sql FROM sqlite_master WHERE type='table' AND name='goose_db_version'").Scan(&createSQL) + if err != nil { + if err == sql.ErrNoRows { + return nil + } + return err + } + if strings.Contains(strings.ToUpper(createSQL), "AUTOINCREMENT") { + return nil // Already upgraded + } + + if _, err := tx.ExecContext(ctx, "ALTER TABLE goose_db_version RENAME TO old_goose_db_version"); err != nil { + return fmt.Errorf("rename goose_db_version failed: %w", err) + } + + newTableSQL := `CREATE TABLE goose_db_version ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + version_id INTEGER NOT NULL, + is_applied INTEGER NOT NULL, + tstamp TIMESTAMP DEFAULT CURRENT_TIMESTAMP + );` + if _, err := tx.ExecContext(ctx, newTableSQL); err != nil { + return fmt.Errorf("create new goose_db_version failed: %w", err) + } + + copySQL := `INSERT INTO goose_db_version (id, version_id, is_applied, tstamp) + SELECT id, version_id, is_applied, tstamp FROM old_goose_db_version;` + if _, err := tx.ExecContext(ctx, copySQL); err != nil { + return fmt.Errorf("copy goose_db_version data failed: %w", err) + } + + if _, err := tx.ExecContext(ctx, "DROP TABLE old_goose_db_version"); err != nil { + return fmt.Errorf("drop old_goose_db_version failed: %w", err) + } + + return nil +} diff --git a/internal/db/migrator/202606200006_migrate_legacy_data.go b/internal/db/migrator/202606200006_migrate_legacy_data.go new file mode 100644 index 00000000..21f0b616 --- /dev/null +++ b/internal/db/migrator/202606200006_migrate_legacy_data.go @@ -0,0 +1,538 @@ +// Copyright 2026 Arctel.net +// SPDX-License-Identifier: Apache-2.0 + +package migrator + +import ( + "context" + "database/sql" + "fmt" + + "github.com/pressly/goose/v3" +) + +func init() { + goose.AddMigrationContext(up202606200006, down202606200006) +} + +type migrationTask struct { + legacyName string + sqliteSQL string + pgSQL string +} + +func up202606200006(ctx context.Context, tx *sql.Tx) error { + dialect := gooseDialect() + tableExistsQuery := tableExistsSQL(dialect) + + checkTable := func(name string) (bool, error) { + var count int + err := tx.QueryRowContext(ctx, tableExistsQuery, name).Scan(&count) + return count > 0, err + } + + // Helper to migrate table if it exists + migrateTable := func(task migrationTask) error { + exists, err := checkTable(task.legacyName) + if err != nil { + return err + } + if !exists { + return nil + } + var sqlToRun string + if dialect == dialectSqlite { + sqlToRun = task.sqliteSQL + } else { + sqlToRun = task.pgSQL + } + if _, err := tx.ExecContext(ctx, sqlToRun); err != nil { + return fmt.Errorf("execute migration query for %s failed: %w", task.legacyName, err) + } + return nil + } + + tasks := getMigrationTasks() + for _, task := range tasks { + if err := migrateTable(task); err != nil { + return err + } + } + + // Drop legacy tables that exist + legacyTables := []string{ + "legacy_users", + "legacy_auth_sources", + "legacy_external_accounts", + "legacy_options", + "legacy_origins", + "legacy_apply_logs", + "legacy_proxy_routes", + "legacy_nodes", + "legacy_waf_rule_groups", + "legacy_waf_rule_group_bindings", + "legacy_waf_ip_groups", + "legacy_tls_certificates", + "legacy_managed_domains", + "legacy_dns_accounts", + "legacy_acme_accounts", + "legacy_config_versions", + "legacy_pages_projects", + "legacy_pages_deployments", + "legacy_pages_deployment_files", + } + + for _, table := range legacyTables { + exists, err := checkTable(table) + if err != nil { + return err + } + if exists { + dropSQL := fmt.Sprintf("DROP TABLE IF EXISTS %s", table) + if dialect == dialectPostgres { + dropSQL += cascadeSuffix + } + if _, err := tx.ExecContext(ctx, dropSQL); err != nil { + return fmt.Errorf("drop legacy table %s failed: %w", table, err) + } + } + } + + return nil +} + +func getMigrationTasks() []migrationTask { + var tasks []migrationTask + tasks = append(tasks, getUserAndOptionTasks()...) + tasks = append(tasks, getOriginAndRouteTasks()...) + tasks = append(tasks, getWafTasks()...) + tasks = append(tasks, getTLSTasks()...) + tasks = append(tasks, getPagesAndConfigTasks()...) + return tasks +} + +func getUserAndOptionTasks() []migrationTask { + return []migrationTask{ + { + legacyName: "legacy_users", + sqliteSQL: `INSERT OR REPLACE INTO w_users (id, username, password, nickname, email, is_active, is_admin, created_at, updated_at) + SELECT id, username, password, display_name, email, + CASE WHEN status = 1 THEN 1 ELSE 0 END, + CASE WHEN role = 100 THEN 1 ELSE 0 END, + CURRENT_TIMESTAMP, CURRENT_TIMESTAMP + FROM legacy_users;`, + pgSQL: `INSERT INTO w_users (id, username, password, nickname, email, is_active, is_admin, created_at, updated_at) + SELECT id, username, password, display_name, email, + CASE WHEN status = 1 THEN TRUE ELSE FALSE END, + CASE WHEN role = 100 THEN TRUE ELSE FALSE END, + CURRENT_TIMESTAMP, CURRENT_TIMESTAMP + FROM legacy_users + ON CONFLICT (id) DO UPDATE SET + username = EXCLUDED.username, + password = EXCLUDED.password, + nickname = EXCLUDED.nickname, + email = EXCLUDED.email, + is_active = EXCLUDED.is_active, + is_admin = EXCLUDED.is_admin, + updated_at = EXCLUDED.updated_at;`, + }, + { + legacyName: "legacy_auth_sources", + sqliteSQL: `INSERT OR REPLACE INTO w_auth_sources (id, name, type, display_name, is_active, client_id, client_secret, openid_discovery_url, scopes, icon_url, created_at, updated_at) + SELECT id, name, type, display_name, is_active, client_id, client_secret, openid_discovery_url, scopes, icon_url, created_at, updated_at + FROM legacy_auth_sources;`, + pgSQL: `INSERT INTO w_auth_sources (id, name, type, display_name, is_active, client_id, client_secret, openid_discovery_url, scopes, icon_url, created_at, updated_at) + SELECT id, name, type, display_name, is_active::boolean, client_id, client_secret, openid_discovery_url, scopes, icon_url, created_at, updated_at + FROM legacy_auth_sources + ON CONFLICT (id) DO UPDATE SET + name = EXCLUDED.name, + type = EXCLUDED.type, + display_name = EXCLUDED.display_name, + is_active = EXCLUDED.is_active, + client_id = EXCLUDED.client_id, + client_secret = EXCLUDED.client_secret, + openid_discovery_url = EXCLUDED.openid_discovery_url, + scopes = EXCLUDED.scopes, + icon_url = EXCLUDED.icon_url, + updated_at = EXCLUDED.updated_at;`, + }, + { + legacyName: "legacy_external_accounts", + sqliteSQL: `INSERT OR REPLACE INTO w_external_accounts (id, auth_source_id, user_id, external_id, external_username, email, created_at, updated_at) + SELECT id, auth_source_id, user_id, external_id, external_username, email, created_at, updated_at + FROM legacy_external_accounts;`, + pgSQL: `INSERT INTO w_external_accounts (id, auth_source_id, user_id, external_id, external_username, email, created_at, updated_at) + SELECT id, auth_source_id, user_id, external_id, external_username, email, created_at, updated_at + FROM legacy_external_accounts + ON CONFLICT (id) DO UPDATE SET + auth_source_id = EXCLUDED.auth_source_id, + user_id = EXCLUDED.user_id, + external_id = EXCLUDED.external_id, + external_username = EXCLUDED.external_username, + email = EXCLUDED.email, + updated_at = EXCLUDED.updated_at;`, + }, + { + legacyName: "legacy_options", + sqliteSQL: `INSERT OR REPLACE INTO of_options (key, value) + SELECT key, COALESCE(value, '') FROM legacy_options;`, + pgSQL: `INSERT INTO of_options (key, value) + SELECT key, COALESCE(value, '') FROM legacy_options + ON CONFLICT (key) DO UPDATE SET + value = EXCLUDED.value;`, + }, + } +} + +func getOriginAndRouteTasks() []migrationTask { + return []migrationTask{ + { + legacyName: "legacy_origins", + sqliteSQL: `INSERT OR REPLACE INTO of_origins (id, name, address, remark, created_at, updated_at) + SELECT id, name, address, COALESCE(remark, ''), created_at, updated_at + FROM legacy_origins;`, + pgSQL: `INSERT INTO of_origins (id, name, address, remark, created_at, updated_at) + SELECT id, name, address, COALESCE(remark, ''), created_at, updated_at + FROM legacy_origins + ON CONFLICT (id) DO UPDATE SET + name = EXCLUDED.name, + address = EXCLUDED.address, + remark = EXCLUDED.remark, + updated_at = EXCLUDED.updated_at;`, + }, + { + legacyName: "legacy_apply_logs", + sqliteSQL: `INSERT OR REPLACE INTO of_apply_logs (id, node_id, version, result, message, checksum, main_config_checksum, route_config_checksum, support_file_count, created_at) + SELECT id, node_id, version, result, message, COALESCE(checksum, ''), COALESCE(main_config_checksum, ''), COALESCE(route_config_checksum, ''), COALESCE(support_file_count, 0), created_at + FROM legacy_apply_logs;`, + pgSQL: `INSERT INTO of_apply_logs (id, node_id, version, result, message, checksum, main_config_checksum, route_config_checksum, support_file_count, created_at) + SELECT id, node_id, version, result, message, COALESCE(checksum, ''), COALESCE(main_config_checksum, ''), COALESCE(route_config_checksum, ''), COALESCE(support_file_count, 0), created_at + FROM legacy_apply_logs + ON CONFLICT (id) DO UPDATE SET + node_id = EXCLUDED.node_id, + version = EXCLUDED.version, + result = EXCLUDED.result, + message = EXCLUDED.message, + checksum = EXCLUDED.checksum, + main_config_checksum = EXCLUDED.main_config_checksum, + route_config_checksum = EXCLUDED.route_config_checksum, + support_file_count = EXCLUDED.support_file_count;`, + }, + { + legacyName: "legacy_proxy_routes", + sqliteSQL: `INSERT OR REPLACE INTO of_proxy_routes (id, site_name, domain, domains, origin_id, origin_url, origin_host, upstreams, enabled, enable_https, cert_id, cert_ids, domain_cert_ids, redirect_http, limit_conn_per_server, limit_conn_per_ip, limit_rate, cache_enabled, cache_policy, cache_rules, custom_headers, basic_auth_enabled, basic_auth_username, basic_auth_password, remark, upstream_type, tunnel_node_id, tunnel_target_addr, tunnel_target_protocol, pages_project_id, created_at, updated_at) + SELECT id, COALESCE(site_name, ''), domain, COALESCE(domains, '[]'), origin_id, origin_url, COALESCE(origin_host, ''), COALESCE(upstreams, '[]'), enabled, enable_https, cert_id, COALESCE(cert_ids, '[]'), COALESCE(domain_cert_ids, '[]'), redirect_http, limit_conn_per_server, limit_conn_per_ip, COALESCE(limit_rate, ''), cache_enabled, COALESCE(cache_policy, ''), COALESCE(cache_rules, '[]'), COALESCE(custom_headers, '[]'), basic_auth_enabled, COALESCE(basic_auth_username, ''), COALESCE(basic_auth_password, ''), COALESCE(remark, ''), COALESCE(upstream_type, 'direct'), tunnel_node_id, COALESCE(tunnel_target_addr, ''), COALESCE(tunnel_target_protocol, ''), pages_project_id, created_at, updated_at + FROM legacy_proxy_routes;`, + pgSQL: `INSERT INTO of_proxy_routes (id, site_name, domain, domains, origin_id, origin_url, origin_host, upstreams, enabled, enable_https, cert_id, cert_ids, domain_cert_ids, redirect_http, limit_conn_per_server, limit_conn_per_ip, limit_rate, cache_enabled, cache_policy, cache_rules, custom_headers, basic_auth_enabled, basic_auth_username, basic_auth_password, remark, upstream_type, tunnel_node_id, tunnel_target_addr, tunnel_target_protocol, pages_project_id, created_at, updated_at) + SELECT id, COALESCE(site_name, ''), domain, COALESCE(domains, '[]'), origin_id, origin_url, COALESCE(origin_host, ''), COALESCE(upstreams, '[]'), enabled, enable_https, cert_id, COALESCE(cert_ids, '[]'), COALESCE(domain_cert_ids, '[]'), redirect_http, limit_conn_per_server, limit_conn_per_ip, COALESCE(limit_rate, ''), cache_enabled, COALESCE(cache_policy, ''), COALESCE(cache_rules, '[]'), COALESCE(custom_headers, '[]'), basic_auth_enabled, COALESCE(basic_auth_username, ''), COALESCE(basic_auth_password, ''), COALESCE(remark, ''), COALESCE(upstream_type, 'direct'), tunnel_node_id, COALESCE(tunnel_target_addr, ''), COALESCE(tunnel_target_protocol, ''), pages_project_id, created_at, updated_at + FROM legacy_proxy_routes + ON CONFLICT (id) DO UPDATE SET + site_name = EXCLUDED.site_name, + domain = EXCLUDED.domain, + domains = EXCLUDED.domains, + origin_id = EXCLUDED.origin_id, + origin_url = EXCLUDED.origin_url, + origin_host = EXCLUDED.origin_host, + upstreams = EXCLUDED.upstreams, + enabled = EXCLUDED.enabled, + enable_https = EXCLUDED.enable_https, + cert_id = EXCLUDED.cert_id, + cert_ids = EXCLUDED.cert_ids, + domain_cert_ids = EXCLUDED.domain_cert_ids, + redirect_http = EXCLUDED.redirect_http, + limit_conn_per_server = EXCLUDED.limit_conn_per_server, + limit_conn_per_ip = EXCLUDED.limit_conn_per_ip, + limit_rate = EXCLUDED.limit_rate, + cache_enabled = EXCLUDED.cache_enabled, + cache_policy = EXCLUDED.cache_policy, + cache_rules = EXCLUDED.cache_rules, + custom_headers = EXCLUDED.custom_headers, + basic_auth_enabled = EXCLUDED.basic_auth_enabled, + basic_auth_username = EXCLUDED.basic_auth_username, + basic_auth_password = EXCLUDED.basic_auth_password, + remark = EXCLUDED.remark, + upstream_type = EXCLUDED.upstream_type, + tunnel_node_id = EXCLUDED.tunnel_node_id, + tunnel_target_addr = EXCLUDED.tunnel_target_addr, + tunnel_target_protocol = EXCLUDED.tunnel_target_protocol, + pages_project_id = EXCLUDED.pages_project_id, + updated_at = EXCLUDED.updated_at;`, + }, + { + legacyName: "legacy_nodes", + sqliteSQL: `INSERT OR REPLACE INTO of_nodes (id, node_id, name, ip, ip_manual_override, geo_name, geo_latitude, geo_longitude, geo_manual_override, access_token, auto_update_enabled, update_requested, update_channel, update_tag, restart_openresty_requested, version, ext_version, openresty_status, openresty_message, status, current_version, last_seen_at, last_error, created_at, updated_at, node_type, relay_bind_port, relay_vhost_http_port, relay_auth_token, relay_agent_access_addr, relay_client_access_addr, relay_client_proxy_url, capabilities_json, relay_status, relay_web_server_enabled) + SELECT id, node_id, name, COALESCE(ip, ''), COALESCE(ip_manual_override, 0), COALESCE(geo_name, ''), geo_latitude, geo_longitude, COALESCE(geo_manual_override, 0), COALESCE(access_token, ''), COALESCE(auto_update_enabled, 0), COALESCE(update_requested, 0), COALESCE(update_channel, 'stable'), COALESCE(update_tag, ''), COALESCE(restart_openresty_requested, 0), COALESCE(version, ''), COALESCE(ext_version, ''), COALESCE(openresty_status, 'unknown'), openresty_message, COALESCE(status, 'offline'), COALESCE(current_version, ''), last_seen_at, last_error, created_at, updated_at, COALESCE(node_type, 'edge_node'), COALESCE(relay_bind_port, 0), COALESCE(relay_vhost_http_port, 0), COALESCE(relay_auth_token, ''), COALESCE(relay_agent_access_addr, ''), COALESCE(relay_client_access_addr, ''), COALESCE(relay_client_proxy_url, ''), COALESCE(capabilities_json, '[]'), COALESCE(relay_status, 'unknown'), COALESCE(relay_web_server_enabled, 0) + FROM legacy_nodes;`, + pgSQL: `INSERT INTO of_nodes (id, node_id, name, ip, ip_manual_override, geo_name, geo_latitude, geo_longitude, geo_manual_override, access_token, auto_update_enabled, update_requested, update_channel, update_tag, restart_openresty_requested, version, ext_version, openresty_status, openresty_message, status, current_version, last_seen_at, last_error, created_at, updated_at, node_type, relay_bind_port, relay_vhost_http_port, relay_auth_token, relay_agent_access_addr, relay_client_access_addr, relay_client_proxy_url, capabilities_json, relay_status, relay_web_server_enabled) + SELECT id, node_id, name, COALESCE(ip, ''), COALESCE(ip_manual_override, FALSE), COALESCE(geo_name, ''), geo_latitude, geo_longitude, COALESCE(geo_manual_override, FALSE), COALESCE(access_token, ''), COALESCE(auto_update_enabled, FALSE), COALESCE(update_requested, FALSE), COALESCE(update_channel, 'stable'), COALESCE(update_tag, ''), COALESCE(restart_openresty_requested, FALSE), COALESCE(version, ''), COALESCE(ext_version, ''), COALESCE(openresty_status, 'unknown'), openresty_message, COALESCE(status, 'offline'), COALESCE(current_version, ''), last_seen_at, last_error, created_at, updated_at, COALESCE(node_type, 'edge_node'), COALESCE(relay_bind_port, 0), COALESCE(relay_vhost_http_port, 0), COALESCE(relay_auth_token, ''), COALESCE(relay_agent_access_addr, ''), COALESCE(relay_client_access_addr, ''), COALESCE(relay_client_proxy_url, ''), COALESCE(capabilities_json, '[]'), COALESCE(relay_status, 'unknown'), COALESCE(relay_web_server_enabled, FALSE) + FROM legacy_nodes + ON CONFLICT (id) DO UPDATE SET + node_id = EXCLUDED.node_id, + name = EXCLUDED.name, + ip = EXCLUDED.ip, + ip_manual_override = EXCLUDED.ip_manual_override, + geo_name = EXCLUDED.geo_name, + geo_latitude = EXCLUDED.geo_latitude, + geo_longitude = EXCLUDED.geo_longitude, + geo_manual_override = EXCLUDED.geo_manual_override, + access_token = EXCLUDED.access_token, + auto_update_enabled = EXCLUDED.auto_update_enabled, + update_requested = EXCLUDED.update_requested, + update_channel = EXCLUDED.update_channel, + update_tag = EXCLUDED.update_tag, + restart_openresty_requested = EXCLUDED.restart_openresty_requested, + version = EXCLUDED.version, + ext_version = EXCLUDED.ext_version, + openresty_status = EXCLUDED.openresty_status, + openresty_message = EXCLUDED.openresty_message, + status = EXCLUDED.status, + current_version = EXCLUDED.current_version, + last_seen_at = EXCLUDED.last_seen_at, + last_error = EXCLUDED.last_error, + updated_at = EXCLUDED.updated_at, + node_type = EXCLUDED.node_type, + relay_bind_port = EXCLUDED.relay_bind_port, + relay_vhost_http_port = EXCLUDED.relay_vhost_http_port, + relay_auth_token = EXCLUDED.relay_auth_token, + relay_agent_access_addr = EXCLUDED.relay_agent_access_addr, + relay_client_access_addr = EXCLUDED.relay_client_access_addr, + relay_client_proxy_url = EXCLUDED.relay_client_proxy_url, + capabilities_json = EXCLUDED.capabilities_json, + relay_status = EXCLUDED.relay_status, + relay_web_server_enabled = EXCLUDED.relay_web_server_enabled;`, + }, + } +} + +func getWafTasks() []migrationTask { + return []migrationTask{ + { + legacyName: "legacy_waf_rule_groups", + sqliteSQL: `INSERT OR REPLACE INTO of_waf_rule_groups (id, name, enabled, is_global, block_status_code, block_response_body, ip_whitelist, ip_blacklist, ip_whitelist_groups, ip_blacklist_groups, country_whitelist, country_blacklist, region_whitelist, region_blacklist, pow_enabled, pow_config, remark, created_at, updated_at) + SELECT id, name, COALESCE(enabled, 1), COALESCE(is_global, 0), COALESCE(block_status_code, 418), COALESCE(block_response_body, ''), COALESCE(ip_whitelist, '[]'), COALESCE(ip_blacklist, '[]'), COALESCE(ip_whitelist_groups, '[]'), COALESCE(ip_blacklist_groups, '[]'), COALESCE(country_whitelist, '[]'), COALESCE(country_blacklist, '[]'), COALESCE(region_whitelist, '[]'), COALESCE(region_blacklist, '[]'), COALESCE(pow_enabled, 0), COALESCE(pow_config, '{}'), COALESCE(remark, ''), created_at, updated_at + FROM legacy_waf_rule_groups;`, + pgSQL: `INSERT INTO of_waf_rule_groups (id, name, enabled, is_global, block_status_code, block_response_body, ip_whitelist, ip_blacklist, ip_whitelist_groups, ip_blacklist_groups, country_whitelist, country_blacklist, region_whitelist, region_blacklist, pow_enabled, pow_config, remark, created_at, updated_at) + SELECT id, name, COALESCE(enabled::boolean, TRUE), COALESCE(is_global::boolean, FALSE), COALESCE(block_status_code, 418), COALESCE(block_response_body, ''), COALESCE(ip_whitelist, '[]'), COALESCE(ip_blacklist, '[]'), COALESCE(ip_whitelist_groups, '[]'), COALESCE(ip_blacklist_groups, '[]'), COALESCE(country_whitelist, '[]'), COALESCE(country_blacklist, '[]'), COALESCE(region_whitelist, '[]'), COALESCE(region_blacklist, '[]'), COALESCE(pow_enabled::boolean, FALSE), COALESCE(pow_config, '{}'), COALESCE(remark, ''), created_at, updated_at + FROM legacy_waf_rule_groups + ON CONFLICT (id) DO UPDATE SET + name = EXCLUDED.name, + enabled = EXCLUDED.enabled, + is_global = EXCLUDED.is_global, + block_status_code = EXCLUDED.block_status_code, + block_response_body = EXCLUDED.block_response_body, + ip_whitelist = EXCLUDED.ip_whitelist, + ip_blacklist = EXCLUDED.ip_blacklist, + ip_whitelist_groups = EXCLUDED.ip_whitelist_groups, + ip_blacklist_groups = EXCLUDED.ip_blacklist_groups, + country_whitelist = EXCLUDED.country_whitelist, + country_blacklist = EXCLUDED.country_blacklist, + region_whitelist = EXCLUDED.region_whitelist, + region_blacklist = EXCLUDED.region_blacklist, + pow_enabled = EXCLUDED.pow_enabled, + pow_config = EXCLUDED.pow_config, + remark = EXCLUDED.remark, + updated_at = EXCLUDED.updated_at;`, + }, + { + legacyName: "legacy_waf_rule_group_bindings", + sqliteSQL: `INSERT OR REPLACE INTO of_waf_rule_group_bindings (id, rule_group_id, proxy_route_id, created_at) + SELECT id, rule_group_id, proxy_route_id, created_at + FROM legacy_waf_rule_group_bindings;`, + pgSQL: `INSERT INTO of_waf_rule_group_bindings (id, rule_group_id, proxy_route_id, created_at) + SELECT id, rule_group_id, proxy_route_id, created_at + FROM legacy_waf_rule_group_bindings + ON CONFLICT (id) DO UPDATE SET + rule_group_id = EXCLUDED.rule_group_id, + proxy_route_id = EXCLUDED.proxy_route_id;`, + }, + { + legacyName: "legacy_waf_ip_groups", + sqliteSQL: `INSERT OR REPLACE INTO of_waf_ip_groups (id, name, type, enabled, ip_list, auto_config, ext_ips, subscription_url, subscription_format, subscription_mapping_rule, sync_interval_minutes, last_synced_at, next_sync_at, last_sync_status, last_sync_message, remark, created_at, updated_at) + SELECT id, name, type, COALESCE(enabled, 1), COALESCE(ip_list, '[]'), COALESCE(auto_config, '{}'), COALESCE(ext_ips, '[]'), COALESCE(subscription_url, ''), COALESCE(subscription_format, 'text'), COALESCE(subscription_mapping_rule, ''), COALESCE(sync_interval_minutes, 1440), last_synced_at, next_sync_at, COALESCE(last_sync_status, ''), COALESCE(last_sync_message, ''), COALESCE(remark, ''), created_at, updated_at + FROM legacy_waf_ip_groups;`, + pgSQL: `INSERT INTO of_waf_ip_groups (id, name, type, enabled, ip_list, auto_config, ext_ips, subscription_url, subscription_format, subscription_mapping_rule, sync_interval_minutes, last_synced_at, next_sync_at, last_sync_status, last_sync_message, remark, created_at, updated_at) + SELECT id, name, type, COALESCE(enabled::boolean, TRUE), COALESCE(ip_list, '[]'), COALESCE(auto_config, '{}'), COALESCE(ext_ips, '[]'), COALESCE(subscription_url, ''), COALESCE(subscription_format, 'text'), COALESCE(subscription_mapping_rule, ''), COALESCE(sync_interval_minutes, 1440), last_synced_at, next_sync_at, COALESCE(last_sync_status, ''), COALESCE(last_sync_message, ''), COALESCE(remark, ''), created_at, updated_at + FROM legacy_waf_ip_groups + ON CONFLICT (id) DO UPDATE SET + name = EXCLUDED.name, + type = EXCLUDED.type, + enabled = EXCLUDED.enabled, + ip_list = EXCLUDED.ip_list, + auto_config = EXCLUDED.auto_config, + ext_ips = EXCLUDED.ext_ips, + subscription_url = EXCLUDED.subscription_url, + subscription_format = EXCLUDED.subscription_format, + subscription_mapping_rule = EXCLUDED.subscription_mapping_rule, + sync_interval_minutes = EXCLUDED.sync_interval_minutes, + last_synced_at = EXCLUDED.last_synced_at, + next_sync_at = EXCLUDED.next_sync_at, + last_sync_status = EXCLUDED.last_sync_status, + last_sync_message = EXCLUDED.last_sync_message, + remark = EXCLUDED.remark, + updated_at = EXCLUDED.updated_at;`, + }, + } +} + +func getTLSTasks() []migrationTask { + return []migrationTask{ + { + legacyName: "legacy_tls_certificates", + sqliteSQL: `INSERT OR REPLACE INTO of_tls_certificates (id, name, cert_pem, key_pem, not_before, not_after, remark, provider, acme_account_id, dns_account_id, key_algorithm, auto_renew, primary_domain, other_domains, disable_cname, skip_dns, dns1, dns2, apply_status, apply_message, created_at, updated_at) + SELECT id, name, cert_pem, key_pem, not_before, not_after, COALESCE(remark, ''), COALESCE(provider, 'upload'), COALESCE(acme_account_id, 0), COALESCE(dns_account_id, 0), COALESCE(key_algorithm, ''), COALESCE(auto_renew, 0), COALESCE(primary_domain, ''), COALESCE(other_domains, ''), COALESCE(disable_cname, 0), COALESCE(skip_dns, 0), COALESCE(dns1, ''), COALESCE(dns2, ''), COALESCE(apply_status, 'ready'), COALESCE(apply_message, ''), created_at, updated_at + FROM legacy_tls_certificates;`, + pgSQL: `INSERT INTO of_tls_certificates (id, name, cert_pem, key_pem, not_before, not_after, remark, provider, acme_account_id, dns_account_id, key_algorithm, auto_renew, primary_domain, other_domains, disable_cname, skip_dns, dns1, dns2, apply_status, apply_message, created_at, updated_at) + SELECT id, name, cert_pem, key_pem, not_before, not_after, COALESCE(remark, ''), COALESCE(provider, 'upload'), COALESCE(acme_account_id, 0), COALESCE(dns_account_id, 0), COALESCE(key_algorithm, ''), COALESCE(auto_renew, FALSE), COALESCE(primary_domain, ''), COALESCE(other_domains, ''), COALESCE(disable_cname, FALSE), COALESCE(skip_dns, FALSE), COALESCE(dns1, ''), COALESCE(dns2, ''), COALESCE(apply_status, 'ready'), COALESCE(apply_message, ''), created_at, updated_at + FROM legacy_tls_certificates + ON CONFLICT (id) DO UPDATE SET + name = EXCLUDED.name, + cert_pem = EXCLUDED.cert_pem, + key_pem = EXCLUDED.key_pem, + not_before = EXCLUDED.not_before, + not_after = EXCLUDED.not_after, + remark = EXCLUDED.remark, + provider = EXCLUDED.provider, + acme_account_id = EXCLUDED.acme_account_id, + dns_account_id = EXCLUDED.dns_account_id, + key_algorithm = EXCLUDED.key_algorithm, + auto_renew = EXCLUDED.auto_renew, + primary_domain = EXCLUDED.primary_domain, + other_domains = EXCLUDED.other_domains, + disable_cname = EXCLUDED.disable_cname, + skip_dns = EXCLUDED.skip_dns, + dns1 = EXCLUDED.dns1, + dns2 = EXCLUDED.dns2, + apply_status = EXCLUDED.apply_status, + apply_message = EXCLUDED.apply_message, + updated_at = EXCLUDED.updated_at;`, + }, + { + legacyName: "legacy_managed_domains", + sqliteSQL: `INSERT OR REPLACE INTO of_managed_domains (id, domain, cert_id, enabled, remark, created_at, updated_at) + SELECT id, domain, cert_id, COALESCE(enabled, 1), COALESCE(remark, ''), created_at, updated_at + FROM legacy_managed_domains;`, + pgSQL: `INSERT INTO of_managed_domains (id, domain, cert_id, enabled, remark, created_at, updated_at) + SELECT id, domain, cert_id, COALESCE(enabled, TRUE), COALESCE(remark, ''), created_at, updated_at + FROM legacy_managed_domains + ON CONFLICT (id) DO UPDATE SET + domain = EXCLUDED.domain, + cert_id = EXCLUDED.cert_id, + enabled = EXCLUDED.enabled, + remark = EXCLUDED.remark, + updated_at = EXCLUDED.updated_at;`, + }, + { + legacyName: "legacy_dns_accounts", + sqliteSQL: `INSERT OR REPLACE INTO of_dns_accounts (id, name, type, authorization, created_at, updated_at) + SELECT id, name, type, authorization, created_at, updated_at + FROM legacy_dns_accounts;`, + pgSQL: `INSERT INTO of_dns_accounts (id, name, type, "authorization", created_at, updated_at) + SELECT id, name, type, "authorization", created_at, updated_at + FROM legacy_dns_accounts + ON CONFLICT (id) DO UPDATE SET + name = EXCLUDED.name, + type = EXCLUDED.type, + "authorization" = EXCLUDED."authorization", + updated_at = EXCLUDED.updated_at;`, + }, + { + legacyName: "legacy_acme_accounts", + sqliteSQL: `INSERT OR REPLACE INTO of_acme_accounts (id, email, url, private_key, created_at, updated_at) + SELECT id, COALESCE(email, ''), COALESCE(url, ''), COALESCE(private_key, ''), created_at, updated_at + FROM legacy_acme_accounts;`, + pgSQL: `INSERT INTO of_acme_accounts (id, email, url, private_key, created_at, updated_at) + SELECT id, COALESCE(email, ''), COALESCE(url, ''), COALESCE(private_key, ''), created_at, updated_at + FROM legacy_acme_accounts + ON CONFLICT (id) DO UPDATE SET + email = EXCLUDED.email, + url = EXCLUDED.url, + private_key = EXCLUDED.private_key, + updated_at = EXCLUDED.updated_at;`, + }, + } +} + +func getPagesAndConfigTasks() []migrationTask { + return []migrationTask{ + { + legacyName: "legacy_config_versions", + sqliteSQL: `INSERT OR REPLACE INTO of_config_versions (id, version, snapshot_json, main_config, rendered_config, support_files_json, checksum, is_active, created_by, created_at) + SELECT id, version, snapshot_json, COALESCE(main_config, ''), rendered_config, COALESCE(support_files_json, '[]'), checksum, COALESCE(is_active, 0), created_by, created_at + FROM legacy_config_versions;`, + pgSQL: `INSERT INTO of_config_versions (id, version, snapshot_json, main_config, rendered_config, support_files_json, checksum, is_active, created_by, created_at) + SELECT id, version, snapshot_json, COALESCE(main_config, ''), rendered_config, COALESCE(support_files_json, '[]'), checksum, COALESCE(is_active, FALSE), created_by, created_at + FROM legacy_config_versions + ON CONFLICT (id) DO UPDATE SET + version = EXCLUDED.version, + snapshot_json = EXCLUDED.snapshot_json, + main_config = EXCLUDED.main_config, + rendered_config = EXCLUDED.rendered_config, + support_files_json = EXCLUDED.support_files_json, + checksum = EXCLUDED.checksum, + is_active = EXCLUDED.is_active, + created_by = EXCLUDED.created_by;`, + }, + { + legacyName: "legacy_pages_projects", + sqliteSQL: `INSERT OR REPLACE INTO of_pages_projects (id, name, slug, description, enabled, spa_fallback_enabled, spa_fallback_path, api_proxy_enabled, api_proxy_path, api_proxy_pass, api_proxy_rewrite, active_deployment_id, root_dir, entry_file, created_at, updated_at) + SELECT id, name, slug, COALESCE(description, ''), COALESCE(enabled, 1), COALESCE(spa_fallback_enabled, 0), COALESCE(spa_fallback_path, '/index.html'), COALESCE(api_proxy_enabled, 0), COALESCE(api_proxy_path, ''), COALESCE(api_proxy_pass, ''), COALESCE(api_proxy_rewrite, ''), active_deployment_id, COALESCE(root_dir, ''), COALESCE(entry_file, 'index.html'), created_at, updated_at + FROM legacy_pages_projects;`, + pgSQL: `INSERT INTO of_pages_projects (id, name, slug, description, enabled, spa_fallback_enabled, spa_fallback_path, api_proxy_enabled, api_proxy_path, api_proxy_pass, api_proxy_rewrite, active_deployment_id, root_dir, entry_file, created_at, updated_at) + SELECT id, name, slug, COALESCE(description, ''), COALESCE(enabled, TRUE), COALESCE(spa_fallback_enabled, FALSE), COALESCE(spa_fallback_path, '/index.html'), COALESCE(api_proxy_enabled, FALSE), COALESCE(api_proxy_path, ''), COALESCE(api_proxy_pass, ''), COALESCE(api_proxy_rewrite, ''), active_deployment_id, COALESCE(root_dir, ''), COALESCE(entry_file, 'index.html'), created_at, updated_at + FROM legacy_pages_projects + ON CONFLICT (id) DO UPDATE SET + name = EXCLUDED.name, + slug = EXCLUDED.slug, + description = EXCLUDED.description, + enabled = EXCLUDED.enabled, + spa_fallback_enabled = EXCLUDED.spa_fallback_enabled, + spa_fallback_path = EXCLUDED.spa_fallback_path, + api_proxy_enabled = EXCLUDED.api_proxy_enabled, + api_proxy_path = EXCLUDED.api_proxy_path, + api_proxy_pass = EXCLUDED.api_proxy_pass, + api_proxy_rewrite = EXCLUDED.api_proxy_rewrite, + active_deployment_id = EXCLUDED.active_deployment_id, + root_dir = EXCLUDED.root_dir, + entry_file = EXCLUDED.entry_file, + updated_at = EXCLUDED.updated_at;`, + }, + { + legacyName: "legacy_pages_deployments", + sqliteSQL: `INSERT OR REPLACE INTO of_pages_deployments (id, project_id, deployment_number, checksum, status, artifact_path, file_count, total_size, created_by, created_at, activated_at) + SELECT id, project_id, deployment_number, checksum, COALESCE(status, 'uploaded'), artifact_path, COALESCE(file_count, 0), COALESCE(total_size, 0), COALESCE(created_by, ''), created_at, activated_at + FROM legacy_pages_deployments;`, + pgSQL: `INSERT INTO of_pages_deployments (id, project_id, deployment_number, checksum, status, artifact_path, file_count, total_size, created_by, created_at, activated_at) + SELECT id, project_id, deployment_number, checksum, COALESCE(status, 'uploaded'), artifact_path, COALESCE(file_count, 0), COALESCE(total_size, 0), COALESCE(created_by, ''), created_at, activated_at + FROM legacy_pages_deployments + ON CONFLICT (id) DO UPDATE SET + project_id = EXCLUDED.project_id, + deployment_number = EXCLUDED.deployment_number, + checksum = EXCLUDED.checksum, + status = EXCLUDED.status, + artifact_path = EXCLUDED.artifact_path, + file_count = EXCLUDED.file_count, + total_size = EXCLUDED.total_size, + created_by = EXCLUDED.created_by, + activated_at = EXCLUDED.activated_at;`, + }, + { + legacyName: "legacy_pages_deployment_files", + sqliteSQL: `INSERT OR REPLACE INTO of_pages_deployment_files (id, deployment_id, path, size, checksum, created_at) + SELECT id, deployment_id, path, COALESCE(size, 0), checksum, created_at + FROM legacy_pages_deployment_files;`, + pgSQL: `INSERT INTO of_pages_deployment_files (id, deployment_id, path, size, checksum, created_at) + SELECT id, deployment_id, path, COALESCE(size, 0), checksum, created_at + FROM legacy_pages_deployment_files + ON CONFLICT (id) DO UPDATE SET + deployment_id = EXCLUDED.deployment_id, + path = EXCLUDED.path, + size = EXCLUDED.size, + checksum = EXCLUDED.checksum;`, + }, + } +} + +func down202606200006(_ context.Context, _ *sql.Tx) error { + // Down migration is a no-op + return nil +} diff --git a/internal/db/migrator/bridge_test.go b/internal/db/migrator/bridge_test.go new file mode 100644 index 00000000..f30880f4 --- /dev/null +++ b/internal/db/migrator/bridge_test.go @@ -0,0 +1,270 @@ +// Copyright 2026 Arctel.net +// SPDX-License-Identifier: Apache-2.0 + +package migrator + +import ( + "testing" + "time" + + "github.com/Rain-kl/Wavelet/internal/config" + "github.com/Rain-kl/Wavelet/internal/db" + "github.com/alicebob/miniredis/v2" + "github.com/glebarez/sqlite" + "github.com/redis/go-redis/v9" + "github.com/redis/go-redis/v9/maintnotifications" + "gorm.io/gorm" +) + +func TestMigrateUpgradesLegacyDatabaseAndPreservesData(t *testing.T) { + // 1. Initialize an empty in-memory SQLite DB + sqliteDB, err := gorm.Open(sqlite.Open(":memory:"), &gorm.Config{ + DisableForeignKeyConstraintWhenMigrating: true, + }) + if err != nil { + t.Fatalf("gorm.Open(sqlite) error = %v", err) + } + + rawDB, err := sqliteDB.DB() + if err != nil { + t.Fatalf("sqliteDB.DB() error = %v", err) + } + + // 2. Set up the legacy schema manually + setupLegacySchema := []string{ + `CREATE TABLE users ( + id integer(64) NOT NULL, + username text, + password text NOT NULL, + display_name text, + role integer(64), + status integer(64), + token text, + email text, + github_id text, + wechat_id text, + CONSTRAINT users_pkey PRIMARY KEY (id) + );`, + `CREATE TABLE options ( + key text NOT NULL, + value text, + CONSTRAINT options_pkey PRIMARY KEY (key) + );`, + `CREATE TABLE origins ( + id integer(64) NOT NULL, + name text(255) NOT NULL, + address text(255) NOT NULL, + remark text(255), + created_at text(6), + updated_at text(6), + CONSTRAINT origins_pkey PRIMARY KEY (id) + );`, + `CREATE TABLE apply_logs ( + id integer(64) NOT NULL, + node_id text(64) NOT NULL, + version text(32) NOT NULL, + result text(32) NOT NULL, + message text, + checksum text(64) NOT NULL, + main_config_checksum text(64) NOT NULL, + route_config_checksum text(64) NOT NULL, + support_file_count integer(64) NOT NULL, + created_at text(6), + CONSTRAINT apply_logs_pkey PRIMARY KEY (id) + );`, + // Legacy log tables (which should be dropped) + `CREATE TABLE node_access_logs_00 ( + id integer(64) NOT NULL, + node_id text(64) NOT NULL, + logged_at text(6) NOT NULL + );`, + `CREATE TABLE node_system_profiles ( + id integer(64) NOT NULL, + node_id text(64) NOT NULL, + hostname text(255) + );`, + } + + for _, query := range setupLegacySchema { + if _, err := rawDB.Exec(query); err != nil { + t.Fatalf("Exec setup query failed: %v", err) + } + } + + // 3. Populate with legacy mock data + insertMockData := []string{ + // Legacy admin user + `INSERT INTO users (id, username, password, display_name, role, status, email) + VALUES (1, 'ryan', 'hashed_pass_123', 'Root User ryan', 100, 1, 'ryan@example.com');`, + // Legacy normal user + `INSERT INTO users (id, username, password, display_name, role, status, email) + VALUES (2, 'jack', 'hashed_pass_456', 'Normal User jack', 10, 1, 'jack@example.com');`, + // Legacy config options + `INSERT INTO options (key, value) VALUES ('GitHubClientId', 'Ov23lixMKW');`, + `INSERT INTO options (key, value) VALUES ('GitHubClientSecret', 'secret_abc_123');`, + // Legacy origin + `INSERT INTO origins (id, name, address, remark, created_at, updated_at) + VALUES (10, 'my_origin', '127.0.0.1:8080', 'original origin', '2026-06-01 12:00:00', '2026-06-01 12:00:00');`, + // Legacy apply log + `INSERT INTO apply_logs (id, node_id, version, result, message, checksum, main_config_checksum, route_config_checksum, support_file_count, created_at) + VALUES (100, 'node_1', 'v1.0.0', 'success', 'applied successfully', 'hash1', 'hash2', 'hash3', 2, '2026-06-01 12:00:00');`, + } + + for _, query := range insertMockData { + if _, err := rawDB.Exec(query); err != nil { + t.Fatalf("Exec mock data query failed: %v", err) + } + } + + // 4. Set up mock services (Redis & Config) + mr, err := miniredis.Run() + if err != nil { + t.Fatalf("miniredis.Run() error = %v", err) + } + redisClient := redis.NewClient(&redis.Options{ + Addr: mr.Addr(), + MaintNotificationsConfig: &maintnotifications.Config{ + Mode: maintnotifications.ModeDisabled, + }, + }) + + previousDBEnabled := config.Config.Database.Enabled + previousRedis := db.Redis + config.Config.Database.Enabled = false + db.SetDB(sqliteDB) + db.Redis = redisClient + + t.Cleanup(func() { + config.Config.Database.Enabled = previousDBEnabled + db.SetDB(nil) + db.Redis = previousRedis + _ = redisClient.Close() + mr.Close() + }) + + // 5. Run the Goose migrations + Migrate() + + // 6. Assertions on migrated data schema & records + + // Check that the w_users table exists and contains correct columns/records + var usersCount int64 + if err := sqliteDB.Table("w_users").Count(&usersCount).Error; err != nil { + t.Fatalf("count w_users error: %v", err) + } + if usersCount < 2 { + t.Errorf("expected at least 2 users, got %d", usersCount) + } + + // Verify user 1 (ryan) details + type userStruct struct { + ID int64 + Username string + Nickname string + Email string + IsActive bool + IsAdmin bool + } + var ryanUser userStruct + if err := sqliteDB.Table("w_users").Where("id = ?", 1).First(&ryanUser).Error; err != nil { + t.Fatalf("query ryan user failed: %v", err) + } + if ryanUser.Username != "ryan" || ryanUser.Nickname != "Root User ryan" || !ryanUser.IsActive || !ryanUser.IsAdmin { + t.Errorf("ryan user migrated incorrectly: %+v", ryanUser) + } + + // Verify user 2 (jack) details + var jackUser userStruct + if err := sqliteDB.Table("w_users").Where("id = ?", 2).First(&jackUser).Error; err != nil { + t.Fatalf("query jack user failed: %v", err) + } + if jackUser.Username != "jack" || jackUser.Nickname != "Normal User jack" || !jackUser.IsActive || jackUser.IsAdmin { + t.Errorf("jack user migrated incorrectly: %+v", jackUser) + } + + // Verify options migrated to of_options + var githubClientIDVal string + if err := sqliteDB.Table("of_options").Where("key = ?", "GitHubClientId").Select("value").Scan(&githubClientIDVal).Error; err != nil { + t.Fatalf("query GitHubClientId failed: %v", err) + } + if githubClientIDVal != "Ov23lixMKW" { + t.Errorf("GitHubClientId value incorrect: got %s, want Ov23lixMKW", githubClientIDVal) + } + + var githubClientSecretVal string + if err := sqliteDB.Table("of_options").Where("key = ?", "GitHubClientSecret").Select("value").Scan(&githubClientSecretVal).Error; err != nil { + t.Fatalf("query GitHubClientSecret failed: %v", err) + } + if githubClientSecretVal != "secret_abc_123" { + t.Errorf("GitHubClientSecret value incorrect: got %s, want secret_abc_123", githubClientSecretVal) + } + + // Verify origin migrated to of_origins + type originStruct struct { + ID int64 + Name string + Address string + Remark string + } + var testOrigin originStruct + if err := sqliteDB.Table("of_origins").Where("id = ?", 10).First(&testOrigin).Error; err != nil { + t.Fatalf("query my_origin failed: %v", err) + } + if testOrigin.Name != "my_origin" || testOrigin.Address != "127.0.0.1:8080" || testOrigin.Remark != "original origin" { + t.Errorf("origin migrated incorrectly: %+v", testOrigin) + } + + // Verify apply log migrated to of_apply_logs + type applyLogStruct struct { + ID int64 + NodeID string + Version string + Result string + Message string + CreatedAt time.Time + } + var testApplyLog applyLogStruct + if err := sqliteDB.Table("of_apply_logs").Where("id = ?", 100).First(&testApplyLog).Error; err != nil { + t.Fatalf("query apply_logs failed: %v", err) + } + if testApplyLog.NodeID != "node_1" || testApplyLog.Version != "v1.0.0" || testApplyLog.Result != "success" || testApplyLog.Message != "applied successfully" { + t.Errorf("apply log migrated incorrectly: %+v", testApplyLog) + } + + // 7. Verify that all legacy_ prefix tables have been dropped + legacyTablesToCheck := []string{ + "legacy_users", + "legacy_options", + "legacy_origins", + "legacy_apply_logs", + "legacy_proxy_routes", + "legacy_nodes", + "legacy_waf_rule_groups", + "legacy_waf_rule_group_bindings", + "legacy_waf_ip_groups", + "legacy_tls_certificates", + "legacy_managed_domains", + "legacy_dns_accounts", + "legacy_acme_accounts", + "legacy_config_versions", + "legacy_pages_projects", + "legacy_pages_deployments", + "legacy_pages_deployment_files", + "users", + "options", + "origins", + "apply_logs", + "node_access_logs_00", + "node_system_profiles", + } + + for _, table := range legacyTablesToCheck { + var count int + if err := sqliteDB.Raw("SELECT COUNT(*) FROM sqlite_master WHERE type='table' AND name=?", table).Scan(&count).Error; err != nil { + t.Fatalf("check table %s existence failed: %v", table, err) + } + if count > 0 { + t.Errorf("expected table %s to be dropped, but it still exists", table) + } + } +} diff --git a/internal/db/migrator/clickhouse.go b/internal/db/migrator/clickhouse.go index a48a1a22..d301ba50 100644 --- a/internal/db/migrator/clickhouse.go +++ b/internal/db/migrator/clickhouse.go @@ -5,8 +5,10 @@ package migrator import ( + "context" "database/sql" "embed" + "io/fs" "log" "time" @@ -55,13 +57,25 @@ func MigrateClickHouse() { BlockBufferSize: cfg.BlockBufferSize, }) - goose.SetBaseFS(clickhouseMigrationFS) - if err := goose.SetDialect("clickhouse"); err != nil { + subFS, err := fs.Sub(clickhouseMigrationFS, "goose/clickhouse") + if err != nil { closeClickHouseDB(sqlDB) - log.Fatalf("[ClickHouse] set goose dialect failed: %v\n", err) + log.Fatalf("[ClickHouse] get sub fs failed: %v\n", err) } - goose.SetTableName(clickhouseGooseVersionTable) - if err := goose.Up(sqlDB, clickhouseMigrationDir); err != nil { + + provider, err := goose.NewProvider( + "clickhouse", + sqlDB, + subFS, + goose.WithTableName(clickhouseGooseVersionTable), + goose.WithDisableGlobalRegistry(true), + ) + if err != nil { + closeClickHouseDB(sqlDB) + log.Fatalf("[ClickHouse] create goose provider failed: %v\n", err) + } + + if _, err := provider.Up(context.Background()); err != nil { closeClickHouseDB(sqlDB) log.Fatalf("[ClickHouse] goose migrate failed: %v\n", err) } @@ -74,4 +88,4 @@ func closeClickHouseDB(sqlDB *sql.DB) { if err := sqlDB.Close(); err != nil { log.Printf("[ClickHouse] close sql db failed: %v\n", err) } -} \ No newline at end of file +} diff --git a/internal/db/migrator/migrator.go b/internal/db/migrator/migrator.go index a28d3b4f..76315080 100644 --- a/internal/db/migrator/migrator.go +++ b/internal/db/migrator/migrator.go @@ -29,11 +29,31 @@ func dbType() string { return "PostgreSQL" } +const ( + dialectSqlite = "sqlite3" + dialectPostgres = "postgres" + cascadeSuffix = " CASCADE" +) + func gooseDialect() string { if !config.Config.Database.Enabled { - return "sqlite3" + return dialectSqlite } - return "postgres" + return dialectPostgres +} + +func tableExistsSQL(dialect string) string { + if dialect == dialectSqlite { + return "SELECT COUNT(*) FROM sqlite_master WHERE type='table' AND name=?" + } + return "SELECT COUNT(*) FROM information_schema.tables WHERE table_schema = 'public' AND table_name = $1" +} + +func tablesWithPrefixSQL(dialect string) string { + if dialect == dialectSqlite { + return "SELECT name FROM sqlite_master WHERE type='table' AND name LIKE ?" + } + return "SELECT table_name FROM information_schema.tables WHERE table_schema='public' AND table_name LIKE $1" } func migrationDir() string {