Files
creator-hub/internal/creator/integration_test.go
T
rogee 998871e3c4
douyin-release-gate / verify (push) Failing after 4m57s
fix(creator): 强制同步抢锁失败分类提示——租约被占/竞品停用分别返回可读 409
- ClaimCompetitorSync 失败后经 SyncClaimFailure 复查竞品状态,
  区分 ErrSyncInProgress(该竞品正在同步中)与 ErrSyncDisabled(已停用)
- creatorError 新增两分类映射 409;通用冲突文案保持不变
- 前端竞品强制同步去掉固定'强制同步失败'后缀,直接展示分类提示
- 回归:TestCompetitorSyncClaimFailureClassification(PG 集成)
2026-09-29 15:11:23 +08:00

721 lines
36 KiB
Go
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
package creator
import (
"context"
"database/sql"
"errors"
"fmt"
"net/url"
"os"
"testing"
"time"
"git.ipao.vip/rogee/creator-hub/internal/account"
hub "git.ipao.vip/rogee/creator-hub/internal/environment"
_ "github.com/jackc/pgx/v5/stdlib"
)
type integrationCredentialBridge struct{}
func (integrationCredentialBridge) Store(context.Context, account.CredentialReference, string, string) error {
return nil
}
func (integrationCredentialBridge) Delete(context.Context, account.CredentialReference, string) error {
return nil
}
type integrationAnalyzer struct{}
func (integrationAnalyzer) MatchTheme(context.Context, string, string, string) (bool, string, error) {
return true, "主题匹配", nil
}
func (integrationAnalyzer) MatchLead(context.Context, string, string, string) (bool, string, error) {
return true, "具备线索意向", nil
}
type integrationExecutor struct{}
func (integrationExecutor) Execute(context.Context, ActionRequest) (ActionResult, error) {
return ActionResult{State: "succeeded", Evidence: map[string]string{"platform_id": "creator-it"}}, nil
}
type integrationCollector struct {
work WorkInput
comment CommentInput
}
func (c integrationCollector) VerifyIdentity(context.Context, string) error { return nil }
func (c integrationCollector) ListWorks(_ context.Context, _, cursor string) (WorkPage, error) {
if cursor != "" {
return WorkPage{Items: nil, HasMore: false}, nil
}
return WorkPage{Items: []WorkInput{c.work}, HasMore: true, NextCursor: "1"}, nil
}
func (c integrationCollector) ListTopLevelComments(_ context.Context, _, cursor string) (CommentPage, error) {
if cursor != "" {
return CommentPage{Items: nil, HasMore: false}, nil
}
return CommentPage{Items: []CommentInput{c.comment}, HasMore: false}, nil
}
func openCreatorIntegrationStore(t *testing.T) (*Store, *account.Store, context.Context) {
t.Helper()
databaseURL := os.Getenv("CREATORHUB_POSTGRES_TEST_URL")
if databaseURL == "" {
t.Skip("set CREATORHUB_POSTGRES_TEST_URL to run CreatorHub PostgreSQL integration coverage")
}
ctx := context.Background()
admin, err := sql.Open("pgx", databaseURL)
if err != nil {
t.Fatal(err)
}
t.Cleanup(func() { _ = admin.Close() })
if err := admin.PingContext(ctx); err != nil {
t.Fatal(err)
}
schema := fmt.Sprintf("creatorhub_it_%d", time.Now().UnixNano())
if _, err := admin.ExecContext(ctx, "CREATE SCHEMA "+schema); err != nil {
t.Fatal(err)
}
t.Cleanup(func() {
if _, err := admin.ExecContext(ctx, "DROP SCHEMA "+schema+" CASCADE"); err != nil {
t.Errorf("drop test schema: %v", err)
}
})
parsed, err := url.Parse(databaseURL)
if err != nil {
t.Fatal(err)
}
query := parsed.Query()
query.Set("search_path", schema)
parsed.RawQuery = query.Encode()
testURL := parsed.String()
phaseAStore, err := account.Open(ctx, testURL)
if err != nil {
t.Fatal(err)
}
t.Cleanup(func() { _ = phaseAStore.Close() })
hubStore, err := hub.Open(ctx, testURL)
if err != nil {
t.Fatal(err)
}
t.Cleanup(func() { _ = hubStore.Close() })
store, err := Open(ctx, testURL)
if err != nil {
t.Fatal(err)
}
t.Cleanup(func() { _ = store.Close() })
return store, phaseAStore, ctx
}
func createIntegrationAccount(t *testing.T, ctx context.Context, phaseAStore *account.Store, suffix string) string {
t.Helper()
id := "cit" + suffix
account := account.Account{
ID: id,
Name: "Creator integration " + suffix,
Platform: PlatformDouyin,
PlatformAccountKey: "sec_uid_" + id,
Tags: []string{"integration"},
CredentialReference: account.CredentialReference{ID: "credential_" + id, Provider: "os_keyring"},
CredentialKey: "creatorhub/" + id,
}
if err := phaseAStore.CreateAccount(ctx, account, integrationCredentialBridge{}); err != nil {
t.Fatal(err)
}
return id
}
func TestCompetitorShareJobLifecycle(t *testing.T) {
store, _, ctx := openCreatorIntegrationStore(t)
stamp := fmt.Sprintf("%d", time.Now().UnixNano())
job, err := store.CreateCompetitorShareJob(ctx, CompetitorShareJobInput{
Platform: PlatformDouyin,
ShareURL: "https://v.douyin.com/share-" + stamp + "/",
Tags: []string{"重点监测"},
})
if err != nil || job.Status != CompetitorShareJobQueued || job.Attempts != 0 {
t.Fatalf("create share job: job=%+v err=%v", job, err)
}
due, err := store.ListDueCompetitorShareJobs(ctx, time.Now().UTC())
if err != nil || len(due) != 1 || due[0].ID != job.ID {
t.Fatalf("list due share jobs: jobs=%+v err=%v", due, err)
}
claimed, leaseToken, ok, err := store.ClaimCompetitorShareJob(ctx, job.ID, time.Now().UTC())
if err != nil || !ok || leaseToken == "" || claimed.Attempts != 1 || claimed.Status != CompetitorShareJobProcessing {
t.Fatalf("claim share job: job=%+v token=%q claimed=%v err=%v", claimed, leaseToken, ok, err)
}
if err := store.MarkCompetitorShareJob(ctx, job.ID, leaseToken, CompetitorShareJobQueued, "", "temporary failure"); err != nil {
t.Fatal(err)
}
claimed, leaseToken, ok, err = store.ClaimCompetitorShareJob(ctx, job.ID, time.Now().UTC())
if err != nil || !ok || leaseToken == "" || claimed.Attempts != 2 {
t.Fatalf("claim retry share job: job=%+v token=%q claimed=%v err=%v", claimed, leaseToken, ok, err)
}
if err := store.MarkCompetitorShareJob(ctx, job.ID, leaseToken, CompetitorShareJobFailed, "", "final failure"); err != nil {
t.Fatal(err)
}
failed, err := store.GetCompetitorShareJob(ctx, job.ID)
if err != nil || failed.Status != CompetitorShareJobFailed || failed.Attempts != 2 || failed.FailureReason != "final failure" {
t.Fatalf("failed share job: job=%+v err=%v", failed, err)
}
requeued, err := store.RetryCompetitorShareJob(ctx, job.ID)
if err != nil || requeued.Status != CompetitorShareJobQueued || requeued.Attempts != 0 || requeued.FailureReason != "final failure" {
t.Fatalf("requeue share job: job=%+v err=%v", requeued, err)
}
}
// 强制同步抢锁失败的分类:租约被占 → ErrSyncInProgress;竞品停用 → ErrSyncDisabled。
func TestCompetitorSyncClaimFailureClassification(t *testing.T) {
store, _, ctx := openCreatorIntegrationStore(t)
stamp := fmt.Sprintf("%d", time.Now().UnixNano())
competitor, err := store.CreateCompetitor(ctx, CompetitorInput{
Platform: PlatformDouyin, PlatformAccountKey: "sec_uid_claim_" + stamp, UniqueID: "claim_" + stamp,
Nickname: "Claim", HomepageURL: "https://www.douyin.com/user/sec_uid_claim_" + stamp,
})
if err != nil {
t.Fatal(err)
}
now := time.Now().UTC()
// 首次 force 抢锁成功;租约未释放前第二次 force 抢锁失败 → 正在同步中。
token, claimed, err := store.ClaimCompetitorSync(ctx, competitor.ID, true, now)
if err != nil || !claimed {
t.Fatalf("first force claim must succeed: err=%v claimed=%v", err, claimed)
}
if _, claimed, err := store.ClaimCompetitorSync(ctx, competitor.ID, true, now.Add(time.Second)); err != nil || claimed {
t.Fatalf("second force claim must fail while lease held: err=%v claimed=%v", err, claimed)
}
if err := store.SyncClaimFailure(ctx, competitor.ID); !errors.Is(err, ErrSyncInProgress) {
t.Fatalf("held lease must classify as sync in progress: err=%v", err)
}
// 释放后停用 → 已停用。
if err := store.MarkCompetitorSync(ctx, competitor.ID, token, "idle", "", "", nil); err != nil {
t.Fatal(err)
}
if _, err := store.SetCompetitorEnabled(ctx, competitor.ID, false); err != nil {
t.Fatal(err)
}
if _, claimed, err := store.ClaimCompetitorSync(ctx, competitor.ID, true, now.Add(2*time.Second)); err != nil || claimed {
t.Fatalf("disabled competitor must not be claimable: err=%v claimed=%v", err, claimed)
}
if err := store.SyncClaimFailure(ctx, competitor.ID); !errors.Is(err, ErrSyncDisabled) {
t.Fatalf("disabled competitor must classify as disabled: err=%v", err)
}
}
func TestCompetitorUpsertIsIdempotent(t *testing.T) {
store, _, ctx := openCreatorIntegrationStore(t)
stamp := fmt.Sprintf("%d", time.Now().UnixNano())
input := CompetitorInput{
Platform: PlatformDouyin,
PlatformAccountKey: "sec_uid_upsert_" + stamp,
UniqueID: "upsert_" + stamp,
Nickname: "Competitor",
HomepageURL: "https://www.douyin.com/user/sec_uid_upsert_" + stamp,
Tags: []string{"重点监测"},
}
first, err := store.UpsertCompetitor(ctx, input)
if err != nil {
t.Fatal(err)
}
second, err := store.UpsertCompetitor(ctx, CompetitorInput{
Platform: input.Platform,
PlatformAccountKey: input.PlatformAccountKey,
Nickname: "Updated Competitor",
HomepageURL: input.HomepageURL,
})
if err != nil || second.ID != first.ID || second.Nickname != "Updated Competitor" || len(second.Tags) != 1 {
t.Fatalf("upsert competitor: first=%+v second=%+v err=%v", first, second, err)
}
}
func TestCreatorPostgresContentAndWorkflow(t *testing.T) {
store, phaseAStore, ctx := openCreatorIntegrationStore(t)
stamp := fmt.Sprintf("%d", time.Now().UnixNano())
bigID := createIntegrationAccount(t, ctx, phaseAStore, "big"+stamp)
smallID := createIntegrationAccount(t, ctx, phaseAStore, "small"+stamp)
if err := store.EnsureAccountProfile(ctx, bigID); err != nil {
t.Fatal(err)
}
if err := store.EnsureAccountProfile(ctx, smallID); err != nil {
t.Fatal(err)
}
profiles, err := store.ListAccountProfiles(ctx)
if err != nil || len(profiles) != 2 {
t.Fatalf("list account profiles: profiles=%+v err=%v", profiles, err)
}
updatedTags, err := store.UpdateAccountTags(ctx, bigID, []string{"主账号"})
if err != nil || len(updatedTags) != 1 || updatedTags[0] != "主账号" {
t.Fatalf("update account tags: tags=%v err=%v", updatedTags, err)
}
if _, err := store.UpdateAccountProfile(ctx, bigID, AccountProfileUpdate{RealNameStatus: "unknown", BusinessStatus: "normal", BigAccount: true, ReplyRequirements: "保持准确", CooldownSeconds: 86400}); err != nil {
t.Fatal(err)
}
if _, err := store.UpdateAccountProfile(ctx, smallID, AccountProfileUpdate{RealNameStatus: "unknown", BusinessStatus: "normal", CooldownSeconds: 86400}); err != nil {
t.Fatal(err)
}
if _, err := store.RecordVerifiedLoginResult(ctx, bigID, "sec_uid_"+bigID); err != nil {
t.Fatal(err)
}
if _, err := store.RecordVerifiedLoginResult(ctx, smallID, "sec_uid_"+smallID); err != nil {
t.Fatal(err)
}
now := time.Now().UTC().Truncate(time.Microsecond)
ownedDue, err := store.ListDueOwnedAccounts(ctx, now, 1800)
if err != nil || len(ownedDue) != 2 || ownedDue[0] != bigID || ownedDue[1] != smallID {
t.Fatalf("list due owned accounts: accounts=%+v err=%v", ownedDue, err)
}
published := now.Add(-2 * time.Hour)
likes, comments, shares := int64(10), int64(2), int64(1)
plays := int64(777)
workKey := "creator-it-work-" + stamp
work, inserted, err := store.UpsertWork(ctx, WorkInput{Platform: PlatformDouyin, WorkKey: workKey, SourceType: SourceOwned, SourceID: bigID, Title: "Topic", Body: "Need consultation", PublishedAt: &published, PublishedAtStatus: "verified", OriginalURL: "https://www.douyin.com/video/" + workKey, Likes: &likes, CommentsCount: &comments, Shares: &shares, PlayCount: &plays}, now)
if err != nil || !inserted {
t.Fatalf("insert work: work=%+v inserted=%v err=%v", work, inserted, err)
}
if work.PlayCount == nil || *work.PlayCount != 777 {
t.Fatalf("play_count not persisted on insert: %+v", work)
}
// COALESCE 语义:二次 upsert 未带 play_count 时保留旧值。
if _, _, err := store.UpsertWork(ctx, WorkInput{Platform: PlatformDouyin, WorkKey: workKey, SourceType: SourceOwned, SourceID: bigID, Title: "Topic-renamed", PublishedAt: &published, PublishedAtStatus: "verified"}, now); err != nil {
t.Fatal(err)
}
retained, err := store.GetWorkByKey(ctx, PlatformDouyin, workKey)
if err != nil || retained.PlayCount == nil || *retained.PlayCount != 777 {
t.Fatalf("play_count must be retained via COALESCE: work=%+v err=%v", retained, err)
}
// 负播放量拒绝。
if _, _, err := store.UpsertWork(ctx, WorkInput{Platform: PlatformDouyin, WorkKey: workKey + "-neg", SourceType: SourceOwned, SourceID: bigID, Title: "neg", PlayCount: &[]int64{-1}[0]}, now); !errors.Is(err, ErrInvalid) {
t.Fatalf("negative play_count must be rejected: err=%v", err)
}
competitor, err := store.CreateCompetitor(ctx, CompetitorInput{Platform: PlatformDouyin, PlatformAccountKey: "sec_uid_competitor_" + stamp, UniqueID: "competitor_" + stamp, Nickname: "Competitor", HomepageURL: "https://www.douyin.com/user/sec_uid_competitor_" + stamp, Tags: []string{"重点监测"}})
if err != nil {
t.Fatal(err)
}
if competitor.UniqueID != "competitor_"+stamp || len(competitor.Tags) != 1 || competitor.Tags[0] != "重点监测" {
t.Fatalf("create competitor tags: %+v", competitor.Tags)
}
updatedCompetitor, err := store.UpdateCompetitorTags(ctx, competitor.ID, []string{"已分类"})
if err != nil || len(updatedCompetitor.Tags) != 1 || updatedCompetitor.Tags[0] != "已分类" {
t.Fatalf("update competitor tags: competitor=%+v err=%v", updatedCompetitor, err)
}
competitors, err := store.ListCompetitorsWithProfile(ctx, PlatformDouyin)
if err != nil || len(competitors) != 1 || len(competitors[0].Tags) != 1 || competitors[0].Tags[0] != "已分类" || competitors[0].WorkCount != 0 {
t.Fatalf("list competitors: competitors=%+v err=%v", competitors, err)
}
if _, err := store.SetCompetitorEnabled(ctx, competitor.ID, false); err != nil {
t.Fatal(err)
}
if _, err := store.SetCompetitorEnabled(ctx, competitor.ID, true); err != nil {
t.Fatal(err)
}
dueNow := time.Now().UTC().Add(time.Minute)
if due, err := store.ListDueCompetitors(ctx, dueNow); err != nil || len(due) != 1 {
t.Fatalf("list due competitors: due=%+v err=%v", due, err)
}
leaseToken, claimed, err := store.ClaimCompetitorSync(ctx, competitor.ID, false, dueNow)
if err != nil || !claimed || leaseToken == "" {
t.Fatalf("claim competitor sync: token=%q claimed=%v err=%v", leaseToken, claimed, err)
}
if err := store.MarkCompetitorSync(ctx, competitor.ID, leaseToken, "idle", "", "", nil); err != nil {
t.Fatal(err)
}
// blocked 同步也必须排下一次重试:否则 next_sync_at 置空后定时获取对该账号永久失效。
blockedToken, blockedClaimed, err := store.ClaimCompetitorSync(ctx, competitor.ID, true, dueNow)
if err != nil || !blockedClaimed {
t.Fatalf("claim competitor sync for blocked: err=%v claimed=%v", err, blockedClaimed)
}
blockedNext := dueNow.Add(2 * time.Hour)
if err := store.MarkCompetitorSync(ctx, competitor.ID, blockedToken, "blocked", "", "runtime unavailable", &blockedNext); err != nil {
t.Fatal(err)
}
if due, err := store.ListDueCompetitors(ctx, blockedNext); err != nil || len(due) != 1 || due[0].ID != competitor.ID {
t.Fatalf("blocked competitor must reschedule: due=%+v err=%v", due, err)
}
if due, err := store.ListDueCompetitors(ctx, dueNow); err != nil || len(due) != 0 {
t.Fatalf("blocked competitor must not fire before retry time: due=%+v err=%v", due, err)
}
// 强制同步语义:next_sync_at 未到期时 force claim 仍立即执行;完成后 MarkCompetitorSync 刷新 next_sync_at。
if _, claimed, err := store.ClaimCompetitorSync(ctx, competitor.ID, false, dueNow); err != nil || claimed {
t.Fatalf("non-forced claim must wait for next_sync_at: claimed=%v err=%v", claimed, err)
}
forceToken, forceClaimed, err := store.ClaimCompetitorSync(ctx, competitor.ID, true, dueNow)
if err != nil || !forceClaimed {
t.Fatalf("forced claim must run before next_sync_at: err=%v claimed=%v", err, forceClaimed)
}
forceNext := dueNow.Add(30 * time.Minute)
if err := store.MarkCompetitorSync(ctx, competitor.ID, forceToken, "idle", "", "", &forceNext); err != nil {
t.Fatal(err)
}
// Postgres timestamp 微秒精度会截断纳秒,比较用秒级容差。
if refreshed, err := store.GetCompetitor(ctx, competitor.ID); err != nil || refreshed.NextSyncAt == nil || refreshed.NextSyncAt.Sub(forceNext).Abs() > time.Second {
t.Fatalf("forced sync must refresh next_sync_at: next=%v err=%v", refreshed.NextSyncAt, err)
}
// 单来源内联模型:同 work_key 的二次 upsert 只合并内容字段,来源保持首个写入方。
work, inserted, err = store.UpsertWork(ctx, WorkInput{Platform: PlatformDouyin, WorkKey: workKey, SourceType: SourceCompetitor, SourceID: competitor.ID, Title: "", Body: "", PublishedAt: nil, Likes: nil, CommentsCount: nil, Shares: nil}, now)
if err != nil || inserted || len(work.Sources) != 1 || work.SourceType != SourceOwned || work.SourceID != bigID {
t.Fatalf("duplicate work upsert must keep the first inline source: work=%+v inserted=%v err=%v", work, inserted, err)
}
competitorWork, inserted, err := store.UpsertWork(ctx, WorkInput{Platform: PlatformDouyin, WorkKey: workKey + "-competitor", SourceType: SourceCompetitor, SourceID: competitor.ID, Title: "竞品作品", Body: "body", PublishedAt: &published, PublishedAtStatus: "verified"}, now)
if err != nil || !inserted {
t.Fatalf("insert competitor work: work=%+v inserted=%v err=%v", competitorWork, inserted, err)
}
filtered, err := store.ListWorks(ctx, WorkFilter{Platform: PlatformDouyin, SourceType: SourceCompetitor})
if err != nil || len(filtered) != 1 || filtered[0].ID != competitorWork.ID {
t.Fatalf("competitor source filter failed: works=%+v err=%v", filtered, err)
}
// 画像聚合视图:作品数与最近发布时间来自 creator_work 聚合。
views, err := store.ListCompetitorsWithProfile(ctx, PlatformDouyin)
if err != nil || len(views) != 1 {
t.Fatalf("list competitor views: views=%+v err=%v", views, err)
}
if views[0].WorkCount != 1 || views[0].LatestPublishedAt == nil {
t.Fatalf("competitor view aggregates: work_count=%d latest=%v", views[0].WorkCount, views[0].LatestPublishedAt)
}
if _, err := store.UpdateCompetitorProfile(ctx, competitor.ID, "https://example.invalid/avatar.jpg", 1200, 300, 88); err != nil {
t.Fatal(err)
}
views, err = store.ListCompetitorsWithProfile(ctx, "")
if err != nil || len(views) != 1 || views[0].FollowerCount == nil || *views[0].FollowerCount != 1200 ||
views[0].FollowingCount == nil || *views[0].FollowingCount != 300 || views[0].AwemeCount == nil || *views[0].AwemeCount != 88 ||
views[0].AvatarURL != "https://example.invalid/avatar.jpg" {
t.Fatalf("competitor profile fields: %+v err=%v", views, err)
}
comment, inserted, err := store.SaveComment(ctx, CommentInput{Platform: PlatformDouyin, CommentKey: "creator-it-comment-" + stamp, WorkID: work.ID, Content: "咨询价格", CommentType: "top_level"})
if err != nil || !inserted || comment.AuthorUID != "" {
t.Fatalf("save comment with missing UID: comment=%+v inserted=%v err=%v", comment, inserted, err)
}
if _, duplicate, err := store.SaveComment(ctx, CommentInput{Platform: PlatformDouyin, CommentKey: comment.CommentKey, WorkID: work.ID, Content: "咨询价格", CommentType: "top_level"}); err != nil || duplicate {
t.Fatalf("comment deduplication failed: duplicate=%v err=%v", duplicate, err)
}
collects := int64(4)
playSnap := int64(5566)
if _, err := store.RecordMetric(ctx, MetricInput{WorkID: work.ID, CollectedAt: now, Likes: &likes, CommentsCount: &comments, Shares: &shares, CollectCount: &collects, PlayCount: &playSnap}, Settings{LookbackDays: 30, NewWorkIntervalSeconds: 1800, MetricInitialIntervalSeconds: 3600, MetricMultiplier: 2, MetricMaxIntervalSeconds: 86400, MetricAgeSeconds: 30 * 24 * 60 * 60}, now); err != nil {
t.Fatal(err)
}
metrics, err := store.ListMetrics(ctx, work.ID)
if err != nil || len(metrics) != 1 {
t.Fatalf("initial metric snapshot failed: metrics=%+v err=%v", metrics, err)
}
if metrics[0].CollectCount == nil || *metrics[0].CollectCount != 4 {
t.Fatalf("collect_count not persisted: %+v", metrics[0])
}
if metrics[0].PlayCount == nil || *metrics[0].PlayCount != 5566 {
t.Fatalf("play_count not persisted in metric point: %+v", metrics[0])
}
// 增速计算:独立作品(避免干扰上方计划推进断言),两个采样点。
growthPublished := now.Add(-2 * time.Hour)
growthLikes0, growthLikes1 := int64(100), int64(400)
growthWork, _, err := store.UpsertWork(ctx, WorkInput{Platform: PlatformDouyin, WorkKey: "creator-it-growth-work-" + stamp, SourceType: SourceOwned, SourceID: smallID, Title: "Growth", Body: "body", PublishedAt: &growthPublished, PublishedAtStatus: "verified", Likes: &growthLikes0}, growthPublished.Add(time.Hour))
if err != nil {
t.Fatal(err)
}
growthSettings := Settings{LookbackDays: 30, NewWorkIntervalSeconds: 1800, MetricInitialIntervalSeconds: 3600, MetricMultiplier: 2, MetricMaxIntervalSeconds: 86400, MetricAgeSeconds: 30 * 24 * 60 * 60}
growthPoint0 := now.Add(-90 * time.Minute) // 首点(晚于发布时间)
if _, err := store.recordMetricWithPlan(ctx, MetricInput{WorkID: growthWork.ID, CollectedAt: growthPoint0, Likes: &growthLikes0}, growthSettings, growthPoint0); err != nil {
t.Fatal(err)
}
growthPoint1 := growthPoint0.Add(2 * time.Hour) // 第二个计划点:next = 首点+2h
if _, err := store.recordMetricWithPlan(ctx, MetricInput{WorkID: growthWork.ID, CollectedAt: growthPoint1, Likes: &growthLikes1}, growthSettings, growthPoint1); err != nil {
t.Fatal(err)
}
growth, err := store.LoadWorksGrowth(ctx, []string{growthWork.ID, work.ID}, WorkGrowthFilter{Hours: 24})
if err != nil {
t.Fatal(err)
}
stat, ok := growth[growthWork.ID]
if !ok || stat.Likes == nil || *stat.Likes != 300 || stat.Coverage != "partial" {
t.Fatalf("growth not computed: stat=%+v ok=%v", stat, ok)
}
if stat, ok := growth[work.ID]; !ok || stat.Likes == nil || *stat.Likes != 0 || stat.Coverage != "partial" {
t.Fatalf("single-point work growth: %+v ok=%v", stat, ok)
}
if _, err := store.RecordMetric(ctx, MetricInput{WorkID: work.ID, CollectedAt: now.Add(-time.Minute), Likes: &likes, CommentsCount: &comments, Shares: &shares}, Settings{LookbackDays: 30, NewWorkIntervalSeconds: 1800, MetricInitialIntervalSeconds: 3600, MetricMultiplier: 2, MetricMaxIntervalSeconds: 86400, MetricAgeSeconds: 30 * 24 * 60 * 60}, now); !errors.Is(err, ErrConflict) {
t.Fatalf("expected early metric point to be rejected, got %v", err)
}
if _, err := store.RecordMetric(ctx, MetricInput{WorkID: work.ID, CollectedAt: now.Add(time.Minute), Likes: &likes, CommentsCount: &comments, Shares: &shares}, Settings{LookbackDays: 30, NewWorkIntervalSeconds: 1800, MetricInitialIntervalSeconds: 3600, MetricMultiplier: 2, MetricMaxIntervalSeconds: 86400, MetricAgeSeconds: 30 * 24 * 60 * 60}, now.Add(-time.Minute)); !errors.Is(err, ErrInvalid) {
t.Fatalf("expected future metric point to be rejected, got %v", err)
}
rule, err := store.CreateRule(ctx, LeadRuleInput{Name: "咨询线索", Enabled: true, SourceType: SourceOwned, Topic: "咨询", IncludeKeywords: []string{"咨询"}, ExcludeKeywords: []string{"招聘"}, AIRequirement: "判断购买意向"})
if err != nil {
t.Fatal(err)
}
rules, err := store.ListRules(ctx, false)
if err != nil || len(rules) != 1 {
t.Fatalf("list rules: rules=%+v err=%v", rules, err)
}
updatedRule, err := store.UpdateRule(ctx, rule.ID, LeadRuleInput{Name: rule.Name, Enabled: true, SourceType: SourceOwned, Topic: rule.Topic, IncludeKeywords: []string{"咨询"}, ExcludeKeywords: []string{"招聘"}, AIRequirement: rule.AIRequirement})
if err != nil || updatedRule.ID != rule.ID {
t.Fatalf("update rule: rule=%+v err=%v", updatedRule, err)
}
result, err := store.AnalyzeComment(ctx, comment.ID, rule.ID, integrationAnalyzer{})
if err != nil || result.Status != "lead" {
t.Fatalf("analyze comment: result=%+v err=%v", result, err)
}
ruleResults, err := store.ListRuleResults(ctx, comment.ID, rule.ID)
if err != nil || len(ruleResults) != 1 {
t.Fatalf("list rule results: results=%+v err=%v", ruleResults, err)
}
leads, err := store.ListLeads(ctx, PlatformDouyin)
if err != nil || len(leads) != 1 {
t.Fatalf("list leads: leads=%+v err=%v", leads, err)
}
if _, err := store.SetRuleEnabled(ctx, rule.ID, false); err != nil {
t.Fatal(err)
}
collector := integrationCollector{work: WorkInput{Platform: PlatformDouyin, WorkKey: "creator-it-collected-" + stamp, SourceType: SourceOwned, SourceID: bigID, Title: "Collected", Body: "body", PublishedAt: &published, PublishedAtStatus: "verified", Likes: &likes}, comment: CommentInput{Platform: PlatformDouyin, CommentKey: "creator-it-collected-comment-" + stamp, WorkID: "", Content: "hello", CommentType: "top_level"}}
collector.comment.WorkID = ""
report, err := store.CollectSource(ctx, PlatformDouyin, SourceOwned, bigID, collector, now)
if err != nil || !report.PaginationComplete || report.WorksSeen != 1 || report.CommentsSaved != 2 {
t.Fatalf("collect source workflow failed: report=%+v err=%v", report, err)
}
}
func TestCreatorPostgresAccountMetricSnapshot(t *testing.T) {
store, phaseAStore, ctx := openCreatorIntegrationStore(t)
stamp := fmt.Sprintf("%d", time.Now().UnixNano())
accountID := createIntegrationAccount(t, ctx, phaseAStore, "metric"+stamp)
now := time.Now().UTC().Truncate(time.Second)
follower, following, favorited, aweme := int64(1200), int64(56), int64(56789), int64(34)
if err := store.RecordAccountMetric(ctx, AccountMetricInput{AccountID: accountID, CollectedAt: now, FollowerCount: &follower, FollowingCount: &following, TotalFavorited: &favorited, AwemeCount: &aweme}); err != nil {
t.Fatal(err)
}
points, err := store.ListAccountMetrics(ctx, accountID)
if err != nil || len(points) != 1 {
t.Fatalf("list account metrics: points=%+v err=%v", points, err)
}
if points[0].FollowerCount == nil || *points[0].FollowerCount != 1200 || points[0].TotalFavorited == nil || *points[0].TotalFavorited != 56789 {
t.Fatalf("account metric snapshot: %+v", points[0])
}
// UNIQUE(account_id, collected_at) 覆盖语义:同一时刻重采更新而非新增。
newFollower := int64(1300)
if err := store.RecordAccountMetric(ctx, AccountMetricInput{AccountID: accountID, CollectedAt: now, FollowerCount: &newFollower}); err != nil {
t.Fatal(err)
}
points, err = store.ListAccountMetrics(ctx, accountID)
if err != nil || len(points) != 1 || points[0].FollowerCount == nil || *points[0].FollowerCount != 1300 {
t.Fatalf("account metric upsert must overwrite same timestamp: points=%+v err=%v", points, err)
}
// 不同时刻追加为独立点。
if err := store.RecordAccountMetric(ctx, AccountMetricInput{AccountID: accountID, CollectedAt: now.Add(30 * time.Minute), FollowerCount: &newFollower}); err != nil {
t.Fatal(err)
}
points, err = store.ListAccountMetrics(ctx, accountID)
if err != nil || len(points) != 2 || !points[0].CollectedAt.Before(points[1].CollectedAt) {
t.Fatalf("account metric timeline must append: points=%+v err=%v", points, err)
}
// 负值拒绝。
bad := int64(-1)
if err := store.RecordAccountMetric(ctx, AccountMetricInput{AccountID: accountID, CollectedAt: now, FollowerCount: &bad}); !errors.Is(err, ErrInvalid) {
t.Fatalf("negative follower_count must be rejected: err=%v", err)
}
// 账号不存在时 INSERT..SELECT 无行写入(旧 FK 拒绝语义由幂等无行替代)。
if err := store.RecordAccountMetric(ctx, AccountMetricInput{AccountID: "account-missing-" + fmt.Sprint(stamp), CollectedAt: now}); err == nil {
t.Fatal("missing account must not record metrics")
}
}
func TestCreatorPostgresAccountMonitorViewAndCollectionStatus(t *testing.T) {
store, phaseAStore, ctx := openCreatorIntegrationStore(t)
stamp := fmt.Sprintf("%d", time.Now().UnixNano())
accountID := createIntegrationAccount(t, ctx, phaseAStore, "metric"+stamp)
now := time.Now().UTC().Truncate(time.Second)
// 无 checkpoint:采集状态为 pending。
status, err := store.GetAccountCollectionStatus(ctx, accountID)
if err != nil || status.Works.Status != "pending" || status.Comments.Status != "pending" {
t.Fatalf("pending collection status: %+v err=%v", status, err)
}
// 无作品:监控视图作品数为 0。
views, err := store.ListAccountMonitorViews(ctx)
if err != nil {
t.Fatal(err)
}
viewFor := func(id string) *AccountMonitorView {
for index := range views {
if views[index].ID == id {
return &views[index]
}
}
return nil
}
view := viewFor(accountID)
if view == nil || view.WorkCount != 0 || view.LatestPublishedAt != nil {
t.Fatalf("monitor view without works: %+v", view)
}
// 写入 checkpoint(succeeded + 窗口)与 owned 作品后:状态带下次窗口,视图聚合作品数。
lease, err := store.beginCheckpoint(ctx, SourceOwned, accountID, "works", now.Add(-24*time.Hour), now)
if err != nil {
t.Fatal(err)
}
if err := store.finishCheckpoint(ctx, SourceOwned, accountID, "works", lease, "succeeded", ""); err != nil {
t.Fatal(err)
}
published := now.Add(-time.Hour)
if _, _, err := store.UpsertWork(ctx, WorkInput{Platform: PlatformDouyin, WorkKey: "monitor-view-" + stamp, SourceType: SourceOwned, SourceID: accountID, Title: "MonitorView", PublishedAt: &published, PublishedAtStatus: "verified"}, now); err != nil {
t.Fatal(err)
}
status, err = store.GetAccountCollectionStatus(ctx, accountID)
if err != nil || status.Works.Status != "succeeded" || status.Works.LastCompletedAt == nil ||
status.Works.NextWindowStart == nil || status.Works.NextWindowEnd == nil {
t.Fatalf("succeeded collection status must carry next window: %+v err=%v", status, err)
}
if !status.Works.NextWindowStart.Before(*status.Works.NextWindowEnd) {
t.Fatalf("next window start must precede end: %+v", status.Works)
}
view = viewFor(accountID)
if view == nil || view.WorkCount != 1 || view.LatestPublishedAt == nil {
// 插入作品后重拉视图:聚合来自 creator_work_source 实时 JOIN。
views, err = store.ListAccountMonitorViews(ctx)
if err != nil {
t.Fatal(err)
}
view = viewFor(accountID)
if view == nil || view.WorkCount != 1 || view.LatestPublishedAt == nil {
t.Fatalf("monitor view aggregates: %+v", view)
}
}
// failed checkpoint 带错误信息(普通错误落在 failed,不是 blocked)。
lease, err = store.beginCheckpoint(ctx, SourceOwned, accountID, "comments", now.Add(-24*time.Hour), now)
if err != nil {
t.Fatal(err)
}
if err := store.failCheckpoint(ctx, SourceOwned, accountID, "comments", lease, errors.New("boom")); err == nil {
t.Fatal("failCheckpoint must return the primary error")
}
status, err = store.GetAccountCollectionStatus(ctx, accountID)
if err != nil || status.Comments.Status != "failed" || status.Comments.LastError == "" {
t.Fatalf("failed comments checkpoint surfaces error: %+v err=%v", status.Comments, err)
}
// 空账号 ID 拒绝。
if _, err := store.GetAccountCollectionStatus(ctx, " "); !errors.Is(err, ErrInvalid) {
t.Fatalf("blank account id must be rejected: err=%v", err)
}
}
func TestCreatorPostgresCollectionAllowsMutedReadAccount(t *testing.T) {
store, phaseAStore, ctx := openCreatorIntegrationStore(t)
stamp := time.Now().UnixNano()
accountID := createIntegrationAccount(t, ctx, phaseAStore, "muted"+fmt.Sprint(stamp))
if err := store.EnsureAccountProfile(ctx, accountID); err != nil {
t.Fatal(err)
}
if _, err := store.RecordVerifiedLoginResult(ctx, accountID, "sec_uid_"+accountID); err != nil {
t.Fatal(err)
}
now := time.Now().UTC()
if _, err := store.db.ExecContext(ctx, `UPDATE creator_account_profile p SET business_status='muted' FROM social_account account WHERE account.id = p.account_id AND account.account_id=$1`, accountID); err != nil {
t.Fatal(err)
}
due, err := store.ListDueOwnedAccounts(ctx, now, 1800)
if err != nil || len(due) != 1 || due[0] != accountID {
t.Fatalf("muted account was not eligible for read collection: accounts=%v err=%v", due, err)
}
if _, err := store.db.ExecContext(ctx, `UPDATE creator_account_profile p SET business_status='banned' FROM social_account account WHERE account.id = p.account_id AND account.account_id=$1`, accountID); err != nil {
t.Fatal(err)
}
due, err = store.ListDueOwnedAccounts(ctx, now, 1800)
if err != nil || len(due) != 0 {
t.Fatalf("banned account was eligible for read collection: accounts=%v err=%v", due, err)
}
}
func TestCreatorPostgresSettingsResetCollectionCheckpoint(t *testing.T) {
store, phaseAStore, ctx := openCreatorIntegrationStore(t)
stamp := time.Now().UnixNano()
accountID := createIntegrationAccount(t, ctx, phaseAStore, "settings"+fmt.Sprint(stamp))
before := time.Now().UTC()
lease, err := store.beginCheckpoint(ctx, SourceOwned, accountID, "works", before.Add(-48*time.Hour), before.Add(-47*time.Hour))
if err != nil {
t.Fatalf("begin checkpoint: %v", err)
}
if err := store.finishCheckpoint(ctx, SourceOwned, accountID, "works", lease, "succeeded", ""); err != nil {
t.Fatalf("finish checkpoint: %v", err)
}
if _, err := store.db.ExecContext(ctx, `UPDATE creator_collection_checkpoint SET window_start='2020-01-01T00:00:00Z',window_end='2020-01-02T00:00:00Z',cursor='old-cursor' WHERE source_type=$1 AND source_id=$2 AND collection_kind='works'`, SourceOwned, accountID); err != nil {
t.Fatal(err)
}
settings, err := store.GetSettings(ctx)
if err != nil {
t.Fatal(err)
}
settings.LookbackDays = 3
if _, err := store.UpdateSettings(ctx, SettingsUpdate{LookbackDays: settings.LookbackDays, NewWorkIntervalSeconds: settings.NewWorkIntervalSeconds, MetricInitialIntervalSeconds: settings.MetricInitialIntervalSeconds, MetricMultiplier: settings.MetricMultiplier, MetricMaxIntervalSeconds: settings.MetricMaxIntervalSeconds, MetricAgeSeconds: settings.MetricAgeSeconds, AIProvider: settings.AIProvider, AIModel: settings.AIModel, AIConfigured: settings.AIConfigured}); err != nil {
t.Fatalf("update settings: %v", err)
}
checkpoint, err := store.checkpoint(ctx, SourceOwned, accountID, "works")
if err != nil {
t.Fatal(err)
}
if checkpoint.Cursor != "" || checkpoint.Status != "idle" || checkpoint.WindowStart.Before(time.Now().UTC().Add(-4*24*time.Hour)) {
t.Fatalf("settings did not reset checkpoint: %+v", checkpoint)
}
lease, err = store.beginCheckpoint(ctx, SourceOwned, accountID, "works", time.Now().UTC().Add(-time.Hour), time.Now().UTC())
if err != nil {
t.Fatalf("begin running checkpoint: %v", err)
}
settings.LookbackDays = 4
if _, err := store.UpdateSettings(ctx, SettingsUpdate{LookbackDays: settings.LookbackDays, NewWorkIntervalSeconds: settings.NewWorkIntervalSeconds, MetricInitialIntervalSeconds: settings.MetricInitialIntervalSeconds, MetricMultiplier: settings.MetricMultiplier, MetricMaxIntervalSeconds: settings.MetricMaxIntervalSeconds, MetricAgeSeconds: settings.MetricAgeSeconds, AIProvider: settings.AIProvider, AIModel: settings.AIModel, AIConfigured: settings.AIConfigured}); !errors.Is(err, ErrConflict) {
t.Fatalf("settings changed during running checkpoint: err=%v", err)
}
if err := store.finishCheckpoint(ctx, SourceOwned, accountID, "works", lease, "succeeded", ""); err != nil {
t.Fatal(err)
}
}
func TestCreatorPostgresMetricPlanFollowsPublishedAt(t *testing.T) {
store, phaseAStore, ctx := openCreatorIntegrationStore(t)
stamp := fmt.Sprintf("%d", time.Now().UnixNano())
bigID := createIntegrationAccount(t, ctx, phaseAStore, "metric"+stamp)
if err := store.EnsureAccountProfile(ctx, bigID); err != nil {
t.Fatal(err)
}
published := time.Now().UTC().Add(-time.Hour).Truncate(time.Microsecond)
likes, comments, shares := int64(1), int64(1), int64(1)
work, _, err := store.UpsertWork(ctx, WorkInput{Platform: PlatformDouyin, WorkKey: "creator-it-metric-work-" + stamp, SourceType: SourceOwned, SourceID: bigID, Title: "Metric", Body: "body", PublishedAt: &published, PublishedAtStatus: "verified", Likes: &likes, CommentsCount: &comments, Shares: &shares}, published.Add(time.Hour))
if err != nil {
t.Fatal(err)
}
settings, err := store.GetSettings(ctx)
if err != nil {
t.Fatal(err)
}
settings.MetricInitialIntervalSeconds = 3600
settings.MetricMultiplier = 2
settings.MetricMaxIntervalSeconds = 55 * 3600
settings.MetricAgeSeconds = 100 * 3600
point := func(at time.Time) MetricInput {
return MetricInput{WorkID: work.ID, CollectedAt: at, Likes: &likes, CommentsCount: &comments, Shares: &shares}
}
if _, err := store.recordMetricWithPlan(ctx, point(published.Add(2*time.Hour)), settings, published.Add(2*time.Hour)); err != nil {
t.Fatal(err)
}
var next time.Time
var interval int64
if err := store.db.QueryRowContext(ctx, `SELECT metric_plan_next_at,metric_plan_interval_seconds FROM creator_work WHERE work_id=$1`, work.ID).Scan(&next, &interval); err != nil {
t.Fatal(err)
}
if !next.Equal(published.Add(3*time.Hour)) || interval != 4*3600 {
t.Fatalf("first metric plan drifted: next=%s interval=%d", next, interval)
}
if _, err := store.recordMetricWithPlan(ctx, point(published.Add(3*time.Hour)), settings, published.Add(3*time.Hour)); err != nil {
t.Fatal(err)
}
if err := store.db.QueryRowContext(ctx, `SELECT metric_plan_next_at,metric_plan_interval_seconds FROM creator_work WHERE work_id=$1`, work.ID).Scan(&next, &interval); err != nil {
t.Fatal(err)
}
if !next.Equal(published.Add(7*time.Hour)) || interval != 8*3600 {
t.Fatalf("second metric plan drifted: next=%s interval=%d", next, interval)
}
}