diff --git a/internal/controlplane/api/account_environment.go b/internal/controlplane/api/account_environment.go new file mode 100644 index 0000000..254dfc5 --- /dev/null +++ b/internal/controlplane/api/account_environment.go @@ -0,0 +1,77 @@ +package api + +import ( + "context" + "errors" + "hash/fnv" + + hub "git.ipao.vip/rogee/creator-hub/internal/environment" +) + +// 账号即环境:账号创建即绑定运行环境。 +// 规则:alias = 账号文本 ID、唯一网关自动取、出口直连、seed 由账号 ID 派生。 + +// deriveSeed 由账号文本 ID 确定性派生 fingerprint seed(1001..2147483647)。 +// 数字主键落地后改为 seed = 账号 bigint id + 1000,仅替换本函数。 +func deriveSeed(accountID string) int64 { + hasher := fnv.New64a() + _, _ = hasher.Write([]byte(accountID)) + return int64(hasher.Sum64()%(2147483647-1000)) + 1000 +} + +// soleGateway 返回平台唯一网关;零配置绑定要求恰有一个网关: +// 缺失返回 ErrNotFound,多个返回 ErrConflict(无法自动选择)。 +func soleGateway(ctx context.Context, store HubStore) (hub.Gateway, error) { + gateways, err := store.ListGateways(ctx) + if err != nil { + return hub.Gateway{}, err + } + if len(gateways) == 0 { + return hub.Gateway{}, hub.ErrNotFound + } + if len(gateways) > 1 { + return hub.Gateway{}, hub.ErrConflict + } + return gateways[0], nil +} + +// ensureAccountEnvironment 幂等补建账号环境:已绑定(任意出口)原样返回; +// 未绑定时以 alias=账号 ID、派生 seed、直连出口创建。 +func ensureAccountEnvironment(ctx context.Context, store HubStore, accountID string) (hub.EnvironmentContext, bool, error) { + environment, err := store.GetEnvironmentContextForAccount(ctx, accountID) + if err == nil { + return environment, false, nil + } + if !errors.Is(err, hub.ErrNotFound) { + return hub.EnvironmentContext{}, false, err + } + gateway, err := soleGateway(ctx, store) + if err != nil { + return hub.EnvironmentContext{}, false, err + } + return store.CreateBoundEnv(ctx, hub.Env{ + Alias: accountID, Name: accountID, Gateway: gateway.Name, + Fingerprint: hub.Fingerprint{Seed: deriveSeed(accountID)}, + }, accountID, "") +} + +// startAccountEnvironment 启动账号环境并记录审计对;环境缺失时幂等补建(创建即绑定的自愈路径)。 +func startAccountEnvironment(ctx context.Context, store HubStore, accountID string) error { + environment, err := store.GetEnvironmentContextForAccount(ctx, accountID) + if errors.Is(err, hub.ErrNotFound) { + environment, _, err = ensureAccountEnvironment(ctx, store, accountID) + } + if err != nil { + return err + } + action := actionForEnvironment("start", environment) + if err := store.AppendEnvironmentAction(ctx, "environment_action_requested", action); err != nil { + return err + } + finish := func(outcome, reason string, current hub.EnvironmentContext) error { + action.Outcome, action.ReasonCode = outcome, reason + action.BindingVersion, action.NetworkExitID = current.BindingVersion, current.Exit.ID + return store.AppendEnvironmentAction(ctx, "environment_action_finished", action) + } + return startBrowserRuntime(ctx, store, defaultNetworkExitProbe(), environment, finish) +} diff --git a/internal/controlplane/api/account_environment_unit_test.go b/internal/controlplane/api/account_environment_unit_test.go new file mode 100644 index 0000000..cb10056 --- /dev/null +++ b/internal/controlplane/api/account_environment_unit_test.go @@ -0,0 +1,140 @@ +package api + +import ( + "context" + "database/sql" + "encoding/json" + "net/http" + "net/http/httptest" + "os" + "testing" + + accountdomain "git.ipao.vip/rogee/creator-hub/internal/account" + hub "git.ipao.vip/rogee/creator-hub/internal/environment" + "github.com/gofiber/fiber/v3" +) + +func TestDeriveSeed(t *testing.T) { + // 确定性:同 ID 派生同 seed + if deriveSeed("account-0123456789abcdef01234567") != deriveSeed("account-0123456789abcdef01234567") { + t.Fatal("deriveSeed must be deterministic") + } + // 值域:1001..2147483647(seed 上限约束,偏移 1000 对齐未来 bigint id + 1000) + for _, accountID := range []string{"account-a", "account-000000000000000000000000", "account-zzzzzzzzzzzzzzzzzzzzzzzz", ""} { + seed := deriveSeed(accountID) + if seed < 1001 || seed > 2147483647 { + t.Fatalf("seed out of range for %q: %d", accountID, seed) + } + } + // 不同 ID 派生不同 seed(固定向量,防回归) + if deriveSeed("account-a") == deriveSeed("account-b") { + t.Fatal("distinct accounts must derive distinct seeds") + } + if got := deriveSeed("account-a"); got != 1816671480 { + t.Fatalf("deriveSeed vector drifted: %d", got) + } +} + +func TestAccountEnvironmentAutoBindingAndStart(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() + databaseURL = isolatedControlPlaneDatabaseURL(t, databaseURL) + accountStore, err := accountdomain.Open(ctx, databaseURL) + if err != nil { + t.Fatal(err) + } + t.Cleanup(func() { _ = accountStore.Close() }) + hubStore, err := hub.Open(ctx, databaseURL) + if err != nil { + t.Fatal(err) + } + t.Cleanup(func() { _ = hubStore.Close() }) + gateway := &fakeGateway{token: "unit-test-gateway-token"} + gatewayServer := httptest.NewServer(gateway.handler(t)) + t.Cleanup(gatewayServer.Close) + app := fiber.New() + RegisterAccountRoutes(app, accountStore, hubStore, &testCredentialBridge{values: map[string]string{}}) + + // 网关缺失:创建账号即绑定失败 → 503 environment_binding_failed,但账号已存在(可补建重试) + response := do(app, http.MethodPost, "/api/phase-a/accounts", + `{"name":"测试账号","platform":"douyin","platform_account_key":"key-binding-1","cookies":"sessionid=1"}`) + if response.Code != http.StatusServiceUnavailable { + t.Fatalf("expected 503 when no gateway exists, got %d: %s", response.Code, response.Body.String()) + } + var failure map[string]string + if err := json.Unmarshal(response.Body.Bytes(), &failure); err != nil || failure["reason_code"] != "environment_binding_failed" || failure["account_id"] == "" { + t.Fatalf("binding failure payload: %s err=%v", response.Body.String(), err) + } + accountID := failure["account_id"] + if _, err := accountStore.GetAccount(ctx, accountID); err != nil { + t.Fatalf("account must exist after failed binding: %v", err) + } + if response := do(app, http.MethodPost, "/api/phase-a/accounts/"+accountID+"/environment", ""); response.Code != http.StatusNotFound { + t.Fatalf("expected 404 environment rebind without gateway, got %d: %s", response.Code, response.Body.String()) + } + + // 注册唯一网关后:补建成功,alias=账号 ID、出口直连、seed 派生 + if _, err := hubStore.CreateGateway(ctx, "gw-main", gatewayServer.URL, gateway.token); err != nil { + t.Fatal(err) + } + response = do(app, http.MethodPost, "/api/phase-a/accounts/"+accountID+"/environment", "") + if response.Code != http.StatusOK { + t.Fatalf("expected 200 environment rebind, got %d: %s", response.Code, response.Body.String()) + } + var bound struct { + Alias string `json:"alias"` + Gateway string `json:"gateway"` + Created bool `json:"created"` + } + if err := json.Unmarshal(response.Body.Bytes(), &bound); err != nil || bound.Alias != accountID || bound.Gateway != "gw-main" || !bound.Created { + t.Fatalf("rebind payload: %s err=%v", response.Body.String(), err) + } + environment, err := hubStore.GetEnvironmentContext(ctx, accountID) + if err != nil || environment.Fingerprint.Seed != deriveSeed(accountID) || environment.Exit.ID != "" { + t.Fatalf("auto-bound environment: %#v err=%v", environment, err) + } + // 幂等:重复补建返回既有环境 + response = do(app, http.MethodPost, "/api/phase-a/accounts/"+accountID+"/environment", "") + var rebound struct { + Created bool `json:"created"` + } + if err := json.Unmarshal(response.Body.Bytes(), &rebound); err != nil || response.Code != http.StatusOK || rebound.Created { + t.Fatalf("rebind must be idempotent: %d %s err=%v", response.Code, response.Body.String(), err) + } + + // start = resume + 启动环境:激活 runtime 并落审计对 + response = do(app, http.MethodPost, "/api/phase-a/accounts/"+accountID+"/start", "") + if response.Code != http.StatusNoContent { + t.Fatalf("expected 204 start, got %d: %s", response.Code, response.Body.String()) + } + environment, err = hubStore.GetEnvironmentContext(ctx, accountID) + if err != nil || environment.RuntimeID == "" { + t.Fatalf("start must activate the environment runtime: %#v err=%v", environment, err) + } + auditDB, err := sql.Open("pgx", databaseURL) + if err != nil { + t.Fatal(err) + } + t.Cleanup(func() { _ = auditDB.Close() }) + var startAudits int + if err := auditDB.QueryRowContext(ctx, + `SELECT count(*) FROM audit_event WHERE account_id = $1 AND action = 'start' AND reason_code IN ('action_requested','environment_started')`, accountID).Scan(&startAudits); err != nil || startAudits != 2 { + t.Fatalf("start audit pair missing: rows=%d err=%v", startAudits, err) + } + + // start 冲突分支:吊销账号后 start → 409 readiness blocked + if response := do(app, http.MethodPost, "/api/phase-a/accounts/"+accountID+"/revoke", ""); response.Code != http.StatusNoContent { + t.Fatalf("revoke failed: %d: %s", response.Code, response.Body.String()) + } + response = do(app, http.MethodPost, "/api/phase-a/accounts/"+accountID+"/start", "") + if response.Code != http.StatusConflict { + t.Fatalf("expected 409 start on revoked account, got %d: %s", response.Code, response.Body.String()) + } + var conflict map[string]string + if err := json.Unmarshal(response.Body.Bytes(), &conflict); err != nil || conflict["reason_code"] != "account_revoked" || conflict["readiness"] != "blocked" { + t.Fatalf("start conflict payload: %s err=%v", response.Body.String(), err) + } +} diff --git a/internal/controlplane/api/accounts_operations.go b/internal/controlplane/api/accounts_operations.go index 27e070e..c572ab2 100644 --- a/internal/controlplane/api/accounts_operations.go +++ b/internal/controlplane/api/accounts_operations.go @@ -22,8 +22,8 @@ type accountRequest struct { Cookies string `json:"cookies"` } -// RegisterAccountRoutes exposes account lifecycle routes (create/list/detail/pause/resume/revoke + audit). -func RegisterAccountRoutes(app *fiber.App, store *accountdomain.Store, runtimeStore RuntimeStopStore, credentials accountdomain.CredentialBridge) { +// RegisterAccountRoutes exposes account lifecycle routes (create with auto-binding/list/detail/补建/start/pause/resume/revoke + audit). +func RegisterAccountRoutes(app *fiber.App, store *accountdomain.Store, runtimeStore HubStore, credentials accountdomain.CredentialBridge) { app.Post("/api/phase-a/accounts", func(c fiber.Ctx) error { var input accountRequest if err := decodePhaseA(c, &input); err != nil { @@ -52,6 +52,14 @@ func RegisterAccountRoutes(app *fiber.App, store *accountdomain.Store, runtimeSt } return phaseAError(c, err) } + // 账号即环境:创建即绑定(幂等)。绑定失败透传原因与账号 ID,客户端可用幂等补建端点重试。 + if runtimeStore != nil { + if _, _, err := ensureAccountEnvironment(c.Context(), runtimeStore, accountID); err != nil { + return c.Status(fiber.StatusServiceUnavailable).JSON(map[string]string{ + "error": err.Error(), "reason_code": "environment_binding_failed", "account_id": accountID, + }) + } + } return c.Status(fiber.StatusCreated).JSON(account) }) @@ -103,6 +111,45 @@ func RegisterAccountRoutes(app *fiber.App, store *accountdomain.Store, runtimeSt return c.SendStatus(fiber.StatusNoContent) }) + // 账号即环境:幂等补建运行环境(网关缺失/多网关透传错误)。 + app.Post("/api/phase-a/accounts/:id/environment", func(c fiber.Ctx) error { + if runtimeStore == nil { + return c.Status(fiber.StatusServiceUnavailable).JSON(map[string]string{"error": "environment store unavailable"}) + } + unlock, err := lockAccountResources(c.Context(), runtimeStore, c.Params("id")) + if err != nil { + return hubError(c, err) + } + defer unlock() + environment, created, err := ensureAccountEnvironment(c.Context(), runtimeStore, c.Params("id")) + if err != nil { + return hubError(c, err) + } + return c.JSON(map[string]any{"alias": environment.Alias, "gateway": environment.Gateway, "created": created}) + }) + + // 账号即环境:start = resume + 启动环境;停止沿用 pause(已联动停环境)。 + app.Post("/api/phase-a/accounts/:id/start", func(c fiber.Ctx) error { + if runtimeStore == nil { + return c.Status(fiber.StatusServiceUnavailable).JSON(map[string]string{"error": "environment store unavailable"}) + } + unlock, err := lockAccountResources(c.Context(), runtimeStore, c.Params("id")) + if err != nil { + return hubError(c, err) + } + defer unlock() + if err := store.ResumeAccount(c.Context(), c.Params("id")); err != nil { + if errors.Is(err, accountdomain.ErrConflict) { + return accountResumeConflict(c, store, runtimeStore, c.Params("id")) + } + return phaseAError(c, err) + } + if err := startAccountEnvironment(c.Context(), runtimeStore, c.Params("id")); err != nil { + return hubError(c, err) + } + return c.SendStatus(fiber.StatusNoContent) + }) + app.Post("/api/phase-a/accounts/:id/revoke", func(c fiber.Ctx) error { unlock, err := lockAccountResources(c.Context(), runtimeStore, c.Params("id")) if err != nil { diff --git a/internal/controlplane/app/app_test.go b/internal/controlplane/app/app_test.go index 44cc34b..e2c3341 100644 --- a/internal/controlplane/app/app_test.go +++ b/internal/controlplane/app/app_test.go @@ -362,6 +362,8 @@ func controlPlaneRouteMatrix() []controlPlaneRouteCase { {http.MethodGet, "/api/phase-a/accounts", "/api/phase-a/accounts", "", http.StatusOK}, {http.MethodGet, "/api/phase-a/accounts/:id", "/api/phase-a/accounts/missing", "", http.StatusNotFound}, {http.MethodPost, "/api/phase-a/accounts/:id/pause", "/api/phase-a/accounts/missing/pause", "", http.StatusNotFound}, + {http.MethodPost, "/api/phase-a/accounts/:id/environment", "/api/phase-a/accounts/missing/environment", "", http.StatusNotFound}, + {http.MethodPost, "/api/phase-a/accounts/:id/start", "/api/phase-a/accounts/missing/start", "", http.StatusNotFound}, {http.MethodPost, "/api/phase-a/accounts/:id/resume", "/api/phase-a/accounts/missing/resume", "", http.StatusNotFound}, {http.MethodPost, "/api/phase-a/accounts/:id/revoke", "/api/phase-a/accounts/missing/revoke", "", http.StatusNotFound},