rssd/downloader/downloader.go
Greg Pomerantz c562ecc411 Schedule upcoming streams at their start time; fix is_upcoming check
yt-dlp reports scheduled premieres as live_status 'is_upcoming' (not
'upcoming'), and exposes release_timestamp with the stream's scheduled
start time. For an upcoming stream with a known start, reschedule the
job at start+2min without touching the wait counter, no matter how far
away the start is — the bounded 10-minute wait clock (maxLiveWaits,
~4h) now only applies while a stream is actually running. Upcoming
streams with an unknown start time keep the 10-minute wait behavior.
2026-09-28 17:54:25 -04:00

421 lines
13 KiB
Go

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()
}