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 } }