diff --git a/README.md b/README.md index d21381f..c92bf84 100644 --- a/README.md +++ b/README.md @@ -78,6 +78,12 @@ source search or payload download. Live acceptance probes must remain metadata-only, use credentials supplied at runtime, and never persist signed URLs or account data. +LibGen's authenticated `/v1/health` endpoint checks its five known mirrors +concurrently, once each, with a five-second probe timeout and a one-second +session-cleanup timeout. Slow mirrors report `unavailable`; process health stays +separate from source availability. Cancelling the request cancels and cleans up +the outstanding probes. Search and download-resolution timeouts are unchanged. + Pull requests run four stable aggregate checks: `CI Required`, `Security Required`, `Workflow Hygiene Required`, and `Container Security Required`. They run on GitHub-hosted runners with read-only @@ -249,8 +255,8 @@ either registry; both names resolve to the same signed digest. | Provider | GHCR | Docker Hub | | --- | --- | --- | | GetComics | `ghcr.io/pullboxapp/pullbox-provider-getcomics:1.0.3` | `docker.io/pullbox/pullbox-provider-getcomics:1.0.3` | -| Anna's Archive | `ghcr.io/pullboxapp/pullbox-provider-annas-archive:1.0.3` | `docker.io/pullbox/pullbox-provider-annas-archive:1.0.3` | -| LibGen | `ghcr.io/pullboxapp/pullbox-provider-libgen:1.0.1` | `docker.io/pullbox/pullbox-provider-libgen:1.0.1` | +| Anna's Archive | `ghcr.io/pullboxapp/pullbox-provider-annas-archive:1.0.4` | `docker.io/pullbox/pullbox-provider-annas-archive:1.0.4` | +| LibGen | `ghcr.io/pullboxapp/pullbox-provider-libgen:1.0.2` | `docker.io/pullbox/pullbox-provider-libgen:1.0.2` | Pin a numbered version or the immutable digest in production. `latest` tracks only the newest stable provider release; prerelease and manual `edge` builds do diff --git a/providers/annas_archive/pyproject.toml b/providers/annas_archive/pyproject.toml index 52dbaef..15416f5 100644 --- a/providers/annas_archive/pyproject.toml +++ b/providers/annas_archive/pyproject.toml @@ -4,13 +4,13 @@ build-backend = "hatchling.build" [project] name = "pullbox-provider-annas-archive" -version = "1.0.3" +version = "1.0.4" description = "Optional Anna's Archive discovery provider for Pullbox" requires-python = ">=3.12" license = "GPL-3.0-or-later" dependencies = [ "pullbox-direct-provider-contract==1.0.0", - "pullbox-provider-libgen==1.0.1", + "pullbox-provider-libgen==1.0.2", "uvicorn>=0.34,<1", ] diff --git a/providers/libgen/pyproject.toml b/providers/libgen/pyproject.toml index c227c46..e5228ae 100644 --- a/providers/libgen/pyproject.toml +++ b/providers/libgen/pyproject.toml @@ -4,7 +4,7 @@ build-backend = "hatchling.build" [project] name = "pullbox-provider-libgen" -version = "1.0.1" +version = "1.0.2" description = "Optional LibGen discovery provider for Pullbox" requires-python = ">=3.12" license = "GPL-3.0-or-later" diff --git a/providers/libgen/src/pullbox_provider_libgen/service.py b/providers/libgen/src/pullbox_provider_libgen/service.py index 3090179..5a67f95 100644 --- a/providers/libgen/src/pullbox_provider_libgen/service.py +++ b/providers/libgen/src/pullbox_provider_libgen/service.py @@ -8,6 +8,7 @@ import socket import time from collections.abc import Awaitable, Callable, Mapping, Sequence +from contextlib import suppress from typing import Protocol from urllib.parse import urlencode, urlsplit, urlunsplit @@ -58,6 +59,9 @@ "home.arpa", ) _MAX_METADATA_BYTES = 512 * 1024 +_SOURCE_HEALTH_TIMEOUT_SECONDS = 5.0 +_SOURCE_HEALTH_CLOSE_TIMEOUT_SECONDS = 1.0 +_DETACHED_HEALTH_CLOSE_TASKS: set[asyncio.Task[None]] = set() _LOGGER = structlog.get_logger(__name__) SourceResolver = Callable[[str, int], Awaitable[Sequence[str]]] @@ -79,6 +83,41 @@ class LibGenSourceOriginError(ValueError): """The configured LibGen source origin is unsafe or unavailable.""" +def _consume_health_close_result(task: asyncio.Task[None]) -> None: + """Retrieve a detached close task's result without delaying health responses.""" + _DETACHED_HEALTH_CLOSE_TASKS.discard(task) + with suppress(asyncio.CancelledError, Exception): + task.result() + + +def _cancel_health_close_task(task: asyncio.Task[None]) -> None: + """Cancel cleanup without trusting the transport to honor cancellation.""" + _DETACHED_HEALTH_CLOSE_TASKS.add(task) + task.cancel() + task.add_done_callback(_consume_health_close_result) + + +async def _close_health_session(session: SourceSession) -> bool: + """Close one health session within its budget, even if the request is cancelled.""" + close_task = asyncio.create_task(session.aclose()) + deadline = asyncio.get_running_loop().time() + _SOURCE_HEALTH_CLOSE_TIMEOUT_SECONDS + try: + try: + async with asyncio.timeout_at(deadline): + await asyncio.shield(close_task) + except asyncio.CancelledError: + try: + async with asyncio.timeout_at(deadline): + await asyncio.shield(close_task) + except TimeoutError: + _cancel_health_close_task(close_task) + raise + except TimeoutError: + _cancel_health_close_task(close_task) + return False + return True + + async def validate_source_origin( raw_url: str, *, @@ -237,27 +276,37 @@ def __init__( ) async def source_health(self) -> dict[str, ProviderStatus]: - health: dict[str, ProviderStatus] = {} - for origin in KNOWN_SOURCE_URLS: - session = self._session_factory(origin, None) - try: + # Probe each fixed mirror once; a slow mirror must not delay the others. + async with asyncio.TaskGroup() as group: + probes = { + origin: group.create_task(self._source_health_probe(origin)) + for origin in KNOWN_SOURCE_URLS + } + return { + urlsplit(origin).hostname or origin: task.result() for origin, task in probes.items() + } + + async def _source_health_probe(self, origin: str) -> ProviderStatus: + session = self._session_factory(origin, None) + try: + async with asyncio.timeout(_SOURCE_HEALTH_TIMEOUT_SECONDS): await session.fetch_text(f"{origin}/index.php") - except BrowserChallengeRequiredError: - status = ProviderStatus.CHALLENGE_REQUIRED - except LibGenSourceError as exc: - status = ( - ProviderStatus.RATE_LIMITED - if exc.code == "source_rate_limited" - else ProviderStatus.UNAVAILABLE - ) - except ProviderResolverError: + except BrowserChallengeRequiredError: + status = ProviderStatus.CHALLENGE_REQUIRED + except LibGenSourceError as exc: + status = ( + ProviderStatus.RATE_LIMITED + if exc.code == "source_rate_limited" + else ProviderStatus.UNAVAILABLE + ) + except (ProviderResolverError, TimeoutError): + status = ProviderStatus.UNAVAILABLE + else: + status = ProviderStatus.HEALTHY + finally: + if not await _close_health_session(session): status = ProviderStatus.UNAVAILABLE - else: - status = ProviderStatus.HEALTHY - finally: - await session.aclose() - health[urlsplit(origin).hostname or origin] = status - return health + return status async def search( self, diff --git a/tests/unit/test_libgen_health.py b/tests/unit/test_libgen_health.py new file mode 100644 index 0000000..fbebee9 --- /dev/null +++ b/tests/unit/test_libgen_health.py @@ -0,0 +1,252 @@ +from __future__ import annotations + +import asyncio + +import httpx +import pytest +from pullbox_provider_contract.models import ProviderStatus, ResolverProfile +from pullbox_provider_libgen import service as service_module +from pullbox_provider_libgen.app import create_app +from pullbox_provider_libgen.service import KNOWN_SOURCE_URLS, LibGenProviderService + +from tests.conftest import TEST_TOKEN + + +class _HealthSession: + def __init__( + self, + *, + slow_fetch: bool = False, + slow_close: bool = False, + close_release: asyncio.Event | None = None, + suppress_close_cancellation: bool = False, + ) -> None: + self.slow_fetch = slow_fetch + self.slow_close = slow_close + self.close_release = close_release + self.suppress_close_cancellation = suppress_close_cancellation + self.started = asyncio.Event() + self.close_started = asyncio.Event() + self.closed_event = asyncio.Event() + self.fetch_count = 0 + self.fetch_cancelled = False + self.close_called = False + self.close_cancelled = False + self.closed = False + + async def fetch_text(self, _url: str, *, max_bytes: int = 0) -> str: + self.fetch_count += 1 + self.started.set() + if self.slow_fetch: + try: + await asyncio.Event().wait() + except asyncio.CancelledError: + self.fetch_cancelled = True + raise + return "source is reachable" + + async def resolve_redirect(self, _url: str) -> str: + raise AssertionError("Health checks must not resolve downloads") + + async def aclose(self) -> None: + self.close_called = True + self.close_started.set() + if self.slow_close: + try: + await asyncio.Event().wait() + except asyncio.CancelledError: + self.close_cancelled = True + if self.suppress_close_cancellation: + assert self.close_release is not None + await self.close_release.wait() + self.closed = True + self.closed_event.set() + return + raise + if self.close_release is not None: + try: + await self.close_release.wait() + except asyncio.CancelledError: + self.close_cancelled = True + raise + self.closed = True + self.closed_event.set() + + +@pytest.fixture +def short_health_budget(monkeypatch: pytest.MonkeyPatch) -> None: + monkeypatch.setattr(service_module, "_SOURCE_HEALTH_TIMEOUT_SECONDS", 0.05, raising=False) + monkeypatch.setattr(service_module, "_SOURCE_HEALTH_CLOSE_TIMEOUT_SECONDS", 0.05, raising=False) + + +def _service(sessions: dict[str, _HealthSession]) -> LibGenProviderService: + def factory(origin: str, profile: ResolverProfile | None) -> _HealthSession: + assert profile is None + return sessions[origin] + + return LibGenProviderService(session_factory=factory) + + +@pytest.mark.parametrize("slow_origins", [(KNOWN_SOURCE_URLS[0],), KNOWN_SOURCE_URLS]) +async def test_source_health_bounds_slow_mirrors_and_preserves_fast_results( + short_health_budget: None, slow_origins: tuple[str, ...] +) -> None: + sessions = { + origin: _HealthSession(slow_fetch=origin in slow_origins) for origin in KNOWN_SOURCE_URLS + } + try: + async with asyncio.timeout(1): + health = await _service(sessions).source_health() + except TimeoutError: + pytest.fail("Source health exceeded its budget while waiting for an upstream mirror") + + assert health == { + origin.removeprefix("https://"): ( + ProviderStatus.UNAVAILABLE if origin in slow_origins else ProviderStatus.HEALTHY + ) + for origin in KNOWN_SOURCE_URLS + } + assert all(session.fetch_count == 1 and session.closed for session in sessions.values()) + assert all(sessions[origin].fetch_cancelled for origin in slow_origins) + + +async def test_source_health_bounds_session_cleanup(short_health_budget: None) -> None: + sessions = {origin: _HealthSession(slow_close=True) for origin in KNOWN_SOURCE_URLS} + try: + async with asyncio.timeout(1): + health = await _service(sessions).source_health() + except TimeoutError: + pytest.fail("Source health exceeded its budget while closing an upstream session") + + assert set(health.values()) == {ProviderStatus.UNAVAILABLE} + assert all(session.close_called and session.close_cancelled for session in sessions.values()) + + +async def test_source_health_does_not_wait_for_cleanup_that_suppresses_cancellation( + short_health_budget: None, +) -> None: + close_release = asyncio.Event() + sessions = { + origin: _HealthSession( + slow_close=True, + close_release=close_release, + suppress_close_cancellation=True, + ) + for origin in KNOWN_SOURCE_URLS + } + task = asyncio.create_task(_service(sessions).source_health()) + try: + await asyncio.wait_for( + asyncio.gather(*(session.close_started.wait() for session in sessions.values())), + timeout=1, + ) + await asyncio.sleep(0.15) + assert task.done(), "Source health waited indefinitely for cancellation-resistant cleanup" + finally: + close_release.set() + await asyncio.wait_for(task, timeout=1) + await asyncio.wait_for( + asyncio.gather(*(session.closed_event.wait() for session in sessions.values())), + timeout=1, + ) + + assert all(session.close_cancelled and session.closed for session in sessions.values()) + + +async def test_source_health_retains_detached_cleanup_until_completion( + short_health_budget: None, +) -> None: + close_release = asyncio.Event() + sessions = { + origin: _HealthSession( + slow_close=True, + close_release=close_release, + suppress_close_cancellation=True, + ) + for origin in KNOWN_SOURCE_URLS + } + task = asyncio.create_task(_service(sessions).source_health()) + try: + await asyncio.wait_for( + asyncio.gather(*(session.close_started.wait() for session in sessions.values())), + timeout=1, + ) + await asyncio.sleep(0.15) + assert task.done(), "Source health waited indefinitely for detached cleanup" + assert len(service_module._DETACHED_HEALTH_CLOSE_TASKS) == len(sessions) + assert all( + not close_task.done() for close_task in service_module._DETACHED_HEALTH_CLOSE_TASKS + ) + finally: + close_release.set() + await asyncio.wait_for(task, timeout=1) + await asyncio.wait_for( + asyncio.gather(*(session.closed_event.wait() for session in sessions.values())), + timeout=1, + ) + + await asyncio.sleep(0) + assert not service_module._DETACHED_HEALTH_CLOSE_TASKS + + +async def test_cancelling_health_cleans_up_all_probes(short_health_budget: None) -> None: + sessions = {origin: _HealthSession(slow_fetch=True) for origin in KNOWN_SOURCE_URLS} + task = asyncio.create_task(_service(sessions).source_health()) + try: + await asyncio.wait_for( + asyncio.gather(*(session.started.wait() for session in sessions.values())), timeout=1 + ) + except TimeoutError: + pytest.fail("A slow first mirror prevented the remaining health probes from starting") + finally: + task.cancel() + with pytest.raises(asyncio.CancelledError): + await task + + assert all(session.fetch_cancelled and session.closed for session in sessions.values()) + + +async def test_cancelling_health_during_cleanup_allows_bounded_close_to_finish( + short_health_budget: None, +) -> None: + close_release = asyncio.Event() + sessions = {origin: _HealthSession(close_release=close_release) for origin in KNOWN_SOURCE_URLS} + task = asyncio.create_task(_service(sessions).source_health()) + await asyncio.wait_for( + asyncio.gather(*(session.close_started.wait() for session in sessions.values())), timeout=1 + ) + + task.cancel() + await asyncio.sleep(0.01) + close_release.set() + with pytest.raises(asyncio.CancelledError): + await task + + assert all(session.closed and not session.close_cancelled for session in sessions.values()) + + +async def test_health_api_reports_healthy_process_when_all_mirrors_time_out( + short_health_budget: None, +) -> None: + sessions = {origin: _HealthSession(slow_fetch=True) for origin in KNOWN_SOURCE_URLS} + app = create_app(bearer_token=TEST_TOKEN, service=_service(sessions)) + async with httpx.AsyncClient( + transport=httpx.ASGITransport(app=app), base_url="http://provider.test" + ) as client: + try: + async with asyncio.timeout(1): + response = await client.get( + "/v1/health", headers={"Authorization": f"Bearer {TEST_TOKEN}"} + ) + except TimeoutError: + pytest.fail("Health API did not respond when upstream mirrors were unavailable") + + assert response.status_code == 200 + payload = response.json() + assert payload["process_status"] == "healthy" + assert payload["source_status"] == "unavailable" + assert payload["diagnostics"] == { + "source": "unavailable", + **{f"source.{origin.removeprefix('https://')}": "unavailable" for origin in sessions}, + } + assert all(session.closed for session in sessions.values()) diff --git a/tests/unit/test_provider_release_metadata.py b/tests/unit/test_provider_release_metadata.py index 4597ec4..33aa798 100644 --- a/tests/unit/test_provider_release_metadata.py +++ b/tests/unit/test_provider_release_metadata.py @@ -17,9 +17,9 @@ def test_security_release_versions_and_bundled_libgen_dependency_agree() -> None for name in ("getcomics", "annas_archive", "libgen") } assert projects["getcomics"]["version"] == "1.0.3" - assert projects["annas_archive"]["version"] == "1.0.3" - assert projects["libgen"]["version"] == "1.0.1" - assert "pullbox-provider-libgen==1.0.1" in projects["annas_archive"]["dependencies"] + assert projects["annas_archive"]["version"] == "1.0.4" + assert projects["libgen"]["version"] == "1.0.2" + assert "pullbox-provider-libgen==1.0.2" in projects["annas_archive"]["dependencies"] def _load_module(): @@ -59,12 +59,12 @@ def test_tag_release_maps_getcomics_to_both_registry_names() -> None: def test_tag_release_maps_annas_archive_to_both_registry_names() -> None: release = _load_module().resolve_release( repository_owner="pullboxapp", - tag="annas-archive-v1.0.3", + tag="annas-archive-v1.0.4", ) assert release.provider == "annas-archive" - assert release.version == "1.0.3" - assert release.release_tag == "annas-archive-v1.0.3" + assert release.version == "1.0.4" + assert release.release_tag == "annas-archive-v1.0.4" assert release.is_release is True assert release.is_prerelease is False assert release.dockerfile == "docker/Dockerfile.annas-archive" @@ -75,12 +75,12 @@ def test_tag_release_maps_annas_archive_to_both_registry_names() -> None: def test_tag_release_maps_libgen_to_both_registry_names() -> None: release = _load_module().resolve_release( repository_owner="pullboxapp", - tag="libgen-v1.0.1", + tag="libgen-v1.0.2", ) assert release.provider == "libgen" - assert release.version == "1.0.1" - assert release.release_tag == "libgen-v1.0.1" + assert release.version == "1.0.2" + assert release.release_tag == "libgen-v1.0.2" assert release.is_release is True assert release.is_prerelease is False assert release.dockerfile == "docker/Dockerfile.libgen"