Files
J621/backend/apps/library/downloads.py
T
JakeBreath c061d2681b Fix Download to Library crashing with a NameError
run_download_task sets e621_match_status from MediaItem, but the module
only imported DownloadTask — every download that carried e621 metadata
failed right after indexing, leaving the file in the library unlinked.
Verified the runner end-to-end with a stubbed fetch.
2026-09-17 19:28:19 -05:00

158 lines
5.1 KiB
Python

"""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()