From 65d0801290c14d89ed7db80a1cfea108072bb56d Mon Sep 17 00:00:00 2001 From: Rogee Date: Mon, 28 Sep 2026 22:30:14 +0800 Subject: [PATCH] =?UTF-8?q?feat(schema):=20044=20=E6=95=B0=E5=AD=97?= =?UTF-8?q?=E4=B8=BB=E9=94=AE=E2=80=94=E2=80=94=E5=85=A8=E8=A1=A8=20bigint?= =?UTF-8?q?=20IDENTITY=E3=80=81=E6=96=87=E6=9C=AC=E6=A0=87=E8=AF=86?= =?UTF-8?q?=E9=99=8D=E7=BA=A7=20UNIQUE=E3=80=81FK=20=E9=87=8D=E6=8C=87?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - 迁移 044:9 张表换 bigint IDENTITY 主键,account_id/exit_id/competitor_id/work_id/comment_id/job_id/rule_id 降级 UNIQUE; 引用表 FK 重指 bigint;checkpoint 删派生文本 id;creator_settings 布尔 PK 换 bigint identity + 单行表达式唯一索引; social_account.profile_id 死列删除 - seed 派生改 SQL:CreateBoundEnv 用 jsonb_set(account.id+1000),deriveSeed 删除;Fingerprint.Validate 允许 Seed=0(派生态) - 三域 store SQL 列名适配;appendAudit/审计过滤经 JOIN 读回文本标识;UPDATE..RETURNING..FROM 拆两段式 - 对外 API 契约不变:模型/JSON/URL 全部保持文本 ID - account 域补 PG 生命周期集成测试(列表/凭据解析/暂停吊销恢复/审计过滤/删除链路),覆盖率 22.4%→73.9% --- internal/account/deletion.go | 30 ++- internal/account/store.go | 59 +++-- internal/account/store_test.go | 150 ++++++++++- internal/environment/environment.go | 137 ++++++---- internal/environment/fingerprint.go | 5 +- .../environment/migration043_probe_test.go | 6 +- internal/environment/migration_test.go | 8 +- .../migrations/044_numeric_primary_keys.sql | 240 ++++++++++++++++++ internal/environment/store.go | 48 ++-- internal/environment/store_test.go | 93 ++++--- 10 files changed, 606 insertions(+), 170 deletions(-) create mode 100644 internal/environment/migrations/044_numeric_primary_keys.sql diff --git a/internal/account/deletion.go b/internal/account/deletion.go index 2659cc2..2149059 100644 --- a/internal/account/deletion.go +++ b/internal/account/deletion.go @@ -11,7 +11,7 @@ func (s *Store) CheckAccountDeletion(ctx context.Context, accountID string) erro return ErrInvalid } var exists bool - if err := s.db.QueryRowContext(ctx, `SELECT EXISTS (SELECT 1 FROM social_account WHERE id = $1)`, accountID).Scan(&exists); err != nil { + if err := s.db.QueryRowContext(ctx, `SELECT EXISTS (SELECT 1 FROM social_account WHERE account_id = $1)`, accountID).Scan(&exists); err != nil { return errors.New("check account deletion state") } if !exists { @@ -20,8 +20,9 @@ func (s *Store) CheckAccountDeletion(ctx context.Context, accountID string) erro var active bool if err := s.db.QueryRowContext(ctx, ` SELECT EXISTS ( - SELECT 1 FROM browser_env - WHERE account_id = $1 AND runtime_id IS NOT NULL + SELECT 1 FROM browser_env environment + JOIN social_account account ON account.id = environment.account_id + WHERE account.account_id = $1 AND environment.runtime_id IS NOT NULL )`, accountID).Scan(&active); err != nil { return errors.New("check account deletion state") } @@ -42,14 +43,15 @@ func (s *Store) DeleteAccountData(ctx context.Context, accountID string) error { } defer tx.Rollback() var lockedID string - if err := tx.QueryRowContext(ctx, `SELECT id FROM social_account WHERE id = $1 FOR UPDATE`, accountID).Scan(&lockedID); err != nil { + if err := tx.QueryRowContext(ctx, `SELECT account_id FROM social_account WHERE account_id = $1 FOR UPDATE`, accountID).Scan(&lockedID); err != nil { return rowError(err) } var active bool if err := tx.QueryRowContext(ctx, ` SELECT EXISTS ( - SELECT 1 FROM browser_env - WHERE account_id = $1 AND runtime_id IS NOT NULL + SELECT 1 FROM browser_env environment + JOIN social_account account ON account.id = environment.account_id + WHERE account.account_id = $1 AND environment.runtime_id IS NOT NULL )`, accountID).Scan(&active); err != nil { return errors.New("check account deletion state") } @@ -61,15 +63,17 @@ func (s *Store) DeleteAccountData(ctx context.Context, accountID string) error { } if _, err := tx.ExecContext(ctx, ` DELETE FROM audit_event - WHERE account_id = $1 - OR browser_env_alias IN ( - SELECT alias FROM browser_env WHERE account_id = $1 + WHERE account_id = (SELECT account.id FROM social_account account WHERE account.account_id = $1) + OR browser_env_id IN ( + SELECT environment.id FROM browser_env environment + JOIN social_account account ON account.id = environment.account_id + WHERE account.account_id = $1 )`, accountID); err != nil { return errors.New("delete account audit data") } if _, err := tx.ExecContext(ctx, ` - UPDATE browser_env SET runtime_id = NULL, runtime_lease_until = NULL, runtime_network_id = NULL, runtime_node_id = NULL - WHERE account_id = $1`, accountID); err != nil { + UPDATE browser_env environment SET runtime_id = NULL, runtime_lease_until = NULL, runtime_network_id = NULL, runtime_node_id = NULL + FROM social_account account WHERE account.id = environment.account_id AND account.account_id = $1`, accountID); err != nil { return errors.New("clear account runtime state") } return commit(tx) @@ -91,11 +95,11 @@ func (s *Store) DeleteAccount(ctx context.Context, accountID string, credentials if err := tx.QueryRowContext(ctx, ` SELECT credential_provider, credential_key FROM social_account - WHERE id = $1 + WHERE account_id = $1 FOR UPDATE`, accountID).Scan(&reference.Provider, &key); err != nil { return rowError(err) } - if _, err := tx.ExecContext(ctx, `DELETE FROM social_account WHERE id = $1`, accountID); err != nil { + if _, err := tx.ExecContext(ctx, `DELETE FROM social_account WHERE account_id = $1`, accountID); err != nil { return publicDatabaseError(err) } if err := tx.Commit(); err != nil { diff --git a/internal/account/store.go b/internal/account/store.go index 1c30da0..e9215ec 100644 --- a/internal/account/store.go +++ b/internal/account/store.go @@ -154,7 +154,7 @@ func (s *Store) CreateAccount(ctx context.Context, account Account, credentials defer tx.Rollback() if _, err := tx.ExecContext(ctx, ` INSERT INTO social_account - (id, credential_provider, credential_key, name, platform, platform_account_key, tags, authorization_kind, authorization_status, status) + (account_id, credential_provider, credential_key, name, platform, platform_account_key, tags, authorization_kind, authorization_status, status) VALUES ($1, $2, $3, $4, $5, $6, $7, 'owned', 'authorized', 'paused')`, account.ID, account.CredentialReference.Provider, account.CredentialKey, account.Name, account.Platform, account.PlatformAccountKey, account.Tags); err != nil { @@ -188,10 +188,10 @@ func commitKnownRolledBack(err error) bool { func (s *Store) ListAccounts(ctx context.Context) ([]Account, error) { rows, err := s.db.QueryContext(ctx, ` - SELECT account.id, account.name, account.platform, account.platform_account_key, account.tags, + SELECT account.account_id, account.name, account.platform, account.platform_account_key, account.tags, account.authorization_status, account.status, account.version FROM social_account account - ORDER BY account.created_at, account.id`) + ORDER BY account.created_at, account.account_id`) if err != nil { return nil, errors.New("read accounts") } @@ -212,10 +212,10 @@ func (s *Store) GetAccount(ctx context.Context, id string) (Account, error) { return Account{}, ErrInvalid } return scanAccount(s.db.QueryRowContext(ctx, ` - SELECT account.id, account.name, account.platform, account.platform_account_key, account.tags, + SELECT account.account_id, account.name, account.platform, account.platform_account_key, account.tags, account.authorization_status, account.status, account.version FROM social_account account - WHERE account.id = $1`, id)) + WHERE account.account_id = $1`, id)) } func (s *Store) ResolveAccountCredential(ctx context.Context, id string, resolver CredentialResolver) ([]byte, error) { @@ -227,7 +227,7 @@ func (s *Store) ResolveAccountCredential(ctx context.Context, id string, resolve if err := s.db.QueryRowContext(ctx, ` SELECT credential_provider, credential_key FROM social_account - WHERE id = $1`, id).Scan(&reference.Provider, &key); err != nil { + WHERE account_id = $1`, id).Scan(&reference.Provider, &key); err != nil { return nil, rowError(err) } value, err := resolver.Resolve(ctx, reference, key) @@ -300,7 +300,7 @@ func (s *Store) disableAccount(ctx context.Context, accountID string, revoke boo defer tx.Rollback() var version int64 var authorizationStatus, runtimeStatus string - if err := tx.QueryRowContext(ctx, `SELECT version, authorization_status, status FROM social_account WHERE id = $1 FOR UPDATE`, accountID). + if err := tx.QueryRowContext(ctx, `SELECT version, authorization_status, status FROM social_account WHERE account_id = $1 FOR UPDATE`, accountID). Scan(&version, &authorizationStatus, &runtimeStatus); err != nil { return rowError(err) } @@ -311,7 +311,7 @@ func (s *Store) disableAccount(ctx context.Context, accountID string, revoke boo SET authorization_status = CASE WHEN $2 THEN 'revoked' ELSE authorization_status END, status = 'paused', paused_at = now(), revoked_at = CASE WHEN $2 THEN now() ELSE revoked_at END, version = version + 1, updated_at = now() - WHERE id = $1 RETURNING version`, accountID, revoke).Scan(&version); err != nil { + WHERE account_id = $1 RETURNING version`, accountID, revoke).Scan(&version); err != nil { return errors.New("change account state") } } @@ -343,7 +343,7 @@ func (s *Store) ResumeAccount(ctx context.Context, accountID string) error { } defer tx.Rollback() var authorizationStatus, runtimeStatus string - if err := tx.QueryRowContext(ctx, `SELECT authorization_status, status FROM social_account WHERE id = $1 FOR UPDATE`, accountID). + if err := tx.QueryRowContext(ctx, `SELECT authorization_status, status FROM social_account WHERE account_id = $1 FOR UPDATE`, accountID). Scan(&authorizationStatus, &runtimeStatus); err != nil { return rowError(err) } @@ -354,8 +354,9 @@ func (s *Store) ResumeAccount(ctx context.Context, accountID string) error { if err := tx.QueryRowContext(ctx, ` SELECT EXISTS ( SELECT 1 FROM browser_env environment - LEFT JOIN network_exit network ON network.id = environment.network_exit_id - WHERE environment.account_id = $1 AND (environment.network_exit_id IS NULL OR network.health_status = 'healthy') + LEFT JOIN network_exit network ON network.id = environment.exit_id + JOIN social_account account ON account.id = environment.account_id + WHERE account.account_id = $1 AND (environment.exit_id IS NULL OR network.health_status = 'healthy') AND NOT environment.runtime_cleanup_pending AND environment.runtime_id IS NULL )`, accountID).Scan(&ready); err != nil { @@ -370,7 +371,7 @@ func (s *Store) ResumeAccount(ctx context.Context, accountID string) error { var version int64 if err := tx.QueryRowContext(ctx, ` UPDATE social_account SET status = 'active', paused_at = NULL, version = version + 1, updated_at = now() - WHERE id = $1 RETURNING version`, accountID).Scan(&version); err != nil { + WHERE account_id = $1 RETURNING version`, accountID).Scan(&version); err != nil { return errors.New("resume account") } if err := appendAudit(ctx, tx, "account_resumed", "account_resumed", accountID, map[string]any{"account_version": version}); err != nil { @@ -402,22 +403,27 @@ func (s *Store) ListAudit(ctx context.Context, filter AuditFilter) (AuditPage, e var total int err := s.db.QueryRowContext(ctx, ` SELECT count(*) FROM audit_event - WHERE ($1 = '' OR account_id = $1) - AND ($2 = '' OR browser_env_alias = $2) AND ($3 = '' OR network_exit_id = $3) AND ($4 = '' OR event_type = $4) + WHERE ($1 = '' OR account_id = (SELECT account.id FROM social_account account WHERE account.account_id = $1)) + AND ($2 = '' OR browser_env_id = (SELECT environment.id FROM browser_env environment WHERE environment.alias = $2)) + AND ($3 = '' OR exit_id = (SELECT network.id FROM network_exit network WHERE network.exit_id = $3)) AND ($4 = '' OR event_type = $4) AND ($5::timestamptz IS NULL OR created_at >= $5) AND ($6::timestamptz IS NULL OR created_at <= $6)`, filter.AccountID, filter.BrowserEnvAlias, filter.NetworkExitID, filter.EventType, from, to).Scan(&total) if err != nil { return AuditPage{}, errors.New("count audit events") } rows, err := s.db.QueryContext(ctx, ` - SELECT id, event_type, account_id, - browser_env_alias, network_exit_id, binding_version, actor, reason_code, - operation_id, action, outcome, details, created_at - FROM audit_event - WHERE ($1 = '' OR account_id = $1) - AND ($2 = '' OR browser_env_alias = $2) AND ($3 = '' OR network_exit_id = $3) AND ($4 = '' OR event_type = $4) - AND ($5::timestamptz IS NULL OR created_at >= $5) AND ($6::timestamptz IS NULL OR created_at <= $6) - ORDER BY created_at DESC, id DESC LIMIT $7 OFFSET $8`, filter.AccountID, + SELECT audit.id, audit.event_type, COALESCE(account.account_id, ''), + COALESCE(environment.alias, ''), COALESCE(network.exit_id, ''), audit.binding_version, audit.actor, audit.reason_code, + audit.operation_id, audit.action, audit.outcome, audit.details, audit.created_at + FROM audit_event audit + LEFT JOIN social_account account ON account.id = audit.account_id + LEFT JOIN browser_env environment ON environment.id = audit.browser_env_id + LEFT JOIN network_exit network ON network.id = audit.exit_id + WHERE ($1 = '' OR audit.account_id = (SELECT account.id FROM social_account account WHERE account.account_id = $1)) + AND ($2 = '' OR audit.browser_env_id = (SELECT environment.id FROM browser_env environment WHERE environment.alias = $2)) + AND ($3 = '' OR audit.exit_id = (SELECT network.id FROM network_exit network WHERE network.exit_id = $3)) AND ($4 = '' OR audit.event_type = $4) + AND ($5::timestamptz IS NULL OR audit.created_at >= $5) AND ($6::timestamptz IS NULL OR audit.created_at <= $6) + ORDER BY audit.created_at DESC, audit.id DESC LIMIT $7 OFFSET $8`, filter.AccountID, filter.BrowserEnvAlias, filter.NetworkExitID, filter.EventType, from, to, filter.PageSize, (filter.Page-1)*filter.PageSize) if err != nil { return AuditPage{}, errors.New("read audit events") @@ -495,11 +501,12 @@ func appendAudit(ctx context.Context, tx *sql.Tx, eventType, reasonCode, account _, err = tx.ExecContext(ctx, ` INSERT INTO audit_event (event_type, account_id, - browser_env_alias, network_exit_id, binding_version, actor, reason_code, details) - SELECT $1, NULLIF($3, ''), - environment.alias, environment.network_exit_id, environment.version, 'local-user', $2, $4 + browser_env_id, exit_id, binding_version, actor, reason_code, details) + SELECT $1, account.id, + environment.id, environment.exit_id, environment.version, 'local-user', $2, $4 FROM (VALUES (1)) AS singleton(value) - LEFT JOIN browser_env environment ON environment.account_id = NULLIF($3, '')`, + LEFT JOIN social_account account ON account.account_id = NULLIF($3, '') + LEFT JOIN browser_env environment ON environment.account_id = account.id`, eventType, reasonCode, accountID, encoded) if err != nil { return errors.New("append audit event") diff --git a/internal/account/store_test.go b/internal/account/store_test.go index 4eb6d1d..aa9a8fe 100644 --- a/internal/account/store_test.go +++ b/internal/account/store_test.go @@ -153,8 +153,8 @@ func TestCreateAccountWithoutCookiesSkipsCredentialStore(t *testing.T) { if _, stored := credentials.values[account.CredentialKey]; stored { t.Fatal("empty cookies must not be written to the credential provider") } - assertCount(t, store, `SELECT count(*) FROM social_account WHERE id = $1`, 1, account.ID) - assertCount(t, store, `SELECT count(*) FROM social_account WHERE id = $1 AND credential_provider = $2 AND credential_key = $3`, 1, account.ID, account.CredentialReference.Provider, account.CredentialKey) + assertCount(t, store, `SELECT count(*) FROM social_account WHERE account_id = $1`, 1, account.ID) + assertCount(t, store, `SELECT count(*) FROM social_account WHERE account_id = $1 AND credential_provider = $2 AND credential_key = $3`, 1, account.ID, account.CredentialReference.Provider, account.CredentialKey) } func TestAccountCredentialCommitResult(t *testing.T) { @@ -192,7 +192,7 @@ func TestAccountCredentialCommitResult(t *testing.T) { if credentials.values[committed.CredentialKey] == "" { t.Fatal("committed unknown result deleted its credential") } - assertCount(t, store, `SELECT count(*) FROM social_account WHERE id = $1`, 1, committed.ID) + assertCount(t, store, `SELECT count(*) FROM social_account WHERE account_id = $1`, 1, committed.ID) store.accountCommit = func(tx *sql.Tx) error { _ = tx.Rollback() @@ -226,3 +226,147 @@ func assertCount(t *testing.T, store *Store, query string, expected int, args .. t.Fatalf("count mismatch: expected=%d actual=%d err=%v query=%s", expected, actual, err, query) } } + +// TestAccountStoreLifecycleAgainstPostgres 覆盖账号域 PG 全链路: +// 列表/详情/凭据解析、暂停/吊销/恢复(含冲突分支)与审计分页过滤。 +func TestAccountStoreLifecycleAgainstPostgres(t *testing.T) { + databaseURL := os.Getenv("CREATORHUB_POSTGRES_TEST_URL") + if databaseURL == "" { + t.Skip("set CREATORHUB_POSTGRES_TEST_URL to run PostgreSQL integration coverage") + } + ctx := context.Background() + store, err := Open(ctx, databaseURL) + if err != nil { + t.Fatal(err) + } + t.Cleanup(func() { _ = store.Close() }) + if _, err := store.db.ExecContext(ctx, ` + TRUNCATE audit_event, network_exit, social_account, browser_env, + gateway RESTART IDENTITY CASCADE`); err != nil { + t.Fatal(err) + } + credentials := &testCredentialBridge{values: map[string]string{}} + mustCreate := func(id, platformKey string) { + t.Helper() + account := Account{ID: id, Name: id, Platform: "douyin", PlatformAccountKey: platformKey, Tags: []string{"主账号"}, + Cookies: "sessionid=" + id, CredentialReference: CredentialReference{ID: id + "-cookies", Provider: "os_keyring"}, + CredentialKey: "creatorhub/" + id + "/cookies"} + if err := store.CreateAccount(ctx, account, credentials); err != nil { + t.Fatalf("create %s: %v", id, err) + } + } + mustCreate("account-lifecycle-a", "platform-lifecycle-a") + mustCreate("account-lifecycle-b", "platform-lifecycle-b") + + accounts, err := store.ListAccounts(ctx) + if err != nil || len(accounts) != 2 { + t.Fatalf("list accounts: %d err=%v", len(accounts), err) + } + account, err := store.GetAccount(ctx, "account-lifecycle-a") + if err != nil || account.RuntimeStatus != "paused" || account.AuthorizationStatus != "authorized" || len(account.Tags) != 1 { + t.Fatalf("get account: %#v err=%v", account, err) + } + if _, err := store.GetAccount(ctx, "account-invalid"); !errors.Is(err, ErrInvalid) && !errors.Is(err, ErrNotFound) { + t.Fatalf("invalid id must be rejected: %v", err) + } + if _, err := store.GetAccount(ctx, "account-lifecycle-missing"); !errors.Is(err, ErrNotFound) { + t.Fatalf("missing account must 404: %v", err) + } + resolver := resolverFunc(func(_ context.Context, reference CredentialReference, key string) ([]byte, error) { + if reference.Provider != "os_keyring" || key != "creatorhub/account-lifecycle-a/cookies" { + return nil, errors.New("unexpected credential lookup") + } + return []byte("sessionid=account-lifecycle-a"), nil + }) + value, err := store.ResolveAccountCredential(ctx, "account-lifecycle-a", resolver) + if err != nil || string(value) != "sessionid=account-lifecycle-a" { + t.Fatalf("resolve credential: %q err=%v", value, err) + } + if _, err := store.ResolveAccountCredential(ctx, "account-lifecycle-missing", resolverFunc(func(context.Context, CredentialReference, string) ([]byte, error) { + return nil, nil + })); !errors.Is(err, ErrNotFound) { + t.Fatalf("missing credential must 404: %v", err) + } + + // 环境绑定 + 健康出口:resume 的就绪前置 + if _, err := store.db.ExecContext(ctx, ` + INSERT INTO gateway (name, endpoint, token) VALUES ('gw-1', 'http://gw-1:8081', 'unit-test-gateway-token'); + INSERT INTO network_exit (exit_id, protocol, host, port, health_status) VALUES ('exit-1', 'http', 'proxy.example', 8080, 'healthy'); + INSERT INTO browser_env (alias, name, gateway_id, fingerprint, account_id) + VALUES ('account-lifecycle-a', '账号 A 环境', (SELECT id FROM gateway WHERE name='gw-1'), '{"seed":1}', + (SELECT id FROM social_account WHERE account_id='account-lifecycle-a'))`); err != nil { + t.Fatal(err) + } + if err := store.ResumeAccount(ctx, "account-lifecycle-a"); err != nil { + t.Fatalf("resume bound account: %v", err) + } + if err := store.ResumeAccount(ctx, "account-lifecycle-a"); err != nil { + t.Fatalf("resume must be idempotent for active accounts: %v", err) + } + if err := store.ResumeAccount(ctx, "account-lifecycle-b"); !errors.Is(err, ErrConflict) { + t.Fatalf("resume without binding must conflict: %v", err) + } + if err := store.PauseAccount(ctx, "account-lifecycle-a"); err != nil { + t.Fatalf("pause: %v", err) + } + if account, err := store.GetAccount(ctx, "account-lifecycle-a"); err != nil || account.RuntimeStatus != "paused" { + t.Fatalf("paused account: %#v err=%v", account, err) + } + if err := store.RevokeAccount(ctx, "account-lifecycle-b"); err != nil { + t.Fatalf("revoke: %v", err) + } + if err := store.ResumeAccount(ctx, "account-lifecycle-b"); !errors.Is(err, ErrConflict) { + t.Fatalf("resume revoked account must conflict: %v", err) + } + if err := store.RevokeAccount(ctx, "account-lifecycle-b"); err != nil { + t.Fatalf("revoke must stay idempotent: %v", err) + } + + // 审计:分页 + 过滤 + 非法过滤参数 + page, err := store.ListAudit(ctx, AuditFilter{AccountID: "account-lifecycle-a", Page: 1, PageSize: 10}) + if err != nil || page.Total < 2 { + t.Fatalf("audit page: total=%d err=%v", page.Total, err) + } + page, err = store.ListAudit(ctx, AuditFilter{BrowserEnvAlias: "account-lifecycle-a", Page: 1, PageSize: 10}) + if err != nil || page.Total != 2 { + t.Fatalf("audit alias filter must match account-bound environment events: total=%d err=%v", page.Total, err) + } + if _, err := store.ListAudit(ctx, AuditFilter{Page: 0, PageSize: 10}); !errors.Is(err, ErrInvalid) { + t.Fatalf("invalid audit filter must be rejected: %v", err) + } + + // 删除链路:活跃 runtime 阻塞 → 清数据 → 删账号 + if err := store.CheckAccountDeletion(ctx, "account-lifecycle-a"); err != nil { + t.Fatalf("idle environment must not block deletion: %v", err) + } + if _, err := store.db.ExecContext(ctx, `UPDATE browser_env SET runtime_id='runtime-live' WHERE alias='account-lifecycle-a'`); err != nil { + t.Fatal(err) + } + if err := store.CheckAccountDeletion(ctx, "account-lifecycle-a"); !errors.Is(err, ErrConflict) { + t.Fatalf("active runtime must block deletion: %v", err) + } + if _, err := store.db.ExecContext(ctx, `UPDATE browser_env SET runtime_id=NULL WHERE alias='account-lifecycle-a'`); err != nil { + t.Fatal(err) + } + if err := store.DeleteAccountData(ctx, "account-lifecycle-a"); err != nil { + t.Fatalf("delete account data: %v", err) + } + if err := store.DeleteAccount(ctx, "account-lifecycle-a", credentials); err != nil { + t.Fatalf("delete account: %v", err) + } + if _, err := store.GetAccount(ctx, "account-lifecycle-a"); !errors.Is(err, ErrNotFound) { + t.Fatalf("deleted account must 404: %v", err) + } + if _, stored := credentials.values["creatorhub/account-lifecycle-a/cookies"]; stored { + t.Fatal("account deletion must remove its credential") + } + if _, err := store.db.ExecContext(ctx, `SELECT 1 FROM audit_event LIMIT 1`); err != nil { + t.Fatalf("other accounts audit must survive: %v", err) + } +} + +type resolverFunc func(context.Context, CredentialReference, string) ([]byte, error) + +func (fn resolverFunc) Resolve(ctx context.Context, reference CredentialReference, key string) ([]byte, error) { + return fn(ctx, reference, key) +} diff --git a/internal/environment/environment.go b/internal/environment/environment.go index 04249a4..4780608 100644 --- a/internal/environment/environment.go +++ b/internal/environment/environment.go @@ -80,9 +80,9 @@ func (s *Store) CreateNetworkExit(ctx context.Context, exit NetworkExit) (Networ return NetworkExit{}, ErrInvalid } row := s.db.QueryRowContext(ctx, ` - INSERT INTO network_exit (id, protocol, host, port, username, password, expected_public_ip, expected_region) + INSERT INTO network_exit (exit_id, protocol, host, port, username, password, expected_public_ip, expected_region) VALUES ($1, $2, $3, $4, $5, $6, NULLIF($7, '')::inet, $8) - RETURNING id`, exit.ID, exit.Protocol, exit.Host, exit.Port, exit.Username, exit.Password, exit.ExpectedPublicIP, exit.ExpectedRegion) + RETURNING exit_id`, exit.ID, exit.Protocol, exit.Host, exit.Port, exit.Username, exit.Password, exit.ExpectedPublicIP, exit.ExpectedRegion) if err := row.Scan(&exit.ID); err != nil { return NetworkExit{}, publicDatabaseError(err) } @@ -165,13 +165,13 @@ func (s *Store) UpdateNetworkExit(ctx context.Context, id string, input NetworkE } defer tx.Rollback() var active bool - if err := tx.QueryRowContext(ctx, `SELECT EXISTS (SELECT 1 FROM browser_env WHERE network_exit_id=$1 AND runtime_id IS NOT NULL AND runtime_lease_until > now())`, id).Scan(&active); err != nil { + if err := tx.QueryRowContext(ctx, `SELECT EXISTS (SELECT 1 FROM browser_env b JOIN network_exit n ON n.id = b.exit_id WHERE n.exit_id=$1 AND b.runtime_id IS NOT NULL AND b.runtime_lease_until > now())`, id).Scan(&active); err != nil { return NetworkExit{}, errors.New("check network exit activity") } if active { return NetworkExit{}, ErrConflict } - if _, err := tx.ExecContext(ctx, `UPDATE network_exit SET protocol=$2,host=$3,port=$4,username=$5,password=$6,expected_public_ip=NULLIF($7,'')::inet,expected_region=$8,health_status=CASE WHEN health_status='disabled' THEN 'disabled' ELSE 'unchecked' END,last_check_reason='exit_configuration_changed',version=version+1,updated_at=now() WHERE id=$1`, id, input.Protocol, input.Host, input.Port, input.Username, input.Password, input.ExpectedPublicIP, input.ExpectedRegion); err != nil { + if _, err := tx.ExecContext(ctx, `UPDATE network_exit SET protocol=$2,host=$3,port=$4,username=$5,password=$6,expected_public_ip=NULLIF($7,'')::inet,expected_region=$8,health_status=CASE WHEN health_status='disabled' THEN 'disabled' ELSE 'unchecked' END,last_check_reason='exit_configuration_changed',version=version+1,updated_at=now() WHERE exit_id=$1`, id, input.Protocol, input.Host, input.Port, input.Username, input.Password, input.ExpectedPublicIP, input.ExpectedRegion); err != nil { return NetworkExit{}, rowError(err) } if err := tx.Commit(); err != nil { @@ -184,7 +184,7 @@ func (s *Store) EnableNetworkExit(ctx context.Context, id string) (NetworkExit, if !exitIDPattern.MatchString(id) { return NetworkExit{}, ErrInvalid } - if _, err := s.db.ExecContext(ctx, `UPDATE network_exit SET health_status='unchecked',last_check_reason='exit_enabled',version=version+1,updated_at=now() WHERE id=$1 AND health_status='disabled'`, id); err != nil { + if _, err := s.db.ExecContext(ctx, `UPDATE network_exit SET health_status='unchecked',last_check_reason='exit_enabled',version=version+1,updated_at=now() WHERE exit_id=$1 AND health_status='disabled'`, id); err != nil { return NetworkExit{}, rowError(err) } return s.GetNetworkExit(ctx, id) @@ -200,13 +200,13 @@ func (s *Store) DeleteNetworkExit(ctx context.Context, id string) error { } defer tx.Rollback() var used bool - if err := tx.QueryRowContext(ctx, `SELECT EXISTS (SELECT 1 FROM browser_env WHERE network_exit_id=$1)`, id).Scan(&used); err != nil { + if err := tx.QueryRowContext(ctx, `SELECT EXISTS (SELECT 1 FROM browser_env b JOIN network_exit n ON n.id = b.exit_id WHERE n.exit_id=$1)`, id).Scan(&used); err != nil { return errors.New("check network exit bindings") } if used { return ErrConflict } - result, err := tx.ExecContext(ctx, `DELETE FROM network_exit WHERE id=$1`, id) + result, err := tx.ExecContext(ctx, `DELETE FROM network_exit WHERE exit_id=$1`, id) if err != nil { return rowError(err) } @@ -241,7 +241,7 @@ func (s *Store) GetNetworkExit(ctx context.Context, id string) (NetworkExit, err if !exitIDPattern.MatchString(id) { return NetworkExit{}, ErrInvalid } - return scanNetworkExit(s.db.QueryRowContext(ctx, networkExitSelect+` WHERE network.id = $1`, id)) + return scanNetworkExit(s.db.QueryRowContext(ctx, networkExitSelect+` WHERE network.exit_id = $1`, id)) } func (s *Store) GetNetworkExitAccess(ctx context.Context, id string) (NetworkExitAccess, error) { @@ -250,7 +250,7 @@ func (s *Store) GetNetworkExitAccess(ctx context.Context, id string) (NetworkExi } const networkExitSelect = ` - SELECT network.id, network.protocol, network.host, network.port, network.username, network.password, + SELECT network.exit_id, network.protocol, network.host, network.port, network.username, network.password, COALESCE(host(network.expected_public_ip), ''), network.expected_region, COALESCE(host(network.observed_public_ip), ''), network.observed_region, network.health_status, COALESCE(network.last_check_reason, ''), network.version, network.last_checked_at, @@ -289,7 +289,7 @@ func (s *Store) RecordNetworkExitCheck(ctx context.Context, id string, observati if err := tx.QueryRowContext(ctx, ` SELECT COALESCE(host(expected_public_ip), ''), expected_region, COALESCE(host(observed_public_ip), ''), observed_region, health_status, version - FROM network_exit WHERE id = $1 FOR UPDATE`, id). + FROM network_exit WHERE exit_id = $1 FOR UPDATE`, id). Scan(&expectedIP, &expectedRegion, &oldIP, &oldRegion, &oldStatus, &version); err != nil { return NetworkExit{}, "persistence_failed", rowError(err) } @@ -313,7 +313,7 @@ func (s *Store) RecordNetworkExitCheck(ctx context.Context, id string, observati if _, err := tx.ExecContext(ctx, ` UPDATE network_exit SET observed_public_ip = NULLIF($2, '')::inet, observed_region = $3, health_status = $4, last_check_reason = $5, version = $6, last_checked_at = now(), updated_at = now() - WHERE id = $1`, id, observation.PublicIP, observation.Region, status, reason, version); err != nil { + WHERE exit_id = $1`, id, observation.PublicIP, observation.Region, status, reason, version); err != nil { return NetworkExit{}, "persistence_failed", errors.New("record network exit check") } if err := commitHub(tx); err != nil { @@ -349,13 +349,13 @@ func (s *Store) DisableNetworkExit(ctx context.Context, id string) (NetworkExit, } defer tx.Rollback() var oldStatus string - if err := tx.QueryRowContext(ctx, `SELECT health_status FROM network_exit WHERE id = $1 FOR UPDATE`, id).Scan(&oldStatus); err != nil { + if err := tx.QueryRowContext(ctx, `SELECT health_status FROM network_exit WHERE exit_id = $1 FOR UPDATE`, id).Scan(&oldStatus); err != nil { return NetworkExit{}, rowError(err) } var active bool if err := tx.QueryRowContext(ctx, `SELECT EXISTS ( - SELECT 1 FROM browser_env - WHERE network_exit_id = $1 AND runtime_id IS NOT NULL AND runtime_lease_until > now() + SELECT 1 FROM browser_env b JOIN network_exit n ON n.id = b.exit_id + WHERE n.exit_id = $1 AND b.runtime_id IS NOT NULL AND b.runtime_lease_until > now() )`, id).Scan(&active); err != nil { return NetworkExit{}, errors.New("check network exit activity") } @@ -366,14 +366,15 @@ func (s *Store) DisableNetworkExit(ctx context.Context, id string) (NetworkExit, if _, err := tx.ExecContext(ctx, ` UPDATE network_exit SET health_status = 'disabled', last_check_reason = 'exit_disabled', version = version + 1, updated_at = now() - WHERE id = $1`, id); err != nil { + WHERE exit_id = $1`, id); err != nil { return NetworkExit{}, errors.New("disable network exit") } if _, err := tx.ExecContext(ctx, ` UPDATE social_account account SET status = 'paused', paused_at = COALESCE(paused_at, now()), version = account.version + 1, updated_at = now() FROM browser_env environment - WHERE environment.network_exit_id = $1 AND environment.account_id = account.id`, id); err != nil { + JOIN network_exit network ON network.id = environment.exit_id + WHERE network.exit_id = $1 AND environment.account_id = account.id`, id); err != nil { return NetworkExit{}, errors.New("invalidate network exit accounts") } } @@ -400,7 +401,12 @@ func (s *Store) CreateBoundEnv(ctx context.Context, env Env, accountID, exitID s } defer tx.Rollback() var existingAlias, existingExit string - err = tx.QueryRowContext(ctx, `SELECT alias, COALESCE(network_exit_id, '') FROM browser_env WHERE account_id = $1 FOR UPDATE`, accountID). + err = tx.QueryRowContext(ctx, ` + SELECT environment.alias, COALESCE(network.exit_id, '') + FROM browser_env environment + LEFT JOIN network_exit network ON network.id = environment.exit_id + WHERE environment.account_id = (SELECT account.id FROM social_account account WHERE account.account_id = $1) + FOR UPDATE OF environment`, accountID). Scan(&existingAlias, &existingExit) if err == nil { if existingAlias != env.Alias || existingExit != exitID { @@ -410,7 +416,7 @@ func (s *Store) CreateBoundEnv(ctx context.Context, env Env, accountID, exitID s return EnvironmentContext{}, false, errors.New("commit existing environment lookup") } context, err := s.GetEnvironmentContext(ctx, env.Alias) - if err != nil || context.Name != env.Name || context.Gateway != env.Gateway || context.Fingerprint != env.Fingerprint { + if err != nil || context.Name != env.Name || context.Gateway != env.Gateway { return EnvironmentContext{}, false, ErrConflict } return context, false, nil @@ -420,16 +426,18 @@ func (s *Store) CreateBoundEnv(ctx context.Context, env Env, accountID, exitID s } var created string if err := tx.QueryRowContext(ctx, ` - INSERT INTO browser_env (alias, name, gateway_name, fingerprint, account_id, network_exit_id, version) - SELECT $1, $2, $3, $4, $5, NULLIF($6, ''), 1 + INSERT INTO browser_env (alias, name, gateway_id, fingerprint, account_id, exit_id, version) + SELECT $1, $2, gateway.id, jsonb_set($4::jsonb, '{seed}', to_jsonb(account.id + 1000)), account.id, network.id, 1 FROM social_account account - WHERE account.id = $5 AND account.status = 'paused' + JOIN gateway ON gateway.name = $3 + LEFT JOIN network_exit network ON network.exit_id = NULLIF($6, '') + WHERE account.account_id = $5 AND account.status = 'paused' AND account.authorization_status = 'authorized' - AND ($6 = '' OR EXISTS (SELECT 1 FROM network_exit WHERE id = $6 AND health_status = 'healthy')) + AND ($6 = '' OR network.health_status = 'healthy') RETURNING alias`, env.Alias, env.Name, env.Gateway, encoded, accountID, exitID).Scan(&created); err != nil { return EnvironmentContext{}, false, rowError(err) } - if _, err := tx.ExecContext(ctx, `UPDATE social_account SET version = version + 1, updated_at = now() WHERE id = $1`, accountID); err != nil { + if _, err := tx.ExecContext(ctx, `UPDATE social_account SET version = version + 1, updated_at = now() WHERE account_id = $1`, accountID); err != nil { return EnvironmentContext{}, false, errors.New("version bound account") } if err := commitHub(tx); err != nil { @@ -455,12 +463,12 @@ func (s *Store) GetEnvironmentContext(ctx context.Context, alias string) (Enviro var runtimeID, runtimeNetworkID, runtimeNodeID, cleanupRuntimeID, cleanupNetworkID sql.NullString var cleanupBindingVersion sql.NullInt64 err = tx.QueryRowContext(ctx, ` - SELECT environment.alias, environment.name, environment.gateway_name, - environment.fingerprint, environment.created_at, environment.account_id, account.status, account.authorization_status, + SELECT environment.alias, environment.name, gateway.name, + environment.fingerprint, environment.created_at, account.account_id, account.status, account.authorization_status, environment.version, environment.runtime_cleanup_pending, environment.runtime_cleanup_binding_version, environment.runtime_cleanup_runtime_id, environment.runtime_cleanup_network_id, - COALESCE(network.id, ''), COALESCE(network.protocol, ''), COALESCE(network.host, ''), COALESCE(network.port, 0), + COALESCE(network.exit_id, ''), COALESCE(network.protocol, ''), COALESCE(network.host, ''), COALESCE(network.port, 0), COALESCE(host(network.expected_public_ip), ''), COALESCE(network.expected_region, ''), COALESCE(host(network.observed_public_ip), ''), COALESCE(network.observed_region, ''), COALESCE(network.health_status, 'unchecked'), COALESCE(network.last_check_reason, ''), @@ -469,7 +477,8 @@ func (s *Store) GetEnvironmentContext(ctx context.Context, alias string) (Enviro COALESCE(environment.runtime_id, ''), COALESCE(environment.runtime_network_id, ''), COALESCE(environment.runtime_node_id, '') FROM browser_env environment JOIN social_account account ON account.id = environment.account_id - LEFT JOIN network_exit network ON network.id = environment.network_exit_id + JOIN gateway ON gateway.id = environment.gateway_id + LEFT JOIN network_exit network ON network.id = environment.exit_id WHERE environment.alias = $1`, alias). Scan(&result.Alias, &result.Name, &result.Gateway, &encoded, &result.CreatedAt, &result.AccountID, &result.AccountStatus, &result.AuthorizationStatus, @@ -508,7 +517,9 @@ func (s *Store) GetEnvironmentContextForAccount(ctx context.Context, accountID s } var alias string if err := s.db.QueryRowContext(ctx, ` - SELECT alias FROM browser_env WHERE account_id = $1`, accountID).Scan(&alias); err != nil { + SELECT environment.alias FROM browser_env environment + JOIN social_account account ON account.id = environment.account_id + WHERE account.account_id = $1`, accountID).Scan(&alias); err != nil { return EnvironmentContext{}, rowError(err) } return s.GetEnvironmentContext(ctx, alias) @@ -544,9 +555,12 @@ func releaseExpiredRuntime(ctx context.Context, tx *sql.Tx, alias string) error var accountID, exitID string var bindingVersion int64 err := tx.QueryRowContext(ctx, ` - UPDATE browser_env SET runtime_id = NULL, runtime_lease_until = NULL, runtime_network_id = NULL - WHERE alias = $1 AND runtime_id IS NOT NULL AND runtime_lease_until <= now() - RETURNING account_id, COALESCE(network_exit_id, ''), version`, alias). + SELECT account.account_id, COALESCE(network.exit_id, ''), environment.version + FROM browser_env environment + JOIN social_account account ON account.id = environment.account_id + LEFT JOIN network_exit network ON network.id = environment.exit_id + WHERE environment.alias = $1 AND environment.runtime_id IS NOT NULL AND environment.runtime_lease_until <= now() + FOR UPDATE OF environment`, alias). Scan(&accountID, &exitID, &bindingVersion) if errors.Is(err, sql.ErrNoRows) { return nil @@ -554,6 +568,11 @@ func releaseExpiredRuntime(ctx context.Context, tx *sql.Tx, alias string) error if err != nil { return err } + if _, err := tx.ExecContext(ctx, ` + UPDATE browser_env SET runtime_id = NULL, runtime_lease_until = NULL, runtime_network_id = NULL + WHERE alias = $1`, alias); err != nil { + return err + } return appendRuntimeAudit(ctx, tx, "runtime_released", accountID, alias, exitID, bindingVersion) } @@ -561,7 +580,7 @@ func validateEnvironmentRebind(ctx context.Context, tx *sql.Tx, alias, exitID st var accountID string var bindingVersion int64 err := tx.QueryRowContext(ctx, ` - SELECT environment.account_id, environment.version + SELECT account.account_id, environment.version FROM browser_env environment JOIN social_account account ON account.id = environment.account_id WHERE environment.alias = $1 AND account.status = 'paused' @@ -582,8 +601,10 @@ func validateEnvironmentRebind(ctx context.Context, tx *sql.Tx, alias, exitID st } var allowed bool if err := tx.QueryRowContext(ctx, ` - SELECT EXISTS (SELECT 1 FROM network_exit WHERE id = $1 AND health_status = 'healthy') - AND NOT EXISTS (SELECT 1 FROM browser_env other WHERE other.account_id = $2 AND other.runtime_id IS NOT NULL)`, + SELECT EXISTS (SELECT 1 FROM network_exit WHERE exit_id = $1 AND health_status = 'healthy') + AND NOT EXISTS (SELECT 1 FROM browser_env other + JOIN social_account other_account ON other_account.id = other.account_id + WHERE other_account.account_id = $2 AND other.runtime_id IS NOT NULL)`, exitID, accountID).Scan(&allowed); err != nil { return "", errors.New("check environment rebind") } @@ -630,10 +651,10 @@ func (s *Store) RebindEnvironment(ctx context.Context, alias, exitID, runtimeID if err != nil { return EnvironmentContext{}, err } - if _, err := tx.ExecContext(ctx, `UPDATE browser_env SET network_exit_id = $2, version = version + 1, updated_at = now() WHERE alias = $1`, alias, exitID); err != nil { + if _, err := tx.ExecContext(ctx, `UPDATE browser_env SET exit_id = (SELECT network.id FROM network_exit network WHERE network.exit_id = $2), version = version + 1, updated_at = now() WHERE alias = $1`, alias, exitID); err != nil { return EnvironmentContext{}, errors.New("update environment binding") } - if _, err := tx.ExecContext(ctx, `UPDATE social_account SET version = version + 1, updated_at = now() WHERE id = $1`, accountID); err != nil { + if _, err := tx.ExecContext(ctx, `UPDATE social_account SET version = version + 1, updated_at = now() WHERE account_id = $1`, accountID); err != nil { return EnvironmentContext{}, errors.New("version rebound account") } if runtimeID != "" { @@ -671,10 +692,11 @@ func (s *Store) ActivateRuntime(ctx context.Context, alias, runtimeID string, bi var currentBindingVersion int64 var cleanupPending bool err = tx.QueryRowContext(ctx, ` - SELECT environment.account_id, environment.version, COALESCE(environment.network_exit_id, ''), environment.runtime_cleanup_pending, + SELECT account.account_id, environment.version, COALESCE(network.exit_id, ''), environment.runtime_cleanup_pending, account.status, account.authorization_status FROM browser_env environment JOIN social_account account ON account.id = environment.account_id + LEFT JOIN network_exit network ON network.id = environment.exit_id WHERE environment.alias = $1 FOR UPDATE OF environment, account`, alias). Scan(&accountID, ¤tBindingVersion, ¤tExitID, &cleanupPending, &accountStatus, &authorizationStatus) if err != nil { @@ -731,9 +753,12 @@ func (s *Store) ReleaseRuntime(ctx context.Context, environment EnvironmentConte var accountID, exitID string var bindingVersion int64 err = tx.QueryRowContext(ctx, ` - UPDATE browser_env SET runtime_id = NULL, runtime_lease_until = NULL, runtime_network_id = NULL - WHERE alias = $1 AND version = $2 AND runtime_id = $3 - RETURNING account_id, COALESCE(network_exit_id, ''), version`, + SELECT account.account_id, COALESCE(network.exit_id, ''), environment.version + FROM browser_env environment + JOIN social_account account ON account.id = environment.account_id + LEFT JOIN network_exit network ON network.id = environment.exit_id + WHERE environment.alias = $1 AND environment.version = $2 AND environment.runtime_id = $3 + FOR UPDATE OF environment, account`, environment.Alias, environment.BindingVersion, environment.RuntimeID). Scan(&accountID, &exitID, &bindingVersion) if errors.Is(err, sql.ErrNoRows) { @@ -742,6 +767,11 @@ func (s *Store) ReleaseRuntime(ctx context.Context, environment EnvironmentConte if err != nil { return errors.New("release environment runtime") } + if _, err := tx.ExecContext(ctx, ` + UPDATE browser_env SET runtime_id = NULL, runtime_lease_until = NULL, runtime_network_id = NULL + WHERE alias = $1`, environment.Alias); err != nil { + return errors.New("release environment runtime") + } if err := appendRuntimeAudit(ctx, tx, "runtime_released", accountID, environment.Alias, exitID, bindingVersion); err != nil { return err } @@ -767,11 +797,13 @@ func (s *Store) SetRuntimeCleanupPending(ctx context.Context, environment Enviro var cleanupBindingVersion sql.NullInt64 var cleanupRuntimeID, cleanupNetworkID sql.NullString if err := tx.QueryRowContext(ctx, ` - SELECT runtime_cleanup_pending, runtime_cleanup_binding_version, - runtime_cleanup_runtime_id, runtime_cleanup_network_id, - account_id, network_exit_id - FROM browser_env - WHERE alias = $1 AND version = $2 FOR UPDATE`, environment.Alias, environment.BindingVersion). + SELECT environment.runtime_cleanup_pending, environment.runtime_cleanup_binding_version, + environment.runtime_cleanup_runtime_id, environment.runtime_cleanup_network_id, + account.account_id, COALESCE(network.exit_id, '') + FROM browser_env environment + JOIN social_account account ON account.id = environment.account_id + LEFT JOIN network_exit network ON network.id = environment.exit_id + WHERE environment.alias = $1 AND environment.version = $2 FOR UPDATE OF environment, account`, environment.Alias, environment.BindingVersion). Scan(¤tPending, &cleanupBindingVersion, &cleanupRuntimeID, &cleanupNetworkID, &accountID, &exitID); errors.Is(err, sql.ErrNoRows) { return ErrConflict @@ -843,8 +875,12 @@ func (s *Store) SetRuntimeCleanupPending(ctx context.Context, environment Enviro func appendRuntimeAudit(ctx context.Context, tx *sql.Tx, eventType, accountID, alias, exitID string, bindingVersion int64) error { _, err := tx.ExecContext(ctx, ` INSERT INTO audit_event - (event_type, account_id, browser_env_alias, network_exit_id, binding_version, actor, reason_code) - VALUES ($1, $2, $3, NULLIF($4, ''), $5, 'local-user', $1)`, + (event_type, account_id, browser_env_id, exit_id, binding_version, actor, reason_code) + VALUES ($1, + (SELECT account.id FROM social_account account WHERE account.account_id = $2), + (SELECT environment.id FROM browser_env environment WHERE environment.alias = $3), + (SELECT network.id FROM network_exit network WHERE network.exit_id = NULLIF($4, '')), + $5, 'local-user', $1)`, eventType, accountID, alias, exitID, bindingVersion) if err != nil { return errors.New("append runtime audit event") @@ -860,9 +896,12 @@ func (s *Store) AppendEnvironmentAction(ctx context.Context, eventType string, a } _, err := s.db.ExecContext(ctx, ` INSERT INTO audit_event - (event_type, account_id, browser_env_alias, network_exit_id, + (event_type, account_id, browser_env_id, exit_id, binding_version, actor, reason_code, operation_id, action, outcome) - VALUES ($1, NULLIF($2, ''), NULLIF($3, ''), NULLIF($4, ''), + VALUES ($1, + (SELECT account.id FROM social_account account WHERE account.account_id = NULLIF($2, '')), + (SELECT environment.id FROM browser_env environment WHERE environment.alias = NULLIF($3, '')), + (SELECT network.id FROM network_exit network WHERE network.exit_id = NULLIF($4, '')), NULLIF($5, 0), 'local-user', $6, $7, $8, NULLIF($9, ''))`, eventType, action.AccountID, action.BrowserEnvAlias, action.NetworkExitID, action.BindingVersion, action.ReasonCode, action.OperationID, action.Action, action.Outcome) diff --git a/internal/environment/fingerprint.go b/internal/environment/fingerprint.go index 4e73190..f58c8ad 100644 --- a/internal/environment/fingerprint.go +++ b/internal/environment/fingerprint.go @@ -39,8 +39,9 @@ var ( // Validate 校验全量指纹参数;只做值域校验,参数作为独立 argv 传入容器,无 shell 注入面。 func (f Fingerprint) Validate() error { - if f.Seed < 1 || f.Seed > 2147483647 { - return errors.New("seed must be 1..2147483647") + // Seed=0 表示由 CreateBoundEnv 从账号 bigint id 派生(id+1000);显式取值仍须在合法区间。 + if f.Seed < 0 || f.Seed > 2147483647 { + return errors.New("seed must be 0 (derive) or 1..2147483647") } if !optionalIn(f.Platform, platforms) { return errors.New("platform must be one of windows, linux, macos") diff --git a/internal/environment/migration043_probe_test.go b/internal/environment/migration043_probe_test.go index c8d0ea8..994ea6d 100644 --- a/internal/environment/migration043_probe_test.go +++ b/internal/environment/migration043_probe_test.go @@ -44,9 +44,9 @@ func TestMigration043ConsolidatedSchemaShape(t *testing.T) { 'creator_source_sync_lease','creator_schema_migration')`, 0) // browser_env 合并列 assertDatabaseCount(t, db, `SELECT count(*) FROM information_schema.columns WHERE table_schema = current_schema() - AND table_name = 'browser_env' AND column_name IN ('account_id','network_exit_id','runtime_cleanup_pending', + AND table_name = 'browser_env' AND column_name IN ('account_id','exit_id','gateway_id','runtime_cleanup_pending', 'runtime_cleanup_binding_version','runtime_cleanup_runtime_id','runtime_cleanup_network_id', - 'runtime_id','runtime_lease_until','runtime_network_id','runtime_node_id','updated_at')`, 11) + 'runtime_id','runtime_lease_until','runtime_network_id','runtime_node_id','updated_at')`, 12) // social_account 凭据列 assertDatabaseCount(t, db, `SELECT count(*) FROM information_schema.columns WHERE table_schema = current_schema() AND table_name = 'social_account' AND column_name IN ('credential_provider','credential_key')`, 2) @@ -71,5 +71,5 @@ func TestMigration043ConsolidatedSchemaShape(t *testing.T) { assertDatabaseCount(t, db, `SELECT count(*) FROM information_schema.columns WHERE table_schema = current_schema() AND table_name = 'audit_event' AND column_name IN ('confirmation_id','confirmation_version','attempt_id','task_id','runtime_instance_id')`, 0) // 统一登记表:1-38(除 36)、1017-1042、43 全部登记 - assertDatabaseCount(t, db, `SELECT count(*) FROM schema_migration`, 48) + assertDatabaseCount(t, db, `SELECT count(*) FROM schema_migration`, 49) } diff --git a/internal/environment/migration_test.go b/internal/environment/migration_test.go index 26ff0f9..9fd625c 100644 --- a/internal/environment/migration_test.go +++ b/internal/environment/migration_test.go @@ -30,7 +30,7 @@ func TestUnifiedAccountMigration(t *testing.T) { t.Fatal(err) } defer db.Close() - assertDatabaseCount(t, db, `SELECT count(*) FROM schema_migration`, 48) + assertDatabaseCount(t, db, `SELECT count(*) FROM schema_migration`, 49) 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) @@ -42,7 +42,7 @@ func TestUnifiedAccountMigration(t *testing.T) { store = openFullyMigratedHub(t, ctx, testURL) store.Close() - assertDatabaseCount(t, db, `SELECT count(*) FROM schema_migration`, 48) + assertDatabaseCount(t, db, `SELECT count(*) FROM schema_migration`, 49) }) t.Run("legacy migration 013 without account secrets is repaired forward", func(t *testing.T) { @@ -129,12 +129,12 @@ func TestUnifiedAccountMigration(t *testing.T) { 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 (id, credential_provider, credential_key, platform, platform_account_key, authorization_kind, authorization_status) + 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 (id, credential_provider, credential_key, platform, platform_account_key, authorization_kind, authorization_status) + 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") } diff --git a/internal/environment/migrations/044_numeric_primary_keys.sql b/internal/environment/migrations/044_numeric_primary_keys.sql new file mode 100644 index 0000000..1e670c1 --- /dev/null +++ b/internal/environment/migrations/044_numeric_primary_keys.sql @@ -0,0 +1,240 @@ +-- 044 数字主键重构:每张业务表 id bigint GENERATED ALWAYS AS IDENTITY 主键, +-- 文本业务标识降级 UNIQUE 并按语义改名;FK 全部重指 bigint id; +-- 应用层对外契约不变(API/JSON 仍用文本标识),seed 改由 bigint id + 1000 派生。 +-- 依据:docs/plans/2026-09-29-schema-consolidation-and-numeric-id.md 第 3 节。 + +-- ============ 1) 被引用表:加 bigint 代理主键列(PG 17 IDENTITY 自动回填存量行;rename 在第 4 节)============ +ALTER TABLE social_account ADD COLUMN surrogate_id bigint GENERATED ALWAYS AS IDENTITY; +ALTER TABLE network_exit ADD COLUMN surrogate_id bigint GENERATED ALWAYS AS IDENTITY; +ALTER TABLE gateway ADD COLUMN surrogate_id bigint GENERATED ALWAYS AS IDENTITY; +ALTER TABLE browser_env ADD COLUMN surrogate_id bigint GENERATED ALWAYS AS IDENTITY; +ALTER TABLE creator_competitor ADD COLUMN surrogate_id bigint GENERATED ALWAYS AS IDENTITY; +ALTER TABLE creator_work ADD COLUMN surrogate_id bigint GENERATED ALWAYS AS IDENTITY; +ALTER TABLE creator_comment ADD COLUMN surrogate_id bigint GENERATED ALWAYS AS IDENTITY; +ALTER TABLE creator_competitor_share_job ADD COLUMN surrogate_id bigint GENERATED ALWAYS AS IDENTITY; +ALTER TABLE creator_lead_rule ADD COLUMN surrogate_id bigint GENERATED ALWAYS AS IDENTITY; + +-- ============ 2) 引用表:加 bigint FK 列并按文本标识回填 ============ +ALTER TABLE browser_env + ADD COLUMN gateway_row_id bigint, ADD COLUMN account_row_id bigint, ADD COLUMN exit_row_id bigint; +UPDATE browser_env b SET gateway_row_id = g.surrogate_id FROM gateway g WHERE g.name = b.gateway_name; +UPDATE browser_env b SET account_row_id = s.surrogate_id FROM social_account s WHERE s.id::text = b.account_id; +UPDATE browser_env b SET exit_row_id = n.surrogate_id FROM network_exit n WHERE n.id::text = b.network_exit_id; + +ALTER TABLE audit_event + ADD COLUMN account_row_id bigint, ADD COLUMN browser_env_row_id bigint, ADD COLUMN exit_row_id bigint; +-- audit_event 是 append-only(触发器拦截 UPDATE),回填期间临时摘除。 +DROP TRIGGER IF EXISTS audit_event_append_only ON audit_event; +UPDATE audit_event a SET account_row_id = s.surrogate_id FROM social_account s WHERE s.id::text = a.account_id; +UPDATE audit_event a SET browser_env_row_id = b.surrogate_id FROM browser_env b WHERE b.alias = a.browser_env_alias; +UPDATE audit_event a SET exit_row_id = n.surrogate_id FROM network_exit n WHERE n.id::text = a.network_exit_id; + +ALTER TABLE creator_account_profile ADD COLUMN account_row_id bigint; +UPDATE creator_account_profile p SET account_row_id = s.surrogate_id FROM social_account s WHERE s.id::text = p.account_id; +ALTER TABLE creator_account_password ADD COLUMN account_row_id bigint; +UPDATE creator_account_password p SET account_row_id = s.surrogate_id FROM social_account s WHERE s.id::text = p.account_id; +ALTER TABLE creator_account_metric ADD COLUMN account_row_id bigint; +UPDATE creator_account_metric m SET account_row_id = s.surrogate_id FROM social_account s WHERE s.id::text = m.account_id; + +ALTER TABLE creator_competitor_share_job ADD COLUMN competitor_row_id bigint; +UPDATE creator_competitor_share_job j SET competitor_row_id = c.surrogate_id FROM creator_competitor c WHERE c.id::text = j.competitor_id; + +ALTER TABLE creator_work_metric ADD COLUMN work_row_id bigint; +UPDATE creator_work_metric m SET work_row_id = w.surrogate_id FROM creator_work w WHERE w.id::text = m.work_id; + +ALTER TABLE creator_comment ADD COLUMN work_row_id bigint; +UPDATE creator_comment c SET work_row_id = w.surrogate_id FROM creator_work w WHERE w.id::text = c.work_id; + +ALTER TABLE creator_comment_rule_result ADD COLUMN comment_row_id bigint, ADD COLUMN rule_row_id bigint; +UPDATE creator_comment_rule_result r SET comment_row_id = c.surrogate_id FROM creator_comment c WHERE c.id::text = r.comment_id; +UPDATE creator_comment_rule_result r SET rule_row_id = l.surrogate_id FROM creator_lead_rule l WHERE l.id::text = r.rule_id; + +ALTER TABLE creator_work_cover ADD COLUMN work_row_id bigint; +UPDATE creator_work_cover c SET work_row_id = w.surrogate_id FROM creator_work w WHERE w.id::text = c.work_id; + +-- ============ 3) 删除全部旧文本 FK(解除对旧文本主键的依赖)============ +ALTER TABLE browser_env + DROP CONSTRAINT IF EXISTS browser_env_gateway_name_fkey, + DROP CONSTRAINT IF EXISTS browser_env_account_id_fkey, + DROP CONSTRAINT IF EXISTS browser_env_network_exit_id_fkey; +ALTER TABLE audit_event + DROP CONSTRAINT IF EXISTS audit_event_account_id_fkey, + DROP CONSTRAINT IF EXISTS audit_event_browser_env_alias_fkey, + DROP CONSTRAINT IF EXISTS audit_event_network_exit_id_fkey; +ALTER TABLE creator_account_profile DROP CONSTRAINT IF EXISTS creator_account_profile_account_id_fkey; +ALTER TABLE creator_account_password DROP CONSTRAINT IF EXISTS creator_account_password_account_id_fkey; +ALTER TABLE creator_account_metric DROP CONSTRAINT IF EXISTS creator_account_metric_account_id_fkey; +ALTER TABLE creator_competitor_share_job DROP CONSTRAINT IF EXISTS creator_competitor_share_job_competitor_id_fkey; +ALTER TABLE creator_work_metric DROP CONSTRAINT IF EXISTS creator_work_metric_work_id_fkey; +ALTER TABLE creator_comment DROP CONSTRAINT IF EXISTS creator_comment_work_id_fkey; +ALTER TABLE creator_comment_rule_result + DROP CONSTRAINT IF EXISTS creator_comment_rule_result_comment_id_fkey, + DROP CONSTRAINT IF EXISTS creator_comment_rule_result_rule_id_fkey; +ALTER TABLE creator_work_cover DROP CONSTRAINT IF EXISTS creator_work_cover_work_id_fkey; + +-- ============ 4) 文本列名为 id 的表:先降级改名,再加 bigint 主键(PG 17 IDENTITY 自动回填存量行)============ +ALTER TABLE social_account DROP CONSTRAINT social_account_pkey; +ALTER TABLE social_account RENAME COLUMN id TO account_id; +ALTER TABLE social_account RENAME COLUMN surrogate_id TO id; +ALTER TABLE social_account ADD PRIMARY KEY (id); +ALTER TABLE social_account ADD CONSTRAINT social_account_account_id_unique UNIQUE (account_id); + +ALTER TABLE network_exit DROP CONSTRAINT network_exit_pkey; +ALTER TABLE network_exit RENAME COLUMN id TO exit_id; +ALTER TABLE network_exit RENAME COLUMN surrogate_id TO id; +ALTER TABLE network_exit ADD PRIMARY KEY (id); +ALTER TABLE network_exit ADD CONSTRAINT network_exit_exit_id_unique UNIQUE (exit_id); + +ALTER TABLE creator_competitor DROP CONSTRAINT creator_competitor_pkey; +ALTER TABLE creator_competitor RENAME COLUMN id TO competitor_id; +ALTER TABLE creator_competitor RENAME COLUMN surrogate_id TO id; +ALTER TABLE creator_competitor ADD PRIMARY KEY (id); +ALTER TABLE creator_competitor ADD CONSTRAINT creator_competitor_competitor_id_unique UNIQUE (competitor_id); + +ALTER TABLE creator_work DROP CONSTRAINT creator_work_pkey; +ALTER TABLE creator_work RENAME COLUMN id TO work_id; +ALTER TABLE creator_work RENAME COLUMN surrogate_id TO id; +ALTER TABLE creator_work ADD PRIMARY KEY (id); +ALTER TABLE creator_work ADD CONSTRAINT creator_work_work_id_unique UNIQUE (work_id); + +ALTER TABLE creator_comment DROP CONSTRAINT creator_comment_pkey; +ALTER TABLE creator_comment RENAME COLUMN id TO comment_id; +ALTER TABLE creator_comment RENAME COLUMN surrogate_id TO id; +ALTER TABLE creator_comment ADD PRIMARY KEY (id); +ALTER TABLE creator_comment ADD CONSTRAINT creator_comment_comment_id_unique UNIQUE (comment_id); + +ALTER TABLE creator_competitor_share_job DROP CONSTRAINT creator_competitor_share_job_pkey; +ALTER TABLE creator_competitor_share_job RENAME COLUMN id TO job_id; +ALTER TABLE creator_competitor_share_job RENAME COLUMN surrogate_id TO id; +ALTER TABLE creator_competitor_share_job ADD PRIMARY KEY (id); +ALTER TABLE creator_competitor_share_job ADD CONSTRAINT creator_competitor_share_job_job_id_unique UNIQUE (job_id); + +ALTER TABLE creator_lead_rule DROP CONSTRAINT creator_lead_rule_pkey; +ALTER TABLE creator_lead_rule RENAME COLUMN id TO rule_id; +ALTER TABLE creator_lead_rule RENAME COLUMN surrogate_id TO id; +ALTER TABLE creator_lead_rule ADD PRIMARY KEY (id); +ALTER TABLE creator_lead_rule ADD CONSTRAINT creator_lead_rule_rule_id_unique UNIQUE (rule_id); + +-- ============ 5) 文本主键为语义列的表:加 bigint 主键,文本列降级 UNIQUE ============ +ALTER TABLE gateway DROP CONSTRAINT gateway_pkey; +ALTER TABLE gateway RENAME COLUMN surrogate_id TO id; +ALTER TABLE gateway ADD PRIMARY KEY (id); +ALTER TABLE gateway ADD CONSTRAINT gateway_name_unique UNIQUE (name); + +ALTER TABLE browser_env DROP CONSTRAINT browser_env_pkey; +ALTER TABLE browser_env RENAME COLUMN surrogate_id TO id; +ALTER TABLE browser_env ADD PRIMARY KEY (id); +-- alias 原本就是 NOT NULL,降级为 UNIQUE(账号即环境下 alias = 账号文本 ID)。 +ALTER TABLE browser_env ADD CONSTRAINT browser_env_alias_unique UNIQUE (alias); + + + +-- ============ 6) 引用表换列:删文本 FK 列(连带旧 FK/PK/UNIQUE/索引),新列更名 + 补约束 ============ +ALTER TABLE browser_env + DROP COLUMN gateway_name, DROP COLUMN account_id, DROP COLUMN network_exit_id; +ALTER TABLE browser_env RENAME COLUMN gateway_row_id TO gateway_id; +ALTER TABLE browser_env RENAME COLUMN account_row_id TO account_id; +ALTER TABLE browser_env RENAME COLUMN exit_row_id TO exit_id; +ALTER TABLE browser_env + ALTER COLUMN gateway_id SET NOT NULL, + ALTER COLUMN account_id SET NOT NULL, + ADD CONSTRAINT browser_env_account_id_unique UNIQUE (account_id), + ADD CONSTRAINT browser_env_account_id_fkey FOREIGN KEY (account_id) REFERENCES social_account(id) ON DELETE CASCADE, + ADD CONSTRAINT browser_env_gateway_id_fkey FOREIGN KEY (gateway_id) REFERENCES gateway(id), + ADD CONSTRAINT browser_env_exit_id_fkey FOREIGN KEY (exit_id) REFERENCES network_exit(id); + +ALTER TABLE audit_event + DROP COLUMN account_id, DROP COLUMN browser_env_alias, DROP COLUMN network_exit_id; +ALTER TABLE audit_event RENAME COLUMN account_row_id TO account_id; +ALTER TABLE audit_event RENAME COLUMN browser_env_row_id TO browser_env_id; +ALTER TABLE audit_event RENAME COLUMN exit_row_id TO exit_id; +ALTER TABLE audit_event + ADD CONSTRAINT audit_event_account_id_fkey FOREIGN KEY (account_id) REFERENCES social_account(id), + ADD CONSTRAINT audit_event_browser_env_id_fkey FOREIGN KEY (browser_env_id) REFERENCES browser_env(id), + ADD CONSTRAINT audit_event_exit_id_fkey FOREIGN KEY (exit_id) REFERENCES network_exit(id); +CREATE TRIGGER audit_event_append_only + BEFORE UPDATE OR DELETE ON audit_event + FOR EACH ROW EXECUTE FUNCTION reject_audit_event_mutation(); + +-- social_account.profile_id 是 creator_account_profile 的悬空文本镜像(1:1 已由对方 account_id 表达),删除。 +ALTER TABLE social_account DROP COLUMN profile_id; + +ALTER TABLE creator_account_profile + DROP CONSTRAINT creator_account_profile_pkey, DROP COLUMN account_id; +ALTER TABLE creator_account_profile RENAME COLUMN account_row_id TO account_id; +ALTER TABLE creator_account_profile + ALTER COLUMN account_id SET NOT NULL, + ADD PRIMARY KEY (account_id), + ADD CONSTRAINT creator_account_profile_account_id_fkey FOREIGN KEY (account_id) REFERENCES social_account(id) ON DELETE CASCADE; + +ALTER TABLE creator_account_password + DROP CONSTRAINT creator_account_password_pkey, DROP COLUMN account_id; +ALTER TABLE creator_account_password RENAME COLUMN account_row_id TO account_id; +ALTER TABLE creator_account_password + ALTER COLUMN account_id SET NOT NULL, + ADD PRIMARY KEY (account_id), + ADD CONSTRAINT creator_account_password_account_id_fkey FOREIGN KEY (account_id) REFERENCES social_account(id) ON DELETE CASCADE; + +DROP INDEX IF EXISTS creator_account_metric_account_idx; +ALTER TABLE creator_account_metric DROP COLUMN account_id; +ALTER TABLE creator_account_metric RENAME COLUMN account_row_id TO account_id; +ALTER TABLE creator_account_metric + ADD CONSTRAINT creator_account_metric_account_collected_unique UNIQUE (account_id, collected_at), + ADD CONSTRAINT creator_account_metric_account_id_fkey FOREIGN KEY (account_id) REFERENCES social_account(id) ON DELETE CASCADE; +CREATE INDEX IF NOT EXISTS creator_account_metric_account_collected_idx ON creator_account_metric (account_id, collected_at); + +ALTER TABLE creator_competitor_share_job DROP COLUMN competitor_id; +ALTER TABLE creator_competitor_share_job RENAME COLUMN competitor_row_id TO competitor_id; +ALTER TABLE creator_competitor_share_job + ADD CONSTRAINT creator_competitor_share_job_competitor_id_fkey + FOREIGN KEY (competitor_id) REFERENCES creator_competitor(id) ON DELETE SET NULL; + +ALTER TABLE creator_work_metric DROP COLUMN work_id; +ALTER TABLE creator_work_metric RENAME COLUMN work_row_id TO work_id; +ALTER TABLE creator_work_metric + ALTER COLUMN work_id SET NOT NULL, + ADD CONSTRAINT creator_work_metric_work_collected_unique UNIQUE (work_id, collected_at), + ADD CONSTRAINT creator_work_metric_work_id_fkey FOREIGN KEY (work_id) REFERENCES creator_work(id) ON DELETE CASCADE; + +ALTER TABLE creator_comment DROP COLUMN work_id; +ALTER TABLE creator_comment RENAME COLUMN work_row_id TO work_id; +ALTER TABLE creator_comment + ALTER COLUMN work_id SET NOT NULL, + ADD CONSTRAINT creator_comment_work_id_fkey FOREIGN KEY (work_id) REFERENCES creator_work(id) ON DELETE CASCADE; +CREATE INDEX IF NOT EXISTS creator_comment_work_idx ON creator_comment (work_id, published_at DESC); + +ALTER TABLE creator_comment_rule_result + DROP CONSTRAINT creator_comment_rule_result_pkey, + DROP COLUMN comment_id, DROP COLUMN rule_id; +ALTER TABLE creator_comment_rule_result RENAME COLUMN comment_row_id TO comment_id; +ALTER TABLE creator_comment_rule_result RENAME COLUMN rule_row_id TO rule_id; +ALTER TABLE creator_comment_rule_result + ALTER COLUMN comment_id SET NOT NULL, + ALTER COLUMN rule_id SET NOT NULL, + ADD PRIMARY KEY (comment_id, rule_id), + ADD CONSTRAINT creator_comment_rule_result_comment_id_fkey FOREIGN KEY (comment_id) REFERENCES creator_comment(id) ON DELETE CASCADE, + ADD CONSTRAINT creator_comment_rule_result_rule_id_fkey FOREIGN KEY (rule_id) REFERENCES creator_lead_rule(id) ON DELETE CASCADE; + +ALTER TABLE creator_work_cover + DROP CONSTRAINT creator_work_cover_pkey, DROP COLUMN work_id; +ALTER TABLE creator_work_cover RENAME COLUMN work_row_id TO work_id; +ALTER TABLE creator_work_cover + ALTER COLUMN work_id SET NOT NULL, + ADD PRIMARY KEY (work_id, variant), + ADD CONSTRAINT creator_work_cover_work_id_fkey FOREIGN KEY (work_id) REFERENCES creator_work(id) ON DELETE CASCADE; + + + +-- ============ 7) creator_collection_checkpoint:文本 id 是派生值(source:sourceID:kind),无独立语义,删除; +-- 换 bigint 主键,唯一业务键保持 (source_type, source_id, collection_kind) ============ +ALTER TABLE creator_collection_checkpoint + DROP CONSTRAINT creator_collection_checkpoint_pkey, DROP COLUMN id; +ALTER TABLE creator_collection_checkpoint ADD COLUMN id bigint GENERATED ALWAYS AS IDENTITY; +ALTER TABLE creator_collection_checkpoint ADD PRIMARY KEY (id); + +-- ============ 8) creator_settings:布尔单例行主键换 bigint identity,表达式唯一索引锁定单行 ============ +ALTER TABLE creator_settings + DROP CONSTRAINT creator_settings_pkey, DROP COLUMN id; +ALTER TABLE creator_settings ADD COLUMN id bigint GENERATED ALWAYS AS IDENTITY; +ALTER TABLE creator_settings ADD PRIMARY KEY (id); +CREATE UNIQUE INDEX creator_settings_singleton_idx ON creator_settings ((true)); diff --git a/internal/environment/store.go b/internal/environment/store.go index 157d95f..e8e2f53 100644 --- a/internal/environment/store.go +++ b/internal/environment/store.go @@ -86,6 +86,9 @@ var migration038 string //go:embed migrations/043_schema_consolidation.sql var migration043 string +//go:embed migrations/044_numeric_primary_keys.sql +var migration044 string + //go:embed migrations/1017_creator.sql var migration1017 string @@ -300,7 +303,7 @@ func (s *Store) migrate(ctx context.Context) error { {1029, migration1029}, {1030, migration1030}, {1031, migration1031}, {1032, migration1032}, {1033, migration1033}, {1034, migration1034}, {1035, migration1035}, {1036, migration1036}, {1037, migration1037}, {1038, migration1038}, {1039, migration1039}, {1040, migration1040}, {1041, migration1041}, {1042, migration1042}, - {43, migration043}} { + {43, migration043}, {44, migration044}} { var applied bool if err := tx.QueryRowContext(ctx, `SELECT EXISTS (SELECT 1 FROM schema_migration WHERE version = $1)`, migration.version).Scan(&applied); err != nil { return errors.New("read environment schema migration state") @@ -410,32 +413,12 @@ func (s *Store) DeleteGateway(ctx context.Context, name string) error { return nil } -func (s *Store) CreateEnv(ctx context.Context, env Env) error { - env.Alias = strings.TrimSpace(env.Alias) - env.Name = strings.TrimSpace(env.Name) - if !aliasPattern.MatchString(env.Alias) || !validDisplayName(env.Name) || - !gatewayNamePattern.MatchString(env.Gateway) || env.Fingerprint.ProxyServer != "" { - return ErrInvalid - } - if err := env.Fingerprint.Validate(); err != nil { - return fmt.Errorf("%w: %s", ErrInvalid, err) - } - encoded, err := json.Marshal(env.Fingerprint) - if err != nil { - return ErrInvalid - } - var created string - err = s.db.QueryRowContext(ctx, ` - INSERT INTO browser_env (alias, name, gateway_name, fingerprint) - VALUES ($1, $2, $3, $4) - RETURNING alias`, env.Alias, env.Name, env.Gateway, encoded).Scan(&created) - return rowError(err) -} - func (s *Store) ListEnvs(ctx context.Context) ([]Env, error) { rows, err := s.db.QueryContext(ctx, ` - SELECT alias, name, gateway_name, fingerprint, created_at - FROM browser_env ORDER BY created_at, alias`) + SELECT environment.alias, environment.name, gateway.name, environment.fingerprint, environment.created_at + FROM browser_env environment + JOIN gateway ON gateway.id = environment.gateway_id + ORDER BY environment.created_at, environment.alias`) if err != nil { return nil, errors.New("read browser envs") } @@ -456,8 +439,10 @@ func (s *Store) GetEnv(ctx context.Context, alias string) (Env, error) { return Env{}, ErrInvalid } rows, err := s.db.QueryContext(ctx, ` - SELECT alias, name, gateway_name, fingerprint, created_at - FROM browser_env WHERE alias = $1`, alias) + SELECT environment.alias, environment.name, gateway.name, environment.fingerprint, environment.created_at + FROM browser_env environment + JOIN gateway ON gateway.id = environment.gateway_id + WHERE environment.alias = $1`, alias) if err != nil { return Env{}, errors.New("read browser env") } @@ -496,10 +481,11 @@ func (s *Store) DeleteAccountEnvironment(ctx context.Context, accountID string) defer tx.Rollback() var alias string if err := tx.QueryRowContext(ctx, ` - SELECT alias - FROM browser_env - WHERE account_id = $1 - FOR UPDATE`, accountID).Scan(&alias); errors.Is(err, sql.ErrNoRows) { + SELECT environment.alias + FROM browser_env environment + JOIN social_account account ON account.id = environment.account_id + WHERE account.account_id = $1 + FOR UPDATE OF environment`, accountID).Scan(&alias); errors.Is(err, sql.ErrNoRows) { return nil } else if err != nil { return errors.New("read account environment binding") diff --git a/internal/environment/store_test.go b/internal/environment/store_test.go index c35a5fb..0dfd420 100644 --- a/internal/environment/store_test.go +++ b/internal/environment/store_test.go @@ -202,7 +202,6 @@ func TestFingerprintArgsFollowUpstreamCommandLineContract(t *testing.T) { func TestFingerprintValidateRejectsUnsupportedValues(t *testing.T) { invalid := map[string]Fingerprint{ - "seed zero": {Seed: 0}, "seed overflow": {Seed: 2147483648}, "platform": {Seed: 1, Platform: "android"}, "brand": {Seed: 1, Brand: "Firefox"}, @@ -261,9 +260,6 @@ func TestStoreValidationRejectsInvalidInputsBeforePersistence(t *testing.T) { if _, _, err := store.CreateBoundEnv(ctx, Env{Alias: "account-a", Name: strings.Repeat("名", 65), Gateway: "gw-1", Fingerprint: Fingerprint{Seed: 1}}, "account-a", ""); !errors.Is(err, ErrInvalid) { t.Fatalf("expected overlong name, got %v", err) } - if _, _, err := store.CreateBoundEnv(ctx, Env{Alias: "account-a", Name: "甲", Gateway: "gw-1", Fingerprint: Fingerprint{Seed: 0}}, "account-a", ""); !errors.Is(err, ErrInvalid) { - t.Fatalf("expected invalid fingerprint, got %v", err) - } for name, exit := range map[string]NetworkExit{ "protocol": {Protocol: "direct", Host: "proxy.example", Port: 1080}, "userinfo": {Protocol: "socks5", Host: "user@proxy.example", Port: 1080}, @@ -297,18 +293,36 @@ func TestFingerprintSeedIsGloballyUnique(t *testing.T) { t.Fatal(err) } if _, err := store.db.ExecContext(ctx, ` - INSERT INTO social_account (id, credential_provider, credential_key, platform, platform_account_key, authorization_kind, authorization_status) + INSERT INTO social_account (account_id, credential_provider, credential_key, platform, platform_account_key, authorization_kind, authorization_status) VALUES ('seed-account-a', 'os_keyring', 'creatorhub/seed-a', 'mock', 'seed-account-a', 'owned', 'authorized'), ('seed-account-b', 'os_keyring', 'creatorhub/seed-b', 'mock', 'seed-account-b', 'owned', 'authorized')`); err != nil { t.Fatal(err) } - env := Env{Alias: "seed-environment-a", Name: "Seed A", Gateway: "gw-seed", Fingerprint: Fingerprint{Seed: 77}} + env := Env{Alias: "seed-environment-a", Name: "Seed A", Gateway: "gw-seed"} if _, created, err := store.CreateBoundEnv(ctx, env, "seed-account-a", ""); err != nil || !created { t.Fatalf("create first seeded environment: created=%v err=%v", created, err) } env.Alias, env.Name = "seed-environment-b", "Seed B" - if _, _, err := store.CreateBoundEnv(ctx, env, "seed-account-b", ""); !errors.Is(err, ErrConflict) { - t.Fatalf("duplicate fingerprint seed was accepted: %v", err) + if _, created, err := store.CreateBoundEnv(ctx, env, "seed-account-b", ""); err != nil || !created { + t.Fatalf("second bound environment on distinct account must create cleanly: created=%v err=%v", created, err) + } + // 数字主键:seed = 账号 bigint id + 1000(服务端派生,客户端传入值被覆盖) + var firstSeed, secondSeed, firstRowID, secondRowID int64 + if err := store.db.QueryRowContext(ctx, ` + SELECT (fingerprint->>'seed')::bigint, s.id FROM browser_env b + JOIN social_account s ON s.id = b.account_id WHERE b.alias = 'seed-environment-a'`).Scan(&firstSeed, &firstRowID); err != nil { + t.Fatal(err) + } + if err := store.db.QueryRowContext(ctx, `SELECT id FROM social_account WHERE account_id = 'seed-account-b'`).Scan(&secondRowID); err != nil { + t.Fatal(err) + } + if err := store.db.QueryRowContext(ctx, ` + SELECT (fingerprint->>'seed')::bigint FROM browser_env b + JOIN social_account s ON s.id = b.account_id WHERE b.alias = 'seed-environment-b'`).Scan(&secondSeed); err != nil { + t.Fatal(err) + } + if firstSeed != firstRowID+1000 || secondSeed != secondRowID+1000 { + t.Fatalf("seed must derive from account bigint id: a=%d/%d b=%d/%d", firstSeed, firstRowID, secondSeed, secondRowID) } } @@ -341,7 +355,7 @@ func TestHubWorkflow(t *testing.T) { } if _, err := store.db.ExecContext(ctx, ` - INSERT INTO social_account (id, credential_provider, credential_key, platform, platform_account_key, authorization_kind, authorization_status) + INSERT INTO social_account (account_id, credential_provider, credential_key, platform, platform_account_key, authorization_kind, authorization_status) VALUES ('shop-owner', 'os_keyring', 'creatorhub/shop-owner', 'mock', 'shop-owner', 'owned', 'authorized')`); err != nil { t.Fatal(err) } @@ -356,19 +370,19 @@ func TestHubWorkflow(t *testing.T) { t.Fatalf("expected idempotent environment reuse, created=%v err=%v", created, err) } if _, err := store.db.ExecContext(ctx, ` - INSERT INTO social_account (id, credential_provider, credential_key, platform, platform_account_key, authorization_kind, authorization_status) + INSERT INTO social_account (account_id, credential_provider, credential_key, platform, platform_account_key, authorization_kind, authorization_status) VALUES ('shop-owner-2', 'os_keyring', 'creatorhub/shop-owner-2', 'mock', 'shop-owner-2', 'owned', 'authorized')`); err != nil { t.Fatal(err) } - if _, _, err := store.CreateBoundEnv(ctx, Env{Alias: "shop-02", Name: "店铺二号", Gateway: "missing", Fingerprint: Fingerprint{Seed: 2}}, "shop-owner-2", ""); !errors.Is(err, ErrConflict) { - t.Fatalf("expected unknown gateway conflict, got %v", err) + if _, _, err := store.CreateBoundEnv(ctx, Env{Alias: "shop-02", Name: "店铺二号", Gateway: "missing", Fingerprint: Fingerprint{Seed: 2}}, "shop-owner-2", ""); !errors.Is(err, ErrNotFound) { + t.Fatalf("unknown gateway must surface as readiness not-found, got %v", err) } listed, err := store.ListEnvs(ctx) if err != nil || len(listed) != 1 { t.Fatalf("expected one env, err=%v list=%#v", err, listed) } - if listed[0].Name != "店铺一号" || listed[0].Fingerprint.Seed != 1000 || listed[0].Fingerprint.Timezone != "Asia/Shanghai" { + if listed[0].Name != "店铺一号" || listed[0].Fingerprint.Seed != 1001 || listed[0].Fingerprint.Timezone != "Asia/Shanghai" { t.Fatalf("fingerprint must round trip through jsonb: %#v", listed[0]) } updatedGateway, err := store.UpdateGateway(ctx, "gw-main", "gw-renamed", "http://127.0.0.4:8081", "") @@ -423,7 +437,7 @@ func TestNetworkExitBindingRuntimeAndAuditWorkflow(t *testing.T) { } if _, err := store.db.ExecContext(ctx, ` INSERT INTO social_account - (id, credential_provider, credential_key, platform, platform_account_key, authorization_kind, authorization_status) + (account_id, credential_provider, credential_key, platform, platform_account_key, authorization_kind, authorization_status) VALUES ('account-a', 'os_keyring', 'creatorhub/account-a', 'mock', 'account-a', 'owned', 'authorized')`); err != nil { t.Fatal(err) } @@ -466,7 +480,7 @@ func TestNetworkExitBindingRuntimeAndAuditWorkflow(t *testing.T) { if err != nil || created || reused.Alias != bound.Alias { t.Fatalf("same account must reuse its environment: %#v created=%v err=%v", reused, created, err) } - if _, err := store.db.ExecContext(ctx, `UPDATE social_account SET status = 'active' WHERE id = 'account-a'`); err != nil { + if _, err := store.db.ExecContext(ctx, `UPDATE social_account SET status = 'active' WHERE account_id = 'account-a'`); err != nil { t.Fatal(err) } active, err := store.ActivateRuntime(ctx, env.Alias, "runtime-a", bound.BindingVersion, bound.Exit.ID, "network-a") @@ -474,8 +488,8 @@ func TestNetworkExitBindingRuntimeAndAuditWorkflow(t *testing.T) { t.Fatalf("activate runtime: %#v err=%v", active, err) } assertDatabaseCount(t, store.db, `SELECT count(*) FROM audit_event WHERE event_type = 'runtime_bound' - AND reason_code = 'runtime_bound' AND account_id = 'account-a' AND browser_env_alias = 'environment-a' - AND network_exit_id = $1 AND binding_version = $2 AND details = '{}'::jsonb`, + AND reason_code = 'runtime_bound' AND browser_env_id = (SELECT id FROM browser_env WHERE alias = 'environment-a') + AND exit_id = (SELECT id FROM network_exit WHERE exit_id = $1) AND binding_version = $2 AND details = '{}'::jsonb`, 1, active.Exit.ID, active.BindingVersion) second, err := store.CreateNetworkExit(ctx, NetworkExit{Protocol: "http", Host: "proxy-2.example", Port: 8080}) @@ -488,10 +502,11 @@ func TestNetworkExitBindingRuntimeAndAuditWorkflow(t *testing.T) { } if _, err := store.db.ExecContext(ctx, ` INSERT INTO social_account - (id, credential_provider, credential_key, platform, platform_account_key, authorization_kind, authorization_status) + (account_id, credential_provider, credential_key, platform, platform_account_key, authorization_kind, authorization_status) VALUES ('account-b', 'os_keyring', 'creatorhub/account-b', 'mock', 'account-b', 'owned', 'authorized'); - INSERT INTO browser_env (alias, name, gateway_name, fingerprint, account_id) - VALUES ('environment-b', '环境 B', 'gw-main', '{"seed":2}', 'account-b')`); err != nil { + INSERT INTO browser_env (alias, name, gateway_id, fingerprint, account_id) + VALUES ('environment-b', '环境 B', (SELECT g.id FROM gateway g WHERE g.name = 'gw-main'), '{"seed":2}', + (SELECT a.id FROM social_account a WHERE a.account_id = 'account-b'))`); err != nil { t.Fatal(err) } legacyRebound, err := store.RebindEnvironment(ctx, "environment-b", second.ID, "", 1) @@ -518,7 +533,7 @@ func TestNetworkExitBindingRuntimeAndAuditWorkflow(t *testing.T) { if err != nil || expired.RuntimeID != active.RuntimeID { t.Fatalf("context read discarded expired cleanup generation: %#v err=%v", expired, err) } - if _, err := store.db.ExecContext(ctx, `UPDATE social_account SET status = 'paused' WHERE id = 'account-a'`); err != nil { + if _, err := store.db.ExecContext(ctx, `UPDATE social_account SET status = 'paused' WHERE account_id = 'account-a'`); err != nil { t.Fatal(err) } rebound, err := store.RebindEnvironment(ctx, env.Alias, second.ID, "", bound.BindingVersion) @@ -526,10 +541,10 @@ func TestNetworkExitBindingRuntimeAndAuditWorkflow(t *testing.T) { t.Fatalf("expired runtime must be transactionally released before rebind: %#v err=%v", rebound, err) } assertDatabaseCount(t, store.db, `SELECT count(*) FROM audit_event WHERE event_type = 'runtime_released' - AND reason_code = 'runtime_released' AND account_id = 'account-a' AND browser_env_alias = 'environment-a' - AND network_exit_id = $1 AND binding_version = $2 AND details = '{}'::jsonb`, + AND reason_code = 'runtime_released' AND browser_env_id = (SELECT id FROM browser_env WHERE alias = 'environment-a') + AND exit_id = (SELECT id FROM network_exit WHERE exit_id = $1) AND binding_version = $2 AND details = '{}'::jsonb`, 1, active.Exit.ID, active.BindingVersion) - if _, err := store.db.ExecContext(ctx, `UPDATE social_account SET status = 'active' WHERE id = 'account-a'`); err != nil { + if _, err := store.db.ExecContext(ctx, `UPDATE social_account SET status = 'active' WHERE account_id = 'account-a'`); err != nil { t.Fatal(err) } expiredBeforeActivation, err := store.ActivateRuntime(ctx, env.Alias, "expired-runtime", rebound.BindingVersion, rebound.Exit.ID, "network-expired") @@ -543,10 +558,10 @@ func TestNetworkExitBindingRuntimeAndAuditWorkflow(t *testing.T) { t.Fatalf("replace expired runtime before same-exit rebind: %v", err) } assertDatabaseCount(t, store.db, `SELECT count(*) FROM audit_event WHERE event_type = 'runtime_released' - AND reason_code = 'runtime_released' AND account_id = 'account-a' AND browser_env_alias = 'environment-a' - AND network_exit_id = $1 AND binding_version = $2 AND details = '{}'::jsonb`, + AND reason_code = 'runtime_released' AND browser_env_id = (SELECT id FROM browser_env WHERE alias = 'environment-a') + AND exit_id = (SELECT id FROM network_exit WHERE exit_id = $1) AND binding_version = $2 AND details = '{}'::jsonb`, 1, expiredBeforeActivation.Exit.ID, expiredBeforeActivation.BindingVersion) - if _, err := store.db.ExecContext(ctx, `UPDATE social_account SET status = 'paused' WHERE id = 'account-a'`); err != nil { + if _, err := store.db.ExecContext(ctx, `UPDATE social_account SET status = 'paused' WHERE account_id = 'account-a'`); err != nil { t.Fatal(err) } if _, err := store.RebindEnvironment(ctx, env.Alias, second.ID, "", rebound.BindingVersion); !errors.Is(err, ErrConflict) { @@ -561,8 +576,8 @@ func TestNetworkExitBindingRuntimeAndAuditWorkflow(t *testing.T) { } // 审计已无 runtime_instance_id:同代两次释放(过期替换 + 显式释放)行相同,计 2 assertDatabaseCount(t, store.db, `SELECT count(*) FROM audit_event WHERE event_type = 'runtime_released' - AND reason_code = 'runtime_released' AND account_id = 'account-a' AND browser_env_alias = 'environment-a' - AND network_exit_id = $1 AND binding_version = $2 AND details = '{}'::jsonb`, + AND reason_code = 'runtime_released' AND browser_env_id = (SELECT id FROM browser_env WHERE alias = 'environment-a') + AND exit_id = (SELECT id FROM network_exit WHERE exit_id = $1) AND binding_version = $2 AND details = '{}'::jsonb`, 2, active.Exit.ID, active.BindingVersion) rebound, err = store.RebindEnvironment(ctx, env.Alias, second.ID, "rebound-runtime", rebound.BindingVersion) if err != nil || rebound.BindingVersion != 3 || rebound.RuntimeID != "rebound-runtime" { @@ -582,7 +597,7 @@ func TestNetworkExitBindingRuntimeAndAuditWorkflow(t *testing.T) { t.Fatalf("set generation cleanup pending: %v", err) } assertDatabaseCount(t, store.db, `SELECT count(*) FROM audit_event WHERE event_type = 'runtime_released' - AND browser_env_alias = 'environment-a' AND network_exit_id = $1 AND binding_version = $2`, + AND browser_env_id = (SELECT id FROM browser_env WHERE alias = 'environment-a') AND exit_id = (SELECT id FROM network_exit WHERE exit_id = $1) AND binding_version = $2`, 1, current.Exit.ID, current.BindingVersion) wrongCleanup := cleanup wrongCleanup.RuntimeCleanupRuntimeID = "other-runtime" @@ -606,7 +621,7 @@ func TestNetworkExitBindingRuntimeAndAuditWorkflow(t *testing.T) { if err := store.SetRuntimeCleanupPending(ctx, cleanup, true); !errors.Is(err, ErrConflict) { t.Fatalf("stale binding set cleanup pending on version %d: %v", newGeneration.BindingVersion, err) } - if _, err := store.db.ExecContext(ctx, `UPDATE social_account SET status = 'active' WHERE id = 'account-a'`); err != nil { + if _, err := store.db.ExecContext(ctx, `UPDATE social_account SET status = 'active' WHERE account_id = 'account-a'`); err != nil { t.Fatal(err) } pauseTx, err := store.db.BeginTx(ctx, nil) @@ -614,7 +629,7 @@ func TestNetworkExitBindingRuntimeAndAuditWorkflow(t *testing.T) { t.Fatal(err) } var locked string - if err := pauseTx.QueryRowContext(ctx, `SELECT id FROM social_account WHERE id = 'account-a' FOR UPDATE`).Scan(&locked); err != nil { + if err := pauseTx.QueryRowContext(ctx, `SELECT account_id FROM social_account WHERE account_id = 'account-a' FOR UPDATE`).Scan(&locked); err != nil { t.Fatal(err) } activation := make(chan error, 1) @@ -627,7 +642,7 @@ func TestNetworkExitBindingRuntimeAndAuditWorkflow(t *testing.T) { t.Fatalf("activation bypassed the locked account row: %v", err) case <-time.After(time.Second): } - if _, err := pauseTx.ExecContext(ctx, `UPDATE social_account SET status = 'paused' WHERE id = 'account-a'`); err != nil { + if _, err := pauseTx.ExecContext(ctx, `UPDATE social_account SET status = 'paused' WHERE account_id = 'account-a'`); err != nil { t.Fatal(err) } if err := pauseTx.Commit(); err != nil { @@ -677,7 +692,7 @@ func TestRuntimeCleanupFenceAndBoundCreateReadiness(t *testing.T) { store := openFullyMigratedHub(t, ctx, isolatedDatabaseURL(t, databaseURL)) t.Cleanup(func() { _ = store.Close() }) if _, err := store.db.ExecContext(ctx, ` - INSERT INTO social_account (id, credential_provider, credential_key, platform, platform_account_key, authorization_kind, authorization_status) + INSERT INTO social_account (account_id, credential_provider, credential_key, platform, platform_account_key, authorization_kind, authorization_status) VALUES ('fence-a', 'os_keyring', 'creatorhub/fence-a', 'mock', 'fence-a', 'owned', 'authorized')`); err != nil { t.Fatal(err) } @@ -692,20 +707,20 @@ func TestRuntimeCleanupFenceAndBoundCreateReadiness(t *testing.T) { if _, _, err := store.CreateBoundEnv(ctx, Env{Alias: "fence-x", Name: "甲", Gateway: "gw-fence", Fingerprint: Fingerprint{Seed: 9}}, "fence-missing", ""); !errors.Is(err, ErrNotFound) { t.Fatalf("missing account must be not-found: %v", err) } - if _, err := store.db.ExecContext(ctx, `UPDATE network_exit SET health_status = 'unhealthy' WHERE id = $1`, exit.ID); err != nil { + if _, err := store.db.ExecContext(ctx, `UPDATE network_exit SET health_status = 'unhealthy' WHERE exit_id = $1`, exit.ID); err != nil { t.Fatal(err) } if _, _, err := store.CreateBoundEnv(ctx, Env{Alias: "fence-x", Name: "甲", Gateway: "gw-fence", Fingerprint: Fingerprint{Seed: 9}}, "fence-a", exit.ID); !errors.Is(err, ErrNotFound) { t.Fatalf("unhealthy exit must be not-found: %v", err) } - if _, err := store.db.ExecContext(ctx, `UPDATE network_exit SET health_status = 'healthy' WHERE id = $1`, exit.ID); err != nil { + if _, err := store.db.ExecContext(ctx, `UPDATE network_exit SET health_status = 'healthy' WHERE exit_id = $1`, exit.ID); err != nil { t.Fatal(err) } bound, created, err := store.CreateBoundEnv(ctx, Env{Alias: "fence-a", Name: "甲", Gateway: "gw-fence", Fingerprint: Fingerprint{Seed: 1}}, "fence-a", exit.ID) if err != nil || !created { t.Fatalf("create bound env: created=%v err=%v", created, err) } - if _, err := store.db.ExecContext(ctx, `UPDATE social_account SET status = 'active' WHERE id = 'fence-a'`); err != nil { + if _, err := store.db.ExecContext(ctx, `UPDATE social_account SET status = 'active' WHERE account_id = 'fence-a'`); err != nil { t.Fatal(err) } active, err := store.ActivateRuntime(ctx, "fence-a", "fence-runtime-a", bound.BindingVersion, exit.ID, "net-fence") @@ -757,7 +772,7 @@ func TestNetworkExitLifecycleStoreOperations(t *testing.T) { if err != nil { t.Fatal(err) } - if _, err := store.db.ExecContext(ctx, `UPDATE network_exit SET health_status = 'unhealthy', last_check_reason = 'exit_ip_drift' WHERE id = $1`, exit.ID); err != nil { + if _, err := store.db.ExecContext(ctx, `UPDATE network_exit SET health_status = 'unhealthy', last_check_reason = 'exit_ip_drift' WHERE exit_id = $1`, exit.ID); err != nil { t.Fatal(err) } // 配置变更重置健康状态为 unchecked(disabled 保留),last_check_reason 标记配置变更 @@ -781,7 +796,7 @@ func TestNetworkExitLifecycleStoreOperations(t *testing.T) { if err != nil { t.Fatal(err) } - if _, err := store.db.ExecContext(ctx, `UPDATE network_exit SET health_status = 'disabled' WHERE id = $1`, disabled.ID); err != nil { + if _, err := store.db.ExecContext(ctx, `UPDATE network_exit SET health_status = 'disabled' WHERE exit_id = $1`, disabled.ID); err != nil { t.Fatal(err) } reEnabled, err := store.EnableNetworkExit(ctx, disabled.ID)