diff --git a/docs/development/IMPORT_REVIEW_RECOVERY.md b/docs/development/IMPORT_REVIEW_RECOVERY.md index 56f091ce..8d76350a 100644 --- a/docs/development/IMPORT_REVIEW_RECOVERY.md +++ b/docs/development/IMPORT_REVIEW_RECOVERY.md @@ -56,7 +56,8 @@ 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. + This includes references confirmed missing by a later source recheck. It does + not delete a review record or touch Mylar's database. - **Skip one-page archives** excludes one-page image archives while leaving the source files intact. A one-page archive may be cover art, a damaged archive, or an intentional one-page comic, so Pullbox does not delete it automatically. @@ -65,7 +66,8 @@ ineligible. uses a red warning modal and an actor-bound signed preview, and requires configured Trash plus source write permission. Reference-only Mylar roots cannot use it; the non-destructive skip remains available instead. -- **Skip unusable files** excludes empty, unsupported, and page-less files. +- **Skip unusable files** excludes empty, unsupported, and page-less files, + including those confirmed unusable by a later source recheck. - **Allow oversized files once** retries only decompression-size blocks marked overrideable. It does not change the global archive safety policy or approve dangerous archive content. @@ -93,12 +95,54 @@ ineligible. validation and import rules. Successful files and source paths are untouched; unresolved files remain in Follow-up. This is not a blanket repair of stale IDs or a replacement for manual review when trusted evidence disagrees. - -For older jobs, **Retry failed** also corrects the misleading series-level -`No eligible files available for import` outcome when the series already has -successful imported files and no failed or ready files remain. It retains the -prior error in diagnostics, rebuilds the counters, and leaves unresolved file -decisions visible. A status-only correction does not launch another import. +- **Recheck deferred files** checks the remaining unmatched files in a resumable + background pass. The preview counts distinct physical paths; its signed scope + still includes every underlying record. Repeated records are consolidated + only when their path, size, portable timestamp, and available content digests + agree. The retained record links back to every superseded record. Files already + registered at the same path and exact issue are recognized without importing + them again. A different file for an owned issue remains a review decision; + it is never automatically substituted for the owned copy. + +The deferred pass uses complete local catalogs first. Exact issue identity may +correct stale Mylar ownership only when the file's title, issue number, type, +and embedded identity agree with the target. Conflicting embedded IDs remain +blocked. Filename-only recovery additionally requires an exact series title or +alias, one issue target, matching issue type, a publication year within one year +of the target issue date, and no pack, volume, or identity conflict. It does not +relax the ordinary search matcher. + +Missing candidate catalogs are fetched once per candidate series per pass, not +once per file. Only trusted saved Mylar, ComicInfo, and sidecar series IDs are +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. + +Recovered files run through normal Step 4 safety, current-source validation, +ownership checks, and the original copy or keep-in-place settings. Only newly +prepared recovery groups execute, not unrelated ready files or Story Arcs. +Cancellation stops this pass without rolling back the original import or +completed recovery files. Manual choices, skips, ambiguous targets, and safety +blocks remain protected. Empty missing-location groups with no file records +are archived with their evidence retained. These rules are shared by Mylar and +folder imports; no source file is moved, renamed, or deleted by reconciliation. + +For older jobs, **Retry failed** first repairs terminal bookkeeping before it +retries file work. A series with any successfully imported file retains its +imported outcome even when another file failed. A series that failed because it +has no ComicVine identity moves to Follow-up instead of being retried without a +target. Files that previously failed to resolve to a library issue are retried +only when their saved identity resolves to one unambiguous issue in the known +Pullbox series; contradictory provider IDs fail closed. Every other unresolved +target becomes an explicit Follow-up decision. Prior errors remain in +diagnostics, counters are rebuilt, and source files remain unchanged. A +status-only correction does not launch another import. + +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. 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/services/import_completed_cleanup.py b/src/pullbox/services/import_completed_cleanup.py index fffb89ca..ed444d36 100644 --- a/src/pullbox/services/import_completed_cleanup.py +++ b/src/pullbox/services/import_completed_cleanup.py @@ -9,6 +9,7 @@ from hashlib import sha256 from itertools import batched from typing import TYPE_CHECKING, Any, Final +from uuid import uuid4 from itsdangerous import BadSignature, SignatureExpired, URLSafeTimedSerializer from sqlalchemy import and_, case, exists, func, or_, select, update @@ -33,8 +34,10 @@ from pullbox.models.series import Series from pullbox.services.audit_service import AuditService from pullbox.services.import_counters import recompute_file_counters, recompute_series_counters +from pullbox.services.import_deferred_recovery import load_empty_stale_series from pullbox.services.import_known_series_recovery import load_known_series_recovery from pullbox.services.import_review_actions import apply_safety_allow_once_to_file +from pullbox.services.import_review_recheck import retryable_failed_source_filters from pullbox.services.import_safety_diagnostics import ImportSafetyCategory from pullbox.services.import_story_arc_resolution import ( refresh_story_arc_entries_for_import_files, @@ -66,6 +69,7 @@ class CompletedImportCleanupAction(enum.StrEnum): ACCEPT_RECOMMENDED_CONFLICTS = "accept_recommended_conflicts" RESOLVE_MIXED_FOLDER_FILES = "resolve_mixed_folder_files" RECOVER_KNOWN_SERIES = "recover_known_series" + RECHECK_DEFERRED_FILES = "recheck_deferred_files" @dataclass(frozen=True, slots=True) @@ -163,10 +167,6 @@ def _source_revalidation_category_expression() -> Any: return ImportedFile.diagnostics["source_revalidation"]["category"].as_string() -def _source_revalidation_retryable_expression() -> Any: - return ImportedFile.diagnostics["source_revalidation"]["retryable"].as_boolean() - - def _safety_filter(*categories: ImportSafetyCategory) -> Any: return _category_expression().in_([category.value for category in categories]) @@ -245,12 +245,21 @@ def _eligible_conflict_groups(job_id: int) -> Any: def _file_filters(job_id: int, action: CompletedImportCleanupAction) -> tuple[Any, ...]: filters: list[Any] = [ImportedFile.import_job_id == job_id] - if action is CompletedImportCleanupAction.DISMISS_MISSING_REFERENCES: - filters.extend( - [ - ImportedFile.status == ImportedFileStatus.SAFETY_BLOCKED, - _safety_filter(ImportSafetyCategory.SOURCE_MISSING), - ] + if action is CompletedImportCleanupAction.RECHECK_DEFERRED_FILES: + filters.append(ImportedFile.status == ImportedFileStatus.NO_MATCH) + elif action is CompletedImportCleanupAction.DISMISS_MISSING_REFERENCES: + filters.append( + or_( + and_( + ImportedFile.status == ImportedFileStatus.SAFETY_BLOCKED, + _safety_filter(ImportSafetyCategory.SOURCE_MISSING), + ), + and_( + ImportedFile.status == ImportedFileStatus.FAILED, + _source_revalidation_category_expression() + == ImportSafetyCategory.SOURCE_MISSING.value, + ), + ) ) elif action is CompletedImportCleanupAction.SKIP_PROBABLE_COVERS: filters.extend( @@ -260,15 +269,24 @@ def _file_filters(job_id: int, action: CompletedImportCleanupAction) -> tuple[An ] ) elif action is CompletedImportCleanupAction.SKIP_UNUSABLE_FILES: - filters.extend( - [ - ImportedFile.status == ImportedFileStatus.SAFETY_BLOCKED, - _safety_filter( - ImportSafetyCategory.ZERO_BYTE, - ImportSafetyCategory.ARCHIVE_NO_PAGES, - ImportSafetyCategory.UNSUPPORTED_FILE_TYPE, + unusable_categories = ( + ImportSafetyCategory.ZERO_BYTE, + ImportSafetyCategory.ARCHIVE_NO_PAGES, + ImportSafetyCategory.UNSUPPORTED_FILE_TYPE, + ) + filters.append( + or_( + and_( + ImportedFile.status == ImportedFileStatus.SAFETY_BLOCKED, + _safety_filter(*unusable_categories), ), - ] + and_( + ImportedFile.status == ImportedFileStatus.FAILED, + _source_revalidation_category_expression().in_( + [category.value for category in unusable_categories] + ), + ), + ) ) elif action is CompletedImportCleanupAction.ALLOW_OVERSIZED_FILES: filters.extend( @@ -293,11 +311,7 @@ def _file_filters(job_id: int, action: CompletedImportCleanupAction) -> tuple[An ImportedFile.status == ImportedFileStatus.SAFETY_BLOCKED, _category_expression().in_(retryable_categories), ), - and_( - ImportedFile.status == ImportedFileStatus.FAILED, - _source_revalidation_category_expression().in_(retryable_categories), - _source_revalidation_retryable_expression().is_(True), - ), + and_(*retryable_failed_source_filters(job_id)), ) ) elif action is CompletedImportCleanupAction.NORMALIZE_ALREADY_OWNED: @@ -697,6 +711,11 @@ async def _load_snapshot( scope_digest=digest.hexdigest(), ) filters = _file_filters(job_id, action) + stale_series = ( + await load_empty_stale_series(session, job_id) + if action is CompletedImportCleanupAction.RECHECK_DEFERRED_FILES + else [] + ) aggregate = ( await session.execute( select( @@ -708,43 +727,59 @@ async def _load_snapshot( ) ).one() file_count = int(aggregate[0] or 0) - if file_count == 0: + if file_count == 0 and not stale_series: return CompletedImportCleanupSnapshot(0, 0, None, None, None, sha256().hexdigest()) digest = sha256() group_ids: set[int] = set() - result = await session.stream( - select( - ImportedFile.id, - ImportedFile.conflict_group_id, - ImportedFile.updated_at, + if file_count: + result = await session.stream( + select( + ImportedFile.id, + ImportedFile.conflict_group_id, + ImportedFile.updated_at, + ) + .where(*filters) + .order_by(ImportedFile.id) + .execution_options(yield_per=20_000) ) - .where(*filters) - .order_by(ImportedFile.id) - .execution_options(yield_per=20_000) - ) - try: - async for rows in result.partitions(20_000): - for file_id, conflict_group_id, updated_at in rows: - digest_line = ( - f"{int(file_id)}|{int(conflict_group_id or 0)}|{updated_at.isoformat()}\n" - ) - digest.update(digest_line.encode()) - if conflict_group_id is not None: - group_ids.add(int(conflict_group_id)) - finally: - await result.close() + try: + async for rows in result.partitions(20_000): + for file_id, conflict_group_id, updated_at in rows: + digest_line = ( + f"file|{int(file_id)}|{int(conflict_group_id or 0)}|" + f"{updated_at.isoformat()}\n" + ) + digest.update(digest_line.encode()) + if conflict_group_id is not None: + group_ids.add(int(conflict_group_id)) + finally: + await result.close() + for item in stale_series: + digest.update(f"series|{int(item.id)}|{item.updated_at.isoformat()}\n".encode()) affected_count = ( len(group_ids) if action is CompletedImportCleanupAction.ACCEPT_RECOMMENDED_CONFLICTS else file_count ) + if action is CompletedImportCleanupAction.RECHECK_DEFERRED_FILES: + affected_count = len(stale_series) + int( + await session.scalar( + select(func.count(func.distinct(ImportedFile.file_path))).where(*filters) + ) + or 0 + ) + updated_at_values = [ + item.updated_at.isoformat(timespec="microseconds") for item in stale_series + ] + if aggregate[3] is not None: + updated_at_values.append(aggregate[3].isoformat(timespec="microseconds")) return CompletedImportCleanupSnapshot( affected_count=affected_count, affected_file_count=file_count, - min_file_id=int(aggregate[1]), - max_file_id=int(aggregate[2]), - max_updated_at=aggregate[3].isoformat(timespec="microseconds"), + min_file_id=int(aggregate[1]) if aggregate[1] is not None else None, + max_file_id=int(aggregate[2]) if aggregate[2] is not None else None, + max_updated_at=max(updated_at_values, default=None), scope_digest=digest.hexdigest(), ) @@ -765,6 +800,14 @@ async def count_completed_import_cleanup_scope( file_count = int( (await session.scalar(select(func.count(ImportedFile.id)).where(*filters))) or 0 ) + if action is CompletedImportCleanupAction.RECHECK_DEFERRED_FILES: + paths = int( + await session.scalar( + select(func.count(func.distinct(ImportedFile.file_path))).where(*filters) + ) + or 0 + ) + return paths + len(await load_empty_stale_series(session, job_id)), file_count if action is not CompletedImportCleanupAction.ACCEPT_RECOMMENDED_CONFLICTS: return file_count, file_count group_count = int( @@ -818,6 +861,12 @@ async def list_completed_import_cleanup_files( total_pages=total_pages, ) filters = _file_filters(job_id, action) + if action is CompletedImportCleanupAction.RECHECK_DEFERRED_FILES: + filters = ( + ImportedFile.id.in_( + select(func.min(ImportedFile.id)).where(*filters).group_by(ImportedFile.file_path) + ), + ) total = int((await session.scalar(select(func.count(ImportedFile.id)).where(*filters))) or 0) total_pages = max(1, (total + normalized_page_size - 1) // normalized_page_size) normalized_page = min(normalized_page, total_pages) @@ -865,10 +914,17 @@ async def list_completed_import_cleanup_examples( ).all() } return tuple(_safe_example_name(names_by_id[file_id]) for file_id in file_ids) + filters = _file_filters(job_id, action) + if action is CompletedImportCleanupAction.RECHECK_DEFERRED_FILES: + filters = ( + ImportedFile.id.in_( + select(func.min(ImportedFile.id)).where(*filters).group_by(ImportedFile.file_path) + ), + ) names = ( await session.scalars( select(ImportedFile.file_name) - .where(*_file_filters(job_id, action)) + .where(*filters) .order_by(ImportedFile.id) .limit(min(max(1, int(limit)), 10)) ) @@ -1563,9 +1619,33 @@ async def apply_completed_import_cleanup( if current_snapshot != preview_snapshot: raise ValidationError("The cleanup scope changed. Preview the action again.") - if action is CompletedImportCleanupAction.RECOVER_KNOWN_SERIES: - affected_series_ids = await _apply_known_series_recovery(session, job) + if action is CompletedImportCleanupAction.RECHECK_DEFERRED_FILES: + from pullbox.services.import_retry_helpers import require_retained_import_destination + + require_retained_import_destination(job) + stale_series_ids = [int(item.id) for item in await load_empty_stale_series(session, job_id)] + job.progress_snapshot = { + **dict(job.progress_snapshot or {}), + "deferred_recovery": { + "state": "queued", + "run_id": uuid4().hex, + "series_ids": [], + "stale_series_ids": stale_series_ids, + "actor_id": actor_id, + }, + "mode": "import", + "phase": "deferred_recovery", + "progress": 0, + "message": "Queued deferred file recovery...", + } + job.status = ImportJobStatus.IMPORTING + job.error_message = None + affected_series_ids = set() affected_file_ids: tuple[int, ...] = () + requires_import_retry = True + elif action is CompletedImportCleanupAction.RECOVER_KNOWN_SERIES: + affected_series_ids = await _apply_known_series_recovery(session, job) + affected_file_ids = () requires_import_retry = await _prepare_series_for_retry(session, job, affected_series_ids) elif action is CompletedImportCleanupAction.ACCEPT_RECOMMENDED_CONFLICTS: affected_series_ids = await _apply_recommended_conflicts(session, job) @@ -1584,8 +1664,9 @@ async def apply_completed_import_cleanup( session, job, affected_series_ids ) - await recompute_file_counters(session, job, series_ids=sorted(affected_series_ids)) - await recompute_series_counters(session, job) + if action is not CompletedImportCleanupAction.RECHECK_DEFERRED_FILES: + await recompute_file_counters(session, job, series_ids=sorted(affected_series_ids)) + await recompute_series_counters(session, job) result = CompletedImportCleanupResult( job_id=job.id, action=action, @@ -1599,7 +1680,11 @@ async def apply_completed_import_cleanup( ), ) item_unit = ( - "group" if action is CompletedImportCleanupAction.ACCEPT_RECOMMENDED_CONFLICTS else "file" + "group" + if action is CompletedImportCleanupAction.ACCEPT_RECOMMENDED_CONFLICTS + else "follow-up item" + if action is CompletedImportCleanupAction.RECHECK_DEFERRED_FILES + else "file" ) session.add( ImportJobLog( diff --git a/src/pullbox/services/import_deferred_recovery.py b/src/pullbox/services/import_deferred_recovery.py new file mode 100644 index 00000000..3c3fdb77 --- /dev/null +++ b/src/pullbox/services/import_deferred_recovery.py @@ -0,0 +1,675 @@ +"""Reconcile completed import decisions using physical files and exact identities.""" + +from __future__ import annotations + +import re +from collections import defaultdict +from dataclasses import dataclass +from datetime import UTC, datetime +from itertools import batched +from typing import TYPE_CHECKING, Any + +from sqlalchemy import select + +from pullbox.core.name_matcher import NameMatcher +from pullbox.core.release_parser import normalize_issue_number, parse_release_title +from pullbox.core.source_metadata import _extract_issue_id_from_notes, _extract_issue_id_from_web +from pullbox.models.import_job import ( + ImportedFile, + ImportedFileStatus, + ImportedSeries, + ImportJob, + ImportJobLog, + ImportSeriesStatus, +) +from pullbox.models.issue import Issue +from pullbox.models.library import LibraryFile +from pullbox.models.series import IssueCatalogState, Series +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 sqlalchemy.ext.asyncio import AsyncSession + from sqlalchemy.sql import Select + + +def positive_id(value: object) -> int | None: + try: + result = int(str(value)) + except (TypeError, ValueError): + return None + return result if result > 0 else None + + +def provider_ids(file: ImportedFile) -> set[int]: + """Retain disagreements instead of silently preferring one saved identity.""" + summary = file.diagnostics.get("target_issue_summary") + summary = summary if isinstance(summary, dict) else {} + source = file.diagnostics.get("source_metadata") + source = source if isinstance(source, dict) else {} + comicinfo = source.get("comicinfo") + comicinfo = comicinfo if isinstance(comicinfo, dict) else {} + embedded_ids = { + value + for value in ( + _extract_issue_id_from_web(comicinfo.get("web")), + _extract_issue_id_from_notes(comicinfo.get("notes")), + ) + if value is not None + } + if len(embedded_ids) > 1: + return embedded_ids + signals = file.diagnostics.get("metadata_signals") + signals = signals if isinstance(signals, dict) else {} + source_id = file.comicvine_issue_id + if len(embedded_ids) == 1 and signals.get("comicvine_issue_id") == "mylar3": + source_id = next(iter(embedded_ids)) + values = (source_id, file.matched_issue_cv_id, summary.get("provider_id"), *embedded_ids) + return {positive_id(value) or -1 for value in values if value is not None} + + +def _unresolved_identity_conflicts(file: ImportedFile) -> bool: + source = file.diagnostics.get("source_metadata") + source = source if isinstance(source, dict) else {} + conflicts = [ + *(source.get("identity_conflicts") or []), + *(file.diagnostics.get("identity_conflicts") or []), + ] + signals = file.diagnostics.get("metadata_signals") + signals = signals if isinstance(signals, dict) else {} + ids = provider_ids(file) + for conflict in conflicts: + if not isinstance(conflict, dict): + return True + if conflict.get("field") == "comicvine_series_id": + continue + if ( + conflict.get("field") == "comicvine_issue_id" + and signals.get("comicvine_issue_id") == "mylar3" + and len(ids) == 1 + and positive_id(conflict.get("first")) == file.comicvine_issue_id + and positive_id(conflict.get("conflicting")) in ids + ): + continue + return True + return False + + +def protected_file(file: ImportedFile, item: ImportedSeries) -> bool: + diagnostics = dict(file.diagnostics or {}) + return bool( + file.status is not ImportedFileStatus.NO_MATCH + or file.library_file_id is not None + or file.include_in_import + or item.user_selected_cv_id is not None + or (file.match_method or "").startswith(("manual", "orphan_recovery")) + or any( + diagnostics.get(key) + for key in ( + "safety_block", + "safety_exception", + "source_revalidation", + "safety_review", + ) + ) + or _unresolved_identity_conflicts(file) + or ( + file.conflict_group_id is not None + and not ( + file.is_preferred + and file.match_confidence == "high" + and file.matched_issue_cv_id is not None + ) + ) + or len(provider_ids(file)) > 1 + or -1 in provider_ids(file) + ) + + +def same_source(left: ImportedFile, right: ImportedFile | LibraryFile) -> bool: + if left.file_path != right.file_path or left.file_size != right.file_size: + return False + first, second = dict(left.source_signature or {}), dict(right.source_signature or {}) + # Persisted path alone is not proof that a file survived unchanged. + # Device IDs change when container mounts are recreated. Size and mtime are portable. + keys = ("size", "size_bytes", "mtime_ns", "content_digest", "content_digest_algorithm") + common = [key for key in keys if key in first and key in second] + return bool("mtime_ns" in common and all(first[key] == second[key] for key in common)) + + +def _titles(series: Series) -> set[str]: + return {NameMatcher.normalize(name) for name in (series.title, *(series.alternate_names or []))} + + +def _source_agrees( + file: ImportedFile, + item: ImportedSeries, + series: Series, + issue: Issue, + *, + require_issue_number: bool = True, +) -> bool: + metadata = source_metadata_for_import_file(item, file) + if metadata.series_name and NameMatcher.normalize(metadata.series_name) not in _titles(series): + return False + issue_numbers: list[float] = [] + if metadata.issue_number is not None: + issue_numbers.append(metadata.issue_number) + hint = metadata.diagnostics.get("archive_entry_issue_hint") + if isinstance(hint, dict) and hint.get("confidence") == "strong": + hint_number = normalize_issue_number(hint.get("issue_number")) + if hint_number is None: + return False + issue_numbers.append(hint_number) + if (require_issue_number and not issue_numbers) or any( + number != issue.issue_number for number in issue_numbers + ): + return False + # Type evidence from a release or ComicInfo must not turn an Annual into #1. + raw_type = file.diagnostics.get("source_issue_type") + return not raw_type or raw_type == issue.issue_type.value + + +def strict_filename_target( + file: ImportedFile, item: ImportedSeries, series: Series, issues: list[Issue] +) -> Issue | None: + """Reparse only after identifying the series, requiring independent agreement.""" + if series.issue_catalog_state is not IssueCatalogState.COMPLETE: + return None + source = file.diagnostics.get("source_metadata") + if isinstance(source, dict) and source.get("identity_conflicts"): + return None + name = re.sub(r"\(converted\)", "", file.file_name, flags=re.IGNORECASE) + name = re.sub(r"(?<=\D)\.(\d+)\.(?=\s|\()", r" \1 ", name) + parsed = parse_release_title( + name, expected_series=(series.title, *(series.alternate_names or [])) + ) + if ( + parsed is None + or parsed.is_pack + or parsed.volume is not None + or parsed.issue_number is None + or parsed.year is None + or NameMatcher.normalize(parsed.series_name or "") not in _titles(series) + or (file.parsed_series and NameMatcher.normalize(file.parsed_series) not in _titles(series)) + or ( + file.parsed_issue_number is not None and file.parsed_issue_number != parsed.issue_number + ) + ): + return None + candidates = [issue for issue in issues if issue.issue_number == parsed.issue_number] + if len(candidates) != 1: + return None + issue = candidates[0] + if ( + issue.release_date is None + or abs(issue.release_date.year - parsed.year) > 1 + or parsed.issue_type != issue.issue_type + or (file.parsed_year is not None and abs(file.parsed_year - issue.release_date.year) > 1) + or not _source_agrees(file, item, series, issue, require_issue_number=False) + ): + return None + return issue + + +@dataclass(frozen=True) +class DeferredRecoveryPlan: + file_id: int + action: str + canonical_file_id: int | None = None + issue_id: int | None = None + reason: str = "" + + +async def load_deferred_rows( + session: AsyncSession, job_id: int +) -> list[tuple[ImportedFile, ImportedSeries]]: + rows: list[tuple[ImportedFile, ImportedSeries]] = [] + cursor = 0 + while True: + batch = ( + await session.execute( + select(ImportedFile, ImportedSeries) + .join(ImportedSeries, ImportedSeries.id == ImportedFile.import_series_id) + .where( + ImportedFile.import_job_id == job_id, + ImportedFile.status == ImportedFileStatus.NO_MATCH, + ImportedFile.id > cursor, + ) + .order_by(ImportedFile.id) + .limit(500) + ) + ).all() + if not batch: + return rows + rows.extend((file, item) for file, item in batch) + cursor = batch[-1][0].id + + +async def plan_deferred_recovery( + session: AsyncSession, job_id: int, *, running: bool = False +) -> tuple[DeferredRecoveryPlan, ...]: + """Return a read-only plan for deterministic local recovery.""" + job = await session.get(ImportJob, job_id) + if job is None or (not running and not allows_terminal_import_recovery(job)): + return () + rows = await load_deferred_rows(session, job_id) + groups: dict[str, list[tuple[ImportedFile, ImportedSeries]]] = defaultdict(list) + for file, item in rows: + groups[file.file_path].append((file, item)) + cv_ids = set().union(*(provider_ids(file) for file, _ in rows)) if rows else set() + local_series_ids = {item.series_id for _, item in rows if item.series_id is not None} + targets: dict[int, tuple[Issue, Series]] = {} + issues_by_series: dict[int, list[Issue]] = defaultdict(list) + series_by_id: dict[int, Series] = {} + for ids in batched(sorted(local_series_ids), 300): + for issue, series in ( + await session.execute( + select(Issue, Series) + .join(Series, Series.id == Issue.series_id) + .where(Series.id.in_(ids)) + ) + ).all(): + issues_by_series[series.id].append(issue) + series_by_id[series.id] = series + if issue.comicvine_id is not None: + targets[issue.comicvine_id] = issue, series + for ids in batched(sorted(cv_ids - targets.keys()), 300): + for issue, series in ( + await session.execute( + select(Issue, Series) + .join(Series, Series.id == Issue.series_id) + .where(Issue.comicvine_id.in_(ids)) + ) + ).all(): + if issue.comicvine_id is not None: + targets[issue.comicvine_id] = issue, series + + registered: dict[str, LibraryFile] = {} + owned: dict[int, LibraryFile] = {} + for ids in batched( + sorted( + {issue.id for issue, _ in targets.values()} + | {issue.id for issues in issues_by_series.values() for issue in issues} + ), + 300, + ): + for library in await session.scalars( + select(LibraryFile).where(LibraryFile.issue_id.in_(ids)) + ): + if library.issue_id is not None: + owned[library.issue_id] = library + registered[library.file_path] = library + for paths in batched(sorted(groups), 300): + for library in await session.scalars( + select(LibraryFile).where(LibraryFile.file_path.in_(paths)) + ): + registered[library.file_path] = library + + plans: list[DeferredRecoveryPlan] = [] + for path, cohort in groups.items(): + if any(protected_file(file, item) for file, item in cohort): + continue + identities = set().union(*(provider_ids(file) for file, _ in cohort)) + if len(identities) > 1: + continue + # Prefer the row whose parent already owns the target catalog. + cohort.sort(key=lambda pair: (pair[1].series_id is None, pair[0].id)) + file, item = cohort[0] + if any(not same_source(file, other) for other, _ in cohort[1:]): + continue + provider_id = next(iter(identities), None) + target = targets.get(provider_id) if provider_id is not None else None + if target is not None: + issue, series = target + if any(not _source_agrees(other, parent, series, issue) for other, parent in cohort): + continue + if any(other.matched_issue_id not in (None, issue.id) for other, _ in cohort): + continue + elif provider_id is None and item.series_id in series_by_id: + series = series_by_id[item.series_id] + issue = strict_filename_target(file, item, series, issues_by_series[series.id]) + target = (issue, series) if issue is not None else None + + registered_file = registered.get(path) + if registered_file is not None: + if ( + target is None + or registered_file.issue_id != target[0].id + or not same_source(file, registered_file) + ): + continue + plans.append(DeferredRecoveryPlan(file.id, "already_registered", issue_id=target[0].id)) + elif target is not None: + issue = target[0] + plans.append( + DeferredRecoveryPlan( + file.id, + "owned_variant" if issue.id in owned else "exact_target", + issue_id=issue.id, + reason="provider_id" if provider_id else "strict_filename", + ) + ) + for other, _ in cohort[1:]: + plans.append( + DeferredRecoveryPlan(other.id, "duplicate_reference", canonical_file_id=file.id) + ) + + # Two distinct physical files targeting one unowned issue still need a choice. + pending: dict[int, list[int]] = defaultdict(list) + for index, plan in enumerate(plans): + if plan.action == "exact_target" and plan.issue_id is not None: + pending[plan.issue_id].append(index) + ambiguous = {index for indices in pending.values() if len(indices) > 1 for index in indices} + return tuple( + sorted( + (plan for index, plan in enumerate(plans) if index not in ambiguous), + key=lambda plan: plan.file_id, + ) + ) + + +def candidate_series_ids(file: ImportedFile, item: ImportedSeries) -> set[int]: + """Use trusted saved provider identities as candidates, never as proof of ownership.""" + diagnostics = dict(file.diagnostics or {}) + signals = diagnostics.get("metadata_signals") + signals = signals if isinstance(signals, dict) else {} + ids: set[int] = set() + if signals.get("comicvine_series_id") in {"mylar3", "comicinfo", "sidecar", "folder_sidecar"}: + cv_id = positive_id(diagnostics.get("comicvine_series_id")) + if cv_id is not None: + ids.add(cv_id) + source = diagnostics.get("source_metadata") + if isinstance(source, dict): + for conflict in source.get("identity_conflicts") or []: + if isinstance(conflict, dict) and conflict.get("field") == "comicvine_series_id": + ids.update( + value + for raw in (conflict.get("first"), conflict.get("conflicting")) + if (value := positive_id(raw)) is not None + ) + candidate = dict(item.diagnostics or {}).get("selected_candidate") + if isinstance(candidate, dict) and candidate.get("match_method") in { + "mylar3_cv_id", + "comicinfo_cv_id", + "folder_cv_id", + }: + cv_id = positive_id(candidate.get("cv_id")) + if cv_id is not None: + ids.add(cv_id) + retained = diagnostics.get("deferred_recovery_candidates") + if isinstance(retained, list): + ids.update(value for raw in retained if (value := positive_id(raw)) is not None) + return ids + + +def issue_summary(issue: Issue) -> dict[str, Any]: + return { + "provider_id": str(issue.comicvine_id) if issue.comicvine_id else None, + "issue_number": issue.issue_number, + "issue_number_text": issue.issue_number_text or str(issue.issue_number), + "title": issue.title, + "release_date": issue.release_date.isoformat() if issue.release_date else None, + "cover_url": issue.cover_url, + "issue_type": issue.issue_type.value, + } + + +def apply_proven_identity( + file: ImportedFile, + *, + issue_cv_id: int | None, + series_cv_id: int | None, + summary: dict[str, Any], +) -> None: + """Retain superseded evidence while replacing only an already-proven identity.""" + diagnostics = dict(file.diagnostics or {}) + previous = dict(diagnostics) + source = dict(diagnostics.get("source_metadata") or {}) + source.pop("identity_conflicts", None) + diagnostics.pop("identity_conflicts", None) + diagnostics.update( + source_metadata=source, + comicvine_series_id=series_cv_id, + kind="deferred_exact_identity", + target_issue_summary=summary, + deferred_recovery_previous_diagnostics=previous, + ) + file.diagnostics = diagnostics + file.comicvine_issue_id = issue_cv_id + + +async def apply_deferred_recovery( + session: AsyncSession, job: ImportJob, *, running: bool = False +) -> dict[str, int]: + """Apply locally proven decisions; materialization remains ordinary Step 4 work.""" + from pullbox.services.import_story_arc_resolution import ( + refresh_story_arc_entries_for_import_files, + ) + + if not running and not allows_terminal_import_recovery(job): + return {} + plans = await plan_deferred_recovery(session, job.id, running=running) + counts: dict[str, int] = defaultdict(int) + affected: set[int] = set() + targets: dict[int, ImportedSeries] = {} + for batch in batched(plans, 300): + files = { + file.id: file + for file in await session.scalars( + select(ImportedFile).where(ImportedFile.id.in_([plan.file_id for plan in batch])) + ) + } + for plan in batch: + file = files[plan.file_id] + parent = await session.get(ImportedSeries, file.import_series_id) + assert parent is not None + file.diagnostics = { + **file.diagnostics, + "deferred_recovery_candidates": sorted(candidate_series_ids(file, parent)), + } + affected.add(parent.id) + evidence = { + "action": plan.action, + "reason": plan.reason, + "source_import_series_id": parent.id, + "previous_error": file.error_message, + "source_preserved": True, + "resolved_at": datetime.now(UTC).isoformat(), + } + if plan.action == "duplicate_reference": + canonical = await session.get(ImportedFile, plan.canonical_file_id) + assert canonical is not None + canonical_parent = await session.get(ImportedSeries, canonical.import_series_id) + assert canonical_parent is not None + candidates = candidate_series_ids(file, parent) | candidate_series_ids( + canonical, canonical_parent + ) + canonical.diagnostics = { + **canonical.diagnostics, + "deferred_recovery_candidates": sorted(candidates), + } + file.duplicate_of_file_id = canonical.id + file.status = ImportedFileStatus.SKIPPED + evidence["canonical_file_id"] = canonical.id + else: + issue = await session.get(Issue, plan.issue_id) + assert issue is not None + series = await session.get(Series, issue.series_id) + assert series is not None + file.matched_issue_id = issue.id + file.matched_issue_cv_id = issue.comicvine_id + if plan.action == "already_registered": + file.status = ImportedFileStatus.ALREADY_OWNED + elif plan.action == "owned_variant": + file.status = ImportedFileStatus.CONFLICT + file.error_message = "This issue is already owned. Review this alternate file." + else: + target = targets.get(series.id) + if target is None: + # Isolate selected files from unrelated ready rows in the old parent. + target = ImportedSeries( + import_job_id=job.id, + raw_series_name=series.title, + raw_year=series.year_start, + cv_id=series.comicvine_id, + cv_title=series.title, + cv_year=series.year_start, + cv_match_method="deferred_exact_identity", + cv_match_score=1.0, + series_id=series.id, + has_files=True, + status=ImportSeriesStatus.CONFIRMED, + selected_for_import=True, + diagnostics={"kind": "deferred_recovery", "source_preserved": True}, + ) + session.add(target) + await session.flush() + targets[series.id] = target + file.import_series_id = target.id + affected.add(target.id) + file.status = ImportedFileStatus.CONFIRMED + file.include_in_import = True + file.match_method = "completed_import_exact_target" + file.match_confidence = "high" + file.parsed_issue_number = issue.issue_number + file.issue_number_raw = issue.issue_number_text + file.conflict_group_id = None + apply_proven_identity( + file, + issue_cv_id=issue.comicvine_id, + series_cv_id=series.comicvine_id, + summary=issue_summary(issue), + ) + evidence["target_issue_id"] = issue.id + if plan.action != "exact_target": + file.include_in_import = False + if plan.action != "owned_variant": + file.error_message = None + file.diagnostics = {**file.diagnostics, "deferred_recovery": evidence} + counts[plan.action] += 1 + await refresh_story_arc_entries_for_import_files( + session, import_job_id=job.id, import_file_ids=list(files) + ) + await session.flush() + + stale_series_ids: tuple[int, ...] | None = None + if running: + state = dict(dict(job.progress_snapshot or {}).get("deferred_recovery") or {}) + raw_ids = state.get("stale_series_ids") + stale_series_ids = ( + tuple(int(value) for value in raw_ids) + if isinstance(raw_ids, list) and all(isinstance(value, int) for value in raw_ids) + else () + ) + counts["stale_series"] = await archive_empty_stale_series( + session, + job.id, + series_ids=stale_series_ids, + ) + await refresh_recovered_groups(session, job, affected) + snapshot = dict(job.progress_snapshot or {}) + recovery = dict(snapshot.get("deferred_recovery") or {}) + recovery["series_ids"] = sorted( + set(recovery.get("series_ids", [])) | {item.id for item in targets.values()} + ) + job.progress_snapshot = {**snapshot, "deferred_recovery": recovery} + session.add( + ImportJobLog( + import_job_id=job.id, + level="INFO", + event="import_deferred_recovery_local", + message="Reconciled deferred file records.", + data={"counts": dict(counts), "source_preserved": True}, + ) + ) + await session.flush() + return dict(counts) + + +async def refresh_recovered_groups( + session: AsyncSession, job: ImportJob, affected: set[int] +) -> None: + """Remove finished matching groups while keeping their audit records.""" + from pullbox.services.import_counters import recompute_file_counters, recompute_series_counters + + if affected: + await recompute_file_counters(session, job, series_ids=sorted(affected)) + # Groups containing only completed decisions no longer need a matching task. + actionable = { + ImportedFileStatus.NO_MATCH, + ImportedFileStatus.CONFLICT, + ImportedFileStatus.FAILED, + ImportedFileStatus.SAFETY_BLOCKED, + ImportedFileStatus.SAFETY_APPROVED, + ImportedFileStatus.MATCHED, + ImportedFileStatus.CONFIRMED, + ImportedFileStatus.PENDING, + } + for ids in batched(sorted(affected), 300): + open_ids = set( + await session.scalars( + select(ImportedFile.import_series_id) + .where( + ImportedFile.import_series_id.in_(ids), ImportedFile.status.in_(actionable) + ) + .distinct() + ) + ) + for item in await session.scalars( + select(ImportedSeries).where(ImportedSeries.id.in_(ids)) + ): + if item.id not in open_ids and item.status in { + ImportSeriesStatus.NO_MATCH, + ImportSeriesStatus.RECOVERY_PENDING, + }: + item.status = ImportSeriesStatus.SKIPPED + item.selected_for_import = False + item.diagnostics = { + **item.diagnostics, + "follow_up_resolved": "all_files_handled", + } + await recompute_series_counters(session, job) + + +def empty_stale_series_query(job_id: int) -> Select[tuple[ImportedSeries]]: + """Return empty missing-location rows eligible for follow-up archival.""" + return select(ImportedSeries).where( + ImportedSeries.import_job_id == job_id, + ImportedSeries.status.in_( + (ImportSeriesStatus.NO_MATCH, ImportSeriesStatus.RECOVERY_PENDING) + ), + ImportedSeries.user_selected_cv_id.is_(None), + ImportedSeries.diagnostics["reason"].as_string().in_(("path_missing", "source_missing")), + ~select(ImportedFile.id).where(ImportedFile.import_series_id == ImportedSeries.id).exists(), + ) + + +async def load_empty_stale_series( + session: AsyncSession, + job_id: int, +) -> list[ImportedSeries]: + """Load empty stale series deterministically for signed cleanup previews.""" + return list(await session.scalars(empty_stale_series_query(job_id).order_by(ImportedSeries.id))) + + +async def archive_empty_stale_series( + session: AsyncSession, + job_id: int, + *, + series_ids: tuple[int, ...] | None = None, +) -> int: + """Retain empty missing Mylar locations in history rather than active matching.""" + query = empty_stale_series_query(job_id) + if series_ids is not None: + query = query.where(ImportedSeries.id.in_(series_ids)) + items = list(await session.scalars(query.order_by(ImportedSeries.id))) + for item in items: + item.status = ImportSeriesStatus.SKIPPED + item.selected_for_import = False + item.diagnostics = { + **item.diagnostics, + "follow_up_resolved": "stale_empty_reference", + "archived_at": datetime.now(UTC).isoformat(), + } + return len(items) diff --git a/src/pullbox/services/import_deferred_recovery_execution.py b/src/pullbox/services/import_deferred_recovery_execution.py new file mode 100644 index 00000000..e762fb2b --- /dev/null +++ b/src/pullbox/services/import_deferred_recovery_execution.py @@ -0,0 +1,377 @@ +"""Durable background preparation for completed-import file recovery.""" + +from __future__ import annotations + +from collections import Counter, defaultdict +from dataclasses import asdict +from datetime import UTC, datetime +from typing import TYPE_CHECKING, Any + +from sqlalchemy import select + +from pullbox.core.exceptions import JobPausedError, NotFoundError, ProviderError, ValidationError +from pullbox.core.name_matcher import NameMatcher +from pullbox.models.import_job import ( + ImportControlRequest, + ImportedFile, + ImportedFileStatus, + ImportedSeries, + ImportJob, + ImportJobLog, + ImportJobStatus, + ImportSeriesStatus, +) +from pullbox.models.issue import Issue +from pullbox.schemas.import_job import ImportProgressEvent +from pullbox.services.import_counters import recompute_file_counters, recompute_series_counters +from pullbox.services.import_deferred_recovery import ( + apply_deferred_recovery, + apply_proven_identity, + candidate_series_ids, + load_deferred_rows, + positive_id, + protected_file, + provider_ids, + 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 + +if TYPE_CHECKING: + from collections.abc import Awaitable, Callable + + from sqlalchemy.ext.asyncio import AsyncSession + + from pullbox.services.metadata_service import MetadataService + + +def recovery_state(job: ImportJob) -> dict[str, Any]: + return dict(dict(job.progress_snapshot or {}).get("deferred_recovery") or {}) + + +def save_recovery_state(job: ImportJob, state: dict[str, Any]) -> None: + job.progress_snapshot = {**dict(job.progress_snapshot or {}), "deferred_recovery": state} + + +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) + if state.get("state") not in {"queued", "catalogs", "prepared"}: + return False + ids = state.get("series_ids", []) + for item in await session.scalars( + select(ImportedSeries).where( + ImportedSeries.import_job_id == job.id, ImportedSeries.id.in_(ids) + ) + ): + pending = list( + await session.scalars( + select(ImportedFile).where( + ImportedFile.import_series_id == item.id, + ImportedFile.status.in_( + (ImportedFileStatus.CONFIRMED, ImportedFileStatus.MATCHED) + ), + ) + ) + ) + for file in pending: + file.status = ImportedFileStatus.NO_MATCH + file.include_in_import = False + if pending: + item.status = ImportSeriesStatus.RECOVERY_PENDING + item.selected_for_import = False + await recompute_file_counters(session, job, series_ids=ids) + await recompute_series_counters(session, job) + state["state"] = "cancelled" + save_recovery_state(job, state) + job.status = ImportJobStatus.COMPLETED + job.control_request = ImportControlRequest.NONE + job.error_message = None + job.progress_snapshot = { + **job.progress_snapshot, + "status": "completed", + "phase": "done", + "progress": 100, + "message": "Recovery stopped. Completed imports and source files were preserved.", + } + session.add( + ImportJobLog( + import_job_id=job.id, + level="INFO", + event="import_deferred_recovery_cancelled", + message="Recovery stopped without rolling back the original import.", + data={"source_preserved": True}, + ) + ) + await session.commit() + return True + + +async def _catalog_candidates(session: AsyncSession, job: ImportJob) -> dict[str, list[int]]: + candidates: dict[str, set[int]] = defaultdict(set) + rows = await load_deferred_rows(session, job.id) + for file, item in rows: + if protected_file(file, item) or len(provider_ids(file)) != 1: + continue + issue_cv_id = next(iter(provider_ids(file))) + for series_cv_id in candidate_series_ids(file, item): + candidates[str(series_cv_id)].add(issue_cv_id) + # Never fetch catalogs for an issue whose canonical local ownership is already known. + from itertools import batched + + known: set[int] = set() + for ids in batched(sorted(set().union(*candidates.values()) if candidates else set()), 300): + known.update( + value + for value in await session.scalars( + select(Issue.comicvine_id).where(Issue.comicvine_id.in_(ids)) + ) + if value is not None + ) + return {key: sorted(values - known) for key, values in candidates.items() if values - known} + + +def _catalog_target_agrees( + file: ImportedFile, item: ImportedSeries, target: dict[str, Any] +) -> bool: + metadata = source_metadata_for_import_file(item, file) + summary = target["summary"] + source_title = metadata.series_name or "" + hint = metadata.diagnostics.get("archive_entry_issue_hint") + if ( + isinstance(hint, dict) + and hint.get("confidence") == "strong" + and hint.get("issue_number") != summary["issue_number"] + ): + return False + return bool( + source_title + and NameMatcher.normalize(source_title) == NameMatcher.normalize(target["title"]) + and metadata.issue_number is not None + and metadata.issue_number == summary["issue_number"] + and ( + not file.diagnostics.get("source_issue_type") + or file.diagnostics["source_issue_type"] == summary["issue_type"] + ) + and file.matched_issue_id is None + ) + + +async def _prepare_catalog_targets(session: AsyncSession, job: ImportJob) -> int: + """Only stage unique issue membership with agreeing file evidence.""" + state = recovery_state(job) + matches = state.get("matches", {}) + rows = await load_deferred_rows(session, job.id) + eligible: list[tuple[ImportedFile, ImportedSeries, dict[str, Any]]] = [] + for file, item in rows: + ids = provider_ids(file) + if protected_file(file, item) or len(ids) != 1: + continue + options = matches.get(str(next(iter(ids))), []) + # Membership in conflicting candidate catalogs is ambiguous even when titles agree. + if len(options) != 1 or not _catalog_target_agrees(file, item, options[0]): + continue + eligible.append((file, item, options[0])) + counts = Counter(str(target["summary"]["provider_id"]) for _, _, target in eligible) + targets: dict[int, ImportedSeries] = {} + affected: set[int] = set() + ready = 0 + for file, original, target in eligible: + issue_cv_id = positive_id(target["summary"]["provider_id"]) + if issue_cv_id is None or counts[str(issue_cv_id)] != 1: + continue + if await session.scalar(select(Issue.id).where(Issue.comicvine_id == issue_cv_id)): + continue + target_cv_id = int(target["cv_id"]) + target_item = targets.get(target_cv_id) + if target_item is None: + target_item = ImportedSeries( + import_job_id=job.id, + raw_series_name=target["title"], + raw_year=target["year"], + cv_id=target_cv_id, + cv_title=target["title"], + cv_year=target["year"], + cv_match_score=1.0, + cv_match_method="deferred_catalog_identity", + status=ImportSeriesStatus.CONFIRMED, + selected_for_import=True, + has_files=True, + diagnostics={"kind": "deferred_recovery", "source_preserved": True}, + ) + session.add(target_item) + await session.flush() + targets[target_cv_id] = target_item + affected.update((original.id, target_item.id)) + file.import_series_id = target_item.id + file.status = ImportedFileStatus.CONFIRMED + file.matched_issue_cv_id = issue_cv_id + file.include_in_import = True + file.match_confidence = "high" + file.match_method = "completed_import_exact_target" + file.error_message = None + file.conflict_group_id = None + apply_proven_identity( + file, issue_cv_id=issue_cv_id, series_cv_id=target_cv_id, summary=target["summary"] + ) + file.diagnostics = { + **file.diagnostics, + "target_issue_summary": target["summary"], + "deferred_recovery": { + "action": "catalog_identity", + "source_import_series_id": original.id, + "target_series_cv_id": target_cv_id, + "source_preserved": True, + "resolved_at": datetime.now(UTC).isoformat(), + }, + } + ready += 1 + if affected: + from pullbox.services.import_story_arc_resolution import ( + refresh_story_arc_entries_for_import_files, + ) + + await refresh_recovered_groups(session, job, affected) + await refresh_story_arc_entries_for_import_files( + session, + import_job_id=job.id, + import_file_ids=[ + file.id for file, _, _ in eligible if file.status is ImportedFileStatus.CONFIRMED + ], + ) + state["series_ids"] = sorted( + set(state.get("series_ids", [])) | {item.id for item in targets.values()} + ) + save_recovery_state(job, state) + await session.flush() + return ready + + +async def prepare_deferred_recovery( + session: AsyncSession, + job_id: int, + *, + metadata_service: MetadataService, + progress_callback: Callable[[ImportProgressEvent], Awaitable[None]] | None = None, +) -> bool: + """Resume a saved recovery request, preparing exact files for normal execution.""" + job = await session.get(ImportJob, job_id) + if job is None: + raise NotFoundError("ImportJob", job_id) + state = recovery_state(job) + if not state or state.get("state") in {"prepared", "completed", "cancelled"}: + return False + 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, + 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, + ) + + if state.get("state") == "queued": + await report(0, 1, "Reconciling deferred files with the existing library...") + local_counts = await apply_deferred_recovery(session, job, running=True) + state = recovery_state(job) + state.update( + state="catalogs", + local_counts=local_counts, + candidates=await _catalog_candidates(session, job), + completed=[], + matches={}, + ) + save_recovery_state(job, state) + await session.commit() + + candidates = state["candidates"] + completed = set(state.get("completed", [])) + for cv_id_text, needed_ids in sorted(candidates.items()): + if cv_id_text in completed: + continue + await report( + len(completed), + len(candidates), + f"Checking series catalog {len(completed) + 1} of {len(candidates)}...", + ) + # The progress commit releases the writer lock before provider I/O. + cv_id = int(cv_id_text) + try: + series = await metadata_service.get_series_metadata(cv_id) + summaries = await metadata_service.get_issue_summaries_for_series(cv_id) + except NotFoundError: + series = None + summaries = [] + except ProviderError as exc: + job.error_message = ( + "Metadata is temporarily unavailable. Resume recovery when it is available." + ) + await report(len(completed), len(candidates), job.error_message) + raise JobPausedError(job.error_message) from exc + await raise_if_job_cancelled(session, job_id) + if series is not None and positive_id(series.provider_id) == cv_id: + needed = set(needed_ids) + for summary in summaries: + if positive_id(summary.provider_id) not in needed: + continue + entries = state["matches"].setdefault(str(summary.provider_id), []) + entries.append( + { + "cv_id": cv_id, + "title": series.title, + "year": series.year_start, + "summary": asdict(summary), + } + ) + completed.add(cv_id_text) + state["completed"] = sorted(completed) + save_recovery_state(job, state) + await session.commit() + + await report(len(completed), max(len(candidates), 1), "Preparing verified files for import...") + catalog_count = await _prepare_catalog_targets(session, job) + state = recovery_state(job) + state.update(state="prepared", catalog_files_prepared=catalog_count) + job.error_message = None + if not state.get("series_ids"): + state["state"] = "completed" + job.status = ImportJobStatus.COMPLETED + job.progress_snapshot = { + **job.progress_snapshot, + "status": "completed", + "phase": "done", + "progress": 100, + "message": "Recovery completed. Remaining files still need review.", + } + save_recovery_state(job, state) + session.add( + ImportJobLog( + import_job_id=job_id, + level="INFO", + event="import_deferred_recovery_prepared", + message="Deferred recovery preparation completed.", + data={ + "local_counts": state.get("local_counts", {}), + "catalogs_checked": len(completed), + "catalog_files_prepared": catalog_count, + "source_preserved": True, + }, + ) + ) + await session.commit() + return True diff --git a/src/pullbox/services/import_job_controls.py b/src/pullbox/services/import_job_controls.py index 941655dd..fdd87546 100644 --- a/src/pullbox/services/import_job_controls.py +++ b/src/pullbox/services/import_job_controls.py @@ -141,6 +141,8 @@ async def cancel_job( if ( job.status in {ImportJobStatus.PAUSED, ImportJobStatus.STALLED} and job.import_started_at is None + and dict(dict(job.progress_snapshot or {}).get("deferred_recovery") or {}).get("state") + not in {"queued", "catalogs", "prepared"} ): await discard_unpublished_import_story_arc_sync_work(session, (job_id,)) await session.delete(job) @@ -316,13 +318,19 @@ async def request_cancel( if ( job.status in {ImportJobStatus.PAUSED, ImportJobStatus.STALLED} and job.import_started_at is None + and dict(dict(job.progress_snapshot or {}).get("deferred_recovery") or {}).get("state") + not in {"queued", "catalogs", "prepared"} ): await session.delete(job) await session.flush() return job def _apply_cancel(target: ImportJob) -> None: - if _is_story_arc_placement_wait(target) or ( + recovery = dict(dict(target.progress_snapshot or {}).get("deferred_recovery") or {}) + if recovery.get("state") in {"queued", "catalogs", "prepared"}: + target.status = ImportJobStatus.CANCELLING + target.control_request = ImportControlRequest.CANCEL + elif _is_story_arc_placement_wait(target) or ( target.status in {ImportJobStatus.PAUSED, ImportJobStatus.STALLED} and target.import_started_at is not None ): diff --git a/src/pullbox/services/import_job_execution.py b/src/pullbox/services/import_job_execution.py index 238675a0..3eab15b2 100644 --- a/src/pullbox/services/import_job_execution.py +++ b/src/pullbox/services/import_job_execution.py @@ -113,6 +113,7 @@ auto_confirm_trusted_logical_story_arcs, ) from pullbox.services.import_workflow_state import ( + deferred_recovery_scope, emit_live_progress, ) @@ -285,6 +286,12 @@ async def execute_import_job( confirmed_ids={item.id for item in confirmed_items}, ) execution_items = _build_execution_item_plans(confirmed_items, duplicate_items) + recovery_scope = deferred_recovery_scope(job) + if recovery_scope is not None: + authorized_ids = set(recovery_scope) + confirmed_items = [item for item in confirmed_items if item.id in authorized_ids] + duplicate_items = [item for item in duplicate_items if item.id in authorized_ids] + execution_items = _build_execution_item_plans(confirmed_items, duplicate_items) await log_event( session, job_id, @@ -682,6 +689,43 @@ def queue_catalog_hydration(series_id: int, search_on_add: bool) -> None: job.series_failed = failed_count job.total_files_imported = total_files_imported job.total_files_failed = total_files_failed + if recovery_scope is not None: + from pullbox.services.import_counters import ( + recompute_file_counters, + recompute_series_counters, + ) + + await recompute_file_counters(session, job, series_ids=list(recovery_scope)) + await recompute_series_counters(session, job) + snapshot = dict(job.progress_snapshot or {}) + recovery = dict(snapshot.get("deferred_recovery") or {}) + recovery["state"] = "completed" + job.status = ImportJobStatus.COMPLETED + job.progress_snapshot = { + **snapshot, + "deferred_recovery": recovery, + "status": "completed", + "phase": "done", + "progress": 100, + "message": "Deferred file recovery completed. Remaining decisions are in Follow-up.", + } + await log_event( + session, + job_id, + "INFO", + "import_deferred_recovery_completed", + message="Completed the scoped recovery without executing unrelated review groups.", + series_ids=list(recovery_scope), + ) + for request in pending_catalog_hydrations: + _schedule_catalog_hydration( + session, + series_service=series_service, + series_id=request.series_id, + search_on_add=request.search_on_add, + ) + await session.flush() + return job, story_arc_materialization = await _execute_story_arc_materialization( session, job, diff --git a/src/pullbox/services/import_managed_copy_preflight.py b/src/pullbox/services/import_managed_copy_preflight.py index f00b8c06..d7de457f 100644 --- a/src/pullbox/services/import_managed_copy_preflight.py +++ b/src/pullbox/services/import_managed_copy_preflight.py @@ -24,6 +24,7 @@ ImportedFileStatus, ImportedSeries, ImportFileHandlingMode, + ImportJob, ImportJobStatus, ImportSeriesStatus, ) @@ -35,8 +36,6 @@ if TYPE_CHECKING: from sqlalchemy.ext.asyncio import AsyncSession - from pullbox.models.import_job import ImportJob - _ONE_GIB = 1024**3 _CAPACITY_SNAPSHOT_KEY = "managed_copy_capacity" _STORY_ARC_PAGE_SIZE = 250 @@ -166,6 +165,10 @@ async def selected_managed_copy_source_bytes( selected_duplicate_series = ( ImportedSeries.status == ImportSeriesStatus.DUPLICATE ) & ImportedFile.include_in_import.is_(True) + from pullbox.services.import_workflow_state import deferred_recovery_scope + + job = await session.get(ImportJob, job_id) + scope = deferred_recovery_scope(job) if job is not None else None total = await session.scalar( sa_select(sa_func.coalesce(sa_func.sum(ImportedFile.file_size), 0)) .join(ImportedSeries, ImportedFile.import_series_id == ImportedSeries.id) @@ -173,6 +176,7 @@ async def selected_managed_copy_source_bytes( ImportedFile.import_job_id == job_id, ImportedFile.status.in_([ImportedFileStatus.MATCHED, ImportedFileStatus.CONFIRMED]), sa_or(selected_new_series, selected_duplicate_series), + *([ImportedSeries.id.in_(scope)] if scope is not None else []), ) ) return max(int(total or 0), 0) @@ -381,6 +385,9 @@ async def estimate_conversion_workspace_source_bytes( selected_duplicate_series = ( ImportedSeries.status == ImportSeriesStatus.DUPLICATE ) & ImportedFile.include_in_import.is_(True) + from pullbox.services.import_workflow_state import deferred_recovery_scope + + scope = deferred_recovery_scope(job) sizes = list( ( await session.scalars( @@ -392,6 +399,7 @@ async def estimate_conversion_workspace_source_bytes( [ImportedFileStatus.MATCHED, ImportedFileStatus.CONFIRMED] ), sa_or(selected_new_series, selected_duplicate_series), + *([ImportedSeries.id.in_(scope)] if scope is not None else []), sa_func.lower(ImportedFile.file_format) != "cbz", ) .order_by(ImportedFile.file_size.desc(), ImportedFile.id.asc()) @@ -418,7 +426,13 @@ async def validate_managed_copy_preflight( job_root = await _resolve_job_managed_root(session, job) selected_bytes_by_root[job_root.id] = selected_source_bytes - story_arc_bytes = await selected_story_arc_copy_source_bytes_by_root(session, job.id) + from pullbox.services.import_workflow_state import deferred_recovery_scope + + story_arc_bytes = ( + await selected_story_arc_copy_source_bytes_by_root(session, job.id) + if deferred_recovery_scope(job) is None + else {} + ) for root_id, selected_source_bytes in story_arc_bytes.items(): selected_bytes_by_root[root_id] = ( selected_bytes_by_root.get(root_id, 0) + selected_source_bytes diff --git a/src/pullbox/services/import_orphan_recovery_context.py b/src/pullbox/services/import_orphan_recovery_context.py index d2f1d27d..2aa39d54 100644 --- a/src/pullbox/services/import_orphan_recovery_context.py +++ b/src/pullbox/services/import_orphan_recovery_context.py @@ -19,7 +19,7 @@ ) from pullbox.models.issue import Issue from pullbox.models.library import LibraryRoot -from pullbox.services.import_orphans import is_active_orphan_row +from pullbox.services.import_orphans import is_active_orphan_row, requires_orphan_issue_decision if TYPE_CHECKING: from sqlalchemy.ext.asyncio import AsyncSession @@ -186,10 +186,7 @@ def _build_recovery_file_rows( ): suggested_issue_cv_id = issue_number_to_cv_ids[imp_file.parsed_issue_number][0] - decision_locked = imp_file.status in { - ImportedFileStatus.IMPORTED, - ImportedFileStatus.SKIPPED, - } + decision_locked = not requires_orphan_issue_decision(imp_file) if decision_locked: files_completed += 1 else: diff --git a/src/pullbox/services/import_orphans.py b/src/pullbox/services/import_orphans.py index 041c734f..9dce5308 100644 --- a/src/pullbox/services/import_orphans.py +++ b/src/pullbox/services/import_orphans.py @@ -2,6 +2,7 @@ from __future__ import annotations +from datetime import UTC, datetime from typing import TYPE_CHECKING, Any, Protocol from sqlalchemy import and_, or_ @@ -9,6 +10,7 @@ from sqlalchemy import select as sa_select from pullbox.core.exceptions import NotFoundError, ValidationError +from pullbox.core.release_parser import parse_release_title from pullbox.models.import_job import ( ImportedFile, ImportedFileStatus, @@ -17,14 +19,19 @@ ImportJobStatus, ImportSeriesStatus, ) +from pullbox.models.issue import Issue, IssueType from pullbox.models.series import Series from pullbox.services.import_counters import recompute_file_counters, recompute_series_counters +from pullbox.services.import_file_issue_signals import ( + candidate_issue_number, + candidate_issue_number_text, + comicinfo_issue_number, + filename_issue_number, + volume_issue_number, +) from pullbox.services.import_file_resolution import load_issue_lookup_for_series from pullbox.services.import_job_actions import build_series_created_action_payload -from pullbox.services.import_job_execution_items import ( - ensure_target_issue_summary_for_import_file, - retain_imported_series_outcome, -) +from pullbox.services.import_job_execution_items import ensure_target_issue_summary_for_import_file 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 @@ -36,7 +43,6 @@ from sqlalchemy.sql import Select from sqlalchemy.sql.elements import ColumnElement - from pullbox.models.issue import Issue from pullbox.providers.base import SeriesMetadata from pullbox.schemas.import_job import OrphanRecoveryDecision, RecoverOrphanRequest @@ -87,6 +93,231 @@ async def add_from_comicvine( ImportSeriesStatus.NO_MATCH, ImportSeriesStatus.RECOVERY_PENDING, ) +_UNRESOLVED_TARGET_ERROR = "Could not resolve to a library issue" + + +def requires_orphan_issue_decision(file: ImportedFile) -> bool: + """Return whether Follow-up should ask for an issue assignment or skip.""" + return file.status in { + ImportedFileStatus.PENDING, + ImportedFileStatus.MATCHED, + ImportedFileStatus.CONFIRMED, + ImportedFileStatus.CONFLICT, + ImportedFileStatus.NO_MATCH, + } + + +def _saved_target_provider_ids(file: ImportedFile) -> set[int]: + diagnostics = dict(file.diagnostics or {}) + raw_summary = diagnostics.get("target_issue_summary") + summary = raw_summary if isinstance(raw_summary, dict) else {} + raw_values: tuple[object, ...] = ( + file.matched_issue_cv_id, + file.comicvine_issue_id, + summary.get("provider_id"), + ) + provider_ids: set[int] = set() + for raw_value in raw_values: + if raw_value is None: + continue + try: + provider_ids.add(int(str(raw_value))) + except ValueError: + provider_ids.add(-1) + return provider_ids + + +def _source_issue_type_matches(file: ImportedFile, issue: Issue) -> bool: + diagnostics = dict(file.diagnostics or {}) + source_metadata = diagnostics.get("source_metadata") + filename_parse = ( + source_metadata.get("filename_parse") if isinstance(source_metadata, dict) else None + ) + raw_issue_type = diagnostics.get("source_issue_type") or ( + filename_parse.get("issue_type") if isinstance(filename_parse, dict) else None + ) + if raw_issue_type: + try: + return IssueType(str(raw_issue_type)) is issue.issue_type + except ValueError: + return False + parsed = parse_release_title(file.file_name or "") + return parsed is None or parsed.issue_type is issue.issue_type + + +def _failed_target_issue_number(file: ImportedFile) -> float | None: + """Return affirmative issue-number evidence without treating volume as issue.""" + exact_number = candidate_issue_number_text(file) + if exact_number is not None: + return candidate_issue_number(file) + filename_number = filename_issue_number(file) + if filename_number is not None: + return filename_number + if file.parsed_issue_number is not None and volume_issue_number(file) is None: + return file.parsed_issue_number + return comicinfo_issue_number(file) + + +async def _resolve_proven_failed_target( + session: AsyncSession, + item: ImportedSeries, + file: ImportedFile, + *, + cv_id_to_issue: dict[int, Issue], + exact_number_to_issue: dict[str, Issue], + number_to_issue: dict[float, Issue], +) -> Issue | None: + """Resolve only an exact target inside the already identified local series.""" + if item.series_id is None: + return None + if file.matched_issue_id is not None: + saved_issue = await session.get(Issue, file.matched_issue_id) + if saved_issue is not None and saved_issue.series_id == item.series_id: + return saved_issue + return None + + provider_ids = _saved_target_provider_ids(file) + if provider_ids: + if len(provider_ids) != 1: + return None + return cv_id_to_issue.get(next(iter(provider_ids))) + + exact_number = candidate_issue_number_text(file) + if exact_number is not None: + issue = exact_number_to_issue.get(exact_number) + else: + issue_number = _failed_target_issue_number(file) + issue = number_to_issue.get(issue_number) if issue_number is not None else None + return issue if issue is not None and _source_issue_type_matches(file, issue) else None + + +def _prepare_exact_target_retry(file: ImportedFile, issue: Issue) -> None: + diagnostics = dict(file.diagnostics or {}) + diagnostics["previous_import_error"] = file.error_message + diagnostics["completed_import_follow_up"] = { + "resolution": "exact_issue_target", + "resolved_at": datetime.now(UTC).isoformat(), + "source_preserved": True, + } + file.status = ImportedFileStatus.CONFIRMED + file.include_in_import = True + file.matched_issue_id = issue.id + file.matched_issue_cv_id = issue.comicvine_id + file.match_confidence = "high" + file.match_method = "completed_import_exact_target" + file.error_message = None + file.diagnostics = diagnostics + + +def _defer_target_to_follow_up(file: ImportedFile) -> None: + diagnostics = dict(file.diagnostics or {}) + diagnostics["previous_import_error"] = file.error_message + diagnostics["completed_import_follow_up"] = { + "resolution": "issue_decision_required", + "resolved_at": datetime.now(UTC).isoformat(), + "source_preserved": True, + } + file.status = ImportedFileStatus.NO_MATCH + file.include_in_import = False + file.matched_issue_id = None + file.error_message = "Choose the correct issue in Follow-up." + file.diagnostics = diagnostics + + +async def _prepare_terminal_follow_up( + session: AsyncSession, + job: ImportJob, +) -> tuple[set[int], set[int], int]: + """Repair legacy terminal outcomes before retrying only proven work.""" + affected_series_ids: set[int] = set() + retry_series_ids: set[int] = set() + target_rows = list( + ( + await session.execute( + sa_select(ImportedSeries, ImportedFile) + .join(ImportedFile, ImportedFile.import_series_id == ImportedSeries.id) + .where( + ImportedSeries.import_job_id == job.id, + ImportedFile.status == ImportedFileStatus.FAILED, + ImportedFile.error_message == _UNRESOLVED_TARGET_ERROR, + ) + .order_by(ImportedFile.id) + ) + ).all() + ) + issue_lookups: dict[int, tuple[dict[int, Issue], dict[str, Issue], dict[float, Issue]]] = {} + for item, file in target_rows: + if item.series_id is not None and item.series_id not in issue_lookups: + issue_lookups[item.series_id] = await load_issue_lookup_for_series( + session, + item.series_id, + ) + lookup = issue_lookups.get(item.series_id or -1, ({}, {}, {})) + issue = await _resolve_proven_failed_target( + session, + item, + file, + cv_id_to_issue=lookup[0], + exact_number_to_issue=lookup[1], + number_to_issue=lookup[2], + ) + if issue is None: + _defer_target_to_follow_up(file) + else: + _prepare_exact_target_retry(file, issue) + retry_series_ids.add(item.id) + affected_series_ids.add(item.id) + + imported_series_ids = set( + await session.scalars( + sa_select(ImportedFile.import_series_id) + .where( + ImportedFile.import_job_id == job.id, + ImportedFile.status == ImportedFileStatus.IMPORTED, + ) + .distinct() + ) + ) + failed_series = list( + await session.scalars( + sa_select(ImportedSeries).where( + ImportedSeries.import_job_id == job.id, + ImportedSeries.status == ImportSeriesStatus.FAILED, + ) + ) + ) + normalized_series_count = 0 + for item in failed_series: + previous_error = item.error_message + if item.id in imported_series_ids: + item.status = ImportSeriesStatus.IMPORTED + elif item.error_message == "No ComicVine ID available": + item.status = ImportSeriesStatus.NO_MATCH + item.selected_for_import = False + elif item.error_message == "No eligible files available for import" or ( + item.id in affected_series_ids + ): + item.status = ImportSeriesStatus.RECOVERY_PENDING + item.selected_for_import = False + else: + continue + item.error_message = None + item.diagnostics = { + **dict(item.diagnostics or {}), + "previous_series_error": previous_error, + "completed_import_follow_up": { + "resolution": "series_status_repaired", + "resolved_at": datetime.now(UTC).isoformat(), + "source_preserved": True, + }, + } + normalized_series_count += 1 + affected_series_ids.add(item.id) + + if affected_series_ids: + await recompute_file_counters(session, job, series_ids=sorted(affected_series_ids)) + await recompute_series_counters(session, job) + return affected_series_ids, retry_series_ids, normalized_series_count def _active_issue_recovery_clause() -> ColumnElement[bool]: @@ -231,11 +462,7 @@ async def recover_orphan( .order_by(ImportedFile.id.asc()) ) files = list(files_result.scalars().all()) - active_files = [ - imp_file - for imp_file in files - if imp_file.status not in {ImportedFileStatus.IMPORTED, ImportedFileStatus.SKIPPED} - ] + active_files = [imp_file for imp_file in files if requires_orphan_issue_decision(imp_file)] decision_by_id = {decision.imported_file_id: decision for decision in request.decisions} missing_ids = [imp_file.id for imp_file in active_files if imp_file.id not in decision_by_id] @@ -544,42 +771,30 @@ async def retry_failed_series( job, file_ids=normalized_file_ids, ) - repaired_series_ids: list[int] = [] + repaired_series_ids: set[int] = set() + prepared_target_series_ids: set[int] = set() + normalized_series_count = 0 if normalized_file_ids is None: - exhausted = await session.scalars( - sa_select(ImportedSeries).where( - ImportedSeries.import_job_id == job_id, - ImportedSeries.status == ImportSeriesStatus.FAILED, - ImportedSeries.error_message == "No eligible files available for import", - ~sa_select(ImportedFile.id) - .where( - ImportedFile.import_series_id == ImportedSeries.id, - ImportedFile.status.in_( - ( - ImportedFileStatus.MATCHED, - ImportedFileStatus.CONFIRMED, - ImportedFileStatus.FAILED, - ) - ), - ) - .exists(), - ) + ( + repaired_series_ids, + prepared_target_series_ids, + normalized_series_count, + ) = await _prepare_terminal_follow_up( + session, + job, ) - for item in exhausted: - if await retain_imported_series_outcome(session, item): - repaired_series_ids.append(item.id) if repaired_series_ids: - await recompute_file_counters(session, job, series_ids=repaired_series_ids) - await recompute_series_counters(session, job) await log_event( session, job_id, "INFO", - "import_partial_success_restored", + "import_terminal_follow_up_prepared", message=( - "Restored successful series outcomes; unresolved files remain in Follow-up." + "Restored successful outcomes and moved unresolved identities to Follow-up." ), - repaired_series_count=len(repaired_series_ids), + repaired_series_count=normalized_series_count, + affected_series_count=len(repaired_series_ids), + retry_series_count=len(prepared_target_series_ids), ) identity_blocked_count = 0 @@ -623,6 +838,11 @@ async def retry_failed_series( if not isinstance(dict(imp_file.diagnostics or {}).get("source_revalidation"), dict) ] retry_items_by_id = {item.id: item for item in [*failed_items, *failed_file_items]} + if prepared_target_series_ids: + prepared_items = await session.scalars( + sa_select(ImportedSeries).where(ImportedSeries.id.in_(prepared_target_series_ids)) + ) + retry_items_by_id.update({item.id: item for item in prepared_items}) else: scoped_items = await session.execute( sa_select(ImportedSeries, ImportedFile) @@ -664,7 +884,11 @@ async def retry_failed_series( retry_series_ids = [item.id for item in retry_items] for item in retry_items: - if item.status in {ImportSeriesStatus.FAILED, ImportSeriesStatus.IMPORTED}: + if item.status in { + ImportSeriesStatus.FAILED, + ImportSeriesStatus.IMPORTED, + ImportSeriesStatus.RECOVERY_PENDING, + }: item.status = ImportSeriesStatus.CONFIRMED item.error_message = None diff --git a/src/pullbox/services/import_review_recheck.py b/src/pullbox/services/import_review_recheck.py index 9219f8b9..ea598a62 100644 --- a/src/pullbox/services/import_review_recheck.py +++ b/src/pullbox/services/import_review_recheck.py @@ -8,7 +8,7 @@ from pathlib import Path from typing import TYPE_CHECKING, Any -from sqlalchemy import exists, func, or_, select +from sqlalchemy import and_, exists, func, or_, select from pullbox.core.exceptions import NotFoundError, ValidationError from pullbox.core.file_safety import ( @@ -38,7 +38,10 @@ ) from pullbox.models.library import LibraryRoot from pullbox.services.import_content_inspection import inspect_import_content -from pullbox.services.import_safety_diagnostics import build_import_safety_diagnostics +from pullbox.services.import_safety_diagnostics import ( + ImportSafetyCategory, + 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 @@ -49,6 +52,44 @@ from sqlalchemy.ext.asyncio import AsyncSession +_TRANSIENT_SOURCE_RECHECK_CATEGORIES = ( + ImportSafetyCategory.PERMISSION_UNREADABLE.value, + ImportSafetyCategory.ARCHIVE_INSPECTION_FAILED.value, +) +_TRANSIENT_SOURCE_RECHECK_CODES = ( + "permission_unreadable", + "permission_denied", + "source_unreadable", + "unreadable", + "archive_inspection_failed", + "corrupt_archive", + "inspection_failed", + "source_changed", + "source_signature_missing", + "source_signature_unsupported", + "source_unavailable", +) + + +def retryable_failed_source_filters(job_id: int) -> tuple[Any, ...]: + """Select only transient source failures that another inspection can resolve.""" + category = ImportedFile.diagnostics["source_revalidation"]["category"].as_string() + code = ImportedFile.diagnostics["source_revalidation"]["code"].as_string() + return ( + ImportedFile.import_job_id == job_id, + ImportedFile.status == ImportedFileStatus.FAILED, + ImportedFile.diagnostics["source_revalidation"]["retryable"].as_boolean().is_(True), + or_( + category.in_(_TRANSIENT_SOURCE_RECHECK_CATEGORIES), + code.in_(_TRANSIENT_SOURCE_RECHECK_CODES), + and_( + category == ImportSafetyCategory.SOURCE_CHANGED.value, + code.is_(None), + ), + ), + ) + + def _preserve_mylar_identity(base: SourceMetadata, fresh: SourceMetadata) -> SourceMetadata: updates: dict[str, Any] = {} signals = dict(fresh.signals) @@ -208,32 +249,15 @@ def _apply_completed_file_recheck( diagnostics = dict(file.diagnostics or {}) previous_signature = dict(file.source_signature or {}) source = {**metadata.diagnostics, **content} - block = source.pop("file_safety", None) - identity_conflicts = source.get("identity_conflicts") - saved_target_conflicts = _saved_target_identity_conflicts( + block = _completed_file_recheck_block( file, metadata, + source, reviewed_series_cv_id=reviewed_series_cv_id, ) - if block is None and isinstance(identity_conflicts, list) and identity_conflicts: - block = build_import_safety_diagnostics( - "The current source identity conflicts with the file reviewed during import.", - kind="source_revalidation", - code="source_identity_changed", - source="completed_import_recheck", - overrideable_hint=False, - ) - if block is None and saved_target_conflicts: - block = build_import_safety_diagnostics( - "The replacement source does not match the issue reviewed during import.", - kind="source_revalidation", - code="source_identity_changed", - source="completed_import_recheck", - overrideable_hint=False, - ) - block["identity_conflicts"] = saved_target_conflicts + source.pop("file_safety", None) checked_at = datetime.now(UTC).isoformat() - if isinstance(block, dict): + if block is not None: diagnostics["source_revalidation"] = { **block, "kind": "source_revalidation", @@ -280,6 +304,42 @@ def _apply_completed_file_recheck( return True +def _completed_file_recheck_block( + file: ImportedFile, + metadata: SourceMetadata, + source: dict[str, Any], + *, + reviewed_series_cv_id: int | None, +) -> dict[str, Any] | None: + """Return the final safety block after archive and identity checks.""" + raw_block = source.get("file_safety") + block = dict(raw_block) if isinstance(raw_block, dict) else None + identity_conflicts = source.get("identity_conflicts") + saved_target_conflicts = _saved_target_identity_conflicts( + file, + metadata, + reviewed_series_cv_id=reviewed_series_cv_id, + ) + if block is None and isinstance(identity_conflicts, list) and identity_conflicts: + block = build_import_safety_diagnostics( + "The current source identity conflicts with the file reviewed during import.", + kind="source_revalidation", + code="source_identity_changed", + source="completed_import_recheck", + overrideable_hint=False, + ) + if block is None and saved_target_conflicts: + block = build_import_safety_diagnostics( + "The replacement source does not match the issue reviewed during import.", + kind="source_revalidation", + code="source_identity_changed", + source="completed_import_recheck", + overrideable_hint=False, + ) + block["identity_conflicts"] = saved_target_conflicts + return block + + def _saved_target_identity_conflicts( file: ImportedFile, metadata: SourceMetadata, @@ -357,9 +417,7 @@ async def _retry_source_roots( else: candidates.extend(Path(value) for value in dict(job.mylar3_path_map or {}).values()) signature_query = select(ImportedFile.source_signature).where( - ImportedFile.import_job_id == job.id, - ImportedFile.status == ImportedFileStatus.FAILED, - ImportedFile.diagnostics["source_revalidation"]["retryable"].as_boolean().is_(True), + *retryable_failed_source_filters(job.id) ) if file_ids is not None: signature_query = signature_query.where(ImportedFile.id.in_(file_ids)) @@ -407,11 +465,7 @@ async def prepare_retryable_failed_sources_for_retry( file_ids: Sequence[int] | None = None, ) -> dict[str, int]: """Revalidate changed failed sources as part of the in-app retry action.""" - retryable_filters = [ - ImportedFile.import_job_id == job.id, - ImportedFile.status == ImportedFileStatus.FAILED, - ImportedFile.diagnostics["source_revalidation"]["retryable"].as_boolean().is_(True), - ] + retryable_filters = list(retryable_failed_source_filters(job.id)) if file_ids is not None: retryable_filters.append(ImportedFile.id.in_(file_ids)) retryable_count = int( @@ -480,10 +534,8 @@ async def prepare_completed_import_file_recheck( select(ImportedFile, ImportedSeries) .join(ImportedSeries, ImportedSeries.id == ImportedFile.import_series_id) .where( - ImportedFile.import_job_id == job_id, - ImportedFile.status == ImportedFileStatus.FAILED, + *retryable_failed_source_filters(job_id), ImportedFile.id > cursor, - ImportedFile.diagnostics["source_revalidation"]["retryable"].as_boolean().is_(True), ) .order_by(ImportedFile.id) .limit(250) @@ -509,17 +561,28 @@ async def prepare_completed_import_file_recheck( sidecars=sidecars, ) report["files_checked"] += 1 - blocked = "file_safety" in content - report["blocked_files"] += int(blocked) - report["files_prepared"] += int(not blocked) if apply: - _apply_completed_file_recheck( + ready_for_retry = _apply_completed_file_recheck( imported_file, metadata, content, signature, reviewed_series_cv_id=imported_series.cv_id, ) + else: + source = {**metadata.diagnostics, **content} + ready_for_retry = ( + _completed_file_recheck_block( + imported_file, + metadata, + source, + reviewed_series_cv_id=imported_series.cv_id, + ) + is None + ) + blocked = not ready_for_retry + report["blocked_files"] += int(blocked) + report["files_prepared"] += int(not blocked) if apply: await session.flush() diff --git a/src/pullbox/services/import_workflow_state.py b/src/pullbox/services/import_workflow_state.py index 513bb2fa..cada4113 100644 --- a/src/pullbox/services/import_workflow_state.py +++ b/src/pullbox/services/import_workflow_state.py @@ -35,12 +35,23 @@ SCAN_PROGRESS_FILE_MATCH_END = 99 WORKFLOW_SNAPSHOT_VERSION = 2 _PERSISTENT_IMPORT_CONTEXT_KEYS = ( + "deferred_recovery", "clean_library_adoption", "clean_library_adoption_prepared", "clean_library_source_snapshot", "source_import_job_id", ) ImportProgressMode = Literal["scan", "import", "rollback"] + + +def deferred_recovery_scope(job: ImportJob) -> tuple[int, ...] | None: + """A prepared follow-up authorizes only its own series groups, not the old review.""" + state = dict(dict(job.progress_snapshot or {}).get("deferred_recovery") or {}) + if state.get("state") != "prepared": + return None + return tuple(int(value) for value in state.get("series_ids", [])) + + _INVENTORY_PROGRESS_THRESHOLDS: tuple[tuple[int, int], ...] = ( (1, 1), (10, 2), diff --git a/src/pullbox/tasks/import_task.py b/src/pullbox/tasks/import_task.py index 2a15e766..cbdcade6 100644 --- a/src/pullbox/tasks/import_task.py +++ b/src/pullbox/tasks/import_task.py @@ -27,6 +27,10 @@ ) from pullbox.schemas.import_job import ImportProgressEvent from pullbox.services.import_counters import job_stats +from pullbox.services.import_deferred_recovery_execution import ( + cancel_deferred_preparation, + prepare_deferred_recovery, +) from pullbox.services.import_job_execution_progress import ( reconcile_durable_import_execution_counters, ) @@ -573,6 +577,9 @@ async def _finalize_cancel(self, session: AsyncSession, job_id: int) -> None: job = await session.get(ImportJob, job_id) if job is None: return + if await cancel_deferred_preparation(session, job): + purge_import_runtime_state(job_id) + return if job.import_started_at is None: await session.delete(job) await session.commit() @@ -696,6 +703,15 @@ async def progress_callback(event: ImportProgressEvent) -> None: ) await session.commit() elif job.status == ImportJobStatus.IMPORTING: + if dict(job.progress_snapshot or {}).get("deferred_recovery"): + await prepare_deferred_recovery( + session, + job_id, + metadata_service=service._metadata_service, + progress_callback=progress_callback, + ) + if job.status != ImportJobStatus.IMPORTING: + return await prepare_clean_library_import( session, job_id, @@ -815,6 +831,16 @@ async def progress_callback(event: ImportProgressEvent) -> None: progress_callback=progress_callback, ) else: + if dict(job.progress_snapshot or {}).get("deferred_recovery"): + await prepare_deferred_recovery( + session, + job_id, + metadata_service=service._metadata_service, + progress_callback=progress_callback, + ) + if job.status != ImportJobStatus.IMPORTING: + await session.commit() + return await prepare_clean_library_import( session, job_id, @@ -846,7 +872,9 @@ async def progress_callback(event: ImportProgressEvent) -> None: await session.rollback() job = await session.get(ImportJob, job_id) if job is not None: - if job.import_started_at is None: + if await cancel_deferred_preparation(session, job): + purge_import_runtime_state(job_id) + elif job.import_started_at is None: terminal_event_override = ImportProgressEvent( job_id=job_id, status=ImportJobStatus.CANCELLED, diff --git a/src/pullbox/ui/import_results_context.py b/src/pullbox/ui/import_results_context.py index 19559de4..617243c9 100644 --- a/src/pullbox/ui/import_results_context.py +++ b/src/pullbox/ui/import_results_context.py @@ -325,6 +325,16 @@ async def _load_files_for_status( } _CLEANUP_ACTION_PRESENTATION = { + CompletedImportCleanupAction.RECHECK_DEFERRED_FILES: { + "label": "Recheck deferred files", + "description": ( + "Group repeated file records, recognize completed imports, and recover exact issue " + "matches. Missing series catalogs are checked in the background. " + "Files that still need a decision remain here." + ), + "button_label": "Recheck files", + "tone": "warning", + }, CompletedImportCleanupAction.DISMISS_MISSING_REFERENCES: { "label": "Dismiss stale Mylar references", "description": ( @@ -427,6 +437,8 @@ async def _load_cleanup_action_summaries( "item_unit": ( "group" if action is CompletedImportCleanupAction.ACCEPT_RECOMMENDED_CONFLICTS + else "follow-up item" + if action is CompletedImportCleanupAction.RECHECK_DEFERRED_FILES else "file" ), "examples": summary.examples, diff --git a/src/pullbox/ui/static/js/pullbox.js b/src/pullbox/ui/static/js/pullbox.js index 20544cee..d8f04d44 100644 --- a/src/pullbox/ui/static/js/pullbox.js +++ b/src/pullbox/ui/static/js/pullbox.js @@ -8606,7 +8606,12 @@ function importResultsData(config) { if (!previewResponse.ok) { throw new Error(preview.detail || "Could not preview this cleanup action."); } - var unit = preview.item_unit === "group" ? "conflict group" : "file"; + var unit = + preview.item_unit === "group" + ? "conflict group" + : preview.item_unit === "follow-up item" + ? "follow-up item" + : "file"; var count = Math.max(0, Number(preview.affected_count) || 0); var examples = Array.isArray(preview.examples) ? preview.examples.slice(0, 3) : []; var message = @@ -8617,13 +8622,20 @@ function importResultsData(config) { unit + (count === 1 ? "" : "s") + ". Source files will remain unchanged."; + if (action === "recheck_deferred_files") { + message = + "Check " + count.toLocaleString() + " deferred file paths and stale series records in the background. " + + "Pullbox will reconcile repeated records and import only proven matches using " + + "this import's original file settings. Uncertain matches stay in Follow-up. " + + "Source files will remain unchanged."; + } if (examples.length) { message += " Examples: " + examples.join(", ") + "."; } var confirmed = await pbConfirm({ title: label + "?", message: message, - confirmText: "Apply cleanup", + confirmText: action === "recheck_deferred_files" ? "Recheck files" : "Apply cleanup", destructive: false, }); if (!confirmed) { diff --git a/tests/api/test_import_completed_cleanup_api.py b/tests/api/test_import_completed_cleanup_api.py index 76c3aa6c..f52c93bf 100644 --- a/tests/api/test_import_completed_cleanup_api.py +++ b/tests/api/test_import_completed_cleanup_api.py @@ -50,6 +50,58 @@ def _csrf_header_for(client: AsyncClient) -> dict[str, str]: return {"X-CSRF-Token": csrf} +@pytest.mark.asyncio +async def test_deferred_file_recheck_api_queues_only_after_signed_confirmation( + authenticated_client: AsyncClient, + sec_db: async_sessionmaker[AsyncSession], + monkeypatch: pytest.MonkeyPatch, +) -> None: + from pullbox.models.import_job import ImportFileHandlingMode + from tests.unit.test_import_deferred_recovery import add_file, seed + + async with sec_db() as session: + job, item, _, _, _ = await seed(session) + job.file_handling_mode = ImportFileHandlingMode.IN_PLACE + job.move_to_library = False + await add_file(session, job, item) + await add_file(session, job, item) + await session.commit() + job_id = job.id + calls = [] + monkeypatch.setattr( + "pullbox.api.v1.import_completed_cleanup.trigger_import_execute", calls.append + ) + url = f"/api/v1/import/{job_id}/cleanup/recheck_deferred_files" + preview = await authenticated_client.get(url + "/preview") + assert preview.status_code == 200 + assert preview.json()["affected_count"] == 1 + assert calls == [] + response = await authenticated_client.post( + url, + headers=_csrf_header_for(authenticated_client), + json={"preview_token": preview.json()["preview_token"], "confirmation": "APPLY CLEANUP"}, + ) + assert response.status_code == 200 + assert response.json()["requires_import_retry"] is True + assert calls == [job_id] + async with sec_db() as session: + job = await session.get(ImportJob, job_id) + assert job is not None + assert job.status is ImportJobStatus.IMPORTING + assert job.progress_snapshot["deferred_recovery"]["state"] == "queued" + assert ( + await session.scalar( + select(func.count()) + .select_from(ImportedFile) + .where( + ImportedFile.import_job_id == job_id, + ImportedFile.status == ImportedFileStatus.NO_MATCH, + ) + ) + == 2 + ) + + @pytest.mark.asyncio async def test_known_series_recovery_api_previews_then_queues_background_import( authenticated_client: AsyncClient, diff --git a/tests/ui/test_import_results_context.py b/tests/ui/test_import_results_context.py index afc7cf4b..b8de3534 100644 --- a/tests/ui/test_import_results_context.py +++ b/tests/ui/test_import_results_context.py @@ -134,7 +134,14 @@ async def test_load_import_results_context_splits_unmatched_queue_counts(db_sess assert context["no_match_count"] == 1 assert context["unmatched_queue_count"] == 2 assert context["cleanup_needs_review_count"] == 0 - assert context["follow_up_group_count"] == 1 + assert context["follow_up_group_count"] == 2 + recheck = next( + item + for item in context["cleanup_action_summaries"] + if item["action"] == "recheck_deferred_files" + ) + assert recheck["affected_count"] == 2 + assert recheck["button_label"] == "Recheck files" @pytest.mark.asyncio diff --git a/tests/unit/test_import_completed_cleanup.py b/tests/unit/test_import_completed_cleanup.py index 66c707a0..1eff7fa0 100644 --- a/tests/unit/test_import_completed_cleanup.py +++ b/tests/unit/test_import_completed_cleanup.py @@ -136,6 +136,98 @@ async def test_missing_references_can_be_dismissed_without_deleting_records( assert all(row.error_message is None for row in rows) +@pytest.mark.asyncio +async def test_rechecked_missing_references_use_the_existing_safe_cleanup( + db_session: AsyncSession, +) -> None: + job, imported_series = await _seed_job(db_session) + failed = ImportedFile( + import_job_id=job.id, + import_series_id=imported_series.id, + file_path="/comics/stale-reference.cbz", + file_name="stale-reference.cbz", + file_size=1024, + file_format="cbz", + status=ImportedFileStatus.FAILED, + diagnostics={ + "source_revalidation": { + "kind": "source_revalidation", + "category": ImportSafetyCategory.SOURCE_MISSING.value, + "code": ImportSafetyCategory.SOURCE_MISSING.value, + "reason": "The recorded file is missing.", + "retryable": True, + } + }, + error_message="The recorded file is missing.", + ) + db_session.add(failed) + await db_session.commit() + + preview = await preview_completed_import_cleanup( + db_session, + job.id, + CompletedImportCleanupAction.DISMISS_MISSING_REFERENCES, + actor_id=42, + ) + await apply_completed_import_cleanup( + db_session, + job.id, + CompletedImportCleanupAction.DISMISS_MISSING_REFERENCES, + actor_id=42, + preview_token=preview.preview_token, + ) + + assert preview.affected_file_count == 1 + assert failed.status is ImportedFileStatus.SKIPPED + assert failed.diagnostics["completed_import_cleanup"]["source_preserved"] is True + + +@pytest.mark.asyncio +async def test_rechecked_zero_byte_files_use_the_existing_safe_cleanup( + db_session: AsyncSession, +) -> None: + job, imported_series = await _seed_job(db_session) + failed = ImportedFile( + import_job_id=job.id, + import_series_id=imported_series.id, + file_path="/comics/empty.cbz", + file_name="empty.cbz", + file_size=0, + file_format="cbz", + status=ImportedFileStatus.FAILED, + diagnostics={ + "source_revalidation": { + "kind": "source_revalidation", + "category": ImportSafetyCategory.ZERO_BYTE.value, + "code": "zero_byte_file", + "reason": "The file is empty.", + "retryable": True, + } + }, + error_message="The file is empty.", + ) + db_session.add(failed) + await db_session.commit() + + preview = await preview_completed_import_cleanup( + db_session, + job.id, + CompletedImportCleanupAction.SKIP_UNUSABLE_FILES, + actor_id=42, + ) + await apply_completed_import_cleanup( + db_session, + job.id, + CompletedImportCleanupAction.SKIP_UNUSABLE_FILES, + actor_id=42, + preview_token=preview.preview_token, + ) + + assert preview.affected_file_count == 1 + assert failed.status is ImportedFileStatus.SKIPPED + assert failed.diagnostics["completed_import_cleanup"]["source_preserved"] is True + + @pytest.mark.asyncio async def test_probable_covers_are_separate_from_oversized_files( db_session: AsyncSession, @@ -1083,6 +1175,46 @@ async def test_retry_source_inspection_includes_failed_completed_rechecks( assert failed.diagnostics["source_revalidation"]["retryable"] is True +@pytest.mark.parametrize( + "code", + ("source_identity_changed", "source_root_changed", "source_root_unconfirmed"), +) +async def test_retry_source_inspection_excludes_identity_and_root_drift( + db_session: AsyncSession, + code: str, +) -> None: + job, imported_series = await _seed_job(db_session) + failed = ImportedFile( + import_job_id=job.id, + import_series_id=imported_series.id, + file_path=f"/comics/{code}.cbz", + file_name=f"{code}.cbz", + file_size=1024, + file_format="cbz", + status=ImportedFileStatus.FAILED, + diagnostics={ + "source_revalidation": { + "category": ImportSafetyCategory.SOURCE_CHANGED.value, + "code": code, + "retryable": True, + } + }, + error_message="The saved source identity changed.", + ) + db_session.add(failed) + await db_session.commit() + + summary = await summarize_completed_import_cleanup_scope( + db_session, + job.id, + CompletedImportCleanupAction.RETRY_SOURCE_INSPECTION, + ) + + assert summary.affected_count == 0 + assert summary.affected_file_count == 0 + assert summary.examples == () + + @pytest.mark.asyncio async def test_cleanup_preview_cannot_be_reused_after_scope_changes( db_session: AsyncSession, diff --git a/tests/unit/test_import_deferred_recovery.py b/tests/unit/test_import_deferred_recovery.py new file mode 100644 index 00000000..53b21d6c --- /dev/null +++ b/tests/unit/test_import_deferred_recovery.py @@ -0,0 +1,336 @@ +"""Deferred recovery respects physical identity, source evidence, and operator decisions.""" + +from datetime import UTC, date, datetime + +import pytest + +from pullbox.models.import_job import ( + ImportedFile, + ImportedFileStatus, + ImportedSeries, + ImportJob, + ImportJobStatus, + ImportSeriesStatus, + ImportSourceType, +) +from pullbox.models.issue import Issue, IssueType +from pullbox.models.library import FileFormat, LibraryFile, LibraryRoot +from pullbox.models.series import IssueCatalogState, Series +from pullbox.services.import_deferred_recovery import plan_deferred_recovery + + +async def seed(session, *, source_type=ImportSourceType.MYLAR3): + root = LibraryRoot(name="Comics", path="/comics") + job = ImportJob( + source_path="/imports", source_type=source_type, status=ImportJobStatus.COMPLETED + ) + session.add_all([root, job]) + await session.flush() + series = Series( + title="Batman", + sort_title="batman", + year_start=2016, + comicvine_id=100, + library_root_id=root.id, + issue_catalog_state=IssueCatalogState.COMPLETE, + ) + session.add(series) + await session.flush() + issue = Issue( + series_id=series.id, + comicvine_id=1001, + issue_number=104, + issue_number_text="104", + issue_type=IssueType.ISSUE, + release_date=date(2021, 1, 1), + ) + item = ImportedSeries( + import_job_id=job.id, + raw_series_name="Batman", + raw_year=2016, + cv_id=100, + series_id=series.id, + status=ImportSeriesStatus.IMPORTED, + ) + session.add_all([issue, item]) + await session.flush() + return job, item, series, issue, root + + +async def add_file(session, job, item, **overrides): + values = dict( + import_job_id=job.id, + import_series_id=item.id, + file_path="/comics/Batman/Batman 104 (2021).cbz", + file_name="Batman 104 (2021).cbz", + file_size=1024, + file_format="cbz", + parsed_series="Batman", + parsed_issue_number=104, + parsed_year=2021, + status=ImportedFileStatus.NO_MATCH, + comicvine_issue_id=1001, + diagnostics={"source_issue_type": "issue"}, + source_signature={"size_bytes": 1024, "mtime_ns": 123456}, + ) + values.update(overrides) + result = ImportedFile(**values) + session.add(result) + await session.flush() + return result + + +async def register(session, file, issue, root): + library = LibraryFile( + file_path=file.file_path, + file_name=file.file_name, + file_size=file.file_size, + file_format=FileFormat.CBZ, + file_modified_at=datetime.now(UTC), + issue_id=issue.id, + library_root_id=root.id, + source_signature=file.source_signature, + ) + session.add(library) + await session.flush() + return library + + +async def test_conflicting_saved_content_digest_never_collapses_files(db_session): + job, item, _, issue, root = await seed(db_session) + file = await add_file(db_session, job, item) + file.source_signature = {**file.source_signature, "content_digest": "one"} + registered = await register(db_session, file, issue, root) + registered.source_signature = {**registered.source_signature, "content_digest": "two"} + assert await plan_deferred_recovery(db_session, job.id) == () + + +@pytest.mark.parametrize("source_type", list(ImportSourceType)) +async def test_exact_registered_path_is_handled_without_reimport(db_session, source_type): + job, item, _, issue, root = await seed(db_session, source_type=source_type) + file = await add_file(db_session, job, item) + await register(db_session, file, issue, root) + await db_session.commit() + + plans = await plan_deferred_recovery(db_session, job.id) + + assert [(p.file_id, p.action, p.issue_id) for p in plans] == [ + (file.id, "already_registered", issue.id) + ] + assert file.status is ImportedFileStatus.NO_MATCH + assert not db_session.dirty + + +async def test_one_physical_file_with_different_mylar_parents_has_one_decision(db_session): + job, item, _, issue, _ = await seed(db_session) + wrong = ImportedSeries( + import_job_id=job.id, + raw_series_name="Superman", + cv_id=200, + status=ImportSeriesStatus.NO_MATCH, + ) + db_session.add(wrong) + await db_session.flush() + file = await add_file(db_session, job, item) + twin = await add_file( + db_session, + job, + wrong, + diagnostics={ + "comicvine_series_id": 200, + "metadata_signals": {"comicvine_series_id": "mylar3"}, + }, + ) + plans = await plan_deferred_recovery(db_session, job.id) + assert [(p.file_id, p.action) for p in plans] == [ + (file.id, "exact_target"), + (twin.id, "duplicate_reference"), + ] + assert plans[1].canonical_file_id == file.id + assert plans[0].issue_id == issue.id + + +async def test_distinct_variant_paths_are_left_for_one_duplicate_decision(db_session): + job, item, _, issue, root = await seed(db_session) + owned = await add_file(db_session, job, item, status=ImportedFileStatus.IMPORTED) + await register(db_session, owned, issue, root) + variant = await add_file(db_session, job, item, file_path="/comics/Batman/104 variant.cbz") + plans = await plan_deferred_recovery(db_session, job.id) + assert [(p.file_id, p.action) for p in plans] == [(variant.id, "owned_variant")] + + +async def test_stale_provider_id_without_issue_number_evidence_remains_deferred(db_session): + job, item, _, _, _ = await seed(db_session) + await add_file( + db_session, + job, + item, + file_name="Batman.cbz", + parsed_issue_number=None, + issue_number_raw=None, + ) + + assert await plan_deferred_recovery(db_session, job.id) == () + + +@pytest.mark.parametrize( + "protection", ["safety", "manual", "skip", "conflicting_id", "wrong_title", "changed"] +) +async def test_recovery_preserves_protected_or_contradictory_evidence(db_session, protection): + job, item, _, issue, root = await seed(db_session) + file = await add_file(db_session, job, item) + if protection == "safety": + file.diagnostics = {"safety_block": {"category": "dangerous_path_or_payload"}} + elif protection == "manual": + file.match_method = "manual_issue" + elif protection == "skip": + file.status = ImportedFileStatus.SKIPPED + elif protection == "conflicting_id": + file.matched_issue_cv_id = 9999 + elif protection == "wrong_title": + file.parsed_series = "Superman" + else: + await register(db_session, file, issue, root) + file.source_signature = {"size_bytes": 1024, "mtime_ns": 9999} + assert await plan_deferred_recovery(db_session, job.id) == () + + +@pytest.mark.parametrize( + "name,number", + [ + ("batman.104. (2021).cbz", None), + ("Batman 104 (2021) (converted).cbz", None), + ("Batman 104 (2021).cbz", 104), + ], +) +async def test_contextual_reparse_requires_exact_series_type_number_and_date( + db_session, name, number +): + job, item, _, issue, _ = await seed(db_session) + file = await add_file( + db_session, job, item, comicvine_issue_id=None, file_name=name, parsed_issue_number=number + ) + plans = await plan_deferred_recovery(db_session, job.id) + assert [(p.file_id, p.action, p.issue_id) for p in plans] == [ + (file.id, "exact_target", issue.id) + ] + + +@pytest.mark.parametrize( + "name,year,issue_type", + [ + ("Batman Annual 104 (2021).cbz", 2021, "annual"), + ("Batman 104 (1990).cbz", 1990, "issue"), + ("Batman 104-105 (2021).cbz", 2021, "issue"), + ("Batman Vol. 104 (2021).cbz", 2021, "tpb"), + ("Absolute Batman 104 (2021).cbz", 2021, "issue"), + ], +) +async def test_number_only_recovery_does_not_cross_type_title_pack_or_year( + db_session, name, year, issue_type +): + job, item, _, _, _ = await seed(db_session) + await add_file( + db_session, + job, + item, + comicvine_issue_id=None, + file_name=name, + parsed_series=None, + parsed_year=year, + diagnostics={"source_issue_type": issue_type}, + ) + assert await plan_deferred_recovery(db_session, job.id) == () + + +async def test_identical_unresolved_path_groups_once_without_inventing_a_target(db_session): + job, item, _, _, _ = await seed(db_session) + file = await add_file( + db_session, + job, + item, + comicvine_issue_id=None, + parsed_series="Unknown", + file_name="Unknown.cbz", + parsed_issue_number=None, + ) + twin = await add_file( + db_session, + job, + item, + comicvine_issue_id=None, + parsed_series="Unknown", + file_name="Unknown.cbz", + parsed_issue_number=None, + ) + plans = await plan_deferred_recovery(db_session, job.id) + assert [(p.file_id, p.action, p.canonical_file_id) for p in plans] == [ + (twin.id, "duplicate_reference", file.id) + ] + + +async def test_registered_identity_survives_container_device_number_change(db_session): + job, item, _, issue, root = await seed(db_session) + file = await add_file( + db_session, + job, + item, + source_signature={ + "size": 1024, + "mtime_ns": 123456, + "device": 98, + "inode": 123, + }, + ) + library = await register(db_session, file, issue, root) + library.source_signature = {**file.source_signature, "device": 99} + plans = await plan_deferred_recovery(db_session, job.id) + assert plans[0].action == "already_registered" + + +async def test_stale_mylar_id_can_be_reconciled_by_embedded_id_and_registered_path(db_session): + job, item, _, issue, root = await seed(db_session) + file = await add_file( + db_session, + job, + item, + comicvine_issue_id=999, + diagnostics={ + "metadata_signals": {"comicvine_issue_id": "mylar3", "series_name": "comicinfo"}, + "source_metadata": { + "archive_metadata_loaded": True, + "comicinfo": { + "series": "Batman", + "number": "104", + "web": "https://comicvine.gamespot.com/batman/4000-1001/", + }, + "identity_conflicts": [ + {"field": "comicvine_issue_id", "first": 999, "conflicting": 1001} + ], + }, + }, + ) + await register(db_session, file, issue, root) + plans = await plan_deferred_recovery(db_session, job.id) + assert [(p.action, p.issue_id) for p in plans] == [("already_registered", issue.id)] + + +async def test_conflicting_embedded_ids_remain_reviewable(db_session): + job, item, _, issue, root = await seed(db_session) + file = await add_file( + db_session, + job, + item, + diagnostics={ + "source_metadata": { + "comicinfo": { + "series": "Batman", + "number": "104", + "web": "https://comicvine.gamespot.com/batman/4000-1001/", + "notes": "[cv_issue_id:999]", + } + }, + }, + ) + await register(db_session, file, issue, root) + assert await plan_deferred_recovery(db_session, job.id) == () diff --git a/tests/unit/test_import_deferred_recovery_execution.py b/tests/unit/test_import_deferred_recovery_execution.py new file mode 100644 index 00000000..cb636106 --- /dev/null +++ b/tests/unit/test_import_deferred_recovery_execution.py @@ -0,0 +1,435 @@ +"""Recovery stages work without changing source files or reviving unrelated decisions.""" + +from unittest.mock import AsyncMock + +import pytest + +from pullbox.core.exceptions import JobPausedError, ProviderError +from pullbox.models.import_job import ( + ImportedFileStatus, + ImportedSeries, + ImportJobStatus, + ImportSeriesStatus, +) +from pullbox.providers.base import IssueSummary, SeriesMetadata +from pullbox.services.import_deferred_recovery import ( + apply_deferred_recovery, + plan_deferred_recovery, +) +from pullbox.services.import_deferred_recovery_execution import prepare_deferred_recovery +from tests.unit.test_import_deferred_recovery import add_file, register, seed + + +async def test_apply_is_repeatable_and_prepares_only_its_exact_scope(db_session): + job, item, _, _, _ = await seed(db_session) + file = await add_file(db_session, job, item) + unrelated = await add_file( + db_session, job, item, file_path="/comics/Batman/105.cbz", status=ImportedFileStatus.MATCHED + ) + original_path = file.file_path + counts = await apply_deferred_recovery(db_session, job) + assert counts["exact_target"] == 1 + assert file.status is ImportedFileStatus.CONFIRMED + assert file.import_series_id != item.id + assert file.file_path == original_path + assert unrelated.status is ImportedFileStatus.MATCHED + assert unrelated.import_series_id == item.id + assert ( + await db_session.get(ImportedSeries, file.import_series_id) + ).status is ImportSeriesStatus.CONFIRMED + assert await plan_deferred_recovery(db_session, job.id) == () + assert (await apply_deferred_recovery(db_session, job)).get("exact_target", 0) == 0 + + +async def test_duplicate_reference_keeps_canonical_link_and_evidence(db_session): + job, item, _, issue, root = await seed(db_session) + file = await add_file(db_session, job, item) + twin = await add_file(db_session, job, item) + await register(db_session, file, issue, root) + await apply_deferred_recovery(db_session, job) + assert file.status is ImportedFileStatus.ALREADY_OWNED + assert twin.status is ImportedFileStatus.SKIPPED + assert twin.duplicate_of_file_id == file.id + assert twin.diagnostics["deferred_recovery"]["source_preserved"] + assert job.total_files_no_match == 0 + + +async def test_only_stale_series_without_any_file_records_are_archived(db_session): + job, item, _, _, _ = await seed(db_session) + item.status = ImportSeriesStatus.NO_MATCH + item.diagnostics = {"reason": "path_missing"} + stale = ImportedSeries( + import_job_id=job.id, + raw_series_name="Stale", + status=ImportSeriesStatus.NO_MATCH, + diagnostics={"reason": "path_missing"}, + ) + db_session.add(stale) + await add_file( + db_session, + job, + item, + comicvine_issue_id=None, + parsed_series="Unidentified", + file_name="Unidentified.cbz", + parsed_issue_number=None, + ) + result = await apply_deferred_recovery(db_session, job) + assert result["stale_series"] == 1 + assert stale.status is ImportSeriesStatus.SKIPPED + assert item.status is ImportSeriesStatus.NO_MATCH + + +async def test_distinct_unowned_candidates_are_not_silently_selected(db_session): + job, item, _, _, _ = await seed(db_session) + await add_file(db_session, job, item) + await add_file(db_session, job, item, file_path="/comics/Batman/104 alternate.cbz") + assert await plan_deferred_recovery(db_session, job.id) == () + + +async def test_background_recovery_fetches_each_candidate_catalog_once_and_resumes(db_session): + job, item, _, _, _ = await seed(db_session) + item.series_id = None + item.cv_id = None + item.status = ImportSeriesStatus.NO_MATCH + file = await add_file( + db_session, + job, + item, + comicvine_issue_id=7001, + diagnostics={ + "comicvine_series_id": 700, + "metadata_signals": {"comicvine_series_id": "mylar3"}, + }, + ) + twin = await add_file( + db_session, job, item, comicvine_issue_id=7001, diagnostics=file.diagnostics + ) + 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", + year_start=2016, + issue_count=1, + sort_title="batman", + year_end=None, + status=None, + publisher=None, + description=None, + cover_url=None, + comicvine_url=None, + ) + provider.get_issue_summaries_for_series.return_value = [ + IssueSummary( + provider_id="7001", + issue_number=104, + issue_number_text="104", + title=None, + release_date="2021-01-01", + cover_url=None, + issue_type="issue", + ) + ] + + assert await prepare_deferred_recovery(db_session, job.id, metadata_service=provider) + assert file.status is ImportedFileStatus.CONFIRMED + assert twin.status is ImportedFileStatus.SKIPPED + assert file.diagnostics["target_issue_summary"]["provider_id"] == "7001" + target = await db_session.get(ImportedSeries, file.import_series_id) + assert target.cv_id == 700 + assert item.status is ImportSeriesStatus.SKIPPED + assert file.file_path == twin.file_path + provider.get_issue_summaries_for_series.assert_awaited_once_with(700) + provider.get_series_metadata.assert_awaited_once_with(700) + assert not await prepare_deferred_recovery(db_session, job.id, metadata_service=provider) + assert provider.get_issue_summaries_for_series.await_count == 1 + + +async def test_catalog_target_respects_strong_archive_number_evidence(db_session): + from pullbox.services.import_deferred_recovery_execution import _catalog_target_agrees + + job, item, _, _, _ = await seed(db_session) + file = await add_file( + db_session, + job, + item, + diagnostics={ + "source_metadata": { + "archive_entry_issue_hint": {"confidence": "strong", "issue_number": 105} + } + }, + ) + assert not _catalog_target_agrees( + file, + item, + { + "title": "Batman", + "summary": {"provider_id": "1001", "issue_number": 104, "issue_type": "issue"}, + }, + ) + + +async def test_local_background_recovery_requires_no_provider_calls(db_session): + job, item, _, _, _ = await seed(db_session) + file = await add_file(db_session, job, item) + job.status = ImportJobStatus.IMPORTING + job.progress_snapshot = {"deferred_recovery": {"state": "queued"}} + provider = AsyncMock() + assert await prepare_deferred_recovery(db_session, job.id, metadata_service=provider) + assert file.status is ImportedFileStatus.CONFIRMED + provider.get_series_metadata.assert_not_awaited() + provider.get_issue_summaries_for_series.assert_not_awaited() + + +async def test_catalog_failure_checkpoints_without_repeating_completed_requests(db_session): + job, item, _, _, _ = await seed(db_session) + item.series_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"}, + "source_metadata": { + "identity_conflicts": [ + {"field": "comicvine_series_id", "first": 700, "conflicting": 800}, + ] + }, + }, + ) + job.status = ImportJobStatus.IMPORTING + job.progress_snapshot = {"deferred_recovery": {"state": "queued"}} + provider = AsyncMock() + + async def get_series(cv_id): + assert not db_session.in_transaction() + return SeriesMetadata( + provider_id=str(cv_id), + 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.side_effect = [ + [], + ProviderError("comicvine", "offline"), + ] + with pytest.raises(JobPausedError): + await prepare_deferred_recovery(db_session, job.id, metadata_service=provider) + assert job.progress_snapshot["deferred_recovery"]["completed"] == ["700"] + provider.get_issue_summaries_for_series.side_effect = None + provider.get_issue_summaries_for_series.return_value = [] + await prepare_deferred_recovery(db_session, job.id, metadata_service=provider) + assert [call.args[0] for call in provider.get_issue_summaries_for_series.await_args_list] == [ + 700, + 800, + 800, + ] + assert job.status is ImportJobStatus.COMPLETED + + +async def test_preparation_cancel_preserves_original_import(db_session): + from pullbox.services.import_deferred_recovery_execution import cancel_deferred_preparation + + job, item, _, _, _ = await seed(db_session) + original = await add_file(db_session, job, item, status=ImportedFileStatus.IMPORTED) + file = await add_file(db_session, job, item, file_path="/comics/Batman/other.cbz") + job.progress_snapshot = {"deferred_recovery": {"state": "queued"}} + await apply_deferred_recovery(db_session, job) + job.status = ImportJobStatus.CANCELLING + assert await cancel_deferred_preparation(db_session, job) + assert job.status is ImportJobStatus.COMPLETED + assert original.status is ImportedFileStatus.IMPORTED + assert file.status is ImportedFileStatus.NO_MATCH + + +@pytest.mark.parametrize("started", [True, False]) +async def test_paused_recovery_cancel_never_requests_original_rollback(db_session, started): + from datetime import UTC, datetime + + from pullbox.services.import_job_controls import request_cancel + + job, _, _, _, _ = await seed(db_session) + job.status = ImportJobStatus.PAUSED + job.import_started_at = datetime.now(UTC) if started else None + job.progress_snapshot = {"deferred_recovery": {"state": "catalogs"}} + await db_session.commit() + await request_cancel(db_session, job.id, log_event=AsyncMock()) + assert job.status is ImportJobStatus.CANCELLING + + +async def test_recovery_execution_leaves_unrelated_ready_groups_and_arcs_untouched( + db_session, monkeypatch +): + from pullbox.models.import_job import ImportFileHandlingMode + from pullbox.services import import_job_execution as execution + + job, item, _, _, _ = await seed(db_session) + job.file_handling_mode = ImportFileHandlingMode.IN_PLACE + job.move_to_library = False + file = await add_file(db_session, job, item) + await apply_deferred_recovery(db_session, job) + unrelated = ImportedSeries( + import_job_id=job.id, + raw_series_name="Not selected", + cv_id=999, + status=ImportSeriesStatus.CONFIRMED, + ) + db_session.add(unrelated) + await db_session.flush() + await add_file( + db_session, + job, + unrelated, + status=ImportedFileStatus.CONFIRMED, + file_path="/comics/other.cbz", + ) + job.progress_snapshot = { + **job.progress_snapshot, + "deferred_recovery": {**job.progress_snapshot["deferred_recovery"], "state": "prepared"}, + } + execute_group = AsyncMock(return_value=(1, 0, 1, 0, True)) + arcs = AsyncMock(side_effect=AssertionError("Unrelated story arcs must not execute")) + monkeypatch.setattr(execution, "_execute_new_series", execute_group) + monkeypatch.setattr(execution, "_execute_story_arc_materialization", arcs) + await execution.execute_import_job( + db_session, + job.id, + series_service=AsyncMock(), + process_series_files=AsyncMock(), + raise_if_cancelled=AsyncMock(), + record_action=AsyncMock(), + log_event=AsyncMock(), + emit_progress=AsyncMock(), + estimate_remaining_seconds=lambda *a, **kw: None, + maybe_slow_item_delay=AsyncMock(), + ) + assert execute_group.await_count == 1 + assert execute_group.await_args.args[2].id == file.import_series_id + assert unrelated.status is ImportSeriesStatus.CONFIRMED + assert job.progress_snapshot["deferred_recovery"]["state"] == "completed" + assert job.status is ImportJobStatus.COMPLETED + + +async def test_copy_capacity_is_scoped_to_recovery_files(db_session): + from pullbox.services.import_managed_copy_preflight import selected_managed_copy_source_bytes + + job, item, _, _, _ = await seed(db_session) + await add_file(db_session, job, item) + await apply_deferred_recovery(db_session, job) + item.status = ImportSeriesStatus.CONFIRMED + await add_file( + db_session, + job, + item, + file_path="/comics/unrelated.cbz", + status=ImportedFileStatus.CONFIRMED, + file_size=99999, + ) + job.progress_snapshot = { + **job.progress_snapshot, + "deferred_recovery": {**job.progress_snapshot["deferred_recovery"], "state": "prepared"}, + } + assert await selected_managed_copy_source_bytes(db_session, job.id) == 1024 + + +async def test_cleanup_preview_counts_physical_paths_and_queues_without_inline_work(db_session): + from pullbox.models.import_job import ImportFileHandlingMode + from pullbox.models.user import User + from pullbox.services.import_completed_cleanup import ( + CompletedImportCleanupAction, + apply_completed_import_cleanup, + preview_completed_import_cleanup, + summarize_completed_import_cleanup_scope, + ) + + db_session.add(User(id=42, username="recovery-test", password_hash="unused")) + job, item, _, _, _ = await seed(db_session) + job.file_handling_mode = ImportFileHandlingMode.IN_PLACE + job.move_to_library = False + files = [await add_file(db_session, job, item), await add_file(db_session, job, item)] + await db_session.commit() + action = CompletedImportCleanupAction.RECHECK_DEFERRED_FILES + preview = await preview_completed_import_cleanup(db_session, job.id, action, actor_id=42) + assert preview.affected_count == 1 + assert preview.affected_file_count == 2 + summary = await summarize_completed_import_cleanup_scope(db_session, job.id, action) + assert summary.affected_count == 1 + result = await apply_completed_import_cleanup( + db_session, job.id, action, actor_id=42, preview_token=preview.preview_token + ) + assert result.requires_import_retry + assert job.status is ImportJobStatus.IMPORTING + assert job.progress_snapshot["deferred_recovery"]["state"] == "queued" + assert all(file.status is ImportedFileStatus.NO_MATCH for file in files) + + +async def test_cleanup_preview_signs_and_bounds_empty_stale_series(db_session): + from pullbox.models.import_job import ImportFileHandlingMode + from pullbox.models.user import User + from pullbox.services.import_completed_cleanup import ( + CompletedImportCleanupAction, + apply_completed_import_cleanup, + preview_completed_import_cleanup, + ) + + db_session.add(User(id=42, username="recovery-test", password_hash="unused")) + job, item, _, _, _ = await seed(db_session) + job.file_handling_mode = ImportFileHandlingMode.IN_PLACE + job.move_to_library = False + await add_file(db_session, job, item) + stale = ImportedSeries( + import_job_id=job.id, + raw_series_name="Missing before preview", + status=ImportSeriesStatus.NO_MATCH, + diagnostics={"reason": "path_missing"}, + ) + db_session.add(stale) + await db_session.commit() + + action = CompletedImportCleanupAction.RECHECK_DEFERRED_FILES + preview = await preview_completed_import_cleanup(db_session, job.id, action, actor_id=42) + + assert preview.affected_count == 2 + assert preview.affected_file_count == 1 + result = await apply_completed_import_cleanup( + db_session, job.id, action, actor_id=42, preview_token=preview.preview_token + ) + assert result.affected_count == 2 + assert job.progress_snapshot["deferred_recovery"]["stale_series_ids"] == [stale.id] + + late_stale = ImportedSeries( + import_job_id=job.id, + raw_series_name="Missing after preview", + status=ImportSeriesStatus.NO_MATCH, + diagnostics={"reason": "source_missing"}, + ) + db_session.add(late_stale) + await db_session.flush() + + await apply_deferred_recovery(db_session, job, running=True) + + assert stale.status is ImportSeriesStatus.SKIPPED + assert late_stale.status is ImportSeriesStatus.NO_MATCH + + +async def test_active_import_is_not_changed_by_offline_recovery(db_session): + job, item, _, _, _ = await seed(db_session) + item.status = ImportSeriesStatus.NO_MATCH + item.diagnostics = {"reason": "path_missing"} + job.status = ImportJobStatus.IMPORTING + assert await apply_deferred_recovery(db_session, job) == {} + assert item.status is ImportSeriesStatus.NO_MATCH diff --git a/tests/unit/test_import_orphans.py b/tests/unit/test_import_orphans.py index 8c5b92b4..1eade473 100644 --- a/tests/unit/test_import_orphans.py +++ b/tests/unit/test_import_orphans.py @@ -180,6 +180,23 @@ def test_apply_orphan_recovery_decisions_skips_file() -> None: assert imp_file.diagnostics["resolution"] == "skipped" +def test_source_failure_is_not_an_issue_recovery_decision() -> None: + from pullbox.services.import_orphans import requires_orphan_issue_decision + + source_failure = ImportedFile( + status=ImportedFileStatus.FAILED, + diagnostics={ + "source_revalidation": { + "category": "source_missing", + "retryable": True, + } + }, + ) + + assert requires_orphan_issue_decision(source_failure) is False + assert requires_orphan_issue_decision(ImportedFile(status=ImportedFileStatus.NO_MATCH)) is True + + def test_summarize_orphan_recovery_marks_imported_when_no_files_remaining() -> None: from pullbox.services.import_orphans import summarize_orphan_recovery_result @@ -790,6 +807,229 @@ async def test_retry_failed_repairs_exhausted_partial_success_without_reimport(d assert item.diagnostics["previous_series_error"] == "No eligible files available for import" +async def test_retry_failed_restores_partial_success_with_unresolved_target_to_follow_up( + db_session, +): + service = _make_service() + job = await _create_job_row(db_session, series_failed=1) + item = await _create_imported_series( + db_session, job, name="Hellblazer", status=ImportSeriesStatus.FAILED + ) + series = Series(title="Hellblazer", sort_title="hellblazer", comicvine_id=4008) + db_session.add(series) + await db_session.flush() + item.series_id = series.id + item.cv_id = 4008 + item.error_message = "No eligible files available for import" + imported = ImportedFile( + import_job_id=job.id, + import_series_id=item.id, + file_path="/comics/hellblazer-001.cbz", + file_name="Hellblazer 001.cbz", + file_size=1024, + file_format="cbz", + status=ImportedFileStatus.IMPORTED, + ) + unresolved = ImportedFile( + import_job_id=job.id, + import_series_id=item.id, + file_path="/comics/hellblazer-special.cbz", + file_name="Hellblazer Special.cbz", + file_size=1024, + file_format="cbz", + status=ImportedFileStatus.FAILED, + error_message="Could not resolve to a library issue", + diagnostics={"kind": "file_conflict"}, + ) + db_session.add_all([imported, unresolved]) + await db_session.flush() + + updated_job, count = await service.retry_failed_series(db_session, job.id) + + assert count == 0 + assert updated_job.status is ImportJobStatus.COMPLETED + assert item.status is ImportSeriesStatus.IMPORTED + assert unresolved.status is ImportedFileStatus.NO_MATCH + assert unresolved.include_in_import is False + assert item.files_imported == 1 + assert item.files_no_match == 1 + assert job.series_failed == 0 + + +async def test_retry_failed_routes_series_without_comicvine_id_to_follow_up(db_session): + service = _make_service() + job = await _create_job_row(db_session, series_failed=1) + item = await _create_imported_series( + db_session, job, name="Unknown anthology", status=ImportSeriesStatus.FAILED + ) + item.error_message = "No ComicVine ID available" + db_session.add( + ImportedFile( + import_job_id=job.id, + import_series_id=item.id, + file_path="/comics/unknown-001.cbz", + file_name="Unknown anthology 001.cbz", + file_size=1024, + file_format="cbz", + status=ImportedFileStatus.MATCHED, + ) + ) + await db_session.flush() + + updated_job, count = await service.retry_failed_series(db_session, job.id) + + assert count == 0 + assert updated_job.status is ImportJobStatus.COMPLETED + assert item.status is ImportSeriesStatus.NO_MATCH + assert item.error_message is None + assert item.diagnostics["previous_series_error"] == "No ComicVine ID available" + assert job.series_failed == 0 + assert job.series_no_match == 1 + + +async def test_retry_failed_requeues_unique_local_issue_target(db_session): + service = _make_service() + job = await _create_job_row(db_session) + item = await _create_imported_series( + db_session, job, name="2000AD", status=ImportSeriesStatus.IMPORTED + ) + await _create_series_with_issue( + db_session, + series_id=160, + issue_id=3218, + title="2000AD", + issue_number=2487.0, + ) + item.series_id = 160 + item.cv_id = 1783 + failed_file = ImportedFile( + import_job_id=job.id, + import_series_id=item.id, + file_path="/comics/2000AD-2487.cbz", + file_name="2000AD 2487.cbz", + file_size=1024, + file_format="cbz", + status=ImportedFileStatus.FAILED, + parsed_issue_number=2487.0, + match_confidence="high", + match_method="issue_number", + error_message="Could not resolve to a library issue", + diagnostics={"kind": "file_conflict"}, + ) + db_session.add(failed_file) + await db_session.flush() + + updated_job, count = await service.retry_failed_series(db_session, job.id) + + assert count == 1 + assert updated_job.status is ImportJobStatus.IMPORTING + assert item.status is ImportSeriesStatus.CONFIRMED + assert failed_file.status is ImportedFileStatus.CONFIRMED + assert failed_file.matched_issue_id == 3218 + assert failed_file.include_in_import is True + assert failed_file.match_method == "completed_import_exact_target" + + +async def test_retry_failed_does_not_use_issue_number_after_provider_id_conflict(db_session): + service = _make_service() + job = await _create_job_row(db_session) + item = await _create_imported_series( + db_session, job, name="Aquaman", status=ImportSeriesStatus.IMPORTED + ) + await _create_series_with_issue( + db_session, + series_id=160, + issue_id=3218, + title="Aquaman", + issue_number=1.0, + ) + item.series_id = 160 + item.cv_id = 91738 + failed_file = ImportedFile( + import_job_id=job.id, + import_series_id=item.id, + file_path="/comics/aquaman-volume-1.cbz", + file_name="Aquaman v01 - The Drowning.cbz", + file_size=1024, + file_format="cbz", + status=ImportedFileStatus.FAILED, + parsed_issue_number=1.0, + matched_issue_cv_id=1089296, + match_confidence="high", + match_method="issue_number", + error_message="Could not resolve to a library issue", + diagnostics={ + "kind": "file_conflict", + "target_issue_summary": { + "provider_id": "1089296", + "issue_number": 1.0, + }, + }, + ) + db_session.add(failed_file) + await db_session.flush() + + updated_job, count = await service.retry_failed_series(db_session, job.id) + + assert count == 0 + assert updated_job.status is ImportJobStatus.COMPLETED + assert item.status is ImportSeriesStatus.IMPORTED + assert failed_file.status is ImportedFileStatus.NO_MATCH + assert failed_file.matched_issue_id is None + assert failed_file.include_in_import is False + + +async def test_retry_failed_does_not_treat_volume_as_standard_issue_one(db_session): + service = _make_service() + job = await _create_job_row(db_session) + item = await _create_imported_series( + db_session, job, name="Aquaman", status=ImportSeriesStatus.IMPORTED + ) + await _create_series_with_issue( + db_session, + series_id=160, + issue_id=3218, + title="Aquaman", + issue_number=1.0, + ) + item.series_id = 160 + item.cv_id = 91738 + failed_file = ImportedFile( + import_job_id=job.id, + import_series_id=item.id, + file_path="/comics/aquaman-volume-1.cbz", + file_name="Aquaman v01 - The Drowning.cbz", + file_size=1024, + file_format="cbz", + status=ImportedFileStatus.FAILED, + match_confidence="high", + match_method="issue_number", + error_message="Could not resolve to a library issue", + diagnostics={ + "kind": "file_conflict", + "source_issue_type": IssueType.TPB.value, + "source_metadata": { + "filename_parse": { + "issue_number": None, + "issue_type": IssueType.TPB.value, + "volume": "v01", + } + }, + }, + ) + db_session.add(failed_file) + await db_session.flush() + + updated_job, count = await service.retry_failed_series(db_session, job.id) + + assert count == 0 + assert updated_job.status is ImportJobStatus.COMPLETED + assert item.status is ImportSeriesStatus.IMPORTED + assert failed_file.status is ImportedFileStatus.NO_MATCH + assert failed_file.matched_issue_id is None + assert failed_file.include_in_import is False + + async def test_retry_failed_does_not_requeue_unresolved_legacy_series_identity(db_session): service = _make_service() job = await _create_job_row(db_session, series_failed=1) diff --git a/tests/unit/test_import_review_recheck.py b/tests/unit/test_import_review_recheck.py index bc175659..6eb5073a 100644 --- a/tests/unit/test_import_review_recheck.py +++ b/tests/unit/test_import_review_recheck.py @@ -25,6 +25,7 @@ from pullbox.services.import_review_recheck import ( prepare_completed_import_file_recheck, prepare_import_recheck, + prepare_retryable_failed_sources_for_retry, prepare_review_recheck, ) @@ -298,7 +299,7 @@ async def test_completed_recheck_keeps_missing_source_blocked(db_session, tmp_pa missing.diagnostics = { **missing.diagnostics, "source_revalidation": { - "code": "source_missing", + "code": "source_changed", "retryable": True, }, } @@ -321,6 +322,48 @@ async def test_completed_recheck_keeps_missing_source_blocked(db_session, tmp_pa assert missing.diagnostics["source_revalidation"]["code"] == "source_missing" +@pytest.mark.parametrize( + "category", + ( + "source_missing", + "zero_byte", + "archive_no_pages", + "unsupported_file_type", + "source_identity_changed", + "outside_approved_root", + ), +) +async def test_retry_failed_does_not_reinspect_terminal_source_failures( + db_session, + tmp_path, + category, +): + job, item, files = await _fixture(db_session, tmp_path, ImportSourceType.MYLAR3) + job.status = ImportJobStatus.COMPLETED + item.status = ImportSeriesStatus.IMPORTED + failed = files[0] + failed.status = ImportedFileStatus.FAILED + failed.diagnostics = { + **failed.diagnostics, + "source_revalidation": { + "category": category, + "code": category, + "retryable": True, + }, + } + await db_session.flush() + + report = await prepare_retryable_failed_sources_for_retry(db_session, job) + + assert report == { + "files_checked": 0, + "files_prepared": 0, + "blocked_files": 0, + "skipped_files": 0, + } + assert failed.status is ImportedFileStatus.FAILED + + @pytest.mark.parametrize( "source_type", [ImportSourceType.MYLAR3, ImportSourceType.FILESYSTEM], @@ -553,6 +596,48 @@ async def test_retry_failed_rejects_changed_source_with_conflicting_identity( assert changed.diagnostics["source_revalidation"]["code"] == "source_identity_changed" +async def test_completed_recheck_counts_identity_conflict_as_blocked( + db_session, + tmp_path, +): + 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}, + } + (tmp_path / "Firefly (2018)" / "cvinfo").write_text( + "https://comicvine.gamespot.com/other/4050-123456/" + ) + path = Path(changed.file_path) + with zipfile.ZipFile(path, "w") as archive: + archive.writestr("1.jpg", b"replacement image") + archive.writestr("2.jpg", b"replacement image") + await db_session.flush() + + 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": 1, + "skipped_files": 0, + } + assert changed.diagnostics["source_revalidation"]["code"] == "source_identity_changed" + + async def test_retry_failed_rejects_replacement_for_different_saved_issue( db_session, tmp_path,