Files
J621/backend/apps/library/downloads.py
T
JakeBreath 09405d1a0f Match local files to e621: MD5 lookups, manual links, batch scans
- MediaItem gains e621_match_status (unknown/matched/not_found/deleted)
  and e621_checked_at, backfilled for existing matched items.
- Server-side e621 client (apps/library/e621.py) using the user's stored
  credentials, throttled to 2 req/s, with typed errors.
- Matching service: MD5 lookup, manual post linking (flags MD5
  mismatches), unlink, metadata refresh, deleted-post detection.
- Detail actions POST /api/files/J-x/match/ and /unlink/ (uploader or
  staff only).
- Background library scans: MatchTask + /api/matches/ with missing/all
  scopes, progress polling, cancel and stale-task reaping; the scan
  counts toward the footer's Active Workers. Same pass available as
  manage.py match_e621 for cron.
- Library gains not_found/deleted status filters; the detail page adds
  an e621 match card (check / link by post ID / unlink) and the metadata
  card warns when a post was deleted on e621.
2026-09-17 13:41:08 -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
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()