Compare commits

...

1 Commits

Author SHA1 Message Date
9503fd1cfc US06-02: Transfer, Verify, and Remove Active Sources 2026-08-16 18:45:52 +02:00
14 changed files with 2167 additions and 24 deletions

View File

@@ -0,0 +1,96 @@
"""Archive plans, per-item transfer journal, and archived asset location (US06-02).
Revision ID: 0012_archive_plans
Revises: 0011_archive_locations
Create Date: 2026-08-16
The journal is what makes removing an original recoverable: every item records its
source, destination, expected hash, and the state it had reached before the process
died. Assets gain the archive location and relative path so an offline original is
still explained rather than looking missing.
"""
import sqlalchemy as sa
from alembic import op
revision = "0012_archive_plans"
down_revision = "0011_archive_locations"
branch_labels = None
depends_on = None
def upgrade() -> None:
op.create_table(
"archive_plans",
sa.Column("id", sa.String(), primary_key=True),
sa.Column(
"location_id",
sa.String(),
sa.ForeignKey("archive_locations.id"),
nullable=False,
index=True,
),
# The preflight token this plan was approved against.
sa.Column("token", sa.String(), nullable=False),
sa.Column("albums", sa.String(), nullable=True), # JSON array
# planned | applying | complete | failed
sa.Column("state", sa.String(), nullable=False, server_default="planned"),
sa.Column("schema_version", sa.Integer(), nullable=False, server_default="1"),
sa.Column("asset_count", sa.Integer(), nullable=False, server_default="0"),
sa.Column("byte_size", sa.Integer(), nullable=False, server_default="0"),
# Bumped on every claim and used as the fencing token.
sa.Column("version", sa.Integer(), nullable=False, server_default="1"),
sa.Column("worker_id", sa.String(), nullable=True),
sa.Column("completed_at", sa.DateTime(timezone=True), nullable=True),
sa.Column(
"created_at", sa.DateTime(timezone=True), nullable=False, server_default=sa.func.now()
),
sa.Column(
"updated_at", sa.DateTime(timezone=True), nullable=False, server_default=sa.func.now()
),
)
op.create_table(
"archive_operations",
sa.Column("id", sa.String(), primary_key=True),
sa.Column(
"plan_id",
sa.String(),
sa.ForeignKey("archive_plans.id", ondelete="CASCADE"),
nullable=False,
index=True,
),
sa.Column("sequence", sa.Integer(), nullable=False),
sa.Column("album", sa.String(), nullable=False),
sa.Column("asset_id", sa.String(), sa.ForeignKey("assets.id"), nullable=False, index=True),
sa.Column("source_path", sa.String(), nullable=False),
sa.Column("destination_path", sa.String(), nullable=False),
# Relative to the location root: the medium can be mounted anywhere later.
sa.Column("archive_path", sa.String(), nullable=False),
sa.Column("expected_sha256", sa.String(), nullable=False),
sa.Column("byte_size", sa.Integer(), nullable=True),
sa.Column("same_filesystem", sa.Boolean(), nullable=True),
# planned | transferring | verified | removing | complete | failed
sa.Column("journal_state", sa.String(), nullable=False, server_default="planned"),
sa.Column("attempt_count", sa.Integer(), nullable=False, server_default="0"),
sa.Column("fencing_token", sa.Integer(), nullable=True),
sa.Column("worker_id", sa.String(), nullable=True),
sa.Column("verified_at", sa.DateTime(timezone=True), nullable=True),
sa.Column("removed_at", sa.DateTime(timezone=True), nullable=True),
sa.Column("error_code", sa.String(), nullable=True),
sa.Column("error_message", sa.String(), nullable=True),
sa.Column(
"updated_at", sa.DateTime(timezone=True), nullable=False, server_default=sa.func.now()
),
sa.UniqueConstraint("plan_id", "sequence", name="uq_archive_operations_plan_sequence"),
)
# Plain columns: SQLite cannot ALTER a table to add a foreign key, and the
# relationship is enforced by the service that writes them.
op.add_column("assets", sa.Column("archive_location_id", sa.String(), nullable=True))
op.add_column("assets", sa.Column("archive_path", sa.String(), nullable=True))
def downgrade() -> None:
op.drop_column("assets", "archive_path")
op.drop_column("assets", "archive_location_id")
op.drop_table("archive_operations")
op.drop_table("archive_plans")

View File

@@ -1,9 +1,10 @@
"""Archive location and preflight API (US06-01).
"""Archive location, preflight, and plan API (US06-01, US06-02).
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
the token a later archive plan must present (US06-02). Neither endpoint moves or
removes a single library file.
the token an archive plan must present. Creating a plan writes only database rows —
the transfer itself runs on the durable ``archive`` lane, never in the request
thread, because it removes originals.
"""
from __future__ import annotations
@@ -12,12 +13,17 @@ 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.services.archives import ArchiveError, ArchiveService
from photo_pipeline.services.archive_transfer import ArchiveTransferService
from photo_pipeline.services.jobs import JobBlocked, JobService
router = APIRouter(tags=["archives"])
# Which failures are the caller's request (422) and which are a missing thing (404).
NOT_FOUND_CODES = {"unknown_location"}
# Which failures are the caller's request (422), a missing thing (404), or state
# that changed under the caller (409).
NOT_FOUND_CODES = {"unknown_location", "unknown_plan"}
CONFLICT_CODES = {"stale_token", "stale_plan", "archive_pending"}
class RegisterLocationRequest(BaseModel):
@@ -31,12 +37,28 @@ class PreflightRequest(BaseModel):
albums: list[str] | None = None
class CreatePlanRequest(PreflightRequest):
# The token of the preflight the user approved; a stale one is refused.
token: str
def _service(request: Request) -> ArchiveService:
return ArchiveService(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
)
def _error(error: ArchiveError) -> JSONResponse:
status = 404 if error.code in NOT_FOUND_CODES else 422
if error.code in NOT_FOUND_CODES:
status = 404
elif error.code in CONFLICT_CODES:
status = 409
else:
status = 422
return JSONResponse(
status_code=status, content={"error": {"code": error.code, "message": str(error)}}
)
@@ -61,3 +83,67 @@ def preflight(body: PreflightRequest, request: Request):
return _service(request).preflight(body.location_id, body.albums)
except ArchiveError as error:
return _error(error)
@router.post("/archive-plans", status_code=201)
def create_plan(body: CreatePlanRequest, request: Request):
"""Turn an approved preflight into a durable, journaled plan. Nothing moves."""
try:
return _transfers(request).create(body.location_id, body.albums, token=body.token)
except ArchiveError as error:
return _error(error)
@router.get("/archive-plans")
def list_plans(request: Request) -> dict:
return {"plans": _transfers(request).list()}
@router.get("/archive-plans/{plan_id}")
def get_plan(plan_id: str, request: Request):
plan = _transfers(request).get(plan_id)
if plan is None:
return _error(ArchiveError("unknown_plan", f"unknown archive plan {plan_id}"))
return plan
@router.post("/archive-plans/{plan_id}/apply")
def apply_plan(plan_id: str, request: Request):
"""Queue the transfer on the archiver lane. The worker removes the sources."""
service = _transfers(request)
plan = service.get(plan_id)
if plan is None:
return _error(ArchiveError("unknown_plan", f"unknown archive plan {plan_id}"))
if service.journal.blocks_mutation():
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(
ARCHIVE_PLAN,
lock=ARCHIVE_LOCK,
# One queued attempt per plan version: a double-clicked apply reuses it.
idempotency_key=f"archive:{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("/archive-recovery")
def recovery_status(request: Request) -> dict:
"""What an interrupted transfer left behind, straight from journal + disk."""
return _transfers(request).recovery_status()
@router.post("/archive-recovery/resolve")
def resolve_recovery(request: Request) -> dict:
return _transfers(request).recover()

View File

@@ -1,10 +1,12 @@
"""Domain job handlers: safety scoring, content analysis, uploads (US02-06, US05-02).
"""Domain job handlers: safety scoring, content analysis, uploads, archive
transfers (US02-06, US05-02, US06-02).
Importing this module registers the ``safety_score``, ``analysis``, and
``upload_batch`` 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, and an upload batch refuses to re-run an attempt whose outcome is unknown.
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
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.
Providers/models are the service defaults here (real NsfwModel / vision provider);
tests exercise the services directly with injected fakes rather than the worker.
@@ -17,12 +19,12 @@ from photo_pipeline.jobs.handlers import Cancelled, JobContext, register
SAFETY_SCORE = "safety_score"
ANALYSIS = "analysis"
UPLOAD_BATCH = "upload_batch"
ARCHIVE_PLAN = "archive_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).
UPLOAD_LOCK = "upload"
# The archiver lane: one archive/restore plan at a time (concept §16). No handler
# runs on it yet (US06-02); preflight already refuses to plan around a held lease.
# The archiver lane: one archive/restore plan at a time (concept §16).
ARCHIVE_LOCK = "archive"
@@ -54,6 +56,21 @@ def _upload_batch_item(batch_id: str, ctx: JobContext) -> None:
raise RuntimeError(f"upload batch {batch_id} is {batch['state']}: {batch['error_code']}")
def _archive_plan_item(plan_id: str, ctx: JobContext) -> None:
"""One item = one archive plan. Item-level failures stay in the journal (the
source is then still there); only an unusable plan fails the job."""
from photo_pipeline.config import Config
from photo_pipeline.services.archive_transfer import ArchiveTransferService
config = ctx.config if ctx.config is not None else Config.from_env()
result = ArchiveTransferService(ctx.session_factory, config=config).apply(
plan_id, worker_id=ctx.worker_id
)
if result["failed"]:
raise RuntimeError(f"archive 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)

View File

@@ -5,7 +5,7 @@ Alembic environment relies on.
"""
from photo_pipeline.models.albums import AlbumProposal
from photo_pipeline.models.archives import ArchiveLocation
from photo_pipeline.models.archives import ArchiveLocation, ArchiveOperation, ArchivePlan
from photo_pipeline.models.assets import Asset, AssetPath
from photo_pipeline.models.duplicates import (
DuplicateCluster,
@@ -21,6 +21,8 @@ from photo_pipeline.models.workflow import AnalysisResult, SafetyReview
__all__ = [
"AlbumProposal",
"ArchiveLocation",
"ArchiveOperation",
"ArchivePlan",
"Asset",
"AssetPath",
"DuplicateCluster",

View File

@@ -1,4 +1,4 @@
"""Archive location persistence (US06-01).
"""Archive location, plan, and transfer-journal persistence (US06-01, US06-02).
An archive location is a *medium*, not a path. External disks get mounted at
different mountpoints, and a different disk can be mounted at the same one, so a
@@ -9,13 +9,27 @@ location therefore owns a marker file written onto the medium itself; its
``capabilities`` and ``state`` are the last probe result, kept so the UI can list
locations without touching a sleeping disk. Preflight always re-probes — a stored
state is a hint, never evidence.
An ``ArchivePlan`` is one approved preflight turned into durable work, and each
``ArchiveOperation`` is one file's crash-safe journal row (US06-02). The row records
what the transfer *intends* to do before it does it — source, destination, expected
hash — because after a crash that intent plus the files on disk is the only evidence
available for deciding whether an original may be removed.
"""
from __future__ import annotations
from datetime import datetime
from sqlalchemy import DateTime, String, func
from sqlalchemy import (
Boolean,
DateTime,
ForeignKey,
Integer,
String,
UniqueConstraint,
func,
)
from sqlalchemy.orm import Mapped, mapped_column
from photo_pipeline.db import Base
@@ -40,3 +54,69 @@ class ArchiveLocation(Base):
updated_at: Mapped[datetime] = mapped_column(
DateTime(timezone=True), nullable=False, server_default=func.now(), onupdate=func.now()
)
class ArchivePlan(Base):
__tablename__ = "archive_plans"
id: Mapped[str] = mapped_column(String, primary_key=True)
location_id: Mapped[str] = mapped_column(
ForeignKey("archive_locations.id"), nullable=False, index=True
)
# 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
# planned | applying | complete | failed
state: Mapped[str] = mapped_column(String, nullable=False, default="planned")
schema_version: Mapped[int] = mapped_column(Integer, nullable=False, default=1)
asset_count: Mapped[int] = mapped_column(Integer, nullable=False, default=0)
byte_size: Mapped[int] = mapped_column(Integer, nullable=False, default=0)
# Bumped on every claim and used as the fencing token, so a superseded attempt
# cannot commit.
version: Mapped[int] = mapped_column(Integer, nullable=False, default=1)
worker_id: Mapped[str | None] = mapped_column(String)
completed_at: Mapped[datetime | None] = mapped_column(DateTime(timezone=True))
created_at: Mapped[datetime] = mapped_column(
DateTime(timezone=True), nullable=False, server_default=func.now()
)
updated_at: Mapped[datetime] = mapped_column(
DateTime(timezone=True), nullable=False, server_default=func.now(), onupdate=func.now()
)
class ArchiveOperation(Base):
__tablename__ = "archive_operations"
__table_args__ = (
UniqueConstraint("plan_id", "sequence", name="uq_archive_operations_plan_sequence"),
)
id: Mapped[str] = mapped_column(String, primary_key=True)
plan_id: Mapped[str] = mapped_column(
ForeignKey("archive_plans.id", ondelete="CASCADE"), nullable=False, index=True
)
sequence: Mapped[int] = mapped_column(Integer, nullable=False)
album: Mapped[str] = mapped_column(String, nullable=False)
asset_id: Mapped[str] = mapped_column(ForeignKey("assets.id"), nullable=False, index=True)
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.
archive_path: Mapped[str] = mapped_column(String, nullable=False)
expected_sha256: Mapped[str] = mapped_column(String, nullable=False)
byte_size: Mapped[int | None] = mapped_column(Integer)
same_filesystem: Mapped[bool | None] = mapped_column(Boolean)
# planned | transferring | verified | removing | complete | failed
journal_state: Mapped[str] = mapped_column(String, nullable=False, default="planned")
attempt_count: Mapped[int] = mapped_column(Integer, nullable=False, default=0)
fencing_token: Mapped[int | None] = mapped_column(Integer)
worker_id: Mapped[str | None] = mapped_column(String)
verified_at: Mapped[datetime | None] = mapped_column(DateTime(timezone=True))
removed_at: Mapped[datetime | None] = mapped_column(DateTime(timezone=True))
error_code: Mapped[str | None] = mapped_column(String)
error_message: Mapped[str | None] = mapped_column(String)
updated_at: Mapped[datetime] = mapped_column(
DateTime(timezone=True), nullable=False, server_default=func.now(), onupdate=func.now()
)

View File

@@ -35,7 +35,14 @@ class Asset(Base):
discovered_at: Mapped[datetime] = mapped_column(DateTime(timezone=True), nullable=False)
missing_at: Mapped[datetime | None] = mapped_column(DateTime(timezone=True))
# active | archiving | archived_online | archived_offline | restoring |
# missing_unexpected (concept §9). ``current_path`` is NULL once archived; the
# original is then explained by the location plus its relative archive path.
availability_state: Mapped[str] = mapped_column(String, nullable=False, default="active")
# Not a declared foreign key: SQLite cannot add one to an existing table, so the
# 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)
# Duplicate canonical link: NULL when the asset is itself canonical or undecided.
canonical_asset_id: Mapped[str | None] = mapped_column(ForeignKey("assets.id"))

View File

@@ -0,0 +1,337 @@
"""Archive transfer journal — the durable record of every per-file transition
(US06-02).
Archiving is the only stage that deletes an original, so the journal exists to make
one question answerable after any crash: *may this source file be removed?* Intent
is written before the mutation it describes, and the recorded state plus the real
files on disk are the sole basis for answering it later. This module owns the state
machine, the durable writes, and the evidence table; it never touches a photo
(:mod:`photo_pipeline.services.archive_transfer` does).
Per-item state machine (concept §9 "Transfer and removal semantics"):
```
planned → transferring → verified → removing → complete
↘ ↘ ↘ failed
```
- ``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
recorded, and the manifest entry is durable. Only from here may a source go.
- ``removing`` — the source removal is committed to; the source may already be
gone while the database still points at it.
- ``complete`` — source absent, database updated, availability recorded.
``classify`` labels each incomplete item from the journal plus disk evidence:
- ``resumable`` — nothing was published; the source is intact, so applying again is
safe.
- ``forward`` — the archived copy exists and matches its recorded hash, so the
remaining steps (manifest, removal, bookkeeping) can be finished deterministically.
- ``manual`` — the evidence contradicts the journal (missing archive copy, wrong
bytes, source and archive both gone). Nothing is guessed and nothing is removed;
the item blocks unrelated mutations until a human decides.
"""
from __future__ import annotations
from datetime import datetime, timezone
from pathlib import Path
from sqlalchemy import select
from sqlalchemy.orm import sessionmaker
from photo_pipeline.models import ArchiveOperation, ArchivePlan
from photo_pipeline.services.hashing import sha256_file
class ArchiveState:
PLANNED = "planned"
TRANSFERRING = "transferring"
VERIFIED = "verified"
REMOVING = "removing"
COMPLETE = "complete"
FAILED = "failed"
ALLOWED_TRANSITIONS = {
ArchiveState.PLANNED: {ArchiveState.TRANSFERRING, ArchiveState.FAILED},
# From `transferring` the outcome is unknown until evidence is gathered, so it
# may resolve forward, back to planned (proven nothing was published), or fail.
ArchiveState.TRANSFERRING: {
ArchiveState.VERIFIED,
ArchiveState.PLANNED,
ArchiveState.FAILED,
},
ArchiveState.VERIFIED: {ArchiveState.REMOVING, ArchiveState.FAILED},
# No path back: once the source may be gone, only finishing is safe.
ArchiveState.REMOVING: {ArchiveState.COMPLETE, ArchiveState.FAILED},
ArchiveState.COMPLETE: set(),
# A retry re-enters `transferring`, which rechecks every precondition from
# scratch; recovery may also reset a failed item to `planned`.
ArchiveState.FAILED: {ArchiveState.PLANNED, ArchiveState.TRANSFERRING},
}
TERMINAL_STATES = frozenset({ArchiveState.COMPLETE})
# States where this item may already have touched the filesystem.
UNSAFE_STATES = frozenset({ArchiveState.TRANSFERRING, ArchiveState.VERIFIED, ArchiveState.REMOVING})
RESUMABLE = "resumable"
FORWARD = "forward"
MANUAL = "manual"
class JournalError(RuntimeError):
pass
class InvalidTransition(JournalError):
pass
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 _now() -> datetime:
return datetime.now(timezone.utc)
class ArchiveJournal:
def __init__(self, session_factory: sessionmaker) -> None:
self._session_factory = session_factory
# ── intent ────────────────────────────────────────────────────────────────
def begin(self, operation_id: str, *, worker_id: str, fencing_token: int) -> dict:
"""Record the intent to transfer **before** touching the filesystem."""
with self._session_factory() as session:
row = self._require(session, operation_id)
if row.fencing_token is not None and fencing_token < row.fencing_token:
raise JournalConflict(
f"stale fencing token {fencing_token} (current {row.fencing_token})"
)
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
):
raise InvalidTransition(f"{row.journal_state} -> {ArchiveState.TRANSFERRING}")
if row.journal_state != ArchiveState.TRANSFERRING:
row.attempt_count += 1
row.journal_state = ArchiveState.TRANSFERRING
row.worker_id = worker_id
row.fencing_token = fencing_token
row.error_code = row.error_message = None
row.updated_at = _now()
session.commit()
return _operation_dict(row)
# ── transitions ───────────────────────────────────────────────────────────
def transition(
self,
operation_id: str,
target: str,
*,
fencing_token: int | None = None,
error: tuple[str, str] | None = None,
same_filesystem: bool | None = None,
) -> dict:
"""Move one operation to ``target``, enforcing the state machine.
Re-entering the state an operation already holds is a no-op, which is what
makes recovery idempotent across repeated restarts.
"""
with self._session_factory() as session:
row = self._require(session, operation_id)
if fencing_token is not None and row.fencing_token is not None:
if fencing_token < row.fencing_token:
raise JournalConflict(
f"stale fencing token {fencing_token} (current {row.fencing_token})"
)
if same_filesystem is not None:
row.same_filesystem = same_filesystem
if row.journal_state == target:
session.commit()
return _operation_dict(row) # idempotent
if not can_transition(row.journal_state, target):
raise InvalidTransition(f"{row.journal_state} -> {target}")
row.journal_state = target
row.updated_at = _now()
if target == ArchiveState.VERIFIED:
row.verified_at = _now()
if target == ArchiveState.COMPLETE:
row.removed_at = _now()
if error:
row.error_code, row.error_message = error[0], error[1][:500]
elif target != ArchiveState.FAILED:
row.error_code = row.error_message = None
session.commit()
return _operation_dict(row)
# ── reads ─────────────────────────────────────────────────────────────────
def get(self, operation_id: str) -> dict | None:
with self._session_factory() as session:
row = session.get(ArchiveOperation, operation_id)
return _operation_dict(row) if row else None
def operations(self, plan_id: str) -> list[dict]:
with self._session_factory() as session:
rows = session.scalars(
select(ArchiveOperation)
.where(ArchiveOperation.plan_id == plan_id)
.order_by(ArchiveOperation.sequence)
)
return [_operation_dict(row) for row in rows]
def incomplete(self) -> list[dict]:
"""Every operation left in a non-terminal, non-planned state — the work a
restart has to reason about."""
with self._session_factory() as session:
rows = session.scalars(
select(ArchiveOperation)
.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]
# ── startup classification ────────────────────────────────────────────────
def classify(self, operation_id: str) -> dict:
"""Classify one incomplete operation from the journal plus disk evidence.
Hashes the archived copy when one exists: "a file is at the destination" is
not evidence that the *right* bytes are, and only the right bytes justify
removing an original. Never mutates anything.
"""
row = self.get(operation_id)
if row is None:
raise JournalError(f"unknown archive operation {operation_id!r}")
source = Path(row["source_path"])
destination = Path(row["destination_path"])
source_exists = source.exists()
destination_exists = destination.exists()
destination_matches = (
destination_exists and sha256_file(destination) == row["expected_sha256"]
)
classification, reason = _classify(
row["journal_state"], source_exists, destination_exists, destination_matches
)
return {
"operation_id": operation_id,
"plan_id": row["plan_id"],
"album": row["album"],
"asset_id": row["asset_id"],
"source_path": row["source_path"],
"destination_path": row["destination_path"],
"journal_state": row["journal_state"],
"classification": classification,
"reason": reason,
"source_exists": source_exists,
"destination_exists": destination_exists,
"destination_matches": destination_matches,
}
def classify_all(self) -> list[dict]:
return [self.classify(row["id"]) for row in self.incomplete()]
def blocks_mutation(self) -> bool:
"""True when any item may have the library half-archived."""
return any(row["journal_state"] in UNSAFE_STATES for row in self.incomplete())
# ── plan-level ────────────────────────────────────────────────────────────
def plan_state(self, plan_id: str) -> str:
"""Derive the plan's state from its items, so the summary can never disagree
with the journal."""
states = {row["journal_state"] for row in self.operations(plan_id)}
if not states:
return "planned"
if states <= {ArchiveState.COMPLETE}:
return "complete"
if states & {ArchiveState.FAILED}:
return "failed"
if states & UNSAFE_STATES:
return "applying"
return "planned"
def sync_plan_state(self, plan_id: str) -> str:
state = self.plan_state(plan_id)
with self._session_factory() as session:
plan = session.get(ArchivePlan, plan_id)
if plan is None:
raise JournalError(f"unknown archive plan {plan_id!r}")
if plan.state != state:
plan.state = state
plan.version += 1
plan.updated_at = _now()
if state == "complete" and plan.completed_at is None:
plan.completed_at = _now()
session.commit()
return state
@staticmethod
def _require(session, operation_id: str) -> ArchiveOperation:
row = session.get(ArchiveOperation, operation_id)
if row is None:
raise JournalError(f"unknown archive operation {operation_id!r}")
return row
def _classify(
state: str, source_exists: bool, destination_exists: bool, destination_matches: bool
) -> tuple[str, str]:
"""The evidence table. Kept a pure function so every combination is testable."""
if destination_exists and not destination_matches and state != ArchiveState.PLANNED:
# Someone else's file, or a partial/edited copy: never overwrite it, and
# never treat it as the durable archive that justifies a deletion.
return MANUAL, "the archived path holds bytes that are not the recorded ones"
if state in (ArchiveState.TRANSFERRING, ArchiveState.FAILED):
if destination_matches:
return FORWARD, "the archived copy is durable; finish the remaining steps"
if source_exists:
return RESUMABLE, "nothing was published; the source is intact"
return MANUAL, "neither the source nor a verified archive copy is present"
if state in (ArchiveState.VERIFIED, ArchiveState.REMOVING):
if destination_matches:
return FORWARD, "the archived copy is durable; finish the remaining steps"
return MANUAL, f"journal says {state} but the archived copy is missing"
return MANUAL, f"unhandled journal state {state}"
def _operation_dict(row: ArchiveOperation) -> dict:
return {
"id": row.id,
"plan_id": row.plan_id,
"sequence": row.sequence,
"album": row.album,
"asset_id": row.asset_id,
"source_path": row.source_path,
"destination_path": row.destination_path,
"archive_path": row.archive_path,
"expected_sha256": row.expected_sha256,
"byte_size": row.byte_size,
"same_filesystem": row.same_filesystem,
"journal_state": row.journal_state,
"attempt_count": row.attempt_count,
"fencing_token": row.fencing_token,
"worker_id": row.worker_id,
"verified_at": row.verified_at.isoformat() if row.verified_at else None,
"removed_at": row.removed_at.isoformat() if row.removed_at else None,
"error_code": row.error_code,
"error_message": row.error_message,
}

View File

@@ -0,0 +1,605 @@
"""Transfer an approved archive plan, verify it, and remove the active sources
(US06-02).
This is the only module that deletes originals from the photo library, so every
step exists to make one promise keepable: **a source is removed only after the
archived bytes are durable and proven identical.** The journal (US06-02,
:mod:`photo_pipeline.services.archive_journal`) records intent before each mutation;
this module performs the mutations and the recovery that reads that intent back.
Per file the sequence is:
```
journal.begin (transferring) ← intent persisted BEFORE any disk change
recheck preconditions ← source hash, free destination, no symlink
copy to a temporary file ← same directory, so the publish is atomic
fsync, close, hash it back ← read from disk; the write is not the evidence
atomically publish ← rename onto the final archive path
append the manifest entry ← durable on the medium itself, fsynced
journal → verified
journal → removing ← intent to delete, persisted first
re-verify the archive copy, remove the source, verify its absence
current_path=NULL, close the path occurrence, availability + location recorded
journal → complete
```
Same-filesystem albums may skip the copy and use an atomic ``rename`` instead
(concept §9), but only when the shared device is proven at run time — never from the
plan's stored guess — and the published file is hashed afterwards exactly as in the
copy path.
Rules that are never relaxed:
- An occupied destination is never overwritten; the item fails with the source
untouched.
- A source whose bytes no longer match the plan is never archived and never removed.
- A crash resolves from journal + disk evidence only: an archive copy that is
missing or hashes differently blocks the item for a human instead of being
retried or, worse, treated as a successful archive.
- Recovery is idempotent — repeated passes converge on the same state.
"""
from __future__ import annotations
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.models import ArchiveLocation, ArchiveOperation, ArchivePlan, Asset, AssetPath
from photo_pipeline.services.archive_journal import (
MANUAL,
RESUMABLE,
ArchiveJournal,
ArchiveState,
)
from photo_pipeline.services.archives import MARKER_NAME, ArchiveError, ArchiveService
from photo_pipeline.services.hashing import sha256_file
from photo_pipeline.services.rename_apply import PreconditionFailed, maybe_fault
# 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
# read back if the database is lost.
MANIFEST_NAME = "archive-manifest.jsonl"
MANIFEST_VERSION = 1
TEMP_SUFFIX = ".part"
TEMP_PREFIX = ".archive-"
# Fault barrier between "the source is gone" and "the database knows it" — not a
# journal state, but the transition crash tests care about most.
SOURCE_REMOVED = "source_removed"
APPLYABLE_PLAN_STATES = frozenset({"planned", "applying", "failed", "complete"})
def _now() -> datetime:
return datetime.now(timezone.utc)
class ArchiveTransferService:
def __init__(self, session_factory: sessionmaker, *, config: Config) -> None:
self._session_factory = session_factory
self._config = config
self.journal = ArchiveJournal(session_factory)
# ── plans ─────────────────────────────────────────────────────────────────
def create(self, location_id: str, albums: list[str] | None = None, *, token: str) -> dict:
"""Turn an approved preflight into a durable plan.
The token is re-derived from a fresh preflight, so a plan can only be
created for the exact scope, bytes, and destination the user approved.
"""
preflight = ArchiveService(self._session_factory, config=self._config).preflight(
location_id, albums
)
if not token or token != preflight["token"]:
raise ArchiveError("stale_token", "the archive preflight changed since it was approved")
if preflight["state"] != "ready":
codes = ", ".join(sorted({issue["code"] for issue in preflight["blockers"]})) or "-"
raise ArchiveError("blocked", f"the archive scope is blocked: {codes}")
root = Path(preflight["location"]["root"])
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(albums) if albums is not None else None,
state="planned",
schema_version=MANIFEST_VERSION,
asset_count=preflight["totals"]["assets"],
byte_size=preflight["totals"]["bytes"],
)
)
session.flush() # the plan row must exist before its items reference it
sequence = 0
for album in preflight["albums"]:
destination_dir = Path(album["destination"])
for asset in album["assets"]:
source = Path(asset["current_path"])
destination = destination_dir / source.name
session.add(
ArchiveOperation(
id=str(uuid.uuid4()),
plan_id=plan_id,
sequence=sequence,
album=album["album"],
asset_id=asset["asset_id"],
source_path=str(source),
destination_path=str(destination),
archive_path=str(destination.relative_to(root)),
expected_sha256=asset["current_sha256"],
byte_size=asset["byte_size"],
# Recorded as a preview only; the real decision is made
# against the devices at apply time.
same_filesystem=album["transfer_method"] == "move",
journal_state=ArchiveState.PLANNED,
)
)
sequence += 1
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:
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).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 = "archive",
) -> dict:
"""Archive every item of a plan, then report what happened.
Items are independent: one failure records its reason and leaves that
source in place; the rest of the album continues.
"""
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']}")
# One archiver lane: never start while another plan may be half-archived
# (concept §16 lock hierarchy).
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"])
archived = failed = skipped = 0
for operation in self.journal.operations(plan_id):
if operation["journal_state"] == ArchiveState.COMPLETE:
skipped += 1 # repeated apply is a no-op for finished work
continue
try:
if operation["journal_state"] in (ArchiveState.VERIFIED, ArchiveState.REMOVING):
# The bytes are already archived; never transfer them twice.
self._finish(operation, location, token=token, worker_id=worker_id)
else:
self._archive_one(operation, location, token=token, worker_id=worker_id)
archived += 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, "archive_error", str(error))
failed += 1
self._prune_empty_sources(plan_id)
state = self.journal.sync_plan_state(plan_id)
return {
"plan_id": plan_id,
"archived": archived,
"failed": failed,
"skipped": skipped,
"state": state,
}
def _prune_empty_sources(self, plan_id: str) -> None:
"""Drop an album folder once every one of its files is archived.
``rmdir`` only: a folder that still holds anything at all — an unarchived
file, someone else's file, a subfolder — is left exactly as it is.
"""
folders: dict[Path, set[str]] = {}
for operation in self.journal.operations(plan_id):
folders.setdefault(Path(operation["source_path"]).parent, set()).add(
operation["journal_state"]
)
for folder, states in folders.items():
if states == {ArchiveState.COMPLETE}:
try:
folder.rmdir()
except OSError:
pass # not empty, or gone already; either way, leave it alone
def _archive_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 — after this point a crash is recoverable from evidence.
self.journal.begin(operation["id"], worker_id=worker_id, fencing_token=token)
maybe_fault(ArchiveState.TRANSFERRING)
# 2. Recheck immediately before mutating; the plan's snapshot is not trusted.
self._recheck(operation, source, destination, location)
destination.parent.mkdir(parents=True, exist_ok=True)
# 3. Transfer. Same-filesystem is an optimisation, so it has to be proven
# here rather than believed from the plan.
same_filesystem = _same_filesystem(source, destination.parent)
if same_filesystem:
os.rename(source, destination)
else:
self._copy_and_publish(operation, source, destination)
_fsync_dir(destination.parent)
# 4. The published file is the archive only once it hashes as recorded.
if sha256_file(destination) != operation["expected_sha256"]:
raise PreconditionFailed(
"archive_mismatch", f"{destination} does not hold the expected bytes"
)
_append_manifest(destination.parent, _manifest_entry(operation, location))
self.journal.transition(
operation["id"],
ArchiveState.VERIFIED,
fencing_token=token,
same_filesystem=same_filesystem,
)
maybe_fault(ArchiveState.VERIFIED)
# 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."""
destination = Path(operation["destination_path"])
source = Path(operation["source_path"])
state = operation["journal_state"]
if state == ArchiveState.FAILED:
# The archive copy is durable even though the attempt ended badly:
# re-enter the transfer state so the remaining steps can run.
self.journal.transition(operation["id"], ArchiveState.TRANSFERRING, fencing_token=token)
state = ArchiveState.TRANSFERRING
if state == ArchiveState.TRANSFERRING:
if sha256_file(destination) != operation["expected_sha256"]:
raise PreconditionFailed(
"archive_mismatch", f"{destination} does not hold the expected bytes"
)
_append_manifest(destination.parent, _manifest_entry(operation, location))
self.journal.transition(operation["id"], ArchiveState.VERIFIED, fencing_token=token)
state = ArchiveState.VERIFIED
if state == ArchiveState.VERIFIED:
self.journal.transition(operation["id"], ArchiveState.REMOVING, fencing_token=token)
maybe_fault(ArchiveState.REMOVING)
state = ArchiveState.REMOVING
if state == ArchiveState.REMOVING:
# Re-verify the archived bytes immediately before deleting the original:
# this check is the entire justification for the removal.
if not destination.exists() or sha256_file(destination) != operation["expected_sha256"]:
raise PreconditionFailed(
"archive_unverified", f"{destination} is not a verified archive copy"
)
if source.exists():
if source.is_symlink():
raise PreconditionFailed("symlink", f"{source} became a symlink")
if sha256_file(source) != operation["expected_sha256"]:
raise PreconditionFailed(
"source_changed", f"{source} changed; it is not ours to remove"
)
source.unlink()
if source.exists():
raise PreconditionFailed("removal_failed", f"{source} is still present")
# The dangerous window: active storage no longer holds the file while the
# database still points at it.
maybe_fault(SOURCE_REMOVED)
self._record_archived(operation, location, 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:
if not source.exists():
raise PreconditionFailed("source_missing", f"source {source} disappeared")
if source.is_symlink() or destination.is_symlink():
raise PreconditionFailed("symlink", "refusing to archive through a symlink")
if destination.exists():
raise PreconditionFailed("destination_exists", f"destination {destination} is occupied")
root = Path(location["root"])
if root not in destination.parents:
raise PreconditionFailed(
"destination_escape", f"{destination} is outside the archive location {root}"
)
if not root.is_dir() or not (root / MARKER_NAME).exists():
raise PreconditionFailed("location_offline", f"{root} is not the archive medium")
if sha256_file(source) != operation["expected_sha256"]:
raise PreconditionFailed(
"source_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.current_path != str(source):
raise PreconditionFailed(
"asset_moved", f"asset {operation['asset_id']} is no longer at {source}"
)
# ── database ──────────────────────────────────────────────────────────────
def _record_archived(self, operation: dict, location: dict, destination: Path) -> None:
"""The original is gone from active storage: drop ``current_path``, close its
occurrence, record where the bytes now live, and set availability."""
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"
)
if asset.current_path:
for row in session.scalars(
select(AssetPath).where(
AssetPath.asset_id == asset.id,
AssetPath.path == asset.current_path,
AssetPath.valid_until.is_(None),
)
):
row.valid_until = now
recorded = session.scalar(
select(AssetPath).where(
AssetPath.asset_id == asset.id, AssetPath.path == str(destination)
)
)
if recorded is None: # idempotent: recovery may replay this
session.add(
AssetPath(
asset_id=asset.id,
path=str(destination),
valid_from=now,
reason="archive",
)
)
asset.current_path = None
asset.availability_state = (
"archived_online" if destination.exists() else "archived_offline"
)
asset.archive_location_id = location["id"]
asset.archive_path = operation["archive_path"]
asset.state_version += 1
asset.updated_at = now
session.commit()
# ── recovery ──────────────────────────────────────────────────────────────
def recover(self, *, worker_id: str = "archive-recovery") -> dict:
"""Resolve every incomplete item from journal + disk evidence.
Idempotent: running it repeatedly converges. Ambiguous (``manual``) work is
left exactly as found and keeps blocking unrelated mutations.
"""
results = {"resumed": 0, "completed": 0, "manual": 0}
touched: set[str] = set()
for verdict in self.journal.classify_all():
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:
# Nothing was published: discard the debris and let a later apply
# retry the item cleanly.
_clean_temp_files(Path(operation["destination_path"]).parent)
self.journal.transition(operation["id"], ArchiveState.PLANNED, fencing_token=token)
results["resumed"] += 1
continue
location = self._location(self._require_plan(operation["plan_id"])["location_id"])
try:
self._finish(operation, location, token=token, worker_id=worker_id)
results["completed"] += 1
except PreconditionFailed as error:
self._fail(operation, token, error.code, str(error))
results["manual"] += 1
for plan_id in touched:
self._prune_empty_sources(plan_id)
self.journal.sync_plan_state(plan_id)
return results
def recovery_status(self) -> dict:
verdicts = self.journal.classify_all()
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:
with self._session_factory() as session:
plan = session.get(ArchivePlan, plan_id)
if plan is None:
raise ArchiveError("unknown_plan", f"unknown archive plan {plan_id!r}")
return _plan_dict(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:
"""Bump the plan version and use it as this attempt's fencing token, so a
worker from a superseded attempt cannot commit."""
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 _same_filesystem(source: Path, destination_dir: Path) -> bool:
"""Proven at run time from the actual devices, never from the plan's preview."""
try:
return source.stat().st_dev == destination_dir.stat().st_dev
except OSError:
return False
def _fsync_dir(path: Path) -> None:
"""Make the directory entry itself durable, so the published name survives a
power loss and not just the file's data."""
fd = os.open(path, os.O_RDONLY)
try:
os.fsync(fd)
except OSError:
pass # some filesystems refuse directory fsync; the data is already synced
finally:
os.close(fd)
def _manifest_entry(operation: dict, location: dict) -> dict:
return {
"schema_version": MANIFEST_VERSION,
"plan_id": operation["plan_id"],
"asset_id": operation["asset_id"],
"album": operation["album"],
"archive_path": operation["archive_path"],
"source_path": operation["source_path"],
"sha256": operation["expected_sha256"],
"byte_size": operation["byte_size"],
"media_id": location["media_id"],
"archived_at": _now().isoformat(),
}
def _append_manifest(directory: Path, entry: dict) -> None:
"""Append one durable manifest line, skipping an entry that is already there.
The manifest is written before the source is removed, so it is the medium's own
record of what it holds even if the database is lost.
"""
path = directory / MANIFEST_NAME
# ponytail: rereads the album manifest per file (O(n²) lines for one album).
# Keep an in-memory index per plan if an album ever holds enough files to matter.
if path.exists():
for line in path.read_text(encoding="utf-8").splitlines():
try:
existing = json.loads(line)
except ValueError:
continue
if (existing.get("asset_id"), existing.get("sha256")) == (
entry["asset_id"],
entry["sha256"],
):
return
with open(path, "a", encoding="utf-8") as handle:
handle.write(json.dumps(entry, sort_keys=True) + "\n")
handle.flush()
os.fsync(handle.fileno())
_fsync_dir(directory)
def read_manifest(directory: Path) -> list[dict]:
"""Every manifest entry an archived album directory holds."""
path = directory / MANIFEST_NAME
if not path.exists():
return []
entries = []
for line in path.read_text(encoding="utf-8").splitlines():
try:
entries.append(json.loads(line))
except ValueError:
continue
return entries
def _clean_temp_files(directory: Path) -> None:
"""Remove this application's own abandoned transfer temporaries — never any
other file (concept §17: startup cleans only recognised stale temporaries)."""
if not directory.is_dir():
return
for temp in directory.glob(f"{TEMP_PREFIX}*{TEMP_SUFFIX}"):
temp.unlink(missing_ok=True)
def _plan_dict(plan: ArchivePlan) -> dict:
return {
"id": plan.id,
"location_id": plan.location_id,
"token": plan.token,
"albums": json.loads(plan.albums) if plan.albums else None,
"state": plan.state,
"schema_version": plan.schema_version,
"asset_count": plan.asset_count,
"byte_size": plan.byte_size,
"version": plan.version,
"worker_id": plan.worker_id,
"completed_at": plan.completed_at.isoformat() if plan.completed_at else None,
"created_at": plan.created_at.isoformat() if plan.created_at else None,
}

View File

@@ -24,6 +24,7 @@ Preflight proves, per concept §9 "Archive preflight":
probed by writing them, not assumed.
Blocker codes: ``no_library_root``, ``location_offline``, ``wrong_volume``,
``archive_pending``,
``unsafe_destination``, ``destination_not_writable``, ``manifest_unwritable``,
``insufficient_capacity``, ``backup_unavailable``, ``lock_conflict``,
``rename_pending``, ``empty_scope``, ``destination_collision``,
@@ -57,6 +58,7 @@ from photo_pipeline.jobs.domain_handlers import ARCHIVE_LOCK, LIBRARY_WRITE_LOCK
from photo_pipeline.models import ArchiveLocation, Asset, UploadBatch, UploadItem
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.hashing import sha256_file
from photo_pipeline.services.jobs import JobService
from photo_pipeline.services.rename_journal import RenameJournal
@@ -291,6 +293,13 @@ class ArchiveService:
blockers.append(
_issue("rename_pending", "an unresolved rename must be recovered before archiving")
)
if ArchiveJournal(self._session_factory).blocks_mutation():
blockers.append(
_issue(
"archive_pending",
"an unresolved archive transfer must be recovered before archiving again",
)
)
return blockers
def _capacity(self, required: int, probe: dict) -> dict:

View File

@@ -83,12 +83,13 @@ def _now() -> datetime:
return datetime.now(timezone.utc)
def _maybe_fault(state: str) -> None:
def maybe_fault(state: str) -> None:
"""Test-only crash barrier (concept §18 fault injection).
When ``PHOTO_PIPELINE_FAULT_AFTER`` names a journal state, the process dies
abruptly the moment that state has been persisted — modelling a real kill at
exactly that transition. Never set outside tests.
exactly that transition. Never set outside tests. Shared with the archive
transfer journal (US06-02), which uses the same env var and its own state names.
"""
if os.environ.get("PHOTO_PIPELINE_FAULT_AFTER") == state:
os._exit(9)
@@ -172,7 +173,7 @@ class RenameApplyService:
# 1. Intent first — after this point a crash is recoverable from evidence.
self.journal.begin(operation["id"], worker_id=worker_id, fencing_token=token)
_maybe_fault(JournalState.MOVING)
maybe_fault(JournalState.MOVING)
# 2. Recheck preconditions immediately before mutating, never trusting the
# plan's snapshot: files can change between preview and confirmation.
@@ -187,21 +188,21 @@ class RenameApplyService:
os.rename(source, destination)
self.journal.transition(operation["id"], JournalState.MOVED, fencing_token=token)
_maybe_fault(JournalState.MOVED)
maybe_fault(JournalState.MOVED)
# 4. Database: stable IDs keep their identity, paths are re-pointed and the
# old occurrence is closed — all in one transaction.
self._reconcile_paths(operation, source, destination)
self.journal.transition(operation["id"], JournalState.DATABASE_UPDATED, fencing_token=token)
_maybe_fault(JournalState.DATABASE_UPDATED)
maybe_fault(JournalState.DATABASE_UPDATED)
# 5. Postconditions: the bytes really are at the new paths.
self._verify(operation, destination)
self.journal.transition(operation["id"], JournalState.VERIFIED, fencing_token=token)
_maybe_fault(JournalState.VERIFIED)
maybe_fault(JournalState.VERIFIED)
self.journal.transition(operation["id"], JournalState.COMPLETE, fencing_token=token)
_maybe_fault(JournalState.COMPLETE)
maybe_fault(JournalState.COMPLETE)
def _recheck(self, operation: dict, source: Path, destination: Path) -> None:
if not source.exists():

View File

@@ -0,0 +1,313 @@
"""Recovering interrupted archive transfers (US06-02).
The crash tests are real: a child process applies an archive plan and is killed by
the ``PHOTO_PIPELINE_FAULT_AFTER`` barrier at each persisted transition in turn —
including the moment immediately after an active source has been unlinked. The
parent then reopens the database and asserts the one invariant archiving exists to
uphold: **no verified file is ever lost, and no source is removed without a durable,
byte-identical archive copy.**
Both transfer paths are exercised: the same-filesystem atomic move and (with the
child forcing the device comparison) the cross-filesystem copy/verify/publish.
"""
import os
import subprocess
import sys
import uuid
from datetime import datetime, timezone
from pathlib import Path
import pytest
from sqlalchemy import select
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, UploadBatch, UploadItem
from photo_pipeline.services.archive_journal import FORWARD, MANUAL, RESUMABLE, ArchiveState
from photo_pipeline.services.archive_transfer import (
MANIFEST_NAME,
ArchiveTransferService,
read_manifest,
)
from photo_pipeline.services.archives import ArchiveService
from photo_pipeline.services.hashing import sha256_file
pytestmark = pytest.mark.phase_f
REPO = Path(__file__).resolve().parents[2]
NOW = datetime(2026, 1, 1, tzinfo=timezone.utc)
# Barriers, in the order the transfer persists them. ``source_removed`` is the one
# that matters most: the original is already gone at that point.
BARRIERS = [
ArchiveState.TRANSFERRING,
ArchiveState.VERIFIED,
ArchiveState.REMOVING,
"source_removed",
ArchiveState.COMPLETE,
]
# A child that applies the plan and dies at the configured barrier. ``force_copy``
# makes it take the cross-filesystem path without a second real volume.
APPLY_SCRIPT = """
import sys
sys.path.insert(0, {repo!r})
from photo_pipeline.config import Config
from photo_pipeline.db import create_db_engine, create_session_factory
from photo_pipeline.services import archive_transfer
db_url, data_dir, lib, plan_id, force_copy = sys.argv[1:6]
if force_copy == "1":
archive_transfer._same_filesystem = lambda *args: False
config = Config.from_env(
{{"PHOTO_PIPELINE_DATA_DIR": data_dir, "PHOTO_PIPELINE_LIBRARY_ROOTS": lib}}
)
sf = create_session_factory(create_db_engine(db_url))
archive_transfer.ArchiveTransferService(sf, config=config).apply(plan_id)
"""
# ── 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 _album(sf, lib, album="rome", names=("a.jpg", "b.jpg")):
folder = lib / album
folder.mkdir(parents=True, exist_ok=True)
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 name in names:
path = folder / name
path.write_bytes(f"{album}/{name} content".encode() * 8)
asset_id = str(uuid.uuid4())
session.add(
Asset(
id=asset_id,
original_path=str(path),
current_path=str(path),
discovered_at=NOW,
hash_version=1,
byte_size=path.stat().st_size,
current_sha256=sha256_file(path),
)
)
session.add(
UploadItem(
batch_id=batch_id,
asset_id=asset_id,
path=str(path),
sha256=sha256_file(path),
sha1="0" * 40,
state="sent",
outcome="uploaded",
)
)
session.commit()
return folder
def _plan(sf, config, archive, albums=None):
location = ArchiveService(sf, config=config).register("external", str(archive))
token = ArchiveService(sf, config=config).preflight(location["id"], albums)["token"]
service = ArchiveTransferService(sf, config=config)
return service, service.create(location["id"], albums, token=token)
def _crash_during_apply(config, tmp_path, lib, plan, barrier, *, force_copy=False):
script = tmp_path / f"apply_{barrier}.py"
script.write_text(APPLY_SCRIPT.format(repo=str(REPO)))
env = dict(os.environ)
env["PHOTO_PIPELINE_FAULT_AFTER"] = barrier
result = subprocess.run(
[
sys.executable,
str(script),
config.database_url,
str(tmp_path / "data"),
str(lib),
plan["id"],
"1" if force_copy else "0",
],
env=env,
capture_output=True,
)
assert result.returncode in (9, -9), (
f"child should have been killed at {barrier}, got {result.returncode}: "
f"{result.stderr.decode(errors='replace')[-400:]}"
)
def _contents(*roots):
return sorted(
path.read_bytes()
for root in roots
for path in root.rglob("*")
if path.is_file() and path.name != MANIFEST_NAME and not path.name.startswith(".")
)
def _reopen(config):
return create_session_factory(create_db_engine(config.database_url))
# ── fault injection at every persisted transition ────────────────────────────
@pytest.mark.parametrize("barrier", BARRIERS)
@pytest.mark.parametrize("force_copy", [False, True], ids=["move", "copy"])
def test_a_crash_at_every_transition_loses_nothing_and_recovers(tmp_path, barrier, force_copy):
config, sf, lib, archive = _env(tmp_path)
_album(sf, lib)
service, plan = _plan(sf, config, archive)
before = _contents(lib)
_crash_during_apply(config, tmp_path, lib, plan, barrier, force_copy=force_copy)
# Every file still exists somewhere: the crash may not have cost a single byte.
assert set(before) <= set(_contents(lib, archive)), f"content lost at {barrier}"
reopened = _reopen(config)
recovery = ArchiveTransferService(reopened, config=config)
# Whatever the crash interrupted is recoverable from evidence — never ambiguous.
# (A crash on the last barrier can land on a terminal item and leave none.)
verdicts = recovery.journal.classify_all()
assert all(v["classification"] in {RESUMABLE, FORWARD} for v in verdicts), verdicts
recovery.recover()
recovery.apply(plan["id"]) # finish whatever the crash never started
assert _contents(archive) == before
assert _contents(lib) == []
assert recovery.journal.blocks_mutation() is False
assert recovery.journal.plan_state(plan["id"]) == "complete"
@pytest.mark.parametrize("barrier", [ArchiveState.VERIFIED, "source_removed"])
def test_recovery_is_idempotent_across_repeated_restarts(tmp_path, barrier):
config, sf, lib, archive = _env(tmp_path)
_album(sf, lib)
service, plan = _plan(sf, config, archive)
_crash_during_apply(config, tmp_path, lib, plan, barrier)
reopened = _reopen(config)
recovery = ArchiveTransferService(reopened, config=config)
recovery.recover()
settled = (_contents(lib), _contents(archive), read_manifest(archive / "rome"))
for _ in range(2):
recovery.recover()
assert (_contents(lib), _contents(archive), read_manifest(archive / "rome")) == settled
def test_a_crash_before_anything_was_published_leaves_the_source_intact(tmp_path):
config, sf, lib, archive = _env(tmp_path)
folder = _album(sf, lib, names=("a.jpg",))
service, plan = _plan(sf, config, archive)
_crash_during_apply(config, tmp_path, lib, plan, ArchiveState.TRANSFERRING)
recovery = ArchiveTransferService(_reopen(config), config=config)
verdict = recovery.journal.classify_all()[0]
assert verdict["classification"] == RESUMABLE
assert (folder / "a.jpg").exists()
assert not (archive / "rome" / "a.jpg").exists()
assert recovery.recover() == {"resumed": 1, "completed": 0, "manual": 0}
assert recovery.journal.operations(plan["id"])[0]["journal_state"] == ArchiveState.PLANNED
def test_a_crash_after_the_source_was_removed_finishes_the_bookkeeping(tmp_path):
"""The dangerous window: the original is gone and the database still points at
it. Recovery must complete the record, never re-transfer or report loss."""
config, sf, lib, archive = _env(tmp_path)
folder = _album(sf, lib, names=("a.jpg",))
service, plan = _plan(sf, config, archive)
expected = service.journal.operations(plan["id"])[0]["expected_sha256"]
_crash_during_apply(config, tmp_path, lib, plan, "source_removed")
reopened = _reopen(config)
recovery = ArchiveTransferService(reopened, config=config)
verdict = recovery.journal.classify_all()[0]
assert verdict["classification"] == FORWARD
assert not (folder / "a.jpg").exists()
assert sha256_file(archive / "rome" / "a.jpg") == expected
with reopened() as session:
stranded = session.scalar(select(Asset))
assert stranded.current_path is not None, "the crash happened before the DB update"
assert recovery.recover() == {"resumed": 0, "completed": 1, "manual": 0}
with reopened() as session:
asset = session.scalar(select(Asset))
assert asset.current_path is None
assert asset.availability_state == "archived_online"
assert asset.archive_path == "rome/a.jpg"
assert recovery.journal.operations(plan["id"])[0]["journal_state"] == ArchiveState.COMPLETE
assert len(read_manifest(archive / "rome")) == 1, "the manifest is not duplicated"
def test_an_archive_copy_that_changed_after_the_crash_blocks_for_a_human(tmp_path):
"""Wrong bytes at the destination can never justify deleting the original, and
are never overwritten either."""
config, sf, lib, archive = _env(tmp_path)
folder = _album(sf, lib, names=("a.jpg",))
service, plan = _plan(sf, config, archive)
_crash_during_apply(config, tmp_path, lib, plan, ArchiveState.VERIFIED, force_copy=True)
(archive / "rome" / "a.jpg").write_bytes(b"tampered with while the app was down")
recovery = ArchiveTransferService(_reopen(config), config=config)
verdict = recovery.journal.classify_all()[0]
assert verdict["classification"] == MANUAL
assert recovery.recover() == {"resumed": 0, "completed": 0, "manual": 1}
assert (folder / "a.jpg").exists(), "the source is kept while the archive is unproven"
assert (archive / "rome" / "a.jpg").read_bytes() == b"tampered with while the app was down"
# And the unresolved item keeps blocking further archiving.
assert recovery.journal.blocks_mutation() is True
report = ArchiveService(_reopen(config), config=config).preflight(plan["location_id"])
assert "archive_pending" in {issue["code"] for issue in report["blockers"]}
def test_a_verified_archive_whose_copy_vanished_is_never_reported_as_archived(tmp_path):
config, sf, lib, archive = _env(tmp_path)
folder = _album(sf, lib, names=("a.jpg",))
service, plan = _plan(sf, config, archive)
_crash_during_apply(config, tmp_path, lib, plan, ArchiveState.REMOVING, force_copy=True)
(archive / "rome" / "a.jpg").unlink() # the medium lost it
recovery = ArchiveTransferService(_reopen(config), config=config)
assert recovery.journal.classify_all()[0]["classification"] == MANUAL
assert recovery.recover()["manual"] == 1
assert (folder / "a.jpg").exists()
assert recovery.journal.operations(plan["id"])[0]["journal_state"] != ArchiveState.COMPLETE

View File

@@ -0,0 +1,480 @@
"""Transferring, verifying, and removing active sources (US06-02).
Every case here asks the same question the service exists to answer: could an
original leave active storage without a durable, byte-identical archive copy? The
files are real, the hashes are real, and each failure path asserts that the source
is still exactly where it was.
Cross-filesystem behaviour is forced by patching the device comparison rather than
by requiring a second real filesystem in CI — the copy/verify/publish code that runs
is the production one.
"""
import json
import uuid
from datetime import datetime, timezone
import pytest
from fastapi.testclient import TestClient
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, UploadBatch, UploadItem
from photo_pipeline.services import archive_transfer as transfer_module
from photo_pipeline.services.archive_journal import ArchiveState
from photo_pipeline.services.archive_transfer import (
MANIFEST_NAME,
ArchiveTransferService,
read_manifest,
)
from photo_pipeline.services.archives import ArchiveError, ArchiveService
from photo_pipeline.services.hashing import sha256_file
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 _album(sf, lib, album="rome", names=("a.jpg", "b.jpg")):
"""A real album whose assets carry the verified upload evidence archiving needs."""
folder = lib / album
folder.mkdir(parents=True, exist_ok=True)
ids = []
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 name in names:
path = folder / name
path.write_bytes(f"{album}/{name} bytes".encode() * 8)
asset_id = str(uuid.uuid4())
ids.append(asset_id)
session.add(
Asset(
id=asset_id,
original_path=str(path),
current_path=str(path),
discovered_at=NOW,
hash_version=1,
byte_size=path.stat().st_size,
current_sha256=sha256_file(path),
)
)
session.add(AssetPath(asset_id=asset_id, path=str(path), valid_from=NOW))
session.add(
UploadItem(
batch_id=batch_id,
asset_id=asset_id,
path=str(path),
sha256=sha256_file(path),
sha1="0" * 40,
state="sent",
outcome="uploaded",
)
)
session.commit()
return folder, ids
def _location(sf, config, archive, name="external"):
return ArchiveService(sf, config=config).register(name, str(archive))
def _plan(sf, config, location_id, albums=None):
service = ArchiveTransferService(sf, config=config)
token = ArchiveService(sf, config=config).preflight(location_id, albums)["token"]
return service, service.create(location_id, albums, token=token)
def _contents(*roots):
"""Every file's bytes under the given roots, ignoring our own bookkeeping."""
return sorted(
path.read_bytes()
for root in roots
for path in root.rglob("*")
if path.is_file() and path.name != MANIFEST_NAME and not path.name.startswith(".")
)
def _states(service, plan_id):
return [row["journal_state"] for row in service.journal.operations(plan_id)]
def _asset(sf, asset_id):
with sf() as session:
return session.get(Asset, asset_id)
# ── planning ─────────────────────────────────────────────────────────────────
def test_a_plan_records_every_file_with_its_expected_hash_and_moves_nothing(tmp_path):
config, sf, lib, archive = _env(tmp_path)
folder, ids = _album(sf, lib)
location = _location(sf, config, archive)
before = _contents(lib)
service, plan = _plan(sf, config, location["id"])
assert plan["state"] == "planned" and plan["asset_count"] == 2
assert sorted(op["asset_id"] for op in plan["operations"]) == sorted(ids)
for operation in plan["operations"]:
source = folder / operation["archive_path"].split("/")[-1]
assert operation["expected_sha256"] == sha256_file(source)
assert operation["archive_path"].startswith("rome/")
assert operation["journal_state"] == ArchiveState.PLANNED
assert _contents(lib) == before
assert _contents(archive) == []
def test_a_stale_token_cannot_create_a_plan(tmp_path):
config, sf, lib, archive = _env(tmp_path)
folder, _ = _album(sf, lib)
location = _location(sf, config, archive)
token = ArchiveService(sf, config=config).preflight(location["id"])["token"]
(folder / "a.jpg").write_bytes(b"edited after approval")
with pytest.raises(ArchiveError) as error:
ArchiveTransferService(sf, config=config).create(location["id"], token=token)
assert error.value.code == "stale_token"
def test_a_blocked_scope_cannot_create_a_plan(tmp_path):
"""No verified upload means Immich may not hold these bytes; archiving would
remove the only copy."""
config, sf, lib, archive = _env(tmp_path)
folder = lib / "rome"
folder.mkdir()
path = folder / "a.jpg"
path.write_bytes(b"never uploaded")
with sf() as session:
session.add(
Asset(
id=str(uuid.uuid4()),
original_path=str(path),
current_path=str(path),
discovered_at=NOW,
hash_version=1,
byte_size=path.stat().st_size,
current_sha256=sha256_file(path),
)
)
session.commit()
location = _location(sf, config, archive)
token = ArchiveService(sf, config=config).preflight(location["id"])["token"]
with pytest.raises(ArchiveError) as error:
ArchiveTransferService(sf, config=config).create(location["id"], token=token)
assert error.value.code == "blocked"
assert path.exists()
# ── the transfer ─────────────────────────────────────────────────────────────
@pytest.mark.parametrize("same_filesystem", [True, False])
def test_archiving_verifies_the_copy_before_the_source_is_removed(
tmp_path, monkeypatch, same_filesystem
):
config, sf, lib, archive = _env(tmp_path)
folder, ids = _album(sf, lib)
location = _location(sf, config, archive)
service, plan = _plan(sf, config, location["id"])
expected = {op["asset_id"]: op["expected_sha256"] for op in plan["operations"]}
before = _contents(lib)
if not same_filesystem:
# Force the copy-verify-publish path without needing a second real volume.
monkeypatch.setattr(transfer_module, "_same_filesystem", lambda *_: False)
result = service.apply(plan["id"])
assert result == {
"plan_id": plan["id"],
"archived": 2,
"failed": 0,
"skipped": 0,
"state": "complete",
}
# The bytes moved: nothing is left in the library, everything is in the archive.
assert _contents(archive) == before
assert _contents(lib) == []
assert not folder.exists(), "an emptied album folder is not left behind"
for asset_id, digest in expected.items():
archived = archive / "rome" / _asset(sf, asset_id).archive_path.split("/")[-1]
assert sha256_file(archived) == digest
assert _states(service, plan["id"]) == [ArchiveState.COMPLETE] * 2
# No transfer temporaries survive either path.
assert not list((archive / "rome").glob(".archive-*"))
@pytest.mark.parametrize("same_filesystem", [True, False])
def test_the_manifest_on_the_medium_matches_the_archived_bytes(
tmp_path, monkeypatch, same_filesystem
):
"""The manifest is the medium's own record: it must be usable to verify the
archive with no database at all."""
config, sf, lib, archive = _env(tmp_path)
_album(sf, lib)
location = _location(sf, config, archive)
service, plan = _plan(sf, config, location["id"])
if not same_filesystem:
monkeypatch.setattr(transfer_module, "_same_filesystem", lambda *_: False)
service.apply(plan["id"])
entries = read_manifest(archive / "rome")
assert len(entries) == 2
for entry in entries:
archived = archive / entry["archive_path"]
assert sha256_file(archived) == entry["sha256"]
assert entry["byte_size"] == archived.stat().st_size
assert entry["media_id"] == location["media_id"]
assert entry["plan_id"] == plan["id"] and entry["album"] == "rome"
def test_archived_assets_keep_their_identity_and_gain_their_new_location(tmp_path):
config, sf, lib, archive = _env(tmp_path)
folder, ids = _album(sf, lib)
location = _location(sf, config, archive)
service, plan = _plan(sf, config, location["id"])
service.apply(plan["id"])
with sf() as session:
for asset_id in ids:
asset = session.get(Asset, asset_id)
assert asset is not None, "archiving never deletes the record"
assert asset.current_path is None
assert asset.availability_state == "archived_online"
assert asset.archive_location_id == location["id"]
assert (archive / asset.archive_path).exists()
occurrences = session.scalars(
select(AssetPath).where(AssetPath.asset_id == asset_id)
).all()
active = [row for row in occurrences if row.valid_until is None]
assert [row.path for row in active] == [str(archive / asset.archive_path)]
closed = [row for row in occurrences if row.valid_until is not None]
assert [row.path for row in closed] == [str(folder / asset.archive_path.split("/")[-1])]
def test_applying_the_same_plan_again_changes_nothing(tmp_path):
config, sf, lib, archive = _env(tmp_path)
_album(sf, lib)
location = _location(sf, config, archive)
service, plan = _plan(sf, config, location["id"])
service.apply(plan["id"])
archived = _contents(archive)
manifest = read_manifest(archive / "rome")
again = service.apply(plan["id"])
assert again["skipped"] == 2 and again["archived"] == 0 and again["failed"] == 0
assert _contents(archive) == archived
assert read_manifest(archive / "rome") == manifest
# ── unexpected changes stop the item ─────────────────────────────────────────
def test_an_occupied_destination_is_never_overwritten(tmp_path):
config, sf, lib, archive = _env(tmp_path)
folder, _ = _album(sf, lib, names=("a.jpg",))
location = _location(sf, config, archive)
service, plan = _plan(sf, config, location["id"])
(archive / "rome").mkdir()
(archive / "rome" / "a.jpg").write_bytes(b"someone else's file")
result = service.apply(plan["id"])
assert result["failed"] == 1 and result["state"] == "failed"
assert (archive / "rome" / "a.jpg").read_bytes() == b"someone else's file"
assert (folder / "a.jpg").exists(), "the source must survive a refused transfer"
assert service.journal.operations(plan["id"])[0]["error_code"] == "destination_exists"
def test_a_source_edited_after_planning_is_neither_archived_nor_removed(tmp_path):
config, sf, lib, archive = _env(tmp_path)
folder, _ = _album(sf, lib, names=("a.jpg",))
location = _location(sf, config, archive)
service, plan = _plan(sf, config, location["id"])
(folder / "a.jpg").write_bytes(b"edited between approval and apply")
result = service.apply(plan["id"])
assert result["failed"] == 1
assert (folder / "a.jpg").read_bytes() == b"edited between approval and apply"
assert not (archive / "rome" / "a.jpg").exists()
assert service.journal.operations(plan["id"])[0]["error_code"] == "source_changed"
def test_a_copy_that_lands_with_the_wrong_bytes_is_not_published(tmp_path, monkeypatch):
"""The read-back hash — not the fact that a write returned — is what proves the
archive copy."""
config, sf, lib, archive = _env(tmp_path)
folder, _ = _album(sf, lib, names=("a.jpg",))
location = _location(sf, config, archive)
service, plan = _plan(sf, config, location["id"])
monkeypatch.setattr(transfer_module, "_same_filesystem", lambda *_: False)
monkeypatch.setattr(
transfer_module.shutil,
"copyfileobj",
lambda src, dst, length=0: dst.write(b"corrupted in flight"),
)
result = service.apply(plan["id"])
assert result["failed"] == 1
assert service.journal.operations(plan["id"])[0]["error_code"] == "copy_mismatch"
assert (folder / "a.jpg").exists()
assert not (archive / "rome" / "a.jpg").exists()
assert not list((archive / "rome").glob(".archive-*")), "the failed copy is cleaned up"
def test_a_source_replaced_by_a_symlink_is_refused(tmp_path):
config, sf, lib, archive = _env(tmp_path)
folder, _ = _album(sf, lib, names=("a.jpg",))
location = _location(sf, config, archive)
service, plan = _plan(sf, config, location["id"])
elsewhere = tmp_path / "elsewhere.jpg"
elsewhere.write_bytes(b"not a library file")
(folder / "a.jpg").unlink()
(folder / "a.jpg").symlink_to(elsewhere)
result = service.apply(plan["id"])
assert result["failed"] == 1
assert elsewhere.exists() and (folder / "a.jpg").is_symlink()
assert not (archive / "rome" / "a.jpg").exists()
def test_one_failed_file_does_not_stop_the_rest_of_the_album(tmp_path):
config, sf, lib, archive = _env(tmp_path)
folder, _ = _album(sf, lib, names=("a.jpg", "b.jpg"))
location = _location(sf, config, archive)
service, plan = _plan(sf, config, location["id"])
(folder / "a.jpg").write_bytes(b"edited after approval")
result = service.apply(plan["id"])
assert (result["archived"], result["failed"]) == (1, 1)
assert (folder / "a.jpg").exists() and not (folder / "b.jpg").exists()
assert (archive / "rome" / "b.jpg").exists()
assert folder.exists(), "a folder that still holds an unarchived file stays"
def test_an_unresolved_transfer_blocks_the_next_preflight(tmp_path, monkeypatch):
config, sf, lib, archive = _env(tmp_path)
_album(sf, lib, "rome", names=("a.jpg",))
_album(sf, lib, "paris", names=("c.jpg",))
location = _location(sf, config, archive)
service, plan = _plan(sf, config, location["id"], ["rome"])
# Leave one item stuck mid-transfer, as a killed process would.
service.journal.begin(
service.journal.operations(plan["id"])[0]["id"], worker_id="w", fencing_token=1
)
report = ArchiveService(sf, config=config).preflight(location["id"], ["paris"])
assert report["state"] == "blocked"
assert "archive_pending" in {issue["code"] for issue in report["blockers"]}
# ── API surface ──────────────────────────────────────────────────────────────
def test_api_creates_a_plan_and_queues_it_on_the_archiver_lane(tmp_path):
config, sf, lib, archive = _env(tmp_path)
_album(sf, lib)
with TestClient(create_app(config)) as client:
location = client.post(
"/api/v1/archive-locations", json={"name": "external", "root": str(archive)}
).json()
token = client.post(
"/api/v1/archive-preflight", json={"location_id": location["id"]}
).json()["token"]
created = client.post(
"/api/v1/archive-plans", json={"location_id": location["id"], "token": token}
)
plan_id = created.json()["id"]
fetched = client.get(f"/api/v1/archive-plans/{plan_id}")
applied = client.post(f"/api/v1/archive-plans/{plan_id}/apply")
listed = client.get("/api/v1/archive-plans")
assert created.status_code == 201 and created.json()["asset_count"] == 2
assert fetched.status_code == 200 and len(fetched.json()["operations"]) == 2
assert applied.status_code == 200 and applied.json()["job"]["state"] == "queued"
assert applied.json()["job"]["lock_key"] == "archive"
assert [row["id"] for row in listed.json()["plans"]] == [plan_id]
# Queuing alone must not have touched a single file.
assert _contents(lib) and _contents(archive) == []
def test_api_refuses_a_stale_token_and_an_unknown_plan(tmp_path):
config, sf, lib, archive = _env(tmp_path)
_album(sf, lib)
location = _location(sf, config, archive)
with TestClient(create_app(config)) as client:
stale = client.post(
"/api/v1/archive-plans",
json={"location_id": location["id"], "token": "v1:not-the-real-token"},
)
unknown = client.get("/api/v1/archive-plans/nope")
unknown_apply = client.post("/api/v1/archive-plans/nope/apply")
assert stale.status_code == 409 and stale.json()["error"]["code"] == "stale_token"
assert unknown.status_code == 404 and unknown_apply.status_code == 404
def test_api_reports_recovery_state_for_an_interrupted_transfer(tmp_path):
config, sf, lib, archive = _env(tmp_path)
_album(sf, lib, names=("a.jpg",))
location = _location(sf, config, archive)
service, plan = _plan(sf, config, location["id"])
service.journal.begin(
service.journal.operations(plan["id"])[0]["id"], worker_id="w", fencing_token=1
)
with TestClient(create_app(config)) as client:
status = client.get("/api/v1/archive-recovery").json()
resolved = client.post("/api/v1/archive-recovery/resolve").json()
assert status["blocks_mutation"] is True
assert status["operations"][0]["classification"] == "resumable"
assert resolved == {"resumed": 1, "completed": 0, "manual": 0}
assert json.loads(json.dumps(resolved)) # plain JSON, nothing exotic

View File

@@ -126,6 +126,11 @@
],
"US06-01": [
"tests/integration/test_archive_preflight.py"
],
"US06-02": [
"tests/unit/test_archive_journal_states.py",
"tests/integration/test_archive_transfer.py",
"tests/integration/test_archive_recovery.py"
]
}
}

View File

@@ -0,0 +1,105 @@
"""Archive journal state machine and evidence table (US06-02).
The evidence table decides whether an original may be deleted, so every
combination of journal state and disk reality is asserted here as a pure function —
no database, no files. A wrong cell in this table is data loss.
"""
import pytest
from photo_pipeline.services.archive_journal import (
ALLOWED_TRANSITIONS,
FORWARD,
MANUAL,
RESUMABLE,
TERMINAL_STATES,
UNSAFE_STATES,
ArchiveState,
_classify,
can_transition,
)
pytestmark = pytest.mark.phase_f
# ── state machine ────────────────────────────────────────────────────────────
def test_the_happy_path_is_the_only_way_forward():
assert can_transition(ArchiveState.PLANNED, ArchiveState.TRANSFERRING)
assert can_transition(ArchiveState.TRANSFERRING, ArchiveState.VERIFIED)
assert can_transition(ArchiveState.VERIFIED, ArchiveState.REMOVING)
assert can_transition(ArchiveState.REMOVING, ArchiveState.COMPLETE)
# No shortcut may skip verification before a source is removed.
assert not can_transition(ArchiveState.TRANSFERRING, ArchiveState.REMOVING)
assert not can_transition(ArchiveState.PLANNED, ArchiveState.VERIFIED)
assert not can_transition(ArchiveState.VERIFIED, ArchiveState.COMPLETE)
def test_removal_never_goes_backwards():
"""Once the source may be gone, retrying the transfer would archive nothing and
could overwrite the copy that is now the only one."""
assert ALLOWED_TRANSITIONS[ArchiveState.REMOVING] == {
ArchiveState.COMPLETE,
ArchiveState.FAILED,
}
assert not can_transition(ArchiveState.REMOVING, ArchiveState.TRANSFERRING)
assert not can_transition(ArchiveState.REMOVING, ArchiveState.PLANNED)
def test_complete_is_terminal_and_failed_can_be_retried():
assert ALLOWED_TRANSITIONS[ArchiveState.COMPLETE] == set()
assert TERMINAL_STATES == {ArchiveState.COMPLETE}
assert can_transition(ArchiveState.FAILED, ArchiveState.TRANSFERRING)
assert can_transition(ArchiveState.FAILED, ArchiveState.PLANNED)
def test_every_state_that_can_touch_the_disk_is_marked_unsafe():
assert UNSAFE_STATES == {
ArchiveState.TRANSFERRING,
ArchiveState.VERIFIED,
ArchiveState.REMOVING,
}
assert ArchiveState.PLANNED not in UNSAFE_STATES
# ── evidence table ───────────────────────────────────────────────────────────
@pytest.mark.parametrize(
"state,source,destination,matches,expected",
[
# Nothing published yet: the source is still the only copy.
(ArchiveState.TRANSFERRING, True, False, False, RESUMABLE),
(ArchiveState.FAILED, True, False, False, RESUMABLE),
# The archive copy is durable and correct: finish the remaining steps.
(ArchiveState.TRANSFERRING, True, True, True, FORWARD),
(ArchiveState.TRANSFERRING, False, True, True, FORWARD),
(ArchiveState.VERIFIED, True, True, True, FORWARD),
(ArchiveState.REMOVING, False, True, True, FORWARD),
(ArchiveState.FAILED, True, True, True, FORWARD),
# Wrong bytes at the destination: never overwrite, never remove.
(ArchiveState.TRANSFERRING, True, True, False, MANUAL),
(ArchiveState.VERIFIED, True, True, False, MANUAL),
(ArchiveState.REMOVING, False, True, False, MANUAL),
# The journal claims an archived copy that is not there.
(ArchiveState.VERIFIED, True, False, False, MANUAL),
(ArchiveState.REMOVING, False, False, False, MANUAL),
# Neither copy exists — never silently accepted as success.
(ArchiveState.TRANSFERRING, False, False, False, MANUAL),
(ArchiveState.FAILED, False, False, False, MANUAL),
],
)
def test_classification_of_every_evidence_combination(
state, source, destination, matches, expected
):
classification, reason = _classify(state, source, destination, matches)
assert classification == expected, reason
def test_a_source_that_is_gone_without_an_archive_copy_is_never_called_recoverable():
"""The one combination that must always stop: the original left active storage
and nothing verifiable took its place."""
for state in (ArchiveState.TRANSFERRING, ArchiveState.VERIFIED, ArchiveState.REMOVING):
assert _classify(state, False, False, False)[0] == MANUAL