Files
creator-hub/internal/controlplane/api/creator_internal_coverage_test.go
T

299 lines
12 KiB
Go

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