Fix download status flapping: only the 'completed' event is terminal
Production logs showed a job oscillating downloading -> completed -> downloading -> completed -> postprocessing -> completed. Cause: yt-dlp reports status='finished' via progress-hook ticks once per stream when downloading separate video+audio for muxing (video lands, audio is still in flight), and MeTube forwards that through 'updated' Socket.IO events too -- we were mapping any status='finished' to our "completed", regardless of which event carried it. Only the dedicated 'completed' event (and the 'done' bucket of MeTube's /history, for startup reconciliation) is now treated as authoritative for terminal status; a transient 'finished' arriving via 'added'/'updated' (or history's queue/pending) is a no-op for status, matching the progress_percent update it also carries. Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com>
This commit is contained in:
parent
fe16c08daa
commit
e333296170
2 changed files with 80 additions and 11 deletions
|
|
@ -106,18 +106,27 @@ async def handle_metube_event(db: Session, event_name: str, raw_payload) -> None
|
||||||
if job is None:
|
if job is None:
|
||||||
return
|
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()
|
db.commit()
|
||||||
logger.info("download_job %s updated to status=%s via '%s' event", job.id, job.status, event_name)
|
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"):
|
if info.get("id"):
|
||||||
job.metube_job_id = info["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)
|
job.started_at = datetime.now(timezone.utc)
|
||||||
|
|
||||||
percent = info.get("percent")
|
percent = info.get("percent")
|
||||||
|
|
@ -128,6 +137,10 @@ def _apply_metube_info(job: DownloadJob, info: dict) -> None:
|
||||||
job.status = "failed"
|
job.status = "failed"
|
||||||
job.error_message = info.get("msg") or info.get("error") or "Ошибка загрузки"
|
job.error_message = info.get("msg") or info.get("error") or "Ошибка загрузки"
|
||||||
elif our_status == "completed":
|
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")
|
filename = info.get("filename")
|
||||||
if filename:
|
if filename:
|
||||||
job.metube_filename = filename
|
job.metube_filename = filename
|
||||||
|
|
@ -135,7 +148,7 @@ def _apply_metube_info(job: DownloadJob, info: dict) -> None:
|
||||||
job.status = "completed"
|
job.status = "completed"
|
||||||
job.progress_percent = 100
|
job.progress_percent = 100
|
||||||
job.completed_at = job.completed_at or datetime.now(timezone.utc)
|
job.completed_at = job.completed_at or datetime.now(timezone.utc)
|
||||||
else:
|
elif our_status != "unknown":
|
||||||
job.status = our_status
|
job.status = our_status
|
||||||
|
|
||||||
|
|
||||||
|
|
@ -158,22 +171,25 @@ def reconcile_on_startup(db: Session) -> None:
|
||||||
db.commit()
|
db.commit()
|
||||||
return
|
return
|
||||||
|
|
||||||
by_url: dict[str, dict] = {}
|
# (item, authoritative) -- only the 'done' bucket is a terminal, trustworthy
|
||||||
for bucket in ("queue", "pending", "done"):
|
# 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 []:
|
for item in history.get(bucket, []) or []:
|
||||||
url = item.get("url")
|
url = item.get("url")
|
||||||
if 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()}
|
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:
|
for job in jobs:
|
||||||
video = videos.get(job.video_id)
|
video = videos.get(job.video_id)
|
||||||
info = by_url.get(video.youtube_url) if video else None
|
match = by_url.get(video.youtube_url) if video else None
|
||||||
if info is None:
|
if match is None:
|
||||||
job.status = "unknown"
|
job.status = "unknown"
|
||||||
continue
|
continue
|
||||||
_apply_metube_info(job, info)
|
info, authoritative = match
|
||||||
|
_apply_metube_info(job, info, authoritative=authoritative)
|
||||||
|
|
||||||
db.commit()
|
db.commit()
|
||||||
logger.info("Reconciled %d download job(s) against MeTube history", len(jobs))
|
logger.info("Reconciled %d download job(s) against MeTube history", len(jobs))
|
||||||
|
|
|
||||||
|
|
@ -138,6 +138,59 @@ async def test_handle_metube_event_completed_builds_media_url(monkeypatch, db_se
|
||||||
assert job.completed_at is not None
|
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
|
@pytest.mark.asyncio
|
||||||
async def test_handle_metube_event_error_marks_failed(db_session):
|
async def test_handle_metube_event_error_marks_failed(db_session):
|
||||||
video = _seed_video(db_session)
|
video = _seed_video(db_session)
|
||||||
|
|
|
||||||
Loading…
Add table
Add a link
Reference in a new issue