Files
J621/backend/apps/library/services.py
T
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

419 lines
13 KiB
Python

import hashlib
import json
import logging
import mimetypes
import os
import re
import shutil
import subprocess
from pathlib import Path
import imagehash
from django.conf import settings
from django.core import signing
from django.http import FileResponse, Http404, HttpResponse
from django.utils.text import get_valid_filename
from PIL import Image
from .models import MediaItem, MediaLocation
logger = logging.getLogger(__name__)
HASH_FIELDS = ("ahash", "dhash", "phash", "whash")
IMAGE_EXTENSIONS = {".jpg", ".jpeg", ".png", ".gif", ".apng", ".webp"}
ALLOWED_EXTENSIONS = {
".jpg",
".jpeg",
".png",
".gif",
".apng",
".webp",
".mp4",
".webm",
}
VIDEO_EXTENSIONS = {".mp4", ".webm"}
UPLOAD_FILE_SALT = "j621.upload-file"
MEDIA_FILE_SALT = "j621.media-file"
CHUNK_SIZE = 1024 * 1024
RANGE_RE = re.compile(r"bytes=(\d*)-(\d*)$")
def compute_md5(path):
digest = hashlib.md5()
with open(path, "rb") as handle:
for chunk in iter(lambda: handle.read(CHUNK_SIZE), b""):
digest.update(chunk)
return digest.hexdigest()
def index_file(path, folder):
"""Index one file into MediaItem/MediaLocation.
Returns (item, created_item, location, created_location).
"""
path = Path(path)
folder = Path(folder)
stat = path.stat()
md5 = compute_md5(path)
item, created_item = MediaItem.objects.get_or_create(
md5=md5, defaults={"size": stat.st_size}
)
if not created_item and item.size != stat.st_size:
item.size = stat.st_size
item.save(update_fields=["size", "updated_at"])
location, created_location = MediaLocation.objects.update_or_create(
item=item,
path=str(path),
defaults={
"rel_path": str(path.relative_to(folder)),
"mtime": stat.st_mtime,
},
)
return item, created_item, location, created_location
def signed_media_url(item, user, action="raw"):
"""Media URL that <img>/<video> tags can load for a signed-in user."""
base = f"/api/files/J-{item.id}/{action}/"
if user is None or not getattr(user, "is_authenticated", False):
return base
signature = signing.dumps(
{"item": item.id, "user": user.id, "action": action},
salt=MEDIA_FILE_SALT,
)
return f"{base}?sig={signature}"
def rename_location_to_j_id(item, location):
"""Name a freshly indexed copy J-<id>.<ext> inside its own folder."""
path = Path(location.path)
if not path.exists():
return location
if path.stem == f"J-{item.id}":
return location
target = unique_destination(path.parent, f"J-{item.id}{path.suffix}")
path.rename(target)
parent = Path(location.rel_path).parent
location.path = str(target)
location.rel_path = (
str(parent / target.name) if str(parent) != "." else target.name
)
location.save(update_fields=["path", "rel_path"])
return location
def compute_visual_hashes(path):
"""Perceptual hashes for an image file (empty dict for other files)."""
path = Path(path)
if path.suffix.lower() not in IMAGE_EXTENSIONS:
return {}
try:
with Image.open(path) as image:
converted = image.convert("RGB")
return {
"ahash": str(imagehash.average_hash(converted, hash_size=8)),
"dhash": str(imagehash.dhash(converted, hash_size=8)),
"phash": str(imagehash.phash(converted, hash_size=8)),
"whash": str(imagehash.whash(converted, hash_size=8)),
}
except Exception: # noqa: BLE001 - hashing must never break indexing
logger.exception("Could not compute visual hashes for %s", path)
return {}
def ensure_visual_hashes(item):
"""Fill in missing perceptual hashes for a media item."""
if all(getattr(item, field) for field in HASH_FIELDS):
return item
location = item.locations.first()
if location is None:
return item
hashes = compute_visual_hashes(location.path)
if not hashes:
return item
for field, value in hashes.items():
setattr(item, field, value)
item.save(update_fields=[*hashes.keys(), "updated_at"])
return item
def parse_tags(raw):
"""Normalize a comma-separated string or JSON list into a list of tags."""
if raw is None:
return None
if isinstance(raw, (list, tuple)):
values = raw
else:
text = str(raw).strip()
if not text:
return None
if text.startswith("["):
try:
values = json.loads(text)
except json.JSONDecodeError:
values = text.split(",")
else:
values = text.split(",")
cleaned = []
for value in values:
tag = str(value).strip()
if tag and tag not in cleaned:
cleaned.append(tag[:100])
return cleaned
def unique_destination(folder, filename):
"""Return a non-existing path inside folder for the given filename."""
folder = Path(folder)
folder.mkdir(parents=True, exist_ok=True)
filename = get_valid_filename(Path(filename).name) or "upload"
candidate = folder / filename
stem, suffix = candidate.stem, candidate.suffix
counter = 1
while candidate.exists():
candidate = folder / f"{stem}-{counter}{suffix}"
counter += 1
return candidate
class RangeFileWrapper:
"""Iterate over a limited byte range of an open file."""
def __init__(self, file, length, chunk_size=64 * 1024):
self.file = file
self.remaining = length
self.chunk_size = chunk_size
def __iter__(self):
return self
def __next__(self):
if self.remaining <= 0:
raise StopIteration
data = self.file.read(min(self.chunk_size, self.remaining))
if not data:
raise StopIteration
self.remaining -= len(data)
return data
def close(self):
self.file.close()
def serve_file(request, path, download=False):
"""Serve a file with HTTP range support (needed for video seeking)."""
path = Path(path)
if not path.is_file():
raise Http404
size = path.stat().st_size
content_type = mimetypes.guess_type(str(path))[0] or "application/octet-stream"
range_header = request.headers.get("Range", "").strip()
if range_header:
match = RANGE_RE.match(range_header)
if match:
start_raw, end_raw = match.groups()
if start_raw == "" and end_raw:
length = min(int(end_raw), size)
start, end = size - length, size - 1
else:
start = int(start_raw or 0)
end = min(int(end_raw) if end_raw else size - 1, size - 1)
if start >= size or start > end:
response = HttpResponse(status=416)
response["Content-Range"] = f"bytes */{size}"
return response
length = end - start + 1
handle = open(path, "rb")
handle.seek(start)
response = FileResponse(
RangeFileWrapper(handle, length),
status=206,
content_type=content_type,
)
response["Content-Length"] = str(length)
response["Content-Range"] = f"bytes {start}-{end}/{size}"
response["Accept-Ranges"] = "bytes"
return response
response = FileResponse(
open(path, "rb"),
content_type=content_type,
as_attachment=download,
filename=path.name,
)
response["Accept-Ranges"] = "bytes"
return response
class DownloadCancelled(Exception):
"""Raised when a streamed download is cancelled by the user."""
def download_file(
url,
destination,
progress_callback=None,
should_cancel=None,
read_timeout=60,
):
"""Stream a remote file into destination (used by Download to Library).
The read timeout bounds how long a stalled connection can block the
worker: without it a hung socket would keep a job "downloading" forever
and the cancel flag could never be observed.
"""
import requests
headers = {"User-Agent": settings.USER_AGENT}
with requests.get(
url, headers=headers, stream=True, timeout=(10, read_timeout)
) as response:
response.raise_for_status()
total = int(response.headers.get("content-length") or 0)
downloaded = 0
with open(destination, "wb") as handle:
for chunk in response.iter_content(chunk_size=CHUNK_SIZE):
if should_cancel is not None and should_cancel():
raise DownloadCancelled("download cancelled")
if chunk:
handle.write(chunk)
downloaded += len(chunk)
if progress_callback is not None:
progress_callback(downloaded, total)
if progress_callback is not None:
progress_callback(downloaded, total or downloaded)
E621_DESCRIPTION_LIMIT = 20000
def trim_e621_post(post):
"""Keep a compact, size-bounded copy of an e621 post payload."""
if not isinstance(post, dict):
return None
file_data = post.get("file") or {}
relationships = post.get("relationships") or {}
score = post.get("score") or {}
categories = post.get("tags") or {}
return {
"id": post.get("id"),
"created_at": post.get("created_at"),
"rating": post.get("rating"),
"tags": {
str(category): [str(tag)[:200] for tag in tags][:500]
for category, tags in categories.items()
if isinstance(tags, list)
},
"score": {
"up": score.get("up"),
"down": score.get("down"),
"total": score.get("total"),
},
"fav_count": post.get("fav_count"),
"comment_count": post.get("comment_count"),
"sources": [
str(source)[:500]
for source in (post.get("sources") or [])
if source
][:20],
"description": str(post.get("description") or "")[:E621_DESCRIPTION_LIMIT],
"pools": [
int(pool)
for pool in (post.get("pools") or [])
if str(pool).isdigit()
][:50],
"relationships": {
"parent_id": relationships.get("parent_id"),
"has_children": relationships.get("has_children"),
},
"file": {
"md5": file_data.get("md5"),
"ext": file_data.get("ext"),
"size": file_data.get("size"),
"width": file_data.get("width"),
"height": file_data.get("height"),
"url": file_data.get("url"),
},
"preview": {
"url": (post.get("preview") or {}).get("url"),
},
"uploader_name": post.get("uploader_name"),
}
def sanitize_iqdb_results(results):
"""Keep a compact, size-bounded copy of IQDB candidates."""
cleaned = []
for result in results[:10]:
if not isinstance(result, dict):
continue
score = result.get("score")
post_id = result.get("post_id")
tags_preview = result.get("tags_preview")
entry = {
"post_id": post_id if isinstance(post_id, int) else None,
"score": float(score) if isinstance(score, (int, float)) else None,
"preview_url": str(result.get("preview_url") or "")[:500] or None,
"rating": str(result.get("rating") or "")[:1] or None,
"md5": str(result.get("md5") or "")[:32] or None,
"score_total": (
result.get("score_total")
if isinstance(result.get("score_total"), int)
else None
),
"fav_count": (
result.get("fav_count")
if isinstance(result.get("fav_count"), int)
else None
),
"width": (
result.get("width")
if isinstance(result.get("width"), int)
else None
),
"height": (
result.get("height")
if isinstance(result.get("height"), int)
else None
),
"tags_preview": (
[str(tag)[:100] for tag in tags_preview][:8]
if isinstance(tags_preview, list)
else []
),
}
if entry["post_id"] or entry["preview_url"]:
cleaned.append(entry)
return cleaned
def generate_video_thumbnail(md5, path):
"""Extract a JPEG thumbnail from a video, cached under MEDIA_ROOT/thumbs."""
if not shutil.which("ffmpeg"):
return None
thumbs_dir = Path(settings.MEDIA_ROOT) / "thumbs"
thumbs_dir.mkdir(parents=True, exist_ok=True)
target = thumbs_dir / f"{md5}.jpg"
if target.exists() and target.stat().st_mtime >= os.path.getmtime(path):
return target
command = [
"ffmpeg",
"-y",
"-ss",
"0.5",
"-i",
str(path),
"-frames:v",
"1",
"-vf",
"scale=480:-2",
"-loglevel",
"error",
str(target),
]
try:
subprocess.run(command, check=True, capture_output=True, timeout=60)
except (subprocess.SubprocessError, OSError):
return None
return target if target.exists() else None