fix(creator): 竞品同步 blocked 状态也排下次重试,避免定时获取永久失效

This commit is contained in:
2026-09-27 13:22:54 +08:00
parent e8205678ff
commit 7cb2a0a2b4
3 changed files with 40 additions and 3 deletions
+25 -3
View File
@@ -1778,8 +1778,20 @@ func syncCreatorCompetitorWithClaim(ctx context.Context, store *creator.Store, p
if !claimed {
return creator.CollectionReport{}, creator.ErrConflict
}
// blocked(环境暂不可用)也必须排下一次重试:否则 next_sync_at 置空后定时获取对该账号永久失效。
nextRetry := func() *time.Time {
nextBase := now
if competitor.NextSyncAt != nil {
nextBase = competitor.NextSyncAt.UTC()
}
next := creator.NextFixedRun(nextBase, time.Now().UTC(), time.Duration(settings.NewWorkIntervalSeconds)*time.Second)
if next.IsZero() {
return nil
}
return &next
}
blocked := func(blockErr error) (creator.CollectionReport, error) {
markErr := store.MarkCompetitorSync(ctx, competitorID, leaseToken, "blocked", "", blockErr.Error(), nil)
markErr := store.MarkCompetitorSync(ctx, competitorID, leaseToken, "blocked", "", blockErr.Error(), nextRetry())
return creator.CollectionReport{}, errors.Join(blockErr, markErr)
}
if competitor.Platform != creator.PlatformDouyin {
@@ -1834,7 +1846,7 @@ func syncCreatorCompetitorWithClaim(ctx context.Context, store *creator.Store, p
next := creator.NextFixedRun(nextBase, time.Now().UTC(), time.Duration(settings.NewWorkIntervalSeconds)*time.Second)
status := "failed"
if errors.Is(collectErr, creator.ErrUnavailable) || errors.Is(collectErr, creator.ErrConflict) {
status, next = "blocked", time.Time{}
status = "blocked"
}
var nextAt *time.Time
if !next.IsZero() {
@@ -1928,7 +1940,17 @@ func runCreatorScheduleOnce(ctx context.Context, store *creator.Store, phaseASto
continue
}
if claimed {
if markErr := store.MarkCompetitorSync(ctx, competitor.ID, leaseToken, "blocked", "", err.Error(), nil); markErr != nil {
// blocked 也排下一次重试(按新作品间隔),避免无采集账号时永久卡死。
nextBase := now
if competitor.NextSyncAt != nil {
nextBase = competitor.NextSyncAt.UTC()
}
next := creator.NextFixedRun(nextBase, now, time.Duration(settings.NewWorkIntervalSeconds)*time.Second)
var nextAt *time.Time
if !next.IsZero() {
nextAt = &next
}
if markErr := store.MarkCompetitorSync(ctx, competitor.ID, leaseToken, "blocked", "", err.Error(), nextAt); markErr != nil {
logrus.WithError(markErr).WithField("competitor_id", competitor.ID).Warn("creator competitor sync block update failed")
}
}
+15
View File
@@ -289,6 +289,21 @@ func TestCreatorPostgresContentAndWorkflow(t *testing.T) {
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)
}
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)