Compare commits
1 Commits
chore/E09-
...
us/US06-03
| Author | SHA1 | Date | |
|---|---|---|---|
| 85fe7e6642 |
29
migrations/versions/0013_protected_thumbnails.py
Normal file
29
migrations/versions/0013_protected_thumbnails.py
Normal 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")
|
||||||
@@ -11,7 +11,7 @@ from __future__ import annotations
|
|||||||
|
|
||||||
from datetime import datetime
|
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 sqlalchemy.orm import Mapped, mapped_column
|
||||||
|
|
||||||
from photo_pipeline.db import Base
|
from photo_pipeline.db import Base
|
||||||
@@ -29,6 +29,9 @@ class Thumbnail(Base):
|
|||||||
width: Mapped[int | None] = mapped_column(Integer)
|
width: Mapped[int | None] = mapped_column(Integer)
|
||||||
height: Mapped[int | None] = mapped_column(Integer)
|
height: Mapped[int | None] = mapped_column(Integer)
|
||||||
format: Mapped[str | None] = mapped_column(String)
|
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(
|
created_at: Mapped[datetime] = mapped_column(
|
||||||
DateTime(timezone=True), nullable=False, server_default=func.now()
|
DateTime(timezone=True), nullable=False, server_default=func.now()
|
||||||
)
|
)
|
||||||
|
|||||||
@@ -60,8 +60,10 @@ from photo_pipeline.services.archive_journal import (
|
|||||||
ArchiveState,
|
ArchiveState,
|
||||||
)
|
)
|
||||||
from photo_pipeline.services.archives import MARKER_NAME, ArchiveError, ArchiveService
|
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.hashing import sha256_file
|
||||||
from photo_pipeline.services.rename_apply import PreconditionFailed, maybe_fault
|
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
|
# 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
|
# before its source is removed. It lives with the bytes so the archive can still be
|
||||||
@@ -335,6 +337,7 @@ class ArchiveTransferService:
|
|||||||
raise PreconditionFailed(
|
raise PreconditionFailed(
|
||||||
"archive_unverified", f"{destination} is not a verified archive copy"
|
"archive_unverified", f"{destination} is not a verified archive copy"
|
||||||
)
|
)
|
||||||
|
self._require_evidence(operation["asset_id"], destination)
|
||||||
if source.exists():
|
if source.exists():
|
||||||
if source.is_symlink():
|
if source.is_symlink():
|
||||||
raise PreconditionFailed("symlink", f"{source} became a 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}"
|
"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 ──────────────────────────────────────────────────────────────
|
# ── database ──────────────────────────────────────────────────────────────
|
||||||
|
|
||||||
def _record_archived(self, operation: dict, location: dict, destination: Path) -> None:
|
def _record_archived(self, operation: dict, location: dict, destination: Path) -> None:
|
||||||
|
|||||||
@@ -28,7 +28,11 @@ Blocker codes: ``no_library_root``, ``location_offline``, ``wrong_volume``,
|
|||||||
``unsafe_destination``, ``destination_not_writable``, ``manifest_unwritable``,
|
``unsafe_destination``, ``destination_not_writable``, ``manifest_unwritable``,
|
||||||
``insufficient_capacity``, ``backup_unavailable``, ``lock_conflict``,
|
``insufficient_capacity``, ``backup_unavailable``, ``lock_conflict``,
|
||||||
``rename_pending``, ``empty_scope``, ``destination_collision``,
|
``rename_pending``, ``empty_scope``, ``destination_collision``,
|
||||||
``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
|
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
|
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.path_policy import PathPolicyError, is_excluded, normalize_root, resolve_within
|
||||||
from photo_pipeline.services.albums import album_label
|
from photo_pipeline.services.albums import album_label
|
||||||
from photo_pipeline.services.archive_journal import ArchiveJournal
|
from photo_pipeline.services.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.hashing import sha256_file
|
||||||
from photo_pipeline.services.jobs import JobService
|
from photo_pipeline.services.jobs import JobService
|
||||||
from photo_pipeline.services.rename_journal import RenameJournal
|
from photo_pipeline.services.rename_journal import RenameJournal
|
||||||
|
from photo_pipeline.services.thumbnails import ThumbnailService
|
||||||
from photo_pipeline.services.upload_reports import VERIFIED
|
from photo_pipeline.services.upload_reports import VERIFIED
|
||||||
|
|
||||||
PREFLIGHT_VERSION = 1
|
PREFLIGHT_VERSION = 1
|
||||||
TOKEN_PREFIX = f"v{PREFLIGHT_VERSION}"
|
TOKEN_PREFIX = f"v{PREFLIGHT_VERSION}"
|
||||||
MARKER_NAME = ".photo-pipeline-archive.json"
|
|
||||||
MANIFEST_NAME = "archive-manifest.json"
|
MANIFEST_NAME = "archive-manifest.json"
|
||||||
|
|
||||||
# Upload outcomes that prove Immich holds these exact bytes. ``skipped``/``failed``/
|
# Upload outcomes that prove Immich holds these exact bytes. ``skipped``/``failed``/
|
||||||
@@ -163,6 +169,9 @@ class ArchiveService:
|
|||||||
location.capabilities = json.dumps(probe["capabilities"])
|
location.capabilities = json.dumps(probe["capabilities"])
|
||||||
reports.append(self._location_report(location, probe=probe))
|
reports.append(self._location_report(location, probe=probe))
|
||||||
session.commit()
|
session.commit()
|
||||||
|
# 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
|
return reports
|
||||||
|
|
||||||
# ── preflight ─────────────────────────────────────────────────────────────
|
# ── preflight ─────────────────────────────────────────────────────────────
|
||||||
@@ -420,7 +429,10 @@ class ArchiveService:
|
|||||||
|
|
||||||
def _album(self, name: str, rows: list[dict], root: Path, *, reachable: bool) -> dict:
|
def _album(self, name: str, rows: list[dict], root: Path, *, reachable: bool) -> dict:
|
||||||
folder = Path(rows[0]["path"]).parent
|
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"]]
|
blocked = [item for item in items if item["blockers"]]
|
||||||
blockers: list[dict] = []
|
blockers: list[dict] = []
|
||||||
|
|
||||||
@@ -456,6 +468,30 @@ class ArchiveService:
|
|||||||
"assets": items,
|
"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 ────────────────────────────────────────────────────────────────
|
# ── internals ────────────────────────────────────────────────────────────────
|
||||||
|
|
||||||
@@ -512,13 +548,6 @@ def _transfer_method(folder: Path, root: Path) -> str:
|
|||||||
return "copy_verify_remove"
|
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:
|
def _probe_write(path: Path, payload: bytes, *, keep: bool = False) -> str | None:
|
||||||
"""Write ``payload`` to ``path``; return the failure detail or ``None``."""
|
"""Write ``payload`` to ``path``; return the failure detail or ``None``."""
|
||||||
try:
|
try:
|
||||||
|
|||||||
124
photo_pipeline/services/availability.py
Normal file
124
photo_pipeline/services/availability.py
Normal 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)
|
||||||
@@ -11,6 +11,12 @@ Detection runs in two categories:
|
|||||||
band (NEAR/SIMILAR). These are review candidates: never decided automatically, and
|
band (NEAR/SIMILAR). These are review candidates: never decided automatically, and
|
||||||
negative-linked pairs are suppressed so a rejected pair is not re-suggested.
|
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,
|
Decisions (``canonical`` / ``not_duplicate`` / ``deferred``) persist with evidence,
|
||||||
use optimistic version checks, are reversible, and can never form a canonical cycle.
|
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
|
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 sqlalchemy.orm import sessionmaker
|
||||||
|
|
||||||
from photo_pipeline.models import (
|
from photo_pipeline.models import (
|
||||||
|
ArchiveLocation,
|
||||||
Asset,
|
Asset,
|
||||||
DuplicateCluster,
|
DuplicateCluster,
|
||||||
DuplicateMember,
|
DuplicateMember,
|
||||||
DuplicateNegativeLink,
|
DuplicateNegativeLink,
|
||||||
|
Thumbnail,
|
||||||
)
|
)
|
||||||
from photo_pipeline.services import hashing
|
from photo_pipeline.services import availability, hashing
|
||||||
|
|
||||||
NEAR_MAX = 5
|
NEAR_MAX = 5
|
||||||
SIMILAR_MAX = 10
|
SIMILAR_MAX = 10
|
||||||
@@ -127,18 +135,18 @@ class DuplicateService:
|
|||||||
|
|
||||||
# ── perceptual hash backfill ───────────────────────────────────────────
|
# ── perceptual hash backfill ───────────────────────────────────────────
|
||||||
def ensure_phashes(self) -> int:
|
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
|
updated = 0
|
||||||
with self._session_factory() as session:
|
with self._session_factory() as session:
|
||||||
assets = session.execute(
|
assets = session.execute(select(Asset)).scalars()
|
||||||
select(Asset).where(
|
|
||||||
Asset.availability_state == "active",
|
|
||||||
Asset.current_path.isnot(None),
|
|
||||||
)
|
|
||||||
).scalars()
|
|
||||||
for asset in assets:
|
for asset in assets:
|
||||||
if asset.phash is not None and asset.phash_version == hashing.PHASH_VERSION:
|
if asset.phash is not None and asset.phash_version == hashing.PHASH_VERSION:
|
||||||
continue
|
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:
|
if value is not None:
|
||||||
asset.phash = value
|
asset.phash = value
|
||||||
asset.phash_version = hashing.PHASH_VERSION
|
asset.phash_version = hashing.PHASH_VERSION
|
||||||
@@ -146,20 +154,40 @@ class DuplicateService:
|
|||||||
session.commit()
|
session.commit()
|
||||||
return updated
|
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 ──────────────────────────────────────────────────────────
|
# ── detection ──────────────────────────────────────────────────────────
|
||||||
def detect(self) -> DetectionReport:
|
def detect(self) -> DetectionReport:
|
||||||
self.ensure_phashes()
|
self.ensure_phashes()
|
||||||
now = datetime.now(timezone.utc)
|
now = datetime.now(timezone.utc)
|
||||||
report = DetectionReport()
|
report = DetectionReport()
|
||||||
with self._session_factory() as session:
|
with self._session_factory() as session:
|
||||||
assets = list(
|
# Every known asset stays in the indexes, archived or not: a copy of an
|
||||||
session.execute(
|
# archived original must be recognised as a duplicate rather than
|
||||||
select(Asset).where(
|
# treated as a new photo (concept §9, invariant 12).
|
||||||
Asset.availability_state == "active",
|
assets = list(session.execute(select(Asset)).scalars())
|
||||||
Asset.current_path.isnot(None),
|
|
||||||
)
|
|
||||||
).scalars()
|
|
||||||
)
|
|
||||||
by_id = {a.id: a for a in assets}
|
by_id = {a.id: a for a in assets}
|
||||||
negatives = {
|
negatives = {
|
||||||
_pair(link.asset_a, link.asset_b)
|
_pair(link.asset_a, link.asset_b)
|
||||||
@@ -431,10 +459,19 @@ class DuplicateService:
|
|||||||
|
|
||||||
@staticmethod
|
@staticmethod
|
||||||
def _recommend_canonical(ids, by_id) -> str:
|
def _recommend_canonical(ids, by_id) -> str:
|
||||||
# ponytail: largest file, path as deterministic tie-break. The concept's
|
# ponytail: largest file, then the archived copy, then path as a
|
||||||
# richer policy (resolution, least recompression, metadata richness) lands
|
# deterministic tie-break. Archived wins ties because it is the reviewed,
|
||||||
# with the review UI story.
|
# uploaded original — a fresh active copy must not demote it to a variant.
|
||||||
return max(ids, key=lambda i: (by_id[i].byte_size or 0, by_id[i].current_path or ""))
|
# 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):
|
def _apply_canonical(self, session, cluster, ids, canonical_id):
|
||||||
for member in session.execute(
|
for member in session.execute(
|
||||||
@@ -507,9 +544,15 @@ class DuplicateService:
|
|||||||
"current_path": asset.current_path if asset else None,
|
"current_path": asset.current_path if asset else None,
|
||||||
"byte_size": asset.byte_size if asset else None,
|
"byte_size": asset.byte_size if asset else None,
|
||||||
"phash": asset.phash if asset else None,
|
"phash": asset.phash if asset else None,
|
||||||
|
**self._offline_evidence(session, asset),
|
||||||
}
|
}
|
||||||
)
|
)
|
||||||
members.sort(key=lambda m: m["asset_id"])
|
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 {
|
return {
|
||||||
"id": cluster.id,
|
"id": cluster.id,
|
||||||
"method": cluster.method,
|
"method": cluster.method,
|
||||||
@@ -519,9 +562,54 @@ class DuplicateService:
|
|||||||
"canonical_asset_id": cluster.canonical_asset_id,
|
"canonical_asset_id": cluster.canonical_asset_id,
|
||||||
"version": cluster.version,
|
"version": cluster.version,
|
||||||
"requires_confirmation": cluster.method == Method.PERCEPTUAL.value,
|
"requires_confirmation": cluster.method == Method.PERCEPTUAL.value,
|
||||||
|
"mount_required": mount_required,
|
||||||
"members": members,
|
"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 ────────────────────────────────────────────────────────────
|
# ── decisions ────────────────────────────────────────────────────────────
|
||||||
def decide(
|
def decide(
|
||||||
self,
|
self,
|
||||||
|
|||||||
@@ -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`` — a known active asset whose file is gone (kept, flagged).
|
||||||
|
|
||||||
Missing files are never pruned (that would break identity); the asset is retained
|
Missing files are never pruned (that would break identity); the asset is retained
|
||||||
with ``missing_at`` set. Archived assets are left untouched. Rescanning unchanged
|
with ``missing_at`` set and its availability becomes ``missing_unexpected`` —
|
||||||
input makes no durable change.
|
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
|
Extracted from photo_analyzer.discover_photos/reconcile_moved/prune_missing
|
||||||
(see donor_ledger.yaml: pa-discovery, pa-prune-missing).
|
(see donor_ledger.yaml: pa-discovery, pa-prune-missing).
|
||||||
@@ -30,12 +33,12 @@ from enum import Enum
|
|||||||
from pathlib import Path
|
from pathlib import Path
|
||||||
from typing import Iterable
|
from typing import Iterable
|
||||||
|
|
||||||
from sqlalchemy import func, select
|
from sqlalchemy import func, or_, select
|
||||||
from sqlalchemy.orm import Session, sessionmaker
|
from sqlalchemy.orm import Session, sessionmaker
|
||||||
|
|
||||||
from photo_pipeline import path_policy
|
from photo_pipeline import path_policy
|
||||||
from photo_pipeline.models import Asset, AssetPath
|
from photo_pipeline.models import Asset, AssetPath
|
||||||
from photo_pipeline.services import hashing
|
from photo_pipeline.services import availability, hashing
|
||||||
|
|
||||||
|
|
||||||
class Occurrence(str, Enum):
|
class Occurrence(str, Enum):
|
||||||
@@ -60,6 +63,8 @@ def _asset_dict(asset: Asset) -> dict:
|
|||||||
"id": asset.id,
|
"id": asset.id,
|
||||||
"current_path": asset.current_path,
|
"current_path": asset.current_path,
|
||||||
"availability_state": asset.availability_state,
|
"availability_state": asset.availability_state,
|
||||||
|
"archive_location_id": asset.archive_location_id,
|
||||||
|
"archive_path": asset.archive_path,
|
||||||
"byte_size": asset.byte_size,
|
"byte_size": asset.byte_size,
|
||||||
"current_sha256": asset.current_sha256,
|
"current_sha256": asset.current_sha256,
|
||||||
"pixel_sha256": asset.pixel_sha256,
|
"pixel_sha256": asset.pixel_sha256,
|
||||||
@@ -101,7 +106,9 @@ class InventoryService:
|
|||||||
result.asset_ids[str(path)] = asset.id
|
result.asset_ids[str(path)] = asset.id
|
||||||
|
|
||||||
for asset in assets:
|
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
|
continue
|
||||||
if asset.current_path and asset.current_path not in discovered_paths:
|
if asset.current_path and asset.current_path not in discovered_paths:
|
||||||
if not Path(asset.current_path).exists():
|
if not Path(asset.current_path).exists():
|
||||||
@@ -109,10 +116,18 @@ class InventoryService:
|
|||||||
asset.missing_at = now
|
asset.missing_at = now
|
||||||
asset.state_version += 1
|
asset.state_version += 1
|
||||||
asset.updated_at = now
|
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
|
result.occurrences[asset.current_path] = Occurrence.MISSING.value
|
||||||
|
|
||||||
session.commit()
|
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()))
|
result.counts = dict(Counter(result.occurrences.values()))
|
||||||
return result
|
return result
|
||||||
|
|
||||||
@@ -132,7 +147,11 @@ class InventoryService:
|
|||||||
if availability:
|
if availability:
|
||||||
stmt = stmt.where(Asset.availability_state == availability)
|
stmt = stmt.where(Asset.availability_state == availability)
|
||||||
if query:
|
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()))
|
total = session.scalar(select(func.count()).select_from(stmt.subquery()))
|
||||||
rows = session.execute(
|
rows = session.execute(
|
||||||
stmt.order_by(Asset.current_path).limit(limit).offset(offset)
|
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)
|
self._open_path(session, existing.id, path_str, now, occ.value)
|
||||||
if existing.missing_at is not None:
|
if existing.missing_at is not None:
|
||||||
existing.missing_at = None
|
existing.missing_at = None
|
||||||
|
existing.availability_state = availability.ACTIVE
|
||||||
existing.state_version += 1
|
existing.state_version += 1
|
||||||
existing.updated_at = now
|
existing.updated_at = now
|
||||||
return existing, occ
|
return existing, occ
|
||||||
@@ -189,6 +209,7 @@ class InventoryService:
|
|||||||
moved_from.current_path = path_str
|
moved_from.current_path = path_str
|
||||||
moved_from.byte_size = size
|
moved_from.byte_size = size
|
||||||
moved_from.missing_at = None
|
moved_from.missing_at = None
|
||||||
|
moved_from.availability_state = availability.ACTIVE
|
||||||
moved_from.state_version += 1
|
moved_from.state_version += 1
|
||||||
moved_from.updated_at = now
|
moved_from.updated_at = now
|
||||||
self._open_path(session, moved_from.id, path_str, now, Occurrence.MOVED.value)
|
self._open_path(session, moved_from.id, path_str, now, Occurrence.MOVED.value)
|
||||||
|
|||||||
@@ -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
|
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.
|
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
|
Reuses photo_analyzer.prepare_image decode/resize/HEIC handling, adding the missing
|
||||||
EXIF-orientation step, WebP output, and a managed cache (donor_ledger.yaml:
|
EXIF-orientation step, WebP output, and a managed cache (donor_ledger.yaml:
|
||||||
pa-imaging).
|
pa-imaging).
|
||||||
@@ -26,6 +31,7 @@ from sqlalchemy.orm import sessionmaker
|
|||||||
from photo_pipeline import path_policy
|
from photo_pipeline import path_policy
|
||||||
from photo_pipeline.config import Config
|
from photo_pipeline.config import Config
|
||||||
from photo_pipeline.models import Asset, Thumbnail
|
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.
|
# Best-effort HEIC support: registered only if the optional decoder is installed.
|
||||||
try: # pragma: no cover - depends on an optional native dependency
|
try: # pragma: no cover - depends on an optional native dependency
|
||||||
@@ -38,6 +44,8 @@ except Exception: # pragma: no cover
|
|||||||
SIZES = (256, 512, 1280)
|
SIZES = (256, 512, 1280)
|
||||||
THUMB_VERSION = 1
|
THUMB_VERSION = 1
|
||||||
THUMB_FORMAT = "webp"
|
THUMB_FORMAT = "webp"
|
||||||
|
# The size kept as durable comparison evidence for archived assets (concept §9).
|
||||||
|
PROTECTED_SIZE = 1280
|
||||||
|
|
||||||
|
|
||||||
class ThumbnailError(RuntimeError):
|
class ThumbnailError(RuntimeError):
|
||||||
@@ -86,7 +94,12 @@ class ThumbnailService:
|
|||||||
self._config = config
|
self._config = config
|
||||||
self._cache_dir = config.thumbnail_cache_dir
|
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:
|
if size not in SIZES:
|
||||||
raise InvalidSize(f"size must be one of {SIZES}")
|
raise InvalidSize(f"size must be one of {SIZES}")
|
||||||
|
|
||||||
@@ -94,9 +107,7 @@ class ThumbnailService:
|
|||||||
asset = session.get(Asset, asset_id)
|
asset = session.get(Asset, asset_id)
|
||||||
if asset is None:
|
if asset is None:
|
||||||
raise ThumbnailNotFound(f"unknown asset {asset_id}")
|
raise ThumbnailNotFound(f"unknown asset {asset_id}")
|
||||||
if asset.availability_state != "active" or not asset.current_path:
|
archived = asset.availability_state in availability.ARCHIVED
|
||||||
raise ThumbnailUnavailable(f"asset {asset_id} has no active file")
|
|
||||||
self._validate_path(asset.current_path)
|
|
||||||
cache_key = self._cache_key(asset, size)
|
cache_key = self._cache_key(asset, size)
|
||||||
|
|
||||||
row = session.get(Thumbnail, cache_key)
|
row = session.get(Thumbnail, cache_key)
|
||||||
@@ -107,9 +118,18 @@ class ThumbnailService:
|
|||||||
)
|
)
|
||||||
if row.path and Path(row.path).exists():
|
if row.path and Path(row.path).exists():
|
||||||
_touch(row.path)
|
_touch(row.path)
|
||||||
|
if protected and not row.protected:
|
||||||
|
self._protect(cache_key)
|
||||||
return Path(row.path)
|
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).
|
# Rendering happens outside the DB session (no transaction held during I/O).
|
||||||
try:
|
try:
|
||||||
@@ -119,10 +139,61 @@ class ThumbnailService:
|
|||||||
self._record_error(cache_key, asset_id, size, error.code)
|
self._record_error(cache_key, asset_id, size, error.code)
|
||||||
raise
|
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"])
|
self._enforce_quota(keep=rendered["path"])
|
||||||
return Path(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 ──────────────────────────────────────────────────────────
|
# ── path safety ──────────────────────────────────────────────────────────
|
||||||
def _validate_path(self, current_path: str) -> None:
|
def _validate_path(self, current_path: str) -> None:
|
||||||
path = Path(current_path)
|
path = Path(current_path)
|
||||||
@@ -184,7 +255,9 @@ class ThumbnailService:
|
|||||||
}
|
}
|
||||||
|
|
||||||
# ── persistence ────────────────────────────────────────────────────────────
|
# ── 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:
|
with self._session_factory() as session:
|
||||||
session.merge(
|
session.merge(
|
||||||
Thumbnail(
|
Thumbnail(
|
||||||
@@ -197,6 +270,7 @@ class ThumbnailService:
|
|||||||
width=rendered["width"],
|
width=rendered["width"],
|
||||||
height=rendered["height"],
|
height=rendered["height"],
|
||||||
format=rendered["format"],
|
format=rendered["format"],
|
||||||
|
protected=protected,
|
||||||
)
|
)
|
||||||
)
|
)
|
||||||
try:
|
try:
|
||||||
@@ -230,6 +304,10 @@ class ThumbnailService:
|
|||||||
if total <= quota:
|
if total <= quota:
|
||||||
return
|
return
|
||||||
files.sort(key=lambda f: f.stat().st_mtime) # least-recently-used first
|
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
|
keep_path = str(Path(keep)) if keep else None
|
||||||
evicted: list[str] = []
|
evicted: list[str] = []
|
||||||
for f in files:
|
for f in files:
|
||||||
@@ -247,6 +325,16 @@ class ThumbnailService:
|
|||||||
if evicted:
|
if evicted:
|
||||||
self._forget(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:
|
def _forget(self, paths: list[str]) -> None:
|
||||||
with self._session_factory() as session:
|
with self._session_factory() as session:
|
||||||
rows = session.execute(
|
rows = session.execute(
|
||||||
|
|||||||
@@ -167,7 +167,8 @@ def test_missing_file_is_flagged_not_deleted(
|
|||||||
assets = assets_by_id(make_factory(db_url))
|
assets = assets_by_id(make_factory(db_url))
|
||||||
assert missing_id in assets # not pruned
|
assert missing_id in assets # not pruned
|
||||||
assert assets[missing_id].missing_at is not None
|
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(
|
def test_reappearing_file_clears_missing(
|
||||||
@@ -186,6 +187,7 @@ def test_reappearing_file_clears_missing(
|
|||||||
inventory.scan(lib)
|
inventory.scan(lib)
|
||||||
assets = assets_by_id(make_factory(db_url))
|
assets = assets_by_id(make_factory(db_url))
|
||||||
assert assets[asset_id].missing_at is None
|
assert assets[asset_id].missing_at is None
|
||||||
|
assert assets[asset_id].availability_state == "active"
|
||||||
|
|
||||||
|
|
||||||
def test_identity_and_state_durable_across_restart(
|
def test_identity_and_state_durable_across_restart(
|
||||||
|
|||||||
377
tests/integration/test_offline_assets.py
Normal file
377
tests/integration/test_offline_assets.py
Normal 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()}
|
||||||
@@ -131,6 +131,9 @@
|
|||||||
"tests/unit/test_archive_journal_states.py",
|
"tests/unit/test_archive_journal_states.py",
|
||||||
"tests/integration/test_archive_transfer.py",
|
"tests/integration/test_archive_transfer.py",
|
||||||
"tests/integration/test_archive_recovery.py"
|
"tests/integration/test_archive_recovery.py"
|
||||||
|
],
|
||||||
|
"US06-03": [
|
||||||
|
"tests/integration/test_offline_assets.py"
|
||||||
]
|
]
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user