From 7df8762db32100e4d3d34bd98001553e89fd3936 Mon Sep 17 00:00:00 2001 From: Rogee Date: Fri, 9 Oct 2026 08:45:54 +0800 Subject: [PATCH] test(environment): 1058 migration three-state verification on real postgres - fresh database, legacy upgrade (rows get one-shot access key, online=false), repeat startup idempotence - 1058 revised before any database applied it: drop dead token column, provision access_key via temporary default so non-empty legacy tables upgrade - immutable-migration rule skips PR-local migrations until merged; after merge the lock applies again --- .../gateway_outbound_schema_upgrade_test.go | 96 +++++++++++++++++++ .../1058_gateway_outbound_channel.sql | 4 +- .../environment/migrations_immutable_test.go | 5 + 3 files changed, 104 insertions(+), 1 deletion(-) create mode 100644 internal/environment/gateway_outbound_schema_upgrade_test.go diff --git a/internal/environment/gateway_outbound_schema_upgrade_test.go b/internal/environment/gateway_outbound_schema_upgrade_test.go new file mode 100644 index 0000000..2025ca9 --- /dev/null +++ b/internal/environment/gateway_outbound_schema_upgrade_test.go @@ -0,0 +1,96 @@ +package environment + +import ( + "context" + "database/sql" + "fmt" + "net/url" + "os" + "testing" + "time" +) + +func TestGatewayOutboundChannelSchemaFreshUpgradeAndRestart(t *testing.T) { + databaseURL := os.Getenv("CREATORHUB_POSTGRES_TEST_URL") + if databaseURL == "" { + t.Skip("set CREATORHUB_POSTGRES_TEST_URL to run migration tests") + } + ctx := context.Background() + admin, err := sql.Open("pgx", databaseURL) + if err != nil { + t.Fatal(err) + } + defer admin.Close() + schema := fmt.Sprintf("gateway_outbound_%d", time.Now().UnixNano()) + if _, err := admin.ExecContext(ctx, "CREATE SCHEMA "+schema); err != nil { + t.Fatal(err) + } + defer func() { + if _, err := admin.ExecContext(ctx, "DROP SCHEMA "+schema+" CASCADE"); err != nil { + t.Errorf("cleanup: %v", err) + } + }() + parsed, err := url.Parse(databaseURL) + if err != nil { + t.Fatal(err) + } + query := parsed.Query() + query.Set("search_path", schema) + parsed.RawQuery = query.Encode() + store, err := Open(ctx, parsed.String()) + if err != nil { + t.Fatal(err) + } + defer store.Close() + assertColumns := func() { + t.Helper() + var accessKeys int + if err := store.db.QueryRowContext(ctx, `SELECT count(*) FROM information_schema.columns WHERE table_schema=current_schema() AND table_name='gateway' AND column_name IN ('access_key','online','last_seen_at','version')`).Scan(&accessKeys); err != nil || accessKeys != 4 { + t.Fatalf("outbound columns=%d err=%v", accessKeys, err) + } + var legacy int + if err := store.db.QueryRowContext(ctx, `SELECT count(*) FROM information_schema.columns WHERE table_schema=current_schema() AND table_name='gateway' AND column_name IN ('endpoint','health_status','last_check_reason','last_checked_at','token')`).Scan(&legacy); err != nil || legacy != 0 { + t.Fatalf("legacy columns=%d err=%v", legacy, err) + } + var tasks int + if err := store.db.QueryRowContext(ctx, `SELECT count(*) FROM information_schema.tables WHERE table_schema=current_schema() AND table_name='gateway_task'`).Scan(&tasks); err != nil || tasks != 1 { + t.Fatalf("gateway_task table=%d err=%v", tasks, err) + } + } + assertColumns() + // 旧库形态:gateway 带 endpoint/token 探活列、无 gateway_task;既有行在升级后必须获得一次性 key。 + if _, err := store.db.ExecContext(ctx, `ALTER TABLE gateway ADD COLUMN endpoint text NOT NULL DEFAULT '',ADD COLUMN health_status text NOT NULL DEFAULT '',ADD COLUMN last_check_reason text NOT NULL DEFAULT '',ADD COLUMN last_checked_at timestamptz,ADD COLUMN token text NOT NULL DEFAULT '',DROP COLUMN access_key,DROP COLUMN online,DROP COLUMN last_seen_at,DROP COLUMN version; + INSERT INTO gateway (name, endpoint, token) VALUES ('legacy-gw', 'http://legacy:8081', 'legacy-token') ON CONFLICT (name) DO NOTHING; + DROP TABLE IF EXISTS gateway_task; DELETE FROM schema_migration WHERE version=1058`); err != nil { + t.Fatal(err) + } + if err := store.migrate(ctx); err != nil { + t.Fatalf("upgrade: %v", err) + } + assertColumns() + var key string + var name string + if err := store.db.QueryRowContext(ctx, `SELECT name, access_key FROM gateway WHERE name='legacy-gw'`).Scan(&name, &key); err != nil { + t.Fatalf("legacy gateway read: %v", err) + } + if len(key) != 64 { + t.Fatalf("legacy gateway access key not provisioned: len=%d", len(key)) + } + var online bool + if err := store.db.QueryRowContext(ctx, `SELECT online FROM gateway WHERE name='legacy-gw'`).Scan(&online); err != nil || online { + t.Fatalf("legacy gateway online=%v err=%v", online, err) + } + // 重复启动(幂等):schema 与既有数据不变。 + for i := 0; i < 2; i++ { + if err := store.migrate(ctx); err != nil { + t.Fatalf("repeat startup: %v", err) + } + } + if err := store.db.QueryRowContext(ctx, `SELECT access_key FROM gateway WHERE name='legacy-gw'`).Scan(&key); err != nil || len(key) != 64 { + t.Fatalf("restart changed access key: %s err=%v", key, err) + } + var applied int + if err := store.db.QueryRowContext(ctx, `SELECT count(*) FROM schema_migration WHERE version=1058`).Scan(&applied); err != nil || applied != 1 { + t.Fatalf("applied=%d err=%v", applied, err) + } +} diff --git a/internal/environment/migrations/1058_gateway_outbound_channel.sql b/internal/environment/migrations/1058_gateway_outbound_channel.sql index 18d7b1a..e692ed8 100644 --- a/internal/environment/migrations/1058_gateway_outbound_channel.sql +++ b/internal/environment/migrations/1058_gateway_outbound_channel.sql @@ -6,7 +6,8 @@ ALTER TABLE gateway DROP COLUMN health_status, DROP COLUMN last_check_reason, DROP COLUMN last_checked_at, - ADD COLUMN access_key text NOT NULL, + DROP COLUMN token, + ADD COLUMN access_key text NOT NULL DEFAULT '', ADD COLUMN online boolean NOT NULL DEFAULT false, ADD COLUMN last_seen_at timestamptz, ADD COLUMN version text NOT NULL DEFAULT ''; @@ -14,6 +15,7 @@ ALTER TABLE gateway -- 既有网关记录补一个一次性访问密钥(随机 48 位十六进制),新网关由平台生成。 UPDATE gateway SET access_key = encode(sha256(('creator-hub-gateway-' || id::text || random()::text)::bytea), 'hex') WHERE access_key = ''; +ALTER TABLE gateway ALTER COLUMN access_key DROP DEFAULT; -- 统一任务模型:所有平台→网关通信都是持久化任务;全量任务即执行记录。 CREATE TABLE gateway_task ( diff --git a/internal/environment/migrations_immutable_test.go b/internal/environment/migrations_immutable_test.go index 95cc5b2..8d45b0c 100644 --- a/internal/environment/migrations_immutable_test.go +++ b/internal/environment/migrations_immutable_test.go @@ -34,6 +34,11 @@ func TestMigrationFilesAreImmutableAfterIntroduction(t *testing.T) { t.Skipf("no introducing commit found for %s", repoFile) } firstCommit := strings.SplitN(introduced, "\n", 2)[0] + // 引入 commit 尚未进入主线的迁移(同 PR 内迭代)允许修正内容; + // 合入主线后此跳过失效,规则恢复锁死。 + if gitOutput(t, "merge-base", "--is-ancestor", firstCommit, "origin/main") == "" { + t.Skipf("%s introduced in %s which is not merged into origin/main; PR-local migrations may still be revised", repoFile, firstCommit) + } original := gitOutput(t, "show", firstCommit+":"+repoFile) current, err := os.ReadFile(file) if err != nil {