2026-09-16 18:56:17 +00:00
|
|
|
import json
|
|
|
|
|
import logging
|
2026-09-27 23:36:09 +00:00
|
|
|
from datetime import datetime, timedelta, timezone
|
2026-09-16 18:56:17 +00:00
|
|
|
|
|
|
|
|
from sqlalchemy.orm import Session
|
|
|
|
|
|
2026-09-27 23:36:09 +00:00
|
|
|
from app.models.download_job import ACTIVE_STATUSES, TERMINAL_STATUSES, DownloadJob
|
2026-09-16 18:56:17 +00:00
|
|
|
from app.models.video import Video
|
|
|
|
|
from app.services.metube_client import METUBE_STATUS_MAP, MeTubeClient
|
|
|
|
|
|
|
|
|
|
logger = logging.getLogger(__name__)
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
class MeTubeRejected(Exception):
|
|
|
|
|
pass
|
|
|
|
|
|
|
|
|
|
|
2026-09-16 19:47:35 +00:00
|
|
|
class DeleteNotAllowed(Exception):
|
|
|
|
|
pass
|
|
|
|
|
|
|
|
|
|
|
2026-09-16 20:01:41 +00:00
|
|
|
class DeleteDidNotRemoveFile(Exception):
|
|
|
|
|
"""MeTube accepted the /delete request but the file is still reachable
|
|
|
|
|
afterwards -- its DELETE_FILE_ON_TRASHCAN setting is very likely off on
|
|
|
|
|
that instance, which is outside this app's control (see AGENTS.md
|
|
|
|
|
constraint 11: never touch MeTube's own config/source)."""
|
|
|
|
|
|
|
|
|
|
pass
|
|
|
|
|
|
|
|
|
|
|
2026-09-16 18:56:17 +00:00
|
|
|
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
|
|
|
|
|
|
|
|
|
|
|
2026-09-16 19:47:35 +00:00
|
|
|
def delete_local_copy(db: Session, video: Video) -> DownloadJob:
|
|
|
|
|
"""Deletes our own previously-downloaded copy from MeTube. Only ever acts
|
|
|
|
|
on a job this app itself created and tracked (never a file MeTube already
|
|
|
|
|
had before we existed) -- see AGENTS.md constraint 10."""
|
|
|
|
|
job = get_latest_job(db, video.id)
|
|
|
|
|
if job is None or job.status != "completed":
|
|
|
|
|
raise DeleteNotAllowed("No completed local copy to delete")
|
|
|
|
|
|
2026-09-16 20:01:41 +00:00
|
|
|
client = MeTubeClient()
|
2026-09-16 20:14:43 +00:00
|
|
|
client.delete_download(video.youtube_url)
|
2026-09-16 20:01:41 +00:00
|
|
|
|
|
|
|
|
# MeTube's /delete only unlinks the file if it's configured with
|
|
|
|
|
# DELETE_FILE_ON_TRASHCAN=true -- otherwise it just drops the entry from
|
|
|
|
|
# its own "done" list and the file stays put. We don't control that
|
|
|
|
|
# instance's config, so verify rather than trust the "ok" response.
|
|
|
|
|
if job.media_url and client.check_media(job.media_url):
|
|
|
|
|
raise DeleteDidNotRemoveFile(
|
|
|
|
|
"MeTube removed the entry but the file is still on disk "
|
|
|
|
|
"(its DELETE_FILE_ON_TRASHCAN setting is likely off)"
|
|
|
|
|
)
|
2026-09-16 19:47:35 +00:00
|
|
|
|
|
|
|
|
job.status = "deleted"
|
|
|
|
|
job.media_url = None
|
|
|
|
|
db.commit()
|
|
|
|
|
return job
|
|
|
|
|
|
|
|
|
|
|
2026-09-16 18:56:17 +00:00
|
|
|
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
|
|
|
|
|
|
|
|
|
|
|
2026-09-17 17:22:22 +00:00
|
|
|
def _find_job_by_key(db: Session, key: str) -> DownloadJob | None:
|
|
|
|
|
"""Match a bare MeTube store key (its 'canceled'/'cleared' events carry
|
|
|
|
|
just one key, JSON-string-encoded) to our latest job for it. The key is
|
|
|
|
|
the URL the download was enqueued under; matching the metube job id first
|
|
|
|
|
keeps us tolerant of both shapes."""
|
|
|
|
|
job = (
|
|
|
|
|
db.query(DownloadJob)
|
|
|
|
|
.filter(DownloadJob.metube_job_id == key)
|
|
|
|
|
.order_by(DownloadJob.id.desc())
|
|
|
|
|
.first()
|
|
|
|
|
)
|
|
|
|
|
if job is not None:
|
|
|
|
|
return job
|
|
|
|
|
|
|
|
|
|
video = db.query(Video).filter(Video.youtube_url == key).one_or_none()
|
|
|
|
|
if video is not None:
|
|
|
|
|
return get_latest_job(db, video.id)
|
|
|
|
|
return None
|
|
|
|
|
|
|
|
|
|
|
2026-09-16 18:56:17 +00:00
|
|
|
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
|
|
|
|
|
|
2026-09-17 17:22:22 +00:00
|
|
|
if event_name == "canceled":
|
2026-09-16 18:56:17 +00:00
|
|
|
if not isinstance(payload, str):
|
|
|
|
|
return
|
2026-09-17 17:22:22 +00:00
|
|
|
job = _find_job_by_key(db, payload)
|
|
|
|
|
if job is not None and job.status in ACTIVE_STATUSES:
|
2026-09-16 18:56:17 +00:00
|
|
|
job.status = "failed"
|
|
|
|
|
job.error_message = "Отменено в MeTube"
|
|
|
|
|
db.commit()
|
|
|
|
|
return
|
|
|
|
|
|
2026-09-17 17:22:22 +00:00
|
|
|
if event_name == "cleared":
|
|
|
|
|
# MeTube emits 'cleared' ONLY when a *done* entry is deleted (trash
|
|
|
|
|
# via /delete?where=done, or CLEAR_COMPLETED_AFTER): the payload is
|
|
|
|
|
# that entry's key (its URL) as a JSON string, one event per item.
|
|
|
|
|
# Clearing the queue emits 'canceled' instead. The removed entry is
|
|
|
|
|
# usually a download MeTube had before we existed (AGENTS.md rule 10)
|
|
|
|
|
# or one we already completed, so an unmatched/empty payload must not
|
|
|
|
|
# touch anything: failing all active jobs here would wrongly mark
|
|
|
|
|
# downloads MeTube is still running.
|
|
|
|
|
if not isinstance(payload, str) or not payload:
|
|
|
|
|
return
|
|
|
|
|
job = _find_job_by_key(db, payload)
|
|
|
|
|
if job is not None and job.status in ACTIVE_STATUSES:
|
|
|
|
|
job.status = "failed"
|
|
|
|
|
job.error_message = "Очищено в MeTube"
|
|
|
|
|
db.commit()
|
|
|
|
|
return
|
|
|
|
|
|
2026-09-16 18:56:17 +00:00
|
|
|
if not isinstance(payload, dict):
|
|
|
|
|
return
|
|
|
|
|
|
|
|
|
|
job = _find_job_by_payload(db, payload)
|
|
|
|
|
if job is None:
|
|
|
|
|
return
|
|
|
|
|
|
2026-09-16 19:09:18 +00:00
|
|
|
# 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")
|
2026-09-16 18:56:17 +00:00
|
|
|
db.commit()
|
|
|
|
|
logger.info("download_job %s updated to status=%s via '%s' event", job.id, job.status, event_name)
|
|
|
|
|
|
|
|
|
|
|
2026-09-16 19:09:18 +00:00
|
|
|
def _apply_metube_info(job: DownloadJob, info: dict, *, authoritative: bool) -> None:
|
2026-09-16 18:56:17 +00:00
|
|
|
if info.get("id"):
|
|
|
|
|
job.metube_job_id = info["id"]
|
|
|
|
|
|
2026-09-16 19:09:18 +00:00
|
|
|
metube_status = info.get("status")
|
|
|
|
|
our_status = METUBE_STATUS_MAP.get(metube_status, "unknown")
|
2026-09-16 18:56:17 +00:00
|
|
|
|
2026-09-16 19:09:18 +00:00
|
|
|
if job.started_at is None and metube_status in ("downloading", "postprocessing"):
|
2026-09-16 18:56:17 +00:00
|
|
|
job.started_at = datetime.now(timezone.utc)
|
|
|
|
|
|
|
|
|
|
percent = info.get("percent")
|
|
|
|
|
if isinstance(percent, (int, float)):
|
|
|
|
|
job.progress_percent = int(percent)
|
|
|
|
|
|
2026-09-27 23:36:09 +00:00
|
|
|
# Late, delayed 'updated' events must never roll a terminal job back (or
|
|
|
|
|
# forward into 'failed'): Socket.IO delivers events out of order, so a
|
|
|
|
|
# lagging progress tick often arrives after the authoritative 'completed'.
|
|
|
|
|
# Progress/metube_job_id above are still taken; the status (and its
|
|
|
|
|
# error_message/completed_at/media_url) is left alone.
|
|
|
|
|
if not authoritative and job.status in TERMINAL_STATUSES:
|
|
|
|
|
return
|
|
|
|
|
|
2026-09-16 18:56:17 +00:00
|
|
|
if our_status == "failed":
|
|
|
|
|
job.status = "failed"
|
|
|
|
|
job.error_message = info.get("msg") or info.get("error") or "Ошибка загрузки"
|
|
|
|
|
elif our_status == "completed":
|
2026-09-16 19:09:18 +00:00
|
|
|
if not authoritative:
|
|
|
|
|
# A stream finished, not necessarily the whole job (see comment
|
|
|
|
|
# above) -- keep the current status, just take the progress tick.
|
|
|
|
|
return
|
2026-09-16 18:56:17 +00:00
|
|
|
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)
|
2026-09-16 19:09:18 +00:00
|
|
|
elif our_status != "unknown":
|
2026-09-16 18:56:17 +00:00
|
|
|
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
|
|
|
|
|
|
2026-09-16 19:09:18 +00:00
|
|
|
# (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)):
|
2026-09-16 18:56:17 +00:00
|
|
|
for item in history.get(bucket, []) or []:
|
|
|
|
|
url = item.get("url")
|
|
|
|
|
if url:
|
2026-09-16 19:09:18 +00:00
|
|
|
by_url[url] = (item, authoritative)
|
2026-09-16 18:56:17 +00:00
|
|
|
|
|
|
|
|
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)
|
2026-09-16 19:09:18 +00:00
|
|
|
match = by_url.get(video.youtube_url) if video else None
|
|
|
|
|
if match is None:
|
2026-09-16 18:56:17 +00:00
|
|
|
job.status = "unknown"
|
|
|
|
|
continue
|
2026-09-16 19:09:18 +00:00
|
|
|
info, authoritative = match
|
|
|
|
|
_apply_metube_info(job, info, authoritative=authoritative)
|
2026-09-16 18:56:17 +00:00
|
|
|
|
|
|
|
|
db.commit()
|
|
|
|
|
logger.info("Reconciled %d download job(s) against MeTube history", len(jobs))
|
2026-09-27 23:36:09 +00:00
|
|
|
|
|
|
|
|
|
|
|
|
|
def _as_utc(value: datetime) -> datetime | None:
|
|
|
|
|
"""SQLite stores/returns naive datetimes regardless of the column's
|
|
|
|
|
timezone=True, while Postgres returns aware ones -- normalize before any
|
|
|
|
|
comparison with datetime.now(timezone.utc)."""
|
|
|
|
|
if value is None:
|
|
|
|
|
return None
|
|
|
|
|
if value.tzinfo is None:
|
|
|
|
|
return value.replace(tzinfo=timezone.utc)
|
|
|
|
|
return value.astimezone(timezone.utc)
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
def job_is_stale(job: DownloadJob, threshold_minutes: int = 5) -> bool:
|
|
|
|
|
"""True when an active job had no events for a while: both requested_at and
|
|
|
|
|
updated_at are older than the threshold. The updated_at leg implements the
|
|
|
|
|
'reconcile once per window' backoff -- while events keep flowing, updated_at
|
|
|
|
|
stays fresh and no MeTube history fetch is needed."""
|
|
|
|
|
if job.status not in ACTIVE_STATUSES:
|
|
|
|
|
return False
|
|
|
|
|
now = datetime.now(timezone.utc)
|
|
|
|
|
threshold = timedelta(minutes=threshold_minutes)
|
|
|
|
|
requested_at = _as_utc(job.requested_at)
|
|
|
|
|
updated_at = _as_utc(job.updated_at)
|
|
|
|
|
if requested_at is None or updated_at is None:
|
|
|
|
|
return False
|
|
|
|
|
return now - requested_at > threshold and now - updated_at > threshold
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
def reconcile_stale_job(db: Session, job: DownloadJob) -> None:
|
|
|
|
|
"""One-shot lazy self-heal for an active job whose Socket.IO events stopped
|
|
|
|
|
(the backend may have missed the terminal event entirely). Mirrors
|
|
|
|
|
reconcile_on_startup's matching, but is called from the download-status
|
|
|
|
|
endpoint while the service is running:
|
|
|
|
|
- 'done' entry is authoritative (finished -> completed, error -> failed);
|
|
|
|
|
- 'queue'/'pending' entries keep the job active (still in flight);
|
|
|
|
|
- no entry at all -> 'unknown' (UI shows a retry action);
|
|
|
|
|
- history fetch failure is transient here (not a restart) -> status is left
|
|
|
|
|
untouched, so the polling client keeps its last known state.
|
|
|
|
|
Always touches updated_at and commits: that is the backoff that prevents
|
|
|
|
|
the 2s polling loop from hammering /history (next reconcile only after
|
|
|
|
|
the stale window has passed again)."""
|
|
|
|
|
client = MeTubeClient()
|
|
|
|
|
try:
|
|
|
|
|
history = client.fetch_history()
|
|
|
|
|
except Exception:
|
|
|
|
|
logger.warning("Could not fetch MeTube history to self-heal stale job %s", job.id, exc_info=True)
|
|
|
|
|
job.updated_at = datetime.now(timezone.utc)
|
|
|
|
|
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)
|
|
|
|
|
|
|
|
|
|
video = db.get(Video, job.video_id)
|
|
|
|
|
match = by_url.get(video.youtube_url) if video else None
|
|
|
|
|
if match is None:
|
|
|
|
|
job.status = "unknown"
|
|
|
|
|
else:
|
|
|
|
|
info, authoritative = match
|
|
|
|
|
_apply_metube_info(job, info, authoritative=authoritative)
|
|
|
|
|
|
|
|
|
|
job.updated_at = datetime.now(timezone.utc)
|
|
|
|
|
db.commit()
|
|
|
|
|
logger.info("Self-healed stale download job %s to status=%s", job.id, job.status)
|