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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
8 changes: 8 additions & 0 deletions docs/development/IMPORT_PERFORMANCE.md
Original file line number Diff line number Diff line change
Expand Up @@ -12,6 +12,14 @@ Archive page-name matching runs off the event loop and uses a bounded local
cache for repeated page-title parsing. It still counts every page toward
consensus and does not reuse safety decisions between scans.

Completed-import source recovery inspects each bounded page before its database
mutation phase, then commits before inspecting the next page. Slow files
therefore do not hold SQLite's writer lock, and results do not accumulate in an
unbounded in-memory collection.
Deferred catalog recovery publishes live provider progress without rewriting
the full recovery snapshot, then writes one JSON-safe checkpoint after each
completed catalog.

`PULLBOX_IMPORT_SCAN_WORKER_COUNT=0` selects automatic inspection concurrency,
up to four workers. CPU affinity, cgroup v2 CPU quotas and parent limits,
cgroup memory headroom, and OS available memory cap the budget. Common cgroup
Expand Down
8 changes: 8 additions & 0 deletions docs/development/IMPORT_REVIEW_RECOVERY.md
Original file line number Diff line number Diff line change
Expand Up @@ -118,6 +118,10 @@ candidates; membership in one catalog plus agreeing file evidence is required
before staging a target. Completed catalog checks are checkpointed. A provider
failure pauses the pass, and Resume continues without repeating completed
checks. Network requests do not hold a database write transaction.
Catalog summaries are converted to JSON-safe checkpoint payloads before they
are stored. Live provider progress does not rewrite the full durable recovery
snapshot; each completed catalog produces one durable checkpoint, so a worker
restart resumes after the last completed catalog without replaying it.

Recovered files run through normal Step 4 safety, current-source validation,
ownership checks, and the original copy or keep-in-place settings. Only newly
Expand All @@ -143,6 +147,10 @@ A completed source recheck reports a file ready only after both archive safety
and saved target identity checks pass. Missing, empty, or otherwise blocked
sources are counted as blocked even when the archive-level inspection itself
completed successfully.
Completed-import source rechecks inspect one bounded page before writing its
refreshed evidence, then commit that page before reading more archives. This
keeps slow archive I/O outside SQLite's single-writer window, bounds memory, and
leaves completed pages durable if a later source needs another attempt.

Recovery queries must not expand an entire library into SQL bind parameters.
Mixed-folder lookups join existing references and discard exact same-title
Expand Down
3 changes: 3 additions & 0 deletions src/pullbox/api/v1/import_completed_cleanup.py
Original file line number Diff line number Diff line change
Expand Up @@ -117,6 +117,9 @@ async def apply_completed_import_cleanup_route(
)

if action is CompletedImportCleanupAction.RETRY_SOURCE_INSPECTION:
# Persist the signed cleanup transition before slow archive inspection.
# The recheck then starts without an existing SQLite writer lock.
await session.commit()
service = await build_import_service(session)
_job, retrying_count = await service.retry_failed_series(
session,
Expand Down
12 changes: 12 additions & 0 deletions src/pullbox/services/diagnostic_db_snapshot.py
Original file line number Diff line number Diff line change
Expand Up @@ -29,6 +29,18 @@ def create_sanitized_db_copy(db_path: Path) -> bytes | None:
src.backup(dst)
src.close()

tables = {
str(row[0])
for row in dst.execute(
"SELECT name FROM sqlite_schema WHERE type='table'"
).fetchall()
}
if "audit_logs" in tables:
# Keep operational evidence while removing its private account link.
dst.execute("UPDATE audit_logs SET user_id = NULL")
if "issue_reader_states" in tables:
# Per-user reading state has no meaning once users are removed.
dst.execute("DELETE FROM issue_reader_states")
dst.execute("DELETE FROM users")
dst.execute("DELETE FROM api_keys")

Expand Down
88 changes: 66 additions & 22 deletions src/pullbox/services/import_deferred_recovery_execution.py
Original file line number Diff line number Diff line change
Expand Up @@ -35,13 +35,18 @@
refresh_recovered_groups,
)
from pullbox.services.import_source_metadata import source_metadata_for_import_file
from pullbox.services.import_workflow_state import emit_progress, raise_if_job_cancelled
from pullbox.services.import_workflow_state import (
emit_live_progress,
emit_progress,
raise_if_job_cancelled,
)

if TYPE_CHECKING:
from collections.abc import Awaitable, Callable

from sqlalchemy.ext.asyncio import AsyncSession

from pullbox.providers.base import IssueSummary
from pullbox.services.metadata_service import MetadataService


Expand All @@ -53,6 +58,15 @@ def save_recovery_state(job: ImportJob, state: dict[str, Any]) -> None:
job.progress_snapshot = {**dict(job.progress_snapshot or {}), "deferred_recovery": state}


def _catalog_summary_payload(summary: IssueSummary) -> dict[str, Any]:
"""Return a durable JSON-safe catalog checkpoint payload."""
payload = asdict(summary)
source_cutoff_at = payload.get("source_cutoff_at")
if isinstance(source_cutoff_at, datetime):
payload["source_cutoff_at"] = source_cutoff_at.isoformat()
return payload


async def cancel_deferred_preparation(session: AsyncSession, job: ImportJob) -> bool:
"""Stop this recovery pass, retaining the original and any completed imports."""
state = recovery_state(job)
Expand Down Expand Up @@ -264,25 +278,49 @@ async def prepare_deferred_recovery(
if job.status is not ImportJobStatus.IMPORTING:
raise ValidationError("Deferred recovery must run inside the import worker.")

async def report(current: int, total: int, message: str) -> None:
await raise_if_job_cancelled(session, job_id)
await emit_progress(
session,
revision_state = {"value": int(job.progress_revision or 0)}

def progress_event(current: int, total: int, message: str) -> ImportProgressEvent:
return ImportProgressEvent(
job_id=job_id,
status=ImportJobStatus.IMPORTING,
mode="import",
phase="deferred_recovery",
progress=round(15 * current / max(total, 1)),
message=message,
current_file_stage="deferred_recovery",
current_file_progress_current=current,
current_file_progress_total=total,
current_file_progress_pct=round(100 * current / max(total, 1)),
current_file_progress_unit="catalogs",
)

async def report(
current: int,
total: int,
message: str,
*,
durable: bool = True,
check_control: bool = True,
) -> None:
if check_control:
await raise_if_job_cancelled(session, job_id)
event = progress_event(current, total, message)
if durable:
event.progress_revision = revision_state["value"] + 1
await emit_progress(session, job, event, progress_callback)
revision_state["value"] = event.progress_revision
return
# A read-only control check still opens a transaction. Close it before
# provider I/O, then publish live progress without rewriting the large
# durable recovery checkpoint a second time.
await session.commit()
await emit_live_progress(
job,
ImportProgressEvent(
job_id=job_id,
status=ImportJobStatus.IMPORTING,
mode="import",
phase="deferred_recovery",
progress=round(15 * current / max(total, 1)),
message=message,
current_file_stage="deferred_recovery",
current_file_progress_current=current,
current_file_progress_total=total,
current_file_progress_pct=round(100 * current / max(total, 1)),
current_file_progress_unit="catalogs",
),
progress_callback,
event,
progress_callback=progress_callback,
revision_state=revision_state,
started_at=job.import_started_at,
)

if state.get("state") == "queued":
Expand All @@ -308,8 +346,9 @@ async def report(current: int, total: int, message: str) -> None:
len(completed),
len(candidates),
f"Checking series catalog {len(completed) + 1} of {len(candidates)}...",
durable=False,
)
# The progress commit releases the writer lock before provider I/O.
# The live progress report closes its read transaction before provider I/O.
cv_id = int(cv_id_text)
try:
series = await metadata_service.get_series_metadata(cv_id)
Expand All @@ -335,13 +374,18 @@ async def report(current: int, total: int, message: str) -> None:
"cv_id": cv_id,
"title": series.title,
"year": series.year_start,
"summary": asdict(summary),
"summary": _catalog_summary_payload(summary),
}
)
completed.add(cv_id_text)
state["completed"] = sorted(completed)
save_recovery_state(job, state)
await session.commit()
await report(
len(completed),
len(candidates),
f"Checked series catalog {len(completed)} of {len(candidates)}.",
check_control=False,
)

await report(len(completed), max(len(candidates), 1), "Preparing verified files for import...")
catalog_count = await _prepare_catalog_targets(session, job)
Expand Down
75 changes: 64 additions & 11 deletions src/pullbox/services/import_review_recheck.py
Original file line number Diff line number Diff line change
Expand Up @@ -547,6 +547,14 @@ async def prepare_completed_import_file_recheck(
rows = list((await session.execute(query)).all())
if not rows:
break
inspected: list[
tuple[
int,
SourceMetadata,
dict[str, Any],
dict[str, int | str],
]
] = []
for imported_file, imported_series in rows:
cursor = int(imported_file.id)
metadata, content, signature = await asyncio.to_thread(
Expand All @@ -562,13 +570,7 @@ async def prepare_completed_import_file_recheck(
)
report["files_checked"] += 1
if apply:
ready_for_retry = _apply_completed_file_recheck(
imported_file,
metadata,
content,
signature,
reviewed_series_cv_id=imported_series.cv_id,
)
inspected.append((int(imported_file.id), metadata, content, signature))
else:
source = {**metadata.diagnostics, **content}
ready_for_retry = (
Expand All @@ -580,11 +582,62 @@ async def prepare_completed_import_file_recheck(
)
is None
)
blocked = not ready_for_retry
report["blocked_files"] += int(blocked)
report["files_prepared"] += int(not blocked)
blocked = not ready_for_retry
report["blocked_files"] += int(blocked)
report["files_prepared"] += int(not blocked)

if apply:
await session.flush()
# Inspect the complete bounded page before taking SQLite's writer
# lock. Reload and lock the job and eligible rows so a concurrent
# retry cannot have its newer import evidence overwritten.
current_job = await session.scalar(
select(ImportJob)
.where(ImportJob.id == job_id)
.with_for_update()
.execution_options(populate_existing=True)
)
if current_job is None:
raise NotFoundError("ImportJob", job_id)
if not allows_terminal_import_recovery(current_job):
raise ValidationError(
"Job changed while failed sources were being inspected; retry the recheck"
)

inspected_by_id = {
file_id: (metadata, content, signature)
for file_id, metadata, content, signature in inspected
}
apply_query = (
select(ImportedFile, ImportedSeries)
.join(ImportedSeries, ImportedSeries.id == ImportedFile.import_series_id)
.where(
*retryable_failed_source_filters(job_id),
ImportedFile.id.in_(inspected_by_id),
)
.order_by(ImportedFile.id)
.with_for_update()
.execution_options(populate_existing=True)
)
current_rows = list((await session.execute(apply_query)).all())
report["skipped_files"] += len(inspected) - len(current_rows)
for imported_file, imported_series in current_rows:
metadata, content, signature = inspected_by_id[int(imported_file.id)]
ready_for_retry = _apply_completed_file_recheck(
Comment thread
DeusExTaco marked this conversation as resolved.
imported_file,
metadata,
content,
signature,
reviewed_series_cv_id=imported_series.cv_id,
)
blocked = not ready_for_retry
report["blocked_files"] += int(blocked)
report["files_prepared"] += int(not blocked)
# Commit before reading the next page of source files.
await session.commit()
rows.clear()
current_rows.clear()
inspected.clear()
inspected_by_id.clear()

if apply and report["files_checked"]:
session.add(
Expand Down
89 changes: 89 additions & 0 deletions tests/api/test_import_completed_cleanup_api.py
Original file line number Diff line number Diff line change
Expand Up @@ -102,6 +102,95 @@ async def test_deferred_file_recheck_api_queues_only_after_signed_confirmation(
)


@pytest.mark.asyncio
async def test_retry_source_inspection_commits_cleanup_before_recheck(
db_session: AsyncSession,
monkeypatch: pytest.MonkeyPatch,
) -> None:
from starlette.requests import Request

from pullbox.api.v1.import_completed_cleanup import apply_completed_import_cleanup_route
from pullbox.models.user import User
from pullbox.schemas.import_completed_cleanup import CompletedImportCleanupApplyRequest
from pullbox.services.import_completed_cleanup import (
CompletedImportCleanupAction,
preview_completed_import_cleanup,
)

user = User(id=42, username="recovery-test", password_hash="unused")
job = ImportJob(
source_path="/imports/mylar.db",
source_type=ImportSourceType.MYLAR3,
status=ImportJobStatus.COMPLETED,
)
db_session.add_all([user, job])
await db_session.flush()
imported_series = ImportedSeries(
import_job_id=job.id,
raw_series_name="Temporarily unavailable",
status=ImportSeriesStatus.IMPORTED,
)
db_session.add(imported_series)
await db_session.flush()
block = build_import_safety_diagnostics(
ImportSafetyCategory.ARCHIVE_INSPECTION_FAILED.value,
code=ImportSafetyCategory.ARCHIVE_INSPECTION_FAILED.value,
)
db_session.add(
ImportedFile(
import_job_id=job.id,
import_series_id=imported_series.id,
file_path="/comics/temporarily-unavailable.cbz",
file_name="temporarily-unavailable.cbz",
file_size=1024,
file_format="cbz",
status=ImportedFileStatus.SAFETY_BLOCKED,
diagnostics={"safety_block": block},
)
)
await db_session.commit()
job_id = int(job.id)
action = CompletedImportCleanupAction.RETRY_SOURCE_INSPECTION
preview = await preview_completed_import_cleanup(db_session, job_id, action, actor_id=user.id)

transaction_states: list[bool] = []

class RetryService:
async def retry_failed_series(
self,
session: AsyncSession,
requested_job_id: int,
*,
file_ids: list[int] | None = None,
) -> tuple[ImportJob, int]:
transaction_states.append(session.in_transaction())
job = await session.get(ImportJob, requested_job_id)
assert job is not None
return job, 0

async def build_retry_service(_session: AsyncSession) -> RetryService:
return RetryService()

monkeypatch.setattr(
"pullbox.composition.services.build_import_service",
build_retry_service,
)
response = await apply_completed_import_cleanup_route(
job_id,
action,
CompletedImportCleanupApplyRequest(
preview_token=preview.preview_token,
confirmation="APPLY CLEANUP",
),
Request({"type": "http", "client": ("127.0.0.1", 12345), "headers": []}),
user,
db_session,
)

assert response.requires_import_retry is False
assert transaction_states == [False]


@pytest.mark.asyncio
async def test_known_series_recovery_api_previews_then_queues_background_import(
authenticated_client: AsyncClient,
Expand Down
Loading
Loading