myYouTube/backend/app/services/download_jobs.py

196 lines
6.8 KiB
Python
Raw Normal View History

import json
import logging
from datetime import datetime, timezone
from sqlalchemy.orm import Session
from app.models.download_job import ACTIVE_STATUSES, DownloadJob
from app.models.video import Video
from app.services.metube_client import METUBE_STATUS_MAP, MeTubeClient
logger = logging.getLogger(__name__)
class MeTubeRejected(Exception):
pass
def get_latest_job(db: Session, video_id: int) -> DownloadJob | None:
return (
db.query(DownloadJob)
.filter(DownloadJob.video_id == video_id)
.order_by(DownloadJob.requested_at.desc(), DownloadJob.id.desc())
.first()
)
def latest_jobs_map(db: Session, video_ids: list[int]) -> dict[int, DownloadJob]:
if not video_ids:
return {}
rows = (
db.query(DownloadJob)
.filter(DownloadJob.video_id.in_(video_ids))
.order_by(DownloadJob.requested_at.desc(), DownloadJob.id.desc())
.all()
)
result: dict[int, DownloadJob] = {}
for row in rows:
result.setdefault(row.video_id, row)
return result
def request_download(db: Session, video: Video) -> DownloadJob:
existing = get_latest_job(db, video.id)
if existing is not None and existing.status in ACTIVE_STATUSES:
return existing
result = MeTubeClient().enqueue_video(video.youtube_url, video.youtube_video_id)
if result.get("status") == "error":
raise MeTubeRejected(result.get("msg") or "MeTube rejected the download")
job = DownloadJob(video_id=video.id, status="queued")
db.add(job)
db.commit()
db.refresh(job)
return job
def _find_job_by_payload(db: Session, payload: dict) -> DownloadJob | None:
metube_id = payload.get("id")
if metube_id:
job = (
db.query(DownloadJob)
.filter(DownloadJob.metube_job_id == metube_id)
.order_by(DownloadJob.id.desc())
.first()
)
if job is not None:
return job
url = payload.get("url")
if url:
video = db.query(Video).filter(Video.youtube_url == url).one_or_none()
if video is not None:
job = get_latest_job(db, video.id)
if job is not None:
return job
return None
async def handle_metube_event(db: Session, event_name: str, raw_payload) -> None:
try:
payload = json.loads(raw_payload) if isinstance(raw_payload, str) else raw_payload
except (TypeError, ValueError):
logger.warning("Could not parse MeTube event payload for %s: %r", event_name, raw_payload)
return
if event_name in ("canceled", "cleared"):
if not isinstance(payload, str):
return
job = (
db.query(DownloadJob)
.filter(DownloadJob.metube_job_id == payload)
.order_by(DownloadJob.id.desc())
.first()
)
if job is not None and event_name == "canceled" and job.status in ACTIVE_STATUSES:
job.status = "failed"
job.error_message = "Отменено в MeTube"
db.commit()
return
if not isinstance(payload, dict):
return
job = _find_job_by_payload(db, payload)
if job is None:
return
# 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, *, authoritative: bool) -> None:
if info.get("id"):
job.metube_job_id = info["id"]
metube_status = info.get("status")
our_status = METUBE_STATUS_MAP.get(metube_status, "unknown")
if job.started_at is None and metube_status in ("downloading", "postprocessing"):
job.started_at = datetime.now(timezone.utc)
percent = info.get("percent")
if isinstance(percent, (int, float)):
job.progress_percent = int(percent)
if our_status == "failed":
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
job.media_url = MeTubeClient().build_media_url(filename)
job.status = "completed"
job.progress_percent = 100
job.completed_at = job.completed_at or datetime.now(timezone.utc)
elif our_status != "unknown":
job.status = our_status
def reconcile_on_startup(db: Session) -> None:
"""Section 19: after a restart we don't trust in-flight jobs until we've
checked MeTube's current state. Completed jobs with a media_url are left
alone — they don't need MeTube's history to remain trustworthy."""
stale_statuses = list(ACTIVE_STATUSES) + ["unknown"]
jobs = db.query(DownloadJob).filter(DownloadJob.status.in_(stale_statuses)).all()
if not jobs:
return
client = MeTubeClient()
try:
history = client.fetch_history()
except Exception:
logger.warning("Could not fetch MeTube history for reconciliation", exc_info=True)
for job in jobs:
job.status = "unknown"
db.commit()
return
# (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, 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)
match = by_url.get(video.youtube_url) if video else None
if match is None:
job.status = "unknown"
continue
info, authoritative = match
_apply_metube_info(job, info, authoritative=authoritative)
db.commit()
logger.info("Reconciled %d download job(s) against MeTube history", len(jobs))