Compare commits
1 Commits
us/US07-07
...
us/US05-06
| Author | SHA1 | Date | |
|---|---|---|---|
| fe027eada9 |
33
README.md
33
README.md
@@ -136,3 +136,36 @@ 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 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.
|
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.
|
||||||
|
|||||||
@@ -96,6 +96,16 @@ class Worker:
|
|||||||
return
|
return
|
||||||
try:
|
try:
|
||||||
if cancelled or snapshot["state"] == JobState.CANCELLING:
|
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(
|
self.service.transition(
|
||||||
job_id, JobState.CANCELLED, worker_id=self.worker_id, fencing_token=token
|
job_id, JobState.CANCELLED, worker_id=self.worker_id, fencing_token=token
|
||||||
)
|
)
|
||||||
|
|||||||
@@ -31,4 +31,5 @@ markers = [
|
|||||||
"phase_b: Phase B end-to-end acceptance (US02-07) — API, worker-recovery, and browser journeys",
|
"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_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_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",
|
||||||
]
|
]
|
||||||
|
|||||||
@@ -9,14 +9,18 @@ and worker, never mocked inside a test.
|
|||||||
|
|
||||||
from __future__ import annotations
|
from __future__ import annotations
|
||||||
|
|
||||||
|
import json
|
||||||
import os
|
import os
|
||||||
import socket
|
import socket
|
||||||
|
import stat
|
||||||
import subprocess
|
import subprocess
|
||||||
import sys
|
import sys
|
||||||
|
import threading
|
||||||
import time
|
import time
|
||||||
import uuid
|
import uuid
|
||||||
from dataclasses import dataclass
|
from dataclasses import dataclass
|
||||||
from datetime import datetime, timezone
|
from datetime import datetime, timezone
|
||||||
|
from http.server import BaseHTTPRequestHandler, HTTPServer
|
||||||
from pathlib import Path
|
from pathlib import Path
|
||||||
|
|
||||||
import httpx
|
import httpx
|
||||||
@@ -258,3 +262,182 @@ def approve_album(base: str, *, album: str = "rome", name: str) -> None:
|
|||||||
json={**payload, "expected_version": current["version"]},
|
json={**payload, "expected_version": current["version"]},
|
||||||
timeout=10,
|
timeout=10,
|
||||||
).raise_for_status()
|
).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()
|
||||||
|
|||||||
568
tests/e2e/test_phase_e_pipeline.py
Normal file
568
tests/e2e/test_phase_e_pipeline.py
Normal file
@@ -0,0 +1,568 @@
|
|||||||
|
"""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,6 +13,7 @@ REPO = Path(__file__).resolve().parents[2]
|
|||||||
MAP = json.loads((REPO / "tests" / "story_traceability.json").read_text())["stories"]
|
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_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_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():
|
def test_all_phase_a_stories_are_mapped():
|
||||||
@@ -25,6 +26,12 @@ def test_all_phase_d_stories_are_mapped():
|
|||||||
assert PHASE_D_STORIES <= set(MAP)
|
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():
|
def test_every_mapped_test_file_exists_and_is_nonempty():
|
||||||
for story, files in MAP.items():
|
for story, files in MAP.items():
|
||||||
assert files, f"{story} maps to no tests"
|
assert files, f"{story} maps to no tests"
|
||||||
|
|||||||
@@ -13,180 +13,31 @@ API key is a sentinel string, so the last test can prove it never reached the pa
|
|||||||
|
|
||||||
from __future__ import annotations
|
from __future__ import annotations
|
||||||
|
|
||||||
import json
|
|
||||||
import stat
|
|
||||||
import threading
|
|
||||||
from http.server import BaseHTTPRequestHandler, HTTPServer
|
|
||||||
from pathlib import Path
|
|
||||||
|
|
||||||
import httpx
|
import httpx
|
||||||
import pytest
|
import pytest
|
||||||
from playwright.sync_api import expect
|
from playwright.sync_api import expect
|
||||||
|
|
||||||
from tests.e2e._pipeline_harness import (
|
from tests.e2e._pipeline_harness import (
|
||||||
Server,
|
SENTINEL_KEY,
|
||||||
|
SILENT_UPLOADER,
|
||||||
|
UPLOADER_VERSION,
|
||||||
|
UploadStack,
|
||||||
|
mark_upload_ready,
|
||||||
seed_album,
|
seed_album,
|
||||||
session_factory,
|
session_factory,
|
||||||
start_worker,
|
|
||||||
wait_until,
|
wait_until,
|
||||||
)
|
)
|
||||||
|
|
||||||
|
pytestmark = pytest.mark.phase_e
|
||||||
|
|
||||||
TIMEOUT = 10
|
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
|
@pytest.fixture
|
||||||
def stack(tmp_path):
|
def stack(tmp_path):
|
||||||
seeded = seed_album(tmp_path)
|
seeded = seed_album(tmp_path)
|
||||||
_mark_upload_ready(seeded)
|
mark_upload_ready(seeded)
|
||||||
running = Stack(tmp_path, seeded, FakeImmich())
|
running = UploadStack(tmp_path, seeded)
|
||||||
try:
|
try:
|
||||||
yield running
|
yield running
|
||||||
finally:
|
finally:
|
||||||
@@ -236,7 +87,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):
|
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)
|
stack.start(worker=False)
|
||||||
_open(page, stack)
|
_open(page, stack)
|
||||||
|
|
||||||
@@ -422,8 +273,7 @@ def test_an_interrupted_attempt_is_shown_as_uncertain_after_a_restart(page, stac
|
|||||||
session.get(UploadBatch, created["id"]).state = "running"
|
session.get(UploadBatch, created["id"]).state = "running"
|
||||||
session.commit()
|
session.commit()
|
||||||
|
|
||||||
stack.server.stop()
|
stack.restart_server() # the same port, so recovery runs in a genuinely fresh process
|
||||||
stack.server.start() # the same port, so recovery runs in a genuinely fresh process
|
|
||||||
page.reload()
|
page.reload()
|
||||||
|
|
||||||
expect(page.get_by_test_id("detail-state")).to_have_text("unknown_requires_verification")
|
expect(page.get_by_test_id("detail-state")).to_have_text("unknown_requires_verification")
|
||||||
|
|||||||
@@ -38,6 +38,8 @@ from photo_pipeline.services.upload_batches import (
|
|||||||
)
|
)
|
||||||
from photo_pipeline.services.uploads import UploadService
|
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)
|
NOW = datetime(2026, 1, 1, tzinfo=timezone.utc)
|
||||||
SENTINEL_KEY = "immich-sentinel-9f3a2b"
|
SENTINEL_KEY = "immich-sentinel-9f3a2b"
|
||||||
UPLOADER_VERSION = "immich-go 0.21.0"
|
UPLOADER_VERSION = "immich-go 0.21.0"
|
||||||
|
|||||||
@@ -29,6 +29,8 @@ from photo_pipeline.models import (
|
|||||||
from photo_pipeline.services.hashing import sha256_file
|
from photo_pipeline.services.hashing import sha256_file
|
||||||
from photo_pipeline.services.uploads import UploadError, UploadService
|
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)
|
NOW = datetime(2026, 1, 1, tzinfo=timezone.utc)
|
||||||
# A sentinel credential: every assertion below proves it never leaves configuration.
|
# A sentinel credential: every assertion below proves it never leaves configuration.
|
||||||
SENTINEL_KEY = "immich-sentinel-9f3a2b"
|
SENTINEL_KEY = "immich-sentinel-9f3a2b"
|
||||||
|
|||||||
@@ -32,6 +32,8 @@ from photo_pipeline.services.upload_reports import (
|
|||||||
)
|
)
|
||||||
from photo_pipeline.services.uploads import UploadService
|
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)
|
NOW = datetime(2026, 1, 1, tzinfo=timezone.utc)
|
||||||
SUPPORTED_VERSION = "immich-go 0.21.0"
|
SUPPORTED_VERSION = "immich-go 0.21.0"
|
||||||
|
|
||||||
|
|||||||
@@ -40,6 +40,8 @@ from photo_pipeline.services.upload_verification import (
|
|||||||
)
|
)
|
||||||
from photo_pipeline.services.uploads import UploadService
|
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)
|
NOW = datetime(2026, 1, 1, tzinfo=timezone.utc)
|
||||||
SENTINEL_KEY = "immich-sentinel-9f3a2b"
|
SENTINEL_KEY = "immich-sentinel-9f3a2b"
|
||||||
UPLOADER_VERSION = "immich-go 0.21.0"
|
UPLOADER_VERSION = "immich-go 0.21.0"
|
||||||
|
|||||||
@@ -90,6 +90,26 @@ def test_cooperative_cancellation_leaves_items_resumable(sf, jobs):
|
|||||||
assert by_state.get(ItemState.QUEUED) == 1 # "b" left resumable
|
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):
|
def test_fencing_rejects_superseded_worker(sf, jobs):
|
||||||
job = jobs.enqueue("scan", items=["a"])
|
job = jobs.enqueue("scan", items=["a"])
|
||||||
stale = jobs.claim(["scan"], "old")
|
stale = jobs.claim(["scan"], "old")
|
||||||
|
|||||||
@@ -120,6 +120,9 @@
|
|||||||
],
|
],
|
||||||
"US05-05": [
|
"US05-05": [
|
||||||
"tests/e2e/test_uploads_ui.py"
|
"tests/e2e/test_uploads_ui.py"
|
||||||
|
],
|
||||||
|
"US05-06": [
|
||||||
|
"tests/e2e/test_phase_e_pipeline.py"
|
||||||
]
|
]
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user