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
11 changes: 11 additions & 0 deletions docs/development/ARCHITECTURE_OVERVIEW.md
Original file line number Diff line number Diff line change
Expand Up @@ -583,6 +583,17 @@ effects such as:
- File matched handling.
- Series added handling.

Utility cancellation is shared by the serial utility queue. Repeated requests
are idempotent, and cancellation takes precedence over a pending pause. No new
batch can be leased after cancellation commits. Busy workers observe a per-pool
signal; conversion and metadata stages on disposable outputs use interruptible
archive subprocesses. Source replacement and its result journal finish before
the queue releases its execution slot. Other non-interruptible item operations
finish at their safe item boundary rather than being killed mid-mutation.
Terminal job state and activity progress are committed together. Startup recovery
finishes abandoned `CANCELLING` jobs as `CANCELLED`, preserving completed results
and rollback records without replaying unfinished files.

Composition helpers make event-bus intent explicit:

- `build_domain_event_bus()` returns the shared application event bus.
Expand Down
17 changes: 15 additions & 2 deletions docs/development/IMPORT_PERFORMANCE.md
Original file line number Diff line number Diff line change
Expand Up @@ -18,7 +18,11 @@ 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.
completed catalog. Once preparation has durably stored its decisions on file
rows, consumed catalog and title-lookup caches are removed from the progress
snapshot. Recovery scope, summary counts, and restart state remain. Older
prepared/completed snapshots are compacted on their next progress checkpoint;
unfinished catalog preparation retains its caches for resume.

`PULLBOX_IMPORT_SCAN_WORKER_COUNT=0` selects automatic inspection concurrency,
up to four workers. CPU affinity, cgroup v2 CPU quotas and parent limits,
Expand All @@ -34,13 +38,22 @@ Automatic mode deliberately does not saturate high-core machines: local ZIP
header parsing did not improve with more than 2-4 workers. Docker Desktop's
VM resources are the relevant limits, not the Mac's advertised RAM.

Step 4 keeps `PULLBOX_IMPORT_FILE_WORKER_COUNT=2` and its existing temporary-space
Managed Step 4 keeps `PULLBOX_IMPORT_FILE_WORKER_COUNT=2` and its existing temporary-space
preflight, target-collision serialization, per-worker sessions, and rollback
journal. It now bounds submitted tasks as well as active workers, rather than
creating one waiting task per file. Exiting or canceling either worker pool
drains active work before the job can transition; no orphan filesystem work
may continue after cancellation is reported complete.

SQLite in-place registration uses one file worker, with an isolated transaction
per file. It retries only transient SQLite lock failures, rolling back and
rechecking control requests and source identity in a fresh session. Exhausted
retries leave the import resumable as database-busy, not a failed match. Source
archives, skips, safety blocks, and confirmed targets are not loosened. Progress
inside the file transaction is live-only; the terminal file checkpoint is
persisted after commit, avoiding a competing progress writer. Archive scanning
and managed copy/conversion keep their existing bounded parallelism.

## Progress and Evidence

Unknown inventory totals are indeterminate. Completed series report 100% for
Expand Down
37 changes: 37 additions & 0 deletions docs/development/IMPORT_REVIEW_RECOVERY.md
Original file line number Diff line number Diff line change
Expand Up @@ -34,6 +34,26 @@ recovery is available in Import Follow-up and does not require an
offline command. It works for Mylar and folder imports. This is not a full
rescan, a database restore, or an import.

## Managed File Publication

Mylar and folder imports publish completed sibling staging files with an atomic
no-overwrite rename where available (Linux `RENAME_NOREPLACE`, macOS
`RENAME_EXCL`, or Windows exclusive rename). Copy mode therefore does not
require hard-link support. If exclusive rename is unsupported, the existing
link/unlink publication remains a fallback. Pullbox never falls back to an
overwriting rename or exposes a partially copied final file. Existing targets,
including dangling symlinks, remain protected; journaled stages and source
restoration retain their existing recovery rules.

Confirmation and execution preflight exercise publication and collision
protection with disposable files under each selected managed destination.
Failure returns an actionable storage message before bulk work starts. The
probe runs off the event loop, does not change import source files, and does not
run for reference-only in-place imports. A root probe cannot guarantee that
every existing subdirectory has the same permissions or that a mount will stay
available; execution failures retain the destination and underlying filesystem
error for retry. No database migration is required.

## Mylar Inventory And Missing References

Copy and keep-in-place scans retain the same Mylar file inventory, including
Expand Down Expand Up @@ -221,6 +241,23 @@ checkpoint. Resume preserves completed repairs, and repeated runs do not create
another file registration. This is logical library repair, not authorization to
reorganize the user's filesystem.

Recovery reads the original per-file ComicInfo title and issue designation even
when older imports saved the parent folder's name as the parsed title. Explicit
reading-order prefixes and `(converted)` annotations are removed for candidate
lookups only; filenames remain unchanged. Deferred and referenced files share
the same per-file catalog checks, including publication dates and issue types.
A date accepted for one file never authorizes another file in the same batch.
Conflicting embedded identities, manual decisions, unsafe sources and duplicate
copies remain protected. Unresolved catalog candidates retain their proposed
series, issue and reason in `mixed_folder_recovery` diagnostics.

An issue mistakenly registered under the parent series can itself be reassigned
only when its embedded ComicVine issue ID agrees with the unique catalog target,
it is the current referenced issue, and no other file or imported registration
owns it. The issue ID and reader history are preserved. The old series identity
is retained in the cleanup audit. A provisional issue is not offered under a
folder series when independent file evidence identifies an unrelated series.

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
Expand Down
99 changes: 99 additions & 0 deletions src/pullbox/core/file_publication.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,99 @@
"""Publish completed staging files without replacing an existing destination."""

from __future__ import annotations

import ctypes
import errno
import os
import stat
import sys
import tempfile
from pathlib import Path

_UNSUPPORTED_RENAME = {errno.ENOSYS, errno.EINVAL, errno.ENOTSUP, errno.EOPNOTSUPP}
_AT_FDCWD = -100
_RENAME_NOREPLACE = 1
_RENAME_EXCL = 4


def publish_file_without_overwrite(stage: Path, target: Path) -> None:
"""Consume a stage using an atomic destination claim, never plain POSIX rename.

Native exclusive rename avoids requiring hard-link support on managed-copy
destinations. Older filesystems retain the link/unlink fallback. Neither
path copies bytes into a visible, partially written final file.
"""
mode = stage.lstat().st_mode
if not (stat.S_ISREG(mode) or stat.S_ISLNK(mode)):
raise OSError(errno.EINVAL, "Publication requires a regular file or symlink", str(stage))
if _native_rename_without_overwrite(stage, target):
return
if stat.S_ISLNK(mode):
os.symlink(os.readlink(stage), target)
else:
os.link(stage, target, follow_symlinks=False)
stage.unlink()


def _native_rename_without_overwrite(stage: Path, target: Path) -> bool:
if sys.platform == "win32":
# Windows rename fails if the destination exists; POSIX rename does not.
os.rename(stage, target)
return True
if sys.platform not in {"linux", "darwin"}:
return False

source_bytes, target_bytes = os.fsencode(stage), os.fsencode(target)
if b"\0" in source_bytes or b"\0" in target_bytes:
raise ValueError("embedded null byte")
libc = ctypes.CDLL(None, use_errno=True)
symbol = "renameat2" if sys.platform == "linux" else "renamex_np"
rename = getattr(libc, symbol, None)
if rename is None:
return False
rename.restype = ctypes.c_int
if sys.platform == "linux":
rename.argtypes = [
ctypes.c_int,
ctypes.c_char_p,
ctypes.c_int,
ctypes.c_char_p,
ctypes.c_uint,
]
result = rename(_AT_FDCWD, source_bytes, _AT_FDCWD, target_bytes, _RENAME_NOREPLACE)
else:
rename.argtypes = [ctypes.c_char_p, ctypes.c_char_p, ctypes.c_uint]
result = rename(source_bytes, target_bytes, _RENAME_EXCL)
if result == 0:
return True
error_number = ctypes.get_errno()
if error_number in _UNSUPPORTED_RENAME:
return False
raise OSError(error_number, os.strerror(error_number), str(stage), None, str(target))


def probe_file_publication(root: Path) -> None:
"""Exercise publication and collision protection using only disposable files."""
with tempfile.TemporaryDirectory(prefix=".pullbox-publication-probe-", dir=root) as directory:
stage = Path(directory) / "stage"
target = Path(directory) / "target"
stage.write_bytes(b"original")
publish_file_without_overwrite(stage, target)
stage.write_bytes(b"replacement")
try:
publish_file_without_overwrite(stage, target)
except FileExistsError:
if target.read_bytes() == b"original" and stage.read_bytes() == b"replacement":
return
raise OSError(errno.ENOTSUP, "Storage did not preserve an existing destination", str(root))


def publication_failure_message(directory: Path, error: OSError) -> str:
"""Actionable destination guidance shared by preflight and execution failures."""
return (
f"Cannot safely publish imported files in {directory}: {error.strerror or str(error)}. "
"Check the destination mount, free space, and the container user's file permissions. "
"The destination must support exclusive rename or hard links without replacing "
"existing files. Correct the storage settings or choose another managed library root, "
"then retry."
)
72 changes: 43 additions & 29 deletions src/pullbox/services/import_deferred_recovery_execution.py
Original file line number Diff line number Diff line change
Expand Up @@ -36,6 +36,12 @@
provider_ids,
refresh_recovered_groups,
)
from pullbox.services.import_recovery_checkpoint import compact_recovery_state
from pullbox.services.import_recovery_identity import (
catalog_file_identity,
catalog_target_agrees,
record_catalog_review,
)
from pullbox.services.import_reference_recovery import (
reference_candidates,
repair_catalog_references,
Expand All @@ -61,7 +67,10 @@ def recovery_state(job: ImportJob) -> dict[str, Any]:


def save_recovery_state(job: ImportJob, state: dict[str, Any]) -> None:
job.progress_snapshot = {**dict(job.progress_snapshot or {}), "deferred_recovery": state}
job.progress_snapshot = {
**dict(job.progress_snapshot or {}),
"deferred_recovery": compact_recovery_state(state),
}


def _catalog_summary_payload(summary: IssueSummary) -> dict[str, Any]:
Expand All @@ -76,45 +85,31 @@ def _catalog_summary_payload(summary: IssueSummary) -> dict[str, Any]:
def _filename_catalog_identity(
file: ImportedFile,
item: ImportedSeries,
) -> dict[str, str | int] | None:
) -> dict[str, Any] | None:
"""Return bounded filename identity eligible for local-catalog recovery."""
if protected_file(file, item) or provider_ids(file):
if protected_file(file, item):
return None
diagnostics = dict(file.diagnostics or {})
source = source_metadata_for_import_file(item, file).diagnostics
if source.get("identity_conflicts"):
identity = catalog_file_identity(file)
if identity is None:
return None
parsed = source.get("filename_parse")
parsed = parsed if isinstance(parsed, dict) else {}
source_title = str(parsed.get("series_name") or "").strip()
raw_number = parsed.get("issue_number_text") or parsed.get("issue_number")
if not source_title or raw_number is None:
ids = provider_ids(file)
if ids and ids != {identity.get("issue_cv_id")}:
return None
try:
issue_number = normalize_issue_number_text(str(raw_number))
except ValueError:
return None
issue_type = str(parsed.get("issue_type") or diagnostics.get("source_issue_type") or "issue")
if issue_type in {"annual", "special"}:
label = issue_type.title()
if not NameMatcher.normalize(source_title).endswith(f" {NameMatcher.normalize(label)}"):
source_title = f"{source_title} {label}"
normalized_title = NameMatcher.normalize(source_title)
source_title = identity["query"]
normalized_title = identity["key"]
parent_title = NameMatcher.normalize(item.cv_title or item.raw_series_name)
if not normalized_title or normalized_title == parent_title:
return None
title_match = NameMatcher().match(source_title, item.cv_title or item.raw_series_name)
corroborated = diagnostics.get("conflict_type") == "corroborated_file_series_mismatch"
if title_match.is_match and not corroborated and issue_type not in {"annual", "special"}:
if (
title_match.is_match
and not corroborated
and identity["issue_type"] not in {"annual", "special"}
):
return None
year = parsed.get("year") or file.parsed_year
return {
"key": normalized_title,
"query": source_title,
"issue_number": issue_number,
"year": int(year) if isinstance(year, int | float) else 0,
"issue_type": issue_type,
}
return identity


async def _catalog_title_candidates(
Expand Down Expand Up @@ -431,7 +426,17 @@ async def _prepare_title_catalog_targets(session: AsyncSession, job: ImportJob)
continue
options_by_number = title_matches.get(str(identity["key"]), {})
options = options_by_number.get(str(identity["issue_number"]), [])
options = [option for option in options if catalog_target_agrees(identity, option)]
if len(options) != 1:
record_catalog_review(
file,
identity,
(
"Multiple catalog issues agree; choose the correct series and issue."
if options
else "No catalog issue agrees with this file's title, number, type and year."
),
)
continue
eligible.append((file, item, options[0]))

Expand Down Expand Up @@ -465,6 +470,15 @@ async def _prepare_title_catalog_targets(session: AsyncSession, job: ImportJob)
or counts[str(issue_cv_id)] != 1
or issue_cv_id in existing_issue_ids
):
identity = _filename_catalog_identity(file, original)
if identity is not None:
record_catalog_review(
file,
identity,
"Multiple files claim this issue; choose which copy to keep."
if counts[str(issue_cv_id)] != 1
else "This issue already has a catalog entry; review its existing ownership.",
)
continue
target_item = targets.get(target_cv_id)
if target_item is None:
Expand Down
Loading
Loading