Compare commits
3 Commits
us/US06-02
...
us/US06-04
| Author | SHA1 | Date | |
|---|---|---|---|
| 0982ebda70 | |||
| 439eb9971e | |||
| b90d1718be |
29
migrations/versions/0013_protected_thumbnails.py
Normal file
29
migrations/versions/0013_protected_thumbnails.py
Normal file
@@ -0,0 +1,29 @@
|
||||
"""Protected thumbnails (US06-03).
|
||||
|
||||
Revision ID: 0013_protected_thumbnails
|
||||
Revises: 0012_archive_plans
|
||||
Create Date: 2026-08-16
|
||||
|
||||
A protected thumbnail is the durable comparison preview of an asset whose
|
||||
original has left active storage. It is evidence rather than cache, so the LRU
|
||||
quota must not evict it: the archive medium may be offline when it is needed.
|
||||
"""
|
||||
|
||||
import sqlalchemy as sa
|
||||
from alembic import op
|
||||
|
||||
revision = "0013_protected_thumbnails"
|
||||
down_revision = "0012_archive_plans"
|
||||
branch_labels = None
|
||||
depends_on = None
|
||||
|
||||
|
||||
def upgrade() -> None:
|
||||
op.add_column(
|
||||
"thumbnails",
|
||||
sa.Column("protected", sa.Boolean(), nullable=False, server_default=sa.false()),
|
||||
)
|
||||
|
||||
|
||||
def downgrade() -> None:
|
||||
op.drop_column("thumbnails", "protected")
|
||||
37
migrations/versions/0014_restore_plans.py
Normal file
37
migrations/versions/0014_restore_plans.py
Normal 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")
|
||||
@@ -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."""
|
||||
|
||||
@@ -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)
|
||||
|
||||
@@ -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.
|
||||
|
||||
@@ -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"))
|
||||
|
||||
|
||||
@@ -11,7 +11,7 @@ from __future__ import annotations
|
||||
|
||||
from datetime import datetime
|
||||
|
||||
from sqlalchemy import DateTime, ForeignKey, Integer, String, func
|
||||
from sqlalchemy import Boolean, DateTime, ForeignKey, Integer, String, func
|
||||
from sqlalchemy.orm import Mapped, mapped_column
|
||||
|
||||
from photo_pipeline.db import Base
|
||||
@@ -29,6 +29,9 @@ class Thumbnail(Base):
|
||||
width: Mapped[int | None] = mapped_column(Integer)
|
||||
height: Mapped[int | None] = mapped_column(Integer)
|
||||
format: Mapped[str | None] = mapped_column(String)
|
||||
# Durable comparison evidence for an archived asset: never evicted by the LRU
|
||||
# quota, because the original may be on a medium that is no longer reachable.
|
||||
protected: Mapped[bool] = mapped_column(Boolean, nullable=False, default=False)
|
||||
created_at: Mapped[datetime] = mapped_column(
|
||||
DateTime(timezone=True), nullable=False, server_default=func.now()
|
||||
)
|
||||
|
||||
@@ -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:
|
||||
rows = session.scalars(
|
||||
select(ArchiveOperation)
|
||||
.where(
|
||||
stmt = select(ArchiveOperation).where(
|
||||
ArchiveOperation.journal_state.not_in([*TERMINAL_STATES, ArchiveState.PLANNED])
|
||||
)
|
||||
.order_by(ArchiveOperation.plan_id, ArchiveOperation.sequence)
|
||||
if direction is not None:
|
||||
stmt = stmt.where(ArchiveOperation.direction == direction)
|
||||
rows = session.scalars(
|
||||
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,
|
||||
|
||||
@@ -54,14 +54,17 @@ 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,
|
||||
ArchiveState,
|
||||
)
|
||||
from photo_pipeline.services.archives import MARKER_NAME, ArchiveError, ArchiveService
|
||||
from photo_pipeline.services.duplicates import DuplicateService
|
||||
from photo_pipeline.services.hashing import sha256_file
|
||||
from photo_pipeline.services.rename_apply import PreconditionFailed, maybe_fault
|
||||
from photo_pipeline.services.thumbnails import ThumbnailService
|
||||
|
||||
# The per-medium manifest: one JSON line per archived file, appended and fsynced
|
||||
# before its source is removed. It lives with the bytes so the archive can still be
|
||||
@@ -113,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"],
|
||||
@@ -130,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"],
|
||||
@@ -159,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 ─────────────────────────────────────────────────────────────────
|
||||
@@ -260,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.
|
||||
@@ -280,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."""
|
||||
@@ -335,6 +323,7 @@ class ArchiveTransferService:
|
||||
raise PreconditionFailed(
|
||||
"archive_unverified", f"{destination} is not a verified archive copy"
|
||||
)
|
||||
self._require_evidence(operation["asset_id"], destination)
|
||||
if source.exists():
|
||||
if source.is_symlink():
|
||||
raise PreconditionFailed("symlink", f"{source} became a symlink")
|
||||
@@ -377,6 +366,28 @@ class ArchiveTransferService:
|
||||
"asset_moved", f"asset {operation['asset_id']} is no longer at {source}"
|
||||
)
|
||||
|
||||
def _require_evidence(self, asset_id: str, source: Path) -> dict:
|
||||
"""Review evidence must be durable before the original goes.
|
||||
|
||||
The perceptual hash keeps the asset in the fuzzy index once its bytes are
|
||||
unreachable, and the protected preview is what duplicate review can still
|
||||
look at. Both are read from the freshly verified archive copy, which holds
|
||||
exactly the bytes being archived. A file that cannot be decoded has neither
|
||||
— recorded, not fatal, since its exact hashes remain — but failing to
|
||||
produce a preview from a decodable original stops the removal (concept §9).
|
||||
"""
|
||||
DuplicateService(self._session_factory).ensure_phash(asset_id, source=source)
|
||||
preview = ThumbnailService(self._session_factory, self._config).ensure_protected(
|
||||
asset_id, source=source
|
||||
)
|
||||
if preview["state"] == "unavailable":
|
||||
raise PreconditionFailed(
|
||||
"preview_unavailable",
|
||||
f"no durable comparison preview for asset {asset_id} "
|
||||
f"({preview['error_code']})",
|
||||
)
|
||||
return preview
|
||||
|
||||
# ── database ──────────────────────────────────────────────────────────────
|
||||
|
||||
def _record_archived(self, operation: dict, location: dict, destination: Path) -> None:
|
||||
@@ -432,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
|
||||
@@ -459,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],
|
||||
@@ -503,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:
|
||||
@@ -593,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,
|
||||
|
||||
@@ -28,7 +28,11 @@ Blocker codes: ``no_library_root``, ``location_offline``, ``wrong_volume``,
|
||||
``unsafe_destination``, ``destination_not_writable``, ``manifest_unwritable``,
|
||||
``insufficient_capacity``, ``backup_unavailable``, ``lock_conflict``,
|
||||
``rename_pending``, ``empty_scope``, ``destination_collision``,
|
||||
``upload_unverified``, ``bytes_changed``, ``file_missing``.
|
||||
``upload_unverified``, ``bytes_changed``, ``file_missing``, ``preview_unavailable``.
|
||||
|
||||
Preflight also *creates* the durable comparison preview of every asset in scope
|
||||
(US06-03): it is the evidence duplicate review falls back on once the original is
|
||||
on a medium that may be offline, so it has to exist before the original leaves.
|
||||
|
||||
Like upload preflight, the confirmation token is *derived* from the report rather
|
||||
than stored: any change to the scope, the bytes, the destination, or the blockers
|
||||
@@ -59,14 +63,16 @@ from photo_pipeline.models import ArchiveLocation, Asset, UploadBatch, UploadIte
|
||||
from photo_pipeline.path_policy import PathPolicyError, is_excluded, normalize_root, resolve_within
|
||||
from photo_pipeline.services.albums import album_label
|
||||
from photo_pipeline.services.archive_journal import ArchiveJournal
|
||||
from photo_pipeline.services.availability import MARKER_NAME, read_marker as _read_marker
|
||||
from photo_pipeline.services.availability import refresh as refresh_availability
|
||||
from photo_pipeline.services.hashing import sha256_file
|
||||
from photo_pipeline.services.jobs import JobService
|
||||
from photo_pipeline.services.rename_journal import RenameJournal
|
||||
from photo_pipeline.services.thumbnails import ThumbnailService
|
||||
from photo_pipeline.services.upload_reports import VERIFIED
|
||||
|
||||
PREFLIGHT_VERSION = 1
|
||||
TOKEN_PREFIX = f"v{PREFLIGHT_VERSION}"
|
||||
MARKER_NAME = ".photo-pipeline-archive.json"
|
||||
MANIFEST_NAME = "archive-manifest.json"
|
||||
|
||||
# Upload outcomes that prove Immich holds these exact bytes. ``skipped``/``failed``/
|
||||
@@ -163,6 +169,9 @@ class ArchiveService:
|
||||
location.capabilities = json.dumps(probe["capabilities"])
|
||||
reports.append(self._location_report(location, probe=probe))
|
||||
session.commit()
|
||||
# A medium that just appeared or vanished changes what is readable, so the
|
||||
# archived assets are re-derived from the same probe (US06-03).
|
||||
refresh_availability(self._session_factory)
|
||||
return reports
|
||||
|
||||
# ── preflight ─────────────────────────────────────────────────────────────
|
||||
@@ -420,7 +429,10 @@ class ArchiveService:
|
||||
|
||||
def _album(self, name: str, rows: list[dict], root: Path, *, reachable: bool) -> dict:
|
||||
folder = Path(rows[0]["path"]).parent
|
||||
items = sorted((_item(row) for row in rows), key=lambda item: item["current_path"])
|
||||
items = sorted(
|
||||
(self._with_preview(_item(row)) for row in rows),
|
||||
key=lambda item: item["current_path"],
|
||||
)
|
||||
blocked = [item for item in items if item["blockers"]]
|
||||
blockers: list[dict] = []
|
||||
|
||||
@@ -456,6 +468,30 @@ class ArchiveService:
|
||||
"assets": items,
|
||||
}
|
||||
|
||||
def _with_preview(self, item: dict) -> dict:
|
||||
"""Create the durable comparison preview while the original is still here.
|
||||
|
||||
This is the last moment it can be made: once the file is archived and the
|
||||
medium leaves, only the retained preview can answer "is this new photo the
|
||||
same picture?". An original that cannot be decoded at all has no preview to
|
||||
keep — its hashes and metadata stay the evidence — but a preview that fails
|
||||
for any other reason blocks the archive (concept §9).
|
||||
"""
|
||||
preview = self._previews().ensure_protected(item["asset_id"])
|
||||
item["preview"] = preview
|
||||
if preview["state"] == "unavailable" and not item["blockers"]:
|
||||
item["blockers"].append(
|
||||
_issue(
|
||||
"preview_unavailable",
|
||||
f"a durable comparison preview of {item['current_path']} could not be "
|
||||
f"created ({preview['error_code']})",
|
||||
)
|
||||
)
|
||||
return item
|
||||
|
||||
def _previews(self) -> ThumbnailService:
|
||||
return ThumbnailService(self._session_factory, self._config)
|
||||
|
||||
|
||||
# ── internals ────────────────────────────────────────────────────────────────
|
||||
|
||||
@@ -512,13 +548,6 @@ def _transfer_method(folder: Path, root: Path) -> str:
|
||||
return "copy_verify_remove"
|
||||
|
||||
|
||||
def _read_marker(root: Path) -> dict | None:
|
||||
try:
|
||||
return json.loads((root / MARKER_NAME).read_text(encoding="utf-8"))
|
||||
except (OSError, ValueError):
|
||||
return None
|
||||
|
||||
|
||||
def _probe_write(path: Path, payload: bytes, *, keep: bool = False) -> str | None:
|
||||
"""Write ``payload`` to ``path``; return the failure detail or ``None``."""
|
||||
try:
|
||||
|
||||
124
photo_pipeline/services/availability.py
Normal file
124
photo_pipeline/services/availability.py
Normal file
@@ -0,0 +1,124 @@
|
||||
"""Where an asset's bytes are right now (US06-03).
|
||||
|
||||
Archiving removes the original from the active library but never removes the
|
||||
asset: its identity, hashes, decisions, and evidence stay. This module is the one
|
||||
place that answers "can these bytes be read, and if not, why" so inventory,
|
||||
duplicate review, thumbnails, and the archive service all give the same answer.
|
||||
|
||||
States (concept §9):
|
||||
|
||||
- ``active`` — the original is in the active library;
|
||||
- ``archived_online`` — the recorded medium is mounted and holds the file;
|
||||
- ``archived_offline`` — archived, but the medium is not available right now;
|
||||
- ``missing_unexpected`` — neither an active path nor the recorded archive
|
||||
location explains the absence. This is the state that must never be confused
|
||||
with ``archived_offline``: an unmounted disk is normal, a mounted disk with a
|
||||
hole in it is not.
|
||||
|
||||
A medium is identified by its marker file, never by its mountpoint, so a
|
||||
different disk mounted at the recorded root is offline rather than accepted.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import json
|
||||
from collections import Counter
|
||||
from datetime import datetime, timezone
|
||||
from pathlib import Path
|
||||
|
||||
from sqlalchemy import select
|
||||
from sqlalchemy.orm import Session, sessionmaker
|
||||
|
||||
from photo_pipeline.models import ArchiveLocation, Asset
|
||||
|
||||
ACTIVE = "active"
|
||||
ARCHIVED_ONLINE = "archived_online"
|
||||
ARCHIVED_OFFLINE = "archived_offline"
|
||||
MISSING_UNEXPECTED = "missing_unexpected"
|
||||
ARCHIVED = (ARCHIVED_ONLINE, ARCHIVED_OFFLINE)
|
||||
|
||||
MARKER_NAME = ".photo-pipeline-archive.json"
|
||||
|
||||
|
||||
def read_marker(root: Path) -> dict | None:
|
||||
"""The medium's identity marker, or ``None`` when it is not readable."""
|
||||
try:
|
||||
return json.loads((root / MARKER_NAME).read_text(encoding="utf-8"))
|
||||
except (OSError, ValueError):
|
||||
return None
|
||||
|
||||
|
||||
def location_online(location: ArchiveLocation) -> bool:
|
||||
"""True only when the *recorded* medium is mounted at its root."""
|
||||
marker = read_marker(Path(location.root))
|
||||
return bool(marker) and marker.get("media_id") == location.media_id
|
||||
|
||||
|
||||
def archive_file(session: Session, asset: Asset) -> Path | None:
|
||||
"""The archived file's absolute path, whether or not the medium is mounted."""
|
||||
if not asset.archive_location_id or not asset.archive_path:
|
||||
return None
|
||||
location = session.get(ArchiveLocation, asset.archive_location_id)
|
||||
if location is None:
|
||||
return None
|
||||
return Path(location.root) / asset.archive_path
|
||||
|
||||
|
||||
def readable_path(session: Session, asset: Asset) -> Path | None:
|
||||
"""A path whose bytes can be read now: the active file, else the archive copy."""
|
||||
if asset.current_path and Path(asset.current_path).exists():
|
||||
return Path(asset.current_path)
|
||||
archived = archive_file(session, asset)
|
||||
if archived is None:
|
||||
return None
|
||||
location = session.get(ArchiveLocation, asset.archive_location_id)
|
||||
if not location_online(location) or not archived.exists():
|
||||
return None
|
||||
return archived
|
||||
|
||||
|
||||
def state_of(session: Session, asset: Asset, *, online: dict[str, bool] | None = None) -> str:
|
||||
"""The availability this asset's storage actually justifies right now."""
|
||||
if asset.current_path:
|
||||
return ACTIVE if Path(asset.current_path).exists() else MISSING_UNEXPECTED
|
||||
if not asset.archive_location_id:
|
||||
return MISSING_UNEXPECTED if asset.availability_state != ACTIVE else ACTIVE
|
||||
location = session.get(ArchiveLocation, asset.archive_location_id)
|
||||
if location is None:
|
||||
return MISSING_UNEXPECTED
|
||||
reachable = (
|
||||
online[location.id] if online and location.id in online else location_online(location)
|
||||
)
|
||||
if not reachable:
|
||||
return ARCHIVED_OFFLINE
|
||||
archived = archive_file(session, asset)
|
||||
# The medium is mounted and identified: the file is either there, or it is
|
||||
# genuinely gone — that is not "offline", it needs a human.
|
||||
return ARCHIVED_ONLINE if archived and archived.exists() else MISSING_UNEXPECTED
|
||||
|
||||
|
||||
def refresh(session_factory: sessionmaker) -> dict[str, int]:
|
||||
"""Re-derive availability for every archived asset from the media themselves.
|
||||
|
||||
Only archived assets are probed: whether an *active* file is present is the
|
||||
inventory scan's job and costs one stat per library file. Each medium is
|
||||
probed once, not once per asset.
|
||||
"""
|
||||
counts: Counter[str] = Counter()
|
||||
now = datetime.now(timezone.utc)
|
||||
with session_factory() as session:
|
||||
online = {
|
||||
location.id: location_online(location)
|
||||
for location in session.scalars(select(ArchiveLocation))
|
||||
}
|
||||
for asset in session.scalars(
|
||||
select(Asset).where(Asset.archive_location_id.is_not(None))
|
||||
):
|
||||
state = state_of(session, asset, online=online)
|
||||
counts[state] += 1
|
||||
if state != asset.availability_state:
|
||||
asset.availability_state = state
|
||||
asset.state_version += 1
|
||||
asset.updated_at = now
|
||||
session.commit()
|
||||
return dict(counts)
|
||||
@@ -11,6 +11,12 @@ Detection runs in two categories:
|
||||
band (NEAR/SIMILAR). These are review candidates: never decided automatically, and
|
||||
negative-linked pairs are suppressed so a rejected pair is not re-suggested.
|
||||
|
||||
Archived assets stay in both indexes (US06-03): a new active copy of an archived
|
||||
original is recognised through its hashes even while the medium is offline, and
|
||||
cluster review falls back to the retained protected preview plus hash evidence.
|
||||
An exact/pixel match links straight to the archived canonical; a perceptual match
|
||||
is a review candidate that names the medium to mount for a pixel-level decision.
|
||||
|
||||
Decisions (``canonical`` / ``not_duplicate`` / ``deferred``) persist with evidence,
|
||||
use optimistic version checks, are reversible, and can never form a canonical cycle.
|
||||
A new content-identical member of an already-decided cluster inherits the established
|
||||
@@ -34,12 +40,14 @@ from sqlalchemy import func, select
|
||||
from sqlalchemy.orm import sessionmaker
|
||||
|
||||
from photo_pipeline.models import (
|
||||
ArchiveLocation,
|
||||
Asset,
|
||||
DuplicateCluster,
|
||||
DuplicateMember,
|
||||
DuplicateNegativeLink,
|
||||
Thumbnail,
|
||||
)
|
||||
from photo_pipeline.services import hashing
|
||||
from photo_pipeline.services import availability, hashing
|
||||
|
||||
NEAR_MAX = 5
|
||||
SIMILAR_MAX = 10
|
||||
@@ -127,18 +135,18 @@ class DuplicateService:
|
||||
|
||||
# ── perceptual hash backfill ───────────────────────────────────────────
|
||||
def ensure_phashes(self) -> int:
|
||||
"""Hash whatever is readable now — an archived asset keeps the hash it
|
||||
already has, and gains one whenever its medium happens to be mounted."""
|
||||
updated = 0
|
||||
with self._session_factory() as session:
|
||||
assets = session.execute(
|
||||
select(Asset).where(
|
||||
Asset.availability_state == "active",
|
||||
Asset.current_path.isnot(None),
|
||||
)
|
||||
).scalars()
|
||||
assets = session.execute(select(Asset)).scalars()
|
||||
for asset in assets:
|
||||
if asset.phash is not None and asset.phash_version == hashing.PHASH_VERSION:
|
||||
continue
|
||||
value = hashing.safe_phash(asset.current_path)
|
||||
source = availability.readable_path(session, asset)
|
||||
if source is None:
|
||||
continue
|
||||
value = hashing.safe_phash(str(source))
|
||||
if value is not None:
|
||||
asset.phash = value
|
||||
asset.phash_version = hashing.PHASH_VERSION
|
||||
@@ -146,20 +154,40 @@ class DuplicateService:
|
||||
session.commit()
|
||||
return updated
|
||||
|
||||
def ensure_phash(self, asset_id: str, *, source=None) -> str | None:
|
||||
"""Backfill one asset's perceptual hash while its bytes are still readable.
|
||||
|
||||
Archiving calls this before the original leaves — passing the archive copy
|
||||
as ``source``, since the database does not point at it yet — because an
|
||||
asset without a pHash silently drops out of the fuzzy index the moment its
|
||||
medium is away.
|
||||
"""
|
||||
with self._session_factory() as session:
|
||||
asset = session.get(Asset, asset_id)
|
||||
if asset is None:
|
||||
return None
|
||||
if asset.phash is not None and asset.phash_version == hashing.PHASH_VERSION:
|
||||
return asset.phash
|
||||
source = source or availability.readable_path(session, asset)
|
||||
if source is None:
|
||||
return None
|
||||
value = hashing.safe_phash(str(source))
|
||||
if value is not None:
|
||||
asset.phash = value
|
||||
asset.phash_version = hashing.PHASH_VERSION
|
||||
session.commit()
|
||||
return value
|
||||
|
||||
# ── detection ──────────────────────────────────────────────────────────
|
||||
def detect(self) -> DetectionReport:
|
||||
self.ensure_phashes()
|
||||
now = datetime.now(timezone.utc)
|
||||
report = DetectionReport()
|
||||
with self._session_factory() as session:
|
||||
assets = list(
|
||||
session.execute(
|
||||
select(Asset).where(
|
||||
Asset.availability_state == "active",
|
||||
Asset.current_path.isnot(None),
|
||||
)
|
||||
).scalars()
|
||||
)
|
||||
# Every known asset stays in the indexes, archived or not: a copy of an
|
||||
# archived original must be recognised as a duplicate rather than
|
||||
# treated as a new photo (concept §9, invariant 12).
|
||||
assets = list(session.execute(select(Asset)).scalars())
|
||||
by_id = {a.id: a for a in assets}
|
||||
negatives = {
|
||||
_pair(link.asset_a, link.asset_b)
|
||||
@@ -431,10 +459,19 @@ class DuplicateService:
|
||||
|
||||
@staticmethod
|
||||
def _recommend_canonical(ids, by_id) -> str:
|
||||
# ponytail: largest file, path as deterministic tie-break. The concept's
|
||||
# richer policy (resolution, least recompression, metadata richness) lands
|
||||
# with the review UI story.
|
||||
return max(ids, key=lambda i: (by_id[i].byte_size or 0, by_id[i].current_path or ""))
|
||||
# ponytail: largest file, then the archived copy, then path as a
|
||||
# deterministic tie-break. Archived wins ties because it is the reviewed,
|
||||
# uploaded original — a fresh active copy must not demote it to a variant.
|
||||
# The concept's richer policy (resolution, least recompression, metadata
|
||||
# richness) lands with the review UI story.
|
||||
return max(
|
||||
ids,
|
||||
key=lambda i: (
|
||||
by_id[i].byte_size or 0,
|
||||
by_id[i].availability_state in availability.ARCHIVED,
|
||||
by_id[i].current_path or by_id[i].archive_path or "",
|
||||
),
|
||||
)
|
||||
|
||||
def _apply_canonical(self, session, cluster, ids, canonical_id):
|
||||
for member in session.execute(
|
||||
@@ -507,9 +544,15 @@ class DuplicateService:
|
||||
"current_path": asset.current_path if asset else None,
|
||||
"byte_size": asset.byte_size if asset else None,
|
||||
"phash": asset.phash if asset else None,
|
||||
**self._offline_evidence(session, asset),
|
||||
}
|
||||
)
|
||||
members.sort(key=lambda m: m["asset_id"])
|
||||
# A full-resolution comparison of an offline original is impossible; the
|
||||
# UI asks for that named medium instead of guessing (concept §9).
|
||||
mount_required = sorted(
|
||||
{m["archive_location"] for m in members if m["requires_mount"]}
|
||||
)
|
||||
return {
|
||||
"id": cluster.id,
|
||||
"method": cluster.method,
|
||||
@@ -519,9 +562,54 @@ class DuplicateService:
|
||||
"canonical_asset_id": cluster.canonical_asset_id,
|
||||
"version": cluster.version,
|
||||
"requires_confirmation": cluster.method == Method.PERCEPTUAL.value,
|
||||
"mount_required": mount_required,
|
||||
"members": members,
|
||||
}
|
||||
|
||||
def _offline_evidence(self, session, asset: Asset | None) -> dict:
|
||||
"""What review can still rely on when a member's original is not readable."""
|
||||
if asset is None:
|
||||
return {
|
||||
"availability_state": None,
|
||||
"archive_location": None,
|
||||
"archive_location_id": None,
|
||||
"archive_path": None,
|
||||
"preview": {"state": "missing", "protected": False},
|
||||
"requires_mount": False,
|
||||
}
|
||||
location = (
|
||||
session.get(ArchiveLocation, asset.archive_location_id)
|
||||
if asset.archive_location_id
|
||||
else None
|
||||
)
|
||||
preview = self._preview_evidence(session, asset.id)
|
||||
archived = asset.availability_state in availability.ARCHIVED
|
||||
return {
|
||||
"availability_state": asset.availability_state,
|
||||
"archive_location": location.name if location else None,
|
||||
"archive_location_id": asset.archive_location_id,
|
||||
"archive_path": asset.archive_path,
|
||||
"preview": preview,
|
||||
# Offline archived members can still be compared through their retained
|
||||
# preview and hash evidence; only pixel-level review needs the medium.
|
||||
"requires_mount": archived
|
||||
and asset.availability_state == availability.ARCHIVED_OFFLINE
|
||||
and bool(location),
|
||||
}
|
||||
|
||||
@staticmethod
|
||||
def _preview_evidence(session, asset_id: str) -> dict:
|
||||
rows = list(
|
||||
session.execute(select(Thumbnail).where(Thumbnail.asset_id == asset_id)).scalars()
|
||||
)
|
||||
ready = [r for r in rows if r.state == "ready" and r.path]
|
||||
if ready:
|
||||
best = max(ready, key=lambda r: (bool(r.protected), r.size or 0))
|
||||
return {"state": "ready", "protected": bool(best.protected), "size": best.size}
|
||||
if rows:
|
||||
return {"state": "unsupported", "protected": False, "size": rows[0].size}
|
||||
return {"state": "missing", "protected": False, "size": None}
|
||||
|
||||
# ── decisions ────────────────────────────────────────────────────────────
|
||||
def decide(
|
||||
self,
|
||||
|
||||
@@ -13,8 +13,11 @@ renames. Every discovered or absent path is classified as one occurrence:
|
||||
- ``missing`` — a known active asset whose file is gone (kept, flagged).
|
||||
|
||||
Missing files are never pruned (that would break identity); the asset is retained
|
||||
with ``missing_at`` set. Archived assets are left untouched. Rescanning unchanged
|
||||
input makes no durable change.
|
||||
with ``missing_at`` set and its availability becomes ``missing_unexpected`` —
|
||||
nothing explains where the bytes went. Archived assets are left untouched: their
|
||||
absence from the active roots is expected, and each scan re-derives whether their
|
||||
medium is reachable (:mod:`photo_pipeline.services.availability`). Rescanning
|
||||
unchanged input makes no durable change.
|
||||
|
||||
Extracted from photo_analyzer.discover_photos/reconcile_moved/prune_missing
|
||||
(see donor_ledger.yaml: pa-discovery, pa-prune-missing).
|
||||
@@ -30,12 +33,12 @@ from enum import Enum
|
||||
from pathlib import Path
|
||||
from typing import Iterable
|
||||
|
||||
from sqlalchemy import func, select
|
||||
from sqlalchemy import func, or_, select
|
||||
from sqlalchemy.orm import Session, sessionmaker
|
||||
|
||||
from photo_pipeline import path_policy
|
||||
from photo_pipeline.models import Asset, AssetPath
|
||||
from photo_pipeline.services import hashing
|
||||
from photo_pipeline.services import availability, hashing
|
||||
|
||||
|
||||
class Occurrence(str, Enum):
|
||||
@@ -60,6 +63,8 @@ def _asset_dict(asset: Asset) -> dict:
|
||||
"id": asset.id,
|
||||
"current_path": asset.current_path,
|
||||
"availability_state": asset.availability_state,
|
||||
"archive_location_id": asset.archive_location_id,
|
||||
"archive_path": asset.archive_path,
|
||||
"byte_size": asset.byte_size,
|
||||
"current_sha256": asset.current_sha256,
|
||||
"pixel_sha256": asset.pixel_sha256,
|
||||
@@ -101,7 +106,9 @@ class InventoryService:
|
||||
result.asset_ids[str(path)] = asset.id
|
||||
|
||||
for asset in assets:
|
||||
if asset.availability_state != "active" or asset.id in seen_ids:
|
||||
# Archived assets are explained by their location, not by the active
|
||||
# roots: a scan must never prune or flag them (concept §9).
|
||||
if asset.availability_state in availability.ARCHIVED or asset.id in seen_ids:
|
||||
continue
|
||||
if asset.current_path and asset.current_path not in discovered_paths:
|
||||
if not Path(asset.current_path).exists():
|
||||
@@ -109,10 +116,18 @@ class InventoryService:
|
||||
asset.missing_at = now
|
||||
asset.state_version += 1
|
||||
asset.updated_at = now
|
||||
# Nothing explains this absence — it is not an offline medium.
|
||||
if asset.availability_state != availability.MISSING_UNEXPECTED:
|
||||
asset.availability_state = availability.MISSING_UNEXPECTED
|
||||
asset.state_version += 1
|
||||
asset.updated_at = now
|
||||
result.occurrences[asset.current_path] = Occurrence.MISSING.value
|
||||
|
||||
session.commit()
|
||||
|
||||
# Media may have been mounted or removed since the last scan.
|
||||
availability.refresh(self._session_factory)
|
||||
|
||||
result.counts = dict(Counter(result.occurrences.values()))
|
||||
return result
|
||||
|
||||
@@ -132,7 +147,11 @@ class InventoryService:
|
||||
if availability:
|
||||
stmt = stmt.where(Asset.availability_state == availability)
|
||||
if query:
|
||||
stmt = stmt.where(Asset.current_path.like(f"%{query}%"))
|
||||
like = f"%{query}%"
|
||||
# An archived asset has no active path; it is searched where it lives.
|
||||
stmt = stmt.where(
|
||||
or_(Asset.current_path.like(like), Asset.archive_path.like(like))
|
||||
)
|
||||
total = session.scalar(select(func.count()).select_from(stmt.subquery()))
|
||||
rows = session.execute(
|
||||
stmt.order_by(Asset.current_path).limit(limit).offset(offset)
|
||||
@@ -169,6 +188,7 @@ class InventoryService:
|
||||
self._open_path(session, existing.id, path_str, now, occ.value)
|
||||
if existing.missing_at is not None:
|
||||
existing.missing_at = None
|
||||
existing.availability_state = availability.ACTIVE
|
||||
existing.state_version += 1
|
||||
existing.updated_at = now
|
||||
return existing, occ
|
||||
@@ -189,6 +209,7 @@ class InventoryService:
|
||||
moved_from.current_path = path_str
|
||||
moved_from.byte_size = size
|
||||
moved_from.missing_at = None
|
||||
moved_from.availability_state = availability.ACTIVE
|
||||
moved_from.state_version += 1
|
||||
moved_from.updated_at = now
|
||||
self._open_path(session, moved_from.id, path_str, now, Occurrence.MOVED.value)
|
||||
|
||||
648
photo_pipeline/services/restores.py
Normal file
648
photo_pipeline/services/restores.py
Normal 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}"
|
||||
@@ -8,6 +8,11 @@ an EXIF-only edit reuses the file while a real pixel change invalidates it; writ
|
||||
are atomic and the cache is bounded by an LRU quota. Failures are persisted as typed
|
||||
errors so a broken original is not retried on every request.
|
||||
|
||||
An archived asset is served from its medium when that medium is mounted, and from
|
||||
its *protected* preview when it is not (US06-03). Protected previews are evidence,
|
||||
not cache: the quota never evicts them, because the original they describe may be
|
||||
unreachable when duplicate review needs it.
|
||||
|
||||
Reuses photo_analyzer.prepare_image decode/resize/HEIC handling, adding the missing
|
||||
EXIF-orientation step, WebP output, and a managed cache (donor_ledger.yaml:
|
||||
pa-imaging).
|
||||
@@ -26,6 +31,7 @@ from sqlalchemy.orm import sessionmaker
|
||||
from photo_pipeline import path_policy
|
||||
from photo_pipeline.config import Config
|
||||
from photo_pipeline.models import Asset, Thumbnail
|
||||
from photo_pipeline.services import availability
|
||||
|
||||
# Best-effort HEIC support: registered only if the optional decoder is installed.
|
||||
try: # pragma: no cover - depends on an optional native dependency
|
||||
@@ -38,6 +44,8 @@ except Exception: # pragma: no cover
|
||||
SIZES = (256, 512, 1280)
|
||||
THUMB_VERSION = 1
|
||||
THUMB_FORMAT = "webp"
|
||||
# The size kept as durable comparison evidence for archived assets (concept §9).
|
||||
PROTECTED_SIZE = 1280
|
||||
|
||||
|
||||
class ThumbnailError(RuntimeError):
|
||||
@@ -86,7 +94,12 @@ class ThumbnailService:
|
||||
self._config = config
|
||||
self._cache_dir = config.thumbnail_cache_dir
|
||||
|
||||
def generate(self, asset_id: str, size: int) -> Path:
|
||||
def generate(
|
||||
self, asset_id: str, size: int, *, protected: bool = False, source: Path | None = None
|
||||
) -> Path:
|
||||
"""Render (or reuse) a preview. ``source`` overrides where the bytes are read
|
||||
from — the archiver passes its verified archive copy, which the database does
|
||||
not yet point at while the transfer is still in flight."""
|
||||
if size not in SIZES:
|
||||
raise InvalidSize(f"size must be one of {SIZES}")
|
||||
|
||||
@@ -94,9 +107,7 @@ class ThumbnailService:
|
||||
asset = session.get(Asset, asset_id)
|
||||
if asset is None:
|
||||
raise ThumbnailNotFound(f"unknown asset {asset_id}")
|
||||
if asset.availability_state != "active" or not asset.current_path:
|
||||
raise ThumbnailUnavailable(f"asset {asset_id} has no active file")
|
||||
self._validate_path(asset.current_path)
|
||||
archived = asset.availability_state in availability.ARCHIVED
|
||||
cache_key = self._cache_key(asset, size)
|
||||
|
||||
row = session.get(Thumbnail, cache_key)
|
||||
@@ -107,9 +118,18 @@ class ThumbnailService:
|
||||
)
|
||||
if row.path and Path(row.path).exists():
|
||||
_touch(row.path)
|
||||
if protected and not row.protected:
|
||||
self._protect(cache_key)
|
||||
return Path(row.path)
|
||||
|
||||
source = asset.current_path
|
||||
# An archived original is read from its medium; when that medium is not
|
||||
# mounted the retained preview above is the only evidence there is.
|
||||
source = source or availability.readable_path(session, asset)
|
||||
if source is None:
|
||||
raise ThumbnailUnavailable(f"asset {asset_id} has no readable file")
|
||||
source = str(source)
|
||||
if source == asset.current_path:
|
||||
self._validate_path(source) # archive roots lie outside the library
|
||||
|
||||
# Rendering happens outside the DB session (no transaction held during I/O).
|
||||
try:
|
||||
@@ -119,10 +139,61 @@ class ThumbnailService:
|
||||
self._record_error(cache_key, asset_id, size, error.code)
|
||||
raise
|
||||
|
||||
self._record_ready(cache_key, asset_id, size, rendered)
|
||||
# Archived assets keep their preview permanently: it is the comparison
|
||||
# evidence that survives the original leaving active storage.
|
||||
self._record_ready(cache_key, asset_id, size, rendered, protected=protected or archived)
|
||||
self._enforce_quota(keep=rendered["path"])
|
||||
return Path(rendered["path"])
|
||||
|
||||
def ensure_protected(self, asset_id: str, *, source: Path | None = None) -> dict:
|
||||
"""Produce (or confirm) the durable comparison preview for an asset.
|
||||
|
||||
Returns evidence rather than raising, because the caller — archive
|
||||
preflight and the transfer itself — decides what an unrenderable original
|
||||
means. ``unsupported`` is a recorded property of the file, not a failure of
|
||||
the policy: its hashes and metadata remain the comparison evidence.
|
||||
"""
|
||||
try:
|
||||
path = self.generate(asset_id, PROTECTED_SIZE, protected=True, source=source)
|
||||
except tuple(_PERSISTED_ERRORS) as error:
|
||||
return {"state": "unsupported", "error_code": error.code, "path": None}
|
||||
except ThumbnailError as error:
|
||||
return {"state": "unavailable", "error_code": error.code, "path": None}
|
||||
return {"state": "ready", "error_code": None, "path": str(path)}
|
||||
|
||||
def evidence(self, asset_id: str) -> dict:
|
||||
"""What durable preview this asset has right now, without rendering."""
|
||||
with self._session_factory() as session:
|
||||
rows = list(
|
||||
session.execute(
|
||||
select(Thumbnail).where(Thumbnail.asset_id == asset_id)
|
||||
).scalars()
|
||||
)
|
||||
for row in rows:
|
||||
if row.state == "ready" and row.path and Path(row.path).exists():
|
||||
return {
|
||||
"state": "ready",
|
||||
"protected": bool(row.protected),
|
||||
"size": row.size,
|
||||
"error_code": None,
|
||||
}
|
||||
for row in rows:
|
||||
if row.state == "error":
|
||||
return {
|
||||
"state": "unsupported",
|
||||
"protected": False,
|
||||
"size": row.size,
|
||||
"error_code": row.error_code,
|
||||
}
|
||||
return {"state": "missing", "protected": False, "size": None, "error_code": None}
|
||||
|
||||
def _protect(self, cache_key: str) -> None:
|
||||
with self._session_factory() as session:
|
||||
row = session.get(Thumbnail, cache_key)
|
||||
if row is not None:
|
||||
row.protected = True
|
||||
session.commit()
|
||||
|
||||
# ── path safety ──────────────────────────────────────────────────────────
|
||||
def _validate_path(self, current_path: str) -> None:
|
||||
path = Path(current_path)
|
||||
@@ -184,7 +255,9 @@ class ThumbnailService:
|
||||
}
|
||||
|
||||
# ── persistence ────────────────────────────────────────────────────────────
|
||||
def _record_ready(self, cache_key: str, asset_id: str, size: int, rendered: dict) -> None:
|
||||
def _record_ready(
|
||||
self, cache_key: str, asset_id: str, size: int, rendered: dict, *, protected: bool = False
|
||||
) -> None:
|
||||
with self._session_factory() as session:
|
||||
session.merge(
|
||||
Thumbnail(
|
||||
@@ -197,6 +270,7 @@ class ThumbnailService:
|
||||
width=rendered["width"],
|
||||
height=rendered["height"],
|
||||
format=rendered["format"],
|
||||
protected=protected,
|
||||
)
|
||||
)
|
||||
try:
|
||||
@@ -230,6 +304,10 @@ class ThumbnailService:
|
||||
if total <= quota:
|
||||
return
|
||||
files.sort(key=lambda f: f.stat().st_mtime) # least-recently-used first
|
||||
# Protected previews are evidence, not cache: an archived original cannot be
|
||||
# re-rendered once its medium is away, so eviction never touches them.
|
||||
protected = self._protected_paths()
|
||||
files = [f for f in files if str(f) not in protected]
|
||||
keep_path = str(Path(keep)) if keep else None
|
||||
evicted: list[str] = []
|
||||
for f in files:
|
||||
@@ -247,6 +325,16 @@ class ThumbnailService:
|
||||
if evicted:
|
||||
self._forget(evicted)
|
||||
|
||||
def _protected_paths(self) -> set[str]:
|
||||
with self._session_factory() as session:
|
||||
return {
|
||||
row.path
|
||||
for row in session.execute(
|
||||
select(Thumbnail).where(Thumbnail.protected.is_(True))
|
||||
).scalars()
|
||||
if row.path
|
||||
}
|
||||
|
||||
def _forget(self, paths: list[str]) -> None:
|
||||
with self._session_factory() as session:
|
||||
rows = session.execute(
|
||||
|
||||
@@ -167,7 +167,8 @@ def test_missing_file_is_flagged_not_deleted(
|
||||
assets = assets_by_id(make_factory(db_url))
|
||||
assert missing_id in assets # not pruned
|
||||
assert assets[missing_id].missing_at is not None
|
||||
assert assets[missing_id].availability_state == "active"
|
||||
# Nothing explains the absence: this is not an offline archive medium (US06-03).
|
||||
assert assets[missing_id].availability_state == "missing_unexpected"
|
||||
|
||||
|
||||
def test_reappearing_file_clears_missing(
|
||||
@@ -186,6 +187,7 @@ def test_reappearing_file_clears_missing(
|
||||
inventory.scan(lib)
|
||||
assets = assets_by_id(make_factory(db_url))
|
||||
assert assets[asset_id].missing_at is None
|
||||
assert assets[asset_id].availability_state == "active"
|
||||
|
||||
|
||||
def test_identity_and_state_durable_across_restart(
|
||||
|
||||
377
tests/integration/test_offline_assets.py
Normal file
377
tests/integration/test_offline_assets.py
Normal file
@@ -0,0 +1,377 @@
|
||||
"""Offline identity and review evidence (US06-03).
|
||||
|
||||
An archived photo is not gone: it keeps its identity, its hashes, and enough
|
||||
evidence to be recognised in a duplicate cluster while its medium sits in a
|
||||
drawer. Every case here archives a *real* album through the real transfer, then
|
||||
takes the medium away by removing its marker — the same thing the service sees
|
||||
when an external disk is unplugged — and asks whether the application still tells
|
||||
the truth about where the bytes are.
|
||||
|
||||
The distinction that matters throughout: an unmounted medium is
|
||||
``archived_offline`` (expected, harmless), a mounted medium with a hole in it is
|
||||
``missing_unexpected`` (needs a human). Confusing the two is how an archive
|
||||
quietly loses a photo.
|
||||
"""
|
||||
|
||||
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, Thumbnail, UploadBatch, UploadItem
|
||||
from photo_pipeline.services.archive_transfer import ArchiveTransferService
|
||||
from photo_pipeline.services.archives import MARKER_NAME, ArchiveService
|
||||
from photo_pipeline.services.availability import (
|
||||
ACTIVE,
|
||||
ARCHIVED_OFFLINE,
|
||||
ARCHIVED_ONLINE,
|
||||
MISSING_UNEXPECTED,
|
||||
)
|
||||
from photo_pipeline.services.duplicates import ClusterState, DuplicateService, Method
|
||||
from photo_pipeline.services.hashing import sha256_file
|
||||
from photo_pipeline.services.inventory import InventoryService
|
||||
from photo_pipeline.services.thumbnails import PROTECTED_SIZE, ThumbnailService, ThumbnailUnavailable
|
||||
|
||||
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=(256, 192)):
|
||||
"""A deterministic, decodable photo — previews and pHashes must be real."""
|
||||
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(6):
|
||||
x0 = int(rng.integers(0, w - 60))
|
||||
y0 = int(rng.integers(0, h - 60))
|
||||
base[y0 : y0 + 60, x0 : x0 + 60] = rng.integers(0, 256, 3)
|
||||
grad = np.linspace(0, 120, w, dtype=np.uint8)
|
||||
base[:, :, 0] = np.clip(base[:, :, 0].astype(int) + grad[None, :], 0, 255)
|
||||
Image.fromarray(base).save(path, quality=95)
|
||||
return path
|
||||
|
||||
|
||||
def resized_copy(src, dst, scale=0.5):
|
||||
with Image.open(src) as image:
|
||||
image.resize(
|
||||
(int(image.width * scale), int(image.height * scale)), Image.LANCZOS
|
||||
).save(dst, quality=95)
|
||||
return dst
|
||||
|
||||
|
||||
def _uploaded(sf, lib, album="rome", seeds=(1, 2)):
|
||||
"""A scanned album carrying the verified upload evidence archiving requires."""
|
||||
folder = lib / album
|
||||
for index, seed in enumerate(seeds):
|
||||
structured(folder / f"{index}.jpg", seed)
|
||||
result = 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 result.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",
|
||||
)
|
||||
)
|
||||
session.commit()
|
||||
return folder, result.asset_ids
|
||||
|
||||
|
||||
def _archive(sf, config, archive, albums=None):
|
||||
service = ArchiveService(sf, config=config)
|
||||
location = service.register("external", str(archive))
|
||||
token = service.preflight(location["id"], albums)["token"]
|
||||
transfers = ArchiveTransferService(sf, config=config)
|
||||
plan = transfers.create(location["id"], albums, token=token)
|
||||
transfers.apply(plan["id"])
|
||||
return location
|
||||
|
||||
|
||||
def _unmount(archive):
|
||||
"""Take the medium away the way a real one goes: its marker stops answering."""
|
||||
(archive / MARKER_NAME).rename(archive / f"{MARKER_NAME}.away")
|
||||
|
||||
|
||||
def _remount(archive):
|
||||
(archive / f"{MARKER_NAME}.away").rename(archive / MARKER_NAME)
|
||||
|
||||
|
||||
def _assets(sf):
|
||||
with sf() as session:
|
||||
return {asset.id: asset for asset in session.scalars(select(Asset))}
|
||||
|
||||
|
||||
# ── availability ─────────────────────────────────────────────────────────────
|
||||
|
||||
|
||||
def test_archived_assets_report_online_offline_and_missing(tmp_path):
|
||||
config, sf, lib, archive = _env(tmp_path)
|
||||
_, ids = _uploaded(sf, lib)
|
||||
_archive(sf, config, archive)
|
||||
inventory = InventoryService(sf)
|
||||
|
||||
states = {a.availability_state for a in _assets(sf).values()}
|
||||
assert states == {ARCHIVED_ONLINE}
|
||||
|
||||
_unmount(archive)
|
||||
inventory.scan(lib) # the album folder is gone from the active roots
|
||||
assert {a.availability_state for a in _assets(sf).values()} == {ARCHIVED_OFFLINE}
|
||||
assert all(a.missing_at is None for a in _assets(sf).values()) # not "missing"
|
||||
|
||||
_remount(archive)
|
||||
inventory.scan(lib)
|
||||
assert {a.availability_state for a in _assets(sf).values()} == {ARCHIVED_ONLINE}
|
||||
|
||||
# Mounted medium, absent file: that is not an offline archive, it needs a human.
|
||||
victim = sorted(ids.values())[0]
|
||||
with sf() as session:
|
||||
asset = session.get(Asset, victim)
|
||||
(archive / asset.archive_path).unlink()
|
||||
inventory.scan(lib)
|
||||
assert _assets(sf)[victim].availability_state == MISSING_UNEXPECTED
|
||||
|
||||
|
||||
def test_scan_never_prunes_or_flags_offline_assets(tmp_path):
|
||||
config, sf, lib, archive = _env(tmp_path)
|
||||
_, ids = _uploaded(sf, lib)
|
||||
_archive(sf, config, archive)
|
||||
_unmount(archive)
|
||||
|
||||
before = _assets(sf)
|
||||
result = InventoryService(sf).scan(lib)
|
||||
|
||||
assert result.counts.get("missing") is None
|
||||
after = _assets(sf)
|
||||
assert set(after) == set(before) == set(ids.values())
|
||||
for asset in after.values():
|
||||
assert asset.availability_state == ARCHIVED_OFFLINE
|
||||
assert asset.current_sha256 and asset.pixel_sha256 # hashes retained
|
||||
assert asset.archive_location_id and asset.archive_path
|
||||
|
||||
|
||||
def test_active_missing_file_is_missing_unexpected_not_offline(tmp_path):
|
||||
config, sf, lib, _archive_root = _env(tmp_path)
|
||||
structured(lib / "loose" / "a.jpg", 7)
|
||||
ids = InventoryService(sf).scan(lib).asset_ids
|
||||
asset_id = next(iter(ids.values()))
|
||||
|
||||
(lib / "loose" / "a.jpg").unlink()
|
||||
InventoryService(sf).scan(lib)
|
||||
|
||||
asset = _assets(sf)[asset_id]
|
||||
assert asset.availability_state == MISSING_UNEXPECTED
|
||||
assert asset.missing_at is not None
|
||||
|
||||
|
||||
# ── deduplication against archived originals ─────────────────────────────────
|
||||
|
||||
|
||||
def test_exact_copy_of_offline_asset_links_to_archived_canonical(tmp_path):
|
||||
config, sf, lib, archive = _env(tmp_path)
|
||||
folder, ids = _uploaded(sf, lib, seeds=(1,))
|
||||
archived_id = next(iter(ids.values()))
|
||||
original = next(iter(ids))
|
||||
kept = tmp_path / "kept.jpg"
|
||||
shutil.copy2(original, kept)
|
||||
_archive(sf, config, archive)
|
||||
_unmount(archive)
|
||||
|
||||
# The same photo turns up again in the active library while the disk is away.
|
||||
(lib / "inbox").mkdir(parents=True, exist_ok=True)
|
||||
shutil.copy2(kept, lib / "inbox" / "again.jpg")
|
||||
scan = InventoryService(sf).scan(lib)
|
||||
new_id = scan.asset_ids[str(lib / "inbox" / "again.jpg")]
|
||||
assert scan.occurrences[str(lib / "inbox" / "again.jpg")] == "copied"
|
||||
|
||||
clusters = DuplicateService(sf).detect().clusters
|
||||
exact = [c for c in clusters if c["method"] == Method.EXACT.value]
|
||||
assert len(exact) == 1
|
||||
cluster = exact[0]
|
||||
# Byte-identical: linked directly, and the archived original stays canonical.
|
||||
assert cluster["state"] == ClusterState.DECIDED.value
|
||||
assert cluster["canonical_asset_id"] == archived_id
|
||||
assert _assets(sf)[new_id].canonical_asset_id == archived_id
|
||||
|
||||
|
||||
def test_fuzzy_copy_of_offline_asset_requires_review_and_names_the_medium(tmp_path):
|
||||
config, sf, lib, archive = _env(tmp_path)
|
||||
folder, ids = _uploaded(sf, lib, seeds=(1,))
|
||||
archived_id = next(iter(ids.values()))
|
||||
original = next(iter(ids))
|
||||
variant_source = resized_copy(original, tmp_path / "small.jpg")
|
||||
_archive(sf, config, archive)
|
||||
_unmount(archive)
|
||||
|
||||
(lib / "inbox").mkdir(parents=True, exist_ok=True)
|
||||
shutil.copy2(variant_source, lib / "inbox" / "small.jpg")
|
||||
InventoryService(sf).scan(lib)
|
||||
|
||||
duplicates = DuplicateService(sf)
|
||||
clusters = duplicates.detect().clusters
|
||||
perceptual = [c for c in clusters if c["method"] == Method.PERCEPTUAL.value]
|
||||
assert len(perceptual) == 1
|
||||
|
||||
detail = duplicates.get_cluster(perceptual[0]["id"])
|
||||
assert detail["state"] == ClusterState.OPEN.value # never auto-decided
|
||||
assert detail["requires_confirmation"] is True
|
||||
assert detail["mount_required"] == ["external"] # full-resolution needs the disk
|
||||
archived = next(m for m in detail["members"] if m["asset_id"] == archived_id)
|
||||
assert archived["availability_state"] == ARCHIVED_OFFLINE
|
||||
assert archived["current_path"] is None
|
||||
assert archived["archive_path"] and archived["evidence"]["phash"]
|
||||
# The retained preview is what makes the offline member reviewable at all.
|
||||
assert archived["preview"] == {"state": "ready", "protected": True, "size": PROTECTED_SIZE}
|
||||
|
||||
|
||||
# ── protected review evidence ────────────────────────────────────────────────
|
||||
|
||||
|
||||
def test_protected_preview_survives_quota_and_serves_while_offline(tmp_path):
|
||||
config, sf, lib, archive = _env(tmp_path)
|
||||
_, ids = _uploaded(sf, lib, seeds=(1,))
|
||||
asset_id = next(iter(ids.values()))
|
||||
_archive(sf, config, archive)
|
||||
_unmount(archive)
|
||||
|
||||
thumbnails = ThumbnailService(sf, config)
|
||||
served = thumbnails.generate(asset_id, PROTECTED_SIZE)
|
||||
assert served.exists() # rendered before the original left, not from the medium
|
||||
|
||||
# An aggressive quota may empty the cache, but not this evidence.
|
||||
tight = ThumbnailService(sf, config.model_copy(update={"thumbnail_cache_quota_bytes": 1}))
|
||||
tight._enforce_quota()
|
||||
assert served.exists()
|
||||
with sf() as session:
|
||||
row = session.scalar(select(Thumbnail).where(Thumbnail.asset_id == asset_id))
|
||||
assert row.protected is True
|
||||
assert thumbnails.evidence(asset_id)["state"] == "ready"
|
||||
|
||||
|
||||
def test_offline_asset_without_preview_reports_unavailable(tmp_path):
|
||||
"""No preview and no medium is an honest 409, never a wrong picture."""
|
||||
config, sf, lib, archive = _env(tmp_path)
|
||||
_, ids = _uploaded(sf, lib, seeds=(1,))
|
||||
asset_id = next(iter(ids.values()))
|
||||
_archive(sf, config, archive)
|
||||
with sf() as session:
|
||||
for row in session.scalars(select(Thumbnail).where(Thumbnail.asset_id == asset_id)):
|
||||
session.delete(row)
|
||||
session.commit()
|
||||
_unmount(archive)
|
||||
|
||||
with pytest.raises(ThumbnailUnavailable):
|
||||
ThumbnailService(sf, config).generate(asset_id, 256)
|
||||
|
||||
_remount(archive) # mounted again: the archived original is readable
|
||||
assert ThumbnailService(sf, config).generate(asset_id, 256).exists()
|
||||
|
||||
|
||||
def test_offline_asset_is_browsable_through_the_api(tmp_path):
|
||||
"""The browser sees an archived asset, its medium, and its preview — offline."""
|
||||
config, sf, lib, archive = _env(tmp_path)
|
||||
_, ids = _uploaded(sf, lib, seeds=(1,))
|
||||
asset_id = next(iter(ids.values()))
|
||||
_archive(sf, config, archive)
|
||||
_unmount(archive)
|
||||
|
||||
with TestClient(create_app(config)) as client:
|
||||
client.post("/api/v1/inventory/scan")
|
||||
listed = client.get("/api/v1/inventory/assets", params={"availability": ARCHIVED_OFFLINE})
|
||||
assert listed.status_code == 200
|
||||
item = next(row for row in listed.json()["items"] if row["id"] == asset_id)
|
||||
assert item["current_path"] is None
|
||||
assert item["archive_path"] == "rome/0.jpg"
|
||||
assert item["missing"] is False
|
||||
|
||||
# Searching by the archived path still finds it.
|
||||
found = client.get("/api/v1/inventory/assets", params={"q": "rome"}).json()
|
||||
assert [row["id"] for row in found["items"]] == [asset_id]
|
||||
|
||||
preview = client.get(f"/api/v1/assets/{asset_id}/thumbnail", params={"size": 1280})
|
||||
assert preview.status_code == 200
|
||||
assert preview.headers["content-type"] == "image/webp"
|
||||
|
||||
|
||||
def test_offline_state_is_stable_across_restart(tmp_path):
|
||||
config, sf, lib, archive = _env(tmp_path)
|
||||
_, ids = _uploaded(sf, lib, seeds=(1, 2))
|
||||
_archive(sf, config, archive)
|
||||
_unmount(archive)
|
||||
InventoryService(sf).scan(lib)
|
||||
before = {
|
||||
asset_id: (
|
||||
asset.availability_state,
|
||||
asset.archive_path,
|
||||
asset.current_sha256,
|
||||
asset.phash,
|
||||
)
|
||||
for asset_id, asset in _assets(sf).items()
|
||||
}
|
||||
|
||||
# Restart: a fresh engine and session factory against the same database.
|
||||
restarted = create_session_factory(create_db_engine(config.database_url))
|
||||
after = {
|
||||
asset_id: (
|
||||
asset.availability_state,
|
||||
asset.archive_path,
|
||||
asset.current_sha256,
|
||||
asset.phash,
|
||||
)
|
||||
for asset_id, asset in _assets(restarted).items()
|
||||
}
|
||||
assert after == before
|
||||
assert set(after) == set(ids.values())
|
||||
|
||||
# And the medium coming back is picked up by the restarted process.
|
||||
_remount(archive)
|
||||
assert ArchiveService(restarted, config=config).locations()[0]["state"] == "online"
|
||||
assert {a.availability_state for a in _assets(restarted).values()} == {ARCHIVED_ONLINE}
|
||||
assert ACTIVE not in {a.availability_state for a in _assets(restarted).values()}
|
||||
417
tests/integration/test_restore.py
Normal file
417
tests/integration/test_restore.py
Normal 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
|
||||
)
|
||||
@@ -131,6 +131,12 @@
|
||||
"tests/unit/test_archive_journal_states.py",
|
||||
"tests/integration/test_archive_transfer.py",
|
||||
"tests/integration/test_archive_recovery.py"
|
||||
],
|
||||
"US06-03": [
|
||||
"tests/integration/test_offline_assets.py"
|
||||
],
|
||||
"US06-04": [
|
||||
"tests/integration/test_restore.py"
|
||||
]
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user