US06-03: Preserve Offline Identity and Review Evidence #79

Merged
domverse merged 1 commits from us/US06-03-preserve-offline-identity-and-review-evidence into main 2026-08-16 19:21:30 +02:00
11 changed files with 835 additions and 46 deletions
Showing only changes of commit 85fe7e6642 - Show all commits

View File

@@ -0,0 +1,29 @@
"""Protected thumbnails (US06-03).
Revision ID: 0013_protected_thumbnails
Revises: 0012_archive_plans
Create Date: 2026-08-16
A protected thumbnail is the durable comparison preview of an asset whose
original has left active storage. It is evidence rather than cache, so the LRU
quota must not evict it: the archive medium may be offline when it is needed.
"""
import sqlalchemy as sa
from alembic import op
revision = "0013_protected_thumbnails"
down_revision = "0012_archive_plans"
branch_labels = None
depends_on = None
def upgrade() -> None:
op.add_column(
"thumbnails",
sa.Column("protected", sa.Boolean(), nullable=False, server_default=sa.false()),
)
def downgrade() -> None:
op.drop_column("thumbnails", "protected")

View File

@@ -11,7 +11,7 @@ from __future__ import annotations
from datetime import datetime
from sqlalchemy import DateTime, ForeignKey, Integer, String, func
from sqlalchemy import Boolean, DateTime, ForeignKey, Integer, String, func
from sqlalchemy.orm import Mapped, mapped_column
from photo_pipeline.db import Base
@@ -29,6 +29,9 @@ class Thumbnail(Base):
width: Mapped[int | None] = mapped_column(Integer)
height: Mapped[int | None] = mapped_column(Integer)
format: Mapped[str | None] = mapped_column(String)
# Durable comparison evidence for an archived asset: never evicted by the LRU
# quota, because the original may be on a medium that is no longer reachable.
protected: Mapped[bool] = mapped_column(Boolean, nullable=False, default=False)
created_at: Mapped[datetime] = mapped_column(
DateTime(timezone=True), nullable=False, server_default=func.now()
)

View File

@@ -60,8 +60,10 @@ from photo_pipeline.services.archive_journal import (
ArchiveState,
)
from photo_pipeline.services.archives import MARKER_NAME, ArchiveError, ArchiveService
from photo_pipeline.services.duplicates import DuplicateService
from photo_pipeline.services.hashing import sha256_file
from photo_pipeline.services.rename_apply import PreconditionFailed, maybe_fault
from photo_pipeline.services.thumbnails import ThumbnailService
# The per-medium manifest: one JSON line per archived file, appended and fsynced
# before its source is removed. It lives with the bytes so the archive can still be
@@ -335,6 +337,7 @@ class ArchiveTransferService:
raise PreconditionFailed(
"archive_unverified", f"{destination} is not a verified archive copy"
)
self._require_evidence(operation["asset_id"], destination)
if source.exists():
if source.is_symlink():
raise PreconditionFailed("symlink", f"{source} became a symlink")
@@ -377,6 +380,28 @@ class ArchiveTransferService:
"asset_moved", f"asset {operation['asset_id']} is no longer at {source}"
)
def _require_evidence(self, asset_id: str, source: Path) -> dict:
"""Review evidence must be durable before the original goes.
The perceptual hash keeps the asset in the fuzzy index once its bytes are
unreachable, and the protected preview is what duplicate review can still
look at. Both are read from the freshly verified archive copy, which holds
exactly the bytes being archived. A file that cannot be decoded has neither
— recorded, not fatal, since its exact hashes remain — but failing to
produce a preview from a decodable original stops the removal (concept §9).
"""
DuplicateService(self._session_factory).ensure_phash(asset_id, source=source)
preview = ThumbnailService(self._session_factory, self._config).ensure_protected(
asset_id, source=source
)
if preview["state"] == "unavailable":
raise PreconditionFailed(
"preview_unavailable",
f"no durable comparison preview for asset {asset_id} "
f"({preview['error_code']})",
)
return preview
# ── database ──────────────────────────────────────────────────────────────
def _record_archived(self, operation: dict, location: dict, destination: Path) -> None:

View File

@@ -28,7 +28,11 @@ Blocker codes: ``no_library_root``, ``location_offline``, ``wrong_volume``,
``unsafe_destination``, ``destination_not_writable``, ``manifest_unwritable``,
``insufficient_capacity``, ``backup_unavailable``, ``lock_conflict``,
``rename_pending``, ``empty_scope``, ``destination_collision``,
``upload_unverified``, ``bytes_changed``, ``file_missing``.
``upload_unverified``, ``bytes_changed``, ``file_missing``, ``preview_unavailable``.
Preflight also *creates* the durable comparison preview of every asset in scope
(US06-03): it is the evidence duplicate review falls back on once the original is
on a medium that may be offline, so it has to exist before the original leaves.
Like upload preflight, the confirmation token is *derived* from the report rather
than stored: any change to the scope, the bytes, the destination, or the blockers
@@ -59,14 +63,16 @@ from photo_pipeline.models import ArchiveLocation, Asset, UploadBatch, UploadIte
from photo_pipeline.path_policy import PathPolicyError, is_excluded, normalize_root, resolve_within
from photo_pipeline.services.albums import album_label
from photo_pipeline.services.archive_journal import ArchiveJournal
from photo_pipeline.services.availability import MARKER_NAME, read_marker as _read_marker
from photo_pipeline.services.availability import refresh as refresh_availability
from photo_pipeline.services.hashing import sha256_file
from photo_pipeline.services.jobs import JobService
from photo_pipeline.services.rename_journal import RenameJournal
from photo_pipeline.services.thumbnails import ThumbnailService
from photo_pipeline.services.upload_reports import VERIFIED
PREFLIGHT_VERSION = 1
TOKEN_PREFIX = f"v{PREFLIGHT_VERSION}"
MARKER_NAME = ".photo-pipeline-archive.json"
MANIFEST_NAME = "archive-manifest.json"
# Upload outcomes that prove Immich holds these exact bytes. ``skipped``/``failed``/
@@ -163,7 +169,10 @@ class ArchiveService:
location.capabilities = json.dumps(probe["capabilities"])
reports.append(self._location_report(location, probe=probe))
session.commit()
return reports
# A medium that just appeared or vanished changes what is readable, so the
# archived assets are re-derived from the same probe (US06-03).
refresh_availability(self._session_factory)
return reports
# ── preflight ─────────────────────────────────────────────────────────────
@@ -420,7 +429,10 @@ class ArchiveService:
def _album(self, name: str, rows: list[dict], root: Path, *, reachable: bool) -> dict:
folder = Path(rows[0]["path"]).parent
items = sorted((_item(row) for row in rows), key=lambda item: item["current_path"])
items = sorted(
(self._with_preview(_item(row)) for row in rows),
key=lambda item: item["current_path"],
)
blocked = [item for item in items if item["blockers"]]
blockers: list[dict] = []
@@ -456,6 +468,30 @@ class ArchiveService:
"assets": items,
}
def _with_preview(self, item: dict) -> dict:
"""Create the durable comparison preview while the original is still here.
This is the last moment it can be made: once the file is archived and the
medium leaves, only the retained preview can answer "is this new photo the
same picture?". An original that cannot be decoded at all has no preview to
keep — its hashes and metadata stay the evidence — but a preview that fails
for any other reason blocks the archive (concept §9).
"""
preview = self._previews().ensure_protected(item["asset_id"])
item["preview"] = preview
if preview["state"] == "unavailable" and not item["blockers"]:
item["blockers"].append(
_issue(
"preview_unavailable",
f"a durable comparison preview of {item['current_path']} could not be "
f"created ({preview['error_code']})",
)
)
return item
def _previews(self) -> ThumbnailService:
return ThumbnailService(self._session_factory, self._config)
# ── internals ────────────────────────────────────────────────────────────────
@@ -512,13 +548,6 @@ def _transfer_method(folder: Path, root: Path) -> str:
return "copy_verify_remove"
def _read_marker(root: Path) -> dict | None:
try:
return json.loads((root / MARKER_NAME).read_text(encoding="utf-8"))
except (OSError, ValueError):
return None
def _probe_write(path: Path, payload: bytes, *, keep: bool = False) -> str | None:
"""Write ``payload`` to ``path``; return the failure detail or ``None``."""
try:

View File

@@ -0,0 +1,124 @@
"""Where an asset's bytes are right now (US06-03).
Archiving removes the original from the active library but never removes the
asset: its identity, hashes, decisions, and evidence stay. This module is the one
place that answers "can these bytes be read, and if not, why" so inventory,
duplicate review, thumbnails, and the archive service all give the same answer.
States (concept §9):
- ``active`` — the original is in the active library;
- ``archived_online`` — the recorded medium is mounted and holds the file;
- ``archived_offline`` — archived, but the medium is not available right now;
- ``missing_unexpected`` — neither an active path nor the recorded archive
location explains the absence. This is the state that must never be confused
with ``archived_offline``: an unmounted disk is normal, a mounted disk with a
hole in it is not.
A medium is identified by its marker file, never by its mountpoint, so a
different disk mounted at the recorded root is offline rather than accepted.
"""
from __future__ import annotations
import json
from collections import Counter
from datetime import datetime, timezone
from pathlib import Path
from sqlalchemy import select
from sqlalchemy.orm import Session, sessionmaker
from photo_pipeline.models import ArchiveLocation, Asset
ACTIVE = "active"
ARCHIVED_ONLINE = "archived_online"
ARCHIVED_OFFLINE = "archived_offline"
MISSING_UNEXPECTED = "missing_unexpected"
ARCHIVED = (ARCHIVED_ONLINE, ARCHIVED_OFFLINE)
MARKER_NAME = ".photo-pipeline-archive.json"
def read_marker(root: Path) -> dict | None:
"""The medium's identity marker, or ``None`` when it is not readable."""
try:
return json.loads((root / MARKER_NAME).read_text(encoding="utf-8"))
except (OSError, ValueError):
return None
def location_online(location: ArchiveLocation) -> bool:
"""True only when the *recorded* medium is mounted at its root."""
marker = read_marker(Path(location.root))
return bool(marker) and marker.get("media_id") == location.media_id
def archive_file(session: Session, asset: Asset) -> Path | None:
"""The archived file's absolute path, whether or not the medium is mounted."""
if not asset.archive_location_id or not asset.archive_path:
return None
location = session.get(ArchiveLocation, asset.archive_location_id)
if location is None:
return None
return Path(location.root) / asset.archive_path
def readable_path(session: Session, asset: Asset) -> Path | None:
"""A path whose bytes can be read now: the active file, else the archive copy."""
if asset.current_path and Path(asset.current_path).exists():
return Path(asset.current_path)
archived = archive_file(session, asset)
if archived is None:
return None
location = session.get(ArchiveLocation, asset.archive_location_id)
if not location_online(location) or not archived.exists():
return None
return archived
def state_of(session: Session, asset: Asset, *, online: dict[str, bool] | None = None) -> str:
"""The availability this asset's storage actually justifies right now."""
if asset.current_path:
return ACTIVE if Path(asset.current_path).exists() else MISSING_UNEXPECTED
if not asset.archive_location_id:
return MISSING_UNEXPECTED if asset.availability_state != ACTIVE else ACTIVE
location = session.get(ArchiveLocation, asset.archive_location_id)
if location is None:
return MISSING_UNEXPECTED
reachable = (
online[location.id] if online and location.id in online else location_online(location)
)
if not reachable:
return ARCHIVED_OFFLINE
archived = archive_file(session, asset)
# The medium is mounted and identified: the file is either there, or it is
# genuinely gone — that is not "offline", it needs a human.
return ARCHIVED_ONLINE if archived and archived.exists() else MISSING_UNEXPECTED
def refresh(session_factory: sessionmaker) -> dict[str, int]:
"""Re-derive availability for every archived asset from the media themselves.
Only archived assets are probed: whether an *active* file is present is the
inventory scan's job and costs one stat per library file. Each medium is
probed once, not once per asset.
"""
counts: Counter[str] = Counter()
now = datetime.now(timezone.utc)
with session_factory() as session:
online = {
location.id: location_online(location)
for location in session.scalars(select(ArchiveLocation))
}
for asset in session.scalars(
select(Asset).where(Asset.archive_location_id.is_not(None))
):
state = state_of(session, asset, online=online)
counts[state] += 1
if state != asset.availability_state:
asset.availability_state = state
asset.state_version += 1
asset.updated_at = now
session.commit()
return dict(counts)

View File

@@ -11,6 +11,12 @@ Detection runs in two categories:
band (NEAR/SIMILAR). These are review candidates: never decided automatically, and
negative-linked pairs are suppressed so a rejected pair is not re-suggested.
Archived assets stay in both indexes (US06-03): a new active copy of an archived
original is recognised through its hashes even while the medium is offline, and
cluster review falls back to the retained protected preview plus hash evidence.
An exact/pixel match links straight to the archived canonical; a perceptual match
is a review candidate that names the medium to mount for a pixel-level decision.
Decisions (``canonical`` / ``not_duplicate`` / ``deferred``) persist with evidence,
use optimistic version checks, are reversible, and can never form a canonical cycle.
A new content-identical member of an already-decided cluster inherits the established
@@ -34,12 +40,14 @@ from sqlalchemy import func, select
from sqlalchemy.orm import sessionmaker
from photo_pipeline.models import (
ArchiveLocation,
Asset,
DuplicateCluster,
DuplicateMember,
DuplicateNegativeLink,
Thumbnail,
)
from photo_pipeline.services import hashing
from photo_pipeline.services import availability, hashing
NEAR_MAX = 5
SIMILAR_MAX = 10
@@ -127,18 +135,18 @@ class DuplicateService:
# ── perceptual hash backfill ───────────────────────────────────────────
def ensure_phashes(self) -> int:
"""Hash whatever is readable now — an archived asset keeps the hash it
already has, and gains one whenever its medium happens to be mounted."""
updated = 0
with self._session_factory() as session:
assets = session.execute(
select(Asset).where(
Asset.availability_state == "active",
Asset.current_path.isnot(None),
)
).scalars()
assets = session.execute(select(Asset)).scalars()
for asset in assets:
if asset.phash is not None and asset.phash_version == hashing.PHASH_VERSION:
continue
value = hashing.safe_phash(asset.current_path)
source = availability.readable_path(session, asset)
if source is None:
continue
value = hashing.safe_phash(str(source))
if value is not None:
asset.phash = value
asset.phash_version = hashing.PHASH_VERSION
@@ -146,20 +154,40 @@ class DuplicateService:
session.commit()
return updated
def ensure_phash(self, asset_id: str, *, source=None) -> str | None:
"""Backfill one asset's perceptual hash while its bytes are still readable.
Archiving calls this before the original leaves — passing the archive copy
as ``source``, since the database does not point at it yet — because an
asset without a pHash silently drops out of the fuzzy index the moment its
medium is away.
"""
with self._session_factory() as session:
asset = session.get(Asset, asset_id)
if asset is None:
return None
if asset.phash is not None and asset.phash_version == hashing.PHASH_VERSION:
return asset.phash
source = source or availability.readable_path(session, asset)
if source is None:
return None
value = hashing.safe_phash(str(source))
if value is not None:
asset.phash = value
asset.phash_version = hashing.PHASH_VERSION
session.commit()
return value
# ── detection ──────────────────────────────────────────────────────────
def detect(self) -> DetectionReport:
self.ensure_phashes()
now = datetime.now(timezone.utc)
report = DetectionReport()
with self._session_factory() as session:
assets = list(
session.execute(
select(Asset).where(
Asset.availability_state == "active",
Asset.current_path.isnot(None),
)
).scalars()
)
# Every known asset stays in the indexes, archived or not: a copy of an
# archived original must be recognised as a duplicate rather than
# treated as a new photo (concept §9, invariant 12).
assets = list(session.execute(select(Asset)).scalars())
by_id = {a.id: a for a in assets}
negatives = {
_pair(link.asset_a, link.asset_b)
@@ -431,10 +459,19 @@ class DuplicateService:
@staticmethod
def _recommend_canonical(ids, by_id) -> str:
# ponytail: largest file, path as deterministic tie-break. The concept's
# richer policy (resolution, least recompression, metadata richness) lands
# with the review UI story.
return max(ids, key=lambda i: (by_id[i].byte_size or 0, by_id[i].current_path or ""))
# ponytail: largest file, then the archived copy, then path as a
# deterministic tie-break. Archived wins ties because it is the reviewed,
# uploaded original — a fresh active copy must not demote it to a variant.
# The concept's richer policy (resolution, least recompression, metadata
# richness) lands with the review UI story.
return max(
ids,
key=lambda i: (
by_id[i].byte_size or 0,
by_id[i].availability_state in availability.ARCHIVED,
by_id[i].current_path or by_id[i].archive_path or "",
),
)
def _apply_canonical(self, session, cluster, ids, canonical_id):
for member in session.execute(
@@ -507,9 +544,15 @@ class DuplicateService:
"current_path": asset.current_path if asset else None,
"byte_size": asset.byte_size if asset else None,
"phash": asset.phash if asset else None,
**self._offline_evidence(session, asset),
}
)
members.sort(key=lambda m: m["asset_id"])
# A full-resolution comparison of an offline original is impossible; the
# UI asks for that named medium instead of guessing (concept §9).
mount_required = sorted(
{m["archive_location"] for m in members if m["requires_mount"]}
)
return {
"id": cluster.id,
"method": cluster.method,
@@ -519,9 +562,54 @@ class DuplicateService:
"canonical_asset_id": cluster.canonical_asset_id,
"version": cluster.version,
"requires_confirmation": cluster.method == Method.PERCEPTUAL.value,
"mount_required": mount_required,
"members": members,
}
def _offline_evidence(self, session, asset: Asset | None) -> dict:
"""What review can still rely on when a member's original is not readable."""
if asset is None:
return {
"availability_state": None,
"archive_location": None,
"archive_location_id": None,
"archive_path": None,
"preview": {"state": "missing", "protected": False},
"requires_mount": False,
}
location = (
session.get(ArchiveLocation, asset.archive_location_id)
if asset.archive_location_id
else None
)
preview = self._preview_evidence(session, asset.id)
archived = asset.availability_state in availability.ARCHIVED
return {
"availability_state": asset.availability_state,
"archive_location": location.name if location else None,
"archive_location_id": asset.archive_location_id,
"archive_path": asset.archive_path,
"preview": preview,
# Offline archived members can still be compared through their retained
# preview and hash evidence; only pixel-level review needs the medium.
"requires_mount": archived
and asset.availability_state == availability.ARCHIVED_OFFLINE
and bool(location),
}
@staticmethod
def _preview_evidence(session, asset_id: str) -> dict:
rows = list(
session.execute(select(Thumbnail).where(Thumbnail.asset_id == asset_id)).scalars()
)
ready = [r for r in rows if r.state == "ready" and r.path]
if ready:
best = max(ready, key=lambda r: (bool(r.protected), r.size or 0))
return {"state": "ready", "protected": bool(best.protected), "size": best.size}
if rows:
return {"state": "unsupported", "protected": False, "size": rows[0].size}
return {"state": "missing", "protected": False, "size": None}
# ── decisions ────────────────────────────────────────────────────────────
def decide(
self,

View File

@@ -13,8 +13,11 @@ renames. Every discovered or absent path is classified as one occurrence:
- ``missing`` — a known active asset whose file is gone (kept, flagged).
Missing files are never pruned (that would break identity); the asset is retained
with ``missing_at`` set. Archived assets are left untouched. Rescanning unchanged
input makes no durable change.
with ``missing_at`` set and its availability becomes ``missing_unexpected`` —
nothing explains where the bytes went. Archived assets are left untouched: their
absence from the active roots is expected, and each scan re-derives whether their
medium is reachable (:mod:`photo_pipeline.services.availability`). Rescanning
unchanged input makes no durable change.
Extracted from photo_analyzer.discover_photos/reconcile_moved/prune_missing
(see donor_ledger.yaml: pa-discovery, pa-prune-missing).
@@ -30,12 +33,12 @@ from enum import Enum
from pathlib import Path
from typing import Iterable
from sqlalchemy import func, select
from sqlalchemy import func, or_, select
from sqlalchemy.orm import Session, sessionmaker
from photo_pipeline import path_policy
from photo_pipeline.models import Asset, AssetPath
from photo_pipeline.services import hashing
from photo_pipeline.services import availability, hashing
class Occurrence(str, Enum):
@@ -60,6 +63,8 @@ def _asset_dict(asset: Asset) -> dict:
"id": asset.id,
"current_path": asset.current_path,
"availability_state": asset.availability_state,
"archive_location_id": asset.archive_location_id,
"archive_path": asset.archive_path,
"byte_size": asset.byte_size,
"current_sha256": asset.current_sha256,
"pixel_sha256": asset.pixel_sha256,
@@ -101,7 +106,9 @@ class InventoryService:
result.asset_ids[str(path)] = asset.id
for asset in assets:
if asset.availability_state != "active" or asset.id in seen_ids:
# Archived assets are explained by their location, not by the active
# roots: a scan must never prune or flag them (concept §9).
if asset.availability_state in availability.ARCHIVED or asset.id in seen_ids:
continue
if asset.current_path and asset.current_path not in discovered_paths:
if not Path(asset.current_path).exists():
@@ -109,10 +116,18 @@ class InventoryService:
asset.missing_at = now
asset.state_version += 1
asset.updated_at = now
# Nothing explains this absence — it is not an offline medium.
if asset.availability_state != availability.MISSING_UNEXPECTED:
asset.availability_state = availability.MISSING_UNEXPECTED
asset.state_version += 1
asset.updated_at = now
result.occurrences[asset.current_path] = Occurrence.MISSING.value
session.commit()
# Media may have been mounted or removed since the last scan.
availability.refresh(self._session_factory)
result.counts = dict(Counter(result.occurrences.values()))
return result
@@ -132,7 +147,11 @@ class InventoryService:
if availability:
stmt = stmt.where(Asset.availability_state == availability)
if query:
stmt = stmt.where(Asset.current_path.like(f"%{query}%"))
like = f"%{query}%"
# An archived asset has no active path; it is searched where it lives.
stmt = stmt.where(
or_(Asset.current_path.like(like), Asset.archive_path.like(like))
)
total = session.scalar(select(func.count()).select_from(stmt.subquery()))
rows = session.execute(
stmt.order_by(Asset.current_path).limit(limit).offset(offset)
@@ -169,6 +188,7 @@ class InventoryService:
self._open_path(session, existing.id, path_str, now, occ.value)
if existing.missing_at is not None:
existing.missing_at = None
existing.availability_state = availability.ACTIVE
existing.state_version += 1
existing.updated_at = now
return existing, occ
@@ -189,6 +209,7 @@ class InventoryService:
moved_from.current_path = path_str
moved_from.byte_size = size
moved_from.missing_at = None
moved_from.availability_state = availability.ACTIVE
moved_from.state_version += 1
moved_from.updated_at = now
self._open_path(session, moved_from.id, path_str, now, Occurrence.MOVED.value)

View File

@@ -8,6 +8,11 @@ an EXIF-only edit reuses the file while a real pixel change invalidates it; writ
are atomic and the cache is bounded by an LRU quota. Failures are persisted as typed
errors so a broken original is not retried on every request.
An archived asset is served from its medium when that medium is mounted, and from
its *protected* preview when it is not (US06-03). Protected previews are evidence,
not cache: the quota never evicts them, because the original they describe may be
unreachable when duplicate review needs it.
Reuses photo_analyzer.prepare_image decode/resize/HEIC handling, adding the missing
EXIF-orientation step, WebP output, and a managed cache (donor_ledger.yaml:
pa-imaging).
@@ -26,6 +31,7 @@ from sqlalchemy.orm import sessionmaker
from photo_pipeline import path_policy
from photo_pipeline.config import Config
from photo_pipeline.models import Asset, Thumbnail
from photo_pipeline.services import availability
# Best-effort HEIC support: registered only if the optional decoder is installed.
try: # pragma: no cover - depends on an optional native dependency
@@ -38,6 +44,8 @@ except Exception: # pragma: no cover
SIZES = (256, 512, 1280)
THUMB_VERSION = 1
THUMB_FORMAT = "webp"
# The size kept as durable comparison evidence for archived assets (concept §9).
PROTECTED_SIZE = 1280
class ThumbnailError(RuntimeError):
@@ -86,7 +94,12 @@ class ThumbnailService:
self._config = config
self._cache_dir = config.thumbnail_cache_dir
def generate(self, asset_id: str, size: int) -> Path:
def generate(
self, asset_id: str, size: int, *, protected: bool = False, source: Path | None = None
) -> Path:
"""Render (or reuse) a preview. ``source`` overrides where the bytes are read
from — the archiver passes its verified archive copy, which the database does
not yet point at while the transfer is still in flight."""
if size not in SIZES:
raise InvalidSize(f"size must be one of {SIZES}")
@@ -94,9 +107,7 @@ class ThumbnailService:
asset = session.get(Asset, asset_id)
if asset is None:
raise ThumbnailNotFound(f"unknown asset {asset_id}")
if asset.availability_state != "active" or not asset.current_path:
raise ThumbnailUnavailable(f"asset {asset_id} has no active file")
self._validate_path(asset.current_path)
archived = asset.availability_state in availability.ARCHIVED
cache_key = self._cache_key(asset, size)
row = session.get(Thumbnail, cache_key)
@@ -107,9 +118,18 @@ class ThumbnailService:
)
if row.path and Path(row.path).exists():
_touch(row.path)
if protected and not row.protected:
self._protect(cache_key)
return Path(row.path)
source = asset.current_path
# An archived original is read from its medium; when that medium is not
# mounted the retained preview above is the only evidence there is.
source = source or availability.readable_path(session, asset)
if source is None:
raise ThumbnailUnavailable(f"asset {asset_id} has no readable file")
source = str(source)
if source == asset.current_path:
self._validate_path(source) # archive roots lie outside the library
# Rendering happens outside the DB session (no transaction held during I/O).
try:
@@ -119,10 +139,61 @@ class ThumbnailService:
self._record_error(cache_key, asset_id, size, error.code)
raise
self._record_ready(cache_key, asset_id, size, rendered)
# Archived assets keep their preview permanently: it is the comparison
# evidence that survives the original leaving active storage.
self._record_ready(cache_key, asset_id, size, rendered, protected=protected or archived)
self._enforce_quota(keep=rendered["path"])
return Path(rendered["path"])
def ensure_protected(self, asset_id: str, *, source: Path | None = None) -> dict:
"""Produce (or confirm) the durable comparison preview for an asset.
Returns evidence rather than raising, because the caller — archive
preflight and the transfer itself — decides what an unrenderable original
means. ``unsupported`` is a recorded property of the file, not a failure of
the policy: its hashes and metadata remain the comparison evidence.
"""
try:
path = self.generate(asset_id, PROTECTED_SIZE, protected=True, source=source)
except tuple(_PERSISTED_ERRORS) as error:
return {"state": "unsupported", "error_code": error.code, "path": None}
except ThumbnailError as error:
return {"state": "unavailable", "error_code": error.code, "path": None}
return {"state": "ready", "error_code": None, "path": str(path)}
def evidence(self, asset_id: str) -> dict:
"""What durable preview this asset has right now, without rendering."""
with self._session_factory() as session:
rows = list(
session.execute(
select(Thumbnail).where(Thumbnail.asset_id == asset_id)
).scalars()
)
for row in rows:
if row.state == "ready" and row.path and Path(row.path).exists():
return {
"state": "ready",
"protected": bool(row.protected),
"size": row.size,
"error_code": None,
}
for row in rows:
if row.state == "error":
return {
"state": "unsupported",
"protected": False,
"size": row.size,
"error_code": row.error_code,
}
return {"state": "missing", "protected": False, "size": None, "error_code": None}
def _protect(self, cache_key: str) -> None:
with self._session_factory() as session:
row = session.get(Thumbnail, cache_key)
if row is not None:
row.protected = True
session.commit()
# ── path safety ──────────────────────────────────────────────────────────
def _validate_path(self, current_path: str) -> None:
path = Path(current_path)
@@ -184,7 +255,9 @@ class ThumbnailService:
}
# ── persistence ────────────────────────────────────────────────────────────
def _record_ready(self, cache_key: str, asset_id: str, size: int, rendered: dict) -> None:
def _record_ready(
self, cache_key: str, asset_id: str, size: int, rendered: dict, *, protected: bool = False
) -> None:
with self._session_factory() as session:
session.merge(
Thumbnail(
@@ -197,6 +270,7 @@ class ThumbnailService:
width=rendered["width"],
height=rendered["height"],
format=rendered["format"],
protected=protected,
)
)
try:
@@ -230,6 +304,10 @@ class ThumbnailService:
if total <= quota:
return
files.sort(key=lambda f: f.stat().st_mtime) # least-recently-used first
# Protected previews are evidence, not cache: an archived original cannot be
# re-rendered once its medium is away, so eviction never touches them.
protected = self._protected_paths()
files = [f for f in files if str(f) not in protected]
keep_path = str(Path(keep)) if keep else None
evicted: list[str] = []
for f in files:
@@ -247,6 +325,16 @@ class ThumbnailService:
if evicted:
self._forget(evicted)
def _protected_paths(self) -> set[str]:
with self._session_factory() as session:
return {
row.path
for row in session.execute(
select(Thumbnail).where(Thumbnail.protected.is_(True))
).scalars()
if row.path
}
def _forget(self, paths: list[str]) -> None:
with self._session_factory() as session:
rows = session.execute(

View File

@@ -167,7 +167,8 @@ def test_missing_file_is_flagged_not_deleted(
assets = assets_by_id(make_factory(db_url))
assert missing_id in assets # not pruned
assert assets[missing_id].missing_at is not None
assert assets[missing_id].availability_state == "active"
# Nothing explains the absence: this is not an offline archive medium (US06-03).
assert assets[missing_id].availability_state == "missing_unexpected"
def test_reappearing_file_clears_missing(
@@ -186,6 +187,7 @@ def test_reappearing_file_clears_missing(
inventory.scan(lib)
assets = assets_by_id(make_factory(db_url))
assert assets[asset_id].missing_at is None
assert assets[asset_id].availability_state == "active"
def test_identity_and_state_durable_across_restart(

View File

@@ -0,0 +1,377 @@
"""Offline identity and review evidence (US06-03).
An archived photo is not gone: it keeps its identity, its hashes, and enough
evidence to be recognised in a duplicate cluster while its medium sits in a
drawer. Every case here archives a *real* album through the real transfer, then
takes the medium away by removing its marker — the same thing the service sees
when an external disk is unplugged — and asks whether the application still tells
the truth about where the bytes are.
The distinction that matters throughout: an unmounted medium is
``archived_offline`` (expected, harmless), a mounted medium with a hole in it is
``missing_unexpected`` (needs a human). Confusing the two is how an archive
quietly loses a photo.
"""
import shutil
import uuid
from datetime import datetime, timezone
import numpy as np
import pytest
from fastapi.testclient import TestClient
from PIL import Image
from sqlalchemy import select
from photo_pipeline.api.app import create_app
from photo_pipeline.config import Config
from photo_pipeline.db import create_db_engine, create_session_factory, run_migrations
from photo_pipeline.models import Asset, Thumbnail, UploadBatch, UploadItem
from photo_pipeline.services.archive_transfer import ArchiveTransferService
from photo_pipeline.services.archives import MARKER_NAME, ArchiveService
from photo_pipeline.services.availability import (
ACTIVE,
ARCHIVED_OFFLINE,
ARCHIVED_ONLINE,
MISSING_UNEXPECTED,
)
from photo_pipeline.services.duplicates import ClusterState, DuplicateService, Method
from photo_pipeline.services.hashing import sha256_file
from photo_pipeline.services.inventory import InventoryService
from photo_pipeline.services.thumbnails import PROTECTED_SIZE, ThumbnailService, ThumbnailUnavailable
pytestmark = pytest.mark.phase_f # part of the Phase F acceptance gate (US06-06)
NOW = datetime(2026, 1, 1, tzinfo=timezone.utc)
# ── environment ──────────────────────────────────────────────────────────────
def _env(tmp_path):
(tmp_path / "data").mkdir(exist_ok=True)
lib = tmp_path / "lib"
lib.mkdir(exist_ok=True)
archive = tmp_path / "archive"
archive.mkdir(exist_ok=True)
config = Config.from_env(
{
"PHOTO_PIPELINE_DATA_DIR": str(tmp_path / "data"),
"PHOTO_PIPELINE_LIBRARY_ROOTS": str(lib),
"PHOTO_PIPELINE_ARCHIVE_FREE_SPACE_RESERVE_BYTES": "0",
}
)
run_migrations(config.database_url)
return config, create_session_factory(create_db_engine(config.database_url)), lib, archive
def structured(path, seed, size=(256, 192)):
"""A deterministic, decodable photo — previews and pHashes must be real."""
path.parent.mkdir(parents=True, exist_ok=True)
rng = np.random.default_rng(seed)
w, h = size
base = np.zeros((h, w, 3), dtype=np.uint8)
for _ in range(6):
x0 = int(rng.integers(0, w - 60))
y0 = int(rng.integers(0, h - 60))
base[y0 : y0 + 60, x0 : x0 + 60] = rng.integers(0, 256, 3)
grad = np.linspace(0, 120, w, dtype=np.uint8)
base[:, :, 0] = np.clip(base[:, :, 0].astype(int) + grad[None, :], 0, 255)
Image.fromarray(base).save(path, quality=95)
return path
def resized_copy(src, dst, scale=0.5):
with Image.open(src) as image:
image.resize(
(int(image.width * scale), int(image.height * scale)), Image.LANCZOS
).save(dst, quality=95)
return dst
def _uploaded(sf, lib, album="rome", seeds=(1, 2)):
"""A scanned album carrying the verified upload evidence archiving requires."""
folder = lib / album
for index, seed in enumerate(seeds):
structured(folder / f"{index}.jpg", seed)
result = InventoryService(sf).scan(lib)
with sf() as session:
batch_id = str(uuid.uuid4())
session.add(
UploadBatch(
id=batch_id,
album=album,
folder=str(folder),
album_name=album,
state="succeeded",
preflight_token="v1:test",
outcome_state="verified",
created_at=NOW,
)
)
for path, asset_id in result.asset_ids.items():
session.add(
UploadItem(
batch_id=batch_id,
asset_id=asset_id,
path=path,
sha256=sha256_file(path),
sha1="0" * 40,
state="sent",
outcome="uploaded",
)
)
session.commit()
return folder, result.asset_ids
def _archive(sf, config, archive, albums=None):
service = ArchiveService(sf, config=config)
location = service.register("external", str(archive))
token = service.preflight(location["id"], albums)["token"]
transfers = ArchiveTransferService(sf, config=config)
plan = transfers.create(location["id"], albums, token=token)
transfers.apply(plan["id"])
return location
def _unmount(archive):
"""Take the medium away the way a real one goes: its marker stops answering."""
(archive / MARKER_NAME).rename(archive / f"{MARKER_NAME}.away")
def _remount(archive):
(archive / f"{MARKER_NAME}.away").rename(archive / MARKER_NAME)
def _assets(sf):
with sf() as session:
return {asset.id: asset for asset in session.scalars(select(Asset))}
# ── availability ─────────────────────────────────────────────────────────────
def test_archived_assets_report_online_offline_and_missing(tmp_path):
config, sf, lib, archive = _env(tmp_path)
_, ids = _uploaded(sf, lib)
_archive(sf, config, archive)
inventory = InventoryService(sf)
states = {a.availability_state for a in _assets(sf).values()}
assert states == {ARCHIVED_ONLINE}
_unmount(archive)
inventory.scan(lib) # the album folder is gone from the active roots
assert {a.availability_state for a in _assets(sf).values()} == {ARCHIVED_OFFLINE}
assert all(a.missing_at is None for a in _assets(sf).values()) # not "missing"
_remount(archive)
inventory.scan(lib)
assert {a.availability_state for a in _assets(sf).values()} == {ARCHIVED_ONLINE}
# Mounted medium, absent file: that is not an offline archive, it needs a human.
victim = sorted(ids.values())[0]
with sf() as session:
asset = session.get(Asset, victim)
(archive / asset.archive_path).unlink()
inventory.scan(lib)
assert _assets(sf)[victim].availability_state == MISSING_UNEXPECTED
def test_scan_never_prunes_or_flags_offline_assets(tmp_path):
config, sf, lib, archive = _env(tmp_path)
_, ids = _uploaded(sf, lib)
_archive(sf, config, archive)
_unmount(archive)
before = _assets(sf)
result = InventoryService(sf).scan(lib)
assert result.counts.get("missing") is None
after = _assets(sf)
assert set(after) == set(before) == set(ids.values())
for asset in after.values():
assert asset.availability_state == ARCHIVED_OFFLINE
assert asset.current_sha256 and asset.pixel_sha256 # hashes retained
assert asset.archive_location_id and asset.archive_path
def test_active_missing_file_is_missing_unexpected_not_offline(tmp_path):
config, sf, lib, _archive_root = _env(tmp_path)
structured(lib / "loose" / "a.jpg", 7)
ids = InventoryService(sf).scan(lib).asset_ids
asset_id = next(iter(ids.values()))
(lib / "loose" / "a.jpg").unlink()
InventoryService(sf).scan(lib)
asset = _assets(sf)[asset_id]
assert asset.availability_state == MISSING_UNEXPECTED
assert asset.missing_at is not None
# ── deduplication against archived originals ─────────────────────────────────
def test_exact_copy_of_offline_asset_links_to_archived_canonical(tmp_path):
config, sf, lib, archive = _env(tmp_path)
folder, ids = _uploaded(sf, lib, seeds=(1,))
archived_id = next(iter(ids.values()))
original = next(iter(ids))
kept = tmp_path / "kept.jpg"
shutil.copy2(original, kept)
_archive(sf, config, archive)
_unmount(archive)
# The same photo turns up again in the active library while the disk is away.
(lib / "inbox").mkdir(parents=True, exist_ok=True)
shutil.copy2(kept, lib / "inbox" / "again.jpg")
scan = InventoryService(sf).scan(lib)
new_id = scan.asset_ids[str(lib / "inbox" / "again.jpg")]
assert scan.occurrences[str(lib / "inbox" / "again.jpg")] == "copied"
clusters = DuplicateService(sf).detect().clusters
exact = [c for c in clusters if c["method"] == Method.EXACT.value]
assert len(exact) == 1
cluster = exact[0]
# Byte-identical: linked directly, and the archived original stays canonical.
assert cluster["state"] == ClusterState.DECIDED.value
assert cluster["canonical_asset_id"] == archived_id
assert _assets(sf)[new_id].canonical_asset_id == archived_id
def test_fuzzy_copy_of_offline_asset_requires_review_and_names_the_medium(tmp_path):
config, sf, lib, archive = _env(tmp_path)
folder, ids = _uploaded(sf, lib, seeds=(1,))
archived_id = next(iter(ids.values()))
original = next(iter(ids))
variant_source = resized_copy(original, tmp_path / "small.jpg")
_archive(sf, config, archive)
_unmount(archive)
(lib / "inbox").mkdir(parents=True, exist_ok=True)
shutil.copy2(variant_source, lib / "inbox" / "small.jpg")
InventoryService(sf).scan(lib)
duplicates = DuplicateService(sf)
clusters = duplicates.detect().clusters
perceptual = [c for c in clusters if c["method"] == Method.PERCEPTUAL.value]
assert len(perceptual) == 1
detail = duplicates.get_cluster(perceptual[0]["id"])
assert detail["state"] == ClusterState.OPEN.value # never auto-decided
assert detail["requires_confirmation"] is True
assert detail["mount_required"] == ["external"] # full-resolution needs the disk
archived = next(m for m in detail["members"] if m["asset_id"] == archived_id)
assert archived["availability_state"] == ARCHIVED_OFFLINE
assert archived["current_path"] is None
assert archived["archive_path"] and archived["evidence"]["phash"]
# The retained preview is what makes the offline member reviewable at all.
assert archived["preview"] == {"state": "ready", "protected": True, "size": PROTECTED_SIZE}
# ── protected review evidence ────────────────────────────────────────────────
def test_protected_preview_survives_quota_and_serves_while_offline(tmp_path):
config, sf, lib, archive = _env(tmp_path)
_, ids = _uploaded(sf, lib, seeds=(1,))
asset_id = next(iter(ids.values()))
_archive(sf, config, archive)
_unmount(archive)
thumbnails = ThumbnailService(sf, config)
served = thumbnails.generate(asset_id, PROTECTED_SIZE)
assert served.exists() # rendered before the original left, not from the medium
# An aggressive quota may empty the cache, but not this evidence.
tight = ThumbnailService(sf, config.model_copy(update={"thumbnail_cache_quota_bytes": 1}))
tight._enforce_quota()
assert served.exists()
with sf() as session:
row = session.scalar(select(Thumbnail).where(Thumbnail.asset_id == asset_id))
assert row.protected is True
assert thumbnails.evidence(asset_id)["state"] == "ready"
def test_offline_asset_without_preview_reports_unavailable(tmp_path):
"""No preview and no medium is an honest 409, never a wrong picture."""
config, sf, lib, archive = _env(tmp_path)
_, ids = _uploaded(sf, lib, seeds=(1,))
asset_id = next(iter(ids.values()))
_archive(sf, config, archive)
with sf() as session:
for row in session.scalars(select(Thumbnail).where(Thumbnail.asset_id == asset_id)):
session.delete(row)
session.commit()
_unmount(archive)
with pytest.raises(ThumbnailUnavailable):
ThumbnailService(sf, config).generate(asset_id, 256)
_remount(archive) # mounted again: the archived original is readable
assert ThumbnailService(sf, config).generate(asset_id, 256).exists()
def test_offline_asset_is_browsable_through_the_api(tmp_path):
"""The browser sees an archived asset, its medium, and its preview — offline."""
config, sf, lib, archive = _env(tmp_path)
_, ids = _uploaded(sf, lib, seeds=(1,))
asset_id = next(iter(ids.values()))
_archive(sf, config, archive)
_unmount(archive)
with TestClient(create_app(config)) as client:
client.post("/api/v1/inventory/scan")
listed = client.get("/api/v1/inventory/assets", params={"availability": ARCHIVED_OFFLINE})
assert listed.status_code == 200
item = next(row for row in listed.json()["items"] if row["id"] == asset_id)
assert item["current_path"] is None
assert item["archive_path"] == "rome/0.jpg"
assert item["missing"] is False
# Searching by the archived path still finds it.
found = client.get("/api/v1/inventory/assets", params={"q": "rome"}).json()
assert [row["id"] for row in found["items"]] == [asset_id]
preview = client.get(f"/api/v1/assets/{asset_id}/thumbnail", params={"size": 1280})
assert preview.status_code == 200
assert preview.headers["content-type"] == "image/webp"
def test_offline_state_is_stable_across_restart(tmp_path):
config, sf, lib, archive = _env(tmp_path)
_, ids = _uploaded(sf, lib, seeds=(1, 2))
_archive(sf, config, archive)
_unmount(archive)
InventoryService(sf).scan(lib)
before = {
asset_id: (
asset.availability_state,
asset.archive_path,
asset.current_sha256,
asset.phash,
)
for asset_id, asset in _assets(sf).items()
}
# Restart: a fresh engine and session factory against the same database.
restarted = create_session_factory(create_db_engine(config.database_url))
after = {
asset_id: (
asset.availability_state,
asset.archive_path,
asset.current_sha256,
asset.phash,
)
for asset_id, asset in _assets(restarted).items()
}
assert after == before
assert set(after) == set(ids.values())
# And the medium coming back is picked up by the restarted process.
_remount(archive)
assert ArchiveService(restarted, config=config).locations()[0]["state"] == "online"
assert {a.availability_state for a in _assets(restarted).values()} == {ARCHIVED_ONLINE}
assert ACTIVE not in {a.availability_state for a in _assets(restarted).values()}

View File

@@ -131,6 +131,9 @@
"tests/unit/test_archive_journal_states.py",
"tests/integration/test_archive_transfer.py",
"tests/integration/test_archive_recovery.py"
],
"US06-03": [
"tests/integration/test_offline_assets.py"
]
}
}