Compare commits
1 Commits
us/US09-03
...
us/US06-04
| Author | SHA1 | Date | |
|---|---|---|---|
| 0982ebda70 |
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
|
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
|
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 fastapi.responses import JSONResponse
|
||||||
from pydantic import BaseModel
|
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.archives import ArchiveError, ArchiveService
|
||||||
from photo_pipeline.services.archive_transfer import ArchiveTransferService
|
from photo_pipeline.services.archive_transfer import ArchiveTransferService
|
||||||
from photo_pipeline.services.jobs import JobBlocked, JobService
|
from photo_pipeline.services.jobs import JobBlocked, JobService
|
||||||
|
from photo_pipeline.services.restores import RestoreService
|
||||||
|
|
||||||
router = APIRouter(tags=["archives"])
|
router = APIRouter(tags=["archives"])
|
||||||
|
|
||||||
@@ -42,10 +43,24 @@ class CreatePlanRequest(PreflightRequest):
|
|||||||
token: str
|
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:
|
def _service(request: Request) -> ArchiveService:
|
||||||
return ArchiveService(request.app.state.session_factory, config=request.app.state.config)
|
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:
|
def _transfers(request: Request) -> ArchiveTransferService:
|
||||||
return ArchiveTransferService(
|
return ArchiveTransferService(
|
||||||
request.app.state.session_factory, config=request.app.state.config
|
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}
|
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")
|
@router.get("/archive-recovery")
|
||||||
def recovery_status(request: Request) -> dict:
|
def recovery_status(request: Request) -> dict:
|
||||||
"""What an interrupted transfer left behind, straight from journal + disk."""
|
"""What an interrupted transfer left behind, straight from journal + disk."""
|
||||||
|
|||||||
@@ -1,9 +1,9 @@
|
|||||||
"""Domain job handlers: safety scoring, content analysis, uploads, archive
|
"""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``,
|
Importing this module registers the ``safety_score``, ``analysis``,
|
||||||
``upload_batch``, and ``archive_plan`` job types so the generic worker can run them
|
``upload_batch``, ``archive_plan``, and ``restore_plan`` job types so the generic
|
||||||
per item. Each handler delegates to its service, which owns the real work and the
|
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
|
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
|
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.
|
outcome is unknown, and an archive plan skips items it already completed.
|
||||||
@@ -20,6 +20,7 @@ SAFETY_SCORE = "safety_score"
|
|||||||
ANALYSIS = "analysis"
|
ANALYSIS = "analysis"
|
||||||
UPLOAD_BATCH = "upload_batch"
|
UPLOAD_BATCH = "upload_batch"
|
||||||
ARCHIVE_PLAN = "archive_plan"
|
ARCHIVE_PLAN = "archive_plan"
|
||||||
|
RESTORE_PLAN = "restore_plan"
|
||||||
# Both mutate the library's metadata/derived state; one at a time (concept §one job).
|
# Both mutate the library's metadata/derived state; one at a time (concept §one job).
|
||||||
LIBRARY_WRITE_LOCK = "library_write"
|
LIBRARY_WRITE_LOCK = "library_write"
|
||||||
# The uploader lane: one album batch at a time (concept §16).
|
# 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")
|
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(SAFETY_SCORE, _safety_score_item)
|
||||||
register(ANALYSIS, _analysis_item)
|
register(ANALYSIS, _analysis_item)
|
||||||
register(UPLOAD_BATCH, _upload_batch_item)
|
register(UPLOAD_BATCH, _upload_batch_item)
|
||||||
register(ARCHIVE_PLAN, _archive_plan_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.
|
# The preflight token this plan was approved against; re-verified before apply.
|
||||||
token: Mapped[str] = mapped_column(String, nullable=False)
|
token: Mapped[str] = mapped_column(String, nullable=False)
|
||||||
albums: Mapped[str | None] = mapped_column(String) # JSON array
|
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
|
# planned | applying | complete | failed
|
||||||
state: Mapped[str] = mapped_column(String, nullable=False, default="planned")
|
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)
|
album: Mapped[str] = mapped_column(String, nullable=False)
|
||||||
asset_id: Mapped[str] = mapped_column(ForeignKey("assets.id"), nullable=False, index=True)
|
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)
|
source_path: Mapped[str] = mapped_column(String, nullable=False)
|
||||||
destination_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.
|
# 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).
|
# link is written and read by the archive service (US06-02).
|
||||||
archive_location_id: Mapped[str | None] = mapped_column(String)
|
archive_location_id: Mapped[str | None] = mapped_column(String)
|
||||||
archive_path: 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.
|
# Duplicate canonical link: NULL when the asset is itself canonical or undecided.
|
||||||
canonical_asset_id: Mapped[str | None] = mapped_column(ForeignKey("assets.id"))
|
canonical_asset_id: Mapped[str | None] = mapped_column(ForeignKey("assets.id"))
|
||||||
|
|
||||||
|
|||||||
@@ -15,6 +15,11 @@ planned → transferring → verified → removing → complete
|
|||||||
↘ ↘ ↘ failed
|
↘ ↘ ↘ 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
|
- ``transferring`` — intent recorded; a temporary copy may exist, the destination
|
||||||
may or may not have been published. Nothing has been removed.
|
may or may not have been published. Nothing has been removed.
|
||||||
- ``verified`` — the archived bytes exist at their final path, hash exactly as
|
- ``verified`` — the archived bytes exist at their final path, hash exactly as
|
||||||
@@ -73,10 +78,22 @@ ALLOWED_TRANSITIONS = {
|
|||||||
ArchiveState.FAILED: {ArchiveState.PLANNED, ArchiveState.TRANSFERRING},
|
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})
|
TERMINAL_STATES = frozenset({ArchiveState.COMPLETE})
|
||||||
# States where this item may already have touched the filesystem.
|
# States where this item may already have touched the filesystem.
|
||||||
UNSAFE_STATES = frozenset({ArchiveState.TRANSFERRING, ArchiveState.VERIFIED, ArchiveState.REMOVING})
|
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"
|
RESUMABLE = "resumable"
|
||||||
FORWARD = "forward"
|
FORWARD = "forward"
|
||||||
MANUAL = "manual"
|
MANUAL = "manual"
|
||||||
@@ -94,8 +111,9 @@ class JournalConflict(JournalError):
|
|||||||
"""Fencing check failed; a newer owner has taken over this operation."""
|
"""Fencing check failed; a newer owner has taken over this operation."""
|
||||||
|
|
||||||
|
|
||||||
def can_transition(current: str, target: str) -> bool:
|
def can_transition(current: str, target: str, direction: str = ARCHIVE) -> bool:
|
||||||
return target in ALLOWED_TRANSITIONS.get(current, set())
|
table = RESTORE_TRANSITIONS if direction == RESTORE else ALLOWED_TRANSITIONS
|
||||||
|
return target in table.get(current, set())
|
||||||
|
|
||||||
|
|
||||||
def _now() -> datetime:
|
def _now() -> datetime:
|
||||||
@@ -119,7 +137,7 @@ class ArchiveJournal:
|
|||||||
if row.journal_state in TERMINAL_STATES:
|
if row.journal_state in TERMINAL_STATES:
|
||||||
raise InvalidTransition(f"{row.journal_state} is terminal")
|
raise InvalidTransition(f"{row.journal_state} is terminal")
|
||||||
if row.journal_state != ArchiveState.TRANSFERRING and not can_transition(
|
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}")
|
raise InvalidTransition(f"{row.journal_state} -> {ArchiveState.TRANSFERRING}")
|
||||||
if row.journal_state != ArchiveState.TRANSFERRING:
|
if row.journal_state != ArchiveState.TRANSFERRING:
|
||||||
@@ -160,7 +178,7 @@ class ArchiveJournal:
|
|||||||
if row.journal_state == target:
|
if row.journal_state == target:
|
||||||
session.commit()
|
session.commit()
|
||||||
return _operation_dict(row) # idempotent
|
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}")
|
raise InvalidTransition(f"{row.journal_state} -> {target}")
|
||||||
|
|
||||||
row.journal_state = target
|
row.journal_state = target
|
||||||
@@ -192,16 +210,18 @@ class ArchiveJournal:
|
|||||||
)
|
)
|
||||||
return [_operation_dict(row) for row in rows]
|
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
|
"""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:
|
with self._session_factory() as session:
|
||||||
|
stmt = select(ArchiveOperation).where(
|
||||||
|
ArchiveOperation.journal_state.not_in([*TERMINAL_STATES, ArchiveState.PLANNED])
|
||||||
|
)
|
||||||
|
if direction is not None:
|
||||||
|
stmt = stmt.where(ArchiveOperation.direction == direction)
|
||||||
rows = session.scalars(
|
rows = session.scalars(
|
||||||
select(ArchiveOperation)
|
stmt.order_by(ArchiveOperation.plan_id, ArchiveOperation.sequence)
|
||||||
.where(
|
|
||||||
ArchiveOperation.journal_state.not_in([*TERMINAL_STATES, ArchiveState.PLANNED])
|
|
||||||
)
|
|
||||||
.order_by(ArchiveOperation.plan_id, ArchiveOperation.sequence)
|
|
||||||
)
|
)
|
||||||
return [_operation_dict(row) for row in rows]
|
return [_operation_dict(row) for row in rows]
|
||||||
|
|
||||||
@@ -231,6 +251,7 @@ class ArchiveJournal:
|
|||||||
return {
|
return {
|
||||||
"operation_id": operation_id,
|
"operation_id": operation_id,
|
||||||
"plan_id": row["plan_id"],
|
"plan_id": row["plan_id"],
|
||||||
|
"direction": row["direction"],
|
||||||
"album": row["album"],
|
"album": row["album"],
|
||||||
"asset_id": row["asset_id"],
|
"asset_id": row["asset_id"],
|
||||||
"source_path": row["source_path"],
|
"source_path": row["source_path"],
|
||||||
@@ -243,8 +264,8 @@ class ArchiveJournal:
|
|||||||
"destination_matches": destination_matches,
|
"destination_matches": destination_matches,
|
||||||
}
|
}
|
||||||
|
|
||||||
def classify_all(self) -> list[dict]:
|
def classify_all(self, *, direction: str | None = None) -> list[dict]:
|
||||||
return [self.classify(row["id"]) for row in self.incomplete()]
|
return [self.classify(row["id"]) for row in self.incomplete(direction=direction)]
|
||||||
|
|
||||||
def blocks_mutation(self) -> bool:
|
def blocks_mutation(self) -> bool:
|
||||||
"""True when any item may have the library half-archived."""
|
"""True when any item may have the library half-archived."""
|
||||||
@@ -317,6 +338,7 @@ def _operation_dict(row: ArchiveOperation) -> dict:
|
|||||||
return {
|
return {
|
||||||
"id": row.id,
|
"id": row.id,
|
||||||
"plan_id": row.plan_id,
|
"plan_id": row.plan_id,
|
||||||
|
"direction": row.direction,
|
||||||
"sequence": row.sequence,
|
"sequence": row.sequence,
|
||||||
"album": row.album,
|
"album": row.album,
|
||||||
"asset_id": row.asset_id,
|
"asset_id": row.asset_id,
|
||||||
|
|||||||
@@ -54,6 +54,7 @@ from sqlalchemy.orm import sessionmaker
|
|||||||
from photo_pipeline.config import Config
|
from photo_pipeline.config import Config
|
||||||
from photo_pipeline.models import ArchiveLocation, ArchiveOperation, ArchivePlan, Asset, AssetPath
|
from photo_pipeline.models import ArchiveLocation, ArchiveOperation, ArchivePlan, Asset, AssetPath
|
||||||
from photo_pipeline.services.archive_journal import (
|
from photo_pipeline.services.archive_journal import (
|
||||||
|
ARCHIVE,
|
||||||
MANUAL,
|
MANUAL,
|
||||||
RESUMABLE,
|
RESUMABLE,
|
||||||
ArchiveJournal,
|
ArchiveJournal,
|
||||||
@@ -115,6 +116,7 @@ class ArchiveTransferService:
|
|||||||
location_id=location_id,
|
location_id=location_id,
|
||||||
token=token,
|
token=token,
|
||||||
albums=json.dumps(albums) if albums is not None else None,
|
albums=json.dumps(albums) if albums is not None else None,
|
||||||
|
direction=ARCHIVE,
|
||||||
state="planned",
|
state="planned",
|
||||||
schema_version=MANIFEST_VERSION,
|
schema_version=MANIFEST_VERSION,
|
||||||
asset_count=preflight["totals"]["assets"],
|
asset_count=preflight["totals"]["assets"],
|
||||||
@@ -132,6 +134,7 @@ class ArchiveTransferService:
|
|||||||
ArchiveOperation(
|
ArchiveOperation(
|
||||||
id=str(uuid.uuid4()),
|
id=str(uuid.uuid4()),
|
||||||
plan_id=plan_id,
|
plan_id=plan_id,
|
||||||
|
direction=ARCHIVE,
|
||||||
sequence=sequence,
|
sequence=sequence,
|
||||||
album=album["album"],
|
album=album["album"],
|
||||||
asset_id=asset["asset_id"],
|
asset_id=asset["asset_id"],
|
||||||
@@ -161,7 +164,11 @@ class ArchiveTransferService:
|
|||||||
|
|
||||||
def list(self) -> list[dict]:
|
def list(self) -> list[dict]:
|
||||||
with self._session_factory() as session:
|
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]
|
return [_plan_dict(row) for row in rows]
|
||||||
|
|
||||||
# ── apply ─────────────────────────────────────────────────────────────────
|
# ── apply ─────────────────────────────────────────────────────────────────
|
||||||
@@ -262,7 +269,7 @@ class ArchiveTransferService:
|
|||||||
if same_filesystem:
|
if same_filesystem:
|
||||||
os.rename(source, destination)
|
os.rename(source, destination)
|
||||||
else:
|
else:
|
||||||
self._copy_and_publish(operation, source, destination)
|
copy_verify_publish(source, destination, operation["expected_sha256"])
|
||||||
_fsync_dir(destination.parent)
|
_fsync_dir(destination.parent)
|
||||||
|
|
||||||
# 4. The published file is the archive only once it hashes as recorded.
|
# 4. The published file is the archive only once it hashes as recorded.
|
||||||
@@ -282,27 +289,6 @@ class ArchiveTransferService:
|
|||||||
# 5. Only now may the active source go.
|
# 5. Only now may the active source go.
|
||||||
self._finish(self.journal.get(operation["id"]), location, token=token, worker_id=worker_id)
|
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:
|
def _finish(self, operation: dict, location: dict, *, token: int, worker_id: str) -> None:
|
||||||
"""Drive an item whose archive copy is durable through removal and
|
"""Drive an item whose archive copy is durable through removal and
|
||||||
bookkeeping. Every step is idempotent, so recovery may replay it."""
|
bookkeeping. Every step is idempotent, so recovery may replay it."""
|
||||||
@@ -457,7 +443,7 @@ class ArchiveTransferService:
|
|||||||
"""
|
"""
|
||||||
results = {"resumed": 0, "completed": 0, "manual": 0}
|
results = {"resumed": 0, "completed": 0, "manual": 0}
|
||||||
touched: set[str] = set()
|
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"])
|
operation = self.journal.get(verdict["operation_id"])
|
||||||
touched.add(operation["plan_id"])
|
touched.add(operation["plan_id"])
|
||||||
token = (operation["fencing_token"] or 0) + 1
|
token = (operation["fencing_token"] or 0) + 1
|
||||||
@@ -484,7 +470,7 @@ class ArchiveTransferService:
|
|||||||
return results
|
return results
|
||||||
|
|
||||||
def recovery_status(self) -> dict:
|
def recovery_status(self) -> dict:
|
||||||
verdicts = self.journal.classify_all()
|
verdicts = self.journal.classify_all(direction=ARCHIVE)
|
||||||
return {
|
return {
|
||||||
"operations": verdicts,
|
"operations": verdicts,
|
||||||
"manual": [v for v in verdicts if v["classification"] == MANUAL],
|
"manual": [v for v in verdicts if v["classification"] == MANUAL],
|
||||||
@@ -528,6 +514,33 @@ class ArchiveTransferService:
|
|||||||
# ── module helpers ───────────────────────────────────────────────────────────
|
# ── 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:
|
def _same_filesystem(source: Path, destination_dir: Path) -> bool:
|
||||||
"""Proven at run time from the actual devices, never from the plan's preview."""
|
"""Proven at run time from the actual devices, never from the plan's preview."""
|
||||||
try:
|
try:
|
||||||
@@ -618,6 +631,7 @@ def _plan_dict(plan: ArchivePlan) -> dict:
|
|||||||
"id": plan.id,
|
"id": plan.id,
|
||||||
"location_id": plan.location_id,
|
"location_id": plan.location_id,
|
||||||
"token": plan.token,
|
"token": plan.token,
|
||||||
|
"direction": plan.direction,
|
||||||
"albums": json.loads(plan.albums) if plan.albums else None,
|
"albums": json.loads(plan.albums) if plan.albums else None,
|
||||||
"state": plan.state,
|
"state": plan.state,
|
||||||
"schema_version": plan.schema_version,
|
"schema_version": plan.schema_version,
|
||||||
|
|||||||
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}"
|
||||||
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
|
||||||
|
)
|
||||||
@@ -134,6 +134,9 @@
|
|||||||
],
|
],
|
||||||
"US06-03": [
|
"US06-03": [
|
||||||
"tests/integration/test_offline_assets.py"
|
"tests/integration/test_offline_assets.py"
|
||||||
|
],
|
||||||
|
"US06-04": [
|
||||||
|
"tests/integration/test_restore.py"
|
||||||
]
|
]
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user