Compare commits

...

1 Commits

Author SHA1 Message Date
a18d40955c US06-01: Configure and Preflight Archive Destinations 2026-08-16 17:37:47 +02:00
11 changed files with 1350 additions and 0 deletions

View File

@@ -0,0 +1,40 @@
"""Archive locations (US06-01).
Revision ID: 0011_archive_locations
Revises: 0010_upload_verification
Create Date: 2026-08-16
Configured archive destinations with their stable media identity, last probed
capabilities, and state.
"""
import sqlalchemy as sa
from alembic import op
revision = "0011_archive_locations"
down_revision = "0010_upload_verification"
branch_labels = None
depends_on = None
def upgrade() -> None:
op.create_table(
"archive_locations",
sa.Column("id", sa.String(), primary_key=True),
sa.Column("name", sa.String(), nullable=False, unique=True),
sa.Column("root", sa.String(), nullable=False),
sa.Column("media_id", sa.String(), nullable=False, unique=True),
sa.Column("capabilities", sa.String(), nullable=True),
sa.Column("state", sa.String(), nullable=False, server_default="offline"),
sa.Column("last_seen_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()
),
)
def downgrade() -> None:
op.drop_table("archive_locations")

View File

@@ -17,6 +17,7 @@ from fastapi.staticfiles import StaticFiles
from photo_pipeline.api.routes import (
albums,
analysis,
archives,
duplicates,
health,
inventory,
@@ -73,6 +74,7 @@ def create_app(config: Config | None = None) -> FastAPI:
app.include_router(albums.router, prefix="/api/v1")
app.include_router(renames.router, prefix="/api/v1")
app.include_router(uploads.router, prefix="/api/v1")
app.include_router(archives.router, prefix="/api/v1")
# Static single-page app (hash-routed). Mounted last so /api/v1 wins.
if FRONTEND_DIR.is_dir():
app.mount("/app", StaticFiles(directory=FRONTEND_DIR, html=True), name="app")

View File

@@ -0,0 +1,63 @@
"""Archive location and preflight API (US06-01).
Registering a location writes a marker onto the medium; preflight is a command
rather than a read, because it probes the destination, hashes the scope, and issues
the token a later archive plan must present (US06-02). Neither endpoint moves or
removes a single library file.
"""
from __future__ import annotations
from fastapi import APIRouter, Request
from fastapi.responses import JSONResponse
from pydantic import BaseModel
from photo_pipeline.services.archives import ArchiveError, ArchiveService
router = APIRouter(tags=["archives"])
# Which failures are the caller's request (422) and which are a missing thing (404).
NOT_FOUND_CODES = {"unknown_location"}
class RegisterLocationRequest(BaseModel):
name: str
root: str
class PreflightRequest(BaseModel):
location_id: str
# ``None`` means every album; an explicit list scopes the check.
albums: list[str] | None = None
def _service(request: Request) -> ArchiveService:
return ArchiveService(request.app.state.session_factory, config=request.app.state.config)
def _error(error: ArchiveError) -> JSONResponse:
status = 404 if error.code in NOT_FOUND_CODES else 422
return JSONResponse(
status_code=status, content={"error": {"code": error.code, "message": str(error)}}
)
@router.post("/archive-locations", status_code=201)
def register_location(body: RegisterLocationRequest, request: Request):
try:
return _service(request).register(body.name, body.root)
except ArchiveError as error:
return _error(error)
@router.get("/archive-locations")
def list_locations(request: Request) -> dict:
return {"locations": _service(request).locations()}
@router.post("/archive-preflight")
def preflight(body: PreflightRequest, request: Request):
try:
return _service(request).preflight(body.location_id, body.albums)
except ArchiveError as error:
return _error(error)

View File

@@ -35,6 +35,9 @@ class Config(BaseModel):
thumbnail_cache_quota_bytes: int = 500_000_000
thumbnail_max_pixels: int = 100_000_000
# Free space an archive destination must keep beyond the transfer itself.
archive_free_space_reserve_bytes: int = 1_000_000_000
vision_api_key: SecretStr | None = None
immich_api_key: SecretStr | None = None
immich_server_url: str = ""

View File

@@ -21,6 +21,9 @@ UPLOAD_BATCH = "upload_batch"
LIBRARY_WRITE_LOCK = "library_write"
# The uploader lane: one album batch at a time (concept §16).
UPLOAD_LOCK = "upload"
# The archiver lane: one archive/restore plan at a time (concept §16). No handler
# runs on it yet (US06-02); preflight already refuses to plan around a held lease.
ARCHIVE_LOCK = "archive"
def _safety_score_item(asset_id: str, ctx: JobContext) -> None:

View File

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

View File

@@ -0,0 +1,42 @@
"""Archive location persistence (US06-01).
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
recorded root alone can never prove "these bytes went to that volume". Each
location therefore owns a marker file written onto the medium itself; its
``media_id`` is the stable identity, and the root is only where it was last seen.
``capabilities`` and ``state`` are the last probe result, kept so the UI can list
locations without touching a sleeping disk. Preflight always re-probes — a stored
state is a hint, never evidence.
"""
from __future__ import annotations
from datetime import datetime
from sqlalchemy import DateTime, String, func
from sqlalchemy.orm import Mapped, mapped_column
from photo_pipeline.db import Base
class ArchiveLocation(Base):
__tablename__ = "archive_locations"
id: Mapped[str] = mapped_column(String, primary_key=True)
name: Mapped[str] = mapped_column(String, nullable=False, unique=True)
root: Mapped[str] = mapped_column(String, nullable=False)
# Written into the marker file on the medium; proves the right volume is mounted.
media_id: Mapped[str] = mapped_column(String, nullable=False, unique=True)
capabilities: Mapped[str | None] = mapped_column(String) # JSON, last probe
# online | offline | wrong_volume | unwritable — the last probe's verdict.
state: Mapped[str] = mapped_column(String, nullable=False, default="offline")
last_seen_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()
)

View File

@@ -0,0 +1,599 @@
"""ArchiveService — destinations and archive preflight (US06-01).
Archive is the only stage that *removes* originals from the active library, so
this service does the opposite of removing anything: it registers destinations and
proves, before a single byte moves, that an album could be archived safely. The
transfer itself is US06-02.
An archive location is a medium, not a path (see :class:`ArchiveLocation`). A
marker file on the medium carries its ``media_id``, so a disk mounted at the
recorded root but holding a different marker is ``wrong_volume`` rather than
silently accepted — the classic "the external disk came back at the same
mountpoint" data-loss path.
Preflight proves, per concept §9 "Archive preflight":
- the album's upload is *verified*, not merely process-successful, and its bytes on
disk still hash to exactly what was uploaded;
- the destination medium is mounted, is the right one, is writable, lies outside
every library root and every ``_IGNORE/`` tree, and has room for the scope plus a
configured reserve;
- nothing already occupies the destination;
- no rename/upload/archive lease is held, and no rename is half-applied;
- a database backup and the archive manifest can really be written — both are
probed by writing them, not assumed.
Blocker codes: ``no_library_root``, ``location_offline``, ``wrong_volume``,
``unsafe_destination``, ``destination_not_writable``, ``manifest_unwritable``,
``insufficient_capacity``, ``backup_unavailable``, ``lock_conflict``,
``rename_pending``, ``empty_scope``, ``destination_collision``,
``upload_unverified``, ``bytes_changed``, ``file_missing``.
Like upload preflight, the confirmation token is *derived* from the report rather
than stored: any change to the scope, the bytes, the destination, or the blockers
produces a different token, so a stale browser confirmation can never apply. Values
that drift without meaning anything (free space, backup size, timestamps) are left
out of the digest.
"""
from __future__ import annotations
import hashlib
import json
import os
import shutil
import sqlite3
import uuid
from contextlib import closing
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.integrations import immich_go_report as report_parser
from photo_pipeline.jobs.domain_handlers import ARCHIVE_LOCK, LIBRARY_WRITE_LOCK, UPLOAD_LOCK
from photo_pipeline.models import ArchiveLocation, Asset, UploadBatch, UploadItem
from photo_pipeline.path_policy import PathPolicyError, is_excluded, normalize_root, resolve_within
from photo_pipeline.services.albums import album_label
from photo_pipeline.services.hashing import sha256_file
from photo_pipeline.services.jobs import JobService
from photo_pipeline.services.rename_journal import RenameJournal
from photo_pipeline.services.upload_reports import VERIFIED
PREFLIGHT_VERSION = 1
TOKEN_PREFIX = f"v{PREFLIGHT_VERSION}"
MARKER_NAME = ".photo-pipeline-archive.json"
MANIFEST_NAME = "archive-manifest.json"
# Upload outcomes that prove Immich holds these exact bytes. ``skipped``/``failed``/
# ``unknown`` never qualify: archiving on them would remove the only copy.
ARCHIVED_OUTCOMES = frozenset(
{report_parser.UPLOADED, report_parser.UPGRADED, report_parser.DUPLICATE}
)
LOCKS = (LIBRARY_WRITE_LOCK, UPLOAD_LOCK, ARCHIVE_LOCK)
class ArchiveError(RuntimeError):
"""The request cannot be carried out (unknown location/album, unsafe root)."""
def __init__(self, code: str, message: str) -> None:
super().__init__(message)
self.code = code
def _now() -> datetime:
return datetime.now(timezone.utc)
def _issue(code: str, message: str) -> dict:
return {"code": code, "message": message}
class ArchiveService:
def __init__(self, session_factory: sessionmaker, *, config: Config) -> None:
self._session_factory = session_factory
self._config = config
self._roots = tuple(normalize_root(root) for root in config.library_roots)
# ── locations ─────────────────────────────────────────────────────────────
def register(self, name: str, root: str) -> dict:
"""Register an archive destination and stamp its medium with a marker.
The marker is what makes the location identifiable later, so registering is
the one archive operation that writes to the destination up front.
"""
name = (name or "").strip()
if not name:
raise ArchiveError("name_required", "an archive location needs a name")
path = Path(root).expanduser()
if not path.is_dir():
raise ArchiveError("root_missing", f"{path} is not an existing directory")
path = normalize_root(path)
unsafe = self._unsafe_destination(path)
if unsafe:
raise ArchiveError("unsafe_destination", unsafe)
marker = _read_marker(path)
with self._session_factory() as session:
if session.scalar(select(ArchiveLocation).where(ArchiveLocation.name == name)):
raise ArchiveError("duplicate_name", f"an archive location named {name!r} exists")
if marker and session.scalar(
select(ArchiveLocation).where(ArchiveLocation.media_id == marker.get("media_id"))
):
raise ArchiveError(
"already_registered", f"{path} already belongs to another archive location"
)
media_id = marker.get("media_id") if marker else str(uuid.uuid4())
error = _probe_write(
path / MARKER_NAME,
json.dumps({"media_id": media_id, "name": name}, indent=2).encode("utf-8"),
keep=True,
)
if error:
raise ArchiveError("destination_not_writable", error)
location = ArchiveLocation(
id=str(uuid.uuid4()),
name=name,
root=str(path),
media_id=media_id,
state="online",
last_seen_at=_now(),
)
location.capabilities = json.dumps(_capabilities(path))
session.add(location)
session.commit()
return self._location_report(location, probe=_probe_location(location))
def locations(self) -> list[dict]:
"""Every configured location with a fresh probe of its medium."""
with self._session_factory() as session:
rows = list(session.scalars(select(ArchiveLocation).order_by(ArchiveLocation.name)))
reports = []
for location in rows:
probe = _probe_location(location)
location.state = probe["state"]
if probe["state"] == "online":
location.last_seen_at = _now()
location.capabilities = json.dumps(probe["capabilities"])
reports.append(self._location_report(location, probe=probe))
session.commit()
return reports
# ── preflight ─────────────────────────────────────────────────────────────
def preflight(self, location_id: str, albums: list[str] | None = None) -> dict:
"""Validate an archive scope against a destination and issue its token.
Read-only with respect to the library: it hashes files, probes the
destination with its own temporary files, and writes nothing else.
"""
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}")
probe = _probe_location(location)
location.state = probe["state"]
if probe["state"] == "online":
location.last_seen_at = _now()
location.capabilities = json.dumps(probe["capabilities"])
report = {
"schema_version": PREFLIGHT_VERSION,
"location": self._location_report(location, probe=probe),
"blockers": [],
}
root = Path(location.root)
session.commit()
report["blockers"] += self._destination_blockers(root, probe)
report["blockers"] += self._lock_blockers()
report["albums"] = self._albums(albums, root, reachable=probe["state"] == "online")
report["totals"] = _totals(report["albums"])
report["capacity"] = self._capacity(report["totals"]["bytes"], probe)
if not report["capacity"]["sufficient"]:
report["blockers"].append(
_issue(
"insufficient_capacity",
f"{report['totals']['bytes']} B plus a "
f"{self._config.archive_free_space_reserve_bytes} B reserve do not fit in "
f"{report['capacity']['free_bytes']} B of free space",
)
)
report["backup"] = self._backup_probe()
if not report["backup"]["ok"]:
report["blockers"].append(
_issue(
"backup_unavailable",
f"a database backup could not be written: {report['backup']['detail']}",
)
)
report["manifest"] = self._manifest_probe(
root, report["albums"], writable=probe["writable"]
)
if not report["manifest"]["ok"]:
report["blockers"].append(
_issue(
"manifest_unwritable",
f"the archive manifest could not be written: {report['manifest']['detail']}",
)
)
if not report["albums"]:
report["blockers"].append(
_issue("empty_scope", "no canonical, active assets are in the selected scope")
)
report["state"] = (
"ready"
if not report["blockers"] and all(a["state"] == "ready" for a in report["albums"])
else "blocked"
)
report["token"] = _token(report)
report["generated_at"] = _now().isoformat()
return report
def verify_token(self, token: str, location_id: str, albums: list[str] | None = None) -> bool:
"""True when ``token`` still describes this scope and this destination.
Recomputed, never looked up: an edited source file, a swapped medium, or a
newly occupied destination invalidates it without anything writing to the
database.
"""
return bool(token) and token == self.preflight(location_id, albums)["token"]
# ── destination ───────────────────────────────────────────────────────────
def _unsafe_destination(self, root: Path) -> str | None:
"""Why this root may never hold archived originals, or ``None``."""
if is_excluded(root):
return f"{root} is inside an excluded (_IGNORE/) tree"
for library in self._roots:
if root == library or library in root.parents or root in library.parents:
return f"{root} overlaps the active library root {library}"
return None
def _destination_blockers(self, root: Path, probe: dict) -> list[dict]:
blockers: list[dict] = []
if not self._roots:
blockers.append(_issue("no_library_root", "no library root is configured"))
if probe["state"] == "offline":
blockers.append(
_issue("location_offline", f"the archive medium is not mounted at {root}")
)
elif probe["state"] == "wrong_volume":
blockers.append(
_issue(
"wrong_volume",
f"{root} holds a different archive medium ({probe['detail']})",
)
)
unsafe = self._unsafe_destination(root)
if unsafe:
blockers.append(_issue("unsafe_destination", unsafe))
if probe["state"] == "unwritable":
blockers.append(
_issue("destination_not_writable", f"{root} is not writable: {probe['detail']}")
)
return blockers
def _lock_blockers(self) -> list[dict]:
"""Archive is blocked by any lease that may still be moving bytes or metadata."""
blockers: list[dict] = []
jobs = JobService(self._session_factory)
for lock in LOCKS:
held = jobs.blockers(lock)
if held:
blockers.append(
_issue("lock_conflict", f"the {lock} lane is busy: job {held[0]['id']}")
)
if RenameJournal(self._session_factory).blocks_mutation():
blockers.append(
_issue("rename_pending", "an unresolved rename must be recovered before archiving")
)
return blockers
def _capacity(self, required: int, probe: dict) -> dict:
reserve = self._config.archive_free_space_reserve_bytes
free = probe["free_bytes"]
return {
"required_bytes": required,
"reserve_bytes": reserve,
"free_bytes": free,
"sufficient": free is not None and free >= required + reserve,
}
def _backup_probe(self) -> dict:
"""Write a real online backup of the database, then discard it.
A backup that is merely assumed to be possible is worth nothing on the day
the archive removes the originals, so this actually runs SQLite's backup API.
"""
source = self._config.database_path
target = source.parent / f".archive-preflight-backup-{uuid.uuid4()}.db"
try:
with closing(sqlite3.connect(source)) as src, closing(sqlite3.connect(target)) as dst:
src.backup(dst)
size = target.stat().st_size
except (sqlite3.Error, OSError) as error:
return {"ok": False, "bytes": None, "detail": str(error)}
finally:
target.unlink(missing_ok=True)
return {"ok": True, "bytes": size, "detail": None}
def _manifest_probe(self, root: Path, albums: list[dict], *, writable: bool) -> dict:
"""Prove the manifest can be created by writing this exact content and
removing it again. The real manifest is written by the transfer (US06-02)."""
manifest = {
"schema_version": PREFLIGHT_VERSION,
"albums": [
{
"album": album["album"],
"destination": album["destination"],
"files": [
{
"asset_id": asset["asset_id"],
"source": asset["current_path"],
"sha256": asset["current_sha256"],
"byte_size": asset["byte_size"],
}
for asset in album["assets"]
],
}
for album in albums
],
}
payload = json.dumps(manifest, indent=2, sort_keys=True).encode("utf-8")
if not writable:
return {"ok": False, "bytes": len(payload), "detail": "the destination is unavailable"}
error = _probe_write(root / f".{MANIFEST_NAME}.probe-{uuid.uuid4()}", payload)
return {"ok": error is None, "bytes": len(payload), "detail": error}
def _location_report(self, location: ArchiveLocation, *, probe: dict) -> dict:
return {
"id": location.id,
"name": location.name,
"root": location.root,
"media_id": location.media_id,
"state": probe["state"],
"writable": probe["writable"],
"device_id": probe["device_id"],
"detail": probe["detail"],
"last_seen_at": location.last_seen_at.isoformat() if location.last_seen_at else None,
}
# ── scope ─────────────────────────────────────────────────────────────────
def _albums(self, requested: list[str] | None, root: Path, *, reachable: bool) -> list[dict]:
by_album = self._scope()
if requested is not None:
unknown = sorted(set(requested) - set(by_album))
if unknown:
raise ArchiveError("unknown_album", f"unknown album(s): {', '.join(unknown)}")
by_album = {name: by_album[name] for name in sorted(set(requested))}
return [
self._album(name, rows, root, reachable=reachable)
for name, rows in sorted(by_album.items())
]
def _scope(self) -> dict[str, list[dict]]:
"""Canonical, active assets grouped by album, each with its upload evidence."""
with self._session_factory() as session:
assets = list(
session.scalars(
select(Asset).where(
Asset.canonical_asset_id.is_(None),
Asset.availability_state == "active",
Asset.current_path.is_not(None),
)
)
)
uploads: dict[str, UploadItem] = {}
for item, batch in session.execute(
select(UploadItem, UploadBatch)
.join(UploadBatch, UploadBatch.id == UploadItem.batch_id)
.order_by(UploadBatch.created_at)
):
if _proves_upload(item, batch):
uploads[item.asset_id] = item # the latest verified batch wins
by_album: dict[str, list[dict]] = {}
for asset in assets:
by_album.setdefault(album_label(asset.current_path, self._roots), []).append(
{
"asset_id": asset.id,
"path": asset.current_path,
"byte_size": asset.byte_size,
"upload": uploads.get(asset.id),
}
)
return by_album
def _album(self, name: str, rows: list[dict], root: Path, *, reachable: bool) -> dict:
folder = Path(rows[0]["path"]).parent
items = sorted((_item(row) for row in rows), key=lambda item: item["current_path"])
blocked = [item for item in items if item["blockers"]]
blockers: list[dict] = []
destination = root / name
try:
resolve_within(root, destination)
except PathPolicyError as error:
blockers.append(_issue("unsafe_destination", str(error)))
if reachable and destination.exists() and any(destination.iterdir()):
blockers.append(
_issue("destination_collision", f"{destination} already exists and is not empty")
)
if blocked:
blockers.append(
_issue(
"partial_scope",
f"{len(blocked)} of {len(items)} asset(s) are not archivable; an album is "
"archived whole or not at all",
)
)
return {
"album": name,
"folder": str(folder),
"destination": str(destination),
# Same filesystem means the transfer can be an atomic move; anything else
# is copy-verify-remove (concept §9).
"transfer_method": _transfer_method(folder, root),
"asset_count": len(items),
"blocked_count": len(blocked),
"reclaimable_bytes": sum(item["byte_size"] or 0 for item in items),
"state": "blocked" if blockers else "ready",
"blockers": blockers,
"assets": items,
}
# ── internals ────────────────────────────────────────────────────────────────
def _item(row: dict) -> dict:
"""One asset's archivability: verified upload plus the bytes on disk right now."""
path = Path(row["path"])
upload: UploadItem | None = row["upload"]
blockers: list[dict] = []
current_sha256 = None
if not path.exists():
blockers.append(_issue("file_missing", f"{path} is missing"))
else:
# ponytail: full re-hash of the scope. Gate on (size, mtime_ns) first if a
# large album makes this slow — the hash stays the authority.
current_sha256 = sha256_file(path)
if upload is None:
blockers.append(
_issue("upload_unverified", "a verified Immich upload of these bytes is required")
)
elif current_sha256 is not None and upload.sha256 and current_sha256 != upload.sha256:
blockers.append(
_issue("bytes_changed", f"{path} changed since it was uploaded; re-upload it first")
)
return {
"asset_id": row["asset_id"],
"current_path": str(path),
"byte_size": row["byte_size"],
"current_sha256": current_sha256,
"uploaded_sha256": upload.sha256 if upload else None,
"blockers": blockers,
}
def _proves_upload(item: UploadItem, batch: UploadBatch) -> bool:
"""Whether this upload item is evidence that Immich holds these exact bytes."""
return (
batch.outcome_state == VERIFIED
and not batch.stale_bytes
and not item.changed_after_upload
and item.outcome in ARCHIVED_OUTCOMES
)
def _transfer_method(folder: Path, root: Path) -> str:
try:
if folder.stat().st_dev == root.stat().st_dev:
return "move"
except OSError:
pass
return "copy_verify_remove"
def _read_marker(root: Path) -> dict | None:
try:
return json.loads((root / MARKER_NAME).read_text(encoding="utf-8"))
except (OSError, ValueError):
return None
def _probe_write(path: Path, payload: bytes, *, keep: bool = False) -> str | None:
"""Write ``payload`` to ``path``; return the failure detail or ``None``."""
try:
path.write_bytes(payload)
except OSError as error:
return str(error)
if not keep:
try:
path.unlink()
except OSError as error:
return str(error)
return None
def _capabilities(root: Path) -> dict:
usage = shutil.disk_usage(root)
return {
"device_id": root.stat().st_dev,
"total_bytes": usage.total,
"writable": os.access(root, os.W_OK),
}
def _probe_location(location: ArchiveLocation) -> dict:
"""Is the right medium mounted, and can it take bytes right now?"""
root = Path(location.root)
blank = {"device_id": None, "free_bytes": None, "total_bytes": None, "capabilities": {}}
if not root.is_dir():
return {"state": "offline", "writable": False, "detail": f"{root} is not mounted", **blank}
marker = _read_marker(root)
if marker is None:
return {
"state": "offline",
"writable": False,
"detail": f"no archive marker found at {root}",
**blank,
}
if marker.get("media_id") != location.media_id:
return {
"state": "wrong_volume",
"writable": False,
"detail": f"marker media_id {marker.get('media_id')!r}",
**blank,
}
capabilities = _capabilities(root)
usage = shutil.disk_usage(root)
# os.access lies on some filesystems; a real write is the only proof.
detail = _probe_write(root / f".archive-write-probe-{uuid.uuid4()}", b"")
return {
"state": "online" if detail is None else "unwritable",
"writable": detail is None,
"detail": detail,
"device_id": capabilities["device_id"],
"free_bytes": usage.free,
"total_bytes": usage.total,
"capabilities": capabilities,
}
def _totals(albums: list[dict]) -> dict:
return {
"albums": len(albums),
"ready_albums": sum(1 for album in albums if album["state"] == "ready"),
"assets": sum(album["asset_count"] for album in albums),
"blocked": sum(album["blocked_count"] for album in albums),
"bytes": sum(album["reclaimable_bytes"] for album in albums),
}
def _token(report: dict) -> str:
"""Digest of everything the report asserts about the scope and the destination.
Values that drift without changing what would happen — free space, backup size,
timestamps — are excluded so the same situation always yields the same token.
"""
payload = {key: value for key, value in report.items() if key not in ("generated_at", "token")}
payload["location"] = {
key: value for key, value in payload["location"].items() if key != "last_seen_at"
}
payload["capacity"] = {
key: value for key, value in payload["capacity"].items() if key != "free_bytes"
}
payload["backup"] = {key: value for key, value in payload["backup"].items() if key != "bytes"}
digest = hashlib.sha256(
json.dumps(payload, sort_keys=True, ensure_ascii=False, default=str).encode("utf-8")
).hexdigest()
return f"{TOKEN_PREFIX}:{digest}"

View File

@@ -32,4 +32,5 @@ markers = [
"phase_c: Phase C end-to-end acceptance (US03-05) — album proposal API and browser journeys",
"phase_d: Phase D end-to-end acceptance (US04-06) — guarded rename API, fault, and browser journeys",
"phase_e: Phase E end-to-end acceptance (US05-06) — upload preflight, uploader, and browser journeys",
"phase_f: Phase F end-to-end acceptance (US06-06) — archive destination, transfer, and restore journeys",
]

View File

@@ -0,0 +1,592 @@
"""Archive destinations and preflight (US06-01).
Archive is the only stage that removes originals, so every case here asks the same
question: would this preflight let an album leave active storage when it should
not? The destinations are real directories on real filesystems — mounted, missing,
swapped for another medium, read-only, or full — and preflight itself must stay
non-destructive: the library snapshot is asserted unchanged.
"""
import json
import os
import stat
import uuid
from datetime import datetime, timedelta, timezone
import pytest
from fastapi.testclient import TestClient
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.jobs.domain_handlers import ARCHIVE_LOCK, UPLOAD_LOCK
from photo_pipeline.models import (
Asset,
RenameOperation,
RenamePlan,
UploadBatch,
UploadItem,
)
from photo_pipeline.services.archives import MARKER_NAME, ArchiveError, ArchiveService
from photo_pipeline.services.hashing import sha256_file
from photo_pipeline.services.jobs import JobService
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, *, reserve=0):
(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": str(reserve),
}
)
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"),
*,
uploaded=True,
outcome="uploaded",
outcome_state="verified",
stale_bytes=False,
):
"""A real album folder whose assets carry their upload evidence."""
folder = lib / album
folder.mkdir(parents=True, exist_ok=True)
ids = []
with sf() as session:
batch_id = str(uuid.uuid4())
if uploaded:
session.add(
UploadBatch(
id=batch_id,
album=album,
folder=str(folder),
album_name=album,
state="succeeded",
preflight_token="v1:test",
outcome_state=outcome_state,
stale_bytes=stale_bytes,
created_at=NOW,
)
)
for name in names:
path = folder / name
path.write_bytes(name.encode() * 16)
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),
)
)
if uploaded:
session.add(
UploadItem(
batch_id=batch_id,
asset_id=asset_id,
path=str(path),
sha256=sha256_file(path),
sha1="0" * 40,
state="sent",
outcome=outcome,
)
)
session.commit()
return folder, ids
def _service(sf, config):
return ArchiveService(sf, config=config)
def _location(sf, config, archive, name="external"):
return _service(sf, config).register(name, str(archive))
def _snapshot(lib):
return {
str(p.relative_to(lib)): (p.read_bytes() if p.is_file() else None)
for p in sorted(lib.rglob("*"))
}
def _codes(report):
return (
{issue["code"] for issue in report["blockers"]}
| {issue["code"] for album in report["albums"] for issue in album["blockers"]}
| {
issue["code"]
for album in report["albums"]
for asset in album["assets"]
for issue in asset["blockers"]
}
)
# ── locations ────────────────────────────────────────────────────────────────
def test_registering_a_location_stamps_the_medium_with_its_identity(tmp_path):
config, sf, _, archive = _env(tmp_path)
location = _location(sf, config, archive)
marker = json.loads((archive / MARKER_NAME).read_text())
assert marker["media_id"] == location["media_id"]
assert location["state"] == "online" and location["writable"] is True
assert location["root"] == str(archive.resolve())
listed = _service(sf, config).locations()
assert [(row["id"], row["media_id"], row["state"]) for row in listed] == [
(location["id"], location["media_id"], "online")
]
def test_a_second_location_cannot_claim_the_same_medium(tmp_path):
config, sf, _, archive = _env(tmp_path)
_location(sf, config, archive)
with pytest.raises(ArchiveError) as error:
_location(sf, config, archive, name="second")
assert error.value.code == "already_registered"
@pytest.mark.parametrize("inside", ["", "sub"])
def test_a_destination_inside_the_library_is_refused(tmp_path, inside):
"""The library may never archive into itself: the 'reclaimed' bytes would still
be in the active tree, and a later scan would rediscover them."""
config, sf, lib, _ = _env(tmp_path)
root = lib / inside if inside else lib
root.mkdir(exist_ok=True)
with pytest.raises(ArchiveError) as error:
_service(sf, config).register("bad", str(root))
assert error.value.code == "unsafe_destination"
def test_an_ignored_destination_is_refused(tmp_path):
config, sf, _, _ = _env(tmp_path)
root = tmp_path / "_IGNORE" / "archive"
root.mkdir(parents=True)
with pytest.raises(ArchiveError) as error:
_service(sf, config).register("ignored", str(root))
assert error.value.code == "unsafe_destination"
def test_listing_reports_an_unmounted_medium_as_offline(tmp_path):
config, sf, _, archive = _env(tmp_path)
_location(sf, config, archive)
(archive / MARKER_NAME).unlink()
assert [row["state"] for row in _service(sf, config).locations()] == ["offline"]
# ── happy path ───────────────────────────────────────────────────────────────
def test_ready_preflight_previews_scope_method_and_reclaimable_bytes(tmp_path):
config, sf, lib, archive = _env(tmp_path)
folder, ids = _album(sf, lib)
location = _location(sf, config, archive)
before = _snapshot(lib)
report = _service(sf, config).preflight(location["id"])
assert report["state"] == "ready" and report["blockers"] == []
album = report["albums"][0]
assert album["album"] == "rome" and album["folder"] == str(folder)
assert album["destination"] == str(archive.resolve() / "rome")
assert album["transfer_method"] in ("move", "copy_verify_remove")
assert album["reclaimable_bytes"] == sum(p.stat().st_size for p in folder.iterdir())
assert sorted(a["asset_id"] for a in album["assets"]) == sorted(ids)
assert report["totals"]["bytes"] == album["reclaimable_bytes"]
assert report["capacity"]["sufficient"] is True
# Both must be proven by writing, not assumed.
assert report["backup"]["ok"] is True and report["backup"]["bytes"] > 0
assert report["manifest"]["ok"] is True
assert report["token"].startswith("v1:")
assert _snapshot(lib) == before, "preflight must not touch the library"
assert not list(archive.glob("*probe*")), "probe files must be cleaned up"
def test_same_filesystem_destination_is_previewed_as_a_move(tmp_path):
config, sf, lib, archive = _env(tmp_path)
_album(sf, lib)
location = _location(sf, config, archive)
report = _service(sf, config).preflight(location["id"])
# tmp_path is one filesystem, so this is the same-filesystem case by construction.
assert report["albums"][0]["transfer_method"] == "move"
def test_scoping_to_one_album_excludes_the_others(tmp_path):
config, sf, lib, archive = _env(tmp_path)
_album(sf, lib, "rome")
_album(sf, lib, "paris", names=("c.jpg",))
location = _location(sf, config, archive)
report = _service(sf, config).preflight(location["id"], ["paris"])
assert [album["album"] for album in report["albums"]] == ["paris"]
assert report["totals"]["assets"] == 1
def test_unknown_album_and_unknown_location_are_refused(tmp_path):
config, sf, lib, archive = _env(tmp_path)
_album(sf, lib)
location = _location(sf, config, archive)
service = _service(sf, config)
with pytest.raises(ArchiveError) as unknown_album:
service.preflight(location["id"], ["atlantis"])
with pytest.raises(ArchiveError) as unknown_location:
service.preflight("nope")
assert unknown_album.value.code == "unknown_album"
assert unknown_location.value.code == "unknown_location"
def test_empty_scope_is_a_blocker(tmp_path):
config, sf, _, archive = _env(tmp_path)
location = _location(sf, config, archive)
report = _service(sf, config).preflight(location["id"])
assert report["state"] == "blocked" and "empty_scope" in _codes(report)
# ── destination ──────────────────────────────────────────────────────────────
def test_offline_medium_blocks(tmp_path):
config, sf, lib, archive = _env(tmp_path)
_album(sf, lib)
location = _location(sf, config, archive)
(archive / MARKER_NAME).unlink() # the disk went away
report = _service(sf, config).preflight(location["id"])
assert report["state"] == "blocked" and "location_offline" in _codes(report)
assert report["location"]["state"] == "offline"
def test_a_different_medium_at_the_same_mountpoint_blocks(tmp_path):
"""The mountpoint is right, the disk is not — never write the archive here."""
config, sf, lib, archive = _env(tmp_path)
_album(sf, lib)
location = _location(sf, config, archive)
(archive / MARKER_NAME).write_text(json.dumps({"media_id": "some-other-disk"}))
report = _service(sf, config).preflight(location["id"])
assert "wrong_volume" in _codes(report)
assert report["location"]["state"] == "wrong_volume"
def test_read_only_destination_blocks_and_cannot_write_the_manifest(tmp_path):
config, sf, lib, archive = _env(tmp_path)
_album(sf, lib)
location = _location(sf, config, archive)
mode = archive.stat().st_mode
archive.chmod(mode & ~stat.S_IWUSR & ~stat.S_IWGRP & ~stat.S_IWOTH)
try:
report = _service(sf, config).preflight(location["id"])
finally:
archive.chmod(mode)
assert {"destination_not_writable", "manifest_unwritable"} <= _codes(report)
assert report["manifest"]["ok"] is False
@pytest.mark.skipif(os.geteuid() == 0, reason="root ignores directory permissions")
def test_read_only_destination_is_detected_by_a_real_write(tmp_path):
config, sf, lib, archive = _env(tmp_path)
_album(sf, lib)
location = _location(sf, config, archive)
mode = archive.stat().st_mode
archive.chmod(stat.S_IRUSR | stat.S_IXUSR)
try:
report = _service(sf, config).preflight(location["id"])
finally:
archive.chmod(mode)
assert report["location"]["writable"] is False
def test_insufficient_capacity_blocks(tmp_path):
"""The reserve is what stops an archive from filling its own destination."""
config, sf, lib, archive = _env(tmp_path, reserve=10**15)
_album(sf, lib)
location = _location(sf, config, archive)
report = _service(sf, config).preflight(location["id"])
assert report["state"] == "blocked" and "insufficient_capacity" in _codes(report)
assert report["capacity"]["sufficient"] is False
assert report["capacity"]["reserve_bytes"] == 10**15
def test_an_occupied_destination_blocks_that_album(tmp_path):
config, sf, lib, archive = _env(tmp_path)
_album(sf, lib)
location = _location(sf, config, archive)
(archive / "rome").mkdir()
(archive / "rome" / "a.jpg").write_bytes(b"something already here")
report = _service(sf, config).preflight(location["id"])
assert "destination_collision" in _codes(report)
assert report["albums"][0]["state"] == "blocked"
def test_a_destination_moved_into_the_library_blocks_even_though_it_registered(tmp_path):
"""Registration validated the root once; preflight validates it again, because a
mountpoint can be moved after the fact."""
config, sf, lib, archive = _env(tmp_path)
_album(sf, lib)
location = _location(sf, config, archive)
with sf() as session:
from photo_pipeline.models import ArchiveLocation
session.get(ArchiveLocation, location["id"]).root = str(lib / "inside")
session.commit()
(lib / "inside").mkdir()
(lib / "inside" / MARKER_NAME).write_text(json.dumps({"media_id": location["media_id"]}))
report = _service(sf, config).preflight(location["id"])
assert "unsafe_destination" in _codes(report)
# ── source readiness ─────────────────────────────────────────────────────────
@pytest.mark.parametrize(
"kwargs,code",
[
({"uploaded": False}, "upload_unverified"),
({"outcome_state": "requires_verification"}, "upload_unverified"),
({"outcome": "failed"}, "upload_unverified"),
({"outcome": "skipped"}, "upload_unverified"),
({"stale_bytes": True}, "upload_unverified"),
],
)
def test_an_unverified_upload_blocks_the_album(tmp_path, kwargs, code):
config, sf, lib, archive = _env(tmp_path)
_album(sf, lib, **kwargs)
location = _location(sf, config, archive)
report = _service(sf, config).preflight(location["id"])
assert report["state"] == "blocked" and code in _codes(report)
assert report["albums"][0]["blocked_count"] == 2
@pytest.mark.parametrize("outcome", ["upgraded", "duplicate"])
def test_upgraded_and_duplicate_uploads_are_evidence_enough(tmp_path, outcome):
"""Immich already holds these exact bytes; that is what archiving requires."""
config, sf, lib, archive = _env(tmp_path)
_album(sf, lib, outcome=outcome)
location = _location(sf, config, archive)
assert _service(sf, config).preflight(location["id"])["state"] == "ready"
def test_bytes_changed_since_upload_block_the_album(tmp_path):
config, sf, lib, archive = _env(tmp_path)
folder, _ = _album(sf, lib)
location = _location(sf, config, archive)
(folder / "a.jpg").write_bytes(b"edited after the upload")
report = _service(sf, config).preflight(location["id"])
assert {"bytes_changed", "partial_scope"} <= _codes(report)
assert report["albums"][0]["blocked_count"] == 1
def test_a_missing_source_file_blocks_the_album(tmp_path):
config, sf, lib, archive = _env(tmp_path)
folder, _ = _album(sf, lib)
location = _location(sf, config, archive)
(folder / "a.jpg").unlink()
assert "file_missing" in _codes(_service(sf, config).preflight(location["id"]))
# ── leases ───────────────────────────────────────────────────────────────────
@pytest.mark.parametrize("lock", [UPLOAD_LOCK, ARCHIVE_LOCK])
def test_a_held_lease_blocks_archiving(tmp_path, lock):
config, sf, lib, archive = _env(tmp_path)
_album(sf, lib)
location = _location(sf, config, archive)
JobService(sf).enqueue("upload_batch", lock=lock, items=["x"])
report = _service(sf, config).preflight(location["id"])
assert report["state"] == "blocked" and "lock_conflict" in _codes(report)
def test_a_half_applied_rename_blocks_archiving(tmp_path):
config, sf, lib, archive = _env(tmp_path)
folder, _ = _album(sf, lib)
location = _location(sf, config, archive)
with sf() as session:
plan_id = str(uuid.uuid4())
session.add(RenamePlan(id=plan_id, state="applying", operation_count=1))
session.flush()
session.add(
RenameOperation(
id=str(uuid.uuid4()),
plan_id=plan_id,
sequence=0,
operation="move_folder",
source_path=str(folder),
destination_path=str(lib / "2019 Rome"),
journal_state="moving",
)
)
session.commit()
report = _service(sf, config).preflight(location["id"])
assert report["state"] == "blocked" and "rename_pending" in _codes(report)
# ── token ────────────────────────────────────────────────────────────────────
def test_token_is_stable_while_nothing_relevant_changes(tmp_path):
config, sf, lib, archive = _env(tmp_path)
_album(sf, lib)
location = _location(sf, config, archive)
service = _service(sf, config)
first = service.preflight(location["id"])["token"]
assert service.preflight(location["id"])["token"] == first
assert service.verify_token(first, location["id"]) is True
def test_an_edited_source_makes_the_token_stale(tmp_path):
config, sf, lib, archive = _env(tmp_path)
folder, _ = _album(sf, lib)
location = _location(sf, config, archive)
service = _service(sf, config)
token = service.preflight(location["id"])["token"]
(folder / "b.jpg").write_bytes(b"edited outside the app")
assert service.verify_token(token, location["id"]) is False
def test_a_changed_destination_makes_the_token_stale(tmp_path):
config, sf, lib, archive = _env(tmp_path)
_album(sf, lib)
location = _location(sf, config, archive)
service = _service(sf, config)
token = service.preflight(location["id"])["token"]
(archive / "rome").mkdir()
(archive / "rome" / "a.jpg").write_bytes(b"appeared after approval")
assert service.verify_token(token, location["id"]) is False
def test_a_token_from_another_scope_or_medium_is_rejected(tmp_path):
config, sf, lib, archive = _env(tmp_path)
_album(sf, lib, "rome")
_album(sf, lib, "paris", names=("c.jpg",))
other = tmp_path / "archive2"
other.mkdir()
location = _location(sf, config, archive)
second = _location(sf, config, other, name="second")
service = _service(sf, config)
rome = service.preflight(location["id"], ["rome"])["token"]
assert service.verify_token(rome, location["id"], ["paris"]) is False
assert service.verify_token(rome, second["id"], ["rome"]) is False
assert service.verify_token("v1:not-a-real-token", location["id"]) is False
assert service.verify_token("", location["id"]) is False
def test_the_token_survives_free_space_and_timestamp_drift(tmp_path):
"""Free space changes constantly on a live disk; a token that expired on every
byte written elsewhere would train users to ignore it."""
config, sf, lib, archive = _env(tmp_path)
_album(sf, lib)
location = _location(sf, config, archive)
service = _service(sf, config)
token = service.preflight(location["id"])["token"]
(tmp_path / "unrelated.bin").write_bytes(b"0" * 100_000)
with sf() as session:
from photo_pipeline.models import ArchiveLocation
session.get(ArchiveLocation, location["id"]).last_seen_at = NOW - timedelta(days=5)
session.commit()
assert service.verify_token(token, location["id"]) is True
# ── API surface ──────────────────────────────────────────────────────────────
def test_api_registers_a_location_and_returns_a_preflight_report(tmp_path):
config, sf, lib, archive = _env(tmp_path)
_album(sf, lib)
with TestClient(create_app(config)) as client:
created = client.post(
"/api/v1/archive-locations", json={"name": "external", "root": str(archive)}
)
listed = client.get("/api/v1/archive-locations")
report = client.post(
"/api/v1/archive-preflight", json={"location_id": created.json()["id"]}
)
assert created.status_code == 201
assert [row["name"] for row in listed.json()["locations"]] == ["external"]
assert report.status_code == 200
assert report.json()["state"] == "ready" and report.json()["token"].startswith("v1:")
def test_api_rejects_an_unknown_location_and_an_unsafe_root(tmp_path):
config, sf, lib, _ = _env(tmp_path)
with TestClient(create_app(config)) as client:
unknown = client.post("/api/v1/archive-preflight", json={"location_id": "nope"})
unsafe = client.post("/api/v1/archive-locations", json={"name": "bad", "root": str(lib)})
assert unknown.status_code == 404 and unknown.json()["error"]["code"] == "unknown_location"
assert unsafe.status_code == 422 and unsafe.json()["error"]["code"] == "unsafe_destination"

View File

@@ -123,6 +123,9 @@
],
"US05-06": [
"tests/e2e/test_phase_e_pipeline.py"
],
"US06-01": [
"tests/integration/test_archive_preflight.py"
]
}
}