US06-04: Plan and Execute Safe Restores #80

Merged
domverse merged 1 commits from us/US06-04-plan-and-execute-safe-restores into main 2026-08-16 19:49:22 +02:00
10 changed files with 1294 additions and 43 deletions

View File

@@ -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")

View File

@@ -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."""

View File

@@ -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)

View File

@@ -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.

View File

@@ -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"))

View File

@@ -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,

View File

@@ -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,

View File

@@ -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}"

View File

@@ -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
)

View File

@@ -134,6 +134,9 @@
],
"US06-03": [
"tests/integration/test_offline_assets.py"
],
"US06-04": [
"tests/integration/test_restore.py"
]
}
}