From 1dd902723807a0c214fb798692e463ddf78c6130 Mon Sep 17 00:00:00 2001 From: Adam Hernandez Date: Sun, 13 Sep 2026 15:23:36 -0700 Subject: [PATCH 1/4] feat: add verified local Comic Vine catalog downloads and matching --- .gitignore | 1 + CHANGELOG.md | 5 + docs/development/LOCAL_CATALOG_V2.md | 138 +++++++ pyproject.toml | 1 + src/pullbox/api/v1/catalog.py | 41 +++ src/pullbox/api/v1/router.py | 2 + src/pullbox/app.py | 13 + src/pullbox/composition/services.py | 6 + src/pullbox/services/catalog/__init__.py | 1 + src/pullbox/services/catalog/contract.py | 165 +++++++++ src/pullbox/services/catalog/database.py | 235 ++++++++++++ src/pullbox/services/catalog/lookup.py | 103 ++++++ src/pullbox/services/catalog/reader.py | 193 ++++++++++ src/pullbox/services/catalog/retention.py | 47 +++ src/pullbox/services/catalog/service.py | 344 ++++++++++++++++++ src/pullbox/services/catalog/storage.py | 116 ++++++ src/pullbox/services/import_cv_search.py | 5 + src/pullbox/services/import_provider_cache.py | 4 +- src/pullbox/services/import_service.py | 2 +- .../services/import_service_matching.py | 6 +- src/pullbox/services/metadata_service.py | 73 +++- src/pullbox/services/series_service.py | 3 +- src/pullbox/tasks/__init__.py | 1 + src/pullbox/tasks/catalog_task.py | 19 + src/pullbox/ui/comicvine_provider.py | 8 + src/pullbox/ui/import_orphaned_routes.py | 16 +- src/pullbox/ui/import_routes.py | 16 +- src/pullbox/ui/series_routes.py | 11 +- .../components/comicvine_search_loading.html | 4 +- .../partials/add_series_header_metrics.html | 2 +- .../partials/add_series_results.html | 13 +- .../templates/partials/settings_catalog.html | 75 ++++ .../templates/partials/settings_metadata.html | 5 +- tests/catalog_fixtures.py | 141 +++++++ tests/e2e/test_accessibility.py | 1 + tests/ui/test_catalog_controls.py | 90 +++++ tests/unit/test_catalog_contract.py | 95 +++++ tests/unit/test_catalog_controls.py | 23 ++ tests/unit/test_catalog_database.py | 75 ++++ tests/unit/test_catalog_integration.py | 155 ++++++++ tests/unit/test_catalog_reader.py | 63 ++++ tests/unit/test_catalog_service.py | 204 +++++++++++ tests/unit/test_catalog_storage.py | 29 ++ tests/unit/test_catalog_task.py | 13 + tests/unit/test_import_service.py | 7 +- 45 files changed, 2535 insertions(+), 35 deletions(-) create mode 100644 docs/development/LOCAL_CATALOG_V2.md create mode 100644 src/pullbox/api/v1/catalog.py create mode 100644 src/pullbox/services/catalog/__init__.py create mode 100644 src/pullbox/services/catalog/contract.py create mode 100644 src/pullbox/services/catalog/database.py create mode 100644 src/pullbox/services/catalog/lookup.py create mode 100644 src/pullbox/services/catalog/reader.py create mode 100644 src/pullbox/services/catalog/retention.py create mode 100644 src/pullbox/services/catalog/service.py create mode 100644 src/pullbox/services/catalog/storage.py create mode 100644 src/pullbox/tasks/catalog_task.py create mode 100644 src/pullbox/ui/templates/partials/settings_catalog.html create mode 100644 tests/catalog_fixtures.py create mode 100644 tests/ui/test_catalog_controls.py create mode 100644 tests/unit/test_catalog_contract.py create mode 100644 tests/unit/test_catalog_controls.py create mode 100644 tests/unit/test_catalog_database.py create mode 100644 tests/unit/test_catalog_integration.py create mode 100644 tests/unit/test_catalog_reader.py create mode 100644 tests/unit/test_catalog_service.py create mode 100644 tests/unit/test_catalog_storage.py create mode 100644 tests/unit/test_catalog_task.py diff --git a/.gitignore b/.gitignore index 4e4fbd22..ab7dc7ba 100644 --- a/.gitignore +++ b/.gitignore @@ -113,6 +113,7 @@ node_modules/ # Playwright test artifacts test-results/ .playwright-cli/ +output/playwright/ # Local security artifacts bandit-report.json diff --git a/CHANGELOG.md b/CHANGELOG.md index be4344f0..b8948cfa 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -9,6 +9,11 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 ### Added +- Added an optional local Comic Vine catalog for series searches and import + matching without live metadata requests. Download it from Metadata settings + with progress, resumable transfers, signed verification, and daily updates. + Failed updates preserve the installed catalog; full metadata enrichment still + uses the user's Comic Vine API key. - Added import source-layout previews for series folders, publisher/series folders, and custom folder and issue naming patterns. - Added independent options to keep existing files in place and use an approved diff --git a/docs/development/LOCAL_CATALOG_V2.md b/docs/development/LOCAL_CATALOG_V2.md new file mode 100644 index 00000000..d14e507a --- /dev/null +++ b/docs/development/LOCAL_CATALOG_V2.md @@ -0,0 +1,138 @@ +# Local Comic Vine catalog v2 client + +## Scope + +The catalog is an optional, derived SQLite database on the **Pullbox server**, +separate from the user's library database. Settings → Metadata → Download catalog +starts the first download. Browsers display status and progress; they do not +download or unzip the database onto the browser device. + +The Pullbox Data API serves only signed publications and catalog artifacts. This +feature adds no remote search or individual metadata endpoints. ComicRack/v1, +MEGA, arbitrary file imports, and user-configurable signing keys are not supported. + +After installation: + +- Add Series, automatic import matching, manual match search, orphan recovery, + explicit Comic Vine ID lookup, and basic issue-list hydration use the catalog. +- A local miss does not silently fall back to online requests. The user can check + for a catalog update; existing review and override decisions remain intact. +- Without an installed catalog, existing Comic Vine discovery behavior remains. +- Full metadata refresh and post-import ComicInfo enrichment retain their direct + Comic Vine provider path, including its batch cache and rate controls. They + still need the user's API key. Basic local matching does not need that key. +- Catalog installation changes no library rows or files. Imports still require + the existing review/confirmation steps before Step 4 materializes files. +- New basic records use `metadata_source=pullbox_catalog`. Series freshness uses + the publication's source cutoff, not download time. A completed full Comic Vine + series refresh is not demoted or overwritten by basic catalog hydration. + Existing live Comic Vine issue records likewise retain their identity and + fields; basic hydration adds missing issues without reverting live metadata. + +## Download and activation contract + +`services/catalog/contract.py` owns the Ed25519 public verification key ring. +The release currently trusts `catalog-2026-09`, verified against the publisher's +public key and production publication. A future signing-key rotation must ship +the next public key in a Pullbox release before the publisher switches. Unknown +keys fail closed. No signing secret is distributed with Pullbox. + +1. Fetch `/api/v2/catalog/latest` from `PULLBOX_DATA_API_BASE_URL` (defaults to + `https://api.pullbox.app`), using the cached ETag when available. +2. Verify the signature over compact sorted JSON, payload SHA-256, schema, + versions, lineage, sizes, and exact same-API download coordinates. Reverify + the cached signed envelope after a 304. Do not follow redirects. +3. Stream the required snapshot or cumulative patch into a checksum-named + `.part` file. Interrupted transfers resume with Range/If-Range. A server that + returns 200 starts a fresh transfer; 206 must match the expected byte range. +4. Verify compressed byte count and SHA-256 **before decompression**. Decompress + with a bounded window and output ceiling, checking free space as it expands. +5. Validate the SQLite application/user versions, dataset identity, source + cutoff, row counts, logical content hash, foreign keys, integrity, and FTS + index. Reject unsupported tables, views, or triggers. +6. Reconstruct each daily cumulative patch from a **fresh copy of its immutable + weekly base**, never from yesterday's patched database. Apply child-first + deletes and parent-first upserts transactionally, replace the dataset + manifest, rebuild FTS, and verify the target logical hash. +7. Flush and atomically move the validated generation into place, then atomically + replace `active.json`. Retain the previous reference and file. Readers open a + read-only immutable generation for each query; long disk work is offloaded + from the event loop. Cancellation waits for an owned disk operation to finish + before releasing the update lock. + +When a new weekly base is published, it is downloaded in full. Between weekly +bases, only the latest cumulative patch is needed. If no current patch exists, +the signed publication's full snapshot is installed. A healthy newer local +version is never downgraded by a stale API response. + +Limits: manifest 1 MiB, compressed artifact 2 GiB, decompressed artifact 4 GiB, +zstd window 128 MiB; search input 256 characters / 16 words; up to 1,000 results. +Disk-space checks reserve a safety margin, and patching checks space for copying +the weekly base. The byte-progress bar applies to transfer; unpacking, +verification, and installation are separate indeterminate stages. + +## Persistent storage and recovery + +Under `/catalog/`: + +```text +active.json installed generation and source cutoff +previous.json last installed generation before a successful update +state.json opt-in, update preference, coarse status, timestamps +manifest.json signed publication cache and ETag +update.lock cross-process advisory update lock +bases/.db immutable weekly snapshots +versions/.db validated reconstructed daily generations +downloads/.part resumable compressed transfer +staging/catalog-*.db uncommitted work, never used for searches +``` + +An async lock and filesystem lock prevent simultaneous installers. Failed +downloads or patches leave the installed catalog active. Retry resumes a matching +partial transfer. Checksum failures discard the bad transfer. Retrying repairs a +missing or invalid file at the currently published version. Unknown signing keys +or formats require a client update, not bypassing verification. An unreadable +active catalog reports an actionable error rather than making live metadata +requests. Symlinked catalog storage is rejected. + +Cleanup runs under the update lock: abandon unfinished staging files, remove +compressed downloads after success, and expire unreferenced generations/partials +older than two days. Keep active, previous, and both referenced weekly bases. +Only recognized catalog-owned filenames are eligible; no library paths are used. +The grace period accommodates readers that started before activation. + +## Scheduling and local control API + +`catalog_update` appears as **Local Catalog Update** in the normal task system. +It checks daily at 06:30 in the scheduler timezone, with up to 30 minutes of jitter. +Automatic runs do nothing until the user requests the first download or when +automatic updates are disabled. A startup check catches overdue work; a 23-hour +freshness guard allows the next day's jitter to be earlier than yesterday's. +Manual checks bypass the age guard. Startup checks use the same installer lock +and status but do not create a separate scheduler history entry. + +| Local route | Access | Result | +|---|---|---| +| `GET /api/v1/catalog` | Authenticated | Safe status, version, cutoff, byte progress, timestamps, error | +| `POST /api/v1/catalog/sync` | Interactive operator + CSRF | 202 queued/already queued/already running; 503 if unavailable | +| `PATCH /api/v1/catalog/preferences` | Interactive operator + CSRF | `{ "automatic_updates": true/false }`; updated status | + +These routes control this Pullbox instance; they are not new Pullbox Data API +distribution routes. Machine API keys cannot trigger downloads or change the +preference. Settings polls only while the page is mounted and live updates are +enabled; a download continues when the page is closed. Errors never include +credentials or raw response bodies. + +## Verification + +`tests/unit/test_catalog_*` covers signed-manifest failures, hash and lineage +checks, corrupted/unsupported SQLite artifacts, cumulative reversion, interrupted +transfers, installation failure preservation, overlap, repair, cleanup, realistic +zstd windows, search escaping, exact issue IDs and fractions, source provenance, +and live enrichment separation. `tests/ui/test_catalog_controls.py` covers +settings, no-key local search, session/CSRF and machine-key boundaries. + +The external-artifact SQLite adapter is deliberately separate from the ORM +application database. Its SQL identifiers come exclusively from a closed contract +allowlist; all search terms, IDs and filters are bound parameters. Existing +application database access remains SQLAlchemy-based. diff --git a/pyproject.toml b/pyproject.toml index 0c89c99a..a1f74e8f 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -46,6 +46,7 @@ dependencies = [ "pillow~=12.3", "pdf2image~=1.17", "tzlocal~=5.3", + "zstandard>=0.25,<1.0", ] [project.optional-dependencies] diff --git a/src/pullbox/api/v1/catalog.py b/src/pullbox/api/v1/catalog.py new file mode 100644 index 00000000..7f23c501 --- /dev/null +++ b/src/pullbox/api/v1/catalog.py @@ -0,0 +1,41 @@ +"""Instance-local catalog controls. No metadata proxy routes.""" + +from fastapi import APIRouter, HTTPException +from pydantic import BaseModel + +from pullbox.api.deps import AuthenticatedUser, InteractiveOperatorUser +from pullbox.services.catalog.service import CatalogStatus + +router = APIRouter(prefix="/catalog", tags=["catalog"]) + + +@router.get("") +async def catalog_status(user: AuthenticatedUser) -> CatalogStatus: + from pullbox.services.catalog.service import get_catalog_service + + return get_catalog_service().status() + + +@router.post("/sync", status_code=202) +async def sync_catalog(user: InteractiveOperatorUser) -> dict[str, str]: + from pullbox.core.scheduler import get_scheduler + + status = get_scheduler().run_task_now("catalog_update") + if status is None: + raise HTTPException(503, "The catalog task is unavailable. Restart Pullbox and retry.") + return {"status": status} + + +class CatalogPreferences(BaseModel): + automatic_updates: bool + + +@router.patch("/preferences") +async def catalog_preferences( + body: CatalogPreferences, user: InteractiveOperatorUser +) -> CatalogStatus: + from pullbox.services.catalog.service import get_catalog_service + + service = get_catalog_service() + await service.set_automatic_updates(body.automatic_updates) + return service.status() diff --git a/src/pullbox/api/v1/router.py b/src/pullbox/api/v1/router.py index 8761b951..4aae8f42 100644 --- a/src/pullbox/api/v1/router.py +++ b/src/pullbox/api/v1/router.py @@ -6,6 +6,7 @@ from pullbox.api.v1.audit import router as audit_router from pullbox.api.v1.auth import router as auth_router from pullbox.api.v1.blocklist import router as blocklist_router +from pullbox.api.v1.catalog import router as catalog_router from pullbox.api.v1.clients import router as clients_router from pullbox.api.v1.config import router as config_router from pullbox.api.v1.covers import router as covers_router @@ -36,6 +37,7 @@ v1_router = APIRouter(prefix="/api/v1") v1_router.include_router(activity_router) +v1_router.include_router(catalog_router) v1_router.include_router(audit_router) v1_router.include_router(blocklist_router) v1_router.include_router(auth_router) diff --git a/src/pullbox/app.py b/src/pullbox/app.py index 29a91c76..877f58d5 100644 --- a/src/pullbox/app.py +++ b/src/pullbox/app.py @@ -808,6 +808,19 @@ async def _startup_update_check() -> None: exc_info=True, ) + async def _startup_catalog_check() -> None: + from pullbox.services.catalog.contract import CatalogError + from pullbox.services.catalog.service import get_catalog_service + + try: + await get_catalog_service().sync() + except CatalogError: + logger.warning("startup_catalog_check_failed") + + catalog_startup_task = asyncio.create_task(_startup_catalog_check()) + _startup_background_tasks.add(catalog_startup_task) + catalog_startup_task.add_done_callback(_startup_background_tasks.discard) + if settings.startup_update_check_enabled: startup_update_task = asyncio.create_task(_startup_update_check()) _startup_background_tasks.add(startup_update_task) diff --git a/src/pullbox/composition/services.py b/src/pullbox/composition/services.py index 22f57019..d37f5213 100644 --- a/src/pullbox/composition/services.py +++ b/src/pullbox/composition/services.py @@ -206,6 +206,8 @@ def _datanodes_login_failure(exc: Exception) -> ArtifactHostResolutionError: async def build_metadata_service(session: AsyncSession) -> MetadataService: """Construct a MetadataService using persisted ComicVine settings.""" + from pullbox.services.catalog.reader import get_catalog_reader + settings = get_settings() api_key = await get_comicvine_api_key(session) provider = ComicVineProvider(api_key=api_key) @@ -215,6 +217,7 @@ async def build_metadata_service(session: AsyncSession) -> MetadataService: provider=provider, covers_dir=covers_dir, refresh_days=settings.metadata_refresh_days, + catalog=get_catalog_reader(), ) @@ -239,6 +242,8 @@ async def build_import_service( min_burst_limit: int | None = None, ) -> ImportService: """Construct an ImportService using persisted ComicVine settings.""" + from pullbox.services.catalog.reader import get_catalog_reader + settings = get_settings() api_key = await get_comicvine_api_key(session) persisted_rate_config = await session.get(SystemConfig, "comicvine_rate_limit_per_second") @@ -271,6 +276,7 @@ async def build_import_service( provider, covers_dir=await resolve_covers_dir(session), refresh_days=settings.metadata_refresh_days, + catalog=get_catalog_reader(), ) event_bus = build_scoped_event_bus() series_svc = SeriesService(metadata_svc, event_bus) diff --git a/src/pullbox/services/catalog/__init__.py b/src/pullbox/services/catalog/__init__.py new file mode 100644 index 00000000..2742c6ce --- /dev/null +++ b/src/pullbox/services/catalog/__init__.py @@ -0,0 +1 @@ +"""Verified local Comic Vine catalog downloads and queries.""" diff --git a/src/pullbox/services/catalog/contract.py b/src/pullbox/services/catalog/contract.py new file mode 100644 index 00000000..de358742 --- /dev/null +++ b/src/pullbox/services/catalog/contract.py @@ -0,0 +1,165 @@ +"""Verify the publisher before accepting any download coordinates.""" + +from __future__ import annotations + +import base64 +import binascii +import hashlib +import json +import re +from dataclasses import dataclass +from datetime import datetime +from typing import TYPE_CHECKING, Any + +from cryptography.exceptions import InvalidSignature +from cryptography.hazmat.primitives.asymmetric.ed25519 import Ed25519PublicKey + +if TYPE_CHECKING: + from collections.abc import Mapping + +MAX_MANIFEST_BYTES = 1024 * 1024 +MAX_ARTIFACT_BYTES = 2 * 1024**3 +VERSION = re.compile(r"\d{8}T\d{6}Z") +SHA256 = re.compile(r"[0-9a-f]{64}") + +# Public verification keys only. Ship a new key before the publisher rotates to it. +TRUSTED_KEYS = { + "catalog-2026-09": Ed25519PublicKey.from_public_bytes( + base64.b64decode("SwsNUSLDGU8PvkvImuLSRopRgKMdQM2Odr2JA5bA+tg=") + ), +} + + +class CatalogError(ValueError): + """A safe, actionable catalog failure.""" + + +@dataclass(frozen=True) +class Artifact: + version: str + download_path: str + sha256: str + size_bytes: int + base_version: str | None = None + + +@dataclass(frozen=True) +class Publication: + latest_version: str + full_snapshot: Artifact + snapshots: tuple[Artifact, ...] + patches: tuple[Artifact, ...] + + def latest_patch(self) -> Artifact | None: + return next( + ( + patch + for patch in self.patches + if patch.version == self.latest_version + and patch.base_version == self.full_snapshot.version + ), + None, + ) + + +def valid_version(value: object) -> str: + if not isinstance(value, str) or VERSION.fullmatch(value) is None: + raise CatalogError("Catalog version is not supported. Update Pullbox and try again.") + try: + datetime.strptime(value, "%Y%m%dT%H%M%SZ") + except ValueError as exc: + raise CatalogError("Catalog version is invalid.") from exc + return value + + +def _object(value: object) -> dict[str, Any]: + if not isinstance(value, dict): + raise CatalogError("Catalog manifest is invalid.") + return value + + +def _artifact(value: object, *, patch: bool = False) -> Artifact: + item = _object(value) + version = valid_version(item.get("target_version" if patch else "version")) + base = valid_version(item.get("base_version")) if patch else None + path = ( + f"/api/v2/catalog/patches/{base}/{version}" + if patch + else f"/api/v2/catalog/snapshots/{version}" + ) + checksum, size = item.get("sha256"), item.get("size_bytes") + if ( + item.get("download_path") != path + or not isinstance(checksum, str) + or SHA256.fullmatch(checksum) is None + or type(size) is not int + or not 0 < size <= MAX_ARTIFACT_BYTES + or (base is not None and base >= version) + ): + raise CatalogError("Catalog artifact information is invalid.") + return Artifact(version, path, checksum, size, base) + + +def verify_manifest( + raw: bytes, + keys: Mapping[str, Ed25519PublicKey] = TRUSTED_KEYS, +) -> Publication: + """Parse only bounded JSON and verify its canonical signed payload first.""" + if len(raw) > MAX_MANIFEST_BYTES: + raise CatalogError("Catalog manifest is too large.") + try: + document = _object(json.loads(raw)) + payload, signature = _object(document.get("payload")), _object(document.get("signature")) + key_id = signature.get("key_id") + if not isinstance(key_id, str) or key_id not in keys: + raise CatalogError("Unknown catalog signing key. Update Pullbox and try again.") + canonical = json.dumps( + payload, sort_keys=True, separators=(",", ":"), allow_nan=False + ).encode() + if ( + signature.get("algorithm") != "Ed25519" + or signature.get("payload_sha256") != hashlib.sha256(canonical).hexdigest() + ): + raise CatalogError("Catalog signature verification failed.") + encoded = signature.get("value") + if not isinstance(encoded, str): + raise CatalogError("Catalog signature verification failed.") + keys[key_id].verify( + base64.b64decode(encoded + "=" * (-len(encoded) % 4), altchars=b"-_", validate=True), + canonical, + ) + if ( + payload.get("format_id") != "pullbox-catalog-v2-publication" + or payload.get("schema_version") != "1" + ): + raise CatalogError("Catalog format is not supported. Update Pullbox and try again.") + full = _artifact(payload.get("full_snapshot")) + snapshots_raw, patches_raw = payload.get("snapshots", []), payload.get("patches", []) + if not isinstance(snapshots_raw, list) or not isinstance(patches_raw, list): + raise CatalogError("Catalog manifest is invalid.") + snapshots = tuple(_artifact(item) for item in snapshots_raw) + patches = tuple(_artifact(item, patch=True) for item in patches_raw) + bases = {full.version, *(entry.version for entry in snapshots)} + coordinates = [item.download_path for item in (full, *snapshots, *patches)] + latest = valid_version(payload.get("latest_version")) + if ( + len(coordinates) != len(set(coordinates)) + or any(item.version >= full.version for item in snapshots) + or any(item.base_version not in bases for item in patches) + or any( + item.base_version != full.version and item.version >= full.version + for item in patches + ) + or latest != max([full.version, *(item.version for item in patches)]) + ): + raise CatalogError("Catalog update lineage is invalid.") + result = Publication(latest, full, snapshots, patches) + if latest != full.version and result.latest_patch() is None: + raise CatalogError("Catalog latest update is unavailable.") + return result + except (InvalidSignature, binascii.Error) as exc: + raise CatalogError("Catalog signature verification failed.") from exc + except CatalogError: + raise + except (ValueError, TypeError, KeyError, RecursionError) as exc: + raise CatalogError("Catalog manifest is invalid.") from exc diff --git a/src/pullbox/services/catalog/database.py b/src/pullbox/services/catalog/database.py new file mode 100644 index 00000000..e15250f3 --- /dev/null +++ b/src/pullbox/services/catalog/database.py @@ -0,0 +1,235 @@ +"""Validate signed SQLite artifacts using the producer's fixed schema and hash. + +This database is an external read-only artifact, separate from the ORM application +database. SQL identifiers below are contract constants, never user input. +""" + +from __future__ import annotations + +import hashlib +import json +import os +import shutil +import sqlite3 +from contextlib import closing +from datetime import datetime +from typing import TYPE_CHECKING, Any + +from pullbox.services.catalog.contract import CatalogError, valid_version + +if TYPE_CHECKING: + from pathlib import Path + +TABLES = { + "publishers": ("id,name", "id", "publisher"), + "series": ("id,name,start_year,publisher_id,issue_count,cover_url", "id", "series"), + "series_aliases": ( + "series_id,position,alias,normalized_alias", + "series_id,position", + "series_alias", + ), + "issues": ( + "id,series_id,issue_number,normalized_issue_number,sort_number,title,cover_date,store_date,cover_url", + "id", + "issue", + ), +} +FTS_TABLES = { + "series_fts", + *(f"series_fts_{suffix}" for suffix in ("config", "content", "data", "docsize", "idx")), +} +REQUIRED_TABLES = {*TABLES, "series_fts", "dataset_manifest", "schema_migrations"} +# Closed SQL allowlist, generated only from contract constants. Request input +# never becomes an identifier, including in the external artifact database. +CONTENT_QUERIES = { + table: f"SELECT {columns} FROM {table} ORDER BY {order}" + for table, (columns, order, _) in TABLES.items() +} +COUNT_QUERIES = {table: f"SELECT COUNT(*) FROM {table}" for table in TABLES} +PATCH_DELETE_QUERIES = { + table: f"SELECT {keys} FROM {prefix}_deletes" for table, (_, keys, prefix) in TABLES.items() +} +PATCH_UPSERT_QUERIES = { + table: f"SELECT {columns} FROM {prefix}_upserts" + for table, (columns, _, prefix) in TABLES.items() +} +MANIFEST_QUERIES = { + "dataset_manifest": "SELECT key,value FROM dataset_manifest", + "patch_manifest": "SELECT key,value FROM patch_manifest", + "target_dataset_manifest": "SELECT key,value FROM target_dataset_manifest", +} +REBUILD_FTS = """ +INSERT INTO series_fts(rowid,series_id,name,aliases) +SELECT s.id,s.id,s.name,COALESCE(GROUP_CONCAT(a.alias,' '),'') +FROM series s LEFT JOIN series_aliases a ON a.series_id=s.id +GROUP BY s.id,s.name ORDER BY s.id +""" + + +def safe_path(path: Path) -> Path: + """Reject symlink components in app-owned catalog paths.""" + if any(part.is_symlink() for part in (path, *path.parents)): + raise CatalogError("Catalog storage contains a symbolic link. Check the data volume.") + return path + + +def open_readonly(path: Path) -> sqlite3.Connection: + safe_path(path) + db = sqlite3.connect(path.absolute().as_uri() + "?mode=ro&immutable=1", uri=True) + db.execute("PRAGMA query_only=ON") + db.execute("PRAGMA trusted_schema=OFF") + db.execute("PRAGMA cache_size=-8192") + return db + + +def file_sha256(path: Path) -> str: + safe_path(path) + with path.open("rb") as stream: + return hashlib.file_digest(stream, "sha256").hexdigest() + + +def read_manifest(db: sqlite3.Connection, table: str = "dataset_manifest") -> dict[str, Any]: + if table not in MANIFEST_QUERIES: + raise CatalogError("Catalog manifest table is invalid.") + return {str(k): json.loads(v) for k, v in db.execute(MANIFEST_QUERIES[table])} + + +def logical_hash(db: sqlite3.Connection) -> str: + digest = hashlib.sha256() + for table in TABLES: + for row in db.execute(CONTENT_QUERIES[table]): + digest.update(table.encode() + b"\0") + digest.update( + json.dumps(row, ensure_ascii=False, separators=(",", ":")).encode() + b"\n" + ) + return digest.hexdigest() + + +def _check_database(db: sqlite3.Connection, application_id: int, allowed: set[str]) -> None: + if ( + db.execute("PRAGMA application_id").fetchone()[0] != application_id + or db.execute("PRAGMA user_version").fetchone()[0] != 1 + ): + raise CatalogError("Catalog database format is not supported. Update Pullbox.") + if db.execute("PRAGMA integrity_check").fetchone()[0] != "ok": + raise CatalogError("Catalog database integrity check failed. Retry the download.") + objects = db.execute( + "SELECT name,type FROM sqlite_master WHERE type IN ('table','view','trigger')" + ).fetchall() + if any( + kind != "table" or name not in allowed + for name, kind in objects + if not name.startswith("sqlite_") + ): + raise CatalogError("Catalog contains unsupported database objects.") + + +def validate_snapshot(path: Path, version: str) -> dict[str, Any]: + """Recompute logical content before allowing a new catalog to become active.""" + valid_version(version) + try: + with closing(open_readonly(path)) as db: + _check_database(db, 0x50424332, REQUIRED_TABLES | FTS_TABLES) + if db.execute("PRAGMA foreign_key_check").fetchone() is not None: + raise CatalogError("Catalog has invalid record relationships.") + manifest = read_manifest(db) + for key, value in { + "format_id": "pullbox-catalog-v2", + "schema_version": "1", + "compatibility_version": "pullbox-only", + "application_id": 0x50424332, + "dataset_version": version, + }.items(): + if manifest.get(key) != value: + raise CatalogError("Catalog identity does not match its signed publication.") + cutoff = datetime.fromisoformat(str(manifest.get("source_cutoff_at", ""))) + if cutoff.tzinfo is None or cutoff.strftime("%Y%m%dT%H%M%SZ") != version: + raise CatalogError("Catalog source timestamp is invalid.") + counts = {table: db.execute(COUNT_QUERIES[table]).fetchone()[0] for table in TABLES} + if counts != manifest.get("counts") or logical_hash(db) != manifest.get( + "content_sha256" + ): + raise CatalogError("Catalog content checksum failed. Retry the download.") + if db.execute("SELECT COUNT(*) FROM series_fts").fetchone()[0] != counts["series"]: + raise CatalogError("Catalog search index is incomplete.") + if db.execute( + "SELECT 1 FROM series s LEFT JOIN series_fts f ON f.rowid=s.id " + "WHERE f.series_id IS NULL OR f.series_id != s.id OR f.name != s.name LIMIT 1" + ).fetchone(): + raise CatalogError("Catalog search index is inconsistent.") + db.execute("SELECT version,applied_at FROM schema_migrations LIMIT 1").fetchall() + return manifest + except (sqlite3.Error, OSError, json.JSONDecodeError, TypeError, ValueError) as exc: + if isinstance(exc, CatalogError): + raise + raise CatalogError("Catalog database could not be validated. Retry the download.") from exc + + +def apply_catalog_patch( + base: Path, patch: Path, output: Path, base_version: str, target_version: str +) -> dict[str, Any]: + """Always reconstruct from the immutable weekly base, including reverted rows.""" + validate_snapshot(base, base_version) + safe_path(output) + try: + with closing(open_readonly(patch)) as changes: + allowed = {"patch_manifest", "target_dataset_manifest"} + allowed.update( + f"{prefix}_{kind}" + for _, _, prefix in TABLES.values() + for kind in ("upserts", "deletes") + ) + _check_database(changes, 0x50425044, allowed) + meta = read_manifest(changes, "patch_manifest") + expected = { + "format_id": "pullbox-catalog-v2-patch", + "schema_version": "1", + "base_version": base_version, + "target_version": target_version, + "base_snapshot_sha256": file_sha256(base), + } + if any(meta.get(k) != v for k, v in expected.items()): + raise CatalogError("Catalog patch does not match the retained weekly base.") + with base.open("rb") as source, output.open("xb") as destination: + shutil.copyfileobj(source, destination, 1024 * 1024) + with closing(sqlite3.connect(output)) as db: + db.execute("PRAGMA trusted_schema=OFF") + db.execute("PRAGMA foreign_keys=ON") + db.execute("PRAGMA journal_mode=DELETE") + db.execute("PRAGMA synchronous=FULL") + with db: + db.execute("PRAGMA defer_foreign_keys=ON") + for table, (_, keys, _prefix) in reversed(TABLES.items()): + where = " AND ".join(f"{key}=?" for key in keys.split(",")) + db.executemany( + f"DELETE FROM {table} WHERE {where}", + changes.execute(PATCH_DELETE_QUERIES[table]), + ) + for table, (columns, keys, _prefix) in TABLES.items(): + cols = columns.split(",") + assignments = ",".join( + f"{column}=excluded.{column}" + for column in cols + if column not in keys.split(",") + ) + db.executemany( + f"INSERT INTO {table} ({columns}) " + f"VALUES ({','.join('?' for _ in cols)}) " + f"ON CONFLICT ({keys}) DO UPDATE SET {assignments}", + changes.execute(PATCH_UPSERT_QUERIES[table]), + ) + db.execute("DELETE FROM dataset_manifest") + db.executemany( + "INSERT INTO dataset_manifest VALUES (?,?)", + changes.execute("SELECT key,value FROM target_dataset_manifest"), + ) + db.execute("DELETE FROM series_fts") + db.execute(REBUILD_FTS) + result = validate_snapshot(output, target_version) + if result["content_sha256"] != meta.get("target_content_sha256"): + raise CatalogError("Catalog patch target checksum failed.") + with output.open("rb") as stream: + os.fsync(stream.fileno()) + return result + except (sqlite3.Error, OSError, json.JSONDecodeError) as exc: + raise CatalogError("Catalog patch could not be applied. Retry the update.") from exc diff --git a/src/pullbox/services/catalog/lookup.py b/src/pullbox/services/catalog/lookup.py new file mode 100644 index 00000000..0607ac59 --- /dev/null +++ b/src/pullbox/services/catalog/lookup.py @@ -0,0 +1,103 @@ +"""Catalog discovery adapter for existing search and import matching contracts.""" + +from __future__ import annotations + +from typing import TYPE_CHECKING, Any + +from pullbox.services.catalog.contract import CatalogError + +if TYPE_CHECKING: + from pullbox.providers.base import ( + IssueMetadata, + IssueSummary, + SeriesMetadata, + SeriesSearchResult, + ) + from pullbox.services.catalog.reader import CatalogReader + + +class CatalogLookupService: + """Basic catalog discovery only; full metadata refresh still owns its provider.""" + + is_local_catalog = True + + def __init__(self, reader: CatalogReader) -> None: + self.reader = reader + + async def search_series( + self, + query: str, + year: int | None = None, + *, + limit: int = 1000, + offset: int = 0, + suppress_errors: bool = False, + ) -> list[SeriesSearchResult]: + return await self.reader.search(query, year, limit, offset) + + async def search_series_page( + self, + query: str, + year: int | None = None, + *, + limit: int = 100, + offset: int = 0, + suppress_errors: bool = False, + ) -> tuple[list[SeriesSearchResult], int]: + rows = await self.reader.search(query, year, limit, offset) + return rows, len(rows) + + async def search_series_globally( + self, + query: str, + *, + max_results: int = 1000, + page_size: int = 100, + suppress_errors: bool = False, + ) -> tuple[list[SeriesSearchResult], int]: + rows = await self.reader.search(query, limit=max_results) + return rows, len(rows) + + async def get_series(self, series_provider_id: str) -> SeriesMetadata: + result = await self.reader.series(int(series_provider_id)) + if result is None: + raise CatalogError( + "This series is not in the local catalog. Check for a catalog update." + ) + return result + + async def get_series_cached(self, series_provider_id: str) -> SeriesMetadata | None: + return await self.reader.series(int(series_provider_id)) + + async def get_issues_for_series(self, series_provider_id: str) -> list[IssueSummary]: + await self.get_series(series_provider_id) + return await self.reader.issues(int(series_provider_id)) + + async def get_issues_for_series_by_numbers( + self, series_provider_id: str, issue_numbers: list[float] + ) -> list[IssueSummary]: + numbers = set(issue_numbers) + return [ + issue + for issue in await self.get_issues_for_series(series_provider_id) + if issue.issue_number in numbers + ] + + async def get_issue(self, issue_provider_id: str) -> IssueMetadata: + result = await self.reader.issue(int(issue_provider_id)) + if result is None: + raise CatalogError( + "This issue is not in the local catalog. Check for a catalog update." + ) + return result + + async def close(self) -> None: + """Queries own and close their connections individually.""" + + +def catalog_or_provider(provider: Any) -> Any: + """Select the local discovery source without altering full refresh providers.""" + from pullbox.services.catalog.reader import get_catalog_reader + + reader = get_catalog_reader() + return CatalogLookupService(reader) if reader.available else provider diff --git a/src/pullbox/services/catalog/reader.py b/src/pullbox/services/catalog/reader.py new file mode 100644 index 00000000..304a9c70 --- /dev/null +++ b/src/pullbox/services/catalog/reader.py @@ -0,0 +1,193 @@ +"""Bounded SQLite reads from an immutable catalog generation.""" + +from __future__ import annotations + +import re +import sqlite3 +import threading +from contextlib import closing +from dataclasses import dataclass +from datetime import datetime +from functools import lru_cache +from typing import TYPE_CHECKING, Any + +import structlog + +from pullbox.config import get_settings +from pullbox.core.issue_numbers import parse_issue_number_text +from pullbox.core.naming import detect_issue_type_from_metadata_title +from pullbox.providers.base import IssueMetadata, IssueSummary, SeriesMetadata, SeriesSearchResult +from pullbox.services.catalog.contract import CatalogError, valid_version +from pullbox.services.catalog.database import open_readonly, safe_path, validate_snapshot +from pullbox.services.catalog.storage import disk_work, load_json + +if TYPE_CHECKING: + from pathlib import Path + +logger = structlog.get_logger(__name__) +SERIES_COLUMNS = "s.id,s.name,s.start_year,p.name,s.issue_count,s.cover_url" +SERIES_FROM = "series s LEFT JOIN publishers p ON p.id=s.publisher_id" +ISSUE_COLUMNS = ( + "id,series_id,issue_number,normalized_issue_number,sort_number," + "title,cover_date,store_date,cover_url" +) + + +@dataclass(frozen=True) +class CatalogSeriesMetadata(SeriesMetadata): + source_cutoff_at: datetime | None = None + + +@dataclass(frozen=True) +class CatalogIssueSummary(IssueSummary): + source_cutoff_at: datetime | None = None + + +class CatalogReader: + """Pin a file per query; verify a generation once before serving its rows.""" + + def __init__(self, root: Path) -> None: + self.root = root + self._validated: set[tuple[str, int, int, int]] = set() + self._validation_lock = threading.Lock() + + @property + def available(self) -> bool: + return (self.root / "active.json").exists() + + def _generation(self) -> tuple[Path, datetime]: + reference = load_json(self.root / "active.json") + version = valid_version(reference.get("version")) + relative = reference.get("path") + if relative not in {f"bases/{version}.db", f"versions/{version}.db"}: + raise CatalogError("Catalog active file is invalid. Retry the catalog update.") + path = safe_path(self.root / str(relative)) + stat = path.stat() + identity = (str(path), stat.st_ino, stat.st_mtime_ns, stat.st_size) + with self._validation_lock: + if identity not in self._validated: + validate_snapshot(path, version) + if len(self._validated) >= 8: + self._validated.clear() + self._validated.add(identity) + return path, datetime.fromisoformat(str(reference["source_cutoff_at"])) + + def _query(self, sql: str, params: tuple[object, ...]) -> tuple[list[Any], datetime]: + try: + path, cutoff = self._generation() + with closing(open_readonly(path)) as db: + return db.execute(sql, params).fetchall(), cutoff + except (sqlite3.Error, OSError, KeyError) as exc: + raise CatalogError( + "The local catalog could not be read. Retry its update in Metadata settings." + ) from exc + + async def search( + self, query: str, year: int | None = None, limit: int = 1000, offset: int = 0 + ) -> list[SeriesSearchResult]: + terms = re.findall(r"\w+", query[:256], flags=re.UNICODE)[:16] + if not terms: + return [] + expression = " AND ".join(f'"{term}"*' for term in terms) + rows, _ = await disk_work( + self._query, + f"SELECT {SERIES_COLUMNS} FROM series_fts f JOIN series s ON s.id=f.rowid " + "LEFT JOIN publishers p ON p.id=s.publisher_id " + "WHERE series_fts MATCH ? AND (? IS NULL OR s.start_year=?) " + "ORDER BY rank,s.id LIMIT ? OFFSET ?", + (expression, year, year, max(1, min(limit, 1000)), max(0, min(offset, 10000))), + ) + return [ + SeriesSearchResult( + str(r[0]), + r[1], + r[2], + r[3], + r[4], + None, + r[5], + None, + f"https://comicvine.gamespot.com/volume/4050-{r[0]}/", + ) + for r in rows + ] + + async def series(self, series_id: int) -> CatalogSeriesMetadata | None: + rows, cutoff = await disk_work( + self._query, f"SELECT {SERIES_COLUMNS} FROM {SERIES_FROM} WHERE s.id=?", (series_id,) + ) + if not rows: + return None + r = rows[0] + return CatalogSeriesMetadata( + str(r[0]), + r[1], + r[1], + r[2], + None, + None, + r[3], + None, + r[5], + r[4], + f"https://comicvine.gamespot.com/volume/4050-{r[0]}/", + cutoff, + ) + + async def issues(self, series_id: int) -> list[IssueSummary]: + rows, cutoff = await disk_work( + self._query, + f"SELECT {ISSUE_COLUMNS} FROM issues WHERE series_id=? " + "ORDER BY CAST(sort_number AS REAL),issue_number,id", + (series_id,), + ) + return [self._summary(row, cutoff) for row in rows] + + @staticmethod + def _summary(row: Any, cutoff: datetime) -> CatalogIssueSummary: + try: + number, exact = parse_issue_number_text(str(row[2] or row[3] or "0")) + except ValueError: + number, exact = 0.0, None + return CatalogIssueSummary( + str(row[0]), + number, + row[5], + row[6], + row[8], + detect_issue_type_from_metadata_title(row[5] or ""), + exact, + cutoff, + ) + + async def issue(self, issue_id: int) -> IssueMetadata | None: + """Return only the basic identity fields used during import file matching.""" + rows, cutoff = await disk_work( + self._query, f"SELECT {ISSUE_COLUMNS} FROM issues WHERE id=?", (issue_id,) + ) + if not rows: + return None + row = rows[0] + summary = self._summary(row, cutoff) + return IssueMetadata( + summary.provider_id, + str(row[1]), + summary.issue_number, + summary.title, + None, + summary.release_date, + row[7], + summary.cover_url, + None, + f"https://comicvine.gamespot.com/issue/4000-{row[0]}/", + issue_number_text=summary.issue_number_text, + ) + + +@lru_cache(maxsize=4) +def _reader(root: Path) -> CatalogReader: + return CatalogReader(root) + + +def get_catalog_reader() -> CatalogReader: + return _reader(get_settings().data_dir / "catalog") diff --git a/src/pullbox/services/catalog/retention.py b/src/pullbox/services/catalog/retention.py new file mode 100644 index 00000000..1a70c046 --- /dev/null +++ b/src/pullbox/services/catalog/retention.py @@ -0,0 +1,47 @@ +"""Bounded cleanup of catalog-owned artifacts, never of library data.""" + +from __future__ import annotations + +import re +import time +from typing import TYPE_CHECKING + +from pullbox.services.catalog.database import safe_path +from pullbox.services.catalog.storage import load_json + +if TYPE_CHECKING: + from pathlib import Path + + +def cleanup(root: Path, *, completed: bool = False) -> None: + """Caller holds the update lock; keep active, previous and their weekly bases. + + Old generations get a two-day grace period so in-flight readers can finish. + Only fixed-format owned files are removed; unknown files are left alone. + """ + keep = set() + for pointer in ("active.json", "previous.json"): + ref = load_json(root / pointer) + keep.add(str(ref.get("path", ""))) + keep.add(f"bases/{ref.get('base_version', '')}.db") + for directory, pattern in ( + ("staging", r"catalog-[A-Za-z0-9_-]+(?:\.patch)?\.db(?:-journal)?"), + ("downloads", r"[a-f0-9]{64}\.part"), + ("bases", r"\d{8}T\d{6}Z\.db"), + ("versions", r"\d{8}T\d{6}Z\.db"), + ): + folder = safe_path(root / directory) + if not folder.exists(): + continue + for path in folder.iterdir(): + safe_path(path) + if not path.is_file() or not re.fullmatch(pattern, path.name): + continue + relative = str(path.relative_to(root)) + expired = path.stat().st_mtime < time.time() - 2 * 86400 + if ( + directory == "staging" + or (directory == "downloads" and completed) + or (expired and relative not in keep) + ): + path.unlink() diff --git a/src/pullbox/services/catalog/service.py b/src/pullbox/services/catalog/service.py new file mode 100644 index 00000000..e0b274d5 --- /dev/null +++ b/src/pullbox/services/catalog/service.py @@ -0,0 +1,344 @@ +"""One resumable, verified catalog update at a time, independent of library writes.""" + +from __future__ import annotations + +import asyncio +import fcntl +import os +import shutil +from datetime import UTC, datetime, timedelta +from functools import lru_cache +from typing import TYPE_CHECKING, Any + +import httpx +import structlog +from pydantic import BaseModel, ValidationError + +from pullbox.config import get_settings +from pullbox.services.catalog.contract import ( + MAX_MANIFEST_BYTES, + TRUSTED_KEYS, + Artifact, + CatalogError, + Publication, + valid_version, + verify_manifest, +) +from pullbox.services.catalog.database import ( + apply_catalog_patch, + file_sha256, + safe_path, + validate_snapshot, +) +from pullbox.services.catalog.retention import cleanup +from pullbox.services.catalog.storage import ( + activate_file, + atomic_json, + decompress, + disk_work, + load_json, + stage_path, +) + +if TYPE_CHECKING: + from collections.abc import Mapping + from pathlib import Path + + from cryptography.hazmat.primitives.asymmetric.ed25519 import Ed25519PublicKey + +logger = structlog.get_logger(__name__) +BUSY_PHASES = {"checking", "downloading", "verifying", "decompressing", "installing"} + + +class CatalogStatus(BaseModel): + phase: str = "not_downloaded" + requested: bool = False + automatic_updates: bool = True + installed_version: str | None = None + source_cutoff_at: str | None = None + target_version: str | None = None + bytes_downloaded: int = 0 + bytes_total: int = 0 + catalog_size_bytes: int = 0 + last_checked_at: datetime | None = None + last_updated_at: datetime | None = None + attempt_started_at: datetime | None = None + error: str | None = None + + +class CatalogService: + """Use only signed API coordinates and activate after all validation succeeds.""" + + def __init__( + self, + root: Path, + base_url: str, + *, + keys: Mapping[str, Ed25519PublicKey] | None = None, + transport: httpx.AsyncBaseTransport | None = None, + ) -> None: + self.root = root + url = httpx.URL(base_url) + if ( + url.scheme not in {"http", "https"} + or not url.host + or url.userinfo + or url.query + or url.fragment + ): + raise CatalogError("The Pullbox API address is invalid.") + self.base_url = str(url).rstrip("/") + self.keys = TRUSTED_KEYS if keys is None else keys + self.transport = transport + self._lock = asyncio.Lock() + self._state_lock = asyncio.Lock() + try: + self._state = CatalogStatus.model_validate(load_json(root / "state.json")) + if self._state.phase in BUSY_PHASES: + self._state.phase = "interrupted" + self._state.error = "The previous download was interrupted. It can be resumed." + active = load_json(root / "active.json") + if active: + self._state.installed_version = str(active["version"]) + self._state.source_cutoff_at = str(active["source_cutoff_at"]) + except (CatalogError, ValidationError, KeyError): + self._state = CatalogStatus( + phase="failed", error="Catalog state could not be read. Check the data volume." + ) + + def status(self) -> CatalogStatus: + return self._state.model_copy(deep=True) + + async def set_automatic_updates(self, enabled: bool) -> None: + self._state.automatic_updates = enabled + await self._save() + + async def _save(self) -> None: + async with self._state_lock: + await disk_work( + atomic_json, self.root / "state.json", self._state.model_dump(mode="json") + ) + + async def sync(self, *, manual: bool = False) -> bool: + if self._lock.locked(): + return False + if not manual: + if not self._state.requested or not self._state.automatic_updates: + return False + checked = self._state.last_checked_at + # Allow the daily scheduler's jitter to move earlier than yesterday. + if checked and checked > datetime.now(UTC) - timedelta(hours=23): + return False + async with self._lock: + safe_path(self.root).mkdir(parents=True, exist_ok=True) + lock_path = safe_path(self.root / "update.lock") + with lock_path.open("a+b") as lock: + try: + fcntl.flock(lock, fcntl.LOCK_EX | fcntl.LOCK_NB) + except BlockingIOError: + return False + try: + self._state.attempt_started_at = datetime.now(UTC) + await disk_work(cleanup, self.root) + self._state.requested = True + self._state.phase, self._state.error = "checking", None + await self._save() + async with httpx.AsyncClient( + timeout=httpx.Timeout(60, connect=15), + follow_redirects=False, + transport=self.transport, + ) as client: + changed = await self._update(client) + await disk_work(lambda: cleanup(self.root, completed=True)) + self._state.phase = "current" + self._state.last_checked_at = datetime.now(UTC) + await self._save() + logger.info( + "catalog_update_complete", + version=self._state.installed_version, + changed=changed, + ) + return changed + except asyncio.CancelledError: + self._state.phase = "interrupted" + await self._save() + raise + except (CatalogError, httpx.HTTPError, OSError) as exc: + message = ( + str(exc) + if isinstance(exc, CatalogError) + else "Catalog download failed. Check the API connection and " + "available disk space, then retry." + ) + self._state.phase, self._state.error = "failed", message + await self._save() + logger.warning("catalog_update_failed", reason=message) + raise CatalogError(message) from exc + finally: + fcntl.flock(lock, fcntl.LOCK_UN) + + async def _manifest(self, client: httpx.AsyncClient) -> Publication: + cache = await disk_work(load_json, self.root / "manifest.json") + headers = {"Accept-Encoding": "identity"} + cached_raw = cache.get("raw") + if isinstance(cached_raw, str) and isinstance(cache.get("etag"), str): + verify_manifest(cached_raw.encode(), self.keys) + headers["If-None-Match"] = cache["etag"] + async with client.stream( + "GET", self.base_url + "/api/v2/catalog/latest", headers=headers + ) as response: + if response.status_code == 304 and isinstance(cached_raw, str): + return verify_manifest(cached_raw.encode(), self.keys) + response.raise_for_status() + raw = bytearray() + async for chunk in response.aiter_bytes(): + raw.extend(chunk) + if len(raw) > MAX_MANIFEST_BYTES: + raise CatalogError("Catalog manifest is too large.") + publication = verify_manifest(bytes(raw), self.keys) + await disk_work( + atomic_json, + self.root / "manifest.json", + {"raw": raw.decode(), "etag": response.headers.get("etag")}, + ) + return publication + + async def _download(self, client: httpx.AsyncClient, artifact: Artifact) -> Path: + directory = safe_path(self.root / "downloads") + directory.mkdir(parents=True, exist_ok=True) + path = safe_path(directory / f"{artifact.sha256}.part") + offset = path.stat().st_size if path.exists() else 0 + if offset > artifact.size_bytes: + path.unlink() + offset = 0 + if shutil.disk_usage(directory).free < artifact.size_bytes - offset + 64 * 1024 * 1024: + raise CatalogError("Not enough disk space to download the catalog.") + self._state.phase = "downloading" + self._state.bytes_downloaded, self._state.bytes_total = offset, artifact.size_bytes + await self._save() + if offset < artifact.size_bytes: + headers = {"Accept-Encoding": "identity"} + if offset: + headers.update({"Range": f"bytes={offset}-", "If-Range": f'"{artifact.sha256}"'}) + async with client.stream( + "GET", self.base_url + artifact.download_path, headers=headers + ) as response: + response.raise_for_status() + if response.status_code == 206: + expected = f"bytes {offset}-{artifact.size_bytes - 1}/{artifact.size_bytes}" + if response.headers.get("content-range") != expected: + raise CatalogError( + "Catalog resume response is invalid. Retry the download." + ) + elif response.status_code == 200: + offset = 0 + else: + raise CatalogError("Catalog download response is invalid.") + with path.open("ab" if offset else "wb") as stream: + async for chunk in response.aiter_bytes(256 * 1024): + offset += len(chunk) + if offset > artifact.size_bytes: + raise CatalogError("Catalog download exceeded its signed size.") + await disk_work(stream.write, chunk) + self._state.bytes_downloaded = offset + await disk_work(stream.flush) + await disk_work(os.fsync, stream.fileno()) + self._state.phase = "verifying" + if ( + path.stat().st_size != artifact.size_bytes + or await disk_work(file_sha256, path) != artifact.sha256 + ): + path.unlink(missing_ok=True) + raise CatalogError("Catalog download checksum failed. Retry the download.") + return path + + async def _snapshot(self, client: httpx.AsyncClient, artifact: Artifact) -> Path: + target = safe_path(self.root / "bases" / f"{artifact.version}.db") + if target.exists(): + try: + await disk_work(validate_snapshot, target, artifact.version) + return target + except CatalogError: + logger.warning("catalog_weekly_base_invalid", version=artifact.version) + archive = await self._download(client, artifact) + stage = stage_path(self.root, ".db") + try: + self._state.phase = "decompressing" + await disk_work(decompress, archive, stage) + self._state.phase = "verifying" + await disk_work(validate_snapshot, stage, artifact.version) + await disk_work(activate_file, stage, target) + finally: + stage.unlink(missing_ok=True) + return target + + async def _update(self, client: httpx.AsyncClient) -> bool: + publication = await self._manifest(client) + active = await disk_work(load_json, self.root / "active.json") + self._state.target_version = publication.latest_version + if active and str(active.get("version", "")) >= publication.latest_version: + version = valid_version(active.get("version")) + relative = active.get("path") + if relative not in {f"bases/{version}.db", f"versions/{version}.db"}: + raise CatalogError("Catalog active file is invalid. Check the data volume.") + try: + await disk_work(validate_snapshot, self.root / str(relative), version) + return False + except CatalogError: + if version > publication.latest_version: + raise CatalogError( + "The API has an older catalog. The installed version was not replaced." + ) from None + active = {} # Do not replace a good previous reference with a broken one. + base = await self._snapshot(client, publication.full_snapshot) + patch = publication.latest_patch() + target = base + if patch: + if shutil.disk_usage(self.root).free < base.stat().st_size * 2 + 64 * 1024 * 1024: + raise CatalogError("Not enough disk space to apply the catalog update.") + archive = await self._download(client, patch) + unpacked, stage = stage_path(self.root, ".patch.db"), stage_path(self.root, ".db") + try: + self._state.phase = "decompressing" + await disk_work(decompress, archive, unpacked) + self._state.phase = "installing" + await disk_work( + apply_catalog_patch, + base, + unpacked, + stage, + publication.full_snapshot.version, + patch.version, + ) + target = self.root / "versions" / f"{patch.version}.db" + await disk_work(activate_file, stage, target) + finally: + unpacked.unlink(missing_ok=True) + stage.unlink(missing_ok=True) + self._state.phase = "verifying" + manifest = await disk_work(validate_snapshot, target, publication.latest_version) + self._state.phase = "installing" + reference: dict[str, Any] = { + "version": publication.latest_version, + "base_version": publication.full_snapshot.version, + "path": str(target.relative_to(self.root)), + "source_cutoff_at": manifest["source_cutoff_at"], + } + if active: + await disk_work(atomic_json, self.root / "previous.json", active) + await disk_work(atomic_json, self.root / "active.json", reference) + self._state.installed_version = publication.latest_version + self._state.source_cutoff_at = str(manifest["source_cutoff_at"]) + self._state.catalog_size_bytes = target.stat().st_size + self._state.last_updated_at = datetime.now(UTC) + return True + + +@lru_cache(maxsize=4) +def _service(root: Path, base_url: str) -> CatalogService: + return CatalogService(root, base_url) + + +def get_catalog_service() -> CatalogService: + settings = get_settings() + return _service(settings.data_dir / "catalog", settings.pullbox_data_api_base_url) diff --git a/src/pullbox/services/catalog/storage.py b/src/pullbox/services/catalog/storage.py new file mode 100644 index 00000000..3a733789 --- /dev/null +++ b/src/pullbox/services/catalog/storage.py @@ -0,0 +1,116 @@ +"""Bounded artifact decompression and atomic catalog state on the data volume.""" + +from __future__ import annotations + +import asyncio +import json +import os +import shutil +import tempfile +from pathlib import Path +from typing import TYPE_CHECKING, Any, TypeVar + +import zstandard + +from pullbox.services.catalog.contract import CatalogError +from pullbox.services.catalog.database import safe_path + +if TYPE_CHECKING: + from collections.abc import Callable + +T = TypeVar("T") +MAX_DATABASE_BYTES = 4 * 1024**3 + + +async def disk_work[T](func: Callable[..., T], *args: Any) -> T: + """Finish an owned disk operation before cancellation releases the update lock.""" + task = asyncio.create_task(asyncio.to_thread(func, *args)) + try: + return await asyncio.shield(task) + except asyncio.CancelledError: + try: + await task + finally: + raise + + +def atomic_json(path: Path, data: dict[str, Any]) -> None: + safe_path(path) + path.parent.mkdir(parents=True, exist_ok=True) + descriptor, name = tempfile.mkstemp(prefix=".catalog-", dir=path.parent) + stage = Path(name) + try: + with os.fdopen(descriptor, "w") as stream: + json.dump(data, stream, sort_keys=True, separators=(",", ":")) + stream.flush() + os.fsync(stream.fileno()) + os.replace(stage, path) + sync_directory(path.parent) + finally: + stage.unlink(missing_ok=True) + + +def sync_directory(path: Path) -> None: + descriptor = os.open(path, os.O_RDONLY) + try: + os.fsync(descriptor) + finally: + os.close(descriptor) + + +def load_json(path: Path) -> dict[str, Any]: + safe_path(path) + if not path.exists(): + return {} + if path.stat().st_size > 1024 * 1024: + raise CatalogError("Catalog state is invalid. Check the data volume.") + try: + result = json.loads(path.read_bytes()) + except (ValueError, OSError) as exc: + raise CatalogError("Catalog state could not be read. Check the data volume.") from exc + if not isinstance(result, dict): + raise CatalogError("Catalog state is invalid.") + return result + + +def decompress(archive: Path, output: Path) -> None: + safe_path(archive) + safe_path(output) + total = 0 + try: + with archive.open("rb") as source, output.open("xb") as destination: + # The C backend passes this limit to ZSTD_DCtx_setMaxWindowSize in bytes. + with zstandard.ZstdDecompressor(max_window_size=128 * 1024 * 1024).stream_reader( + source + ) as reader: + while chunk := reader.read(1024 * 1024): + total += len(chunk) + if total > MAX_DATABASE_BYTES: + raise CatalogError("Catalog exceeds the supported storage size.") + if ( + total % (32 * 1024 * 1024) == 0 + and shutil.disk_usage(output.parent).free < 64 * 1024 * 1024 + ): + raise CatalogError("Not enough disk space to install the catalog.") + destination.write(chunk) + destination.flush() + os.fsync(destination.fileno()) + except zstandard.ZstdError as exc: + raise CatalogError("Catalog decompression failed. Retry the download.") from exc + + +def stage_path(root: Path, suffix: str) -> Path: + directory = safe_path(root / "staging") + directory.mkdir(parents=True, exist_ok=True) + descriptor, name = tempfile.mkstemp(prefix="catalog-", suffix=suffix, dir=directory) + os.close(descriptor) + path = Path(name) + path.unlink() + return path + + +def activate_file(stage: Path, target: Path) -> None: + safe_path(target) + target.parent.mkdir(parents=True, exist_ok=True) + os.replace(stage, target) + sync_directory(target.parent) diff --git a/src/pullbox/services/import_cv_search.py b/src/pullbox/services/import_cv_search.py index 0c049278..26f0714f 100644 --- a/src/pullbox/services/import_cv_search.py +++ b/src/pullbox/services/import_cv_search.py @@ -29,6 +29,11 @@ async def search_with_retry( year: int | None, ) -> list[SeriesSearchResult]: """Search ComicVine with retry on transient provider failures.""" + if getattr(provider, "is_local_catalog", False) is True: + results, _ = await provider.search_series_globally( + query, max_results=_IMPORT_GLOBAL_SEARCH_LIMIT + ) + return list(results) last_provider_error: str | None = None saw_successful_response = False for attempt in range(_MAX_RETRIES): diff --git a/src/pullbox/services/import_provider_cache.py b/src/pullbox/services/import_provider_cache.py index dd557516..5655e768 100644 --- a/src/pullbox/services/import_provider_cache.py +++ b/src/pullbox/services/import_provider_cache.py @@ -44,8 +44,10 @@ def build_import_scan_metadata_provider( provider: Any, ) -> CachedImportMetadataProvider: """Return the Step 2 provider stack: persistent cache, then per-job cache.""" + from pullbox.services.catalog.lookup import catalog_or_provider + return CachedImportMetadataProvider( - build_persistent_import_metadata_provider(session, provider) + catalog_or_provider(build_persistent_import_metadata_provider(session, provider)) ) diff --git a/src/pullbox/services/import_service.py b/src/pullbox/services/import_service.py index faec8d8f..769b327c 100644 --- a/src/pullbox/services/import_service.py +++ b/src/pullbox/services/import_service.py @@ -868,7 +868,7 @@ async def override_cv_id( async def _fetch_series_metadata_for_override(self, cv_id: int) -> SeriesMetadata: """Fetch ComicVine metadata for a manual imported-series override.""" - return await self._metadata_service._provider.get_series(str(cv_id)) + return await self._metadata_service.get_series_metadata(cv_id) async def rematch_imported_series_files( self, diff --git a/src/pullbox/services/import_service_matching.py b/src/pullbox/services/import_service_matching.py index cd9593c5..74bcce80 100644 --- a/src/pullbox/services/import_service_matching.py +++ b/src/pullbox/services/import_service_matching.py @@ -493,9 +493,13 @@ async def _run_file_matching( def _metadata_provider_for_job(self: ImportServiceMatchingContext, job_id: int) -> Any: """Return the job-scoped cached provider when Step 2 is actively scanning.""" + from pullbox.services.catalog.lookup import catalog_or_provider + if self._metadata_service is None: return None - return self._scan_provider_cache_by_job.get(job_id, self._metadata_service._provider) + return self._scan_provider_cache_by_job.get( + job_id, catalog_or_provider(self._metadata_service._provider) + ) async def _record_duplicate_copy_cluster( self: ImportServiceMatchingContext, diff --git a/src/pullbox/services/metadata_service.py b/src/pullbox/services/metadata_service.py index 08a9c3be..8f1a7e48 100644 --- a/src/pullbox/services/metadata_service.py +++ b/src/pullbox/services/metadata_service.py @@ -36,6 +36,7 @@ ) from pullbox.providers.base import SeriesMetadata from pullbox.providers.metadata.comicvine import ComicVineError +from pullbox.services.catalog.reader import CatalogIssueSummary, CatalogSeriesMetadata from pullbox.services.cover_cache_service import purge_series_cover_cache if TYPE_CHECKING: @@ -43,6 +44,7 @@ from pullbox.providers.base import IssueMetadata, IssueSummary from pullbox.providers.metadata.comicvine import ComicVineProvider + from pullbox.services.catalog.reader import CatalogReader logger = structlog.get_logger(__name__) @@ -129,10 +131,13 @@ def __init__( provider: ComicVineProvider, covers_dir: Path, refresh_days: int = 30, + *, + catalog: CatalogReader | None = None, ) -> None: self._provider = provider self._covers_dir = covers_dir self._refresh_days = refresh_days + self._catalog = catalog async def fetch_series( self, @@ -150,7 +155,7 @@ async def fetch_series( log = logger.bind(comicvine_id=comicvine_id) log.debug("metadata_fetch_series") - meta = await self.get_series_metadata(comicvine_id) + meta = await self.get_series_metadata(comicvine_id, use_catalog=False) series = await self.upsert_series_metadata( session, comicvine_id, @@ -176,6 +181,11 @@ async def upsert_series_metadata( """ log = logger.bind(comicvine_id=comicvine_id) + source = "pullbox_catalog" if isinstance(meta, CatalogSeriesMetadata) else "comicvine" + refreshed_at = ( + meta.source_cutoff_at if isinstance(meta, CatalogSeriesMetadata) else datetime.now(UTC) + ) + publisher_id = None if meta.publisher: publisher_id = await self._ensure_publisher(session, meta.publisher) @@ -184,6 +194,13 @@ async def upsert_series_metadata( await session.execute(select(Series).where(Series.comicvine_id == comicvine_id)) ).scalar_one_or_none() + if ( + existing + and isinstance(meta, CatalogSeriesMetadata) + and existing.metadata_source == "comicvine" + ): + # Basic catalog hydration must never replace a completed full refresh. + return existing if existing: existing.title = meta.title existing.sort_title = meta.sort_title or meta.title @@ -203,8 +220,8 @@ async def upsert_series_metadata( existing.comicvine_url = meta.comicvine_url existing.cover_url = meta.cover_url existing.publisher_id = publisher_id - existing.metadata_last_refreshed = datetime.now(UTC) - existing.metadata_source = "comicvine" + existing.metadata_last_refreshed = refreshed_at + existing.metadata_source = source series = existing log.debug("metadata_series_updated", series_id=series.id) else: @@ -220,8 +237,8 @@ async def upsert_series_metadata( comicvine_url=meta.comicvine_url, cover_url=meta.cover_url, publisher_id=publisher_id, - metadata_last_refreshed=datetime.now(UTC), - metadata_source="comicvine", + metadata_last_refreshed=refreshed_at, + metadata_source=source, ) session.add(series) await session.flush() @@ -239,11 +256,18 @@ async def upsert_series_metadata( async def get_series_metadata( self, comicvine_id: int, + *, + use_catalog: bool = True, ) -> SeriesMetadata: """Fetch provider series metadata without creating or updating local rows.""" log = logger.bind(comicvine_id=comicvine_id) log.debug("metadata_get_series_metadata") + if use_catalog and self._catalog is not None and self._catalog.available: + from pullbox.services.catalog.lookup import CatalogLookupService + + return await CatalogLookupService(self._catalog).get_series(str(comicvine_id)) + try: return await self._provider.get_series(str(comicvine_id)) except ComicVineError as exc: @@ -254,6 +278,10 @@ async def get_series_metadata_batch( comicvine_ids: list[int], ) -> dict[int, SeriesMetadata]: """Fetch multiple provider series profiles through the optional bulk contract.""" + if self._catalog is not None and self._catalog.available: + return { + key: await self.get_series_metadata(key) for key in dict.fromkeys(comicvine_ids) + } batch_fetch = getattr(type(self._provider), "get_series_batch", None) try: if callable(batch_fetch): @@ -274,6 +302,8 @@ async def get_cached_series_metadata( comicvine_id: int, ) -> SeriesMetadata | None: """Return fresh cached series metadata without starting a provider request.""" + if self._catalog is not None and self._catalog.available: + return await self._catalog.series(comicvine_id) cached_lookup = getattr(type(self._provider), "get_series_cached", None) if cached_lookup is None: return None @@ -283,11 +313,20 @@ async def get_cached_series_metadata( async def get_issue_summaries_for_series( self, comicvine_id: int, + *, + use_catalog: bool = True, ) -> list[IssueSummary]: """Fetch provider issue summaries for a series without touching local issues.""" log = logger.bind(comicvine_id=comicvine_id) log.debug("metadata_get_issue_summaries_for_series") + if use_catalog and self._catalog is not None and self._catalog.available: + from pullbox.services.catalog.lookup import CatalogLookupService + + return await CatalogLookupService(self._catalog).get_issues_for_series( + str(comicvine_id) + ) + try: return await self._provider.get_issues_for_series(str(comicvine_id)) except ComicVineError as exc: @@ -298,6 +337,11 @@ async def get_issue_catalog_batch( comicvine_ids: list[int], ) -> dict[int, list[IssueSummary]]: """Fetch multiple full issue catalogs through the optional bulk contract.""" + if self._catalog is not None and self._catalog.available: + return { + key: await self.get_issue_summaries_for_series(key) + for key in dict.fromkeys(comicvine_ids) + } batch_fetch = getattr(type(self._provider), "get_issue_catalog_batch", None) try: if callable(batch_fetch): @@ -637,7 +681,9 @@ async def fetch_issues_for_series( if not series or not series.comicvine_id: raise NotFoundError("Series", series_id) - summaries = await self.get_issue_summaries_for_series(series.comicvine_id) + summaries = await self.get_issue_summaries_for_series( + series.comicvine_id, use_catalog=False + ) return await self.upsert_issue_summaries( session, series, @@ -706,6 +752,7 @@ async def upsert_issue_summaries( } summary_evidence_types: list[IssueType] = [] for summary in summaries: + source = "pullbox_catalog" if isinstance(summary, CatalogIssueSummary) else "comicvine" provider_issue_id = int(summary.provider_id) exact_issue_number_text = _exact_issue_number_text( summary.issue_number, @@ -727,6 +774,15 @@ async def upsert_issue_summaries( ) existing_by_provider = None assign_provider_issue_id = False + if ( + source == "pullbox_catalog" + and existing_by_provider is not None + and existing_by_provider.metadata_source == "comicvine" + ): + # A basic snapshot cannot establish that live fields are stale. + # Keep its identity and fields, and do not infer parent type here. + summary_evidence_types.append(IssueType.ISSUE) + continue if existing_by_provider is not None and existing_by_provider is not existing: if existing is None: old_issue_number = existing_by_provider.issue_number @@ -792,7 +848,8 @@ async def upsert_issue_summaries( ) if not preserve_explicit_import_type: existing.issue_type = detected_type - existing.metadata_source = "comicvine" + if source != "pullbox_catalog" or existing.metadata_source != "comicvine": + existing.metadata_source = source else: issue = Issue( series_id=series_id, @@ -803,7 +860,7 @@ async def upsert_issue_summaries( release_date=_parse_date(summary.release_date), cover_url=summary.cover_url, issue_type=detected_type, - metadata_source="comicvine", + metadata_source=source, ) session.add(issue) created.append(issue) diff --git a/src/pullbox/services/series_service.py b/src/pullbox/services/series_service.py index f1e9d065..2b279538 100644 --- a/src/pullbox/services/series_service.py +++ b/src/pullbox/services/series_service.py @@ -275,7 +275,8 @@ async def add_from_import_review_targeted( series.issue_catalog_error = None series.issue_catalog_last_synced_at = None series.issue_catalog_last_checked_at = None - series.metadata_source = "comicvine_partial" + if series.metadata_source != "pullbox_catalog": + series.metadata_source = "comicvine_partial" diagnostics = dict(import_series.diagnostics or {}) diagnostics.pop("series_folder_ownership", None) diff --git a/src/pullbox/tasks/__init__.py b/src/pullbox/tasks/__init__.py index 4f1c631d..86cc5373 100644 --- a/src/pullbox/tasks/__init__.py +++ b/src/pullbox/tasks/__init__.py @@ -6,6 +6,7 @@ from pullbox.tasks import backup_task as backup_task from pullbox.tasks import blocklist_task as blocklist_task +from pullbox.tasks import catalog_task as catalog_task from pullbox.tasks import cover_backfill_task as cover_backfill_task from pullbox.tasks import dashboard_task as dashboard_task from pullbox.tasks import database_maintenance_task as database_maintenance_task diff --git a/src/pullbox/tasks/catalog_task.py b/src/pullbox/tasks/catalog_task.py new file mode 100644 index 00000000..74dfafde --- /dev/null +++ b/src/pullbox/tasks/catalog_task.py @@ -0,0 +1,19 @@ +"""Catalog update scheduler integration.""" + +from pullbox.core.scheduler import get_current_task_trigger_type, scheduled_task + + +@scheduled_task( + task_id="catalog_update", + display_name="Local Catalog Update", + trigger="cron", + hour=6, + minute=30, + jitter=1800, + misfire_grace_time=3600, +) +async def update_catalog() -> None: + """Check daily after opt-in; manual runs also allow the first download.""" + from pullbox.services.catalog.service import get_catalog_service + + await get_catalog_service().sync(manual=get_current_task_trigger_type() == "manual") diff --git a/src/pullbox/ui/comicvine_provider.py b/src/pullbox/ui/comicvine_provider.py index 5372b0df..d1d32b0e 100644 --- a/src/pullbox/ui/comicvine_provider.py +++ b/src/pullbox/ui/comicvine_provider.py @@ -22,9 +22,17 @@ async def open_comicvine_ui_provider( session: AsyncSession, *, session_factory: async_sessionmaker[AsyncSession] | None = None, + prefer_catalog: bool = False, ) -> AsyncIterator[Any]: """Open the configured Comic Vine provider for one UI operation.""" from pullbox.core.comicvine_key import get_comicvine_api_key + from pullbox.services.catalog.lookup import CatalogLookupService + from pullbox.services.catalog.reader import get_catalog_reader + + reader = get_catalog_reader() + if prefer_catalog and reader.available: + yield CatalogLookupService(reader) + return api_key = await get_comicvine_api_key(session) await session.rollback() diff --git a/src/pullbox/ui/import_orphaned_routes.py b/src/pullbox/ui/import_orphaned_routes.py index d01f3ae1..8fbab0f1 100644 --- a/src/pullbox/ui/import_orphaned_routes.py +++ b/src/pullbox/ui/import_orphaned_routes.py @@ -168,11 +168,19 @@ async def import_orphaned_cv_search( if query_text: parsed_query = parse_comicvine_series_query(query_text) api_key = await get_comicvine_api_key(session) - if api_key: + from pullbox.services.catalog.lookup import CatalogLookupService + from pullbox.services.catalog.reader import get_catalog_reader + + catalog = get_catalog_reader() + if api_key or catalog.available: try: - provider = wrap_comicvine_provider_for_ui_cache( - ComicVineProvider(api_key=api_key, rate_limit=10), - request, + provider = ( + CatalogLookupService(catalog) + if catalog.available + else wrap_comicvine_provider_for_ui_cache( + ComicVineProvider(api_key=api_key, rate_limit=10), + request, + ) ) cv_results, _total_results = await provider.search_series_globally( parsed_query.title_query, diff --git a/src/pullbox/ui/import_routes.py b/src/pullbox/ui/import_routes.py index 7ccd07be..ed5a313d 100644 --- a/src/pullbox/ui/import_routes.py +++ b/src/pullbox/ui/import_routes.py @@ -1482,11 +1482,19 @@ async def import_cv_search( if query_text: parsed_query = parse_comicvine_series_query(query_text) api_key = await get_comicvine_api_key(session) - if api_key: + from pullbox.services.catalog.lookup import CatalogLookupService + from pullbox.services.catalog.reader import get_catalog_reader + + catalog = get_catalog_reader() + if api_key or catalog.available: try: - provider = wrap_comicvine_provider_for_ui_cache( - ComicVineProvider(api_key=api_key, rate_limit=10), - request, + provider = ( + CatalogLookupService(catalog) + if catalog.available + else wrap_comicvine_provider_for_ui_cache( + ComicVineProvider(api_key=api_key, rate_limit=10), + request, + ) ) cv_results, _total_results = await provider.search_series_globally( parsed_query.title_query, diff --git a/src/pullbox/ui/series_routes.py b/src/pullbox/ui/series_routes.py index e0ad5b5e..cd6d9b07 100644 --- a/src/pullbox/ui/series_routes.py +++ b/src/pullbox/ui/series_routes.py @@ -601,6 +601,9 @@ async def load_add_series_search_context( search_mode: str | None = None, session_factory: async_sessionmaker[AsyncSession] | None = None, ) -> dict[str, object]: + from pullbox.services.catalog.reader import get_catalog_reader + + local_catalog = get_catalog_reader().available per_page = ADD_SERIES_PER_PAGE normalized_query = (query or "").strip() normalized_sort = normalize_add_series_sort(sort) @@ -615,6 +618,7 @@ async def load_add_series_search_context( ) or 0 base_context: dict[str, object] = { + "local_catalog": local_catalog, "search_query": normalized_query, "add_series_sort": normalized_sort, "is_preview_search": False, @@ -661,6 +665,7 @@ async def load_add_series_search_context( async with open_comicvine_ui_provider( session, session_factory=session_factory, + prefer_catalog=True, ) as provider: if preview_mode: cv_results, _total_results = await provider.search_series_page( @@ -695,7 +700,11 @@ async def load_add_series_search_context( ) except Exception: logger.exception("comicvine_search_failed", query=normalized_query) - base_context["search_error"] = "ComicVine search failed. Check your API key in settings." + base_context["search_error"] = ( + "Local catalog search failed. Check its status in Metadata settings." + if local_catalog + else "ComicVine search failed. Check your API key in settings." + ) return base_context in_library_count = sum(1 for item in search_results if bool(item.get("already_added"))) diff --git a/src/pullbox/ui/templates/components/comicvine_search_loading.html b/src/pullbox/ui/templates/components/comicvine_search_loading.html index 8f05169b..a913ebf1 100644 --- a/src/pullbox/ui/templates/components/comicvine_search_loading.html +++ b/src/pullbox/ui/templates/components/comicvine_search_loading.html @@ -1,4 +1,4 @@ -{% macro comicvine_search_loading(element_id, testid) -%} +{% macro comicvine_search_loading(element_id, testid, local_catalog=false) -%}
-

Searching ComicVine

+

{{ 'Searching local catalog' if local_catalog else 'Searching ComicVine' }}

Large catalogs can take a moment.

diff --git a/src/pullbox/ui/templates/partials/add_series_header_metrics.html b/src/pullbox/ui/templates/partials/add_series_header_metrics.html index 4d477a78..e7382d21 100644 --- a/src/pullbox/ui/templates/partials/add_series_header_metrics.html +++ b/src/pullbox/ui/templates/partials/add_series_header_metrics.html @@ -7,7 +7,7 @@

ADD SERIES

- Search ComicVine · review · add to library + {{ 'Search local catalog' if local_catalog else 'Search ComicVine' }} · review · add to library

diff --git a/src/pullbox/ui/templates/partials/add_series_results.html b/src/pullbox/ui/templates/partials/add_series_results.html index 2591b8d0..3fdd745a 100644 --- a/src/pullbox/ui/templates/partials/add_series_results.html +++ b/src/pullbox/ui/templates/partials/add_series_results.html @@ -5,7 +5,7 @@ data-testid="add-series-results" class="add-series-results-shell space-y-3" > - {{ comicvine_search_loading("add-series-results-loading", "add-series-results-loading") }} + {{ comicvine_search_loading("add-series-results-loading", "add-series-results-loading", local_catalog) }} {% if is_preview_search and search_query %}
Quick preview

{% if search_shown_count %} - Showing {{ search_shown_count }} fast ComicVine match{% if search_shown_count != 1 %}es{% endif %}. + Showing {{ search_shown_count }} fast {{ 'local catalog' if local_catalog else 'ComicVine' }} match{% if search_shown_count != 1 %}es{% endif %}. {% else %} No quick matches were found in the fast preview. {% endif %} - Press Enter or search all results to load the full ComicVine result set. + Press Enter or search all results to load the full {{ 'local catalog' if local_catalog else 'ComicVine' }} result set.

- Search all ComicVine results + Search all {{ 'local catalog' if local_catalog else 'ComicVine' }} results {% endif %} @@ -66,9 +66,10 @@ -

No ComicVine matches

+

No {{ 'local catalog' if local_catalog else 'ComicVine' }} matches

Try the exact title, add the year, or remove punctuation and subtitles from the search. + {% if local_catalog %}You can also check for a catalog update in Metadata settings. No online search was made.{% endif %}

{% else %} @@ -76,7 +77,7 @@ -

Search ComicVine

+

{{ 'Search local catalog' if local_catalog else 'Search ComicVine' }}

Type a series name above to find and add it to your library.

diff --git a/src/pullbox/ui/templates/partials/settings_catalog.html b/src/pullbox/ui/templates/partials/settings_catalog.html new file mode 100644 index 00000000..c275cb22 --- /dev/null +++ b/src/pullbox/ui/templates/partials/settings_catalog.html @@ -0,0 +1,75 @@ +{% call section_card(title="Local Comic Vine catalog", eyebrow="Offline search and imports", title_style="plain", body_class="p-5") %} +
+

Download a verified catalog to your Pullbox server for local series searches and import matching. Daily updates download only changes when possible. Full metadata refreshes still use your Comic Vine API key.

+

Catalog files stay separate from your library database. A download never moves or imports comics.

+
+ + · Data through +
+
+ + +

You can leave this page while Pullbox works.

+
+ +
+
Installed version:
+
Catalog size:
+
Last checked:
+
+ + +

Local searches do not automatically fall back to online searches when no match is found.

+
+{% endcall %} + diff --git a/src/pullbox/ui/templates/partials/settings_metadata.html b/src/pullbox/ui/templates/partials/settings_metadata.html index dd0bc84d..67b9bc97 100644 --- a/src/pullbox/ui/templates/partials/settings_metadata.html +++ b/src/pullbox/ui/templates/partials/settings_metadata.html @@ -2,6 +2,7 @@ {% from "components/settings_shell.html" import field_note, inline_alert, section_card, settings_footer, settings_row, ui_icon %}
+ {% include "partials/settings_catalog.html" %} {% call section_card( title="ComicVine access", eyebrow="API credentials", @@ -89,14 +90,14 @@
{% call settings_row(label="Refresh Interval") %}
- + days
{% endcall %} {% call settings_row(label="Import Match Burst", help="How aggressively bulk import matching may burst ComicVine requests before pacing them out. Lower is safer.") %}
- + requests
{% endcall %} diff --git a/tests/catalog_fixtures.py b/tests/catalog_fixtures.py new file mode 100644 index 00000000..b884ada1 --- /dev/null +++ b/tests/catalog_fixtures.py @@ -0,0 +1,141 @@ +"""Small deterministic fixtures for the published catalog v2 wire contract.""" + +import hashlib +import json +import sqlite3 + +TABLES = { + "publishers": ("id,name", "id"), + "series": ("id,name,start_year,publisher_id,issue_count,cover_url", "id"), + "series_aliases": ("series_id,position,alias,normalized_alias", "series_id,position"), + "issues": ( + "id,series_id,issue_number,normalized_issue_number,sort_number,title,cover_date,store_date,cover_url", + "id", + ), +} +SCHEMA = """ +PRAGMA application_id=1346519858; +PRAGMA user_version=1; +CREATE TABLE publishers(id INTEGER PRIMARY KEY, name TEXT); +CREATE TABLE series(id INTEGER PRIMARY KEY, name TEXT NOT NULL, start_year INTEGER, + publisher_id INTEGER REFERENCES publishers(id), issue_count INTEGER, cover_url TEXT); +CREATE TABLE series_aliases(series_id INTEGER REFERENCES series(id), position INTEGER, + alias TEXT NOT NULL, normalized_alias TEXT NOT NULL, PRIMARY KEY(series_id,position)); +CREATE TABLE issues(id INTEGER PRIMARY KEY, series_id INTEGER REFERENCES series(id), + issue_number TEXT, normalized_issue_number TEXT, sort_number TEXT, title TEXT, + cover_date TEXT, store_date TEXT, cover_url TEXT); +CREATE VIRTUAL TABLE series_fts USING fts5(series_id UNINDEXED,name,aliases, + tokenize='porter unicode61 remove_diacritics 1'); +CREATE TABLE dataset_manifest(key TEXT PRIMARY KEY,value TEXT NOT NULL); +CREATE TABLE schema_migrations(version TEXT PRIMARY KEY,applied_at TEXT NOT NULL); +""" + + +def content_hash(db): + digest = hashlib.sha256() + for table, (columns, order) in TABLES.items(): + for row in db.execute(f"SELECT {columns} FROM {table} ORDER BY {order}"): + digest.update(table.encode() + b"\0") + digest.update( + json.dumps(row, ensure_ascii=False, separators=(",", ":")).encode() + b"\n" + ) + return digest.hexdigest() + + +def build_snapshot(path, version="20260913T050000Z", name="Batman", *, extra_issue=False): + with sqlite3.connect(path) as db: + db.executescript(SCHEMA) + db.execute("INSERT INTO publishers VALUES (1,'DC')") + db.execute("INSERT INTO series VALUES (10,?,2016,1,1,NULL)", (name,)) + db.execute("INSERT INTO series_aliases VALUES (10,0,'The Dark Knight','the dark knight')") + db.execute( + "INSERT INTO issues VALUES (100,10,'½','0.5','0.5','Rebirth','2016-08-01',NULL,NULL)" + ) + if extra_issue: + db.execute("INSERT INTO issues VALUES (101,10,'2','2','2','Second',NULL,NULL,NULL)") + db.execute( + "INSERT INTO series_fts(rowid,series_id,name,aliases) " + "VALUES (10,10,?,'The Dark Knight')", + (name,), + ) + manifest = { + "application_id": 0x50424332, + "format_id": "pullbox-catalog-v2", + "schema_version": "1", + "compatibility_version": "pullbox-only", + "dataset_version": version, + "source_cutoff_at": f"{version[:4]}-{version[4:6]}-{version[6:8]}T05:00:00+00:00", + "counts": { + table: db.execute(f"SELECT COUNT(*) FROM {table}").fetchone()[0] for table in TABLES + }, + "content_sha256": content_hash(db), + } + db.executemany( + "INSERT INTO dataset_manifest VALUES (?,?)", + [(k, json.dumps(v)) for k, v in manifest.items()], + ) + db.execute("INSERT INTO schema_migrations VALUES ('1','2026-09-13T05:00:00+00:00')") + return manifest + + +def build_patch(path, base, target): + names = { + "publishers": "publisher", + "series": "series", + "series_aliases": "series_alias", + "issues": "issue", + } + with ( + sqlite3.connect(path) as db, + sqlite3.connect(base) as before, + sqlite3.connect(target) as after, + ): + db.execute("PRAGMA application_id=1346523204") + db.execute("PRAGMA user_version=1") + for table, (columns, key_columns) in TABLES.items(): + before.execute("SELECT sql FROM sqlite_master WHERE name=?", (table,)).fetchone()[0] + # Patch upserts use the same columns but do not carry snapshot foreign keys. + cols = columns.split(",") + keys = key_columns.split(",") + prefix = names[table] + db.execute(f"CREATE TABLE {prefix}_upserts ({','.join(cols)})") + db.execute(f"CREATE TABLE {prefix}_deletes ({key_columns})") + old = { + tuple(row[cols.index(k)] for k in keys): row + for row in before.execute(f"SELECT {columns} FROM {table}") + } + new = { + tuple(row[cols.index(k)] for k in keys): row + for row in after.execute(f"SELECT {columns} FROM {table}") + } + for key, row in new.items(): + if old.get(key) != row: + db.execute( + f"INSERT INTO {prefix}_upserts VALUES ({','.join('?' for _ in cols)})", row + ) + for key in old.keys() - new.keys(): + db.execute( + f"INSERT INTO {prefix}_deletes VALUES ({','.join('?' for _ in keys)})", key + ) + db.execute("CREATE TABLE target_dataset_manifest(key TEXT PRIMARY KEY,value TEXT)") + db.executemany( + "INSERT INTO target_dataset_manifest VALUES (?,?)", + after.execute("SELECT * FROM dataset_manifest"), + ) + before_meta = { + k: json.loads(v) for k, v in before.execute("SELECT * FROM dataset_manifest") + } + after_meta = {k: json.loads(v) for k, v in after.execute("SELECT * FROM dataset_manifest")} + manifest = { + "format_id": "pullbox-catalog-v2-patch", + "schema_version": "1", + "base_version": before_meta["dataset_version"], + "target_version": after_meta["dataset_version"], + "base_snapshot_sha256": hashlib.sha256(base.read_bytes()).hexdigest(), + "target_content_sha256": after_meta["content_sha256"], + } + db.execute("CREATE TABLE patch_manifest(key TEXT PRIMARY KEY,value TEXT)") + db.executemany( + "INSERT INTO patch_manifest VALUES (?,?)", + [(k, json.dumps(v)) for k, v in manifest.items()], + ) diff --git a/tests/e2e/test_accessibility.py b/tests/e2e/test_accessibility.py index d2b26b0b..d46032ce 100644 --- a/tests/e2e/test_accessibility.py +++ b/tests/e2e/test_accessibility.py @@ -33,6 +33,7 @@ def test_login_page_has_no_wcag_aa_violations( [ ("/", "[data-testid='dashboard-page']"), ("/settings?tab=general", "[data-testid='settings-page']"), + ("/settings?tab=metadata", "[data-testid='settings-page']"), ("/settings?tab=resolvers", "[data-testid='settings-resolvers-card']"), ("/security?tab=authentication", "[data-testid='security-page']"), ("/system?tab=tasks", "[data-testid='system-page']"), diff --git a/tests/ui/test_catalog_controls.py b/tests/ui/test_catalog_controls.py new file mode 100644 index 00000000..33571c53 --- /dev/null +++ b/tests/ui/test_catalog_controls.py @@ -0,0 +1,90 @@ +"""Catalog settings and controls exercise the real authentication middleware.""" + +import os +import sys +from unittest.mock import Mock + +import pytest + +from pullbox.services.auth_service import SESSION_COOKIE_NAME, AuthService +from tests.unit.test_catalog_reader import installed_reader +from tests.unit.test_catalog_service import Server + +sys.path.insert(0, os.path.join(os.path.dirname(__file__), "..")) +pytest_plugins = ["conftest_security"] + + +@pytest.fixture +def local_catalog(tmp_path, monkeypatch): + service = Server(tmp_path).service(tmp_path / "local") + monkeypatch.setattr("pullbox.services.catalog.service.get_catalog_service", lambda: service) + return service + + +async def test_catalog_status_is_authenticated(unauthenticated_client, sec_user, local_catalog): + response = await unauthenticated_client.get("/api/v1/catalog") + assert response.status_code == 401 + + +async def test_catalog_writes_require_csrf(authenticated_client, local_catalog): + response = await authenticated_client.post("/api/v1/catalog/sync") + assert response.status_code == 403 + response = await authenticated_client.patch( + "/api/v1/catalog/preferences", json={"automatic_updates": False} + ) + assert response.status_code == 403 + assert local_catalog.status().requested is False + + +async def test_operator_can_queue_and_save_preference( + authenticated_client, local_catalog, monkeypatch +): + scheduler = Mock() + scheduler.run_task_now.return_value = "queued" + monkeypatch.setattr("pullbox.core.scheduler.get_scheduler", lambda: scheduler) + token = authenticated_client.cookies.get(SESSION_COOKIE_NAME) + headers = {"X-CSRF-Token": AuthService.get_csrf_token_from_session(token)} + response = await authenticated_client.post("/api/v1/catalog/sync", headers=headers) + assert response.status_code == 202 + scheduler.run_task_now.assert_called_once_with("catalog_update") + response = await authenticated_client.patch( + "/api/v1/catalog/preferences", headers=headers, json={"automatic_updates": False} + ) + assert response.status_code == 200 + assert local_catalog.status().automatic_updates is False + assert local_catalog.status().requested is False + + +async def test_metadata_settings_explains_local_storage_and_progress(authenticated_client): + response = await authenticated_client.get("/settings?tab=metadata") + assert response.status_code == 200 + assert 'id="catalog-download-progress"' in response.text + assert 'role="status" aria-live="polite"' in response.text + assert "window.pullboxLiveUpdatesEnabled()" in response.text + assert "Download catalog" in response.text + assert "A download never moves or imports comics." in response.text + + +async def test_add_series_labels_local_search_without_api_key( + authenticated_client, tmp_path, monkeypatch +): + reader = installed_reader(tmp_path) + monkeypatch.setattr("pullbox.services.catalog.reader.get_catalog_reader", lambda: reader) + response = await authenticated_client.get("/series/add?q=Batman") + assert response.status_code == 200 + assert "Search local catalog" in response.text + assert "Batman" in response.text + assert "Check your API key" not in response.text + + +async def test_api_keys_cannot_start_catalog_download( + unauthenticated_client, sec_api_key, local_catalog +): + readable = await unauthenticated_client.get( + "/api/v1/catalog", headers={"X-Api-Key": sec_api_key} + ) + assert readable.status_code == 200 + response = await unauthenticated_client.post( + "/api/v1/catalog/sync", headers={"X-Api-Key": sec_api_key} + ) + assert response.status_code == 401 diff --git a/tests/unit/test_catalog_contract.py b/tests/unit/test_catalog_contract.py new file mode 100644 index 00000000..b6dec120 --- /dev/null +++ b/tests/unit/test_catalog_contract.py @@ -0,0 +1,95 @@ +"""The publication signature is the trust boundary for all download coordinates.""" + +import base64 +import hashlib +import json + +import pytest +from cryptography.hazmat.primitives.asymmetric.ed25519 import Ed25519PrivateKey + +from pullbox.services.catalog.contract import CatalogError, verify_manifest + + +def signed_publication(payload, key): + canonical = json.dumps(payload, sort_keys=True, separators=(",", ":")).encode() + return json.dumps( + { + "payload": payload, + "signature": { + "algorithm": "Ed25519", + "key_id": "test", + "payload_sha256": hashlib.sha256(canonical).hexdigest(), + "value": base64.urlsafe_b64encode(key.sign(canonical)).decode().rstrip("="), + }, + } + ).encode() + + +def publication_payload(): + return { + "format_id": "pullbox-catalog-v2-publication", + "schema_version": "1", + "published_at": "2026-09-13T05:00:00+00:00", + "latest_version": "20260913T050000Z", + "full_snapshot": { + "version": "20260913T050000Z", + "size_bytes": 123, + "sha256": "a" * 64, + "download_path": "/api/v2/catalog/snapshots/20260913T050000Z", + }, + "snapshots": [], + "patches": [], + "retention": {"snapshot_count": 6, "patch_days": 35}, + } + + +def test_verifies_the_signed_artifact_coordinates(): + key = Ed25519PrivateKey.generate() + result = verify_manifest( + signed_publication(publication_payload(), key), {"test": key.public_key()} + ) + assert result is not None + assert result.latest_version == "20260913T050000Z" + assert result.full_snapshot.size_bytes == 123 + + +def test_rejects_tampering_before_accepting_a_download(): + key = Ed25519PrivateKey.generate() + document = json.loads(signed_publication(publication_payload(), key)) + document["payload"]["full_snapshot"]["size_bytes"] = 1 + with pytest.raises(CatalogError, match="signature"): + verify_manifest(json.dumps(document).encode(), {"test": key.public_key()}) + + +def test_rejects_unknown_signing_key(): + key = Ed25519PrivateKey.generate() + with pytest.raises(CatalogError, match="signing key"): + verify_manifest(signed_publication(publication_payload(), key), {}) + + +@pytest.mark.parametrize("invalid_number", [float("nan"), float("inf")]) +def test_non_json_numbers_report_a_safe_manifest_error(invalid_number): + key = Ed25519PrivateKey.generate() + payload = publication_payload() + payload["full_snapshot"]["size_bytes"] = invalid_number + with pytest.raises(CatalogError, match="manifest is invalid"): + verify_manifest(signed_publication(payload, key), {"test": key.public_key()}) + + +@pytest.mark.parametrize( + "change", + [ + {"download_path": "https://elsewhere.example/private"}, + {"download_path": "/api/v2/catalog/snapshots/../../secret"}, + {"size_bytes": -1}, + {"size_bytes": True}, + {"sha256": "bad"}, + {"version": "garbage"}, + ], +) +def test_rejects_invalid_signed_artifacts(change): + key = Ed25519PrivateKey.generate() + payload = publication_payload() + payload["full_snapshot"].update(change) + with pytest.raises(CatalogError): + verify_manifest(signed_publication(payload, key), {"test": key.public_key()}) diff --git a/tests/unit/test_catalog_controls.py b/tests/unit/test_catalog_controls.py new file mode 100644 index 00000000..63634afb --- /dev/null +++ b/tests/unit/test_catalog_controls.py @@ -0,0 +1,23 @@ +"""Catalog controls use the existing operator queue and expose safe status.""" + +from unittest.mock import Mock + +from pullbox.api.v1.catalog import catalog_status, sync_catalog +from pullbox.services.catalog.service import CatalogStatus + + +async def test_status_reads_local_service_only(monkeypatch): + service = Mock() + service.status.return_value = CatalogStatus( + phase="current", installed_version="20260913T050000Z" + ) + monkeypatch.setattr("pullbox.services.catalog.service.get_catalog_service", lambda: service) + assert (await catalog_status(Mock())).installed_version == "20260913T050000Z" + + +async def test_download_is_queued_not_awaited_in_request(monkeypatch): + scheduler = Mock() + scheduler.run_task_now.return_value = "queued" + monkeypatch.setattr("pullbox.core.scheduler.get_scheduler", lambda: scheduler) + assert await sync_catalog(Mock()) == {"status": "queued"} + scheduler.run_task_now.assert_called_once_with("catalog_update") diff --git a/tests/unit/test_catalog_database.py b/tests/unit/test_catalog_database.py new file mode 100644 index 00000000..677d6c11 --- /dev/null +++ b/tests/unit/test_catalog_database.py @@ -0,0 +1,75 @@ +"""Reject invalid catalogs and reproduce cumulative updates from the weekly base.""" + +import sqlite3 + +import pytest + +from pullbox.services.catalog.contract import CatalogError +from pullbox.services.catalog.database import apply_catalog_patch, validate_snapshot +from tests.catalog_fixtures import build_patch, build_snapshot + + +def test_validates_snapshot_content_and_identity(tmp_path): + path = tmp_path / "base.db" + expected = build_snapshot(path) + result = validate_snapshot(path, "20260913T050000Z") + assert result == expected + + +@pytest.mark.parametrize( + "sql", + [ + "UPDATE series SET name='corrupt'", + "PRAGMA application_id=0", + "PRAGMA user_version=99", + "DELETE FROM series_fts", + "CREATE TABLE private_payload(secret TEXT)", + "UPDATE issues SET series_id=999", + "CREATE TRIGGER extra AFTER INSERT ON series BEGIN DELETE FROM issues; END", + ], +) +def test_rejects_invalid_snapshot(tmp_path, sql): + path = tmp_path / "base.db" + build_snapshot(path) + with sqlite3.connect(path) as db: + db.execute(sql) + with pytest.raises(CatalogError): + validate_snapshot(path, "20260913T050000Z") + + +def test_rejects_symlink_snapshot(tmp_path): + base = tmp_path / "base.db" + build_snapshot(base) + link = tmp_path / "link.db" + link.symlink_to(base) + with pytest.raises(CatalogError): + validate_snapshot(link, "20260913T050000Z") + + +def test_daily_updates_rebuild_from_base_including_reverted_changes(tmp_path): + base, day1, day2 = (tmp_path / name for name in ("base.db", "day1.db", "day2.db")) + build_snapshot(base) + build_snapshot(day1, "20260914T050000Z", "Changed", extra_issue=True) + expected = build_snapshot(day2, "20260915T050000Z") + for target, version in ((day1, "20260914T050000Z"), (day2, "20260915T050000Z")): + patch, output = tmp_path / f"{version}.patch", tmp_path / f"{version}.db" + build_patch(patch, base, target) + apply_catalog_patch(base, patch, output, "20260913T050000Z", version) + assert output.exists() + assert validate_snapshot(output, "20260915T050000Z") == expected + with sqlite3.connect(output) as db: + assert db.execute("SELECT name FROM series").fetchone()[0] == "Batman" + assert db.execute("SELECT COUNT(*) FROM issues").fetchone()[0] == 1 + + +def test_wrong_base_fails_without_modifying_it(tmp_path): + base, target, patch = (tmp_path / name for name in ("base.db", "target.db", "patch.db")) + build_snapshot(base) + build_snapshot(target, "20260914T050000Z") + build_patch(patch, base, target) + original = base.read_bytes() + with pytest.raises(CatalogError): + apply_catalog_patch( + base, patch, tmp_path / "out.db", "20260912T050000Z", "20260914T050000Z" + ) + assert base.read_bytes() == original diff --git a/tests/unit/test_catalog_integration.py b/tests/unit/test_catalog_integration.py new file mode 100644 index 00000000..1cfd9fb5 --- /dev/null +++ b/tests/unit/test_catalog_integration.py @@ -0,0 +1,155 @@ +"""Catalog-backed discovery never labels basic data as a full provider refresh.""" + +from unittest.mock import AsyncMock + +from sqlalchemy import select + +from pullbox.models.issue import Issue, IssueType +from pullbox.services.catalog.lookup import CatalogLookupService +from pullbox.services.import_cv_search import search_with_retry +from pullbox.services.import_provider_cache import build_import_scan_metadata_provider +from pullbox.services.metadata_service import MetadataService +from pullbox.services.series_service import SeriesService +from pullbox.ui.comicvine_provider import open_comicvine_ui_provider +from tests.unit.test_catalog_reader import installed_reader + + +async def test_metadata_hydration_uses_catalog_without_provider_calls(tmp_path): + reader = installed_reader(tmp_path) + live = AsyncMock() + service = MetadataService(live, tmp_path) + service._catalog = reader + meta = await service.get_series_metadata(10) + issues = await service.get_issue_summaries_for_series(10) + assert meta.title == "Batman" + assert len(issues) == 1 + live.get_series.assert_not_awaited() + live.get_issues_for_series.assert_not_awaited() + + +async def test_empty_local_search_does_not_retry_or_sleep(tmp_path, monkeypatch): + provider = CatalogLookupService(installed_reader(tmp_path)) + sleep = AsyncMock() + monkeypatch.setattr("pullbox.services.import_cv_search.asyncio.sleep", sleep) + assert await search_with_retry(provider, "No matching book", None) == [] + sleep.assert_not_awaited() + + +async def test_ui_search_works_without_comicvine_key(tmp_path, monkeypatch): + reader = installed_reader(tmp_path) + monkeypatch.setattr("pullbox.services.catalog.reader.get_catalog_reader", lambda: reader) + monkeypatch.setattr( + "pullbox.core.comicvine_key.get_comicvine_api_key", AsyncMock(return_value="") + ) + async with open_comicvine_ui_provider(AsyncMock(), prefer_catalog=True) as lookup: + results, _ = await lookup.search_series_globally("Batman") + assert results[0].provider_id == "10" + + +async def test_catalog_provenance_and_cutoff_persist_for_new_records(tmp_path, db_session): + reader = installed_reader(tmp_path) + service = MetadataService(AsyncMock(), tmp_path) + service._catalog = reader + meta = await reader.series(10) + series = await service.upsert_series_metadata(db_session, 10, meta) + assert series.metadata_source == "pullbox_catalog" + assert series.metadata_last_refreshed == meta.source_cutoff_at + issues = await service.upsert_issue_summaries(db_session, series, await reader.issues(10)) + assert issues[0].metadata_source == "pullbox_catalog" + + +async def test_post_import_enrichment_keeps_live_batching(tmp_path): + reader = installed_reader(tmp_path) + + class Live: + def __init__(self): + self.calls = [] + + async def get_issue_batch(self, ids): + self.calls.append(ids) + return {} + + live = Live() + service = MetadataService(live, tmp_path) + service._catalog = reader + assert await service.prefetch_issue_metadata_batch([100]) == {} + assert live.calls == [["100"]] + + +async def test_basic_catalog_does_not_demote_full_metadata(tmp_path, db_session): + reader = installed_reader(tmp_path) + service = MetadataService(AsyncMock(), tmp_path) + meta = await reader.series(10) + series = await service.upsert_series_metadata(db_session, 10, meta) + series.metadata_source = "comicvine" + series.description = "Complete live description" + series.title = "Newer live title" + original_date = series.metadata_last_refreshed + await service.upsert_series_metadata(db_session, 10, meta) + assert series.metadata_source == "comicvine" + assert series.title == "Newer live title" + assert series.description == "Complete live description" + assert series.metadata_last_refreshed == original_date + + +async def test_basic_catalog_preserves_existing_live_issue_fields(tmp_path, db_session): + reader = installed_reader(tmp_path) + service = MetadataService(AsyncMock(), tmp_path) + series = await service.upsert_series_metadata(db_session, 10, await reader.series(10)) + summaries = await reader.issues(10) + issue = (await service.upsert_issue_summaries(db_session, series, summaries))[0] + issue.metadata_source = "comicvine" + issue.title = "Newer live issue title" + issue.issue_type = IssueType.SPECIAL + issue.issue_number = 1.0 + issue.issue_number_text = "1" + await db_session.flush() + assert await service.upsert_issue_summaries(db_session, series, summaries) == [] + assert issue.title == "Newer live issue title" + assert issue.issue_type == IssueType.SPECIAL + assert issue.issue_number_text == "1" + assert issue.metadata_source == "comicvine" + + +async def test_full_issue_refresh_bypasses_basic_catalog(tmp_path, db_session): + reader = installed_reader(tmp_path) + live = AsyncMock() + live.get_issues_for_series.return_value = [] + service = MetadataService(live, tmp_path, catalog=reader) + series = await service.upsert_series_metadata(db_session, 10, await reader.series(10)) + await service.fetch_issues_for_series(db_session, series.id) + live.get_issues_for_series.assert_awaited_once_with("10") + + +async def test_scan_provider_stack_keeps_search_and_identity_local( + tmp_path, db_session, monkeypatch +): + reader = installed_reader(tmp_path) + monkeypatch.setattr("pullbox.services.catalog.reader.get_catalog_reader", lambda: reader) + live = AsyncMock() + provider = build_import_scan_metadata_provider(db_session, live) + first = await search_with_retry(provider, "Batman", None) + second = await search_with_retry(provider, "Batman", None) + assert first == second + assert first[0].provider_id == "10" + assert (await provider.get_series("10")).title == "Batman" + assert (await provider.get_issue("100")).series_provider_id == "10" + assert (await provider.get_issues_for_series("10"))[0].provider_id == "100" + assert await search_with_retry(provider, "Missing series", None) == [] + assert provider.cache_metrics()["memory_hits"]["search_series_globally"] == 1 + assert live.mock_calls == [] + + +async def test_add_series_materializes_basic_catalog_without_live_metadata(tmp_path, db_session): + reader = installed_reader(tmp_path) + live = AsyncMock() + metadata = MetadataService(live, tmp_path, catalog=reader) + service = SeriesService(metadata, AsyncMock()) + series = await service.add_from_comicvine(db_session, 10) + issue = await db_session.scalar(select(Issue).where(Issue.series_id == series.id)) + assert series.comicvine_id == 10 + assert series.metadata_source == "pullbox_catalog" + assert issue.comicvine_id == 100 + assert issue.issue_number_text == "0.5" + assert issue.metadata_source == "pullbox_catalog" + assert live.mock_calls == [] diff --git a/tests/unit/test_catalog_reader.py b/tests/unit/test_catalog_reader.py new file mode 100644 index 00000000..4b1a1434 --- /dev/null +++ b/tests/unit/test_catalog_reader.py @@ -0,0 +1,63 @@ +"""Local matching uses aliases, exact IDs, and normalized issue designations.""" + +import json + +import pytest + +from pullbox.services.catalog.contract import CatalogError +from pullbox.services.catalog.reader import CatalogReader +from tests.catalog_fixtures import build_snapshot + + +def installed_reader(tmp_path): + root = tmp_path / "catalog" + (root / "bases").mkdir(parents=True) + build_snapshot(root / "bases/20260913T050000Z.db") + (root / "active.json").write_text( + json.dumps( + { + "path": "bases/20260913T050000Z.db", + "version": "20260913T050000Z", + "base_version": "20260913T050000Z", + "source_cutoff_at": "2026-09-13T05:00:00+00:00", + } + ) + ) + return CatalogReader(root) + + +async def test_alias_search_returns_existing_comicvine_identity(tmp_path): + reader = installed_reader(tmp_path) + results = await reader.search("Dark Knight") + assert len(results) == 1 + assert results[0].provider_id == "10" + assert results[0].title == "Batman" + assert results[0].publisher == "DC" + assert await reader.search("Batman", year=2024) == [] + + +async def test_search_input_is_not_fts_syntax_or_sql(tmp_path): + reader = installed_reader(tmp_path) + assert await reader.search('" OR * )') == [] + assert await reader.search("'; DROP TABLE series; --") == [] + assert (await reader.series(10)).title == "Batman" + + +async def test_issue_catalog_preserves_fractional_identity_and_cutoff(tmp_path): + reader = installed_reader(tmp_path) + issues = await reader.issues(10) + assert len(issues) == 1 + assert issues[0].provider_id == "100" + assert issues[0].issue_number == 0.5 + assert issues[0].issue_number_text == "0.5" + assert issues[0].source_cutoff_at.isoformat() == "2026-09-13T05:00:00+00:00" + assert await reader.series(999) is None + + +async def test_rejects_active_pointer_outside_catalog(tmp_path): + reader = installed_reader(tmp_path) + (reader.root / "active.json").write_text( + json.dumps({"path": "../../user.db", "version": "20260913T050000Z"}) + ) + with pytest.raises(CatalogError): + await reader.search("Batman") diff --git a/tests/unit/test_catalog_service.py b/tests/unit/test_catalog_service.py new file mode 100644 index 00000000..1e05bb38 --- /dev/null +++ b/tests/unit/test_catalog_service.py @@ -0,0 +1,204 @@ +"""Downloads remain resumable and a failed update preserves the active catalog.""" + +import asyncio +import hashlib + +import httpx +import pytest +import zstandard +from cryptography.hazmat.primitives.asymmetric.ed25519 import Ed25519PrivateKey + +from pullbox.services.catalog.contract import CatalogError +from pullbox.services.catalog.service import CatalogService +from tests.catalog_fixtures import build_patch, build_snapshot +from tests.unit.test_catalog_contract import publication_payload, signed_publication + + +class Server: + def __init__(self, tmp_path): + self.key = Ed25519PrivateKey.generate() + self.payload = publication_payload() + self.requests = [] + self.artifacts = {} + self.etag = '"publication-1"' + self.fail_download = False + self.base = tmp_path / "producer.db" + build_snapshot(self.base) + self.publish(self.base) + + def publish(self, path, *, base=None, version="20260913T050000Z"): + data = zstandard.ZstdCompressor().compress(path.read_bytes()) + url = ( + f"/api/v2/catalog/patches/{base}/{version}" + if base + else f"/api/v2/catalog/snapshots/{version}" + ) + item = { + "sha256": hashlib.sha256(data).hexdigest(), + "size_bytes": len(data), + "download_path": url, + } + item.update( + {"base_version": base, "target_version": version} if base else {"version": version} + ) + if base: + self.payload["patches"] = [item] + else: + self.payload["full_snapshot"] = item + self.payload["latest_version"] = version + self.artifacts[url] = data + self.etag = f'"{version}"' + + async def handle(self, request): + self.requests.append(request) + if request.url.path.endswith("/latest"): + if request.headers.get("if-none-match") == self.etag: + return httpx.Response(304) + return httpx.Response( + 200, content=signed_publication(self.payload, self.key), headers={"etag": self.etag} + ) + if self.fail_download: + return httpx.Response(503) + data = self.artifacts[request.url.path] + start = int(request.headers.get("range", "bytes=0-").split("=")[1].split("-")[0]) + headers = {"etag": f'"{hashlib.sha256(data).hexdigest()}"'} + if start: + headers["content-range"] = f"bytes {start}-{len(data) - 1}/{len(data)}" + return httpx.Response(206 if start else 200, content=data[start:], headers=headers) + + def service(self, root): + return CatalogService( + root, + "https://catalog.example", + keys={"test": self.key.public_key()}, + transport=httpx.MockTransport(self.handle), + ) + + +async def test_first_download_requires_opt_in_then_installs_verified_catalog(tmp_path): + server = Server(tmp_path) + service = server.service(tmp_path / "local") + assert await service.sync() is False + assert server.requests == [] + assert await service.sync(manual=True) is True + state = service.status() + assert state.installed_version == "20260913T050000Z" + assert state.phase == "current" + assert (service.root / "active.json").is_file() + count = len(server.requests) + assert await service.sync() is False + assert len(server.requests) == count + + +async def test_manual_retry_repairs_missing_installed_file(tmp_path): + server = Server(tmp_path) + service = server.service(tmp_path / "local") + await service.sync(manual=True) + installed = service.root / "bases/20260913T050000Z.db" + installed.unlink() + assert await service.sync(manual=True) is True + assert installed.is_file() + + +async def test_cancellation_releases_lock_and_allows_retry(tmp_path): + server = Server(tmp_path) + service = server.service(tmp_path / "local") + started = asyncio.Event() + original = server.handle + + async def waiting(request): + started.set() + await asyncio.Event().wait() + + service.transport = httpx.MockTransport(waiting) + task = asyncio.create_task(service.sync(manual=True)) + await started.wait() + task.cancel() + with pytest.raises(asyncio.CancelledError): + await task + assert service.status().phase == "interrupted" + assert not (service.root / "active.json").exists() + service.transport = httpx.MockTransport(original) + assert await service.sync(manual=True) is True + + +async def test_disabling_daily_updates_survives_process_restart(tmp_path): + server = Server(tmp_path) + service = server.service(tmp_path / "local") + await service.sync(manual=True) + await service.set_automatic_updates(False) + restored = server.service(service.root) + assert restored.status().automatic_updates is False + server.requests.clear() + assert await restored.sync() is False + assert server.requests == [] + + +async def test_success_discards_compressed_downloads_and_stale_staging(tmp_path): + server = Server(tmp_path) + service = server.service(tmp_path / "local") + staging = service.root / "staging" + staging.mkdir(parents=True) + (staging / "catalog-abandoned.db").write_bytes(b"incomplete") + await service.sync(manual=True) + assert list((service.root / "downloads").iterdir()) == [] + assert list(staging.iterdir()) == [] + + +async def test_resumes_partial_artifact_using_signed_checksum_as_identity(tmp_path): + server = Server(tmp_path) + root = tmp_path / "local" + download = root / "downloads" + download.mkdir(parents=True) + artifact = server.payload["full_snapshot"] + data = server.artifacts[artifact["download_path"]] + (download / f"{artifact['sha256']}.part").write_bytes(data[:100]) + assert await server.service(root).sync(manual=True) is True + assert server.requests[-1].headers["range"] == "bytes=100-" + + +async def test_installs_cumulative_patch_without_redownloading_weekly_base(tmp_path): + server = Server(tmp_path) + service = server.service(tmp_path / "local") + assert await service.sync(manual=True) is True + target, patch = tmp_path / "next.db", tmp_path / "patch.db" + build_snapshot(target, "20260914T050000Z", "Updated") + build_patch(patch, server.base, target) + server.publish(patch, base="20260913T050000Z", version="20260914T050000Z") + assert await service.sync(manual=True) is True + assert service.status().installed_version == "20260914T050000Z" + assert len([r for r in server.requests if "/snapshots/" in r.url.path]) == 1 + + +async def test_failed_update_keeps_active_catalog_and_safe_error(tmp_path): + server = Server(tmp_path) + service = server.service(tmp_path / "local") + assert await service.sync(manual=True) is True + active = (service.root / "active.json").read_bytes() + target = tmp_path / "next.db" + build_snapshot(target, "20260920T050000Z") + server.publish(target, version="20260920T050000Z") + server.fail_download = True + with pytest.raises(CatalogError): + await service.sync(manual=True) + assert (service.root / "active.json").read_bytes() == active + assert service.status().installed_version == "20260913T050000Z" + assert service.status().phase == "failed" + + +async def test_corrupt_download_never_becomes_active(tmp_path): + server = Server(tmp_path) + artifact = server.payload["full_snapshot"] + server.artifacts[artifact["download_path"]] = b"x" * artifact["size_bytes"] + service = server.service(tmp_path / "local") + with pytest.raises(CatalogError, match="checksum"): + await service.sync(manual=True) + assert not (service.root / "active.json").exists() + + +async def test_manual_and_scheduled_attempts_do_not_overlap(tmp_path): + server = Server(tmp_path) + service = server.service(tmp_path / "local") + results = await asyncio.gather(service.sync(manual=True), service.sync(manual=True)) + assert sorted(results) == [False, True] + assert len([r for r in server.requests if "/snapshots/" in r.url.path]) == 1 diff --git a/tests/unit/test_catalog_storage.py b/tests/unit/test_catalog_storage.py new file mode 100644 index 00000000..1c13e190 --- /dev/null +++ b/tests/unit/test_catalog_storage.py @@ -0,0 +1,29 @@ +"""Realistic zstd windows and expansion limits protect the installed catalog.""" + +import pytest +import zstandard + +from pullbox.services.catalog import storage +from pullbox.services.catalog.contract import CatalogError + + +def test_streams_production_sized_zstd_window(tmp_path): + data = b"catalog" * 500_000 + archive = tmp_path / "snapshot.zst" + archive.write_bytes(zstandard.ZstdCompressor(level=3).compress(data)) + output = tmp_path / "snapshot.db" + error = None + try: + storage.decompress(archive, output) + except CatalogError as exc: + error = str(exc) + assert error is None + assert output.read_bytes() == data + + +def test_enforces_expansion_ceiling(tmp_path, monkeypatch): + archive = tmp_path / "snapshot.zst" + archive.write_bytes(zstandard.ZstdCompressor().compress(b"a" * 10000)) + monkeypatch.setattr(storage, "MAX_DATABASE_BYTES", 100) + with pytest.raises(CatalogError, match="storage size"): + storage.decompress(archive, tmp_path / "snapshot.db") diff --git a/tests/unit/test_catalog_task.py b/tests/unit/test_catalog_task.py new file mode 100644 index 00000000..97cb4d29 --- /dev/null +++ b/tests/unit/test_catalog_task.py @@ -0,0 +1,13 @@ +"""Catalog scheduling respects explicit first-download consent.""" + +from unittest.mock import AsyncMock + +from pullbox.services.catalog import service +from pullbox.tasks.catalog_task import update_catalog + + +async def test_task_delegates_automatic_policy(monkeypatch): + catalog = AsyncMock() + monkeypatch.setattr(service, "get_catalog_service", lambda: catalog) + await update_catalog() + catalog.sync.assert_awaited_once_with(manual=False) diff --git a/tests/unit/test_import_service.py b/tests/unit/test_import_service.py index 0032f85a..5590f4ca 100644 --- a/tests/unit/test_import_service.py +++ b/tests/unit/test_import_service.py @@ -5079,7 +5079,7 @@ async def test_override_sets_user_cv_id( issue_count=85, comicvine_url="https://comicvine.gamespot.com/batman/4050-12345/", ) - mock_metadata_service._provider.get_series = AsyncMock(return_value=meta) + mock_metadata_service.get_series_metadata = AsyncMock(return_value=meta) svc = ImportService( series_service=mock_series_service, @@ -5104,6 +5104,7 @@ async def test_override_sets_user_cv_id( assert updated.status == ImportSeriesStatus.MATCHED assert updated.cv_match_method == "user_override" assert updated.diagnostics == {} + mock_metadata_service.get_series_metadata.assert_awaited_once_with(12345) @pytest.mark.asyncio async def test_override_recomputes_pending_file_matches( @@ -5129,7 +5130,7 @@ async def test_override_recomputes_pending_file_matches( issue_count=10, comicvine_url="https://comicvine.gamespot.com/absolute-martian-manhunter/4050-162966/", ) - mock_metadata_service._provider.get_series = AsyncMock(return_value=meta) + mock_metadata_service.get_series_metadata = AsyncMock(return_value=meta) mock_metadata_service._provider.get_issues_for_series = AsyncMock( return_value=[ IssueSummary( @@ -5199,7 +5200,7 @@ async def test_override_reconsolidates_matching_logical_group( issue_count=4, comicvine_url="https://comicvine.gamespot.com/chicken-devil/4050-139451/", ) - mock_metadata_service._provider.get_series = AsyncMock(return_value=meta) + mock_metadata_service.get_series_metadata = AsyncMock(return_value=meta) mock_metadata_service._provider.get_issues_for_series = AsyncMock( return_value=[ IssueSummary( From a3c502be93c75f952342b1635ee23b5a5ee6581c Mon Sep 17 00:00:00 2001 From: Adam Hernandez Date: Sun, 13 Sep 2026 17:32:51 -0700 Subject: [PATCH 2/4] fix: harden local catalog recovery and settings order Keep matching profiles when a catalog batch contains missing IDs. Recover installed catalog identity independently of status state and retain update opt-in after status loss. Restore ComicVine access before the local catalog card. --- src/pullbox/services/catalog/service.py | 19 ++++++--- src/pullbox/services/metadata_service.py | 9 +++-- .../templates/partials/settings_metadata.html | 3 +- tests/ui/test_catalog_controls.py | 11 +++++ tests/unit/test_catalog_integration.py | 13 ++++++ tests/unit/test_catalog_service.py | 40 +++++++++++++++++++ 6 files changed, 86 insertions(+), 9 deletions(-) diff --git a/src/pullbox/services/catalog/service.py b/src/pullbox/services/catalog/service.py index e0b274d5..99defaf9 100644 --- a/src/pullbox/services/catalog/service.py +++ b/src/pullbox/services/catalog/service.py @@ -97,14 +97,23 @@ def __init__( if self._state.phase in BUSY_PHASES: self._state.phase = "interrupted" self._state.error = "The previous download was interrupted. It can be resumed." - active = load_json(root / "active.json") - if active: - self._state.installed_version = str(active["version"]) - self._state.source_cutoff_at = str(active["source_cutoff_at"]) - except (CatalogError, ValidationError, KeyError): + except (CatalogError, ValidationError): self._state = CatalogStatus( phase="failed", error="Catalog state could not be read. Check the data volume." ) + # The active generation survives a missing or damaged optional status file. + try: + active = load_json(root / "active.json") + if active: + version, cutoff = str(active["version"]), str(active["source_cutoff_at"]) + self._state.installed_version = version + self._state.source_cutoff_at = cutoff + self._state.requested = True + if self._state.phase == "not_downloaded": + self._state.phase = "current" + except (CatalogError, KeyError): + self._state.phase = "failed" + self._state.error = "Catalog state could not be read. Check the data volume." def status(self) -> CatalogStatus: return self._state.model_copy(deep=True) diff --git a/src/pullbox/services/metadata_service.py b/src/pullbox/services/metadata_service.py index 8f1a7e48..6bc8aa82 100644 --- a/src/pullbox/services/metadata_service.py +++ b/src/pullbox/services/metadata_service.py @@ -279,9 +279,12 @@ async def get_series_metadata_batch( ) -> dict[int, SeriesMetadata]: """Fetch multiple provider series profiles through the optional bulk contract.""" if self._catalog is not None and self._catalog.available: - return { - key: await self.get_series_metadata(key) for key in dict.fromkeys(comicvine_ids) - } + profiles: dict[int, SeriesMetadata] = {} + for key in dict.fromkeys(comicvine_ids): + profile = await self._catalog.series(key) + if profile is not None: + profiles[key] = profile + return profiles batch_fetch = getattr(type(self._provider), "get_series_batch", None) try: if callable(batch_fetch): diff --git a/src/pullbox/ui/templates/partials/settings_metadata.html b/src/pullbox/ui/templates/partials/settings_metadata.html index 67b9bc97..67c69d71 100644 --- a/src/pullbox/ui/templates/partials/settings_metadata.html +++ b/src/pullbox/ui/templates/partials/settings_metadata.html @@ -2,7 +2,6 @@ {% from "components/settings_shell.html" import field_note, inline_alert, section_card, settings_footer, settings_row, ui_icon %}
- {% include "partials/settings_catalog.html" %} {% call section_card( title="ComicVine access", eyebrow="API credentials", @@ -80,6 +79,8 @@ {% endcall %} + {% include "partials/settings_catalog.html" %} + {% call section_card( title="Refresh policy", eyebrow="Sync cadence", diff --git a/tests/ui/test_catalog_controls.py b/tests/ui/test_catalog_controls.py index 33571c53..3f47699e 100644 --- a/tests/ui/test_catalog_controls.py +++ b/tests/ui/test_catalog_controls.py @@ -65,6 +65,17 @@ async def test_metadata_settings_explains_local_storage_and_progress(authenticat assert "A download never moves or imports comics." in response.text +@pytest.mark.parametrize("url", ["/settings?tab=metadata", "/htmx/settings/metadata"]) +async def test_metadata_settings_places_access_before_local_catalog(authenticated_client, url): + response = await authenticated_client.get(url) + assert response.status_code == 200 + assert ( + response.text.index("ComicVine access") + < response.text.index("Local Comic Vine catalog") + < response.text.index("Refresh policy") + ) + + async def test_add_series_labels_local_search_without_api_key( authenticated_client, tmp_path, monkeypatch ): diff --git a/tests/unit/test_catalog_integration.py b/tests/unit/test_catalog_integration.py index 1cfd9fb5..7798ef6c 100644 --- a/tests/unit/test_catalog_integration.py +++ b/tests/unit/test_catalog_integration.py @@ -2,6 +2,7 @@ from unittest.mock import AsyncMock +import pytest from sqlalchemy import select from pullbox.models.issue import Issue, IssueType @@ -27,6 +28,18 @@ async def test_metadata_hydration_uses_catalog_without_provider_calls(tmp_path): live.get_issues_for_series.assert_not_awaited() +@pytest.mark.parametrize("ids", [[999, 10, 10], [10, 999]]) +async def test_catalog_profile_batch_preserves_matches_when_a_series_is_missing(tmp_path, ids): + live = AsyncMock() + service = MetadataService(live, tmp_path, catalog=installed_reader(tmp_path)) + + profiles = await service.get_series_metadata_batch(ids) + + assert list(profiles) == [10] + assert profiles[10].title == "Batman" + assert live.mock_calls == [] + + async def test_empty_local_search_does_not_retry_or_sleep(tmp_path, monkeypatch): provider = CatalogLookupService(installed_reader(tmp_path)) sleep = AsyncMock() diff --git a/tests/unit/test_catalog_service.py b/tests/unit/test_catalog_service.py index 1e05bb38..0c48cf3e 100644 --- a/tests/unit/test_catalog_service.py +++ b/tests/unit/test_catalog_service.py @@ -134,6 +134,46 @@ async def test_disabling_daily_updates_survives_process_restart(tmp_path): assert server.requests == [] +@pytest.mark.parametrize("saved_state", ["{", '{"requested": {}}', None]) +async def test_installed_catalog_survives_missing_or_invalid_status_state(tmp_path, saved_state): + server = Server(tmp_path) + service = server.service(tmp_path / "local") + await service.sync(manual=True) + installed = service.status() + active = (service.root / "active.json").read_bytes() + state_path = service.root / "state.json" + if saved_state is None: + state_path.unlink() + else: + state_path.write_text(saved_state) + server.requests.clear() + + restored = server.service(service.root) + + assert restored.status().installed_version == installed.installed_version + assert restored.status().source_cutoff_at == installed.source_cutoff_at + assert restored.status().requested is True + assert await restored.sync() is False # Current generation needs no download. + assert [request.url.path for request in server.requests] == ["/api/v2/catalog/latest"] + assert restored.status().phase == "current" + assert restored.status().error is None + assert (service.root / "active.json").read_bytes() == active + + +async def test_invalid_status_without_an_installed_catalog_does_not_opt_in(tmp_path): + server = Server(tmp_path) + root = tmp_path / "local" + root.mkdir() + (root / "state.json").write_text("{") + + restored = server.service(root) + + assert restored.status().requested is False + assert restored.status().installed_version is None + assert await restored.sync() is False + assert server.requests == [] + + async def test_success_discards_compressed_downloads_and_stale_staging(tmp_path): server = Server(tmp_path) service = server.service(tmp_path / "local") From a468d14a3bdda38d58b8612494d70cfaa73dfe5a Mon Sep 17 00:00:00 2001 From: Adam Hernandez Date: Sun, 13 Sep 2026 18:32:14 -0700 Subject: [PATCH 3/4] fix: preserve new focus choices during story arc reordering --- src/pullbox/ui/static/js/story-arc-preview.js | 4 + tests/e2e/test_story_arc_catalog_page.py | 73 +++++++++++++++++++ 2 files changed, 77 insertions(+) diff --git a/src/pullbox/ui/static/js/story-arc-preview.js b/src/pullbox/ui/static/js/story-arc-preview.js index f3df4771..9ebab321 100644 --- a/src/pullbox/ui/static/js/story-arc-preview.js +++ b/src/pullbox/ui/static/js/story-arc-preview.js @@ -41,6 +41,7 @@ function storyArcPreview() { const index = this.members.findIndex(member => member.provider_id === providerId); const target = index + direction; if (index < 0 || target < 0 || target >= this.members.length) return; + const previousFocus = document.activeElement; // Replace the array in one pass so keyed rows never see duplicate IDs. const reordered = [...this.members]; [reordered[index], reordered[target]] = [reordered[target], reordered[index]]; @@ -50,6 +51,9 @@ function storyArcPreview() { const moved = this.members[target]; this.reorderAnnouncement = `${moved.series_name} #${moved.issue_number} moved to position ${target + 1}.`; this.$nextTick(() => { + // A newer keyboard/pointer focus choice takes precedence over restoring + // focus lost while Alpine moves or replaces the keyed row. + if (document.activeElement !== previousFocus && document.activeElement !== document.body) return; const row = this.$root.querySelector(`[data-provider-issue-id="${providerId}"]`); const button = row?.querySelector(`[data-order-direction="${direction < 0 ? 'up' : 'down'}"]`); const focusTarget = button?.disabled ? row.querySelector('[data-order-direction]:not(:disabled)') : button; diff --git a/tests/e2e/test_story_arc_catalog_page.py b/tests/e2e/test_story_arc_catalog_page.py index 8de7f18a..bdc08448 100644 --- a/tests/e2e/test_story_arc_catalog_page.py +++ b/tests/e2e/test_story_arc_catalog_page.py @@ -579,6 +579,79 @@ def test_preview_pagination_retry_and_add_keep_all_page_choices( assert errors == [] +@pytest.mark.parametrize("provider_id, initial_position", [("124", 24), ("125", 25)]) +@pytest.mark.parametrize("next_control", ["up", "monitor"]) +def test_preview_reorder_does_not_override_new_keyboard_focus( + authed_page: Page, + seeded_server: str, + catalog_provider: CatalogProvider, + provider_id: str, + initial_position: int, + next_control: str, +) -> None: + page = authed_page + catalog_provider.metadata = replace( + catalog_provider.metadata, + issue_provider_ids=tuple(str(number) for number in range(101, 152)), + declared_issue_count=51, + ) + page.goto(f"{seeded_server}/story-arcs/catalog/177") + row = page.locator(f'[data-provider-issue-id="{provider_id}"]') + expect(row).to_be_visible() + expect(row.locator("[data-reading-position]")).to_have_text(str(initial_position)) + next_selector = ( + f'[data-provider-issue-id="{provider_id}"] [data-order-direction="up"]' + if next_control == "up" + else "#catalog-monitored" + ) + # Put the user's next focus change between Alpine's DOM update and its + # deferred focus restoration, including a move across the page boundary. + rendered_position = row.evaluate( + """async (row, nextSelector) => { + const id = row.dataset.providerIssueId; + const down = row.querySelector('[data-order-direction="down"]'); + down.focus(); + down.click(); + await Promise.resolve(); + const moved = document.querySelector(`[data-provider-issue-id="${id}"]`); + document.querySelector(nextSelector).focus(); + const position = moved.querySelector('[data-reading-position]').textContent; + await Alpine.nextTick(); + return position; + }""", + next_selector, + ) + assert rendered_position == str(initial_position + 1) + expect(page.locator(next_selector)).to_be_focused() + if next_control == "monitor": + page.keyboard.press("Space") + expect(page.get_by_role("switch", name="Monitor this story arc")).to_be_checked() + expect(row.locator("[data-reading-position]")).to_have_text(str(initial_position + 1)) + return + page.keyboard.press("Enter") + expect(row.locator("[data-reading-position]")).to_have_text(str(initial_position)) + + +def test_preview_reorder_restores_focus_to_enabled_direction_at_list_ends( + authed_page: Page, seeded_server: str, catalog_provider: CatalogProvider +) -> None: + page = authed_page + page.goto(f"{seeded_server}/story-arcs/catalog/178") + row = page.locator('[data-provider-issue-id="101"]') + expect(row).to_be_visible() + down = row.locator('[data-order-direction="down"]') + up = row.locator('[data-order-direction="up"]') + expect(row.locator("[data-reading-position]")).to_have_text("1") + down.press("Enter") + expect(row.locator("[data-reading-position]")).to_have_text("2") + expect(down).to_be_disabled() + expect(up).to_be_focused() + page.keyboard.press("Enter") + expect(row.locator("[data-reading-position]")).to_have_text("1") + expect(up).to_be_disabled() + expect(down).to_be_focused() + + def test_preview_reorder_single_member_and_pending_retry_are_disabled( authed_page: Page, seeded_server: str, catalog_provider: CatalogProvider ) -> None: From 959768a3cc3b59915f8910d4e9579d50713e6ffe Mon Sep 17 00:00:00 2001 From: Adam Hernandez Date: Sun, 13 Sep 2026 19:09:30 -0700 Subject: [PATCH 4/4] fix: prevent stale dropdown focus restoration --- src/pullbox/ui/static/js/pullbox.js | 9 +++------ tests/e2e/test_story_arc_catalog_page.py | 22 ++++++++++++++++++++++ 2 files changed, 25 insertions(+), 6 deletions(-) diff --git a/src/pullbox/ui/static/js/pullbox.js b/src/pullbox/ui/static/js/pullbox.js index e3c00b97..20544cee 100644 --- a/src/pullbox/ui/static/js/pullbox.js +++ b/src/pullbox/ui/static/js/pullbox.js @@ -14887,6 +14887,9 @@ function dropdownSelectData(config) { this.currentLabel = option.label; this.syncInput(); this.close(false); + if (this.$refs && this.$refs.trigger) { + this.$refs.trigger.focus(); + } this.runChangeExpression(); @@ -14906,12 +14909,6 @@ function dropdownSelectData(config) { ); } - var self = this; - this.$nextTick(function () { - if (self.$refs && self.$refs.trigger) { - self.$refs.trigger.focus(); - } - }); }, selectActive: function () { diff --git a/tests/e2e/test_story_arc_catalog_page.py b/tests/e2e/test_story_arc_catalog_page.py index bdc08448..bc6d2301 100644 --- a/tests/e2e/test_story_arc_catalog_page.py +++ b/tests/e2e/test_story_arc_catalog_page.py @@ -652,6 +652,28 @@ def test_preview_reorder_restores_focus_to_enabled_direction_at_list_ends( expect(down).to_be_focused() +def test_dropdown_selection_does_not_restore_focus_after_user_moves_on( + authed_page: Page, seeded_server: str, catalog_provider: CatalogProvider +) -> None: + page = authed_page + page.goto(f"{seeded_server}/story-arcs/catalog/179") + root = page.get_by_label("Library root for new series") + root.click() + option = page.locator("[data-dropdown-select-panel]:visible [data-dropdown-option]").nth(1) + move_down = page.locator('[data-provider-issue-id="101"] [data-order-direction="down"]') + option.evaluate( + """async (option, targetSelector) => { + option.click(); + document.querySelector(targetSelector).focus(); + await Alpine.nextTick(); + }""", + '[data-provider-issue-id="101"] [data-order-direction="down"]', + ) + expect(move_down).to_be_focused() + page.keyboard.press("Enter") + expect(page.locator('[data-provider-issue-id="101"] [data-reading-position]')).to_have_text("2") + + def test_preview_reorder_single_member_and_pending_retry_are_disabled( authed_page: Page, seeded_server: str, catalog_provider: CatalogProvider ) -> None: