Compare commits

...
3 Commits
Author SHA1 Message Date
mathiasandClaude Opus 4.8 5219561a91 feat(discovery): per-channel caption-availability memory (ADR-024)
CI / Lint / Test / Vet (push) Successful in 30s
CI / Build & Import (push) Successful in 11s
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) <noreply@anthropic.com>
2026-06-10 20:27:25 +02:00
mathiasandClaude Opus 4.8 9db06d8a63 feat(discovery): drop Shorts and livestreams before the caption fetch (ADR-023)
CI / Lint / Test / Vet (push) Successful in 12s
CI / Build & Import (push) Successful in 10s
The scarce resource is the per-IP timedtext caption fetch (ADR-014); the pilot's
candidate set was mostly Shorts/clips/livestreams, each burning a fetch (a "none"
result is a completed fetch — it costs budget even when it yields nothing).

NewVideos now enriches candidates with one cheap Data API videos.list call
(contentDetails.duration + snippet.liveBroadcastContent — the quota API, a
DIFFERENT limit from the timedtext 429) and drops, before returning: videos
shorter than TAPIR_MIN_VIDEO_SECONDS (default 60) and any live/upcoming
broadcast. Dropped videos are never persisted, so the list declutters too.

Degrade-open: MinVideoSeconds=0 disables it (no quota call); a videos.list error
returns candidates unfiltered so discovery never breaks on a metadata hiccup. The
paste-a-URL path (VideoByID) is not filtered — an explicit request is honoured.

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
2026-06-10 19:34:45 +02:00
mathiasandClaude Opus 4.8 1aa8a97f95 fix(scheduler): rotate over connected users only for true fair share
CI / Lint / Test / Vet (push) Successful in 12s
CI / Build & Import (push) Successful in 10s
The lead-user rotation rotated the full ListAllUsers set, so a connectionless
orphan identity ate a rotation slot — collapsing onto the next real user and
skewing the lead share (two real users got 2/3 vs 1/3 instead of 50/50). Filter
to connected users BEFORE rotating so the rotation is over exactly the users that
consume the caption budget. A dead identity can no longer skew fairness.

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
2026-06-10 18:28:54 +02:00
16 changed files with 692 additions and 43 deletions
+68
View File
@@ -855,6 +855,74 @@ reply still parses; a transcript within budget is unchanged).
--- ---
## ADR-023 — Drop Shorts/livestreams at discovery to protect the caption budget
**Status:** Accepted (2026-06-10). **Builds on ADR-014** (per-IP caption rate limit is the
binding constraint) and the ADR-022 live-run findings.
**Context.** The scarce resource is the unofficial timedtext caption fetch (per-egress-IP
429, ~3 successful/pass). The first multi-user run showed the candidate set was mostly noise —
Shorts, sub-minute clips, and live broadcasts — each of which still consumes a caption-fetch
attempt (and a "none" result is a *completed* fetch, so it costs budget even when it yields
nothing). Spending the rate-limited budget on content the user will not read is the waste to
cut first; it is cheaper and lower-risk than raising the ceiling (multi-IP, Whisper).
Duration and live status are NOT in the playlistItems discovery response, but they ARE in the
Data API `videos.list` (contentDetails.duration + snippet.liveBroadcastContent) — the official
**quota-based** API (1 unit/call, 50 ids/call), which is a *different* limit from the timedtext
429. So one cheap quota call buys a filter that saves many expensive throttled fetches.
**Decision.**
1. `NewVideos` enriches its candidates with a single `videos.list` call and drops, before
returning: videos shorter than `TAPIR_MIN_VIDEO_SECONDS` (default 60) and any `live`/
`upcoming` broadcast. Dropped videos are never persisted, so they also declutter the list.
2. The filter is **degrade-open**: `MinVideoSeconds=0` disables it (no quota call), and a
`videos.list` error returns the candidates unfiltered — discovery must never break because a
metadata call hiccuped (worst case = pre-ADR-023 behaviour).
3. The paste-a-URL path (`VideoByID`) is **not** filtered — an explicit user request for a
specific video (even a Short) is honoured.
**Reversibility.** Pure discovery-time filter + config. `TAPIR_MIN_VIDEO_SECONDS=0` restores
the old behaviour. No schema change, no effect on already-stored videos.
**Quota note.** Per-channel enrichment adds ~1 unit/channel/pass. At pilot scale (≤3 users)
this is well under the 10k/day cap; at larger scale, batch `videos.list` across channels
(50 ids/call) by collecting all discovered ids per pass before enriching.
---
## 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 ## Rejected alternatives
Approaches considered during the 2026-06-02 planning + grill session and **deliberately not Approaches considered during the 2026-06-02 planning + grill session and **deliberately not
+2 -1
View File
@@ -136,7 +136,8 @@ func cmdRun(ctx context.Context, log *slog.Logger) error {
r := runner.New(engine.Source, st, engine, cfg.UserID, log, r := runner.New(engine.Source, st, engine, cfg.UserID, log,
runner.WithBackoff(cfg.FetchBackoff), 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, log.Info("starting run", "user", cfg.UserID, "model", cfg.SummarizerModel,
"gateway", cfg.GatewayURL, "poll_interval", cfg.PollInterval, "fetch_backoff", cfg.FetchBackoff, "gateway", cfg.GatewayURL, "poll_interval", cfg.PollInterval, "fetch_backoff", cfg.FetchBackoff,
+1
View File
@@ -88,6 +88,7 @@ func buildProcessor(cfg config.Config, st *store.Store) (*usecase.Engine, error)
ClientSecret: cfg.YTClientSecret, ClientSecret: cfg.YTClientSecret,
TokenSecretRef: cfg.YTTokenRef, TokenSecretRef: cfg.YTTokenRef,
PreferredLanguages: []string{"en"}, PreferredLanguages: []string{"en"},
MinVideoSeconds: cfg.MinVideoSeconds,
}, secretStore) }, secretStore)
sum := buildSummarizer(cfg) sum := buildSummarizer(cfg)
+29 -14
View File
@@ -33,6 +33,7 @@ func buildUserRunner(cfg config.Config, st *store.Store, secretStore ports.Secre
ClientSecret: cfg.YTClientSecret, ClientSecret: cfg.YTClientSecret,
TokenSecretRef: web.YouTubeTokenRef(userID), TokenSecretRef: web.YouTubeTokenRef(userID),
PreferredLanguages: []string{"en"}, PreferredLanguages: []string{"en"},
MinVideoSeconds: cfg.MinVideoSeconds,
}, secretStore) }, secretStore)
engine := usecase.NewEngine(src, buildSummarizer(cfg), st) engine := usecase.NewEngine(src, buildSummarizer(cfg), st)
@@ -44,7 +45,8 @@ func buildUserRunner(cfg config.Config, st *store.Store, secretStore ports.Secre
return runner.New(src, st, engine, userID, log, return runner.New(src, st, engine, userID, log,
runner.WithBackoff(cfg.FetchBackoff), 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 // userLister enumerates every registered user and reports a user's video
@@ -73,22 +75,17 @@ func runDiscoveryPass(
return runner.Stats{} return runner.Stats{}
} }
// Rotate who goes first each pass. Caption fetches share one per-egress-IP // Keep only users with a video connection. A pass for a connectionless user
// rate budget (ADR-014); whoever runs first each pass spends the pre-throttle // (e.g. a stale Dex-era orphan identity) only tries to resolve a token that
// window, so a FIXED user order permanently starves whoever is last (a new // was never minted, logging a spurious "ref not found" every tick. Filtering
// pilot user got 0 fetches for 12h while the first-listed user got all of // here — BEFORE rotation — also keeps fairness honest: rotation is over the
// them). Rotation gives every user the lead slot in turn. // users that actually consume the caption budget, so a dead identity can't eat
users = rotateUsers(users, pass) // a rotation slot and skew the lead share.
var connected []store.UserIdentity
log.Info("scheduler: starting discovery pass", "users", len(users))
var total runner.Stats
for _, u := range users { for _, u := range users {
if ctx.Err() != nil { if ctx.Err() != nil {
break // shutting down: stop enumerating return runner.Stats{} // shutting down
} }
// Skip users with no video connection. A discovery pass for them only
// attempts to resolve a token that was never minted, logging a spurious
// "ref not found" every tick (e.g. stale Dex-era orphan identities).
conns, err := lister.ConnectionsForUser(ctx, u.UserID) conns, err := lister.ConnectionsForUser(ctx, u.UserID)
if err != nil { if err != nil {
log.Warn("scheduler: list connections failed", "user", u.UserID, "err", err) log.Warn("scheduler: list connections failed", "user", u.UserID, "err", err)
@@ -98,6 +95,22 @@ func runDiscoveryPass(
log.Debug("scheduler: skipping user with no video connections", "user", u.UserID) log.Debug("scheduler: skipping user with no video connections", "user", u.UserID)
continue continue
} }
connected = append(connected, u)
}
// Rotate who goes first each pass. Caption fetches share one per-egress-IP
// rate budget (ADR-014); whoever runs first each pass spends the pre-throttle
// window, so a FIXED order permanently starves whoever is last (a new pilot
// user got 0 fetches for 12h while the first-listed user got all of them).
// Rotation over the connected set gives each real user the lead in turn.
connected = rotateUsers(connected, pass)
log.Info("scheduler: starting discovery pass", "users", len(connected))
var total runner.Stats
for _, u := range connected {
if ctx.Err() != nil {
break // shutting down: stop enumerating
}
stats, err := runUser(ctx, u.UserID) stats, err := runUser(ctx, u.UserID)
total = sumStats(total, stats) total = sumStats(total, stats)
if err != nil { if err != nil {
@@ -109,6 +122,7 @@ func runDiscoveryPass(
"skipped_seen", total.SkippedSeen, "skipped_no_text", total.SkippedNoText, "skipped_seen", total.SkippedSeen, "skipped_no_text", total.SkippedNoText,
"skipped_manual", total.SkippedManual, "skipped_too_old", total.SkippedTooOld, "skipped_manual", total.SkippedManual, "skipped_too_old", total.SkippedTooOld,
"skipped_rate_limited", total.SkippedRateLimited, "skipped_rate_limited", total.SkippedRateLimited,
"skipped_no_caption_channel", total.SkippedNoCaptionChannel,
"channel_unavailable", total.ChannelUnavailable, "errors", total.Errors) "channel_unavailable", total.ChannelUnavailable, "errors", total.Errors)
return total return total
} }
@@ -176,6 +190,7 @@ func sumStats(a, b runner.Stats) runner.Stats {
SkippedManual: a.SkippedManual + b.SkippedManual, SkippedManual: a.SkippedManual + b.SkippedManual,
SkippedTooOld: a.SkippedTooOld + b.SkippedTooOld, SkippedTooOld: a.SkippedTooOld + b.SkippedTooOld,
SkippedRateLimited: a.SkippedRateLimited + b.SkippedRateLimited, SkippedRateLimited: a.SkippedRateLimited + b.SkippedRateLimited,
SkippedNoCaptionChannel: a.SkippedNoCaptionChannel + b.SkippedNoCaptionChannel,
ChannelUnavailable: a.ChannelUnavailable + b.ChannelUnavailable, ChannelUnavailable: a.ChannelUnavailable + b.ChannelUnavailable,
Errors: a.Errors + b.Errors, Errors: a.Errors + b.Errors,
} }
+15
View File
@@ -128,6 +128,21 @@ func TestDiscoveryPassRotatesLeadUser(t *testing.T) {
require.Equal(t, 3, rc.count("c")) require.Equal(t, 3, rc.count("c"))
} }
// A connectionless orphan must not consume a rotation slot: rotation is over the
// connected users only, so two real users alternate the lead 50/50 even with a
// dead identity listed between them.
func TestDiscoveryPassRotationIgnoresConnectionlessUsers(t *testing.T) {
lister := fakeLister{users: usersN("a", "orphan", "c"), noConn: map[string]bool{"orphan": true}}
rc := newCountingRunUser()
runDiscoveryPass(context.Background(), 0, lister, rc.run, quietLog())
runDiscoveryPass(context.Background(), 1, lister, rc.run, quietLog())
require.Equal(t, []string{"a", "c", "c", "a"}, rc.runOrder(),
"only connected users rotate; the orphan never runs and never holds a slot")
require.Equal(t, 0, rc.count("orphan"))
}
func TestDiscoveryPassSkipsUsersWithoutConnections(t *testing.T) { func TestDiscoveryPassSkipsUsersWithoutConnections(t *testing.T) {
// b never connected a video source (e.g. a stale Dex-era orphan identity). // b never connected a video source (e.g. a stale Dex-era orphan identity).
// It must be skipped silently — not run and logged as a token error every pass. // It must be skipped silently — not run and logged as a token error every pass.
+10
View File
@@ -39,6 +39,16 @@ it** — endpoints and aliases drift, and this file is a snapshot (2026-06-06),
- `TAPIR_MAX_TRANSCRIPT_CHARS` — transcript truncation budget sent to the model. **Default - `TAPIR_MAX_TRANSCRIPT_CHARS` — transcript truncation budget sent to the model. **Default
`18000`** (~fits an 8k-context model). `0` disables truncation. Prevents the context-overflow `18000`** (~fits an 8k-context model). `0` disables truncation. Prevents the context-overflow
HTTP 400 a long transcript caused on `phi4-mini`. HTTP 400 a long transcript caused on `phi4-mini`.
- **Discovery low-value filter (ADR-023).** `TAPIR_MIN_VIDEO_SECONDS`**default `60`**. At
discovery, `NewVideos` enriches candidates with one cheap `videos.list` call (quota API, NOT
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 - **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 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):** parser treats an empty summary as an error for exactly this reason. **Done (2026-06-02, Worker F):**
@@ -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
})
}
@@ -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")
}
+10 -3
View File
@@ -53,7 +53,9 @@ func TestMigration010LoginEventsUpDown(t *testing.T) {
require.True(t, loginEventsExists(t), "login_events must exist at latest migration") require.True(t, loginEventsExists(t), "login_events must exist at latest migration")
m := fileMigrator(t) 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.NoError(t, m.Steps(-1), "down 015 reshapes transcripts, login_events intact")
require.True(t, loginEventsExists(t), "015 down leaves 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") 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.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.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") 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") require.Equal(t, "true", autoSummarizeDefault(t), "011 sets the default to TRUE")
m := fileMigrator(t) 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 015 reshapes transcripts")
require.NoError(t, m.Steps(-1), "down 014 drops channel_title") require.NoError(t, m.Steps(-1), "down 014 drops channel_title")
require.NoError(t, m.Steps(-1), "down 013 drops channel_errors") 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 013 creates channel_errors")
require.NoError(t, m.Steps(1), "up 014 recreates channel_title") 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 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. // 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") require.True(t, channelTitleExists(t), "channel_title exists at latest migration")
m := fileMigrator(t) 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.NoError(t, m.Steps(-1), "down 015 reshapes transcripts, channel_title intact")
require.True(t, channelTitleExists(t), "015 down leaves 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") 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.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.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 // TestMigration012FixAutoSummarizeRLS proves 012 runs cleanly and flips any
@@ -0,0 +1 @@
DROP TABLE channel_caption_state;
@@ -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);
+103 -1
View File
@@ -66,6 +66,13 @@ type Config struct {
// poll. Zero means defaultMaxVideos. // poll. Zero means defaultMaxVideos.
MaxVideosPerSubscription int MaxVideosPerSubscription int
// MinVideoSeconds drops videos shorter than this from discovery (Shorts/clips,
// ADR-023). NewVideos enriches candidates with a single cheap videos.list call
// (contentDetails.duration + snippet.liveBroadcastContent) and filters before
// returning, so the scarce caption-fetch budget is never spent on them. Live
// and upcoming broadcasts are dropped too. Zero disables the filter.
MinVideoSeconds int
// BaseURL overrides the Data API root. Empty means defaultBaseURL. // BaseURL overrides the Data API root. Empty means defaultBaseURL.
BaseURL string BaseURL string
@@ -242,7 +249,97 @@ func (a *Adapter) NewVideos(ctx context.Context, sub domain.Subscription) ([]dom
break break
} }
} }
return videos, nil
// Drop Shorts/sub-minute clips and live/upcoming broadcasts before they ever
// reach the rate-limited caption path (ADR-023). One cheap videos.list call
// (quota API, not the timedtext throttle) supplies duration + live status.
return a.filterLowValue(ctx, client, videos), nil
}
// filterLowValue removes videos shorter than cfg.MinVideoSeconds and any live or
// upcoming broadcast, using a single videos.list lookup for duration +
// liveBroadcastContent. The filter is best-effort: if MinVideoSeconds is 0 (off)
// or the lookup fails, the input is returned unfiltered — discovery must not break
// because a metadata call hiccuped; the worst case is the pre-ADR-023 behaviour.
func (a *Adapter) filterLowValue(ctx context.Context, client *http.Client, videos []domain.Video) []domain.Video {
if a.cfg.MinVideoSeconds <= 0 || len(videos) == 0 {
return videos
}
ids := make([]string, 0, len(videos))
for _, v := range videos {
ids = append(ids, v.ProviderVideoID)
}
q := url.Values{
"part": {"contentDetails,snippet"},
"id": {strings.Join(ids, ",")},
}
var resp videoListResponse
if err := a.getJSON(ctx, client, "/videos", q, &resp); err != nil {
// Degrade open: keep the candidates rather than lose discovery.
return videos
}
type meta struct {
seconds int
live string
}
byID := make(map[string]meta, len(resp.Items))
for _, it := range resp.Items {
byID[it.ID] = meta{seconds: parseISO8601Seconds(it.ContentDetails.Duration), live: it.Snippet.LiveBroadcastContent}
}
kept := videos[:0]
for _, v := range videos {
m, ok := byID[v.ProviderVideoID]
if !ok {
kept = append(kept, v) // unknown metadata: keep, let the fetch decide
continue
}
if m.live != "" && m.live != "none" {
continue // live or upcoming broadcast
}
if m.seconds > 0 && m.seconds < a.cfg.MinVideoSeconds {
continue // Short / sub-threshold clip
}
kept = append(kept, v)
}
return kept
}
// parseISO8601Seconds parses an ISO 8601 duration as returned by the YouTube Data
// API (e.g. "PT1H2M3S", "PT45S", "PT3M") into seconds. Only the hour/minute/second
// components YouTube emits are handled; an unparseable or zero value returns 0,
// which the caller treats as "unknown" (not filtered on duration).
func parseISO8601Seconds(d string) int {
if !strings.HasPrefix(d, "PT") {
return 0
}
d = d[2:]
total, num := 0, 0
seen := false
for _, r := range d {
switch {
case r >= '0' && r <= '9':
num = num*10 + int(r-'0')
seen = true
case r == 'H':
total += num * 3600
num, seen = 0, false
case r == 'M':
total += num * 60
num, seen = 0, false
case r == 'S':
total += num
num, seen = 0, false
default:
return 0 // unexpected component (days/weeks) — treat as unknown
}
}
if seen {
return 0 // trailing digits without a unit: malformed
}
return total
} }
// VideoByID fetches a single video's metadata (videos.list, snippet) for an // VideoByID fetches a single video's metadata (videos.list, snippet) for an
@@ -377,11 +474,16 @@ type playlistItemListResponse struct {
type videoListResponse struct { type videoListResponse struct {
Items []struct { Items []struct {
ID string `json:"id"`
Snippet struct { Snippet struct {
Title string `json:"title"` Title string `json:"title"`
ChannelTitle string `json:"channelTitle"` ChannelTitle string `json:"channelTitle"`
PublishedAt time.Time `json:"publishedAt"` PublishedAt time.Time `json:"publishedAt"`
LiveBroadcastContent string `json:"liveBroadcastContent"`
} `json:"snippet"` } `json:"snippet"`
ContentDetails struct {
Duration string `json:"duration"` // ISO 8601, e.g. "PT1M30S"
} `json:"contentDetails"`
} `json:"items"` } `json:"items"`
} }
+85
View File
@@ -186,6 +186,91 @@ func TestNewVideosCapsAtMax(t *testing.T) {
} }
} }
// TestNewVideosFiltersShortsAndLive: with MinVideoSeconds set, discovery enriches
// candidates via videos.list and drops sub-threshold clips (Shorts) and
// live/upcoming broadcasts before they reach the rate-limited caption path.
func TestNewVideosFiltersShortsAndLive(t *testing.T) {
a, _ := newTestAdapter(t, func(w http.ResponseWriter, r *http.Request) {
switch r.URL.Path {
case "/playlistItems":
_, _ = w.Write([]byte(`{
"items": [
{"snippet": {"title": "Real Talk", "publishedAt": "2026-06-03T10:00:00Z", "resourceId": {"videoId": "long1"}}},
{"snippet": {"title": "A Short", "publishedAt": "2026-06-03T09:00:00Z", "resourceId": {"videoId": "short1"}}},
{"snippet": {"title": "Live Now", "publishedAt": "2026-06-03T08:00:00Z", "resourceId": {"videoId": "live1"}}}
]
}`))
case "/videos":
if got := r.URL.Query().Get("part"); got != "contentDetails,snippet" {
t.Errorf("videos.list part=%q, want contentDetails,snippet", got)
}
_, _ = w.Write([]byte(`{
"items": [
{"id": "long1", "contentDetails": {"duration": "PT12M30S"}, "snippet": {"liveBroadcastContent": "none"}},
{"id": "short1", "contentDetails": {"duration": "PT45S"}, "snippet": {"liveBroadcastContent": "none"}},
{"id": "live1", "contentDetails": {"duration": "PT0S"}, "snippet": {"liveBroadcastContent": "live"}}
]
}`))
default:
t.Errorf("unexpected path %q", r.URL.Path)
}
})
a.cfg.MinVideoSeconds = 60
vids, err := a.NewVideos(context.Background(), domain.Subscription{ID: "s1", UserID: "u1", ChannelID: "UC_acme"})
if err != nil {
t.Fatalf("NewVideos: %v", err)
}
if len(vids) != 1 || vids[0].ProviderVideoID != "long1" {
t.Fatalf("expected only long1 to survive the filter, got %+v", vids)
}
}
// TestNewVideosNoFilterWhenDisabled: MinVideoSeconds=0 keeps the pre-ADR-023
// behaviour — no videos.list call, no filtering.
func TestNewVideosNoFilterWhenDisabled(t *testing.T) {
a, _ := newTestAdapter(t, func(w http.ResponseWriter, r *http.Request) {
if r.URL.Path == "/videos" {
t.Errorf("videos.list must not be called when MinVideoSeconds is 0")
}
_, _ = w.Write([]byte(`{"items": [
{"snippet": {"title": "A Short", "publishedAt": "2026-06-03T09:00:00Z", "resourceId": {"videoId": "short1"}}}
]}`))
})
a.cfg.MinVideoSeconds = 0
vids, err := a.NewVideos(context.Background(), domain.Subscription{ID: "s1", UserID: "u1", ChannelID: "UC_acme"})
if err != nil {
t.Fatalf("NewVideos: %v", err)
}
if len(vids) != 1 {
t.Fatalf("filter disabled must keep all videos, got %d", len(vids))
}
}
func TestParseISO8601Seconds(t *testing.T) {
cases := []struct {
in string
want int
}{
{"PT45S", 45},
{"PT1M30S", 90},
{"PT3M", 180},
{"PT1H2M3S", 3723},
{"PT2H", 7200},
{"PT0S", 0},
{"", 0},
{"garbage", 0},
{"P1D", 0}, // days component not handled → unknown
{"PT10", 0}, // trailing digits without a unit → malformed
}
for _, c := range cases {
if got := parseISO8601Seconds(c.in); got != c.want {
t.Errorf("parseISO8601Seconds(%q) = %d, want %d", c.in, got, c.want)
}
}
}
// TestUploadsPlaylistID covers the zero-cost UC->UU derivation, including // TestUploadsPlaylistID covers the zero-cost UC->UU derivation, including
// non-standard ids that must fall through unchanged (handled via fallback). // non-standard ids that must fall through unchanged (handled via fallback).
func TestUploadsPlaylistID(t *testing.T) { func TestUploadsPlaylistID(t *testing.T) {
+41
View File
@@ -46,6 +46,20 @@ type Config struct {
// MaxTranscriptChars bounds the transcript text sent to the model so a long // MaxTranscriptChars bounds the transcript text sent to the model so a long
// transcript does not overflow a small-context primary. 0 disables truncation. // transcript does not overflow a small-context primary. 0 disables truncation.
MaxTranscriptChars int MaxTranscriptChars int
// MinVideoSeconds drops videos shorter than this from discovery (Shorts and
// other sub-minute clips that are noise and waste the scarce caption-fetch
// 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 // SummarizerTimeout bounds a single completion call. Thinking models are
// slow, so the default is generous. // slow, so the default is generous.
SummarizerTimeout time.Duration SummarizerTimeout time.Duration
@@ -138,6 +152,9 @@ const (
defaultCloudFallbackModel = "berget/mistral-small" defaultCloudFallbackModel = "berget/mistral-small"
defaultSummaryMaxTokens = 1500 defaultSummaryMaxTokens = 1500
defaultMaxTranscriptChars = 18000 defaultMaxTranscriptChars = 18000
defaultMinVideoSeconds = 60
defaultCaptionlessThreshold = 5
defaultCaptionlessWindow = 14 * 24 * time.Hour
defaultSummarizerTimeout = 5 * time.Minute defaultSummarizerTimeout = 5 * time.Minute
defaultYTTokenRef = "youtube/refresh_token" defaultYTTokenRef = "youtube/refresh_token"
defaultYTConnectRedirectURL = "https://tapir.d-ma.be/oauth/youtube/callback" defaultYTConnectRedirectURL = "https://tapir.d-ma.be/oauth/youtube/callback"
@@ -230,6 +247,30 @@ func Load() (Config, error) {
} }
c.MaxTranscriptChars = maxChars c.MaxTranscriptChars = maxChars
minVideo, err := intOr("TAPIR_MIN_VIDEO_SECONDS", defaultMinVideoSeconds)
if err != nil {
return Config{}, err
}
if minVideo < 0 {
minVideo = 0
}
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) onboard, err := intOr("TAPIR_ONBOARD_SUMMARIZE_COUNT", defaultOnboardSummarizeCount)
if err != nil { if err != nil {
return Config{}, err return Config{}, err
+52 -1
View File
@@ -28,6 +28,7 @@ import (
// and is queued for transcript fetch + summarization in this pass. // and is queued for transcript fetch + summarization in this pass.
type passCandidate struct { type passCandidate struct {
v domain.Video 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 pos int // discovery position — used as a stable tiebreak when published_at ties
} }
@@ -74,6 +75,14 @@ type VideoStore interface {
// Called when NewVideos returns domain.ErrChannelUnavailable; best-effort, errors // Called when NewVideos returns domain.ErrChannelUnavailable; best-effort, errors
// are logged and never abort the pass. // are logged and never abort the pass.
UpsertChannelError(ctx context.Context, userID, channelID, channelTitle string) error 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 // 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 backoff time.Duration // rate-limit retry window; 0 = always retry
autoWindow time.Duration // recency bound for auto-summarize; 0 = no bound autoWindow time.Duration // recency bound for auto-summarize; 0 = no bound
now func() time.Time // injectable clock (tests); defaults to time.Now 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 // 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). // 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 } } 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. // 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 { func New(src ports.VideoSource, store VideoStore, engine Processor, userID string, log *slog.Logger, opts ...Option) *Runner {
if log == nil { if log == nil {
@@ -138,6 +158,7 @@ type Stats struct {
SkippedManual int // discovered but not queued, in manual mode SkippedManual int // discovered but not queued, in manual mode
SkippedTooOld int // auto mode: published outside the recency window (not requested) SkippedTooOld int // auto mode: published outside the recency window (not requested)
SkippedRateLimited int // 429'd previously and still inside the backoff window SkippedRateLimited int // 429'd previously and still inside the backoff window
SkippedNoCaptionChannel int // channel suppressed as caption-less (ADR-024)
Errors int Errors int
ChannelUnavailable int // channels that returned HTTP 404 (deleted/private) ChannelUnavailable int // channels that returned HTTP 404 (deleted/private)
} }
@@ -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) subs, err := r.src.ListSubscriptions(ctx, r.userID)
if err != nil { if err != nil {
return stats, fmt.Errorf("runner: list subscriptions: %w", err) return stats, fmt.Errorf("runner: list subscriptions: %w", err)
@@ -273,6 +305,14 @@ func (r *Runner) RunOnce(ctx context.Context) (Stats, error) {
continue 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. // Still inside the rate-limit backoff window: skip without fetching.
if at, ok := rateLimited[id]; ok && r.now().Sub(at) < r.backoff { if at, ok := rateLimited[id]; ok && r.now().Sub(at) < r.backoff {
stats.SkippedRateLimited++ stats.SkippedRateLimited++
@@ -280,7 +320,7 @@ func (r *Runner) RunOnce(ctx context.Context) (Stats, error) {
continue continue
} }
candidates = append(candidates, passCandidate{v: v, pos: pos}) candidates = append(candidates, passCandidate{v: v, channelID: sub.ChannelID, pos: 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)) errs = append(errs, fmt.Errorf("set none status %q: %w", c.v.ProviderVideoID, err))
stats.Errors++ 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) r.log.Info("skipped video (no transcript)", "video", c.v.ProviderVideoID, "title", c.v.Title)
case res.Summary != nil: case res.Summary != nil:
stats.Summarized++ 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)) errs = append(errs, fmt.Errorf("set fetched status %q: %w", c.v.ProviderVideoID, err))
stats.Errors++ 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 // In manual mode the video was explicitly queued; clear the flag so
// it is not re-summarized and the UI drops the "Queued" chip. // it is not re-summarized and the UI drops the "Queued" chip.
if !auto { 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_seen", stats.SkippedSeen, "skipped_no_text", stats.SkippedNoText,
"skipped_manual", stats.SkippedManual, "skipped_too_old", stats.SkippedTooOld, "skipped_manual", stats.SkippedManual, "skipped_too_old", stats.SkippedTooOld,
"skipped_rate_limited", stats.SkippedRateLimited, "skipped_rate_limited", stats.SkippedRateLimited,
"skipped_no_caption_channel", stats.SkippedNoCaptionChannel,
"channel_unavailable", stats.ChannelUnavailable, "errors", stats.Errors) "channel_unavailable", stats.ChannelUnavailable, "errors", stats.Errors)
if err != nil { if err != nil {
r.log.Warn("run pass had errors", "err", err) r.log.Warn("run pass had errors", "err", err)
+72
View File
@@ -51,6 +51,13 @@ type fakeStore struct {
cleared []string cleared []string
rateLimited map[string]time.Time // id -> when 429'd (seeds the backoff window) rateLimited map[string]time.Time // id -> when 429'd (seeds the backoff window)
statuses map[string]string // id -> last SetTranscriptStatus value 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) { 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) 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 { func (f *fakeStore) SetTranscriptStatus(_ context.Context, _, videoID, status string) error {
if f.statuses == nil { if f.statuses == nil {
f.statuses = map[string]string{} f.statuses = map[string]string{}
@@ -284,6 +307,55 @@ func TestRunOnce_AutoMode_OldVideoRequestedBypassesWindow(t *testing.T) {
require.Len(t, sink.delivered, 1) 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 — // TestRunOnce_AutoWindowZero_SummarizesOld: a zero window disables the bound —
// the pre-recency behaviour (summarize every unseen video) is preserved. // the pre-recency behaviour (summarize every unseen video) is preserved.
func TestRunOnce_AutoWindowZero_SummarizesOld(t *testing.T) { func TestRunOnce_AutoWindowZero_SummarizesOld(t *testing.T) {