Compare commits
1 Commits
us/US06-01
...
us/US05-05
| Author | SHA1 | Date | |
|---|---|---|---|
| c860c6d412 |
33
README.md
33
README.md
@@ -136,36 +136,3 @@ work_item/scripts/python -m pytest tests/e2e -m phase_d -q
|
||||
|
||||
The fault barrier is test-only configuration; without `PHOTO_PIPELINE_FAULT_AFTER`
|
||||
the apply path has no crash points. Phases A–C remain green in the full run above.
|
||||
|
||||
### Phase E acceptance gate
|
||||
|
||||
Phase E (Epic E05: Immich upload) is the one stage the application cannot take back,
|
||||
so its gate runs the fake-uploader suites, the black-box upload API journeys, and the
|
||||
browser suite as a single command:
|
||||
|
||||
```bash
|
||||
work_item/scripts/python -m pytest -m phase_e -q
|
||||
```
|
||||
|
||||
- `tests/integration/test_upload_*.py` drive a **real executable** standing in for
|
||||
`immich-go` through the real adapter and `subprocess` — argument construction,
|
||||
output bounding, report parsing, verification, and killing a running process.
|
||||
- `tests/e2e/test_phase_e_pipeline.py` drives a real server and a real durable worker
|
||||
over HTTP: credential failure and an unreachable server, preflight blockers and the
|
||||
explicitly approved partial scope, a new album, an exact duplicate, an upgrade, a
|
||||
retryable failure and its successful retry, a lost acceptance response, verification
|
||||
against Immich, an inconclusive answer resolved by an operator with evidence, bytes
|
||||
edited after upload, cancellation and resume, and an interrupted attempt recovered
|
||||
across a restart.
|
||||
- **EXIF precedes upload** is asserted, not assumed: an album without its verified
|
||||
safety and analysis checkpoints cannot be approved, and the uploader's own argv log
|
||||
proves it was never executed. Each finished upload re-hashes the files in the folder
|
||||
the uploader was handed and requires the persisted SHA-256/SHA-1 to match.
|
||||
- **No secret is retained.** The API key is a sentinel string; after a full upload and
|
||||
verification it must appear in the uploader's argv and nowhere else — not in the
|
||||
database, the retained report, or any response the browser can read.
|
||||
- `tests/e2e/test_uploads_ui.py` covers the browser journeys (preflight preview,
|
||||
confirmation, progress, stopping a run, verification, manual resolution, stale
|
||||
bytes, and recovery after a restart).
|
||||
|
||||
Phases A–D remain green in the full run above.
|
||||
|
||||
@@ -1,40 +0,0 @@
|
||||
"""Archive locations (US06-01).
|
||||
|
||||
Revision ID: 0011_archive_locations
|
||||
Revises: 0010_upload_verification
|
||||
Create Date: 2026-08-16
|
||||
|
||||
Configured archive destinations with their stable media identity, last probed
|
||||
capabilities, and state.
|
||||
"""
|
||||
|
||||
import sqlalchemy as sa
|
||||
from alembic import op
|
||||
|
||||
revision = "0011_archive_locations"
|
||||
down_revision = "0010_upload_verification"
|
||||
branch_labels = None
|
||||
depends_on = None
|
||||
|
||||
|
||||
def upgrade() -> None:
|
||||
op.create_table(
|
||||
"archive_locations",
|
||||
sa.Column("id", sa.String(), primary_key=True),
|
||||
sa.Column("name", sa.String(), nullable=False, unique=True),
|
||||
sa.Column("root", sa.String(), nullable=False),
|
||||
sa.Column("media_id", sa.String(), nullable=False, unique=True),
|
||||
sa.Column("capabilities", sa.String(), nullable=True),
|
||||
sa.Column("state", sa.String(), nullable=False, server_default="offline"),
|
||||
sa.Column("last_seen_at", sa.DateTime(timezone=True), nullable=True),
|
||||
sa.Column(
|
||||
"created_at", sa.DateTime(timezone=True), nullable=False, server_default=sa.func.now()
|
||||
),
|
||||
sa.Column(
|
||||
"updated_at", sa.DateTime(timezone=True), nullable=False, server_default=sa.func.now()
|
||||
),
|
||||
)
|
||||
|
||||
|
||||
def downgrade() -> None:
|
||||
op.drop_table("archive_locations")
|
||||
@@ -17,7 +17,6 @@ from fastapi.staticfiles import StaticFiles
|
||||
from photo_pipeline.api.routes import (
|
||||
albums,
|
||||
analysis,
|
||||
archives,
|
||||
duplicates,
|
||||
health,
|
||||
inventory,
|
||||
@@ -74,7 +73,6 @@ def create_app(config: Config | None = None) -> FastAPI:
|
||||
app.include_router(albums.router, prefix="/api/v1")
|
||||
app.include_router(renames.router, prefix="/api/v1")
|
||||
app.include_router(uploads.router, prefix="/api/v1")
|
||||
app.include_router(archives.router, prefix="/api/v1")
|
||||
# Static single-page app (hash-routed). Mounted last so /api/v1 wins.
|
||||
if FRONTEND_DIR.is_dir():
|
||||
app.mount("/app", StaticFiles(directory=FRONTEND_DIR, html=True), name="app")
|
||||
|
||||
@@ -1,63 +0,0 @@
|
||||
"""Archive location and preflight API (US06-01).
|
||||
|
||||
Registering a location writes a marker onto the medium; preflight is a command
|
||||
rather than a read, because it probes the destination, hashes the scope, and issues
|
||||
the token a later archive plan must present (US06-02). Neither endpoint moves or
|
||||
removes a single library file.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
from fastapi import APIRouter, Request
|
||||
from fastapi.responses import JSONResponse
|
||||
from pydantic import BaseModel
|
||||
|
||||
from photo_pipeline.services.archives import ArchiveError, ArchiveService
|
||||
|
||||
router = APIRouter(tags=["archives"])
|
||||
|
||||
# Which failures are the caller's request (422) and which are a missing thing (404).
|
||||
NOT_FOUND_CODES = {"unknown_location"}
|
||||
|
||||
|
||||
class RegisterLocationRequest(BaseModel):
|
||||
name: str
|
||||
root: str
|
||||
|
||||
|
||||
class PreflightRequest(BaseModel):
|
||||
location_id: str
|
||||
# ``None`` means every album; an explicit list scopes the check.
|
||||
albums: list[str] | None = None
|
||||
|
||||
|
||||
def _service(request: Request) -> ArchiveService:
|
||||
return ArchiveService(request.app.state.session_factory, config=request.app.state.config)
|
||||
|
||||
|
||||
def _error(error: ArchiveError) -> JSONResponse:
|
||||
status = 404 if error.code in NOT_FOUND_CODES else 422
|
||||
return JSONResponse(
|
||||
status_code=status, content={"error": {"code": error.code, "message": str(error)}}
|
||||
)
|
||||
|
||||
|
||||
@router.post("/archive-locations", status_code=201)
|
||||
def register_location(body: RegisterLocationRequest, request: Request):
|
||||
try:
|
||||
return _service(request).register(body.name, body.root)
|
||||
except ArchiveError as error:
|
||||
return _error(error)
|
||||
|
||||
|
||||
@router.get("/archive-locations")
|
||||
def list_locations(request: Request) -> dict:
|
||||
return {"locations": _service(request).locations()}
|
||||
|
||||
|
||||
@router.post("/archive-preflight")
|
||||
def preflight(body: PreflightRequest, request: Request):
|
||||
try:
|
||||
return _service(request).preflight(body.location_id, body.albums)
|
||||
except ArchiveError as error:
|
||||
return _error(error)
|
||||
@@ -35,9 +35,6 @@ class Config(BaseModel):
|
||||
thumbnail_cache_quota_bytes: int = 500_000_000
|
||||
thumbnail_max_pixels: int = 100_000_000
|
||||
|
||||
# Free space an archive destination must keep beyond the transfer itself.
|
||||
archive_free_space_reserve_bytes: int = 1_000_000_000
|
||||
|
||||
vision_api_key: SecretStr | None = None
|
||||
immich_api_key: SecretStr | None = None
|
||||
immich_server_url: str = ""
|
||||
|
||||
@@ -21,9 +21,6 @@ UPLOAD_BATCH = "upload_batch"
|
||||
LIBRARY_WRITE_LOCK = "library_write"
|
||||
# The uploader lane: one album batch at a time (concept §16).
|
||||
UPLOAD_LOCK = "upload"
|
||||
# The archiver lane: one archive/restore plan at a time (concept §16). No handler
|
||||
# runs on it yet (US06-02); preflight already refuses to plan around a held lease.
|
||||
ARCHIVE_LOCK = "archive"
|
||||
|
||||
|
||||
def _safety_score_item(asset_id: str, ctx: JobContext) -> None:
|
||||
|
||||
@@ -96,16 +96,6 @@ class Worker:
|
||||
return
|
||||
try:
|
||||
if cancelled or snapshot["state"] == JobState.CANCELLING:
|
||||
# A handler may stop on its own — an upload batch cancelled through
|
||||
# its own API never touches the job — so the job can still be
|
||||
# ``running`` here. Record the request before the outcome: a stop is
|
||||
# always observable as cancelling → cancelled, and ``running ->
|
||||
# cancelled`` is not a legal jump. Without this hop the transition
|
||||
# is rejected and the job keeps its lock forever.
|
||||
if snapshot["state"] == JobState.RUNNING:
|
||||
self.service.transition(
|
||||
job_id, JobState.CANCELLING, worker_id=self.worker_id, fencing_token=token
|
||||
)
|
||||
self.service.transition(
|
||||
job_id, JobState.CANCELLED, worker_id=self.worker_id, fencing_token=token
|
||||
)
|
||||
|
||||
@@ -5,7 +5,6 @@ Alembic environment relies on.
|
||||
"""
|
||||
|
||||
from photo_pipeline.models.albums import AlbumProposal
|
||||
from photo_pipeline.models.archives import ArchiveLocation
|
||||
from photo_pipeline.models.assets import Asset, AssetPath
|
||||
from photo_pipeline.models.duplicates import (
|
||||
DuplicateCluster,
|
||||
@@ -20,7 +19,6 @@ from photo_pipeline.models.workflow import AnalysisResult, SafetyReview
|
||||
|
||||
__all__ = [
|
||||
"AlbumProposal",
|
||||
"ArchiveLocation",
|
||||
"Asset",
|
||||
"AssetPath",
|
||||
"DuplicateCluster",
|
||||
|
||||
@@ -1,42 +0,0 @@
|
||||
"""Archive location persistence (US06-01).
|
||||
|
||||
An archive location is a *medium*, not a path. External disks get mounted at
|
||||
different mountpoints, and a different disk can be mounted at the same one, so a
|
||||
recorded root alone can never prove "these bytes went to that volume". Each
|
||||
location therefore owns a marker file written onto the medium itself; its
|
||||
``media_id`` is the stable identity, and the root is only where it was last seen.
|
||||
|
||||
``capabilities`` and ``state`` are the last probe result, kept so the UI can list
|
||||
locations without touching a sleeping disk. Preflight always re-probes — a stored
|
||||
state is a hint, never evidence.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
from datetime import datetime
|
||||
|
||||
from sqlalchemy import DateTime, String, func
|
||||
from sqlalchemy.orm import Mapped, mapped_column
|
||||
|
||||
from photo_pipeline.db import Base
|
||||
|
||||
|
||||
class ArchiveLocation(Base):
|
||||
__tablename__ = "archive_locations"
|
||||
|
||||
id: Mapped[str] = mapped_column(String, primary_key=True)
|
||||
name: Mapped[str] = mapped_column(String, nullable=False, unique=True)
|
||||
root: Mapped[str] = mapped_column(String, nullable=False)
|
||||
# Written into the marker file on the medium; proves the right volume is mounted.
|
||||
media_id: Mapped[str] = mapped_column(String, nullable=False, unique=True)
|
||||
capabilities: Mapped[str | None] = mapped_column(String) # JSON, last probe
|
||||
# online | offline | wrong_volume | unwritable — the last probe's verdict.
|
||||
state: Mapped[str] = mapped_column(String, nullable=False, default="offline")
|
||||
last_seen_at: Mapped[datetime | None] = mapped_column(DateTime(timezone=True))
|
||||
|
||||
created_at: Mapped[datetime] = mapped_column(
|
||||
DateTime(timezone=True), nullable=False, server_default=func.now()
|
||||
)
|
||||
updated_at: Mapped[datetime] = mapped_column(
|
||||
DateTime(timezone=True), nullable=False, server_default=func.now(), onupdate=func.now()
|
||||
)
|
||||
@@ -1,599 +0,0 @@
|
||||
"""ArchiveService — destinations and archive preflight (US06-01).
|
||||
|
||||
Archive is the only stage that *removes* originals from the active library, so
|
||||
this service does the opposite of removing anything: it registers destinations and
|
||||
proves, before a single byte moves, that an album could be archived safely. The
|
||||
transfer itself is US06-02.
|
||||
|
||||
An archive location is a medium, not a path (see :class:`ArchiveLocation`). A
|
||||
marker file on the medium carries its ``media_id``, so a disk mounted at the
|
||||
recorded root but holding a different marker is ``wrong_volume`` rather than
|
||||
silently accepted — the classic "the external disk came back at the same
|
||||
mountpoint" data-loss path.
|
||||
|
||||
Preflight proves, per concept §9 "Archive preflight":
|
||||
|
||||
- the album's upload is *verified*, not merely process-successful, and its bytes on
|
||||
disk still hash to exactly what was uploaded;
|
||||
- the destination medium is mounted, is the right one, is writable, lies outside
|
||||
every library root and every ``_IGNORE/`` tree, and has room for the scope plus a
|
||||
configured reserve;
|
||||
- nothing already occupies the destination;
|
||||
- no rename/upload/archive lease is held, and no rename is half-applied;
|
||||
- a database backup and the archive manifest can really be written — both are
|
||||
probed by writing them, not assumed.
|
||||
|
||||
Blocker codes: ``no_library_root``, ``location_offline``, ``wrong_volume``,
|
||||
``unsafe_destination``, ``destination_not_writable``, ``manifest_unwritable``,
|
||||
``insufficient_capacity``, ``backup_unavailable``, ``lock_conflict``,
|
||||
``rename_pending``, ``empty_scope``, ``destination_collision``,
|
||||
``upload_unverified``, ``bytes_changed``, ``file_missing``.
|
||||
|
||||
Like upload preflight, the confirmation token is *derived* from the report rather
|
||||
than stored: any change to the scope, the bytes, the destination, or the blockers
|
||||
produces a different token, so a stale browser confirmation can never apply. Values
|
||||
that drift without meaning anything (free space, backup size, timestamps) are left
|
||||
out of the digest.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import hashlib
|
||||
import json
|
||||
import os
|
||||
import shutil
|
||||
import sqlite3
|
||||
import uuid
|
||||
from contextlib import closing
|
||||
from datetime import datetime, timezone
|
||||
from pathlib import Path
|
||||
|
||||
from sqlalchemy import select
|
||||
from sqlalchemy.orm import sessionmaker
|
||||
|
||||
from photo_pipeline.config import Config
|
||||
from photo_pipeline.integrations import immich_go_report as report_parser
|
||||
from photo_pipeline.jobs.domain_handlers import ARCHIVE_LOCK, LIBRARY_WRITE_LOCK, UPLOAD_LOCK
|
||||
from photo_pipeline.models import ArchiveLocation, Asset, UploadBatch, UploadItem
|
||||
from photo_pipeline.path_policy import PathPolicyError, is_excluded, normalize_root, resolve_within
|
||||
from photo_pipeline.services.albums import album_label
|
||||
from photo_pipeline.services.hashing import sha256_file
|
||||
from photo_pipeline.services.jobs import JobService
|
||||
from photo_pipeline.services.rename_journal import RenameJournal
|
||||
from photo_pipeline.services.upload_reports import VERIFIED
|
||||
|
||||
PREFLIGHT_VERSION = 1
|
||||
TOKEN_PREFIX = f"v{PREFLIGHT_VERSION}"
|
||||
MARKER_NAME = ".photo-pipeline-archive.json"
|
||||
MANIFEST_NAME = "archive-manifest.json"
|
||||
|
||||
# Upload outcomes that prove Immich holds these exact bytes. ``skipped``/``failed``/
|
||||
# ``unknown`` never qualify: archiving on them would remove the only copy.
|
||||
ARCHIVED_OUTCOMES = frozenset(
|
||||
{report_parser.UPLOADED, report_parser.UPGRADED, report_parser.DUPLICATE}
|
||||
)
|
||||
|
||||
LOCKS = (LIBRARY_WRITE_LOCK, UPLOAD_LOCK, ARCHIVE_LOCK)
|
||||
|
||||
|
||||
class ArchiveError(RuntimeError):
|
||||
"""The request cannot be carried out (unknown location/album, unsafe root)."""
|
||||
|
||||
def __init__(self, code: str, message: str) -> None:
|
||||
super().__init__(message)
|
||||
self.code = code
|
||||
|
||||
|
||||
def _now() -> datetime:
|
||||
return datetime.now(timezone.utc)
|
||||
|
||||
|
||||
def _issue(code: str, message: str) -> dict:
|
||||
return {"code": code, "message": message}
|
||||
|
||||
|
||||
class ArchiveService:
|
||||
def __init__(self, session_factory: sessionmaker, *, config: Config) -> None:
|
||||
self._session_factory = session_factory
|
||||
self._config = config
|
||||
self._roots = tuple(normalize_root(root) for root in config.library_roots)
|
||||
|
||||
# ── locations ─────────────────────────────────────────────────────────────
|
||||
|
||||
def register(self, name: str, root: str) -> dict:
|
||||
"""Register an archive destination and stamp its medium with a marker.
|
||||
|
||||
The marker is what makes the location identifiable later, so registering is
|
||||
the one archive operation that writes to the destination up front.
|
||||
"""
|
||||
name = (name or "").strip()
|
||||
if not name:
|
||||
raise ArchiveError("name_required", "an archive location needs a name")
|
||||
path = Path(root).expanduser()
|
||||
if not path.is_dir():
|
||||
raise ArchiveError("root_missing", f"{path} is not an existing directory")
|
||||
path = normalize_root(path)
|
||||
unsafe = self._unsafe_destination(path)
|
||||
if unsafe:
|
||||
raise ArchiveError("unsafe_destination", unsafe)
|
||||
|
||||
marker = _read_marker(path)
|
||||
with self._session_factory() as session:
|
||||
if session.scalar(select(ArchiveLocation).where(ArchiveLocation.name == name)):
|
||||
raise ArchiveError("duplicate_name", f"an archive location named {name!r} exists")
|
||||
if marker and session.scalar(
|
||||
select(ArchiveLocation).where(ArchiveLocation.media_id == marker.get("media_id"))
|
||||
):
|
||||
raise ArchiveError(
|
||||
"already_registered", f"{path} already belongs to another archive location"
|
||||
)
|
||||
media_id = marker.get("media_id") if marker else str(uuid.uuid4())
|
||||
error = _probe_write(
|
||||
path / MARKER_NAME,
|
||||
json.dumps({"media_id": media_id, "name": name}, indent=2).encode("utf-8"),
|
||||
keep=True,
|
||||
)
|
||||
if error:
|
||||
raise ArchiveError("destination_not_writable", error)
|
||||
location = ArchiveLocation(
|
||||
id=str(uuid.uuid4()),
|
||||
name=name,
|
||||
root=str(path),
|
||||
media_id=media_id,
|
||||
state="online",
|
||||
last_seen_at=_now(),
|
||||
)
|
||||
location.capabilities = json.dumps(_capabilities(path))
|
||||
session.add(location)
|
||||
session.commit()
|
||||
return self._location_report(location, probe=_probe_location(location))
|
||||
|
||||
def locations(self) -> list[dict]:
|
||||
"""Every configured location with a fresh probe of its medium."""
|
||||
with self._session_factory() as session:
|
||||
rows = list(session.scalars(select(ArchiveLocation).order_by(ArchiveLocation.name)))
|
||||
reports = []
|
||||
for location in rows:
|
||||
probe = _probe_location(location)
|
||||
location.state = probe["state"]
|
||||
if probe["state"] == "online":
|
||||
location.last_seen_at = _now()
|
||||
location.capabilities = json.dumps(probe["capabilities"])
|
||||
reports.append(self._location_report(location, probe=probe))
|
||||
session.commit()
|
||||
return reports
|
||||
|
||||
# ── preflight ─────────────────────────────────────────────────────────────
|
||||
|
||||
def preflight(self, location_id: str, albums: list[str] | None = None) -> dict:
|
||||
"""Validate an archive scope against a destination and issue its token.
|
||||
|
||||
Read-only with respect to the library: it hashes files, probes the
|
||||
destination with its own temporary files, and writes nothing else.
|
||||
"""
|
||||
with self._session_factory() as session:
|
||||
location = session.get(ArchiveLocation, location_id)
|
||||
if location is None:
|
||||
raise ArchiveError("unknown_location", f"unknown archive location {location_id!r}")
|
||||
probe = _probe_location(location)
|
||||
location.state = probe["state"]
|
||||
if probe["state"] == "online":
|
||||
location.last_seen_at = _now()
|
||||
location.capabilities = json.dumps(probe["capabilities"])
|
||||
report = {
|
||||
"schema_version": PREFLIGHT_VERSION,
|
||||
"location": self._location_report(location, probe=probe),
|
||||
"blockers": [],
|
||||
}
|
||||
root = Path(location.root)
|
||||
session.commit()
|
||||
|
||||
report["blockers"] += self._destination_blockers(root, probe)
|
||||
report["blockers"] += self._lock_blockers()
|
||||
report["albums"] = self._albums(albums, root, reachable=probe["state"] == "online")
|
||||
report["totals"] = _totals(report["albums"])
|
||||
report["capacity"] = self._capacity(report["totals"]["bytes"], probe)
|
||||
if not report["capacity"]["sufficient"]:
|
||||
report["blockers"].append(
|
||||
_issue(
|
||||
"insufficient_capacity",
|
||||
f"{report['totals']['bytes']} B plus a "
|
||||
f"{self._config.archive_free_space_reserve_bytes} B reserve do not fit in "
|
||||
f"{report['capacity']['free_bytes']} B of free space",
|
||||
)
|
||||
)
|
||||
report["backup"] = self._backup_probe()
|
||||
if not report["backup"]["ok"]:
|
||||
report["blockers"].append(
|
||||
_issue(
|
||||
"backup_unavailable",
|
||||
f"a database backup could not be written: {report['backup']['detail']}",
|
||||
)
|
||||
)
|
||||
report["manifest"] = self._manifest_probe(
|
||||
root, report["albums"], writable=probe["writable"]
|
||||
)
|
||||
if not report["manifest"]["ok"]:
|
||||
report["blockers"].append(
|
||||
_issue(
|
||||
"manifest_unwritable",
|
||||
f"the archive manifest could not be written: {report['manifest']['detail']}",
|
||||
)
|
||||
)
|
||||
if not report["albums"]:
|
||||
report["blockers"].append(
|
||||
_issue("empty_scope", "no canonical, active assets are in the selected scope")
|
||||
)
|
||||
report["state"] = (
|
||||
"ready"
|
||||
if not report["blockers"] and all(a["state"] == "ready" for a in report["albums"])
|
||||
else "blocked"
|
||||
)
|
||||
report["token"] = _token(report)
|
||||
report["generated_at"] = _now().isoformat()
|
||||
return report
|
||||
|
||||
def verify_token(self, token: str, location_id: str, albums: list[str] | None = None) -> bool:
|
||||
"""True when ``token`` still describes this scope and this destination.
|
||||
|
||||
Recomputed, never looked up: an edited source file, a swapped medium, or a
|
||||
newly occupied destination invalidates it without anything writing to the
|
||||
database.
|
||||
"""
|
||||
return bool(token) and token == self.preflight(location_id, albums)["token"]
|
||||
|
||||
# ── destination ───────────────────────────────────────────────────────────
|
||||
|
||||
def _unsafe_destination(self, root: Path) -> str | None:
|
||||
"""Why this root may never hold archived originals, or ``None``."""
|
||||
if is_excluded(root):
|
||||
return f"{root} is inside an excluded (_IGNORE/) tree"
|
||||
for library in self._roots:
|
||||
if root == library or library in root.parents or root in library.parents:
|
||||
return f"{root} overlaps the active library root {library}"
|
||||
return None
|
||||
|
||||
def _destination_blockers(self, root: Path, probe: dict) -> list[dict]:
|
||||
blockers: list[dict] = []
|
||||
if not self._roots:
|
||||
blockers.append(_issue("no_library_root", "no library root is configured"))
|
||||
if probe["state"] == "offline":
|
||||
blockers.append(
|
||||
_issue("location_offline", f"the archive medium is not mounted at {root}")
|
||||
)
|
||||
elif probe["state"] == "wrong_volume":
|
||||
blockers.append(
|
||||
_issue(
|
||||
"wrong_volume",
|
||||
f"{root} holds a different archive medium ({probe['detail']})",
|
||||
)
|
||||
)
|
||||
unsafe = self._unsafe_destination(root)
|
||||
if unsafe:
|
||||
blockers.append(_issue("unsafe_destination", unsafe))
|
||||
if probe["state"] == "unwritable":
|
||||
blockers.append(
|
||||
_issue("destination_not_writable", f"{root} is not writable: {probe['detail']}")
|
||||
)
|
||||
return blockers
|
||||
|
||||
def _lock_blockers(self) -> list[dict]:
|
||||
"""Archive is blocked by any lease that may still be moving bytes or metadata."""
|
||||
blockers: list[dict] = []
|
||||
jobs = JobService(self._session_factory)
|
||||
for lock in LOCKS:
|
||||
held = jobs.blockers(lock)
|
||||
if held:
|
||||
blockers.append(
|
||||
_issue("lock_conflict", f"the {lock} lane is busy: job {held[0]['id']}")
|
||||
)
|
||||
if RenameJournal(self._session_factory).blocks_mutation():
|
||||
blockers.append(
|
||||
_issue("rename_pending", "an unresolved rename must be recovered before archiving")
|
||||
)
|
||||
return blockers
|
||||
|
||||
def _capacity(self, required: int, probe: dict) -> dict:
|
||||
reserve = self._config.archive_free_space_reserve_bytes
|
||||
free = probe["free_bytes"]
|
||||
return {
|
||||
"required_bytes": required,
|
||||
"reserve_bytes": reserve,
|
||||
"free_bytes": free,
|
||||
"sufficient": free is not None and free >= required + reserve,
|
||||
}
|
||||
|
||||
def _backup_probe(self) -> dict:
|
||||
"""Write a real online backup of the database, then discard it.
|
||||
|
||||
A backup that is merely assumed to be possible is worth nothing on the day
|
||||
the archive removes the originals, so this actually runs SQLite's backup API.
|
||||
"""
|
||||
source = self._config.database_path
|
||||
target = source.parent / f".archive-preflight-backup-{uuid.uuid4()}.db"
|
||||
try:
|
||||
with closing(sqlite3.connect(source)) as src, closing(sqlite3.connect(target)) as dst:
|
||||
src.backup(dst)
|
||||
size = target.stat().st_size
|
||||
except (sqlite3.Error, OSError) as error:
|
||||
return {"ok": False, "bytes": None, "detail": str(error)}
|
||||
finally:
|
||||
target.unlink(missing_ok=True)
|
||||
return {"ok": True, "bytes": size, "detail": None}
|
||||
|
||||
def _manifest_probe(self, root: Path, albums: list[dict], *, writable: bool) -> dict:
|
||||
"""Prove the manifest can be created by writing this exact content and
|
||||
removing it again. The real manifest is written by the transfer (US06-02)."""
|
||||
manifest = {
|
||||
"schema_version": PREFLIGHT_VERSION,
|
||||
"albums": [
|
||||
{
|
||||
"album": album["album"],
|
||||
"destination": album["destination"],
|
||||
"files": [
|
||||
{
|
||||
"asset_id": asset["asset_id"],
|
||||
"source": asset["current_path"],
|
||||
"sha256": asset["current_sha256"],
|
||||
"byte_size": asset["byte_size"],
|
||||
}
|
||||
for asset in album["assets"]
|
||||
],
|
||||
}
|
||||
for album in albums
|
||||
],
|
||||
}
|
||||
payload = json.dumps(manifest, indent=2, sort_keys=True).encode("utf-8")
|
||||
if not writable:
|
||||
return {"ok": False, "bytes": len(payload), "detail": "the destination is unavailable"}
|
||||
error = _probe_write(root / f".{MANIFEST_NAME}.probe-{uuid.uuid4()}", payload)
|
||||
return {"ok": error is None, "bytes": len(payload), "detail": error}
|
||||
|
||||
def _location_report(self, location: ArchiveLocation, *, probe: dict) -> dict:
|
||||
return {
|
||||
"id": location.id,
|
||||
"name": location.name,
|
||||
"root": location.root,
|
||||
"media_id": location.media_id,
|
||||
"state": probe["state"],
|
||||
"writable": probe["writable"],
|
||||
"device_id": probe["device_id"],
|
||||
"detail": probe["detail"],
|
||||
"last_seen_at": location.last_seen_at.isoformat() if location.last_seen_at else None,
|
||||
}
|
||||
|
||||
# ── scope ─────────────────────────────────────────────────────────────────
|
||||
|
||||
def _albums(self, requested: list[str] | None, root: Path, *, reachable: bool) -> list[dict]:
|
||||
by_album = self._scope()
|
||||
if requested is not None:
|
||||
unknown = sorted(set(requested) - set(by_album))
|
||||
if unknown:
|
||||
raise ArchiveError("unknown_album", f"unknown album(s): {', '.join(unknown)}")
|
||||
by_album = {name: by_album[name] for name in sorted(set(requested))}
|
||||
return [
|
||||
self._album(name, rows, root, reachable=reachable)
|
||||
for name, rows in sorted(by_album.items())
|
||||
]
|
||||
|
||||
def _scope(self) -> dict[str, list[dict]]:
|
||||
"""Canonical, active assets grouped by album, each with its upload evidence."""
|
||||
with self._session_factory() as session:
|
||||
assets = list(
|
||||
session.scalars(
|
||||
select(Asset).where(
|
||||
Asset.canonical_asset_id.is_(None),
|
||||
Asset.availability_state == "active",
|
||||
Asset.current_path.is_not(None),
|
||||
)
|
||||
)
|
||||
)
|
||||
uploads: dict[str, UploadItem] = {}
|
||||
for item, batch in session.execute(
|
||||
select(UploadItem, UploadBatch)
|
||||
.join(UploadBatch, UploadBatch.id == UploadItem.batch_id)
|
||||
.order_by(UploadBatch.created_at)
|
||||
):
|
||||
if _proves_upload(item, batch):
|
||||
uploads[item.asset_id] = item # the latest verified batch wins
|
||||
|
||||
by_album: dict[str, list[dict]] = {}
|
||||
for asset in assets:
|
||||
by_album.setdefault(album_label(asset.current_path, self._roots), []).append(
|
||||
{
|
||||
"asset_id": asset.id,
|
||||
"path": asset.current_path,
|
||||
"byte_size": asset.byte_size,
|
||||
"upload": uploads.get(asset.id),
|
||||
}
|
||||
)
|
||||
return by_album
|
||||
|
||||
def _album(self, name: str, rows: list[dict], root: Path, *, reachable: bool) -> dict:
|
||||
folder = Path(rows[0]["path"]).parent
|
||||
items = sorted((_item(row) for row in rows), key=lambda item: item["current_path"])
|
||||
blocked = [item for item in items if item["blockers"]]
|
||||
blockers: list[dict] = []
|
||||
|
||||
destination = root / name
|
||||
try:
|
||||
resolve_within(root, destination)
|
||||
except PathPolicyError as error:
|
||||
blockers.append(_issue("unsafe_destination", str(error)))
|
||||
if reachable and destination.exists() and any(destination.iterdir()):
|
||||
blockers.append(
|
||||
_issue("destination_collision", f"{destination} already exists and is not empty")
|
||||
)
|
||||
if blocked:
|
||||
blockers.append(
|
||||
_issue(
|
||||
"partial_scope",
|
||||
f"{len(blocked)} of {len(items)} asset(s) are not archivable; an album is "
|
||||
"archived whole or not at all",
|
||||
)
|
||||
)
|
||||
return {
|
||||
"album": name,
|
||||
"folder": str(folder),
|
||||
"destination": str(destination),
|
||||
# Same filesystem means the transfer can be an atomic move; anything else
|
||||
# is copy-verify-remove (concept §9).
|
||||
"transfer_method": _transfer_method(folder, root),
|
||||
"asset_count": len(items),
|
||||
"blocked_count": len(blocked),
|
||||
"reclaimable_bytes": sum(item["byte_size"] or 0 for item in items),
|
||||
"state": "blocked" if blockers else "ready",
|
||||
"blockers": blockers,
|
||||
"assets": items,
|
||||
}
|
||||
|
||||
|
||||
# ── internals ────────────────────────────────────────────────────────────────
|
||||
|
||||
|
||||
def _item(row: dict) -> dict:
|
||||
"""One asset's archivability: verified upload plus the bytes on disk right now."""
|
||||
path = Path(row["path"])
|
||||
upload: UploadItem | None = row["upload"]
|
||||
blockers: list[dict] = []
|
||||
current_sha256 = None
|
||||
|
||||
if not path.exists():
|
||||
blockers.append(_issue("file_missing", f"{path} is missing"))
|
||||
else:
|
||||
# ponytail: full re-hash of the scope. Gate on (size, mtime_ns) first if a
|
||||
# large album makes this slow — the hash stays the authority.
|
||||
current_sha256 = sha256_file(path)
|
||||
|
||||
if upload is None:
|
||||
blockers.append(
|
||||
_issue("upload_unverified", "a verified Immich upload of these bytes is required")
|
||||
)
|
||||
elif current_sha256 is not None and upload.sha256 and current_sha256 != upload.sha256:
|
||||
blockers.append(
|
||||
_issue("bytes_changed", f"{path} changed since it was uploaded; re-upload it first")
|
||||
)
|
||||
|
||||
return {
|
||||
"asset_id": row["asset_id"],
|
||||
"current_path": str(path),
|
||||
"byte_size": row["byte_size"],
|
||||
"current_sha256": current_sha256,
|
||||
"uploaded_sha256": upload.sha256 if upload else None,
|
||||
"blockers": blockers,
|
||||
}
|
||||
|
||||
|
||||
def _proves_upload(item: UploadItem, batch: UploadBatch) -> bool:
|
||||
"""Whether this upload item is evidence that Immich holds these exact bytes."""
|
||||
return (
|
||||
batch.outcome_state == VERIFIED
|
||||
and not batch.stale_bytes
|
||||
and not item.changed_after_upload
|
||||
and item.outcome in ARCHIVED_OUTCOMES
|
||||
)
|
||||
|
||||
|
||||
def _transfer_method(folder: Path, root: Path) -> str:
|
||||
try:
|
||||
if folder.stat().st_dev == root.stat().st_dev:
|
||||
return "move"
|
||||
except OSError:
|
||||
pass
|
||||
return "copy_verify_remove"
|
||||
|
||||
|
||||
def _read_marker(root: Path) -> dict | None:
|
||||
try:
|
||||
return json.loads((root / MARKER_NAME).read_text(encoding="utf-8"))
|
||||
except (OSError, ValueError):
|
||||
return None
|
||||
|
||||
|
||||
def _probe_write(path: Path, payload: bytes, *, keep: bool = False) -> str | None:
|
||||
"""Write ``payload`` to ``path``; return the failure detail or ``None``."""
|
||||
try:
|
||||
path.write_bytes(payload)
|
||||
except OSError as error:
|
||||
return str(error)
|
||||
if not keep:
|
||||
try:
|
||||
path.unlink()
|
||||
except OSError as error:
|
||||
return str(error)
|
||||
return None
|
||||
|
||||
|
||||
def _capabilities(root: Path) -> dict:
|
||||
usage = shutil.disk_usage(root)
|
||||
return {
|
||||
"device_id": root.stat().st_dev,
|
||||
"total_bytes": usage.total,
|
||||
"writable": os.access(root, os.W_OK),
|
||||
}
|
||||
|
||||
|
||||
def _probe_location(location: ArchiveLocation) -> dict:
|
||||
"""Is the right medium mounted, and can it take bytes right now?"""
|
||||
root = Path(location.root)
|
||||
blank = {"device_id": None, "free_bytes": None, "total_bytes": None, "capabilities": {}}
|
||||
if not root.is_dir():
|
||||
return {"state": "offline", "writable": False, "detail": f"{root} is not mounted", **blank}
|
||||
marker = _read_marker(root)
|
||||
if marker is None:
|
||||
return {
|
||||
"state": "offline",
|
||||
"writable": False,
|
||||
"detail": f"no archive marker found at {root}",
|
||||
**blank,
|
||||
}
|
||||
if marker.get("media_id") != location.media_id:
|
||||
return {
|
||||
"state": "wrong_volume",
|
||||
"writable": False,
|
||||
"detail": f"marker media_id {marker.get('media_id')!r}",
|
||||
**blank,
|
||||
}
|
||||
capabilities = _capabilities(root)
|
||||
usage = shutil.disk_usage(root)
|
||||
# os.access lies on some filesystems; a real write is the only proof.
|
||||
detail = _probe_write(root / f".archive-write-probe-{uuid.uuid4()}", b"")
|
||||
return {
|
||||
"state": "online" if detail is None else "unwritable",
|
||||
"writable": detail is None,
|
||||
"detail": detail,
|
||||
"device_id": capabilities["device_id"],
|
||||
"free_bytes": usage.free,
|
||||
"total_bytes": usage.total,
|
||||
"capabilities": capabilities,
|
||||
}
|
||||
|
||||
|
||||
def _totals(albums: list[dict]) -> dict:
|
||||
return {
|
||||
"albums": len(albums),
|
||||
"ready_albums": sum(1 for album in albums if album["state"] == "ready"),
|
||||
"assets": sum(album["asset_count"] for album in albums),
|
||||
"blocked": sum(album["blocked_count"] for album in albums),
|
||||
"bytes": sum(album["reclaimable_bytes"] for album in albums),
|
||||
}
|
||||
|
||||
|
||||
def _token(report: dict) -> str:
|
||||
"""Digest of everything the report asserts about the scope and the destination.
|
||||
|
||||
Values that drift without changing what would happen — free space, backup size,
|
||||
timestamps — are excluded so the same situation always yields the same token.
|
||||
"""
|
||||
payload = {key: value for key, value in report.items() if key not in ("generated_at", "token")}
|
||||
payload["location"] = {
|
||||
key: value for key, value in payload["location"].items() if key != "last_seen_at"
|
||||
}
|
||||
payload["capacity"] = {
|
||||
key: value for key, value in payload["capacity"].items() if key != "free_bytes"
|
||||
}
|
||||
payload["backup"] = {key: value for key, value in payload["backup"].items() if key != "bytes"}
|
||||
digest = hashlib.sha256(
|
||||
json.dumps(payload, sort_keys=True, ensure_ascii=False, default=str).encode("utf-8")
|
||||
).hexdigest()
|
||||
return f"{TOKEN_PREFIX}:{digest}"
|
||||
@@ -31,6 +31,4 @@ markers = [
|
||||
"phase_b: Phase B end-to-end acceptance (US02-07) — API, worker-recovery, and browser journeys",
|
||||
"phase_c: Phase C end-to-end acceptance (US03-05) — album proposal API and browser journeys",
|
||||
"phase_d: Phase D end-to-end acceptance (US04-06) — guarded rename API, fault, and browser journeys",
|
||||
"phase_e: Phase E end-to-end acceptance (US05-06) — upload preflight, uploader, and browser journeys",
|
||||
"phase_f: Phase F end-to-end acceptance (US06-06) — archive destination, transfer, and restore journeys",
|
||||
]
|
||||
|
||||
@@ -9,18 +9,14 @@ and worker, never mocked inside a test.
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import json
|
||||
import os
|
||||
import socket
|
||||
import stat
|
||||
import subprocess
|
||||
import sys
|
||||
import threading
|
||||
import time
|
||||
import uuid
|
||||
from dataclasses import dataclass
|
||||
from datetime import datetime, timezone
|
||||
from http.server import BaseHTTPRequestHandler, HTTPServer
|
||||
from pathlib import Path
|
||||
|
||||
import httpx
|
||||
@@ -262,182 +258,3 @@ def approve_album(base: str, *, album: str = "rome", name: str) -> None:
|
||||
json={**payload, "expected_version": current["version"]},
|
||||
timeout=10,
|
||||
).raise_for_status()
|
||||
|
||||
|
||||
# ── Phase E: an upload-ready album, a fake Immich, and a real fake uploader ───
|
||||
|
||||
SENTINEL_KEY = "immich-sentinel-9f3a2b"
|
||||
UPLOADER_VERSION = "immich-go 0.21.0" # a pinned family, so reports are parsable
|
||||
|
||||
# Uploader bodies for the pinned ``text-v1`` grammar. ``$6`` is the folder argument
|
||||
# of ``upload from-folder``.
|
||||
REPORTING_UPLOADER = (
|
||||
'echo "INFO uploaded $6/a.jpg"\n'
|
||||
'echo "INFO server has the same file $6/b.jpg"\n'
|
||||
'echo "Uploaded 1, duplicates 1"\n'
|
||||
"exit 0\n"
|
||||
)
|
||||
# Exits cleanly but says nothing about any file: the process succeeded, the
|
||||
# per-file outcome is unknown.
|
||||
SILENT_UPLOADER = "exit 0\n"
|
||||
|
||||
|
||||
def _immich_handler(state: dict):
|
||||
class Handler(BaseHTTPRequestHandler):
|
||||
def do_GET(self): # noqa: N802 (BaseHTTPRequestHandler API)
|
||||
self._json(200, {"res": "pong"})
|
||||
|
||||
def do_POST(self): # noqa: N802
|
||||
length = int(self.headers.get("Content-Length", 0))
|
||||
payload = json.loads(self.rfile.read(length) or b"{}")
|
||||
if state["mode"] == "broken":
|
||||
self.send_error(500, "bulk-upload-check is unavailable")
|
||||
return
|
||||
reject = state["mode"] == "present"
|
||||
self._json(
|
||||
200,
|
||||
{
|
||||
"results": [
|
||||
{
|
||||
"id": asset["id"],
|
||||
"action": "reject" if reject else "accept",
|
||||
"reason": "duplicate" if reject else None,
|
||||
}
|
||||
for asset in payload.get("assets", [])
|
||||
]
|
||||
},
|
||||
)
|
||||
|
||||
def _json(self, code: int, body: dict) -> None:
|
||||
raw = json.dumps(body).encode()
|
||||
self.send_response(code)
|
||||
self.send_header("Content-Type", "application/json")
|
||||
self.send_header("Content-Length", str(len(raw)))
|
||||
self.end_headers()
|
||||
self.wfile.write(raw)
|
||||
|
||||
def log_message(self, *args):
|
||||
pass
|
||||
|
||||
return Handler
|
||||
|
||||
|
||||
class FakeImmich:
|
||||
"""An Immich that answers ping, and says whether it holds the exact bytes.
|
||||
|
||||
``mode`` is what the next verification will find: ``present`` (the server
|
||||
deduplicates them, so it has them), ``absent`` (it would accept them, so it does
|
||||
not), or ``broken`` (no usable answer at all).
|
||||
"""
|
||||
|
||||
def __init__(self) -> None:
|
||||
self.state = {"mode": "present"}
|
||||
self._server = HTTPServer(("127.0.0.1", 0), _immich_handler(self.state))
|
||||
threading.Thread(target=self._server.serve_forever, daemon=True).start()
|
||||
self.url = f"http://127.0.0.1:{self._server.server_port}"
|
||||
self._running = True
|
||||
|
||||
def mode(self, mode: str) -> None:
|
||||
self.state["mode"] = mode
|
||||
|
||||
def stop(self) -> None:
|
||||
"""Idempotent, so a test may take Immich away mid-journey."""
|
||||
if not self._running:
|
||||
return
|
||||
self._running = False
|
||||
self._server.shutdown()
|
||||
self._server.server_close()
|
||||
|
||||
|
||||
def fake_uploader(tmp_path: Path, body: str) -> Path:
|
||||
"""A real executable standing in for immich-go.
|
||||
|
||||
``--version`` answers like the real tool; any other invocation appends its
|
||||
complete argv to ``immich-go.argv`` — which is how a test proves the uploader
|
||||
ran, what folder it was handed, or that it never ran at all.
|
||||
"""
|
||||
path = tmp_path / "immich-go"
|
||||
path.write_text(
|
||||
"#!/bin/sh\n"
|
||||
f'if [ "$1" = "--version" ]; then echo "{UPLOADER_VERSION}"; exit 0; fi\n'
|
||||
f'printf "%s\\n" "$*" >> "{tmp_path / "immich-go.argv"}"\n'
|
||||
f"{body}"
|
||||
)
|
||||
path.chmod(path.stat().st_mode | stat.S_IEXEC | stat.S_IXGRP | stat.S_IXOTH)
|
||||
return path
|
||||
|
||||
|
||||
def uploader_argv(tmp_path: Path) -> list[str]:
|
||||
"""Every upload invocation the fake uploader saw, oldest first."""
|
||||
log = tmp_path / "immich-go.argv"
|
||||
return log.read_text().splitlines() if log.exists() else []
|
||||
|
||||
|
||||
def mark_upload_ready(seeded: Seeded, *, unverified: tuple[str, ...] = ()) -> None:
|
||||
"""Give every seeded photo the verified EXIF checkpoints upload requires.
|
||||
|
||||
``unverified`` names stems whose analysis checkpoint stays incomplete, which is
|
||||
what makes an album partially blocked.
|
||||
"""
|
||||
from sqlalchemy import select
|
||||
|
||||
from photo_pipeline.models import AnalysisResult, SafetyReview
|
||||
|
||||
blocked = {seeded.asset_ids[stem] for stem in unverified}
|
||||
with session_factory(seeded) as sf:
|
||||
with sf() as session:
|
||||
for review in session.scalars(select(SafetyReview)):
|
||||
review.exif_verified_at = NOW
|
||||
for analysis in session.scalars(select(AnalysisResult)):
|
||||
analysis.exif_written_at = None if analysis.asset_id in blocked else NOW
|
||||
session.commit()
|
||||
|
||||
|
||||
class UploadStack:
|
||||
"""A seeded, upload-ready library plus the server, worker, and fake Immich."""
|
||||
|
||||
def __init__(self, tmp_path: Path, seeded: Seeded) -> None:
|
||||
self.tmp_path = tmp_path
|
||||
self.seeded = seeded
|
||||
self.immich = FakeImmich()
|
||||
self.server: Server | None = None
|
||||
self.worker: subprocess.Popen | None = None
|
||||
self.base = ""
|
||||
|
||||
def start(
|
||||
self,
|
||||
*,
|
||||
uploader: str = REPORTING_UPLOADER,
|
||||
worker: bool = True,
|
||||
credentials: bool = True,
|
||||
) -> "UploadStack":
|
||||
env = {
|
||||
"PHOTO_PIPELINE_IMMICH_SERVER_URL": self.immich.url if credentials else "",
|
||||
"PHOTO_PIPELINE_IMMICH_GO_BINARY": str(fake_uploader(self.tmp_path, uploader)),
|
||||
}
|
||||
if credentials:
|
||||
env["PHOTO_PIPELINE_IMMICH_API_KEY"] = SENTINEL_KEY
|
||||
self.server = Server(self.seeded, extra_env=env).start()
|
||||
self.base = self.server.base
|
||||
if worker:
|
||||
self.worker = start_worker(self.seeded, extra_env=env)
|
||||
return self
|
||||
|
||||
def restart_server(self) -> None:
|
||||
"""A genuinely fresh process against the same database and library."""
|
||||
self.server.stop()
|
||||
self.server.start()
|
||||
|
||||
def batches(self) -> list[dict]:
|
||||
return httpx.get(f"{self.base}/api/v1/upload-batches", timeout=20).json()["batches"]
|
||||
|
||||
def argv(self) -> list[str]:
|
||||
return uploader_argv(self.tmp_path)
|
||||
|
||||
def stop(self) -> None:
|
||||
if self.worker is not None:
|
||||
self.worker.kill()
|
||||
self.worker.wait(timeout=10)
|
||||
if self.server is not None:
|
||||
self.server.stop()
|
||||
self.immich.stop()
|
||||
|
||||
@@ -1,568 +0,0 @@
|
||||
"""Phase E end-to-end acceptance (US05-06): Immich upload, black box.
|
||||
|
||||
Every journey drives a real ``photo_pipeline serve`` child process and a real durable
|
||||
worker over HTTP — preflight, approve, upload, duplicate, upgrade, fail, retry, lose
|
||||
the acceptance response, verify, cancel, crash, restart. Nothing external is mocked
|
||||
inside the application: ``immich-go`` is a real executable that records the argv it
|
||||
was handed, and Immich is a real HTTP server answering the same ``ping`` and
|
||||
``bulk-upload-check`` endpoints the adapter calls in production.
|
||||
|
||||
Two invariants are asserted in every relevant journey, because they are what make an
|
||||
irreversible stage safe:
|
||||
|
||||
- **EXIF precedes upload.** An album without its verified safety and analysis
|
||||
checkpoints cannot be approved, and the uploader's argv log proves it was never
|
||||
even executed.
|
||||
- **The persisted hashes are the submitted bytes.** After each upload the recorded
|
||||
SHA-256/SHA-1 of every item is recomputed from the files in the folder the uploader
|
||||
was actually given.
|
||||
|
||||
The API key is a sentinel string, so the last journey can prove it reached the
|
||||
uploader and nothing else that was retained.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import hashlib
|
||||
from pathlib import Path
|
||||
|
||||
import httpx
|
||||
import pytest
|
||||
|
||||
from tests.e2e._pipeline_harness import (
|
||||
SENTINEL_KEY,
|
||||
SILENT_UPLOADER,
|
||||
UploadStack,
|
||||
mark_upload_ready,
|
||||
seed_album,
|
||||
session_factory,
|
||||
wait_until,
|
||||
)
|
||||
|
||||
pytestmark = pytest.mark.phase_e
|
||||
|
||||
TIMEOUT = 20
|
||||
ALBUM = "rome"
|
||||
TERMINAL = {"succeeded", "failed", "cancelled", "unknown_requires_verification"}
|
||||
|
||||
# Uploader bodies in the pinned ``text-v1`` grammar. ``$6`` is the folder argument of
|
||||
# ``upload from-folder``, so each line names the real path of a real file.
|
||||
ALL_NEW = (
|
||||
'echo "INFO uploaded $6/a.jpg"\n'
|
||||
'echo "INFO uploaded $6/b.jpg"\n'
|
||||
'echo "Uploaded 2"\n'
|
||||
"exit 0\n"
|
||||
)
|
||||
EXACT_DUPLICATES = (
|
||||
'echo "INFO server has the same file $6/a.jpg"\n'
|
||||
'echo "INFO server has the same file $6/b.jpg"\n'
|
||||
'echo "Duplicates 2"\n'
|
||||
"exit 0\n"
|
||||
)
|
||||
UPGRADES = (
|
||||
'echo "INFO server has an older file $6/a.jpg"\n'
|
||||
'echo "INFO server has an older file $6/b.jpg"\n'
|
||||
'echo "Upgraded 2"\n'
|
||||
"exit 0\n"
|
||||
)
|
||||
|
||||
|
||||
def _once_then(tmp_path: Path, first: str, rest: str) -> str:
|
||||
"""An uploader that behaves one way on its first attempt and another afterwards.
|
||||
|
||||
The flag file is the attempt counter, so retry and resume journeys are
|
||||
deterministic without any test reaching into the running application.
|
||||
"""
|
||||
flag = tmp_path / "first-attempt.flag"
|
||||
return f'if [ ! -f "{flag}" ]; then\n touch "{flag}"\n{first}fi\n{rest}'
|
||||
|
||||
|
||||
FAILING_FIRST = (' echo "ERROR error uploading $6/a.jpg: connection reset"\n exit 1\n', ALL_NEW)
|
||||
SLOW_FIRST = (' echo "INFO starting"\n sleep 30\n exit 0\n', ALL_NEW)
|
||||
|
||||
|
||||
@pytest.fixture
|
||||
def stack(tmp_path):
|
||||
"""An analysed album whose EXIF checkpoints are already verified."""
|
||||
seeded = seed_album(tmp_path)
|
||||
mark_upload_ready(seeded)
|
||||
running = UploadStack(tmp_path, seeded)
|
||||
try:
|
||||
yield running
|
||||
finally:
|
||||
running.stop()
|
||||
|
||||
|
||||
@pytest.fixture
|
||||
def unfinished(tmp_path):
|
||||
"""The same album *before* its EXIF checkpoints were written."""
|
||||
seeded = seed_album(tmp_path)
|
||||
running = UploadStack(tmp_path, seeded)
|
||||
try:
|
||||
yield running
|
||||
finally:
|
||||
running.stop()
|
||||
|
||||
|
||||
# ── helpers ──────────────────────────────────────────────────────────────────
|
||||
|
||||
|
||||
def _preflight(base: str, **body) -> dict:
|
||||
response = httpx.post(f"{base}/api/v1/upload-preflight", json=body, timeout=TIMEOUT)
|
||||
response.raise_for_status()
|
||||
return response.json()
|
||||
|
||||
|
||||
def _create(base: str, report: dict, **body) -> httpx.Response:
|
||||
return httpx.post(
|
||||
f"{base}/api/v1/upload-batches",
|
||||
json={"token": report["token"], **body},
|
||||
timeout=TIMEOUT,
|
||||
)
|
||||
|
||||
|
||||
def _get(base: str, batch_id: str) -> dict:
|
||||
return httpx.get(f"{base}/api/v1/upload-batches/{batch_id}", timeout=TIMEOUT).json()
|
||||
|
||||
|
||||
def _start(base: str, batch_id: str) -> httpx.Response:
|
||||
return httpx.post(f"{base}/api/v1/upload-batches/{batch_id}/start", timeout=TIMEOUT)
|
||||
|
||||
|
||||
def _start_accepted(base: str, batch_id: str) -> httpx.Response:
|
||||
"""Start, waiting out the uploader lane the previous attempt still holds.
|
||||
|
||||
A stopped attempt releases its job a moment after the batch itself reaches
|
||||
``cancelled``; ``lock_held`` is that gap, not a refusal of this batch.
|
||||
"""
|
||||
|
||||
def _attempt():
|
||||
response = _start(base, batch_id)
|
||||
if response.status_code == 409 and response.json()["error"]["code"] == "lock_held":
|
||||
return None
|
||||
response.raise_for_status()
|
||||
return response
|
||||
|
||||
return wait_until(_attempt)
|
||||
|
||||
|
||||
def _verify(base: str, batch_id: str) -> dict:
|
||||
response = httpx.post(f"{base}/api/v1/upload-batches/{batch_id}/verify", timeout=TIMEOUT)
|
||||
response.raise_for_status()
|
||||
return response.json()
|
||||
|
||||
|
||||
def _await_state(base: str, batch_id: str, states: set[str], *, timeout: float = 60) -> dict:
|
||||
return wait_until(
|
||||
lambda: (lambda b: b if b.get("state") in states else None)(_get(base, batch_id)),
|
||||
timeout=timeout,
|
||||
)
|
||||
|
||||
|
||||
def _await_report(base: str, batch_id: str, states: set[str] = TERMINAL) -> dict:
|
||||
"""Wait for a finished attempt *and* the report that explains it.
|
||||
|
||||
The batch state is recorded a moment before its report is parsed — the outcome of
|
||||
the process and the outcome of each file are deliberately separate facts — so a
|
||||
journey that reads per-item evidence must wait for the second one too.
|
||||
"""
|
||||
return wait_until(
|
||||
lambda: (lambda b: b if b.get("state") in states and b.get("parsed_at") else None)(
|
||||
_get(base, batch_id)
|
||||
),
|
||||
timeout=60,
|
||||
)
|
||||
|
||||
|
||||
def _approve(stack, **body) -> dict:
|
||||
"""Preflight, approve exactly that report, and return the created batch."""
|
||||
report = _preflight(stack.base, **body)
|
||||
response = _create(stack.base, report, **body)
|
||||
response.raise_for_status()
|
||||
return response.json()["batches"][0]
|
||||
|
||||
|
||||
def _upload(stack, **body) -> dict:
|
||||
"""The whole approved journey, up to whatever terminal state it reaches."""
|
||||
batch = _approve(stack, **body)
|
||||
_start(stack.base, batch["id"]).raise_for_status()
|
||||
return _await_report(stack.base, batch["id"])
|
||||
|
||||
|
||||
def _outcomes(batch: dict) -> dict[str, str]:
|
||||
return {Path(item["path"]).name: item["outcome"] for item in batch["items"]}
|
||||
|
||||
|
||||
def _uploaded_folder(stack) -> Path:
|
||||
"""The folder the uploader was actually handed, from its own argv log."""
|
||||
invocations = stack.argv()
|
||||
assert invocations, "the uploader was never executed"
|
||||
return Path(invocations[-1].split()[-1])
|
||||
|
||||
|
||||
def _assert_hashes_match_submitted_bytes(stack, batch: dict) -> None:
|
||||
folder = _uploaded_folder(stack)
|
||||
for item in batch["items"]:
|
||||
submitted = folder / Path(item["path"]).name
|
||||
raw = submitted.read_bytes()
|
||||
assert item["sha256"] == hashlib.sha256(raw).hexdigest(), submitted
|
||||
assert item["sha1"] == hashlib.sha1(raw).hexdigest(), submitted # noqa: S324 — Immich's
|
||||
|
||||
|
||||
# ── credentials ──────────────────────────────────────────────────────────────
|
||||
|
||||
|
||||
def test_missing_credentials_block_the_preflight_and_no_upload_is_attempted(stack):
|
||||
stack.start(worker=False, credentials=False)
|
||||
|
||||
report = _preflight(stack.base)
|
||||
|
||||
assert report["state"] == "blocked"
|
||||
assert [issue["code"] for issue in report["blockers"]] == ["credentials_missing"]
|
||||
assert report["credentials"]["api_key_configured"] is False
|
||||
# A blocked scope still issues a token; approving it is what is refused.
|
||||
refused = _create(stack.base, report)
|
||||
assert refused.status_code == 422
|
||||
assert refused.json()["error"]["code"] == "not_ready"
|
||||
assert stack.batches() == []
|
||||
assert stack.argv() == [], "the uploader must not run without credentials"
|
||||
|
||||
|
||||
def test_a_server_that_stops_answering_blocks_the_preflight(stack):
|
||||
stack.start(worker=False)
|
||||
assert _preflight(stack.base)["state"] == "ready"
|
||||
|
||||
stack.immich.stop() # Immich goes away between one preview and the next
|
||||
|
||||
report = _preflight(stack.base)
|
||||
assert report["state"] == "blocked"
|
||||
assert [issue["code"] for issue in report["blockers"]] == ["server_unreachable"]
|
||||
assert report["server"]["reachable"] is False
|
||||
assert _create(stack.base, report).status_code == 422
|
||||
assert stack.argv() == []
|
||||
|
||||
|
||||
# ── EXIF precedes upload ─────────────────────────────────────────────────────
|
||||
|
||||
|
||||
def test_exif_checkpoints_must_be_written_before_anything_is_uploaded(unfinished):
|
||||
stack = unfinished
|
||||
stack.start()
|
||||
|
||||
report = _preflight(stack.base)
|
||||
assert report["state"] == "blocked"
|
||||
codes = {
|
||||
issue["code"]
|
||||
for album in report["albums"]
|
||||
for asset in album["assets"]
|
||||
for issue in asset["blockers"]
|
||||
}
|
||||
assert codes == {"safety_exif_unverified"}
|
||||
assert _create(stack.base, report).status_code == 422
|
||||
assert stack.argv() == [], "the uploader must not run before the EXIF checkpoints"
|
||||
|
||||
mark_upload_ready(stack.seeded) # the checkpoints are written
|
||||
|
||||
batch = _upload(stack)
|
||||
assert batch["state"] == "succeeded"
|
||||
assert stack.argv(), "the same scope uploads once its checkpoints exist"
|
||||
# The recorded checkpoints predate the attempt that was allowed to run.
|
||||
from sqlalchemy import select
|
||||
|
||||
from photo_pipeline.models import AnalysisResult, SafetyReview
|
||||
|
||||
with session_factory(stack.seeded) as sf, sf() as session:
|
||||
checkpoints = [
|
||||
*[row.exif_verified_at for row in session.scalars(select(SafetyReview))],
|
||||
*[row.exif_written_at for row in session.scalars(select(AnalysisResult))],
|
||||
]
|
||||
started = batch["started_at"]
|
||||
assert checkpoints and all(stamp.isoformat() < started for stamp in checkpoints)
|
||||
|
||||
|
||||
def test_an_unfinished_photo_blocks_its_album_until_a_partial_upload_is_approved(tmp_path):
|
||||
seeded = seed_album(tmp_path)
|
||||
mark_upload_ready(seeded, unverified=("b",))
|
||||
stack = UploadStack(tmp_path, seeded)
|
||||
try:
|
||||
stack.start()
|
||||
report = _preflight(stack.base)
|
||||
assert report["state"] == "blocked"
|
||||
assert [issue["code"] for issue in report["albums"][0]["blockers"]] == ["partial_scope"]
|
||||
assert report["totals"] == {
|
||||
"albums": 1,
|
||||
"ready_albums": 0,
|
||||
"assets": 2,
|
||||
"eligible": 1,
|
||||
"blocked": 1,
|
||||
}
|
||||
assert _create(stack.base, report).status_code == 422
|
||||
|
||||
partial = _preflight(stack.base, allow_partial=True)
|
||||
assert partial["state"] == "ready"
|
||||
assert partial["token"] != report["token"]
|
||||
# A full-scope approval can never be replayed as a partial one.
|
||||
replayed = _create(stack.base, report, allow_partial=True)
|
||||
assert replayed.status_code == 409
|
||||
assert replayed.json()["error"]["code"] == "stale_preflight"
|
||||
|
||||
batch = _upload(stack, allow_partial=True)
|
||||
assert batch["state"] == "succeeded"
|
||||
assert batch["allow_partial"] is True
|
||||
assert [Path(item["path"]).name for item in batch["items"]] == ["a.jpg"]
|
||||
finally:
|
||||
stack.stop()
|
||||
|
||||
|
||||
# ── outcomes ─────────────────────────────────────────────────────────────────
|
||||
|
||||
|
||||
def test_a_new_album_uploads_and_persists_the_hashes_that_were_submitted(stack):
|
||||
stack.start(uploader=ALL_NEW)
|
||||
|
||||
batch = _upload(stack)
|
||||
|
||||
assert batch["state"] == "succeeded"
|
||||
assert batch["outcome_state"] == "verified"
|
||||
assert _outcomes(batch) == {"a.jpg": "uploaded", "b.jpg": "uploaded"}
|
||||
assert batch["outcome_counts"]["uploaded"] == 2
|
||||
assert batch["report_counts"] == {"uploaded": 2}
|
||||
assert batch["parser"] == "text-v1"
|
||||
_assert_hashes_match_submitted_bytes(stack, batch)
|
||||
# One album, one invocation, scoped to that album's own folder.
|
||||
assert len(stack.argv()) == 1
|
||||
assert f"--album-name={ALBUM}" in stack.argv()[0]
|
||||
assert _uploaded_folder(stack) == stack.seeded.lib / ALBUM
|
||||
|
||||
|
||||
def test_an_exact_duplicate_is_recorded_as_a_duplicate_not_a_new_asset(stack):
|
||||
stack.start(uploader=EXACT_DUPLICATES)
|
||||
|
||||
batch = _upload(stack)
|
||||
|
||||
assert batch["state"] == "succeeded"
|
||||
assert batch["outcome_state"] == "verified"
|
||||
assert set(_outcomes(batch).values()) == {"duplicate"}
|
||||
assert batch["outcome_counts"]["uploaded"] == 0
|
||||
_assert_hashes_match_submitted_bytes(stack, batch)
|
||||
|
||||
|
||||
def test_a_better_copy_is_recorded_as_an_upgrade(stack):
|
||||
stack.start(uploader=UPGRADES)
|
||||
|
||||
batch = _upload(stack)
|
||||
|
||||
assert batch["state"] == "succeeded"
|
||||
assert batch["outcome_state"] == "verified"
|
||||
assert set(_outcomes(batch).values()) == {"upgraded"}
|
||||
assert batch["report_counts"] == {"upgraded": 2}
|
||||
|
||||
|
||||
# ── failure and retry ────────────────────────────────────────────────────────
|
||||
|
||||
|
||||
def test_a_plain_uploader_failure_is_retryable_and_the_retry_succeeds(stack, tmp_path):
|
||||
stack.start(uploader=_once_then(tmp_path, *FAILING_FIRST))
|
||||
|
||||
failed = _upload(stack)
|
||||
|
||||
assert failed["state"] == "failed"
|
||||
assert failed["error_code"] == "uploader_failed"
|
||||
assert failed["exit_code"] == 1
|
||||
assert _outcomes(failed)["a.jpg"] == "failed"
|
||||
# Nothing uncertain happened, so the batch is offered again rather than blocked.
|
||||
assert failed["retry_blockers"] == []
|
||||
|
||||
_start_accepted(stack.base, failed["id"])
|
||||
retried = _await_report(stack.base, failed["id"], {"succeeded"})
|
||||
|
||||
assert retried["attempt_count"] == 2
|
||||
assert _outcomes(retried) == {"a.jpg": "uploaded", "b.jpg": "uploaded"}
|
||||
assert len(stack.argv()) == 2
|
||||
|
||||
|
||||
# ── uncertainty ──────────────────────────────────────────────────────────────
|
||||
|
||||
|
||||
def test_a_lost_acceptance_response_stays_uncertain_and_is_never_retried_blindly(stack):
|
||||
stack.start(uploader=SILENT_UPLOADER)
|
||||
|
||||
batch = _upload(stack)
|
||||
|
||||
# The process succeeded; what happened to each file is simply not known.
|
||||
assert batch["state"] == "succeeded"
|
||||
assert batch["outcome_state"] == "requires_verification"
|
||||
assert set(_outcomes(batch).values()) == {"unknown"}
|
||||
refused = _start(stack.base, batch["id"])
|
||||
assert refused.status_code == 409
|
||||
assert refused.json()["error"]["code"] == "not_runnable"
|
||||
assert len(stack.argv()) == 1, "a blind retry must not reach the uploader"
|
||||
|
||||
|
||||
def test_verification_asks_immich_and_resolves_every_uncertain_item(stack):
|
||||
stack.start(uploader=SILENT_UPLOADER)
|
||||
stack.immich.mode("present") # Immich holds exactly the bytes that were sent
|
||||
batch = _upload(stack)
|
||||
|
||||
result = _verify(stack.base, batch["id"])
|
||||
|
||||
assert result["counts"] == {"present": 2}
|
||||
assert result["outcome_state"] == "verified"
|
||||
verified = _get(stack.base, batch["id"])
|
||||
assert set(_outcomes(verified).values()) == {"uploaded"}
|
||||
history = httpx.get(
|
||||
f"{stack.base}/api/v1/upload-batches/{batch['id']}/verifications", timeout=TIMEOUT
|
||||
).json()["verifications"]
|
||||
assert {entry["source"] for entry in history} == {"immich_api"}
|
||||
|
||||
|
||||
def test_an_unusable_answer_stays_uncertain_until_an_operator_records_evidence(stack):
|
||||
stack.start(uploader=SILENT_UPLOADER)
|
||||
stack.immich.mode("broken") # answers, but nothing this adapter will interpret
|
||||
batch = _upload(stack)
|
||||
|
||||
result = _verify(stack.base, batch["id"])
|
||||
|
||||
# No answer is never "no": the items stay uncertain rather than being called failed.
|
||||
assert result["counts"] == {"inconclusive": 2}
|
||||
assert result["outcome_state"] == "requires_verification"
|
||||
assert set(_outcomes(_get(stack.base, batch["id"])).values()) == {"unknown"}
|
||||
|
||||
asset_id = batch["items"][0]["asset_id"]
|
||||
incomplete = httpx.post(
|
||||
f"{stack.base}/api/v1/upload-batches/{batch['id']}/resolve",
|
||||
json={"asset_id": asset_id, "outcome": "uploaded", "evidence": "", "actor": "dom"},
|
||||
timeout=TIMEOUT,
|
||||
)
|
||||
assert incomplete.status_code == 422
|
||||
assert incomplete.json()["error"]["code"] == "evidence_required"
|
||||
|
||||
resolved = httpx.post(
|
||||
f"{stack.base}/api/v1/upload-batches/{batch['id']}/resolve",
|
||||
json={
|
||||
"asset_id": asset_id,
|
||||
"outcome": "uploaded",
|
||||
"evidence": "found it in Immich by checksum",
|
||||
"actor": "dom",
|
||||
},
|
||||
timeout=TIMEOUT,
|
||||
)
|
||||
resolved.raise_for_status()
|
||||
assert _outcomes(_get(stack.base, batch["id"]))["a.jpg"] == "uploaded"
|
||||
manual = httpx.get(
|
||||
f"{stack.base}/api/v1/upload-batches/{batch['id']}/verifications", timeout=TIMEOUT
|
||||
).json()["verifications"][-1]
|
||||
assert manual["source"] == "operator"
|
||||
assert manual["actor"] == "dom"
|
||||
|
||||
|
||||
def test_bytes_edited_after_the_upload_are_flagged_and_block_another_run(stack):
|
||||
stack.start(uploader=ALL_NEW)
|
||||
batch = _upload(stack)
|
||||
assert batch["state"] == "succeeded"
|
||||
|
||||
(stack.seeded.lib / ALBUM / "a.jpg").write_bytes(b"edited after the upload")
|
||||
|
||||
result = _verify(stack.base, batch["id"])
|
||||
assert result["stale_bytes"] is True
|
||||
changed = [item for item in result["items"] if item["changed_after_upload"]]
|
||||
assert [Path(item["path"]).name for item in changed] == ["a.jpg"]
|
||||
refused = _start(stack.base, batch["id"])
|
||||
assert refused.status_code == 409
|
||||
assert refused.json()["error"]["code"] == "changed_after_upload"
|
||||
# The preflight agrees: those bytes are no longer approved for any new upload.
|
||||
report = _preflight(stack.base)
|
||||
assert report["state"] == "blocked"
|
||||
assert "bytes_changed" in {
|
||||
issue["code"]
|
||||
for album in report["albums"]
|
||||
for asset in album["assets"]
|
||||
for issue in asset["blockers"]
|
||||
}
|
||||
|
||||
|
||||
# ── cancellation and resume ──────────────────────────────────────────────────
|
||||
|
||||
|
||||
def test_a_running_album_can_be_stopped_and_run_again_from_that_boundary(stack, tmp_path):
|
||||
stack.start(uploader=_once_then(tmp_path, *SLOW_FIRST))
|
||||
batch = _approve(stack)
|
||||
_start(stack.base, batch["id"]).raise_for_status()
|
||||
_await_state(stack.base, batch["id"], {"running"})
|
||||
|
||||
httpx.post(
|
||||
f"{stack.base}/api/v1/upload-batches/{batch['id']}/cancel", timeout=TIMEOUT
|
||||
).raise_for_status()
|
||||
cancelled = _await_state(stack.base, batch["id"], {"cancelled"})
|
||||
|
||||
# A stopped album is a clean boundary, not an uncertain one.
|
||||
assert cancelled["retry_blockers"] == []
|
||||
_start_accepted(stack.base, batch["id"])
|
||||
resumed = _await_report(stack.base, batch["id"], {"succeeded"})
|
||||
assert resumed["attempt_count"] == 2
|
||||
assert _outcomes(resumed) == {"a.jpg": "uploaded", "b.jpg": "uploaded"}
|
||||
|
||||
|
||||
def test_an_interrupted_attempt_is_uncertain_after_a_restart_and_stays_blocked(stack, tmp_path):
|
||||
"""The worker vanishes mid-upload: Immich may hold the files, so the outcome is
|
||||
unknown. Startup recovery must say so, refuse a retry, and survive the restart."""
|
||||
stack.start(uploader=_once_then(tmp_path, *SLOW_FIRST))
|
||||
batch = _approve(stack)
|
||||
_start(stack.base, batch["id"]).raise_for_status()
|
||||
_await_state(stack.base, batch["id"], {"running"})
|
||||
|
||||
stack.worker.kill() # no chance to record any outcome
|
||||
stack.worker.wait(timeout=10)
|
||||
stack.restart_server()
|
||||
|
||||
recovered = _get(stack.base, batch["id"])
|
||||
assert recovered["state"] == "unknown_requires_verification"
|
||||
assert recovered["error_code"] == "interrupted"
|
||||
assert recovered["attempt_count"] == 1
|
||||
refused = _start(stack.base, batch["id"])
|
||||
assert refused.status_code == 409
|
||||
assert refused.json()["error"]["code"] == "requires_verification"
|
||||
|
||||
# Verification is the only way out, and it is what makes the batch certain again.
|
||||
stack.immich.mode("present")
|
||||
result = _verify(stack.base, batch["id"])
|
||||
assert result["state"] == "succeeded"
|
||||
assert result["outcome_state"] == "verified"
|
||||
|
||||
stack.restart_server() # the resolution is durable, not in-process memory
|
||||
after = _get(stack.base, batch["id"])
|
||||
assert after["state"] == "succeeded"
|
||||
assert set(_outcomes(after).values()) == {"uploaded"}
|
||||
|
||||
|
||||
# ── privacy ──────────────────────────────────────────────────────────────────
|
||||
|
||||
|
||||
def test_no_secret_appears_in_any_retained_artifact(stack):
|
||||
stack.start(uploader=ALL_NEW)
|
||||
planned = _approve(stack)
|
||||
job = _start(stack.base, planned["id"]).json()["job"]
|
||||
batch = _await_report(stack.base, planned["id"])
|
||||
stack.immich.mode("present")
|
||||
_verify(stack.base, batch["id"])
|
||||
|
||||
# The uploader really was given the key…
|
||||
assert f"--api-key={SENTINEL_KEY}" in stack.argv()[0]
|
||||
# …and it is in nothing that was kept: not the database, not the retained report,
|
||||
# not any response the browser can read.
|
||||
retained = [path for path in stack.seeded.data.rglob("*") if path.is_file()]
|
||||
assert any(path.suffix == ".log" for path in retained), "the report was not retained"
|
||||
for path in retained:
|
||||
assert SENTINEL_KEY.encode() not in path.read_bytes(), path
|
||||
for url in (
|
||||
f"{stack.base}/api/v1/upload-batches",
|
||||
f"{stack.base}/api/v1/upload-batches/{batch['id']}",
|
||||
f"{stack.base}/api/v1/upload-batches/{batch['id']}/verifications",
|
||||
f"{stack.base}/api/v1/jobs/{job['id']}",
|
||||
f"{stack.base}/api/v1/jobs/{job['id']}/events",
|
||||
):
|
||||
assert SENTINEL_KEY not in httpx.get(url, timeout=TIMEOUT).text, url
|
||||
assert SENTINEL_KEY not in httpx.post(
|
||||
f"{stack.base}/api/v1/upload-preflight", json={}, timeout=TIMEOUT
|
||||
).text
|
||||
assert "--api-key=***" in " ".join(_get(stack.base, batch["id"])["command"])
|
||||
@@ -13,7 +13,6 @@ REPO = Path(__file__).resolve().parents[2]
|
||||
MAP = json.loads((REPO / "tests" / "story_traceability.json").read_text())["stories"]
|
||||
PHASE_A_STORIES = {f"US01-0{n}" for n in range(1, 8)}
|
||||
PHASE_D_STORIES = {f"US04-0{n}" for n in range(1, 7)}
|
||||
PHASE_E_STORIES = {f"US05-0{n}" for n in range(1, 7)}
|
||||
|
||||
|
||||
def test_all_phase_a_stories_are_mapped():
|
||||
@@ -26,12 +25,6 @@ def test_all_phase_d_stories_are_mapped():
|
||||
assert PHASE_D_STORIES <= set(MAP)
|
||||
|
||||
|
||||
def test_all_phase_e_stories_are_mapped():
|
||||
"""US05-06 acceptance: every upload story, preflight through browser, is tied to
|
||||
automated tests — upload is the one stage the app cannot take back."""
|
||||
assert PHASE_E_STORIES <= set(MAP)
|
||||
|
||||
|
||||
def test_every_mapped_test_file_exists_and_is_nonempty():
|
||||
for story, files in MAP.items():
|
||||
assert files, f"{story} maps to no tests"
|
||||
|
||||
@@ -13,31 +13,180 @@ API key is a sentinel string, so the last test can prove it never reached the pa
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import json
|
||||
import stat
|
||||
import threading
|
||||
from http.server import BaseHTTPRequestHandler, HTTPServer
|
||||
from pathlib import Path
|
||||
|
||||
import httpx
|
||||
import pytest
|
||||
from playwright.sync_api import expect
|
||||
|
||||
from tests.e2e._pipeline_harness import (
|
||||
SENTINEL_KEY,
|
||||
SILENT_UPLOADER,
|
||||
UPLOADER_VERSION,
|
||||
UploadStack,
|
||||
mark_upload_ready,
|
||||
Server,
|
||||
seed_album,
|
||||
session_factory,
|
||||
start_worker,
|
||||
wait_until,
|
||||
)
|
||||
|
||||
pytestmark = pytest.mark.phase_e
|
||||
|
||||
TIMEOUT = 10
|
||||
SENTINEL_KEY = "immich-sentinel-9f3a2b"
|
||||
UPLOADER_VERSION = "immich-go 0.21.0" # a pinned family, so reports are parsable
|
||||
|
||||
# A report the pinned text-v1 grammar understands: one new file, one the server
|
||||
# already holds. ``$6`` is the folder argument of ``upload from-folder``.
|
||||
REPORTING_UPLOADER = (
|
||||
'echo "INFO uploaded $6/a.jpg"\n'
|
||||
'echo "INFO server has the same file $6/b.jpg"\n'
|
||||
'echo "Uploaded 1, duplicates 1"\n'
|
||||
"exit 0\n"
|
||||
)
|
||||
# Exits cleanly but says nothing about any file: the process succeeded, the
|
||||
# per-file outcome is unknown.
|
||||
SILENT_UPLOADER = "exit 0\n"
|
||||
|
||||
|
||||
# ── fake Immich ──────────────────────────────────────────────────────────────
|
||||
|
||||
|
||||
def _handler(state: dict):
|
||||
class Handler(BaseHTTPRequestHandler):
|
||||
def do_GET(self): # noqa: N802 (BaseHTTPRequestHandler API)
|
||||
self._json(200, {"res": "pong"})
|
||||
|
||||
def do_POST(self): # noqa: N802
|
||||
length = int(self.headers.get("Content-Length", 0))
|
||||
payload = json.loads(self.rfile.read(length) or b"{}")
|
||||
if state["mode"] == "broken":
|
||||
self.send_error(500, "bulk-upload-check is unavailable")
|
||||
return
|
||||
reject = state["mode"] == "present"
|
||||
self._json(
|
||||
200,
|
||||
{
|
||||
"results": [
|
||||
{
|
||||
"id": asset["id"],
|
||||
"action": "reject" if reject else "accept",
|
||||
"reason": "duplicate" if reject else None,
|
||||
}
|
||||
for asset in payload.get("assets", [])
|
||||
]
|
||||
},
|
||||
)
|
||||
|
||||
def _json(self, code: int, body: dict) -> None:
|
||||
raw = json.dumps(body).encode()
|
||||
self.send_response(code)
|
||||
self.send_header("Content-Type", "application/json")
|
||||
self.send_header("Content-Length", str(len(raw)))
|
||||
self.end_headers()
|
||||
self.wfile.write(raw)
|
||||
|
||||
def log_message(self, *args):
|
||||
pass
|
||||
|
||||
return Handler
|
||||
|
||||
|
||||
class FakeImmich:
|
||||
"""An Immich that answers ping, and says whether it holds the exact bytes.
|
||||
|
||||
``mode`` is what the next verification will find: ``present`` (the server
|
||||
deduplicates them, so it has them), ``absent`` (it would accept them, so it does
|
||||
not), or ``broken`` (no usable answer at all).
|
||||
"""
|
||||
|
||||
def __init__(self) -> None:
|
||||
self.state = {"mode": "present"}
|
||||
self._server = HTTPServer(("127.0.0.1", 0), _handler(self.state))
|
||||
threading.Thread(target=self._server.serve_forever, daemon=True).start()
|
||||
self.url = f"http://127.0.0.1:{self._server.server_port}"
|
||||
|
||||
def mode(self, mode: str) -> None:
|
||||
self.state["mode"] = mode
|
||||
|
||||
def stop(self) -> None:
|
||||
self._server.shutdown()
|
||||
self._server.server_close()
|
||||
|
||||
|
||||
# ── stack ────────────────────────────────────────────────────────────────────
|
||||
|
||||
|
||||
def _uploader(tmp_path: Path, body: str) -> Path:
|
||||
path = tmp_path / "immich-go"
|
||||
path.write_text(
|
||||
f'#!/bin/sh\nif [ "$1" = "--version" ]; then echo "{UPLOADER_VERSION}"; exit 0; fi\n{body}'
|
||||
)
|
||||
path.chmod(path.stat().st_mode | stat.S_IEXEC | stat.S_IXGRP | stat.S_IXOTH)
|
||||
return path
|
||||
|
||||
|
||||
def _mark_upload_ready(seeded, *, unverified: tuple[str, ...] = ()) -> None:
|
||||
"""Give every seeded photo the verified EXIF checkpoints upload requires.
|
||||
|
||||
``unverified`` names stems whose analysis checkpoint stays incomplete, which is
|
||||
what makes an album partially blocked.
|
||||
"""
|
||||
from datetime import datetime, timezone
|
||||
|
||||
from sqlalchemy import select
|
||||
|
||||
from photo_pipeline.models import AnalysisResult, SafetyReview
|
||||
|
||||
now = datetime(2026, 1, 1, tzinfo=timezone.utc)
|
||||
blocked = {seeded.asset_ids[stem] for stem in unverified}
|
||||
with session_factory(seeded) as sf:
|
||||
with sf() as session:
|
||||
for review in session.scalars(select(SafetyReview)):
|
||||
review.exif_verified_at = now
|
||||
for analysis in session.scalars(select(AnalysisResult)):
|
||||
analysis.exif_written_at = None if analysis.asset_id in blocked else now
|
||||
session.commit()
|
||||
|
||||
|
||||
class Stack:
|
||||
"""A seeded, upload-ready library plus the server, worker, and fake Immich."""
|
||||
|
||||
def __init__(self, tmp_path: Path, seeded, immich: FakeImmich) -> None:
|
||||
self.tmp_path = tmp_path
|
||||
self.seeded = seeded
|
||||
self.immich = immich
|
||||
self.server: Server | None = None
|
||||
self.worker = None
|
||||
|
||||
def start(self, *, uploader: str = REPORTING_UPLOADER, worker: bool = True) -> "Stack":
|
||||
env = {
|
||||
"PHOTO_PIPELINE_IMMICH_SERVER_URL": self.immich.url,
|
||||
"PHOTO_PIPELINE_IMMICH_API_KEY": SENTINEL_KEY,
|
||||
"PHOTO_PIPELINE_IMMICH_GO_BINARY": str(_uploader(self.tmp_path, uploader)),
|
||||
}
|
||||
self.server = Server(self.seeded, extra_env=env).start()
|
||||
self.base = self.server.base
|
||||
if worker:
|
||||
self.worker = start_worker(self.seeded, extra_env=env)
|
||||
return self
|
||||
|
||||
def batches(self) -> list[dict]:
|
||||
return httpx.get(f"{self.base}/api/v1/upload-batches", timeout=TIMEOUT).json()["batches"]
|
||||
|
||||
def stop(self) -> None:
|
||||
if self.worker is not None:
|
||||
self.worker.terminate()
|
||||
self.worker.wait(timeout=10)
|
||||
if self.server is not None:
|
||||
self.server.stop()
|
||||
self.immich.stop()
|
||||
|
||||
|
||||
@pytest.fixture
|
||||
def stack(tmp_path):
|
||||
seeded = seed_album(tmp_path)
|
||||
mark_upload_ready(seeded)
|
||||
running = UploadStack(tmp_path, seeded)
|
||||
_mark_upload_ready(seeded)
|
||||
running = Stack(tmp_path, seeded, FakeImmich())
|
||||
try:
|
||||
yield running
|
||||
finally:
|
||||
@@ -87,7 +236,7 @@ def test_the_preview_shows_scope_configuration_and_an_exact_confirmation(page, s
|
||||
|
||||
|
||||
def test_an_unfinished_photo_blocks_its_album_and_the_confirmation(page, stack):
|
||||
mark_upload_ready(stack.seeded, unverified=("b",))
|
||||
_mark_upload_ready(stack.seeded, unverified=("b",))
|
||||
stack.start(worker=False)
|
||||
_open(page, stack)
|
||||
|
||||
@@ -273,7 +422,8 @@ def test_an_interrupted_attempt_is_shown_as_uncertain_after_a_restart(page, stac
|
||||
session.get(UploadBatch, created["id"]).state = "running"
|
||||
session.commit()
|
||||
|
||||
stack.restart_server() # the same port, so recovery runs in a genuinely fresh process
|
||||
stack.server.stop()
|
||||
stack.server.start() # the same port, so recovery runs in a genuinely fresh process
|
||||
page.reload()
|
||||
|
||||
expect(page.get_by_test_id("detail-state")).to_have_text("unknown_requires_verification")
|
||||
|
||||
@@ -1,592 +0,0 @@
|
||||
"""Archive destinations and preflight (US06-01).
|
||||
|
||||
Archive is the only stage that removes originals, so every case here asks the same
|
||||
question: would this preflight let an album leave active storage when it should
|
||||
not? The destinations are real directories on real filesystems — mounted, missing,
|
||||
swapped for another medium, read-only, or full — and preflight itself must stay
|
||||
non-destructive: the library snapshot is asserted unchanged.
|
||||
"""
|
||||
|
||||
import json
|
||||
import os
|
||||
import stat
|
||||
import uuid
|
||||
from datetime import datetime, timedelta, timezone
|
||||
|
||||
import pytest
|
||||
from fastapi.testclient import TestClient
|
||||
|
||||
from photo_pipeline.api.app import create_app
|
||||
from photo_pipeline.config import Config
|
||||
from photo_pipeline.db import create_db_engine, create_session_factory, run_migrations
|
||||
from photo_pipeline.jobs.domain_handlers import ARCHIVE_LOCK, UPLOAD_LOCK
|
||||
from photo_pipeline.models import (
|
||||
Asset,
|
||||
RenameOperation,
|
||||
RenamePlan,
|
||||
UploadBatch,
|
||||
UploadItem,
|
||||
)
|
||||
from photo_pipeline.services.archives import MARKER_NAME, ArchiveError, ArchiveService
|
||||
from photo_pipeline.services.hashing import sha256_file
|
||||
from photo_pipeline.services.jobs import JobService
|
||||
|
||||
pytestmark = pytest.mark.phase_f # part of the Phase F acceptance gate (US06-06)
|
||||
|
||||
NOW = datetime(2026, 1, 1, tzinfo=timezone.utc)
|
||||
|
||||
|
||||
# ── environment ──────────────────────────────────────────────────────────────
|
||||
|
||||
|
||||
def _env(tmp_path, *, reserve=0):
|
||||
(tmp_path / "data").mkdir(exist_ok=True)
|
||||
lib = tmp_path / "lib"
|
||||
lib.mkdir(exist_ok=True)
|
||||
archive = tmp_path / "archive"
|
||||
archive.mkdir(exist_ok=True)
|
||||
config = Config.from_env(
|
||||
{
|
||||
"PHOTO_PIPELINE_DATA_DIR": str(tmp_path / "data"),
|
||||
"PHOTO_PIPELINE_LIBRARY_ROOTS": str(lib),
|
||||
"PHOTO_PIPELINE_ARCHIVE_FREE_SPACE_RESERVE_BYTES": str(reserve),
|
||||
}
|
||||
)
|
||||
run_migrations(config.database_url)
|
||||
return config, create_session_factory(create_db_engine(config.database_url)), lib, archive
|
||||
|
||||
|
||||
def _album(
|
||||
sf,
|
||||
lib,
|
||||
album="rome",
|
||||
names=("a.jpg", "b.jpg"),
|
||||
*,
|
||||
uploaded=True,
|
||||
outcome="uploaded",
|
||||
outcome_state="verified",
|
||||
stale_bytes=False,
|
||||
):
|
||||
"""A real album folder whose assets carry their upload evidence."""
|
||||
folder = lib / album
|
||||
folder.mkdir(parents=True, exist_ok=True)
|
||||
ids = []
|
||||
with sf() as session:
|
||||
batch_id = str(uuid.uuid4())
|
||||
if uploaded:
|
||||
session.add(
|
||||
UploadBatch(
|
||||
id=batch_id,
|
||||
album=album,
|
||||
folder=str(folder),
|
||||
album_name=album,
|
||||
state="succeeded",
|
||||
preflight_token="v1:test",
|
||||
outcome_state=outcome_state,
|
||||
stale_bytes=stale_bytes,
|
||||
created_at=NOW,
|
||||
)
|
||||
)
|
||||
for name in names:
|
||||
path = folder / name
|
||||
path.write_bytes(name.encode() * 16)
|
||||
asset_id = str(uuid.uuid4())
|
||||
ids.append(asset_id)
|
||||
session.add(
|
||||
Asset(
|
||||
id=asset_id,
|
||||
original_path=str(path),
|
||||
current_path=str(path),
|
||||
discovered_at=NOW,
|
||||
hash_version=1,
|
||||
byte_size=path.stat().st_size,
|
||||
current_sha256=sha256_file(path),
|
||||
)
|
||||
)
|
||||
if uploaded:
|
||||
session.add(
|
||||
UploadItem(
|
||||
batch_id=batch_id,
|
||||
asset_id=asset_id,
|
||||
path=str(path),
|
||||
sha256=sha256_file(path),
|
||||
sha1="0" * 40,
|
||||
state="sent",
|
||||
outcome=outcome,
|
||||
)
|
||||
)
|
||||
session.commit()
|
||||
return folder, ids
|
||||
|
||||
|
||||
def _service(sf, config):
|
||||
return ArchiveService(sf, config=config)
|
||||
|
||||
|
||||
def _location(sf, config, archive, name="external"):
|
||||
return _service(sf, config).register(name, str(archive))
|
||||
|
||||
|
||||
def _snapshot(lib):
|
||||
return {
|
||||
str(p.relative_to(lib)): (p.read_bytes() if p.is_file() else None)
|
||||
for p in sorted(lib.rglob("*"))
|
||||
}
|
||||
|
||||
|
||||
def _codes(report):
|
||||
return (
|
||||
{issue["code"] for issue in report["blockers"]}
|
||||
| {issue["code"] for album in report["albums"] for issue in album["blockers"]}
|
||||
| {
|
||||
issue["code"]
|
||||
for album in report["albums"]
|
||||
for asset in album["assets"]
|
||||
for issue in asset["blockers"]
|
||||
}
|
||||
)
|
||||
|
||||
|
||||
# ── locations ────────────────────────────────────────────────────────────────
|
||||
|
||||
|
||||
def test_registering_a_location_stamps_the_medium_with_its_identity(tmp_path):
|
||||
config, sf, _, archive = _env(tmp_path)
|
||||
|
||||
location = _location(sf, config, archive)
|
||||
|
||||
marker = json.loads((archive / MARKER_NAME).read_text())
|
||||
assert marker["media_id"] == location["media_id"]
|
||||
assert location["state"] == "online" and location["writable"] is True
|
||||
assert location["root"] == str(archive.resolve())
|
||||
listed = _service(sf, config).locations()
|
||||
assert [(row["id"], row["media_id"], row["state"]) for row in listed] == [
|
||||
(location["id"], location["media_id"], "online")
|
||||
]
|
||||
|
||||
|
||||
def test_a_second_location_cannot_claim_the_same_medium(tmp_path):
|
||||
config, sf, _, archive = _env(tmp_path)
|
||||
_location(sf, config, archive)
|
||||
|
||||
with pytest.raises(ArchiveError) as error:
|
||||
_location(sf, config, archive, name="second")
|
||||
|
||||
assert error.value.code == "already_registered"
|
||||
|
||||
|
||||
@pytest.mark.parametrize("inside", ["", "sub"])
|
||||
def test_a_destination_inside_the_library_is_refused(tmp_path, inside):
|
||||
"""The library may never archive into itself: the 'reclaimed' bytes would still
|
||||
be in the active tree, and a later scan would rediscover them."""
|
||||
config, sf, lib, _ = _env(tmp_path)
|
||||
root = lib / inside if inside else lib
|
||||
root.mkdir(exist_ok=True)
|
||||
|
||||
with pytest.raises(ArchiveError) as error:
|
||||
_service(sf, config).register("bad", str(root))
|
||||
|
||||
assert error.value.code == "unsafe_destination"
|
||||
|
||||
|
||||
def test_an_ignored_destination_is_refused(tmp_path):
|
||||
config, sf, _, _ = _env(tmp_path)
|
||||
root = tmp_path / "_IGNORE" / "archive"
|
||||
root.mkdir(parents=True)
|
||||
|
||||
with pytest.raises(ArchiveError) as error:
|
||||
_service(sf, config).register("ignored", str(root))
|
||||
|
||||
assert error.value.code == "unsafe_destination"
|
||||
|
||||
|
||||
def test_listing_reports_an_unmounted_medium_as_offline(tmp_path):
|
||||
config, sf, _, archive = _env(tmp_path)
|
||||
_location(sf, config, archive)
|
||||
(archive / MARKER_NAME).unlink()
|
||||
|
||||
assert [row["state"] for row in _service(sf, config).locations()] == ["offline"]
|
||||
|
||||
|
||||
# ── happy path ───────────────────────────────────────────────────────────────
|
||||
|
||||
|
||||
def test_ready_preflight_previews_scope_method_and_reclaimable_bytes(tmp_path):
|
||||
config, sf, lib, archive = _env(tmp_path)
|
||||
folder, ids = _album(sf, lib)
|
||||
location = _location(sf, config, archive)
|
||||
before = _snapshot(lib)
|
||||
|
||||
report = _service(sf, config).preflight(location["id"])
|
||||
|
||||
assert report["state"] == "ready" and report["blockers"] == []
|
||||
album = report["albums"][0]
|
||||
assert album["album"] == "rome" and album["folder"] == str(folder)
|
||||
assert album["destination"] == str(archive.resolve() / "rome")
|
||||
assert album["transfer_method"] in ("move", "copy_verify_remove")
|
||||
assert album["reclaimable_bytes"] == sum(p.stat().st_size for p in folder.iterdir())
|
||||
assert sorted(a["asset_id"] for a in album["assets"]) == sorted(ids)
|
||||
assert report["totals"]["bytes"] == album["reclaimable_bytes"]
|
||||
assert report["capacity"]["sufficient"] is True
|
||||
# Both must be proven by writing, not assumed.
|
||||
assert report["backup"]["ok"] is True and report["backup"]["bytes"] > 0
|
||||
assert report["manifest"]["ok"] is True
|
||||
assert report["token"].startswith("v1:")
|
||||
assert _snapshot(lib) == before, "preflight must not touch the library"
|
||||
assert not list(archive.glob("*probe*")), "probe files must be cleaned up"
|
||||
|
||||
|
||||
def test_same_filesystem_destination_is_previewed_as_a_move(tmp_path):
|
||||
config, sf, lib, archive = _env(tmp_path)
|
||||
_album(sf, lib)
|
||||
location = _location(sf, config, archive)
|
||||
|
||||
report = _service(sf, config).preflight(location["id"])
|
||||
|
||||
# tmp_path is one filesystem, so this is the same-filesystem case by construction.
|
||||
assert report["albums"][0]["transfer_method"] == "move"
|
||||
|
||||
|
||||
def test_scoping_to_one_album_excludes_the_others(tmp_path):
|
||||
config, sf, lib, archive = _env(tmp_path)
|
||||
_album(sf, lib, "rome")
|
||||
_album(sf, lib, "paris", names=("c.jpg",))
|
||||
location = _location(sf, config, archive)
|
||||
|
||||
report = _service(sf, config).preflight(location["id"], ["paris"])
|
||||
|
||||
assert [album["album"] for album in report["albums"]] == ["paris"]
|
||||
assert report["totals"]["assets"] == 1
|
||||
|
||||
|
||||
def test_unknown_album_and_unknown_location_are_refused(tmp_path):
|
||||
config, sf, lib, archive = _env(tmp_path)
|
||||
_album(sf, lib)
|
||||
location = _location(sf, config, archive)
|
||||
service = _service(sf, config)
|
||||
|
||||
with pytest.raises(ArchiveError) as unknown_album:
|
||||
service.preflight(location["id"], ["atlantis"])
|
||||
with pytest.raises(ArchiveError) as unknown_location:
|
||||
service.preflight("nope")
|
||||
|
||||
assert unknown_album.value.code == "unknown_album"
|
||||
assert unknown_location.value.code == "unknown_location"
|
||||
|
||||
|
||||
def test_empty_scope_is_a_blocker(tmp_path):
|
||||
config, sf, _, archive = _env(tmp_path)
|
||||
location = _location(sf, config, archive)
|
||||
|
||||
report = _service(sf, config).preflight(location["id"])
|
||||
|
||||
assert report["state"] == "blocked" and "empty_scope" in _codes(report)
|
||||
|
||||
|
||||
# ── destination ──────────────────────────────────────────────────────────────
|
||||
|
||||
|
||||
def test_offline_medium_blocks(tmp_path):
|
||||
config, sf, lib, archive = _env(tmp_path)
|
||||
_album(sf, lib)
|
||||
location = _location(sf, config, archive)
|
||||
(archive / MARKER_NAME).unlink() # the disk went away
|
||||
|
||||
report = _service(sf, config).preflight(location["id"])
|
||||
|
||||
assert report["state"] == "blocked" and "location_offline" in _codes(report)
|
||||
assert report["location"]["state"] == "offline"
|
||||
|
||||
|
||||
def test_a_different_medium_at_the_same_mountpoint_blocks(tmp_path):
|
||||
"""The mountpoint is right, the disk is not — never write the archive here."""
|
||||
config, sf, lib, archive = _env(tmp_path)
|
||||
_album(sf, lib)
|
||||
location = _location(sf, config, archive)
|
||||
(archive / MARKER_NAME).write_text(json.dumps({"media_id": "some-other-disk"}))
|
||||
|
||||
report = _service(sf, config).preflight(location["id"])
|
||||
|
||||
assert "wrong_volume" in _codes(report)
|
||||
assert report["location"]["state"] == "wrong_volume"
|
||||
|
||||
|
||||
def test_read_only_destination_blocks_and_cannot_write_the_manifest(tmp_path):
|
||||
config, sf, lib, archive = _env(tmp_path)
|
||||
_album(sf, lib)
|
||||
location = _location(sf, config, archive)
|
||||
mode = archive.stat().st_mode
|
||||
archive.chmod(mode & ~stat.S_IWUSR & ~stat.S_IWGRP & ~stat.S_IWOTH)
|
||||
try:
|
||||
report = _service(sf, config).preflight(location["id"])
|
||||
finally:
|
||||
archive.chmod(mode)
|
||||
|
||||
assert {"destination_not_writable", "manifest_unwritable"} <= _codes(report)
|
||||
assert report["manifest"]["ok"] is False
|
||||
|
||||
|
||||
@pytest.mark.skipif(os.geteuid() == 0, reason="root ignores directory permissions")
|
||||
def test_read_only_destination_is_detected_by_a_real_write(tmp_path):
|
||||
config, sf, lib, archive = _env(tmp_path)
|
||||
_album(sf, lib)
|
||||
location = _location(sf, config, archive)
|
||||
mode = archive.stat().st_mode
|
||||
archive.chmod(stat.S_IRUSR | stat.S_IXUSR)
|
||||
try:
|
||||
report = _service(sf, config).preflight(location["id"])
|
||||
finally:
|
||||
archive.chmod(mode)
|
||||
|
||||
assert report["location"]["writable"] is False
|
||||
|
||||
|
||||
def test_insufficient_capacity_blocks(tmp_path):
|
||||
"""The reserve is what stops an archive from filling its own destination."""
|
||||
config, sf, lib, archive = _env(tmp_path, reserve=10**15)
|
||||
_album(sf, lib)
|
||||
location = _location(sf, config, archive)
|
||||
|
||||
report = _service(sf, config).preflight(location["id"])
|
||||
|
||||
assert report["state"] == "blocked" and "insufficient_capacity" in _codes(report)
|
||||
assert report["capacity"]["sufficient"] is False
|
||||
assert report["capacity"]["reserve_bytes"] == 10**15
|
||||
|
||||
|
||||
def test_an_occupied_destination_blocks_that_album(tmp_path):
|
||||
config, sf, lib, archive = _env(tmp_path)
|
||||
_album(sf, lib)
|
||||
location = _location(sf, config, archive)
|
||||
(archive / "rome").mkdir()
|
||||
(archive / "rome" / "a.jpg").write_bytes(b"something already here")
|
||||
|
||||
report = _service(sf, config).preflight(location["id"])
|
||||
|
||||
assert "destination_collision" in _codes(report)
|
||||
assert report["albums"][0]["state"] == "blocked"
|
||||
|
||||
|
||||
def test_a_destination_moved_into_the_library_blocks_even_though_it_registered(tmp_path):
|
||||
"""Registration validated the root once; preflight validates it again, because a
|
||||
mountpoint can be moved after the fact."""
|
||||
config, sf, lib, archive = _env(tmp_path)
|
||||
_album(sf, lib)
|
||||
location = _location(sf, config, archive)
|
||||
with sf() as session:
|
||||
from photo_pipeline.models import ArchiveLocation
|
||||
|
||||
session.get(ArchiveLocation, location["id"]).root = str(lib / "inside")
|
||||
session.commit()
|
||||
(lib / "inside").mkdir()
|
||||
(lib / "inside" / MARKER_NAME).write_text(json.dumps({"media_id": location["media_id"]}))
|
||||
|
||||
report = _service(sf, config).preflight(location["id"])
|
||||
|
||||
assert "unsafe_destination" in _codes(report)
|
||||
|
||||
|
||||
# ── source readiness ─────────────────────────────────────────────────────────
|
||||
|
||||
|
||||
@pytest.mark.parametrize(
|
||||
"kwargs,code",
|
||||
[
|
||||
({"uploaded": False}, "upload_unverified"),
|
||||
({"outcome_state": "requires_verification"}, "upload_unverified"),
|
||||
({"outcome": "failed"}, "upload_unverified"),
|
||||
({"outcome": "skipped"}, "upload_unverified"),
|
||||
({"stale_bytes": True}, "upload_unverified"),
|
||||
],
|
||||
)
|
||||
def test_an_unverified_upload_blocks_the_album(tmp_path, kwargs, code):
|
||||
config, sf, lib, archive = _env(tmp_path)
|
||||
_album(sf, lib, **kwargs)
|
||||
location = _location(sf, config, archive)
|
||||
|
||||
report = _service(sf, config).preflight(location["id"])
|
||||
|
||||
assert report["state"] == "blocked" and code in _codes(report)
|
||||
assert report["albums"][0]["blocked_count"] == 2
|
||||
|
||||
|
||||
@pytest.mark.parametrize("outcome", ["upgraded", "duplicate"])
|
||||
def test_upgraded_and_duplicate_uploads_are_evidence_enough(tmp_path, outcome):
|
||||
"""Immich already holds these exact bytes; that is what archiving requires."""
|
||||
config, sf, lib, archive = _env(tmp_path)
|
||||
_album(sf, lib, outcome=outcome)
|
||||
location = _location(sf, config, archive)
|
||||
|
||||
assert _service(sf, config).preflight(location["id"])["state"] == "ready"
|
||||
|
||||
|
||||
def test_bytes_changed_since_upload_block_the_album(tmp_path):
|
||||
config, sf, lib, archive = _env(tmp_path)
|
||||
folder, _ = _album(sf, lib)
|
||||
location = _location(sf, config, archive)
|
||||
(folder / "a.jpg").write_bytes(b"edited after the upload")
|
||||
|
||||
report = _service(sf, config).preflight(location["id"])
|
||||
|
||||
assert {"bytes_changed", "partial_scope"} <= _codes(report)
|
||||
assert report["albums"][0]["blocked_count"] == 1
|
||||
|
||||
|
||||
def test_a_missing_source_file_blocks_the_album(tmp_path):
|
||||
config, sf, lib, archive = _env(tmp_path)
|
||||
folder, _ = _album(sf, lib)
|
||||
location = _location(sf, config, archive)
|
||||
(folder / "a.jpg").unlink()
|
||||
|
||||
assert "file_missing" in _codes(_service(sf, config).preflight(location["id"]))
|
||||
|
||||
|
||||
# ── leases ───────────────────────────────────────────────────────────────────
|
||||
|
||||
|
||||
@pytest.mark.parametrize("lock", [UPLOAD_LOCK, ARCHIVE_LOCK])
|
||||
def test_a_held_lease_blocks_archiving(tmp_path, lock):
|
||||
config, sf, lib, archive = _env(tmp_path)
|
||||
_album(sf, lib)
|
||||
location = _location(sf, config, archive)
|
||||
JobService(sf).enqueue("upload_batch", lock=lock, items=["x"])
|
||||
|
||||
report = _service(sf, config).preflight(location["id"])
|
||||
|
||||
assert report["state"] == "blocked" and "lock_conflict" in _codes(report)
|
||||
|
||||
|
||||
def test_a_half_applied_rename_blocks_archiving(tmp_path):
|
||||
config, sf, lib, archive = _env(tmp_path)
|
||||
folder, _ = _album(sf, lib)
|
||||
location = _location(sf, config, archive)
|
||||
with sf() as session:
|
||||
plan_id = str(uuid.uuid4())
|
||||
session.add(RenamePlan(id=plan_id, state="applying", operation_count=1))
|
||||
session.flush()
|
||||
session.add(
|
||||
RenameOperation(
|
||||
id=str(uuid.uuid4()),
|
||||
plan_id=plan_id,
|
||||
sequence=0,
|
||||
operation="move_folder",
|
||||
source_path=str(folder),
|
||||
destination_path=str(lib / "2019 Rome"),
|
||||
journal_state="moving",
|
||||
)
|
||||
)
|
||||
session.commit()
|
||||
|
||||
report = _service(sf, config).preflight(location["id"])
|
||||
|
||||
assert report["state"] == "blocked" and "rename_pending" in _codes(report)
|
||||
|
||||
|
||||
# ── token ────────────────────────────────────────────────────────────────────
|
||||
|
||||
|
||||
def test_token_is_stable_while_nothing_relevant_changes(tmp_path):
|
||||
config, sf, lib, archive = _env(tmp_path)
|
||||
_album(sf, lib)
|
||||
location = _location(sf, config, archive)
|
||||
service = _service(sf, config)
|
||||
|
||||
first = service.preflight(location["id"])["token"]
|
||||
|
||||
assert service.preflight(location["id"])["token"] == first
|
||||
assert service.verify_token(first, location["id"]) is True
|
||||
|
||||
|
||||
def test_an_edited_source_makes_the_token_stale(tmp_path):
|
||||
config, sf, lib, archive = _env(tmp_path)
|
||||
folder, _ = _album(sf, lib)
|
||||
location = _location(sf, config, archive)
|
||||
service = _service(sf, config)
|
||||
token = service.preflight(location["id"])["token"]
|
||||
|
||||
(folder / "b.jpg").write_bytes(b"edited outside the app")
|
||||
|
||||
assert service.verify_token(token, location["id"]) is False
|
||||
|
||||
|
||||
def test_a_changed_destination_makes_the_token_stale(tmp_path):
|
||||
config, sf, lib, archive = _env(tmp_path)
|
||||
_album(sf, lib)
|
||||
location = _location(sf, config, archive)
|
||||
service = _service(sf, config)
|
||||
token = service.preflight(location["id"])["token"]
|
||||
|
||||
(archive / "rome").mkdir()
|
||||
(archive / "rome" / "a.jpg").write_bytes(b"appeared after approval")
|
||||
|
||||
assert service.verify_token(token, location["id"]) is False
|
||||
|
||||
|
||||
def test_a_token_from_another_scope_or_medium_is_rejected(tmp_path):
|
||||
config, sf, lib, archive = _env(tmp_path)
|
||||
_album(sf, lib, "rome")
|
||||
_album(sf, lib, "paris", names=("c.jpg",))
|
||||
other = tmp_path / "archive2"
|
||||
other.mkdir()
|
||||
location = _location(sf, config, archive)
|
||||
second = _location(sf, config, other, name="second")
|
||||
service = _service(sf, config)
|
||||
|
||||
rome = service.preflight(location["id"], ["rome"])["token"]
|
||||
|
||||
assert service.verify_token(rome, location["id"], ["paris"]) is False
|
||||
assert service.verify_token(rome, second["id"], ["rome"]) is False
|
||||
assert service.verify_token("v1:not-a-real-token", location["id"]) is False
|
||||
assert service.verify_token("", location["id"]) is False
|
||||
|
||||
|
||||
def test_the_token_survives_free_space_and_timestamp_drift(tmp_path):
|
||||
"""Free space changes constantly on a live disk; a token that expired on every
|
||||
byte written elsewhere would train users to ignore it."""
|
||||
config, sf, lib, archive = _env(tmp_path)
|
||||
_album(sf, lib)
|
||||
location = _location(sf, config, archive)
|
||||
service = _service(sf, config)
|
||||
token = service.preflight(location["id"])["token"]
|
||||
|
||||
(tmp_path / "unrelated.bin").write_bytes(b"0" * 100_000)
|
||||
with sf() as session:
|
||||
from photo_pipeline.models import ArchiveLocation
|
||||
|
||||
session.get(ArchiveLocation, location["id"]).last_seen_at = NOW - timedelta(days=5)
|
||||
session.commit()
|
||||
|
||||
assert service.verify_token(token, location["id"]) is True
|
||||
|
||||
|
||||
# ── API surface ──────────────────────────────────────────────────────────────
|
||||
|
||||
|
||||
def test_api_registers_a_location_and_returns_a_preflight_report(tmp_path):
|
||||
config, sf, lib, archive = _env(tmp_path)
|
||||
_album(sf, lib)
|
||||
|
||||
with TestClient(create_app(config)) as client:
|
||||
created = client.post(
|
||||
"/api/v1/archive-locations", json={"name": "external", "root": str(archive)}
|
||||
)
|
||||
listed = client.get("/api/v1/archive-locations")
|
||||
report = client.post(
|
||||
"/api/v1/archive-preflight", json={"location_id": created.json()["id"]}
|
||||
)
|
||||
|
||||
assert created.status_code == 201
|
||||
assert [row["name"] for row in listed.json()["locations"]] == ["external"]
|
||||
assert report.status_code == 200
|
||||
assert report.json()["state"] == "ready" and report.json()["token"].startswith("v1:")
|
||||
|
||||
|
||||
def test_api_rejects_an_unknown_location_and_an_unsafe_root(tmp_path):
|
||||
config, sf, lib, _ = _env(tmp_path)
|
||||
|
||||
with TestClient(create_app(config)) as client:
|
||||
unknown = client.post("/api/v1/archive-preflight", json={"location_id": "nope"})
|
||||
unsafe = client.post("/api/v1/archive-locations", json={"name": "bad", "root": str(lib)})
|
||||
|
||||
assert unknown.status_code == 404 and unknown.json()["error"]["code"] == "unknown_location"
|
||||
assert unsafe.status_code == 422 and unsafe.json()["error"]["code"] == "unsafe_destination"
|
||||
@@ -38,8 +38,6 @@ from photo_pipeline.services.upload_batches import (
|
||||
)
|
||||
from photo_pipeline.services.uploads import UploadService
|
||||
|
||||
pytestmark = pytest.mark.phase_e # part of the Phase E acceptance gate (US05-06)
|
||||
|
||||
NOW = datetime(2026, 1, 1, tzinfo=timezone.utc)
|
||||
SENTINEL_KEY = "immich-sentinel-9f3a2b"
|
||||
UPLOADER_VERSION = "immich-go 0.21.0"
|
||||
|
||||
@@ -29,8 +29,6 @@ from photo_pipeline.models import (
|
||||
from photo_pipeline.services.hashing import sha256_file
|
||||
from photo_pipeline.services.uploads import UploadError, UploadService
|
||||
|
||||
pytestmark = pytest.mark.phase_e # part of the Phase E acceptance gate (US05-06)
|
||||
|
||||
NOW = datetime(2026, 1, 1, tzinfo=timezone.utc)
|
||||
# A sentinel credential: every assertion below proves it never leaves configuration.
|
||||
SENTINEL_KEY = "immich-sentinel-9f3a2b"
|
||||
|
||||
@@ -32,8 +32,6 @@ from photo_pipeline.services.upload_reports import (
|
||||
)
|
||||
from photo_pipeline.services.uploads import UploadService
|
||||
|
||||
pytestmark = pytest.mark.phase_e # part of the Phase E acceptance gate (US05-06)
|
||||
|
||||
NOW = datetime(2026, 1, 1, tzinfo=timezone.utc)
|
||||
SUPPORTED_VERSION = "immich-go 0.21.0"
|
||||
|
||||
|
||||
@@ -40,8 +40,6 @@ from photo_pipeline.services.upload_verification import (
|
||||
)
|
||||
from photo_pipeline.services.uploads import UploadService
|
||||
|
||||
pytestmark = pytest.mark.phase_e # part of the Phase E acceptance gate (US05-06)
|
||||
|
||||
NOW = datetime(2026, 1, 1, tzinfo=timezone.utc)
|
||||
SENTINEL_KEY = "immich-sentinel-9f3a2b"
|
||||
UPLOADER_VERSION = "immich-go 0.21.0"
|
||||
|
||||
@@ -90,26 +90,6 @@ def test_cooperative_cancellation_leaves_items_resumable(sf, jobs):
|
||||
assert by_state.get(ItemState.QUEUED) == 1 # "b" left resumable
|
||||
|
||||
|
||||
def test_a_handler_that_stops_itself_releases_the_lock(sf, jobs):
|
||||
"""A handler may stop without anyone cancelling the *job* — an upload batch
|
||||
cancelled through its own API does exactly that. The job is still ``running``
|
||||
when it raises, so it has to reach ``cancelled`` through ``cancelling``; if that
|
||||
hop is skipped the transition is rejected and the lock is held forever."""
|
||||
from photo_pipeline.jobs.handlers import Cancelled
|
||||
|
||||
def handler(item, ctx):
|
||||
raise Cancelled("the work this job wraps was stopped elsewhere")
|
||||
|
||||
worker = Worker(sf, {"scan": handler}, "w1")
|
||||
job = jobs.enqueue("scan", lock="library_write", items=["a"])
|
||||
worker.run_once()
|
||||
|
||||
assert jobs.get(job["id"])["state"] == JobState.CANCELLED
|
||||
assert jobs.progress(job["id"])["by_state"] == {ItemState.QUEUED: 1} # resumable
|
||||
# The lane is free: the next job may be enqueued under the same lock.
|
||||
assert jobs.enqueue("scan", lock="library_write", items=["b"])["state"] == JobState.QUEUED
|
||||
|
||||
|
||||
def test_fencing_rejects_superseded_worker(sf, jobs):
|
||||
job = jobs.enqueue("scan", items=["a"])
|
||||
stale = jobs.claim(["scan"], "old")
|
||||
|
||||
@@ -120,12 +120,6 @@
|
||||
],
|
||||
"US05-05": [
|
||||
"tests/e2e/test_uploads_ui.py"
|
||||
],
|
||||
"US05-06": [
|
||||
"tests/e2e/test_phase_e_pipeline.py"
|
||||
],
|
||||
"US06-01": [
|
||||
"tests/integration/test_archive_preflight.py"
|
||||
]
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user