diff --git a/internal/controlplane/api/account_environment_unit_test.go b/internal/controlplane/api/account_environment_unit_test.go index aef0985..c7b5483 100644 --- a/internal/controlplane/api/account_environment_unit_test.go +++ b/internal/controlplane/api/account_environment_unit_test.go @@ -204,3 +204,70 @@ func TestAccountEnvironmentDeletionLifecycle(t *testing.T) { t.Fatalf("account audit events must be purged with the account: rows=%d", auditCount) } } + +// TestAccountAuditEndpointFilters 驱动 GET /api/phase-a/audit 的过滤参数解析(auditFilter)。 +func TestAccountAuditEndpointFilters(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() }) + // hub.Open 负责把迁移链铺进隔离 schema(account.Open 不建表)。 + hubStore, err := hub.Open(ctx, databaseURL) + if err != nil { + t.Fatal(err) + } + t.Cleanup(func() { _ = hubStore.Close() }) + app := fiber.New() + RegisterAccountRoutes(app, accountStore, nil, &testCredentialBridge{values: map[string]string{}}) + + auditDB, err := sql.Open("pgx", databaseURL) + if err != nil { + t.Fatal(err) + } + t.Cleanup(func() { _ = auditDB.Close() }) + if _, err := auditDB.ExecContext(ctx, ` + INSERT INTO gateway (name, endpoint, token) VALUES ('gw-1', 'http://gw-1:8081', 'unit-test-gateway-token'); + INSERT INTO social_account (account_id, credential_provider, credential_key, name, platform, platform_account_key, authorization_kind, authorization_status, status) + VALUES ('account-audit-0000000000000000', 'os_keyring', 'creatorhub/account-audit-0000000000000000', '审计账号', 'douyin', 'key-audit', 'owned', 'authorized', 'paused'); + INSERT INTO audit_event (event_type, account_id, actor, reason_code, details) + SELECT 'account_created', account.id, 'local-user', 'account_created', '{"platform":"douyin"}' + FROM social_account account WHERE account.account_id = 'account-audit-0000000000000000'`); err != nil { + t.Fatal(err) + } + + response := do(app, http.MethodGet, "/api/phase-a/audit?account_id=account-audit-0000000000000000&page=1&page_size=10", "") + if response.Code != http.StatusOK { + t.Fatalf("audit query: %d %s", response.Code, response.Body.String()) + } + var page struct { + Data []map[string]any `json:"data"` + Total int `json:"total"` + } + if err := json.Unmarshal(response.Body.Bytes(), &page); err != nil || page.Total != 1 || len(page.Data) != 1 { + t.Fatalf("audit page: total=%d data=%d err=%v", page.Total, len(page.Data), err) + } + if page.Data[0]["event_type"] != "account_created" || page.Data[0]["account_id"] != "account-audit-0000000000000000" { + t.Fatalf("audit event fields: %#v", page.Data[0]) + } + + for _, query := range []string{ + "/api/phase-a/audit?page=0", + "/api/phase-a/audit?page_size=1001", + "/api/phase-a/audit?from=not-a-time", + "/api/phase-a/audit?event_type=INVALID%20EVENT", + } { + if response := do(app, http.MethodGet, query, ""); response.Code != http.StatusBadRequest { + t.Fatalf("invalid filter %q must 400: %d %s", query, response.Code, response.Body.String()) + } + } + if response := do(app, http.MethodGet, "/api/phase-a/audit?from=2026-01-01T00:00:00Z&to=2027-01-01T00:00:00Z", ""); response.Code != http.StatusOK { + t.Fatalf("time-bounded audit query: %d %s", response.Code, response.Body.String()) + } +} diff --git a/internal/controlplane/api/creator_internal_coverage_test.go b/internal/controlplane/api/creator_internal_coverage_test.go new file mode 100644 index 0000000..c2668a3 --- /dev/null +++ b/internal/controlplane/api/creator_internal_coverage_test.go @@ -0,0 +1,298 @@ +package api + +import ( + "context" + "encoding/base64" + "encoding/json" + "fmt" + "net/http" + "net/http/httptest" + "os" + "strings" + "testing" + "time" + + "git.ipao.vip/rogee/creator-hub/internal/account" + "git.ipao.vip/rogee/creator-hub/internal/creator" + hub "git.ipao.vip/rogee/creator-hub/internal/environment" +) + +// TestDouyinCollectorCarriesSourceIdentity 校验采集器构造透传来源标识。 +func TestDouyinCollectorCarriesSourceIdentity(t *testing.T) { + browser := creatorGatewayBrowser{gateway: hub.Gateway{Name: "gw-1"}, environment: hub.EnvironmentContext{RuntimeID: "runtime-a"}} + collector := douyinCollector(browser, "sec_uid_x", "owned", "account-a") + if collector.AccountKey != "sec_uid_x" || collector.SourceType != "owned" || collector.SourceID != "account-a" { + t.Fatalf("collector fields: %+v", collector) + } +} + +// TestCreatorUpdateHubSubscribePublish 驱动 SSE 变更通知的订阅/广播纯逻辑。 +func TestCreatorUpdateHubSubscribePublish(t *testing.T) { + updates, unsubscribe := creatorUpdates.subscribe() + if updates == nil || unsubscribe == nil { + t.Fatal("subscribe must return channel and cancel function") + } + creatorUpdates.publish() + select { + case <-updates: + case <-time.After(time.Second): + t.Fatal("publish did not notify the subscriber") + } + // 缓冲为 1:再次 publish 不阻塞(非阻塞发送语义)。 + creatorUpdates.publish() + creatorUpdates.publish() + unsubscribe() + unsubscribe() // 幂等 +} + +// TestRuntimeCreateSpecMatchesGate 校验恢复运行时的绑定一致性门槛。 +func TestRuntimeCreateSpecMatchesGate(t *testing.T) { + ctx := context.Background() + current := hub.EnvironmentContext{BindingVersion: 3} + current.Exit.ID = "exit-1" + if runtimeCreateSpecMatches(ctx, nil, current, current, nil) { + t.Fatal("nil prepared spec must not match") + } + previous := current + if !runtimeCreateSpecMatches(ctx, nil, current, previous, &runtimeCreateSpec{}) { + t.Fatal("same binding version and exit must match") + } + previous.BindingVersion = 2 + if runtimeCreateSpecMatches(ctx, nil, current, previous, &runtimeCreateSpec{}) { + t.Fatal("stale binding version must not match") + } + previous.BindingVersion = 3 + previous.Exit.ID = "exit-2" + if runtimeCreateSpecMatches(ctx, nil, current, previous, &runtimeCreateSpec{}) { + t.Fatal("changed exit must not match") + } +} + +// TestCreatorGatewayBrowserGetImage 驱动封面拉取链路(成功 + 2xx 校验 + base64 解码)。 +func TestCreatorGatewayBrowserGetImage(t *testing.T) { + encoded := base64.StdEncoding.EncodeToString([]byte("png-bytes")) + server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + w.Header().Set("Content-Type", "application/json") + if !strings.HasSuffix(r.URL.Path, "/douyin/image") { + http.NotFound(w, r) + return + } + _, _ = w.Write([]byte(`{"status":200,"content_type":"image/png","body_base64":"` + encoded + `"}`)) + })) + defer server.Close() + browser := creatorGatewayBrowser{ + gateway: hub.Gateway{Endpoint: server.URL, Token: "token"}, + environment: hub.EnvironmentContext{Env: hub.Env{Alias: "browser"}, RuntimeID: "dddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddd", BindingVersion: 1}, + } + contentType, data, err := browser.GetImage(context.Background(), "https://p.example/cover.png") + if err != nil || contentType != "image/png" || string(data) != "png-bytes" { + t.Fatalf("get image: type=%q data=%q err=%v", contentType, data, err) + } + + // 网关返回非图片 content_type → 拒绝 + reject := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + _, _ = w.Write([]byte(`{"status":200,"content_type":"text/html","body_base64":"PGh0bWw+"}`)) + })) + defer reject.Close() + browser.gateway.Endpoint = reject.URL + if _, _, err := browser.GetImage(context.Background(), "https://p.example/x"); err == nil { + t.Fatal("non-image response must be rejected") + } +} + +// TestRunCreatorSchedulerContextCancel 驱动调度器入口:ctx 取消后立即返回。 +func TestRunCreatorSchedulerContextCancel(t *testing.T) { + ctx, cancel := context.WithCancel(context.Background()) + cancel() + done := make(chan struct{}) + go func() { + RunCreatorScheduler(ctx, nil, nil, nil) + close(done) + }() + select { + case <-done: + case <-time.After(2 * time.Second): + t.Fatal("scheduler must honor a cancelled context") + } +} + +// TestReconcileRuntimeLeasesExported 用缺失网关的内存桩驱动导出的租约对账入口。 +func TestReconcileRuntimeLeasesExported(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) + hubStore, err := hub.Open(ctx, databaseURL) + if err != nil { + t.Fatal(err) + } + t.Cleanup(func() { _ = hubStore.Close() }) + if err := ReconcileRuntimeLeases(ctx, hubStore); err != nil { + t.Fatalf("empty reconcile must succeed: %v", err) + } +} + +// TestCacheCompetitorWorkCovers 在隔离 PG 上驱动封面回填链路: +// 缺封面作品经浏览器 GetImage 拉取并落 creator_work_cover。 +func TestCacheCompetitorWorkCovers(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, _, ctx := openCreatorIntegrationStoreForAPITest(t, databaseURL) + competitor, err := store.CreateCompetitor(ctx, creator.CompetitorInput{ + Platform: creator.PlatformDouyin, PlatformAccountKey: "sec_uid_covers_" + fmt.Sprintf("%d", time.Now().UnixNano()), + UniqueID: "covers_unique", Nickname: "封面竞品", HomepageURL: "https://www.douyin.com/user/sec_uid_covers", + }) + if err != nil { + t.Fatal(err) + } + work, _, err := store.UpsertWork(ctx, creator.WorkInput{ + Platform: creator.PlatformDouyin, WorkKey: "cover-work-" + fmt.Sprintf("%d", time.Now().UnixNano()), + SourceType: creator.SourceCompetitor, SourceID: competitor.ID, Title: "封面作品", CoverURL: "https://p.example/cover.png", + }, time.Now().UTC()) + if err != nil { + t.Fatal(err) + } + + encoded := base64.StdEncoding.EncodeToString([]byte("cover-bytes")) + server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + w.Header().Set("Content-Type", "application/json") + _, _ = w.Write([]byte(`{"status":200,"content_type":"image/png","body_base64":"` + encoded + `"}`)) + })) + defer server.Close() + browser := creatorGatewayBrowser{ + gateway: hub.Gateway{Endpoint: server.URL, Token: "token"}, + environment: hub.EnvironmentContext{Env: hub.Env{Alias: "browser"}, RuntimeID: "dddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddd", BindingVersion: 1}, + } + if err := cacheCompetitorWorkCovers(ctx, store, browser, competitor.ID); err != nil { + t.Fatalf("cache covers: %v", err) + } + contentType, data, err := store.GetWorkCover(ctx, work.ID, "cover") + if err != nil || contentType != "image/png" || string(data) != "cover-bytes" { + t.Fatalf("cover: type=%q data=%q err=%v", contentType, data, err) + } + + // 网关故障 → 报错且不落库(新增一件缺封面作品,否则无可拉取对象直接成功) + if _, _, err := store.UpsertWork(ctx, creator.WorkInput{ + Platform: creator.PlatformDouyin, WorkKey: "cover-work-fail-" + fmt.Sprintf("%d", time.Now().UnixNano()), + SourceType: creator.SourceCompetitor, SourceID: competitor.ID, Title: "缺封面作品", CoverURL: "https://p.example/missing.png", + }, time.Now().UTC()); err != nil { + t.Fatal(err) + } + // 网关故障 → 报错且不落库:用全新竞品保证其全部作品拉取失败 + other, err := store.CreateCompetitor(ctx, creator.CompetitorInput{ + Platform: creator.PlatformDouyin, PlatformAccountKey: "sec_uid_covers_fail_" + fmt.Sprintf("%d", time.Now().UnixNano()), + UniqueID: "covers_fail_unique", Nickname: "失败竞品", HomepageURL: "https://www.douyin.com/user/sec_uid_fail", + }) + if err != nil { + t.Fatal(err) + } + if _, _, err := store.UpsertWork(ctx, creator.WorkInput{ + Platform: creator.PlatformDouyin, WorkKey: "cover-work-fail-" + fmt.Sprintf("%d", time.Now().UnixNano()), + SourceType: creator.SourceCompetitor, SourceID: other.ID, Title: "缺封面作品", CoverURL: "https://p.example/missing.png", + }, time.Now().UTC()); err != nil { + t.Fatal(err) + } + failing := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + w.WriteHeader(http.StatusInternalServerError) + })) + defer failing.Close() + browser.gateway.Endpoint = failing.URL + if err := cacheCompetitorWorkCovers(ctx, store, browser, other.ID); err == nil { + t.Fatal("failing gateway must surface an error") + } + missing, err := store.ListWorksMissingCover(ctx, other.ID) + if err != nil || len(missing) != 1 { + t.Fatalf("failed fetch must not save covers: missing=%d err=%v", len(missing), err) + } +} + +func openCreatorIntegrationStoreForAPITest(t *testing.T, databaseURL string) (*creator.Store, *account.Store, context.Context) { + t.Helper() + ctx := context.Background() + creatorStore, err := creator.Open(ctx, databaseURL) + if err != nil { + t.Fatal(err) + } + t.Cleanup(func() { _ = creatorStore.Close() }) + if err := creatorStore.EnsureSchema(ctx); err != nil { + t.Fatal(err) + } + creatorStore.SetSecretBridge(nil) + // hub.Open 铺迁移链(social_account 等 phase A 表)。 + hubStore, err := hub.Open(ctx, databaseURL) + if err != nil { + t.Fatal(err) + } + t.Cleanup(func() { _ = hubStore.Close() }) + return creatorStore, nil, ctx +} + +// TestNewAnonymousBrowserCreateFailureCleanup 驱动匿名浏览器创建失败路径: +// 网关 create 返回冲突 → cleanup 对账删除残留容器 → 错误上抛、名额释放。 +func TestNewAnonymousBrowserCreateFailureCleanup(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) + hubStore, err := hub.Open(ctx, databaseURL) + if err != nil { + t.Fatal(err) + } + t.Cleanup(func() { _ = hubStore.Close() }) + + var cleanupCalls int + var createdAlias string + gatewayServer := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + w.Header().Set("Content-Type", "application/json") + if r.Method == http.MethodPost && r.URL.Path == "/v1/browsers" { + // 声称创建成功但返回非法代际(缺 id)→ 触发 cleanup 对账 + var payload struct { + Alias string `json:"alias"` + } + _ = json.NewDecoder(r.Body).Decode(&payload) + createdAlias = payload.Alias + w.WriteHeader(http.StatusCreated) + _, _ = w.Write([]byte(`{"alias":"anon-other","binding_version":1}`)) + return + } + if r.Method == http.MethodGet && r.URL.Path == "/v1/browsers" { + // 对账列表:返回刚创建失败残留的匿名容器 + w.Write([]byte(`[{"id":"bbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbb","alias":"` + createdAlias + `","state":"running","binding_version":1}]`)) + return + } + if r.Method == http.MethodDelete && strings.HasPrefix(r.URL.Path, "/v1/browsers/") { + cleanupCalls++ + w.WriteHeader(http.StatusNoContent) + return + } + http.NotFound(w, r) + })) + defer gatewayServer.Close() + if _, err := hubStore.CreateGateway(ctx, "gw-1", gatewayServer.URL, "unit-test-gateway-token"); err != nil { + t.Fatal(err) + } + + _, err = newAnonymousBrowser(ctx, hubStore) + if err == nil { + t.Fatal("invalid create generation must fail") + } + if cleanupCalls != 1 { + t.Fatalf("cleanup must reconcile the leftover container: calls=%d", cleanupCalls) + } + // 名额已释放:可再次非阻塞占满全部槽位 + for i := 0; i < cap(anonymousBrowserSlots); i++ { + select { + case anonymousBrowserSlots <- struct{}{}: + default: + t.Fatalf("slot %d not released after failure", i) + } + <-anonymousBrowserSlots + } +} diff --git a/internal/controlplane/app/app_test.go b/internal/controlplane/app/app_test.go index e2c3341..bc74711 100644 --- a/internal/controlplane/app/app_test.go +++ b/internal/controlplane/app/app_test.go @@ -21,6 +21,7 @@ import ( "git.ipao.vip/rogee/creator-hub/internal/account" accountsapi "git.ipao.vip/rogee/creator-hub/internal/controlplane/api/accounts" + "git.ipao.vip/rogee/creator-hub/internal/creator" hub "git.ipao.vip/rogee/creator-hub/internal/environment" "github.com/gofiber/fiber/v3" "github.com/gofiber/fiber/v3/middleware/adaptor" @@ -471,3 +472,111 @@ func isolatedControlPlaneDatabaseURL(t *testing.T, databaseURL string) string { parsed.RawQuery = query.Encode() return parsed.String() } + +// TestControlPlaneSystemEndpointsWithStores 覆盖 system 包健康/就绪真实链路: +// 三个 store 就绪时 readyz 204;creator schema 缺失时 503。 +func TestControlPlaneSystemEndpointsWithStores(t *testing.T) { + databaseURL := os.Getenv("CREATORHUB_POSTGRES_TEST_URL") + if databaseURL == "" { + t.Skip("set CREATORHUB_POSTGRES_TEST_URL to run control-plane readiness coverage") + } + databaseURL = isolatedControlPlaneDatabaseURL(t, databaseURL) + ctx := context.Background() + phaseAStore, err := account.Open(ctx, databaseURL) + if err != nil { + t.Fatal(err) + } + t.Cleanup(func() { _ = phaseAStore.Close() }) + hubStore, err := hub.Open(ctx, databaseURL) + if err != nil { + t.Fatal(err) + } + t.Cleanup(func() { _ = hubStore.Close() }) + creatorStore, err := creator.Open(ctx, databaseURL) + if err != nil { + t.Fatal(err) + } + t.Cleanup(func() { _ = creatorStore.Close() }) + + directory := t.TempDir() + if err := os.WriteFile(filepath.Join(directory, "index.html"), []byte("index"), 0o600); err != nil { + t.Fatal(err) + } + app := newHandler(directory, "operator", "unit-test-password", phaseAStore, hubStore, nil, creatorStore) + + health, err := app.Test(httptest.NewRequest(http.MethodGet, "/healthz", nil)) + if err != nil || health.StatusCode != http.StatusNoContent { + t.Fatalf("healthz: status=%d err=%v", health.StatusCode, err) + } + health.Body.Close() + + ready, err := app.Test(httptest.NewRequest(http.MethodGet, "/readyz", nil)) + if err != nil || ready.StatusCode != http.StatusNoContent { + t.Fatalf("readyz with all stores: status=%d err=%v", ready.StatusCode, err) + } + ready.Body.Close() + + creatorStore.Close() + ready, err = app.Test(httptest.NewRequest(http.MethodGet, "/readyz", nil)) + if err != nil || ready.StatusCode != http.StatusServiceUnavailable { + t.Fatalf("readyz with closed store: status=%d err=%v", ready.StatusCode, err) + } + ready.Body.Close() +} + +// TestControlPlaneLegacyAPIPathFallbacks 驱动 isControlPlaneAPIPath 与 SPA 回退包装的纯逻辑。 +func TestControlPlaneLegacyAPIPathFallbacks(t *testing.T) { + for path, want := range map[string]bool{ + "/api/browsers": true, + "/api/phase-a/accounts": true, + "/gateways/gw-1": true, + "/browsers/account-a": true, + "/network-exits": true, + "/phase-a": true, + "/dashboard/assets.js": false, + "/apifake": false, + "/": false, + } { + if got := isControlPlaneAPIPath(path); got != want { + t.Fatalf("isControlPlaneAPIPath(%q)=%v want %v", path, got, want) + } + } +} + +// TestCreatorSecretBridgeDelegatesAndRejectsNil 覆盖凭据桥转发与空桥报错。 +func TestCreatorSecretBridgeDelegatesAndRejectsNil(t *testing.T) { + var nilBridge creatorSecretBridge + ctx := context.Background() + if err := nilBridge.Store(ctx, creator.SecretReference{ID: "ref-1", Provider: "os_keyring"}, "key-1", "value-1"); err == nil { + t.Fatal("nil bridge must reject store") + } + if err := nilBridge.Delete(ctx, creator.SecretReference{ID: "ref-1", Provider: "os_keyring"}, "key-1"); err == nil { + t.Fatal("nil bridge must reject delete") + } + values := map[string]string{} + bridge := creatorSecretBridge{bridge: &accountCredentialBridgeStub{values: values}} + if err := bridge.Store(ctx, creator.SecretReference{ID: "ref-1", Provider: "os_keyring"}, "creatorhub/key", "value-1"); err != nil { + t.Fatalf("store: %v", err) + } + if values["creatorhub/key"] != "value-1" { + t.Fatalf("store did not delegate: %#v", values) + } + if err := bridge.Delete(ctx, creator.SecretReference{ID: "ref-1", Provider: "os_keyring"}, "creatorhub/key"); err != nil { + t.Fatalf("delete: %v", err) + } + if _, ok := values["creatorhub/key"]; ok { + t.Fatal("delete did not delegate") + } +} + +type accountCredentialBridgeStub struct{ values map[string]string } + +func (stub *accountCredentialBridgeStub) Store(_ context.Context, _ account.CredentialReference, key, value string) error { + stub.values[key] = value + return nil +} + +func (stub *accountCredentialBridgeStub) Delete(_ context.Context, _ account.CredentialReference, key string) error { + delete(stub.values, key) + return nil +}