diff --git a/downloader/downloader.go b/downloader/downloader.go index 9b3a7ca..3ab3f6f 100644 --- a/downloader/downloader.go +++ b/downloader/downloader.go @@ -57,9 +57,13 @@ const liveWaitInterval = 10 * time.Minute // 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") +// 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) { @@ -82,7 +86,7 @@ func (p *Pool) worker(job DownloadJob) { } if err != nil { if errors.Is(err, errLiveWait) { - fmt.Fprintf(os.Stderr, "rssd: %s still live, retrying in %s\n", job.EnclosureURL, liveWaitInterval) + 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) } @@ -97,13 +101,14 @@ 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 + // 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. @@ -163,40 +168,63 @@ 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 { +// 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 + return false, time.Time{} } - return info.LiveStatus == "is_live" || info.LiveStatus == "upcoming" + 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 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) { +// 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++ - t := time.Now() 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) - t = t.Add(liveWaitInterval) } else { j.LiveWaitCount = 0 j.AttemptCount++ delay := p.backoffBase * time.Duration(1<