Files
mule-image/backend/app/services/cleanup.py
root 697343646a feat: prune orphaned photo rows + retry-pending action
Adds /api/v1/library/maintenance/{missing-stats,prune-missing} backed
by a new cleanup helper that deletes Photo rows whose files no longer
exist on disk under a *mounted* source root. Skips photos under
unmounted roots so a temporarily-disconnected drive doesn't get
silently nuked.

Settings panel surfaces the orphan count with a destructive Prune
button, plus a "Kick pending" action that re-queues photos stuck in
processing_status='pending' (typically left behind when the scanner
created the row but the worker never picked up the thumbnail task).

Common trigger: PHOTO_DIRS in .env was repointed at a different
library root, leaving every old row dangling.

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

231 lines
8.4 KiB
Python

"""
One-shot data integrity cleanup for source_roots / folders / photos.
Earlier versions of the scanner stored paths verbatim, so trailing slashes
and redundant separators produced duplicate SourceRoot and Folder rows for
the same physical directory. The watcher also auto-created source roots
when fired with a parent dir. This module merges the duplicates and
re-points photos to the canonical folder so the data lines up with the
post-fix scanner.
Idempotent: safe to run on every backend startup.
"""
import os
import logging
from datetime import datetime
from sqlalchemy import select, update, func
from sqlalchemy.ext.asyncio import AsyncSession
from app.database import AsyncSessionLocal
from app.models import Photo, Folder, SourceRoot
logger = logging.getLogger(__name__)
def _normalize_path(path: str) -> str:
return os.path.normpath(path)
async def _dedupe_source_roots(session: AsyncSession) -> int:
"""Group source roots by normalized path and merge duplicates. Returns
the number of rows deleted."""
result = await session.execute(select(SourceRoot))
rows = result.scalars().all()
groups: dict[str, list[SourceRoot]] = {}
for sr in rows:
norm = _normalize_path(sr.path)
groups.setdefault(norm, []).append(sr)
deleted = 0
for norm, srs in groups.items():
if len(srs) == 1:
# Make sure the canonical row's path is normalized too.
if srs[0].path != norm:
srs[0].path = norm
continue
# Pick the canonical row: prefer one with a non-empty name and the
# earliest added_at (most likely the original).
canonical = sorted(
srs,
key=lambda s: (not bool(s.name), s.added_at or datetime.max),
)[0]
canonical.path = norm
for sr in srs:
if sr.id == canonical.id:
continue
# Re-point folders that referenced the duplicate root.
await session.execute(
update(Folder)
.where(Folder.source_root_id == sr.id)
.values(source_root_id=canonical.id)
)
await session.delete(sr)
deleted += 1
return deleted
async def _dedupe_folders(session: AsyncSession) -> int:
"""Group folders by normalized path and merge duplicates. Returns the
number of rows deleted."""
result = await session.execute(select(Folder))
rows = result.scalars().all()
groups: dict[str, list[Folder]] = {}
for f in rows:
norm = _normalize_path(f.path)
groups.setdefault(norm, []).append(f)
deleted = 0
for norm, folders in groups.items():
if len(folders) == 1:
if folders[0].path != norm:
folders[0].path = norm
continue
# Canonical = the one with the most photos already attached, then
# the lowest-id (deterministic tiebreaker).
canonical = sorted(
folders,
key=lambda f: (-(f.photo_count or 0), f.id),
)[0]
canonical.path = norm
for f in folders:
if f.id == canonical.id:
continue
# Re-point photos to the canonical folder.
await session.execute(
update(Photo)
.where(Photo.folder_id == f.id)
.values(folder_id=canonical.id)
)
await session.delete(f)
deleted += 1
return deleted
async def _recompute_folder_counts(session: AsyncSession) -> None:
"""Set folder.photo_count to the actual non-discarded photo count."""
result = await session.execute(select(Folder))
folders = result.scalars().all()
for f in folders:
count_result = await session.execute(
select(func.count(Photo.id)).where(
Photo.folder_id == f.id,
Photo.is_discarded == False, # noqa: E712
)
)
f.photo_count = int(count_result.scalar() or 0)
async def _warn_stale_source_roots(session: AsyncSession) -> int:
"""Log a warning for any active source root whose path no longer exists
on disk. Doesn't delete — a missing path could be a temporarily
unmounted drive, and silently dropping user data is worse than
surfacing a noisy log line.
"""
result = await session.execute(select(SourceRoot))
rows = result.scalars().all()
stale = 0
for sr in rows:
if not os.path.isdir(sr.path):
stale += 1
logger.warning(
f"Source root '{sr.name}' path is missing on disk: {sr.path} "
f"— is the docker mount still in place? "
f"(Edit docker-compose.yml or PHOTO_DIRS in .env to fix.)"
)
return stale
async def find_missing_photos(session: AsyncSession) -> tuple[list[str], list[str]]:
"""Walk every non-discarded photo and check whether its file is still
on disk. Returns (deletable_ids, skipped_under_unmounted_roots).
Skipped rows are photos whose owning source_root path itself doesn't
resolve — that's almost always an unmounted drive, and silently
deleting those rows would be data loss. The caller can surface the
skip count separately so the user knows the cleanup wasn't a no-op
by accident.
"""
sr_rows = (await session.execute(select(SourceRoot))).scalars().all()
sr_mounted: dict[str, bool] = {sr.id: os.path.isdir(sr.path) for sr in sr_rows}
photos = (await session.execute(
select(Photo.id, Photo.filepath, Photo.folder_id)
.where(Photo.is_discarded.is_(False))
)).all()
# folder -> source_root lookup
folders = (await session.execute(select(Folder.id, Folder.source_root_id))).all()
folder_to_sr = {fid: srid for fid, srid in folders}
deletable: list[str] = []
skipped: list[str] = []
for pid, fp, folder_id in photos:
sr_id = folder_to_sr.get(folder_id)
if sr_id is None or not sr_mounted.get(sr_id, False):
skipped.append(pid)
continue
if not os.path.exists(fp):
deletable.append(pid)
return deletable, skipped
async def prune_missing_photos(dry_run: bool = True) -> dict:
"""Delete photo rows whose files are no longer on disk *and* whose
source root is currently mounted. Common cause: PHOTO_DIRS in .env
was repointed at a different library, leaving every old row orphaned.
Set dry_run=False to actually delete. The default is intentionally
safe so the matching count can be surfaced in the UI before the
user commits to it.
"""
from sqlalchemy import delete
async with AsyncSessionLocal() as session:
try:
deletable, skipped = await find_missing_photos(session)
if not dry_run and deletable:
# Chunked delete to keep the IN clause within SQLite limits.
CHUNK = 500
for i in range(0, len(deletable), CHUNK):
await session.execute(
delete(Photo).where(Photo.id.in_(deletable[i:i + CHUNK]))
)
await session.commit()
logger.info(f"Pruned {len(deletable)} orphaned photo rows")
return {
"would_delete" if dry_run else "deleted": len(deletable),
"skipped_unmounted": len(skipped),
"dry_run": dry_run,
}
except Exception as e:
logger.error(f"prune_missing_photos failed: {e}")
await session.rollback()
raise
async def cleanup_data_integrity() -> dict:
"""Top-level entry point. Runs the dedupe + count refresh in a single
transaction. Returns a small summary dict for logging."""
async with AsyncSessionLocal() as session:
try:
sr_deleted = await _dedupe_source_roots(session)
f_deleted = await _dedupe_folders(session)
await _recompute_folder_counts(session)
stale = await _warn_stale_source_roots(session)
await session.commit()
summary = {
"source_roots_merged": sr_deleted,
"folders_merged": f_deleted,
"source_roots_stale": stale,
}
if sr_deleted or f_deleted:
logger.info(f"Cleanup merged duplicates: {summary}")
return summary
except Exception as e:
logger.error(f"Cleanup failed: {e}")
await session.rollback()
raise