Files
rogee 3e809c4fc7 feat(creator): 监控账号画像采集入库——粉丝/关注/获赞/作品数与作者介绍
抖音 profile/other 响应本就携带画像数据但被丢弃,同步链路也无回填。
- migration 1046: creator_competitor 增加 signature、total_favorited 列及非负约束
- douyin 解析层: TargetProfile 增加签名与计数,total_favorited 兼容字符串大数
- store: UpdateCompetitorProfile 改为画像回填(nil 计数保留原值)
- 同步链路: 作品采集成功后拉取 profile 回填,失败仅记日志不阻断
2026-09-30 16:37:56 +08:00

809 lines
41 KiB
Go
Raw Permalink 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 TestCompetitorProfileUpdateRoundTrip(t *testing.T) {
store, _, ctx := openCreatorIntegrationStore(t)
stamp := fmt.Sprintf("%d", time.Now().UnixNano())
competitor, err := store.UpsertCompetitor(ctx, CompetitorInput{
Platform: PlatformDouyin,
PlatformAccountKey: "sec_uid_profile_" + stamp,
Nickname: "Profile Target",
HomepageURL: "https://www.douyin.com/user/sec_uid_profile_" + stamp,
})
if err != nil {
t.Fatal(err)
}
follower, following, favorited, aweme := int64(123), int64(45), int64(98765), int64(7)
updated, err := store.UpdateCompetitorProfile(ctx, competitor.ID, CompetitorProfileUpdate{
AvatarURL: "https://example.invalid/avatar",
Signature: "第一行介绍\n第二行介绍",
FollowerCount: &follower,
FollowingCount: &following,
TotalFavorited: &favorited,
AwemeCount: &aweme,
})
if err != nil {
t.Fatal(err)
}
if updated.Signature != "第一行介绍\n第二行介绍" || updated.FollowerCount == nil || *updated.FollowerCount != 123 ||
updated.FollowingCount == nil || *updated.FollowingCount != 45 || updated.TotalFavorited == nil || *updated.TotalFavorited != 98765 ||
updated.AwemeCount == nil || *updated.AwemeCount != 7 {
t.Fatalf("profile update: %+v", updated)
}
// nil 计数保留原值,不置 NULL。
reread, err := store.UpdateCompetitorProfile(ctx, competitor.ID, CompetitorProfileUpdate{Signature: "新签名"})
if err != nil || reread.Signature != "新签名" || reread.FollowerCount == nil || *reread.FollowerCount != 123 || reread.TotalFavorited == nil || *reread.TotalFavorited != 98765 {
t.Fatalf("profile partial update: %+v err=%v", reread, err)
}
negative := int64(-1)
if _, err := store.UpdateCompetitorProfile(ctx, competitor.ID, CompetitorProfileUpdate{FollowerCount: &negative}); err == nil {
t.Fatal("negative follower count accepted")
}
if _, err := store.UpdateCompetitorProfile(ctx, "missing_"+stamp, CompetitorProfileUpdate{}); err == nil {
t.Fatal("missing competitor accepted")
}
}
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)
}
follower, following, aweme := int64(1200), int64(300), int64(88)
if _, err := store.UpdateCompetitorProfile(ctx, competitor.ID, CompetitorProfileUpdate{AvatarURL: "https://example.invalid/avatar.jpg", FollowerCount: &follower, FollowingCount: &following, AwemeCount: &aweme}); 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)
}
// 无指标快照:画像指标为空(不伪造成 0)。
if view.FollowerCount != nil || view.FollowingCount != nil || view.TotalFavorited != nil || view.AwemeCount != nil {
t.Fatalf("monitor view without metrics must leave profile counters nil: %+v", view)
}
if view.CommentTotal != 0 {
t.Fatalf("monitor view without works must have zero comments: %+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)
}
}
// 画像指标快照:写入两轮后视图取最新一轮;评论聚合来自 owned 作品最近指标之和。
first, second := int64(100), int64(200)
following := int64(30)
favorited := int64(5000)
aweme := int64(7)
mutual := int64(9)
if err := store.RecordAccountMetric(ctx, AccountMetricInput{AccountID: accountID, CollectedAt: now.Add(-2 * time.Hour), FollowerCount: &first, FollowingCount: &following, TotalFavorited: &favorited, AwemeCount: &aweme, FriendCount: &mutual}); err != nil {
t.Fatal(err)
}
if err := store.RecordAccountMetric(ctx, AccountMetricInput{AccountID: accountID, CollectedAt: now.Add(-time.Hour), FollowerCount: &second}); err != nil {
t.Fatal(err)
}
monitorWork, _, err := store.UpsertWork(ctx, WorkInput{Platform: PlatformDouyin, WorkKey: "monitor-view-" + stamp, SourceType: SourceOwned, SourceID: accountID, Title: "MonitorView", PublishedAt: &published, PublishedAtStatus: "verified"}, now)
if err != nil {
t.Fatal(err)
}
if _, err := store.RecordMetric(ctx, MetricInput{WorkID: monitorWork.ID, CollectedAt: now, CommentsCount: &favorited}, Settings{LookbackDays: 30, NewWorkIntervalSeconds: 1800, MetricInitialIntervalSeconds: 3600, MetricMultiplier: 2, MetricMaxIntervalSeconds: 86400, MetricAgeSeconds: 30 * 24 * 60 * 60}, now); err != nil {
t.Fatal(err)
}
views, err = store.ListAccountMonitorViews(ctx)
if err != nil {
t.Fatal(err)
}
view = viewFor(accountID)
if view == nil || view.FollowerCount == nil || *view.FollowerCount != second {
t.Fatalf("monitor view must carry the latest metric snapshot: %+v", view)
}
if view.FollowingCount == nil || *view.FollowingCount != following || view.TotalFavorited == nil || *view.TotalFavorited != favorited || view.AwemeCount == nil || *view.AwemeCount != aweme {
t.Fatalf("monitor view must fall back to the latest non-null counters: %+v", view)
}
if view.FriendCount == nil || *view.FriendCount != mutual {
t.Fatalf("monitor view must carry the latest mutual friend count: %+v", view)
}
if view.CommentTotal != favorited {
t.Fatalf("monitor view must aggregate owned work comments: got %d want %d", view.CommentTotal, favorited)
}
// 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)
}
}