myYouTube/tests/test_download_jobs.py

657 lines
23 KiB
Python
Raw Permalink Normal View History

import json
from datetime import datetime, timedelta, timezone
import pytest
from app.models.channel import Channel
from app.models.download_job import DownloadJob
from app.models.video import Video
from app.services import download_jobs
def _seed_video(db_session, youtube_video_id="vid1", youtube_channel_id="chanA"):
channel = Channel(youtube_channel_id=youtube_channel_id, title="Channel", subscribed=True)
db_session.add(channel)
db_session.commit()
from datetime import datetime, timezone
video = Video(
youtube_video_id=youtube_video_id,
channel_id=channel.id,
title="Video",
published_at=datetime(2026, 9, 10, tzinfo=timezone.utc),
youtube_url=f"https://www.youtube.com/watch?v={youtube_video_id}",
)
db_session.add(video)
db_session.commit()
db_session.refresh(video)
return video
def test_request_download_enqueues_and_creates_job(monkeypatch, db_session):
video = _seed_video(db_session)
calls = {}
def fake_enqueue(self, youtube_url, custom_name_prefix):
calls["youtube_url"] = youtube_url
calls["custom_name_prefix"] = custom_name_prefix
return {"status": "ok"}
monkeypatch.setattr("app.services.metube_client.MeTubeClient.enqueue_video", fake_enqueue)
job = download_jobs.request_download(db_session, video)
assert job.status == "queued"
assert calls["youtube_url"] == video.youtube_url
assert calls["custom_name_prefix"] == video.youtube_video_id
def test_request_download_is_idempotent_while_active(monkeypatch, db_session):
video = _seed_video(db_session)
monkeypatch.setattr(
"app.services.metube_client.MeTubeClient.enqueue_video", lambda self, u, p: {"status": "ok"}
)
job1 = download_jobs.request_download(db_session, video)
job2 = download_jobs.request_download(db_session, video)
assert job1.id == job2.id
assert db_session.query(DownloadJob).count() == 1
def test_request_download_raises_on_metube_error(monkeypatch, db_session):
video = _seed_video(db_session)
monkeypatch.setattr(
"app.services.metube_client.MeTubeClient.enqueue_video",
lambda self, u, p: {"status": "error", "msg": "boom"},
)
with pytest.raises(download_jobs.MeTubeRejected):
download_jobs.request_download(db_session, video)
assert db_session.query(DownloadJob).count() == 0
def test_request_download_allows_retry_after_failure(monkeypatch, db_session):
video = _seed_video(db_session)
job = DownloadJob(video_id=video.id, status="failed", error_message="oops")
db_session.add(job)
db_session.commit()
monkeypatch.setattr(
"app.services.metube_client.MeTubeClient.enqueue_video", lambda self, u, p: {"status": "ok"}
)
new_job = download_jobs.request_download(db_session, video)
assert new_job.id != job.id
assert db_session.query(DownloadJob).count() == 2
@pytest.mark.asyncio
async def test_handle_metube_event_added_and_updated(db_session):
video = _seed_video(db_session)
job = DownloadJob(video_id=video.id, status="queued")
db_session.add(job)
db_session.commit()
added_payload = json.dumps({"id": "vid1.vid1", "url": video.youtube_url, "status": "pending"})
await download_jobs.handle_metube_event(db_session, "added", added_payload)
db_session.refresh(job)
assert job.metube_job_id == "vid1.vid1"
assert job.status == "queued"
updating_payload = json.dumps(
{"id": "vid1.vid1", "url": video.youtube_url, "status": "downloading", "percent": 42.5}
)
await download_jobs.handle_metube_event(db_session, "updated", updating_payload)
db_session.refresh(job)
assert job.status == "downloading"
assert job.progress_percent == 42
assert job.started_at is not None
@pytest.mark.asyncio
async def test_handle_metube_event_completed_builds_media_url(monkeypatch, db_session):
video = _seed_video(db_session)
job = DownloadJob(video_id=video.id, status="downloading", metube_job_id="vid1.vid1")
db_session.add(job)
db_session.commit()
monkeypatch.setattr(
"app.services.metube_client.MeTubeClient.build_media_url",
lambda self, filename: f"http://metube.local/download/{filename}",
)
payload = json.dumps(
{"id": "vid1.vid1", "url": video.youtube_url, "status": "finished", "filename": "vid1.vid1.mp4"}
)
await download_jobs.handle_metube_event(db_session, "completed", payload)
db_session.refresh(job)
assert job.status == "completed"
assert job.progress_percent == 100
assert job.metube_filename == "vid1.vid1.mp4"
assert job.media_url == "http://metube.local/download/vid1.vid1.mp4"
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)
job = DownloadJob(video_id=video.id, status="downloading", metube_job_id="vid1.vid1")
db_session.add(job)
db_session.commit()
payload = json.dumps({"id": "vid1.vid1", "url": video.youtube_url, "status": "error", "msg": "network blip"})
await download_jobs.handle_metube_event(db_session, "updated", payload)
db_session.refresh(job)
assert job.status == "failed"
assert job.error_message == "network blip"
@pytest.mark.asyncio
async def test_handle_metube_event_ignores_unrelated_download(db_session):
video = _seed_video(db_session)
job = DownloadJob(video_id=video.id, status="queued")
db_session.add(job)
db_session.commit()
payload = json.dumps(
{"id": "other.other", "url": "https://www.youtube.com/watch?v=other", "status": "finished"}
)
await download_jobs.handle_metube_event(db_session, "completed", payload)
db_session.refresh(job)
assert job.status == "queued"
@pytest.mark.asyncio
async def test_handle_metube_event_canceled_marks_failed(db_session):
video = _seed_video(db_session)
job = DownloadJob(video_id=video.id, status="downloading", metube_job_id="vid1.vid1")
db_session.add(job)
db_session.commit()
await download_jobs.handle_metube_event(db_session, "canceled", json.dumps("vid1.vid1"))
db_session.refresh(job)
assert job.status == "failed"
@pytest.mark.asyncio
async def test_handle_metube_event_canceled_matches_by_url(db_session):
"""MeTube keys canceled/cleared events by the download's URL, not its id."""
video = _seed_video(db_session)
job = DownloadJob(video_id=video.id, status="downloading", metube_job_id="vid1.vid1")
db_session.add(job)
db_session.commit()
await download_jobs.handle_metube_event(db_session, "canceled", json.dumps(video.youtube_url))
db_session.refresh(job)
assert job.status == "failed"
assert job.error_message == "Отменено в MeTube"
@pytest.mark.asyncio
async def test_handle_metube_event_cleared_matches_job_by_url(db_session):
video = _seed_video(db_session)
job = DownloadJob(video_id=video.id, status="downloading", metube_job_id="vid1.vid1")
other_video = _seed_video(db_session, youtube_video_id="vid2", youtube_channel_id="chanB")
other_job = DownloadJob(video_id=other_video.id, status="queued")
db_session.add(job)
db_session.add(other_job)
db_session.commit()
await download_jobs.handle_metube_event(db_session, "cleared", json.dumps(video.youtube_url))
db_session.refresh(job)
db_session.refresh(other_job)
assert job.status == "failed"
assert job.error_message == "Очищено в MeTube"
assert other_job.status == "queued"
@pytest.mark.asyncio
async def test_handle_metube_event_cleared_matching_terminal_job_is_noop(db_session):
"""Cleared payload resolving to a finished job must not touch terminal
statuses, nor spill over onto unrelated active jobs."""
video = _seed_video(db_session)
completed = DownloadJob(video_id=video.id, status="completed", media_url="http://x/f.mp4")
other_video = _seed_video(db_session, youtube_video_id="vid2", youtube_channel_id="chanB")
active = DownloadJob(video_id=other_video.id, status="downloading")
db_session.add_all([completed, active])
db_session.commit()
await download_jobs.handle_metube_event(db_session, "cleared", json.dumps(video.youtube_url))
db_session.refresh(completed)
db_session.refresh(active)
assert completed.status == "completed"
assert active.status == "downloading"
@pytest.mark.asyncio
async def test_handle_metube_event_cleared_unmatched_payload_is_noop(db_session):
"""'cleared' only ever refers to the one done-entry being removed (trash
or CLEAR_COMPLETED_AFTER). A payload that doesn't match any of our jobs is
a foreign MeTube download we never tracked -- the jobs MeTube is still
running must stay untouched."""
video1 = _seed_video(db_session)
job1 = DownloadJob(video_id=video1.id, status="downloading")
video2 = _seed_video(db_session, youtube_video_id="vid2", youtube_channel_id="chanB")
job2 = DownloadJob(video_id=video2.id, status="postprocessing")
video3 = _seed_video(db_session, youtube_video_id="vid3", youtube_channel_id="chanC")
completed = DownloadJob(video_id=video3.id, status="completed", media_url="http://x/f.mp4")
unknown = DownloadJob(video_id=video3.id, status="unknown")
db_session.add_all([job1, job2, completed, unknown])
db_session.commit()
await download_jobs.handle_metube_event(db_session, "cleared", json.dumps("https://example.com/other"))
db_session.refresh(job1)
db_session.refresh(job2)
db_session.refresh(completed)
db_session.refresh(unknown)
assert job1.status == "downloading"
assert job2.status == "postprocessing"
assert completed.status == "completed"
assert unknown.status == "unknown"
@pytest.mark.asyncio
async def test_handle_metube_event_cleared_without_payload_is_noop(db_session):
video = _seed_video(db_session)
job = DownloadJob(video_id=video.id, status="queued")
db_session.add(job)
db_session.commit()
await download_jobs.handle_metube_event(db_session, "cleared", None)
db_session.refresh(job)
assert job.status == "queued"
assert job.error_message is None
def test_reconcile_marks_unknown_when_history_unavailable(monkeypatch, db_session):
video = _seed_video(db_session)
job = DownloadJob(video_id=video.id, status="downloading")
db_session.add(job)
db_session.commit()
def raise_error(self):
raise RuntimeError("connection refused")
monkeypatch.setattr("app.services.metube_client.MeTubeClient.fetch_history", raise_error)
download_jobs.reconcile_on_startup(db_session)
db_session.refresh(job)
assert job.status == "unknown"
def test_reconcile_restores_completed_from_history(monkeypatch, db_session):
video = _seed_video(db_session)
job = DownloadJob(video_id=video.id, status="downloading")
db_session.add(job)
db_session.commit()
monkeypatch.setattr(
"app.services.metube_client.MeTubeClient.fetch_history",
lambda self: {
"queue": [],
"pending": [],
"done": [{"id": "vid1.vid1", "url": video.youtube_url, "status": "finished", "filename": "f.mp4"}],
},
)
monkeypatch.setattr(
"app.services.metube_client.MeTubeClient.build_media_url",
lambda self, filename: f"http://metube.local/download/{filename}",
)
download_jobs.reconcile_on_startup(db_session)
db_session.refresh(job)
assert job.status == "completed"
assert job.media_url == "http://metube.local/download/f.mp4"
def test_reconcile_leaves_completed_jobs_untouched(monkeypatch, db_session):
video = _seed_video(db_session)
job = DownloadJob(video_id=video.id, status="completed", media_url="http://x/f.mp4")
db_session.add(job)
db_session.commit()
def fail_if_called(self):
raise AssertionError("should not fetch history for completed jobs")
monkeypatch.setattr("app.services.metube_client.MeTubeClient.fetch_history", fail_if_called)
download_jobs.reconcile_on_startup(db_session)
db_session.refresh(job)
assert job.status == "completed"
assert job.media_url == "http://x/f.mp4"
# --- Bug 2: non-authoritative 'updated' events must never roll terminal
# statuses back (double progress bar bug) ---------------------------------
@pytest.mark.asyncio
async def test_updated_event_does_not_roll_back_completed_job(db_session):
"""A late 'updated' event arriving after the authoritative 'completed' must
not flip the job back to downloading/postprocessing/failed -- that used to
re-trigger the progress bar after completion. Progress ticks are still
taken; status and its payload (media_url/completed_at/error_message) are
left alone."""
video = _seed_video(db_session)
completed_at = datetime(2026, 9, 20, tzinfo=timezone.utc)
job = DownloadJob(
video_id=video.id,
status="completed",
metube_job_id="vid1.vid1",
media_url="http://metube.local/download/f.mp4",
progress_percent=100,
completed_at=completed_at,
)
db_session.add(job)
db_session.commit()
await download_jobs.handle_metube_event(
db_session,
"updated",
json.dumps({"id": "vid1.vid1", "url": video.youtube_url, "status": "downloading", "percent": 12}),
)
db_session.refresh(job)
assert job.status == "completed"
assert job.progress_percent == 12
assert job.media_url == "http://metube.local/download/f.mp4"
assert job.completed_at is not None
assert job.error_message is None
await download_jobs.handle_metube_event(
db_session,
"updated",
json.dumps({"id": "vid1.vid1", "url": video.youtube_url, "status": "postprocessing", "percent": 55}),
)
db_session.refresh(job)
assert job.status == "completed"
assert job.progress_percent == 55
assert job.media_url == "http://metube.local/download/f.mp4"
# a late 'error' tick must not turn it into failed either
await download_jobs.handle_metube_event(
db_session,
"updated",
json.dumps({"id": "vid1.vid1", "url": video.youtube_url, "status": "error", "msg": "late blip"}),
)
db_session.refresh(job)
assert job.status == "completed"
assert job.media_url == "http://metube.local/download/f.mp4"
assert job.completed_at is not None
assert job.error_message is None
@pytest.mark.asyncio
async def test_updated_event_does_not_roll_back_failed_job(db_session):
video = _seed_video(db_session)
job = DownloadJob(video_id=video.id, status="failed", metube_job_id="vid1.vid1", error_message="boom")
db_session.add(job)
db_session.commit()
await download_jobs.handle_metube_event(
db_session,
"updated",
json.dumps({"id": "vid1.vid1", "url": video.youtube_url, "status": "downloading", "percent": 30}),
)
db_session.refresh(job)
assert job.status == "failed"
assert job.error_message == "boom"
@pytest.mark.asyncio
async def test_updated_event_does_not_touch_deleted_job(db_session):
video = _seed_video(db_session)
job = DownloadJob(video_id=video.id, status="deleted")
db_session.add(job)
db_session.commit()
await download_jobs.handle_metube_event(
db_session,
"updated",
json.dumps({"id": "vid1.vid1", "url": video.youtube_url, "status": "postprocessing", "percent": 10}),
)
db_session.refresh(job)
assert job.status == "deleted"
assert job.media_url is None
@pytest.mark.asyncio
async def test_authoritative_completed_still_finalizes_queued_job(monkeypatch, db_session):
video = _seed_video(db_session)
job = DownloadJob(video_id=video.id, status="queued")
db_session.add(job)
db_session.commit()
monkeypatch.setattr(
"app.services.metube_client.MeTubeClient.build_media_url",
lambda self, filename: f"http://metube.local/download/{filename}",
)
await download_jobs.handle_metube_event(
db_session,
"completed",
json.dumps({"id": "vid1.vid1", "url": video.youtube_url, "status": "finished", "filename": "f.mp4"}),
)
db_session.refresh(job)
assert job.status == "completed"
assert job.media_url == "http://metube.local/download/f.mp4"
# --- Bug 1: stale-job detection and lazy self-heal --------------------------
def test_job_is_stale_requires_both_timestamps_old(db_session):
video = _seed_video(db_session)
job = DownloadJob(video_id=video.id, status="queued")
db_session.add(job)
db_session.commit()
now = datetime.now(timezone.utc)
old = now - timedelta(minutes=10)
job.requested_at = old
job.updated_at = old
assert download_jobs.job_is_stale(job) is True
# events still flowing: fresh updated_at means no reconcile needed
job.updated_at = now
assert download_jobs.job_is_stale(job) is False
# terminal jobs are never reconciled
job.updated_at = old
job.status = "completed"
assert download_jobs.job_is_stale(job) is False
def test_reconcile_stale_job_completed_from_done(monkeypatch, db_session):
video = _seed_video(db_session)
job = DownloadJob(video_id=video.id, status="queued")
db_session.add(job)
db_session.commit()
monkeypatch.setattr(
"app.services.metube_client.MeTubeClient.fetch_history",
lambda self: {
"queue": [],
"pending": [],
"done": [{"id": "vid1.vid1", "url": video.youtube_url, "status": "finished", "filename": "f.mp4"}],
},
)
monkeypatch.setattr(
"app.services.metube_client.MeTubeClient.build_media_url",
lambda self, filename: f"http://metube.local/download/{filename}",
)
download_jobs.reconcile_stale_job(db_session, job)
db_session.refresh(job)
assert job.status == "completed"
assert job.media_url == "http://metube.local/download/f.mp4"
assert job.progress_percent == 100
def test_reconcile_stale_job_failed_from_done_error(monkeypatch, db_session):
video = _seed_video(db_session)
job = DownloadJob(video_id=video.id, status="queued")
db_session.add(job)
db_session.commit()
monkeypatch.setattr(
"app.services.metube_client.MeTubeClient.fetch_history",
lambda self: {
"queue": [],
"pending": [],
"done": [{"id": "vid1.vid1", "url": video.youtube_url, "status": "error", "msg": "HTTP 429 bot check"}],
},
)
download_jobs.reconcile_stale_job(db_session, job)
db_session.refresh(job)
assert job.status == "failed"
assert job.error_message == "HTTP 429 bot check"
def test_reconcile_stale_job_keeps_active_when_still_in_queue(monkeypatch, db_session):
video = _seed_video(db_session)
job = DownloadJob(video_id=video.id, status="queued")
db_session.add(job)
db_session.commit()
monkeypatch.setattr(
"app.services.metube_client.MeTubeClient.fetch_history",
lambda self: {
"queue": [{"id": "vid1.vid1", "url": video.youtube_url, "status": "pending"}],
"pending": [],
"done": [],
},
)
download_jobs.reconcile_stale_job(db_session, job)
db_session.refresh(job)
assert job.status == "queued"
def test_reconcile_stale_job_unknown_when_not_in_history(monkeypatch, db_session):
video = _seed_video(db_session)
job = DownloadJob(video_id=video.id, status="queued")
db_session.add(job)
db_session.commit()
monkeypatch.setattr(
"app.services.metube_client.MeTubeClient.fetch_history",
lambda self: {"queue": [], "pending": [], "done": []},
)
download_jobs.reconcile_stale_job(db_session, job)
db_session.refresh(job)
assert job.status == "unknown"
def test_reconcile_stale_job_fetch_error_leaves_status(monkeypatch, db_session):
video = _seed_video(db_session)
job = DownloadJob(video_id=video.id, status="downloading")
db_session.add(job)
db_session.commit()
def raise_error(self):
raise RuntimeError("connection refused")
monkeypatch.setattr("app.services.metube_client.MeTubeClient.fetch_history", raise_error)
download_jobs.reconcile_stale_job(db_session, job)
db_session.refresh(job)
assert job.status == "downloading"
def test_reconcile_stale_job_touches_updated_at_for_backoff(monkeypatch, db_session):
"""updated_at is bumped on every reconcile attempt (success or failure),
so the 2s polling loop triggers at most one history fetch per window."""
video = _seed_video(db_session)
old = datetime.now(timezone.utc) - timedelta(minutes=10)
job = DownloadJob(video_id=video.id, status="queued", requested_at=old, updated_at=old)
db_session.add(job)
db_session.commit()
monkeypatch.setattr(
"app.services.metube_client.MeTubeClient.fetch_history",
lambda self: {"queue": [], "pending": [], "done": []},
)
assert download_jobs.job_is_stale(job) is True
download_jobs.reconcile_stale_job(db_session, job)
assert download_jobs.job_is_stale(job) is False