diff --git a/migrations/versions/0014_restore_plans.py b/migrations/versions/0014_restore_plans.py new file mode 100644 index 0000000..6fb35eb --- /dev/null +++ b/migrations/versions/0014_restore_plans.py @@ -0,0 +1,37 @@ +"""Restore plans and archive divergence (US06-04). + +Revision ID: 0014_restore_plans +Revises: 0013_protected_thumbnails +Create Date: 2026-08-16 + +Restore reuses the archive plan and journal tables: the crash-safe question is the +same one in the opposite direction (copy, verify, publish, register), so the rows +gain a ``direction`` instead of a parallel pair of tables. ``archive_divergent_at`` +records the moment an archived copy was proven to hold bytes that are not the ones +the database recorded — a restore must never silently accept a different file. +""" + +import sqlalchemy as sa +from alembic import op + +revision = "0014_restore_plans" +down_revision = "0013_protected_thumbnails" +branch_labels = None +depends_on = None + + +def upgrade() -> None: + for table in ("archive_plans", "archive_operations"): + op.add_column( + table, + sa.Column("direction", sa.String(), nullable=False, server_default="archive"), + ) + op.add_column( + "assets", sa.Column("archive_divergent_at", sa.DateTime(timezone=True), nullable=True) + ) + + +def downgrade() -> None: + op.drop_column("assets", "archive_divergent_at") + for table in ("archive_plans", "archive_operations"): + op.drop_column(table, "direction") diff --git a/photo_pipeline/api/routes/archives.py b/photo_pipeline/api/routes/archives.py index 0cd3eca..20d6514 100644 --- a/photo_pipeline/api/routes/archives.py +++ b/photo_pipeline/api/routes/archives.py @@ -1,4 +1,4 @@ -"""Archive location, preflight, and plan API (US06-01, US06-02). +"""Archive location, preflight, plan, and restore API (US06-01, US06-02, US06-04). Registering a location writes a marker onto the medium; preflight is a command rather than a read, because it probes the destination, hashes the scope, and issues @@ -13,10 +13,11 @@ from fastapi import APIRouter, Request from fastapi.responses import JSONResponse from pydantic import BaseModel -from photo_pipeline.jobs.domain_handlers import ARCHIVE_LOCK, ARCHIVE_PLAN +from photo_pipeline.jobs.domain_handlers import ARCHIVE_LOCK, ARCHIVE_PLAN, RESTORE_PLAN from photo_pipeline.services.archives import ArchiveError, ArchiveService from photo_pipeline.services.archive_transfer import ArchiveTransferService from photo_pipeline.services.jobs import JobBlocked, JobService +from photo_pipeline.services.restores import RestoreService router = APIRouter(tags=["archives"]) @@ -42,10 +43,24 @@ class CreatePlanRequest(PreflightRequest): token: str +class RestoreRequest(BaseModel): + location_id: str + # ``None`` means every asset archived at this location. + asset_ids: list[str] | None = None + + +class CreateRestoreRequest(RestoreRequest): + token: str + + def _service(request: Request) -> ArchiveService: return ArchiveService(request.app.state.session_factory, config=request.app.state.config) +def _restores(request: Request) -> RestoreService: + return RestoreService(request.app.state.session_factory, config=request.app.state.config) + + def _transfers(request: Request) -> ArchiveTransferService: return ArchiveTransferService( request.app.state.session_factory, config=request.app.state.config @@ -138,6 +153,76 @@ def apply_plan(plan_id: str, request: Request): return {"plan_id": plan_id, "job": job} +@router.post("/restore-preflight") +def restore_preflight(body: RestoreRequest, request: Request): + """Validate restoring archived assets back into the library. Nothing moves.""" + try: + return _restores(request).preflight(body.location_id, body.asset_ids) + except ArchiveError as error: + return _error(error) + + +@router.post("/restore-plans", status_code=201) +def create_restore_plan(body: CreateRestoreRequest, request: Request): + try: + return _restores(request).create(body.location_id, body.asset_ids, token=body.token) + except ArchiveError as error: + return _error(error) + + +@router.get("/restore-plans") +def list_restore_plans(request: Request) -> dict: + return {"plans": _restores(request).list()} + + +@router.get("/restore-plans/{plan_id}") +def get_restore_plan(plan_id: str, request: Request): + plan = _restores(request).get(plan_id) + if plan is None: + return _error(ArchiveError("unknown_plan", f"unknown restore plan {plan_id}")) + return plan + + +@router.post("/restore-plans/{plan_id}/apply") +def apply_restore_plan(plan_id: str, request: Request): + """Queue the restore on the archiver lane — the same single lane as archiving, + because both move the same originals.""" + service = _restores(request) + plan = service.get(plan_id) + if plan is None: + return _error(ArchiveError("unknown_plan", f"unknown restore plan {plan_id}")) + unresolved = [row for row in service.journal.incomplete() if row["plan_id"] != plan_id] + if unresolved: + return _error( + ArchiveError( + "archive_pending", + f"an unresolved archive operation ({unresolved[0]['id']}) must be recovered", + ) + ) + try: + job = JobService(request.app.state.session_factory).enqueue( + RESTORE_PLAN, + lock=ARCHIVE_LOCK, + idempotency_key=f"restore:{plan_id}:{plan['version']}", + items=[plan_id], + ) + except JobBlocked as error: + return JSONResponse( + status_code=409, content={"error": {"code": error.code, "message": str(error)}} + ) + return {"plan_id": plan_id, "job": job} + + +@router.get("/restore-recovery") +def restore_recovery_status(request: Request) -> dict: + return _restores(request).recovery_status() + + +@router.post("/restore-recovery/resolve") +def resolve_restore_recovery(request: Request) -> dict: + return _restores(request).recover() + + @router.get("/archive-recovery") def recovery_status(request: Request) -> dict: """What an interrupted transfer left behind, straight from journal + disk.""" diff --git a/photo_pipeline/jobs/domain_handlers.py b/photo_pipeline/jobs/domain_handlers.py index 2ef7cff..2ce89ce 100644 --- a/photo_pipeline/jobs/domain_handlers.py +++ b/photo_pipeline/jobs/domain_handlers.py @@ -1,9 +1,9 @@ """Domain job handlers: safety scoring, content analysis, uploads, archive -transfers (US02-06, US05-02, US06-02). +transfers, restores (US02-06, US05-02, US06-02, US06-04). Importing this module registers the ``safety_score``, ``analysis``, -``upload_batch``, and ``archive_plan`` job types so the generic worker can run them -per item. Each handler delegates to its service, which owns the real work and the +``upload_batch``, ``archive_plan``, and ``restore_plan`` job types so the generic +worker can run them per item. Each handler delegates to its service, which owns the real work and the privacy gate. Handlers are idempotent: re-scoring or re-analyzing one asset is safe after an interrupted attempt, an upload batch refuses to re-run an attempt whose outcome is unknown, and an archive plan skips items it already completed. @@ -20,6 +20,7 @@ SAFETY_SCORE = "safety_score" ANALYSIS = "analysis" UPLOAD_BATCH = "upload_batch" ARCHIVE_PLAN = "archive_plan" +RESTORE_PLAN = "restore_plan" # Both mutate the library's metadata/derived state; one at a time (concept §one job). LIBRARY_WRITE_LOCK = "library_write" # The uploader lane: one album batch at a time (concept §16). @@ -70,7 +71,22 @@ def _archive_plan_item(plan_id: str, ctx: JobContext) -> None: raise RuntimeError(f"archive plan {plan_id}: {result['failed']} item(s) failed") +def _restore_plan_item(plan_id: str, ctx: JobContext) -> None: + """One item = one restore plan. A restore removes nothing, so an item failure + simply leaves that asset archived (US06-04).""" + from photo_pipeline.config import Config + from photo_pipeline.services.restores import RestoreService + + config = ctx.config if ctx.config is not None else Config.from_env() + result = RestoreService(ctx.session_factory, config=config).apply( + plan_id, worker_id=ctx.worker_id + ) + if result["failed"]: + raise RuntimeError(f"restore plan {plan_id}: {result['failed']} item(s) failed") + + register(SAFETY_SCORE, _safety_score_item) register(ANALYSIS, _analysis_item) register(UPLOAD_BATCH, _upload_batch_item) register(ARCHIVE_PLAN, _archive_plan_item) +register(RESTORE_PLAN, _restore_plan_item) diff --git a/photo_pipeline/models/archives.py b/photo_pipeline/models/archives.py index c03e8ce..4660a7e 100644 --- a/photo_pipeline/models/archives.py +++ b/photo_pipeline/models/archives.py @@ -66,6 +66,8 @@ class ArchivePlan(Base): # The preflight token this plan was approved against; re-verified before apply. token: Mapped[str] = mapped_column(String, nullable=False) albums: Mapped[str | None] = mapped_column(String) # JSON array + # archive | restore — the same journal read in the opposite direction (US06-04). + direction: Mapped[str] = mapped_column(String, nullable=False, default="archive") # planned | applying | complete | failed state: Mapped[str] = mapped_column(String, nullable=False, default="planned") @@ -100,6 +102,9 @@ class ArchiveOperation(Base): album: Mapped[str] = mapped_column(String, nullable=False) asset_id: Mapped[str] = mapped_column(ForeignKey("assets.id"), nullable=False, index=True) + # archive: library → medium. restore: medium → library (US06-04). ``source_path`` + # and ``destination_path`` always mean "from" and "to" for this direction. + direction: Mapped[str] = mapped_column(String, nullable=False, default="archive") source_path: Mapped[str] = mapped_column(String, nullable=False) destination_path: Mapped[str] = mapped_column(String, nullable=False) # Relative to the location root, because the medium can be mounted elsewhere. diff --git a/photo_pipeline/models/assets.py b/photo_pipeline/models/assets.py index 52075c5..2a9b366 100644 --- a/photo_pipeline/models/assets.py +++ b/photo_pipeline/models/assets.py @@ -43,6 +43,10 @@ class Asset(Base): # link is written and read by the archive service (US06-02). archive_location_id: Mapped[str | None] = mapped_column(String) archive_path: Mapped[str | None] = mapped_column(String) + # Set when the archived copy was proven to hold bytes other than the recorded + # ones (US06-04). Restore refuses such an asset instead of accepting a different + # file; cleared as soon as a verification matches again. + archive_divergent_at: Mapped[datetime | None] = mapped_column(DateTime(timezone=True)) # Duplicate canonical link: NULL when the asset is itself canonical or undecided. canonical_asset_id: Mapped[str | None] = mapped_column(ForeignKey("assets.id")) diff --git a/photo_pipeline/services/archive_journal.py b/photo_pipeline/services/archive_journal.py index 08b3093..b0c9412 100644 --- a/photo_pipeline/services/archive_journal.py +++ b/photo_pipeline/services/archive_journal.py @@ -15,6 +15,11 @@ planned → transferring → verified → removing → complete ↘ ↘ ↘ failed ``` +A restore (US06-04) uses the same rows with ``direction='restore'``: it copies from +the medium back into the library and removes nothing, so it goes ``verified → +complete`` directly. ``source_path``/``destination_path`` always mean "from"/"to", +which is why the evidence table below needs no direction of its own. + - ``transferring`` — intent recorded; a temporary copy may exist, the destination may or may not have been published. Nothing has been removed. - ``verified`` — the archived bytes exist at their final path, hash exactly as @@ -73,10 +78,22 @@ ALLOWED_TRANSITIONS = { ArchiveState.FAILED: {ArchiveState.PLANNED, ArchiveState.TRANSFERRING}, } +# A restore removes nothing, so it has no ``removing`` step: a verified published +# copy is the whole job (US06-04). Keeping this as a separate table means the +# archive direction still cannot reach ``complete`` without going through removal. +RESTORE_TRANSITIONS = { + **ALLOWED_TRANSITIONS, + ArchiveState.VERIFIED: {ArchiveState.COMPLETE, ArchiveState.FAILED}, +} + TERMINAL_STATES = frozenset({ArchiveState.COMPLETE}) # States where this item may already have touched the filesystem. UNSAFE_STATES = frozenset({ArchiveState.TRANSFERRING, ArchiveState.VERIFIED, ArchiveState.REMOVING}) +# Which way the bytes move. Same rows, same evidence table, opposite direction. +ARCHIVE = "archive" +RESTORE = "restore" + RESUMABLE = "resumable" FORWARD = "forward" MANUAL = "manual" @@ -94,8 +111,9 @@ class JournalConflict(JournalError): """Fencing check failed; a newer owner has taken over this operation.""" -def can_transition(current: str, target: str) -> bool: - return target in ALLOWED_TRANSITIONS.get(current, set()) +def can_transition(current: str, target: str, direction: str = ARCHIVE) -> bool: + table = RESTORE_TRANSITIONS if direction == RESTORE else ALLOWED_TRANSITIONS + return target in table.get(current, set()) def _now() -> datetime: @@ -119,7 +137,7 @@ class ArchiveJournal: if row.journal_state in TERMINAL_STATES: raise InvalidTransition(f"{row.journal_state} is terminal") if row.journal_state != ArchiveState.TRANSFERRING and not can_transition( - row.journal_state, ArchiveState.TRANSFERRING + row.journal_state, ArchiveState.TRANSFERRING, row.direction ): raise InvalidTransition(f"{row.journal_state} -> {ArchiveState.TRANSFERRING}") if row.journal_state != ArchiveState.TRANSFERRING: @@ -160,7 +178,7 @@ class ArchiveJournal: if row.journal_state == target: session.commit() return _operation_dict(row) # idempotent - if not can_transition(row.journal_state, target): + if not can_transition(row.journal_state, target, row.direction): raise InvalidTransition(f"{row.journal_state} -> {target}") row.journal_state = target @@ -192,16 +210,18 @@ class ArchiveJournal: ) return [_operation_dict(row) for row in rows] - def incomplete(self) -> list[dict]: + def incomplete(self, *, direction: str | None = None) -> list[dict]: """Every operation left in a non-terminal, non-planned state — the work a - restart has to reason about.""" + restart has to reason about. Without ``direction`` this spans archives and + restores, because either one half-done blocks the other.""" with self._session_factory() as session: + stmt = select(ArchiveOperation).where( + ArchiveOperation.journal_state.not_in([*TERMINAL_STATES, ArchiveState.PLANNED]) + ) + if direction is not None: + stmt = stmt.where(ArchiveOperation.direction == direction) rows = session.scalars( - select(ArchiveOperation) - .where( - ArchiveOperation.journal_state.not_in([*TERMINAL_STATES, ArchiveState.PLANNED]) - ) - .order_by(ArchiveOperation.plan_id, ArchiveOperation.sequence) + stmt.order_by(ArchiveOperation.plan_id, ArchiveOperation.sequence) ) return [_operation_dict(row) for row in rows] @@ -231,6 +251,7 @@ class ArchiveJournal: return { "operation_id": operation_id, "plan_id": row["plan_id"], + "direction": row["direction"], "album": row["album"], "asset_id": row["asset_id"], "source_path": row["source_path"], @@ -243,8 +264,8 @@ class ArchiveJournal: "destination_matches": destination_matches, } - def classify_all(self) -> list[dict]: - return [self.classify(row["id"]) for row in self.incomplete()] + def classify_all(self, *, direction: str | None = None) -> list[dict]: + return [self.classify(row["id"]) for row in self.incomplete(direction=direction)] def blocks_mutation(self) -> bool: """True when any item may have the library half-archived.""" @@ -317,6 +338,7 @@ def _operation_dict(row: ArchiveOperation) -> dict: return { "id": row.id, "plan_id": row.plan_id, + "direction": row.direction, "sequence": row.sequence, "album": row.album, "asset_id": row.asset_id, diff --git a/photo_pipeline/services/archive_transfer.py b/photo_pipeline/services/archive_transfer.py index a8efa1a..754171b 100644 --- a/photo_pipeline/services/archive_transfer.py +++ b/photo_pipeline/services/archive_transfer.py @@ -54,6 +54,7 @@ from sqlalchemy.orm import sessionmaker from photo_pipeline.config import Config from photo_pipeline.models import ArchiveLocation, ArchiveOperation, ArchivePlan, Asset, AssetPath from photo_pipeline.services.archive_journal import ( + ARCHIVE, MANUAL, RESUMABLE, ArchiveJournal, @@ -115,6 +116,7 @@ class ArchiveTransferService: location_id=location_id, token=token, albums=json.dumps(albums) if albums is not None else None, + direction=ARCHIVE, state="planned", schema_version=MANIFEST_VERSION, asset_count=preflight["totals"]["assets"], @@ -132,6 +134,7 @@ class ArchiveTransferService: ArchiveOperation( id=str(uuid.uuid4()), plan_id=plan_id, + direction=ARCHIVE, sequence=sequence, album=album["album"], asset_id=asset["asset_id"], @@ -161,7 +164,11 @@ class ArchiveTransferService: def list(self) -> list[dict]: with self._session_factory() as session: - rows = session.scalars(select(ArchivePlan).order_by(ArchivePlan.created_at)) + rows = session.scalars( + select(ArchivePlan) + .where(ArchivePlan.direction == ARCHIVE) + .order_by(ArchivePlan.created_at) + ) return [_plan_dict(row) for row in rows] # ── apply ───────────────────────────────────────────────────────────────── @@ -262,7 +269,7 @@ class ArchiveTransferService: if same_filesystem: os.rename(source, destination) else: - self._copy_and_publish(operation, source, destination) + copy_verify_publish(source, destination, operation["expected_sha256"]) _fsync_dir(destination.parent) # 4. The published file is the archive only once it hashes as recorded. @@ -282,27 +289,6 @@ class ArchiveTransferService: # 5. Only now may the active source go. self._finish(self.journal.get(operation["id"]), location, token=token, worker_id=worker_id) - def _copy_and_publish(self, operation: dict, source: Path, destination: Path) -> None: - """Cross-filesystem: copy to a temporary file beside the destination, prove - its bytes, then publish it atomically. The source is still untouched.""" - temp = destination.with_name(f"{TEMP_PREFIX}{uuid.uuid4().hex}{TEMP_SUFFIX}") - try: - with open(source, "rb") as src, open(temp, "wb") as out: - shutil.copyfileobj(src, out, 1024 * 1024) - out.flush() - os.fsync(out.fileno()) - if sha256_file(temp) != operation["expected_sha256"]: - raise PreconditionFailed("copy_mismatch", f"{source} copied with wrong bytes") - if destination.exists(): - raise PreconditionFailed( - "destination_exists", f"{destination} appeared during the transfer" - ) - # ponytail: rename after an exists() check. The archiver lane is single - # and local; use O_EXCL/link-based publish if a second writer ever exists. - os.rename(temp, destination) - finally: - temp.unlink(missing_ok=True) - def _finish(self, operation: dict, location: dict, *, token: int, worker_id: str) -> None: """Drive an item whose archive copy is durable through removal and bookkeeping. Every step is idempotent, so recovery may replay it.""" @@ -457,7 +443,7 @@ class ArchiveTransferService: """ results = {"resumed": 0, "completed": 0, "manual": 0} touched: set[str] = set() - for verdict in self.journal.classify_all(): + for verdict in self.journal.classify_all(direction=ARCHIVE): operation = self.journal.get(verdict["operation_id"]) touched.add(operation["plan_id"]) token = (operation["fencing_token"] or 0) + 1 @@ -484,7 +470,7 @@ class ArchiveTransferService: return results def recovery_status(self) -> dict: - verdicts = self.journal.classify_all() + verdicts = self.journal.classify_all(direction=ARCHIVE) return { "operations": verdicts, "manual": [v for v in verdicts if v["classification"] == MANUAL], @@ -528,6 +514,33 @@ class ArchiveTransferService: # ── module helpers ─────────────────────────────────────────────────────────── +def copy_verify_publish(source: Path, destination: Path, expected_sha256: str) -> None: + """Copy to a temporary file beside the destination, prove its bytes, then publish + it atomically. The source is never touched, so a failure costs nothing. + + Shared by archiving (library → medium) and restoring (medium → library, US06-04): + both need the same promise that a published file is either complete and correct + or not there at all. + """ + temp = destination.with_name(f"{TEMP_PREFIX}{uuid.uuid4().hex}{TEMP_SUFFIX}") + try: + with open(source, "rb") as src, open(temp, "wb") as out: + shutil.copyfileobj(src, out, 1024 * 1024) + out.flush() + os.fsync(out.fileno()) + if sha256_file(temp) != expected_sha256: + raise PreconditionFailed("copy_mismatch", f"{source} copied with wrong bytes") + if destination.exists(): + raise PreconditionFailed( + "destination_exists", f"{destination} appeared during the transfer" + ) + # ponytail: rename after an exists() check. The archiver lane is single and + # local; use O_EXCL/link-based publish if a second writer ever exists. + os.rename(temp, destination) + finally: + temp.unlink(missing_ok=True) + + def _same_filesystem(source: Path, destination_dir: Path) -> bool: """Proven at run time from the actual devices, never from the plan's preview.""" try: @@ -618,6 +631,7 @@ def _plan_dict(plan: ArchivePlan) -> dict: "id": plan.id, "location_id": plan.location_id, "token": plan.token, + "direction": plan.direction, "albums": json.loads(plan.albums) if plan.albums else None, "state": plan.state, "schema_version": plan.schema_version, diff --git a/photo_pipeline/services/restores.py b/photo_pipeline/services/restores.py new file mode 100644 index 0000000..d755227 --- /dev/null +++ b/photo_pipeline/services/restores.py @@ -0,0 +1,648 @@ +"""RestoreService — plan and execute safe restores (US06-04). + +Restore is archiving read backwards, with one decisive difference: it removes +nothing. The archived copy stays on its medium, so every failure mode here costs +at most a discarded temporary file. What restore must never do is *lose identity* +— the asset that comes back is the same asset, with its duplicate decision, safety +review, analysis, and upload history intact — or *overwrite* something in the +active library. + +Preflight proves, per concept §9 "Restore": + +- the recorded medium is mounted and is the right one (marker ``media_id``); +- every selected asset is archived, its archive copy exists, and it hashes to + exactly the bytes the database recorded — a mismatch is ``divergent`` and is + refused, never silently accepted as "the file"; +- the destination lies inside the library, outside ``_IGNORE/``, and is free; a + taken path is answered with a collision-free name, never an overwrite; +- the library filesystem has room for the scope plus the configured reserve; +- no rename, archive, or restore lease is holding the lane. + +Blocker codes: ``no_library_root``, ``location_offline``, ``wrong_volume``, +``unsafe_destination``, ``library_not_writable``, ``insufficient_capacity``, +``lock_conflict``, ``rename_pending``, ``archive_pending``, ``empty_scope``, +``not_archived``, ``archive_missing``, ``bytes_changed``. + +Per item the sequence is: + +``` +journal.begin (transferring) ← intent persisted BEFORE any disk change +recheck: medium, hash, free destination, asset still archived +copy to a temporary file beside the destination, fsync, hash it back +atomically publish into the library +journal → verified +current_path = destination, availability = active, path occurrence opened +journal → complete +``` + +Like archiving, the confirmation token is derived from the report, so a changed +scope, a swapped medium, or a destination that filled up invalidates it. +""" + +from __future__ import annotations + +import hashlib +import json +import os +import shutil +import uuid +from datetime import datetime, timezone +from pathlib import Path + +from sqlalchemy import select +from sqlalchemy.orm import sessionmaker + +from photo_pipeline.config import Config +from photo_pipeline.jobs.domain_handlers import ARCHIVE_LOCK, LIBRARY_WRITE_LOCK, UPLOAD_LOCK +from photo_pipeline.models import ArchiveLocation, ArchiveOperation, ArchivePlan, Asset, AssetPath +from photo_pipeline.path_policy import PathPolicyError, is_excluded, normalize_root, resolve_within +from photo_pipeline.services import availability +from photo_pipeline.services.archive_journal import ( + MANUAL, + RESTORE, + RESUMABLE, + ArchiveJournal, + ArchiveState, +) +from photo_pipeline.services.archive_transfer import ( + _clean_temp_files, + _fsync_dir, + _plan_dict, + copy_verify_publish, +) +from photo_pipeline.services.archives import ArchiveError +from photo_pipeline.services.hashing import sha256_file +from photo_pipeline.services.jobs import JobService +from photo_pipeline.services.rename_apply import PreconditionFailed, maybe_fault +from photo_pipeline.services.rename_journal import RenameJournal + +PREFLIGHT_VERSION = 1 +TOKEN_PREFIX = f"r{PREFLIGHT_VERSION}" +# What a restored file is called when its original name is taken. The suffix is +# visible on purpose: a restore that quietly reuses a name is indistinguishable +# from an overwrite. +RESTORED_SUFFIX = "restored" + +LOCKS = (LIBRARY_WRITE_LOCK, UPLOAD_LOCK, ARCHIVE_LOCK) + +APPLYABLE_PLAN_STATES = frozenset({"planned", "applying", "failed", "complete"}) + + +def _now() -> datetime: + return datetime.now(timezone.utc) + + +def _issue(code: str, message: str) -> dict: + return {"code": code, "message": message} + + +class RestoreService: + def __init__(self, session_factory: sessionmaker, *, config: Config) -> None: + self._session_factory = session_factory + self._config = config + self._roots = tuple(normalize_root(root) for root in config.library_roots) + self.journal = ArchiveJournal(session_factory) + + # ── preflight ───────────────────────────────────────────────────────────── + + def preflight(self, location_id: str, asset_ids: list[str] | None = None) -> dict: + """Validate a restore scope and issue its token. Nothing is written.""" + with self._session_factory() as session: + location = session.get(ArchiveLocation, location_id) + if location is None: + raise ArchiveError("unknown_location", f"unknown archive location {location_id!r}") + root = Path(location.root) + online = availability.location_online(location) + marker = availability.read_marker(root) + report = { + "schema_version": PREFLIGHT_VERSION, + "location": { + "id": location.id, + "name": location.name, + "root": str(root), + "media_id": location.media_id, + "state": _location_state(root, marker, location.media_id), + }, + "blockers": [], + } + items = self._items(session, location, asset_ids, reachable=online) + + report["blockers"] += self._destination_blockers(report["location"]["state"], root) + report["blockers"] += self._lock_blockers() + report["items"] = items + report["totals"] = { + "assets": len(items), + "blocked": sum(1 for item in items if item["blockers"]), + "bytes": sum(item["byte_size"] or 0 for item in items), + } + report["capacity"] = self._capacity(report["totals"]["bytes"]) + if not report["capacity"]["sufficient"]: + report["blockers"].append( + _issue( + "insufficient_capacity", + f"{report['totals']['bytes']} B plus a " + f"{self._config.archive_free_space_reserve_bytes} B reserve do not fit in " + f"{report['capacity']['free_bytes']} B of free space", + ) + ) + if not items: + report["blockers"].append( + _issue("empty_scope", "no archived assets are in the selected scope") + ) + report["state"] = ( + "ready" + if not report["blockers"] and not report["totals"]["blocked"] + else "blocked" + ) + report["token"] = _token(report) + report["generated_at"] = _now().isoformat() + return report + + def verify_token(self, token: str, location_id: str, asset_ids: list[str] | None = None) -> bool: + return bool(token) and token == self.preflight(location_id, asset_ids)["token"] + + def _items( + self, session, location: ArchiveLocation, asset_ids: list[str] | None, *, reachable: bool + ) -> list[dict]: + stmt = select(Asset).where(Asset.archive_location_id == location.id) + if asset_ids is None: + # A restored asset keeps its archive link; the default scope is only what + # is still archived, so restoring twice is an empty scope, not a blocker. + stmt = stmt.where(Asset.availability_state.in_(availability.ARCHIVED)) + else: + stmt = stmt.where(Asset.id.in_(asset_ids)) + assets = list(session.scalars(stmt.order_by(Asset.archive_path))) + if asset_ids is not None: + unknown = sorted(set(asset_ids) - {asset.id for asset in assets}) + if unknown: + raise ArchiveError( + "unknown_asset", f"not archived at this location: {', '.join(unknown)}" + ) + taken: set[str] = set() + return [self._item(asset, location, reachable=reachable, taken=taken) for asset in assets] + + def _item(self, asset: Asset, location: ArchiveLocation, *, reachable: bool, taken: set) -> dict: + source = Path(location.root) / (asset.archive_path or "") + blockers: list[dict] = [] + archive_sha256 = None + + if asset.availability_state not in availability.ARCHIVED: + blockers.append( + _issue("not_archived", f"asset {asset.id} is {asset.availability_state}") + ) + if reachable: + if not source.exists(): + blockers.append(_issue("archive_missing", f"{source} is not on the medium")) + else: + archive_sha256 = sha256_file(source) + if asset.current_sha256 and archive_sha256 != asset.current_sha256: + blockers.append( + _issue( + "bytes_changed", + f"{source} holds bytes that are not the recorded ones; " + "the archived copy is divergent", + ) + ) + + destination, destination_blockers = self._destination(asset, taken) + blockers += destination_blockers + if destination is not None: + taken.add(str(destination)) + return { + "asset_id": asset.id, + "archive_path": asset.archive_path, + "source_path": str(source), + "destination_path": str(destination) if destination else None, + "expected_sha256": asset.current_sha256, + "archive_sha256": archive_sha256, + "byte_size": asset.byte_size, + "availability_state": asset.availability_state, + "blockers": blockers, + } + + def _destination(self, asset: Asset, taken: set) -> tuple[Path | None, list[dict]]: + """A free path inside the library that mirrors the archived layout. + + Restoring onto an existing file is never an option, so a taken name is + answered with ``name (restored).ext`` — visible, ordinary, and impossible to + confuse with an overwrite. + """ + if not self._roots: + return None, [_issue("no_library_root", "no library root is configured")] + root = self._roots[0] + try: + candidate = resolve_within(root, root / (asset.archive_path or "")) + except PathPolicyError as error: + return None, [_issue("unsafe_destination", str(error))] + if is_excluded(candidate): + return None, [ + _issue("unsafe_destination", f"{candidate} is inside an excluded (_IGNORE/) tree") + ] + return _free_path(candidate, taken), [] + + def _destination_blockers(self, state: str, root: Path) -> list[dict]: + blockers: list[dict] = [] + if not self._roots: + blockers.append(_issue("no_library_root", "no library root is configured")) + elif not os.access(self._roots[0], os.W_OK): + blockers.append( + _issue("library_not_writable", f"{self._roots[0]} is not writable") + ) + if state == "offline": + blockers.append( + _issue("location_offline", f"the archive medium is not mounted at {root}") + ) + elif state == "wrong_volume": + blockers.append(_issue("wrong_volume", f"{root} holds a different archive medium")) + return blockers + + def _lock_blockers(self) -> list[dict]: + blockers: list[dict] = [] + jobs = JobService(self._session_factory) + for lock in LOCKS: + held = jobs.blockers(lock) + if held: + blockers.append( + _issue("lock_conflict", f"the {lock} lane is busy: job {held[0]['id']}") + ) + if RenameJournal(self._session_factory).blocks_mutation(): + blockers.append( + _issue("rename_pending", "an unresolved rename must be recovered before restoring") + ) + if self.journal.blocks_mutation(): + blockers.append( + _issue( + "archive_pending", + "an unresolved archive or restore must be recovered before restoring", + ) + ) + return blockers + + def _capacity(self, required: int) -> dict: + reserve = self._config.archive_free_space_reserve_bytes + free = shutil.disk_usage(self._roots[0]).free if self._roots else None + return { + "required_bytes": required, + "reserve_bytes": reserve, + "free_bytes": free, + "sufficient": free is not None and free >= required + reserve, + } + + # ── plans ───────────────────────────────────────────────────────────────── + + def create(self, location_id: str, asset_ids: list[str] | None = None, *, token: str) -> dict: + preflight = self.preflight(location_id, asset_ids) + if not token or token != preflight["token"]: + raise ArchiveError("stale_token", "the restore preflight changed since it was approved") + if preflight["state"] != "ready": + codes = ", ".join(sorted({issue["code"] for issue in preflight["blockers"]})) or "-" + blocked = sorted( + {issue["code"] for item in preflight["items"] for issue in item["blockers"]} + ) + raise ArchiveError( + "blocked", f"the restore scope is blocked: {', '.join(blocked) or codes}" + ) + + plan_id = str(uuid.uuid4()) + with self._session_factory() as session: + session.add( + ArchivePlan( + id=plan_id, + location_id=location_id, + token=token, + albums=json.dumps(asset_ids) if asset_ids is not None else None, + direction=RESTORE, + state="planned", + schema_version=PREFLIGHT_VERSION, + asset_count=preflight["totals"]["assets"], + byte_size=preflight["totals"]["bytes"], + ) + ) + session.flush() + for sequence, item in enumerate(preflight["items"]): + session.add( + ArchiveOperation( + id=str(uuid.uuid4()), + plan_id=plan_id, + direction=RESTORE, + sequence=sequence, + album=Path(item["archive_path"]).parent.name or "(root)", + asset_id=item["asset_id"], + source_path=item["source_path"], + destination_path=item["destination_path"], + archive_path=item["archive_path"], + expected_sha256=item["expected_sha256"], + byte_size=item["byte_size"], + journal_state=ArchiveState.PLANNED, + ) + ) + session.commit() + return self.get(plan_id) + + def get(self, plan_id: str) -> dict | None: + with self._session_factory() as session: + plan = session.get(ArchivePlan, plan_id) + if plan is None or plan.direction != RESTORE: + return None + report = _plan_dict(plan) + report["operations"] = self.journal.operations(plan_id) + return report + + def list(self) -> list[dict]: + with self._session_factory() as session: + rows = session.scalars( + select(ArchivePlan) + .where(ArchivePlan.direction == RESTORE) + .order_by(ArchivePlan.created_at) + ) + return [_plan_dict(row) for row in rows] + + # ── apply ───────────────────────────────────────────────────────────────── + + def apply( + self, plan_id: str, *, expected_version: int | None = None, worker_id: str = "restore" + ) -> dict: + plan = self._require_plan(plan_id) + if expected_version is not None and plan["version"] != expected_version: + raise ArchiveError( + "stale_plan", + f"plan {plan_id} is at version {plan['version']}, expected {expected_version}", + ) + if plan["state"] not in APPLYABLE_PLAN_STATES: + raise ArchiveError("invalid_state", f"plan {plan_id} is {plan['state']}") + blocking = [row for row in self.journal.incomplete() if row["plan_id"] != plan_id] + if blocking: + raise ArchiveError( + "archive_pending", + f"another archive operation is unresolved ({blocking[0]['id']}); recover it first", + ) + + token = self._claim_plan(plan_id) + location = self._location(plan["location_id"]) + restored = failed = skipped = 0 + for operation in self.journal.operations(plan_id): + if operation["journal_state"] == ArchiveState.COMPLETE: + skipped += 1 + continue + try: + if operation["journal_state"] == ArchiveState.VERIFIED: + self._finish(operation, token=token) + else: + self._restore_one(operation, location, token=token, worker_id=worker_id) + restored += 1 + except PreconditionFailed as error: + self._fail(operation, token, error.code, str(error)) + failed += 1 + except Exception as error: # unexpected: record and stop touching disk + self._fail(operation, token, "restore_error", str(error)) + failed += 1 + state = self.journal.sync_plan_state(plan_id) + return { + "plan_id": plan_id, + "restored": restored, + "failed": failed, + "skipped": skipped, + "state": state, + } + + def _restore_one(self, operation: dict, location: dict, *, token: int, worker_id: str) -> None: + source = Path(operation["source_path"]) + destination = Path(operation["destination_path"]) + + # 1. Intent first; from here a crash is resolvable from journal + disk. + self.journal.begin(operation["id"], worker_id=worker_id, fencing_token=token) + maybe_fault(ArchiveState.TRANSFERRING) + + # 2. Recheck against the medium and the library as they are right now. + self._recheck(operation, source, destination, location) + destination.parent.mkdir(parents=True, exist_ok=True) + + # 3. Always copy: the archived original stays on its medium. + copy_verify_publish(source, destination, operation["expected_sha256"]) + _fsync_dir(destination.parent) + + if sha256_file(destination) != operation["expected_sha256"]: + raise PreconditionFailed( + "restore_mismatch", f"{destination} does not hold the expected bytes" + ) + self.journal.transition(operation["id"], ArchiveState.VERIFIED, fencing_token=token) + maybe_fault(ArchiveState.VERIFIED) + + self._finish(self.journal.get(operation["id"]), token=token) + + def _finish(self, operation: dict, *, token: int) -> None: + """Publish the restored file to the database. Idempotent, so recovery may + replay it after a crash between the copy and the bookkeeping.""" + destination = Path(operation["destination_path"]) + if not destination.exists() or sha256_file(destination) != operation["expected_sha256"]: + raise PreconditionFailed( + "restore_unverified", f"{destination} is not a verified restored copy" + ) + self._record_restored(operation, destination) + self.journal.transition(operation["id"], ArchiveState.COMPLETE, fencing_token=token) + maybe_fault(ArchiveState.COMPLETE) + + def _recheck(self, operation: dict, source: Path, destination: Path, location: dict) -> None: + root = Path(location["root"]) + if not root.is_dir() or not (root / availability.MARKER_NAME).exists(): + raise PreconditionFailed("location_offline", f"{root} is not the archive medium") + if not source.exists(): + raise PreconditionFailed("archive_missing", f"{source} is not on the medium") + if source.is_symlink() or destination.is_symlink(): + raise PreconditionFailed("symlink", "refusing to restore through a symlink") + if destination.exists(): + # Never overwrite: the plan's free path was taken since it was made. + raise PreconditionFailed( + "destination_exists", f"destination {destination} is occupied" + ) + if not self._inside_library(destination): + raise PreconditionFailed( + "destination_escape", f"{destination} is outside the library roots" + ) + if sha256_file(source) != operation["expected_sha256"]: + self._mark_divergent(operation["asset_id"]) + raise PreconditionFailed( + "bytes_changed", f"{source} changed since the plan was approved" + ) + with self._session_factory() as session: + asset = session.get(Asset, operation["asset_id"]) + if asset is None or asset.availability_state not in availability.ARCHIVED: + raise PreconditionFailed( + "not_archived", f"asset {operation['asset_id']} is no longer archived" + ) + + def _inside_library(self, destination: Path) -> bool: + for root in self._roots: + try: + resolve_within(root, destination) + return True + except PathPolicyError: + continue + return False + + # ── database ────────────────────────────────────────────────────────────── + + def _record_restored(self, operation: dict, destination: Path) -> None: + """The bytes are back in the library: open the new active occurrence and set + availability. Identity, decisions, and history are untouched — that is the + entire point of restoring rather than re-importing.""" + now = _now() + with self._session_factory() as session: + asset = session.get(Asset, operation["asset_id"]) + if asset is None: + raise PreconditionFailed( + "asset_missing", f"asset {operation['asset_id']} no longer exists" + ) + # A restored asset may be returning to a path it once held, so only an + # *open* occurrence counts as already registered — that is what keeps + # recovery idempotent without collapsing the path history. + recorded = session.scalar( + select(AssetPath).where( + AssetPath.asset_id == asset.id, + AssetPath.path == str(destination), + AssetPath.valid_until.is_(None), + ) + ) + if recorded is None: # idempotent: recovery may replay this + session.add( + AssetPath( + asset_id=asset.id, + path=str(destination), + valid_from=now, + reason="restore", + ) + ) + asset.current_path = str(destination) + asset.availability_state = availability.ACTIVE + asset.missing_at = None + # The archive copy stays where it is; keeping the link means a restored + # asset still knows which medium holds its archived bytes. + asset.archive_divergent_at = None + asset.state_version += 1 + asset.updated_at = now + session.commit() + + def _mark_divergent(self, asset_id: str) -> None: + """Record that the archived copy is not the recorded file. Durable, because + the next restore attempt must not rediscover this from scratch.""" + with self._session_factory() as session: + asset = session.get(Asset, asset_id) + if asset is None: + return + asset.archive_divergent_at = _now() + asset.state_version += 1 + session.commit() + + # ── recovery ────────────────────────────────────────────────────────────── + + def recover(self, *, worker_id: str = "restore-recovery") -> dict: + """Resolve every incomplete restore from journal + disk evidence. + + A restore never removed anything, so ``resumable`` simply discards the + temporary debris and re-plans the item; ``forward`` finishes the bookkeeping + for a published file; ``manual`` is left untouched and keeps blocking. + """ + results = {"resumed": 0, "completed": 0, "manual": 0} + touched: set[str] = set() + for verdict in self.journal.classify_all(direction=RESTORE): + operation = self.journal.get(verdict["operation_id"]) + touched.add(operation["plan_id"]) + token = (operation["fencing_token"] or 0) + 1 + if verdict["classification"] == MANUAL: + results["manual"] += 1 + continue + if verdict["classification"] == RESUMABLE: + _clean_temp_files(Path(operation["destination_path"]).parent) + self.journal.transition(operation["id"], ArchiveState.PLANNED, fencing_token=token) + results["resumed"] += 1 + continue + try: + self._finish(operation, token=token) + results["completed"] += 1 + except PreconditionFailed as error: + self._fail(operation, token, error.code, str(error)) + results["manual"] += 1 + for plan_id in touched: + self.journal.sync_plan_state(plan_id) + return results + + def recovery_status(self) -> dict: + verdicts = self.journal.classify_all(direction=RESTORE) + return { + "operations": verdicts, + "manual": [v for v in verdicts if v["classification"] == MANUAL], + "blocks_mutation": self.journal.blocks_mutation(), + } + + # ── helpers ─────────────────────────────────────────────────────────────── + + def _fail(self, operation: dict, token: int, code: str, message: str) -> None: + self.journal.transition( + operation["id"], ArchiveState.FAILED, fencing_token=token, error=(code, message) + ) + + def _require_plan(self, plan_id: str) -> dict: + plan = self.get(plan_id) + if plan is None: + raise ArchiveError("unknown_plan", f"unknown restore plan {plan_id!r}") + return plan + + def _location(self, location_id: str) -> dict: + with self._session_factory() as session: + location = session.get(ArchiveLocation, location_id) + if location is None: + raise ArchiveError("unknown_location", f"unknown archive location {location_id!r}") + return {"id": location.id, "root": location.root, "media_id": location.media_id} + + def _claim_plan(self, plan_id: str) -> int: + with self._session_factory() as session: + plan = session.get(ArchivePlan, plan_id) + plan.version += 1 + plan.state = "applying" + plan.updated_at = _now() + token = plan.version + session.commit() + return token + + +# ── module helpers ─────────────────────────────────────────────────────────── + + +def _location_state(root: Path, marker: dict | None, media_id: str) -> str: + if not root.is_dir() or marker is None: + return "offline" + return "online" if marker.get("media_id") == media_id else "wrong_volume" + + +def _free_path(candidate: Path, taken: set) -> Path: + """``a.jpg`` → ``a (restored).jpg`` → ``a (restored 2).jpg`` … + + ``taken`` holds the destinations already claimed by earlier items of the same + plan, so two restores in one scope cannot plan the same path. + """ + if not candidate.exists() and str(candidate) not in taken: + return candidate + stem, suffix = candidate.stem, candidate.suffix + attempt = 1 + while True: + label = RESTORED_SUFFIX if attempt == 1 else f"{RESTORED_SUFFIX} {attempt}" + alternative = candidate.with_name(f"{stem} ({label}){suffix}") + if not alternative.exists() and str(alternative) not in taken: + return alternative + attempt += 1 + + +def _token(report: dict) -> str: + """Digest of everything the report asserts about the scope and the medium. + + Free space is excluded: it drifts constantly without changing what a restore + would do, and the capacity verdict itself is part of the digest. + """ + payload = {key: value for key, value in report.items() if key not in ("generated_at", "token")} + payload["capacity"] = { + key: value for key, value in payload["capacity"].items() if key != "free_bytes" + } + digest = hashlib.sha256( + json.dumps(payload, sort_keys=True, ensure_ascii=False, default=str).encode("utf-8") + ).hexdigest() + return f"{TOKEN_PREFIX}:{digest}" diff --git a/tests/integration/test_restore.py b/tests/integration/test_restore.py new file mode 100644 index 0000000..c66b51e --- /dev/null +++ b/tests/integration/test_restore.py @@ -0,0 +1,417 @@ +"""Planning and executing safe restores (US06-04). + +Restoring is the one archive operation that can *add* a file to the library, so +every case here asks two questions: did the right bytes come back under the right +identity, and did anything already in the library get touched? The media are real +directories, the hashes are real, and the failure paths assert that the archived +copy is still exactly where it was — a restore that fails must cost nothing. +""" + +import shutil +import uuid +from datetime import datetime, timezone + +import numpy as np +import pytest +from fastapi.testclient import TestClient +from PIL import Image +from sqlalchemy import select + +from photo_pipeline.api.app import create_app +from photo_pipeline.config import Config +from photo_pipeline.db import create_db_engine, create_session_factory, run_migrations +from photo_pipeline.models import Asset, AssetPath, SafetyReview, UploadBatch, UploadItem +from photo_pipeline.services import availability +from photo_pipeline.services.archive_journal import ArchiveState +from photo_pipeline.services.archive_transfer import ArchiveTransferService +from photo_pipeline.services.archives import MARKER_NAME, ArchiveError, ArchiveService +from photo_pipeline.services.hashing import sha256_file +from photo_pipeline.services.inventory import InventoryService +from photo_pipeline.services.restores import RestoreService + +pytestmark = pytest.mark.phase_f # part of the Phase F acceptance gate (US06-06) + +NOW = datetime(2026, 1, 1, tzinfo=timezone.utc) + + +# ── environment ────────────────────────────────────────────────────────────── + + +def _env(tmp_path): + (tmp_path / "data").mkdir(exist_ok=True) + lib = tmp_path / "lib" + lib.mkdir(exist_ok=True) + archive = tmp_path / "archive" + archive.mkdir(exist_ok=True) + config = Config.from_env( + { + "PHOTO_PIPELINE_DATA_DIR": str(tmp_path / "data"), + "PHOTO_PIPELINE_LIBRARY_ROOTS": str(lib), + "PHOTO_PIPELINE_ARCHIVE_FREE_SPACE_RESERVE_BYTES": "0", + } + ) + run_migrations(config.database_url) + return config, create_session_factory(create_db_engine(config.database_url)), lib, archive + + +def structured(path, seed, size=(192, 144)): + path.parent.mkdir(parents=True, exist_ok=True) + rng = np.random.default_rng(seed) + w, h = size + base = np.zeros((h, w, 3), dtype=np.uint8) + for _ in range(5): + x0 = int(rng.integers(0, w - 40)) + y0 = int(rng.integers(0, h - 40)) + base[y0 : y0 + 40, x0 : x0 + 40] = rng.integers(0, 256, 3) + Image.fromarray(base).save(path, quality=95) + return path + + +def _archived(sf, config, lib, archive, album="rome", seeds=(1, 2)): + """A real album taken all the way through archiving, ready to be restored.""" + folder = lib / album + for index, seed in enumerate(seeds): + structured(folder / f"{index}.jpg", seed) + scan = InventoryService(sf).scan(lib) + with sf() as session: + batch_id = str(uuid.uuid4()) + session.add( + UploadBatch( + id=batch_id, + album=album, + folder=str(folder), + album_name=album, + state="succeeded", + preflight_token="v1:test", + outcome_state="verified", + created_at=NOW, + ) + ) + for path, asset_id in scan.asset_ids.items(): + session.add( + UploadItem( + batch_id=batch_id, + asset_id=asset_id, + path=path, + sha256=sha256_file(path), + sha1="0" * 40, + state="sent", + outcome="uploaded", + ) + ) + # A decision that must survive the whole round trip. + session.add( + SafetyReview( + id=str(uuid.uuid4()), + asset_id=asset_id, + decision="sfw", + score=0.01, + reviewer="test", + ) + ) + session.commit() + + service = ArchiveService(sf, config=config) + location = service.register("external", str(archive)) + token = service.preflight(location["id"])["token"] + transfers = ArchiveTransferService(sf, config=config) + plan = transfers.create(location["id"], None, token=token) + transfers.apply(plan["id"]) + return location, scan.asset_ids + + +def _restore(sf, config, location_id, asset_ids=None): + service = RestoreService(sf, config=config) + token = service.preflight(location_id, asset_ids)["token"] + plan = service.create(location_id, asset_ids, token=token) + return service, plan, service.apply(plan["id"]) + + +def _unmount(archive): + (archive / MARKER_NAME).rename(archive / f"{MARKER_NAME}.away") + + +def _assets(sf): + with sf() as session: + return {asset.id: asset for asset in session.scalars(select(Asset))} + + +def _codes(report): + return {issue["code"] for issue in report["blockers"]} | { + issue["code"] for item in report["items"] for issue in item["blockers"] + } + + +# ── preflight ──────────────────────────────────────────────────────────────── + + +def test_preflight_blocks_offline_medium(tmp_path): + config, sf, lib, archive = _env(tmp_path) + location, _ = _archived(sf, config, lib, archive) + _unmount(archive) + + report = RestoreService(sf, config=config).preflight(location["id"]) + assert report["state"] == "blocked" + assert "location_offline" in _codes(report) + + +def test_preflight_blocks_wrong_volume(tmp_path): + config, sf, lib, archive = _env(tmp_path) + location, _ = _archived(sf, config, lib, archive) + (archive / MARKER_NAME).write_text('{"media_id": "someone-elses-disk"}', encoding="utf-8") + + report = RestoreService(sf, config=config).preflight(location["id"]) + assert report["state"] == "blocked" + assert "wrong_volume" in _codes(report) + + +def test_preflight_blocks_changed_archive_bytes(tmp_path): + config, sf, lib, archive = _env(tmp_path) + location, ids = _archived(sf, config, lib, archive, seeds=(1,)) + asset_id = next(iter(ids.values())) + with sf() as session: + archived_file = archive / session.get(Asset, asset_id).archive_path + archived_file.write_bytes(b"not the photo that was archived") + + report = RestoreService(sf, config=config).preflight(location["id"]) + assert report["state"] == "blocked" + assert "bytes_changed" in _codes(report) + with pytest.raises(ArchiveError) as error: + _restore(sf, config, location["id"]) + assert error.value.code == "blocked" + + +def test_preflight_blocks_insufficient_capacity(tmp_path): + config, sf, lib, archive = _env(tmp_path) + location, _ = _archived(sf, config, lib, archive) + greedy = config.model_copy( + update={"archive_free_space_reserve_bytes": 1 << 62} # more than any disk has + ) + + report = RestoreService(sf, config=greedy).preflight(location["id"]) + assert report["state"] == "blocked" + assert "insufficient_capacity" in _codes(report) + + +def test_token_changes_with_the_scope(tmp_path): + config, sf, lib, archive = _env(tmp_path) + location, ids = _archived(sf, config, lib, archive) + service = RestoreService(sf, config=config) + + whole = service.preflight(location["id"])["token"] + partial = service.preflight(location["id"], [sorted(ids.values())[0]])["token"] + assert whole != partial + assert service.verify_token(whole, location["id"]) + assert not service.verify_token(partial, location["id"]) + + +# ── restore ────────────────────────────────────────────────────────────────── + + +def test_restore_returns_bytes_identity_and_decisions(tmp_path): + config, sf, lib, archive = _env(tmp_path) + location, ids = _archived(sf, config, lib, archive) + archived_hashes = { + asset_id: asset.current_sha256 for asset_id, asset in _assets(sf).items() + } + + service, plan, result = _restore(sf, config, location["id"]) + assert (result["restored"], result["failed"], result["state"]) == (2, 0, "complete") + + for asset_id, asset in _assets(sf).items(): + assert asset.availability_state == availability.ACTIVE + assert asset.current_path == str(lib / asset.archive_path) + assert sha256_file(asset.current_path) == archived_hashes[asset_id] + # The archived copy is a copy: restoring never empties the medium. + assert (archive / asset.archive_path).exists() + assert asset.archive_location_id == location["id"] + + with sf() as session: + # Identity and decisions survived: same ids, same reviews, new occurrence. + assert set(ids.values()) == {a.id for a in session.scalars(select(Asset))} + assert {r.decision for r in session.scalars(select(SafetyReview))} == {"sfw"} + occurrences = [ + row.reason + for row in session.scalars( + select(AssetPath).where(AssetPath.asset_id == sorted(ids.values())[0]) + ) + ] + assert "restore" in occurrences + + +def test_restore_never_overwrites_a_collision(tmp_path): + config, sf, lib, archive = _env(tmp_path) + location, ids = _archived(sf, config, lib, archive, seeds=(1,)) + asset_id = next(iter(ids.values())) + with sf() as session: + archive_path = session.get(Asset, asset_id).archive_path + occupied = lib / archive_path + occupied.parent.mkdir(parents=True, exist_ok=True) + occupied.write_bytes(b"a different photo already lives here") + before = occupied.read_bytes() + + service, plan, result = _restore(sf, config, location["id"]) + assert result["failed"] == 0 + + assert occupied.read_bytes() == before # untouched + restored = _assets(sf)[asset_id].current_path + assert restored != str(occupied) + assert "(restored)" in restored + assert sha256_file(restored) == sha256_file(archive / archive_path) + + +def test_apply_refuses_a_destination_taken_after_planning(tmp_path): + config, sf, lib, archive = _env(tmp_path) + location, ids = _archived(sf, config, lib, archive, seeds=(1,)) + service = RestoreService(sf, config=config) + token = service.preflight(location["id"])["token"] + plan = service.create(location["id"], None, token=token) + + # Someone drops a file exactly where the plan intends to publish. + destination = plan["operations"][0]["destination_path"] + from pathlib import Path + + Path(destination).parent.mkdir(parents=True, exist_ok=True) + Path(destination).write_bytes(b"squatter") + + result = service.apply(plan["id"]) + assert result["failed"] == 1 + operation = service.journal.operations(plan["id"])[0] + assert operation["journal_state"] == ArchiveState.FAILED + assert operation["error_code"] == "destination_exists" + assert Path(destination).read_bytes() == b"squatter" + assert _assets(sf)[next(iter(ids.values()))].availability_state == availability.ARCHIVED_ONLINE + + +def test_changed_archive_bytes_mark_the_asset_divergent(tmp_path): + config, sf, lib, archive = _env(tmp_path) + location, ids = _archived(sf, config, lib, archive, seeds=(1,)) + asset_id = next(iter(ids.values())) + service = RestoreService(sf, config=config) + token = service.preflight(location["id"])["token"] + plan = service.create(location["id"], None, token=token) + + # The medium's copy is edited after the plan was approved. + with sf() as session: + archived_file = archive / session.get(Asset, asset_id).archive_path + archived_file.write_bytes(b"edited on the shelf") + + result = service.apply(plan["id"]) + assert result["failed"] == 1 + operation = service.journal.operations(plan["id"])[0] + assert operation["error_code"] == "bytes_changed" + asset = _assets(sf)[asset_id] + assert asset.archive_divergent_at is not None # durable divergence + assert asset.availability_state == availability.ARCHIVED_ONLINE + assert asset.current_path is None # nothing was published + + +# ── interruption and idempotency ───────────────────────────────────────────── + + +def test_interrupted_before_publishing_is_resumable(tmp_path): + config, sf, lib, archive = _env(tmp_path) + location, ids = _archived(sf, config, lib, archive, seeds=(1,)) + service = RestoreService(sf, config=config) + token = service.preflight(location["id"])["token"] + plan = service.create(location["id"], None, token=token) + operation = service.journal.operations(plan["id"])[0] + + # Model a kill right after the intent was written: nothing published yet. + service.journal.begin(operation["id"], worker_id="killed", fencing_token=1) + + status = service.recovery_status() + assert status["operations"][0]["classification"] == "resumable" + assert service.recover() == {"resumed": 1, "completed": 0, "manual": 0} + assert service.journal.operations(plan["id"])[0]["journal_state"] == ArchiveState.PLANNED + + result = service.apply(plan["id"]) + assert result["failed"] == 0 + assert _assets(sf)[next(iter(ids.values()))].availability_state == availability.ACTIVE + + +def test_interrupted_after_publishing_is_finished_by_recovery(tmp_path): + config, sf, lib, archive = _env(tmp_path) + location, ids = _archived(sf, config, lib, archive, seeds=(1,)) + asset_id = next(iter(ids.values())) + service = RestoreService(sf, config=config) + token = service.preflight(location["id"])["token"] + plan = service.create(location["id"], None, token=token) + operation = service.journal.operations(plan["id"])[0] + + # Model a kill between the published copy and the database update. + from pathlib import Path + + destination = Path(operation["destination_path"]) + destination.parent.mkdir(parents=True, exist_ok=True) + shutil.copy2(operation["source_path"], destination) + service.journal.begin(operation["id"], worker_id="killed", fencing_token=1) + service.journal.transition(operation["id"], ArchiveState.VERIFIED, fencing_token=1) + + assert service.recovery_status()["operations"][0]["classification"] == "forward" + assert service.recover()["completed"] == 1 + asset = _assets(sf)[asset_id] + assert asset.availability_state == availability.ACTIVE + assert asset.current_path == str(destination) + + # Repeated recovery and a repeated apply converge on the same state. + assert service.recover() == {"resumed": 0, "completed": 0, "manual": 0} + again = service.apply(plan["id"]) + assert (again["skipped"], again["failed"]) == (1, 0) + assert _assets(sf)[asset_id].current_path == str(destination) + + +def test_restored_state_survives_restart_and_rescan(tmp_path): + config, sf, lib, archive = _env(tmp_path) + location, ids = _archived(sf, config, lib, archive) + _restore(sf, config, location["id"]) + + restarted = create_session_factory(create_db_engine(config.database_url)) + InventoryService(restarted).scan(lib) + + assets = _assets(restarted) + assert set(assets) == set(ids.values()) # no new identities from the rescan + for asset in assets.values(): + assert asset.availability_state == availability.ACTIVE + assert asset.missing_at is None + + # Nothing is archived at that location any more, so there is nothing to restore. + again = RestoreService(restarted, config=config).preflight(location["id"]) + assert _codes(again) == {"empty_scope"} + + +def test_restore_api_round_trip(tmp_path): + config, sf, lib, archive = _env(tmp_path) + location, ids = _archived(sf, config, lib, archive, seeds=(1,)) + + with TestClient(create_app(config)) as client: + report = client.post( + "/api/v1/restore-preflight", json={"location_id": location["id"]} + ).json() + assert report["state"] == "ready" + + stale = client.post( + "/api/v1/restore-plans", + json={"location_id": location["id"], "token": "r1:not-the-token"}, + ) + assert stale.status_code == 409 + + created = client.post( + "/api/v1/restore-plans", + json={"location_id": location["id"], "token": report["token"]}, + ) + assert created.status_code == 201 + plan_id = created.json()["id"] + assert created.json()["direction"] == "restore" + + # The plan is visible and applying it queues work on the archiver lane. + assert client.get(f"/api/v1/restore-plans/{plan_id}").status_code == 200 + queued = client.post(f"/api/v1/restore-plans/{plan_id}/apply") + assert queued.status_code == 200 + assert queued.json()["job"]["job_type"] == "restore_plan" + assert queued.json()["job"]["lock_key"] == "archive" + assert client.get("/api/v1/restore-recovery").json()["manual"] == [] + + assert _assets(sf)[next(iter(ids.values()))].availability_state == ( + availability.ARCHIVED_ONLINE # the worker, not the request, does the work + ) diff --git a/tests/story_traceability.json b/tests/story_traceability.json index 144f4da..fec5eb9 100644 --- a/tests/story_traceability.json +++ b/tests/story_traceability.json @@ -134,6 +134,9 @@ ], "US06-03": [ "tests/integration/test_offline_assets.py" + ], + "US06-04": [ + "tests/integration/test_restore.py" ] } }