diff --git a/backend/app/services/download_jobs.py b/backend/app/services/download_jobs.py index ea3ba42..5342a68 100644 --- a/backend/app/services/download_jobs.py +++ b/backend/app/services/download_jobs.py @@ -106,18 +106,27 @@ async def handle_metube_event(db: Session, event_name: str, raw_payload) -> None if job is None: return - _apply_metube_info(job, payload) + # Only the dedicated 'completed' event is authoritative for terminal + # status. yt-dlp fires per-stream progress hooks with status=='finished' + # when downloading video+audio separately for merging (one stream lands, + # the other is still in flight) -- MeTube forwards that transient state + # through 'updated' events too, so treating it as our own "completed" + # there flaps the job back and forth: downloading -> completed -> + # downloading -> ... -> postprocessing -> completed. Verified against a + # real deploy's logs before this fix. + _apply_metube_info(job, payload, authoritative=event_name == "completed") db.commit() logger.info("download_job %s updated to status=%s via '%s' event", job.id, job.status, event_name) -def _apply_metube_info(job: DownloadJob, info: dict) -> None: +def _apply_metube_info(job: DownloadJob, info: dict, *, authoritative: bool) -> None: if info.get("id"): job.metube_job_id = info["id"] - our_status = METUBE_STATUS_MAP.get(info.get("status"), "unknown") + metube_status = info.get("status") + our_status = METUBE_STATUS_MAP.get(metube_status, "unknown") - if job.started_at is None and our_status in ("downloading", "postprocessing", "completed"): + if job.started_at is None and metube_status in ("downloading", "postprocessing"): job.started_at = datetime.now(timezone.utc) percent = info.get("percent") @@ -128,6 +137,10 @@ def _apply_metube_info(job: DownloadJob, info: dict) -> None: job.status = "failed" job.error_message = info.get("msg") or info.get("error") or "Ошибка загрузки" elif our_status == "completed": + if not authoritative: + # A stream finished, not necessarily the whole job (see comment + # above) -- keep the current status, just take the progress tick. + return filename = info.get("filename") if filename: job.metube_filename = filename @@ -135,7 +148,7 @@ def _apply_metube_info(job: DownloadJob, info: dict) -> None: job.status = "completed" job.progress_percent = 100 job.completed_at = job.completed_at or datetime.now(timezone.utc) - else: + elif our_status != "unknown": job.status = our_status @@ -158,22 +171,25 @@ def reconcile_on_startup(db: Session) -> None: db.commit() return - by_url: dict[str, dict] = {} - for bucket in ("queue", "pending", "done"): + # (item, authoritative) -- only the 'done' bucket is a terminal, trustworthy + # outcome; 'queue'/'pending' are still in flight (see handle_metube_event). + by_url: dict[str, tuple[dict, bool]] = {} + for bucket, authoritative in (("queue", False), ("pending", False), ("done", True)): for item in history.get(bucket, []) or []: url = item.get("url") if url: - by_url[url] = item + by_url[url] = (item, authoritative) videos = {v.id: v for v in db.query(Video).filter(Video.id.in_([j.video_id for j in jobs])).all()} for job in jobs: video = videos.get(job.video_id) - info = by_url.get(video.youtube_url) if video else None - if info is None: + match = by_url.get(video.youtube_url) if video else None + if match is None: job.status = "unknown" continue - _apply_metube_info(job, info) + info, authoritative = match + _apply_metube_info(job, info, authoritative=authoritative) db.commit() logger.info("Reconciled %d download job(s) against MeTube history", len(jobs)) diff --git a/tests/test_download_jobs.py b/tests/test_download_jobs.py index 096fa9d..adad227 100644 --- a/tests/test_download_jobs.py +++ b/tests/test_download_jobs.py @@ -138,6 +138,59 @@ async def test_handle_metube_event_completed_builds_media_url(monkeypatch, db_se assert job.completed_at is not None +@pytest.mark.asyncio +async def test_transient_finished_in_updated_event_does_not_complete_job(db_session): + """Regression test for a real bug seen in production logs: yt-dlp reports + status='finished' via 'updated' events once per stream when merging + separate video+audio (finished for video, then still downloading audio), + which used to flip our job to 'completed' prematurely and then back to + 'downloading', flapping the UI. Only the dedicated 'completed' event may + mark a job completed.""" + video = _seed_video(db_session) + job = DownloadJob(video_id=video.id, status="downloading", metube_job_id="vid1.vid1", progress_percent=80) + db_session.add(job) + db_session.commit() + + # video stream finished (transient) -- must NOT become "completed" + await download_jobs.handle_metube_event( + db_session, + "updated", + json.dumps({"id": "vid1.vid1", "url": video.youtube_url, "status": "finished"}), + ) + db_session.refresh(job) + assert job.status == "downloading" + + # audio stream still downloading + await download_jobs.handle_metube_event( + db_session, + "updated", + json.dumps({"id": "vid1.vid1", "url": video.youtube_url, "status": "downloading", "percent": 10}), + ) + db_session.refresh(job) + assert job.status == "downloading" + assert job.progress_percent == 10 + + # ffmpeg merge + await download_jobs.handle_metube_event( + db_session, + "updated", + json.dumps({"id": "vid1.vid1", "url": video.youtube_url, "status": "postprocessing"}), + ) + db_session.refresh(job) + assert job.status == "postprocessing" + + # only the dedicated 'completed' event finalizes it + await download_jobs.handle_metube_event( + db_session, + "completed", + json.dumps( + {"id": "vid1.vid1", "url": video.youtube_url, "status": "finished", "filename": "vid1.vid1.mp4"} + ), + ) + db_session.refresh(job) + assert job.status == "completed" + + @pytest.mark.asyncio async def test_handle_metube_event_error_marks_failed(db_session): video = _seed_video(db_session)