Files
mule-image/backend/app/routers/library.py
root 42250aa16e feat: worker diagnostics in settings panel
Adds /api/v1/library/maintenance/worker-status (Celery inspect + queue
depths + recent failed photos) and a Workers section in the Settings
dialog so users can debug stuck queues and task failures without
tailing container logs. Auto-polls every 5s while open.

Co-Authored-By: Claude Opus 4.6 (1M context) <noreply@anthropic.com>
2026-04-09 10:11:33 +02:00

404 lines
14 KiB
Python

"""
Library API router for stats, scanning, and maintenance.
The /maintenance/* endpoints are surfaced through the frontend Settings
panel. They're intentionally idempotent and operate by re-queueing the
existing Celery tasks rather than doing any heavy lifting in the
request thread.
"""
import logging
import os
import shutil
from typing import List, Optional
from fastapi import APIRouter, Depends
from pydantic import BaseModel, Field
from sqlalchemy import select, func, update
from sqlalchemy.ext.asyncio import AsyncSession
from app.database import get_db
from app.models import Photo
logger = logging.getLogger(__name__)
router = APIRouter()
# Media types we accept in the regenerate-thumbnails request body. Mirrors
# the values produced by `app.tasks.scan.get_media_type`.
_VALID_MEDIA_TYPES = {'photo', 'raw', 'heic', 'video'}
@router.get("/stats")
async def get_library_stats(db: AsyncSession = Depends(get_db)):
"""Get library statistics + per-section counts. Each section count
matches the filter the sidebar applies when you click it, so the
sidebar badges and the timeline below them stay in sync.
- all_photos: non-discarded photos + videos (matches the All
Photos section's default filter)
- rated: non-discarded with rating >= 1
- duplicates: non-discarded with is_duplicate = true
- discarded: is_discarded = true
- total_size: raw bytes across every row, including discarded
"""
not_discarded = Photo.is_discarded.is_(False)
all_photos_count = (
await db.execute(select(func.count(Photo.id)).where(not_discarded))
).scalar() or 0
rated_count = (
await db.execute(
select(func.count(Photo.id)).where(not_discarded, Photo.rating >= 1)
)
).scalar() or 0
duplicates_count = (
await db.execute(
select(func.count(Photo.id)).where(
not_discarded, Photo.is_duplicate.is_(True)
)
)
).scalar() or 0
discarded_count = (
await db.execute(
select(func.count(Photo.id)).where(Photo.is_discarded.is_(True))
)
).scalar() or 0
# Legacy split (kept for the existing /stats consumers).
photo_count = (
await db.execute(
select(func.count(Photo.id)).where(
Photo.media_type.in_(['photo', 'heic', 'raw'])
)
)
).scalar() or 0
video_count = (
await db.execute(
select(func.count(Photo.id)).where(Photo.media_type == 'video')
)
).scalar() or 0
size = (await db.execute(select(func.sum(Photo.file_size)))).scalar() or 0
return {
"all_photos": all_photos_count,
"rated": rated_count,
"duplicates": duplicates_count,
"discarded": discarded_count,
"total_photos": photo_count,
"total_videos": video_count,
"total_size": size,
"total_size_gb": round(size / (1024**3), 2) if size else 0,
}
@router.post("/scan")
async def trigger_scan():
"""Trigger full library re-scan"""
from app.tasks.scan import scan_all_source_roots
scan_all_source_roots.delay()
return {"status": "success", "message": "Library scan started"}
@router.get("/scan/status")
async def get_scan_status(db: AsyncSession = Depends(get_db)):
"""Get current scan status"""
import redis
from app.config import settings
# Connect to Redis to get scan status
r = redis.Redis.from_url(settings.redis_url)
# Get scan status from Redis (set by worker tasks)
is_scanning = r.get('scan:active') == b'true'
current_folder = r.get('scan:current_folder')
processed_files = int(r.get('scan:processed_files') or 0)
total_files = int(r.get('scan:total_files') or 0)
errors = r.lrange('scan:errors', 0, -1)
return {
"is_scanning": is_scanning,
"current_folder": current_folder.decode() if current_folder else None,
"processed_files": processed_files,
"total_files": total_files,
"errors": [e.decode() for e in errors] if errors else []
}
# ---------------------------------------------------------------------------
# Maintenance endpoints — surfaced via the Settings panel.
# ---------------------------------------------------------------------------
class RegenerateThumbnailsRequest(BaseModel):
"""Optional filters narrowing which photos get re-queued. With both
fields omitted the request resets every photo in the library."""
media_types: Optional[List[str]] = Field(
default=None,
description="Restrict to these media_type values (photo/raw/heic/video).",
)
only_failed: bool = Field(
default=False,
description="If true, only re-queue photos whose processing_status is 'failed'.",
)
@router.get("/maintenance/thumbnail-stats")
async def get_thumbnail_stats(db: AsyncSession = Depends(get_db)):
"""Counts of photos by processing_status, plus a media-type breakdown
so the Settings panel can show the user what's outstanding."""
status_rows = (
await db.execute(
select(Photo.processing_status, func.count(Photo.id)).group_by(
Photo.processing_status
)
)
).all()
media_rows = (
await db.execute(
select(Photo.media_type, func.count(Photo.id)).group_by(Photo.media_type)
)
).all()
by_status = {status or 'unknown': count for status, count in status_rows}
by_media_type = {media or 'unknown': count for media, count in media_rows}
total = sum(by_status.values())
return {
"total": total,
"pending": by_status.get('pending', 0),
"processing": by_status.get('processing', 0),
"completed": by_status.get('completed', 0),
"failed": by_status.get('failed', 0),
"by_media_type": by_media_type,
}
@router.post("/maintenance/regenerate-thumbnails")
async def regenerate_thumbnails(
body: RegenerateThumbnailsRequest,
db: AsyncSession = Depends(get_db),
):
"""Reset matching photos' on-disk thumbnail directories and re-queue
Celery thumbnail generation. Used by the Settings panel for the
'regenerate video thumbnails' / 'regenerate failed' buttons.
Files on disk are removed under /data/thumbs/<photo_id>/ so the next
request to /photos/{id}/thumb/{size} actually re-generates instead of
serving the stale placeholder.
"""
from app.tasks.thumbs import generate_thumbnails
# Validate media_types early so a typo can't silently match nothing.
media_types = body.media_types
if media_types is not None:
invalid = [m for m in media_types if m not in _VALID_MEDIA_TYPES]
if invalid:
return {
"status": "error",
"message": f"Invalid media_types: {invalid}. "
f"Allowed: {sorted(_VALID_MEDIA_TYPES)}",
}
query = select(Photo)
if media_types:
query = query.where(Photo.media_type.in_(media_types))
if body.only_failed:
query = query.where(Photo.processing_status == 'failed')
photos = (await db.execute(query)).scalars().all()
cleared_dirs = 0
file_errors = 0
for photo in photos:
thumb_dir = f"/data/thumbs/{photo.id}"
if os.path.isdir(thumb_dir):
try:
shutil.rmtree(thumb_dir)
cleared_dirs += 1
except OSError as e:
file_errors += 1
logger.warning(f"Could not clear thumb dir {thumb_dir}: {e}")
photo.processing_status = 'pending'
photo.processing_error = None
photo.thumb_small = None
photo.thumb_medium = None
photo.thumb_large = None
await db.commit()
# Queue celery tasks AFTER the commit so the worker sees the reset
# state when it picks the job up.
queued = 0
for photo in photos:
try:
generate_thumbnails.delay(photo.id)
queued += 1
except Exception as e:
logger.warning(f"Could not queue thumbnail job for {photo.id}: {e}")
return {
"status": "success",
"matched": len(photos),
"queued": queued,
"cleared_dirs": cleared_dirs,
"file_errors": file_errors,
"filters": {
"media_types": media_types,
"only_failed": body.only_failed,
},
}
@router.get("/maintenance/worker-status")
async def get_worker_status(db: AsyncSession = Depends(get_db)):
"""Diagnostics for the Celery worker fleet + recent task failures.
Surfaced in the Settings panel so the user can spot a stuck queue or
a worker that's gone away without tailing container logs. Returns:
- workers: list of {name, status, active, concurrency, queues}
derived from celery_app.control.inspect(). `status` is 'online'
when ping succeeds, 'unreachable' otherwise. Empty list means no
workers are responding at all (broker down, container crashed,
wrong queue routing, etc.).
- queues: per-queue depth read from Redis (LLEN of each queue key
used by celery.kombu). Mirrors what tasks are waiting to be
picked up.
- failures: aggregate count of photos with processing_status='failed'
plus the most recent N error messages so the user can see *why*
things failed without opening the DB.
- broker_ok: bool — could we even reach Redis?
"""
from app.tasks.celery import celery_app
from app.config import settings
import redis as _redis
# ----- Celery inspect (workers + active tasks) -------------------------
workers: list[dict] = []
inspect_error: Optional[str] = None
try:
inspect = celery_app.control.inspect(timeout=1.0)
ping = inspect.ping() or {}
active = inspect.active() or {}
reserved = inspect.reserved() or {}
scheduled = inspect.scheduled() or {}
stats = inspect.stats() or {}
active_queues = inspect.active_queues() or {}
worker_names = set(ping) | set(active) | set(stats)
for name in sorted(worker_names):
wstats = stats.get(name) or {}
pool = wstats.get('pool') or {}
workers.append({
"name": name,
"status": "online" if name in ping else "unreachable",
"active": len(active.get(name, []) or []),
"reserved": len(reserved.get(name, []) or []),
"scheduled": len(scheduled.get(name, []) or []),
"concurrency": pool.get('max-concurrency'),
"processed": (wstats.get('total') or {}),
"queues": [q.get('name') for q in (active_queues.get(name) or [])],
"active_tasks": [
{
"id": t.get('id'),
"name": t.get('name'),
"args": t.get('args'),
"time_start": t.get('time_start'),
}
for t in (active.get(name) or [])[:10]
],
})
except Exception as e:
inspect_error = str(e)
logger.warning(f"Celery inspect failed: {e}")
# ----- Broker / queue depth --------------------------------------------
broker_ok = False
queue_depths: dict[str, int] = {}
broker_error: Optional[str] = None
try:
r = _redis.Redis.from_url(settings.redis_url, socket_timeout=1.0)
r.ping()
broker_ok = True
for q in ('default', 'high', 'low'):
try:
queue_depths[q] = int(r.llen(q) or 0)
except Exception:
queue_depths[q] = 0
except Exception as e:
broker_error = str(e)
logger.warning(f"Redis broker unreachable: {e}")
# ----- Recent task failures from the photos table ----------------------
failed_total = (
await db.execute(
select(func.count(Photo.id)).where(Photo.processing_status == 'failed')
)
).scalar() or 0
recent_failed_rows = (
await db.execute(
select(
Photo.id,
Photo.filename,
Photo.media_type,
Photo.processing_error,
Photo.updated_at,
)
.where(Photo.processing_status == 'failed')
.order_by(Photo.updated_at.desc().nullslast())
.limit(20)
)
).all()
recent_failures = [
{
"photo_id": row[0],
"filename": row[1],
"media_type": row[2],
"error": (row[3] or '')[:500],
"updated_at": row[4].isoformat() if row[4] else None,
}
for row in recent_failed_rows
]
# ----- Most recent scan errors (Redis list) ----------------------------
scan_errors: list[str] = []
try:
if broker_ok:
r = _redis.Redis.from_url(settings.redis_url, socket_timeout=1.0)
raw = r.lrange('scan:errors', 0, 19) or []
scan_errors = [e.decode(errors='replace') for e in raw]
except Exception as e:
logger.debug(f"Could not read scan:errors: {e}")
return {
"broker_ok": broker_ok,
"broker_error": broker_error,
"inspect_error": inspect_error,
"workers": workers,
"worker_count": len(workers),
"queues": queue_depths,
"failures": {
"total": failed_total,
"recent": recent_failures,
},
"scan_errors": scan_errors,
}
@router.post("/maintenance/cleanup")
async def run_data_integrity_cleanup():
"""Re-run the source-roots / folders / photos data-integrity cleanup
that normally only runs on backend startup. Idempotent."""
from app.services.cleanup import cleanup_data_integrity
try:
await cleanup_data_integrity()
return {"status": "success"}
except Exception as e:
logger.error(f"Manual cleanup failed: {e}")
return {"status": "error", "message": str(e)}