Compare commits
No commits in common. "master" and "main" have entirely different histories.
3
.gitignore
vendored
3
.gitignore
vendored
|
|
@ -1,3 +0,0 @@
|
||||||
rssd
|
|
||||||
log
|
|
||||||
nohup.out
|
|
||||||
64
cmd/debug/main.go
Normal file
64
cmd/debug/main.go
Normal file
|
|
@ -0,0 +1,64 @@
|
||||||
|
package main
|
||||||
|
|
||||||
|
import (
|
||||||
|
"fmt"
|
||||||
|
"io"
|
||||||
|
"net/http"
|
||||||
|
"os"
|
||||||
|
"time"
|
||||||
|
|
||||||
|
"github.com/mmcdole/gofeed"
|
||||||
|
)
|
||||||
|
|
||||||
|
func main() {
|
||||||
|
client := &http.Client{Timeout: 60 * time.Second}
|
||||||
|
resp, err := client.Get("https://feeds.transistor.fm/oxide-and-friends")
|
||||||
|
if err != nil {
|
||||||
|
fmt.Fprintf(os.Stderr, "Error fetching: %v\n", err)
|
||||||
|
os.Exit(1)
|
||||||
|
}
|
||||||
|
defer resp.Body.Close()
|
||||||
|
|
||||||
|
body, err := io.ReadAll(resp.Body)
|
||||||
|
if err != nil {
|
||||||
|
fmt.Fprintf(os.Stderr, "Error reading body: %v\n", err)
|
||||||
|
os.Exit(1)
|
||||||
|
}
|
||||||
|
|
||||||
|
parser := gofeed.NewParser()
|
||||||
|
feed, err := parser.ParseString(string(body))
|
||||||
|
if err != nil {
|
||||||
|
fmt.Fprintf(os.Stderr, "Error parsing feed: %v\n", err)
|
||||||
|
os.Exit(1)
|
||||||
|
}
|
||||||
|
|
||||||
|
fmt.Printf("Feed title: %s\n", feed.Title)
|
||||||
|
fmt.Printf("Total items: %d\n\n", len(feed.Items))
|
||||||
|
|
||||||
|
for i, item := range feed.Items {
|
||||||
|
if i >= 5 {
|
||||||
|
break
|
||||||
|
}
|
||||||
|
fmt.Printf("--- Item %d ---\n", i+1)
|
||||||
|
fmt.Printf(" Title: %s\n", item.Title)
|
||||||
|
fmt.Printf(" Published (raw): %s\n", item.Published)
|
||||||
|
if item.PublishedParsed != nil {
|
||||||
|
fmt.Printf(" PublishedParsed: %v\n", *item.PublishedParsed)
|
||||||
|
} else {
|
||||||
|
fmt.Printf(" PublishedParsed: <nil>\n")
|
||||||
|
}
|
||||||
|
|
||||||
|
// Check for itunes:pubDate or other date fields
|
||||||
|
for _, enc := range item.Enclosures {
|
||||||
|
fmt.Printf(" Enclosure URL (truncated): %s...\n", enc.URL[:min(len(enc.URL), 80)])
|
||||||
|
}
|
||||||
|
fmt.Println()
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func min(a, b int) int {
|
||||||
|
if a < b {
|
||||||
|
return a
|
||||||
|
}
|
||||||
|
return b
|
||||||
|
}
|
||||||
120
cmd/probe/main.go
Normal file
120
cmd/probe/main.go
Normal file
|
|
@ -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 <start-date YYYY-MM-DD> <url> [url ...]
|
||||||
|
package main
|
||||||
|
|
||||||
|
import (
|
||||||
|
"context"
|
||||||
|
"fmt"
|
||||||
|
"os"
|
||||||
|
"time"
|
||||||
|
|
||||||
|
"rssd/ytdlp"
|
||||||
|
)
|
||||||
|
|
||||||
|
func main() {
|
||||||
|
if len(os.Args) < 3 {
|
||||||
|
fmt.Fprintln(os.Stderr, "usage: probe <start-date YYYY-MM-DD> <url> [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")
|
||||||
|
}
|
||||||
148
config/config.go
Normal file
148
config/config.go
Normal file
|
|
@ -0,0 +1,148 @@
|
||||||
|
package config
|
||||||
|
|
||||||
|
import (
|
||||||
|
"fmt"
|
||||||
|
"os"
|
||||||
|
"path/filepath"
|
||||||
|
"strings"
|
||||||
|
"time"
|
||||||
|
|
||||||
|
"gopkg.in/yaml.v3"
|
||||||
|
)
|
||||||
|
|
||||||
|
// ExpandPath resolves a path that may start with ~ to an absolute path.
|
||||||
|
func ExpandPath(p string) string {
|
||||||
|
if strings.HasPrefix(p, "~") {
|
||||||
|
home, err := os.UserHomeDir()
|
||||||
|
if err != nil {
|
||||||
|
return p // fall back to literal on error
|
||||||
|
}
|
||||||
|
return filepath.Join(home, strings.TrimPrefix(p, "~"))
|
||||||
|
}
|
||||||
|
return p
|
||||||
|
}
|
||||||
|
|
||||||
|
// expandPath resolves a path that may start with ~ to an absolute path.
|
||||||
|
func expandPath(p string) string {
|
||||||
|
return ExpandPath(p)
|
||||||
|
}
|
||||||
|
|
||||||
|
// 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
|
||||||
|
StartDate string `yaml:"start_date,omitempty"` // optional, format "2006-01-02"
|
||||||
|
}
|
||||||
|
|
||||||
|
// PublishedAt returns the parsed start date for this feed, or zero time if not set.
|
||||||
|
func (f Feed) PublishedAt() time.Time {
|
||||||
|
if f.StartDate == "" {
|
||||||
|
return time.Time{}
|
||||||
|
}
|
||||||
|
t, err := time.Parse("2006-01-02", f.StartDate)
|
||||||
|
if err != nil {
|
||||||
|
return time.Time{}
|
||||||
|
}
|
||||||
|
return t
|
||||||
|
}
|
||||||
|
|
||||||
|
// Settings holds global defaults for the daemon. Zero values are replaced by ApplyDefaults.
|
||||||
|
type Settings struct {
|
||||||
|
PollInterval time.Duration `yaml:"poll_interval"`
|
||||||
|
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.
|
||||||
|
type Config struct {
|
||||||
|
Settings Settings `yaml:"settings,omitempty"`
|
||||||
|
Feeds []Feed `yaml:"feeds"`
|
||||||
|
}
|
||||||
|
|
||||||
|
// ApplyDefaults fills in zero-valued settings with sensible defaults and expands ~ paths.
|
||||||
|
func (c *Config) ApplyDefaults() {
|
||||||
|
if c.Settings.PollInterval == 0 {
|
||||||
|
c.Settings.PollInterval = 6 * time.Hour
|
||||||
|
}
|
||||||
|
if c.Settings.MaxRetries == 0 {
|
||||||
|
c.Settings.MaxRetries = 5
|
||||||
|
}
|
||||||
|
if c.Settings.BackoffBase == 0 {
|
||||||
|
c.Settings.BackoffBase = 5 * time.Minute
|
||||||
|
}
|
||||||
|
if c.Settings.DownloadPoolSize == 0 {
|
||||||
|
c.Settings.DownloadPoolSize = 4
|
||||||
|
}
|
||||||
|
|
||||||
|
for i := range c.Feeds {
|
||||||
|
c.Feeds[i].OutputDir = expandPath(c.Feeds[i].OutputDir)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// PollIntervalFor returns the effective poll interval for a feed, falling back to global default.
|
||||||
|
func (c *Config) PollIntervalFor(f Feed) time.Duration {
|
||||||
|
if f.PollInterval != nil && *f.PollInterval > 0 {
|
||||||
|
return *f.PollInterval
|
||||||
|
}
|
||||||
|
return c.Settings.PollInterval
|
||||||
|
}
|
||||||
|
|
||||||
|
// Validate checks that the config is usable.
|
||||||
|
func (c *Config) Validate() error {
|
||||||
|
if len(c.Feeds) == 0 {
|
||||||
|
return fmt.Errorf("no feeds defined")
|
||||||
|
}
|
||||||
|
for i, f := range c.Feeds {
|
||||||
|
if f.URL == "" {
|
||||||
|
return fmt.Errorf("feed[%d]: url is required", i)
|
||||||
|
}
|
||||||
|
if f.OutputDir == "" {
|
||||||
|
return fmt.Errorf("feed[%d] (%s): output_dir is required", i, f.URL)
|
||||||
|
}
|
||||||
|
if c.Settings.PollInterval <= 0 {
|
||||||
|
return fmt.Errorf("settings.poll_interval must be > 0")
|
||||||
|
}
|
||||||
|
if c.Settings.MaxRetries == 0 {
|
||||||
|
return fmt.Errorf("settings.max_retries must be > 0")
|
||||||
|
}
|
||||||
|
if c.Settings.BackoffBase <= 0 {
|
||||||
|
return fmt.Errorf("settings.backoff_base must be > 0")
|
||||||
|
}
|
||||||
|
if c.Settings.DownloadPoolSize == 0 {
|
||||||
|
return fmt.Errorf("settings.download_pool_size must be > 0")
|
||||||
|
}
|
||||||
|
// Validate start_date format if provided.
|
||||||
|
if f.StartDate != "" {
|
||||||
|
if _, err := time.Parse("2006-01-02", f.StartDate); err != nil {
|
||||||
|
return fmt.Errorf("feed[%d] (%s): invalid start_date %q: %w", i, f.URL, f.StartDate, err)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
|
||||||
|
// Load reads and parses a YAML config from the given path.
|
||||||
|
func Load(path string) (*Config, error) {
|
||||||
|
path = expandPath(path)
|
||||||
|
|
||||||
|
data, err := os.ReadFile(path)
|
||||||
|
if err != nil {
|
||||||
|
return nil, fmt.Errorf("read config %s: %w", path, err)
|
||||||
|
}
|
||||||
|
|
||||||
|
var cfg Config
|
||||||
|
if err := yaml.Unmarshal(data, &cfg); err != nil {
|
||||||
|
return nil, fmt.Errorf("parse config: %w", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
cfg.ApplyDefaults()
|
||||||
|
|
||||||
|
if err := cfg.Validate(); err != nil {
|
||||||
|
return nil, fmt.Errorf("validate config: %w", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
return &cfg, nil
|
||||||
|
}
|
||||||
420
downloader/downloader.go
Normal file
420
downloader/downloader.go
Normal file
|
|
@ -0,0 +1,420 @@
|
||||||
|
package downloader
|
||||||
|
|
||||||
|
import (
|
||||||
|
"bytes"
|
||||||
|
"context"
|
||||||
|
"encoding/json"
|
||||||
|
"errors"
|
||||||
|
"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.
|
||||||
|
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.
|
||||||
|
type Pool struct {
|
||||||
|
poolSize int
|
||||||
|
backoffBase time.Duration
|
||||||
|
state *state.Manager
|
||||||
|
}
|
||||||
|
|
||||||
|
// NewPool creates a download pool with the given parameters.
|
||||||
|
func NewPool(poolSize int, backoffBase time.Duration, sm *state.Manager) *Pool {
|
||||||
|
return &Pool{poolSize: poolSize, backoffBase: backoffBase, state: sm}
|
||||||
|
}
|
||||||
|
|
||||||
|
// liveWaitInterval is how often a still-live YouTube stream is re-checked.
|
||||||
|
// While a stream is live, the static formats bestaudio requests may not be
|
||||||
|
// available ("Requested format is not available"); once the stream ends,
|
||||||
|
// YouTube converts it to a regular VOD whose formats download normally.
|
||||||
|
const liveWaitInterval = 10 * time.Minute
|
||||||
|
|
||||||
|
// maxLiveWaits bounds the live wait: after this many 10-minute waits
|
||||||
|
// (~4 hours) the job falls back to a regular retry attempt with the
|
||||||
|
// normal exponential backoff, so a video that never converts still ends
|
||||||
|
// in a permanent failure rather than retrying forever.
|
||||||
|
const maxLiveWaits = 24
|
||||||
|
|
||||||
|
// liveStartBuffer is added to an upcoming stream's scheduled start time
|
||||||
|
// before the next check, so the stream has actually begun by then.
|
||||||
|
const liveStartBuffer = 2 * time.Minute
|
||||||
|
|
||||||
|
// errLiveWait is returned by downloadYT when the video is still live or
|
||||||
|
// upcoming and has been rescheduled; it is not a download failure.
|
||||||
|
var errLiveWait = errors.New("video is live or upcoming, rescheduled")
|
||||||
|
|
||||||
|
// Run starts the worker goroutines. It returns immediately after spawning the workers.
|
||||||
|
func (p *Pool) Run(jobs <-chan DownloadJob) {
|
||||||
|
for i := 0; i < p.poolSize; i++ {
|
||||||
|
go func() {
|
||||||
|
for job := range jobs {
|
||||||
|
p.worker(job)
|
||||||
|
}
|
||||||
|
}()
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func (p *Pool) worker(job DownloadJob) {
|
||||||
|
fmt.Fprintf(os.Stderr, "rssd: worker started download %s\n", job.EnclosureURL)
|
||||||
|
var err error
|
||||||
|
if job.IsYouTube {
|
||||||
|
err = p.downloadYT(job)
|
||||||
|
} else {
|
||||||
|
err = p.download(job)
|
||||||
|
}
|
||||||
|
if err != nil {
|
||||||
|
if errors.Is(err, errLiveWait) {
|
||||||
|
fmt.Fprintf(os.Stderr, "rssd: %s — %v\n", job.EnclosureURL, err)
|
||||||
|
} else {
|
||||||
|
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)
|
||||||
|
}
|
||||||
|
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 + ".<actual-ext>";
|
||||||
|
// state is updated to the real path.
|
||||||
|
func (p *Pool) downloadYT(job DownloadJob) error {
|
||||||
|
// While the video is a live stream (or an upcoming one), the static
|
||||||
|
// formats we request are not available yet. Instead of counting a failed
|
||||||
|
// attempt, reschedule the job: an upcoming stream is re-checked at its
|
||||||
|
// scheduled start time, a running one every liveWaitInterval. Once the
|
||||||
|
// stream ends and converts to a regular VOD the normal download succeeds.
|
||||||
|
if wait, start := p.probeLive(job); wait {
|
||||||
|
next := p.markLiveWait(job.EnclosureURL, start)
|
||||||
|
return fmt.Errorf("%w, next check at %s", errLiveWait, next.Format(time.RFC3339))
|
||||||
|
}
|
||||||
|
|
||||||
|
// 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
|
||||||
|
j.LiveWaitCount = 0
|
||||||
|
})
|
||||||
|
if err := p.state.Save(); err != nil {
|
||||||
|
return fmt.Errorf("save state (downloaded): %w", err)
|
||||||
|
}
|
||||||
|
job.DestPath = actual
|
||||||
|
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
|
||||||
|
// probeLive reports whether the video should not be downloaded yet (it is a
|
||||||
|
// running or upcoming live stream), and, for an upcoming stream, its
|
||||||
|
// scheduled start time (zero if unknown). If the probe itself fails, no
|
||||||
|
// wait is returned so the normal download path (with its own error
|
||||||
|
// handling) can proceed.
|
||||||
|
func (p *Pool) probeLive(job DownloadJob) (wait bool, scheduledStart time.Time) {
|
||||||
|
ctx, cancel := context.WithTimeout(context.Background(), 2*time.Minute)
|
||||||
|
defer cancel()
|
||||||
|
info, err := ytdlp.Info(ctx, job.YTSettings, job.VideoID)
|
||||||
|
if err != nil {
|
||||||
|
return false, time.Time{}
|
||||||
|
}
|
||||||
|
switch info.LiveStatus {
|
||||||
|
case "is_upcoming":
|
||||||
|
if info.ReleaseTimestamp > 0 {
|
||||||
|
return true, time.Unix(info.ReleaseTimestamp, 0)
|
||||||
|
}
|
||||||
|
return true, time.Time{}
|
||||||
|
case "is_live":
|
||||||
|
return true, time.Time{}
|
||||||
|
}
|
||||||
|
return false, time.Time{}
|
||||||
|
}
|
||||||
|
|
||||||
|
// markLiveWait reschedules a job whose video is live or upcoming, and
|
||||||
|
// returns the next check time. An upcoming stream with a known start time
|
||||||
|
// is simply rescheduled past its start (liveStartBuffer) without touching
|
||||||
|
// the wait counter, no matter how far away the start is. A running stream
|
||||||
|
// consumes one of the maxLiveWaits 10-minute waits (which do not count
|
||||||
|
// against the job's retry budget); after maxLiveWaits consecutive waits the
|
||||||
|
// next check falls back to a regular retry attempt.
|
||||||
|
func (p *Pool) markLiveWait(enclosureURL string, scheduledStart time.Time) time.Time {
|
||||||
|
var next time.Time
|
||||||
|
p.state.UpdateJob(enclosureURL, func(j *state.Job) {
|
||||||
|
j.Status = state.StatusRetrying
|
||||||
|
now := time.Now()
|
||||||
|
if scheduledStart.After(now) {
|
||||||
|
next = scheduledStart.Add(liveStartBuffer)
|
||||||
|
j.Error = fmt.Sprintf("video is scheduled to start %s, re-checking after start", scheduledStart.Format("2006-01-02 15:04 MST"))
|
||||||
|
j.NextAttemptAt = &next
|
||||||
|
return
|
||||||
|
}
|
||||||
|
j.LiveWaitCount++
|
||||||
|
if j.LiveWaitCount <= maxLiveWaits {
|
||||||
|
next = now.Add(liveWaitInterval)
|
||||||
|
j.Error = fmt.Sprintf("video is live (wait %d/%d), retrying after stream ends", j.LiveWaitCount, maxLiveWaits)
|
||||||
|
} else {
|
||||||
|
j.LiveWaitCount = 0
|
||||||
|
j.AttemptCount++
|
||||||
|
delay := p.backoffBase * time.Duration(1<<uint(j.AttemptCount-1))
|
||||||
|
next = now.Add(delay)
|
||||||
|
j.Error = fmt.Sprintf("video still live after %d waits; falling back to regular retry", maxLiveWaits)
|
||||||
|
}
|
||||||
|
j.NextAttemptAt = &next
|
||||||
|
})
|
||||||
|
p.state.Save()
|
||||||
|
return next
|
||||||
|
}
|
||||||
|
|
||||||
|
// 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) {
|
||||||
|
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
|
||||||
|
}
|
||||||
|
|
||||||
|
// Create destination file.
|
||||||
|
dstFile, err := os.Create(job.DestPath)
|
||||||
|
if err != nil {
|
||||||
|
p.markRetry(job.EnclosureURL, fmt.Errorf("create dest file: %w", err))
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
|
||||||
|
// Download with timeout.
|
||||||
|
client := &http.Client{Timeout: 30 * time.Minute}
|
||||||
|
req, err := http.NewRequest("GET", job.EnclosureURL, nil)
|
||||||
|
if err != nil {
|
||||||
|
dstFile.Close()
|
||||||
|
os.Remove(job.DestPath)
|
||||||
|
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 {
|
||||||
|
dstFile.Close()
|
||||||
|
os.Remove(job.DestPath)
|
||||||
|
p.markRetry(job.EnclosureURL, fmt.Errorf("fetch: %w", err))
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
defer resp.Body.Close()
|
||||||
|
|
||||||
|
if resp.StatusCode != http.StatusOK {
|
||||||
|
dstFile.Close()
|
||||||
|
os.Remove(job.DestPath)
|
||||||
|
p.markRetry(job.EnclosureURL, fmt.Errorf("HTTP %d", resp.StatusCode))
|
||||||
|
return fmt.Errorf("HTTP %d", resp.StatusCode)
|
||||||
|
}
|
||||||
|
|
||||||
|
// Copy body to file.
|
||||||
|
if _, err := io.Copy(dstFile, resp.Body); err != nil {
|
||||||
|
dstFile.Close()
|
||||||
|
os.Remove(job.DestPath)
|
||||||
|
p.markRetry(job.EnclosureURL, fmt.Errorf("write file: %w", err))
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
if err := dstFile.Close(); err != nil {
|
||||||
|
os.Remove(job.DestPath)
|
||||||
|
p.markRetry(job.EnclosureURL, fmt.Errorf("close file: %w", err))
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
|
||||||
|
// Mark as downloaded.
|
||||||
|
p.state.UpdateJob(job.EnclosureURL, func(j *state.Job) {
|
||||||
|
j.Status = state.StatusDownloaded
|
||||||
|
})
|
||||||
|
if err := p.state.Save(); err != nil {
|
||||||
|
return fmt.Errorf("save state (downloaded): %w", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
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) {
|
||||||
|
j.Status = state.StatusRetrying
|
||||||
|
j.LiveWaitCount = 0
|
||||||
|
j.Error = err.Error()
|
||||||
|
// Exponential backoff: base * 2^attemptCount (attempt already incremented).
|
||||||
|
delay := p.backoffBase * time.Duration(1<<uint(j.AttemptCount-1))
|
||||||
|
t := time.Now().Add(delay)
|
||||||
|
j.NextAttemptAt = &t
|
||||||
|
})
|
||||||
|
p.state.Save()
|
||||||
|
}
|
||||||
16
go.mod
Normal file
16
go.mod
Normal file
|
|
@ -0,0 +1,16 @@
|
||||||
|
module rssd
|
||||||
|
|
||||||
|
go 1.24.2
|
||||||
|
|
||||||
|
require (
|
||||||
|
github.com/PuerkitoBio/goquery v1.8.0 // indirect
|
||||||
|
github.com/andybalholm/cascadia v1.3.1 // indirect
|
||||||
|
github.com/json-iterator/go v1.1.12 // indirect
|
||||||
|
github.com/mmcdole/gofeed v1.3.0 // indirect
|
||||||
|
github.com/mmcdole/goxpp v1.1.1-0.20240225020742-a0c311522b23 // indirect
|
||||||
|
github.com/modern-go/concurrent v0.0.0-20180306012644-bacd9c7ef1dd // indirect
|
||||||
|
github.com/modern-go/reflect2 v1.0.2 // indirect
|
||||||
|
golang.org/x/net v0.4.0 // indirect
|
||||||
|
golang.org/x/text v0.5.0 // indirect
|
||||||
|
gopkg.in/yaml.v3 v3.0.1 // indirect
|
||||||
|
)
|
||||||
34
go.sum
Normal file
34
go.sum
Normal file
|
|
@ -0,0 +1,34 @@
|
||||||
|
github.com/PuerkitoBio/goquery v1.8.0 h1:PJTF7AmFCFKk1N6V6jmKfrNH9tV5pNE6lZMkG0gta/U=
|
||||||
|
github.com/PuerkitoBio/goquery v1.8.0/go.mod h1:ypIiRMtY7COPGk+I/YbZLbxsxn9g5ejnI2HSMtkjZvI=
|
||||||
|
github.com/andybalholm/cascadia v1.3.1 h1:nhxRkql1kdYCc8Snf7D5/D3spOX+dBgjA6u8x004T2c=
|
||||||
|
github.com/andybalholm/cascadia v1.3.1/go.mod h1:R4bJ1UQfqADjvDa4P6HZHLh/3OxWWEqc0Sk8XGwHqvA=
|
||||||
|
github.com/davecgh/go-spew v1.1.0/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38=
|
||||||
|
github.com/davecgh/go-spew v1.1.1/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38=
|
||||||
|
github.com/google/gofuzz v1.0.0/go.mod h1:dBl0BpW6vV/+mYPU4Po3pmUjxk6FQPldtuIdl/M65Eg=
|
||||||
|
github.com/json-iterator/go v1.1.12 h1:PV8peI4a0ysnczrg+LtxykD8LfKY9ML6u2jnxaEnrnM=
|
||||||
|
github.com/json-iterator/go v1.1.12/go.mod h1:e30LSqwooZae/UwlEbR2852Gd8hjQvJoHmT4TnhNGBo=
|
||||||
|
github.com/mmcdole/gofeed v1.3.0 h1:5yn+HeqlcvjMeAI4gu6T+crm7d0anY85+M+v6fIFNG4=
|
||||||
|
github.com/mmcdole/gofeed v1.3.0/go.mod h1:9TGv2LcJhdXePDzxiuMnukhV2/zb6VtnZt1mS+SjkLE=
|
||||||
|
github.com/mmcdole/goxpp v1.1.1-0.20240225020742-a0c311522b23 h1:Zr92CAlFhy2gL+V1F+EyIuzbQNbSgP4xhTODZtrXUtk=
|
||||||
|
github.com/mmcdole/goxpp v1.1.1-0.20240225020742-a0c311522b23/go.mod h1:v+25+lT2ViuQ7mVxcncQ8ch1URund48oH+jhjiwEgS8=
|
||||||
|
github.com/modern-go/concurrent v0.0.0-20180228061459-e0a39a4cb421/go.mod h1:6dJC0mAP4ikYIbvyc7fijjWJddQyLn8Ig3JB5CqoB9Q=
|
||||||
|
github.com/modern-go/concurrent v0.0.0-20180306012644-bacd9c7ef1dd h1:TRLaZ9cD/w8PVh93nsPXa1VrQ6jlwL5oN8l14QlcNfg=
|
||||||
|
github.com/modern-go/concurrent v0.0.0-20180306012644-bacd9c7ef1dd/go.mod h1:6dJC0mAP4ikYIbvyc7fijjWJddQyLn8Ig3JB5CqoB9Q=
|
||||||
|
github.com/modern-go/reflect2 v1.0.2 h1:xBagoLtFs94CBntxluKeaWgTMpvLxC4ur3nMaC9Gz0M=
|
||||||
|
github.com/modern-go/reflect2 v1.0.2/go.mod h1:yWuevngMOJpCy52FWWMvUC8ws7m/LJsjYzDa0/r8luk=
|
||||||
|
github.com/pmezard/go-difflib v1.0.0/go.mod h1:iKH77koFhYxTK1pcRnkKkqfTogsbg7gZNVY4sRDYZ/4=
|
||||||
|
github.com/stretchr/objx v0.1.0/go.mod h1:HFkY916IF+rwdDfMAkV7OtwuqBVzrE8GR6GFx+wExME=
|
||||||
|
github.com/stretchr/testify v1.3.0/go.mod h1:M5WIy9Dh21IEIfnGCwXGc5bZfKNJtfHm1UVUgZn+9EI=
|
||||||
|
golang.org/x/net v0.0.0-20210916014120-12bc252f5db8/go.mod h1:9nx3DQGgdP8bBQD5qxJ1jj9UTztislL4KSBs9R2vV5Y=
|
||||||
|
golang.org/x/net v0.4.0 h1:Q5QPcMlvfxFTAPV0+07Xz/MpK9NTXu2VDUuy0FeMfaU=
|
||||||
|
golang.org/x/net v0.4.0/go.mod h1:MBQ8lrhLObU/6UmLb4fmbmk5OcyYmqtbGd/9yIeKjEE=
|
||||||
|
golang.org/x/sys v0.0.0-20201119102817-f84b799fce68/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs=
|
||||||
|
golang.org/x/sys v0.0.0-20210423082822-04245dca01da/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs=
|
||||||
|
golang.org/x/term v0.0.0-20201126162022-7de9c90e9dd1/go.mod h1:bj7SfCRtBDWHUb9snDiAeCFNEtKQo2Wmx5Cou7ajbmo=
|
||||||
|
golang.org/x/text v0.3.6/go.mod h1:5Zoc/QRtKVWzQhOtBMvqHzDpF6irO9z98xDceosuGiQ=
|
||||||
|
golang.org/x/text v0.5.0 h1:OLmvp0KP+FVG99Ct/qFiL/Fhk4zp4QQnZ7b2U+5piUM=
|
||||||
|
golang.org/x/text v0.5.0/go.mod h1:mrYo+phRRbMaCq/xk9113O4dZlRixOauAjOtrjsXDZ8=
|
||||||
|
golang.org/x/tools v0.0.0-20180917221912-90fa682c2a6e/go.mod h1:n7NCudcB/nEzxVGmLbDWY5pfWTLqBcC2KZ6jyYvM4mQ=
|
||||||
|
gopkg.in/check.v1 v0.0.0-20161208181325-20d25e280405/go.mod h1:Co6ibVJAznAaIkqp8huTwlJQCZ016jof/cbN4VW5Yz0=
|
||||||
|
gopkg.in/yaml.v3 v3.0.1 h1:fxVm/GzAzEWqLHuvctI91KS9hhNmmWOoWu0XTYJS7CA=
|
||||||
|
gopkg.in/yaml.v3 v3.0.1/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM=
|
||||||
922
main.go
922
main.go
|
|
@ -1,495 +1,463 @@
|
||||||
package main
|
package main
|
||||||
|
|
||||||
import (
|
import (
|
||||||
"bufio"
|
"context"
|
||||||
"bytes"
|
|
||||||
"fmt"
|
"fmt"
|
||||||
"io"
|
|
||||||
"log"
|
|
||||||
"net/url"
|
"net/url"
|
||||||
"os"
|
"os"
|
||||||
"path"
|
"os/signal"
|
||||||
"sort"
|
"path/filepath"
|
||||||
"strconv"
|
"strings"
|
||||||
"sync"
|
"syscall"
|
||||||
"time"
|
"time"
|
||||||
|
|
||||||
"github.com/BurntSushi/toml"
|
"rssd/config"
|
||||||
"github.com/cavaliercoder/grab"
|
"rssd/downloader"
|
||||||
homedir "github.com/mitchellh/go-homedir"
|
"rssd/poller"
|
||||||
"github.com/mmcdole/gofeed"
|
"rssd/state"
|
||||||
|
"rssd/ytdlp"
|
||||||
)
|
)
|
||||||
|
|
||||||
var (
|
const (
|
||||||
confFile,dataFile,queueFile,dstDir string
|
stateFile = "state.json"
|
||||||
)
|
)
|
||||||
|
|
||||||
func init() {
|
|
||||||
os.Chdir(path.Join()) // go to the root directory
|
|
||||||
homeDir,err := homedir.Dir()
|
|
||||||
if err != nil {
|
|
||||||
log.Fatal("Cannot locate user's home directory")
|
|
||||||
}
|
|
||||||
confDir := path.Join(homeDir,".config","rssd")
|
|
||||||
confFile = path.Join(confDir,"rssd.conf")
|
|
||||||
dataFile = path.Join(confDir,"podcasts.conf")
|
|
||||||
queueFile = path.Join(confDir,"queue.conf")
|
|
||||||
}
|
|
||||||
|
|
||||||
type Config struct {
|
|
||||||
Workers int
|
|
||||||
DestDir string
|
|
||||||
Urls []string
|
|
||||||
}
|
|
||||||
|
|
||||||
type Item struct {
|
|
||||||
Title, Description, Url, Filename string
|
|
||||||
Length int
|
|
||||||
Published time.Time
|
|
||||||
Podcast *Podcast
|
|
||||||
}
|
|
||||||
|
|
||||||
type Podcast struct {
|
|
||||||
Title, Description, Url string
|
|
||||||
Items []Item
|
|
||||||
}
|
|
||||||
|
|
||||||
type pcList struct {
|
|
||||||
Podcasts []Podcast
|
|
||||||
}
|
|
||||||
|
|
||||||
func newpcList(confs ...string) (ret *pcList) {
|
|
||||||
ret = &pcList{}
|
|
||||||
ret.Podcasts = make([]Podcast,0)
|
|
||||||
if len(confs) > 0 {
|
|
||||||
if _, err := toml.DecodeFile(dataFile, &ret); err != nil {
|
|
||||||
log.Print("Error reading podcast list:",err)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
return
|
|
||||||
}
|
|
||||||
|
|
||||||
func (p *pcList) Find(x *Podcast) (int, bool) {
|
|
||||||
for i,y := range p.Podcasts {
|
|
||||||
if y.Title == x.Title {
|
|
||||||
return i, true
|
|
||||||
}
|
|
||||||
}
|
|
||||||
return 0, false
|
|
||||||
}
|
|
||||||
|
|
||||||
func (p *pcList) Add(x *Podcast) {
|
|
||||||
if x == nil {
|
|
||||||
return
|
|
||||||
}
|
|
||||||
if i,ok := p.Find(x); ok == true {
|
|
||||||
log.Print(" Existing podcast")
|
|
||||||
p.Podcasts[i].Merge(x)
|
|
||||||
} else {
|
|
||||||
log.Print(" New podcast")
|
|
||||||
p.Podcasts = append((*p).Podcasts,*x)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
func (p *Podcast) Merge(x *Podcast) {
|
|
||||||
for _,item := range x.Items {
|
|
||||||
if !p.Has(item) {
|
|
||||||
p.Items = append(p.Items,item)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
func (p *Podcast) Has(i Item) bool {
|
|
||||||
for _,x := range p.Items {
|
|
||||||
if x.Title == i.Title {
|
|
||||||
return true
|
|
||||||
}
|
|
||||||
}
|
|
||||||
return false
|
|
||||||
}
|
|
||||||
|
|
||||||
type Selector func(*gofeed.Item) bool
|
|
||||||
|
|
||||||
func AllSelectors(ss ...Selector) Selector {
|
|
||||||
return func(i *gofeed.Item) bool {
|
|
||||||
for _, s := range ss {
|
|
||||||
if !s(i) {
|
|
||||||
return false
|
|
||||||
}
|
|
||||||
}
|
|
||||||
return true
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
func AnySelector(ss ...Selector) Selector {
|
|
||||||
return func(i *gofeed.Item) bool {
|
|
||||||
for _, s := range ss {
|
|
||||||
if s(i) {
|
|
||||||
return true
|
|
||||||
}
|
|
||||||
}
|
|
||||||
return false
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
func newerThan(t time.Time) Selector {
|
|
||||||
return func(i *gofeed.Item) bool {
|
|
||||||
if i.PublishedParsed.After(t) {
|
|
||||||
return true
|
|
||||||
}
|
|
||||||
return false
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
func daysAgo(x int) Selector {
|
|
||||||
d := time.Now()
|
|
||||||
return newerThan(time.Date(d.Year(),d.Month(),d.Day()-x,0,0,0,0,time.Local))
|
|
||||||
}
|
|
||||||
|
|
||||||
func toPodcast(sel Selector, u string, feed *gofeed.Feed) (ret *Podcast) {
|
|
||||||
ret = &Podcast{
|
|
||||||
Title: feed.Title,
|
|
||||||
Description: feed.Description,
|
|
||||||
Url: u,
|
|
||||||
Items: []Item{},
|
|
||||||
}
|
|
||||||
for _, i := range feed.Items {
|
|
||||||
if sel(i) {
|
|
||||||
fn := i.PublishedParsed.Format("20060102--") + i.Title
|
|
||||||
it := Item{
|
|
||||||
Title: i.Title,
|
|
||||||
Description: i.Description,
|
|
||||||
Filename: path.Join(ret.Title, fn),
|
|
||||||
Published: *i.PublishedParsed,
|
|
||||||
}
|
|
||||||
for _, n := range i.Enclosures {
|
|
||||||
if n.Type == "audio/mpeg" {
|
|
||||||
u,err := url.Parse(n.URL)
|
|
||||||
if err != nil {
|
|
||||||
it.Url = n.URL
|
|
||||||
} else {
|
|
||||||
it.Url = fmt.Sprintf("%s://%s%s",u.Scheme,u.Host,u.Path)
|
|
||||||
}
|
|
||||||
if l, err := strconv.Atoi(n.Length); err == nil {
|
|
||||||
it.Length = l
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
ret.Items = append(ret.Items,it)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
return
|
|
||||||
}
|
|
||||||
|
|
||||||
func readFeed(u string,sel Selector) *Podcast {
|
|
||||||
fp := gofeed.NewParser()
|
|
||||||
feed, err := fp.ParseURL(u)
|
|
||||||
if err != nil {
|
|
||||||
log.Print(err)
|
|
||||||
return nil
|
|
||||||
}
|
|
||||||
return toPodcast(sel,u,feed)
|
|
||||||
}
|
|
||||||
|
|
||||||
type dlItem struct {
|
|
||||||
Item *Item
|
|
||||||
Filename string
|
|
||||||
Downloading bool
|
|
||||||
Complete bool
|
|
||||||
}
|
|
||||||
|
|
||||||
type ByDate []*dlItem
|
|
||||||
func (a ByDate) Len() int { return len(a) }
|
|
||||||
func (a ByDate) Swap(i, j int) { a[i], a[j] = a[j], a[i] }
|
|
||||||
func (a ByDate) Less(i, j int) bool { return a[j].Item.Published.After(a[i].Item.Published) }
|
|
||||||
|
|
||||||
type dlQueue struct {
|
|
||||||
Items []*dlItem
|
|
||||||
reqch chan *grab.Request
|
|
||||||
respch chan *grab.Response
|
|
||||||
sync.Mutex
|
|
||||||
wg sync.WaitGroup
|
|
||||||
wake chan struct{}
|
|
||||||
}
|
|
||||||
|
|
||||||
func (q *dlQueue) Sort() {
|
|
||||||
sort.Sort(ByDate(q.Items))
|
|
||||||
}
|
|
||||||
|
|
||||||
func (q *dlQueue) Find(x *Item) (int, bool) {
|
|
||||||
for i,y := range q.Items {
|
|
||||||
if y.Item.Title == x.Title {
|
|
||||||
return i, true
|
|
||||||
}
|
|
||||||
}
|
|
||||||
return 0, false
|
|
||||||
}
|
|
||||||
|
|
||||||
func (q *dlQueue) FindFilename(x string) (int, bool) {
|
|
||||||
for i,y := range q.Items {
|
|
||||||
if y.Filename == x {
|
|
||||||
return i, true
|
|
||||||
}
|
|
||||||
}
|
|
||||||
return 0, false
|
|
||||||
}
|
|
||||||
|
|
||||||
func (q *dlQueue) Waiting() *dlQueue {
|
|
||||||
ret := &dlQueue{
|
|
||||||
Items: make([]*dlItem,0),
|
|
||||||
reqch: q.reqch,
|
|
||||||
respch: q.respch,
|
|
||||||
Mutex: q.Mutex,
|
|
||||||
wg: q.wg,
|
|
||||||
wake: q.wake,
|
|
||||||
}
|
|
||||||
for _,i := range q.Items {
|
|
||||||
if i.Downloading == false {
|
|
||||||
ret.Items = append(ret.Items,i)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
return ret
|
|
||||||
}
|
|
||||||
|
|
||||||
func (q *dlQueue) Add(i *Item) {
|
|
||||||
if _, ok := q.Find(i); ok == true {
|
|
||||||
return
|
|
||||||
}
|
|
||||||
di := &dlItem{Item: &Item{} }
|
|
||||||
*di.Item = *i
|
|
||||||
di.Filename = path.Join(dstDir,i.Filename) + path.Ext(i.Url)
|
|
||||||
q.Items = append(q.Items,di)
|
|
||||||
}
|
|
||||||
|
|
||||||
func NewQueue() *dlQueue {
|
|
||||||
ret := &dlQueue{
|
|
||||||
Items: make([]*dlItem,0),
|
|
||||||
reqch: make(chan *grab.Request),
|
|
||||||
respch: make(chan *grab.Response),
|
|
||||||
wake: make(chan struct{}),
|
|
||||||
}
|
|
||||||
return ret
|
|
||||||
}
|
|
||||||
|
|
||||||
func (q *dlQueue) Load() {
|
|
||||||
q.Lock()
|
|
||||||
if _, err := toml.DecodeFile(queueFile, q); err != nil {
|
|
||||||
log.Print("Cannot read queue status file:",err)
|
|
||||||
}
|
|
||||||
q.Unlock()
|
|
||||||
}
|
|
||||||
|
|
||||||
func (q *dlQueue) Save() {
|
|
||||||
of,err := os.Create(queueFile)
|
|
||||||
if err != nil {
|
|
||||||
log.Print("dlQueue.Save(): Cannot open output file")
|
|
||||||
} else {
|
|
||||||
var buf bytes.Buffer
|
|
||||||
w := bufio.NewWriter(&buf)
|
|
||||||
enc := toml.NewEncoder(w)
|
|
||||||
log.Printf("dlQueue.Save(): encoding %d entries\n",len(q.Items))
|
|
||||||
q.Lock() // do not lock around IO
|
|
||||||
enc.Encode(q)
|
|
||||||
q.Unlock()
|
|
||||||
io.Copy(of,&buf)
|
|
||||||
of.Close()
|
|
||||||
log.Print("dlQueue.Save(): done")
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
type Daemon struct {
|
|
||||||
conf Config
|
|
||||||
g *grab.Client
|
|
||||||
pl *pcList
|
|
||||||
queue *dlQueue
|
|
||||||
sync.Mutex
|
|
||||||
workers int
|
|
||||||
dlwake chan struct{}
|
|
||||||
}
|
|
||||||
|
|
||||||
func NewDaemon(conf Config, pl *pcList) *Daemon {
|
|
||||||
ret := &Daemon{
|
|
||||||
conf: conf,
|
|
||||||
g: grab.NewClient(),
|
|
||||||
pl: pl,
|
|
||||||
queue: NewQueue(),
|
|
||||||
dlwake: make(chan struct{}),
|
|
||||||
}
|
|
||||||
ret.queue.Load()
|
|
||||||
return ret
|
|
||||||
}
|
|
||||||
|
|
||||||
func (d *Daemon) Update(urls []string) {
|
|
||||||
sel := daysAgo(30)
|
|
||||||
for _,url := range urls {
|
|
||||||
log.Print(" -> ",url)
|
|
||||||
f := readFeed(url,sel) // do not lock around IO
|
|
||||||
d.Lock()
|
|
||||||
d.pl.Add(f)
|
|
||||||
d.Unlock()
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
func (d *Daemon) Monitor() {
|
|
||||||
status := func(resp *grab.Response) {
|
|
||||||
log.Printf(" %s: %v bytes (%.2f%%)\n",
|
|
||||||
resp.Filename,
|
|
||||||
resp.BytesComplete(),
|
|
||||||
100*resp.Progress())
|
|
||||||
}
|
|
||||||
mon := func(resp *grab.Response) {
|
|
||||||
t := time.NewTicker(5 * time.Second)
|
|
||||||
defer t.Stop()
|
|
||||||
Loop:
|
|
||||||
for {
|
|
||||||
select {
|
|
||||||
case <-t.C:
|
|
||||||
status(resp)
|
|
||||||
case <-resp.Done:
|
|
||||||
status(resp)
|
|
||||||
d.queue.Lock()
|
|
||||||
if i, ok := d.queue.FindFilename(resp.Filename); ok {
|
|
||||||
d.queue.Items[i].Complete = true
|
|
||||||
d.queue.Unlock()
|
|
||||||
d.queue.Save()
|
|
||||||
} else {
|
|
||||||
d.queue.Unlock()
|
|
||||||
log.Printf("Error finding queue entry for %s\n",resp.Filename)
|
|
||||||
}
|
|
||||||
break Loop
|
|
||||||
}
|
|
||||||
}
|
|
||||||
if err := resp.Err(); err != nil {
|
|
||||||
log.Printf("Download failed for %s (%s)\n",resp.Filename,err)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
for r := range d.queue.respch {
|
|
||||||
go mon(r)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
func (d *Daemon) StartDownloader() {
|
|
||||||
log.Print("Downloader(): spawning workers")
|
|
||||||
for i := 0; i < d.workers; i++ {
|
|
||||||
d.queue.wg.Add(1)
|
|
||||||
go func() {
|
|
||||||
d.g.DoChannel(d.queue.reqch,d.queue.respch)
|
|
||||||
d.queue.wg.Done()
|
|
||||||
}()
|
|
||||||
}
|
|
||||||
log.Print("Downloader(): starting monitor")
|
|
||||||
go d.Monitor()
|
|
||||||
go d.QueueUpdater()
|
|
||||||
go d.queue.Downloader()
|
|
||||||
}
|
|
||||||
|
|
||||||
func (d *Daemon) QueueUpdater() {
|
|
||||||
t := time.NewTicker(30 * time.Minute)
|
|
||||||
defer t.Stop()
|
|
||||||
for {
|
|
||||||
// lock Podcast list and update download queue
|
|
||||||
log.Print("QueueUpdater(): Updating download queue")
|
|
||||||
d.queue.Lock()
|
|
||||||
for _,p := range d.pl.Podcasts {
|
|
||||||
for _,i := range p.Items {
|
|
||||||
d.queue.Add(&i)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
log.Print("QueueUpdater(): Done updating download queue")
|
|
||||||
d.queue.Unlock()
|
|
||||||
d.queue.Save()
|
|
||||||
d.queue.wake<- struct{}{}
|
|
||||||
log.Print("QueueUpdater(): Sleeping")
|
|
||||||
select {
|
|
||||||
case <-t.C:
|
|
||||||
continue
|
|
||||||
case <-d.dlwake:
|
|
||||||
continue
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
func (q *dlQueue) Downloader() {
|
|
||||||
t := time.NewTicker(30 * time.Minute)
|
|
||||||
defer t.Stop()
|
|
||||||
LOOP:
|
|
||||||
for {
|
|
||||||
// launch requests for files we are not yet downloading
|
|
||||||
q.Lock()
|
|
||||||
waiting := q.Waiting()
|
|
||||||
q.Unlock()
|
|
||||||
log.Print("Download queue length: ",len(waiting.Items))
|
|
||||||
for _,i := range waiting.Items {
|
|
||||||
if !i.Downloading && !i.Complete {
|
|
||||||
req,err := grab.NewRequest(i.Filename,i.Item.Url)
|
|
||||||
if err != nil {
|
|
||||||
log.Print("Request error: ",err)
|
|
||||||
continue
|
|
||||||
}
|
|
||||||
i.Downloading = true
|
|
||||||
t := time.Now()
|
|
||||||
q.reqch <- req
|
|
||||||
if time.Now().After(t.Add(5 * time.Second)) {
|
|
||||||
continue LOOP // refresh list
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
log.Print("Downloader(): sleeping")
|
|
||||||
select {
|
|
||||||
case <-t.C:
|
|
||||||
continue
|
|
||||||
case <-q.wake:
|
|
||||||
continue
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
func (d *Daemon) Updater() {
|
|
||||||
log.Print("Updater(): starting")
|
|
||||||
d.StartDownloader()
|
|
||||||
for {
|
|
||||||
time.Sleep(1 * time.Second)
|
|
||||||
d.Update(d.conf.Urls)
|
|
||||||
of,err := os.Create(dataFile)
|
|
||||||
if err != nil {
|
|
||||||
log.Print("Updater(): Cannot open output file")
|
|
||||||
} else {
|
|
||||||
enc := toml.NewEncoder(of)
|
|
||||||
log.Print("Updater(): writing output")
|
|
||||||
d.Lock()
|
|
||||||
enc.Encode(d.pl)
|
|
||||||
d.Unlock()
|
|
||||||
of.Close()
|
|
||||||
}
|
|
||||||
d.dlwake <-struct{}{}
|
|
||||||
time.Sleep(30 * time.Minute)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
func (d *Daemon) Start() {
|
|
||||||
go d.Updater()
|
|
||||||
}
|
|
||||||
|
|
||||||
func main() {
|
func main() {
|
||||||
log.Print("rssd")
|
fmt.Fprintf(os.Stderr, "rssd: starting\n")
|
||||||
log.Print("reading configuration")
|
|
||||||
var conf Config
|
// Load configuration.
|
||||||
var err error
|
cfgPath := os.Getenv("RSSD_CONFIG")
|
||||||
if _, err = toml.DecodeFile(confFile, &conf); err != nil {
|
if cfgPath == "" {
|
||||||
log.Fatal("Error reading config file:",err)
|
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)
|
||||||
}
|
}
|
||||||
|
|
||||||
pl := newpcList(dataFile)
|
// Load state — store alongside the config file.
|
||||||
d := NewDaemon(conf,pl)
|
stateDir := filepath.Dir(cfgPath)
|
||||||
if d.conf.Workers != 0 {
|
sm, err := state.NewManager(filepath.Join(stateDir, stateFile))
|
||||||
d.workers = d.conf.Workers
|
|
||||||
} else {
|
|
||||||
d.workers = 3
|
|
||||||
}
|
|
||||||
dstDir,err = homedir.Expand(conf.DestDir)
|
|
||||||
if err != nil {
|
if err != nil {
|
||||||
log.Fatal("Error locating DestDir.")
|
fmt.Fprintf(os.Stderr, "rssd: fatal: %v\n", err)
|
||||||
|
os.Exit(1)
|
||||||
}
|
}
|
||||||
d.Start()
|
if err := sm.Load(); err != nil {
|
||||||
select { }
|
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 - <sanitized title> [<id>]" (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)
|
||||||
}
|
}
|
||||||
|
|
|
||||||
129
poller/poller.go
Normal file
129
poller/poller.go
Normal file
|
|
@ -0,0 +1,129 @@
|
||||||
|
package poller
|
||||||
|
|
||||||
|
import (
|
||||||
|
"fmt"
|
||||||
|
"io"
|
||||||
|
"net/http"
|
||||||
|
"time"
|
||||||
|
|
||||||
|
"github.com/mmcdole/gofeed"
|
||||||
|
)
|
||||||
|
|
||||||
|
// EnclosureInfo carries an enclosure URL and its associated item's metadata.
|
||||||
|
type EnclosureInfo struct {
|
||||||
|
URL string
|
||||||
|
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
|
||||||
|
}
|
||||||
|
|
||||||
|
// FeedResult holds the result of polling a single feed.
|
||||||
|
type FeedResult struct {
|
||||||
|
URL string
|
||||||
|
NewEnclosures []EnclosureInfo
|
||||||
|
LastModified string
|
||||||
|
ETag string
|
||||||
|
NotModified bool
|
||||||
|
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
|
||||||
|
parser *gofeed.Parser
|
||||||
|
}
|
||||||
|
|
||||||
|
// NewFetcher creates a new Fetcher with sensible defaults.
|
||||||
|
func NewFetcher() *Fetcher {
|
||||||
|
return &Fetcher{
|
||||||
|
httpClient: &http.Client{
|
||||||
|
Timeout: 30 * time.Second,
|
||||||
|
},
|
||||||
|
parser: gofeed.NewParser(),
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// Poll fetches and parses an RSS/Atom feed, returning new enclosure info.
|
||||||
|
// lastModified and etag are used for conditional requests (If-None-Match / If-Modified-Since).
|
||||||
|
// Returns a FeedResult with NewEnclosures populated only if the feed changed.
|
||||||
|
func (f *Fetcher) Poll(url, lastModified, etag string) (*FeedResult, error) {
|
||||||
|
req, err := http.NewRequest("GET", url, nil)
|
||||||
|
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)
|
||||||
|
}
|
||||||
|
if lastModified != "" {
|
||||||
|
req.Header.Set("If-Modified-Since", lastModified)
|
||||||
|
}
|
||||||
|
|
||||||
|
resp, err := f.httpClient.Do(req)
|
||||||
|
if err != nil {
|
||||||
|
return nil, fmt.Errorf("fetch %s: %w", url, err)
|
||||||
|
}
|
||||||
|
defer resp.Body.Close()
|
||||||
|
|
||||||
|
// Handle 304 Not Modified — feed hasn't changed.
|
||||||
|
if resp.StatusCode == http.StatusNotModified {
|
||||||
|
return &FeedResult{
|
||||||
|
URL: url,
|
||||||
|
NotModified: true,
|
||||||
|
}, nil
|
||||||
|
}
|
||||||
|
|
||||||
|
if resp.StatusCode != http.StatusOK {
|
||||||
|
return nil, fmt.Errorf("fetch %s: HTTP %d", url, resp.StatusCode)
|
||||||
|
}
|
||||||
|
|
||||||
|
body, err := io.ReadAll(resp.Body)
|
||||||
|
if err != nil {
|
||||||
|
return nil, fmt.Errorf("read body %s: %w", url, err)
|
||||||
|
}
|
||||||
|
|
||||||
|
feed, err := f.parser.ParseString(string(body))
|
||||||
|
if err != nil {
|
||||||
|
return nil, fmt.Errorf("parse feed %s: %w", url, err)
|
||||||
|
}
|
||||||
|
|
||||||
|
// Extract enclosure info from feed items.
|
||||||
|
var enclosures []EnclosureInfo
|
||||||
|
for _, item := range feed.Items {
|
||||||
|
for _, enc := range item.Enclosures {
|
||||||
|
if enc.URL != "" {
|
||||||
|
encInfo := EnclosureInfo{
|
||||||
|
URL: enc.URL,
|
||||||
|
Title: item.Title,
|
||||||
|
MimeType: enc.Type,
|
||||||
|
}
|
||||||
|
if item.PublishedParsed != nil {
|
||||||
|
encInfo.PublishedAt = *item.PublishedParsed
|
||||||
|
}
|
||||||
|
enclosures = append(enclosures, encInfo)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
result := &FeedResult{
|
||||||
|
URL: url,
|
||||||
|
NewEnclosures: enclosures,
|
||||||
|
NotModified: false,
|
||||||
|
}
|
||||||
|
|
||||||
|
// Capture cache headers for the next conditional request.
|
||||||
|
if v := resp.Header.Get("Last-Modified"); v != "" {
|
||||||
|
result.LastModified = v
|
||||||
|
}
|
||||||
|
if v := resp.Header.Get("ETag"); v != "" {
|
||||||
|
result.ETag = v
|
||||||
|
}
|
||||||
|
|
||||||
|
return result, nil
|
||||||
|
}
|
||||||
290
state/state.go
Normal file
290
state/state.go
Normal file
|
|
@ -0,0 +1,290 @@
|
||||||
|
package state
|
||||||
|
|
||||||
|
import (
|
||||||
|
"encoding/json"
|
||||||
|
"fmt"
|
||||||
|
"os"
|
||||||
|
"path/filepath"
|
||||||
|
"sync"
|
||||||
|
"time"
|
||||||
|
)
|
||||||
|
|
||||||
|
// State is the full persisted state of the daemon.
|
||||||
|
type State struct {
|
||||||
|
Feeds []FeedState `json:"feeds"`
|
||||||
|
Jobs []Job `json:"jobs"`
|
||||||
|
}
|
||||||
|
|
||||||
|
// FeedState holds per-feed metadata from HTTP responses.
|
||||||
|
type FeedState struct {
|
||||||
|
URL string `json:"url"`
|
||||||
|
LastModified string `json:"last_modified,omitempty"`
|
||||||
|
ETag string `json:"etag,omitempty"`
|
||||||
|
}
|
||||||
|
|
||||||
|
// Job represents a download job with retry state.
|
||||||
|
type Job struct {
|
||||||
|
FeedURL string `json:"feed_url"`
|
||||||
|
EnclosureURL string `json:"enclosure_url"`
|
||||||
|
DestPath string `json:"dest_path"`
|
||||||
|
Status string `json:"status"`
|
||||||
|
AttemptCount int `json:"attempt_count"`
|
||||||
|
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)
|
||||||
|
LiveWaitCount int `json:"live_wait_count,omitempty"` // consecutive 10-min waits while a video is still live (does not count against AttemptCount)
|
||||||
|
}
|
||||||
|
|
||||||
|
// JobStatus values.
|
||||||
|
const (
|
||||||
|
StatusPending = "pending"
|
||||||
|
StatusDownloading = "downloading"
|
||||||
|
StatusRetrying = "retrying"
|
||||||
|
StatusDownloaded = "downloaded"
|
||||||
|
StatusFailed = "failed"
|
||||||
|
StatusSkipped = "skipped" // resolved but excluded by start_date; never downloaded
|
||||||
|
)
|
||||||
|
|
||||||
|
// Manager provides thread-safe access to persisted state.
|
||||||
|
type Manager struct {
|
||||||
|
mu sync.Mutex
|
||||||
|
path string
|
||||||
|
state *State
|
||||||
|
}
|
||||||
|
|
||||||
|
// NewManager creates a new state manager pointing at the given file path.
|
||||||
|
func NewManager(path string) (*Manager, error) {
|
||||||
|
dir := filepath.Dir(path)
|
||||||
|
if err := os.MkdirAll(dir, 0755); err != nil {
|
||||||
|
return nil, fmt.Errorf("create state directory: %w", err)
|
||||||
|
}
|
||||||
|
return &Manager{path: path, state: &State{}}, nil
|
||||||
|
}
|
||||||
|
|
||||||
|
// Load reads the state file from disk. If it doesn't exist or is empty, returns an empty state.
|
||||||
|
func (m *Manager) Load() error {
|
||||||
|
m.mu.Lock()
|
||||||
|
defer m.mu.Unlock()
|
||||||
|
|
||||||
|
data, err := os.ReadFile(m.path)
|
||||||
|
if err != nil {
|
||||||
|
if os.IsNotExist(err) {
|
||||||
|
m.state = &State{}
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
return fmt.Errorf("read state: %w", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
if len(data) == 0 {
|
||||||
|
m.state = &State{}
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
|
||||||
|
var s State
|
||||||
|
if err := json.Unmarshal(data, &s); err != nil {
|
||||||
|
return fmt.Errorf("parse state: %w", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
m.state = &s
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
|
||||||
|
// Save atomically writes the current state to disk with pretty-printed JSON.
|
||||||
|
func (m *Manager) Save() error {
|
||||||
|
m.mu.Lock()
|
||||||
|
defer m.mu.Unlock()
|
||||||
|
|
||||||
|
data, err := json.MarshalIndent(m.state, "", " ")
|
||||||
|
if err != nil {
|
||||||
|
return fmt.Errorf("marshal state: %w", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
dir := filepath.Dir(m.path)
|
||||||
|
tmpFile, err := os.CreateTemp(dir, ".rssd-state-*.tmp")
|
||||||
|
if err != nil {
|
||||||
|
return fmt.Errorf("create temp file: %w", err)
|
||||||
|
}
|
||||||
|
tmpName := tmpFile.Name()
|
||||||
|
|
||||||
|
if _, err := tmpFile.Write(data); err != nil {
|
||||||
|
tmpFile.Close()
|
||||||
|
os.Remove(tmpName)
|
||||||
|
return fmt.Errorf("write temp state: %w", err)
|
||||||
|
}
|
||||||
|
if err := tmpFile.Close(); err != nil {
|
||||||
|
os.Remove(tmpName)
|
||||||
|
return fmt.Errorf("close temp state: %w", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
if err := os.Rename(tmpName, m.path); err != nil {
|
||||||
|
os.Remove(tmpName)
|
||||||
|
return fmt.Errorf("rename state file: %w", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
|
||||||
|
// JobSnapshot is a copy of a Job suitable for read-only inspection.
|
||||||
|
type JobSnapshot struct {
|
||||||
|
FeedURL string
|
||||||
|
EnclosureURL string
|
||||||
|
DestPath string
|
||||||
|
Status string
|
||||||
|
AttemptCount int
|
||||||
|
MaxRetries int
|
||||||
|
NextAttemptAt *time.Time // nil means ready now
|
||||||
|
Error string
|
||||||
|
Duration float64
|
||||||
|
}
|
||||||
|
|
||||||
|
// ForEachJob iterates over all jobs and calls fn for each one.
|
||||||
|
// The callback receives a pointer to a copy of the job (not the live state),
|
||||||
|
// so it must not modify fields. Use UpdateJob or AddJob for mutations.
|
||||||
|
// fn may perform I/O without risk of deadlock — the lock is released before
|
||||||
|
// any callback invocation.
|
||||||
|
func (m *Manager) ForEachJob(fn func(*JobSnapshot)) {
|
||||||
|
m.mu.Lock()
|
||||||
|
jobsCopy := make([]JobSnapshot, len(m.state.Jobs))
|
||||||
|
for i := range m.state.Jobs {
|
||||||
|
j := &m.state.Jobs[i]
|
||||||
|
jobsCopy[i] = JobSnapshot{
|
||||||
|
FeedURL: j.FeedURL,
|
||||||
|
EnclosureURL: j.EnclosureURL,
|
||||||
|
DestPath: j.DestPath,
|
||||||
|
Status: j.Status,
|
||||||
|
AttemptCount: j.AttemptCount,
|
||||||
|
MaxRetries: j.MaxRetries,
|
||||||
|
NextAttemptAt: j.NextAttemptAt,
|
||||||
|
Error: j.Error,
|
||||||
|
Duration: j.Duration,
|
||||||
|
}
|
||||||
|
}
|
||||||
|
m.mu.Unlock()
|
||||||
|
|
||||||
|
for i := range jobsCopy {
|
||||||
|
fn(&jobsCopy[i])
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// State returns a copy of the current in-memory state for iteration.
|
||||||
|
func (m *Manager) State() *State {
|
||||||
|
m.mu.Lock()
|
||||||
|
defer m.mu.Unlock()
|
||||||
|
return m.state
|
||||||
|
}
|
||||||
|
|
||||||
|
// AddJob adds a new download job if one with this enclosure URL doesn't already exist.
|
||||||
|
// 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()
|
||||||
|
|
||||||
|
// Check if this enclosure URL already has a job.
|
||||||
|
for _, j := range m.state.Jobs {
|
||||||
|
if j.EnclosureURL == enclosureURL {
|
||||||
|
return false // already exists
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
job := Job{
|
||||||
|
FeedURL: feedURL,
|
||||||
|
EnclosureURL: enclosureURL,
|
||||||
|
DestPath: destPath,
|
||||||
|
Status: status,
|
||||||
|
AttemptCount: 0,
|
||||||
|
MaxRetries: maxRetries,
|
||||||
|
NextAttemptAt: nil,
|
||||||
|
Error: "",
|
||||||
|
Duration: duration,
|
||||||
|
}
|
||||||
|
|
||||||
|
m.state.Jobs = append(m.state.Jobs, job)
|
||||||
|
return true
|
||||||
|
}
|
||||||
|
|
||||||
|
// UpdateJob finds a job by enclosure URL and applies the update function.
|
||||||
|
// Returns true if the job was found and updated.
|
||||||
|
func (m *Manager) UpdateJob(enclosureURL string, fn func(*Job)) bool {
|
||||||
|
m.mu.Lock()
|
||||||
|
defer m.mu.Unlock()
|
||||||
|
|
||||||
|
for i := range m.state.Jobs {
|
||||||
|
if m.state.Jobs[i].EnclosureURL == enclosureURL {
|
||||||
|
fn(&m.state.Jobs[i])
|
||||||
|
return true
|
||||||
|
}
|
||||||
|
}
|
||||||
|
return false
|
||||||
|
}
|
||||||
|
|
||||||
|
// FindJobByURL returns a pointer to the job for the given enclosure URL, or nil.
|
||||||
|
func (m *Manager) FindJobByURL(enclosureURL string) *Job {
|
||||||
|
m.mu.Lock()
|
||||||
|
defer m.mu.Unlock()
|
||||||
|
|
||||||
|
for i := range m.state.Jobs {
|
||||||
|
if m.state.Jobs[i].EnclosureURL == enclosureURL {
|
||||||
|
return &m.state.Jobs[i]
|
||||||
|
}
|
||||||
|
}
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
|
||||||
|
// HasJob returns true if a job with the given enclosure URL exists.
|
||||||
|
func (m *Manager) HasJob(enclosureURL string) bool {
|
||||||
|
m.mu.Lock()
|
||||||
|
defer m.mu.Unlock()
|
||||||
|
|
||||||
|
for _, j := range m.state.Jobs {
|
||||||
|
if j.EnclosureURL == enclosureURL {
|
||||||
|
return true
|
||||||
|
}
|
||||||
|
}
|
||||||
|
return false
|
||||||
|
}
|
||||||
|
|
||||||
|
// FindFeedByURL returns the FeedState for a given feed URL, or nil.
|
||||||
|
func (m *Manager) FindFeedByURL(feedURL string) *FeedState {
|
||||||
|
m.mu.Lock()
|
||||||
|
defer m.mu.Unlock()
|
||||||
|
|
||||||
|
for i := range m.state.Feeds {
|
||||||
|
if m.state.Feeds[i].URL == feedURL {
|
||||||
|
return &m.state.Feeds[i]
|
||||||
|
}
|
||||||
|
}
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
|
||||||
|
// UpsertFeed adds or updates a feed's metadata (ETag, Last-Modified).
|
||||||
|
func (m *Manager) UpsertFeed(feedURL, lastModified, etag string) {
|
||||||
|
m.mu.Lock()
|
||||||
|
defer m.mu.Unlock()
|
||||||
|
|
||||||
|
for i := range m.state.Feeds {
|
||||||
|
if m.state.Feeds[i].URL == feedURL {
|
||||||
|
m.state.Feeds[i].LastModified = lastModified
|
||||||
|
m.state.Feeds[i].ETag = etag
|
||||||
|
return
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
m.state.Feeds = append(m.state.Feeds, FeedState{
|
||||||
|
URL: feedURL,
|
||||||
|
LastModified: lastModified,
|
||||||
|
ETag: etag,
|
||||||
|
})
|
||||||
|
}
|
||||||
197
ytdlp/ytdlp.go
Normal file
197
ytdlp/ytdlp.go
Normal file
|
|
@ -0,0 +1,197 @@
|
||||||
|
// 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
|
||||||
|
LiveStatus string // "not_live", "is_live", "was_live", "is_upcoming", or "" if unknown
|
||||||
|
// ReleaseTimestamp is the scheduled start time of a live stream in unix
|
||||||
|
// seconds (from liveBroadcastDetails.startTimestamp); 0 if unknown.
|
||||||
|
ReleaseTimestamp int64
|
||||||
|
}
|
||||||
|
|
||||||
|
// 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
|
||||||
|
}
|
||||||
|
}
|
||||||
|
info.LiveStatus, _ = data["live_status"].(string)
|
||||||
|
if ts, ok := num(data["release_timestamp"]); ok {
|
||||||
|
info.ReleaseTimestamp = int64(ts)
|
||||||
|
}
|
||||||
|
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
|
||||||
|
}
|
||||||
Loading…
Reference in New Issue
Block a user