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
This commit is contained in:
@@ -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)
|
||||
}
|
||||
}
|
||||
@@ -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 (
|
||||
|
||||
@@ -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 {
|
||||
|
||||
Reference in New Issue
Block a user