From 9ccc460c6907ba5af42f30161c41201b75ca686a Mon Sep 17 00:00:00 2001
From: Bryan Van Deusen
Date: Mon, 21 Sep 2026 12:46:06 -0400
Subject: [PATCH 1/5] =?UTF-8?q?feat:=20placement=20reconciler=20=E2=80=94?=
=?UTF-8?q?=20plan,=20apply,=20revert=20(4246,=20slice=203a)?=
MIME-Version: 1.0
Content-Type: text/plain; charset=UTF-8
Content-Transfer-Encoding: 8bit
Milestone #421 step 3, reframed on the operator's steer: not a one-off
migration but the system that keeps the tree true. The 33,789 misplaced rows
the survey found are just its first run.
The placement half was already done, verified by reading each writer rather
than assuming: downloads have always written `///`
(gallery_dl.py:523), attach_in_place leaves files where the downloader put
them, and `_copy_to_library` / `_supersede` became canonical in #4244. So
nothing is written off-canon today; what remains is the backlog and a standing
check for future drift.
`LibraryPlacementRun` (migration 0099) holds the plan as JSONB, and that one
structure does three jobs: it is the PREVIEW the operator reads, the list the
APPLY executes (rather than re-deriving the set, so the two cannot disagree),
and — because `from` is retained — the UNDO.
The undo is the point. It makes a 33,789-file operation something to do one
artist at a time, look at in the gallery, and reverse if it reads wrong. That
settles whether artist_id or the folder held the truth (spike #4257) by doing
rather than by arguing it from a 50-row sample.
An applied run is therefore HISTORY, not state — lesson #4226's trap, since
it is the only record of where those files used to be. The model and the
migration both say so: any future retention here may prune ready/cancelled/
error runs, never an applied one.
Everything fails closed. The apply re-checks each row against what the plan
recorded — source still there, destination still free, row still pointing
where the plan said — because a download or a supersede can land in between.
A refusal is recorded with its reason and the run continues; one stale row is
not a reason to abandon the other 33,788. The row is updated only after its
rename lands, so a failed move can never leave `path` naming a file that is
not there.
Writing the collision test caught the code disagreeing with its own comment:
it claimed the first of two rows wanting one destination and skipped the
second, silently picking a winner by iteration order. Now it counts first and
filters after, so genuinely neither is planned.
Thumbnails are sha-addressed, not path-keyed, so they do not move — pinned by
a test.
Co-Authored-By: Claude Opus 5 (1M context)
Claude-Session: https://claude.ai/code/session_01LVjrnpQjRgHdvq95rASoiR
---
.../versions/0099_library_placement_run.py | 99 +++++++++
backend/app/models/__init__.py | 2 +
backend/app/models/library_placement_run.py | 99 +++++++++
backend/app/services/library_layout.py | 174 +++++++++++++++-
tests/test_library_layout.py | 188 ++++++++++++++++++
5 files changed, 560 insertions(+), 2 deletions(-)
create mode 100644 alembic/versions/0099_library_placement_run.py
create mode 100644 backend/app/models/library_placement_run.py
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/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..dc3961b 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,167 @@ 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,
+) -> 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.
+ """
+ if run.status != "ready":
+ raise ValueError(f"run {run.id} is {run.status}, not ready")
+
+ refusals: list[dict] = []
+ moved: list[dict] = []
+ for move in run.moves:
+ 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)
+
+ 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.
+ """
+ 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/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)
--
2.54.0
From abe449b4f272a3179af34c98e6163ab4f9e95129 Mon Sep 17 00:00:00 2001
From: Bryan Van Deusen
Date: Mon, 21 Sep 2026 14:13:25 -0400
Subject: [PATCH 2/5] feat: placement reconciler tasks + API (4246, slice 3b)
MIME-Version: 1.0
Content-Type: text/plain; charset=UTF-8
Content-Transfer-Encoding: 8bit
Three Celery tasks wrapping the 3a service, and the endpoints that drive
them. Routed to `maintenance_long` alongside backups: 33k renames on NFS have
no business in the quick lane where the self-healing sweeps live (the
2026-06-07 starvation).
A durability bug in 3a, found by thinking about what a crash costs rather
than by a failing test: `apply_run` wrote its ledger only at the end, so a
worker dying at row 30,000 of 33,789 would have taken the undo information
for the first 29,999 with it — and that ledger is the ONLY record of where
those files came from. It now persists every 200 moves. Two things fell out
of writing that:
- `_persist` reassigns `run.moves`, so the loop had to snapshot the plan
first rather than iterate the attribute it rewrites.
- the reassignment is itself load-bearing: SQLAlchemy does not track in-place
mutation of a JSONB list, so an `.append()` alone would never reach the
database and the ledger would have stayed silently empty.
Re-running a partially-applied plan is safe — the moved rows no longer match
their `from` and refuse as "row moved since planning" — but `apply_placement`
deliberately has NO autoretry: re-entering a half-applied plan should be the
operator's call after reading what happened, not the queue's.
Endpoints gate on run state as well as the service does, so a stray POST
cannot re-apply an applied run. The list response omits `moves` (an applied
whole-library run carries tens of thousands of entries); the detail endpoint
includes them, because that detail IS the preview read before agreeing.
Co-Authored-By: Claude Opus 5 (1M context)
Claude-Session: https://claude.ai/code/session_01LVjrnpQjRgHdvq95rASoiR
---
backend/app/api/cleanup.py | 90 +++++++++++++-
backend/app/celery_app.py | 3 +
backend/app/services/library_layout.py | 37 +++++-
backend/app/tasks/library_placement.py | 125 +++++++++++++++++++
tests/test_api_placement.py | 163 +++++++++++++++++++++++++
5 files changed, 415 insertions(+), 3 deletions(-)
create mode 100644 backend/app/tasks/library_placement.py
create mode 100644 tests/test_api_placement.py
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/services/library_layout.py b/backend/app/services/library_layout.py
index dc3961b..c14c711 100644
--- a/backend/app/services/library_layout.py
+++ b/backend/app/services/library_layout.py
@@ -317,7 +317,7 @@ def _move_one(src: Path, dest: Path) -> str | None:
def apply_run(
- session: Session, run: LibraryPlacementRun,
+ session: Session, run: LibraryPlacementRun, *, chunk: int = 0,
) -> LibraryPlacementRun:
"""Execute a `ready` run's stored plan: file and row together, per row.
@@ -325,13 +325,38 @@ def apply_run(
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] = []
- for move in run.moves:
+
+ 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"})
@@ -349,6 +374,8 @@ def apply_run(
continue
record.path = move["to"]
moved.append(move)
+ if chunk and done % chunk == 0:
+ _persist()
run.moves = moved
run.refusals = refusals
@@ -372,6 +399,12 @@ def revert_run(
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")
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/tests/test_api_placement.py b/tests/test_api_placement.py
new file mode 100644
index 0000000..4bbe5e6
--- /dev/null
+++ b/tests/test_api_placement.py
@@ -0,0 +1,163 @@
+"""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
+
+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 flush. `db` is the ASYNC session,
+ so flushing here would leave an un-awaited coroutine and the row would
+ never reach the database."""
+ 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.flush()
+
+ 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.flush()
+
+ 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.flush()
+
+ 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.flush()
+
+ 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.flush()
+
+ 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.flush()
+
+ 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"
--
2.54.0
From db633777543606d57d7503feb5545aab440c982d Mon Sep 17 00:00:00 2001
From: Bryan Van Deusen
Date: Mon, 21 Sep 2026 14:15:52 -0400
Subject: [PATCH 3/5] test: placement API tests must commit, not flush (4246)
MIME-Version: 1.0
Content-Type: text/plain; charset=UTF-8
Content-Transfer-Encoding: 8bit
Five of the new tests failed in CI: every one that created a
LibraryPlacementRun and then read it back through the client. The ones that
touched no rows passed.
A flush stays inside the test's own transaction, and the app under test runs
on a separate session and connection — so the endpoint queried a database
where the row did not exist yet and got its 404 / empty list honestly.
`_seed_runs` in test_api_system_backup already commits for this reason; the
idiom was there to copy and I did not look first.
Recorded in the helper's docstring rather than just fixed, since the next
person writing a create-then-fetch API test will reach for flush too.
Co-Authored-By: Claude Opus 5 (1M context)
Claude-Session: https://claude.ai/code/session_01LVjrnpQjRgHdvq95rASoiR
---
tests/test_api_placement.py | 21 ++++++++++++---------
1 file changed, 12 insertions(+), 9 deletions(-)
diff --git a/tests/test_api_placement.py b/tests/test_api_placement.py
index 4bbe5e6..d9f0f49 100644
--- a/tests/test_api_placement.py
+++ b/tests/test_api_placement.py
@@ -35,9 +35,12 @@ def test_placement_runs_on_the_long_maintenance_lane():
def _run(db, status="ready", moves=None, artist_id=None):
- """Adds the row; the caller awaits the flush. `db` is the ASYNC session,
- so flushing here would leave an un-awaited coroutine and the row would
- never reach the database."""
+ """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 []),
@@ -56,7 +59,7 @@ async def test_runs_list_omits_the_moves(client, db):
planned_count=1, moved_count=1,
)
db.add(run)
- await db.flush()
+ await db.commit()
resp = await client.get("/api/cleanup/placement/runs")
assert resp.status_code == 200
@@ -74,7 +77,7 @@ async def test_run_detail_carries_the_moves(client, db):
planned_count=1,
)
db.add(run)
- await db.flush()
+ await db.commit()
resp = await client.get(f"/api/cleanup/placement/runs/{run.id}")
assert resp.status_code == 200
@@ -109,7 +112,7 @@ async def test_plan_accepts_an_artist_scope(client, db, monkeypatch):
)
artist = Artist(name="Conto", slug="conto")
db.add(artist)
- await db.flush()
+ await db.commit()
resp = await client.post(
"/api/cleanup/placement/plan", json={"artist_id": artist.id},
@@ -123,7 +126,7 @@ 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.flush()
+ await db.commit()
resp = await client.post(f"/api/cleanup/placement/runs/{run.id}/apply")
assert resp.status_code == 400
@@ -133,7 +136,7 @@ async def test_apply_refuses_a_run_that_is_not_ready(client, db):
@pytest.mark.asyncio
async def test_revert_refuses_a_run_that_was_never_applied(client, db):
run = _run(db, status="ready")
- await db.flush()
+ await db.commit()
resp = await client.post(f"/api/cleanup/placement/runs/{run.id}/revert")
assert resp.status_code == 400
@@ -150,7 +153,7 @@ async def test_apply_dispatches_for_a_ready_run(client, db, monkeypatch):
lambda run_id: sent.update(run_id=run_id),
)
run = _run(db, status="ready")
- await db.flush()
+ await db.commit()
resp = await client.post(f"/api/cleanup/placement/runs/{run.id}/apply")
assert resp.status_code == 202
--
2.54.0
From a4bdbcaca4d87131ff9bee8f37cb652dbf2b92ee Mon Sep 17 00:00:00 2001
From: Bryan Van Deusen
Date: Mon, 21 Sep 2026 14:21:20 -0400
Subject: [PATCH 4/5] test: import the placement task module so its names reach
celery.tasks (4246)
MIME-Version: 1.0
Content-Type: text/plain; charset=UTF-8
Content-Transfer-Encoding: 8bit
Three registration assertions failed: a task name only enters `celery.tasks`
when its module is imported, and nothing in the test process imported
`library_placement`. The API routes import it lazily inside the handlers, and
the registration test runs before any handler test triggers that.
`include=[...]` is what gets the module imported in a real WORKER, so
production registration was never in question — the test was asserting
something only observable after an import it never performed.
test_tasks_admin already carries the convention verbatim
(`import backend.app.tasks.admin # noqa: F401 — register tasks`); I wrote the
assertion from what I meant instead of copying the idiom next to it. Same
mistake shape as the commit-vs-flush bounce one commit ago: the pattern was
already in the suite both times.
Co-Authored-By: Claude Opus 5 (1M context)
Claude-Session: https://claude.ai/code/session_01LVjrnpQjRgHdvq95rASoiR
---
tests/test_api_placement.py | 1 +
1 file changed, 1 insertion(+)
diff --git a/tests/test_api_placement.py b/tests/test_api_placement.py
index d9f0f49..3fd56b4 100644
--- a/tests/test_api_placement.py
+++ b/tests/test_api_placement.py
@@ -8,6 +8,7 @@ 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
--
2.54.0
From 2dd9b956d519887a62c8d40ab2c6997a249843ee Mon Sep 17 00:00:00 2001
From: Bryan Van Deusen
Date: Mon, 21 Sep 2026 14:32:51 -0400
Subject: [PATCH 5/5] =?UTF-8?q?feat:=20File=20placement=20card=20=E2=80=94?=
=?UTF-8?q?=20survey,=20plan=20per=20artist,=20apply,=20put=20back=20(4246?=
=?UTF-8?q?,=20slice=203c)?=
MIME-Version: 1.0
Content-Type: text/plain; charset=UTF-8
Content-Transfer-Encoding: 8bit
The UI half of the reconciler, in Maintenance. Survey shows how many images
sit in the wrong artist's folder and which folders they are in; each artist
gets its own Plan button; each run can be reviewed, applied, and put back.
Deliberately NOT using useMaintenanceTask. That composable stashes a task id
in localStorage so a result survives navigate-away, which is the right answer
when the only record is a Celery result. Here the runs are database rows — so
a reload, another machine, or coming back tomorrow simply shows the same
state, because the state IS the row. The card polls the runs endpoint instead.
Copy avoids the vocabulary this work has been tripping over: "folder", "put
back", "in the wrong folder" rather than artist_id, revert and canonical. The
one thing the operator most needs to know — nothing here changes who an image
belongs to — is what the blurb says first.
Three things I had assumed and checked instead: MaintenanceTile lives in
common/ not settings/; there is no generic ConfirmDialog (BackupCard uses a
purpose-built modal), so this uses a plain v-dialog; and `loadArtistNames`
did not exist — it does now, mapping id to name so a run row reads "Conto"
rather than "#47". The name deliberately is not denormalised into the run: it
belongs to the artist and would go stale on a rename.
Also extracted `stubFetch` to frontend/test/stubFetch.js. The shape ledger
flagged it as a byte-identical duplicate across five specs and this would
have been the sixth; the others keep their copies until each is next touched.
Co-Authored-By: Claude Opus 5 (1M context)
Claude-Session: https://claude.ai/code/session_01LVjrnpQjRgHdvq95rASoiR
---
.../components/settings/MaintenancePanel.vue | 2 +
.../src/components/settings/PlacementCard.vue | 312 ++++++++++++++++++
frontend/src/stores/cleanup.js | 56 ++++
frontend/test/placement.spec.js | 111 +++++++
frontend/test/stubFetch.js | 24 ++
5 files changed, 505 insertions(+)
create mode 100644 frontend/src/components/settings/PlacementCard.vue
create mode 100644 frontend/test/placement.spec.js
create mode 100644 frontend/test/stubFetch.js
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 @@
+
+
+
+ The library keeps one folder per artist, named after them. Files written
+ under older rules can sit in another artist's folder — this moves them
+ home, updating the record and the file together. Every run can be
+ reverted, so the safe way to use it is one artist at a time: run it,
+ look at the gallery, then continue or put it back.
+
+
+ {{ error }}
+
+
+
+ Check placement
+
+ {{ layout.misplaced_rows.toLocaleString() }} of
+ {{ layout.total_rows.toLocaleString() }} images are in the wrong folder
+
+ — across {{ layout.artists.length }} artists
+
+
+
+
+
+
+
+ | Artist |
+ To move |
+ Currently in |
+ Plan |
+
+
+
+
+ | {{ a.name }} |
+ {{ a.misplaced_rows.toLocaleString() }} |
+ {{ a.stray_dirs.join(', ') }} |
+
+ Plan
+ |
+
+
+
+
+
+
+
+ Runs
+
+
+
+ | When |
+ Scope |
+ Status |
+ Planned |
+ Moved |
+ Refused |
+ Actions |
+
+
+
+
+ |
+ {{ formatRelative(r.started_at) }}
+ |
+ {{ artistName(r.artist_id) }} |
+
+
+ {{ statusIcon(r.status) }}
+
+ {{ r.status }}
+ |
+ {{ r.planned_count.toLocaleString() }} |
+ {{ r.moved_count.toLocaleString() }} |
+
+
+ {{ r.refused_count.toLocaleString() }}
+
+ |
+
+
+
+
+
+
+ |
+
+
+ |
+ No runs yet. Check placement above, then plan one artist.
+ |
+
+
+
+
+
+
+
+
+ Run {{ reviewRun?.id }} — {{ reviewRun?.planned_count?.toLocaleString() }} moves
+
+
+
+ Showing the first {{ REVIEW_LIMIT }}. Each row moves the file and
+ its record together; nothing is overwritten.
+
+
+
+
+ | {{ m.from }} |
+ → {{ m.to }} |
+
+
+
+
+
Refused ({{ reviewRun.refusals.length }})
+
+ Rows the run declined to touch — the source moved, the
+ destination was taken, or the record changed since planning.
+
+
#{{ f.image_id }} — {{ f.reason }}
+
+
+
+
+ Close
+
+
+
+
+
+
+ {{ confirmTitle }}
+ {{ confirmMessage }}
+
+
+ Cancel
+ Go ahead
+
+
+
+
+
+
+
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
+}
--
2.54.0