Compare commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
1665a1e7c4 | ||
|
|
5219561a91 | ||
|
|
9db06d8a63 | ||
|
|
1aa8a97f95 |
+103
@@ -855,6 +855,109 @@ 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).
|
||||||
|
|
||||||
|
---
|
||||||
|
|
||||||
## 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
@@ -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,
|
||||||
|
|||||||
+23
-1
@@ -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)
|
||||||
@@ -131,7 +132,28 @@ func (p *engineProcessor) ProcessVideo(ctx context.Context, userID, videoID stri
|
|||||||
if err != nil {
|
if err != nil {
|
||||||
return fmt.Errorf("process video %q: %w", videoID, err)
|
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 {
|
if err := p.store.ClearSummarizeRequested(ctx, userID, videoID); err != nil {
|
||||||
return fmt.Errorf("clear summarize flag %q: %w", videoID, err)
|
return fmt.Errorf("clear summarize flag %q: %w", videoID, err)
|
||||||
}
|
}
|
||||||
|
|||||||
+38
-23
@@ -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
|
||||||
}
|
}
|
||||||
@@ -169,14 +183,15 @@ func rotateUsers(users []store.UserIdentity, pass int) []store.UserIdentity {
|
|||||||
// per-tick aggregate across all users.
|
// per-tick aggregate across all users.
|
||||||
func sumStats(a, b runner.Stats) runner.Stats {
|
func sumStats(a, b runner.Stats) runner.Stats {
|
||||||
return runner.Stats{
|
return runner.Stats{
|
||||||
Candidates: a.Candidates + b.Candidates,
|
Candidates: a.Candidates + b.Candidates,
|
||||||
Summarized: a.Summarized + b.Summarized,
|
Summarized: a.Summarized + b.Summarized,
|
||||||
SkippedSeen: a.SkippedSeen + b.SkippedSeen,
|
SkippedSeen: a.SkippedSeen + b.SkippedSeen,
|
||||||
SkippedNoText: a.SkippedNoText + b.SkippedNoText,
|
SkippedNoText: a.SkippedNoText + b.SkippedNoText,
|
||||||
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,
|
||||||
ChannelUnavailable: a.ChannelUnavailable + b.ChannelUnavailable,
|
SkippedNoCaptionChannel: a.SkippedNoCaptionChannel + b.SkippedNoCaptionChannel,
|
||||||
Errors: a.Errors + b.Errors,
|
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"))
|
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.
|
||||||
|
|||||||
@@ -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")
|
||||||
|
}
|
||||||
@@ -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);
|
||||||
@@ -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"`
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -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) {
|
||||||
|
|||||||
@@ -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
|
||||||
|
|||||||
+63
-12
@@ -27,8 +27,9 @@ import (
|
|||||||
// passCandidate is a video that passed all pre-filters (seen/manual/backoff)
|
// passCandidate is a video that passed all pre-filters (seen/manual/backoff)
|
||||||
// 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
|
||||||
pos int // discovery position — used as a stable tiebreak when published_at ties
|
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,
|
// 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
|
// 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 {
|
||||||
@@ -131,15 +151,16 @@ func New(src ports.VideoSource, store VideoStore, engine Processor, userID strin
|
|||||||
|
|
||||||
// Stats summarizes one RunOnce pass.
|
// Stats summarizes one RunOnce pass.
|
||||||
type Stats struct {
|
type Stats struct {
|
||||||
Candidates int
|
Candidates int
|
||||||
Summarized int
|
Summarized int
|
||||||
SkippedSeen int
|
SkippedSeen int
|
||||||
SkippedNoText int
|
SkippedNoText int
|
||||||
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
|
||||||
Errors int
|
SkippedNoCaptionChannel int // channel suppressed as caption-less (ADR-024)
|
||||||
ChannelUnavailable int // channels that returned HTTP 404 (deleted/private)
|
Errors int
|
||||||
|
ChannelUnavailable int // channels that returned HTTP 404 (deleted/private)
|
||||||
}
|
}
|
||||||
|
|
||||||
// tooOld reports whether a video published at publishedAt falls outside the
|
// 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)
|
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)
|
||||||
|
|||||||
@@ -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) {
|
||||||
|
|||||||
@@ -495,11 +495,22 @@ func (a *App) handleStatus(w http.ResponseWriter, r *http.Request) {
|
|||||||
return
|
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))
|
a.render(w, r, VideoCard(*row))
|
||||||
return
|
|
||||||
}
|
}
|
||||||
a.render(w, r, processingCard(*row))
|
|
||||||
}
|
}
|
||||||
|
|
||||||
// handleSummarizeMode toggles the user's auto/manual summarization mode. The form
|
// handleSummarizeMode toggles the user's auto/manual summarization mode. The form
|
||||||
|
|||||||
@@ -7,6 +7,7 @@ import (
|
|||||||
"testing"
|
"testing"
|
||||||
"time"
|
"time"
|
||||||
|
|
||||||
|
"github.com/jackc/pgx/v5/pgxpool"
|
||||||
"github.com/stretchr/testify/require"
|
"github.com/stretchr/testify/require"
|
||||||
|
|
||||||
"gitea.d-ma.be/mathias/tapir/internal/web"
|
"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)
|
app.Router().ServeHTTP(rec, req)
|
||||||
return rec
|
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); }
|
.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); } }
|
@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; }
|
.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) {
|
@media (prefers-reduced-motion: reduce) {
|
||||||
.tapir-charm pre { animation: none; }
|
.tapir-charm pre { animation: none; }
|
||||||
.tapir-charm .tapir-f2, .tapir-charm .tapir-f3 { display: none; }
|
.tapir-charm .tapir-f2, .tapir-charm .tapir-f3 { display: none; }
|
||||||
.tapir-charm .tapir-f1 { opacity: 1; }
|
.tapir-charm .tapir-f1 { opacity: 1; }
|
||||||
.tapir-bar-fill { animation: none; clip-path: inset(0 35% 0 0); }
|
.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 */
|
/* summarization mode toggle on the account page */
|
||||||
|
|||||||
@@ -339,7 +339,17 @@ templ TapirSpinner() {
|
|||||||
<pre class="tapir-f3">@templ.Raw(tapirFrameHTML3)</pre>
|
<pre class="tapir-f3">@templ.Raw(tapirFrameHTML3)</pre>
|
||||||
<div class="tapir-bar"><span class="tapir-bar-fill" style={ "color:" + CharmMint }>{ tapirBarFill }</span></div>
|
<div class="tapir-bar"><span class="tapir-bar-fill" style={ "color:" + CharmMint }>{ tapirBarFill }</span></div>
|
||||||
</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
|
// processingCard is the in-flight summarization card. It replaces the Summarize
|
||||||
@@ -363,6 +373,43 @@ templ processingCard(r store.SummaryRow) {
|
|||||||
</li>
|
</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,
|
// DetailPage is the full summary view: text, highlights, takeaways, metadata,
|
||||||
// and the action button group.
|
// and the action button group.
|
||||||
templ DetailPage(r store.SummaryRow) {
|
templ DetailPage(r store.SummaryRow) {
|
||||||
|
|||||||
+401
-218
File diff suppressed because it is too large
Load Diff
Reference in New Issue
Block a user