858 lines
44 KiB
Go
858 lines
44 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)
|
||
plays := int64(777)
|
||
workKey := "creator-it-work-" + stamp
|
||
work, inserted, err := store.UpsertWork(ctx, WorkInput{Platform: PlatformDouyin, WorkKey: workKey, SourceType: SourceOwned, SourceID: bigID, Title: "Topic", Body: "Need consultation", PublishedAt: &published, PublishedAtStatus: "verified", OriginalURL: "https://www.douyin.com/video/" + workKey, Likes: &likes, CommentsCount: &comments, Shares: &shares, PlayCount: &plays}, now)
|
||
if err != nil || !inserted {
|
||
t.Fatalf("insert work: work=%+v inserted=%v err=%v", work, inserted, err)
|
||
}
|
||
if work.PlayCount == nil || *work.PlayCount != 777 {
|
||
t.Fatalf("play_count not persisted on insert: %+v", work)
|
||
}
|
||
// COALESCE 语义:二次 upsert 未带 play_count 时保留旧值。
|
||
if _, _, err := store.UpsertWork(ctx, WorkInput{Platform: PlatformDouyin, WorkKey: workKey, SourceType: SourceOwned, SourceID: bigID, Title: "Topic-renamed", PublishedAt: &published, PublishedAtStatus: "verified"}, now); err != nil {
|
||
t.Fatal(err)
|
||
}
|
||
retained, err := store.GetWorkByKey(ctx, PlatformDouyin, workKey)
|
||
if err != nil || retained.PlayCount == nil || *retained.PlayCount != 777 {
|
||
t.Fatalf("play_count must be retained via COALESCE: work=%+v err=%v", retained, err)
|
||
}
|
||
// 负播放量拒绝。
|
||
if _, _, err := store.UpsertWork(ctx, WorkInput{Platform: PlatformDouyin, WorkKey: workKey + "-neg", SourceType: SourceOwned, SourceID: bigID, Title: "neg", PlayCount: &[]int64{-1}[0]}, now); !errors.Is(err, ErrInvalid) {
|
||
t.Fatalf("negative play_count must be rejected: err=%v", err)
|
||
}
|
||
competitor, err := store.CreateCompetitor(ctx, CompetitorInput{Platform: PlatformDouyin, PlatformAccountKey: "sec_uid_competitor_" + stamp, UniqueID: "competitor_" + stamp, Nickname: "Competitor", HomepageURL: "https://www.douyin.com/user/sec_uid_competitor_" + stamp, Tags: []string{"重点监测"}})
|
||
if err != nil {
|
||
t.Fatal(err)
|
||
}
|
||
if competitor.UniqueID != "competitor_"+stamp || len(competitor.Tags) != 1 || competitor.Tags[0] != "重点监测" {
|
||
t.Fatalf("create competitor tags: %+v", competitor.Tags)
|
||
}
|
||
updatedCompetitor, err := store.UpdateCompetitorTags(ctx, competitor.ID, []string{"已分类"})
|
||
if err != nil || len(updatedCompetitor.Tags) != 1 || updatedCompetitor.Tags[0] != "已分类" {
|
||
t.Fatalf("update competitor tags: competitor=%+v err=%v", updatedCompetitor, err)
|
||
}
|
||
competitors, err := store.ListCompetitorsWithProfile(ctx, PlatformDouyin)
|
||
if err != nil || len(competitors) != 1 || len(competitors[0].Tags) != 1 || competitors[0].Tags[0] != "已分类" || competitors[0].WorkCount != 0 {
|
||
t.Fatalf("list competitors: competitors=%+v err=%v", competitors, err)
|
||
}
|
||
if _, err := store.SetCompetitorEnabled(ctx, competitor.ID, false); err != nil {
|
||
t.Fatal(err)
|
||
}
|
||
if _, err := store.SetCompetitorEnabled(ctx, competitor.ID, true); err != nil {
|
||
t.Fatal(err)
|
||
}
|
||
dueNow := time.Now().UTC().Add(time.Minute)
|
||
if due, err := store.ListDueCompetitors(ctx, dueNow); err != nil || len(due) != 1 {
|
||
t.Fatalf("list due competitors: due=%+v err=%v", due, err)
|
||
}
|
||
leaseToken, claimed, err := store.ClaimCompetitorSync(ctx, competitor.ID, false, dueNow)
|
||
if err != nil || !claimed || leaseToken == "" {
|
||
t.Fatalf("claim competitor sync: token=%q claimed=%v err=%v", leaseToken, claimed, err)
|
||
}
|
||
if err := store.MarkCompetitorSync(ctx, competitor.ID, leaseToken, "idle", "", "", nil); err != nil {
|
||
t.Fatal(err)
|
||
}
|
||
// blocked 同步也必须排下一次重试:否则 next_sync_at 置空后定时获取对该账号永久失效。
|
||
blockedToken, blockedClaimed, err := store.ClaimCompetitorSync(ctx, competitor.ID, true, dueNow)
|
||
if err != nil || !blockedClaimed {
|
||
t.Fatalf("claim competitor sync for blocked: err=%v claimed=%v", err, blockedClaimed)
|
||
}
|
||
blockedNext := dueNow.Add(2 * time.Hour)
|
||
if err := store.MarkCompetitorSync(ctx, competitor.ID, blockedToken, "blocked", "", "runtime unavailable", &blockedNext); err != nil {
|
||
t.Fatal(err)
|
||
}
|
||
if due, err := store.ListDueCompetitors(ctx, blockedNext); err != nil || len(due) != 1 || due[0].ID != competitor.ID {
|
||
t.Fatalf("blocked competitor must reschedule: due=%+v err=%v", due, err)
|
||
}
|
||
if due, err := store.ListDueCompetitors(ctx, dueNow); err != nil || len(due) != 0 {
|
||
t.Fatalf("blocked competitor must not fire before retry time: due=%+v err=%v", due, err)
|
||
}
|
||
// 强制同步语义:next_sync_at 未到期时 force claim 仍立即执行;完成后 MarkCompetitorSync 刷新 next_sync_at。
|
||
if _, claimed, err := store.ClaimCompetitorSync(ctx, competitor.ID, false, dueNow); err != nil || claimed {
|
||
t.Fatalf("non-forced claim must wait for next_sync_at: claimed=%v err=%v", claimed, err)
|
||
}
|
||
forceToken, forceClaimed, err := store.ClaimCompetitorSync(ctx, competitor.ID, true, dueNow)
|
||
if err != nil || !forceClaimed {
|
||
t.Fatalf("forced claim must run before next_sync_at: err=%v claimed=%v", err, forceClaimed)
|
||
}
|
||
forceNext := dueNow.Add(30 * time.Minute)
|
||
if err := store.MarkCompetitorSync(ctx, competitor.ID, forceToken, "idle", "", "", &forceNext); err != nil {
|
||
t.Fatal(err)
|
||
}
|
||
// Postgres timestamp 微秒精度会截断纳秒,比较用秒级容差。
|
||
if refreshed, err := store.GetCompetitor(ctx, competitor.ID); err != nil || refreshed.NextSyncAt == nil || refreshed.NextSyncAt.Sub(forceNext).Abs() > time.Second {
|
||
t.Fatalf("forced sync must refresh next_sync_at: next=%v err=%v", refreshed.NextSyncAt, err)
|
||
}
|
||
work, inserted, err = store.UpsertWork(ctx, WorkInput{Platform: PlatformDouyin, WorkKey: workKey, SourceType: SourceCompetitor, SourceID: competitor.ID, Title: "", Body: "", PublishedAt: nil, Likes: nil, CommentsCount: nil, Shares: nil}, now)
|
||
if err != nil || inserted || len(work.Sources) != 2 {
|
||
t.Fatalf("work source association was not retained: work=%+v inserted=%v err=%v", work, inserted, err)
|
||
}
|
||
filtered, err := store.ListWorks(ctx, WorkFilter{Platform: PlatformDouyin, SourceType: SourceCompetitor})
|
||
if err != nil || len(filtered) != 1 || filtered[0].ID != work.ID {
|
||
t.Fatalf("competitor source filter failed: works=%+v err=%v", filtered, err)
|
||
}
|
||
|
||
// 画像聚合视图:作品数与最近发布时间来自 creator_work 聚合。
|
||
views, err := store.ListCompetitorsWithProfile(ctx, PlatformDouyin)
|
||
if err != nil || len(views) != 1 {
|
||
t.Fatalf("list competitor views: views=%+v err=%v", views, err)
|
||
}
|
||
if views[0].WorkCount != 1 || views[0].LatestPublishedAt == nil {
|
||
t.Fatalf("competitor view aggregates: work_count=%d latest=%v", views[0].WorkCount, views[0].LatestPublishedAt)
|
||
}
|
||
if _, err := store.UpdateCompetitorProfile(ctx, competitor.ID, "https://example.invalid/avatar.jpg", 1200, 300, 88); err != nil {
|
||
t.Fatal(err)
|
||
}
|
||
views, err = store.ListCompetitorsWithProfile(ctx, "")
|
||
if err != nil || len(views) != 1 || views[0].FollowerCount == nil || *views[0].FollowerCount != 1200 ||
|
||
views[0].FollowingCount == nil || *views[0].FollowingCount != 300 || views[0].AwemeCount == nil || *views[0].AwemeCount != 88 ||
|
||
views[0].AvatarURL != "https://example.invalid/avatar.jpg" {
|
||
t.Fatalf("competitor profile fields: %+v err=%v", views, err)
|
||
}
|
||
|
||
comment, inserted, err := store.SaveComment(ctx, CommentInput{Platform: PlatformDouyin, CommentKey: "creator-it-comment-" + stamp, WorkID: work.ID, Content: "咨询价格", CommentType: "top_level"})
|
||
if err != nil || !inserted || comment.AuthorUID != "" {
|
||
t.Fatalf("save comment with missing UID: comment=%+v inserted=%v err=%v", comment, inserted, err)
|
||
}
|
||
if _, duplicate, err := store.SaveComment(ctx, CommentInput{Platform: PlatformDouyin, CommentKey: comment.CommentKey, WorkID: work.ID, Content: "咨询价格", CommentType: "top_level"}); err != nil || duplicate {
|
||
t.Fatalf("comment deduplication failed: duplicate=%v err=%v", duplicate, err)
|
||
}
|
||
collects := int64(4)
|
||
playSnap := int64(5566)
|
||
if _, err := store.RecordMetric(ctx, MetricInput{WorkID: work.ID, CollectedAt: now, Likes: &likes, CommentsCount: &comments, Shares: &shares, CollectCount: &collects, PlayCount: &playSnap}, Settings{LookbackDays: 30, NewWorkIntervalSeconds: 1800, MetricInitialIntervalSeconds: 3600, MetricMultiplier: 2, MetricMaxIntervalSeconds: 86400, MetricAgeSeconds: 30 * 24 * 60 * 60}, now); err != nil {
|
||
t.Fatal(err)
|
||
}
|
||
metrics, err := store.ListMetrics(ctx, work.ID)
|
||
if err != nil || len(metrics) != 1 {
|
||
t.Fatalf("initial metric snapshot failed: metrics=%+v err=%v", metrics, err)
|
||
}
|
||
if metrics[0].CollectCount == nil || *metrics[0].CollectCount != 4 {
|
||
t.Fatalf("collect_count not persisted: %+v", metrics[0])
|
||
}
|
||
if metrics[0].PlayCount == nil || *metrics[0].PlayCount != 5566 {
|
||
t.Fatalf("play_count not persisted in metric point: %+v", metrics[0])
|
||
}
|
||
// 增速计算:独立作品(避免干扰上方计划推进断言),两个采样点。
|
||
growthPublished := now.Add(-2 * time.Hour)
|
||
growthLikes0, growthLikes1 := int64(100), int64(400)
|
||
growthWork, _, err := store.UpsertWork(ctx, WorkInput{Platform: PlatformDouyin, WorkKey: "creator-it-growth-work-" + stamp, SourceType: SourceOwned, SourceID: smallID, Title: "Growth", Body: "body", PublishedAt: &growthPublished, PublishedAtStatus: "verified", Likes: &growthLikes0}, growthPublished.Add(time.Hour))
|
||
if err != nil {
|
||
t.Fatal(err)
|
||
}
|
||
growthSettings := Settings{LookbackDays: 30, NewWorkIntervalSeconds: 1800, MetricInitialIntervalSeconds: 3600, MetricMultiplier: 2, MetricMaxIntervalSeconds: 86400, MetricAgeSeconds: 30 * 24 * 60 * 60}
|
||
growthPoint0 := now.Add(-90 * time.Minute) // 首点(晚于发布时间)
|
||
if _, err := store.recordMetricWithPlan(ctx, MetricInput{WorkID: growthWork.ID, CollectedAt: growthPoint0, Likes: &growthLikes0}, growthSettings, growthPoint0); err != nil {
|
||
t.Fatal(err)
|
||
}
|
||
growthPoint1 := growthPoint0.Add(2 * time.Hour) // 第二个计划点:next = 首点+2h
|
||
if _, err := store.recordMetricWithPlan(ctx, MetricInput{WorkID: growthWork.ID, CollectedAt: growthPoint1, Likes: &growthLikes1}, growthSettings, growthPoint1); err != nil {
|
||
t.Fatal(err)
|
||
}
|
||
growth, err := store.LoadWorksGrowth(ctx, []string{growthWork.ID, work.ID}, WorkGrowthFilter{Hours: 24})
|
||
if err != nil {
|
||
t.Fatal(err)
|
||
}
|
||
stat, ok := growth[growthWork.ID]
|
||
if !ok || stat.Likes == nil || *stat.Likes != 300 || stat.Coverage != "partial" {
|
||
t.Fatalf("growth not computed: stat=%+v ok=%v", stat, ok)
|
||
}
|
||
if stat, ok := growth[work.ID]; !ok || stat.Likes == nil || *stat.Likes != 0 || stat.Coverage != "partial" {
|
||
t.Fatalf("single-point work growth: %+v ok=%v", stat, ok)
|
||
}
|
||
if _, err := store.RecordMetric(ctx, MetricInput{WorkID: work.ID, CollectedAt: now.Add(-time.Minute), Likes: &likes, CommentsCount: &comments, Shares: &shares}, Settings{LookbackDays: 30, NewWorkIntervalSeconds: 1800, MetricInitialIntervalSeconds: 3600, MetricMultiplier: 2, MetricMaxIntervalSeconds: 86400, MetricAgeSeconds: 30 * 24 * 60 * 60}, now); !errors.Is(err, ErrConflict) {
|
||
t.Fatalf("expected early metric point to be rejected, got %v", err)
|
||
}
|
||
if _, err := store.RecordMetric(ctx, MetricInput{WorkID: work.ID, CollectedAt: now.Add(time.Minute), Likes: &likes, CommentsCount: &comments, Shares: &shares}, Settings{LookbackDays: 30, NewWorkIntervalSeconds: 1800, MetricInitialIntervalSeconds: 3600, MetricMultiplier: 2, MetricMaxIntervalSeconds: 86400, MetricAgeSeconds: 30 * 24 * 60 * 60}, now.Add(-time.Minute)); !errors.Is(err, ErrInvalid) {
|
||
t.Fatalf("expected future metric point to be rejected, got %v", err)
|
||
}
|
||
|
||
rule, err := store.CreateRule(ctx, LeadRuleInput{Name: "咨询线索", Enabled: true, SourceType: SourceOwned, Topic: "咨询", IncludeKeywords: []string{"咨询"}, ExcludeKeywords: []string{"招聘"}, AIRequirement: "判断购买意向"})
|
||
if err != nil {
|
||
t.Fatal(err)
|
||
}
|
||
rules, err := store.ListRules(ctx, false)
|
||
if err != nil || len(rules) != 1 {
|
||
t.Fatalf("list rules: rules=%+v err=%v", rules, err)
|
||
}
|
||
updatedRule, err := store.UpdateRule(ctx, rule.ID, LeadRuleInput{Name: rule.Name, Enabled: true, SourceType: SourceOwned, Topic: rule.Topic, IncludeKeywords: []string{"咨询"}, ExcludeKeywords: []string{"招聘"}, AIRequirement: rule.AIRequirement})
|
||
if err != nil || updatedRule.ID != rule.ID {
|
||
t.Fatalf("update rule: rule=%+v err=%v", updatedRule, err)
|
||
}
|
||
result, err := store.AnalyzeComment(ctx, comment.ID, rule.ID, integrationAnalyzer{})
|
||
if err != nil || result.Status != "lead" {
|
||
t.Fatalf("analyze comment: result=%+v err=%v", result, err)
|
||
}
|
||
ruleResults, err := store.ListRuleResults(ctx, comment.ID, rule.ID)
|
||
if err != nil || len(ruleResults) != 1 {
|
||
t.Fatalf("list rule results: results=%+v err=%v", ruleResults, err)
|
||
}
|
||
leads, err := store.ListLeads(ctx, PlatformDouyin)
|
||
if err != nil || len(leads) != 1 {
|
||
t.Fatalf("list leads: leads=%+v err=%v", leads, err)
|
||
}
|
||
if _, err := store.SetRuleEnabled(ctx, rule.ID, false); err != nil {
|
||
t.Fatal(err)
|
||
}
|
||
|
||
material, selected, err := store.SelectMaterial(ctx, work.ID)
|
||
if err != nil || !selected || !material.Selected {
|
||
t.Fatalf("select material: material=%+v selected=%v err=%v", material, selected, err)
|
||
}
|
||
for step, status := range map[string]string{"download": "succeeded", "audio": "no_audio", "transcription": "no_speech"} {
|
||
if _, err := store.SetMaterialStep(ctx, work.ID, step, status, "ref-"+step, ""); err != nil {
|
||
t.Fatal(err)
|
||
}
|
||
}
|
||
if _, err := store.ConfirmRewrite(ctx, work.ID, "仿写同主题但不复制原文"); err != nil {
|
||
t.Fatal(err)
|
||
}
|
||
if _, err := store.SaveRewrite(ctx, work.ID, "新标题", "新口播"); err != nil {
|
||
t.Fatal(err)
|
||
}
|
||
|
||
collector := integrationCollector{work: WorkInput{Platform: PlatformDouyin, WorkKey: "creator-it-collected-" + stamp, SourceType: SourceOwned, SourceID: bigID, Title: "Collected", Body: "body", PublishedAt: &published, PublishedAtStatus: "verified", Likes: &likes}, comment: CommentInput{Platform: PlatformDouyin, CommentKey: "creator-it-collected-comment-" + stamp, WorkID: "", Content: "hello", CommentType: "top_level"}}
|
||
collector.comment.WorkID = ""
|
||
report, err := store.CollectSource(ctx, PlatformDouyin, SourceOwned, bigID, collector, now)
|
||
if err != nil || !report.PaginationComplete || report.WorksSeen != 1 || report.CommentsSaved != 2 {
|
||
t.Fatalf("collect source workflow failed: report=%+v err=%v", report, err)
|
||
}
|
||
}
|
||
|
||
func TestCreatorPostgresAccountMetricSnapshot(t *testing.T) {
|
||
store, phaseAStore, ctx := openCreatorIntegrationStore(t)
|
||
stamp := fmt.Sprintf("%d", time.Now().UnixNano())
|
||
accountID := createIntegrationAccount(t, ctx, phaseAStore, "metric"+stamp)
|
||
now := time.Now().UTC().Truncate(time.Second)
|
||
follower, following, favorited, aweme := int64(1200), int64(56), int64(56789), int64(34)
|
||
if err := store.RecordAccountMetric(ctx, AccountMetricInput{AccountID: accountID, CollectedAt: now, FollowerCount: &follower, FollowingCount: &following, TotalFavorited: &favorited, AwemeCount: &aweme}); err != nil {
|
||
t.Fatal(err)
|
||
}
|
||
points, err := store.ListAccountMetrics(ctx, accountID)
|
||
if err != nil || len(points) != 1 {
|
||
t.Fatalf("list account metrics: points=%+v err=%v", points, err)
|
||
}
|
||
if points[0].FollowerCount == nil || *points[0].FollowerCount != 1200 || points[0].TotalFavorited == nil || *points[0].TotalFavorited != 56789 {
|
||
t.Fatalf("account metric snapshot: %+v", points[0])
|
||
}
|
||
// UNIQUE(account_id, collected_at) 覆盖语义:同一时刻重采更新而非新增。
|
||
newFollower := int64(1300)
|
||
if err := store.RecordAccountMetric(ctx, AccountMetricInput{AccountID: accountID, CollectedAt: now, FollowerCount: &newFollower}); err != nil {
|
||
t.Fatal(err)
|
||
}
|
||
points, err = store.ListAccountMetrics(ctx, accountID)
|
||
if err != nil || len(points) != 1 || points[0].FollowerCount == nil || *points[0].FollowerCount != 1300 {
|
||
t.Fatalf("account metric upsert must overwrite same timestamp: points=%+v err=%v", points, err)
|
||
}
|
||
// 不同时刻追加为独立点。
|
||
if err := store.RecordAccountMetric(ctx, AccountMetricInput{AccountID: accountID, CollectedAt: now.Add(30 * time.Minute), FollowerCount: &newFollower}); err != nil {
|
||
t.Fatal(err)
|
||
}
|
||
points, err = store.ListAccountMetrics(ctx, accountID)
|
||
if err != nil || len(points) != 2 || !points[0].CollectedAt.Before(points[1].CollectedAt) {
|
||
t.Fatalf("account metric timeline must append: points=%+v err=%v", points, err)
|
||
}
|
||
// 负值拒绝。
|
||
bad := int64(-1)
|
||
if err := store.RecordAccountMetric(ctx, AccountMetricInput{AccountID: accountID, CollectedAt: now, FollowerCount: &bad}); !errors.Is(err, ErrInvalid) {
|
||
t.Fatalf("negative follower_count must be rejected: err=%v", err)
|
||
}
|
||
// 账号不存在时 FK 拒绝。
|
||
if err := store.RecordAccountMetric(ctx, AccountMetricInput{AccountID: "account-missing-" + fmt.Sprint(stamp), CollectedAt: now}); err == nil {
|
||
t.Fatal("missing account must violate FK")
|
||
}
|
||
}
|
||
|
||
func TestCreatorPostgresAccountMonitorViewAndCollectionStatus(t *testing.T) {
|
||
store, phaseAStore, ctx := openCreatorIntegrationStore(t)
|
||
stamp := fmt.Sprintf("%d", time.Now().UnixNano())
|
||
accountID := createIntegrationAccount(t, ctx, phaseAStore, "metric"+stamp)
|
||
now := time.Now().UTC().Truncate(time.Second)
|
||
|
||
// 无 checkpoint:采集状态为 pending。
|
||
status, err := store.GetAccountCollectionStatus(ctx, accountID)
|
||
if err != nil || status.Works.Status != "pending" || status.Comments.Status != "pending" {
|
||
t.Fatalf("pending collection status: %+v err=%v", status, err)
|
||
}
|
||
// 无作品:监控视图作品数为 0。
|
||
views, err := store.ListAccountMonitorViews(ctx)
|
||
if err != nil {
|
||
t.Fatal(err)
|
||
}
|
||
viewFor := func(id string) *AccountMonitorView {
|
||
for index := range views {
|
||
if views[index].ID == id {
|
||
return &views[index]
|
||
}
|
||
}
|
||
return nil
|
||
}
|
||
view := viewFor(accountID)
|
||
if view == nil || view.WorkCount != 0 || view.LatestPublishedAt != nil {
|
||
t.Fatalf("monitor view without works: %+v", view)
|
||
}
|
||
|
||
// 写入 checkpoint(succeeded + 窗口)与 owned 作品后:状态带下次窗口,视图聚合作品数。
|
||
lease, err := store.beginCheckpoint(ctx, SourceOwned, accountID, "works", now.Add(-24*time.Hour), now)
|
||
if err != nil {
|
||
t.Fatal(err)
|
||
}
|
||
if err := store.finishCheckpoint(ctx, SourceOwned, accountID, "works", lease, "succeeded", ""); err != nil {
|
||
t.Fatal(err)
|
||
}
|
||
published := now.Add(-time.Hour)
|
||
if _, _, err := store.UpsertWork(ctx, WorkInput{Platform: PlatformDouyin, WorkKey: "monitor-view-" + stamp, SourceType: SourceOwned, SourceID: accountID, Title: "MonitorView", PublishedAt: &published, PublishedAtStatus: "verified"}, now); err != nil {
|
||
t.Fatal(err)
|
||
}
|
||
status, err = store.GetAccountCollectionStatus(ctx, accountID)
|
||
if err != nil || status.Works.Status != "succeeded" || status.Works.LastCompletedAt == nil ||
|
||
status.Works.NextWindowStart == nil || status.Works.NextWindowEnd == nil {
|
||
t.Fatalf("succeeded collection status must carry next window: %+v err=%v", status, err)
|
||
}
|
||
if !status.Works.NextWindowStart.Before(*status.Works.NextWindowEnd) {
|
||
t.Fatalf("next window start must precede end: %+v", status.Works)
|
||
}
|
||
view = viewFor(accountID)
|
||
if view == nil || view.WorkCount != 1 || view.LatestPublishedAt == nil {
|
||
// 插入作品后重拉视图:聚合来自 creator_work_source 实时 JOIN。
|
||
views, err = store.ListAccountMonitorViews(ctx)
|
||
if err != nil {
|
||
t.Fatal(err)
|
||
}
|
||
view = viewFor(accountID)
|
||
if view == nil || view.WorkCount != 1 || view.LatestPublishedAt == nil {
|
||
t.Fatalf("monitor view aggregates: %+v", view)
|
||
}
|
||
}
|
||
// failed checkpoint 带错误信息(普通错误落在 failed,不是 blocked)。
|
||
lease, err = store.beginCheckpoint(ctx, SourceOwned, accountID, "comments", now.Add(-24*time.Hour), now)
|
||
if err != nil {
|
||
t.Fatal(err)
|
||
}
|
||
if err := store.failCheckpoint(ctx, SourceOwned, accountID, "comments", lease, errors.New("boom")); err == nil {
|
||
t.Fatal("failCheckpoint must return the primary error")
|
||
}
|
||
status, err = store.GetAccountCollectionStatus(ctx, accountID)
|
||
if err != nil || status.Comments.Status != "failed" || status.Comments.LastError == "" {
|
||
t.Fatalf("failed comments checkpoint surfaces error: %+v err=%v", status.Comments, err)
|
||
}
|
||
// 空账号 ID 拒绝。
|
||
if _, err := store.GetAccountCollectionStatus(ctx, " "); !errors.Is(err, ErrInvalid) {
|
||
t.Fatalf("blank account id must be rejected: err=%v", err)
|
||
}
|
||
}
|
||
|
||
func TestCreatorPostgresCollectionAllowsMutedReadAccount(t *testing.T) {
|
||
store, phaseAStore, ctx := openCreatorIntegrationStore(t)
|
||
stamp := time.Now().UnixNano()
|
||
accountID := createIntegrationAccount(t, ctx, phaseAStore, "muted"+fmt.Sprint(stamp))
|
||
if err := store.EnsureAccountProfile(ctx, accountID); err != nil {
|
||
t.Fatal(err)
|
||
}
|
||
if _, err := store.RecordVerifiedLoginResult(ctx, accountID, "sec_uid_"+accountID); err != nil {
|
||
t.Fatal(err)
|
||
}
|
||
now := time.Now().UTC()
|
||
if _, err := store.db.ExecContext(ctx, `UPDATE creator_account_profile SET business_status='muted' WHERE account_id=$1`, accountID); err != nil {
|
||
t.Fatal(err)
|
||
}
|
||
due, err := store.ListDueOwnedAccounts(ctx, now, 1800)
|
||
if err != nil || len(due) != 1 || due[0] != accountID {
|
||
t.Fatalf("muted account was not eligible for read collection: accounts=%v err=%v", due, err)
|
||
}
|
||
if _, err := store.db.ExecContext(ctx, `UPDATE creator_account_profile SET business_status='banned' WHERE account_id=$1`, accountID); err != nil {
|
||
t.Fatal(err)
|
||
}
|
||
due, err = store.ListDueOwnedAccounts(ctx, now, 1800)
|
||
if err != nil || len(due) != 0 {
|
||
t.Fatalf("banned account was eligible for read collection: accounts=%v err=%v", due, err)
|
||
}
|
||
}
|
||
|
||
func TestCreatorPostgresSettingsResetCollectionCheckpoint(t *testing.T) {
|
||
store, phaseAStore, ctx := openCreatorIntegrationStore(t)
|
||
stamp := time.Now().UnixNano()
|
||
accountID := createIntegrationAccount(t, ctx, phaseAStore, "settings"+fmt.Sprint(stamp))
|
||
before := time.Now().UTC()
|
||
lease, err := store.beginCheckpoint(ctx, SourceOwned, accountID, "works", before.Add(-48*time.Hour), before.Add(-47*time.Hour))
|
||
if err != nil {
|
||
t.Fatalf("begin checkpoint: %v", err)
|
||
}
|
||
if err := store.finishCheckpoint(ctx, SourceOwned, accountID, "works", lease, "succeeded", ""); err != nil {
|
||
t.Fatalf("finish checkpoint: %v", err)
|
||
}
|
||
if _, err := store.db.ExecContext(ctx, `UPDATE creator_collection_checkpoint SET window_start='2020-01-01T00:00:00Z',window_end='2020-01-02T00:00:00Z',cursor='old-cursor' WHERE source_type=$1 AND source_id=$2 AND collection_kind='works'`, SourceOwned, accountID); err != nil {
|
||
t.Fatal(err)
|
||
}
|
||
settings, err := store.GetSettings(ctx)
|
||
if err != nil {
|
||
t.Fatal(err)
|
||
}
|
||
settings.LookbackDays = 3
|
||
if _, err := store.UpdateSettings(ctx, SettingsUpdate{LookbackDays: settings.LookbackDays, NewWorkIntervalSeconds: settings.NewWorkIntervalSeconds, MetricInitialIntervalSeconds: settings.MetricInitialIntervalSeconds, MetricMultiplier: settings.MetricMultiplier, MetricMaxIntervalSeconds: settings.MetricMaxIntervalSeconds, MetricAgeSeconds: settings.MetricAgeSeconds, AIProvider: settings.AIProvider, AIModel: settings.AIModel, AIConfigured: settings.AIConfigured, TranscriptionProvider: settings.TranscriptionProvider, TranscriptionModel: settings.TranscriptionModel, TranscriptionConfigured: settings.TranscriptionConfigured}); err != nil {
|
||
t.Fatalf("update settings: %v", err)
|
||
}
|
||
checkpoint, err := store.checkpoint(ctx, SourceOwned, accountID, "works")
|
||
if err != nil {
|
||
t.Fatal(err)
|
||
}
|
||
if checkpoint.Cursor != "" || checkpoint.Status != "idle" || checkpoint.WindowStart.Before(time.Now().UTC().Add(-4*24*time.Hour)) {
|
||
t.Fatalf("settings did not reset checkpoint: %+v", checkpoint)
|
||
}
|
||
lease, err = store.beginCheckpoint(ctx, SourceOwned, accountID, "works", time.Now().UTC().Add(-time.Hour), time.Now().UTC())
|
||
if err != nil {
|
||
t.Fatalf("begin running checkpoint: %v", err)
|
||
}
|
||
settings.LookbackDays = 4
|
||
if _, err := store.UpdateSettings(ctx, SettingsUpdate{LookbackDays: settings.LookbackDays, NewWorkIntervalSeconds: settings.NewWorkIntervalSeconds, MetricInitialIntervalSeconds: settings.MetricInitialIntervalSeconds, MetricMultiplier: settings.MetricMultiplier, MetricMaxIntervalSeconds: settings.MetricMaxIntervalSeconds, MetricAgeSeconds: settings.MetricAgeSeconds, AIProvider: settings.AIProvider, AIModel: settings.AIModel, AIConfigured: settings.AIConfigured, TranscriptionProvider: settings.TranscriptionProvider, TranscriptionModel: settings.TranscriptionModel, TranscriptionConfigured: settings.TranscriptionConfigured}); !errors.Is(err, ErrConflict) {
|
||
t.Fatalf("settings changed during running checkpoint: err=%v", err)
|
||
}
|
||
if err := store.finishCheckpoint(ctx, SourceOwned, accountID, "works", lease, "succeeded", ""); err != nil {
|
||
t.Fatal(err)
|
||
}
|
||
}
|
||
|
||
func TestCreatorPostgresMetricPlanFollowsPublishedAt(t *testing.T) {
|
||
store, phaseAStore, ctx := openCreatorIntegrationStore(t)
|
||
stamp := fmt.Sprintf("%d", time.Now().UnixNano())
|
||
bigID := createIntegrationAccount(t, ctx, phaseAStore, "metric"+stamp)
|
||
if err := store.EnsureAccountProfile(ctx, bigID); err != nil {
|
||
t.Fatal(err)
|
||
}
|
||
published := time.Now().UTC().Add(-time.Hour).Truncate(time.Microsecond)
|
||
likes, comments, shares := int64(1), int64(1), int64(1)
|
||
work, _, err := store.UpsertWork(ctx, WorkInput{Platform: PlatformDouyin, WorkKey: "creator-it-metric-work-" + stamp, SourceType: SourceOwned, SourceID: bigID, Title: "Metric", Body: "body", PublishedAt: &published, PublishedAtStatus: "verified", Likes: &likes, CommentsCount: &comments, Shares: &shares}, published.Add(time.Hour))
|
||
if err != nil {
|
||
t.Fatal(err)
|
||
}
|
||
settings, err := store.GetSettings(ctx)
|
||
if err != nil {
|
||
t.Fatal(err)
|
||
}
|
||
settings.MetricInitialIntervalSeconds = 3600
|
||
settings.MetricMultiplier = 2
|
||
settings.MetricMaxIntervalSeconds = 55 * 3600
|
||
settings.MetricAgeSeconds = 100 * 3600
|
||
point := func(at time.Time) MetricInput {
|
||
return MetricInput{WorkID: work.ID, CollectedAt: at, Likes: &likes, CommentsCount: &comments, Shares: &shares}
|
||
}
|
||
if _, err := store.recordMetricWithPlan(ctx, point(published.Add(2*time.Hour)), settings, published.Add(2*time.Hour)); err != nil {
|
||
t.Fatal(err)
|
||
}
|
||
var next time.Time
|
||
var interval int64
|
||
if err := store.db.QueryRowContext(ctx, `SELECT next_plan_at,interval_seconds FROM creator_metric_plan WHERE work_id=$1`, work.ID).Scan(&next, &interval); err != nil {
|
||
t.Fatal(err)
|
||
}
|
||
if !next.Equal(published.Add(3*time.Hour)) || interval != 4*3600 {
|
||
t.Fatalf("first metric plan drifted: next=%s interval=%d", next, interval)
|
||
}
|
||
if _, err := store.recordMetricWithPlan(ctx, point(published.Add(3*time.Hour)), settings, published.Add(3*time.Hour)); err != nil {
|
||
t.Fatal(err)
|
||
}
|
||
if err := store.db.QueryRowContext(ctx, `SELECT next_plan_at,interval_seconds FROM creator_metric_plan WHERE work_id=$1`, work.ID).Scan(&next, &interval); err != nil {
|
||
t.Fatal(err)
|
||
}
|
||
if !next.Equal(published.Add(7*time.Hour)) || interval != 8*3600 {
|
||
t.Fatalf("second metric plan drifted: next=%s interval=%d", next, interval)
|
||
}
|
||
}
|
||
|
||
func prepareIntegrationActionFixture(t *testing.T, store *Store, phaseAStore *account.Store, ctx context.Context, stamp string) (string, string, Work, Comment, Strategy) {
|
||
t.Helper()
|
||
bigID := createIntegrationAccount(t, ctx, phaseAStore, "big"+stamp)
|
||
smallID := createIntegrationAccount(t, ctx, phaseAStore, "small"+stamp)
|
||
for _, accountID := range []string{bigID, smallID} {
|
||
if err := store.EnsureAccountProfile(ctx, accountID); err != nil {
|
||
t.Fatal(err)
|
||
}
|
||
if _, err := store.UpdateAccountProfile(ctx, accountID, AccountProfileUpdate{RealNameStatus: "unknown", BusinessStatus: "normal", CooldownSeconds: 86400}); err != nil {
|
||
t.Fatal(err)
|
||
}
|
||
if _, err := store.RecordVerifiedLoginResult(ctx, accountID, "sec_uid_"+accountID); err != nil {
|
||
t.Fatal(err)
|
||
}
|
||
}
|
||
if _, err := store.SetBigAccount(ctx, bigID, true); err != nil {
|
||
t.Fatal(err)
|
||
}
|
||
if err := store.SetRelation(ctx, bigID, smallID, true); err != nil {
|
||
t.Fatal(err)
|
||
}
|
||
now := time.Now().UTC().Truncate(time.Microsecond)
|
||
published := now.Add(-time.Hour)
|
||
likes, comments, shares := int64(1), int64(1), int64(1)
|
||
work, _, err := store.UpsertWork(ctx, WorkInput{Platform: PlatformDouyin, WorkKey: "creator-it-action-work-" + stamp, SourceType: SourceOwned, SourceID: bigID, Title: "Action", Body: "body", PublishedAt: &published, PublishedAtStatus: "verified", Likes: &likes, CommentsCount: &comments, Shares: &shares}, now)
|
||
if err != nil {
|
||
t.Fatal(err)
|
||
}
|
||
comment, _, err := store.SaveComment(ctx, CommentInput{Platform: PlatformDouyin, CommentKey: "creator-it-action-comment-" + stamp, WorkID: work.ID, AuthorUID: "interactor-" + stamp, Content: "hello", CommentType: "top_level"})
|
||
if err != nil {
|
||
t.Fatal(err)
|
||
}
|
||
strategy, err := store.CreateStrategy(ctx, bigID, StrategyInput{ExecutionAccountID: smallID, Position: 1, Enabled: true, EventTypes: []string{"comment"}, Action: ActionReplyComment, TargetType: "comment", CandidateTexts: []string{"已收到"}})
|
||
if err != nil {
|
||
t.Fatal(err)
|
||
}
|
||
strategies, err := store.ListStrategies(ctx, bigID)
|
||
if err != nil || len(strategies) != 1 {
|
||
t.Fatalf("list strategies: strategies=%+v err=%v", strategies, err)
|
||
}
|
||
return bigID, smallID, work, comment, strategy
|
||
}
|
||
|
||
func TestCreatorPostgresActionsAndMessaging(t *testing.T) {
|
||
store, phaseAStore, ctx := openCreatorIntegrationStore(t)
|
||
stamp := fmt.Sprintf("%d", time.Now().UnixNano())
|
||
bigID, smallID, work, comment, strategy := prepareIntegrationActionFixture(t, store, phaseAStore, ctx, stamp)
|
||
event := InteractionEvent{Platform: PlatformDouyin, ReceivingAccountID: bigID, EventKey: "creator-it-event-" + stamp, EventType: "comment", InteractorUID: comment.AuthorUID, CommentID: comment.ID, WorkID: work.ID}
|
||
automatic, err := store.ProcessAutomaticEvent(ctx, event, integrationExecutor{}, nil)
|
||
if err != nil || automatic.Operation == nil || automatic.Operation.State != "succeeded" || automatic.Event.State != "succeeded" {
|
||
t.Fatalf("automatic action failed: result=%+v err=%v", automatic, err)
|
||
}
|
||
events, err := store.ListEvents(ctx, bigID)
|
||
if err != nil || len(events) != 1 {
|
||
t.Fatalf("list events: events=%+v err=%v", events, err)
|
||
}
|
||
duplicate, err := store.ProcessAutomaticEvent(ctx, event, integrationExecutor{}, nil)
|
||
if err != nil || !duplicate.Duplicate || duplicate.Event.ID != automatic.Event.ID {
|
||
t.Fatalf("automatic event deduplication failed: result=%+v err=%v", duplicate, err)
|
||
}
|
||
cooldownEvent := event
|
||
cooldownEvent.EventKey += "-cooldown"
|
||
blocked, err := store.ProcessAutomaticEvent(ctx, cooldownEvent, integrationExecutor{}, nil)
|
||
if err != nil || blocked.Event.State != "blocked" {
|
||
t.Fatalf("cooldown did not block second event: result=%+v err=%v", blocked, err)
|
||
}
|
||
manualInput := OperationInput{IdempotencyKey: "creator-it-manual-" + stamp, Source: "manual", Action: ActionReplyComment, Platform: PlatformDouyin, AccountID: smallID, TargetUID: comment.AuthorUID, TargetCommentID: comment.ID, Text: "人工回复"}
|
||
manual, inserted, err := store.CreateOperation(ctx, manualInput)
|
||
if err != nil || !inserted {
|
||
t.Fatalf("create manual operation: operation=%+v inserted=%v err=%v", manual, inserted, err)
|
||
}
|
||
manualDuplicate, inserted, err := store.CreateOperation(ctx, manualInput)
|
||
if err != nil || inserted || manualDuplicate.ID != manual.ID {
|
||
t.Fatalf("manual operation idempotency failed: operation=%+v inserted=%v err=%v", manualDuplicate, inserted, err)
|
||
}
|
||
operations, err := store.ListOperations(ctx, smallID)
|
||
if err != nil || len(operations) != 2 {
|
||
t.Fatalf("list operations: operations=%+v err=%v", operations, err)
|
||
}
|
||
manual, err = store.ExecuteManualOperation(ctx, manual.ID, integrationExecutor{})
|
||
if err != nil || manual.State != "succeeded" {
|
||
t.Fatalf("execute manual operation: operation=%+v err=%v", manual, err)
|
||
}
|
||
dmOperation, inserted, err := store.CreateOperation(ctx, OperationInput{IdempotencyKey: "creator-it-dm-" + stamp, Source: "manual", Action: ActionDM, Platform: PlatformDouyin, AccountID: smallID, TargetUID: "peer-" + stamp, Text: "人工私信"})
|
||
if err != nil || !inserted {
|
||
t.Fatalf("create direct-message operation: operation=%+v inserted=%v err=%v", dmOperation, inserted, err)
|
||
}
|
||
dmOperation, err = store.ExecuteManualOperation(ctx, dmOperation.ID, integrationExecutor{})
|
||
if err != nil || dmOperation.State != "succeeded" {
|
||
t.Fatalf("execute direct-message operation: operation=%+v err=%v", dmOperation, err)
|
||
}
|
||
messageAt := time.Now().UTC().Truncate(time.Microsecond)
|
||
message, inserted, err := store.SaveMessage(ctx, MessageInput{Platform: PlatformDouyin, AccountID: smallID, PeerUID: "peer-" + stamp, PeerName: "Peer", PlatformMessageKey: "creator-it-message-" + stamp, Direction: "inbound", MessageType: "text", Text: "hello", MessageAt: &messageAt})
|
||
if err != nil || !inserted {
|
||
t.Fatalf("save message: message=%+v inserted=%v err=%v", message, inserted, err)
|
||
}
|
||
if _, inserted, err := store.SaveMessage(ctx, MessageInput{Platform: PlatformDouyin, AccountID: smallID, PeerUID: "peer-" + stamp, PlatformMessageKey: message.PlatformMessageKey, Direction: "inbound", MessageType: "text", Text: "hello", MessageAt: &messageAt}); err != nil || inserted {
|
||
t.Fatalf("message deduplication failed: inserted=%v err=%v", inserted, err)
|
||
}
|
||
conversations, err := store.ListConversations(ctx, smallID)
|
||
if err != nil || len(conversations) != 1 {
|
||
t.Fatalf("list conversations: conversations=%+v err=%v", conversations, err)
|
||
}
|
||
messages, err := store.ListMessages(ctx, conversations[0].ID)
|
||
if err != nil || len(messages) != 2 {
|
||
t.Fatalf("list messages: messages=%+v err=%v", messages, err)
|
||
}
|
||
if _, err := store.SetEventDisplayed(ctx, automatic.Event.ID, time.Now().UTC()); err != nil {
|
||
t.Fatal(err)
|
||
}
|
||
updated, err := store.UpdateStrategy(ctx, strategy.ID, StrategyInput{ExecutionAccountID: smallID, Position: 2, Enabled: true, EventTypes: []string{"comment", "like"}, Action: ActionReplyComment, TargetType: "comment", CandidateTexts: []string{"已更新"}})
|
||
if err != nil || updated.Position != 2 || len(updated.EventTypes) != 2 {
|
||
t.Fatalf("update strategy: strategy=%+v err=%v", updated, err)
|
||
}
|
||
if _, err := store.SetStrategyEnabled(ctx, strategy.ID, false); err != nil {
|
||
t.Fatal(err)
|
||
}
|
||
if err := store.DeleteStrategy(ctx, strategy.ID); err != nil {
|
||
t.Fatal(err)
|
||
}
|
||
if err := store.SetRelation(ctx, bigID, smallID, false); err != nil {
|
||
t.Fatal(err)
|
||
}
|
||
if _, err := store.SetBigAccount(ctx, bigID, false); err != nil {
|
||
t.Fatal(err)
|
||
}
|
||
|
||
start, end := time.Now().UTC().Add(-time.Hour), time.Now().UTC()
|
||
lease, err := store.beginCheckpoint(ctx, SourceOwned, bigID, "works", start, end)
|
||
if err != nil {
|
||
t.Fatal(err)
|
||
}
|
||
if _, err := store.beginCheckpoint(ctx, SourceOwned, bigID, "works", start, end); !errors.Is(err, ErrConflict) {
|
||
t.Fatalf("expected concurrent checkpoint claim conflict, got %v", err)
|
||
}
|
||
if err := store.finishCheckpoint(ctx, SourceOwned, bigID, "works", "stale-lease", "succeeded", ""); !errors.Is(err, ErrConflict) {
|
||
t.Fatalf("expected stale checkpoint completion conflict, got %v", err)
|
||
}
|
||
if err := store.saveCheckpointCursor(ctx, SourceOwned, bigID, "works", lease, "1"); err != nil {
|
||
t.Fatal(err)
|
||
}
|
||
if err := store.finishCheckpoint(ctx, SourceOwned, bigID, "works", lease, "succeeded", ""); err != nil {
|
||
t.Fatal(err)
|
||
}
|
||
}
|