Files
JakeBreath 27cbfd882c Add job cancellation to the stats dashboard
- Active jobs on /stats get a cancel button wired to the existing
  download/match cancel endpoints, showing "cancelling..." and an inline
  error when the task already finished.
- Cancelling now sets the status immediately, so a task whose runner died
  in a restart stops showing as "downloading".
- Download streams use a bounded read timeout (10 s connect / 60 s read):
  a stalled socket fails within a minute (previously it could block
  forever), and a task cancelled while stalled is marked cancelled rather
  than error.
- The stats job list reaps stale download/match tasks, so phantom jobs
  never appear on the dashboard.
2026-09-17 21:41:54 -05:00

170 lines
5.7 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)
if DownloadTask.objects.filter(id=task_id, cancelled=True).exists():
# Cancel was requested while the socket was stalled; the read
# timeout is what breaks the worker out of it.
DownloadTask.objects.filter(id=task_id).update(
status=DownloadTask.STATUS_CANCELLED,
progress=0,
downloaded=0,
speed=None,
error="",
updated_at=timezone.now(),
)
else:
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()