Files
tapir/internal/runner/runner.go
T
mathias f1e9739900
CI / Lint / Test / Vet (push) Successful in 26s
CI / Build & Import (push) Successful in 11s
feat(store,runner,web): channel unavailability notice (migration 013)
YouTube channels that 404 on playlist discovery (deleted/private) are now:
1. Wrapped in domain.ErrChannelUnavailable by the YouTube adapter (instead of
   a generic error), so the runner can identify them without string-matching.
2. Stored per-user in channel_errors (migration 013, RLS-guarded) via runner's
   new UpsertChannelError path — removed from the generic Errors counter,
   counted separately as ChannelUnavailable.
3. Shown on the account page under "Unavailable channels" with name, chip-warn
   badge, and first-seen date, so users know why some subscribed channels
   produce no videos.
2026-06-06 10:09:52 +02:00

290 lines
11 KiB
Go

// Package runner wires the engine to the durable store for the `tapir run`
// command. It owns the cross-restart dedup the engine core deliberately does
// not: the engine's in-memory processed map is process-lifetime only, so this
// loads the store's SeenVideoIDs and skips videos already summarized in a prior
// run. It also assigns each video its durable store id (UpsertVideo) before
// processing, so the summary's video_id equals the dedup key.
//
// It depends on small local interfaces (VideoStore, Processor), not concrete
// types, so the loop is tested with fakes — no live YouTube, gateway, or PG.
package runner
import (
"context"
"errors"
"fmt"
"log/slog"
"os"
"time"
"gitea.d-ma.be/mathias/tapir/internal/domain"
"gitea.d-ma.be/mathias/tapir/internal/ports"
"gitea.d-ma.be/mathias/tapir/internal/usecase"
)
// VideoStore is the durable persistence the run loop needs: assign a stable id +
// metadata, read the already-summarized set, and (for manual summarization mode)
// read the user's mode + queued videos and clear a video's queue flag once it has
// been summarized. *store.Store satisfies it.
type VideoStore interface {
UpsertVideo(ctx context.Context, v domain.Video) (string, error)
SeenVideoIDs(ctx context.Context, userID string) (map[string]bool, error)
GetAutoSummarize(ctx context.Context, userID string) (bool, error)
RequestedVideoIDs(ctx context.Context, userID string) (map[string]bool, error)
ClearSummarizeRequested(ctx context.Context, userID, videoID string) error
// RateLimitedVideoIDs maps the user's still-throttled videos to when they were
// rate-limited, so the loop can back off without re-hitting the caption endpoint.
RateLimitedVideoIDs(ctx context.Context, userID string) (map[string]time.Time, error)
// SetTranscriptStatus records the outcome of a transcript attempt: "none",
// "rate_limited" (stamps the backoff clock), or "fetched".
SetTranscriptStatus(ctx context.Context, userID, videoID, status string) error
// UpsertChannelError records a channel that returned HTTP 404 (deleted/private).
// 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
}
// Processor runs the core use case for a single video. *usecase.Engine
// satisfies it.
type Processor interface {
ProcessNewVideo(ctx context.Context, v domain.Video) (usecase.ProcessResult, error)
}
// Runner walks a user's subscriptions, persists each candidate video, skips the
// ones already summarized (durably), and processes the rest through the engine.
type Runner struct {
src ports.VideoSource
store VideoStore
engine Processor
userID string
log *slog.Logger
backoff time.Duration // rate-limit retry window; 0 = always retry
now func() time.Time // injectable clock (tests); defaults to time.Now
}
// Option configures a Runner at construction. Variadic so existing call sites
// stay valid as new knobs (backoff, clock) are added.
type Option func(*Runner)
// WithBackoff sets the rate-limit retry window. A video that returned HTTP 429 is
// skipped (no caption fetch) until this much time has passed; 0 = always retry.
func WithBackoff(d time.Duration) Option { return func(r *Runner) { r.backoff = d } }
// WithClock overrides the clock used for backoff comparisons. Tests inject a
// fixed time; production leaves the time.Now default.
func WithClock(now func() time.Time) Option { return func(r *Runner) { r.now = now } }
// 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 {
log = slog.Default()
}
r := &Runner{src: src, store: store, engine: engine, userID: userID, log: log, now: time.Now}
for _, opt := range opts {
opt(r)
}
if r.now == nil {
r.now = time.Now
}
return r
}
// 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
SkippedRateLimited int // 429'd previously and still inside the backoff window
Errors int
ChannelUnavailable int // channels that returned HTTP 404 (deleted/private)
}
// RunOnce performs a single pass over the user's subscriptions. Per-item errors
// are logged and collected (one bad video or channel does not abort the pass)
// and returned joined alongside the Stats gathered.
func (r *Runner) RunOnce(ctx context.Context) (Stats, error) {
var (
stats Stats
errs []error
)
// Throttle transcript fetches: unauthenticated caption scraping gets
// soft-throttled by YouTube under heavy back-to-back volume (captionTracks
// silently stripped from the player response). A small per-video delay keeps
// a full pass under the radar. TAPIR_FETCH_DELAY (Go duration), 0 = off.
fetchDelay, _ := time.ParseDuration(os.Getenv("TAPIR_FETCH_DELAY"))
seen, err := r.store.SeenVideoIDs(ctx, r.userID)
if err != nil {
return stats, fmt.Errorf("runner: load seen videos: %w", err)
}
// Summarization mode (per-user, ADR-012). Auto = summarize every unseen video
// (the original behavior). Manual = still discover/persist videos so the user
// sees them, but only summarize the ones explicitly queued via the web UI
// (summarize_requested). The queued set is loaded once per pass, like seen.
auto, err := r.store.GetAutoSummarize(ctx, r.userID)
if err != nil {
return stats, fmt.Errorf("runner: load summarize mode: %w", err)
}
var requested map[string]bool
if !auto {
requested, err = r.store.RequestedVideoIDs(ctx, r.userID)
if err != nil {
return stats, fmt.Errorf("runner: load requested videos: %w", err)
}
}
// Rate-limit backoff: videos that 429'd on a prior pass, mapped to when. Inside
// the backoff window they are skipped before any caption fetch, so a throttled
// IP is not hammered. Loaded once per pass (like seen/requested). Disabled when
// backoff <= 0 ("always retry").
var rateLimited map[string]time.Time
if r.backoff > 0 {
rateLimited, err = r.store.RateLimitedVideoIDs(ctx, r.userID)
if err != nil {
return stats, fmt.Errorf("runner: load rate-limited videos: %w", err)
}
}
subs, err := r.src.ListSubscriptions(ctx, r.userID)
if err != nil {
return stats, fmt.Errorf("runner: list subscriptions: %w", err)
}
for _, sub := range subs {
vids, err := r.src.NewVideos(ctx, sub)
if err != nil {
var unavail *domain.ErrChannelUnavailable
if errors.As(err, &unavail) {
stats.ChannelUnavailable++
r.log.Warn("channel unavailable (playlist 404)", "channel", sub.ChannelTitle, "channel_id", sub.ChannelID)
if storeErr := r.store.UpsertChannelError(ctx, r.userID, unavail.ChannelID, unavail.ChannelTitle); storeErr != nil {
r.log.Warn("failed to store channel error", "err", storeErr)
}
} else {
errs = append(errs, fmt.Errorf("new videos for %q: %w", sub.ChannelTitle, err))
stats.Errors++
}
continue
}
for _, v := range vids {
stats.Candidates++
v.UserID = r.userID // keep the dedup/FK key consistent with config
id, err := r.store.UpsertVideo(ctx, v)
if err != nil {
errs = append(errs, fmt.Errorf("upsert video %q: %w", v.ProviderVideoID, err))
stats.Errors++
continue
}
v.ID = id
if seen[id] {
stats.SkippedSeen++
continue
}
seen[id] = true // also guard against the same video within this pass
// Manual mode: skip summarization for videos the user has not queued.
// Discovery already happened (UpsertVideo above), so the new video is
// visible in the list; it just isn't summarized until requested.
if !auto && !requested[id] {
stats.SkippedManual++
continue
}
// Still inside the rate-limit backoff window: skip without fetching, so
// we don't re-hit a caption endpoint that just 429'd us. After the window
// expires the video falls through and is retried normally.
if at, ok := rateLimited[id]; ok && r.now().Sub(at) < r.backoff {
stats.SkippedRateLimited++
r.log.Info("skipped video (rate-limited, backing off)", "video", v.ProviderVideoID, "title", v.Title)
continue
}
if fetchDelay > 0 {
time.Sleep(fetchDelay)
}
res, err := r.engine.ProcessNewVideo(ctx, v)
if err != nil {
errs = append(errs, fmt.Errorf("process %q: %w", v.ProviderVideoID, err))
stats.Errors++
continue
}
switch {
case res.Skipped && res.TranscriptSource == string(domain.SourceRateLimited):
// Fresh 429 this pass: persist rate_limited (stamps the backoff clock)
// so the next pass skips it until the window expires.
stats.SkippedRateLimited++
if err := r.store.SetTranscriptStatus(ctx, r.userID, id, "rate_limited"); err != nil {
errs = append(errs, fmt.Errorf("set rate_limited status %q: %w", v.ProviderVideoID, err))
stats.Errors++
}
r.log.Info("skipped video (rate-limited)", "video", v.ProviderVideoID, "title", v.Title)
case res.Skipped:
stats.SkippedNoText++
if err := r.store.SetTranscriptStatus(ctx, r.userID, id, "none"); err != nil {
errs = append(errs, fmt.Errorf("set none status %q: %w", v.ProviderVideoID, err))
stats.Errors++
}
r.log.Info("skipped video (no transcript)", "video", v.ProviderVideoID, "title", v.Title)
case res.Summary != nil:
stats.Summarized++
if err := r.store.SetTranscriptStatus(ctx, r.userID, id, "fetched"); err != nil {
errs = append(errs, fmt.Errorf("set fetched status %q: %w", v.ProviderVideoID, err))
stats.Errors++
}
// In manual mode the video was processed because it was queued;
// clear the flag so it is not re-summarized and the UI drops the
// "Queued" chip. (Auto mode never sets the flag.)
if !auto {
if err := r.store.ClearSummarizeRequested(ctx, r.userID, id); err != nil {
errs = append(errs, fmt.Errorf("clear summarize flag %q: %w", v.ProviderVideoID, err))
stats.Errors++
}
}
r.log.Info("summarized video", "video", v.ProviderVideoID, "title", v.Title,
"provider", res.Summary.AIProvider, "model", res.Summary.AIModel)
}
}
}
return stats, errors.Join(errs...)
}
// Loop runs RunOnce immediately, then on every interval tick until ctx is
// cancelled. A zero or negative interval means a single pass (no loop). Per-pass
// errors are logged, not fatal, so a transient failure doesn't kill the watcher.
func (r *Runner) Loop(ctx context.Context, interval time.Duration) error {
runPass := func() {
stats, err := r.RunOnce(ctx)
r.log.Info("run pass complete",
"candidates", stats.Candidates, "summarized", stats.Summarized,
"skipped_seen", stats.SkippedSeen, "skipped_no_text", stats.SkippedNoText,
"skipped_manual", stats.SkippedManual, "skipped_rate_limited", stats.SkippedRateLimited,
"channel_unavailable", stats.ChannelUnavailable, "errors", stats.Errors)
if err != nil {
r.log.Warn("run pass had errors", "err", err)
}
}
runPass()
if interval <= 0 {
return nil
}
ticker := time.NewTicker(interval)
defer ticker.Stop()
for {
select {
case <-ctx.Done():
return ctx.Err()
case <-ticker.C:
runPass()
}
}
}