mirror of
https://github.com/Rain-kl/OpenFlare.git
synced 2026-09-29 05:56:38 +08:00
[优化] 压缩脚本到 V16
This commit is contained in:
@@ -190,7 +190,7 @@ func TestRegisterShardingAutoMigratesShardTables(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
func TestUpgradeDatabaseSchemaV15ToV16AddsWAFIPGroups(t *testing.T) {
|
||||
func TestUpgradeDatabaseSchemaV15ToV16AppliesCompressedReleaseSchema(t *testing.T) {
|
||||
db := openBareTestSQLiteDB(t, "v16.db")
|
||||
if err := registerSharding(db, "sqlite"); err != nil {
|
||||
t.Fatalf("register sharding: %v", err)
|
||||
@@ -216,6 +216,21 @@ func TestUpgradeDatabaseSchemaV15ToV16AddsWAFIPGroups(t *testing.T) {
|
||||
if !db.Migrator().HasColumn(&WAFRuleGroup{}, "ip_whitelist_groups") {
|
||||
t.Fatal("expected waf_rule_groups.ip_whitelist_groups column")
|
||||
}
|
||||
if !db.Migrator().HasColumn(&Node{}, "access_token") {
|
||||
t.Fatal("expected nodes.access_token column")
|
||||
}
|
||||
if !db.Migrator().HasColumn(&Node{}, "version") {
|
||||
t.Fatal("expected nodes.version column")
|
||||
}
|
||||
if !db.Migrator().HasColumn(&Node{}, "ext_version") {
|
||||
t.Fatal("expected nodes.ext_version column")
|
||||
}
|
||||
if !db.Migrator().HasColumn(&ProxyRoute{}, "tunnel_node_id") {
|
||||
t.Fatal("expected proxy_routes.tunnel_node_id column")
|
||||
}
|
||||
if db.Migrator().HasTable("tunnels") {
|
||||
t.Fatal("expected pre-release tunnels table to be absent")
|
||||
}
|
||||
version, ok, err := loadDatabaseSchemaVersion(db)
|
||||
if err != nil {
|
||||
t.Fatalf("load schema version: %v", err)
|
||||
@@ -522,8 +537,8 @@ func TestEnsureDatabaseSchemaUpToDateAddsNodeIPManualOverride(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
func TestEnsureDatabaseSchemaUpToDateV21BackfillsNodeColumnsWhenNewColumnsAlreadyExist(t *testing.T) {
|
||||
db := openBareTestSQLiteDB(t, "node-v21-existing-target-columns.db")
|
||||
func TestEnsureDatabaseSchemaUpToDateV16BackfillsNodeColumnsWhenNewColumnsAlreadyExist(t *testing.T) {
|
||||
db := openBareTestSQLiteDB(t, "node-v16-existing-target-columns.db")
|
||||
if err := registerSharding(db, "sqlite"); err != nil {
|
||||
t.Fatalf("register sharding: %v", err)
|
||||
}
|
||||
@@ -553,10 +568,10 @@ func TestEnsureDatabaseSchemaUpToDateV21BackfillsNodeColumnsWhenNewColumnsAlread
|
||||
agent_token, agent_version, nginx_version,
|
||||
status, last_seen_at, created_at, updated_at
|
||||
) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)
|
||||
`, "node-v21", "Node v21", "127.0.0.1", "", "", "", "legacy-token", "v2.0.0", "openresty/1.25.3", "offline", now, now, now).Error; err != nil {
|
||||
`, "node-v16", "Node v16", "127.0.0.1", "", "", "", "legacy-token", "v2.0.0", "openresty/1.25.3", "offline", now, now, now).Error; err != nil {
|
||||
t.Fatalf("seed node with legacy columns: %v", err)
|
||||
}
|
||||
if err := saveDatabaseSchemaVersion(db, 20); err != nil {
|
||||
if err := saveDatabaseSchemaVersion(db, 15); err != nil {
|
||||
t.Fatalf("save schema version: %v", err)
|
||||
}
|
||||
|
||||
@@ -565,7 +580,7 @@ func TestEnsureDatabaseSchemaUpToDateV21BackfillsNodeColumnsWhenNewColumnsAlread
|
||||
}
|
||||
|
||||
var node Node
|
||||
if err := db.Where("node_id = ?", "node-v21").First(&node).Error; err != nil {
|
||||
if err := db.Where("node_id = ?", "node-v16").First(&node).Error; err != nil {
|
||||
t.Fatalf("query migrated node: %v", err)
|
||||
}
|
||||
if node.AccessToken != "legacy-token" {
|
||||
|
||||
@@ -30,6 +30,14 @@ func (nodeV15) TableName() string {
|
||||
}
|
||||
|
||||
func migrateV15(ctx Context, db *gorm.DB, backend string) error {
|
||||
if db == nil {
|
||||
return fmt.Errorf("database handle is nil")
|
||||
}
|
||||
if !db.Migrator().HasColumn(&nodeV15{}, "ip_manual_override") {
|
||||
if err := db.Migrator().AddColumn(&nodeV15{}, "IPManualOverride"); err != nil {
|
||||
return fmt.Errorf("add nodes.ip_manual_override: %w", err)
|
||||
}
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
|
||||
@@ -1,23 +1,46 @@
|
||||
// v16 is the first database migration after the V15 formal release baseline.
|
||||
// It folds the previously drafted v16-v21 schema work into a single official
|
||||
// upgrade: tunnel-relay fields, WAF IP groups, current node identity/version
|
||||
// columns, and split node observation tables. The migration also backfills
|
||||
// legacy node columns and removes obsolete pre-release tunnel metadata when
|
||||
// present, so V15 deployments can upgrade directly to the new formal schema.
|
||||
package migrate
|
||||
|
||||
import (
|
||||
"fmt"
|
||||
"log/slog"
|
||||
|
||||
"gorm.io/gorm"
|
||||
)
|
||||
|
||||
type nodeV16 struct {
|
||||
NodeType string `gorm:"column:node_type;not null;default:'edge_node'"`
|
||||
RelayBindPort int `gorm:"column:relay_bind_port"`
|
||||
func (nodeV16) TableName() string {
|
||||
return "nodes"
|
||||
}
|
||||
|
||||
type tunnelV16 struct {
|
||||
ID uint `gorm:"primaryKey"`
|
||||
func (tunnelV16) TableName() string {
|
||||
return "tunnels"
|
||||
}
|
||||
|
||||
type proxyRouteV16 struct {
|
||||
UpstreamType string `gorm:"column:upstream_type;not null;default:'direct'"`
|
||||
TunnelID *uint `gorm:"column:tunnel_id"`
|
||||
func (proxyRouteV16) TableName() string {
|
||||
return "proxy_routes"
|
||||
}
|
||||
|
||||
type nodeV16 struct{}
|
||||
|
||||
type tunnelV16 struct{}
|
||||
|
||||
type proxyRouteV16 struct{}
|
||||
|
||||
type wafIPGroupV16 struct{}
|
||||
|
||||
type wafRuleGroupV16 struct{}
|
||||
|
||||
func (wafIPGroupV16) TableName() string {
|
||||
return "waf_ip_groups"
|
||||
}
|
||||
|
||||
func (wafRuleGroupV16) TableName() string {
|
||||
return "waf_rule_groups"
|
||||
}
|
||||
|
||||
func init() {
|
||||
@@ -33,28 +56,50 @@ func V16() Migration {
|
||||
}
|
||||
}
|
||||
|
||||
func (nodeV16) TableName() string {
|
||||
return "nodes"
|
||||
}
|
||||
|
||||
func (tunnelV16) TableName() string {
|
||||
return "tunnels"
|
||||
}
|
||||
|
||||
func (proxyRouteV16) TableName() string {
|
||||
return "proxy_routes"
|
||||
}
|
||||
|
||||
func migrateV16(ctx Context, db *gorm.DB, backend string) error {
|
||||
if err := db.AutoMigrate(&tunnelV16{}); err != nil {
|
||||
return fmt.Errorf("auto migrate tunnelV16: %w", err)
|
||||
if err := ctx.ApplyCurrentSchema(db, backend); err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
migrator := db.Migrator()
|
||||
if migrator.HasColumn(&nodeV16{}, "agent_token") {
|
||||
if err := db.Exec(`UPDATE nodes SET access_token = agent_token WHERE access_token IS NULL OR access_token = ''`).Error; err != nil {
|
||||
return fmt.Errorf("backfill nodes.access_token from agent_token: %w", err)
|
||||
}
|
||||
}
|
||||
if migrator.HasColumn(&nodeV16{}, "agent_version") {
|
||||
if err := db.Exec(`UPDATE nodes SET version = agent_version WHERE version = '' OR version IS NULL`).Error; err != nil {
|
||||
return fmt.Errorf("backfill nodes.version from agent_version: %w", err)
|
||||
}
|
||||
}
|
||||
if migrator.HasColumn(&nodeV16{}, "nginx_version") {
|
||||
if err := db.Exec(`UPDATE nodes SET ext_version = nginx_version WHERE ext_version IS NULL OR ext_version = ''`).Error; err != nil {
|
||||
return fmt.Errorf("backfill nodes.ext_version from nginx_version: %w", err)
|
||||
}
|
||||
}
|
||||
|
||||
if err := db.Exec("UPDATE nodes SET node_type = 'edge_node' WHERE node_type = '' OR node_type IS NULL").Error; err != nil {
|
||||
return fmt.Errorf("backfill nodes.node_type: %w", err)
|
||||
}
|
||||
if err := db.Exec("UPDATE proxy_routes SET upstream_type = 'direct' WHERE upstream_type = '' OR upstream_type IS NULL").Error; err != nil {
|
||||
return fmt.Errorf("backfill proxy_routes.upstream_type: %w", err)
|
||||
}
|
||||
|
||||
if migrator.HasColumn(&proxyRouteV16{}, "tunnel_id") {
|
||||
if err := db.Model(&proxyRouteV16{}).Where("upstream_type = ?", "tunnel").Update("upstream_type", "direct").Error; err != nil {
|
||||
return fmt.Errorf("reset pre-release tunnel proxy routes: %w", err)
|
||||
}
|
||||
if err := migrator.DropColumn(&proxyRouteV16{}, "tunnel_id"); err != nil {
|
||||
return fmt.Errorf("drop pre-release proxy_routes.tunnel_id: %w", err)
|
||||
}
|
||||
}
|
||||
if migrator.HasTable(&tunnelV16{}) {
|
||||
if err := migrator.DropTable(&tunnelV16{}); err != nil {
|
||||
return fmt.Errorf("drop pre-release tunnels table: %w", err)
|
||||
}
|
||||
slog.Info("dropped pre-release tunnels table during v16 migration")
|
||||
}
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
@@ -62,22 +107,57 @@ func validateV16(ctx Context, db *gorm.DB, backend string) error {
|
||||
if err := ctx.ValidateDatabaseSchemaVersion(db, backend, 15); err != nil {
|
||||
return err
|
||||
}
|
||||
if !db.Migrator().HasColumn(&proxyRouteV16{}, "tunnel_node_id") {
|
||||
if db == nil || !db.Migrator().HasTable(&tunnelV16{}) {
|
||||
return fmt.Errorf("table tunnels is missing")
|
||||
}
|
||||
if !db.Migrator().HasColumn(&proxyRouteV16{}, "tunnel_id") {
|
||||
return fmt.Errorf("column proxy_routes.tunnel_id is missing")
|
||||
if db == nil {
|
||||
return fmt.Errorf("database handle is nil")
|
||||
}
|
||||
|
||||
migrator := db.Migrator()
|
||||
for _, column := range []string{
|
||||
"access_token",
|
||||
"version",
|
||||
"ext_version",
|
||||
"node_type",
|
||||
"relay_bind_port",
|
||||
"relay_vhost_http_port",
|
||||
"relay_auth_token",
|
||||
"relay_agent_access_addr",
|
||||
"relay_client_access_addr",
|
||||
"relay_client_proxy_url",
|
||||
"relay_status",
|
||||
} {
|
||||
if !migrator.HasColumn(&nodeV16{}, column) {
|
||||
return fmt.Errorf("column nodes.%s is missing", column)
|
||||
}
|
||||
}
|
||||
if !db.Migrator().HasColumn(&nodeV16{}, "node_type") {
|
||||
return fmt.Errorf("column nodes.node_type is missing")
|
||||
for _, column := range []string{
|
||||
"upstream_type",
|
||||
"tunnel_node_id",
|
||||
"tunnel_target_addr",
|
||||
"tunnel_target_protocol",
|
||||
} {
|
||||
if !migrator.HasColumn(&proxyRouteV16{}, column) {
|
||||
return fmt.Errorf("column proxy_routes.%s is missing", column)
|
||||
}
|
||||
}
|
||||
if !db.Migrator().HasColumn(&proxyRouteV16{}, "upstream_type") {
|
||||
return fmt.Errorf("column proxy_routes.upstream_type is missing")
|
||||
if migrator.HasColumn(&proxyRouteV16{}, "tunnel_id") {
|
||||
return fmt.Errorf("column proxy_routes.tunnel_id should not exist in v16")
|
||||
}
|
||||
if !db.Migrator().HasColumn(&nodeV16{}, "relay_bind_port") {
|
||||
return fmt.Errorf("column nodes.relay_bind_port is missing")
|
||||
if migrator.HasTable(&tunnelV16{}) {
|
||||
return fmt.Errorf("table tunnels should not exist in v16")
|
||||
}
|
||||
if !migrator.HasTable(&wafIPGroupV16{}) {
|
||||
return fmt.Errorf("table waf_ip_groups is missing")
|
||||
}
|
||||
for _, column := range []string{
|
||||
"ip_whitelist_groups",
|
||||
"ip_blacklist_groups",
|
||||
} {
|
||||
if !migrator.HasColumn(&wafRuleGroupV16{}, column) {
|
||||
return fmt.Errorf("column waf_rule_groups.%s is missing", column)
|
||||
}
|
||||
}
|
||||
if !migrator.HasColumn(&wafIPGroupV16{}, "ext_ips") {
|
||||
return fmt.Errorf("column waf_ip_groups.ext_ips is missing")
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
@@ -1,55 +0,0 @@
|
||||
package migrate
|
||||
|
||||
import (
|
||||
"fmt"
|
||||
|
||||
"gorm.io/gorm"
|
||||
)
|
||||
|
||||
type wafIPGroupV17 struct{}
|
||||
|
||||
type wafRuleGroupV17 struct {
|
||||
IPWhitelistGroups string `gorm:"column:ip_whitelist_groups;type:text;not null;default:'[]'"`
|
||||
IPBlacklistGroups string `gorm:"column:ip_blacklist_groups;type:text;not null;default:'[]'"`
|
||||
}
|
||||
|
||||
func init() {
|
||||
Register(V17())
|
||||
}
|
||||
|
||||
func V17() Migration {
|
||||
return Migration{
|
||||
FromVersion: 16,
|
||||
ToVersion: 17,
|
||||
Migrate: migrateV17,
|
||||
Validate: validateV17,
|
||||
}
|
||||
}
|
||||
|
||||
func (wafIPGroupV17) TableName() string {
|
||||
return "waf_ip_groups"
|
||||
}
|
||||
|
||||
func (wafRuleGroupV17) TableName() string {
|
||||
return "waf_rule_groups"
|
||||
}
|
||||
|
||||
func migrateV17(ctx Context, db *gorm.DB, backend string) error {
|
||||
return nil
|
||||
}
|
||||
|
||||
func validateV17(ctx Context, db *gorm.DB, backend string) error {
|
||||
if err := ctx.ValidateDatabaseSchemaVersion(db, backend, 16); err != nil {
|
||||
return err
|
||||
}
|
||||
if db == nil || !db.Migrator().HasTable(&wafIPGroupV17{}) {
|
||||
return fmt.Errorf("table waf_ip_groups is missing")
|
||||
}
|
||||
if !db.Migrator().HasColumn(&wafRuleGroupV17{}, "ip_whitelist_groups") {
|
||||
return fmt.Errorf("column waf_rule_groups.ip_whitelist_groups is missing")
|
||||
}
|
||||
if !db.Migrator().HasColumn(&wafRuleGroupV17{}, "ip_blacklist_groups") {
|
||||
return fmt.Errorf("column waf_rule_groups.ip_blacklist_groups is missing")
|
||||
}
|
||||
return nil
|
||||
}
|
||||
@@ -1,48 +0,0 @@
|
||||
package migrate
|
||||
|
||||
import (
|
||||
"fmt"
|
||||
|
||||
"gorm.io/gorm"
|
||||
)
|
||||
|
||||
// WAF IP Group schema changes for v18:
|
||||
// We introduce the `ext_ips` column to keep track of captured IPs with their capture timestamps.
|
||||
// If more information needs to be saved for captured IPs, it can be added directly inside this JSON structure.
|
||||
type wafIPGroupV18 struct {
|
||||
ExtIPs string `gorm:"column:ext_ips;type:text;not null;default:'[]'"`
|
||||
}
|
||||
|
||||
func init() {
|
||||
Register(V18())
|
||||
}
|
||||
|
||||
func V18() Migration {
|
||||
return Migration{
|
||||
FromVersion: 17,
|
||||
ToVersion: 18,
|
||||
Migrate: migrateV18,
|
||||
Validate: validateV18,
|
||||
}
|
||||
}
|
||||
|
||||
func (wafIPGroupV18) TableName() string {
|
||||
return "waf_ip_groups"
|
||||
}
|
||||
|
||||
func migrateV18(ctx Context, db *gorm.DB, backend string) error {
|
||||
return nil
|
||||
}
|
||||
|
||||
func validateV18(ctx Context, db *gorm.DB, backend string) error {
|
||||
if err := ctx.ValidateDatabaseSchemaVersion(db, backend, 17); err != nil {
|
||||
return err
|
||||
}
|
||||
if db == nil || !db.Migrator().HasTable(&wafIPGroupV18{}) {
|
||||
return fmt.Errorf("table waf_ip_groups is missing")
|
||||
}
|
||||
if !db.Migrator().HasColumn(&wafIPGroupV18{}, "ext_ips") {
|
||||
return fmt.Errorf("column waf_ip_groups.ext_ips is missing")
|
||||
}
|
||||
return nil
|
||||
}
|
||||
@@ -1,88 +0,0 @@
|
||||
package migrate
|
||||
|
||||
import (
|
||||
"fmt"
|
||||
"log/slog"
|
||||
|
||||
"gorm.io/gorm"
|
||||
)
|
||||
|
||||
type proxyRouteV19 struct {
|
||||
ID uint `gorm:"primaryKey"`
|
||||
TunnelID *uint `gorm:"column:tunnel_id"`
|
||||
TunnelNodeID *uint `gorm:"column:tunnel_node_id"`
|
||||
UpstreamType string `gorm:"column:upstream_type"`
|
||||
}
|
||||
|
||||
type tunnelV19 struct{}
|
||||
|
||||
func (tunnelV19) TableName() string {
|
||||
return "tunnels"
|
||||
}
|
||||
|
||||
func (proxyRouteV19) TableName() string {
|
||||
return "proxy_routes"
|
||||
}
|
||||
|
||||
func init() {
|
||||
Register(V19())
|
||||
}
|
||||
|
||||
func V19() Migration {
|
||||
return Migration{
|
||||
FromVersion: 18,
|
||||
ToVersion: 19,
|
||||
Migrate: migrateV19,
|
||||
Validate: validateV19,
|
||||
}
|
||||
}
|
||||
|
||||
func migrateV19(ctx Context, db *gorm.DB, backend string) error {
|
||||
// Drop tunnels table
|
||||
if db.Migrator().HasTable(&tunnelV19{}) {
|
||||
if err := db.Migrator().DropTable(&tunnelV19{}); err != nil {
|
||||
return fmt.Errorf("failed to drop tunnels table: %w", err)
|
||||
}
|
||||
slog.Info("dropped tunnels table")
|
||||
}
|
||||
|
||||
// Add tunnel_node_id column
|
||||
if !db.Migrator().HasColumn(&proxyRouteV19{}, "tunnel_node_id") {
|
||||
if err := db.Migrator().AddColumn(&proxyRouteV19{}, "TunnelNodeID"); err != nil {
|
||||
return fmt.Errorf("failed to add tunnel_node_id to proxy_routes: %w", err)
|
||||
}
|
||||
slog.Info("added tunnel_node_id column to proxy_routes")
|
||||
}
|
||||
|
||||
// Drop old tunnel_id column
|
||||
if db.Migrator().HasColumn(&proxyRouteV19{}, "tunnel_id") {
|
||||
// Update routes that previously used tunnel to be disabled or direct to prevent dangling refs
|
||||
db.Model(&proxyRouteV19{}).Where("upstream_type = ?", "tunnel").Update("upstream_type", "direct")
|
||||
if err := db.Migrator().DropColumn(&proxyRouteV19{}, "tunnel_id"); err != nil {
|
||||
return fmt.Errorf("failed to drop tunnel_id column from proxy_routes: %w", err)
|
||||
}
|
||||
slog.Info("dropped tunnel_id column from proxy_routes")
|
||||
}
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
func validateV19(ctx Context, db *gorm.DB, backend string) error {
|
||||
if err := ctx.ValidateDatabaseSchemaVersion(db, backend, 18); err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
if db.Migrator().HasTable(&tunnelV19{}) {
|
||||
return fmt.Errorf("table tunnels should be dropped in v19")
|
||||
}
|
||||
|
||||
if !db.Migrator().HasColumn(&proxyRouteV19{}, "tunnel_node_id") {
|
||||
return fmt.Errorf("column proxy_routes.tunnel_node_id is missing")
|
||||
}
|
||||
|
||||
if db.Migrator().HasColumn(&proxyRouteV19{}, "tunnel_id") {
|
||||
return fmt.Errorf("column proxy_routes.tunnel_id should be dropped in v19")
|
||||
}
|
||||
|
||||
return nil
|
||||
}
|
||||
@@ -1,59 +0,0 @@
|
||||
// v20 records low-frequency Relay frps counters on nodes so the management UI
|
||||
// can show whether the relay runtime is alive and reporting tunnel load.
|
||||
package migrate
|
||||
|
||||
import (
|
||||
"fmt"
|
||||
|
||||
"gorm.io/gorm"
|
||||
)
|
||||
|
||||
type nodeV20 struct {
|
||||
ID uint `gorm:"primaryKey"`
|
||||
RelayFrpsConnections int `gorm:"column:relay_frps_connections"`
|
||||
RelayFrpsProxyCount int `gorm:"column:relay_frps_proxy_count"`
|
||||
}
|
||||
|
||||
func (nodeV20) TableName() string {
|
||||
return "nodes"
|
||||
}
|
||||
|
||||
func init() {
|
||||
Register(V20())
|
||||
}
|
||||
|
||||
func V20() Migration {
|
||||
return Migration{
|
||||
FromVersion: 19,
|
||||
ToVersion: 20,
|
||||
Migrate: migrateV20,
|
||||
Validate: validateV20,
|
||||
}
|
||||
}
|
||||
|
||||
func migrateV20(ctx Context, db *gorm.DB, backend string) error {
|
||||
if !db.Migrator().HasColumn(&nodeV20{}, "relay_frps_connections") {
|
||||
if err := db.Migrator().AddColumn(&nodeV20{}, "RelayFrpsConnections"); err != nil {
|
||||
return err
|
||||
}
|
||||
}
|
||||
if !db.Migrator().HasColumn(&nodeV20{}, "relay_frps_proxy_count") {
|
||||
if err := db.Migrator().AddColumn(&nodeV20{}, "RelayFrpsProxyCount"); err != nil {
|
||||
return err
|
||||
}
|
||||
}
|
||||
return validateV20(ctx, db, backend)
|
||||
}
|
||||
|
||||
func validateV20(ctx Context, db *gorm.DB, backend string) error {
|
||||
if err := ctx.ValidateDatabaseSchemaVersion(db, backend, 19); err != nil {
|
||||
return err
|
||||
}
|
||||
if !db.Migrator().HasColumn(&nodeV20{}, "relay_frps_connections") {
|
||||
return fmt.Errorf("column nodes.relay_frps_connections is missing")
|
||||
}
|
||||
if !db.Migrator().HasColumn(&nodeV20{}, "relay_frps_proxy_count") {
|
||||
return fmt.Errorf("column nodes.relay_frps_proxy_count is missing")
|
||||
}
|
||||
return nil
|
||||
}
|
||||
@@ -1,120 +0,0 @@
|
||||
// v21 renames agent_token to access_token, unifies versions, and separates node observabilities.
|
||||
package migrate
|
||||
|
||||
import (
|
||||
"fmt"
|
||||
"log/slog"
|
||||
|
||||
"gorm.io/gorm"
|
||||
)
|
||||
|
||||
type nodeV21 struct{}
|
||||
|
||||
func (nodeV21) TableName() string {
|
||||
return "nodes"
|
||||
}
|
||||
|
||||
func init() {
|
||||
Register(V21())
|
||||
}
|
||||
|
||||
func V21() Migration {
|
||||
return Migration{
|
||||
FromVersion: 20,
|
||||
ToVersion: 21,
|
||||
Migrate: migrateV21,
|
||||
Validate: validateV21,
|
||||
}
|
||||
}
|
||||
|
||||
func migrateV21(ctx Context, db *gorm.DB, backend string) error {
|
||||
slog.Info("starting v21 database migration (Node Optimization & Observation Split)")
|
||||
|
||||
migrator := db.Migrator()
|
||||
dropLegacyNodeColumn := func(col string) error {
|
||||
slog.Info("v21: keeping legacy nodes column to avoid lock-heavy schema rewrite", "backend", backend, "column", col)
|
||||
return nil
|
||||
}
|
||||
|
||||
// Rename agent_token → access_token (target may already exist if AutoMigrate ran earlier)
|
||||
if migrator.HasColumn(&nodeV21{}, "agent_token") {
|
||||
if migrator.HasColumn(&nodeV21{}, "access_token") {
|
||||
slog.Info("v21: access_token column already exists, backfilling from agent_token")
|
||||
if err := db.Exec(`UPDATE nodes SET access_token = agent_token WHERE access_token IS NULL OR access_token = ''`).Error; err != nil {
|
||||
return fmt.Errorf("failed to backfill access_token from agent_token: %w", err)
|
||||
}
|
||||
if err := dropLegacyNodeColumn("agent_token"); err != nil {
|
||||
slog.Warn("failed to drop agent_token after backfill", "error", err)
|
||||
}
|
||||
} else {
|
||||
if err := migrator.RenameColumn(&nodeV21{}, "agent_token", "access_token"); err != nil {
|
||||
return fmt.Errorf("failed to rename agent_token to access_token: %w", err)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// Rename agent_version → version (target may already exist if AutoMigrate ran earlier)
|
||||
if migrator.HasColumn(&nodeV21{}, "agent_version") {
|
||||
if migrator.HasColumn(&nodeV21{}, "version") {
|
||||
// version already created by AutoMigrate with empty default; backfill from agent_version
|
||||
slog.Info("v21: version column already exists, backfilling from agent_version")
|
||||
if err := db.Exec(`UPDATE nodes SET version = agent_version WHERE version = ''`).Error; err != nil {
|
||||
return fmt.Errorf("failed to backfill version from agent_version: %w", err)
|
||||
}
|
||||
if err := dropLegacyNodeColumn("agent_version"); err != nil {
|
||||
slog.Warn("failed to drop agent_version after backfill", "error", err)
|
||||
}
|
||||
} else {
|
||||
if err := migrator.RenameColumn(&nodeV21{}, "agent_version", "version"); err != nil {
|
||||
return fmt.Errorf("failed to rename agent_version to version: %w", err)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// Rename nginx_version → ext_version (target may already exist if AutoMigrate ran earlier)
|
||||
if migrator.HasColumn(&nodeV21{}, "nginx_version") {
|
||||
if migrator.HasColumn(&nodeV21{}, "ext_version") {
|
||||
slog.Info("v21: ext_version column already exists, backfilling from nginx_version")
|
||||
if err := db.Exec(`UPDATE nodes SET ext_version = nginx_version WHERE ext_version IS NULL OR ext_version = ''`).Error; err != nil {
|
||||
return fmt.Errorf("failed to backfill ext_version from nginx_version: %w", err)
|
||||
}
|
||||
if err := dropLegacyNodeColumn("nginx_version"); err != nil {
|
||||
slog.Warn("failed to drop nginx_version after backfill", "error", err)
|
||||
}
|
||||
} else {
|
||||
if err := migrator.RenameColumn(&nodeV21{}, "nginx_version", "ext_version"); err != nil {
|
||||
return fmt.Errorf("failed to rename nginx_version to ext_version: %w", err)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// Drop old merged columns
|
||||
columnsToDrop := []string{
|
||||
"relay_version",
|
||||
"relay_frp_version",
|
||||
"relay_frps_connections",
|
||||
"relay_frps_proxy_count",
|
||||
}
|
||||
|
||||
for _, col := range columnsToDrop {
|
||||
if migrator.HasColumn(&nodeV21{}, col) {
|
||||
if err := dropLegacyNodeColumn(col); err != nil {
|
||||
slog.Warn("failed to drop column in v21 migration", "column", col, "error", err)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
if err := ctx.ApplyCurrentSchemaExcept(db, backend, "nodes"); err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
slog.Info("completed v21 database migration")
|
||||
return validateV21(ctx, db, backend)
|
||||
}
|
||||
|
||||
func validateV21(ctx Context, db *gorm.DB, backend string) error {
|
||||
if err := ctx.ValidateDatabaseSchemaVersion(db, backend, 21); err != nil {
|
||||
return err
|
||||
}
|
||||
return nil
|
||||
}
|
||||
@@ -18,6 +18,9 @@ func V8() Migration {
|
||||
}
|
||||
|
||||
func migrateV8(ctx Context, db *gorm.DB, backend string) error {
|
||||
if err := ctx.ApplyCurrentSchema(db, backend); err != nil {
|
||||
return err
|
||||
}
|
||||
if err := ctx.BackfillOriginsFromProxyRoutes(db); err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
@@ -82,16 +82,6 @@ func (databaseSchemaMigrationContext) ValidateDatabaseSchemaVersion(db *gorm.DB,
|
||||
return validateDatabaseSchemaV15(db, backend)
|
||||
case 16:
|
||||
return validateDatabaseSchemaV16(db, backend)
|
||||
case 17:
|
||||
return validateDatabaseSchemaV17(db, backend)
|
||||
case 18:
|
||||
return validateDatabaseSchemaV18(db, backend)
|
||||
case 19:
|
||||
return validateDatabaseSchemaV19(db, backend)
|
||||
case 20:
|
||||
return validateDatabaseSchemaV20(db, backend)
|
||||
case 21:
|
||||
return validateDatabaseSchemaV21(db, backend)
|
||||
default:
|
||||
return fmt.Errorf("database schema validation for v%d is not defined", version)
|
||||
}
|
||||
@@ -1190,26 +1180,46 @@ func validateDatabaseSchemaV16(db *gorm.DB, backend string) error {
|
||||
if err := validateDatabaseSchemaV15(db, backend); err != nil {
|
||||
return err
|
||||
}
|
||||
if !db.Migrator().HasColumn(&ProxyRoute{}, "tunnel_node_id") {
|
||||
if !db.Migrator().HasTable("tunnels") {
|
||||
return fmt.Errorf("table tunnels is missing")
|
||||
}
|
||||
if !db.Migrator().HasColumn(&ProxyRoute{}, "tunnel_id") {
|
||||
return fmt.Errorf("column proxy_routes.tunnel_id is missing")
|
||||
}
|
||||
if !db.Migrator().HasColumn(&Node{}, "access_token") {
|
||||
return fmt.Errorf("column nodes.access_token is missing")
|
||||
}
|
||||
if !db.Migrator().HasColumn(&Node{}, "version") {
|
||||
return fmt.Errorf("column nodes.version is missing")
|
||||
}
|
||||
if !db.Migrator().HasColumn(&Node{}, "ext_version") {
|
||||
return fmt.Errorf("column nodes.ext_version is missing")
|
||||
}
|
||||
if !db.Migrator().HasColumn(&Node{}, "node_type") {
|
||||
return fmt.Errorf("column nodes.node_type is missing")
|
||||
}
|
||||
if !db.Migrator().HasColumn(&ProxyRoute{}, "upstream_type") {
|
||||
return fmt.Errorf("column proxy_routes.upstream_type is missing")
|
||||
for _, column := range []string{
|
||||
"relay_bind_port",
|
||||
"relay_vhost_http_port",
|
||||
"relay_auth_token",
|
||||
"relay_agent_access_addr",
|
||||
"relay_client_access_addr",
|
||||
"relay_client_proxy_url",
|
||||
"relay_status",
|
||||
} {
|
||||
if !db.Migrator().HasColumn(&Node{}, column) {
|
||||
return fmt.Errorf("column nodes.%s is missing", column)
|
||||
}
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
func validateDatabaseSchemaV17(db *gorm.DB, backend string) error {
|
||||
if err := validateDatabaseSchemaV16(db, backend); err != nil {
|
||||
return err
|
||||
for _, column := range []string{
|
||||
"upstream_type",
|
||||
"tunnel_node_id",
|
||||
"tunnel_target_addr",
|
||||
"tunnel_target_protocol",
|
||||
} {
|
||||
if !db.Migrator().HasColumn(&ProxyRoute{}, column) {
|
||||
return fmt.Errorf("column proxy_routes.%s is missing", column)
|
||||
}
|
||||
}
|
||||
if db.Migrator().HasTable("tunnels") {
|
||||
return fmt.Errorf("table tunnels should not exist in v16")
|
||||
}
|
||||
if db.Migrator().HasColumn(&ProxyRoute{}, "tunnel_id") {
|
||||
return fmt.Errorf("column proxy_routes.tunnel_id should not exist in v16")
|
||||
}
|
||||
if !db.Migrator().HasTable(&WAFIPGroup{}) {
|
||||
return fmt.Errorf("table waf_ip_groups is missing")
|
||||
@@ -1220,48 +1230,12 @@ func validateDatabaseSchemaV17(db *gorm.DB, backend string) error {
|
||||
if !db.Migrator().HasColumn(&WAFRuleGroup{}, "ip_blacklist_groups") {
|
||||
return fmt.Errorf("column waf_rule_groups.ip_blacklist_groups is missing")
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
func validateDatabaseSchemaV18(db *gorm.DB, backend string) error {
|
||||
if err := validateDatabaseSchemaV17(db, backend); err != nil {
|
||||
return err
|
||||
}
|
||||
if !db.Migrator().HasColumn(&WAFIPGroup{}, "ext_ips") {
|
||||
return fmt.Errorf("column waf_ip_groups.ext_ips is missing")
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
func validateDatabaseSchemaV19(db *gorm.DB, backend string) error {
|
||||
if err := validateDatabaseSchemaV18(db, backend); err != nil {
|
||||
return err
|
||||
}
|
||||
if db.Migrator().HasTable("tunnels") {
|
||||
return fmt.Errorf("table tunnels should be dropped in v19")
|
||||
}
|
||||
if !db.Migrator().HasColumn(&ProxyRoute{}, "tunnel_node_id") {
|
||||
return fmt.Errorf("column proxy_routes.tunnel_node_id is missing")
|
||||
}
|
||||
if db.Migrator().HasColumn(&ProxyRoute{}, "tunnel_id") {
|
||||
return fmt.Errorf("column proxy_routes.tunnel_id should be dropped in v19")
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
func validateDatabaseSchemaV20(db *gorm.DB, backend string) error {
|
||||
if err := validateDatabaseSchemaV19(db, backend); err != nil {
|
||||
return err
|
||||
}
|
||||
if !db.Migrator().HasColumn(&Node{}, "relay_frps_connections") {
|
||||
return fmt.Errorf("column nodes.relay_frps_connections is missing")
|
||||
}
|
||||
if !db.Migrator().HasColumn(&Node{}, "relay_frps_proxy_count") {
|
||||
return fmt.Errorf("column nodes.relay_frps_proxy_count is missing")
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
func databaseSchemaMigrations() []databaseSchemaMigration {
|
||||
ctx := databaseSchemaMigrationContext{}
|
||||
migrations := []databaseSchemaMigration{}
|
||||
@@ -1392,22 +1366,6 @@ func initializeFreshDatabaseSchema(db *gorm.DB, backend string) error {
|
||||
return saveDatabaseSchemaVersion(db, currentDatabaseSchemaVersion)
|
||||
}
|
||||
|
||||
func validateDatabaseSchemaV21(db *gorm.DB, backend string) error {
|
||||
if err := validateDatabaseSchemaV19(db, backend); err != nil {
|
||||
return err
|
||||
}
|
||||
if !db.Migrator().HasColumn(&Node{}, "access_token") {
|
||||
return fmt.Errorf("column nodes.access_token is missing")
|
||||
}
|
||||
if !db.Migrator().HasColumn(&Node{}, "version") {
|
||||
return fmt.Errorf("column nodes.version is missing")
|
||||
}
|
||||
if !db.Migrator().HasColumn(&Node{}, "ext_version") {
|
||||
return fmt.Errorf("column nodes.ext_version is missing")
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
func ensureDatabaseSchemaUpToDate(db *gorm.DB, backend string) error {
|
||||
version, exists, err := loadDatabaseSchemaVersion(db)
|
||||
if err != nil {
|
||||
|
||||
Reference in New Issue
Block a user