"""Background 'Download to Library' jobs with progress tracking.""" import logging import threading import time from datetime import timedelta from pathlib import Path from urllib.parse import urlparse from django.conf import settings from django.db import connection from django.utils import timezone from . import services from .models import DownloadTask, MediaItem logger = logging.getLogger(__name__) STALE_AFTER = timedelta(minutes=30) def reap_stale_downloads(): """Mark tasks left hanging by a recycled worker as failed. Gunicorn recycles workers (--max-requests, timeouts); a download thread dies with its worker, so long-stuck tasks are surfaced as errors instead of pretending to run forever. """ cutoff = timezone.now() - STALE_AFTER return DownloadTask.objects.filter( status__in=[DownloadTask.STATUS_PENDING, DownloadTask.STATUS_DOWNLOADING], updated_at__lt=cutoff, ).update( status=DownloadTask.STATUS_ERROR, error="The worker restarted before this download finished.", speed=None, updated_at=timezone.now(), ) def start_download_task(task_id): thread = threading.Thread(target=run_download_task, args=(task_id,), daemon=True) thread.start() def run_download_task(task_id): task = DownloadTask.objects.filter(id=task_id).first() if task is None: return folder = Path(settings.WATCHED_FOLDER) name = ( task.filename or Path(urlparse(task.url).path).name or f"post-{task.post_id or 'download'}" ) destination = services.unique_destination(folder, name) state = { "progress_at": 0.0, "bytes_at": 0, "cancel_at": 0.0, "cancelled": False, } def should_cancel(): now = time.monotonic() if now - state["cancel_at"] >= 1.0: state["cancelled"] = DownloadTask.objects.filter( id=task_id, cancelled=True ).exists() state["cancel_at"] = now return state["cancelled"] def on_progress(downloaded, total): now = time.monotonic() if now - state["progress_at"] < 0.5 and (total == 0 or downloaded < total): return elapsed = max(now - state["progress_at"], 0.001) speed = None if state["bytes_at"] > 0 and downloaded >= state["bytes_at"]: speed = (downloaded - state["bytes_at"]) / elapsed state["progress_at"] = now state["bytes_at"] = downloaded DownloadTask.objects.filter(id=task_id).update( downloaded=downloaded, total=total, progress=int(downloaded * 100 / total) if total else 0, speed=speed, updated_at=timezone.now(), ) DownloadTask.objects.filter(id=task_id).update( status=DownloadTask.STATUS_DOWNLOADING, updated_at=timezone.now() ) try: services.download_file( task.url, destination, progress_callback=on_progress, should_cancel=should_cancel, ) item, _, location, _ = services.index_file(destination, folder) services.rename_location_to_j_id(item, location) services.ensure_visual_hashes(item) update_fields = [] if item.uploaded_by_id is None and task.user_id is not None: item.uploaded_by = task.user update_fields.append("uploaded_by") if task.e621_data: item.e621_post_id = task.post_id item.e621_data = task.e621_data item.e621_match_status = MediaItem.E621_MATCHED item.e621_checked_at = timezone.now() update_fields += [ "e621_post_id", "e621_data", "e621_match_status", "e621_checked_at", ] rating = (task.e621_data or {}).get("rating") if not item.rating and rating in {"s", "q", "e"}: item.rating = rating update_fields.append("rating") if update_fields: item.save(update_fields=update_fields + ["updated_at"]) DownloadTask.objects.filter(id=task_id).update( status=DownloadTask.STATUS_COMPLETE, progress=100, speed=None, library_item=item, updated_at=timezone.now(), ) except services.DownloadCancelled: destination.unlink(missing_ok=True) DownloadTask.objects.filter(id=task_id).update( status=DownloadTask.STATUS_CANCELLED, progress=0, downloaded=0, speed=None, updated_at=timezone.now(), ) except Exception as exc: # noqa: BLE001 - report background failures destination.unlink(missing_ok=True) logger.exception("Download task %s failed", task_id) DownloadTask.objects.filter(id=task_id).update( status=DownloadTask.STATUS_ERROR, error=str(exc), speed=None, updated_at=timezone.now(), ) finally: # Background threads hold their own DB connection; release it so # Gunicorn workers do not leak connections when threads finish. connection.close()