Files
2026-09-21 16:44:09 +08:00

635 lines
32 KiB
Go

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