diff --git a/migrations/versions/0013_protected_thumbnails.py b/migrations/versions/0013_protected_thumbnails.py new file mode 100644 index 0000000..043bfd9 --- /dev/null +++ b/migrations/versions/0013_protected_thumbnails.py @@ -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") diff --git a/photo_pipeline/models/thumbnails.py b/photo_pipeline/models/thumbnails.py index d9cae8e..616d44c 100644 --- a/photo_pipeline/models/thumbnails.py +++ b/photo_pipeline/models/thumbnails.py @@ -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() ) diff --git a/photo_pipeline/services/archive_transfer.py b/photo_pipeline/services/archive_transfer.py index a42a7ed..a8efa1a 100644 --- a/photo_pipeline/services/archive_transfer.py +++ b/photo_pipeline/services/archive_transfer.py @@ -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: diff --git a/photo_pipeline/services/archives.py b/photo_pipeline/services/archives.py index f7de620..36d20e8 100644 --- a/photo_pipeline/services/archives.py +++ b/photo_pipeline/services/archives.py @@ -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: diff --git a/photo_pipeline/services/availability.py b/photo_pipeline/services/availability.py new file mode 100644 index 0000000..6bf3b34 --- /dev/null +++ b/photo_pipeline/services/availability.py @@ -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) diff --git a/photo_pipeline/services/duplicates.py b/photo_pipeline/services/duplicates.py index 80ada60..2ba67b4 100644 --- a/photo_pipeline/services/duplicates.py +++ b/photo_pipeline/services/duplicates.py @@ -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, diff --git a/photo_pipeline/services/inventory.py b/photo_pipeline/services/inventory.py index 07c3a48..83ab24c 100644 --- a/photo_pipeline/services/inventory.py +++ b/photo_pipeline/services/inventory.py @@ -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) diff --git a/photo_pipeline/services/thumbnails.py b/photo_pipeline/services/thumbnails.py index 03e26a0..def3188 100644 --- a/photo_pipeline/services/thumbnails.py +++ b/photo_pipeline/services/thumbnails.py @@ -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( diff --git a/tests/integration/test_inventory_reconcile.py b/tests/integration/test_inventory_reconcile.py index a03f9a1..a2ac5dc 100644 --- a/tests/integration/test_inventory_reconcile.py +++ b/tests/integration/test_inventory_reconcile.py @@ -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( diff --git a/tests/integration/test_offline_assets.py b/tests/integration/test_offline_assets.py new file mode 100644 index 0000000..4c49499 --- /dev/null +++ b/tests/integration/test_offline_assets.py @@ -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()} diff --git a/tests/story_traceability.json b/tests/story_traceability.json index 4a35f07..144f4da 100644 --- a/tests/story_traceability.json +++ b/tests/story_traceability.json @@ -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" ] } }