package creator import ( "context" "database/sql" "encoding/json" "errors" "fmt" "strings" "time" ) type WorkPage struct { Items []WorkInput NextCursor string HasMore bool } type CommentPage struct { Items []CommentInput NextCursor string HasMore bool } type PlatformCollector interface { ListWorks(context.Context, string, string) (WorkPage, error) ListTopLevelComments(context.Context, string, string) (CommentPage, error) } // WorkContextCollector lets a platform retain per-work request context across // the separate works and comments phases, including a resumed collection. type WorkContextCollector interface { SetWorkContext(string, string) error } type CollectionReport struct { SourceType string `json:"source_type"` SourceID string `json:"source_id"` WorksSeen int `json:"works_seen"` WorksSaved int `json:"works_saved"` CommentsSeen int `json:"comments_seen"` CommentsSaved int `json:"comments_saved"` WindowStart time.Time `json:"window_start"` WindowEnd time.Time `json:"window_end"` PaginationComplete bool `json:"pagination_complete"` } func collectPages[T any](ctx context.Context, fetch func(context.Context, string) ([]T, string, bool, error)) ([]T, error) { items := make([]T, 0) seenCursors := map[string]struct{}{"": {}} cursor := "" for page := 0; page < 100; page++ { pageItems, nextCursor, hasMore, err := fetch(ctx, cursor) if err != nil { return nil, err } items = append(items, pageItems...) if !hasMore { return items, nil } if nextCursor == "" { return nil, fmt.Errorf("%w: pagination cursor missing", ErrInvalid) } if _, exists := seenCursors[nextCursor]; exists { return nil, fmt.Errorf("%w: pagination cursor did not advance", ErrInvalid) } seenCursors[nextCursor] = struct{}{} cursor = nextCursor } return nil, fmt.Errorf("%w: pagination exceeded 100 pages", ErrInvalid) } func CollectWorkPages(ctx context.Context, collector PlatformCollector) ([]WorkInput, error) { if collector == nil { return nil, ErrUnavailable } return collectPages(ctx, func(ctx context.Context, cursor string) ([]WorkInput, string, bool, error) { page, err := collector.ListWorks(ctx, "", cursor) if err != nil { return nil, "", false, err } return page.Items, page.NextCursor, page.HasMore, nil }) } func CollectCommentPages(ctx context.Context, collector PlatformCollector, workKey string) ([]CommentInput, error) { if collector == nil || strings.TrimSpace(workKey) == "" { return nil, ErrInvalid } return collectPages(ctx, func(ctx context.Context, cursor string) ([]CommentInput, string, bool, error) { page, err := collector.ListTopLevelComments(ctx, workKey, cursor) if err != nil { return nil, "", false, err } for _, item := range page.Items { if item.CommentType == "reply" { return nil, "", false, fmt.Errorf("%w: reply comment returned by top-level collector", ErrInvalid) } } return page.Items, page.NextCursor, page.HasMore, nil }) } func NewCollectionWindow(now time.Time, lookbackDays int) (time.Time, time.Time, error) { if now.IsZero() || lookbackDays <= 0 || int64(lookbackDays) > maxDurationSeconds/(24*60*60) { return time.Time{}, time.Time{}, ErrInvalid } end := now.UTC() return end.Add(-time.Duration(lookbackDays) * 24 * time.Hour), end, nil } func publishedAtInCollectionWindow(publishedAt *time.Time, start, end time.Time) bool { if publishedAt == nil { return true } published := publishedAt.UTC() return !published.Before(start.UTC()) && !published.After(end.UTC()) } func checkpointID(sourceType, sourceID, kind string) string { return sourceType + ":" + sourceID + ":" + kind } type checkpointState struct { Cursor string LeaseToken string WindowStart time.Time WindowEnd time.Time Status string } func (s *Store) beginCheckpoint(ctx context.Context, sourceType, sourceID, kind string, start, end time.Time) (string, error) { if sourceType != SourceOwned && sourceType != SourceCompetitor || sourceID == "" || kind != "works" && kind != "comments" || start.IsZero() || end.IsZero() || !start.Before(end) { return "", ErrInvalid } leaseToken := newID("lease") var claimedToken string err := s.db.QueryRowContext(ctx, ` INSERT INTO creator_collection_checkpoint (id, source_type, source_id, collection_kind, window_start, window_end, status, lease_token, lease_until, last_started_at, last_error) VALUES ($1, $2, $3, $4, $5, $6, 'running', $7, now() + interval '10 minutes', now(), '') ON CONFLICT (source_type, source_id, collection_kind) DO UPDATE SET window_start = CASE WHEN creator_collection_checkpoint.status IN ('failed', 'running') THEN creator_collection_checkpoint.window_start ELSE EXCLUDED.window_start END, window_end = CASE WHEN creator_collection_checkpoint.status IN ('failed', 'running') THEN creator_collection_checkpoint.window_end ELSE EXCLUDED.window_end END, cursor = CASE WHEN creator_collection_checkpoint.status IN ('failed', 'running') THEN creator_collection_checkpoint.cursor ELSE '' END, status = 'running', lease_token = EXCLUDED.lease_token, lease_until = EXCLUDED.lease_until, last_started_at = EXCLUDED.last_started_at, last_error = '' WHERE creator_collection_checkpoint.status <> 'running' OR creator_collection_checkpoint.lease_until IS NULL OR creator_collection_checkpoint.lease_until <= now() RETURNING lease_token`, checkpointID(sourceType, sourceID, kind), sourceType, sourceID, kind, start.UTC(), end.UTC(), leaseToken).Scan(&claimedToken) if errors.Is(err, sql.ErrNoRows) { return "", ErrConflict } if err != nil { return "", databaseError(err) } return claimedToken, nil } func (s *Store) checkpoint(ctx context.Context, sourceType, sourceID, kind string) (checkpointState, error) { var state checkpointState err := s.db.QueryRowContext(ctx, `SELECT cursor,lease_token,window_start,window_end,status FROM creator_collection_checkpoint WHERE source_type=$1 AND source_id=$2 AND collection_kind=$3`, sourceType, sourceID, kind).Scan(&state.Cursor, &state.LeaseToken, &state.WindowStart, &state.WindowEnd, &state.Status) if err != nil { return checkpointState{}, rowError(err) } state.WindowStart, state.WindowEnd = state.WindowStart.UTC(), state.WindowEnd.UTC() return state, nil } // NextCollectionWindow preserves the fixed schedule grid. A failed attempt // retries its exact window; a completed attempt advances from the prior // scheduled end rather than from the wall-clock completion time. func (s *Store) NextCollectionWindow(ctx context.Context, sourceType, sourceID string, now time.Time, interval time.Duration, lookbackDays int) (time.Time, time.Time, error) { if sourceType != SourceOwned && sourceType != SourceCompetitor || sourceID == "" || now.IsZero() || interval <= 0 { return time.Time{}, time.Time{}, ErrInvalid } var end time.Time var status string err := s.db.QueryRowContext(ctx, `SELECT window_end,status FROM creator_collection_checkpoint WHERE source_type=$1 AND source_id=$2 AND collection_kind='works'`, sourceType, sourceID).Scan(&end, &status) if errors.Is(err, sql.ErrNoRows) { return NewCollectionWindow(now, lookbackDays) } if err != nil { return time.Time{}, time.Time{}, databaseError(err) } end = end.UTC() if status == "failed" || status == "blocked" || status == "running" { return end.Add(-time.Duration(lookbackDays) * 24 * time.Hour), end, nil } nextEnd := NextFixedRun(end, now.UTC(), interval) return nextEnd.Add(-time.Duration(lookbackDays) * 24 * time.Hour), nextEnd, nil } func (s *Store) MarkCollectionBlocked(ctx context.Context, sourceType, sourceID, reason string, now time.Time, lookbackDays int) error { start, end, err := NewCollectionWindow(now, lookbackDays) if err != nil { return err } if len(reason) > 2000 { return ErrInvalid } for _, kind := range []string{"works", "comments"} { _, err := s.db.ExecContext(ctx, ` INSERT INTO creator_collection_checkpoint (id,source_type,source_id,collection_kind,window_start,window_end,status,lease_token,lease_until,last_error) VALUES ($1,$2,$3,$4,$5,$6,'blocked','',NULL,$7) ON CONFLICT (source_type,source_id,collection_kind) DO UPDATE SET status='blocked', lease_token='', lease_until=NULL, last_error=$7 WHERE creator_collection_checkpoint.status <> 'running' OR creator_collection_checkpoint.lease_until IS NULL OR creator_collection_checkpoint.lease_until <= now()`, checkpointID(sourceType, sourceID, kind), sourceType, sourceID, kind, start, end, reason) if err != nil { return databaseError(err) } } return nil } func (s *Store) renewCheckpoint(ctx context.Context, sourceType, sourceID, kind, leaseToken string) error { if leaseToken == "" { return ErrInvalid } result, err := s.db.ExecContext(ctx, `UPDATE creator_collection_checkpoint SET lease_until=now()+interval '10 minutes' WHERE source_type=$1 AND source_id=$2 AND collection_kind=$3 AND lease_token=$4 AND status='running' AND lease_until > now()`, sourceType, sourceID, kind, leaseToken) if err != nil { return databaseError(err) } if affected, err := result.RowsAffected(); err != nil { return err } else if affected != 1 { return ErrConflict } return nil } func (s *Store) saveCheckpointCursor(ctx context.Context, sourceType, sourceID, kind, leaseToken, cursor string) error { if err := s.renewCheckpoint(ctx, sourceType, sourceID, kind, leaseToken); err != nil { return err } result, err := s.db.ExecContext(ctx, `UPDATE creator_collection_checkpoint SET cursor=$5 WHERE source_type=$1 AND source_id=$2 AND collection_kind=$3 AND lease_token=$4 AND status='running' AND lease_until > now()`, sourceType, sourceID, kind, leaseToken, cursor) if err != nil { return databaseError(err) } if affected, err := result.RowsAffected(); err != nil { return err } else if affected != 1 { return ErrConflict } return nil } func (s *Store) finishCheckpoint(ctx context.Context, sourceType, sourceID, kind, leaseToken, status, reason string) error { if status != "succeeded" && status != "failed" && status != "blocked" || leaseToken == "" { return ErrInvalid } result, err := s.db.ExecContext(ctx, `UPDATE creator_collection_checkpoint SET status=$5, cursor=CASE WHEN $5='succeeded' THEN '' ELSE cursor END, lease_until=NULL, last_completed_at=CASE WHEN $5='succeeded' THEN now() ELSE last_completed_at END, last_error=$6 WHERE source_type=$1 AND source_id=$2 AND collection_kind=$3 AND lease_token=$4 AND status='running' AND lease_until > now()`, sourceType, sourceID, kind, leaseToken, status, reason) if err != nil { return databaseError(err) } if affected, err := result.RowsAffected(); err != nil { return err } else if affected != 1 { return ErrConflict } return nil } func (s *Store) failCheckpoint(ctx context.Context, sourceType, sourceID, kind, leaseToken string, primary error) error { status := "failed" if errors.Is(primary, ErrConflict) || errors.Is(primary, ErrUnavailable) { status = "blocked" } if err := s.finishCheckpoint(ctx, sourceType, sourceID, kind, leaseToken, status, primary.Error()); err != nil { return errors.Join(primary, err) } return primary } func collectPagesFromCursor[T any](ctx context.Context, initialCursor string, fetch func(context.Context, string) ([]T, string, bool, error), afterPage func([]T, string, bool) error) ([]T, error) { items := make([]T, 0) seenCursors := map[string]struct{}{initialCursor: {}} cursor := initialCursor for page := 0; page < 100; page++ { pageItems, nextCursor, hasMore, err := fetch(ctx, cursor) if err != nil { return nil, err } items = append(items, pageItems...) if afterPage != nil { if err := afterPage(pageItems, nextCursor, hasMore); err != nil { return nil, err } } if !hasMore { return items, nil } if nextCursor == "" { return nil, fmt.Errorf("%w: pagination cursor missing", ErrInvalid) } if _, exists := seenCursors[nextCursor]; exists { return nil, fmt.Errorf("%w: pagination cursor did not advance", ErrInvalid) } seenCursors[nextCursor] = struct{}{} cursor = nextCursor } return nil, fmt.Errorf("%w: pagination exceeded 100 pages", ErrInvalid) } type commentCheckpointCursor struct { WorkKey string `json:"work_key"` Cursor string `json:"cursor"` } func encodeCommentCheckpoint(workKey, cursor string) (string, error) { value, err := json.Marshal(commentCheckpointCursor{WorkKey: workKey, Cursor: cursor}) if err != nil { return "", fmt.Errorf("encode comment checkpoint: %w", err) } return string(value), nil } func decodeCommentCheckpoint(value string) (commentCheckpointCursor, error) { var result commentCheckpointCursor if value == "" { return result, nil } if err := json.Unmarshal([]byte(value), &result); err != nil || result.WorkKey == "" { return commentCheckpointCursor{}, fmt.Errorf("%w: invalid comment checkpoint", ErrInvalid) } return result, nil } func (s *Store) ListDueOwnedAccounts(ctx context.Context, now time.Time, intervalSeconds int64) ([]string, error) { if now.IsZero() || intervalSeconds <= 0 || intervalSeconds > maxDurationSeconds { return nil, ErrInvalid } cutoff := now.UTC().Add(-time.Duration(intervalSeconds) * time.Second) rows, err := s.db.QueryContext(ctx, ` SELECT account.id FROM social_account account JOIN creator_account_profile profile ON profile.account_id=account.id AND profile.business_status='normal' LEFT JOIN creator_collection_checkpoint works_checkpoint ON works_checkpoint.source_type='owned' AND works_checkpoint.source_id=account.id AND works_checkpoint.collection_kind='works' LEFT JOIN creator_collection_checkpoint comments_checkpoint ON comments_checkpoint.source_type='owned' AND comments_checkpoint.source_id=account.id AND comments_checkpoint.collection_kind='comments' WHERE account.platform IN ('douyin', 'xiaohongshu') AND account.authorization_status='authorized' AND profile.login_status='logged_in' AND profile.big_account=true AND COALESCE(works_checkpoint.status, '') <> 'blocked' AND COALESCE(comments_checkpoint.status, '') <> 'blocked' AND (works_checkpoint.id IS NULL OR comments_checkpoint.id IS NULL OR works_checkpoint.status='failed' OR comments_checkpoint.status='failed' OR works_checkpoint.last_completed_at IS NULL OR comments_checkpoint.last_completed_at IS NULL OR works_checkpoint.last_completed_at <= $1 OR comments_checkpoint.last_completed_at <= $1) ORDER BY account.id`, cutoff) if err != nil { return nil, databaseError(err) } defer rows.Close() accounts := make([]string, 0) for rows.Next() { var accountID string if err := rows.Scan(&accountID); err != nil { return nil, err } accounts = append(accounts, accountID) } return accounts, rows.Err() } func (s *Store) CollectSource(ctx context.Context, platform, sourceType, sourceID string, collector PlatformCollector, now time.Time) (CollectionReport, error) { if !ValidatePlatform(platform) || sourceID == "" || (sourceType != SourceOwned && sourceType != SourceCompetitor) || collector == nil { return CollectionReport{}, ErrInvalid } settings, err := s.GetSettings(ctx) if err != nil { return CollectionReport{}, err } start, end, err := NewCollectionWindow(now, settings.LookbackDays) if err != nil { return CollectionReport{}, err } report := CollectionReport{SourceType: sourceType, SourceID: sourceID, WindowStart: start, WindowEnd: end} worksLease, err := s.beginCheckpoint(ctx, sourceType, sourceID, "works", start, end) if err != nil { return report, err } worksCheckpoint, err := s.checkpoint(ctx, sourceType, sourceID, "works") if err != nil { return report, err } report.WindowStart, report.WindowEnd = worksCheckpoint.WindowStart, worksCheckpoint.WindowEnd seenWorks := map[string]string{} _, err = collectPagesFromCursor(ctx, worksCheckpoint.Cursor, func(ctx context.Context, cursor string) ([]WorkInput, string, bool, error) { page, err := collector.ListWorks(ctx, sourceID, cursor) if err != nil { return nil, "", false, err } return page.Items, page.NextCursor, page.HasMore, nil }, func(pageItems []WorkInput, nextCursor string, hasMore bool) error { if err := s.renewCheckpoint(ctx, sourceType, sourceID, "works", worksLease); err != nil { return err } for _, work := range pageItems { if err := s.renewCheckpoint(ctx, sourceType, sourceID, "works", worksLease); err != nil { return err } report.WorksSeen++ if !publishedAtInCollectionWindow(work.PublishedAt, report.WindowStart, report.WindowEnd) { continue } work.Platform, work.SourceType, work.SourceID = platform, sourceType, sourceID if _, exists := seenWorks[work.WorkKey]; exists { continue } savedWork, _, err := s.UpsertWork(ctx, work, report.WindowEnd) if err != nil { return err } if work.PublishedAt == nil || !work.PublishedAt.After(report.WindowEnd) { if err := s.EnsureMetricPlan(ctx, savedWork.ID, settings); err != nil { return err } } if work.Likes != nil || work.CommentsCount != nil || work.Shares != nil { if _, metricErr := s.RecordMetric(ctx, MetricInput{WorkID: savedWork.ID, CollectedAt: now.UTC(), Likes: work.Likes, CommentsCount: work.CommentsCount, Shares: work.Shares}, settings, now.UTC()); metricErr != nil && !errors.Is(metricErr, ErrConflict) { return metricErr } } seenWorks[work.WorkKey] = savedWork.ID report.WorksSaved++ } if hasMore { return s.saveCheckpointCursor(ctx, sourceType, sourceID, "works", worksLease, nextCursor) } return nil }) if err != nil { return report, s.failCheckpoint(ctx, sourceType, sourceID, "works", worksLease, err) } if err := s.finishCheckpoint(ctx, sourceType, sourceID, "works", worksLease, "succeeded", ""); err != nil { return report, err } commentsLease, err := s.beginCheckpoint(ctx, sourceType, sourceID, "comments", report.WindowStart, report.WindowEnd) if err != nil { return report, err } commentsCheckpoint, err := s.checkpoint(ctx, sourceType, sourceID, "comments") if err != nil { return report, err } report.WindowStart, report.WindowEnd = commentsCheckpoint.WindowStart, commentsCheckpoint.WindowEnd works, err := s.ListWorks(ctx, WorkFilter{Platform: platform, SourceType: sourceType, SourceID: sourceID, PublishedAfter: &report.WindowStart, PublishedBefore: &report.WindowEnd}) if err != nil { return report, err } startIndex, initialCommentCursor := 0, "" if commentsCheckpoint.Cursor != "" { position, err := decodeCommentCheckpoint(commentsCheckpoint.Cursor) if err != nil { return report, s.failCheckpoint(ctx, sourceType, sourceID, "comments", commentsLease, err) } initialCommentCursor = position.Cursor for index := range works { if works[index].WorkKey == position.WorkKey { startIndex = index break } if index == len(works)-1 { return report, s.failCheckpoint(ctx, sourceType, sourceID, "comments", commentsLease, fmt.Errorf("%w: comment checkpoint work missing", ErrInvalid)) } } } for index := startIndex; index < len(works); index++ { work := works[index] if contextual, ok := collector.(WorkContextCollector); ok { if err := contextual.SetWorkContext(work.WorkKey, work.OriginalURL); err != nil { return report, s.failCheckpoint(ctx, sourceType, sourceID, "comments", commentsLease, err) } } cursor := initialCommentCursor initialCommentCursor = "" _, err := collectPagesFromCursor(ctx, cursor, func(ctx context.Context, cursor string) ([]CommentInput, string, bool, error) { page, err := collector.ListTopLevelComments(ctx, work.WorkKey, cursor) if err != nil { return nil, "", false, err } for _, item := range page.Items { if item.CommentType == "reply" { return nil, "", false, fmt.Errorf("%w: reply comment returned by top-level collector", ErrInvalid) } } return page.Items, page.NextCursor, page.HasMore, nil }, func(pageItems []CommentInput, nextCursor string, hasMore bool) error { if err := s.renewCheckpoint(ctx, sourceType, sourceID, "comments", commentsLease); err != nil { return err } for _, comment := range pageItems { if err := s.renewCheckpoint(ctx, sourceType, sourceID, "comments", commentsLease); err != nil { return err } report.CommentsSeen++ if !publishedAtInCollectionWindow(comment.PublishedAt, report.WindowStart, report.WindowEnd) || comment.CommentType == "reply" { continue } comment.Platform, comment.WorkID = platform, work.ID if _, _, err := s.SaveComment(ctx, comment); err != nil { return err } report.CommentsSaved++ } if hasMore { checkpoint, err := encodeCommentCheckpoint(work.WorkKey, nextCursor) if err != nil { return err } return s.saveCheckpointCursor(ctx, sourceType, sourceID, "comments", commentsLease, checkpoint) } return nil }) if err != nil { return report, s.failCheckpoint(ctx, sourceType, sourceID, "comments", commentsLease, err) } if index+1 < len(works) { checkpoint, err := encodeCommentCheckpoint(works[index+1].WorkKey, "") if err != nil { return report, s.failCheckpoint(ctx, sourceType, sourceID, "comments", commentsLease, err) } if err := s.saveCheckpointCursor(ctx, sourceType, sourceID, "comments", commentsLease, checkpoint); err != nil { return report, s.failCheckpoint(ctx, sourceType, sourceID, "comments", commentsLease, err) } } } if err := s.finishCheckpoint(ctx, sourceType, sourceID, "comments", commentsLease, "succeeded", ""); err != nil { return report, err } report.PaginationComplete = true return report, nil }