diff --git a/docs/changelog/index.md b/docs/changelog/index.md index 25cc6352..c69cfe29 100644 --- a/docs/changelog/index.md +++ b/docs/changelog/index.md @@ -26,7 +26,7 @@ sidebar: false - 新增第一阶段 Zone 与正规化 Zone 域名数据库表及路由绑定模型,为后续以稳定 ID 管理网站与域名关联提供基础。 - 新增 Zone 管理 API 与显式历史域名导入命令,使用公共后缀列表验证注册根域和域名归属。 - 管理端网站入口改为 Zone 列表与 `/websites/:zoneId` 详情(概览 / 域名 / 路由 / 证书 / 设置),反代路由通过 Zone 域名选择器绑定。 -- 新增 Zone 域名迁移指南(`docs/guide/zone-domain-migration.md`),覆盖备份、`wavelet migrate-zones`、预览等价性与回滚步骤。 +- 新增 Zone 域名迁移指南(`docs/guide/zone-domain-migration.md`);历史域名在 Server 启动时由 goose 自动导入,无需单独命令。 ### 修改 diff --git a/docs/design/zone-design.md b/docs/design/zone-design.md index e7e2ed19..3635b669 100644 --- a/docs/design/zone-design.md +++ b/docs/design/zone-design.md @@ -84,15 +84,15 @@ WAF、Pages、上游与发布版本仍属于 `proxy_routes`。Zone 概览只聚 本次改造分两个发布阶段,以免 SQL 用错误的“末两段域名”规则处理多级公共后缀。操作细则见 [Zone 域名迁移与发布验收](../guide/zone-domain-migration.md)。 1. **第一阶段 DDL**:PostgreSQL 与 SQLite 同版本 Goose 创建 `of_zones` / `of_zone_domains`;暂时保留 `of_managed_domains` 与路由冗余列。 -2. **数据导入**:`wavelet migrate-zones` 使用 `publicsuffix.EffectiveTLDPlusOne`,从既有路由域名及(无路由域名时的)`managed_domains` 创建 Zone / Zone 域名,按旧 `domain_cert_ids` 位置写入 `zone_domains.cert_id`。冲突时整单回滚并输出报告,禁止在有冲突时发布。 +2. **数据导入(自动)**:Server 启动时 `migrator.Migrate()` 先应用 goose SQL 至 `202607120002`,再自动导入旧路由域名 / `managed_domains`(`publicsuffix` 解析注册根域,写入 `cert_id` 与 `proxy_route_id`),最后继续后续 SQL。冲突时启动失败;修复后重启可幂等重试。无需手动命令。 3. **代码切换**:控制面 API、配置快照、渲染、前端均以 Zone 域名为唯一来源;路由写入仅使用 `zone_domain_ids`。 -4. **第二阶段清理**:独立 Goose 迁移 `202607130001_drop_legacy_route_domain_columns` 删除 `of_managed_domains` 与 `of_proxy_routes` 的 `domain` / `domains` / `cert_id` / `cert_ids` / `domain_cert_ids`。Down 仅恢复开发库空结构,不回填历史数据。 +4. **第二阶段清理**:Goose SQL `202607130001_drop_legacy_route_domain_columns` 删除 `of_managed_domains` 与 `of_proxy_routes` 冗余列。Down 仅恢复开发库空结构,不回填历史数据。 ### 运行时模型边界 * 持久化:域名与证书只存在于 `of_zone_domains`;`of_proxy_routes` 仅保存路由策略(上游、缓存、限流、WAF 绑定键等)。 * 渲染:配置快照在内存中组装临时 `Domains` / `DomainCertIDs` 供 OpenResty 渲染,不写回数据库。 -* 导入命令:第二阶段后若旧列/旧表已不存在则跳过对应源,对已导入 Zone 域名保持幂等。 +* 结构迁移仅使用 `internal/db/migrator/goose/{postgres,sqlite}/*.sql`;启动时自动导入历史域名,第二阶段后旧列不存在则为空操作。 ## 验证 diff --git a/docs/guide/index.md b/docs/guide/index.md index bd60dc4e..580a8f6b 100644 --- a/docs/guide/index.md +++ b/docs/guide/index.md @@ -11,7 +11,7 @@ OpenFlare 是一套自托管的 OpenResty 控制面。它把反向代理网站 1. [快速开始](./quick-start.md):用 Docker Compose 启动 Server,登录管理端,并接入第一个 Agent。 2. [发布第一份配置](./first-site.md):快速新建一条最基础的 HTTP 反代站点规则,并验证节点生效状态。 3. [新建反代配置](./proxy-config.md):一步一步了解如何从证书导入与申请开始,配置 HTTPS 加密与上游源站管理。 -4. [Zone 域名迁移](./zone-domain-migration.md):从旧托管域名/路由内嵌域名升级到 Zone 模型,含备份、导入、预览与回滚。 +4. [Zone 域名迁移](./zone-domain-migration.md):从旧托管域名/路由内嵌域名升级到 Zone 模型(goose 自动导入),含备份、验收与回滚说明。 5. [Pages 静态托管使用](./pages-usage.md):了解静态项目 ZIP 上传限制、SPA Fallback、以及内置 API 反向代理配置。 6. [内网穿透与隧道使用](./tunnel-usage.md):部署 Relay 与 Client,实现安全、无公网 IP 反向穿透。 7. [WAF 安全防护使用](./waf-usage.md):配置 WAF 规则组,掌握 IP 黑白名单、自动/订阅 IP 组、地域限制与 PoW CC 防护。 diff --git a/docs/guide/zone-domain-migration.md b/docs/guide/zone-domain-migration.md index 5d3c31de..311ce243 100644 --- a/docs/guide/zone-domain-migration.md +++ b/docs/guide/zone-domain-migration.md @@ -1,95 +1,49 @@ # Zone 域名迁移与发布验收 -从旧版 `managed_domains` / 反代路由内嵌域名列迁移到 Zone + Zone 域名模型时,按本指南操作。**导入报告存在冲突时禁止继续发布。** +从旧版 `managed_domains` / 反代路由内嵌域名列迁移到 Zone + Zone 域名模型时,数据导入与表结构升级均由 **Server 启动时的 goose 自动迁移**完成,无需单独执行导入命令。 -## 前置条件 +## 升级时发生了什么 -* 已备份 PostgreSQL / SQLite 数据库与当前激活配置版本(可导出管理端「配置版本」中的激活快照)。 -* Server 二进制已升级到包含 `of_zones` / `of_zone_domains` 表迁移的版本。 -* 维护窗口内可暂停非必要配置发布。 +启动(或滚动升级)包含 Zone 改造的 Server 版本时,**无需手动命令**,`migrator.Migrate()` 自动: -## 1. 备份 +1. 应用 goose SQL:创建 `of_zones` / `of_zone_domains`(若尚未存在)。 +2. **自动导入**旧路由域名列(及无路由域名时的 `of_managed_domains`)为 Zone / Zone 域名,并绑定 `proxy_route_id` / `cert_id`(公共后缀列表解析注册根域)。 +3. 继续 goose SQL:删除 `of_managed_domains` 与 `of_proxy_routes` 冗余域名/证书列。 + +导入幂等:已存在的域名会跳过或补绑路由。 + +**若历史数据无法解析(冲突域名、无效根域、证书不存在等),启动失败。** 修复数据或恢复备份后再次启动即可重试。 + +## 建议操作 + +### 1. 升级前备份 ```bash # PostgreSQL 示例 pg_dump "$DATABASE_URL" > openflare-pre-zone-$(date +%Y%m%d).sql -# 或复制备份卷 / 快照;SQLite 则直接复制 data 目录中的库文件 +# 或复制备份卷 / 快照;SQLite 则复制 data 目录中的库文件 ``` -在管理端确认当前**激活版本号**并记下 checksum,便于回滚对比。 +可选:在管理端记下当前**激活配置版本号**与 checksum,便于配置回滚对比。 -## 2. 执行历史导入 +### 2. 升级并启动 Server -```bash -# 二进制名称为 wavelet(或你的部署包中的同名入口) -wavelet migrate-zones -``` +部署新版本并启动即可。观察启动日志中的 goose 成功信息;若出现「迁移 Zone 失败(N 个冲突)」则按日志中的冲突项修复源数据后重启。 -命令会: +### 3. 升级后检查 -1. 先跑 goose 迁移(确保 Zone 表存在)。 -2. 以事务从旧路由域名(及无路由域名时的 `of_managed_domains`)导入 Zone / Zone 域名。 -3. 使用公共后缀列表解析注册根域;冲突时整单回滚并输出报告。 +1. 管理端 **网站** `/websites`:Zone 根域与域名计数是否合理。 +2. Zone 详情:域名、证书、关联路由 ID。 +3. **反代路由**:域名绑定来自 Zone 域名,而非旧手写字段。 -**成功标志:** 进程退出码 0,且日志/标准输出无「冲突」列表。 +### 4. 配置预览与发布 -**失败时:** 阅读冲突项(无法解析的根域、通配符 FQDN、全局域名冲突、证书不存在等),修复源数据后重新执行。`migrate-zones` 幂等:已导入的域名会跳过,不会重复创建。 - -**有冲突时不要发布配置、不要执行第二阶段删列迁移。** - -## 3. 导入后检查 - -1. 打开管理端 **网站** `/websites`:确认 Zone 根域与域名计数合理。 -2. 进入各 Zone 详情:域名、证书绑定、关联路由 ID 是否正确。 -3. 打开 **反代路由**:域名区应展示 Zone 域名绑定,而不是手写域名。 - -## 4. 配置预览与快照等价性 - -在升级前若已导出激活快照,导入后: - -1. 在管理端打开配置差异 / 预览(或调用配置 diff / preview API)。 -2. **逐路由**核对: - * 明确 `server_name` 集合(全部 FQDN) - * 证书支持文件路径 / 证书 ID 与域名对应关系 - * WAF 绑定的 Route ID(`site_name` 与路由 ID 不变) - * Pages 项目引用 -3. **允许**旧快照 JSON 中路由上的冗余 `domain` / `domains` / `cert_ids` 字段消失。 -4. **不允许**数据面语义变化(域名集合、证书覆盖、上游、WAF、Pages 绑定)。 - -不一致时:停止发布,修正 Zone 域名/证书绑定后重新预览。 - -## 5. 发布与回滚 - -1. 预览通过后,在管理端执行**配置发布**,记录新版本号。 -2. 用根域与各子域发起 HTTP(S) 请求,确认节点应用成功。 -3. **回滚:** 在配置版本中重新激活导入前的版本;节点会拉取旧快照。数据库侧若需回退,使用升级前备份恢复(第二阶段删列后 Down 迁移不回填业务数据)。 - -## 6. 第二阶段:删除旧表与冗余列 - -仅在以下条件全部满足后执行: - -* `migrate-zones` 无冲突 -* 至少完成一次预览对比与发布(及必要时的回滚演练) -* 运维确认不再依赖 `of_managed_domains` 与 `of_proxy_routes` 上的 `domain` / `domains` / `cert_*` 列 - -然后升级到包含 `202607130001_drop_legacy_route_domain_columns` 的版本并启动 Server(自动 goose)。该迁移将: - -* 删除 `of_proxy_routes` 的 `domain`、`domains`、`cert_id`、`cert_ids`、`domain_cert_ids` -* 删除表 `of_managed_domains` - -**不可逆业务数据:** Down 仅在开发库重建空结构,不恢复历史域名行。 - -## 7. 命令速查 - -| 步骤 | 命令 / 操作 | -| --- | --- | -| 备份 | `pg_dump` / 复制 SQLite 文件 | -| 导入 | `wavelet migrate-zones` | -| 预览 | 管理端配置差异 / Preview API | -| 发布 | 管理端发布激活版本 | -| 回滚配置 | 管理端激活旧版本 | -| 回滚库 | 恢复备份(勿依赖 Down 填数) | +1. 在管理端查看配置差异 / 预览。 +2. **逐路由**核对:`server_name` 集合、证书路径、WAF Route ID、Pages 引用。 +3. **允许**旧快照 JSON 中路由上的冗余 `domain` / `domains` / `cert_ids` 消失。 +4. **不允许**数据面语义变化。 +5. 预览通过后发布;需要时在配置版本中激活升级前版本做配置回滚。数据库回退请使用升级前备份(Down 迁移不回填业务域名数据)。 ## 相关文档 diff --git a/internal/apps/openflare/routeidentity/identity.go b/internal/apps/openflare/routeidentity/identity.go index 29bfdd09..71443de6 100644 --- a/internal/apps/openflare/routeidentity/identity.go +++ b/internal/apps/openflare/routeidentity/identity.go @@ -36,8 +36,8 @@ func NormalizeDomains(rawDomains []string) ([]string, error) { return normalized, nil } -// DecodeDomains parses legacy route domain fields for the explicit migration -// command. Runtime consumers must read ZoneDomain bindings instead. +// DecodeDomains parses legacy route domain fields for the goose upgrade importer. +// Runtime consumers must read ZoneDomain bindings instead. func DecodeDomains(raw string, fallbackDomain string) ([]string, error) { text := strings.TrimSpace(raw) if text == "" { diff --git a/internal/apps/openflare/zone/legacy_import.go b/internal/apps/openflare/zone/legacy_import.go index 1f665c89..54c70b2d 100644 --- a/internal/apps/openflare/zone/legacy_import.go +++ b/internal/apps/openflare/zone/legacy_import.go @@ -5,14 +5,14 @@ package zone import ( "context" + "database/sql" "encoding/json" + "errors" "fmt" + "strconv" "strings" "github.com/Rain-kl/Wavelet/internal/apps/openflare/routeidentity" - "github.com/Rain-kl/Wavelet/internal/db" - "github.com/Rain-kl/Wavelet/internal/model" - "gorm.io/gorm" ) // ImportReport describes the idempotent legacy migration result. @@ -31,138 +31,244 @@ func (r ImportReport) LogAndReturn(err error) error { } type legacyDomain struct { - Domain string - CertID *uint - Remark string + Domain string + CertID *uint + Remark string + ProxyRouteID *uint } -// legacyRouteRow reads pre-cleanup of_proxy_routes columns via raw scan. -type legacyRouteRow struct { - ID uint `gorm:"column:id"` - Domain string `gorm:"column:domain"` - Domains string `gorm:"column:domains"` - DomainCertIDs string `gorm:"column:domain_cert_ids"` - Remark string `gorm:"column:remark"` -} - -// legacyManagedRow reads of_managed_domains while the table still exists. -type legacyManagedRow struct { - Domain string `gorm:"column:domain"` - CertID *uint `gorm:"column:cert_id"` - Remark string `gorm:"column:remark"` -} - -// ImportLegacy imports legacy proxy-route / managed-domain rows into Zone tables. -// After the phase-2 schema cleanup, missing legacy columns or tables are skipped. +// ImportLegacyTx imports legacy proxy-route / managed-domain rows into Zone tables +// within an existing SQL transaction (goose runs this on Server upgrade). +// postgres selects $n placeholders; otherwise SQLite-style ? is used. +// Missing legacy columns or tables are skipped so re-runs after phase-2 cleanup are no-ops. // -//nolint:cyclop // the transactional importer intentionally validates every legacy source in one pass. -func ImportLegacy(ctx context.Context) (report ImportReport, resultErr error) { - conn := db.DB(ctx) - if conn == nil { - return report, fmt.Errorf("database is not initialized") +//nolint:cyclop,gocyclo // single-pass legacy importer validates every source before write. +func ImportLegacyTx(ctx context.Context, tx *sql.Tx, postgres bool) (report ImportReport, err error) { + if tx == nil { + return report, errors.New("transaction is required") } - resultErr = conn.Transaction(func(tx *gorm.DB) error { - items := make([]legacyDomain, 0) - hasRouteDomains := false + q := func(sqlText string) string { return rebindSQL(sqlText, postgres) } - if tx.Migrator().HasColumn("of_proxy_routes", "domain") && - tx.Migrator().HasColumn("of_proxy_routes", "domains") { - var routes []legacyRouteRow - if err := tx.Table("of_proxy_routes"). - Select("id, domain, domains, domain_cert_ids, remark"). - Find(&routes).Error; err != nil { - return err - } - for _, route := range routes { - domains, err := routeidentity.DecodeDomains(route.Domains, route.Domain) - if err != nil { - report.Conflicts = append(report.Conflicts, fmt.Sprintf("route %d: %v", route.ID, err)) - continue - } - if len(domains) > 0 { - hasRouteDomains = true - } - certIDs := decodeLegacyCertIDs(route.DomainCertIDs, len(domains)) - for i, domain := range domains { - var certID *uint - if i < len(certIDs) && certIDs[i] > 0 { - v := certIDs[i] - certID = &v - } - items = append(items, legacyDomain{Domain: domain, CertID: certID, Remark: route.Remark}) - } + items := make([]legacyDomain, 0) + hasRouteDomains := false + + hasDomainCol, err := hasTableColumn(ctx, tx, q, postgres, "of_proxy_routes", "domain") + if err != nil { + return report, err + } + hasDomainsCol, err := hasTableColumn(ctx, tx, q, postgres, "of_proxy_routes", "domains") + if err != nil { + return report, err + } + if hasDomainCol && hasDomainsCol { + var collectErr error + items, hasRouteDomains, report.Conflicts, collectErr = collectLegacyRouteDomainsImpl(ctx, tx, q) + if collectErr != nil { + return report, collectErr + } + } + + if !hasRouteDomains { + exists, tableErr := hasTable(ctx, tx, q, postgres, "of_managed_domains") + if tableErr != nil { + return report, tableErr + } + if exists { + managed, managedErr := collectLegacyManagedDomains(ctx, tx, q) + if managedErr != nil { + return report, managedErr } + items = append(items, managed...) + } + } + + for _, item := range items { + domain, normErr := normalizeDomain(item.Domain) + if normErr != nil { + report.Conflicts = append(report.Conflicts, fmt.Sprintf("%s: %v", item.Domain, normErr)) + continue + } + root, rootErr := zoneRoot(domain) + if rootErr != nil { + report.Conflicts = append(report.Conflicts, fmt.Sprintf("%s: %v", domain, rootErr)) + continue } - if !hasRouteDomains && tx.Migrator().HasTable("of_managed_domains") { - var legacy []legacyManagedRow - if err := tx.Table("of_managed_domains"). - Select("domain, cert_id, remark"). - Find(&legacy).Error; err != nil { - return err - } - for _, item := range legacy { - items = append(items, legacyDomain(item)) + var existingID uint + var existingZoneDomain string + scanErr := tx.QueryRowContext(ctx, q(` + SELECT zd.id, z.domain + FROM of_zone_domains zd + JOIN of_zones z ON z.id = zd.zone_id + WHERE zd.domain = ? + `), domain).Scan(&existingID, &existingZoneDomain) + if scanErr == nil { + if existingZoneDomain != root { + report.Conflicts = append(report.Conflicts, fmt.Sprintf("%s: global domain conflict", domain)) + } else if item.ProxyRouteID != nil { + if _, bindErr := tx.ExecContext(ctx, q(` + UPDATE of_zone_domains + SET proxy_route_id = COALESCE(proxy_route_id, ?), + cert_id = COALESCE(cert_id, ?) + WHERE id = ? + `), *item.ProxyRouteID, nullableUint(item.CertID), existingID); bindErr != nil { + return report, bindErr + } } + continue + } + if !errors.Is(scanErr, sql.ErrNoRows) { + return report, scanErr } - for _, item := range items { - domain, err := normalizeDomain(item.Domain) - if err != nil { - report.Conflicts = append(report.Conflicts, fmt.Sprintf("%s: %v", item.Domain, err)) - continue - } - root, err := zoneRoot(domain) - if err != nil { - report.Conflicts = append(report.Conflicts, fmt.Sprintf("%s: %v", domain, err)) - continue - } - var existing model.ZoneDomain - err = tx.Where("domain = ?", domain).First(&existing).Error - if err == nil { - var z model.Zone - if tx.First(&z, existing.ZoneID).Error != nil || z.Domain != root { - report.Conflicts = append(report.Conflicts, fmt.Sprintf("%s: global domain conflict", domain)) - } - continue - } - if err != nil && !isNotFound(err) { - return err - } - var zone model.Zone - err = tx.Where("domain = ?", root).First(&zone).Error - if isNotFound(err) { - zone = model.Zone{Domain: root} - if err = tx.Create(&zone).Error; err != nil { - return err - } - report.Zones++ - } else if err != nil { - return err - } - if item.CertID != nil { - var cert model.TLSCertificate - if err = tx.First(&cert, *item.CertID).Error; err != nil { + zoneID, zoneErr := ensureZone(ctx, tx, q, root, &report) + if zoneErr != nil { + return report, zoneErr + } + + if item.CertID != nil { + var certID uint + if certErr := tx.QueryRowContext(ctx, q(`SELECT id FROM of_tls_certificates WHERE id = ?`), *item.CertID). + Scan(&certID); certErr != nil { + if errors.Is(certErr, sql.ErrNoRows) { report.Conflicts = append(report.Conflicts, fmt.Sprintf("%s: %s", domain, errCertificateNotFound)) continue } + return report, certErr } - if err = tx.Create(&model.ZoneDomain{ - ZoneID: zone.ID, - Domain: domain, - CertID: item.CertID, - Remark: item.Remark, - }).Error; err != nil { - return err + } + + if _, insErr := tx.ExecContext(ctx, q(` + INSERT INTO of_zone_domains (zone_id, proxy_route_id, domain, cert_id, remark, created_at, updated_at) + VALUES (?, ?, ?, ?, ?, CURRENT_TIMESTAMP, CURRENT_TIMESTAMP) + `), zoneID, nullableUint(item.ProxyRouteID), domain, nullableUint(item.CertID), item.Remark); insErr != nil { + return report, insErr + } + report.Domains++ + } + + if len(report.Conflicts) > 0 { + return report, fmt.Errorf("legacy data has conflicts") + } + return report, nil +} + +func ensureZone( + ctx context.Context, + tx *sql.Tx, + q func(string) string, + root string, + report *ImportReport, +) (uint, error) { + var zoneID uint + err := tx.QueryRowContext(ctx, q(`SELECT id FROM of_zones WHERE domain = ?`), root).Scan(&zoneID) + if err == nil { + return zoneID, nil + } + if !errors.Is(err, sql.ErrNoRows) { + return 0, err + } + if _, execErr := tx.ExecContext(ctx, q(` + INSERT INTO of_zones (domain, remark, created_at, updated_at) + VALUES (?, '', CURRENT_TIMESTAMP, CURRENT_TIMESTAMP) + `), root); execErr != nil { + return 0, execErr + } + if err := tx.QueryRowContext(ctx, q(`SELECT id FROM of_zones WHERE domain = ?`), root).Scan(&zoneID); err != nil { + return 0, err + } + report.Zones++ + return zoneID, nil +} + +func collectLegacyRouteDomainsImpl( + ctx context.Context, + tx *sql.Tx, + q func(string) string, +) (items []legacyDomain, hasRouteDomains bool, conflicts []string, err error) { + // Probe domain_cert_ids: if SELECT fails, fall back without it. + queryWithCert := q(`SELECT id, domain, domains, COALESCE(domain_cert_ids, '[]'), remark FROM of_proxy_routes`) + rows, err := tx.QueryContext(ctx, queryWithCert) + useCert := true + if err != nil { + useCert = false + rows, err = tx.QueryContext(ctx, q(`SELECT id, domain, domains, remark FROM of_proxy_routes`)) + if err != nil { + return nil, false, nil, err + } + } + defer func() { _ = rows.Close() }() + + for rows.Next() { + var ( + id uint + domain string + domains string + certIDs string + remark string + ) + if useCert { + if err := rows.Scan(&id, &domain, &domains, &certIDs, &remark); err != nil { + return nil, false, nil, err } - report.Domains++ + } else { + if err := rows.Scan(&id, &domain, &domains, &remark); err != nil { + return nil, false, nil, err + } + certIDs = "[]" } - if len(report.Conflicts) > 0 { - return fmt.Errorf("legacy data has conflicts") + decoded, decodeErr := routeidentity.DecodeDomains(domains, domain) + if decodeErr != nil { + conflicts = append(conflicts, fmt.Sprintf("route %d: %v", id, decodeErr)) + continue } - return nil - }) - return report, resultErr + if len(decoded) > 0 { + hasRouteDomains = true + } + ids := decodeLegacyCertIDs(certIDs, len(decoded)) + routeID := id + for i, d := range decoded { + var certID *uint + if i < len(ids) && ids[i] > 0 { + v := ids[i] + certID = &v + } + items = append(items, legacyDomain{ + Domain: d, + CertID: certID, + Remark: remark, + ProxyRouteID: &routeID, + }) + } + } + return items, hasRouteDomains, conflicts, rows.Err() +} + +func collectLegacyManagedDomains(ctx context.Context, tx *sql.Tx, q func(string) string) ([]legacyDomain, error) { + rows, err := tx.QueryContext(ctx, q(`SELECT domain, cert_id, remark FROM of_managed_domains`)) + if err != nil { + return nil, err + } + defer func() { _ = rows.Close() }() + + items := make([]legacyDomain, 0) + for rows.Next() { + var ( + domain string + certID sql.NullInt64 + remark string + ) + if err := rows.Scan(&domain, &certID, &remark); err != nil { + return nil, err + } + item := legacyDomain{Domain: domain, Remark: remark} + if certID.Valid && certID.Int64 > 0 { + v := uint(certID.Int64) + item.CertID = &v + } + items = append(items, item) + } + return items, rows.Err() } func decodeLegacyCertIDs(raw string, count int) []uint { @@ -176,4 +282,78 @@ func decodeLegacyCertIDs(raw string, count int) []uint { return values } -func isNotFound(err error) bool { return err == gorm.ErrRecordNotFound } +func nullableUint(v *uint) any { + if v == nil { + return nil + } + return *v +} + +func rebindSQL(query string, postgres bool) string { + if !postgres { + return query + } + var b strings.Builder + b.Grow(len(query) + len(query)/4) + n := 0 + for i := 0; i < len(query); i++ { + if query[i] == '?' { + n++ + b.WriteByte('$') + b.WriteString(strconv.Itoa(n)) + continue + } + b.WriteByte(query[i]) + } + return b.String() +} + +func hasTable( + ctx context.Context, + tx *sql.Tx, + _ func(string) string, + postgres bool, + table string, +) (bool, error) { + var count int + var err error + if postgres { + err = tx.QueryRowContext(ctx, ` + SELECT COUNT(*) FROM information_schema.tables + WHERE table_schema = 'public' AND table_name = $1 + `, table).Scan(&count) + } else { + err = tx.QueryRowContext(ctx, + `SELECT COUNT(*) FROM sqlite_master WHERE type = 'table' AND name = ?`, table, + ).Scan(&count) + } + if err != nil { + return false, err + } + return count > 0, nil +} + +func hasTableColumn( + ctx context.Context, + tx *sql.Tx, + _ func(string) string, + postgres bool, + table, column string, +) (bool, error) { + var count int + var err error + if postgres { + err = tx.QueryRowContext(ctx, ` + SELECT COUNT(*) FROM information_schema.columns + WHERE table_schema = 'public' AND table_name = $1 AND column_name = $2 + `, table, column).Scan(&count) + } else { + err = tx.QueryRowContext(ctx, + `SELECT COUNT(*) FROM pragma_table_info(?) WHERE name = ?`, table, column, + ).Scan(&count) + } + if err != nil { + return false, err + } + return count > 0, nil +} diff --git a/internal/apps/openflare/zone/legacy_import_test.go b/internal/apps/openflare/zone/legacy_import_test.go new file mode 100644 index 00000000..68fcd198 --- /dev/null +++ b/internal/apps/openflare/zone/legacy_import_test.go @@ -0,0 +1,154 @@ +// Copyright 2026 Arctel.net +// SPDX-License-Identifier: Apache-2.0 + +package zone + +import ( + "context" + "database/sql" + "testing" + + "github.com/Rain-kl/Wavelet/internal/db" + "github.com/glebarez/sqlite" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" + "gorm.io/gorm" +) + +func setupLegacyImportDB(t *testing.T) (*sql.DB, func()) { + t.Helper() + gormDB, err := gorm.Open(sqlite.Open(":memory:"), &gorm.Config{}) + require.NoError(t, err) + sqlDB, err := gormDB.DB() + require.NoError(t, err) + + // Pre-phase-2 schema: legacy route columns + managed domains + zone tables. + stmts := []string{ + `CREATE TABLE of_zones ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + domain TEXT NOT NULL UNIQUE, + remark TEXT NOT NULL DEFAULT '', + created_at DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP, + updated_at DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP + )`, + `CREATE TABLE of_zone_domains ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + zone_id INTEGER NOT NULL, + proxy_route_id INTEGER, + domain TEXT NOT NULL UNIQUE, + cert_id INTEGER, + remark TEXT NOT NULL DEFAULT '', + created_at DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP, + updated_at DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP + )`, + `CREATE TABLE of_proxy_routes ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + site_name TEXT NOT NULL DEFAULT '', + domain TEXT NOT NULL DEFAULT '', + domains TEXT NOT NULL DEFAULT '[]', + domain_cert_ids TEXT NOT NULL DEFAULT '[]', + origin_url TEXT NOT NULL DEFAULT '', + remark TEXT NOT NULL DEFAULT '' + )`, + `CREATE TABLE of_tls_certificates ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + name TEXT NOT NULL DEFAULT '' + )`, + `CREATE TABLE of_managed_domains ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + domain TEXT NOT NULL, + cert_id INTEGER, + remark TEXT NOT NULL DEFAULT '' + )`, + } + for _, stmt := range stmts { + _, err := sqlDB.Exec(stmt) + require.NoError(t, err) + } + + previous := db.DB(context.Background()) + db.SetDB(gormDB) + return sqlDB, func() { + db.SetDB(previous) + _ = sqlDB.Close() + } +} + +func TestImportLegacyTxBindsRouteDomains(t *testing.T) { + sqlDB, cleanup := setupLegacyImportDB(t) + defer cleanup() + ctx := context.Background() + + _, err := sqlDB.Exec(`INSERT INTO of_tls_certificates (id, name) VALUES (7, 'cert')`) + require.NoError(t, err) + _, err = sqlDB.Exec(` + INSERT INTO of_proxy_routes (id, site_name, domain, domains, domain_cert_ids, origin_url, remark) + VALUES (3, 'api', 'api.example.com', '["api.example.com","www.example.com"]', '[7,7]', 'http://origin', 'r') + `) + require.NoError(t, err) + + tx, err := sqlDB.Begin() + require.NoError(t, err) + report, err := ImportLegacyTx(ctx, tx, false) + require.NoError(t, err) + require.NoError(t, tx.Commit()) + + assert.Equal(t, 1, report.Zones) + assert.Equal(t, 2, report.Domains) + + var zoneDomain string + require.NoError(t, sqlDB.QueryRow(`SELECT domain FROM of_zones`).Scan(&zoneDomain)) + assert.Equal(t, "example.com", zoneDomain) + + var count int + require.NoError(t, sqlDB.QueryRow(`SELECT COUNT(*) FROM of_zone_domains WHERE proxy_route_id = 3`).Scan(&count)) + assert.Equal(t, 2, count) + + // Idempotent re-run + tx, err = sqlDB.Begin() + require.NoError(t, err) + report2, err := ImportLegacyTx(ctx, tx, false) + require.NoError(t, err) + require.NoError(t, tx.Commit()) + assert.Equal(t, 0, report2.Domains) +} + +func TestImportLegacyTxNoOpWithoutLegacyColumns(t *testing.T) { + gormDB, err := gorm.Open(sqlite.Open(":memory:"), &gorm.Config{}) + require.NoError(t, err) + sqlDB, err := gormDB.DB() + require.NoError(t, err) + defer sqlDB.Close() + + _, err = sqlDB.Exec(` + CREATE TABLE of_zones ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + domain TEXT NOT NULL UNIQUE, + remark TEXT NOT NULL DEFAULT '', + created_at DATETIME, updated_at DATETIME + ); + CREATE TABLE of_zone_domains ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + zone_id INTEGER NOT NULL, + proxy_route_id INTEGER, + domain TEXT NOT NULL UNIQUE, + cert_id INTEGER, + remark TEXT NOT NULL DEFAULT '', + created_at DATETIME, updated_at DATETIME + ); + CREATE TABLE of_proxy_routes ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + site_name TEXT NOT NULL DEFAULT '', + origin_url TEXT NOT NULL DEFAULT '' + ); + `) + require.NoError(t, err) + + tx, err := sqlDB.Begin() + require.NoError(t, err) + report, err := ImportLegacyTx(context.Background(), tx, false) + require.NoError(t, err) + require.NoError(t, tx.Commit()) + assert.Equal(t, 0, report.Zones) + assert.Equal(t, 0, report.Domains) +} diff --git a/internal/cmd/migrate_zones.go b/internal/cmd/migrate_zones.go deleted file mode 100644 index 92935451..00000000 --- a/internal/cmd/migrate_zones.go +++ /dev/null @@ -1,21 +0,0 @@ -// Copyright 2026 Arctel.net -// SPDX-License-Identifier: Apache-2.0 - -package cmd - -import ( - "context" - - "github.com/Rain-kl/Wavelet/internal/apps/openflare/zone" - "github.com/Rain-kl/Wavelet/internal/db/migrator" - "github.com/spf13/cobra" -) - -var migrateZonesCmd = &cobra.Command{ - Use: "migrate-zones", Short: "导入旧域名数据到 Zone", - PreRun: func(_ *cobra.Command, _ []string) { migrator.Migrate() }, - RunE: func(_ *cobra.Command, _ []string) error { - report, err := zone.ImportLegacy(context.Background()) - return report.LogAndReturn(err) - }, -} diff --git a/internal/cmd/root.go b/internal/cmd/root.go index daecd885..c287cf84 100644 --- a/internal/cmd/root.go +++ b/internal/cmd/root.go @@ -72,7 +72,7 @@ func init() { schedulerCmd.PreRun = migratePreRun // 2. 集中将这些命令注册为真正的子命令,以解决 Cobra 的 unknown command 校验限制 - rootCmd.AddCommand(allCmd, apiCmd, workerCmd, schedulerCmd, migrateZonesCmd) + rootCmd.AddCommand(allCmd, apiCmd, workerCmd, schedulerCmd) } // Execute 执行根命令 diff --git a/internal/db/migrator/202606050001_bridge_legacy.go b/internal/db/migrator/202606050001_bridge_legacy.go deleted file mode 100644 index 26ba4e0e..00000000 --- a/internal/db/migrator/202606050001_bridge_legacy.go +++ /dev/null @@ -1,204 +0,0 @@ -// 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 deleted file mode 100644 index 2f901a36..00000000 --- a/internal/db/migrator/202606200006_migrate_legacy_data.go +++ /dev/null @@ -1,550 +0,0 @@ -// 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) - } - } - } - - if dialect == dialectPostgres { - if _, err := tx.ExecContext(ctx, ` - SELECT setval( - pg_get_serial_sequence('of_waf_rule_group_bindings', 'id'), - GREATEST(COALESCE((SELECT MAX(id) FROM of_waf_rule_group_bindings), 0), 1), - COALESCE((SELECT MAX(id) FROM of_waf_rule_group_bindings), 0) > 0 - ) - `); err != nil { - return fmt.Errorf("sync of_waf_rule_group_bindings sequence failed: %w", 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 deleted file mode 100644 index 75a69b7c..00000000 --- a/internal/db/migrator/bridge_test.go +++ /dev/null @@ -1,261 +0,0 @@ -// 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(GeoIPProvider 会被迁移到 w_system_configs.geoip_provider) - `INSERT INTO options (key, value) VALUES ('GeoIPProvider', 'ip-api');`, - // 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 legacy option migrated all the way to w_system_configs(of_options 已被删除) - var geoIPProviderVal string - if err := sqliteDB.Table("w_system_configs").Where("key = ?", "geoip_provider").Select("value").Scan(&geoIPProviderVal).Error; err != nil { - t.Fatalf("query geoip_provider failed: %v", err) - } - if geoIPProviderVal != "ip-api" { - t.Errorf("geoip_provider value incorrect: got %s, want ip-api", geoIPProviderVal) - } - - // 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/goose/postgres/202606050001_bridge_legacy.sql b/internal/db/migrator/goose/postgres/202606050001_bridge_legacy.sql new file mode 100644 index 00000000..52839fa9 --- /dev/null +++ b/internal/db/migrator/goose/postgres/202606050001_bridge_legacy.sql @@ -0,0 +1,7 @@ +-- +goose Up +-- Historical Wavelet→OpenFlare table rename bridge (formerly Go migration). +-- Fresh installs and current of_* schemas need no action. +SELECT 1; + +-- +goose Down +SELECT 1; diff --git a/internal/db/migrator/goose/postgres/202606200006_migrate_legacy_data.sql b/internal/db/migrator/goose/postgres/202606200006_migrate_legacy_data.sql new file mode 100644 index 00000000..9840f6cd --- /dev/null +++ b/internal/db/migrator/goose/postgres/202606200006_migrate_legacy_data.sql @@ -0,0 +1,7 @@ +-- +goose Up +-- Historical copy from legacy_* tables into of_* / w_* (formerly Go migration). +-- Already-applied environments keep their data; new installs have no legacy_* sources. +SELECT 1; + +-- +goose Down +SELECT 1; diff --git a/internal/db/migrator/goose/postgres/202607120002_import_zone_domains.sql b/internal/db/migrator/goose/postgres/202607120002_import_zone_domains.sql new file mode 100644 index 00000000..97890d8c --- /dev/null +++ b/internal/db/migrator/goose/postgres/202607120002_import_zone_domains.sql @@ -0,0 +1,8 @@ +-- +goose Up +-- Marker: Zone 域名从旧路由列 / of_managed_domains 的导入由 migrator.Migrate() +-- 在全部 goose SQL 应用后自动执行(publicsuffix 根域解析无法用纯 SQL 正确完成)。 +-- 本文件仅占位版本号,保证升级链路有序:建表 → 导入 → 删旧列。 +SELECT 1; + +-- +goose Down +SELECT 1; diff --git a/internal/db/migrator/goose/sqlite/202606050001_bridge_legacy.sql b/internal/db/migrator/goose/sqlite/202606050001_bridge_legacy.sql new file mode 100644 index 00000000..52839fa9 --- /dev/null +++ b/internal/db/migrator/goose/sqlite/202606050001_bridge_legacy.sql @@ -0,0 +1,7 @@ +-- +goose Up +-- Historical Wavelet→OpenFlare table rename bridge (formerly Go migration). +-- Fresh installs and current of_* schemas need no action. +SELECT 1; + +-- +goose Down +SELECT 1; diff --git a/internal/db/migrator/goose/sqlite/202606200006_migrate_legacy_data.sql b/internal/db/migrator/goose/sqlite/202606200006_migrate_legacy_data.sql new file mode 100644 index 00000000..9840f6cd --- /dev/null +++ b/internal/db/migrator/goose/sqlite/202606200006_migrate_legacy_data.sql @@ -0,0 +1,7 @@ +-- +goose Up +-- Historical copy from legacy_* tables into of_* / w_* (formerly Go migration). +-- Already-applied environments keep their data; new installs have no legacy_* sources. +SELECT 1; + +-- +goose Down +SELECT 1; diff --git a/internal/db/migrator/goose/sqlite/202607120002_import_zone_domains.sql b/internal/db/migrator/goose/sqlite/202607120002_import_zone_domains.sql new file mode 100644 index 00000000..97890d8c --- /dev/null +++ b/internal/db/migrator/goose/sqlite/202607120002_import_zone_domains.sql @@ -0,0 +1,8 @@ +-- +goose Up +-- Marker: Zone 域名从旧路由列 / of_managed_domains 的导入由 migrator.Migrate() +-- 在全部 goose SQL 应用后自动执行(publicsuffix 根域解析无法用纯 SQL 正确完成)。 +-- 本文件仅占位版本号,保证升级链路有序:建表 → 导入 → 删旧列。 +SELECT 1; + +-- +goose Down +SELECT 1; diff --git a/internal/db/migrator/migrator.go b/internal/db/migrator/migrator.go index 1a56b21e..a3860269 100644 --- a/internal/db/migrator/migrator.go +++ b/internal/db/migrator/migrator.go @@ -12,6 +12,7 @@ import ( "fmt" "log" + "github.com/Rain-kl/Wavelet/internal/apps/openflare/zone" "github.com/Rain-kl/Wavelet/internal/config" "github.com/Rain-kl/Wavelet/internal/db" "github.com/Rain-kl/Wavelet/internal/repository" @@ -34,7 +35,9 @@ func dbType() string { const ( dialectSqlite = "sqlite3" dialectPostgres = "postgres" - cascadeSuffix = " CASCADE" + // zoneImportSQLVersion is the goose SQL marker after of_zones creation and + // before drop of legacy route domain columns. Zone data import runs here. + zoneImportSQLVersion int64 = 202607120002 ) func gooseDialect() string { @@ -51,7 +54,7 @@ func migrationDir() string { return "goose/postgres" } -// Migrate 执行数据库迁移 +// Migrate 执行数据库迁移:全部结构变更走 goose SQL;Zone 历史域名导入在 SQL 之后自动执行。 func Migrate() { gormDB := db.DB(context.Background()) if gormDB == nil { @@ -70,6 +73,15 @@ func Migrate() { if err := resyncGooseVersionSequence(sqlDB); err != nil { log.Fatalf("[%s] resync goose_db_version sequence failed: %v\n", dbType(), err) } + // 1) SQL up to zone-import marker (includes of_zones DDL; still has legacy columns). + if err := goose.UpTo(sqlDB, migrationDir(), zoneImportSQLVersion); err != nil { + log.Fatalf("[%s] goose migrate (up to zone import) failed: %v\n", dbType(), err) + } + // 2) Auto-import legacy domains (publicsuffix; idempotent; no-op after phase-2 drop). + if err := importZoneDomainsAfterGoose(sqlDB); err != nil { + log.Fatalf("[%s] zone domain import failed: %v\n", dbType(), err) + } + // 3) Remaining SQL (drop legacy columns / managed_domains, later migrations). if err := goose.Up(sqlDB, migrationDir()); err != nil { log.Fatalf("[%s] goose migrate failed: %v\n", dbType(), err) } @@ -79,6 +91,31 @@ func Migrate() { log.Printf("[%s] goose migrate success\n", dbType()) } +// importZoneDomainsAfterGoose 在 SQL 迁移完成后自动导入旧路由/托管域名。 +// 必须使用 publicsuffix 解析注册根域,故不能放在纯 SQL 中;幂等且在旧列删除后为空操作。 +func importZoneDomainsAfterGoose(sqlDB *sql.DB) error { + ctx := context.Background() + tx, err := sqlDB.BeginTx(ctx, nil) + if err != nil { + return fmt.Errorf("begin zone import transaction: %w", err) + } + report, err := zone.ImportLegacyTx(ctx, tx, gooseDialect() == dialectPostgres) + if err != nil { + _ = tx.Rollback() + return report.LogAndReturn(err) + } + if err := tx.Commit(); err != nil { + return fmt.Errorf("commit zone import: %w", err) + } + if report.Zones > 0 || report.Domains > 0 { + log.Printf( + "[%s] imported zone domains automatically: zones=%d domains=%d\n", + dbType(), report.Zones, report.Domains, + ) + } + return nil +} + // resyncGooseVersionSequence 修复 PostgreSQL 下 goose_db_version.id 自增序列落后于 // MAX(id) 的问题(常见于从 dump 恢复或历史迁移以显式 id 复制数据后)。序列落后会 // 导致 goose 记录新版本号时 INSERT 命中 goose_db_version_pkey 唯一约束冲突。 @@ -115,17 +152,3 @@ func clearSystemConfigCache() { log.Printf("[%s] clear system config cache failed: %v\n", dbType(), err) } } - -func tableExistsSQL(dialect string) string { - if dialect == dialectPostgres { - return "SELECT count(*) FROM information_schema.tables WHERE table_schema='public' AND table_name=$1" - } - return "SELECT count(*) FROM sqlite_master WHERE type='table' AND name=?" -} - -func tablesWithPrefixSQL(dialect string) string { - if dialect == dialectPostgres { - return "SELECT table_name FROM information_schema.tables WHERE table_schema='public' AND table_name LIKE $1" - } - return "SELECT name FROM sqlite_master WHERE type='table' AND name LIKE ?" -}