package environment import ( "context" "database/sql" "fmt" "net/url" "os" "strings" "testing" "time" "git.ipao.vip/rogee/creator-hub/internal/account" ) func TestUnifiedAccountMigration(t *testing.T) { databaseURL := os.Getenv("CREATORHUB_POSTGRES_TEST_URL") if databaseURL == "" { t.Skip("set CREATORHUB_POSTGRES_TEST_URL to run PostgreSQL integration coverage") } t.Run("fresh database", func(t *testing.T) { ctx := context.Background() testURL := isolatedDatabaseURL(t, databaseURL) store := openFullyMigratedHub(t, ctx, testURL) store.Close() db, err := sql.Open("pgx", testURL) if err != nil { t.Fatal(err) } defer db.Close() assertDatabaseCount(t, db, `SELECT count(*) FROM schema_migration`, 50) assertDatabaseCount(t, db, `SELECT count(*) FROM information_schema.tables WHERE table_schema = current_schema() AND table_name IN ('social_account', 'browser_env', 'network_exit', 'environment_binding')`, 3) assertDatabaseCount(t, db, `SELECT count(*) FROM information_schema.tables WHERE table_schema = current_schema() AND table_name = 'browser_image'`, 0) assertDatabaseCount(t, db, `SELECT count(*) FROM information_schema.tables WHERE table_schema = current_schema() AND table_name = 'browser_version'`, 0) assertDatabaseCount(t, db, `SELECT count(*) FROM information_schema.columns WHERE table_schema = current_schema() AND table_name = 'social_account' AND column_name = 'cookies'`, 0) assertDatabaseCount(t, db, `SELECT count(*) FROM information_schema.columns WHERE table_schema = current_schema() AND table_name = 'browser_env' AND column_name = 'runtime_cleanup_pending'`, 1) assertDatabaseCount(t, db, `SELECT count(*) FROM information_schema.columns WHERE table_schema = current_schema() AND table_name = 'network_exit' AND column_name IN ('username', 'password')`, 2) assertDatabaseCount(t, db, `SELECT count(*) FROM information_schema.columns WHERE table_schema = current_schema() AND table_name = 'network_exit' AND column_name = 'credential_reference_id'`, 0) store = openFullyMigratedHub(t, ctx, testURL) store.Close() assertDatabaseCount(t, db, `SELECT count(*) FROM schema_migration`, 50) }) t.Run("legacy migration 013 without account secrets is repaired forward", func(t *testing.T) { ctx := context.Background() testURL := isolatedDatabaseURL(t, databaseURL) db := openLegacyAccountCreationSchema(t, ctx, testURL) defer db.Close() store, err := Open(ctx, testURL) if err != nil { t.Fatal(err) } store.Close() assertDatabaseCount(t, db, `SELECT count(*) FROM schema_migration WHERE version = 14`, 1) assertDatabaseCount(t, db, `SELECT count(*) FROM information_schema.columns WHERE table_schema = current_schema() AND table_name = 'social_account' AND column_name = 'cookies'`, 0) assertDatabaseCount(t, db, `SELECT count(*) FROM information_schema.columns WHERE table_schema = current_schema() AND table_name = 'social_account' AND column_name = 'credential_key'`, 1) }) t.Run("legacy migration 013 with account secrets blocks upgrade", func(t *testing.T) { ctx := context.Background() testURL := isolatedDatabaseURL(t, databaseURL) db := openLegacyAccountCreationSchema(t, ctx, testURL) defer db.Close() if _, err := db.Exec(` INSERT INTO social_account (id, name, platform, platform_account_key, tags, cookies, authorization_kind, authorization_status, status) VALUES ('legacy-secret', 'Legacy', 'douyin', 'legacy-secret', ARRAY[]::text[], 'sessionid=migration-secret', 'owned', 'authorized', 'paused')`); err != nil { t.Fatal(err) } _, err := Open(ctx, testURL) if err == nil || !strings.Contains(err.Error(), "apply environment schema migration 14") || strings.Contains(err.Error(), "migration-secret") { t.Fatalf("unsafe legacy migration was not blocked safely: %v", err) } assertDatabaseCount(t, db, `SELECT count(*) FROM schema_migration WHERE version = 14`, 0) }) t.Run("v1 and v2 data", func(t *testing.T) { ctx := context.Background() testURL := isolatedDatabaseURL(t, databaseURL) openLegacyPhaseASchema(t, ctx, testURL) db, err := sql.Open("pgx", testURL) if err != nil { t.Fatal(err) } defer db.Close() if _, err := db.Exec(migration002); err != nil { t.Fatal(err) } if _, err := db.Exec(` INSERT INTO schema_migration (version) VALUES (2); ALTER TABLE browser_image DROP CONSTRAINT browser_image_image_ref_check; INSERT INTO credential_reference (id, provider, reference_key) VALUES ('credential-mapped', 'os_keyring', 'creatorhub/mapped'), ('credential-unbound', 'secret_manager', 'creatorhub/unbound'); INSERT INTO social_account (id, credential_reference_id, profile_id, status) VALUES ('mapped', 'credential-mapped', 'legacy-profile-mapped', 'active'), ('unbound', 'credential-unbound', 'legacy-profile-unbound', 'active'); INSERT INTO gateway (name, endpoint, token) VALUES ('legacy-gateway', 'http://127.0.0.1:8081', 'legacy-gateway-token'); INSERT INTO browser_image (version, image_ref) VALUES ('1', '/opt/creatorhub/browsers/1'); INSERT INTO browser_env (alias, name, gateway_name, image_version, fingerprint) VALUES ('mapped', 'Mapped', 'legacy-gateway', '1', '{"seed":1,"proxy_server":"http://legacy:secret@proxy.example:8080","disable_non_proxied_udp":true}'), ('orphan-env', 'Orphan', 'legacy-gateway', '1', '{"seed":2}'); INSERT INTO runtime_instance (id, account_id, runtime_id, lease_until) VALUES ('instance-mapped', 'mapped', 'runtime-mapped', now() + interval '1 hour'), ('instance-unbound', 'unbound', 'runtime-unbound', now() + interval '1 hour'); INSERT INTO audit_event (event_type, account_id) VALUES ('legacy_event', 'mapped')`); err != nil { t.Fatal(err) } store, err := Open(ctx, testURL) if err != nil { t.Fatal(err) } store.Close() // 043 收敛清空账号域:legacy 绑定/实例数据清空,验证收敛形态与约束仍在 assertDatabaseCount(t, db, `SELECT count(*) FROM social_account`, 0) assertDatabaseCount(t, db, `SELECT count(*) FROM browser_env`, 0) assertDatabaseCount(t, db, `SELECT count(*) FROM audit_event`, 0) assertDatabaseCount(t, db, `SELECT count(*) FROM schema_migration WHERE version IN (3, 4, 5, 6, 7, 8, 9)`, 7) if _, err := db.Exec(` INSERT INTO social_account (account_id, credential_provider, credential_key, platform, platform_account_key, authorization_kind, authorization_status) VALUES ('mapped', 'os_keyring', 'creatorhub/mapped', 'mock', 'mapped', 'owned', 'authorized')`); err != nil { t.Fatal(err) } if _, err := db.Exec(` INSERT INTO social_account (account_id, credential_provider, credential_key, platform, platform_account_key, authorization_kind, authorization_status) VALUES ('duplicate', 'os_keyring', 'creatorhub/duplicate', 'mock', 'mapped', 'owned', 'authorized')`); err == nil { t.Fatal("duplicate platform account must fail") } }) } func openFullyMigratedHub(t *testing.T, ctx context.Context, databaseURL string) *Store { t.Helper() phaseAStore, err := account.Open(ctx, databaseURL) if err != nil { t.Fatal(err) } phaseAStore.Close() store, err := Open(ctx, databaseURL) if err != nil { t.Fatal(err) } return store } func isolatedDatabaseURL(t *testing.T, databaseURL string) string { t.Helper() admin, err := sql.Open("pgx", databaseURL) if err != nil { t.Fatal(err) } t.Cleanup(func() { admin.Close() }) schema := fmt.Sprintf("creatorhub_hh804_%d", time.Now().UnixNano()) if _, err := admin.Exec("CREATE SCHEMA " + schema); err != nil { t.Fatal(err) } t.Cleanup(func() { if _, err := admin.Exec("DROP SCHEMA " + schema + " CASCADE"); err != nil { t.Errorf("drop test schema: %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() return parsed.String() } // openLegacyPhaseASchema 在空 schema 上重放 001 账号基座(legacy 迁移重放测试用;account.Open 已不再执行迁移)。 func openLegacyPhaseASchema(t *testing.T, ctx context.Context, databaseURL string) *sql.DB { t.Helper() db, err := sql.Open("pgx", databaseURL) if err != nil { t.Fatal(err) } if _, err := db.Exec(`CREATE TABLE IF NOT EXISTS schema_migration (version integer PRIMARY KEY, applied_at timestamptz NOT NULL DEFAULT now())`); err != nil { t.Fatal(err) } if _, err := db.Exec(migration001); err != nil { t.Fatal(err) } if _, err := db.Exec(`INSERT INTO schema_migration (version) VALUES (1)`); err != nil { t.Fatal(err) } return db } func openLegacyAccountCreationSchema(t *testing.T, ctx context.Context, databaseURL string) *sql.DB { t.Helper() openLegacyPhaseASchema(t, ctx, databaseURL) db, err := sql.Open("pgx", databaseURL) if err != nil { t.Fatal(err) } for _, migration := range []struct { version int sql string }{{2, migration002}, {3, migration003}, {4, migration004}, {5, migration005}, {6, migration006}, {7, migration007}, {8, migration008}, {9, migration009}, {10, migration010}, {11, migration011}, {12, migration012}} { if _, err := db.Exec(migration.sql); err != nil { t.Fatal(err) } if _, err := db.Exec(`INSERT INTO schema_migration (version) VALUES ($1)`, migration.version); err != nil { t.Fatal(err) } } if _, err := db.Exec(` ALTER TABLE social_account ADD COLUMN name text, ADD COLUMN tags text[], ADD COLUMN cookies text; UPDATE social_account SET name = platform_account_key, tags = ARRAY[]::text[], cookies = ''; ALTER TABLE social_account ALTER COLUMN credential_reference_id DROP NOT NULL, ALTER COLUMN name SET NOT NULL, ALTER COLUMN name SET DEFAULT '未命名账号', ALTER COLUMN tags SET NOT NULL, ALTER COLUMN tags SET DEFAULT ARRAY[]::text[], ALTER COLUMN cookies SET NOT NULL, ALTER COLUMN cookies SET DEFAULT '', ADD CONSTRAINT social_account_name_check CHECK (name = btrim(name) AND length(name) BETWEEN 1 AND 128), ADD CONSTRAINT social_account_tags_check CHECK (cardinality(tags) <= 20), ADD CONSTRAINT social_account_cookies_length_check CHECK (length(cookies) <= 8192); INSERT INTO schema_migration (version) VALUES (13)`); err != nil { t.Fatal(err) } return db } func assertDatabaseCount(t *testing.T, db *sql.DB, query string, want int, args ...any) { t.Helper() var got int if err := db.QueryRow(query, args...).Scan(&got); err != nil || got != want { t.Fatalf("count mismatch: got=%d want=%d err=%v query=%s", got, want, err, query) } }