Files
creator-hub/internal/environment/store_test.go
T
rogee a94bdf7921 feat(gateway): 网关健康状态——后台定时探活只保留最近一次结果并在列表展示
- 迁移 1043:gateway 表加 health_status/last_check_reason/last_checked_at
- Store.RecordGatewayCheck 覆盖写最近一次探活结果,Create/Update/List/Get 全量返回健康字段
- api.RecordGatewayHealthChecks 并发探 /healthz:连接失败或 5xx 判离线,其余判在线(401 等异常写入原因供排查)
- workers 每 30 秒一轮,app 启动接线
- 网关管理列表新增健康状态列:在线/离线/未检查 Tag,原因与检查时间收进 Tooltip
2026-09-29 14:47:49 +08:00

870 lines
40 KiB
Go
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
package environment
import (
"context"
"encoding/json"
"errors"
"fmt"
"os"
"reflect"
"strings"
"sync/atomic"
"testing"
"time"
)
func TestNetworkExitCredentialValidation(t *testing.T) {
valid := NetworkExit{Protocol: "socks5", Host: "proxy.example", Port: 1080, Username: "operator", Password: "plain-password"}
if !validNetworkExit(valid) {
t.Fatal("valid stored credentials were rejected")
}
for name, mutate := range map[string]func(*NetworkExit){
"password without username": func(exit *NetworkExit) { exit.Username = "" },
"username too long": func(exit *NetworkExit) { exit.Username = strings.Repeat("u", 256) },
"password too long": func(exit *NetworkExit) { exit.Password = strings.Repeat("p", 256) },
"username control character": func(exit *NetworkExit) { exit.Username = "operator\n" },
"password control character": func(exit *NetworkExit) { exit.Password = "plain\x7fpassword" },
} {
t.Run(name, func(t *testing.T) {
exit := valid
mutate(&exit)
if validNetworkExit(exit) {
t.Fatalf("invalid credentials were accepted: %#v", exit)
}
})
}
}
func TestEnvironmentLocksCoordinateAcrossStoreInstances(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()
testURL := isolatedDatabaseURL(t, databaseURL)
first := openFullyMigratedHub(t, ctx, testURL)
t.Cleanup(func() { _ = first.Close() })
second, err := Open(ctx, testURL)
if err != nil {
t.Fatal(err)
}
t.Cleanup(func() { _ = second.Close() })
unlockFirst, err := first.LockResources(ctx, []string{"account-a"}, nil, nil)
if err != nil {
t.Fatal(err)
}
firstReleased := false
defer func() {
if !firstReleased {
unlockFirst()
}
}()
differentAlias, err := second.LockResources(ctx, []string{"account-b"}, nil, nil)
if err != nil {
t.Fatalf("different aliases must not share a lock: %v", err)
}
differentAlias()
acquired := make(chan func(), 1)
errors := make(chan error, 1)
started := make(chan struct{})
go func() {
close(started)
unlock, lockErr := second.LockResources(ctx, []string{"account-a"}, nil, nil)
if lockErr != nil {
errors <- lockErr
return
}
acquired <- unlock
}()
<-started
select {
case unlock := <-acquired:
unlock()
t.Fatal("same alias lock did not block across Store instances")
case err := <-errors:
t.Fatal(err)
case <-time.After(50 * time.Millisecond):
}
unlockFirst()
firstReleased = true
select {
case unlock := <-acquired:
unlock()
case err := <-errors:
t.Fatal(err)
case <-time.After(time.Second):
t.Fatal("same alias lock was not released")
}
}
func TestResourceLocksReserveConnectionsForLifecycleQueries(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, cancel := context.WithTimeout(context.Background(), 3*time.Second)
defer cancel()
store := openFullyMigratedHub(t, ctx, isolatedDatabaseURL(t, databaseURL))
t.Cleanup(func() { _ = store.Close() })
const workers = 10
var acquired atomic.Int32
startQueries := make(chan struct{})
done := make(chan error, workers)
for worker := 0; worker < workers; worker++ {
go func(worker int) {
unlock, err := store.LockResources(ctx, []string{fmt.Sprintf("account-%d", worker)}, nil, nil)
if err != nil {
done <- err
return
}
defer unlock()
acquired.Add(1)
<-startQueries
_, err = store.ListEnvs(ctx)
done <- err
}(worker)
}
deadline := time.NewTimer(100 * time.Millisecond)
ticker := time.NewTicker(time.Millisecond)
for acquired.Load() < workers {
select {
case <-ticker.C:
case <-deadline.C:
goto release
}
}
release:
ticker.Stop()
if !deadline.Stop() {
select {
case <-deadline.C:
default:
}
}
close(startQueries)
for worker := 0; worker < workers; worker++ {
if err := <-done; err != nil {
t.Fatalf("locked lifecycle query %d did not complete: %v", worker, err)
}
}
}
func TestFingerprintArgsFollowUpstreamCommandLineContract(t *testing.T) {
full := Fingerprint{
Seed: 2024, Platform: "windows", PlatformVersion: "11.0.0",
Brand: "Edge", BrandVersion: "132.0.6834.159", HardwareConcurrency: 8,
Lang: "zh-CN", AcceptLang: "zh-CN,en-US", Timezone: "Asia/Shanghai",
ProxyServer: "socks5://127.0.0.1:1080", DisableNonProxiedUDP: true, DisableSpoofing: "font,gpu",
}
if err := full.Validate(); err != nil {
t.Fatalf("expected full fingerprint to be valid: %v", err)
}
want := []string{
"--fingerprint=2024",
"--fingerprint-platform=windows",
"--fingerprint-platform-version=11.0.0",
"--fingerprint-brand=Edge",
"--fingerprint-brand-version=132.0.6834.159",
"--fingerprint-hardware-concurrency=8",
"--lang=zh-CN",
"--accept-lang=zh-CN,en-US",
"--timezone=Asia/Shanghai",
"--proxy-server=socks5://127.0.0.1:1080",
"--disable-non-proxied-udp",
"--disable-spoofing=font,gpu",
}
if !reflect.DeepEqual(full.Args(), want) {
t.Fatalf("unexpected args:\n got %v\nwant %v", full.Args(), want)
}
minimal := Fingerprint{Seed: 1}
if err := minimal.Validate(); err != nil {
t.Fatalf("minimal fingerprint must be valid: %v", err)
}
if args := minimal.Args(); len(args) != 1 || args[0] != "--fingerprint=1" {
t.Fatalf("zero-value fields must be omitted: %v", args)
}
encoded, err := json.Marshal(minimal)
if err != nil {
t.Fatal(err)
}
var decoded Fingerprint
if err := json.Unmarshal(encoded, &decoded); err != nil || decoded.Seed != 1 || decoded.Args()[0] != "--fingerprint=1" {
t.Fatalf("fingerprint must survive JSON round trip: %#v %v", decoded, err)
}
}
func TestFingerprintValidateRejectsUnsupportedValues(t *testing.T) {
invalid := map[string]Fingerprint{
"seed overflow": {Seed: 2147483648},
"platform": {Seed: 1, Platform: "android"},
"brand": {Seed: 1, Brand: "Firefox"},
"platform version": {Seed: 1, PlatformVersion: "bad value"},
"brand version": {Seed: 1, BrandVersion: strings.Repeat("x", 33)},
"concurrency": {Seed: 1, HardwareConcurrency: 129},
"lang": {Seed: 1, Lang: "zh CN"},
"accept lang": {Seed: 1, AcceptLang: "zh-CN;drop"},
"timezone": {Seed: 1, Timezone: "Asia/Shanghai\n"},
"proxy scheme": {Seed: 1, ProxyServer: "ftp://proxy:21"},
"proxy host": {Seed: 1, ProxyServer: "http://"},
"proxy userinfo": {Seed: 1, ProxyServer: "socks5://user:password@proxy:1080"},
"spoofing unknown": {Seed: 1, DisableSpoofing: "webrtc"},
"spoofing repeated": {Seed: 1, DisableSpoofing: "font,font"},
}
for name, fingerprint := range invalid {
t.Run(name, func(t *testing.T) {
if err := fingerprint.Validate(); err == nil {
t.Fatalf("expected rejection for %#v", fingerprint)
}
})
}
}
func TestStoreValidationRejectsInvalidInputsBeforePersistence(t *testing.T) {
store := &Store{}
ctx := context.Background()
if _, err := store.CreateGateway(ctx, "bad name!", "http://gw:8081", ""); !errors.Is(err, ErrInvalid) {
t.Fatalf("expected invalid gateway name, got %v", err)
}
if _, err := store.CreateGateway(ctx, "gw-1", "ftp://gw:8081", ""); !errors.Is(err, ErrInvalid) {
t.Fatalf("expected invalid gateway endpoint, got %v", err)
}
if _, err := store.CreateGateway(ctx, "gw-1", "http://gw:8081", "short-token"); !errors.Is(err, ErrInvalid) {
t.Fatalf("expected invalid gateway token, got %v", err)
}
for _, test := range []struct {
name string
currentName string
newName string
endpoint string
token string
}{
{name: "current name", currentName: "bad name!", newName: "gw-2", endpoint: "http://gw:8081"},
{name: "new name", currentName: "gw-1", newName: "bad name!", endpoint: "http://gw:8081"},
{name: "endpoint", currentName: "gw-1", newName: "gw-2", endpoint: "ftp://gw:8081"},
{name: "token", currentName: "gw-1", newName: "gw-2", endpoint: "http://gw:8081", token: "short-token"},
} {
if _, err := store.UpdateGateway(ctx, test.currentName, test.newName, test.endpoint, test.token); !errors.Is(err, ErrInvalid) {
t.Fatalf("expected invalid gateway update %s, got %v", test.name, err)
}
}
if _, _, err := store.CreateBoundEnv(ctx, Env{Alias: "UP", Name: "甲", Gateway: "gw-1", Fingerprint: Fingerprint{Seed: 1}}, "account-a", ""); !errors.Is(err, ErrInvalid) {
t.Fatalf("expected invalid alias, got %v", err)
}
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)
}
for name, exit := range map[string]NetworkExit{
"protocol": {Protocol: "direct", Host: "proxy.example", Port: 1080},
"userinfo": {Protocol: "socks5", Host: "user@proxy.example", Port: 1080},
"URL host": {Protocol: "socks5", Host: "socks5://proxy.example", Port: 1080},
"port": {Protocol: "socks5", Host: "proxy.example", Port: 0},
"ip": {Protocol: "socks5", Host: "proxy.example", Port: 1080, ExpectedPublicIP: "not-an-ip"},
} {
t.Run("network exit "+name, func(t *testing.T) {
if _, err := store.CreateNetworkExit(ctx, exit); !errors.Is(err, ErrInvalid) {
t.Fatalf("expected invalid network exit, got %v", err)
}
})
}
if _, _, err := store.CreateBoundEnv(ctx, Env{Alias: "account-a", Name: "甲", Gateway: "gw-1", Fingerprint: Fingerprint{Seed: 1, ProxyServer: "socks5://proxy.example:1080"}}, "account-a", ""); !errors.Is(err, ErrInvalid) {
t.Fatalf("stored fingerprint proxy must be rejected, got %v", err)
}
}
func TestFingerprintSeedIsGloballyUnique(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 := openFullyMigratedHub(t, ctx, isolatedDatabaseURL(t, databaseURL))
t.Cleanup(func() { _ = store.Close() })
if _, err := store.db.ExecContext(ctx, `TRUNCATE browser_env, social_account, gateway CASCADE`); err != nil {
t.Fatal(err)
}
if _, err := store.CreateGateway(ctx, "gw-seed", "http://127.0.0.1:8081", "unit-test-gateway-token"); err != nil {
t.Fatal(err)
}
if _, err := store.db.ExecContext(ctx, `
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"}
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 _, 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)
}
}
func TestHubWorkflow(t *testing.T) {
databaseURL := os.Getenv("CREATORHUB_POSTGRES_TEST_URL")
if databaseURL == "" {
t.Skip("set CREATORHUB_POSTGRES_TEST_URL to run PostgreSQL integration coverage")
}
databaseURL = isolatedDatabaseURL(t, databaseURL)
ctx := context.Background()
store := openFullyMigratedHub(t, ctx, databaseURL)
t.Cleanup(func() { _ = store.Close() })
if _, err := store.db.ExecContext(ctx, `TRUNCATE browser_env, gateway CASCADE`); err != nil {
t.Fatal(err)
}
gateway, err := store.CreateGateway(ctx, "gw-main", "http://127.0.0.1:8081", "")
if err != nil {
t.Fatal(err)
}
if len(gateway.Token) != 48 {
t.Fatalf("platform must assign a 48-hex-char token: %q", gateway.Token)
}
custom, err := store.CreateGateway(ctx, "gw-custom", "http://127.0.0.3:8081", "operator-provided-token-1234")
if err != nil || custom.Token != "operator-provided-token-1234" {
t.Fatalf("explicit token must be honored: %#v %v", custom, err)
}
if _, err := store.CreateGateway(ctx, "gw-main", "http://127.0.0.2:8081", ""); !errors.Is(err, ErrConflict) {
t.Fatalf("expected duplicate gateway conflict, got %v", err)
}
if _, err := store.db.ExecContext(ctx, `
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)
}
env := Env{
Alias: "shop-01", Name: "店铺一号", Gateway: "gw-main",
Fingerprint: Fingerprint{Seed: 1000, Timezone: "Asia/Shanghai", Lang: "zh-CN"},
}
if _, created, err := store.CreateBoundEnv(ctx, env, "shop-owner", ""); err != nil || !created {
t.Fatalf("create bound env: created=%v err=%v", created, err)
}
if _, created, err := store.CreateBoundEnv(ctx, env, "shop-owner", ""); err != nil || created {
t.Fatalf("expected idempotent environment reuse, created=%v err=%v", created, err)
}
if _, err := store.db.ExecContext(ctx, `
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, 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 != 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", "")
if err != nil || updatedGateway.Name != "gw-renamed" || updatedGateway.Endpoint != "http://127.0.0.4:8081" || updatedGateway.Token != gateway.Token {
t.Fatalf("gateway update did not preserve the token: %#v err=%v", updatedGateway, err)
}
renamedEnv, err := store.GetEnv(ctx, "shop-01")
if err != nil || renamedEnv.Gateway != "gw-renamed" {
t.Fatalf("gateway rename did not cascade to environment: %#v err=%v", renamedEnv, err)
}
if _, err := store.UpdateGateway(ctx, "gw-renamed", "gw-custom", "http://127.0.0.4:8081", ""); !errors.Is(err, ErrConflict) {
t.Fatalf("expected gateway rename conflict, got %v", err)
}
if _, err := store.UpdateGateway(ctx, "missing", "gw-missing", "http://127.0.0.5:8081", ""); !errors.Is(err, ErrNotFound) {
t.Fatalf("expected missing gateway on update, got %v", err)
}
if _, err := store.UpdateGateway(ctx, "gw-renamed", "gw-main", "http://127.0.0.1:8081", ""); err != nil {
t.Fatalf("restore gateway name after cascade check: %v", err)
}
if _, err := store.GetEnv(ctx, "ghost"); !errors.Is(err, ErrNotFound) {
t.Fatalf("expected missing env, got %v", err)
}
if err := store.DeleteGateway(ctx, "gw-main"); !errors.Is(err, ErrConflict) {
t.Fatalf("referenced gateway must not be deletable, got %v", err)
}
if err := store.DeleteEnv(ctx, "shop-01"); err != nil {
t.Fatal(err)
}
if err := store.DeleteEnv(ctx, "shop-01"); !errors.Is(err, ErrNotFound) {
t.Fatalf("expected missing env on double delete, got %v", err)
}
if err := store.DeleteGateway(ctx, "gw-main"); err != nil {
t.Fatal(err)
}
if err := store.DeleteGateway(ctx, "gw-custom"); err != nil {
t.Fatal(err)
}
}
func TestNetworkExitBindingRuntimeAndAuditWorkflow(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 := openFullyMigratedHub(t, ctx, isolatedDatabaseURL(t, databaseURL))
t.Cleanup(func() { _ = store.Close() })
if _, err := store.db.ExecContext(ctx, `TRUNCATE audit_event, network_exit,
social_account, browser_env, gateway CASCADE`); err != nil {
t.Fatal(err)
}
if _, err := store.db.ExecContext(ctx, `
INSERT INTO social_account
(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)
}
if _, err := store.CreateGateway(ctx, "gw-main", "http://127.0.0.1:8081", "unit-test-gateway-token"); err != nil {
t.Fatal(err)
}
exit, err := store.CreateNetworkExit(ctx, NetworkExit{
Protocol: "socks5", Host: "proxy.example", Port: 1080,
Username: "proxy-user", Password: "plain-password",
ExpectedPublicIP: "203.0.113.10", ExpectedRegion: "test-region",
})
if err != nil || exit.HealthStatus != "unchecked" || exit.Username != "proxy-user" || exit.Password != "plain-password" {
t.Fatalf("unexpected network exit: %#v err=%v", exit, err)
}
exported, _ := json.Marshal(exit)
if !strings.Contains(string(exported), `"username":"proxy-user"`) || !strings.Contains(string(exported), `"password":"plain-password"`) {
t.Fatalf("network exit response must include stored credentials: %s", exported)
}
access, err := store.GetNetworkExitAccess(ctx, exit.ID)
if err != nil || access.Username != "proxy-user" || access.Password != "plain-password" {
t.Fatalf("runtime network exit credentials unavailable: %#v err=%v", access, err)
}
exit, reason, err := store.RecordNetworkExitCheck(ctx, exit.ID, ExitObservation{PublicIP: "203.0.113.11", Region: "test-region"}, "")
if err != nil || exit.HealthStatus != "unhealthy" || reason != "exit_ip_drift" {
t.Fatalf("identity drift must make the exit unhealthy: %#v reason=%s err=%v", exit, reason, err)
}
exit, reason, err = store.RecordNetworkExitCheck(ctx, exit.ID, ExitObservation{PublicIP: "203.0.113.10", Region: "test-region"}, "")
if err != nil || exit.HealthStatus != "healthy" || reason != "exit_healthy" {
t.Fatalf("matching identity must make the exit healthy: %#v reason=%s err=%v", exit, reason, err)
}
env := Env{Alias: "environment-a", Name: "环境 A", Gateway: "gw-main", Fingerprint: Fingerprint{Seed: 1}}
bound, created, err := store.CreateBoundEnv(ctx, env, "account-a", exit.ID)
if err != nil || !created || bound.AccountID != "account-a" || bound.Exit.ID != exit.ID {
t.Fatalf("create stable binding: %#v created=%v err=%v", bound, created, err)
}
reused, created, err := store.CreateBoundEnv(ctx, env, "account-a", exit.ID)
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 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")
if err != nil || active.RuntimeID == "" {
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 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})
if err != nil {
t.Fatal(err)
}
second, _, err = store.RecordNetworkExitCheck(ctx, second.ID, ExitObservation{PublicIP: "198.51.100.2", Region: "other"}, "")
if err != nil || second.HealthStatus != "healthy" {
t.Fatalf("prepare second exit: %#v err=%v", second, err)
}
if _, err := store.db.ExecContext(ctx, `
INSERT INTO social_account
(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_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)
if err != nil || legacyRebound.Exit.ID != second.ID {
t.Fatalf("binding without an exit must support explicit rebind: %#v err=%v", legacyRebound, err)
}
if _, err := store.ActivateRuntime(ctx, env.Alias, "runtime-a", bound.BindingVersion+1, second.ID, "network-a"); !errors.Is(err, ErrConflict) {
t.Fatalf("stale binding metadata must not activate a runtime: %v", err)
}
if _, err := store.db.ExecContext(ctx, `UPDATE browser_env SET runtime_lease_until = now() + interval '10 seconds' WHERE alias = $1`, env.Alias); err != nil {
t.Fatal(err)
}
if _, err := store.ActivateRuntime(ctx, env.Alias, "runtime-a", bound.BindingVersion, bound.Exit.ID, "network-a"); err != nil {
t.Fatalf("runtime heartbeat failed: %v", err)
}
var renewed bool
if err := store.db.QueryRowContext(ctx, `SELECT runtime_lease_until > now() + interval '30 seconds' FROM browser_env WHERE alias = $1`, env.Alias).Scan(&renewed); err != nil || !renewed {
t.Fatalf("runtime lease was not renewed: renewed=%v err=%v", renewed, err)
}
if _, err := store.db.ExecContext(ctx, `UPDATE browser_env SET runtime_lease_until = now() - interval '1 second' WHERE alias = $1`, env.Alias); err != nil {
t.Fatal(err)
}
expired, err := store.GetEnvironmentContext(ctx, env.Alias)
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 account_id = 'account-a'`); err != nil {
t.Fatal(err)
}
rebound, err := store.RebindEnvironment(ctx, env.Alias, second.ID, "", bound.BindingVersion)
if err != nil || rebound.Exit.ID != second.ID || rebound.BindingVersion != 2 {
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 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 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")
if err != nil {
t.Fatalf("activate runtime to expire: %v", err)
}
if _, err := store.db.ExecContext(ctx, `UPDATE browser_env SET runtime_lease_until = now() - interval '1 second' WHERE alias = $1`, env.Alias); err != nil {
t.Fatal(err)
}
if _, err := store.ActivateRuntime(ctx, env.Alias, "same-exit-runtime", rebound.BindingVersion, rebound.Exit.ID, "network-same-exit"); err != nil {
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 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 account_id = 'account-a'`); err != nil {
t.Fatal(err)
}
if _, err := store.RebindEnvironment(ctx, env.Alias, second.ID, "", rebound.BindingVersion); !errors.Is(err, ErrConflict) {
t.Fatalf("active runtime must block same-exit rebind: %v", err)
}
active, err = store.GetEnvironmentContext(ctx, env.Alias)
if err != nil {
t.Fatal(err)
}
if err := store.ReleaseRuntime(ctx, active); err != nil {
t.Fatal(err)
}
// 审计已无 runtime_instance_id:同代两次释放(过期替换 + 显式释放)行相同,计 2
assertDatabaseCount(t, store.db, `SELECT count(*) FROM audit_event WHERE event_type = 'runtime_released'
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" {
t.Fatalf("same-exit rebind must atomically CAS the binding and runtime: %#v err=%v", rebound, err)
}
if err := store.ReleaseRuntime(ctx, active); !errors.Is(err, ErrConflict) {
t.Fatalf("stale generation release must conflict: %v", err)
}
current, err := store.GetEnvironmentContext(ctx, env.Alias)
if err != nil || current.RuntimeID != rebound.RuntimeID || current.RuntimeID != "rebound-runtime" {
t.Fatalf("stale release changed the current runtime: %#v err=%v", current, err)
}
cleanup := current
cleanup.RuntimeCleanupBindingVersion = current.BindingVersion
cleanup.RuntimeCleanupRuntimeID = current.RuntimeID
if err := store.SetRuntimeCleanupPending(ctx, cleanup, true); err != nil {
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_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"
if err := store.SetRuntimeCleanupPending(ctx, wrongCleanup, false); !errors.Is(err, ErrConflict) {
t.Fatalf("wrong cleanup generation cleared pending state: %v", err)
}
pending, err := store.GetEnvironmentContext(ctx, env.Alias)
if err != nil || !pending.RuntimeCleanupPending || pending.RuntimeCleanupRuntimeID != current.RuntimeID {
t.Fatalf("cleanup generation was not persisted: %#v err=%v", pending, err)
}
if _, err := store.ActivateRuntime(ctx, env.Alias, "candidate-runtime", current.BindingVersion, current.Exit.ID, "network-candidate"); !errors.Is(err, ErrConflict) {
t.Fatalf("paused account activated a stale request: %v", err)
}
if err := store.SetRuntimeCleanupPending(ctx, pending, false); err != nil {
t.Fatalf("clear matching cleanup generation: %v", err)
}
newGeneration, err := store.RebindEnvironment(ctx, env.Alias, current.Exit.ID, "", current.BindingVersion)
if err != nil {
t.Fatalf("advance binding generation: %v", err)
}
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 account_id = 'account-a'`); err != nil {
t.Fatal(err)
}
pauseTx, err := store.db.BeginTx(ctx, nil)
if err != nil {
t.Fatal(err)
}
var locked string
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)
go func() {
_, err := store.ActivateRuntime(ctx, env.Alias, "racing-runtime", newGeneration.BindingVersion, newGeneration.Exit.ID, "network-racing")
activation <- err
}()
select {
case err := <-activation:
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 account_id = 'account-a'`); err != nil {
t.Fatal(err)
}
if err := pauseTx.Commit(); err != nil {
t.Fatal(err)
}
if err := <-activation; !errors.Is(err, ErrConflict) {
t.Fatalf("activation won the pause race: %v", err)
}
latest, err := store.GetEnvironmentContext(ctx, env.Alias)
if err != nil || latest.RuntimeID != "" {
t.Fatalf("pause race left an active runtime: %#v err=%v", latest, err)
}
action := EnvironmentAction{
OperationID: NewOperationID(), Action: "start", AccountID: rebound.AccountID,
BrowserEnvAlias: rebound.Alias, NetworkExitID: rebound.Exit.ID, BindingVersion: rebound.BindingVersion,
ReasonCode: "action_requested",
}
if err := store.AppendEnvironmentAction(ctx, "environment_action_requested", action); err != nil {
t.Fatal(err)
}
action.Outcome, action.ReasonCode = "succeeded", "environment_started"
if err := store.AppendEnvironmentAction(ctx, "environment_action_finished", action); err != nil {
t.Fatal(err)
}
assertDatabaseCount(t, store.db, `SELECT count(*) FROM audit_event WHERE operation_id = '`+action.OperationID+`'`, 2)
if _, err := store.DisableNetworkExit(ctx, newGeneration.Exit.ID); err != nil {
t.Fatalf("exit without active runtime should be disableable: %v", err)
}
var auditText string
if err := store.db.QueryRowContext(ctx, `SELECT string_agg(row_to_json(event)::text, '') FROM audit_event event`).Scan(&auditText); err != nil {
t.Fatal(err)
}
for _, forbidden := range []string{"creatorhub/proxy-main", "credential-exit", "username", "password"} {
if strings.Contains(auditText, forbidden) {
t.Fatalf("audit leaked sensitive value %q: %s", forbidden, auditText)
}
}
}
func TestRuntimeCleanupFenceAndBoundCreateReadiness(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 := openFullyMigratedHub(t, ctx, isolatedDatabaseURL(t, databaseURL))
t.Cleanup(func() { _ = store.Close() })
if _, err := store.db.ExecContext(ctx, `
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)
}
if _, err := store.CreateGateway(ctx, "gw-fence", "http://127.0.0.1:8081", "unit-test-gateway-token"); err != nil {
t.Fatal(err)
}
exit, err := store.CreateNetworkExit(ctx, NetworkExit{Protocol: "http", Host: "proxy.example", Port: 8080})
if err != nil {
t.Fatal(err)
}
// CreateBoundEnv 就绪门禁:账号缺失/出口不健康 → ErrNotFound
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 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 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 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")
if err != nil {
t.Fatal(err)
}
// 未知代 fence:允许登记待清理并事务性释放活跃实例
unknown := active
unknown.RuntimeCleanupBindingVersion, unknown.RuntimeCleanupRuntimeID, unknown.RuntimeCleanupNetworkID = active.BindingVersion, MissingRuntimeID, ""
if err := store.SetRuntimeCleanupPending(ctx, unknown, true); err != nil {
t.Fatalf("unknown-generation fence must be accepted: %v", err)
}
after, err := store.GetEnvironmentContext(ctx, "fence-a")
if err != nil || !after.RuntimeCleanupPending || after.RuntimeCleanupRuntimeID != MissingRuntimeID || after.RuntimeID != "" {
t.Fatalf("unknown-generation fence state: %#v err=%v", after, err)
}
if err := store.SetRuntimeCleanupPending(ctx, after, false); err != nil {
t.Fatalf("clear matching fence: %v", err)
}
// 重新激活后,异代 fence 请求必须冲突
active2, err := store.ActivateRuntime(ctx, "fence-a", "fence-runtime-b", bound.BindingVersion, exit.ID, "net-fence")
if err != nil {
t.Fatal(err)
}
stale := active2
stale.RuntimeCleanupBindingVersion, stale.RuntimeCleanupRuntimeID, stale.RuntimeCleanupNetworkID = active2.BindingVersion, "fence-runtime-other", ""
if err := store.SetRuntimeCleanupPending(ctx, stale, true); !errors.Is(err, ErrConflict) {
t.Fatalf("stale fence generation must conflict: %v", err)
}
if _, err := store.GetEnvironmentContext(ctx, "fence-a"); err != nil {
t.Fatal(err)
}
emptyFence := active2
emptyFence.RuntimeCleanupRuntimeID = ""
if err := store.SetRuntimeCleanupPending(ctx, emptyFence, true); !errors.Is(err, ErrInvalid) {
t.Fatalf("pending without cleanup runtime id must be invalid: %v", err)
}
}
func TestNetworkExitLifecycleStoreOperations(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 := openFullyMigratedHub(t, ctx, isolatedDatabaseURL(t, databaseURL))
t.Cleanup(func() { _ = store.Close() })
exit, err := store.CreateNetworkExit(ctx, NetworkExit{Protocol: "http", Host: "lifecycle.example", Port: 8080, Username: "user", Password: "pass"})
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 exit_id = $1`, exit.ID); err != nil {
t.Fatal(err)
}
// 配置变更重置健康状态为 unchecked(disabled 保留),last_check_reason 标记配置变更
updated, err := store.UpdateNetworkExit(ctx, exit.ID, NetworkExit{Protocol: "socks5", Host: "lifecycle-2.example", Port: 1080})
if err != nil || updated.Host != "lifecycle-2.example" || updated.Protocol != "socks5" || updated.HealthStatus != "unchecked" || updated.LastCheckReason != "exit_configuration_changed" {
t.Fatalf("update network exit must reset health to unchecked: %#v err=%v", updated, err)
}
if _, err := store.UpdateNetworkExit(ctx, exit.ID, NetworkExit{}); !errors.Is(err, ErrInvalid) {
t.Fatalf("invalid exit update must be rejected: %v", err)
}
if _, err := store.UpdateNetworkExit(ctx, "exit-missing-0001", NetworkExit{Protocol: "http", Host: "x.example", Port: 8080}); !errors.Is(err, ErrNotFound) {
t.Fatalf("missing exit update must be not-found: %v", err)
}
if err := store.DeleteNetworkExit(ctx, exit.ID); err != nil {
t.Fatalf("delete unused exit: %v", err)
}
if _, err := store.GetNetworkExit(ctx, exit.ID); !errors.Is(err, ErrNotFound) {
t.Fatalf("deleted exit must be gone: %v", err)
}
disabled, err := store.CreateNetworkExit(ctx, NetworkExit{Protocol: "http", Host: "disabled.example", Port: 8080})
if err != nil {
t.Fatal(err)
}
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)
if err != nil || reEnabled.HealthStatus != "unchecked" {
t.Fatalf("enable disabled exit: %#v err=%v", reEnabled, err)
}
exits, err := store.ListNetworkExits(ctx)
if err != nil || len(exits) != 1 || exits[0].ID != disabled.ID {
t.Fatalf("list exits: exits=%+v err=%v", exits, err)
}
}
func TestGatewayHealthCheckPersistence(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 := openFullyMigratedHub(t, ctx, isolatedDatabaseURL(t, databaseURL))
t.Cleanup(func() { _ = store.Close() })
gateway, err := store.CreateGateway(ctx, "gw-health", "http://127.0.0.1:8081", "")
if err != nil {
t.Fatal(err)
}
if gateway.HealthStatus != "unchecked" || gateway.LastCheckReason != "" || gateway.LastCheckedAt != nil {
t.Fatalf("new gateway must start unchecked: %#v", gateway)
}
checked, err := store.RecordGatewayCheck(ctx, "gw-health", "healthy", "")
if err != nil {
t.Fatal(err)
}
if checked.HealthStatus != "healthy" || checked.LastCheckReason != "" || checked.LastCheckedAt == nil {
t.Fatalf("healthy check was not persisted: %#v", checked)
}
failed, err := store.RecordGatewayCheck(ctx, "gw-health", "unhealthy", `Get "http://127.0.0.1:8081/healthz": connection refused`)
if err != nil {
t.Fatal(err)
}
if failed.HealthStatus != "unhealthy" || failed.LastCheckReason == "" || failed.LastCheckedAt == nil {
t.Fatalf("unhealthy check was not persisted: %#v", failed)
}
if !failed.LastCheckedAt.After(*checked.LastCheckedAt) && failed.LastCheckedAt.Equal(*checked.LastCheckedAt) {
t.Fatalf("last_checked_at must advance: %v", failed.LastCheckedAt)
}
reread, err := store.GetGateway(ctx, "gw-health")
if err != nil || reread.HealthStatus != "unhealthy" || reread.LastCheckReason != failed.LastCheckReason {
t.Fatalf("health fields must round trip through get/list: %#v err=%v", reread, err)
}
listed, err := store.ListGateways(ctx)
if err != nil || len(listed) != 1 || listed[0].HealthStatus != "unhealthy" {
t.Fatalf("health fields must round trip through list: list=%#v err=%v", listed, err)
}
if _, err := store.RecordGatewayCheck(ctx, "gw-health", "degraded", ""); !errors.Is(err, ErrInvalid) {
t.Fatalf("unknown status must be rejected, got %v", err)
}
if _, err := store.RecordGatewayCheck(ctx, "missing", "healthy", ""); !errors.Is(err, ErrNotFound) {
t.Fatalf("missing gateway must surface not found, got %v", err)
}
// 改名/改地址不触碰健康字段(由探活负责覆盖)。
renamed, err := store.UpdateGateway(ctx, "gw-health", "gw-health-2", "http://127.0.0.2:8081", "")
if err != nil || renamed.HealthStatus != "unhealthy" || renamed.LastCheckReason == "" {
t.Fatalf("update must preserve health fields: %#v err=%v", renamed, err)
}
}