From 0c0225f9c664e45e021697af75508b9a6ba0ba4c Mon Sep 17 00:00:00 2001 From: Mathias Date: Sat, 6 Jun 2026 21:27:38 +0200 Subject: [PATCH] feat(runner): process candidates newest-first within each pass (ADR-018) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Restructures RunOnce from per-channel inline processing to collect-sort-process: Phase 1 — discover, persist (UpsertVideo), apply pre-filters (seen/manual/backoff) and collect surviving candidates with their discovery position. Phase 2 — sort candidates by published_at DESC, NULLS LAST, pos ASC tiebreak so videos with no publish date never jump ahead of dated content. Phase 3 — process in sorted order through the unchanged globalFetchGate. Before (per-channel): chanA=[v-old, v-mid], chanB=[v-new, v-null] → [v-old, v-mid, v-new, v-null] After (newest-first): [v-new, v-mid, v-old, v-null] Same set of videos processed; only the order changes within a pass. All existing behaviour is preserved: failure isolation, backoff skip, manual mode, channel-unavailable, stats. In-memory sort; no new table or persisted queue. The ordering is onboarding prioritisation — new users get summaries of their most recent, relevant videos first; the back-catalogue fills in behind across subsequent passes. Both this background batch and the foreground 'Try now' button honour the shared globalFetchGate: rate limiting is respected, not evaded. --- internal/runner/runner.go | 179 ++++++++++++++++++++++----------- internal/runner/runner_test.go | 92 +++++++++++++++++ 2 files changed, 210 insertions(+), 61 deletions(-) diff --git a/internal/runner/runner.go b/internal/runner/runner.go index 4e4fcae..1d3b17c 100644 --- a/internal/runner/runner.go +++ b/internal/runner/runner.go @@ -10,11 +10,13 @@ package runner import ( + "cmp" "context" "errors" "fmt" "log/slog" "os" + "slices" "time" "gitea.d-ma.be/mathias/tapir/internal/domain" @@ -22,6 +24,36 @@ import ( "gitea.d-ma.be/mathias/tapir/internal/usecase" ) +// passCandidate is a video that passed all pre-filters (seen/manual/backoff) +// and is queued for transcript fetch + summarization in this pass. +type passCandidate struct { + v domain.Video + pos int // discovery position — used as a stable tiebreak when published_at ties +} + +// compareNewestFirst orders candidates by published_at descending, NULLS LAST, +// with pos ascending as a stable tiebreak. Videos with a zero published_at +// (schema 001: nullable) sort after all dated videos regardless of pos. +func compareNewestFirst(a, b passCandidate) int { + aNull := a.v.PublishedAt.IsZero() + bNull := b.v.PublishedAt.IsZero() + switch { + case aNull && bNull: + return cmp.Compare(a.pos, b.pos) + case aNull: + return 1 // a is null → after b + case bNull: + return -1 // b is null → after a + } + if !a.v.PublishedAt.Equal(b.v.PublishedAt) { + if a.v.PublishedAt.After(b.v.PublishedAt) { + return -1 // newer first + } + return 1 + } + return cmp.Compare(a.pos, b.pos) // same timestamp: preserve discovery order +} + // VideoStore is the durable persistence the run loop needs: assign a stable id + // metadata, read the already-summarized set, and (for manual summarization mode) // read the user's mode + queued videos and clear a video's queue flag once it has @@ -101,19 +133,31 @@ type Stats struct { ChannelUnavailable int // channels that returned HTTP 404 (deleted/private) } -// RunOnce performs a single pass over the user's subscriptions. Per-item errors -// are logged and collected (one bad video or channel does not abort the pass) -// and returned joined alongside the Stats gathered. +// RunOnce performs a single pass over the user's subscriptions in three phases: +// +// 1. Discovery: walk all channels, persist each candidate video (UpsertVideo), +// apply pre-filters (seen/manual/backoff) — same as before. +// 2. Sort: order the surviving candidates newest-first (published_at DESC, NULLS +// LAST) so new users get summaries of their most recent, relevant videos first; +// the back-catalogue fills in behind across subsequent passes. +// 3. Process: feed candidates to the engine in sorted order through the shared +// globalFetchGate — the gate is unchanged and still governs honest rate pacing. +// +// All existing behaviour is preserved: per-item failure isolation, the rate-limit +// backoff skip, manual mode, channel-unavailable handling, and stats accounting. +// Only the processing order changes within a pass. func (r *Runner) RunOnce(ctx context.Context) (Stats, error) { var ( - stats Stats - errs []error + stats Stats + errs []error + candidates []passCandidate + pos int ) // Throttle transcript fetches: unauthenticated caption scraping gets // soft-throttled by YouTube under heavy back-to-back volume (captionTracks // silently stripped from the player response). A small per-video delay keeps - // a full pass under the radar. TAPIR_FETCH_DELAY (Go duration), 0 = off. + // the fetch rate polite. TAPIR_FETCH_DELAY (Go duration), 0 = off. fetchDelay, _ := time.ParseDuration(os.Getenv("TAPIR_FETCH_DELAY")) seen, err := r.store.SeenVideoIDs(ctx, r.userID) @@ -121,10 +165,9 @@ func (r *Runner) RunOnce(ctx context.Context) (Stats, error) { return stats, fmt.Errorf("runner: load seen videos: %w", err) } - // Summarization mode (per-user, ADR-012). Auto = summarize every unseen video - // (the original behavior). Manual = still discover/persist videos so the user - // sees them, but only summarize the ones explicitly queued via the web UI - // (summarize_requested). The queued set is loaded once per pass, like seen. + // Summarization mode (per-user, ADR-012). Auto = summarize every unseen video. + // Manual = discover/persist videos (visible in list) but only summarize ones + // explicitly queued via the web UI (summarize_requested). auto, err := r.store.GetAutoSummarize(ctx, r.userID) if err != nil { return stats, fmt.Errorf("runner: load summarize mode: %w", err) @@ -137,10 +180,9 @@ func (r *Runner) RunOnce(ctx context.Context) (Stats, error) { } } - // Rate-limit backoff: videos that 429'd on a prior pass, mapped to when. Inside - // the backoff window they are skipped before any caption fetch, so a throttled - // IP is not hammered. Loaded once per pass (like seen/requested). Disabled when - // backoff <= 0 ("always retry"). + // Rate-limit backoff: videos that 429'd on a prior pass, mapped to when. + // Within the backoff window they are skipped before any caption fetch. + // Disabled when backoff <= 0 ("always retry"). var rateLimited map[string]time.Time if r.backoff > 0 { rateLimited, err = r.store.RateLimitedVideoIDs(ctx, r.userID) @@ -154,6 +196,10 @@ func (r *Runner) RunOnce(ctx context.Context) (Stats, error) { return stats, fmt.Errorf("runner: list subscriptions: %w", err) } + // ── Phase 1: discover, persist, filter ─────────────────────────────────── + // Walk all channels. Persist every video (UpsertVideo) so it appears in the + // list regardless of whether it will be summarized this pass. Apply pre-filters + // and collect surviving candidates with their discovery position. for _, sub := range subs { vids, err := r.src.NewVideos(ctx, sub) if err != nil { @@ -172,7 +218,7 @@ func (r *Runner) RunOnce(ctx context.Context) (Stats, error) { } for _, v := range vids { stats.Candidates++ - v.UserID = r.userID // keep the dedup/FK key consistent with config + v.UserID = r.userID id, err := r.store.UpsertVideo(ctx, v) if err != nil { @@ -186,69 +232,80 @@ func (r *Runner) RunOnce(ctx context.Context) (Stats, error) { stats.SkippedSeen++ continue } - seen[id] = true // also guard against the same video within this pass + seen[id] = true // guard against duplicates within this pass // Manual mode: skip summarization for videos the user has not queued. - // Discovery already happened (UpsertVideo above), so the new video is + // Discovery already happened (UpsertVideo above), so the video is // visible in the list; it just isn't summarized until requested. if !auto && !requested[id] { stats.SkippedManual++ continue } - // Still inside the rate-limit backoff window: skip without fetching, so - // we don't re-hit a caption endpoint that just 429'd us. After the window - // expires the video falls through and is retried normally. + // Still inside the rate-limit backoff window: skip without fetching. if at, ok := rateLimited[id]; ok && r.now().Sub(at) < r.backoff { stats.SkippedRateLimited++ r.log.Info("skipped video (rate-limited, backing off)", "video", v.ProviderVideoID, "title", v.Title) continue } - if fetchDelay > 0 { - time.Sleep(fetchDelay) - } - res, err := r.engine.ProcessNewVideo(ctx, v) - if err != nil { - errs = append(errs, fmt.Errorf("process %q: %w", v.ProviderVideoID, err)) + candidates = append(candidates, passCandidate{v: v, pos: pos}) + pos++ + } + } + + // ── Phase 2: sort newest-first, NULLS LAST ──────────────────────────────── + // Within this pass, process the newest videos first so a new user gets + // summaries of their most recent content quickly; the back-catalogue fills in + // behind across subsequent passes. Both this background batch and the foreground + // "Try now" button honour the shared globalFetchGate — ordering is onboarding + // prioritisation, not rate-limit evasion. + slices.SortStableFunc(candidates, compareNewestFirst) + + // ── Phase 3: process in sorted order ───────────────────────────────────── + for _, c := range candidates { + if fetchDelay > 0 { + time.Sleep(fetchDelay) + } + res, err := r.engine.ProcessNewVideo(ctx, c.v) + if err != nil { + errs = append(errs, fmt.Errorf("process %q: %w", c.v.ProviderVideoID, err)) + stats.Errors++ + continue + } + id := c.v.ID + switch { + case res.Skipped && res.TranscriptSource == string(domain.SourceRateLimited): + // Fresh 429: stamp the backoff clock so the next pass skips it. + stats.SkippedRateLimited++ + if err := r.store.SetTranscriptStatus(ctx, r.userID, id, "rate_limited"); err != nil { + errs = append(errs, fmt.Errorf("set rate_limited status %q: %w", c.v.ProviderVideoID, err)) stats.Errors++ - continue } - switch { - case res.Skipped && res.TranscriptSource == string(domain.SourceRateLimited): - // Fresh 429 this pass: persist rate_limited (stamps the backoff clock) - // so the next pass skips it until the window expires. - stats.SkippedRateLimited++ - if err := r.store.SetTranscriptStatus(ctx, r.userID, id, "rate_limited"); err != nil { - errs = append(errs, fmt.Errorf("set rate_limited status %q: %w", v.ProviderVideoID, err)) - stats.Errors++ - } - r.log.Info("skipped video (rate-limited)", "video", v.ProviderVideoID, "title", v.Title) - case res.Skipped: - stats.SkippedNoText++ - if err := r.store.SetTranscriptStatus(ctx, r.userID, id, "none"); err != nil { - errs = append(errs, fmt.Errorf("set none status %q: %w", v.ProviderVideoID, err)) - stats.Errors++ - } - r.log.Info("skipped video (no transcript)", "video", v.ProviderVideoID, "title", v.Title) - case res.Summary != nil: - stats.Summarized++ - if err := r.store.SetTranscriptStatus(ctx, r.userID, id, "fetched"); err != nil { - errs = append(errs, fmt.Errorf("set fetched status %q: %w", v.ProviderVideoID, err)) - stats.Errors++ - } - // In manual mode the video was processed because it was queued; - // clear the flag so it is not re-summarized and the UI drops the - // "Queued" chip. (Auto mode never sets the flag.) - if !auto { - if err := r.store.ClearSummarizeRequested(ctx, r.userID, id); err != nil { - errs = append(errs, fmt.Errorf("clear summarize flag %q: %w", v.ProviderVideoID, err)) - stats.Errors++ - } - } - r.log.Info("summarized video", "video", v.ProviderVideoID, "title", v.Title, - "provider", res.Summary.AIProvider, "model", res.Summary.AIModel) + r.log.Info("skipped video (rate-limited)", "video", c.v.ProviderVideoID, "title", c.v.Title) + case res.Skipped: + stats.SkippedNoText++ + if err := r.store.SetTranscriptStatus(ctx, r.userID, id, "none"); err != nil { + errs = append(errs, fmt.Errorf("set none status %q: %w", c.v.ProviderVideoID, err)) + stats.Errors++ } + r.log.Info("skipped video (no transcript)", "video", c.v.ProviderVideoID, "title", c.v.Title) + case res.Summary != nil: + stats.Summarized++ + if err := r.store.SetTranscriptStatus(ctx, r.userID, id, "fetched"); err != nil { + errs = append(errs, fmt.Errorf("set fetched status %q: %w", c.v.ProviderVideoID, err)) + stats.Errors++ + } + // In manual mode the video was explicitly queued; clear the flag so + // it is not re-summarized and the UI drops the "Queued" chip. + if !auto { + if err := r.store.ClearSummarizeRequested(ctx, r.userID, id); err != nil { + errs = append(errs, fmt.Errorf("clear summarize flag %q: %w", c.v.ProviderVideoID, err)) + stats.Errors++ + } + } + r.log.Info("summarized video", "video", c.v.ProviderVideoID, "title", c.v.Title, + "provider", res.Summary.AIProvider, "model", res.Summary.AIModel) } } diff --git a/internal/runner/runner_test.go b/internal/runner/runner_test.go index 62928a0..371d3c6 100644 --- a/internal/runner/runner_test.go +++ b/internal/runner/runner_test.go @@ -123,6 +123,12 @@ func vid(provID, title string) domain.Video { return domain.Video{UserID: testUser, Provider: domain.ProviderYouTube, ProviderVideoID: provID, Title: title} } +func vidAt(provID, title string, publishedAt time.Time) domain.Video { + v := vid(provID, title) + v.PublishedAt = publishedAt + return v +} + func quietLogger() *slog.Logger { return slog.New(slog.NewTextHandler(io.Discard, nil)) } @@ -298,3 +304,89 @@ func TestRunOnce_UpsertsEveryCandidate(t *testing.T) { require.NoError(t, err) require.Len(t, st.upserted, 2, "every candidate is upserted, including seen ones") } + +// TestRunOnce_NewestFirstOrdering asserts that within a pass, candidates are +// processed newest-first (published_at DESC, NULLS LAST) across all channels, +// and that the set of processed videos is identical to what per-channel inline +// processing would produce (only the order differs). +// +// Fixture: two channels, four videos with mixed published_at (one NULL). +// +// Per-channel (before): chanA=[v-old, v-mid], chanB=[v-new, v-null] +// → [v-old, v-mid, v-new, v-null] +// Newest-first (after): [v-new, v-mid, v-old, v-null] +func TestRunOnce_NewestFirstOrdering(t *testing.T) { + old := time.Date(2024, 1, 1, 0, 0, 0, 0, time.UTC) + mid := time.Date(2024, 6, 1, 0, 0, 0, 0, time.UTC) + newt := time.Date(2024, 12, 1, 0, 0, 0, 0, time.UTC) + // zero time = NULL published_at (schema 001: nullable) + + src := &fakeSource{ + subs: []domain.Subscription{ + sub("chanA", "Channel A"), + sub("chanB", "Channel B"), + }, + videos: map[string][]domain.Video{ + "chanA": { + vidAt("v-old", "Old Video", old), + vidAt("v-mid", "Mid Video", mid), + }, + "chanB": { + vidAt("v-new", "New Video", newt), + vidAt("v-null", "No Date Video", time.Time{}), // NULL + }, + }, + } + st := &fakeStore{seen: map[string]bool{}, auto: true} + sink := &recordingSink{} + eng := usecase.NewEngine(src, fakeSummarizer{}, sink) + r := runner.New(src, st, eng, testUser, quietLogger()) + + stats, err := r.RunOnce(context.Background()) + require.NoError(t, err) + + // Same set: all 4 candidates processed regardless of order. + require.Equal(t, 4, stats.Candidates) + require.Equal(t, 4, stats.Summarized, "same set of videos processed as per-channel order") + require.Len(t, sink.delivered, 4) + + // Build video-id → delivery-position map. + order := make(map[string]int, len(sink.delivered)) + for i, s := range sink.delivered { + order[s.VideoID] = i + t.Logf("position %d: %s", i, s.VideoID) + } + + require.Less(t, order["id-v-new"], order["id-v-mid"], "newest (Dec) before mid (Jun)") + require.Less(t, order["id-v-mid"], order["id-v-old"], "mid (Jun) before old (Jan)") + require.Less(t, order["id-v-old"], order["id-v-null"], "dated before NULL (NULLS LAST)") +} + +// TestRunOnce_NewestFirstNullsOnly asserts that when all candidates have NULL +// published_at, discovery order (stable) is preserved as the tiebreak. +func TestRunOnce_NewestFirstNullsOnly(t *testing.T) { + src := &fakeSource{ + subs: []domain.Subscription{ + sub("chanA", "Channel A"), + sub("chanB", "Channel B"), + }, + videos: map[string][]domain.Video{ + "chanA": {vidAt("v1", "V1", time.Time{}), vidAt("v2", "V2", time.Time{})}, + "chanB": {vidAt("v3", "V3", time.Time{})}, + }, + } + st := &fakeStore{seen: map[string]bool{}, auto: true} + sink := &recordingSink{} + eng := usecase.NewEngine(src, fakeSummarizer{}, sink) + r := runner.New(src, st, eng, testUser, quietLogger()) + + stats, err := r.RunOnce(context.Background()) + require.NoError(t, err) + require.Equal(t, 3, stats.Summarized, "all null-date videos processed") + + // Discovery order: chanA[v1, v2], chanB[v3] → [v1, v2, v3]. + // All have NULL published_at so the sort is stable; discovery order must hold. + require.Equal(t, "id-v1", sink.delivered[0].VideoID) + require.Equal(t, "id-v2", sink.delivered[1].VideoID) + require.Equal(t, "id-v3", sink.delivered[2].VideoID) +}