From bc26e73a90da07e0bdffe0080dcced68ef4045a5 Mon Sep 17 00:00:00 2001 From: Greg Pomerantz Date: Wed, 23 Sep 2026 09:06:36 -0400 Subject: [PATCH] Add YouTube feed support via yt-dlp, MIME-based extensions, robust state - ytdlp: thin wrapper around yt-dlp CLI for channel/playlist discovery (--flat-playlist) and per-video metadata (--print-json); defaults to player_client=web_embedded to avoid SABR-only format restrictions. - downloader: YouTube path extracts audio via yt-dlp -x, probes duration with ffprobe and rejects clips far shorter than the expected length. - main: YouTube channel/playlist feeds, incremental discovery with pre-start-date boundary for newest-first channels, two-phase record-then-enqueue so discovered jobs survive crashes. - poller: propagate enclosure MIME type; computeDestPath appends the matching file extension. - state: jobs carry expected duration; AddJobWithStatus for skipped items. - cmd/probe: dry-run discovery/date-resolution diagnostics. --- cmd/probe/main.go | 120 +++++++++++++++++ config/config.go | 7 +- downloader/downloader.go | 186 ++++++++++++++++++++++++++- main.go | 269 +++++++++++++++++++++++++++++++++------ poller/poller.go | 15 ++- state/state.go | 31 ++++- ytdlp/ytdlp.go | 189 +++++++++++++++++++++++++++ 7 files changed, 768 insertions(+), 49 deletions(-) create mode 100644 cmd/probe/main.go create mode 100644 ytdlp/ytdlp.go diff --git a/cmd/probe/main.go b/cmd/probe/main.go new file mode 100644 index 0000000..758bbf7 --- /dev/null +++ b/cmd/probe/main.go @@ -0,0 +1,120 @@ +// Command probe performs a dry run of the YouTube discovery + date-resolution +// logic for one or more channel/playlist URLs. It resolves every video's +// upload date and reports how many fall before/after a start date, without +// downloading anything. +// +// Usage: +// +// probe [url ...] +package main + +import ( + "context" + "fmt" + "os" + "time" + + "rssd/ytdlp" +) + +func main() { + if len(os.Args) < 3 { + fmt.Fprintln(os.Stderr, "usage: probe [url ...]") + os.Exit(2) + } + start, err := time.Parse("2006-01-02", os.Args[1]) + if err != nil { + fmt.Fprintln(os.Stderr, "bad start date:", err) + os.Exit(2) + } + + ctx, cancel := context.WithTimeout(context.Background(), 3*time.Hour) + defer cancel() + s := ytdlp.Settings{} + + for _, u := range os.Args[2:] { + probeOne(ctx, s, u, start) + } +} + +func probeOne(ctx context.Context, s ytdlp.Settings, url string, start time.Time) { + fmt.Printf("\n=== %s ===\n", url) + vids, err := ytdlp.Discovery(ctx, s, url) + if err != nil { + fmt.Fprintf(os.Stderr, " discovery error: %v\n", err) + return + } + fmt.Printf(" discovered %d videos\n", len(vids)) + + newestFirst := true + for _, tok := range []string{"list="} { + if contains(url, tok) { + newestFirst = false + } + } + fmt.Printf(" assumed newest-first: %v\n", newestFirst) + + total := 0 + inRange, beforeStart, noDate, failed := 0, 0, 0, 0 + var oldest, newest time.Time + for _, v := range vids { + total++ + info, err := ytdlp.Info(ctx, s, v.ID) + if err != nil { + failed++ + if failed <= 5 { + fmt.Fprintf(os.Stderr, " resolve warn %s: %v\n", v.ID, err) + } + continue + } + if info.UploadDate.IsZero() { + noDate++ + continue + } + if oldest.IsZero() || info.UploadDate.Before(oldest) { + oldest = info.UploadDate + } + if info.UploadDate.After(newest) { + newest = info.UploadDate + } + if info.UploadDate.Before(start) { + beforeStart++ + if newestFirst { + // Channel is newest-first: everything after this is older, so stop. + fmt.Printf(" reached pre-start-date boundary, stopping\n") + break + } + continue + } + inRange++ + if inRange <= 5 { + fmt.Printf(" IN %s %s (%s)\n", info.UploadDate.Format("2006-01-02"), v.ID, trunc(v.Title, 60)) + } + } + fmt.Printf(" resolved=%d inRange(>=%s)=%d beforeStart=%d noDate=%d failed=%d\n", + total, start.Format("2006-01-02"), inRange, beforeStart, noDate, failed) + fmt.Printf(" oldest=%s newest=%s\n", fmtDate(oldest), fmtDate(newest)) +} + +func contains(s, sub string) bool { + for i := 0; i+len(sub) <= len(s); i++ { + if s[i:i+len(sub)] == sub { + return true + } + } + return false +} + +func trunc(s string, n int) string { + if len(s) <= n { + return s + } + return s[:n] + "…" +} + +func fmtDate(t time.Time) string { + if t.IsZero() { + return "-" + } + return t.Format("2006-01-02") +} diff --git a/config/config.go b/config/config.go index 6170c74..f37020d 100644 --- a/config/config.go +++ b/config/config.go @@ -29,9 +29,9 @@ func expandPath(p string) string { // Feed represents a single podcast feed to monitor. type Feed struct { - URL string `yaml:"url"` - OutputDir string `yaml:"output_dir"` - PollInterval *time.Duration `yaml:"poll_interval,omitempty"` // nil means use global default + URL string `yaml:"url"` + OutputDir string `yaml:"output_dir"` + PollInterval *time.Duration `yaml:"poll_interval,omitempty"` // nil means use global default StartDate string `yaml:"start_date,omitempty"` // optional, format "2006-01-02" } @@ -53,6 +53,7 @@ type Settings struct { MaxRetries int `yaml:"max_retries"` BackoffBase time.Duration `yaml:"backoff_base"` DownloadPoolSize int `yaml:"download_pool_size"` + YTDlpPath string `yaml:"yt_dlp_path,omitempty"` // optional; defaults to "yt-dlp" on PATH } // Config is the top-level parsed configuration. diff --git a/downloader/downloader.go b/downloader/downloader.go index cfc095c..0f95fbb 100644 --- a/downloader/downloader.go +++ b/downloader/downloader.go @@ -1,14 +1,22 @@ package downloader import ( + "bytes" + "context" + "encoding/json" "fmt" "io" + "math" "net/http" "os" + "os/exec" "path/filepath" + "strings" "time" + "rssd/poller" "rssd/state" + "rssd/ytdlp" ) // DownloadJob carries all information needed for a single download. @@ -16,6 +24,12 @@ type DownloadJob struct { FeedURL string EnclosureURL string DestPath string + + // YouTube-only fields. IsYouTube selects the yt-dlp download path. + IsYouTube bool + VideoID string // YouTube video ID + Duration float64 // expected video duration in seconds (0 = unknown) + YTSettings ytdlp.Settings } // Pool manages a fixed-size worker pool for downloading files. @@ -43,7 +57,13 @@ func (p *Pool) Run(jobs <-chan DownloadJob) { func (p *Pool) worker(job DownloadJob) { fmt.Fprintf(os.Stderr, "rssd: worker started download %s\n", job.EnclosureURL) - if err := p.download(job); err != nil { + var err error + if job.IsYouTube { + err = p.downloadYT(job) + } else { + err = p.download(job) + } + if err != nil { fmt.Fprintf(os.Stderr, "ERROR: download %s failed: %v\n", job.EnclosureURL, err) } else { fmt.Fprintf(os.Stderr, "OK: downloaded %s -> %s\n", job.EnclosureURL, job.DestPath) @@ -51,6 +71,146 @@ func (p *Pool) worker(job DownloadJob) { fmt.Fprintf(os.Stderr, "rssd: worker finished download %s\n", job.EnclosureURL) } +// downloadYT downloads the native audio stream of a YouTube video via yt-dlp +// (bestaudio, no -x), i.e. exactly what YouTube serves (typically Opus in +// WebM) with no re-encoding. The file lands at destPath + "."; +// state is updated to the real path. +func (p *Pool) downloadYT(job DownloadJob) error { + // Update state to downloading. + if !p.state.UpdateJob(job.EnclosureURL, func(j *state.Job) { + j.Status = state.StatusDownloading + j.AttemptCount++ + }) { + return fmt.Errorf("job not found for %s", job.EnclosureURL) + } + if err := p.state.Save(); err != nil { + return fmt.Errorf("save state (downloading): %w", err) + } + + // Ensure output directory exists. + if err := os.MkdirAll(filepath.Dir(job.DestPath), 0755); err != nil { + p.markRetry(job.EnclosureURL, fmt.Errorf("create dest dir: %w", err)) + return err + } + + ctx, cancel := context.WithTimeout(context.Background(), 60*time.Minute) + defer cancel() + + // Download the native best audio stream. yt-dlp writes the file to + // destPath with the source extension (e.g. .webm, .m4a) appended. + if err := p.runYT(ctx, job, "-f", "bestaudio", "-o", job.DestPath+".%(ext)s"); err != nil { + p.markRetry(job.EnclosureURL, err) + return err + } + + // Find the actual file: destPath with the real extension appended. + actual, err := findDownloadedFile(job.DestPath) + if err != nil { + p.markRetry(job.EnclosureURL, err) + return err + } + + // Duration sanity check: fail fast on videos that only offer short + // clips (e.g. some shorts where only a partial audio stream exists). + if err := checkDuration(actual, job.Duration); err != nil { + os.Remove(actual) + p.markRetry(job.EnclosureURL, err) + return err + } + + // Mark as downloaded and record the real destination path (extension + // was only known after the download). + p.state.UpdateJob(job.EnclosureURL, func(j *state.Job) { + j.Status = state.StatusDownloaded + j.DestPath = actual + }) + if err := p.state.Save(); err != nil { + return fmt.Errorf("save state (downloaded): %w", err) + } + job.DestPath = actual + + return nil +} + +// runYT executes yt-dlp for a single video, downloading audio only. +func (p *Pool) runYT(ctx context.Context, job DownloadJob, args ...string) error { + cmd := ytdlp.CommandFor(ctx, job.YTSettings) + cmd.Args = append(cmd.Args, args...) + cmd.Args = append(cmd.Args, "--no-playlist", job.EnclosureURL) + + var stderr bytes.Buffer + cmd.Stderr = &stderr + if err := cmd.Run(); err != nil { + return fmt.Errorf("yt-dlp download %s: %v: %s", job.VideoID, err, ytdlp.LastLines(stderr.String(), 3)) + } + return nil +} + +// findDownloadedFile locates the file yt-dlp produced for the given output +// base path. yt-dlp appends the source extension, so we look for files +// starting with base + "." in the same directory (preferring the most +// recently modified match). +func findDownloadedFile(base string) (string, error) { + dir := filepath.Dir(base) + entries, err := os.ReadDir(dir) + if err != nil { + return "", fmt.Errorf("read output dir: %w", err) + } + var best string + var bestMod time.Time + for _, e := range entries { + if e.IsDir() { + continue + } + name := e.Name() + if name == filepath.Base(base) { + return filepath.Join(dir, name), nil + } + if strings.HasPrefix(name, filepath.Base(base)+".") { + if best == "" { + best = filepath.Join(dir, name) + info, err := e.Info() + if err == nil { + bestMod = info.ModTime() + } + continue + } + info, err := e.Info() + if err == nil && info.ModTime().After(bestMod) { + best = filepath.Join(dir, name) + bestMod = info.ModTime() + } + } + } + if best == "" { + return "", fmt.Errorf("downloaded file not found for %s", base) + } + return best, nil +} + +// checkDuration verifies the extracted audio's duration is within tolerance +// of the expected video duration. This catches videos where yt-dlp fell back +// to a shorter clip (e.g. some shorts only offer partial audio). +func checkDuration(path string, expected float64) error { + if expected <= 0 { + return nil // no expectation to check against + } + act, err := ffprobeDuration(path) + if err != nil { + // If we can't probe, don't block the download. + return nil + } + // Allow a 10% tolerance plus a 5s floor for rounding/formatting. + tol := expected * 0.10 + if tol < 5 { + tol = 5 + } + if math.Abs(act-expected) > tol { + return fmt.Errorf("duration mismatch: extracted %.0fs vs expected %.0fs (video may be unavailable in full length)", act, expected) + } + return nil +} + func (p *Pool) download(job DownloadJob) error { // Update state to downloading. if !p.state.UpdateJob(job.EnclosureURL, func(j *state.Job) { @@ -85,6 +245,7 @@ func (p *Pool) download(job DownloadJob) error { p.markRetry(job.EnclosureURL, fmt.Errorf("create request: %w", err)) return err } + req.Header.Set("User-Agent", poller.UserAgent) resp, err := client.Do(req) if err != nil { @@ -126,6 +287,29 @@ func (p *Pool) download(job DownloadJob) error { return nil } +// ffprobeDuration returns the duration of a media file in seconds. +func ffprobeDuration(path string) (float64, error) { + out, err := exec.Command("ffprobe", "-v", "error", + "-show_entries", "format=duration", + "-of", "json", path).Output() + if err != nil { + return 0, fmt.Errorf("ffprobe %s: %w", filepath.Base(path), err) + } + var res struct { + Format struct { + Duration string `json:"duration"` + } `json:"format"` + } + if err := json.Unmarshal(out, &res); err != nil { + return 0, err + } + var d float64 + if _, err := fmt.Sscanf(res.Format.Duration, "%g", &d); err != nil { + return 0, err + } + return d, nil +} + // markRetry updates a job to retrying status with the given error and computes backoff delay. func (p *Pool) markRetry(enclosureURL string, err error) { p.state.UpdateJob(enclosureURL, func(j *state.Job) { diff --git a/main.go b/main.go index 7251d5d..83f1631 100644 --- a/main.go +++ b/main.go @@ -1,6 +1,7 @@ package main import ( + "context" "fmt" "net/url" "os" @@ -14,6 +15,7 @@ import ( "rssd/downloader" "rssd/poller" "rssd/state" + "rssd/ytdlp" ) const ( @@ -51,15 +53,14 @@ func main() { pool := downloader.NewPool(cfg.Settings.DownloadPoolSize, cfg.Settings.BackoffBase, sm) go pool.Run(jobs) - // Re-enqueue any "pending" jobs from a previous run that were never picked up by workers. + // Re-enqueue any "pending" jobs from a previous run that were never picked up + // by workers, and any "downloading" jobs that were interrupted mid-download by + // a crash or forced termination (a stuck "downloading" job would otherwise be + // left alone forever, since the retry dispatch only handles "retrying"). sm.ForEachJob(func(j *state.JobSnapshot) { - if j.Status == state.StatusPending { - fmt.Fprintf(os.Stderr, "rssd: re-enqueued pending job %s\n", j.EnclosureURL) - jobs <- downloader.DownloadJob{ - FeedURL: j.FeedURL, - EnclosureURL: j.EnclosureURL, - DestPath: j.DestPath, - } + if j.Status == state.StatusPending || j.Status == state.StatusDownloading { + fmt.Fprintf(os.Stderr, "rssd: re-enqueued interrupted job %s (was %s)\n", j.EnclosureURL, j.Status) + jobs <- jobFromSnapshot(j, cfg) } }) @@ -79,7 +80,7 @@ func main() { defer retryTicker.Stop() go func() { for range retryTicker.C { - dispatchRetries(sm, jobs, cfg.Settings.MaxRetries) + dispatchRetries(sm, jobs, cfg.Settings.MaxRetries, cfg) } }() @@ -87,10 +88,6 @@ func main() { sig := <-sigCh fmt.Fprintf(os.Stderr, "rssd: received %v, shutting down...\n", sig) - // Close jobs channel to stop accepting new work. Workers will finish in-flight downloads - // but we don't wait — just exit and let the next run handle any remaining state. - close(jobs) - // Save final state before exiting so no progress is lost. if err := sm.Save(); err != nil { fmt.Fprintf(os.Stderr, "rssd: error saving final state: %v\n", err) @@ -99,9 +96,23 @@ func main() { } func pollFeed(feed config.Feed, sm *state.Manager, jobs chan<- downloader.DownloadJob, cfg *config.Config) { - fetcher := poller.NewFetcher() interval := cfg.PollIntervalFor(feed) + if isYouTubeURL(feed.URL) { + fmt.Fprintf(os.Stderr, "rssd: polling YouTube %s every %s\n", feed.URL, interval) + + // Poll immediately on startup. + pollYouTube(feed, sm, jobs, cfg) + + ticker := time.NewTicker(interval) + defer ticker.Stop() + for range ticker.C { + pollYouTube(feed, sm, jobs, cfg) + } + return + } + + fetcher := poller.NewFetcher() fmt.Fprintf(os.Stderr, "rssd: polling %s every %s\n", feed.URL, interval) // Poll immediately on startup. @@ -115,6 +126,151 @@ func pollFeed(feed config.Feed, sm *state.Manager, jobs chan<- downloader.Downlo } } +// isYouTubeURL reports whether the feed URL points at a YouTube channel or +// playlist. Playlist URLs carry their identifier in the query string +// (e.g. /playlist?list=...), so both path and query are checked. +func isYouTubeURL(rawURL string) bool { + u, err := url.Parse(rawURL) + if err != nil { + return false + } + host := strings.ToLower(u.Hostname()) + switch host { + case "youtube.com", "www.youtube.com", "m.youtube.com", "music.youtube.com", "youtu.be": + default: + return false + } + p := strings.ToLower(u.Path) + return strings.HasPrefix(p, "/channel/") || strings.HasPrefix(p, "/user/") || + strings.HasPrefix(p, "/@") || u.Query().Has("list") +} + +// ytSettings builds ytdlp settings from the global config. +func ytSettings(cfg *config.Config) ytdlp.Settings { + return ytdlp.Settings{BinPath: cfg.Settings.YTDlpPath} +} + +// pollYouTube enumerates a YouTube channel/playlist, resolves the date and +// title of any videos not yet tracked in state, and enqueues those that pass +// the start_date filter. Existing jobs (any status) make their videos count +// as "already seen", so discovery is incremental after the first run. +func pollYouTube(feed config.Feed, sm *state.Manager, jobs chan<- downloader.DownloadJob, cfg *config.Config) { + // A first-run backfill can resolve and enqueue many videos, so give the + // whole pass a generous deadline; discovery itself is quick. + ctx, cancel := context.WithTimeout(context.Background(), 2*time.Hour) + defer cancel() + + vids, err := ytdlp.Discovery(ctx, ytSettings(cfg), feed.URL) + if err != nil { + fmt.Fprintf(os.Stderr, "rssd: error discovering %s: %v\n", feed.URL, err) + return + } + + // Collect video IDs already tracked for this feed (any status). + seen := make(map[string]bool, len(vids)) + sm.ForEachJob(func(j *state.JobSnapshot) { + if j.FeedURL == feed.URL { + seen[videoIDFromURL(j.EnclosureURL)] = true + } + }) + + start := feed.PublishedAt() + // Channels enumerate newest-first, so the first video older than the start + // date is a stable boundary: everything after it is older too. Playlists + // have no guaranteed order, so they must be walked fully and their + // out-of-range videos recorded as skipped to avoid re-resolving them. + newestFirst := !strings.Contains(strings.ToLower(feed.URL), "list=") + + newCount, enqueued, skipped := 0, 0, 0 + // Per-video resolve timeout. A single pathological video (JS-challenge + // solving, region/members-only) can make yt-dlp hang; bounding each call + // means one bad video can't stall the whole feed walk. + resolveTimeout := 2 * time.Minute + for _, v := range vids { + if seen[v.ID] { + continue + } + newCount++ + + vctx, vcancel := context.WithTimeout(ctx, resolveTimeout) + info, err := ytdlp.Info(vctx, ytSettings(cfg), v.ID) + vcancel() + if err != nil { + fmt.Fprintf(os.Stderr, "rssd: warning: could not resolve %s: %v\n", v.ID, err) + continue + } + + // Out of the start_date window. + if !start.IsZero() && !info.UploadDate.IsZero() && info.UploadDate.Before(start) { + if newestFirst { + // Stop at the first out-of-range video without recording it. It + // stays "new" each poll, so it remains a stable boundary and we + // never walk backward through older videos. + fmt.Fprintf(os.Stderr, "rssd: %s reached pre-start-date boundary, stopping\n", feed.URL) + break + } + // Playlist: record as skipped so it isn't re-resolved next poll. + dp := computeYTDestPath(feed.OutputDir, info.UploadDate, info.Title, v.ID) + if sm.AddJobWithStatus(feed.URL, watchURL(v.ID), dp, state.StatusSkipped, info.Duration) { + skipped++ + } + continue + } + + dp := computeYTDestPath(feed.OutputDir, info.UploadDate, info.Title, v.ID) + if !sm.AddJob(feed.URL, watchURL(v.ID), dp, cfg.Settings.MaxRetries, info.Duration) { + continue + } + // Persist immediately so a crash mid-backfill doesn't lose this. + if err := sm.Save(); err != nil { + fmt.Fprintf(os.Stderr, "rssd: error saving state: %v\n", err) + } + // Enqueue; downloads run in the worker pool and overlap further + // resolution of later videos. + job := downloader.DownloadJob{ + FeedURL: feed.URL, + EnclosureURL: watchURL(v.ID), + DestPath: dp, + IsYouTube: true, + VideoID: v.ID, + Duration: info.Duration, + YTSettings: ytSettings(cfg), + } + fmt.Fprintf(os.Stderr, "rssd: enqueued %s -> %s\n", job.EnclosureURL, job.DestPath) + jobs <- job + enqueued++ + } + + fmt.Fprintf(os.Stderr, "rssd: %s done — %d new, %d enqueued, %d skipped\n", feed.URL, newCount, enqueued, skipped) +} + +func videoIDFromURL(u string) string { + q := strings.TrimPrefix(u, "https://www.youtube.com/watch?v=") + if i := strings.IndexAny(q, "?&"); i >= 0 { + q = q[:i] + } + return q +} + +func watchURL(id string) string { return "https://www.youtube.com/watch?v=" + id } + +// computeYTDestPath constructs the local file base path for a YouTube +// video: "YYYY-MM-DD - []" (no extension). yt-dlp +// appends the native audio extension (usually .webm) and the downloader +// records the final path in state. The bracketed ID disambiguates +// duplicate titles. +func computeYTDestPath(outputDir string, published time.Time, title, id string) string { + base := sanitizeFilename(title) + if base == "" { + base = "untitled" + } + name := fmt.Sprintf("%s [%s]", base, id) + if !published.IsZero() { + name = fmt.Sprintf("%s - %s", published.Format("2006-01-02"), name) + } + return filepath.Join(outputDir, name) +} + func pollAndEnqueue(feed config.Feed, sm *state.Manager, jobs chan<- downloader.DownloadJob, fetcher *poller.Fetcher, cfg *config.Config) { feedState := sm.FindFeedByURL(feed.URL) var lastModified, etag string @@ -142,43 +298,40 @@ func pollAndEnqueue(feed config.Feed, sm *state.Manager, jobs chan<- downloader. // Save updated feed metadata for next conditional request. sm.UpsertFeed(feed.URL, result.LastModified, result.ETag) - // Enqueue new enclosures (filtered by start_date). + // Phase 1: Record all new enclosures in state first. + var newlyAdded []downloader.DownloadJob startDate := feed.PublishedAt() for _, encInfo := range result.NewEnclosures { // Skip items older than the configured start date. if !startDate.IsZero() && (!encInfo.PublishedAt.IsZero() && encInfo.PublishedAt.Before(startDate)) { - fmt.Fprintf(os.Stderr, "rssd: skipping %s (published %s, before start_date %s)\n", - encInfo.URL, - encInfo.PublishedAt.Format("2006-01-02"), - feed.StartDate) continue } destPath := computeDestPath(feed.OutputDir, encInfo) - if sm.AddJob(feed.URL, encInfo.URL, destPath, cfg.Settings.MaxRetries) { - fmt.Fprintf(os.Stderr, "rssd: enqueued %s -> %s\n", encInfo.URL, destPath) - jobs <- downloader.DownloadJob{ + if sm.AddJob(feed.URL, encInfo.URL, destPath, cfg.Settings.MaxRetries, 0) { + newlyAdded = append(newlyAdded, downloader.DownloadJob{ FeedURL: feed.URL, EnclosureURL: encInfo.URL, DestPath: destPath, - } - } else { - job := sm.FindJobByURL(encInfo.URL) - if job != nil { - fmt.Fprintf(os.Stderr, "rssd: skipping %s (already %s)\n", encInfo.URL, job.Status) - } else { - fmt.Fprintf(os.Stderr, "rssd: skipping duplicate %s\n", encInfo.URL) - } + }) } } + // Commit all newly discovered jobs to disk before attempting to enqueue. + // This ensures they survive a crash even if we block on the channel send. if err := sm.Save(); err != nil { fmt.Fprintf(os.Stderr, "rssd: error saving state: %v\n", err) } + + // Phase 2: Enqueue the jobs to workers. + for _, job := range newlyAdded { + fmt.Fprintf(os.Stderr, "rssd: enqueued %s -> %s\n", job.EnclosureURL, job.DestPath) + jobs <- job + } } // dispatchRetries scans all jobs in state and re-enqueues those whose backoff has elapsed. -func dispatchRetries(sm *state.Manager, jobs chan<- downloader.DownloadJob, maxRetries int) { +func dispatchRetries(sm *state.Manager, jobs chan<- downloader.DownloadJob, maxRetries int, cfg *config.Config) { // Collect snapshots (lock released before any callback invocation). var ready []downloader.DownloadJob var failed []string @@ -204,11 +357,7 @@ func dispatchRetries(sm *state.Manager, jobs chan<- downloader.DownloadJob, maxR j.Status = state.StatusPending j.NextAttemptAt = nil }) - ready = append(ready, downloader.DownloadJob{ - FeedURL: snap.FeedURL, - EnclosureURL: snap.EnclosureURL, - DestPath: snap.DestPath, - }) + ready = append(ready, jobFromSnapshot(snap, cfg)) } }) @@ -228,6 +377,23 @@ func dispatchRetries(sm *state.Manager, jobs chan<- downloader.DownloadJob, maxR } } +// jobFromSnapshot rebuilds a DownloadJob from a persisted state snapshot, +// restoring YouTube-specific fields that aren't part of the on-disk schema. +func jobFromSnapshot(s *state.JobSnapshot, cfg *config.Config) downloader.DownloadJob { + j := downloader.DownloadJob{ + FeedURL: s.FeedURL, + EnclosureURL: s.EnclosureURL, + DestPath: s.DestPath, + Duration: s.Duration, + } + if isYouTubeURL(s.FeedURL) { + j.IsYouTube = true + j.VideoID = videoIDFromURL(s.EnclosureURL) + j.YTSettings = ytSettings(cfg) + } + return j +} + // sanitizeFilename replaces characters that are invalid or problematic in filenames. func sanitizeFilename(s string) string { replacer := strings.NewReplacer( @@ -240,8 +406,34 @@ func sanitizeFilename(s string) string { return strings.TrimSpace(s) } +// extFromMimeType maps common MIME types to their file extensions. +// Returns "" for unknown types. +func extFromMimeType(mime string) string { + switch mime { + case "audio/mpeg": + return ".mp3" + case "audio/mp4", "audio/x-m4a": + return ".m4a" + case "audio/wav": + return ".wav" + case "audio/ogg": + return ".ogg" + case "audio/flac": + return ".flac" + case "video/mp4": + return ".mp4" + case "video/webm": + return ".webm" + case "audio/webm": + return ".webm" + default: + return "" + } +} + // computeDestPath constructs the local file path from an enclosure URL and output directory. // The filename is prefixed with YYYYMMDD - for chronological sorting when dates are available. +// A file extension is appended based on the enclosure's MIME type. func computeDestPath(outputDir string, encInfo poller.EnclosureInfo) string { baseName := sanitizeFilename(encInfo.Title) if baseName == "" { @@ -256,6 +448,11 @@ func computeDestPath(outputDir string, encInfo poller.EnclosureInfo) string { } } + // Append file extension from MIME type. + if ext := extFromMimeType(encInfo.MimeType); ext != "" { + baseName += ext + } + // Prepend date prefix for chronological sorting if available. if !encInfo.PublishedAt.IsZero() { datePrefix := encInfo.PublishedAt.Format("2006-01-02") diff --git a/poller/poller.go b/poller/poller.go index 4df5edc..d3fb27a 100644 --- a/poller/poller.go +++ b/poller/poller.go @@ -12,7 +12,8 @@ import ( // EnclosureInfo carries an enclosure URL and its associated item's metadata. type EnclosureInfo struct { URL string - Title string // episode title, if available + Title string // episode title, if available + MimeType string // MIME type from the enclosure (e.g. "audio/mpeg") PublishedAt time.Time // zero value means no date available } @@ -26,6 +27,11 @@ type FeedResult struct { Error error } +// UserAgent is sent on all outbound HTTP requests (feed fetches and +// enclosure downloads). Some hosts (e.g. Buzzsprout) reject Go's default +// "Go-http-client" UA with 403. +const UserAgent = "Mozilla/5.0 (compatible; rssd/1.0; +https://github.com/rssd)" + // Fetcher abstracts the HTTP fetch + parse cycle for a single feed. type Fetcher struct { httpClient *http.Client @@ -50,6 +56,7 @@ func (f *Fetcher) Poll(url, lastModified, etag string) (*FeedResult, error) { if err != nil { return nil, fmt.Errorf("create request: %w", err) } + req.Header.Set("User-Agent", UserAgent) if etag != "" { req.Header.Set("If-None-Match", etag) @@ -91,7 +98,11 @@ func (f *Fetcher) Poll(url, lastModified, etag string) (*FeedResult, error) { for _, item := range feed.Items { for _, enc := range item.Enclosures { if enc.URL != "" { - encInfo := EnclosureInfo{URL: enc.URL, Title: item.Title} + encInfo := EnclosureInfo{ + URL: enc.URL, + Title: item.Title, + MimeType: enc.Type, + } if item.PublishedParsed != nil { encInfo.PublishedAt = *item.PublishedParsed } diff --git a/state/state.go b/state/state.go index 0096117..6b26a7a 100644 --- a/state/state.go +++ b/state/state.go @@ -32,15 +32,17 @@ type Job struct { MaxRetries int `json:"max_retries"` NextAttemptAt *time.Time `json:"next_attempt_at,omitempty"` // nil means ready now Error string `json:"error,omitempty"` + Duration float64 `json:"duration,omitempty"` // expected media duration in seconds (YouTube) } // JobStatus values. const ( - StatusPending = "pending" + StatusPending = "pending" StatusDownloading = "downloading" - StatusRetrying = "retrying" - StatusDownloaded = "downloaded" - StatusFailed = "failed" + StatusRetrying = "retrying" + StatusDownloaded = "downloaded" + StatusFailed = "failed" + StatusSkipped = "skipped" // resolved but excluded by start_date; never downloaded ) // Manager provides thread-safe access to persisted state. @@ -132,6 +134,7 @@ type JobSnapshot struct { MaxRetries int NextAttemptAt *time.Time // nil means ready now Error string + Duration float64 } // ForEachJob iterates over all jobs and calls fn for each one. @@ -153,6 +156,7 @@ func (m *Manager) ForEachJob(fn func(*JobSnapshot)) { MaxRetries: j.MaxRetries, NextAttemptAt: j.NextAttemptAt, Error: j.Error, + Duration: j.Duration, } } m.mu.Unlock() @@ -170,8 +174,20 @@ func (m *Manager) State() *State { } // AddJob adds a new download job if one with this enclosure URL doesn't already exist. -// Returns true if the job was added, false if it already existed. -func (m *Manager) AddJob(feedURL, enclosureURL, destPath string, maxRetries int) bool { +// duration is the expected media duration in seconds (0 = unknown, used for +// YouTube duration sanity checks). Returns true if the job was added. +func (m *Manager) AddJob(feedURL, enclosureURL, destPath string, maxRetries int, duration float64) bool { + return m.addJob(feedURL, enclosureURL, destPath, StatusPending, maxRetries, duration) +} + +// AddJobWithStatus adds a new job with an explicit initial status (used to +// record videos resolved but excluded by the start_date filter, so they are +// not re-resolved on every poll). Returns true if the job was added. +func (m *Manager) AddJobWithStatus(feedURL, enclosureURL, destPath, status string, duration float64) bool { + return m.addJob(feedURL, enclosureURL, destPath, status, 0, duration) +} + +func (m *Manager) addJob(feedURL, enclosureURL, destPath, status string, maxRetries int, duration float64) bool { m.mu.Lock() defer m.mu.Unlock() @@ -186,11 +202,12 @@ func (m *Manager) AddJob(feedURL, enclosureURL, destPath string, maxRetries int) FeedURL: feedURL, EnclosureURL: enclosureURL, DestPath: destPath, - Status: StatusPending, + Status: status, AttemptCount: 0, MaxRetries: maxRetries, NextAttemptAt: nil, Error: "", + Duration: duration, } m.state.Jobs = append(m.state.Jobs, job) diff --git a/ytdlp/ytdlp.go b/ytdlp/ytdlp.go new file mode 100644 index 0000000..bec96e5 --- /dev/null +++ b/ytdlp/ytdlp.go @@ -0,0 +1,189 @@ +// Package ytdlp provides a thin wrapper around the yt-dlp CLI for +// YouTube channel/playlist discovery and per-video metadata resolution. +package ytdlp + +import ( + "bufio" + "bytes" + "context" + "encoding/json" + "fmt" + "os/exec" + "strings" + "time" +) + +// maxJSONLineBytes bounds one line of --print-json output. A full video's +// JSON (including the formats list) routinely exceeds 1 MB for long videos, +// so the limit must be generous. +const maxJSONLineBytes = 16 * 1024 * 1024 + +// Settings controls how yt-dlp is invoked. +type Settings struct { + // BinPath is the path to the yt-dlp executable (or a name on $PATH). + BinPath string + // PlayerClient is passed via --extractor-args "youtube:player_client=..." + // to work around YouTube client restrictions (see DefaultPlayerClient). + PlayerClient string +} + +// DefaultPlayerClient works around YouTube's SABR-only streaming experiment, +// which withholds plain HTTPS URLs from the default web/android clients. +const DefaultPlayerClient = "web_embedded" + +// FlatVideo is one entry from a --flat-playlist enumeration. +type FlatVideo struct { + ID string `json:"id"` + Title string `json:"title"` + Duration *float64 `json:"duration"` // seconds; may be nil (e.g. some shorts) + Channel string `json:"channel"` +} + +// Discovery enumerates all videos in a YouTube channel or playlist using +// `yt-dlp --flat-playlist` (fast; does not resolve each video). +// Returns videos in the order yt-dlp emits them: +// +// - channel: newest first +// - playlist: playlist order (not necessarily chronological) +func Discovery(ctx context.Context, s Settings, sourceURL string) ([]FlatVideo, error) { + args := []string{"--flat-playlist", "--print-json", sourceURL} + stdout, stderr, err := run(ctx, s, args...) + if err != nil { + return nil, fmt.Errorf("yt-dlp discovery %s: %v: %s", sourceURL, err, LastLines(string(stderr), 3)) + } + + var vids []FlatVideo + sc := bufio.NewScanner(bytes.NewReader(stdout)) + sc.Buffer(make([]byte, 0, 64*1024), maxJSONLineBytes) + for sc.Scan() { + line := strings.TrimSpace(sc.Text()) + if line == "" { + continue + } + var v FlatVideo + if err := json.Unmarshal([]byte(line), &v); err != nil { + continue // skip malformed lines + } + if v.ID == "" { + continue + } + vids = append(vids, v) + } + if len(vids) == 0 { + return nil, fmt.Errorf("yt-dlp discovery %s: no videos found", sourceURL) + } + return vids, nil +} + +// VideoInfo carries resolved metadata for a single video. +type VideoInfo struct { + ID string + Title string + Duration float64 // seconds, 0 if unknown + UploadDate time.Time +} + +// Info resolves full metadata for a single video by ID. +func Info(ctx context.Context, s Settings, videoID string) (VideoInfo, error) { + url := "https://www.youtube.com/watch?v=" + videoID + stdout, stderr, err := run(ctx, s, + "--no-download", "--no-playlist", "--print-json", url) + if err != nil { + return VideoInfo{}, fmt.Errorf("yt-dlp info %s: %v: %s", videoID, err, LastLines(string(stderr), 3)) + } + + sc := bufio.NewScanner(bytes.NewReader(stdout)) + sc.Buffer(make([]byte, 0, 64*1024), maxJSONLineBytes) + var data map[string]any + for sc.Scan() { + line := strings.TrimSpace(sc.Text()) + if line == "" { + continue + } + if err := json.Unmarshal([]byte(line), &data); err == nil { + break + } + } + if data == nil { + if err := sc.Err(); err != nil { + return VideoInfo{}, fmt.Errorf("yt-dlp info %s: %v", videoID, err) + } + return VideoInfo{}, fmt.Errorf("yt-dlp info %s: empty or unparseable output", videoID) + } + + info := VideoInfo{ID: videoID} + if t, _ := data["title"].(string); t != "" { + info.Title = t + } + if d, ok := num(data["duration"]); ok && d > 0 { + info.Duration = d + } + // upload_date is "YYYYMMDD" or "NA" (or null). + if ud, ok := data["upload_date"].(string); ok && ud != "" && ud != "NA" { + if t, err := time.Parse("20060102", ud); err == nil { + info.UploadDate = t + } + } + return info, nil +} + +// CommandFor builds an *exec.Cmd running yt-dlp with the standard +// arguments for the given settings (binary + YouTube extractor args). +// Additional arguments may be appended to cmd.Args by the caller. +func CommandFor(ctx context.Context, s Settings) *exec.Cmd { + bin := s.BinPath + if bin == "" { + bin = "yt-dlp" + } + pc := s.PlayerClient + if pc == "" { + pc = DefaultPlayerClient + } + return exec.CommandContext(ctx, bin, + "--no-warnings", "--extractor-args", "youtube:player_client="+pc) +} + +// run executes yt-dlp with the given arguments and returns stdout and stderr. +func run(ctx context.Context, s Settings, args ...string) ([]byte, []byte, error) { + cmd := CommandFor(ctx, s) + cmd.Args = append(cmd.Args, args...) + var stdout, stderr bytes.Buffer + cmd.Stdout = &stdout + cmd.Stderr = &stderr + if err := cmd.Run(); err != nil { + return stdout.Bytes(), stderr.Bytes(), err + } + return stdout.Bytes(), stderr.Bytes(), nil +} + +// LastLines returns the last n non-empty lines of s, joined with " | ". +// Useful for compact error messages from yt-dlp output. +func LastLines(s string, n int) string { + var lines []string + for _, l := range strings.Split(s, "\n") { + if t := strings.TrimSpace(l); t != "" { + lines = append(lines, t) + } + } + if len(lines) > n { + lines = lines[len(lines)-n:] + } + return strings.Join(lines, " | ") +} + +func num(v any) (float64, bool) { + switch x := v.(type) { + case float64: + return x, true + case float32: + return float64(x), true + case int: + return float64(x), true + case int64: + return float64(x), true + case json.Number: + f, err := x.Float64() + return f, err == nil + } + return 0, false +}