diff --git a/backend/app/routers/library.py b/backend/app/routers/library.py index 3c3f52d..a418340 100644 --- a/backend/app/routers/library.py +++ b/backend/app/routers/library.py @@ -252,6 +252,144 @@ async def regenerate_thumbnails( } +@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 diff --git a/frontend/src/components/dialogs/SettingsDialog.tsx b/frontend/src/components/dialogs/SettingsDialog.tsx index 77560cb..93f6247 100644 --- a/frontend/src/components/dialogs/SettingsDialog.tsx +++ b/frontend/src/components/dialogs/SettingsDialog.tsx @@ -8,6 +8,9 @@ import { AlertTriangle, Database, Loader2, + Cpu, + AlertCircle, + CheckCircle2, } from 'lucide-react' import clsx from 'clsx' import { @@ -15,6 +18,7 @@ import { type ThumbnailStats, type LibraryStats, type MediaType, + type WorkerStatus, } from '../../services/api' import { toast } from '../ToastContainer' @@ -37,7 +41,10 @@ interface SettingsDialogProps { export function SettingsDialog({ isOpen, onClose }: SettingsDialogProps) { const [thumbStats, setThumbStats] = useState(null) const [libStats, setLibStats] = useState(null) + const [workerStatus, setWorkerStatus] = useState(null) const [loadingStats, setLoadingStats] = useState(false) + const [loadingWorkers, setLoadingWorkers] = useState(false) + const [showAllErrors, setShowAllErrors] = useState(false) // One key per action so each button has its own spinner without // blocking the others. const [busy, setBusy] = useState>({}) @@ -59,16 +66,38 @@ export function SettingsDialog({ isOpen, onClose }: SettingsDialogProps) { } }, []) - // Esc closes; load stats when opened. + const refreshWorkers = useCallback(async () => { + setLoadingWorkers(true) + try { + const ws = await library.maintenance.workerStatus() + setWorkerStatus(ws) + } catch (e) { + console.error('Failed to load worker status', e) + toast.error('Could not load worker status') + } finally { + setLoadingWorkers(false) + } + }, []) + + // Esc closes; load stats when opened. Workers section auto-polls + // every 5s while the dialog is open so the user sees live worker + // activity without manually hammering the refresh button. useEffect(() => { if (!isOpen) return refreshStats() + refreshWorkers() const handler = (e: KeyboardEvent) => { if (e.key === 'Escape') onClose() } window.addEventListener('keydown', handler) - return () => window.removeEventListener('keydown', handler) - }, [isOpen, onClose, refreshStats]) + const poll = window.setInterval(() => { + refreshWorkers() + }, 5000) + return () => { + window.removeEventListener('keydown', handler) + window.clearInterval(poll) + } + }, [isOpen, onClose, refreshStats, refreshWorkers]) const runAction = useCallback( async ( @@ -254,6 +283,235 @@ export function SettingsDialog({ isOpen, onClose }: SettingsDialogProps) { + {/* ----------------------------------------------------- */} + {/* Worker fleet diagnostics */} + {/* ----------------------------------------------------- */} +
} + title="Workers" + right={ + + } + > + {/* Top-line health */} +
+ 0 + ? 'ok' + : 'warn' + : 'muted' + } + /> + + 0 + ? 'warn' + : 'muted' + } + /> +
+ + {/* Inline error banners for the obvious failure modes */} + {workerStatus?.broker_error && ( + + )} + {workerStatus?.inspect_error && ( + + )} + {workerStatus && + workerStatus.broker_ok && + workerStatus.worker_count === 0 && ( + + )} + + {/* Queue depth */} + {workerStatus && ( +
+
+ Queue depth +
+
+ {Object.entries(workerStatus.queues).map(([name, depth]) => ( +
+ {name} + 0 ? 'text-text' : 'text-text-muted' + )} + > + {depth} + +
+ ))} +
+
+ )} + + {/* Per-worker breakdown */} + {workerStatus && workerStatus.workers.length > 0 && ( +
+
+ Worker fleet +
+ {workerStatus.workers.map((w) => ( +
+
+
+ {w.status === 'online' ? ( + + ) : ( + + )} + {w.name} +
+ + {w.active}/{w.concurrency ?? '?'} active + +
+
+ reserved: {w.reserved} + scheduled: {w.scheduled} + {w.queues.length > 0 && ( + queues: {w.queues.join(', ')} + )} +
+ {w.active_tasks.length > 0 && ( +
+ {w.active_tasks.map((t) => ( +
+ {t.name}{' '} + {Array.isArray(t.args) + ? t.args.map((a) => String(a)).join(', ') + : ''} +
+ ))} +
+ )} +
+ ))} +
+ )} + + {/* Recent task failures */} + {workerStatus && workerStatus.failures.recent.length > 0 && ( +
+
+
+ Recent failures ({workerStatus.failures.total}) +
+ {workerStatus.failures.recent.length > 5 && ( + + )} +
+
+ {(showAllErrors + ? workerStatus.failures.recent + : workerStatus.failures.recent.slice(0, 5) + ).map((f) => ( +
+
+ + {f.filename} + + + {f.media_type} + +
+
+ {f.error || '(no error message)'} +
+
+ ))} +
+
+ )} + + {/* Recent scan errors (Redis list) */} + {workerStatus && workerStatus.scan_errors.length > 0 && ( +
+
+ Scan errors +
+
+ {workerStatus.scan_errors.map((e, i) => ( +
+ {e} +
+ ))} +
+
+ )} + + {!workerStatus && ( +
+ + Loading worker status… +
+ )} +
+ {/* ----------------------------------------------------- */} {/* Data integrity */} {/* ----------------------------------------------------- */} @@ -345,6 +603,20 @@ function Stat({ ) } +function ErrorBanner({ title, detail }: { title: string; detail: string }) { + return ( +
+
+ + {title} +
+
+ {detail} +
+
+ ) +} + function ActionButton({ loading, disabled, diff --git a/frontend/src/services/api.ts b/frontend/src/services/api.ts index 68e4c85..eddfb4f 100644 --- a/frontend/src/services/api.ts +++ b/frontend/src/services/api.ts @@ -230,6 +230,45 @@ export interface RegenerateResult { } } +export interface WorkerInfo { + name: string + status: 'online' | 'unreachable' + active: number + reserved: number + scheduled: number + concurrency: number | null + processed: Record + queues: string[] + active_tasks: Array<{ + id: string + name: string + args: unknown + time_start: number | null + }> +} + +export interface WorkerFailure { + photo_id: string + filename: string + media_type: string + error: string + updated_at: string | null +} + +export interface WorkerStatus { + broker_ok: boolean + broker_error: string | null + inspect_error: string | null + workers: WorkerInfo[] + worker_count: number + queues: Record + failures: { + total: number + recent: WorkerFailure[] + } + scan_errors: string[] +} + export const library = { scan: async () => { const response = await api.post('/library/scan') @@ -265,6 +304,14 @@ export const library = { return response.data }, + /** Celery worker fleet diagnostics + recent task failures. Surfaced + * in the Settings panel so users can debug stuck queues without + * tailing container logs. */ + workerStatus: async (): Promise => { + const response = await api.get('/library/maintenance/worker-status') + return response.data + }, + /** Re-run the source-roots / folders / photos integrity cleanup that * normally runs on backend startup. */ cleanup: async (): Promise<{ status: string; message?: string }> => {