From 8d599b72f352e2f42aa3d228d6fc5df80ba25f32 Mon Sep 17 00:00:00 2001 From: Adam Hernandez Date: Mon, 14 Sep 2026 15:40:43 -0700 Subject: [PATCH 1/2] fix(import): harden recovery checkpoints and retries --- docs/development/IMPORT_PERFORMANCE.md | 8 + docs/development/IMPORT_REVIEW_RECOVERY.md | 8 + .../api/v1/import_completed_cleanup.py | 3 + .../services/diagnostic_db_snapshot.py | 12 ++ .../import_deferred_recovery_execution.py | 88 ++++++++--- src/pullbox/services/import_review_recheck.py | 39 +++-- .../api/test_import_completed_cleanup_api.py | 89 +++++++++++ tests/unit/test_diagnostic_db_snapshot.py | 20 +++ ...test_import_deferred_recovery_execution.py | 102 +++++++++++++ tests/unit/test_import_review_recheck.py | 144 ++++++++++++++++++ 10 files changed, 480 insertions(+), 33 deletions(-) diff --git a/docs/development/IMPORT_PERFORMANCE.md b/docs/development/IMPORT_PERFORMANCE.md index c9bb9078..7de632bb 100644 --- a/docs/development/IMPORT_PERFORMANCE.md +++ b/docs/development/IMPORT_PERFORMANCE.md @@ -12,6 +12,14 @@ Archive page-name matching runs off the event loop and uses a bounded local cache for repeated page-title parsing. It still counts every page toward consensus and does not reuse safety decisions between scans. +Completed-import source recovery inspects each bounded page before its database +mutation phase, then commits before inspecting the next page. Slow files +therefore do not hold SQLite's writer lock, and results do not accumulate in an +unbounded in-memory collection. +Deferred catalog recovery publishes live provider progress without rewriting +the full recovery snapshot, then writes one JSON-safe checkpoint after each +completed catalog. + `PULLBOX_IMPORT_SCAN_WORKER_COUNT=0` selects automatic inspection concurrency, up to four workers. CPU affinity, cgroup v2 CPU quotas and parent limits, cgroup memory headroom, and OS available memory cap the budget. Common cgroup diff --git a/docs/development/IMPORT_REVIEW_RECOVERY.md b/docs/development/IMPORT_REVIEW_RECOVERY.md index 8d76350a..77943e6d 100644 --- a/docs/development/IMPORT_REVIEW_RECOVERY.md +++ b/docs/development/IMPORT_REVIEW_RECOVERY.md @@ -118,6 +118,10 @@ candidates; membership in one catalog plus agreeing file evidence is required before staging a target. Completed catalog checks are checkpointed. A provider failure pauses the pass, and Resume continues without repeating completed checks. Network requests do not hold a database write transaction. +Catalog summaries are converted to JSON-safe checkpoint payloads before they +are stored. Live provider progress does not rewrite the full durable recovery +snapshot; each completed catalog produces one durable checkpoint, so a worker +restart resumes after the last completed catalog without replaying it. Recovered files run through normal Step 4 safety, current-source validation, ownership checks, and the original copy or keep-in-place settings. Only newly @@ -143,6 +147,10 @@ A completed source recheck reports a file ready only after both archive safety and saved target identity checks pass. Missing, empty, or otherwise blocked sources are counted as blocked even when the archive-level inspection itself completed successfully. +Completed-import source rechecks inspect one bounded page before writing its +refreshed evidence, then commit that page before reading more archives. This +keeps slow archive I/O outside SQLite's single-writer window, bounds memory, and +leaves completed pages durable if a later source needs another attempt. Recovery queries must not expand an entire library into SQL bind parameters. Mixed-folder lookups join existing references and discard exact same-title diff --git a/src/pullbox/api/v1/import_completed_cleanup.py b/src/pullbox/api/v1/import_completed_cleanup.py index 42d0645c..0be8755a 100644 --- a/src/pullbox/api/v1/import_completed_cleanup.py +++ b/src/pullbox/api/v1/import_completed_cleanup.py @@ -117,6 +117,9 @@ async def apply_completed_import_cleanup_route( ) if action is CompletedImportCleanupAction.RETRY_SOURCE_INSPECTION: + # Persist the signed cleanup transition before slow archive inspection. + # The recheck then starts without an existing SQLite writer lock. + await session.commit() service = await build_import_service(session) _job, retrying_count = await service.retry_failed_series( session, diff --git a/src/pullbox/services/diagnostic_db_snapshot.py b/src/pullbox/services/diagnostic_db_snapshot.py index 4ab3151d..cf95bdc9 100644 --- a/src/pullbox/services/diagnostic_db_snapshot.py +++ b/src/pullbox/services/diagnostic_db_snapshot.py @@ -29,6 +29,18 @@ def create_sanitized_db_copy(db_path: Path) -> bytes | None: src.backup(dst) src.close() + tables = { + str(row[0]) + for row in dst.execute( + "SELECT name FROM sqlite_schema WHERE type='table'" + ).fetchall() + } + if "audit_logs" in tables: + # Keep operational evidence while removing its private account link. + dst.execute("UPDATE audit_logs SET user_id = NULL") + if "issue_reader_states" in tables: + # Per-user reading state has no meaning once users are removed. + dst.execute("DELETE FROM issue_reader_states") dst.execute("DELETE FROM users") dst.execute("DELETE FROM api_keys") diff --git a/src/pullbox/services/import_deferred_recovery_execution.py b/src/pullbox/services/import_deferred_recovery_execution.py index e762fb2b..f5137dbc 100644 --- a/src/pullbox/services/import_deferred_recovery_execution.py +++ b/src/pullbox/services/import_deferred_recovery_execution.py @@ -35,13 +35,18 @@ refresh_recovered_groups, ) from pullbox.services.import_source_metadata import source_metadata_for_import_file -from pullbox.services.import_workflow_state import emit_progress, raise_if_job_cancelled +from pullbox.services.import_workflow_state import ( + emit_live_progress, + emit_progress, + raise_if_job_cancelled, +) if TYPE_CHECKING: from collections.abc import Awaitable, Callable from sqlalchemy.ext.asyncio import AsyncSession + from pullbox.providers.base import IssueSummary from pullbox.services.metadata_service import MetadataService @@ -53,6 +58,15 @@ def save_recovery_state(job: ImportJob, state: dict[str, Any]) -> None: job.progress_snapshot = {**dict(job.progress_snapshot or {}), "deferred_recovery": state} +def _catalog_summary_payload(summary: IssueSummary) -> dict[str, Any]: + """Return a durable JSON-safe catalog checkpoint payload.""" + payload = asdict(summary) + source_cutoff_at = payload.get("source_cutoff_at") + if isinstance(source_cutoff_at, datetime): + payload["source_cutoff_at"] = source_cutoff_at.isoformat() + return payload + + async def cancel_deferred_preparation(session: AsyncSession, job: ImportJob) -> bool: """Stop this recovery pass, retaining the original and any completed imports.""" state = recovery_state(job) @@ -264,25 +278,49 @@ async def prepare_deferred_recovery( if job.status is not ImportJobStatus.IMPORTING: raise ValidationError("Deferred recovery must run inside the import worker.") - async def report(current: int, total: int, message: str) -> None: - await raise_if_job_cancelled(session, job_id) - await emit_progress( - session, + revision_state = {"value": int(job.progress_revision or 0)} + + def progress_event(current: int, total: int, message: str) -> ImportProgressEvent: + return ImportProgressEvent( + job_id=job_id, + status=ImportJobStatus.IMPORTING, + mode="import", + phase="deferred_recovery", + progress=round(15 * current / max(total, 1)), + message=message, + current_file_stage="deferred_recovery", + current_file_progress_current=current, + current_file_progress_total=total, + current_file_progress_pct=round(100 * current / max(total, 1)), + current_file_progress_unit="catalogs", + ) + + async def report( + current: int, + total: int, + message: str, + *, + durable: bool = True, + check_control: bool = True, + ) -> None: + if check_control: + await raise_if_job_cancelled(session, job_id) + event = progress_event(current, total, message) + if durable: + event.progress_revision = revision_state["value"] + 1 + await emit_progress(session, job, event, progress_callback) + revision_state["value"] = event.progress_revision + return + # A read-only control check still opens a transaction. Close it before + # provider I/O, then publish live progress without rewriting the large + # durable recovery checkpoint a second time. + await session.commit() + await emit_live_progress( job, - ImportProgressEvent( - job_id=job_id, - status=ImportJobStatus.IMPORTING, - mode="import", - phase="deferred_recovery", - progress=round(15 * current / max(total, 1)), - message=message, - current_file_stage="deferred_recovery", - current_file_progress_current=current, - current_file_progress_total=total, - current_file_progress_pct=round(100 * current / max(total, 1)), - current_file_progress_unit="catalogs", - ), - progress_callback, + event, + progress_callback=progress_callback, + revision_state=revision_state, + started_at=job.import_started_at, ) if state.get("state") == "queued": @@ -308,8 +346,9 @@ async def report(current: int, total: int, message: str) -> None: len(completed), len(candidates), f"Checking series catalog {len(completed) + 1} of {len(candidates)}...", + durable=False, ) - # The progress commit releases the writer lock before provider I/O. + # The live progress report closes its read transaction before provider I/O. cv_id = int(cv_id_text) try: series = await metadata_service.get_series_metadata(cv_id) @@ -335,13 +374,18 @@ async def report(current: int, total: int, message: str) -> None: "cv_id": cv_id, "title": series.title, "year": series.year_start, - "summary": asdict(summary), + "summary": _catalog_summary_payload(summary), } ) completed.add(cv_id_text) state["completed"] = sorted(completed) save_recovery_state(job, state) - await session.commit() + await report( + len(completed), + len(candidates), + f"Checked series catalog {len(completed)} of {len(candidates)}.", + check_control=False, + ) await report(len(completed), max(len(candidates), 1), "Preparing verified files for import...") catalog_count = await _prepare_catalog_targets(session, job) diff --git a/src/pullbox/services/import_review_recheck.py b/src/pullbox/services/import_review_recheck.py index ea598a62..cbdeec77 100644 --- a/src/pullbox/services/import_review_recheck.py +++ b/src/pullbox/services/import_review_recheck.py @@ -547,6 +547,15 @@ async def prepare_completed_import_file_recheck( rows = list((await session.execute(query)).all()) if not rows: break + inspected: list[ + tuple[ + ImportedFile, + ImportedSeries, + SourceMetadata, + dict[str, Any], + dict[str, int | str], + ] + ] = [] for imported_file, imported_series in rows: cursor = int(imported_file.id) metadata, content, signature = await asyncio.to_thread( @@ -562,13 +571,7 @@ async def prepare_completed_import_file_recheck( ) report["files_checked"] += 1 if apply: - ready_for_retry = _apply_completed_file_recheck( - imported_file, - metadata, - content, - signature, - reviewed_series_cv_id=imported_series.cv_id, - ) + inspected.append((imported_file, imported_series, metadata, content, signature)) else: source = {**metadata.diagnostics, **content} ready_for_retry = ( @@ -580,11 +583,25 @@ async def prepare_completed_import_file_recheck( ) is None ) - blocked = not ready_for_retry - report["blocked_files"] += int(blocked) - report["files_prepared"] += int(not blocked) + blocked = not ready_for_retry + report["blocked_files"] += int(blocked) + report["files_prepared"] += int(not blocked) + if apply: - await session.flush() + # Inspect the complete bounded page before taking SQLite's writer + # lock, then commit before reading the next page of source files. + for imported_file, imported_series, metadata, content, signature in inspected: + ready_for_retry = _apply_completed_file_recheck( + imported_file, + metadata, + content, + signature, + reviewed_series_cv_id=imported_series.cv_id, + ) + blocked = not ready_for_retry + report["blocked_files"] += int(blocked) + report["files_prepared"] += int(not blocked) + await session.commit() if apply and report["files_checked"]: session.add( diff --git a/tests/api/test_import_completed_cleanup_api.py b/tests/api/test_import_completed_cleanup_api.py index f52c93bf..3e9170ff 100644 --- a/tests/api/test_import_completed_cleanup_api.py +++ b/tests/api/test_import_completed_cleanup_api.py @@ -102,6 +102,95 @@ async def test_deferred_file_recheck_api_queues_only_after_signed_confirmation( ) +@pytest.mark.asyncio +async def test_retry_source_inspection_commits_cleanup_before_recheck( + db_session: AsyncSession, + monkeypatch: pytest.MonkeyPatch, +) -> None: + from starlette.requests import Request + + from pullbox.api.v1.import_completed_cleanup import apply_completed_import_cleanup_route + from pullbox.models.user import User + from pullbox.schemas.import_completed_cleanup import CompletedImportCleanupApplyRequest + from pullbox.services.import_completed_cleanup import ( + CompletedImportCleanupAction, + preview_completed_import_cleanup, + ) + + user = User(id=42, username="recovery-test", password_hash="unused") + job = ImportJob( + source_path="/imports/mylar.db", + source_type=ImportSourceType.MYLAR3, + status=ImportJobStatus.COMPLETED, + ) + db_session.add_all([user, job]) + await db_session.flush() + imported_series = ImportedSeries( + import_job_id=job.id, + raw_series_name="Temporarily unavailable", + status=ImportSeriesStatus.IMPORTED, + ) + db_session.add(imported_series) + await db_session.flush() + block = build_import_safety_diagnostics( + ImportSafetyCategory.ARCHIVE_INSPECTION_FAILED.value, + code=ImportSafetyCategory.ARCHIVE_INSPECTION_FAILED.value, + ) + db_session.add( + ImportedFile( + import_job_id=job.id, + import_series_id=imported_series.id, + file_path="/comics/temporarily-unavailable.cbz", + file_name="temporarily-unavailable.cbz", + file_size=1024, + file_format="cbz", + status=ImportedFileStatus.SAFETY_BLOCKED, + diagnostics={"safety_block": block}, + ) + ) + await db_session.commit() + job_id = int(job.id) + action = CompletedImportCleanupAction.RETRY_SOURCE_INSPECTION + preview = await preview_completed_import_cleanup(db_session, job_id, action, actor_id=user.id) + + transaction_states: list[bool] = [] + + class RetryService: + async def retry_failed_series( + self, + session: AsyncSession, + requested_job_id: int, + *, + file_ids: list[int] | None = None, + ) -> tuple[ImportJob, int]: + transaction_states.append(session.in_transaction()) + job = await session.get(ImportJob, requested_job_id) + assert job is not None + return job, 0 + + async def build_retry_service(_session: AsyncSession) -> RetryService: + return RetryService() + + monkeypatch.setattr( + "pullbox.composition.services.build_import_service", + build_retry_service, + ) + response = await apply_completed_import_cleanup_route( + job_id, + action, + CompletedImportCleanupApplyRequest( + preview_token=preview.preview_token, + confirmation="APPLY CLEANUP", + ), + Request({"type": "http", "client": ("127.0.0.1", 12345), "headers": []}), + user, + db_session, + ) + + assert response.requires_import_retry is False + assert transaction_states == [False] + + @pytest.mark.asyncio async def test_known_series_recovery_api_previews_then_queues_background_import( authenticated_client: AsyncClient, diff --git a/tests/unit/test_diagnostic_db_snapshot.py b/tests/unit/test_diagnostic_db_snapshot.py index 17cc0926..2c028bce 100644 --- a/tests/unit/test_diagnostic_db_snapshot.py +++ b/tests/unit/test_diagnostic_db_snapshot.py @@ -14,6 +14,17 @@ def _create_snapshot_source(path) -> None: # type: ignore[no-untyped-def] """ CREATE TABLE users (id INTEGER PRIMARY KEY, username TEXT); CREATE TABLE api_keys (id INTEGER PRIMARY KEY, token TEXT); + CREATE TABLE issues (id INTEGER PRIMARY KEY); + CREATE TABLE audit_logs ( + id INTEGER PRIMARY KEY, + user_id INTEGER REFERENCES users(id) ON DELETE SET NULL, + detail TEXT + ); + CREATE TABLE issue_reader_states ( + id INTEGER PRIMARY KEY, + user_id INTEGER NOT NULL REFERENCES users(id) ON DELETE CASCADE, + issue_id INTEGER NOT NULL REFERENCES issues(id) ON DELETE CASCADE + ); CREATE TABLE system_config (key TEXT PRIMARY KEY, value TEXT); CREATE TABLE download_client_configs ( id INTEGER PRIMARY KEY, @@ -28,6 +39,9 @@ def _create_snapshot_source(path) -> None: # type: ignore[no-untyped-def] ) conn.execute("INSERT INTO users (username) VALUES ('admin')") conn.execute("INSERT INTO api_keys (token) VALUES ('tok-secret')") + conn.execute("INSERT INTO issues (id) VALUES (10)") + conn.execute("INSERT INTO audit_logs (user_id, detail) VALUES (1, 'preserve me')") + conn.execute("INSERT INTO issue_reader_states (user_id, issue_id) VALUES (1, 10)") conn.execute( "INSERT INTO system_config (key, value) VALUES (?, ?)", ("comicvine_api_key", "cv-secret"), @@ -62,6 +76,12 @@ def test_create_sanitized_db_copy_removes_auth_rows_and_redacts_secrets(tmp_path try: assert conn.execute("SELECT count(*) FROM users").fetchone()[0] == 0 assert conn.execute("SELECT count(*) FROM api_keys").fetchone()[0] == 0 + assert conn.execute("SELECT user_id, detail FROM audit_logs").fetchone() == ( + None, + "preserve me", + ) + assert conn.execute("SELECT count(*) FROM issue_reader_states").fetchone()[0] == 0 + assert conn.execute("PRAGMA foreign_key_check").fetchall() == [] assert ( conn.execute( "SELECT value FROM system_config WHERE key = 'comicvine_api_key'" diff --git a/tests/unit/test_import_deferred_recovery_execution.py b/tests/unit/test_import_deferred_recovery_execution.py index cb636106..64dba9fd 100644 --- a/tests/unit/test_import_deferred_recovery_execution.py +++ b/tests/unit/test_import_deferred_recovery_execution.py @@ -1,8 +1,10 @@ """Recovery stages work without changing source files or reviving unrelated decisions.""" +from datetime import UTC, datetime from unittest.mock import AsyncMock import pytest +from sqlalchemy import event from pullbox.core.exceptions import JobPausedError, ProviderError from pullbox.models.import_job import ( @@ -12,6 +14,7 @@ ImportSeriesStatus, ) from pullbox.providers.base import IssueSummary, SeriesMetadata +from pullbox.services.catalog.reader import CatalogIssueSummary from pullbox.services.import_deferred_recovery import ( apply_deferred_recovery, plan_deferred_recovery, @@ -147,6 +150,105 @@ async def test_background_recovery_fetches_each_candidate_catalog_once_and_resum assert provider.get_issue_summaries_for_series.await_count == 1 +async def test_catalog_checkpoint_serializes_local_catalog_cutoff(db_session): + job, item, _, _, _ = await seed(db_session) + item.series_id = None + item.cv_id = None + item.status = ImportSeriesStatus.NO_MATCH + await add_file( + db_session, + job, + item, + comicvine_issue_id=7001, + diagnostics={ + "comicvine_series_id": 700, + "metadata_signals": {"comicvine_series_id": "mylar3"}, + }, + ) + job.status = ImportJobStatus.IMPORTING + job.progress_snapshot = {"deferred_recovery": {"state": "queued"}} + provider = AsyncMock() + provider.get_series_metadata.return_value = SeriesMetadata( + provider_id="700", + title="Batman", + sort_title="batman", + year_start=2016, + year_end=None, + status=None, + publisher=None, + description=None, + cover_url=None, + issue_count=1, + comicvine_url=None, + ) + provider.get_issue_summaries_for_series.return_value = [ + CatalogIssueSummary( + provider_id="7001", + issue_number=104, + issue_number_text="104", + title=None, + release_date="2021-01-01", + cover_url=None, + issue_type="issue", + source_cutoff_at=datetime(2026, 9, 13, 5, tzinfo=UTC), + ) + ] + + assert await prepare_deferred_recovery(db_session, job.id, metadata_service=provider) + + stored = job.progress_snapshot["deferred_recovery"]["matches"]["7001"][0]["summary"] + assert stored["source_cutoff_at"] == "2026-09-13T05:00:00+00:00" + await db_session.commit() + + +async def test_catalog_recovery_does_not_write_checkpoint_before_provider_io(db_session): + job, _item, _, _, _ = await seed(db_session) + job.status = ImportJobStatus.IMPORTING + job.progress_snapshot = { + "deferred_recovery": { + "state": "catalogs", + "candidates": {"700": [7001]}, + "completed": [], + "matches": {}, + "series_ids": [], + } + } + await db_session.commit() + + updates: list[str] = [] + engine = db_session.bind.sync_engine + + def record_statement(_conn, _cursor, statement, _parameters, _context, _executemany): + if statement.lstrip().upper().startswith("UPDATE IMPORT_JOBS"): + updates.append(statement) + + event.listen(engine, "before_cursor_execute", record_statement) + provider = AsyncMock() + + async def get_series(_cv_id): + assert updates == [] + return SeriesMetadata( + provider_id="700", + title="Batman", + sort_title="batman", + year_start=2016, + year_end=None, + status=None, + publisher=None, + description=None, + cover_url=None, + issue_count=0, + comicvine_url=None, + ) + + provider.get_series_metadata.side_effect = get_series + provider.get_issue_summaries_for_series.return_value = [] + try: + assert await prepare_deferred_recovery(db_session, job.id, metadata_service=provider) + finally: + event.remove(engine, "before_cursor_execute", record_statement) + + async def test_catalog_target_respects_strong_archive_number_evidence(db_session): from pullbox.services.import_deferred_recovery_execution import _catalog_target_agrees diff --git a/tests/unit/test_import_review_recheck.py b/tests/unit/test_import_review_recheck.py index 6eb5073a..24316c7c 100644 --- a/tests/unit/test_import_review_recheck.py +++ b/tests/unit/test_import_review_recheck.py @@ -2,6 +2,7 @@ from __future__ import annotations +import gc import zipfile from pathlib import Path from unittest.mock import AsyncMock @@ -290,6 +291,149 @@ async def test_completed_recheck_repairs_only_changed_failed_sources( assert untouched.diagnostics == before_untouched["diagnostics"] +async def test_completed_recheck_finishes_inspection_before_mutating_rows( + db_session, + tmp_path, + monkeypatch, +): + from pullbox.services import import_review_recheck + + job, item, files = await _fixture(db_session, tmp_path, ImportSourceType.MYLAR3) + job.status = ImportJobStatus.COMPLETED + item.status = ImportSeriesStatus.IMPORTED + changed_files = files[1:] + for changed in changed_files: + changed.status = ImportedFileStatus.FAILED + changed.include_in_import = False + changed.matched_issue_cv_id = 100000 + int(changed.parsed_issue_number or 0) + changed.match_method = "comicvine_issue_id" + changed.match_confidence = "high" + changed.diagnostics = { + **changed.diagnostics, + "target_issue_summary": { + "provider_id": str(changed.matched_issue_cv_id), + "issue_number": changed.parsed_issue_number, + }, + "source_revalidation": {"code": "source_changed", "retryable": True}, + } + path = Path(changed.file_path) + with zipfile.ZipFile(path, "w") as archive: + archive.writestr( + "ComicInfo.xml", + "Firefly" + f"{int(changed.parsed_issue_number or 0)}", + ) + archive.writestr("1.jpg", b"replacement image") + archive.writestr("2.jpg", b"replacement image") + await db_session.flush() + + original_inspect = import_review_recheck.inspect_review_source + inspected = 0 + + def inspect_before_write(*args, **kwargs): # type: ignore[no-untyped-def] + nonlocal inspected + if inspected: + assert "source_recheck" not in changed_files[0].diagnostics + result = original_inspect(*args, **kwargs) + inspected += 1 + return result + + monkeypatch.setattr(import_review_recheck, "inspect_review_source", inspect_before_write) + + report = await prepare_completed_import_file_recheck( + db_session, + job.id, + source_roots=[tmp_path], + apply=True, + accept_replaced_files=True, + ) + + assert report["files_checked"] == 2 + assert report["files_prepared"] == 2 + assert all( + changed.diagnostics["source_recheck"]["ready_for_retry"] for changed in changed_files + ) + + +async def test_completed_recheck_keeps_loaded_file_rows_bounded( + db_session, + tmp_path, + monkeypatch, +): + from pullbox.services import import_review_recheck + + job, item, files = await _fixture(db_session, tmp_path, ImportSourceType.MYLAR3) + job.status = ImportJobStatus.COMPLETED + item.status = ImportSeriesStatus.IMPORTED + for index in range(248): + path = tmp_path / f"retry-{index:03d}.cbz" + files.append( + ImportedFile( + import_job_id=job.id, + import_series_id=item.id, + file_path=str(path), + file_name=path.name, + file_size=1024, + file_format="cbz", + status=ImportedFileStatus.FAILED, + parsed_series="Firefly", + parsed_issue_number=float(index + 10), + matched_issue_cv_id=200000 + index, + diagnostics={ + "target_issue_summary": { + "provider_id": str(200000 + index), + "issue_number": float(index + 10), + }, + "source_revalidation": {"code": "source_changed", "retryable": True}, + }, + ) + ) + for imported_file in files[:3]: + imported_file.status = ImportedFileStatus.FAILED + imported_file.matched_issue_cv_id = 100000 + int(imported_file.parsed_issue_number or 0) + imported_file.diagnostics = { + **imported_file.diagnostics, + "target_issue_summary": { + "provider_id": str(imported_file.matched_issue_cv_id), + "issue_number": imported_file.parsed_issue_number, + }, + "source_revalidation": {"code": "source_changed", "retryable": True}, + } + db_session.add_all(files[3:]) + await db_session.commit() + job_id = int(job.id) + db_session.expunge_all() + + loaded_file_counts: list[int] = [] + + async def tracked_to_thread(function, *args, **kwargs): # type: ignore[no-untyped-def] + gc.collect() + loaded_file_counts.append( + sum(isinstance(value, ImportedFile) for value in db_session.identity_map.values()) + ) + return function(*args, **kwargs) + + def inspect_without_io(path, base, _signature, **_kwargs): # type: ignore[no-untyped-def] + return base, {}, {"size": 1024, "mtime_ns": 1} + + monkeypatch.setattr(import_review_recheck.asyncio, "to_thread", tracked_to_thread) + monkeypatch.setattr(import_review_recheck, "inspect_review_source", inspect_without_io) + monkeypatch.setattr( + import_review_recheck, "_apply_completed_file_recheck", lambda *a, **k: True + ) + + report = await prepare_completed_import_file_recheck( + db_session, + job_id, + source_roots=[tmp_path], + apply=True, + accept_replaced_files=True, + ) + + assert report["files_checked"] == 251 + assert max(loaded_file_counts) <= 250 + + async def test_completed_recheck_keeps_missing_source_blocked(db_session, tmp_path): job, item, files = await _fixture(db_session, tmp_path, ImportSourceType.MYLAR3) job.status = ImportJobStatus.COMPLETED From 9e6ae68abcb7e4f4dae1c548ba95849483068a20 Mon Sep 17 00:00:00 2001 From: Adam Hernandez Date: Mon, 14 Sep 2026 15:54:31 -0700 Subject: [PATCH 2/2] fix(import): revalidate recheck rows before apply --- src/pullbox/services/import_review_recheck.py | 46 +++++++++++-- tests/unit/test_import_review_recheck.py | 64 +++++++++++++++++++ 2 files changed, 105 insertions(+), 5 deletions(-) diff --git a/src/pullbox/services/import_review_recheck.py b/src/pullbox/services/import_review_recheck.py index cbdeec77..29f2d86e 100644 --- a/src/pullbox/services/import_review_recheck.py +++ b/src/pullbox/services/import_review_recheck.py @@ -549,8 +549,7 @@ async def prepare_completed_import_file_recheck( break inspected: list[ tuple[ - ImportedFile, - ImportedSeries, + int, SourceMetadata, dict[str, Any], dict[str, int | str], @@ -571,7 +570,7 @@ async def prepare_completed_import_file_recheck( ) report["files_checked"] += 1 if apply: - inspected.append((imported_file, imported_series, metadata, content, signature)) + inspected.append((int(imported_file.id), metadata, content, signature)) else: source = {**metadata.diagnostics, **content} ready_for_retry = ( @@ -589,8 +588,40 @@ async def prepare_completed_import_file_recheck( if apply: # Inspect the complete bounded page before taking SQLite's writer - # lock, then commit before reading the next page of source files. - for imported_file, imported_series, metadata, content, signature in inspected: + # lock. Reload and lock the job and eligible rows so a concurrent + # retry cannot have its newer import evidence overwritten. + current_job = await session.scalar( + select(ImportJob) + .where(ImportJob.id == job_id) + .with_for_update() + .execution_options(populate_existing=True) + ) + if current_job is None: + raise NotFoundError("ImportJob", job_id) + if not allows_terminal_import_recovery(current_job): + raise ValidationError( + "Job changed while failed sources were being inspected; retry the recheck" + ) + + inspected_by_id = { + file_id: (metadata, content, signature) + for file_id, metadata, content, signature in inspected + } + apply_query = ( + select(ImportedFile, ImportedSeries) + .join(ImportedSeries, ImportedSeries.id == ImportedFile.import_series_id) + .where( + *retryable_failed_source_filters(job_id), + ImportedFile.id.in_(inspected_by_id), + ) + .order_by(ImportedFile.id) + .with_for_update() + .execution_options(populate_existing=True) + ) + current_rows = list((await session.execute(apply_query)).all()) + report["skipped_files"] += len(inspected) - len(current_rows) + for imported_file, imported_series in current_rows: + metadata, content, signature = inspected_by_id[int(imported_file.id)] ready_for_retry = _apply_completed_file_recheck( imported_file, metadata, @@ -601,7 +632,12 @@ async def prepare_completed_import_file_recheck( blocked = not ready_for_retry report["blocked_files"] += int(blocked) report["files_prepared"] += int(not blocked) + # Commit before reading the next page of source files. await session.commit() + rows.clear() + current_rows.clear() + inspected.clear() + inspected_by_id.clear() if apply and report["files_checked"]: session.add( diff --git a/tests/unit/test_import_review_recheck.py b/tests/unit/test_import_review_recheck.py index 24316c7c..cc7ca6d9 100644 --- a/tests/unit/test_import_review_recheck.py +++ b/tests/unit/test_import_review_recheck.py @@ -355,6 +355,70 @@ def inspect_before_write(*args, **kwargs): # type: ignore[no-untyped-def] ) +async def test_completed_recheck_skips_file_completed_during_inspection( + db_session, + tmp_path, + monkeypatch, +): + from pullbox.services import import_review_recheck + + job, item, files = await _fixture(db_session, tmp_path, ImportSourceType.MYLAR3) + job.status = ImportJobStatus.COMPLETED + item.status = ImportSeriesStatus.IMPORTED + changed = files[1] + changed.status = ImportedFileStatus.FAILED + changed.include_in_import = False + changed.matched_issue_cv_id = 100008 + changed.diagnostics = { + **changed.diagnostics, + "target_issue_summary": {"provider_id": "100008", "issue_number": 8.0}, + "source_revalidation": {"code": "source_changed", "retryable": True}, + } + path = Path(changed.file_path) + with zipfile.ZipFile(path, "w") as archive: + archive.writestr( + "ComicInfo.xml", + "Firefly8", + ) + archive.writestr("1.jpg", b"replacement image") + archive.writestr("2.jpg", b"replacement image") + await db_session.flush() + + original_to_thread = import_review_recheck.asyncio.to_thread + transitioned = False + + async def inspect_then_complete(function, *args, **kwargs): # type: ignore[no-untyped-def] + nonlocal transitioned + result = await original_to_thread(function, *args, **kwargs) + if not transitioned: + changed.status = ImportedFileStatus.IMPORTED + changed.error_message = None + changed.diagnostics = {"concurrent_import": {"completed": True}} + await db_session.flush() + transitioned = True + return result + + monkeypatch.setattr(import_review_recheck.asyncio, "to_thread", inspect_then_complete) + + report = await prepare_completed_import_file_recheck( + db_session, + job.id, + source_roots=[tmp_path], + apply=True, + accept_replaced_files=True, + ) + + assert report == { + "files_checked": 1, + "files_prepared": 0, + "blocked_files": 0, + "skipped_files": 1, + } + assert changed.status is ImportedFileStatus.IMPORTED + assert changed.error_message is None + assert changed.diagnostics == {"concurrent_import": {"completed": True}} + + async def test_completed_recheck_keeps_loaded_file_rows_bounded( db_session, tmp_path,