Files
J621/backend/apps/library/downloads.py
T
JakeBreath cd490b0a23 Duplicates, delete & storage, users page with J-ID avatars
Backend:
- Perceptual hashes (aHash/dHash/pHash/wHash via imagehash, no imgdd)
  stored on items, computed on upload/download and by the new
  compute_visual_hashes command
- Duplicates API: exact duplicates (multi-location items), visual matches
  for one item, union-find similarity groups with pagination
- Delete API with ownership/staff checks, per-item and per-copy deletion,
  watched-folder path validation; storage overview and temp cleanup;
  file list accepts j_ids batches
- Staged uploads are flagged visual_match with their library matches
  (threshold via VISUAL_MATCH_THRESHOLD)
- Staff users API: list with upload counts, set role and avatar by J-ID;
  User.avatar FK with signed avatar URLs
- Download threads close their DB connection and stale tasks are reaped,
  keeping behaviour Gunicorn-friendly

Frontend:
- /duplicates: exact duplicate groups with per-copy delete, visual
  similarity controls, search similar to a J-ID, paginated groups with
  selection, bulk delete and dismiss
- /delete: storage cards, delete by J-ID with preview grid, temp cleanup
- /users: staff directory with role selects and avatar J-ID inputs
- Nav + command palette entries; top-bar avatar; upload cards and the
  metadata modal show library visual matches
2026-09-17 12:49:10 -05:00

151 lines
4.9 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
update_fields += ["e621_post_id", "e621_data"]
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()