Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
6 changes: 6 additions & 0 deletions docs/development/IMPORT_REVIEW_RECOVERY.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
14 changes: 10 additions & 4 deletions src/pullbox/services/health_persistence.py
Original file line number Diff line number Diff line change
Expand Up @@ -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 (
Expand Down Expand Up @@ -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:
Expand Down
2 changes: 1 addition & 1 deletion src/pullbox/services/health_service.py
Original file line number Diff line number Diff line change
Expand Up @@ -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()
Expand Down
13 changes: 6 additions & 7 deletions src/pullbox/services/import_completed_cleanup.py
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand All @@ -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
Expand Down Expand Up @@ -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


Expand Down
12 changes: 3 additions & 9 deletions src/pullbox/services/import_known_series_recovery.py
Original file line number Diff line number Diff line change
Expand Up @@ -18,7 +18,6 @@
ImportedFileStatus,
ImportedSeries,
ImportJob,
ImportJobStatus,
ImportSeriesStatus,
ImportSourceType,
)
Expand All @@ -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
Expand Down Expand Up @@ -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] = {}
Expand Down
8 changes: 6 additions & 2 deletions src/pullbox/services/import_orphans.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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)

Expand Down
8 changes: 6 additions & 2 deletions src/pullbox/services/import_review_recheck.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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]
Expand Down
19 changes: 19 additions & 0 deletions src/pullbox/services/import_terminal_recovery.py
Original file line number Diff line number Diff line change
@@ -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
7 changes: 4 additions & 3 deletions src/pullbox/ui/import_results_context.py
Original file line number Diff line number Diff line change
Expand Up @@ -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:
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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,
Expand Down
5 changes: 3 additions & 2 deletions src/pullbox/ui/templates/partials/import_results.html
Original file line number Diff line number Diff line change
Expand Up @@ -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) %}

<div
x-data="importResultsData({
Expand All @@ -49,7 +50,7 @@
<div data-testid="import-follow-up-job-action-bar" class="rounded-xl border border-pb-border bg-pb-card/60 px-4 py-3">
<div class="flex flex-wrap items-center justify-between gap-2">
<div class="flex flex-wrap items-center gap-2">
{% 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 %}
<button
type="button"
data-testid="import-results-retry-action"
Expand Down Expand Up @@ -300,7 +301,7 @@ <h2 class="text-2xl font-semibold tracking-tight text-pb-text">
{% set safety_category_summaries = safety_category_summaries|default([]) %}
{% set cleanup_safe_action_count = cleanup_safe_action_count|default(0) %}
{% set cleanup_needs_review_count = cleanup_needs_review_count|default(0) %}
{% if job.status.value == 'completed' and not job.archived_at and (cleanup_action_summaries or cleanup_needs_review_count > 0) %}
{% if recovery_actions_available and (cleanup_action_summaries or cleanup_needs_review_count > 0) %}
<section data-testid="import-results-recovery-dashboard" class="section-card overflow-hidden">
<div class="section-header">
<div class="section-header-copy">
Expand Down
45 changes: 45 additions & 0 deletions tests/ui/test_import_results_context.py
Original file line number Diff line number Diff line change
Expand Up @@ -181,6 +181,51 @@ async def test_results_count_retryable_failed_inspection_as_safe_next_step(db_se
assert context["cleanup_action_summaries"][0]["action"] == "retry_source_inspection"


@pytest.mark.asyncio
async def test_results_keep_recovery_actions_for_post_completion_failure(db_session) -> None: # type: ignore[no-untyped-def]
from pullbox.ui.import_results_context import load_import_results_context

job = ImportJob(
source_path="/tmp/comics",
source_type=ImportSourceType.FILESYSTEM,
status=ImportJobStatus.FAILED,
import_completed_at=datetime.now(UTC),
error_message="Optional follow-up failed after canonical import completed.",
)
db_session.add(job)
await db_session.flush()
imported_series = ImportedSeries(
import_job_id=job.id,
raw_series_name="Recoverable import",
status=ImportSeriesStatus.IMPORTED,
)
db_session.add(imported_series)
await db_session.flush()
db_session.add(
ImportedFile(
import_job_id=job.id,
import_series_id=imported_series.id,
file_path="/tmp/comics/missing.cbz",
file_name="missing.cbz",
file_size=1024,
file_format="cbz",
status=ImportedFileStatus.SAFETY_BLOCKED,
diagnostics={
"safety_block": build_import_safety_diagnostics(
ImportSafetyCategory.SOURCE_MISSING.value,
code=ImportSafetyCategory.SOURCE_MISSING.value,
)
},
)
)
await db_session.flush()

context = await load_import_results_context(db_session, job)

assert context["recovery_actions_available"] is True
assert context["cleanup_action_summaries"][0]["action"] == "dismiss_missing_references"


@pytest.mark.asyncio
async def test_load_import_results_context_reports_changed_sources_separately(
db_session,
Expand Down
25 changes: 24 additions & 1 deletion tests/ui/test_import_ui_routes.py
Original file line number Diff line number Diff line change
Expand Up @@ -1116,7 +1116,11 @@ def test_failed_results_template_preserves_bounded_safety_actions(self) -> None:
diagnostics={"safety_block": {"overrideable": True, "reason": "Too large"}},
)
html = templates.env.get_template("partials/import_results.html").render(
job=SimpleNamespace(id=35, status=SimpleNamespace(value="failed")),
job=SimpleNamespace(
id=35,
status=SimpleNamespace(value="failed"),
archived_at=None,
),
can_rollback=False,
imported_count=0,
failed_count=1,
Expand All @@ -1142,6 +1146,23 @@ def test_failed_results_template_preserves_bounded_safety_actions(self) -> None:
files_safety_blocked=105,
safety_blocked_files=[blocked_file],
safety_blocked_files_truncated=5,
recovery_actions_available=True,
cleanup_action_summaries=[
{
"action": "dismiss_missing_references",
"label": "Dismiss stale Mylar references",
"description": "Clear stale references without deleting files.",
"button_label": "Dismiss references",
"tone": "neutral",
"affected_count": 1,
"affected_file_count": 1,
"item_unit": "file",
"examples": ("missing.cbz",),
}
],
safety_category_summaries=[],
cleanup_safe_action_count=1,
cleanup_needs_review_count=0,
resume_step=5,
resume_job_id=35,
resume_progress_snapshot={},
Expand All @@ -1152,6 +1173,8 @@ def test_failed_results_template_preserves_bounded_safety_actions(self) -> None:
assert "large-tpb.cbz" in html
assert "5 more safety exceptions are not shown" in html
assert 'data-testid="import-results-safety-retry-91"' in html
assert 'data-testid="import-results-recovery-dashboard"' in html
assert "Dismiss stale Mylar references" in html

def test_results_template_explains_changed_source_failures(self) -> None:
from types import SimpleNamespace
Expand Down
23 changes: 23 additions & 0 deletions tests/unit/test_health_checks.py
Original file line number Diff line number Diff line change
Expand Up @@ -2540,6 +2540,29 @@ async def exploding_check() -> CheckOutcome:
assert outcomes[0].status == HealthStatus.UNHEALTHY
assert "boom" in outcomes[0].message.lower()

@pytest.mark.asyncio
async def test_failed_check_always_rolls_back_before_persistence(
self,
settings: MagicMock,
) -> None:
service = _make_service(settings)
session = MagicMock()
session.is_active = True
session.rollback = AsyncMock()

async def cancelled_database_check() -> CheckOutcome:
raise TimeoutError

outcomes = await service._safe_run(
cancelled_database_check(),
"database",
"connectivity",
session=session,
)

assert outcomes[0].status == HealthStatus.UNHEALTHY
session.rollback.assert_awaited_once_with()

@pytest.mark.asyncio
async def test_actionable_guidance_present(self, db_session: AsyncSession) -> None:
"""Non-healthy outcomes should have non-empty guidance."""
Expand Down
Loading
Loading