US06-02: Transfer, Verify, and Remove Active Sources (#78)

This commit was merged in pull request #78.
This commit is contained in:
2026-08-16 18:47:04 +02:00
parent b4b316acf5
commit b90d1718be
14 changed files with 2167 additions and 24 deletions

View File

@@ -1,9 +1,10 @@
"""Archive location and preflight API (US06-01).
"""Archive location, preflight, and plan API (US06-01, US06-02).
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.
the token an archive plan must present. Creating a plan writes only database rows —
the transfer itself runs on the durable ``archive`` lane, never in the request
thread, because it removes originals.
"""
from __future__ import annotations
@@ -12,12 +13,17 @@ from fastapi import APIRouter, Request
from fastapi.responses import JSONResponse
from pydantic import BaseModel
from photo_pipeline.jobs.domain_handlers import ARCHIVE_LOCK, ARCHIVE_PLAN
from photo_pipeline.services.archives import ArchiveError, ArchiveService
from photo_pipeline.services.archive_transfer import ArchiveTransferService
from photo_pipeline.services.jobs import JobBlocked, JobService
router = APIRouter(tags=["archives"])
# Which failures are the caller's request (422) and which are a missing thing (404).
NOT_FOUND_CODES = {"unknown_location"}
# Which failures are the caller's request (422), a missing thing (404), or state
# that changed under the caller (409).
NOT_FOUND_CODES = {"unknown_location", "unknown_plan"}
CONFLICT_CODES = {"stale_token", "stale_plan", "archive_pending"}
class RegisterLocationRequest(BaseModel):
@@ -31,12 +37,28 @@ class PreflightRequest(BaseModel):
albums: list[str] | None = None
class CreatePlanRequest(PreflightRequest):
# The token of the preflight the user approved; a stale one is refused.
token: str
def _service(request: Request) -> ArchiveService:
return ArchiveService(request.app.state.session_factory, config=request.app.state.config)
def _transfers(request: Request) -> ArchiveTransferService:
return ArchiveTransferService(
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
if error.code in NOT_FOUND_CODES:
status = 404
elif error.code in CONFLICT_CODES:
status = 409
else:
status = 422
return JSONResponse(
status_code=status, content={"error": {"code": error.code, "message": str(error)}}
)
@@ -61,3 +83,67 @@ def preflight(body: PreflightRequest, request: Request):
return _service(request).preflight(body.location_id, body.albums)
except ArchiveError as error:
return _error(error)
@router.post("/archive-plans", status_code=201)
def create_plan(body: CreatePlanRequest, request: Request):
"""Turn an approved preflight into a durable, journaled plan. Nothing moves."""
try:
return _transfers(request).create(body.location_id, body.albums, token=body.token)
except ArchiveError as error:
return _error(error)
@router.get("/archive-plans")
def list_plans(request: Request) -> dict:
return {"plans": _transfers(request).list()}
@router.get("/archive-plans/{plan_id}")
def get_plan(plan_id: str, request: Request):
plan = _transfers(request).get(plan_id)
if plan is None:
return _error(ArchiveError("unknown_plan", f"unknown archive plan {plan_id}"))
return plan
@router.post("/archive-plans/{plan_id}/apply")
def apply_plan(plan_id: str, request: Request):
"""Queue the transfer on the archiver lane. The worker removes the sources."""
service = _transfers(request)
plan = service.get(plan_id)
if plan is None:
return _error(ArchiveError("unknown_plan", f"unknown archive plan {plan_id}"))
if service.journal.blocks_mutation():
unresolved = [row for row in service.journal.incomplete() if row["plan_id"] != plan_id]
if unresolved:
return _error(
ArchiveError(
"archive_pending",
f"an unresolved archive operation ({unresolved[0]['id']}) must be recovered",
)
)
try:
job = JobService(request.app.state.session_factory).enqueue(
ARCHIVE_PLAN,
lock=ARCHIVE_LOCK,
# One queued attempt per plan version: a double-clicked apply reuses it.
idempotency_key=f"archive:{plan_id}:{plan['version']}",
items=[plan_id],
)
except JobBlocked as error:
return JSONResponse(
status_code=409, content={"error": {"code": error.code, "message": str(error)}}
)
return {"plan_id": plan_id, "job": job}
@router.get("/archive-recovery")
def recovery_status(request: Request) -> dict:
"""What an interrupted transfer left behind, straight from journal + disk."""
return _transfers(request).recovery_status()
@router.post("/archive-recovery/resolve")
def resolve_recovery(request: Request) -> dict:
return _transfers(request).recover()