package main import ( "context" "fmt" "net/url" "os" "os/signal" "path/filepath" "strings" "syscall" "time" "rssd/config" "rssd/downloader" "rssd/poller" "rssd/state" "rssd/ytdlp" ) const ( stateFile = "state.json" ) func main() { fmt.Fprintf(os.Stderr, "rssd: starting\n") // Load configuration. cfgPath := os.Getenv("RSSD_CONFIG") if cfgPath == "" { cfgPath = "~/.config/rssd/config.yaml" } cfgPath = config.ExpandPath(cfgPath) // resolve ~ early for consistent path usage cfg, err := config.Load(cfgPath) if err != nil { fmt.Fprintf(os.Stderr, "rssd: fatal: %v\n", err) os.Exit(1) } // Load state — store alongside the config file. stateDir := filepath.Dir(cfgPath) sm, err := state.NewManager(filepath.Join(stateDir, stateFile)) if err != nil { fmt.Fprintf(os.Stderr, "rssd: fatal: %v\n", err) os.Exit(1) } if err := sm.Load(); err != nil { fmt.Fprintf(os.Stderr, "rssd: warning: could not load state (%v), starting fresh\n", err) } // Start download worker pool. jobs := make(chan downloader.DownloadJob) 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, 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 || 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) } }) fmt.Fprintf(os.Stderr, "rssd: workers started, polling feeds\n") sigCh := make(chan os.Signal, 1) signal.Notify(sigCh, syscall.SIGTERM, syscall.SIGINT) fmt.Fprintf(os.Stderr, "rssd: %d feeds configured, polling started\n", len(cfg.Feeds)) // Launch a poller goroutine per feed. for _, feed := range cfg.Feeds { go pollFeed(feed, sm, jobs, cfg) } // Periodically dispatch retrying jobs whose backoff has elapsed. retryTicker := time.NewTicker(1 * time.Minute) defer retryTicker.Stop() go func() { for range retryTicker.C { dispatchRetries(sm, jobs, cfg.Settings.MaxRetries, cfg) } }() // Wait for shutdown signal. sig := <-sigCh fmt.Fprintf(os.Stderr, "rssd: received %v, shutting down...\n", sig) // 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) } fmt.Fprintf(os.Stderr, "rssd: stopped\n") } func pollFeed(feed config.Feed, sm *state.Manager, jobs chan<- downloader.DownloadJob, cfg *config.Config) { 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. pollAndEnqueue(feed, sm, jobs, fetcher, cfg) ticker := time.NewTicker(interval) defer ticker.Stop() for range ticker.C { pollAndEnqueue(feed, sm, jobs, fetcher, cfg) } } // 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 if feedState != nil { lastModified = feedState.LastModified etag = feedState.ETag } result, err := fetcher.Poll(feed.URL, lastModified, etag) if err != nil { fmt.Fprintf(os.Stderr, "rssd: error polling %s: %v\n", feed.URL, err) return } if result.NotModified { fmt.Fprintf(os.Stderr, "rssd: %s unchanged (skipping)\n", feed.URL) return } fmt.Fprintf(os.Stderr, "rssd: %s polled — %d enclosures found\n", feed.URL, len(result.NewEnclosures)) if len(result.NewEnclosures) == 0 { fmt.Fprintf(os.Stderr, "rssd: warning: %s has no enclosures in latest items\n", feed.URL) } // Save updated feed metadata for next conditional request. sm.UpsertFeed(feed.URL, result.LastModified, result.ETag) // 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)) { continue } destPath := computeDestPath(feed.OutputDir, encInfo) 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, }) } } // 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, cfg *config.Config) { // Collect snapshots (lock released before any callback invocation). var ready []downloader.DownloadJob var failed []string sm.ForEachJob(func(snap *state.JobSnapshot) { if snap.Status != state.StatusRetrying { return } if snap.AttemptCount >= maxRetries { // Permanently failed — use UpdateJob to mutate the live state. sm.UpdateJob(snap.EnclosureURL, func(j *state.Job) { j.Status = state.StatusFailed }) failed = append(failed, fmt.Sprintf("%s (after %d attempts: %s)", snap.EnclosureURL, snap.AttemptCount, snap.Error)) return } if snap.NextAttemptAt == nil || time.Now().After(*snap.NextAttemptAt) { // Re-enqueue — use UpdateJob to mutate the live state. sm.UpdateJob(snap.EnclosureURL, func(j *state.Job) { j.Status = state.StatusPending j.NextAttemptAt = nil }) ready = append(ready, jobFromSnapshot(snap, cfg)) } }) if len(failed) > 0 { for _, msg := range failed { fmt.Fprintf(os.Stderr, "rssd: %s permanently failed\n", msg) } } // Send to channel and persist state (lock released). for _, job := range ready { jobs <- job } if err := sm.Save(); err != nil { fmt.Fprintf(os.Stderr, "rssd: error saving state (retry dispatch): %v\n", err) } } // 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( "/", "_", "\\", "_", ":", "_", "<", "_", ">", "_", "|", "_", "\"", "_", "?", "_", "*", "_", ) s = replacer.Replace(s) for strings.Contains(s, "__") { s = strings.ReplaceAll(s, "__", "_") } 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 == "" { // Fallback to enclosure URL basename if title is empty or invalid. parsed, err := url.Parse(encInfo.URL) if err != nil || parsed.Path == "" || parsed.Path == "/" { return filepath.Join(outputDir, strings.ReplaceAll(encInfo.URL, "/", "_")) } baseName = sanitizeFilename(filepath.Base(parsed.Path)) if baseName == "." || baseName == "" { baseName = "unknown" } } // 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") return filepath.Join(outputDir, fmt.Sprintf("%s - %s", datePrefix, baseName)) } return filepath.Join(outputDir, baseName) }