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

246 lines
7.3 KiB
Python

"""Runtime statistics for the staff dashboard.
Everything here is read-only and cheap enough to poll every couple of
seconds: psutil counters, an optional nvidia-smi query, the cached storage
numbers and the tail of the application log file.
"""
import logging
import os
import subprocess
import time
from pathlib import Path
import psutil
from django.conf import settings
logger = logging.getLogger(__name__)
NVIDIA_QUERY = (
"name,utilization.gpu,memory.used,memory.total,temperature.gpu"
)
LOG_TAIL_LINES = 140
RECENT_LIMIT = 6
# Prime the non-blocking CPU counters so the first request already has a delta.
psutil.cpu_percent(interval=None)
psutil.cpu_percent(interval=None, percpu=True)
def _as_int(value):
try:
return int(float(str(value).strip()))
except (TypeError, ValueError):
return None
def system_stats():
"""CPU, memory and uptime counters (non-blocking)."""
memory = psutil.virtual_memory()
swap = psutil.swap_memory()
try:
load = [round(value, 2) for value in psutil.getloadavg()]
except (AttributeError, OSError):
load = []
return {
"cpu": {
"percent": psutil.cpu_percent(interval=None),
"per_cpu": psutil.cpu_percent(interval=None, percpu=True),
"cores": psutil.cpu_count(logical=True),
"physical_cores": psutil.cpu_count(logical=False),
"load": load,
},
"memory": {
"total": memory.total,
"used": memory.used,
"available": memory.available,
"percent": memory.percent,
"swap_total": swap.total,
"swap_used": swap.used,
"swap_percent": swap.percent,
},
"uptime_seconds": int(time.time() - psutil.boot_time()),
}
def gpu_stats():
"""NVIDIA GPUs via nvidia-smi, or an empty list when unavailable."""
try:
result = subprocess.run(
[
"nvidia-smi",
f"--query-gpu={NVIDIA_QUERY}",
"--format=csv,noheader,nounits",
],
capture_output=True,
text=True,
timeout=4,
check=False,
)
except (OSError, subprocess.SubprocessError) as exc:
logger.debug("nvidia-smi unavailable: %s", exc)
return []
if result.returncode != 0:
return []
gpus = []
for line in result.stdout.strip().splitlines():
parts = [part.strip() for part in line.split(",")]
if len(parts) != 5:
continue
name, utilization, memory_used, memory_total, temperature = parts
used_mib = _as_int(memory_used)
total_mib = _as_int(memory_total)
gpus.append(
{
"name": name,
"utilization": _as_int(utilization),
"memory_used": used_mib * 1024 * 1024
if used_mib is not None
else None,
"memory_total": total_mib * 1024 * 1024
if total_mib is not None
else None,
"temperature": _as_int(temperature),
}
)
return gpus
def _download_job(task):
label = (
f"Post #{task.post_id}"
if task.post_id
else (task.filename or "Download")
)
detail = f"{task.downloaded} / {task.total}" if task.total else task.filename
return {
"kind": "download",
"id": str(task.id),
"label": label,
"status": task.status,
"progress": task.progress,
"downloaded": task.downloaded,
"total": task.total,
"detail": detail,
"cancelled": task.cancelled,
"created_at": task.created_at,
"updated_at": task.updated_at,
}
def _match_job(task):
progress = (
int(task.processed * 100 / task.total) if task.total else 0
)
return {
"kind": "match",
"id": str(task.id),
"label": f"e621 match scan ({task.scope})",
"status": task.status,
"progress": progress,
"processed": task.processed,
"total": task.total,
"detail": (
f"{task.processed}/{task.total} · {task.matched} matched · "
f"{task.not_found} not found · {task.deleted} deleted"
),
"cancelled": task.cancelled,
"created_at": task.created_at,
"updated_at": task.updated_at,
}
def jobs_stats():
"""Running/pending jobs plus the most recently finished ones."""
from apps.library.downloads import reap_stale_downloads
from apps.library.matching import reap_stale_match_tasks
from apps.library.models import DownloadTask, MatchTask
# Tasks whose runner died in a restart must not linger as "downloading".
reap_stale_downloads()
reap_stale_match_tasks()
active_statuses_download = [
DownloadTask.STATUS_PENDING,
DownloadTask.STATUS_DOWNLOADING,
]
active_statuses_match = [
MatchTask.STATUS_PENDING,
MatchTask.STATUS_RUNNING,
]
finished_download = [
DownloadTask.STATUS_COMPLETE,
DownloadTask.STATUS_ERROR,
DownloadTask.STATUS_CANCELLED,
]
finished_match = [
MatchTask.STATUS_COMPLETE,
MatchTask.STATUS_ERROR,
MatchTask.STATUS_CANCELLED,
]
active = [
_download_job(task)
for task in DownloadTask.objects.filter(
status__in=active_statuses_download
).order_by("created_at")[:RECENT_LIMIT]
]
active += [
_match_job(task)
for task in MatchTask.objects.filter(
status__in=active_statuses_match
).order_by("created_at")[:RECENT_LIMIT]
]
recent = []
for task in DownloadTask.objects.filter(
status__in=finished_download
).order_by("-updated_at")[:RECENT_LIMIT]:
entry = _download_job(task)
entry["summary"] = (
f"J-{task.library_item_id}"
if task.library_item_id
else (task.error[:160] if task.error else task.status)
)
recent.append(entry)
for task in MatchTask.objects.filter(
status__in=finished_match
).order_by("-updated_at")[:RECENT_LIMIT]:
entry = _match_job(task)
entry["summary"] = task.error[:160] if task.error else entry["detail"]
recent.append(entry)
recent.sort(key=lambda entry: entry["updated_at"], reverse=True)
return {
"active": active,
"recent": recent[:RECENT_LIMIT],
"active_count": len(active),
}
def tail_lines(path, count=LOG_TAIL_LINES):
"""Last `count` lines of a file without reading all of it."""
with open(path, "rb") as handle:
handle.seek(0, os.SEEK_END)
size = handle.tell()
block = 16384
data = b""
while size > 0 and data.count(b"\n") <= count:
step = min(block, size)
size -= step
handle.seek(size)
data = handle.read(step) + data
return data.decode("utf-8", errors="replace").splitlines()[-count:]
def log_tail():
path = Path(getattr(settings, "LOG_FILE", ""))
if not path or not path.exists():
return {"path": str(path), "lines": [], "missing": True}
try:
lines = tail_lines(path)
except OSError as exc:
logger.warning("Could not read the log file: %s", exc)
return {"path": str(path), "lines": [], "missing": True}
return {"path": str(path), "lines": lines, "missing": False}