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..29f2d86e 100644
--- a/src/pullbox/services/import_review_recheck.py
+++ b/src/pullbox/services/import_review_recheck.py
@@ -547,6 +547,14 @@ async def prepare_completed_import_file_recheck(
rows = list((await session.execute(query)).all())
if not rows:
break
+ inspected: list[
+ tuple[
+ int,
+ 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 +570,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((int(imported_file.id), metadata, content, signature))
else:
source = {**metadata.diagnostics, **content}
ready_for_retry = (
@@ -580,11 +582,62 @@ 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. 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,
+ 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)
+ # 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/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..cc7ca6d9 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,213 @@ 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_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,
+ 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