Compare commits
1 Commits
main
...
us/US06-02
| Author | SHA1 | Date | |
|---|---|---|---|
| 9503fd1cfc |
96
migrations/versions/0012_archive_plans.py
Normal file
96
migrations/versions/0012_archive_plans.py
Normal 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")
|
||||||
@@ -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
|
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
|
||||||
the token a later archive plan must present (US06-02). Neither endpoint moves or
|
the token an archive plan must present. Creating a plan writes only database rows —
|
||||||
removes a single library file.
|
the transfer itself runs on the durable ``archive`` lane, never in the request
|
||||||
|
thread, because it removes originals.
|
||||||
"""
|
"""
|
||||||
|
|
||||||
from __future__ import annotations
|
from __future__ import annotations
|
||||||
@@ -12,12 +13,17 @@ 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.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.jobs import JobBlocked, JobService
|
||||||
|
|
||||||
router = APIRouter(tags=["archives"])
|
router = APIRouter(tags=["archives"])
|
||||||
|
|
||||||
# Which failures are the caller's request (422) and which are a missing thing (404).
|
# Which failures are the caller's request (422), a missing thing (404), or state
|
||||||
NOT_FOUND_CODES = {"unknown_location"}
|
# that changed under the caller (409).
|
||||||
|
NOT_FOUND_CODES = {"unknown_location", "unknown_plan"}
|
||||||
|
CONFLICT_CODES = {"stale_token", "stale_plan", "archive_pending"}
|
||||||
|
|
||||||
|
|
||||||
class RegisterLocationRequest(BaseModel):
|
class RegisterLocationRequest(BaseModel):
|
||||||
@@ -31,12 +37,28 @@ class PreflightRequest(BaseModel):
|
|||||||
albums: list[str] | None = None
|
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:
|
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 _transfers(request: Request) -> ArchiveTransferService:
|
||||||
|
return ArchiveTransferService(
|
||||||
|
request.app.state.session_factory, config=request.app.state.config
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
def _error(error: ArchiveError) -> JSONResponse:
|
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(
|
return JSONResponse(
|
||||||
status_code=status, content={"error": {"code": error.code, "message": str(error)}}
|
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)
|
return _service(request).preflight(body.location_id, body.albums)
|
||||||
except ArchiveError as error:
|
except ArchiveError as error:
|
||||||
return _error(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()
|
||||||
|
|||||||
@@ -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
|
Importing this module registers the ``safety_score``, ``analysis``,
|
||||||
``upload_batch`` job types so the generic worker can run them per item. Each handler
|
``upload_batch``, and ``archive_plan`` job types so the generic worker can run them
|
||||||
delegates to its service, which owns the real work and the privacy gate. Handlers
|
per item. Each handler delegates to its service, which owns the real work and the
|
||||||
are idempotent: re-scoring or re-analyzing one asset is safe after an interrupted
|
privacy gate. Handlers are idempotent: re-scoring or re-analyzing one asset is safe
|
||||||
attempt, and an upload batch refuses to re-run an attempt whose outcome is unknown.
|
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);
|
Providers/models are the service defaults here (real NsfwModel / vision provider);
|
||||||
tests exercise the services directly with injected fakes rather than the worker.
|
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"
|
SAFETY_SCORE = "safety_score"
|
||||||
ANALYSIS = "analysis"
|
ANALYSIS = "analysis"
|
||||||
UPLOAD_BATCH = "upload_batch"
|
UPLOAD_BATCH = "upload_batch"
|
||||||
|
ARCHIVE_PLAN = "archive_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).
|
||||||
UPLOAD_LOCK = "upload"
|
UPLOAD_LOCK = "upload"
|
||||||
# The archiver lane: one archive/restore plan at a time (concept §16). No handler
|
# The archiver lane: one archive/restore plan at a time (concept §16).
|
||||||
# runs on it yet (US06-02); preflight already refuses to plan around a held lease.
|
|
||||||
ARCHIVE_LOCK = "archive"
|
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']}")
|
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(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)
|
||||||
|
|||||||
@@ -5,7 +5,7 @@ Alembic environment relies on.
|
|||||||
"""
|
"""
|
||||||
|
|
||||||
from photo_pipeline.models.albums import AlbumProposal
|
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.assets import Asset, AssetPath
|
||||||
from photo_pipeline.models.duplicates import (
|
from photo_pipeline.models.duplicates import (
|
||||||
DuplicateCluster,
|
DuplicateCluster,
|
||||||
@@ -21,6 +21,8 @@ from photo_pipeline.models.workflow import AnalysisResult, SafetyReview
|
|||||||
__all__ = [
|
__all__ = [
|
||||||
"AlbumProposal",
|
"AlbumProposal",
|
||||||
"ArchiveLocation",
|
"ArchiveLocation",
|
||||||
|
"ArchiveOperation",
|
||||||
|
"ArchivePlan",
|
||||||
"Asset",
|
"Asset",
|
||||||
"AssetPath",
|
"AssetPath",
|
||||||
"DuplicateCluster",
|
"DuplicateCluster",
|
||||||
|
|||||||
@@ -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
|
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
|
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
|
``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
|
locations without touching a sleeping disk. Preflight always re-probes — a stored
|
||||||
state is a hint, never evidence.
|
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 __future__ import annotations
|
||||||
|
|
||||||
from datetime import datetime
|
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 sqlalchemy.orm import Mapped, mapped_column
|
||||||
|
|
||||||
from photo_pipeline.db import Base
|
from photo_pipeline.db import Base
|
||||||
@@ -40,3 +54,69 @@ class ArchiveLocation(Base):
|
|||||||
updated_at: Mapped[datetime] = mapped_column(
|
updated_at: Mapped[datetime] = mapped_column(
|
||||||
DateTime(timezone=True), nullable=False, server_default=func.now(), onupdate=func.now()
|
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()
|
||||||
|
)
|
||||||
|
|||||||
@@ -35,7 +35,14 @@ class Asset(Base):
|
|||||||
|
|
||||||
discovered_at: Mapped[datetime] = mapped_column(DateTime(timezone=True), nullable=False)
|
discovered_at: Mapped[datetime] = mapped_column(DateTime(timezone=True), nullable=False)
|
||||||
missing_at: Mapped[datetime | None] = mapped_column(DateTime(timezone=True))
|
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")
|
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.
|
# 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"))
|
||||||
|
|
||||||
|
|||||||
337
photo_pipeline/services/archive_journal.py
Normal file
337
photo_pipeline/services/archive_journal.py
Normal 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,
|
||||||
|
}
|
||||||
605
photo_pipeline/services/archive_transfer.py
Normal file
605
photo_pipeline/services/archive_transfer.py
Normal 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,
|
||||||
|
}
|
||||||
@@ -24,6 +24,7 @@ Preflight proves, per concept §9 "Archive preflight":
|
|||||||
probed by writing them, not assumed.
|
probed by writing them, not assumed.
|
||||||
|
|
||||||
Blocker codes: ``no_library_root``, ``location_offline``, ``wrong_volume``,
|
Blocker codes: ``no_library_root``, ``location_offline``, ``wrong_volume``,
|
||||||
|
``archive_pending``,
|
||||||
``unsafe_destination``, ``destination_not_writable``, ``manifest_unwritable``,
|
``unsafe_destination``, ``destination_not_writable``, ``manifest_unwritable``,
|
||||||
``insufficient_capacity``, ``backup_unavailable``, ``lock_conflict``,
|
``insufficient_capacity``, ``backup_unavailable``, ``lock_conflict``,
|
||||||
``rename_pending``, ``empty_scope``, ``destination_collision``,
|
``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.models import ArchiveLocation, Asset, UploadBatch, UploadItem
|
||||||
from photo_pipeline.path_policy import PathPolicyError, is_excluded, normalize_root, resolve_within
|
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.albums import album_label
|
||||||
|
from photo_pipeline.services.archive_journal import ArchiveJournal
|
||||||
from photo_pipeline.services.hashing import sha256_file
|
from photo_pipeline.services.hashing import sha256_file
|
||||||
from photo_pipeline.services.jobs import JobService
|
from photo_pipeline.services.jobs import JobService
|
||||||
from photo_pipeline.services.rename_journal import RenameJournal
|
from photo_pipeline.services.rename_journal import RenameJournal
|
||||||
@@ -291,6 +293,13 @@ class ArchiveService:
|
|||||||
blockers.append(
|
blockers.append(
|
||||||
_issue("rename_pending", "an unresolved rename must be recovered before archiving")
|
_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
|
return blockers
|
||||||
|
|
||||||
def _capacity(self, required: int, probe: dict) -> dict:
|
def _capacity(self, required: int, probe: dict) -> dict:
|
||||||
|
|||||||
@@ -83,12 +83,13 @@ def _now() -> datetime:
|
|||||||
return datetime.now(timezone.utc)
|
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).
|
"""Test-only crash barrier (concept §18 fault injection).
|
||||||
|
|
||||||
When ``PHOTO_PIPELINE_FAULT_AFTER`` names a journal state, the process dies
|
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
|
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:
|
if os.environ.get("PHOTO_PIPELINE_FAULT_AFTER") == state:
|
||||||
os._exit(9)
|
os._exit(9)
|
||||||
@@ -172,7 +173,7 @@ class RenameApplyService:
|
|||||||
|
|
||||||
# 1. Intent first — after this point a crash is recoverable from evidence.
|
# 1. Intent first — after this point a crash is recoverable from evidence.
|
||||||
self.journal.begin(operation["id"], worker_id=worker_id, fencing_token=token)
|
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
|
# 2. Recheck preconditions immediately before mutating, never trusting the
|
||||||
# plan's snapshot: files can change between preview and confirmation.
|
# plan's snapshot: files can change between preview and confirmation.
|
||||||
@@ -187,21 +188,21 @@ class RenameApplyService:
|
|||||||
os.rename(source, destination)
|
os.rename(source, destination)
|
||||||
|
|
||||||
self.journal.transition(operation["id"], JournalState.MOVED, fencing_token=token)
|
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
|
# 4. Database: stable IDs keep their identity, paths are re-pointed and the
|
||||||
# old occurrence is closed — all in one transaction.
|
# old occurrence is closed — all in one transaction.
|
||||||
self._reconcile_paths(operation, source, destination)
|
self._reconcile_paths(operation, source, destination)
|
||||||
self.journal.transition(operation["id"], JournalState.DATABASE_UPDATED, fencing_token=token)
|
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.
|
# 5. Postconditions: the bytes really are at the new paths.
|
||||||
self._verify(operation, destination)
|
self._verify(operation, destination)
|
||||||
self.journal.transition(operation["id"], JournalState.VERIFIED, fencing_token=token)
|
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)
|
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:
|
def _recheck(self, operation: dict, source: Path, destination: Path) -> None:
|
||||||
if not source.exists():
|
if not source.exists():
|
||||||
|
|||||||
313
tests/integration/test_archive_recovery.py
Normal file
313
tests/integration/test_archive_recovery.py
Normal 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
|
||||||
480
tests/integration/test_archive_transfer.py
Normal file
480
tests/integration/test_archive_transfer.py
Normal 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
|
||||||
@@ -126,6 +126,11 @@
|
|||||||
],
|
],
|
||||||
"US06-01": [
|
"US06-01": [
|
||||||
"tests/integration/test_archive_preflight.py"
|
"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"
|
||||||
]
|
]
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
105
tests/unit/test_archive_journal_states.py
Normal file
105
tests/unit/test_archive_journal_states.py
Normal 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
|
||||||
Reference in New Issue
Block a user