Files
photoanalyzer/tests/e2e/test_phase_b_pipeline.py

255 lines
10 KiB
Python

"""Phase B black-box end-to-end acceptance (US02-07).
Launches the real FastAPI server and durable worker as child processes and drives
them only over HTTP/SSE — no service/repository imports for the behaviour under
test. Covers the Phase B journeys the browser suite does not exhaustively assert:
analysis start/progress, resumable SSE reconnect, the polling fallback,
cancellation, per-asset error inspection, one-mutating-job rejection during
read-only browsing, the NSFW→vision privacy gate, and durability across restart.
External vision is the in-repo recording fake, invoked through the real worker.
"""
from __future__ import annotations
import json
import httpx
import pytest
from tests.e2e._pipeline_harness import Server, seed_library, start_worker, wait_until
pytestmark = pytest.mark.phase_b
TIMEOUT = 10
def _job_state(base: str, job_id: str) -> str:
return httpx.get(f"{base}/api/v1/jobs/{job_id}", timeout=TIMEOUT).json()["state"]
def _sse_events(
base: str,
job_id: str,
*,
after: int = 0,
last_event_id: int | None = None,
stop_after: int | None = None,
) -> list[dict]:
"""Read the resumable SSE stream, returning ``[{seq, type, message}, ...]``. When
``stop_after`` is set, close the stream after that many events to model a client
that drops mid-stream. ``last_event_id`` sends the reconnect header."""
headers = {}
if last_event_id is not None:
headers["Last-Event-ID"] = str(last_event_id)
url = f"{base}/api/v1/jobs/{job_id}/events/stream?after={after}"
events: list[dict] = []
seq: int | None = None
with httpx.stream("GET", url, headers=headers, timeout=TIMEOUT) as response:
response.raise_for_status()
for line in response.iter_lines():
if line.startswith("id:"):
seq = int(line[3:].strip())
elif line.startswith("data:"):
payload = json.loads(line[5:].strip())
if seq is not None: # skip the terminal `event: done` (empty {} data)
events.append({"seq": seq, **payload})
seq = None
if stop_after is not None and len(events) >= stop_after:
return events
elif line.startswith("event: done"):
break
return events
# ── analysis job: start, progress, and completion ────────────────────────────
def test_analysis_job_runs_to_completion_over_real_http(tmp_path):
seeded = seed_library(tmp_path, {"beach": 1}, {"beach": "sfw"})
log = tmp_path / "vision.log"
server = Server(seeded).start()
worker = start_worker(seeded, fake_vision_log=log)
try:
job = httpx.post(f"{server.base}/api/v1/analysis/jobs", timeout=TIMEOUT).json()
wait_until(lambda: _job_state(server.base, job["id"]) == "succeeded", timeout=TIMEOUT)
counts = httpx.get(f"{server.base}/api/v1/analysis/counts", timeout=TIMEOUT).json()
assert counts["analyzed"] == 1
# The fake actually ran through the real worker for the one SFW asset.
assert log.read_text().strip().endswith("beach.jpg")
finally:
worker.terminate()
server.stop()
# ── resumable SSE reconnect and the polling fallback ─────────────────────────
def test_sse_reconnect_resumes_without_gap_or_duplicate(tmp_path):
seeded = seed_library(tmp_path, {"beach": 1, "meadow": 2}, {"beach": "sfw", "meadow": "sfw"})
log = tmp_path / "vision.log"
server = Server(seeded).start()
worker = start_worker(seeded, fake_vision_log=log)
try:
job = httpx.post(f"{server.base}/api/v1/analysis/jobs", timeout=TIMEOUT).json()
wait_until(lambda: _job_state(server.base, job["id"]) == "succeeded", timeout=TIMEOUT)
full = _sse_events(server.base, job["id"])
assert [e["type"] for e in full][:2] == ["queued", "claimed"]
assert full[-1]["type"] == "state:succeeded"
assert len(full) >= 3
# A client that reads only the first event, drops, then reconnects with the
# last id it saw must receive exactly the remaining events — no gap, no dup.
first = _sse_events(server.base, job["id"], stop_after=1)
assert len(first) == 1
rest = _sse_events(server.base, job["id"], last_event_id=first[0]["seq"])
assert first + rest == full
seqs = [e["seq"] for e in full]
assert seqs == sorted(set(seqs)) # strictly increasing, unique
finally:
worker.terminate()
server.stop()
def test_polling_fallback_returns_the_same_durable_events(tmp_path):
seeded = seed_library(tmp_path, {"beach": 1}, {"beach": "sfw"})
log = tmp_path / "vision.log"
server = Server(seeded).start()
worker = start_worker(seeded, fake_vision_log=log)
try:
job = httpx.post(f"{server.base}/api/v1/analysis/jobs", timeout=TIMEOUT).json()
wait_until(lambda: _job_state(server.base, job["id"]) == "succeeded", timeout=TIMEOUT)
sse = _sse_events(server.base, job["id"])
polled = httpx.get(
f"{server.base}/api/v1/jobs/{job['id']}/events?after=0", timeout=TIMEOUT
).json()
assert polled["state"] == "succeeded"
assert [(e["seq"], e["type"]) for e in polled["events"]] == [
(e["seq"], e["type"]) for e in sse
]
finally:
worker.terminate()
server.stop()
# ── cancellation ─────────────────────────────────────────────────────────────
def test_queued_job_cancels_and_runs_no_items(tmp_path):
# No worker: the safety job stays queued so cancellation is deterministic.
seeded = seed_library(tmp_path, {"beach": 1, "city": 2}, {})
server = Server(seeded).start()
try:
job = httpx.post(f"{server.base}/api/v1/safety/jobs", timeout=TIMEOUT).json()
cancelled = httpx.post(
f"{server.base}/api/v1/jobs/{job['id']}/cancel", timeout=TIMEOUT
).json()
assert cancelled["state"] == "cancelled"
snapshot = httpx.get(f"{server.base}/api/v1/jobs/{job['id']}", timeout=TIMEOUT).json()
assert snapshot["progress"]["by_state"].get("succeeded", 0) == 0
finally:
server.stop()
# ── per-asset error inspection without losing completed work ──────────────────
def test_one_analysis_error_is_inspectable_and_others_still_complete(tmp_path):
seeded = seed_library(tmp_path, {"beach": 1, "boom": 2}, {"beach": "sfw", "boom": "sfw"})
log = tmp_path / "vision.log"
server = Server(seeded).start()
worker = start_worker(seeded, fake_vision_log=log)
try:
job = httpx.post(f"{server.base}/api/v1/analysis/jobs", timeout=TIMEOUT).json()
wait_until(lambda: _job_state(server.base, job["id"]) == "succeeded", timeout=TIMEOUT)
counts = httpx.get(f"{server.base}/api/v1/analysis/counts", timeout=TIMEOUT).json()
assert counts["analyzed"] == 1 and counts["error"] == 1 # boom failed, beach kept
boom_id = seeded.asset_ids["boom"]
result = httpx.get(
f"{server.base}/api/v1/analysis/results/{boom_id}", timeout=TIMEOUT
).json()
assert result["status"] == "error"
finally:
worker.terminate()
server.stop()
# ── one mutating job at a time, read-only browsing stays available ───────────
def test_second_mutation_rejected_while_reads_continue(tmp_path):
# No worker: the first safety job holds the library_write lock for the test.
seeded = seed_library(tmp_path, {"beach": 1, "city": 2}, {})
server = Server(seeded).start()
try:
first = httpx.post(f"{server.base}/api/v1/safety/jobs", timeout=TIMEOUT)
assert first.status_code == 200
second = httpx.post(f"{server.base}/api/v1/safety/jobs", timeout=TIMEOUT)
assert second.status_code == 409
assert second.json()["error"]["code"] == "lock_held"
# Read-only browsing is unaffected by the held mutation lock.
for path in ("/api/v1/workflow", "/api/v1/library/assets", "/api/v1/safety/queue"):
assert httpx.get(f"{server.base}{path}", timeout=TIMEOUT).status_code == 200
finally:
server.stop()
# ── privacy gate: NSFW assets never reach the vision provider ────────────────
def test_nsfw_asset_never_reaches_the_vision_provider(tmp_path):
seeded = seed_library(tmp_path, {"beach": 1, "city": 2}, {"beach": "sfw", "city": "nsfw"})
log = tmp_path / "vision.log"
server = Server(seeded).start()
worker = start_worker(seeded, fake_vision_log=log)
try:
job = httpx.post(f"{server.base}/api/v1/analysis/jobs", timeout=TIMEOUT).json()
wait_until(lambda: _job_state(server.base, job["id"]) == "succeeded", timeout=TIMEOUT)
logged = log.read_text()
assert "beach.jpg" in logged # the SFW asset was analysed
assert "city.jpg" not in logged # the NSFW asset never reached the provider
# And the NSFW asset was never even enqueued: no analysis row exists for it.
city_id = seeded.asset_ids["city"]
assert (
httpx.get(
f"{server.base}/api/v1/analysis/results/{city_id}", timeout=TIMEOUT
).status_code
== 404
)
finally:
worker.terminate()
server.stop()
# ── durability: decisions survive a full server restart ──────────────────────
def test_safety_decisions_survive_a_restart(tmp_path):
seeded = seed_library(tmp_path, {"beach": 1, "city": 2}, {"beach": "sfw", "city": "nsfw"})
server = Server(seeded).start()
try:
before = httpx.get(f"{server.base}/api/v1/safety/queue?state=nsfw", timeout=TIMEOUT).json()
assert any(row["asset_id"] == seeded.asset_ids["city"] for row in before["items"])
finally:
server.stop()
# Fresh process against the same data dir: the durable decision is still there.
restarted = Server(seeded).start()
try:
after = httpx.get(
f"{restarted.base}/api/v1/safety/queue?state=nsfw", timeout=TIMEOUT
).json()
assert any(row["asset_id"] == seeded.asset_ids["city"] for row in after["items"])
finally:
restarted.stop()