From 9b37740dbe940776f8caeab71e39a6e3a9a279e4 Mon Sep 17 00:00:00 2001 From: Adam Hernandez Date: Sun, 13 Sep 2026 20:37:46 -0700 Subject: [PATCH] fix(import): restore terminal recovery actions --- docs/development/IMPORT_REVIEW_RECOVERY.md | 6 +++ src/pullbox/services/health_persistence.py | 14 ++++-- src/pullbox/services/health_service.py | 2 +- .../services/import_completed_cleanup.py | 13 +++-- .../services/import_known_series_recovery.py | 12 ++--- src/pullbox/services/import_orphans.py | 8 ++- src/pullbox/services/import_review_recheck.py | 8 ++- .../services/import_terminal_recovery.py | 19 +++++++ src/pullbox/ui/import_results_context.py | 7 +-- .../ui/templates/partials/import_results.html | 5 +- tests/ui/test_import_results_context.py | 45 +++++++++++++++++ tests/ui/test_import_ui_routes.py | 25 +++++++++- tests/unit/test_health_checks.py | 23 +++++++++ tests/unit/test_health_service_persistence.py | 48 +++++++++++++++++- tests/unit/test_import_completed_cleanup.py | 18 +++++++ .../unit/test_import_known_series_recovery.py | 21 ++++++++ tests/unit/test_import_orphans.py | 10 +++- tests/unit/test_import_terminal_recovery.py | 50 +++++++++++++++++++ 18 files changed, 301 insertions(+), 33 deletions(-) create mode 100644 src/pullbox/services/import_terminal_recovery.py create mode 100644 tests/unit/test_import_terminal_recovery.py diff --git a/docs/development/IMPORT_REVIEW_RECOVERY.md b/docs/development/IMPORT_REVIEW_RECOVERY.md index 17f362f1..56f091ce 100644 --- a/docs/development/IMPORT_REVIEW_RECOVERY.md +++ b/docs/development/IMPORT_REVIEW_RECOVERY.md @@ -49,6 +49,12 @@ above-fold count and a direct link to that import's Follow-up workspace. Available recovery actions are intentionally narrow: +A later task or optional Story Arc failure does not hide these actions after +the canonical import has durably completed. Pullbox also requires the job to be +unarchived, idle, and free of rollback work before exposing or applying any +recovery action. An import that failed before durable completion remains +ineligible. + - **Dismiss stale Mylar references** marks missing database references skipped. It does not delete a review record or touch Mylar's database. - **Skip one-page archives** excludes one-page image archives while leaving the diff --git a/src/pullbox/services/health_persistence.py b/src/pullbox/services/health_persistence.py index 8f625965..ed98c57e 100644 --- a/src/pullbox/services/health_persistence.py +++ b/src/pullbox/services/health_persistence.py @@ -11,7 +11,7 @@ import structlog from sqlalchemy import delete, select -from sqlalchemy.exc import OperationalError +from sqlalchemy.exc import OperationalError, PendingRollbackError from sqlalchemy.ext.asyncio import async_sessionmaker from pullbox.core.sqlite_lock import ( @@ -206,16 +206,22 @@ async def persist_health_outcomes( await active_session.commit() await session.commit() return - except OperationalError as exc: + except (OperationalError, PendingRollbackError) as exc: await active_session.rollback() - if not is_sqlite_locked_error(exc) or attempt == lock_retry_attempts: + recoverable = isinstance(exc, PendingRollbackError) or is_sqlite_locked_error(exc) + if not recoverable or attempt == lock_retry_attempts: raise delay_seconds = retry_delay(attempt) logger.warning( - "health_result_persist_retrying_after_sqlite_lock", + ( + "health_result_persist_retrying_after_pending_rollback" + if isinstance(exc, PendingRollbackError) + else "health_result_persist_retrying_after_sqlite_lock" + ), attempt=attempt, max_attempts=lock_retry_attempts, delay_seconds=delay_seconds, + failure_type=type(exc).__name__, ) finally: if manage_commit: diff --git a/src/pullbox/services/health_service.py b/src/pullbox/services/health_service.py index b39cbccf..ccd2a8e2 100644 --- a/src/pullbox/services/health_service.py +++ b/src/pullbox/services/health_service.py @@ -528,7 +528,7 @@ async def _safe_run( @staticmethod async def _rollback_failed_check_session(session: AsyncSession | None) -> None: """Restore a failed check transaction before later checks or persistence.""" - if session is None or session.is_active: + if session is None: return try: await session.rollback() diff --git a/src/pullbox/services/import_completed_cleanup.py b/src/pullbox/services/import_completed_cleanup.py index 35a7b0ce..fffb89ca 100644 --- a/src/pullbox/services/import_completed_cleanup.py +++ b/src/pullbox/services/import_completed_cleanup.py @@ -20,7 +20,6 @@ from pullbox.core.name_matcher import NameMatcher from pullbox.models.audit_log import AuditEventType from pullbox.models.import_job import ( - ImportControlRequest, ImportedFile, ImportedFileStatus, ImportedSeries, @@ -40,6 +39,7 @@ from pullbox.services.import_story_arc_resolution import ( refresh_story_arc_entries_for_import_files, ) +from pullbox.services.import_terminal_recovery import allows_terminal_import_recovery if TYPE_CHECKING: from sqlalchemy.ext.asyncio import AsyncSession @@ -638,12 +638,11 @@ async def _load_completed_job(session: AsyncSession, job_id: int) -> ImportJob: job = await session.get(ImportJob, job_id, populate_existing=True) if job is None: raise NotFoundError("ImportJob", job_id) - if job.status is not ImportJobStatus.COMPLETED: - raise ValidationError("Job must be in COMPLETED state for recovery cleanup") - if job.control_request is not ImportControlRequest.NONE: - raise ValidationError("The import job has a pending control request") - if job.archived_at is not None: - raise ValidationError("Archived import jobs must be restored before cleanup") + if not allows_terminal_import_recovery(job): + raise ValidationError( + "Job must have a COMPLETED canonical import with no pending control or rollback " + "work for recovery cleanup" + ) return job diff --git a/src/pullbox/services/import_known_series_recovery.py b/src/pullbox/services/import_known_series_recovery.py index 59f0946e..d077924d 100644 --- a/src/pullbox/services/import_known_series_recovery.py +++ b/src/pullbox/services/import_known_series_recovery.py @@ -18,7 +18,6 @@ ImportedFileStatus, ImportedSeries, ImportJob, - ImportJobStatus, ImportSeriesStatus, ImportSourceType, ) @@ -31,6 +30,7 @@ build_import_metadata_conflict, source_metadata_for_import_file, ) +from pullbox.services.import_terminal_recovery import allows_terminal_import_recovery if TYPE_CHECKING: from sqlalchemy.ext.asyncio import AsyncSession @@ -213,14 +213,8 @@ async def load_known_series_recovery( job_id: int, ) -> tuple[KnownSeriesRecovery, ...]: """Return only exact recoverable file identities without modifying the session.""" - job = ( - await session.execute( - select(ImportJob.status, ImportJob.source_type).where( - ImportJob.id == job_id, - ) - ) - ).one_or_none() - if job is None or job.status is not ImportJobStatus.COMPLETED: + job = await session.get(ImportJob, job_id) + if job is None or not allows_terminal_import_recovery(job): return () plans: list[KnownSeriesRecovery] = [] matched_local_ids: dict[int, int | None] = {} diff --git a/src/pullbox/services/import_orphans.py b/src/pullbox/services/import_orphans.py index 70f3ad5c..041c734f 100644 --- a/src/pullbox/services/import_orphans.py +++ b/src/pullbox/services/import_orphans.py @@ -27,6 +27,7 @@ ) from pullbox.services.import_retry_helpers import require_retained_import_destination from pullbox.services.import_review_recheck import prepare_retryable_failed_sources_for_retry +from pullbox.services.import_terminal_recovery import allows_terminal_import_recovery if TYPE_CHECKING: from collections.abc import Awaitable, Callable, Sequence @@ -527,8 +528,11 @@ async def retry_failed_series( if job is None: raise NotFoundError("ImportJob", job_id) - if job.status != ImportJobStatus.COMPLETED: - raise ValidationError(f"Job must be in COMPLETED state to retry (current: {job.status})") + if not allows_terminal_import_recovery(job): + raise ValidationError( + "Job must have a COMPLETED canonical import with no pending control or rollback " + f"work to retry (current: {job.status})" + ) require_retained_import_destination(job) diff --git a/src/pullbox/services/import_review_recheck.py b/src/pullbox/services/import_review_recheck.py index 5c4ecaaa..9219f8b9 100644 --- a/src/pullbox/services/import_review_recheck.py +++ b/src/pullbox/services/import_review_recheck.py @@ -41,6 +41,7 @@ from pullbox.services.import_safety_diagnostics import build_import_safety_diagnostics from pullbox.services.import_series_match_state import clear_auto_cv_match_fields from pullbox.services.import_source_metadata import source_metadata_for_import_file +from pullbox.services.import_terminal_recovery import allows_terminal_import_recovery if TYPE_CHECKING: from collections.abc import Sequence @@ -453,8 +454,11 @@ async def prepare_completed_import_file_recheck( job = await session.get(ImportJob, job_id) if job is None: raise NotFoundError("ImportJob", job_id) - if job.status != ImportJobStatus.COMPLETED or job.control_request != ImportControlRequest.NONE: - raise ValidationError("Job must be idle in COMPLETED before a failed-file recheck") + if not allows_terminal_import_recovery(job): + raise ValidationError( + "Job must have a COMPLETED canonical import with no pending control or rollback " + "work before a failed-file recheck" + ) if not source_roots: raise ValidationError("At least one explicit source root is required") roots = [(path.expanduser().absolute(), resolve_preview_source(path)) for path in source_roots] diff --git a/src/pullbox/services/import_terminal_recovery.py b/src/pullbox/services/import_terminal_recovery.py new file mode 100644 index 00000000..fe77456b --- /dev/null +++ b/src/pullbox/services/import_terminal_recovery.py @@ -0,0 +1,19 @@ +"""Shared eligibility rules for safe actions after a durable import completion.""" + +from __future__ import annotations + +from pullbox.models.import_job import ImportControlRequest, ImportJob, ImportJobStatus + + +def allows_terminal_import_recovery(job: ImportJob) -> bool: + """Return whether post-import actions may safely operate on this job.""" + if ( + job.archived_at is not None + or job.control_request is not ImportControlRequest.NONE + or job.story_arc_rollback_waiting_work_id is not None + or dict(job.progress_snapshot or {}).get("mode") == "rollback" + ): + return False + if job.status is ImportJobStatus.COMPLETED: + return True + return job.status is ImportJobStatus.FAILED and job.import_completed_at is not None diff --git a/src/pullbox/ui/import_results_context.py b/src/pullbox/ui/import_results_context.py index 82042c33..19559de4 100644 --- a/src/pullbox/ui/import_results_context.py +++ b/src/pullbox/ui/import_results_context.py @@ -45,6 +45,7 @@ ImportSafetyCategory, import_safety_category_label, ) +from pullbox.services.import_terminal_recovery import allows_terminal_import_recovery from pullbox.services.import_workflow_state import import_control_state_for_job if TYPE_CHECKING: @@ -996,10 +997,9 @@ async def load_import_results_context( safety_category_summaries = ( await _load_safety_category_summaries(session, job_id) if files_safety_blocked > 0 else [] ) + recovery_actions_available = allows_terminal_import_recovery(job) cleanup_action_summaries = ( - await _load_cleanup_action_summaries(session, job_id) - if job.status is ImportJobStatus.COMPLETED and job.archived_at is None - else [] + await _load_cleanup_action_summaries(session, job_id) if recovery_actions_available else [] ) misplaced_source_restore_count = 0 misplaced_source_duplicate_count = 0 @@ -1137,6 +1137,7 @@ async def load_import_results_context( ), "safety_category_summaries": safety_category_summaries, "cleanup_action_summaries": cleanup_action_summaries, + "recovery_actions_available": recovery_actions_available, "recommended_conflict_groups": recommended_conflict_groups, "recommended_conflict_files": recommended_conflict_files, "already_owned_conflict_files": already_owned_conflict_files, diff --git a/src/pullbox/ui/templates/partials/import_results.html b/src/pullbox/ui/templates/partials/import_results.html index 16afc257..b51e62ef 100644 --- a/src/pullbox/ui/templates/partials/import_results.html +++ b/src/pullbox/ui/templates/partials/import_results.html @@ -31,6 +31,7 @@ - rollback_actions_rolled_back: int -#} {% from "components/settings_shell.html" import inline_alert %} +{% set recovery_actions_available = recovery_actions_available|default(job.status.value == 'completed' and not job.archived_at) %}
- {% if (failed_count > 0 or files_failed > 0) and not job.removed_library_root_snapshot %} + {% if recovery_actions_available and (failed_count > 0 or files_failed > 0) and not job.removed_library_root_snapshot %}