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) } } 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) } if err := store.SetRelation(ctx, bigID, smallID, true); err != nil { t.Fatal(err) } relations, err := store.ListRelations(ctx, bigID) if err != nil || len(relations) != 1 { t.Fatalf("list relations: relations=%+v err=%v", relations, err) } if _, err := store.SetBigAccount(ctx, smallID, true); !errors.Is(err, ErrConflict) { t.Fatalf("expected small-account promotion to be rejected, got %v", err) } authorizedID := createIntegrationAccount(t, ctx, phaseAStore, "authorized"+stamp) if _, err := store.db.ExecContext(ctx, `UPDATE social_account SET authorization_kind = 'authorized' WHERE id = $1`, authorizedID); err != nil { t.Fatal(err) } if err := store.SetRelation(ctx, bigID, authorizedID, true); !errors.Is(err, ErrInvalid) { t.Fatalf("expected authorized account relation to be rejected, got %v", 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, 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) != 2 { t.Fatalf("work source association was not retained: work=%+v inserted=%v err=%v", work, inserted, err) } filtered, err := store.ListWorks(ctx, WorkFilter{Platform: PlatformDouyin, SourceType: SourceCompetitor}) if err != nil || len(filtered) != 1 || filtered[0].ID != work.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) } material, selected, err := store.SelectMaterial(ctx, work.ID) if err != nil || !selected || !material.Selected { t.Fatalf("select material: material=%+v selected=%v err=%v", material, selected, err) } for step, status := range map[string]string{"download": "succeeded", "audio": "no_audio", "transcription": "no_speech"} { if _, err := store.SetMaterialStep(ctx, work.ID, step, status, "ref-"+step, ""); err != nil { t.Fatal(err) } } if _, err := store.ConfirmRewrite(ctx, work.ID, "仿写同主题但不复制原文"); err != nil { t.Fatal(err) } if _, err := store.SaveRewrite(ctx, work.ID, "新标题", "新口播"); 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) } // 账号不存在时 FK 拒绝。 if err := store.RecordAccountMetric(ctx, AccountMetricInput{AccountID: "account-missing-" + fmt.Sprint(stamp), CollectedAt: now}); err == nil { t.Fatal("missing account must violate FK") } } 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 SET business_status='muted' WHERE 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 SET business_status='banned' WHERE 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, TranscriptionProvider: settings.TranscriptionProvider, TranscriptionModel: settings.TranscriptionModel, TranscriptionConfigured: settings.TranscriptionConfigured}); 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, TranscriptionProvider: settings.TranscriptionProvider, TranscriptionModel: settings.TranscriptionModel, TranscriptionConfigured: settings.TranscriptionConfigured}); !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 next_plan_at,interval_seconds FROM creator_metric_plan 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 next_plan_at,interval_seconds FROM creator_metric_plan 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) } } func prepareIntegrationActionFixture(t *testing.T, store *Store, phaseAStore *account.Store, ctx context.Context, stamp string) (string, string, Work, Comment, Strategy) { t.Helper() bigID := createIntegrationAccount(t, ctx, phaseAStore, "big"+stamp) smallID := createIntegrationAccount(t, ctx, phaseAStore, "small"+stamp) for _, accountID := range []string{bigID, smallID} { if err := store.EnsureAccountProfile(ctx, accountID); err != nil { t.Fatal(err) } if _, err := store.UpdateAccountProfile(ctx, accountID, AccountProfileUpdate{RealNameStatus: "unknown", BusinessStatus: "normal", CooldownSeconds: 86400}); err != nil { t.Fatal(err) } if _, err := store.RecordVerifiedLoginResult(ctx, accountID, "sec_uid_"+accountID); err != nil { t.Fatal(err) } } if _, err := store.SetBigAccount(ctx, bigID, true); err != nil { t.Fatal(err) } if err := store.SetRelation(ctx, bigID, smallID, true); err != nil { t.Fatal(err) } now := time.Now().UTC().Truncate(time.Microsecond) published := now.Add(-time.Hour) likes, comments, shares := int64(1), int64(1), int64(1) work, _, err := store.UpsertWork(ctx, WorkInput{Platform: PlatformDouyin, WorkKey: "creator-it-action-work-" + stamp, SourceType: SourceOwned, SourceID: bigID, Title: "Action", Body: "body", PublishedAt: &published, PublishedAtStatus: "verified", Likes: &likes, CommentsCount: &comments, Shares: &shares}, now) if err != nil { t.Fatal(err) } comment, _, err := store.SaveComment(ctx, CommentInput{Platform: PlatformDouyin, CommentKey: "creator-it-action-comment-" + stamp, WorkID: work.ID, AuthorUID: "interactor-" + stamp, Content: "hello", CommentType: "top_level"}) if err != nil { t.Fatal(err) } strategy, err := store.CreateStrategy(ctx, bigID, StrategyInput{ExecutionAccountID: smallID, Position: 1, Enabled: true, EventTypes: []string{"comment"}, Action: ActionReplyComment, TargetType: "comment", CandidateTexts: []string{"已收到"}}) if err != nil { t.Fatal(err) } strategies, err := store.ListStrategies(ctx, bigID) if err != nil || len(strategies) != 1 { t.Fatalf("list strategies: strategies=%+v err=%v", strategies, err) } return bigID, smallID, work, comment, strategy } func TestCreatorPostgresActionsAndMessaging(t *testing.T) { store, phaseAStore, ctx := openCreatorIntegrationStore(t) stamp := fmt.Sprintf("%d", time.Now().UnixNano()) bigID, smallID, work, comment, strategy := prepareIntegrationActionFixture(t, store, phaseAStore, ctx, stamp) event := InteractionEvent{Platform: PlatformDouyin, ReceivingAccountID: bigID, EventKey: "creator-it-event-" + stamp, EventType: "comment", InteractorUID: comment.AuthorUID, CommentID: comment.ID, WorkID: work.ID} automatic, err := store.ProcessAutomaticEvent(ctx, event, integrationExecutor{}, nil) if err != nil || automatic.Operation == nil || automatic.Operation.State != "succeeded" || automatic.Event.State != "succeeded" { t.Fatalf("automatic action failed: result=%+v err=%v", automatic, err) } events, err := store.ListEvents(ctx, bigID) if err != nil || len(events) != 1 { t.Fatalf("list events: events=%+v err=%v", events, err) } duplicate, err := store.ProcessAutomaticEvent(ctx, event, integrationExecutor{}, nil) if err != nil || !duplicate.Duplicate || duplicate.Event.ID != automatic.Event.ID { t.Fatalf("automatic event deduplication failed: result=%+v err=%v", duplicate, err) } cooldownEvent := event cooldownEvent.EventKey += "-cooldown" blocked, err := store.ProcessAutomaticEvent(ctx, cooldownEvent, integrationExecutor{}, nil) if err != nil || blocked.Event.State != "blocked" { t.Fatalf("cooldown did not block second event: result=%+v err=%v", blocked, err) } manualInput := OperationInput{IdempotencyKey: "creator-it-manual-" + stamp, Source: "manual", Action: ActionReplyComment, Platform: PlatformDouyin, AccountID: smallID, TargetUID: comment.AuthorUID, TargetCommentID: comment.ID, Text: "人工回复"} manual, inserted, err := store.CreateOperation(ctx, manualInput) if err != nil || !inserted { t.Fatalf("create manual operation: operation=%+v inserted=%v err=%v", manual, inserted, err) } manualDuplicate, inserted, err := store.CreateOperation(ctx, manualInput) if err != nil || inserted || manualDuplicate.ID != manual.ID { t.Fatalf("manual operation idempotency failed: operation=%+v inserted=%v err=%v", manualDuplicate, inserted, err) } operations, err := store.ListOperations(ctx, smallID) if err != nil || len(operations) != 2 { t.Fatalf("list operations: operations=%+v err=%v", operations, err) } manual, err = store.ExecuteManualOperation(ctx, manual.ID, integrationExecutor{}) if err != nil || manual.State != "succeeded" { t.Fatalf("execute manual operation: operation=%+v err=%v", manual, err) } dmOperation, inserted, err := store.CreateOperation(ctx, OperationInput{IdempotencyKey: "creator-it-dm-" + stamp, Source: "manual", Action: ActionDM, Platform: PlatformDouyin, AccountID: smallID, TargetUID: "peer-" + stamp, Text: "人工私信"}) if err != nil || !inserted { t.Fatalf("create direct-message operation: operation=%+v inserted=%v err=%v", dmOperation, inserted, err) } dmOperation, err = store.ExecuteManualOperation(ctx, dmOperation.ID, integrationExecutor{}) if err != nil || dmOperation.State != "succeeded" { t.Fatalf("execute direct-message operation: operation=%+v err=%v", dmOperation, err) } messageAt := time.Now().UTC().Truncate(time.Microsecond) message, inserted, err := store.SaveMessage(ctx, MessageInput{Platform: PlatformDouyin, AccountID: smallID, PeerUID: "peer-" + stamp, PeerName: "Peer", PlatformMessageKey: "creator-it-message-" + stamp, Direction: "inbound", MessageType: "text", Text: "hello", MessageAt: &messageAt}) if err != nil || !inserted { t.Fatalf("save message: message=%+v inserted=%v err=%v", message, inserted, err) } if _, inserted, err := store.SaveMessage(ctx, MessageInput{Platform: PlatformDouyin, AccountID: smallID, PeerUID: "peer-" + stamp, PlatformMessageKey: message.PlatformMessageKey, Direction: "inbound", MessageType: "text", Text: "hello", MessageAt: &messageAt}); err != nil || inserted { t.Fatalf("message deduplication failed: inserted=%v err=%v", inserted, err) } conversations, err := store.ListConversations(ctx, smallID) if err != nil || len(conversations) != 1 { t.Fatalf("list conversations: conversations=%+v err=%v", conversations, err) } messages, err := store.ListMessages(ctx, conversations[0].ID) if err != nil || len(messages) != 2 { t.Fatalf("list messages: messages=%+v err=%v", messages, err) } if _, err := store.SetEventDisplayed(ctx, automatic.Event.ID, time.Now().UTC()); err != nil { t.Fatal(err) } updated, err := store.UpdateStrategy(ctx, strategy.ID, StrategyInput{ExecutionAccountID: smallID, Position: 2, Enabled: true, EventTypes: []string{"comment", "like"}, Action: ActionReplyComment, TargetType: "comment", CandidateTexts: []string{"已更新"}}) if err != nil || updated.Position != 2 || len(updated.EventTypes) != 2 { t.Fatalf("update strategy: strategy=%+v err=%v", updated, err) } if _, err := store.SetStrategyEnabled(ctx, strategy.ID, false); err != nil { t.Fatal(err) } if err := store.DeleteStrategy(ctx, strategy.ID); err != nil { t.Fatal(err) } if err := store.SetRelation(ctx, bigID, smallID, false); err != nil { t.Fatal(err) } if _, err := store.SetBigAccount(ctx, bigID, false); err != nil { t.Fatal(err) } start, end := time.Now().UTC().Add(-time.Hour), time.Now().UTC() lease, err := store.beginCheckpoint(ctx, SourceOwned, bigID, "works", start, end) if err != nil { t.Fatal(err) } if _, err := store.beginCheckpoint(ctx, SourceOwned, bigID, "works", start, end); !errors.Is(err, ErrConflict) { t.Fatalf("expected concurrent checkpoint claim conflict, got %v", err) } if err := store.finishCheckpoint(ctx, SourceOwned, bigID, "works", "stale-lease", "succeeded", ""); !errors.Is(err, ErrConflict) { t.Fatalf("expected stale checkpoint completion conflict, got %v", err) } if err := store.saveCheckpointCursor(ctx, SourceOwned, bigID, "works", lease, "1"); err != nil { t.Fatal(err) } if err := store.finishCheckpoint(ctx, SourceOwned, bigID, "works", lease, "succeeded", ""); err != nil { t.Fatal(err) } }