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.
158 lines
5.1 KiB
Python
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()
|