feat(accounts): 账号即环境——创建即绑定、幂等补建、start 组合接口
- POST /api/phase-a/accounts 创建成功即绑定环境:alias=账号文本 ID、 唯一网关自动取(缺失 404/多网关 409)、出口直连、seed 由账号 ID 派生 (deriveSeed,数字主键落地后改 bigint id+1000);绑定失败 503 透传 account_id - POST /api/phase-a/accounts/:id/environment 幂等补建:已绑定原样返回 - POST /api/phase-a/accounts/:id/start = resume + 启动环境(环境缺失自愈补建, 审计对齐全);吊销账号 start → 409 readiness blocked;停止沿用 pause - 测试:deriveSeed 纯单测(确定性/值域/固定向量)+ PG 集成覆盖 绑定失败恢复、补建幂等、seed 断言、start 激活与冲突分支; RegisterAccountRoutes 参数升级为 HubStore,路由矩阵登记新端点
This commit is contained in:
@@ -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)
|
||||
}
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
@@ -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 {
|
||||
|
||||
@@ -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},
|
||||
|
||||
|
||||
Reference in New Issue
Block a user