Files
creator-hub/internal/creator/integration_test.go
T

858 lines
44 KiB
Go
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
package creator
import (
"context"
"database/sql"
"errors"
"fmt"
"net/url"
"os"
"testing"
"time"
"git.ipao.vip/rogee/creator-hub/internal/account"
hub "git.ipao.vip/rogee/creator-hub/internal/environment"
_ "github.com/jackc/pgx/v5/stdlib"
)
type integrationCredentialBridge struct{}
func (integrationCredentialBridge) Store(context.Context, account.CredentialReference, string, string) error {
return nil
}
func (integrationCredentialBridge) Delete(context.Context, account.CredentialReference, string) error {
return nil
}
type integrationAnalyzer struct{}
func (integrationAnalyzer) MatchTheme(context.Context, string, string, string) (bool, string, error) {
return true, "主题匹配", nil
}
func (integrationAnalyzer) MatchLead(context.Context, string, string, string) (bool, string, error) {
return true, "具备线索意向", nil
}
type integrationExecutor struct{}
func (integrationExecutor) Execute(context.Context, ActionRequest) (ActionResult, error) {
return ActionResult{State: "succeeded", Evidence: map[string]string{"platform_id": "creator-it"}}, nil
}
type integrationCollector struct {
work WorkInput
comment CommentInput
}
func (c integrationCollector) VerifyIdentity(context.Context, string) error { return nil }
func (c integrationCollector) ListWorks(_ context.Context, _, cursor string) (WorkPage, error) {
if cursor != "" {
return WorkPage{Items: nil, HasMore: false}, nil
}
return WorkPage{Items: []WorkInput{c.work}, HasMore: true, NextCursor: "1"}, nil
}
func (c integrationCollector) ListTopLevelComments(_ context.Context, _, cursor string) (CommentPage, error) {
if cursor != "" {
return CommentPage{Items: nil, HasMore: false}, nil
}
return CommentPage{Items: []CommentInput{c.comment}, HasMore: false}, nil
}
func openCreatorIntegrationStore(t *testing.T) (*Store, *account.Store, context.Context) {
t.Helper()
databaseURL := os.Getenv("CREATORHUB_POSTGRES_TEST_URL")
if databaseURL == "" {
t.Skip("set CREATORHUB_POSTGRES_TEST_URL to run CreatorHub PostgreSQL integration coverage")
}
ctx := context.Background()
admin, err := sql.Open("pgx", databaseURL)
if err != nil {
t.Fatal(err)
}
t.Cleanup(func() { _ = admin.Close() })
if err := admin.PingContext(ctx); err != nil {
t.Fatal(err)
}
schema := fmt.Sprintf("creatorhub_it_%d", time.Now().UnixNano())
if _, err := admin.ExecContext(ctx, "CREATE SCHEMA "+schema); err != nil {
t.Fatal(err)
}
t.Cleanup(func() {
if _, err := admin.ExecContext(ctx, "DROP SCHEMA "+schema+" CASCADE"); err != nil {
t.Errorf("drop test schema: %v", err)
}
})
parsed, err := url.Parse(databaseURL)
if err != nil {
t.Fatal(err)
}
query := parsed.Query()
query.Set("search_path", schema)
parsed.RawQuery = query.Encode()
testURL := parsed.String()
phaseAStore, err := account.Open(ctx, testURL)
if err != nil {
t.Fatal(err)
}
t.Cleanup(func() { _ = phaseAStore.Close() })
hubStore, err := hub.Open(ctx, testURL)
if err != nil {
t.Fatal(err)
}
t.Cleanup(func() { _ = hubStore.Close() })
store, err := Open(ctx, testURL)
if err != nil {
t.Fatal(err)
}
t.Cleanup(func() { _ = store.Close() })
return store, phaseAStore, ctx
}
func createIntegrationAccount(t *testing.T, ctx context.Context, phaseAStore *account.Store, suffix string) string {
t.Helper()
id := "cit" + suffix
account := account.Account{
ID: id,
Name: "Creator integration " + suffix,
Platform: PlatformDouyin,
PlatformAccountKey: "sec_uid_" + id,
Tags: []string{"integration"},
CredentialReference: account.CredentialReference{ID: "credential_" + id, Provider: "os_keyring"},
CredentialKey: "creatorhub/" + id,
}
if err := phaseAStore.CreateAccount(ctx, account, integrationCredentialBridge{}); err != nil {
t.Fatal(err)
}
return id
}
func TestCompetitorShareJobLifecycle(t *testing.T) {
store, _, ctx := openCreatorIntegrationStore(t)
stamp := fmt.Sprintf("%d", time.Now().UnixNano())
job, err := store.CreateCompetitorShareJob(ctx, CompetitorShareJobInput{
Platform: PlatformDouyin,
ShareURL: "https://v.douyin.com/share-" + stamp + "/",
Tags: []string{"重点监测"},
})
if err != nil || job.Status != CompetitorShareJobQueued || job.Attempts != 0 {
t.Fatalf("create share job: job=%+v err=%v", job, err)
}
due, err := store.ListDueCompetitorShareJobs(ctx, time.Now().UTC())
if err != nil || len(due) != 1 || due[0].ID != job.ID {
t.Fatalf("list due share jobs: jobs=%+v err=%v", due, err)
}
claimed, leaseToken, ok, err := store.ClaimCompetitorShareJob(ctx, job.ID, time.Now().UTC())
if err != nil || !ok || leaseToken == "" || claimed.Attempts != 1 || claimed.Status != CompetitorShareJobProcessing {
t.Fatalf("claim share job: job=%+v token=%q claimed=%v err=%v", claimed, leaseToken, ok, err)
}
if err := store.MarkCompetitorShareJob(ctx, job.ID, leaseToken, CompetitorShareJobQueued, "", "temporary failure"); err != nil {
t.Fatal(err)
}
claimed, leaseToken, ok, err = store.ClaimCompetitorShareJob(ctx, job.ID, time.Now().UTC())
if err != nil || !ok || leaseToken == "" || claimed.Attempts != 2 {
t.Fatalf("claim retry share job: job=%+v token=%q claimed=%v err=%v", claimed, leaseToken, ok, err)
}
if err := store.MarkCompetitorShareJob(ctx, job.ID, leaseToken, CompetitorShareJobFailed, "", "final failure"); err != nil {
t.Fatal(err)
}
failed, err := store.GetCompetitorShareJob(ctx, job.ID)
if err != nil || failed.Status != CompetitorShareJobFailed || failed.Attempts != 2 || failed.FailureReason != "final failure" {
t.Fatalf("failed share job: job=%+v err=%v", failed, err)
}
requeued, err := store.RetryCompetitorShareJob(ctx, job.ID)
if err != nil || requeued.Status != CompetitorShareJobQueued || requeued.Attempts != 0 || requeued.FailureReason != "final failure" {
t.Fatalf("requeue share job: job=%+v err=%v", requeued, err)
}
}
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)
}
}