diff --git a/alembic/versions/0099_library_placement_run.py b/alembic/versions/0099_library_placement_run.py new file mode 100644 index 0000000..8907a7f --- /dev/null +++ b/alembic/versions/0099_library_placement_run.py @@ -0,0 +1,99 @@ +"""library_placement_run — the placement reconciler's plan/apply/undo ledger. + +Milestone #421 step 3. The survey (#4245) measured 33,789 ImageRecord rows +sitting outside their artist's canonical directory, across 56 artists. This +table holds one run of the sweep that trues them up: the plan, what it did, +and where every file came from. + +## Why the moves live in a table rather than a log line + +`ImageRecord.path` is the only pointer at the bytes, so a move rewrites the +row. Once that write lands, the previous location exists nowhere — unless it +was recorded first. `moves` is that record, which is what makes a 33,789-file +operation something the operator can undo per artist after looking at the +result, rather than a one-way door. + +An `applied` row is therefore HISTORY, not state (lesson #4226). Any future +retention on this table may prune `ready`, `cancelled` and `error` runs; an +`applied` one is only disposable once someone decides undo is no longer +wanted. That is deliberately not a timer's decision, and no pruning is added +here. + +Revision ID: 0099 +Revises: 0098 +Create Date: 2026-09-21 + +""" +from typing import Sequence, Union + +import sqlalchemy as sa +from alembic import op +from sqlalchemy.dialects import postgresql + +revision: str = "0099" +down_revision: Union[str, None] = "0098" +branch_labels: Union[str, Sequence[str], None] = None +depends_on: Union[str, Sequence[str], None] = None + + +def upgrade() -> None: + op.create_table( + "library_placement_run", + sa.Column("id", sa.Integer(), nullable=False), + sa.Column( + "status", sa.String(length=16), server_default="running", + nullable=False, + ), + # SET NULL, not CASCADE: deleting an artist must not destroy the + # record of where their files were moved. + sa.Column("artist_id", sa.Integer(), nullable=True), + sa.Column( + "started_at", sa.DateTime(timezone=True), + server_default=sa.text("now()"), nullable=False, + ), + sa.Column("finished_at", sa.DateTime(timezone=True), nullable=True), + sa.Column( + "planned_count", sa.Integer(), server_default="0", nullable=False, + ), + sa.Column( + "moved_count", sa.Integer(), server_default="0", nullable=False, + ), + sa.Column( + "refused_count", sa.Integer(), server_default="0", nullable=False, + ), + sa.Column( + "moves", postgresql.JSONB(astext_type=sa.Text()), + server_default=sa.text("'[]'::jsonb"), nullable=False, + ), + sa.Column( + "refusals", postgresql.JSONB(astext_type=sa.Text()), + server_default=sa.text("'[]'::jsonb"), nullable=False, + ), + sa.Column("error", sa.Text(), nullable=True), + sa.ForeignKeyConstraint( + ["artist_id"], ["artist.id"], + name="fk_library_placement_run_artist_id", ondelete="SET NULL", + ), + sa.PrimaryKeyConstraint("id"), + ) + op.create_index( + "ix_library_placement_run_status", "library_placement_run", ["status"], + ) + op.create_index( + "ix_library_placement_run_artist_id", "library_placement_run", + ["artist_id"], + ) + + +def downgrade() -> None: + # Dropping this table destroys the only record of where moved files came + # from. That is correct for a downgrade — the code that reads it is going + # away too — but it is worth saying out loud rather than discovering. + op.drop_index( + "ix_library_placement_run_artist_id", + table_name="library_placement_run", + ) + op.drop_index( + "ix_library_placement_run_status", table_name="library_placement_run", + ) + op.drop_table("library_placement_run") diff --git a/backend/app/api/cleanup.py b/backend/app/api/cleanup.py index 44cf6d3..d4b2196 100644 --- a/backend/app/api/cleanup.py +++ b/backend/app/api/cleanup.py @@ -29,7 +29,7 @@ from quart import Blueprint, jsonify, request from sqlalchemy import select from ..extensions import get_session -from ..models import LibraryAuditRun +from ..models import LibraryAuditRun, LibraryPlacementRun from ..services import cleanup_service, library_layout from ._responses import error_response as _bad @@ -218,3 +218,91 @@ async def layout_survey(): ) ) return jsonify({**report.as_dict(), "checked_disk": check_disk}) + + +def _serialize_placement_run(run: LibraryPlacementRun, *, moves: bool = False) -> dict: + """`moves` is opt-in: an applied whole-library run carries tens of + thousands of entries, which is a fine thing to hold in Postgres and a + poor thing to put in every list response.""" + out = { + "id": run.id, + "status": run.status, + "artist_id": run.artist_id, + "started_at": run.started_at.isoformat() if run.started_at else None, + "finished_at": run.finished_at.isoformat() if run.finished_at else None, + "planned_count": run.planned_count, + "moved_count": run.moved_count, + "refused_count": run.refused_count, + "refusals": run.refusals or [], + "error": run.error, + } + if moves: + out["moves"] = run.moves or [] + return out + + +@cleanup_bp.route("/placement/runs", methods=["GET"]) +async def placement_runs(): + """Newest first. Without `moves`, so the list stays small.""" + try: + limit = min(int(request.args.get("limit", "25")), 100) + except ValueError: + return _bad("invalid_limit") + async with get_session() as session: + rows = (await session.execute( + select(LibraryPlacementRun) + .order_by(LibraryPlacementRun.id.desc()).limit(limit) + )).scalars().all() + return jsonify({"runs": [_serialize_placement_run(r) for r in rows]}) + + +@cleanup_bp.route("/placement/runs/", methods=["GET"]) +async def placement_run(run_id: int): + """One run WITH its moves — this is the preview the operator reads before + agreeing, and the record of what happened afterwards.""" + async with get_session() as session: + run = await session.get(LibraryPlacementRun, run_id) + if run is None: + return _bad("not_found", status=404) + return jsonify(_serialize_placement_run(run, moves=True)) + + +@cleanup_bp.route("/placement/plan", methods=["POST"]) +async def placement_plan(): + """Queue a planning run. `artist_id` scopes it to one artist, which is the + intended use: do one, look at the gallery, then continue or revert.""" + body = await request.get_json(silent=True) or {} + artist_id = body.get("artist_id") + if artist_id is not None and not isinstance(artist_id, int): + return _bad("invalid_artist_id") + from ..tasks.library_placement import plan_placement + plan_placement.delay(artist_id) + return jsonify({"status": "dispatched"}), 202 + + +@cleanup_bp.route("/placement/runs//apply", methods=["POST"]) +async def placement_apply(run_id: int): + """Execute a ready run. This renames files and rewrites rows.""" + async with get_session() as session: + run = await session.get(LibraryPlacementRun, run_id) + if run is None: + return _bad("not_found", status=404) + if run.status != "ready": + return _bad("not_ready", detail=f"run is {run.status}") + from ..tasks.library_placement import apply_placement + apply_placement.delay(run_id) + return jsonify({"status": "dispatched"}), 202 + + +@cleanup_bp.route("/placement/runs//revert", methods=["POST"]) +async def placement_revert(run_id: int): + """Put an applied run's files back. The reason the ledger is kept.""" + async with get_session() as session: + run = await session.get(LibraryPlacementRun, run_id) + if run is None: + return _bad("not_found", status=404) + if run.status != "applied": + return _bad("not_applied", detail=f"run is {run.status}") + from ..tasks.library_placement import revert_placement + revert_placement.delay(run_id) + return jsonify({"status": "dispatched"}), 202 diff --git a/backend/app/celery_app.py b/backend/app/celery_app.py index d49c8ef..5f176a2 100644 --- a/backend/app/celery_app.py +++ b/backend/app/celery_app.py @@ -35,6 +35,7 @@ def make_celery() -> Celery: "backend.app.tasks.backup", "backend.app.tasks.admin", "backend.app.tasks.library_audit", + "backend.app.tasks.library_placement", "backend.app.tasks.translation", ], ) @@ -62,6 +63,8 @@ def make_celery() -> Celery: # 2026-06-07: a 2h audit blocked vacuum/backup/normalize for hours). "backend.app.tasks.maintenance.*": {"queue": "maintenance"}, "backend.app.tasks.backup.*": {"queue": "maintenance_long"}, + # 33k renames on NFS: long lane, same as backups. + "backend.app.tasks.library_placement.*": {"queue": "maintenance_long"}, "backend.app.tasks.admin.*": {"queue": "maintenance_long"}, "backend.app.tasks.library_audit.*": {"queue": "maintenance_long"}, # Translation backfill hits the LLM (~1–6s/item) → the long lane so it diff --git a/backend/app/models/__init__.py b/backend/app/models/__init__.py index 9e085f4..b9bcf4a 100644 --- a/backend/app/models/__init__.py +++ b/backend/app/models/__init__.py @@ -22,6 +22,7 @@ from .import_batch import ImportBatch from .import_settings import ImportSettings from .import_task import ImportTask from .library_audit_run import LibraryAuditRun +from .library_placement_run import LibraryPlacementRun from .membership_sync import MembershipSync from .ml_settings import MLSettings from .patreon_failed_media import PatreonFailedMedia @@ -85,6 +86,7 @@ __all__ = [ "ImportTask", "ImportSettings", "LibraryAuditRun", + "LibraryPlacementRun", "MembershipSync", "MLSettings", "HeadAutoApplyRun", diff --git a/backend/app/models/library_placement_run.py b/backend/app/models/library_placement_run.py new file mode 100644 index 0000000..da1fbfc --- /dev/null +++ b/backend/app/models/library_placement_run.py @@ -0,0 +1,99 @@ +"""LibraryPlacementRun — one run of the placement reconciler (milestone #421). + +The library is keyed on the Artist row's `slug`, one directory per artist. +Every writer agrees on that now (`utils.paths.canonical_subdir`, task #4244), +but ~33,789 rows were written under older rules and sit in some other +artist's directory. This row is a run of the sweep that trues them up. + +State machine, mirroring LibraryAuditRun: + + running -> ready -> applied -> reverted + \\-> cancelled + (any) -> error + +## The `moves` column does three jobs + +`moves` is the plan: `[{"image_id": 1, "from": "...", "to": "..."}, ...]`. + + 1. **Preview.** It is what the operator reads before agreeing. + 2. **Apply.** The apply executes THIS list rather than re-deriving the set, + so the preview cannot describe a different set from the apply. That is + rule 93's guarantee reached the way LibraryAuditRun reaches it — the + plan is materialised, not recomputed. + 3. **Revert.** `from` is retained, so a batch that looks wrong in the + gallery goes back where it came from. + +## An applied run IS the undo ledger — it must never be pruned + +This is the trap lesson #4226 names: a record that answers both "what is the +current plan" and "what happened" gets deleted by whatever forgets the first. +A `ready` run is disposable state. An `applied` run is HISTORY, and it is the +only record of where 33,789 files used to be — delete it and the moves become +irreversible. + +No pruning exists for this table today, and that is deliberate. If retention +is ever added here, it may prune `ready`, `cancelled` and `error` runs; an +`applied` run is only safe to drop once someone decides the moves are settled +and undo is no longer wanted, which is an operator decision and not a +timer's. +""" + +from datetime import datetime +from typing import Any + +from sqlalchemy import DateTime, ForeignKey, Integer, String, Text, func, text +from sqlalchemy.dialects.postgresql import JSONB +from sqlalchemy.orm import Mapped, mapped_column + +from .base import Base + + +class LibraryPlacementRun(Base): + __tablename__ = "library_placement_run" + + id: Mapped[int] = mapped_column(Integer, primary_key=True) + status: Mapped[str] = mapped_column( + String(16), nullable=False, default="running", index=True, + server_default="running", + ) + # running | ready | applied | reverted | cancelled | error + + # Scope. NULL = the whole library; set = one artist, which is how this is + # meant to be used — do one artist, look at it in the gallery, continue or + # revert. ondelete SET NULL rather than CASCADE: deleting an artist must + # not destroy the record of where their files were moved. + artist_id: Mapped[int | None] = mapped_column( + ForeignKey("artist.id", ondelete="SET NULL"), nullable=True, index=True, + ) + + started_at: Mapped[datetime] = mapped_column( + DateTime(timezone=True), nullable=False, server_default=func.now(), + ) + finished_at: Mapped[datetime | None] = mapped_column( + DateTime(timezone=True), nullable=True, + ) + + planned_count: Mapped[int] = mapped_column( + Integer, nullable=False, default=0, server_default="0", + ) + moved_count: Mapped[int] = mapped_column( + Integer, nullable=False, default=0, server_default="0", + ) + refused_count: Mapped[int] = mapped_column( + Integer, nullable=False, default=0, server_default="0", + ) + + # [{"image_id": int, "from": str, "to": str}, ...] — see the module + # docstring. This is the plan, the audit trail and the undo, in that order + # of appearance and in one place. + moves: Mapped[list[dict[str, Any]]] = mapped_column( + JSONB, nullable=False, default=list, server_default=text("'[]'::jsonb"), + ) + # [{"image_id": int, "reason": str}, ...] — rows the apply declined to + # touch, with why. A refusal is an expected outcome, not an error: the + # world moves between plan and apply, and every gate fails closed. + refusals: Mapped[list[dict[str, Any]]] = mapped_column( + JSONB, nullable=False, default=list, server_default=text("'[]'::jsonb"), + ) + + error: Mapped[str | None] = mapped_column(Text, nullable=True) diff --git a/backend/app/services/library_layout.py b/backend/app/services/library_layout.py index f1c2eeb..c14c711 100644 --- a/backend/app/services/library_layout.py +++ b/backend/app/services/library_layout.py @@ -16,17 +16,23 @@ writing their own — the house shape for rule 93, snippet #3087. A preview that computes its set differently from the apply is a preview that can lie, and here the apply RENAMES the operator's art. -Nothing in this module writes. It reads rows, and it stats files when asked. +## What writes and what does not + +`survey_layout` and `plan_placement` are read-only — they report and they +record a plan. `apply_run` and `revert_run` are the only functions here that +rename a file or rewrite a row, and each does both for one row at a time, +updating the row only after its rename lands. """ from __future__ import annotations from dataclasses import dataclass, field +from datetime import UTC, datetime from pathlib import Path from sqlalchemy import func, select from sqlalchemy.orm import Session -from ..models import Artist, ImageRecord +from ..models import Artist, ImageRecord, LibraryPlacementRun from ..utils.paths import canonical_subdir # Top-level directories under the images root that are STORES, not artists. @@ -229,3 +235,200 @@ def survey_layout( report.unmovable += layout.unmovable return report + + +# --- the reconciler: plan -> apply -> revert (#4246) ------------------------- +# +# The three verbs share one materialised plan rather than each deriving its +# own set. `plan_placement` writes `LibraryPlacementRun.moves`; `apply_run` +# executes THAT list; `revert_run` walks it backwards. A preview that can +# disagree with its apply is the failure this shape exists to prevent, and +# here the apply renames the operator's art. +# +# Every step fails CLOSED. The world moves between planning and applying — +# a download lands, a supersede rewrites a path, a file is deleted — so the +# apply re-checks each row against what the plan recorded and declines the +# ones that moved on, instead of trusting a plan that may be minutes old. + + +def plan_placement( + session: Session, images_root: Path, *, artist_id: int | None = None, +) -> LibraryPlacementRun: + """Build (and persist) the move plan. Touches no files. + + `artist_id` scopes the run to one artist, which is how this is meant to be + used: do one, look at the gallery, then continue or revert. None plans the + whole library. + """ + stmt = select(Artist).order_by(Artist.slug) + if artist_id is not None: + stmt = stmt.where(Artist.id == artist_id) + artists = session.execute(stmt).scalars().all() + + candidates: list[dict] = [] + wanted: dict[str, int] = {} + for artist in artists: + rows = session.execute( + select(ImageRecord.id, ImageRecord.path) + .where(*_misplaced_conditions(images_root, artist.id, artist.slug)) + ).all() + for row_id, path in rows: + dest = destination_for(path, images_root, artist.slug) + if dest is None or dest.exists(): + continue + key = str(dest) + candidates.append({"image_id": row_id, "from": path, "to": key}) + wanted[key] = wanted.get(key, 0) + 1 + + # Two rows wanting one destination: plan NEITHER. Which of them "wins" is + # not this sweep's call, and planning one of them would silently pick a + # winner by iteration order. Counting first and filtering after is what + # makes that true — claiming as we go would quietly keep whichever came + # first. + moves = [m for m in candidates if wanted[m["to"]] == 1] + + run = LibraryPlacementRun( + status="ready", artist_id=artist_id, moves=moves, + planned_count=len(moves), + ) + session.add(run) + session.flush() + return run + + +def _move_one(src: Path, dest: Path) -> str | None: + """Rename `src` to `dest`. Returns a refusal reason, or None on success. + + A rename within one filesystem, so no copy and no free space needed. The + destination check is not a race-free guarantee — nothing here is — but it + turns the common case of "something already landed there" into a refusal + instead of a silent overwrite. + """ + if not src.exists(): + return "source missing" + if dest.exists(): + return "destination occupied" + try: + dest.parent.mkdir(parents=True, exist_ok=True) + src.rename(dest) + except OSError as exc: + return f"rename failed: {exc}" + return None + + +def apply_run( + session: Session, run: LibraryPlacementRun, *, chunk: int = 0, +) -> LibraryPlacementRun: + """Execute a `ready` run's stored plan: file and row together, per row. + + The row is updated ONLY after its rename lands, so a refused or failed + move can never leave `ImageRecord.path` pointing at a file that is not + there. Refusals are recorded and the run continues — one row that moved + on since planning is not a reason to abandon the other 33,788. + + `chunk` commits progress every N moves. Set it for any real run: the + ledger is the ONLY record of where a file came from, so a worker that + dies at row 30,000 of 33,789 must not take the undo information for the + first 29,999 with it. Left at 0 (tests, small runs) everything persists + in one go at the end. + + Re-running a partially-applied plan is safe rather than clever: the rows + already moved no longer match their `from`, so they refuse as "row moved + since planning" instead of being moved twice. + """ + if run.status != "ready": + raise ValueError(f"run {run.id} is {run.status}, not ready") + + refusals: list[dict] = [] + moved: list[dict] = [] + + def _persist() -> None: + # Reassign rather than mutate: SQLAlchemy does not track in-place + # changes to a JSONB list, so an .append() alone would never reach + # the database and the ledger would silently stay empty. + run.moves = list(moved) + run.refusals = list(refusals) + run.moved_count = len(moved) + run.refused_count = len(refusals) + session.commit() + + # Snapshot the plan before iterating: `_persist` reassigns `run.moves`, + # and iterating the attribute while rewriting it would walk a list that + # changes underneath the loop. + plan = list(run.moves) + for done, move in enumerate(plan, start=1): + record = session.get(ImageRecord, move["image_id"]) + if record is None: + refusals.append({"image_id": move["image_id"], "reason": "row gone"}) + continue + if record.path != move["from"]: + # Something rewrote this row since the plan was built — a + # supersede, or an earlier run. The plan is stale for it. + refusals.append({ + "image_id": move["image_id"], "reason": "row moved since planning", + }) + continue + reason = _move_one(Path(move["from"]), Path(move["to"])) + if reason is not None: + refusals.append({"image_id": move["image_id"], "reason": reason}) + continue + record.path = move["to"] + moved.append(move) + if chunk and done % chunk == 0: + _persist() + + run.moves = moved + run.refusals = refusals + run.moved_count = len(moved) + run.refused_count = len(refusals) + run.status = "applied" + run.finished_at = datetime.now(UTC) + session.flush() + return run + + +def revert_run( + session: Session, run: LibraryPlacementRun, +) -> LibraryPlacementRun: + """Put an applied run's files back where they came from. + + This is why `from` is retained. It is the answer to "do one artist, look + at it, and undo if it reads wrong" — which is a cheaper way to settle + whether artist_id or the folder held the truth (#4257) than arguing it + from a sample. + + Refuses the same way the apply does: a file someone has since moved or + replaced stays where it is, and its row is left alone. + + A revert interrupted half way is resumable by re-running it: the rows + already put back no longer sit at `to`, so they refuse rather than move + twice. Unlike the apply this needs no chunked persistence — it consumes + the ledger rather than producing it, so a crash costs progress, not + information. + """ + if run.status != "applied": + raise ValueError(f"run {run.id} is {run.status}, not applied") + + refusals: list[dict] = [] + reverted = 0 + for move in run.moves: + record = session.get(ImageRecord, move["image_id"]) + if record is None or record.path != move["to"]: + refusals.append({ + "image_id": move["image_id"], "reason": "row changed since apply", + }) + continue + reason = _move_one(Path(move["to"]), Path(move["from"])) + if reason is not None: + refusals.append({"image_id": move["image_id"], "reason": reason}) + continue + record.path = move["from"] + reverted += 1 + + run.refusals = refusals + run.refused_count = len(refusals) + run.moved_count = run.moved_count - reverted + run.status = "reverted" + run.finished_at = datetime.now(UTC) + session.flush() + return run diff --git a/backend/app/tasks/library_placement.py b/backend/app/tasks/library_placement.py new file mode 100644 index 0000000..510795d --- /dev/null +++ b/backend/app/tasks/library_placement.py @@ -0,0 +1,125 @@ +"""Placement reconciler tasks — plan, apply, revert (milestone #421). + +The service (`services.library_layout`) holds the decisions; this module is +only the async wrapper, matching `tasks.library_audit`: run on the +maintenance queue, mark the run `error` with a traceback if anything escapes, +and return a small summary dict so eager-mode tests can assert on it. + +Applying is a long run — 33,789 renames on the operator's library at the time +of writing — so `apply_placement` persists its ledger in chunks rather than +at the end. That ledger is the only record of where each file came from, and +a worker that dies two thirds of the way through must not take the undo +information for the first two thirds with it. +""" + +import logging +import traceback +from datetime import UTC, datetime +from pathlib import Path + +from sqlalchemy.exc import DBAPIError, OperationalError + +from ..celery_app import celery +from ..models import LibraryPlacementRun +from ..services import library_layout +from ._sync_engine import sync_session_factory as _sync_session_factory + +log = logging.getLogger(__name__) + +IMAGES_ROOT = Path("/images") + +# Commit the ledger every this many moves. Small enough that a crash loses +# seconds of work, large enough not to make a COMMIT per rename. +_APPLY_CHUNK = 200 + + +def _fail(session, run_id: int, message: str) -> None: + run = session.get(LibraryPlacementRun, run_id) + if run is not None: + run.status = "error" + run.error = message + run.finished_at = datetime.now(UTC) + session.commit() + + +@celery.task( + name="backend.app.tasks.library_placement.plan_placement", + autoretry_for=(OperationalError, DBAPIError), + retry_backoff=5, retry_backoff_max=60, retry_jitter=True, max_retries=3, + soft_time_limit=900, time_limit=1000, +) +def plan_placement(artist_id: int | None = None) -> dict: + """Build a move plan and leave it `ready` for the operator to read. + + Reads rows and stats destinations; moves nothing. + """ + SessionLocal = _sync_session_factory() + with SessionLocal() as session: + run = library_layout.plan_placement( + session, IMAGES_ROOT, artist_id=artist_id, + ) + session.commit() + return { + "run_id": run.id, "status": run.status, + "planned_count": run.planned_count, + } + + +@celery.task( + name="backend.app.tasks.library_placement.apply_placement", + soft_time_limit=7200, time_limit=7500, +) +def apply_placement(run_id: int) -> dict: + """Execute a `ready` run's stored plan. Renames files and rewrites rows. + + No autoretry: a retry would re-enter a half-applied plan on a schedule + nobody asked for. Re-running IS safe (the applied rows refuse as "row + moved since planning"), but that should be the operator's decision after + reading what happened, not the queue's. + """ + SessionLocal = _sync_session_factory() + with SessionLocal() as session: + run = session.get(LibraryPlacementRun, run_id) + if run is None: + return {"run_id": run_id, "status": "missing"} + if run.status != "ready": + return {"run_id": run_id, "status": run.status, "skipped": True} + try: + library_layout.apply_run(session, run, chunk=_APPLY_CHUNK) + session.commit() + except Exception: + log.exception("placement apply failed for run %s", run_id) + session.rollback() + _fail(session, run_id, traceback.format_exc()) + return {"run_id": run_id, "status": "error"} + return { + "run_id": run_id, "status": run.status, + "moved": run.moved_count, "refused": run.refused_count, + } + + +@celery.task( + name="backend.app.tasks.library_placement.revert_placement", + soft_time_limit=7200, time_limit=7500, +) +def revert_placement(run_id: int) -> dict: + """Put an applied run's files back where they came from.""" + SessionLocal = _sync_session_factory() + with SessionLocal() as session: + run = session.get(LibraryPlacementRun, run_id) + if run is None: + return {"run_id": run_id, "status": "missing"} + if run.status != "applied": + return {"run_id": run_id, "status": run.status, "skipped": True} + try: + library_layout.revert_run(session, run) + session.commit() + except Exception: + log.exception("placement revert failed for run %s", run_id) + session.rollback() + _fail(session, run_id, traceback.format_exc()) + return {"run_id": run_id, "status": "error"} + return { + "run_id": run_id, "status": run.status, + "refused": run.refused_count, + } diff --git a/frontend/src/components/settings/MaintenancePanel.vue b/frontend/src/components/settings/MaintenancePanel.vue index 3931384..0c44eaa 100644 --- a/frontend/src/components/settings/MaintenancePanel.vue +++ b/frontend/src/components/settings/MaintenancePanel.vue @@ -54,6 +54,7 @@ Self-healing and repair: missing files, thumbnails, database upkeep.

+ @@ -80,6 +81,7 @@ import MLBackfillCard from './MLBackfillCard.vue' import ThumbnailBackfillCard from './ThumbnailBackfillCard.vue' import ArchiveReextractCard from './ArchiveReextractCard.vue' import MissingFileRepairCard from './MissingFileRepairCard.vue' +import PlacementCard from './PlacementCard.vue' import GpuTriageCard from './GpuTriageCard.vue' import DbMaintenanceCard from './DbMaintenanceCard.vue' import VideoEmbeddingCard from './VideoEmbeddingCard.vue' diff --git a/frontend/src/components/settings/PlacementCard.vue b/frontend/src/components/settings/PlacementCard.vue new file mode 100644 index 0000000..462c255 --- /dev/null +++ b/frontend/src/components/settings/PlacementCard.vue @@ -0,0 +1,312 @@ + + + diff --git a/frontend/src/stores/cleanup.js b/frontend/src/stores/cleanup.js index 82813a2..73b17d7 100644 --- a/frontend/src/stores/cleanup.js +++ b/frontend/src/stores/cleanup.js @@ -71,10 +71,66 @@ export const useCleanupStore = defineStore('cleanup', () => { return await api.post(`/api/cleanup/audit/${id}/cancel`) } + // --- placement reconciler (milestone #421) -------------------------------- + // + // Runs are server-side rows, so the DATABASE is the durable state here — + // no localStorage resurfacing (useMaintenanceTask) is needed. Reload the + // page, open it on another machine, and the run and its status are simply + // there. That also means a plan survives being walked away from for a day. + + const placementRuns = ref([]) + const layout = ref(null) + + // The survey: which rows sit outside their artist's directory. check_disk + // additionally stats every destination (collisions, missing sources) and + // costs one stat per misplaced row over NFS, so it is opt-in. + async function loadLayout(checkDisk = false) { + layout.value = await api.get('/api/cleanup/layout', { + params: checkDisk ? { check_disk: 1 } : {}, + }) + return layout.value + } + + // id -> name, so a run row can say "Conto" instead of "#47". The runs + // endpoint carries artist_id alone: the name belongs to the artist, and + // denormalising it into every run would go stale the moment one is renamed. + async function loadArtistNames() { + const rows = await api.get('/api/artists/names') + return Object.fromEntries((rows || []).map(a => [a.id, a.name])) + } + + async function loadPlacementRuns(limit = 25) { + const body = await api.get('/api/cleanup/placement/runs', { params: { limit } }) + placementRuns.value = body.runs || [] + return placementRuns.value + } + + // Detail carries `moves` — the plan the operator reads before agreeing. + async function getPlacementRun(id) { + return await api.get(`/api/cleanup/placement/runs/${id}`) + } + + async function planPlacement(artistId = null) { + return await api.post('/api/cleanup/placement/plan', { + body: artistId === null ? {} : { artist_id: artistId }, + }) + } + + async function applyPlacement(id) { + return await api.post(`/api/cleanup/placement/runs/${id}/apply`) + } + + async function revertPlacement(id) { + return await api.post(`/api/cleanup/placement/runs/${id}/revert`) + } + return { defaults, recentRuns, loadDefaults, previewMinDim, deleteMinDim, startAudit, getAudit, loadHistory, latestAuditForRule, applyAudit, cancelAudit, + placementRuns, layout, + loadLayout, loadArtistNames, loadPlacementRuns, getPlacementRun, + planPlacement, applyPlacement, revertPlacement, } }) diff --git a/frontend/test/placement.spec.js b/frontend/test/placement.spec.js new file mode 100644 index 0000000..af5fc1a --- /dev/null +++ b/frontend/test/placement.spec.js @@ -0,0 +1,111 @@ +import { describe, it, expect, beforeEach, afterEach, vi } from 'vitest' +import { setActivePinia, createPinia } from 'pinia' +import { useCleanupStore } from '../src/stores/cleanup.js' +import { stubFetch } from './stubFetch.js' + + +describe('placement reconciler store (milestone #421)', () => { + beforeEach(() => setActivePinia(createPinia())) + afterEach(() => vi.restoreAllMocks()) + + it('loadLayout leaves the disk check off by default', async () => { + const s = useCleanupStore() + let seen = '' + stubFetch((url) => { + seen = url + return { status: 200, body: { total_rows: 10, misplaced_rows: 2, artists: [] } } + }) + await s.loadLayout() + // One stat per misplaced row over NFS is the cost; it must be opt-in. + expect(seen).not.toContain('check_disk') + expect(s.layout.misplaced_rows).toBe(2) + }) + + it('loadLayout asks for the disk check when requested', async () => { + const s = useCleanupStore() + let seen = '' + stubFetch((url) => { + seen = url + return { status: 200, body: { total_rows: 0, misplaced_rows: 0, artists: [] } } + }) + await s.loadLayout(true) + expect(seen).toContain('check_disk=1') + }) + + it('planPlacement scopes to an artist when given one', async () => { + const s = useCleanupStore() + let sent = null + stubFetch((url, init) => { + sent = JSON.parse(init.body) + return { status: 202, body: { status: 'dispatched' } } + }) + await s.planPlacement(47) + expect(sent).toEqual({ artist_id: 47 }) + }) + + it('planPlacement sends no scope for the whole library', async () => { + const s = useCleanupStore() + let sent = null + stubFetch((url, init) => { + sent = JSON.parse(init.body) + return { status: 202, body: { status: 'dispatched' } } + }) + await s.planPlacement() + // Not `{artist_id: null}` — the endpoint rejects a non-integer, and an + // absent key is how "whole library" is spelled. + expect(sent).toEqual({}) + }) + + it('loadPlacementRuns keeps the rows for the table', async () => { + const s = useCleanupStore() + stubFetch(() => ({ + status: 200, + body: { runs: [{ id: 3, status: 'ready', planned_count: 12 }] }, + })) + await s.loadPlacementRuns() + expect(s.placementRuns).toHaveLength(1) + expect(s.placementRuns[0].status).toBe('ready') + }) + + it('getPlacementRun carries the moves — it is the preview', async () => { + const s = useCleanupStore() + stubFetch(() => ({ + status: 200, + body: { + id: 3, status: 'ready', planned_count: 1, + moves: [{ image_id: 9, from: '/images/Conto/x.png', to: '/images/conto/x.png' }], + }, + })) + const run = await s.getPlacementRun(3) + expect(run.moves[0].from).toBe('/images/Conto/x.png') + expect(run.moves[0].to).toBe('/images/conto/x.png') + }) + + it('loadArtistNames maps id to name for the run rows', async () => { + const s = useCleanupStore() + stubFetch(() => ({ + status: 200, + body: [{ id: 47, name: 'Conto', slug: 'conto' }], + })) + expect(await s.loadArtistNames()).toEqual({ 47: 'Conto' }) + }) + + it('loadArtistNames survives an empty roster', async () => { + const s = useCleanupStore() + stubFetch(() => ({ status: 200, body: [] })) + expect(await s.loadArtistNames()).toEqual({}) + }) + + it('applyPlacement and revertPlacement post to their own run', async () => { + const s = useCleanupStore() + const urls = [] + stubFetch((url) => { + urls.push(url) + return { status: 202, body: { status: 'dispatched' } } + }) + await s.applyPlacement(5) + await s.revertPlacement(5) + expect(urls[0]).toContain('/api/cleanup/placement/runs/5/apply') + expect(urls[1]).toContain('/api/cleanup/placement/runs/5/revert') + }) +}) diff --git a/frontend/test/stubFetch.js b/frontend/test/stubFetch.js new file mode 100644 index 0000000..80dc8fd --- /dev/null +++ b/frontend/test/stubFetch.js @@ -0,0 +1,24 @@ +import { vi } from 'vitest' + +// The canonical fetch stub for store specs. +// +// `handler(url, init)` returns `{ status, body }`; body is JSON-encoded, and +// `ok` is derived from the status so a store's error path can be exercised by +// returning 4xx/5xx. Returns the vi.fn so a caller can assert on calls. +// +// Extracted 2026-09-21 from six specs carrying byte-identical copies +// (adminStore, credentials, dbMaintenance, gallery, suggestions, +// galleryRelatedStrip). Those still hold their own; migrate each the next +// time it is touched rather than in one sweep. +export function stubFetch (handler) { + globalThis.fetch = vi.fn(async (url, init) => { + const { status, body } = handler(url, init) + return { + ok: status >= 200 && status < 300, + status, + statusText: String(status), + text: async () => (body == null ? '' : JSON.stringify(body)), + } + }) + return globalThis.fetch +} diff --git a/tests/test_api_placement.py b/tests/test_api_placement.py new file mode 100644 index 0000000..3fd56b4 --- /dev/null +++ b/tests/test_api_placement.py @@ -0,0 +1,167 @@ +"""The placement reconciler's task + API surface (milestone #421, slice 3b). + +The move logic itself is covered in tests/test_library_layout.py; this module +covers the wrapper — that the tasks are registered and routed, that the +endpoints gate on run state, and that a list response stays small. +""" + +import pytest +from sqlalchemy import select + +import backend.app.tasks.library_placement # noqa: F401 — register tasks +from backend.app.celery_app import celery +from backend.app.models import Artist, LibraryPlacementRun + +pytestmark = pytest.mark.integration + +_TASKS = ( + "backend.app.tasks.library_placement.plan_placement", + "backend.app.tasks.library_placement.apply_placement", + "backend.app.tasks.library_placement.revert_placement", +) + + +@pytest.mark.parametrize("name", _TASKS) +def test_placement_tasks_are_registered(name): + assert name in celery.tasks + + +def test_placement_runs_on_the_long_maintenance_lane(): + """33k renames must not sit in the quick lane, which is where the + self-healing sweeps live (the 2026-06-07 starvation).""" + routes = celery.conf.task_routes + assert routes["backend.app.tasks.library_placement.*"] == { + "queue": "maintenance_long" + } + + +def _run(db, status="ready", moves=None, artist_id=None): + """Adds the row; the caller awaits the COMMIT. + + Commit, not flush: the app under test runs on its own session and + connection, so a flush that stays inside this test's transaction is + invisible to the endpoint — the row simply is not there yet. Same reason + `_seed_runs` in test_api_system_backup commits.""" + run = LibraryPlacementRun( + status=status, artist_id=artist_id, moves=moves or [], + planned_count=len(moves or []), + ) + db.add(run) + return run + + +@pytest.mark.asyncio +async def test_runs_list_omits_the_moves(client, db): + """An applied whole-library run carries tens of thousands of entries. + Fine in Postgres, wrong in every list response.""" + run = LibraryPlacementRun( + status="applied", + moves=[{"image_id": 1, "from": "/images/A/x.png", "to": "/images/a/x.png"}], + planned_count=1, moved_count=1, + ) + db.add(run) + await db.commit() + + resp = await client.get("/api/cleanup/placement/runs") + assert resp.status_code == 200 + body = await resp.get_json() + assert body["runs"][0]["planned_count"] == 1 + assert "moves" not in body["runs"][0] + + +@pytest.mark.asyncio +async def test_run_detail_carries_the_moves(client, db): + """The detail IS the preview the operator reads before agreeing.""" + run = LibraryPlacementRun( + status="ready", + moves=[{"image_id": 7, "from": "/images/Conto/x.png", "to": "/images/conto/x.png"}], + planned_count=1, + ) + db.add(run) + await db.commit() + + resp = await client.get(f"/api/cleanup/placement/runs/{run.id}") + assert resp.status_code == 200 + body = await resp.get_json() + assert body["moves"][0]["from"] == "/images/Conto/x.png" + assert body["moves"][0]["to"] == "/images/conto/x.png" + + +@pytest.mark.asyncio +async def test_run_detail_404s_for_an_unknown_run(client): + resp = await client.get("/api/cleanup/placement/runs/999999") + assert resp.status_code == 404 + + +@pytest.mark.asyncio +async def test_plan_rejects_a_non_integer_artist(client): + resp = await client.post( + "/api/cleanup/placement/plan", json={"artist_id": "conto"}, + ) + assert resp.status_code == 400 + assert (await resp.get_json())["error"] == "invalid_artist_id" + + +@pytest.mark.asyncio +async def test_plan_accepts_an_artist_scope(client, db, monkeypatch): + sent = {} + from backend.app.tasks import library_placement + + monkeypatch.setattr( + library_placement.plan_placement, "delay", + lambda artist_id=None: sent.update(artist_id=artist_id), + ) + artist = Artist(name="Conto", slug="conto") + db.add(artist) + await db.commit() + + resp = await client.post( + "/api/cleanup/placement/plan", json={"artist_id": artist.id}, + ) + assert resp.status_code == 202 + assert sent["artist_id"] == artist.id + + +@pytest.mark.asyncio +async def test_apply_refuses_a_run_that_is_not_ready(client, db): + """The gate is here as well as in the service — an applied run must not + be re-applied by a stray POST.""" + run = _run(db, status="applied") + await db.commit() + + resp = await client.post(f"/api/cleanup/placement/runs/{run.id}/apply") + assert resp.status_code == 400 + assert (await resp.get_json())["error"] == "not_ready" + + +@pytest.mark.asyncio +async def test_revert_refuses_a_run_that_was_never_applied(client, db): + run = _run(db, status="ready") + await db.commit() + + resp = await client.post(f"/api/cleanup/placement/runs/{run.id}/revert") + assert resp.status_code == 400 + assert (await resp.get_json())["error"] == "not_applied" + + +@pytest.mark.asyncio +async def test_apply_dispatches_for_a_ready_run(client, db, monkeypatch): + sent = {} + from backend.app.tasks import library_placement + + monkeypatch.setattr( + library_placement.apply_placement, "delay", + lambda run_id: sent.update(run_id=run_id), + ) + run = _run(db, status="ready") + await db.commit() + + resp = await client.post(f"/api/cleanup/placement/runs/{run.id}/apply") + assert resp.status_code == 202 + assert sent["run_id"] == run.id + # Dispatch only — the endpoint must not have moved anything itself. + still = (await db.execute( + select(LibraryPlacementRun.status) + .where(LibraryPlacementRun.id == run.id) + )).scalar_one() + assert still == "ready" diff --git a/tests/test_library_layout.py b/tests/test_library_layout.py index 52cd88d..1594082 100644 --- a/tests/test_library_layout.py +++ b/tests/test_library_layout.py @@ -213,3 +213,191 @@ def test_survey_is_read_only(db_sync, tmp_path): assert db_sync.execute( select(ImageRecord.path).where(ImageRecord.id == rec.id) ).scalar_one() == before + + +# --- plan / apply / revert (#4246) ------------------------------------------ + + +def _staged(db, tmp_path, slug, stray, name="x.png", n=100): + """An artist with one file sitting in `stray`'s directory.""" + artist = _artist(db, slug.title(), slug) + src = tmp_path / stray / name + src.parent.mkdir(parents=True, exist_ok=True) + src.write_bytes(b"pixels") + rec = _image(db, str(src), artist, n) + return artist, rec, src + + +@pytest.mark.integration +def test_plan_records_where_each_file_came_from(db_sync, tmp_path): + from backend.app.services.library_layout import plan_placement + + _, rec, src = _staged(db_sync, tmp_path, "conto", "Conto", n=20) + run = plan_placement(db_sync, tmp_path) + + assert run.status == "ready" + assert run.planned_count == 1 + assert run.moves == [{ + "image_id": rec.id, + "from": str(src), + "to": str(tmp_path / "conto" / "x.png"), + }] + # Planning touches nothing. + assert src.exists() + assert db_sync.get(ImageRecord, rec.id).path == str(src) + + +@pytest.mark.integration +def test_apply_moves_file_and_row_together(db_sync, tmp_path): + from backend.app.services.library_layout import apply_run, plan_placement + + _, rec, src = _staged(db_sync, tmp_path, "conto", "Conto", n=21) + run = apply_run(db_sync, plan_placement(db_sync, tmp_path)) + + dest = tmp_path / "conto" / "x.png" + assert run.status == "applied" + assert run.moved_count == 1 and run.refused_count == 0 + assert dest.exists() and not src.exists() + db_sync.expire_all() + assert db_sync.get(ImageRecord, rec.id).path == str(dest) + + +@pytest.mark.integration +def test_revert_puts_it_back(db_sync, tmp_path): + """The whole reason `from` is retained: do one artist, look, undo.""" + from backend.app.services.library_layout import ( + apply_run, + plan_placement, + revert_run, + ) + + _, rec, src = _staged(db_sync, tmp_path, "conto", "Conto", n=22) + run = revert_run(db_sync, apply_run(db_sync, plan_placement(db_sync, tmp_path))) + + assert run.status == "reverted" + assert src.exists() + assert not (tmp_path / "conto" / "x.png").exists() + db_sync.expire_all() + assert db_sync.get(ImageRecord, rec.id).path == str(src) + + +@pytest.mark.integration +def test_apply_refuses_a_row_that_moved_since_planning(db_sync, tmp_path): + """A supersede or an earlier run can rewrite a path between plan and + apply. The stale entry is declined, not forced.""" + from backend.app.services.library_layout import apply_run, plan_placement + + _, rec, src = _staged(db_sync, tmp_path, "conto", "Conto", n=23) + run = plan_placement(db_sync, tmp_path) + + elsewhere = tmp_path / "conto" / "already-here.png" + elsewhere.parent.mkdir(parents=True, exist_ok=True) + elsewhere.write_bytes(b"pixels") + rec.path = str(elsewhere) + db_sync.flush() + + run = apply_run(db_sync, run) + + assert run.moved_count == 0 and run.refused_count == 1 + assert run.refusals[0]["reason"] == "row moved since planning" + assert src.exists() # untouched + + +@pytest.mark.integration +def test_apply_never_overwrites_an_occupied_destination(db_sync, tmp_path): + from backend.app.services.library_layout import apply_run, plan_placement + + _, rec, src = _staged(db_sync, tmp_path, "conto", "Conto", n=24) + run = plan_placement(db_sync, tmp_path) + + squatter = tmp_path / "conto" / "x.png" + squatter.parent.mkdir(parents=True, exist_ok=True) + squatter.write_bytes(b"someone else") + + run = apply_run(db_sync, run) + + assert run.refused_count == 1 + assert run.refusals[0]["reason"] == "destination occupied" + assert squatter.read_bytes() == b"someone else" + db_sync.expire_all() + assert db_sync.get(ImageRecord, rec.id).path == str(src) + + +@pytest.mark.integration +def test_apply_leaves_the_row_alone_when_the_source_is_gone(db_sync, tmp_path): + from backend.app.services.library_layout import apply_run, plan_placement + + _, rec, src = _staged(db_sync, tmp_path, "conto", "Conto", n=25) + run = plan_placement(db_sync, tmp_path) + src.unlink() + + run = apply_run(db_sync, run) + + assert run.refusals[0]["reason"] == "source missing" + db_sync.expire_all() + # The row still points at the missing file rather than at a file that + # was never created — a broken row is recoverable, a lying one is not. + assert db_sync.get(ImageRecord, rec.id).path == str(src) + + +@pytest.mark.integration +def test_plan_scopes_to_one_artist(db_sync, tmp_path): + """Per-artist scope is what makes this incremental instead of one + irreversible sweep.""" + from backend.app.services.library_layout import plan_placement + + conto, _, _ = _staged(db_sync, tmp_path, "conto", "Conto", n=26) + _staged(db_sync, tmp_path, "maewix", "Maewix", name="y.png", n=27) + + run = plan_placement(db_sync, tmp_path, artist_id=conto.id) + + assert run.planned_count == 1 + assert run.artist_id == conto.id + assert "Conto" in run.moves[0]["from"] + + +@pytest.mark.integration +def test_plan_skips_both_rows_when_two_want_one_destination(db_sync, tmp_path): + """Which of two colliding rows 'wins' is not this sweep's call.""" + from backend.app.services.library_layout import plan_placement + + artist = _artist(db_sync, "Sticky", "sticky") + for stray, n in (("StickySpoodge", 28), ("Stickyspoodge", 29)): + p = tmp_path / stray / "dup.png" + p.parent.mkdir(parents=True, exist_ok=True) + p.write_bytes(b"pixels") + _image(db_sync, str(p), artist, n) + + run = plan_placement(db_sync, tmp_path) + + assert run.planned_count == 0 + + +@pytest.mark.integration +def test_thumbnails_do_not_move(db_sync, tmp_path): + """Thumbs are sha-addressed (`thumbs//.jpg`), not path-keyed, so + a placement move must not touch them. Pinned so nobody 'fixes' it.""" + from backend.app.services.library_layout import apply_run, plan_placement + + artist, rec, _ = _staged(db_sync, tmp_path, "conto", "Conto", n=30) + thumb = tmp_path / "thumbs" / "ab" / "abc.jpg" + thumb.parent.mkdir(parents=True, exist_ok=True) + thumb.write_bytes(b"thumb") + rec.thumbnail_path = str(thumb) + db_sync.flush() + + apply_run(db_sync, plan_placement(db_sync, tmp_path)) + + db_sync.expire_all() + assert thumb.exists() + assert db_sync.get(ImageRecord, rec.id).thumbnail_path == str(thumb) + + +@pytest.mark.integration +def test_apply_refuses_a_run_that_is_not_ready(db_sync, tmp_path): + from backend.app.services.library_layout import apply_run, plan_placement + + _staged(db_sync, tmp_path, "conto", "Conto", n=31) + run = apply_run(db_sync, plan_placement(db_sync, tmp_path)) + with pytest.raises(ValueError): + apply_run(db_sync, run)