diff --git a/downloader/downloader.go b/downloader/downloader.go index 0f95fbb..9b3a7ca 100644 --- a/downloader/downloader.go +++ b/downloader/downloader.go @@ -4,6 +4,7 @@ import ( "bytes" "context" "encoding/json" + "errors" "fmt" "io" "math" @@ -44,6 +45,22 @@ 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 + +// errLiveWait is returned by downloadYT when the video is still live and +// has been rescheduled; it is not a download failure. +var errLiveWait = errors.New("video is live, waiting for stream to end") + // 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++ { @@ -64,7 +81,11 @@ func (p *Pool) worker(job DownloadJob) { err = p.download(job) } if err != nil { - fmt.Fprintf(os.Stderr, "ERROR: download %s failed: %v\n", job.EnclosureURL, err) + if errors.Is(err, errLiveWait) { + fmt.Fprintf(os.Stderr, "rssd: %s still live, retrying in %s\n", job.EnclosureURL, liveWaitInterval) + } 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) } @@ -76,6 +97,15 @@ func (p *Pool) worker(job DownloadJob) { // WebM) with no re-encoding. The file lands at destPath + "."; // state is updated to the real path. func (p *Pool) downloadYT(job DownloadJob) error { + // While the video is still a live stream, the static formats we request + // are not available yet. Instead of counting a failed attempt, reschedule + // the job in liveWaitInterval; once the stream ends and converts to a + // regular VOD the normal download succeeds. + if p.isStillLive(job) { + p.markLiveWait(job.EnclosureURL) + return errLiveWait + } + // Update state to downloading. if !p.state.UpdateJob(job.EnclosureURL, func(j *state.Job) { j.Status = state.StatusDownloading @@ -123,6 +153,7 @@ func (p *Pool) downloadYT(job DownloadJob) error { 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) @@ -132,6 +163,42 @@ func (p *Pool) downloadYT(job DownloadJob) error { return nil } +// isStillLive reports whether the video is a live stream that has not ended +// yet (or is not even started). If the probe itself fails, false is returned +// so the normal download path (with its own error handling) can proceed. +func (p *Pool) isStillLive(job DownloadJob) bool { + 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 + } + return info.LiveStatus == "is_live" || info.LiveStatus == "upcoming" +} + +// markLiveWait reschedules a job whose video is still live. The wait does +// 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) { + p.state.UpdateJob(enclosureURL, func(j *state.Job) { + j.Status = state.StatusRetrying + j.LiveWaitCount++ + t := time.Now() + if j.LiveWaitCount <= maxLiveWaits { + j.Error = fmt.Sprintf("video is live (wait %d/%d), retrying after stream ends", j.LiveWaitCount, maxLiveWaits) + t = t.Add(liveWaitInterval) + } else { + j.LiveWaitCount = 0 + j.AttemptCount++ + delay := p.backoffBase * time.Duration(1<