From 5219561a911df173c5ade0ed45ba732837933ca6 Mon Sep 17 00:00:00 2001 From: Mathias Date: Wed, 10 Jun 2026 20:27:25 +0200 Subject: [PATCH] feat(discovery): per-channel caption-availability memory (ADR-024) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit After transcript caching (ADR-021) and the Shorts filter (ADR-023), the remaining caption waste is the first fetch on every new video of a channel that never has English captions — each costs one rate-limited fetch to resolve to "none", and on a throttled IP churns the backoff machinery first. Remember, per (user, channel), a streak of consecutive no-caption outcomes (channel_caption_state, migration 016, RLS-scoped). Once it reaches TAPIR_CHANNEL_CAPTIONLESS_THRESHOLD (default 5) the channel is suppressed — videos discovered/listed but not caption-fetched — for TAPIR_CHANNEL_CAPTIONLESS_WINDOW (default 14d), then one is re-probed (auto-recovery). A successful fetch resets the streak; a 429 does not count; an explicit manual request bypasses suppression. threshold=0 disables. Co-Authored-By: Claude Opus 4.8 (1M context) --- DECISIONS.md | 32 +++++++ cmd/tapir/main.go | 3 +- cmd/tapir/scheduler.go | 23 ++--- docs/homelab-integration.md | 5 ++ internal/adapters/store/channel_caption.go | 83 +++++++++++++++++++ .../adapters/store/channel_caption_test.go | 67 +++++++++++++++ internal/adapters/store/migrate_test.go | 13 ++- .../016_channel_caption_state.down.sql | 1 + .../016_channel_caption_state.up.sql | 30 +++++++ internal/config/config.go | 25 ++++++ internal/runner/runner.go | 75 ++++++++++++++--- internal/runner/runner_test.go | 72 ++++++++++++++++ 12 files changed, 403 insertions(+), 26 deletions(-) create mode 100644 internal/adapters/store/channel_caption.go create mode 100644 internal/adapters/store/channel_caption_test.go create mode 100644 internal/adapters/store/migrations/016_channel_caption_state.down.sql create mode 100644 internal/adapters/store/migrations/016_channel_caption_state.up.sql diff --git a/DECISIONS.md b/DECISIONS.md index 58cff81..0b1a50d 100644 --- a/DECISIONS.md +++ b/DECISIONS.md @@ -891,6 +891,38 @@ this is well under the 10k/day cap; at larger scale, batch `videos.list` across --- +## ADR-024 — Per-channel caption-availability memory + +**Status:** Accepted (2026-06-10). **Builds on ADR-014** (per-IP caption budget), **ADR-021** +(shared transcript cache), **ADR-023** (Shorts filter). + +**Context.** After ADR-021 caches transcripts and ADR-023 drops Shorts, the remaining caption +waste is the *first* fetch on every new video of a channel that never publishes English captions +(foreign-language news, music, etc.). Each costs one rate-limited fetch to resolve to "none" — +and on a throttled IP that fetch may 429 and churn the backoff machinery before it ever gets a +verdict. A pilot user's feed had several such channels. + +**Decision.** Remember, per `(user, channel)`, a streak of consecutive no-caption outcomes +(`channel_caption_state`, migration 016, RLS-scoped like the rest of the user-owned schema). +Once the streak reaches `TAPIR_CHANNEL_CAPTIONLESS_THRESHOLD` (default 5) the channel is +suppressed — its videos are discovered/listed but not caption-fetched — for +`TAPIR_CHANNEL_CAPTIONLESS_WINDOW` (default 14d), after which one video is re-probed +(auto-recovery for a channel that starts adding captions). A successful fetch resets the streak; +a fresh 429 does NOT count (transient, not a caption verdict). An explicit manual request +bypasses suppression. `threshold = 0` disables the feature. + +**Why per-user, not global.** Caption availability is really a channel property (public), so a +global table would let users share the learning. But subscriptions are per-user (ADR-012) and at +pilot scale users' channel sets barely overlap, so per-user + RLS keeps it consistent with the +existing isolation model with no new non-RLS exception to justify. Promoting to a shared table +(like transcripts, ADR-021) is a future optimisation if channel overlap grows. + +**Reversibility.** Migration 016 is a clean drop; `threshold = 0` disables at runtime. The +memory only ever *suppresses fetches* — it never deletes content or affects already-stored +summaries. + +--- + ## Rejected alternatives Approaches considered during the 2026-06-02 planning + grill session and **deliberately not diff --git a/cmd/tapir/main.go b/cmd/tapir/main.go index 33ef6e8..2e01fc3 100644 --- a/cmd/tapir/main.go +++ b/cmd/tapir/main.go @@ -136,7 +136,8 @@ func cmdRun(ctx context.Context, log *slog.Logger) error { r := runner.New(engine.Source, st, engine, cfg.UserID, log, runner.WithBackoff(cfg.FetchBackoff), - runner.WithAutoWindow(cfg.AutoSummarizeWindow)) + runner.WithAutoWindow(cfg.AutoSummarizeWindow), + runner.WithCaptionMemory(cfg.ChannelCaptionlessThreshold, cfg.ChannelCaptionlessWindow)) log.Info("starting run", "user", cfg.UserID, "model", cfg.SummarizerModel, "gateway", cfg.GatewayURL, "poll_interval", cfg.PollInterval, "fetch_backoff", cfg.FetchBackoff, diff --git a/cmd/tapir/scheduler.go b/cmd/tapir/scheduler.go index d8e6de3..41af918 100644 --- a/cmd/tapir/scheduler.go +++ b/cmd/tapir/scheduler.go @@ -45,7 +45,8 @@ func buildUserRunner(cfg config.Config, st *store.Store, secretStore ports.Secre return runner.New(src, st, engine, userID, log, runner.WithBackoff(cfg.FetchBackoff), - runner.WithAutoWindow(cfg.AutoSummarizeWindow)), nil + runner.WithAutoWindow(cfg.AutoSummarizeWindow), + runner.WithCaptionMemory(cfg.ChannelCaptionlessThreshold, cfg.ChannelCaptionlessWindow)), nil } // userLister enumerates every registered user and reports a user's video @@ -121,6 +122,7 @@ func runDiscoveryPass( "skipped_seen", total.SkippedSeen, "skipped_no_text", total.SkippedNoText, "skipped_manual", total.SkippedManual, "skipped_too_old", total.SkippedTooOld, "skipped_rate_limited", total.SkippedRateLimited, + "skipped_no_caption_channel", total.SkippedNoCaptionChannel, "channel_unavailable", total.ChannelUnavailable, "errors", total.Errors) return total } @@ -181,14 +183,15 @@ func rotateUsers(users []store.UserIdentity, pass int) []store.UserIdentity { // per-tick aggregate across all users. func sumStats(a, b runner.Stats) runner.Stats { return runner.Stats{ - Candidates: a.Candidates + b.Candidates, - Summarized: a.Summarized + b.Summarized, - SkippedSeen: a.SkippedSeen + b.SkippedSeen, - SkippedNoText: a.SkippedNoText + b.SkippedNoText, - SkippedManual: a.SkippedManual + b.SkippedManual, - SkippedTooOld: a.SkippedTooOld + b.SkippedTooOld, - SkippedRateLimited: a.SkippedRateLimited + b.SkippedRateLimited, - ChannelUnavailable: a.ChannelUnavailable + b.ChannelUnavailable, - Errors: a.Errors + b.Errors, + Candidates: a.Candidates + b.Candidates, + Summarized: a.Summarized + b.Summarized, + SkippedSeen: a.SkippedSeen + b.SkippedSeen, + SkippedNoText: a.SkippedNoText + b.SkippedNoText, + SkippedManual: a.SkippedManual + b.SkippedManual, + SkippedTooOld: a.SkippedTooOld + b.SkippedTooOld, + SkippedRateLimited: a.SkippedRateLimited + b.SkippedRateLimited, + SkippedNoCaptionChannel: a.SkippedNoCaptionChannel + b.SkippedNoCaptionChannel, + ChannelUnavailable: a.ChannelUnavailable + b.ChannelUnavailable, + Errors: a.Errors + b.Errors, } } diff --git a/docs/homelab-integration.md b/docs/homelab-integration.md index 2fcbf3f..69c51ba 100644 --- a/docs/homelab-integration.md +++ b/docs/homelab-integration.md @@ -44,6 +44,11 @@ it** — endpoints and aliases drift, and this file is a snapshot (2026-06-06), the timedtext 429 path) and drops videos shorter than this plus any live/upcoming broadcast, so the scarce caption-fetch budget isn't spent on Shorts. `0` disables the filter. The paste-a-URL path is never filtered. +- **Per-channel caption memory (ADR-024).** `TAPIR_CHANNEL_CAPTIONLESS_THRESHOLD` — **default + `5`** consecutive no-caption results before a channel is suppressed (its videos listed but not + caption-fetched). `TAPIR_CHANNEL_CAPTIONLESS_WINDOW` — **default `336h`** (14d) suppression + before one video is re-probed. `THRESHOLD=0` disables. A successful fetch resets the channel; + a 429 does not count; an explicit manual request bypasses suppression. - **Thinking models need an explicit `max_tokens`.** qwen3 / deepseek-r1 spend the budget on reasoning and return **empty content** if `max_tokens` is too low (or unset). The summarizer's parser treats an empty summary as an error for exactly this reason. **Done (2026-06-02, Worker F):** diff --git a/internal/adapters/store/channel_caption.go b/internal/adapters/store/channel_caption.go new file mode 100644 index 0000000..1236e40 --- /dev/null +++ b/internal/adapters/store/channel_caption.go @@ -0,0 +1,83 @@ +package store + +import ( + "context" + "fmt" + "time" + + "github.com/jackc/pgx/v5" +) + +// CaptionlessChannels returns the set of channel ids currently suppressed for the +// user — channels whose recent videos all yielded no captions, within their +// suppression window (ADR-024). The runner skips caption fetches for these +// channels' videos. A channel whose window has expired is not returned, so its +// next video is re-probed (auto-recovery). +func (s *Store) CaptionlessChannels(ctx context.Context, userID string) (map[string]bool, error) { + out := map[string]bool{} + err := s.withUser(ctx, userID, func(tx pgx.Tx) error { + rows, err := tx.Query(ctx, ` + SELECT channel_id FROM channel_caption_state + WHERE user_id = $1 AND captionless_until IS NOT NULL AND captionless_until > now()`, + userID) + if err != nil { + return fmt.Errorf("store: caption-less channels: %w", err) + } + defer rows.Close() + for rows.Next() { + var ch string + if err := rows.Scan(&ch); err != nil { + return fmt.Errorf("store: scan caption-less channel: %w", err) + } + out[ch] = true + } + return rows.Err() + }) + return out, err +} + +// RecordChannelCaptionOutcome updates a channel's caption-availability memory +// after a fetch attempt (ADR-024). hadCaptions resets the channel (consecutive +// count to 0, suppression cleared). Otherwise the consecutive no-caption count is +// incremented; once it reaches threshold the channel is suppressed for window. +// threshold <= 0 is a no-op (feature disabled). An empty channelID is ignored +// (some sources may not carry one). +func (s *Store) RecordChannelCaptionOutcome(ctx context.Context, userID, channelID string, hadCaptions bool, threshold int, window time.Duration) error { + if channelID == "" || threshold <= 0 { + return nil + } + return s.withUser(ctx, userID, func(tx pgx.Tx) error { + if hadCaptions { + _, err := tx.Exec(ctx, ` + INSERT INTO channel_caption_state (user_id, channel_id, consecutive_none, captionless_until, updated_at) + VALUES ($1, $2, 0, NULL, now()) + ON CONFLICT (user_id, channel_id) + DO UPDATE SET consecutive_none = 0, captionless_until = NULL, updated_at = now()`, + userID, channelID) + if err != nil { + return fmt.Errorf("store: reset channel caption state: %w", err) + } + return nil + } + // No captions: increment the streak; suppress once it reaches threshold. + // captionless_until is set from the NEW count inside the same statement so + // the decision is atomic with the increment. + until := time.Now().Add(window) + _, err := tx.Exec(ctx, ` + INSERT INTO channel_caption_state (user_id, channel_id, consecutive_none, captionless_until, updated_at) + VALUES ($1, $2, 1, CASE WHEN 1 >= $3 THEN $4::timestamptz ELSE NULL END, now()) + ON CONFLICT (user_id, channel_id) + DO UPDATE SET + consecutive_none = channel_caption_state.consecutive_none + 1, + captionless_until = CASE + WHEN channel_caption_state.consecutive_none + 1 >= $3 THEN $4::timestamptz + ELSE channel_caption_state.captionless_until + END, + updated_at = now()`, + userID, channelID, threshold, until) + if err != nil { + return fmt.Errorf("store: record channel no-caption: %w", err) + } + return nil + }) +} diff --git a/internal/adapters/store/channel_caption_test.go b/internal/adapters/store/channel_caption_test.go new file mode 100644 index 0000000..8fe4576 --- /dev/null +++ b/internal/adapters/store/channel_caption_test.go @@ -0,0 +1,67 @@ +package store_test + +import ( + "context" + "testing" + "time" + + "github.com/stretchr/testify/require" +) + +func TestChannelCaptionMemory_SuppressesAfterThreshold(t *testing.T) { + ctx := context.Background() + s := newStore(t) + super := rawPool(t) + resetDB(t, super) + seedUser(t, super, userA) + + const threshold = 3 + window := time.Hour + + // Below threshold: not yet suppressed. + require.NoError(t, s.RecordChannelCaptionOutcome(ctx, userA, "chanX", false, threshold, window)) + require.NoError(t, s.RecordChannelCaptionOutcome(ctx, userA, "chanX", false, threshold, window)) + got, err := s.CaptionlessChannels(ctx, userA) + require.NoError(t, err) + require.NotContains(t, got, "chanX", "2 < threshold 3: not suppressed yet") + + // Crossing the threshold suppresses the channel. + require.NoError(t, s.RecordChannelCaptionOutcome(ctx, userA, "chanX", false, threshold, window)) + got, err = s.CaptionlessChannels(ctx, userA) + require.NoError(t, err) + require.Contains(t, got, "chanX", "3 consecutive no-caption results suppress the channel") + + // A successful caption fetch resets it. + require.NoError(t, s.RecordChannelCaptionOutcome(ctx, userA, "chanX", true, threshold, window)) + got, err = s.CaptionlessChannels(ctx, userA) + require.NoError(t, err) + require.NotContains(t, got, "chanX", "a captioned video clears suppression") +} + +func TestChannelCaptionMemory_WindowExpiryReProbes(t *testing.T) { + ctx := context.Background() + s := newStore(t) + super := rawPool(t) + resetDB(t, super) + seedUser(t, super, userA) + + // A negative window means captionless_until lands in the past — modelling an + // elapsed suppression window, which must make the channel eligible again. + require.NoError(t, s.RecordChannelCaptionOutcome(ctx, userA, "chanY", false, 1, -time.Hour)) + got, err := s.CaptionlessChannels(ctx, userA) + require.NoError(t, err) + require.NotContains(t, got, "chanY", "an expired window re-enables the channel for a re-probe") +} + +func TestChannelCaptionMemory_DisabledThresholdIsNoOp(t *testing.T) { + ctx := context.Background() + s := newStore(t) + super := rawPool(t) + resetDB(t, super) + seedUser(t, super, userA) + + require.NoError(t, s.RecordChannelCaptionOutcome(ctx, userA, "chanZ", false, 0, time.Hour)) + got, err := s.CaptionlessChannels(ctx, userA) + require.NoError(t, err) + require.Empty(t, got, "threshold 0 disables the memory — nothing recorded") +} diff --git a/internal/adapters/store/migrate_test.go b/internal/adapters/store/migrate_test.go index e75d023..de48273 100644 --- a/internal/adapters/store/migrate_test.go +++ b/internal/adapters/store/migrate_test.go @@ -53,7 +53,9 @@ func TestMigration010LoginEventsUpDown(t *testing.T) { require.True(t, loginEventsExists(t), "login_events must exist at latest migration") m := fileMigrator(t) - // 011..015 sit above 010; step them down first so 010 is exercised in isolation. + // 011..016 sit above 010; step them down first so 010 is exercised in isolation. + require.NoError(t, m.Steps(-1), "down 016 drops channel_caption_state, login_events intact") + require.True(t, loginEventsExists(t), "016 down leaves login_events intact") require.NoError(t, m.Steps(-1), "down 015 reshapes transcripts, login_events intact") require.True(t, loginEventsExists(t), "015 down leaves login_events intact") require.NoError(t, m.Steps(-1), "down 014 drops channel_title, login_events intact") @@ -68,7 +70,7 @@ func TestMigration010LoginEventsUpDown(t *testing.T) { require.NoError(t, m.Steps(-1), "down 010 must drop login_events") require.False(t, loginEventsExists(t), "login_events must be gone after the down migration") - require.NoError(t, m.Steps(6), "up must recreate 010 then re-apply 011..015") + require.NoError(t, m.Steps(7), "up must recreate 010 then re-apply 011..016") require.True(t, loginEventsExists(t), "login_events must be restored after the up migration") } @@ -91,6 +93,7 @@ func TestMigration011AutoSummarizeDefaultUpDown(t *testing.T) { require.Equal(t, "true", autoSummarizeDefault(t), "011 sets the default to TRUE") m := fileMigrator(t) + require.NoError(t, m.Steps(-1), "down 016 drops channel_caption_state") require.NoError(t, m.Steps(-1), "down 015 reshapes transcripts") require.NoError(t, m.Steps(-1), "down 014 drops channel_title") require.NoError(t, m.Steps(-1), "down 013 drops channel_errors") @@ -104,6 +107,7 @@ func TestMigration011AutoSummarizeDefaultUpDown(t *testing.T) { require.NoError(t, m.Steps(1), "up 013 creates channel_errors") require.NoError(t, m.Steps(1), "up 014 recreates channel_title") require.NoError(t, m.Steps(1), "up 015 reshapes transcripts to shared") + require.NoError(t, m.Steps(1), "up 016 recreates channel_caption_state (HEAD)") } // channelTitleExists reports whether videos.channel_title is present. @@ -123,6 +127,8 @@ func TestMigration014VideoChannelTitleUpDown(t *testing.T) { require.True(t, channelTitleExists(t), "channel_title exists at latest migration") m := fileMigrator(t) + require.NoError(t, m.Steps(-1), "down 016 drops channel_caption_state, channel_title intact") + require.True(t, channelTitleExists(t), "016 down leaves channel_title intact") require.NoError(t, m.Steps(-1), "down 015 reshapes transcripts, channel_title intact") require.True(t, channelTitleExists(t), "015 down leaves channel_title intact") require.NoError(t, m.Steps(-1), "down 014 must drop channel_title") @@ -130,7 +136,8 @@ func TestMigration014VideoChannelTitleUpDown(t *testing.T) { require.NoError(t, m.Steps(1), "up 014 must recreate channel_title") require.True(t, channelTitleExists(t), "channel_title must be restored after the up migration") - require.NoError(t, m.Steps(1), "up 015 restores the shared transcripts shape (HEAD)") + require.NoError(t, m.Steps(1), "up 015 restores the shared transcripts shape") + require.NoError(t, m.Steps(1), "up 016 recreates channel_caption_state (HEAD)") } // TestMigration012FixAutoSummarizeRLS proves 012 runs cleanly and flips any diff --git a/internal/adapters/store/migrations/016_channel_caption_state.down.sql b/internal/adapters/store/migrations/016_channel_caption_state.down.sql new file mode 100644 index 0000000..804b2a6 --- /dev/null +++ b/internal/adapters/store/migrations/016_channel_caption_state.down.sql @@ -0,0 +1 @@ +DROP TABLE channel_caption_state; diff --git a/internal/adapters/store/migrations/016_channel_caption_state.up.sql b/internal/adapters/store/migrations/016_channel_caption_state.up.sql new file mode 100644 index 0000000..b39be83 --- /dev/null +++ b/internal/adapters/store/migrations/016_channel_caption_state.up.sql @@ -0,0 +1,30 @@ +-- Migration 016: per-(user, channel) caption-availability memory (ADR-024). +-- +-- Some channels never publish English captions (foreign-language news, music, +-- etc.). Each of their new videos still costs ONE rate-limited caption fetch +-- (ADR-014) before resolving to "none" — and on a throttled egress IP that fetch +-- may 429 and churn through the backoff machinery first. This table remembers +-- channels that repeatedly yield no captions so discovery can stop attempting +-- their videos, freeing the scarce fetch budget for channels that do have them. +-- +-- consecutive_none counts no-caption outcomes in a row; a successful fetch resets +-- it to 0. Once it crosses the threshold the channel is suppressed until +-- captionless_until, after which one video is re-probed (auto-recovery for a +-- channel that starts adding captions). Per-user + RLS-scoped, consistent with +-- the rest of the user-owned schema (subscriptions are per-user; ADR-012). +CREATE TABLE channel_caption_state ( + user_id UUID NOT NULL REFERENCES users(id) ON DELETE CASCADE, + channel_id TEXT NOT NULL, + consecutive_none INT NOT NULL DEFAULT 0, + captionless_until TIMESTAMPTZ, + updated_at TIMESTAMPTZ NOT NULL DEFAULT now(), + PRIMARY KEY (user_id, channel_id) +); + +CREATE INDEX idx_channel_caption_state_user_id ON channel_caption_state(user_id); + +ALTER TABLE channel_caption_state ENABLE ROW LEVEL SECURITY; +ALTER TABLE channel_caption_state FORCE ROW LEVEL SECURITY; +CREATE POLICY channel_caption_state_isolation ON channel_caption_state + FOR ALL + USING (user_id = current_setting('tapir.current_user_id', true)::uuid); diff --git a/internal/config/config.go b/internal/config/config.go index 4e3334b..79e2dfb 100644 --- a/internal/config/config.go +++ b/internal/config/config.go @@ -52,6 +52,14 @@ type Config struct { // budget, ADR-014/ADR-023). Enforced via a cheap Data API videos.list lookup at // discovery, never the rate-limited caption path. 0 disables the filter. MinVideoSeconds int + + // ChannelCaptionlessThreshold is how many consecutive no-caption results a + // channel may yield before its videos are suppressed from caption fetching + // (ADR-024). 0 disables the per-channel caption memory entirely. + ChannelCaptionlessThreshold int + // ChannelCaptionlessWindow is how long a suppressed channel stays suppressed + // before one video is re-probed (auto-recovery for a channel that adds captions). + ChannelCaptionlessWindow time.Duration // SummarizerTimeout bounds a single completion call. Thinking models are // slow, so the default is generous. SummarizerTimeout time.Duration @@ -145,6 +153,8 @@ const ( defaultSummaryMaxTokens = 1500 defaultMaxTranscriptChars = 18000 defaultMinVideoSeconds = 60 + defaultCaptionlessThreshold = 5 + defaultCaptionlessWindow = 14 * 24 * time.Hour defaultSummarizerTimeout = 5 * time.Minute defaultYTTokenRef = "youtube/refresh_token" defaultYTConnectRedirectURL = "https://tapir.d-ma.be/oauth/youtube/callback" @@ -246,6 +256,21 @@ func Load() (Config, error) { } c.MinVideoSeconds = minVideo + captionThreshold, err := intOr("TAPIR_CHANNEL_CAPTIONLESS_THRESHOLD", defaultCaptionlessThreshold) + if err != nil { + return Config{}, err + } + if captionThreshold < 0 { + captionThreshold = 0 + } + c.ChannelCaptionlessThreshold = captionThreshold + + captionWindow, err := durationOr("TAPIR_CHANNEL_CAPTIONLESS_WINDOW", defaultCaptionlessWindow) + if err != nil { + return Config{}, err + } + c.ChannelCaptionlessWindow = captionWindow + onboard, err := intOr("TAPIR_ONBOARD_SUMMARIZE_COUNT", defaultOnboardSummarizeCount) if err != nil { return Config{}, err diff --git a/internal/runner/runner.go b/internal/runner/runner.go index 1059d47..2c4e673 100644 --- a/internal/runner/runner.go +++ b/internal/runner/runner.go @@ -27,8 +27,9 @@ import ( // 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 + v domain.Video + channelID string // owning channel — keys the caption-availability memory (ADR-024) + pos int // discovery position — used as a stable tiebreak when published_at ties } // compareNewestFirst orders candidates by published_at descending, NULLS LAST, @@ -74,6 +75,14 @@ type VideoStore interface { // Called when NewVideos returns domain.ErrChannelUnavailable; best-effort, errors // are logged and never abort the pass. UpsertChannelError(ctx context.Context, userID, channelID, channelTitle string) error + // CaptionlessChannels returns channel ids currently suppressed because their + // recent videos all yielded no captions (ADR-024). The loop skips caption + // fetches for these channels' (non-requested) videos. + CaptionlessChannels(ctx context.Context, userID string) (map[string]bool, error) + // RecordChannelCaptionOutcome updates a channel's caption memory after a fetch: + // hadCaptions resets it, otherwise the no-caption streak grows and the channel + // is suppressed for window once it reaches threshold. A no-op when threshold<=0. + RecordChannelCaptionOutcome(ctx context.Context, userID, channelID string, hadCaptions bool, threshold int, window time.Duration) error } // Processor runs the core use case for a single video. *usecase.Engine @@ -93,6 +102,9 @@ type Runner struct { backoff time.Duration // rate-limit retry window; 0 = always retry autoWindow time.Duration // recency bound for auto-summarize; 0 = no bound now func() time.Time // injectable clock (tests); defaults to time.Now + + captionThreshold int // consecutive no-caption results before a channel is suppressed; 0 = feature off + captionWindow time.Duration // how long a caption-less channel stays suppressed before re-probe } // Option configures a Runner at construction. Variadic so existing call sites @@ -114,6 +126,14 @@ func WithClock(now func() time.Time) Option { return func(r *Runner) { r.now = n // bypasses the bound. 0 (the default) disables it (summarize every unseen video). func WithAutoWindow(d time.Duration) Option { return func(r *Runner) { r.autoWindow = d } } +// WithCaptionMemory enables per-channel caption-availability suppression +// (ADR-024): after threshold consecutive no-caption results a channel's videos +// are skipped (no caption fetch) for window, then one is re-probed. threshold<=0 +// (the default) disables the feature entirely. +func WithCaptionMemory(threshold int, window time.Duration) Option { + return func(r *Runner) { r.captionThreshold = threshold; r.captionWindow = window } +} + // New builds a Runner. A nil logger falls back to slog.Default. func New(src ports.VideoSource, store VideoStore, engine Processor, userID string, log *slog.Logger, opts ...Option) *Runner { if log == nil { @@ -131,15 +151,16 @@ func New(src ports.VideoSource, store VideoStore, engine Processor, userID strin // Stats summarizes one RunOnce pass. type Stats struct { - Candidates int - Summarized int - SkippedSeen int - SkippedNoText int - SkippedManual int // discovered but not queued, in manual mode - SkippedTooOld int // auto mode: published outside the recency window (not requested) - SkippedRateLimited int // 429'd previously and still inside the backoff window - Errors int - ChannelUnavailable int // channels that returned HTTP 404 (deleted/private) + Candidates int + Summarized int + SkippedSeen int + SkippedNoText int + SkippedManual int // discovered but not queued, in manual mode + SkippedTooOld int // auto mode: published outside the recency window (not requested) + SkippedRateLimited int // 429'd previously and still inside the backoff window + SkippedNoCaptionChannel int // channel suppressed as caption-less (ADR-024) + Errors int + ChannelUnavailable int // channels that returned HTTP 404 (deleted/private) } // tooOld reports whether a video published at publishedAt falls outside the @@ -213,6 +234,17 @@ func (r *Runner) RunOnce(ctx context.Context) (Stats, error) { } } + // Per-channel caption memory (ADR-024): channels whose recent videos all + // yielded no captions are suppressed so their new videos don't burn the scarce + // fetch budget. Loaded only when the feature is enabled (threshold > 0). + var captionless map[string]bool + if r.captionThreshold > 0 { + captionless, err = r.store.CaptionlessChannels(ctx, r.userID) + if err != nil { + return stats, fmt.Errorf("runner: load caption-less channels: %w", err) + } + } + subs, err := r.src.ListSubscriptions(ctx, r.userID) if err != nil { return stats, fmt.Errorf("runner: list subscriptions: %w", err) @@ -273,6 +305,14 @@ func (r *Runner) RunOnce(ctx context.Context) (Stats, error) { continue } + // Caption-less channel (ADR-024): its recent videos all returned no + // captions, so skip the fetch entirely. The video is still listed + // (UpsertVideo above); an explicit manual request bypasses the skip. + if !requested[id] && captionless[sub.ChannelID] { + stats.SkippedNoCaptionChannel++ + continue + } + // Still inside the rate-limit backoff window: skip without fetching. if at, ok := rateLimited[id]; ok && r.now().Sub(at) < r.backoff { stats.SkippedRateLimited++ @@ -280,7 +320,7 @@ func (r *Runner) RunOnce(ctx context.Context) (Stats, error) { continue } - candidates = append(candidates, passCandidate{v: v, pos: pos}) + candidates = append(candidates, passCandidate{v: v, channelID: sub.ChannelID, pos: pos}) pos++ } } @@ -320,6 +360,11 @@ func (r *Runner) RunOnce(ctx context.Context) (Stats, error) { errs = append(errs, fmt.Errorf("set none status %q: %w", c.v.ProviderVideoID, err)) stats.Errors++ } + // No captions: grow this channel's no-caption streak (ADR-024). + if err := r.store.RecordChannelCaptionOutcome(ctx, r.userID, c.channelID, false, r.captionThreshold, r.captionWindow); err != nil { + errs = append(errs, fmt.Errorf("record no-caption %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++ @@ -327,6 +372,11 @@ func (r *Runner) RunOnce(ctx context.Context) (Stats, error) { errs = append(errs, fmt.Errorf("set fetched status %q: %w", c.v.ProviderVideoID, err)) stats.Errors++ } + // Captions present: reset this channel's caption memory (ADR-024). + if err := r.store.RecordChannelCaptionOutcome(ctx, r.userID, c.channelID, true, r.captionThreshold, r.captionWindow); err != nil { + errs = append(errs, fmt.Errorf("record has-caption %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 { @@ -354,6 +404,7 @@ func (r *Runner) Loop(ctx context.Context, interval time.Duration) error { "skipped_seen", stats.SkippedSeen, "skipped_no_text", stats.SkippedNoText, "skipped_manual", stats.SkippedManual, "skipped_too_old", stats.SkippedTooOld, "skipped_rate_limited", stats.SkippedRateLimited, + "skipped_no_caption_channel", stats.SkippedNoCaptionChannel, "channel_unavailable", stats.ChannelUnavailable, "errors", stats.Errors) if err != nil { r.log.Warn("run pass had errors", "err", err) diff --git a/internal/runner/runner_test.go b/internal/runner/runner_test.go index 9d8a278..d50fc6e 100644 --- a/internal/runner/runner_test.go +++ b/internal/runner/runner_test.go @@ -51,6 +51,13 @@ type fakeStore struct { cleared []string rateLimited map[string]time.Time // id -> when 429'd (seeds the backoff window) statuses map[string]string // id -> last SetTranscriptStatus value + captionless map[string]bool // channel ids currently suppressed (ADR-024) + captionRecs []captionRec // RecordChannelCaptionOutcome calls, in order +} + +type captionRec struct { + channelID string + had bool } func (f *fakeStore) UpsertVideo(_ context.Context, v domain.Video) (string, error) { @@ -93,6 +100,22 @@ func (f *fakeStore) RateLimitedVideoIDs(_ context.Context, _ string) (map[string func (f *fakeStore) UpsertChannelError(_ context.Context, _, _, _ string) error { return nil } +func (f *fakeStore) CaptionlessChannels(_ context.Context, _ string) (map[string]bool, error) { + cp := make(map[string]bool, len(f.captionless)) + for k, v := range f.captionless { + cp[k] = v + } + return cp, nil +} + +func (f *fakeStore) RecordChannelCaptionOutcome(_ context.Context, _, channelID string, hadCaptions bool, threshold int, _ time.Duration) error { + if threshold <= 0 { + return nil + } + f.captionRecs = append(f.captionRecs, captionRec{channelID: channelID, had: hadCaptions}) + return nil +} + func (f *fakeStore) SetTranscriptStatus(_ context.Context, _, videoID, status string) error { if f.statuses == nil { f.statuses = map[string]string{} @@ -284,6 +307,55 @@ func TestRunOnce_AutoMode_OldVideoRequestedBypassesWindow(t *testing.T) { require.Len(t, sink.delivered, 1) } +// TestRunOnce_CaptionlessChannelSkipped: a channel flagged caption-less (ADR-024) +// has its videos skipped from fetching but still discovered/listed, while a +// normal channel's video is summarized. +func TestRunOnce_CaptionlessChannelSkipped(t *testing.T) { + src := &fakeSource{ + subs: []domain.Subscription{sub("dead", "Dead Channel"), sub("live", "Live Channel")}, + videos: map[string][]domain.Video{ + "dead": {vid("d1", "Dead One")}, + "live": {vid("l1", "Live One")}, + }, + } + st := &fakeStore{seen: map[string]bool{}, auto: true, captionless: map[string]bool{"dead": true}} + sink := &recordingSink{} + eng := usecase.NewEngine(src, fakeSummarizer{}, sink) + r := runner.New(src, st, eng, testUser, quietLogger(), + runner.WithCaptionMemory(5, 14*24*time.Hour)) + + stats, err := r.RunOnce(context.Background()) + require.NoError(t, err) + require.Equal(t, 1, stats.SkippedNoCaptionChannel, "dead channel's video skipped from fetch") + require.Equal(t, 1, stats.Summarized, "live channel's video still summarized") + require.Len(t, st.upserted, 2, "both videos are still discovered and listed") +} + +// TestRunOnce_RecordsCaptionOutcomes: a no-caption result grows the channel's +// streak (had=false); a successful summary resets it (had=true). +func TestRunOnce_RecordsCaptionOutcomes(t *testing.T) { + src := &fakeSource{ + subs: []domain.Subscription{sub("c1", "Has Caps"), sub("c2", "No Caps")}, + videos: map[string][]domain.Video{ + "c1": {vid("good", "Good")}, + "c2": {vid("bad", "Bad")}, + }, + transcripts: map[string]domain.Transcript{ + "bad": {Source: domain.SourceNone}, // no usable text → engine skips + }, + } + st := &fakeStore{seen: map[string]bool{}, auto: true} + sink := &recordingSink{} + eng := usecase.NewEngine(src, fakeSummarizer{}, sink) + r := runner.New(src, st, eng, testUser, quietLogger(), + runner.WithCaptionMemory(5, 14*24*time.Hour)) + + _, err := r.RunOnce(context.Background()) + require.NoError(t, err) + require.Contains(t, st.captionRecs, captionRec{channelID: "c1", had: true}, "captioned channel reset") + require.Contains(t, st.captionRecs, captionRec{channelID: "c2", had: false}, "no-caption channel streak grown") +} + // TestRunOnce_AutoWindowZero_SummarizesOld: a zero window disables the bound — // the pre-recency behaviour (summarize every unseen video) is preserved. func TestRunOnce_AutoWindowZero_SummarizesOld(t *testing.T) {