Compare commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
fe56e2fe01 | ||
|
|
beeb5bc31b | ||
|
|
09eb31d1fe | ||
|
|
1665a1e7c4 | ||
|
|
5219561a91 | ||
|
|
9db06d8a63 | ||
|
|
1aa8a97f95 |
+133
@@ -855,6 +855,139 @@ 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.
|
||||
|
||||
---
|
||||
|
||||
## ADR-025 — Honest, state-aware foreground summarization status
|
||||
|
||||
**Status:** Accepted (2026-06-10). **Pillar B of the manual-mode UX work** (Pillar A, foreground
|
||||
fetch priority, is a separate follow-up). Builds on ADR-014 (the rate limit the UX must make
|
||||
legible).
|
||||
|
||||
**Context.** Clicking "Summarize" spawned a background goroutine and polled `/status`, which
|
||||
returned only two states: the spinner (in-flight) or the normal card (done). But the web
|
||||
`ProcessVideo` only recorded an outcome on *success* — a 429'd or caption-less click left
|
||||
`transcript_status` unset, so the next poll silently reverted to the "Summarize" button. The
|
||||
user saw either an endless spinner or a button that did nothing useful when clicked again. The
|
||||
binding constraint (YouTube's caption rate limit) was completely invisible.
|
||||
|
||||
**Decision.**
|
||||
1. **Record every outcome on the web path**, mirroring the runner: `ProcessVideo` stamps
|
||||
`rate_limited` / `none` / `fetched`. A rate-limited video keeps its requested flag so the
|
||||
background sweep retries it; `none` and `fetched` are terminal.
|
||||
2. **`/status` is state-aware**: summarized → summary card; in-flight → working spinner;
|
||||
`rate_limited` → a calm "waiting, will retry" card that keeps polling (every 30s) so the
|
||||
summary appears on its own when the retry lands — the user never clicks again;
|
||||
`none` → a terminal "no captions" card with no poll and no dead-end button.
|
||||
3. **Charm status text** (Claude-Code / Crush inspired): the working spinner cycles playful,
|
||||
tapir-themed gerunds ("Chewing the cud…", "Munching leaves…", "Distilling the gist…") via
|
||||
CSS only — no JS, keeping the HTMX/no-JS ethos. Decorative (aria-hidden) with a stable
|
||||
`role=status` line for assistive tech.
|
||||
|
||||
**Principle.** When the system cannot be fast (throttled IP), it is at least honest, and it
|
||||
self-resolves without making the user retry. Honesty is the load-bearing half — Pillar A's
|
||||
priority lane only improves the odds of a fast slot; it cannot beat an already-hot IP.
|
||||
|
||||
**Reversibility.** Pure transport-layer + view change over the unchanged engine/ports. No
|
||||
schema change (reuses `transcript_status` from migration 007).
|
||||
|
||||
---
|
||||
|
||||
## ADR-026 — Foreground caption fetches take priority; the credentials probe is dead
|
||||
|
||||
**Status:** Accepted (2026-06-10). **Pillar A of the manual-mode UX work** (Pillar B was
|
||||
ADR-025). Builds on ADR-014 (the shared per-IP gate).
|
||||
|
||||
**Context.** Every caption fetch — the background sweep and the web click-path — shared one
|
||||
process-wide rate gate equally. So a user waiting on a "Summarize" click competed with the
|
||||
firehose for both pacing and the scarce pre-429 window; on a busy IP the click was slow or
|
||||
429'd while the background churned.
|
||||
|
||||
**Decision.** A context-marked priority lane. The web path
|
||||
(`engineProcessor.ProcessVideo`) wraps its context with `ForegroundContext`; the gate gives
|
||||
foreground fetches a token immediately, while **background fetches yield** — they wait until no
|
||||
foreground fetch is pending before taking a token. Threaded via a context value (not new
|
||||
signatures) and a process-wide `foregroundPending` counter. Clicks are rare and bursty, so the
|
||||
background barely loses throughput; the waiting human gets the next (and cleanest) slot.
|
||||
|
||||
**Credentials probe — rejected, not built.** The idea was to fetch captions with the user's
|
||||
auth in manual mode to dodge 429s. It is a dead end, already settled by ADR-010 and the code:
|
||||
the caption path is *deliberately anonymous* because the InnerTube/timedtext endpoints **reject
|
||||
or break on authenticated requests** (`captions.go`: "no OAuth token — it can break the
|
||||
timedtext endpoint"). The user's OAuth (a Data API credential) does not authenticate InnerTube
|
||||
at all, and the official `captions.download` is owner-only (403 on third-party). So auth cannot
|
||||
help here and can actively hurt. No probe needed — building one would only re-confirm the ADR.
|
||||
|
||||
**Reversibility.** Context-marker + a yield loop in the gate; removing the marker collapses to
|
||||
the prior equal-share behaviour. No schema or API change.
|
||||
|
||||
---
|
||||
|
||||
## Rejected alternatives
|
||||
|
||||
Approaches considered during the 2026-06-02 planning + grill session and **deliberately not
|
||||
|
||||
+2
-1
@@ -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,
|
||||
|
||||
+28
-1
@@ -88,6 +88,7 @@ func buildProcessor(cfg config.Config, st *store.Store) (*usecase.Engine, error)
|
||||
ClientSecret: cfg.YTClientSecret,
|
||||
TokenSecretRef: cfg.YTTokenRef,
|
||||
PreferredLanguages: []string{"en"},
|
||||
MinVideoSeconds: cfg.MinVideoSeconds,
|
||||
}, secretStore)
|
||||
|
||||
sum := buildSummarizer(cfg)
|
||||
@@ -112,6 +113,11 @@ type engineProcessor struct {
|
||||
}
|
||||
|
||||
func (p *engineProcessor) ProcessVideo(ctx context.Context, userID, videoID string) error {
|
||||
// This is the user-initiated (foreground) path — a click on "Summarize",
|
||||
// "Try now", or a pasted URL. Mark the context so the caption gate gives it
|
||||
// priority over the background sweep (ADR-026, Pillar A).
|
||||
ctx = youtube.ForegroundContext(ctx)
|
||||
|
||||
row, err := p.store.GetVideoRow(ctx, userID, videoID)
|
||||
if err != nil {
|
||||
return fmt.Errorf("load video %q: %w", videoID, err)
|
||||
@@ -131,7 +137,28 @@ func (p *engineProcessor) ProcessVideo(ctx context.Context, userID, videoID stri
|
||||
if err != nil {
|
||||
return fmt.Errorf("process video %q: %w", videoID, err)
|
||||
}
|
||||
if res.Summary != nil {
|
||||
|
||||
// Record the outcome so the status endpoint can show honest state (ADR-025):
|
||||
// a 429'd or caption-less click used to leave transcript_status unset, so the
|
||||
// poll silently reverted to the "Summarize" button. Mirror the runner: stamp
|
||||
// rate_limited / none / fetched. A rate-limited video keeps its requested flag
|
||||
// so the background sweep retries it; none and fetched are terminal here.
|
||||
switch {
|
||||
case res.Skipped && res.TranscriptSource == string(domain.SourceRateLimited):
|
||||
if err := p.store.SetTranscriptStatus(ctx, userID, videoID, "rate_limited"); err != nil {
|
||||
return fmt.Errorf("set rate_limited status %q: %w", videoID, err)
|
||||
}
|
||||
case res.Skipped:
|
||||
if err := p.store.SetTranscriptStatus(ctx, userID, videoID, "none"); err != nil {
|
||||
return fmt.Errorf("set none status %q: %w", videoID, err)
|
||||
}
|
||||
if err := p.store.ClearSummarizeRequested(ctx, userID, videoID); err != nil {
|
||||
return fmt.Errorf("clear summarize flag %q: %w", videoID, err)
|
||||
}
|
||||
case res.Summary != nil:
|
||||
if err := p.store.SetTranscriptStatus(ctx, userID, videoID, "fetched"); err != nil {
|
||||
return fmt.Errorf("set fetched status %q: %w", videoID, err)
|
||||
}
|
||||
if err := p.store.ClearSummarizeRequested(ctx, userID, videoID); err != nil {
|
||||
return fmt.Errorf("clear summarize flag %q: %w", videoID, err)
|
||||
}
|
||||
|
||||
+51
-27
@@ -33,6 +33,7 @@ func buildUserRunner(cfg config.Config, st *store.Store, secretStore ports.Secre
|
||||
ClientSecret: cfg.YTClientSecret,
|
||||
TokenSecretRef: web.YouTubeTokenRef(userID),
|
||||
PreferredLanguages: []string{"en"},
|
||||
MinVideoSeconds: cfg.MinVideoSeconds,
|
||||
}, secretStore)
|
||||
|
||||
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,
|
||||
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
|
||||
@@ -73,22 +75,17 @@ func runDiscoveryPass(
|
||||
return runner.Stats{}
|
||||
}
|
||||
|
||||
// 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 user 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 gives every user the lead slot in turn.
|
||||
users = rotateUsers(users, pass)
|
||||
|
||||
log.Info("scheduler: starting discovery pass", "users", len(users))
|
||||
var total runner.Stats
|
||||
// Keep only users with a video connection. A pass for a connectionless user
|
||||
// (e.g. a stale Dex-era orphan identity) only tries to resolve a token that
|
||||
// was never minted, logging a spurious "ref not found" every tick. Filtering
|
||||
// here — BEFORE rotation — also keeps fairness honest: rotation is over the
|
||||
// users that actually consume the caption budget, so a dead identity can't eat
|
||||
// a rotation slot and skew the lead share.
|
||||
var connected []store.UserIdentity
|
||||
for _, u := range users {
|
||||
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)
|
||||
if err != nil {
|
||||
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)
|
||||
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)
|
||||
total = sumStats(total, stats)
|
||||
if err != nil {
|
||||
@@ -109,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
|
||||
}
|
||||
@@ -129,8 +143,18 @@ func runScheduler(
|
||||
return // disabled
|
||||
}
|
||||
|
||||
pass := 0
|
||||
runDiscoveryPass(ctx, pass, lister, runUser, log)
|
||||
// Derive the rotation offset from wall-clock, NOT an in-memory counter. A
|
||||
// counter reset to 0 on every pod restart always hands the lead to the
|
||||
// first-listed user — so frequent deploys re-starve whoever is last (exactly
|
||||
// what happened to the first pilot user during a deploy-heavy session). A
|
||||
// time-based offset advances with real time and is identical across restarts,
|
||||
// so the lead rotates fairly regardless of how often the pod bounces.
|
||||
runPass := func() {
|
||||
pass := int(time.Now().Unix() / int64(interval/time.Second))
|
||||
runDiscoveryPass(ctx, pass, lister, runUser, log)
|
||||
}
|
||||
|
||||
runPass()
|
||||
|
||||
ticker := time.NewTicker(interval)
|
||||
defer ticker.Stop()
|
||||
@@ -139,8 +163,7 @@ func runScheduler(
|
||||
case <-ctx.Done():
|
||||
return
|
||||
case <-ticker.C:
|
||||
pass++
|
||||
runDiscoveryPass(ctx, pass, lister, runUser, log)
|
||||
runPass()
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -169,14 +192,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,
|
||||
}
|
||||
}
|
||||
|
||||
@@ -128,6 +128,21 @@ func TestDiscoveryPassRotatesLeadUser(t *testing.T) {
|
||||
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) {
|
||||
// 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.
|
||||
|
||||
@@ -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
|
||||
`18000`** (~fits an 8k-context model). `0` disables truncation. Prevents the context-overflow
|
||||
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
|
||||
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):**
|
||||
|
||||
@@ -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")
|
||||
}
|
||||
@@ -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
|
||||
|
||||
@@ -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);
|
||||
@@ -4,6 +4,7 @@ import (
|
||||
"context"
|
||||
"fmt"
|
||||
"os"
|
||||
"path/filepath"
|
||||
"testing"
|
||||
|
||||
embeddedpostgres "github.com/fergusstrange/embedded-postgres"
|
||||
@@ -24,11 +25,22 @@ var _ ports.Sink = (*store.Store)(nil)
|
||||
var dsn string
|
||||
|
||||
func TestMain(m *testing.M) {
|
||||
const port = 54329
|
||||
// Port + runtime/data dirs are per-process (PID-derived) so two concurrent
|
||||
// `go test` invocations — e.g. a push-run and a tag-run firing together in CI —
|
||||
// don't collide on a fixed port or a shared data dir (which silently failed
|
||||
// both runs). CachePath is shared so the PG archive is downloaded once, not
|
||||
// per process. Base 54000 keeps this package's range distinct from web's.
|
||||
port := uint32(54000 + os.Getpid()%1000)
|
||||
dsn = fmt.Sprintf("postgres://postgres:postgres@localhost:%d/postgres?sslmode=disable", port)
|
||||
|
||||
rt := filepath.Join(os.TempDir(), fmt.Sprintf("tapir-epg-store-%d", os.Getpid()))
|
||||
pg := embeddedpostgres.NewDatabase(
|
||||
embeddedpostgres.DefaultConfig().Port(port),
|
||||
embeddedpostgres.DefaultConfig().
|
||||
Port(port).
|
||||
RuntimePath(rt).
|
||||
DataPath(filepath.Join(rt, "data")).
|
||||
BinariesPath(filepath.Join(rt, "bin")).
|
||||
CachePath(filepath.Join(os.TempDir(), "tapir-epg-cache")),
|
||||
)
|
||||
if err := pg.Start(); err != nil {
|
||||
fmt.Fprintf(os.Stderr, "embedded-postgres start: %v\n", err)
|
||||
@@ -40,6 +52,7 @@ func TestMain(m *testing.M) {
|
||||
if err := pg.Stop(); err != nil {
|
||||
fmt.Fprintf(os.Stderr, "embedded-postgres stop: %v\n", err)
|
||||
}
|
||||
_ = os.RemoveAll(rt)
|
||||
os.Exit(code)
|
||||
}
|
||||
|
||||
|
||||
@@ -2,6 +2,7 @@ package youtube
|
||||
|
||||
import (
|
||||
"context"
|
||||
"sync/atomic"
|
||||
"time"
|
||||
|
||||
"golang.org/x/time/rate"
|
||||
@@ -29,9 +30,55 @@ func SetFetchRate(interval time.Duration) {
|
||||
globalFetchGate = rate.NewLimiter(rate.Every(interval), 1)
|
||||
}
|
||||
|
||||
// foregroundPending counts in-flight foreground (user-initiated) caption fetches.
|
||||
// The background sweep yields the gate while this is non-zero so a human waiting
|
||||
// on a click gets the next slot — and, on a near-throttled IP, the pre-429 window
|
||||
// — instead of competing equally with the firehose (ADR-026, Pillar A). Clicks are
|
||||
// rare and bursty, so background barely notices; the win to the click is large.
|
||||
var foregroundPending atomic.Int64
|
||||
|
||||
// fgCtxKey marks a context as foreground (user-initiated). Unexported; set via
|
||||
// ForegroundContext and read via isForeground so only this package owns the key.
|
||||
type fgCtxKey struct{}
|
||||
|
||||
// ForegroundContext marks ctx as a user-initiated (foreground) fetch so the gate
|
||||
// gives it priority. The web "Summarize"/paste/retry path wraps its context with
|
||||
// this; the background scheduler leaves it unset.
|
||||
func ForegroundContext(ctx context.Context) context.Context {
|
||||
return context.WithValue(ctx, fgCtxKey{}, true)
|
||||
}
|
||||
|
||||
func isForeground(ctx context.Context) bool {
|
||||
v, _ := ctx.Value(fgCtxKey{}).(bool)
|
||||
return v
|
||||
}
|
||||
|
||||
// fgYieldPoll is how often a background waiter re-checks whether a foreground
|
||||
// fetch is still pending. Short enough to feel immediate, long enough not to spin.
|
||||
const fgYieldPoll = 200 * time.Millisecond
|
||||
|
||||
// WaitFetchGate blocks until the process-wide gate allows one timedtext fetch,
|
||||
// respecting ctx cancellation. Called from httpDo before every live outbound
|
||||
// caption fetch so the scheduler and the click-path share the same egress budget.
|
||||
//
|
||||
// Foreground (user-initiated) fetches take priority: they register as pending and
|
||||
// acquire a token immediately. Background fetches first yield — they wait until no
|
||||
// foreground fetch is pending — so a live click is never stuck behind the
|
||||
// background sweep and gets the cleaner slot against the per-IP limit (ADR-026).
|
||||
func WaitFetchGate(ctx context.Context) error {
|
||||
if isForeground(ctx) {
|
||||
foregroundPending.Add(1)
|
||||
defer foregroundPending.Add(-1)
|
||||
return globalFetchGate.Wait(ctx)
|
||||
}
|
||||
|
||||
// Background: defer to any pending foreground fetch before taking a token.
|
||||
for foregroundPending.Load() > 0 {
|
||||
select {
|
||||
case <-ctx.Done():
|
||||
return ctx.Err()
|
||||
case <-time.After(fgYieldPoll):
|
||||
}
|
||||
}
|
||||
return globalFetchGate.Wait(ctx)
|
||||
}
|
||||
|
||||
@@ -82,3 +82,37 @@ func TestSetFetchRateZeroIsUnlimited(t *testing.T) {
|
||||
require.NoError(t, WaitFetchGate(context.Background()))
|
||||
}
|
||||
}
|
||||
|
||||
func TestForegroundContextMarker(t *testing.T) {
|
||||
require.False(t, isForeground(context.Background()), "plain context is background")
|
||||
require.True(t, isForeground(ForegroundContext(context.Background())), "marked context is foreground")
|
||||
}
|
||||
|
||||
// TestWaitFetchGateForegroundProceedsImmediately: a foreground fetch acquires a
|
||||
// token without yielding, even when background callers exist.
|
||||
func TestWaitFetchGateForegroundProceedsImmediately(t *testing.T) {
|
||||
SetFetchRate(0) // unlimited limiter — isolate the yield logic from pacing
|
||||
foregroundPending.Store(0)
|
||||
t.Cleanup(func() { foregroundPending.Store(0) })
|
||||
|
||||
ctx, cancel := context.WithTimeout(context.Background(), time.Second)
|
||||
defer cancel()
|
||||
require.NoError(t, WaitFetchGate(ForegroundContext(ctx)), "foreground proceeds immediately")
|
||||
}
|
||||
|
||||
// TestWaitFetchGateBackgroundYieldsToForeground: while a foreground fetch is
|
||||
// pending, a background fetch yields (does not take a token) until the foreground
|
||||
// clears — proven by a background wait timing out against its own deadline, then
|
||||
// succeeding once the foreground is done.
|
||||
func TestWaitFetchGateBackgroundYieldsToForeground(t *testing.T) {
|
||||
SetFetchRate(0)
|
||||
foregroundPending.Store(1) // simulate a foreground fetch in flight
|
||||
t.Cleanup(func() { foregroundPending.Store(0) })
|
||||
|
||||
ctx, cancel := context.WithTimeout(context.Background(), 250*time.Millisecond)
|
||||
defer cancel()
|
||||
require.Error(t, WaitFetchGate(ctx), "background yields (blocks) while foreground is pending")
|
||||
|
||||
foregroundPending.Store(0) // foreground done
|
||||
require.NoError(t, WaitFetchGate(context.Background()), "background proceeds once foreground clears")
|
||||
}
|
||||
|
||||
@@ -66,6 +66,13 @@ type Config struct {
|
||||
// poll. Zero means defaultMaxVideos.
|
||||
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 string
|
||||
|
||||
@@ -242,7 +249,97 @@ func (a *Adapter) NewVideos(ctx context.Context, sub domain.Subscription) ([]dom
|
||||
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
|
||||
@@ -377,11 +474,16 @@ type playlistItemListResponse struct {
|
||||
|
||||
type videoListResponse struct {
|
||||
Items []struct {
|
||||
ID string `json:"id"`
|
||||
Snippet struct {
|
||||
Title string `json:"title"`
|
||||
ChannelTitle string `json:"channelTitle"`
|
||||
PublishedAt time.Time `json:"publishedAt"`
|
||||
Title string `json:"title"`
|
||||
ChannelTitle string `json:"channelTitle"`
|
||||
PublishedAt time.Time `json:"publishedAt"`
|
||||
LiveBroadcastContent string `json:"liveBroadcastContent"`
|
||||
} `json:"snippet"`
|
||||
ContentDetails struct {
|
||||
Duration string `json:"duration"` // ISO 8601, e.g. "PT1M30S"
|
||||
} `json:"contentDetails"`
|
||||
} `json:"items"`
|
||||
}
|
||||
|
||||
|
||||
@@ -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
|
||||
// non-standard ids that must fall through unchanged (handled via fallback).
|
||||
func TestUploadsPlaylistID(t *testing.T) {
|
||||
|
||||
@@ -46,6 +46,20 @@ type Config struct {
|
||||
// MaxTranscriptChars bounds the transcript text sent to the model so a long
|
||||
// transcript does not overflow a small-context primary. 0 disables truncation.
|
||||
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
|
||||
// slow, so the default is generous.
|
||||
SummarizerTimeout time.Duration
|
||||
@@ -138,6 +152,9 @@ const (
|
||||
defaultCloudFallbackModel = "berget/mistral-small"
|
||||
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"
|
||||
@@ -230,6 +247,30 @@ func Load() (Config, error) {
|
||||
}
|
||||
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)
|
||||
if err != nil {
|
||||
return Config{}, err
|
||||
|
||||
+63
-12
@@ -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)
|
||||
|
||||
@@ -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) {
|
||||
|
||||
@@ -495,11 +495,22 @@ func (a *App) handleStatus(w http.ResponseWriter, r *http.Request) {
|
||||
return
|
||||
}
|
||||
|
||||
if row.Summarized || !a.Processing.Has(processingKey(userID, videoID)) {
|
||||
// Honest, state-aware status (ADR-025). Order matters: a finished summary wins;
|
||||
// an in-flight goroutine shows the working spinner; a recorded rate-limit shows
|
||||
// the calm "waiting, will retry" card that keeps polling; a recorded "none" is
|
||||
// terminal; anything else falls back to the normal card.
|
||||
switch {
|
||||
case row.Summarized:
|
||||
a.render(w, r, VideoCard(*row))
|
||||
case a.Processing.Has(processingKey(userID, videoID)):
|
||||
a.render(w, r, processingCard(*row))
|
||||
case row.TranscriptStatus == "rate_limited":
|
||||
a.render(w, r, waitingCard(*row))
|
||||
case row.TranscriptStatus == "none":
|
||||
a.render(w, r, noCaptionsCard(*row))
|
||||
default:
|
||||
a.render(w, r, VideoCard(*row))
|
||||
return
|
||||
}
|
||||
a.render(w, r, processingCard(*row))
|
||||
}
|
||||
|
||||
// handleSummarizeMode toggles the user's auto/manual summarization mode. The form
|
||||
|
||||
@@ -7,6 +7,7 @@ import (
|
||||
"net/http"
|
||||
"net/http/httptest"
|
||||
"os"
|
||||
"path/filepath"
|
||||
"strings"
|
||||
"testing"
|
||||
"time"
|
||||
@@ -26,10 +27,22 @@ import (
|
||||
var dsn string
|
||||
|
||||
func TestMain(m *testing.M) {
|
||||
const port = 54330 // distinct from the store package's embedded PG (54329)
|
||||
// Per-process port + dirs so concurrent `go test` runs (e.g. a push-run and a
|
||||
// tag-run in CI) never collide on a fixed port or shared data dir. Base 55000
|
||||
// keeps web's range distinct from the store package (54000). Shared CachePath
|
||||
// downloads the PG archive once.
|
||||
port := uint32(55000 + os.Getpid()%1000)
|
||||
dsn = fmt.Sprintf("postgres://postgres:postgres@localhost:%d/postgres?sslmode=disable", port)
|
||||
|
||||
pg := embeddedpostgres.NewDatabase(embeddedpostgres.DefaultConfig().Port(port))
|
||||
rt := filepath.Join(os.TempDir(), fmt.Sprintf("tapir-epg-web-%d", os.Getpid()))
|
||||
pg := embeddedpostgres.NewDatabase(
|
||||
embeddedpostgres.DefaultConfig().
|
||||
Port(port).
|
||||
RuntimePath(rt).
|
||||
DataPath(filepath.Join(rt, "data")).
|
||||
BinariesPath(filepath.Join(rt, "bin")).
|
||||
CachePath(filepath.Join(os.TempDir(), "tapir-epg-cache")),
|
||||
)
|
||||
if err := pg.Start(); err != nil {
|
||||
fmt.Fprintf(os.Stderr, "embedded-postgres start: %v\n", err)
|
||||
os.Exit(1)
|
||||
@@ -38,6 +51,7 @@ func TestMain(m *testing.M) {
|
||||
if err := pg.Stop(); err != nil {
|
||||
fmt.Fprintf(os.Stderr, "embedded-postgres stop: %v\n", err)
|
||||
}
|
||||
_ = os.RemoveAll(rt)
|
||||
os.Exit(code)
|
||||
}
|
||||
|
||||
|
||||
@@ -7,6 +7,7 @@ import (
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"github.com/jackc/pgx/v5/pgxpool"
|
||||
"github.com/stretchr/testify/require"
|
||||
|
||||
"gitea.d-ma.be/mathias/tapir/internal/web"
|
||||
@@ -129,3 +130,42 @@ func getStatus(t *testing.T, app *web.App, videoID string) *httptest.ResponseRec
|
||||
app.Router().ServeHTTP(rec, req)
|
||||
return rec
|
||||
}
|
||||
|
||||
// setTranscriptStatus stamps videos.transcript_status directly (bypassing RLS via
|
||||
// the super pool) so a test can drive the status endpoint into a given state.
|
||||
func setTranscriptStatus(t *testing.T, p *pgxpool.Pool, videoID, status string) {
|
||||
t.Helper()
|
||||
_, err := p.Exec(context.Background(),
|
||||
`UPDATE videos SET transcript_status = $2 WHERE id = $1`, videoID, status)
|
||||
require.NoError(t, err)
|
||||
}
|
||||
|
||||
// TestStatusRateLimitedShowsWaitingCard: a click that hit YouTube's rate limit
|
||||
// must surface the honest "waiting, will retry" card that keeps polling — not a
|
||||
// silent revert to the Summarize button.
|
||||
func TestStatusRateLimitedShowsWaitingCard(t *testing.T) {
|
||||
app := newApp(t)
|
||||
p := rawPool(t)
|
||||
resetDB(t, p)
|
||||
seedVideo(t, p, videoX, "Throttled Title", "https://x", time.Time{})
|
||||
setTranscriptStatus(t, p, videoX, "rate_limited")
|
||||
|
||||
html := body(t, getStatus(t, app, videoX))
|
||||
require.Contains(t, html, "Waiting on YouTube rate limits", "honest rate-limit copy")
|
||||
require.Contains(t, html, `hx-trigger="every 30s"`, "waiting card keeps polling so it self-resolves")
|
||||
require.NotContains(t, html, "Summarize this video", "must not revert to the Summarize button")
|
||||
}
|
||||
|
||||
// TestStatusNoCaptionsTerminal: a no-captions outcome is terminal — an honest
|
||||
// message, no poll, no button to click back into the same dead end.
|
||||
func TestStatusNoCaptionsTerminal(t *testing.T) {
|
||||
app := newApp(t)
|
||||
p := rawPool(t)
|
||||
resetDB(t, p)
|
||||
seedVideo(t, p, videoX, "Silent Title", "https://x", time.Time{})
|
||||
setTranscriptStatus(t, p, videoX, "none")
|
||||
|
||||
html := body(t, getStatus(t, app, videoX))
|
||||
require.Contains(t, html, "No captions available", "honest terminal copy")
|
||||
require.NotContains(t, html, "hx-trigger", "terminal card must stop polling")
|
||||
}
|
||||
|
||||
@@ -658,11 +658,30 @@ a.btn, a.btn:visited { color: var(--accent-fg); }
|
||||
.tapir-bar-fill { animation: tapir-fill 8s linear infinite; text-shadow: 0 0 6px rgba(14, 249, 182, .7); }
|
||||
@keyframes tapir-fill { 0% { clip-path: inset(0 100% 0 0); } 100% { clip-path: inset(0 0 0 0); } }
|
||||
.tapir-label { color: var(--muted); font-size: .9rem; margin: 0; }
|
||||
/* Cycling status verbs (Claude-Code / Crush style): five gerunds stacked, each
|
||||
visible 1/5 of a 6s loop, cross-faded. The container reserves one line height
|
||||
so the layout does not jump as verbs swap. */
|
||||
.tapir-verbs { position: relative; height: 1.3em; margin: .2em 0 0; color: var(--muted); font-size: .9rem; }
|
||||
.tapir-verbs span { position: absolute; left: 0; top: 0; white-space: nowrap; opacity: 0; animation: tapir-verb 6s steps(1, end) infinite; }
|
||||
.tapir-verbs .tv1 { animation-delay: 0s; }
|
||||
.tapir-verbs .tv2 { animation-delay: 1.2s; }
|
||||
.tapir-verbs .tv3 { animation-delay: 2.4s; }
|
||||
.tapir-verbs .tv4 { animation-delay: 3.6s; }
|
||||
.tapir-verbs .tv5 { animation-delay: 4.8s; }
|
||||
@keyframes tapir-verb { 0%, 19.99% { opacity: 1; } 20%, 100% { opacity: 0; } }
|
||||
/* Resting tapir for the rate-limit waiting state: the panel, one still frame, no
|
||||
animation — calm, not busy, signalling "parked, not stuck". */
|
||||
.tapir-resting pre { position: relative; opacity: 1; animation: none; }
|
||||
.card-waiting { border-style: dashed; opacity: .92; }
|
||||
.card-no-captions .card-state { font-style: italic; }
|
||||
.sr-only { position: absolute; width: 1px; height: 1px; padding: 0; margin: -1px; overflow: hidden; clip: rect(0, 0, 0, 0); white-space: nowrap; border: 0; }
|
||||
@media (prefers-reduced-motion: reduce) {
|
||||
.tapir-charm pre { animation: none; }
|
||||
.tapir-charm .tapir-f2, .tapir-charm .tapir-f3 { display: none; }
|
||||
.tapir-charm .tapir-f1 { opacity: 1; }
|
||||
.tapir-bar-fill { animation: none; clip-path: inset(0 35% 0 0); }
|
||||
.tapir-verbs span { animation: none; }
|
||||
.tapir-verbs .tv1 { opacity: 1; }
|
||||
}
|
||||
|
||||
/* summarization mode toggle on the account page */
|
||||
|
||||
@@ -339,7 +339,17 @@ templ TapirSpinner() {
|
||||
<pre class="tapir-f3">@templ.Raw(tapirFrameHTML3)</pre>
|
||||
<div class="tapir-bar"><span class="tapir-bar-fill" style={ "color:" + CharmMint }>{ tapirBarFill }</span></div>
|
||||
</div>
|
||||
<p class="tapir-label" role="status" aria-live="polite"><em>Summarizing…</em></p>
|
||||
// Claude-Code / Crush-style status: playful gerunds cycle in place (CSS only,
|
||||
// no JS). Decorative — aria-hidden — with one stable status line below for
|
||||
// assistive tech.
|
||||
<p class="tapir-verbs" aria-hidden="true">
|
||||
<span class="tv1"><em>Fetching captions…</em></span>
|
||||
<span class="tv2"><em>Chewing the cud…</em></span>
|
||||
<span class="tv3"><em>Munching leaves…</em></span>
|
||||
<span class="tv4"><em>Distilling the gist…</em></span>
|
||||
<span class="tv5"><em>Summarizing…</em></span>
|
||||
</p>
|
||||
<p class="sr-only" role="status" aria-live="polite">Summarizing…</p>
|
||||
}
|
||||
|
||||
// processingCard is the in-flight summarization card. It replaces the Summarize
|
||||
@@ -363,6 +373,43 @@ templ processingCard(r store.SummaryRow) {
|
||||
</li>
|
||||
}
|
||||
|
||||
// waitingCard is the honest rate-limited state: the click landed but YouTube is
|
||||
// throttling the caption fetch, so the tapir rests and the card keeps polling
|
||||
// (gently, every 30s) until the background retry lands the summary — the user
|
||||
// never has to click again. Replaces the old silent revert to a Summarize button.
|
||||
templ waitingCard(r store.SummaryRow) {
|
||||
<li
|
||||
class="card card-waiting"
|
||||
id={ "video-" + r.VideoID }
|
||||
hx-get={ string(statusURL(r.VideoID)) }
|
||||
hx-trigger="every 30s"
|
||||
hx-swap="outerHTML"
|
||||
>
|
||||
<div class="card-title">{ displayTitle(r) }</div>
|
||||
if cardMeta(r) != "" {
|
||||
<div class="card-meta">{ cardMeta(r) }</div>
|
||||
}
|
||||
<div class="tapir-charm tapir-resting" aria-hidden="true">
|
||||
<pre class="tapir-f1">@templ.Raw(tapirFrameHTML2)</pre>
|
||||
</div>
|
||||
<p class="tapir-label" role="status" aria-live="polite">
|
||||
Waiting on YouTube rate limits. Tapir keeps trying, slowly and politely, and the summary will appear here on its own.
|
||||
</p>
|
||||
</li>
|
||||
}
|
||||
|
||||
// noCaptionsCard is the terminal no-captions state: nothing to summarize, so the
|
||||
// card stops (no poll, no button to click again into the same dead end).
|
||||
templ noCaptionsCard(r store.SummaryRow) {
|
||||
<li class="card card-no-captions" id={ "video-" + r.VideoID }>
|
||||
<div class="card-title">{ displayTitle(r) }</div>
|
||||
if cardMeta(r) != "" {
|
||||
<div class="card-meta">{ cardMeta(r) }</div>
|
||||
}
|
||||
<p class="card-state muted">No captions available, so Tapir cannot summarize this one.</p>
|
||||
</li>
|
||||
}
|
||||
|
||||
// DetailPage is the full summary view: text, highlights, takeaways, metadata,
|
||||
// and the action button group.
|
||||
templ DetailPage(r store.SummaryRow) {
|
||||
|
||||
+401
-218
File diff suppressed because it is too large
Load Diff
Reference in New Issue
Block a user