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) 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}, now) if err != nil || !inserted { t.Fatalf("insert work: work=%+v inserted=%v err=%v", work, inserted, 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.ListCompetitors(ctx, PlatformDouyin) if err != nil || len(competitors) != 1 || len(competitors[0].Tags) != 1 || competitors[0].Tags[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) } 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) } 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) } if _, err := store.RecordMetric(ctx, MetricInput{WorkID: work.ID, CollectedAt: now, Likes: &likes, CommentsCount: &comments, Shares: &shares}, 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 _, 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 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) } }