US06-04: Plan and Execute Safe Restores (#80)

This commit was merged in pull request #80.
This commit is contained in:
2026-08-16 19:49:21 +02:00
parent 439eb9971e
commit fb8567284d
10 changed files with 1294 additions and 43 deletions

View File

@@ -1,9 +1,9 @@
"""Domain job handlers: safety scoring, content analysis, uploads, archive
transfers (US02-06, US05-02, US06-02).
transfers, restores (US02-06, US05-02, US06-02, US06-04).
Importing this module registers the ``safety_score``, ``analysis``,
``upload_batch``, and ``archive_plan`` job types so the generic worker can run them
per item. Each handler delegates to its service, which owns the real work and the
``upload_batch``, ``archive_plan``, and ``restore_plan`` job types so the generic
worker can run them per item. Each handler delegates to its service, which owns the real work and the
privacy gate. Handlers are idempotent: re-scoring or re-analyzing one asset is safe
after an interrupted attempt, an upload batch refuses to re-run an attempt whose
outcome is unknown, and an archive plan skips items it already completed.
@@ -20,6 +20,7 @@ SAFETY_SCORE = "safety_score"
ANALYSIS = "analysis"
UPLOAD_BATCH = "upload_batch"
ARCHIVE_PLAN = "archive_plan"
RESTORE_PLAN = "restore_plan"
# Both mutate the library's metadata/derived state; one at a time (concept §one job).
LIBRARY_WRITE_LOCK = "library_write"
# The uploader lane: one album batch at a time (concept §16).
@@ -70,7 +71,22 @@ def _archive_plan_item(plan_id: str, ctx: JobContext) -> None:
raise RuntimeError(f"archive plan {plan_id}: {result['failed']} item(s) failed")
def _restore_plan_item(plan_id: str, ctx: JobContext) -> None:
"""One item = one restore plan. A restore removes nothing, so an item failure
simply leaves that asset archived (US06-04)."""
from photo_pipeline.config import Config
from photo_pipeline.services.restores import RestoreService
config = ctx.config if ctx.config is not None else Config.from_env()
result = RestoreService(ctx.session_factory, config=config).apply(
plan_id, worker_id=ctx.worker_id
)
if result["failed"]:
raise RuntimeError(f"restore plan {plan_id}: {result['failed']} item(s) failed")
register(SAFETY_SCORE, _safety_score_item)
register(ANALYSIS, _analysis_item)
register(UPLOAD_BATCH, _upload_batch_item)
register(ARCHIVE_PLAN, _archive_plan_item)
register(RESTORE_PLAN, _restore_plan_item)