From 5ce4d7862a5bfeaad3b6afc40f1500bc625d6892 Mon Sep 17 00:00:00 2001 From: Gil Desmarais Date: Sun, 30 Aug 2026 10:25:13 +0200 Subject: [PATCH 1/5] feat(engine): WorkLease owns admit, deadline reclaim, and timeout phase Replace Future.cancel theater and ScrapeProgress with a single lease that admits per-host, kills the session driver on deadline, and builds the timeout envelope from one phase snapshot (SCRAPE_MAX_PER_HOST). --- AGENTS.md | 14 +- app/config.py | 3 + app/domain/scrape_service.py | 69 ++----- app/engine/browser_tier.py | 12 +- app/engine/orchestrator.py | 16 +- app/engine/request_tier.py | 8 +- app/engine/session.py | 24 ++- app/engine/work_lease.py | 191 ++++++++++++++++++ app/infra/scrape_progress.py | 53 ----- tests/api/test_http_contract.py | 10 +- tests/api/test_timeout_http.py | 16 +- tests/domain/test_scrape_service_cancel.py | 42 ++-- tests/domain/test_timeout_error.py | 33 ++- tests/domain/test_work_lease_reclaim.py | 82 ++++++++ tests/engine/test_scraper_engine.py | 2 +- tests/engine/test_timeout_progress.py | 52 ++--- .../test_work_lease_snapshot.py} | 18 +- tests/support/http.py | 8 +- 18 files changed, 434 insertions(+), 219 deletions(-) create mode 100644 app/engine/work_lease.py delete mode 100644 app/infra/scrape_progress.py create mode 100644 tests/domain/test_work_lease_reclaim.py rename tests/{infra/test_scrape_progress.py => engine/test_work_lease_snapshot.py} (63%) diff --git a/AGENTS.md b/AGENTS.md index 74c3c73..7e6992d 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -21,10 +21,11 @@ app/ openapi_examples.py # OpenAPI examples built from Pydantic model instances routes/ # thin HTTP handlers (health, scrape) domain/ - scrape_service.py # request-id resolution, URL guardrails, threadpool execution, status mapping + scrape_service.py # request-id resolution, URL guardrails, WorkLease execution, status mapping engine/ orchestrator.py # ScraperEngine.execute - session.py # ScrapeSession lifecycle + session.py # ScrapeSession lifecycle + lease reclaim hook + work_lease.py # WorkLease admit/deadline/reclaim + phase snapshot; HostConcurrencyGate budget.py # wall-clock budget math shared across tiers (elapsed_ms, step budgets) request_tier.py # HTTP/curl_cffi path browser_tier.py # Chromium path @@ -36,7 +37,8 @@ app/ enums.py # ExecutionMode, NavigationMode, ErrorCategory, ... request.py # ScrapeRequest and validators response.py # ScrapeSuccess, ScrapeError, HealthResponse, ... - infra/ # telemetry, progress, metadata, xhr, runtime cleanup, sentry + infra/ # telemetry, metadata, xhr, runtime cleanup, sentry, challenge detector + # (phase snapshot lives on WorkLease — not a separate progress module) security/ # UrlGuard SSRF guardrails scripts/ bench_scrape.py # TestClient wall-time bench for POST /scrape (request tier) @@ -91,11 +93,15 @@ Conventions: | --- | --- | --- | | `ScraperEngine` | process (app.state) | shared | | `ThreadPoolExecutor` | process (app.state) | shared, sized by `SCRAPE_MAX_WORKERS` | +| `HostConcurrencyGate` | process (ScrapeService) | per-host admit; sized by `SCRAPE_MAX_PER_HOST` | +| `WorkLease` | per request | owns admit → executor run → deadline reclaim (session `force_close`); phase snapshot for timeout envelope; `Future.cancel` is not Chromium reclaim | | `WarmDriverPool` (opt-in) | process (engine.warm_pool) | single spare slot; refill on dedicated daemon thread — never the scrape executor | | `_active_request_ids` | in-process memory | shared; collision guard | | runtime dir `/tmp/scrape/` | per request | isolated; deleted in `finally` | | browser profile | per request (or adopted spare-*) | isolated; no reuse across requests; warm spare dies with the adopting request | -| Botasaurus Driver | may start before assignment | usage stays ≤1 request; closed in session `__exit__`; never returned to the pool | +| Botasaurus Driver | may start before assignment | usage stays ≤1 request; closed in session `__exit__` / lease reclaim; never returned to the pool | + +**Timeout ownership:** `WorkLease` is the single owner of outer deadline, host admit, phase snapshot, and Chromium reclaim. Tiers mark phase on the lease; `WorkLease.timeout_error` is the only outer-timeout envelope builder. Session exit and orphan prune remain safety nets under lease reclaim — not a second reclaim story. Multi-worker uvicorn breaks in-process collision detection unless request ids are sticky to a worker. Default to single-worker for isolation semantics. diff --git a/app/config.py b/app/config.py index e6519e6..05c8988 100644 --- a/app/config.py +++ b/app/config.py @@ -61,6 +61,9 @@ class Settings(BaseSettings): default=30, validation_alias="SCRAPE_WORK_TIMEOUT_SECONDS" ) scrape_max_workers: int = Field(default=4, validation_alias="SCRAPE_MAX_WORKERS") + scrape_max_per_host: int = Field( + default=2, ge=1, validation_alias="SCRAPE_MAX_PER_HOST" + ) scrape_runtime_min_free_bytes: int = Field( default=256 * 1024 * 1024, validation_alias="SCRAPE_RUNTIME_MIN_FREE_BYTES", diff --git a/app/domain/scrape_service.py b/app/domain/scrape_service.py index 2c24581..5865148 100644 --- a/app/domain/scrape_service.py +++ b/app/domain/scrape_service.py @@ -2,7 +2,6 @@ from __future__ import annotations -import asyncio import time from concurrent.futures import ThreadPoolExecutor from dataclasses import dataclass @@ -10,12 +9,10 @@ from app.config import Settings from app.engine import ScraperEngine -from app.engine.budget import elapsed_ms -from app.engine.envelope import TIMEOUT_ERROR_BY_PHASE +from app.engine.work_lease import HostConcurrencyGate, WorkLease from app.exceptions import RequestIdCollisionError from app.infra.ops_telemetry import emit_terminal_telemetry from app.infra.request_id import resolve_request_id -from app.infra.scrape_progress import ScrapeProgress from app.logging_config import get_logger from app.schemas.enums import ErrorCategory, TimeoutPhase from app.schemas.request import ScrapeRequest @@ -37,7 +34,7 @@ class ScrapeOutcome: class ScrapeService: - """Owns request-id resolution, URL guardrails, threadpool execution, status mapping, and telemetry.""" + """Owns request-id resolution, URL guardrails, WorkLease execution, status mapping.""" def __init__( self, @@ -45,10 +42,12 @@ def __init__( settings: Settings, engine: ScraperEngine, executor: ThreadPoolExecutor, + host_gate: HostConcurrencyGate | None = None, ) -> None: self.settings = settings self.engine = engine self.executor = executor + self.host_gate = host_gate or HostConcurrencyGate(settings.scrape_max_per_host) async def process( self, @@ -108,33 +107,6 @@ def _validation_outcome( status_code=validation.status_code, ) - @staticmethod - def build_timeout_error( - url: str, - *, - request_id: str, - started_monotonic: float, - progress: ScrapeProgress, - timeout_seconds: int, - ) -> ScrapeError: - del timeout_seconds # budget length is operational; phase message is the UX - snap = progress.snapshot() - phase = snap.phase - render_ms = elapsed_ms(started_monotonic) - return ScrapeError( - url=url, - error=TIMEOUT_ERROR_BY_PHASE[phase], - error_category=ErrorCategory.TIMEOUT, - diagnostics=ScrapeDiagnostics( - request_id=request_id, - attempts=snap.attempts, - strategy_used=snap.strategy_used, - render_ms=render_ms, - execution_tier=snap.execution_tier, - timeout_phase=phase, - ), - ) - async def _run( self, payload: ScrapeRequest, @@ -142,31 +114,26 @@ async def _run( request_id: str, ) -> ScrapeOutcome: target_url = str(payload.url) - host = payload.url.host + host = payload.url.host or "" + lease = WorkLease( + settings=self.settings, + executor=self.executor, + host_gate=self.host_gate, + ) started_monotonic = time.monotonic() deadline_monotonic = started_monotonic + self.settings.scrape_timeout_seconds - progress = ScrapeProgress() try: - loop = asyncio.get_running_loop() - future = loop.run_in_executor( - self.executor, - partial( + result = await lease.run( + host=host, + work=partial( self.engine.execute, payload, deadline_monotonic, request_id=request_id, - progress=progress, + lease=lease, ), ) - try: - result = await asyncio.wait_for( - future, - timeout=self.settings.scrape_timeout_seconds, - ) - except TimeoutError: - future.cancel() - raise except RequestIdCollisionError: collision_result = ScrapeError( url=target_url, @@ -181,12 +148,10 @@ async def _run( emit_terminal_telemetry(collision_result, http_status=502) return ScrapeOutcome(body=collision_result, status_code=502) except TimeoutError: - timeout_result = self.build_timeout_error( + timeout_result = lease.timeout_error( target_url, request_id=request_id, started_monotonic=started_monotonic, - progress=progress, - timeout_seconds=self.settings.scrape_timeout_seconds, ) phase = timeout_result.diagnostics.timeout_phase or TimeoutPhase.QUEUE logger.warning( @@ -200,7 +165,7 @@ async def _run( emit_terminal_telemetry( timeout_result, http_status=504, - warm_hit=progress.snapshot().warm_hit, + warm_hit=lease.snapshot().warm_hit, ) return ScrapeOutcome(body=timeout_result, status_code=504) @@ -209,7 +174,7 @@ async def _run( emit_terminal_telemetry( result, http_status=status_code, - warm_hit=progress.snapshot().warm_hit, + warm_hit=lease.snapshot().warm_hit, ) logger.info( "scrape_complete request_id=%s host=%s mode=%s tier=%s attempts=%s status=%d error_category=%s", diff --git a/app/engine/browser_tier.py b/app/engine/browser_tier.py index c325d56..ee65d6d 100644 --- a/app/engine/browser_tier.py +++ b/app/engine/browser_tier.py @@ -26,9 +26,9 @@ wait_for_readiness, ) from app.engine.warm_pool import DriverFingerprint +from app.engine.work_lease import WorkLease from app.infra.detector import ChallengeAssessment, ChallengeDetector from app.infra.metadata import MetadataExtractor, MetadataResult -from app.infra.scrape_progress import ScrapeProgress from app.infra.xhr_collector import XhrCollector from app.logging_config import get_logger from app.schemas.enums import ErrorCategory, ExecutionTier, NavigationMode, TimeoutPhase @@ -229,7 +229,7 @@ def run_browser_tier( payload: ScrapeRequest, session: ScrapeSession, started_monotonic: float, - progress: ScrapeProgress, + lease: WorkLease, *, settings: Settings, ) -> ScrapeSuccess | ScrapeError: @@ -237,7 +237,7 @@ def run_browser_tier( target_url = str(payload.url) request_id = session.request_id - progress.mark( + lease.mark( TimeoutPhase.BOOT, execution_tier=ExecutionTier.BROWSER_DRIVER, ) @@ -303,7 +303,7 @@ def run_browser_tier( session.driver = driver session.warm_hit = warm_hit - progress.set_warm_hit(warm_hit) + lease.set_warm_hit(warm_hit) boot_ms = int((time.monotonic() - boot_started) * 1000) logger.info( "scrape_boot request_id=%s warm_hit=%s boot_ms=%d", @@ -320,7 +320,7 @@ def run_browser_tier( configure_driver(driver, payload, target_url, collector=collector) browser_ready_monotonic = time.monotonic() - progress.mark( + lease.mark( TimeoutPhase.WORK, execution_tier=ExecutionTier.BROWSER_DRIVER, ) @@ -328,7 +328,7 @@ def run_browser_tier( for attempt_index, strategy in enumerate(strategies, start=1): attempts = attempt_index has_more = attempt_index < len(strategies) - progress.mark( + lease.mark( TimeoutPhase.WORK, attempts=attempts, strategy_used=strategy, diff --git a/app/engine/orchestrator.py b/app/engine/orchestrator.py index c86c5b3..bec0b1d 100644 --- a/app/engine/orchestrator.py +++ b/app/engine/orchestrator.py @@ -15,12 +15,12 @@ from app.engine.request_tier import run_request_tier from app.engine.session import ScrapeSession from app.engine.warm_pool import DriverFingerprint, WarmDriverPool +from app.engine.work_lease import WorkLease from app.exceptions import RequestIdCollisionError from app.infra.runtime_cleanup import ( prune_orphan_runtime_dirs, runtime_root_low_on_space, ) -from app.infra.scrape_progress import ScrapeProgress from app.logging_config import get_logger from app.schemas.enums import ( ErrorCategory, @@ -89,7 +89,7 @@ def execute( deadline_monotonic: float | None = None, *, request_id: str | None = None, - progress: ScrapeProgress | None = None, + lease: WorkLease | None = None, ) -> ScrapeSuccess | ScrapeError: target_url = str(payload.url) resolved_request_id = request_id or str(uuid.uuid4()) @@ -103,10 +103,10 @@ def execute( ) else: started_monotonic = now - progress = progress or ScrapeProgress() + lease = lease or WorkLease.tracking_only(self.settings) if deadline_monotonic is not None and now >= deadline_monotonic: - progress.mark(TimeoutPhase.QUEUE) + lease.mark(TimeoutPhase.QUEUE) return build_error( target_url, TIMEOUT_ERROR_BY_PHASE[TimeoutPhase.QUEUE], @@ -117,7 +117,7 @@ def execute( warm_fp: DriverFingerprint | None = None try: - with ScrapeSession(self, resolved_request_id) as session: + with ScrapeSession(self, resolved_request_id, lease=lease) as session: should_try_request_tier = ( payload.execution_mode == ExecutionMode.REQUEST or ( @@ -134,7 +134,7 @@ def execute( payload, resolved_request_id, started_monotonic, - progress, + lease, settings=self.settings, ) if request_result is not None: @@ -171,7 +171,7 @@ def execute( ) if remaining_total_seconds(self.settings, started_monotonic) <= 0: - progress.mark( + lease.mark( TimeoutPhase.BOOT, execution_tier=ExecutionTier.BROWSER_DRIVER, ) @@ -191,7 +191,7 @@ def execute( payload, session, started_monotonic, - progress, + lease, settings=self.settings, ) warm_fp = session.warm_fingerprint diff --git a/app/engine/request_tier.py b/app/engine/request_tier.py index 4ddc2f9..84d6ef8 100644 --- a/app/engine/request_tier.py +++ b/app/engine/request_tier.py @@ -12,8 +12,8 @@ remaining_total_seconds, ) from app.engine.envelope import TIMEOUT_ERROR_BY_PHASE, build_error, build_success +from app.engine.work_lease import WorkLease from app.infra.detector import ChallengeDetector -from app.infra.scrape_progress import ScrapeProgress from app.logging_config import get_logger from app.schemas.enums import ErrorCategory, ExecutionMode, ExecutionTier, TimeoutPhase from app.schemas.request import ScrapeRequest @@ -29,7 +29,7 @@ def run_request_tier( payload: ScrapeRequest, request_id: str, started_monotonic: float, - progress: ScrapeProgress, + lease: WorkLease, *, settings: Settings, ) -> ScrapeSuccess | ScrapeError | None: @@ -38,7 +38,7 @@ def run_request_tier( target_url = str(payload.url) remaining_budget = remaining_total_seconds(settings, started_monotonic) if remaining_budget <= 0: - progress.mark( + lease.mark( TimeoutPhase.WORK, attempts=1, execution_tier=ExecutionTier.HTTP_REQUEST, @@ -54,7 +54,7 @@ def run_request_tier( timeout_phase=TimeoutPhase.WORK, ) - progress.mark( + lease.mark( TimeoutPhase.WORK, attempts=1, execution_tier=ExecutionTier.HTTP_REQUEST, diff --git a/app/engine/session.py b/app/engine/session.py index 5c2595e..b02e954 100644 --- a/app/engine/session.py +++ b/app/engine/session.py @@ -12,20 +12,29 @@ if TYPE_CHECKING: from app.engine.orchestrator import ScraperEngine from app.engine.warm_pool import DriverFingerprint + from app.engine.work_lease import WorkLease class ScrapeSession: """Encapsulates per-request concurrency registration and filesystem isolation.""" - def __init__(self, engine: ScraperEngine, request_id: str) -> None: + def __init__( + self, + engine: ScraperEngine, + request_id: str, + *, + lease: WorkLease | None = None, + ) -> None: self.engine = engine self.request_id = request_id + self.lease = lease self.runtime_dir = engine.runtime_root / request_id self.profile_dir = self.runtime_dir / "profile" self.driver: DriverProtocol | None = None self.adopted_profile_dir: Path | None = None self.warm_fingerprint: DriverFingerprint | None = None self.warm_hit: bool | None = None + self._closed = False def __enter__(self) -> ScrapeSession: self.engine.register_request_id(self.request_id) @@ -34,8 +43,18 @@ def __enter__(self) -> ScrapeSession: except Exception: self.engine.unregister_request_id(self.request_id) raise + if self.lease is not None: + self.lease.register_reclaim(self.force_close) return self + def force_close(self) -> None: + """Lease-deadline reclaim: close Chromium once so the worker slot can free.""" + if self._closed: + return + self._closed = True + if self.driver is not None: + call_quietly(self.driver, "close") + def prepare_runtime_dir(self) -> None: """Create the request runtime dir only (warm-path adoption).""" try: @@ -65,8 +84,7 @@ def _make_dirs(self) -> None: def __exit__(self, exc_type: Any, exc_val: Any, exc_tb: Any) -> None: try: - if self.driver is not None: - call_quietly(self.driver, "close") + self.force_close() finally: shutil.rmtree(self.runtime_dir, ignore_errors=True) if self.adopted_profile_dir is not None: diff --git a/app/engine/work_lease.py b/app/engine/work_lease.py new file mode 100644 index 0000000..84d52bb --- /dev/null +++ b/app/engine/work_lease.py @@ -0,0 +1,191 @@ +"""WorkLease: admit, deadline, phase snapshot, and Chromium reclaim ownership.""" + +from __future__ import annotations + +import asyncio +import threading +import time +from collections.abc import Callable +from concurrent.futures import ThreadPoolExecutor +from dataclasses import dataclass, replace +from typing import TypeVar + +from app.config import Settings +from app.engine.budget import elapsed_ms +from app.engine.envelope import TIMEOUT_ERROR_BY_PHASE +from app.logging_config import get_logger +from app.schemas.enums import ( + ErrorCategory, + ExecutionTier, + NavigationMode, + TimeoutPhase, +) +from app.schemas.response import ScrapeDiagnostics, ScrapeError + +logger = get_logger() + +T = TypeVar("T") + + +@dataclass(frozen=True, slots=True) +class LeaseSnapshot: + phase: TimeoutPhase = TimeoutPhase.QUEUE + attempts: int = 0 + strategy_used: NavigationMode | None = None + execution_tier: ExecutionTier | None = None + warm_hit: bool | None = None + + +class HostConcurrencyGate: + """Process-wide per-host admit gate (one fact: max concurrent scrapes per host).""" + + def __init__(self, max_per_host: int) -> None: + if max_per_host < 1: + raise ValueError("max_per_host must be >= 1") + self._max_per_host = max_per_host + self._lock = threading.Lock() + self._semaphores: dict[str, threading.BoundedSemaphore] = {} + + def _semaphore(self, host: str) -> threading.BoundedSemaphore: + with self._lock: + sem = self._semaphores.get(host) + if sem is None: + sem = threading.BoundedSemaphore(self._max_per_host) + self._semaphores[host] = sem + return sem + + def acquire(self, host: str, timeout: float) -> bool: + return self._semaphore(host).acquire(timeout=max(0.0, timeout)) + + def release(self, host: str) -> None: + self._semaphore(host).release() + + +class WorkLease: + """Owns admit → run → terminal timeout; reclaim kills the session driver once. + + Progress phase lives on the lease (replaces ScrapeProgress). Outer deadline + calls register_reclaim hooks (session force-close); Future.cancel is not the + Chromium reclaim story. + """ + + def __init__( + self, + *, + settings: Settings, + executor: ThreadPoolExecutor | None = None, + host_gate: HostConcurrencyGate | None = None, + ) -> None: + self.settings = settings + self.executor = executor + self.host_gate = host_gate + self._lock = threading.Lock() + self._snap = LeaseSnapshot() + self._reclaim_hooks: list[Callable[[], None]] = [] + + @classmethod + def tracking_only(cls, settings: Settings | None = None) -> WorkLease: + """Mark/snapshot-only lease for direct engine.execute tests.""" + from app.config import get_settings + + return cls(settings=settings or get_settings()) + + def mark( + self, + phase: TimeoutPhase, + *, + attempts: int | None = None, + strategy_used: NavigationMode | None = None, + execution_tier: ExecutionTier | None = None, + ) -> None: + fields = { + k: v + for k, v in ( + ("attempts", attempts), + ("strategy_used", strategy_used), + ("execution_tier", execution_tier), + ) + if v is not None + } + with self._lock: + self._snap = replace(self._snap, phase=phase, **fields) + + def set_warm_hit(self, warm_hit: bool) -> None: + with self._lock: + self._snap = replace(self._snap, warm_hit=warm_hit) + + def snapshot(self) -> LeaseSnapshot: + with self._lock: + return self._snap + + def register_reclaim(self, hook: Callable[[], None]) -> None: + with self._lock: + self._reclaim_hooks.append(hook) + + def reclaim(self) -> None: + with self._lock: + hooks = list(self._reclaim_hooks) + for hook in hooks: + try: + hook() + except Exception as exc: + logger.debug("lease_reclaim_hook_failed error=%s", str(exc)) + + def timeout_error( + self, + url: str, + *, + request_id: str, + started_monotonic: float, + ) -> ScrapeError: + """Single timeout envelope builder fed by lease phase.""" + snap = self.snapshot() + phase = snap.phase + return ScrapeError( + url=url, + error=TIMEOUT_ERROR_BY_PHASE[phase], + error_category=ErrorCategory.TIMEOUT, + diagnostics=ScrapeDiagnostics( + request_id=request_id, + attempts=snap.attempts, + strategy_used=snap.strategy_used, + render_ms=elapsed_ms(started_monotonic), + execution_tier=snap.execution_tier, + timeout_phase=phase, + ), + ) + + async def run(self, *, host: str, work: Callable[[], T]) -> T: + """Admit on host gate, run work on executor, reclaim on deadline.""" + if self.executor is None or self.host_gate is None: + raise RuntimeError("WorkLease.run requires executor and host_gate") + + timeout_seconds = self.settings.scrape_timeout_seconds + started_monotonic = time.monotonic() + deadline = started_monotonic + timeout_seconds + + remaining_admit = deadline - time.monotonic() + admitted = await asyncio.to_thread( + self.host_gate.acquire, host, remaining_admit + ) + if not admitted: + self.mark(TimeoutPhase.QUEUE) + raise TimeoutError("host concurrency admit deadline") + + try: + remaining = deadline - time.monotonic() + if remaining <= 0: + self.mark(TimeoutPhase.QUEUE) + raise TimeoutError("scrape deadline before executor submit") + + loop = asyncio.get_running_loop() + awaitable = asyncio.ensure_future(loop.run_in_executor(self.executor, work)) + try: + return await asyncio.wait_for( + asyncio.shield(awaitable), timeout=remaining + ) + except TimeoutError: + self.reclaim() + raise + finally: + self.host_gate.release(host) diff --git a/app/infra/scrape_progress.py b/app/infra/scrape_progress.py deleted file mode 100644 index e47bf8d..0000000 --- a/app/infra/scrape_progress.py +++ /dev/null @@ -1,53 +0,0 @@ -"""In-memory scrape progress tracking for timeout phase diagnostics.""" - -from __future__ import annotations - -import threading -from dataclasses import dataclass, replace - -from app.schemas.enums import ExecutionTier, NavigationMode, TimeoutPhase - - -@dataclass(frozen=True, slots=True) -class ScrapeProgressSnapshot: - phase: TimeoutPhase = TimeoutPhase.QUEUE - attempts: int = 0 - strategy_used: NavigationMode | None = None - execution_tier: ExecutionTier | None = None - warm_hit: bool | None = None - - -class ScrapeProgress: - """Thread-safe scrape stage tracker for handler-timeout diagnostics.""" - - def __init__(self) -> None: - self._lock = threading.Lock() - self._snap = ScrapeProgressSnapshot() - - def mark( - self, - phase: TimeoutPhase, - *, - attempts: int | None = None, - strategy_used: NavigationMode | None = None, - execution_tier: ExecutionTier | None = None, - ) -> None: - fields = { - k: v - for k, v in ( - ("attempts", attempts), - ("strategy_used", strategy_used), - ("execution_tier", execution_tier), - ) - if v is not None - } - with self._lock: - self._snap = replace(self._snap, phase=phase, **fields) - - def set_warm_hit(self, warm_hit: bool) -> None: - with self._lock: - self._snap = replace(self._snap, warm_hit=warm_hit) - - def snapshot(self) -> ScrapeProgressSnapshot: - with self._lock: - return self._snap diff --git a/tests/api/test_http_contract.py b/tests/api/test_http_contract.py index f2a6de4..bf1d766 100644 --- a/tests/api/test_http_contract.py +++ b/tests/api/test_http_contract.py @@ -5,7 +5,7 @@ from app.engine import ( ScraperEngine, ) -from app.infra.scrape_progress import ScrapeProgress +from app.engine.work_lease import WorkLease from app.schemas.enums import ( ExecutionTier, ) @@ -23,9 +23,9 @@ def fake_execute( deadline_monotonic: float | None = None, *, request_id: str | None = None, - progress: ScrapeProgress | None = None, + lease: WorkLease | None = None, ) -> ScrapeSuccess: - del deadline_monotonic, progress + del deadline_monotonic, lease resolved_request_id = request_id or "req-unknown" return ScrapeSuccess( url=str(payload.url), @@ -212,9 +212,9 @@ def fake_execute( deadline_monotonic: float | None = None, *, request_id: str | None = None, - progress: ScrapeProgress | None = None, + lease: WorkLease | None = None, ) -> ScrapeSuccess: - del deadline_monotonic, request_id, progress + del deadline_monotonic, request_id, lease captured["wait"] = payload.wait_timeout_seconds return ScrapeSuccess( url=str(payload.url), diff --git a/tests/api/test_timeout_http.py b/tests/api/test_timeout_http.py index 7ce494f..67ddde0 100644 --- a/tests/api/test_timeout_http.py +++ b/tests/api/test_timeout_http.py @@ -1,11 +1,11 @@ -"""HTTP 504 envelope carries progress-derived timeout diagnostics.""" +"""HTTP 504 envelope carries lease-derived timeout diagnostics.""" from __future__ import annotations import unittest from unittest.mock import patch -from app.infra.scrape_progress import ScrapeProgress +from app.engine.work_lease import WorkLease from app.schemas.enums import ErrorCategory, ExecutionTier, TimeoutPhase from app.schemas.request import ScrapeRequest from app.schemas.response import ScrapeDiagnostics, ScrapeError @@ -15,21 +15,19 @@ class HandlerTimeoutHttpTests(unittest.TestCase): - def test_scrape_handler_timeout_uses_progress(self): + def test_scrape_handler_timeout_uses_lease_phase(self) -> None: def fake_execute( payload: ScrapeRequest, deadline_monotonic: float | None = None, *, request_id: str | None = None, - progress: ScrapeProgress | None = None, + lease: WorkLease | None = None, ) -> ScrapeError: del payload, deadline_monotonic, request_id - assert progress is not None - progress.mark( - TimeoutPhase.BOOT, execution_tier=ExecutionTier.BROWSER_DRIVER - ) + assert lease is not None + lease.mark(TimeoutPhase.BOOT, execution_tier=ExecutionTier.BROWSER_DRIVER) # Sentinel only: the patched wait_for raises TimeoutError, so the - # 504 envelope must come from the handler's own timeout path. + # 504 envelope must come from the lease timeout path. return ScrapeError( url=_URL, error="sentinel-discarded", diff --git a/tests/domain/test_scrape_service_cancel.py b/tests/domain/test_scrape_service_cancel.py index a47ebea..f254869 100644 --- a/tests/domain/test_scrape_service_cancel.py +++ b/tests/domain/test_scrape_service_cancel.py @@ -1,4 +1,4 @@ -"""ScrapeService cancels queued executor work on outer timeout.""" +"""ScrapeService outer timeout maps lease TimeoutError to 504 envelope.""" from __future__ import annotations @@ -9,13 +9,14 @@ from app.config import get_settings from app.domain.scrape_service import ScrapeOutcome, ScrapeService from app.engine import ScraperEngine +from app.engine.envelope import TIMEOUT_ERROR_BY_PHASE from app.schemas.enums import ErrorCategory, TimeoutPhase -from app.schemas.response import ScrapeError +from app.schemas.response import ScrapeDiagnostics, ScrapeError from tests.support.factories import scrape_request -class ScrapeServiceCancelTests(unittest.IsolatedAsyncioTestCase): - async def test_timeout_cancels_executor_future(self) -> None: +class ScrapeServiceLeaseTimeoutTests(unittest.IsolatedAsyncioTestCase): + async def test_timeout_returns_lease_timeout_envelope(self) -> None: settings = get_settings() engine = ScraperEngine(settings=settings) service = ScrapeService( @@ -23,28 +24,39 @@ async def test_timeout_cancels_executor_future(self) -> None: engine=engine, executor=ThreadPoolExecutor(max_workers=1), ) - future = MagicMock() - future.cancel = MagicMock(return_value=True) - async def boom(awaitable: object, timeout: float | None = None) -> None: - del awaitable, timeout + async def boom_run(*, host: str, work: object) -> None: + del host, work raise TimeoutError - with ( - patch("asyncio.get_running_loop") as mock_loop, - patch("asyncio.wait_for", side_effect=boom), - ): - mock_loop.return_value.run_in_executor = MagicMock(return_value=future) + with patch("app.domain.scrape_service.WorkLease") as lease_cls: + lease = MagicMock() + lease.run = boom_run + lease.timeout_error = MagicMock( + return_value=ScrapeError( + url="https://example.com", + error=TIMEOUT_ERROR_BY_PHASE[TimeoutPhase.QUEUE], + error_category=ErrorCategory.TIMEOUT, + diagnostics=ScrapeDiagnostics( + request_id="req-timeout", + attempts=0, + timeout_phase=TimeoutPhase.QUEUE, + ), + ) + ) + lease.snapshot = MagicMock(return_value=MagicMock(warm_hit=None)) + lease_cls.return_value = lease + outcome = await service.process(scrape_request()) - future.cancel.assert_called_once() self.assertIsInstance(outcome, ScrapeOutcome) self.assertEqual(outcome.status_code, 504) self.assertIsInstance(outcome.body, ScrapeError) assert isinstance(outcome.body, ScrapeError) self.assertEqual(outcome.body.error_category, ErrorCategory.TIMEOUT) self.assertEqual(outcome.body.diagnostics.timeout_phase, TimeoutPhase.QUEUE) - self.assertEqual(outcome.body.error, "Scraper at capacity; retry shortly") + self.assertEqual(outcome.body.error, TIMEOUT_ERROR_BY_PHASE[TimeoutPhase.QUEUE]) + lease.timeout_error.assert_called_once() service.executor.shutdown(wait=False, cancel_futures=True) diff --git a/tests/domain/test_timeout_error.py b/tests/domain/test_timeout_error.py index 961d1fe..6075f67 100644 --- a/tests/domain/test_timeout_error.py +++ b/tests/domain/test_timeout_error.py @@ -1,26 +1,23 @@ -"""ScrapeService.build_timeout_error phase and diagnostics mapping.""" +"""WorkLease timeout envelope builder (single owner).""" from __future__ import annotations import time import unittest -from app.domain.scrape_service import ScrapeService from app.engine.envelope import TIMEOUT_ERROR_BY_PHASE -from app.infra.scrape_progress import ScrapeProgress +from app.engine.work_lease import WorkLease from app.schemas.enums import ExecutionTier, NavigationMode, TimeoutPhase _URL = "https://example.com" -class BuildTimeoutErrorTests(unittest.TestCase): - def test_queue_phase_keeps_zero_attempts(self): - result = ScrapeService.build_timeout_error( +class WorkLeaseTimeoutErrorTests(unittest.TestCase): + def test_queue_phase_keeps_zero_attempts(self) -> None: + result = WorkLease.tracking_only().timeout_error( _URL, request_id="req-queue", started_monotonic=time.monotonic(), - progress=ScrapeProgress(), - timeout_seconds=45, ) self.assertEqual(result.error_category.value, "timeout") self.assertEqual(result.error, TIMEOUT_ERROR_BY_PHASE[TimeoutPhase.QUEUE]) @@ -28,33 +25,29 @@ def test_queue_phase_keeps_zero_attempts(self): self.assertEqual(result.diagnostics.attempts, 0) self.assertIsNone(result.diagnostics.strategy_used) - def test_boot_phase_message(self): - progress = ScrapeProgress() - progress.mark(TimeoutPhase.BOOT, execution_tier=ExecutionTier.BROWSER_DRIVER) - result = ScrapeService.build_timeout_error( + def test_boot_phase_message(self) -> None: + lease = WorkLease.tracking_only() + lease.mark(TimeoutPhase.BOOT, execution_tier=ExecutionTier.BROWSER_DRIVER) + result = lease.timeout_error( _URL, request_id="req-boot", started_monotonic=time.monotonic(), - progress=progress, - timeout_seconds=45, ) self.assertEqual(result.error, TIMEOUT_ERROR_BY_PHASE[TimeoutPhase.BOOT]) self.assertEqual(result.diagnostics.timeout_phase, TimeoutPhase.BOOT) - def test_work_phase_preserves_attempts_and_strategy(self): - progress = ScrapeProgress() - progress.mark( + def test_work_phase_preserves_attempts_and_strategy(self) -> None: + lease = WorkLease.tracking_only() + lease.mark( TimeoutPhase.WORK, attempts=2, strategy_used=NavigationMode.GOOGLE_GET, execution_tier=ExecutionTier.BROWSER_DRIVER, ) - result = ScrapeService.build_timeout_error( + result = lease.timeout_error( _URL, request_id="req-work", started_monotonic=time.monotonic() - 1, - progress=progress, - timeout_seconds=45, ) diagnostics = result.diagnostics self.assertEqual(diagnostics.timeout_phase, TimeoutPhase.WORK) diff --git a/tests/domain/test_work_lease_reclaim.py b/tests/domain/test_work_lease_reclaim.py new file mode 100644 index 0000000..1f0dd8d --- /dev/null +++ b/tests/domain/test_work_lease_reclaim.py @@ -0,0 +1,82 @@ +"""WorkLease reclaim on deadline frees host slot; Future.cancel is not reclaim.""" + +from __future__ import annotations + +import threading +import unittest +from concurrent.futures import ThreadPoolExecutor + +from app.config import get_settings +from app.engine.work_lease import HostConcurrencyGate, WorkLease +from app.schemas.enums import TimeoutPhase + + +def _fast_settings(*, max_per_host: int = 1, timeout_seconds: int = 1): + return get_settings().model_copy( + update={ + "scrape_timeout_seconds": timeout_seconds, + "scrape_work_timeout_seconds": timeout_seconds, + "scrape_max_per_host": max_per_host, + } + ) + + +class WorkLeaseReclaimTests(unittest.IsolatedAsyncioTestCase): + async def test_deadline_invokes_reclaim_and_releases_host_slot(self) -> None: + settings = _fast_settings(max_per_host=1, timeout_seconds=1) + gate = HostConcurrencyGate(1) + executor = ThreadPoolExecutor(max_workers=1) + lease = WorkLease(settings=settings, executor=executor, host_gate=gate) + reclaimed = threading.Event() + release_work = threading.Event() + + def hung_work() -> str: + release_work.wait(timeout=5) + return "done" + + lease.register_reclaim(reclaimed.set) + lease.register_reclaim(release_work.set) + + with self.assertRaises(TimeoutError): + await lease.run(host="example.com", work=hung_work) + + self.assertTrue(reclaimed.wait(timeout=1)) + # Host slot must be free so a follow-up admit succeeds immediately. + self.assertTrue(gate.acquire("example.com", timeout=0.1)) + gate.release("example.com") + executor.shutdown(wait=False, cancel_futures=True) + + async def test_host_admit_timeout_marks_queue(self) -> None: + settings = _fast_settings(max_per_host=1, timeout_seconds=1) + gate = HostConcurrencyGate(1) + self.assertTrue(gate.acquire("busy.example", timeout=0.1)) + executor = ThreadPoolExecutor(max_workers=1) + lease = WorkLease(settings=settings, executor=executor, host_gate=gate) + + with self.assertRaises(TimeoutError): + await lease.run(host="busy.example", work=lambda: "never") + + self.assertEqual(lease.snapshot().phase, TimeoutPhase.QUEUE) + gate.release("busy.example") + executor.shutdown(wait=False, cancel_futures=True) + + async def test_successful_run_does_not_reclaim(self) -> None: + settings = _fast_settings(max_per_host=2, timeout_seconds=5) + gate = HostConcurrencyGate(2) + executor = ThreadPoolExecutor(max_workers=1) + lease = WorkLease(settings=settings, executor=executor, host_gate=gate) + reclaim_calls = 0 + + def mark_reclaim() -> None: + nonlocal reclaim_calls + reclaim_calls += 1 + + lease.register_reclaim(mark_reclaim) + result = await lease.run(host="ok.example", work=lambda: "ok") + self.assertEqual(result, "ok") + self.assertEqual(reclaim_calls, 0) + executor.shutdown(wait=False, cancel_futures=True) + + +if __name__ == "__main__": + unittest.main() diff --git a/tests/engine/test_scraper_engine.py b/tests/engine/test_scraper_engine.py index 32ffc39..e3bcc3e 100644 --- a/tests/engine/test_scraper_engine.py +++ b/tests/engine/test_scraper_engine.py @@ -222,7 +222,7 @@ def fake_browser_tier( _payload: object, _session: object, started_monotonic: float, - _progress: object, + _lease: object, *, settings: object, ) -> ScrapeSuccess: diff --git a/tests/engine/test_timeout_progress.py b/tests/engine/test_timeout_progress.py index de043a8..cc49006 100644 --- a/tests/engine/test_timeout_progress.py +++ b/tests/engine/test_timeout_progress.py @@ -1,4 +1,4 @@ -"""Engine progress marking across queue, boot, and work phases.""" +"""Engine lease phase marking across queue, boot, and work phases.""" from __future__ import annotations @@ -12,7 +12,7 @@ from app.config import get_settings from app.engine import ScraperEngine -from app.infra.scrape_progress import ScrapeProgress, ScrapeProgressSnapshot +from app.engine.work_lease import LeaseSnapshot, WorkLease from app.schemas.enums import ( ExecutionMode, ExecutionTier, @@ -29,16 +29,16 @@ class _PhaseProbeDriver(FakeDriver): - """Records progress phase at Driver construction time.""" + """Records lease phase at Driver construction time.""" construction_phase: TimeoutPhase | None = None - progress: ScrapeProgress | None = None + lease: WorkLease | None = None def __init__(self, *args: object, **kwargs: Any) -> None: super().__init__(*args, **kwargs) - progress = type(self).progress + lease = type(self).lease type(self).construction_phase = ( - progress.snapshot().phase if progress is not None else None + lease.snapshot().phase if lease is not None else None ) self.page_html = _HTML self.current_url = f"{_URL}/" @@ -47,7 +47,7 @@ def __init__(self, *args: object, **kwargs: Any) -> None: def _execute( payload: ScrapeRequest, *, - progress: ScrapeProgress, + lease: WorkLease, request_id: str, **patches: Any, ) -> ScrapeSuccess | ScrapeError: @@ -62,12 +62,12 @@ def _execute( stack.enter_context( patch("botasaurus.request.Request", patches["Request"]) ) - return engine.execute(payload, request_id=request_id, progress=progress) + return engine.execute(payload, request_id=request_id, lease=lease) def _snap_eq( test: unittest.TestCase, - snap: ScrapeProgressSnapshot, + snap: LeaseSnapshot, *, phase: TimeoutPhase, attempts: int = 0, @@ -80,9 +80,9 @@ def _snap_eq( test.assertEqual(snap.execution_tier, tier) -class EngineProgressMarkTests(unittest.TestCase): +class EngineLeaseMarkTests(unittest.TestCase): def test_execute_queue_timeout_sets_phase(self) -> None: - progress = ScrapeProgress() + lease = WorkLease.tracking_only() result = ScraperEngine(settings=get_settings()).execute( scrape_request( execution_mode=ExecutionMode.BROWSER, @@ -90,17 +90,17 @@ def test_execute_queue_timeout_sets_phase(self) -> None: ), time.monotonic() - 1, request_id="req-engine-queue", - progress=progress, + lease=lease, ) self.assertIsInstance(result, ScrapeError) assert isinstance(result, ScrapeError) self.assertEqual(result.error_category.value, "timeout") self.assertEqual(result.diagnostics.timeout_phase, TimeoutPhase.QUEUE) - self.assertEqual(progress.snapshot().phase, TimeoutPhase.QUEUE) + self.assertEqual(lease.snapshot().phase, TimeoutPhase.QUEUE) def test_browser_tier_marks_boot_before_driver_then_work(self) -> None: - progress = ScrapeProgress() - _PhaseProbeDriver.progress = progress + lease = WorkLease.tracking_only() + _PhaseProbeDriver.lease = lease _PhaseProbeDriver.construction_phase = None result = _execute( scrape_request( @@ -108,7 +108,7 @@ def test_browser_tier_marks_boot_before_driver_then_work(self) -> None: navigation_mode=NavigationMode.GET, max_retries=0, ), - progress=progress, + lease=lease, request_id="req-boot-mark", Driver=_PhaseProbeDriver, ) @@ -116,7 +116,7 @@ def test_browser_tier_marks_boot_before_driver_then_work(self) -> None: self.assertEqual(_PhaseProbeDriver.construction_phase, TimeoutPhase.BOOT) _snap_eq( self, - progress.snapshot(), + lease.snapshot(), phase=TimeoutPhase.WORK, attempts=1, strategy=NavigationMode.GET, @@ -124,17 +124,17 @@ def test_browser_tier_marks_boot_before_driver_then_work(self) -> None: ) def test_request_tier_marks_work_with_attempt(self) -> None: - progress = ScrapeProgress() + lease = WorkLease.tracking_only() result = _execute( scrape_request(url=_URL, execution_mode=ExecutionMode.REQUEST), - progress=progress, + lease=lease, request_id="req-http-mark", Request=fake_request_cls(html=_HTML, url=f"{_URL}/"), ) self.assertIsInstance(result, ScrapeSuccess) _snap_eq( self, - progress.snapshot(), + lease.snapshot(), phase=TimeoutPhase.WORK, attempts=1, tier=ExecutionTier.HTTP_REQUEST, @@ -148,10 +148,10 @@ def get(self, *_a: object, **_k: object) -> None: def close(self) -> None: return None - progress = ScrapeProgress() + lease = WorkLease.tracking_only() result = _execute( scrape_request(url=_URL, execution_mode=ExecutionMode.REQUEST), - progress=progress, + lease=lease, request_id="req-http-timeout", Request=BoomRequest, ) @@ -160,22 +160,22 @@ def close(self) -> None: self.assertEqual(result.error_category.value, "timeout") self.assertEqual(result.diagnostics.timeout_phase, TimeoutPhase.WORK) self.assertEqual(result.diagnostics.execution_tier, ExecutionTier.HTTP_REQUEST) - self.assertEqual(progress.snapshot().phase, TimeoutPhase.WORK) + self.assertEqual(lease.snapshot().phase, TimeoutPhase.WORK) def test_browser_tier_timeout_exception_sets_phase(self) -> None: class BoomDriver(_PhaseProbeDriver): def get(self, *_a: object, **_k: object) -> None: raise TimeoutError("navigation timeout") - progress = ScrapeProgress() - BoomDriver.progress = progress + lease = WorkLease.tracking_only() + BoomDriver.lease = lease result = _execute( scrape_request( execution_mode=ExecutionMode.BROWSER, navigation_mode=NavigationMode.GET, max_retries=0, ), - progress=progress, + lease=lease, request_id="req-browser-timeout", Driver=BoomDriver, ) diff --git a/tests/infra/test_scrape_progress.py b/tests/engine/test_work_lease_snapshot.py similarity index 63% rename from tests/infra/test_scrape_progress.py rename to tests/engine/test_work_lease_snapshot.py index d600140..38dec60 100644 --- a/tests/infra/test_scrape_progress.py +++ b/tests/engine/test_work_lease_snapshot.py @@ -1,30 +1,30 @@ -"""ScrapeProgress snapshot unit tests.""" +"""WorkLease phase snapshot unit tests (folded from ScrapeProgress).""" from __future__ import annotations import unittest -from app.infra.scrape_progress import ScrapeProgress, ScrapeProgressSnapshot +from app.engine.work_lease import LeaseSnapshot, WorkLease from app.schemas.enums import ExecutionTier, NavigationMode, TimeoutPhase -class ScrapeProgressTests(unittest.TestCase): - def test_snapshot_defaults_to_queue(self): - snap: ScrapeProgressSnapshot = ScrapeProgress().snapshot() +class WorkLeaseSnapshotTests(unittest.TestCase): + def test_snapshot_defaults_to_queue(self) -> None: + snap: LeaseSnapshot = WorkLease.tracking_only().snapshot() self.assertEqual(snap.phase, TimeoutPhase.QUEUE) self.assertEqual(snap.attempts, 0) self.assertIsNone(snap.strategy_used) self.assertIsNone(snap.execution_tier) - def test_mark_updates_snapshot(self): - progress = ScrapeProgress() - progress.mark( + def test_mark_updates_snapshot(self) -> None: + lease = WorkLease.tracking_only() + lease.mark( TimeoutPhase.WORK, attempts=2, strategy_used=NavigationMode.GET, execution_tier=ExecutionTier.BROWSER_DRIVER, ) - snap = progress.snapshot() + snap = lease.snapshot() self.assertEqual(snap.phase, TimeoutPhase.WORK) self.assertEqual(snap.attempts, 2) self.assertEqual(snap.strategy_used, NavigationMode.GET) diff --git a/tests/support/http.py b/tests/support/http.py index 061ba5e..dad961f 100644 --- a/tests/support/http.py +++ b/tests/support/http.py @@ -10,7 +10,7 @@ from app.api.deps import get_engine from app.config import Settings from app.engine import ScraperEngine -from app.infra.scrape_progress import ScrapeProgress +from app.engine.work_lease import WorkLease from app.main import create_app from app.schemas.request import ScrapeRequest from app.schemas.response import ScrapeError, ScrapeSuccess @@ -23,7 +23,7 @@ def __call__( deadline_monotonic: float | None = ..., *, request_id: str | None = ..., - progress: ScrapeProgress | None = ..., + lease: WorkLease | None = ..., ) -> ScrapeSuccess | ScrapeError: ... @@ -43,13 +43,13 @@ def execute( deadline_monotonic: float | None = None, *, request_id: str | None = None, - progress: ScrapeProgress | None = None, + lease: WorkLease | None = None, ) -> ScrapeSuccess | ScrapeError: return self._execute_fn( payload, deadline_monotonic, request_id=request_id, - progress=progress, + lease=lease, ) From e93eafb667607f204bb2bfeed26d52a3875f25e8 Mon Sep 17 00:00:00 2001 From: Gil Desmarais Date: Sun, 30 Aug 2026 10:27:54 +0200 Subject: [PATCH 2/5] fix(engine): fail closed unclean/hung surfaces as challenge_block Delete soft-retry (may_retry_strategies) and prefer-challenge continue paths so unclean or unreadable hung pages always emit challenge_block. --- app/engine/browser_tier.py | 53 ++++++--------------- app/engine/strategies.py | 6 +-- app/infra/detector.py | 14 ------ tests/engine/test_browser_tier_challenge.py | 33 ++++++------- tests/engine/test_scraper_engine.py | 1 - tests/infra/test_challenge_detector.py | 15 ------ 6 files changed, 34 insertions(+), 88 deletions(-) diff --git a/app/engine/browser_tier.py b/app/engine/browser_tier.py index ee65d6d..d76030b 100644 --- a/app/engine/browser_tier.py +++ b/app/engine/browser_tier.py @@ -12,7 +12,6 @@ browser_step_budget_seconds, elapsed_ms, is_timeout_exception, - remaining_work_seconds, ) from app.engine.driver_capabilities import DriverProtocol, call_if_available from app.engine.envelope import TIMEOUT_ERROR_BY_PHASE, build_error, build_success @@ -120,11 +119,8 @@ def _surface_unclean( strategy: NavigationMode, started_monotonic: float, assessment: ChallengeAssessment, - has_more_strategies: bool, - remaining_work: int, - collector: XhrCollector, -) -> ScrapeError | None: - """Log and apply unclean assessment. None ⇒ soft-retry / continue.""" +) -> ScrapeError: + """Fail closed: any unclean assessment is challenge_block (no soft-retry).""" logger.warning( "scrape_challenge_detected request_id=%s host=%s strategy=%s " "attempt=%d marker=%s", @@ -134,12 +130,6 @@ def _surface_unclean( attempts, assessment.detected_marker, ) - if assessment.may_retry_strategies( - has_more=has_more_strategies, - remaining_seconds=remaining_work, - ): - collector.reset() - return None return _challenge_block_error( target_url, request_id=request_id, @@ -150,6 +140,13 @@ def _surface_unclean( ) +_UNREADABLE_SURFACE = ChallengeAssessment( + blocked_detected=True, + challenge_detected=True, + detected_marker="unreadable_surface", +) + + def _boot_storage_error( target_url: str, request_id: str, @@ -345,7 +342,6 @@ def run_browser_tier( strategy=strategy, started_monotonic=started_monotonic, ) - remaining_work = remaining_work_seconds(settings, browser_ready_monotonic) try: navigate(driver, target_url, strategy, step_budget) mid_wait = wait_for_readiness( @@ -354,20 +350,14 @@ def run_browser_tier( timeout_seconds=min(payload.wait_timeout_seconds, step_budget), ) if mid_wait is not None: - blocked = _surface_unclean( + return _surface_unclean( target_url, request_id=request_id, attempts=attempts, strategy=strategy, started_monotonic=started_monotonic, assessment=mid_wait, - has_more_strategies=has_more, - remaining_work=remaining_work, - collector=collector, ) - if blocked is None: - continue - return blocked if payload.scroll: apply_scrolling(driver) @@ -377,20 +367,14 @@ def run_browser_tier( ) if not assessment.is_clean: - blocked = _surface_unclean( + return _surface_unclean( target_url, request_id=request_id, attempts=attempts, strategy=strategy, started_monotonic=started_monotonic, assessment=assessment, - has_more_strategies=has_more, - remaining_work=remaining_work, - collector=collector, ) - if blocked is None: - continue - return blocked return build_success( target_url, @@ -418,25 +402,16 @@ def run_browser_tier( attempt_index, str(exc), ) - # Prefer knowable challenge over timeout when the page is inspectable. assessment = _inspect_assessment(driver) - if assessment is not None and not assessment.is_clean: - blocked = _surface_unclean( + if assessment is None or not assessment.is_clean: + return _surface_unclean( target_url, request_id=request_id, attempts=attempts, strategy=strategy, started_monotonic=started_monotonic, - assessment=assessment, - has_more_strategies=has_more, - remaining_work=remaining_work_seconds( - settings, browser_ready_monotonic - ), - collector=collector, + assessment=assessment or _UNREADABLE_SURFACE, ) - if blocked is None: - continue - return blocked if has_more: collector.reset() continue diff --git a/app/engine/strategies.py b/app/engine/strategies.py index 7577fef..8a9f4c7 100644 --- a/app/engine/strategies.py +++ b/app/engine/strategies.py @@ -131,10 +131,10 @@ def wait_for_readiness( selector: str | None, timeout_seconds: int, ) -> ChallengeAssessment | None: - """Wait for selector / settle; return unclean assessment if challenge appears mid-wait. + """Wait for selector / settle; return unclean assessment if challenge appears. - Selector waits run in ≤2s chunks so a challenge interstitial can fail closed - before the full wait budget burns. Non-selector settles stay short and probe once. + Selector waits run in ≤2s chunks so a challenge interstitial fails closed + before the full wait budget burns. Non-selector settles probe once. """ if not selector: if ( diff --git a/app/infra/detector.py b/app/infra/detector.py index d77dc9d..f8bda28 100644 --- a/app/infra/detector.py +++ b/app/infra/detector.py @@ -28,9 +28,6 @@ (m, m.lower()) for m in _CHALLENGE_MARKERS ) -# Soft strategy retries need this much remaining work budget (seconds). -_MIN_SOFT_RETRY_REMAINING_SECONDS = 5 - _DRIVER_SIGNAL_METHODS: tuple[tuple[str, str], ...] = ( ("is_bot_detected", "botasaurus_driver_bot_detected"), ("is_in_challenge", "botasaurus_driver_challenge"), @@ -48,17 +45,6 @@ class ChallengeAssessment: def is_clean(self) -> bool: return not self.blocked_detected and not self.challenge_detected - def may_retry_strategies( - self, *, has_more: bool, remaining_seconds: float | int - ) -> bool: - """Soft challenge markers may retry strategies; hard HTTP blocks do not. - - Requires enough remaining *work* budget so another strategy can finish. - """ - if remaining_seconds < _MIN_SOFT_RETRY_REMAINING_SECONDS: - return False - return self.challenge_detected and has_more - def to_signal(self) -> ChallengeSignal: """Convert domain assessment to wire ChallengeSignal DTO.""" return ChallengeSignal( diff --git a/tests/engine/test_browser_tier_challenge.py b/tests/engine/test_browser_tier_challenge.py index cf86350..44997e9 100644 --- a/tests/engine/test_browser_tier_challenge.py +++ b/tests/engine/test_browser_tier_challenge.py @@ -1,4 +1,4 @@ -"""Browser-tier challenge-before-timeout and hard-block abort.""" +"""Browser-tier challenge-before-timeout and fail-closed unclean surfaces.""" from __future__ import annotations @@ -141,42 +141,43 @@ def test_hard_block_aborts_remaining_strategies(self) -> None: self.assertTrue(result.diagnostics.challenge.blocked) self.assertFalse(result.diagnostics.challenge.detected) - def test_soft_challenge_retries_strategies_then_blocks(self) -> None: + def test_challenge_fails_closed_without_strategy_retry(self) -> None: result = self._execute( _SoftChallengeDriver, navigation_mode=NavigationMode.AUTO, max_retries=2, - request_id="req-soft-challenge", + request_id="req-challenge-once", ) self.assertEqual(result.error_category, ErrorCategory.CHALLENGE_BLOCK) - self.assertEqual(result.diagnostics.attempts, 3) - self.assertEqual(_SoftChallengeDriver.navigate_calls, 3) + self.assertEqual(result.diagnostics.attempts, 1) + self.assertEqual(_SoftChallengeDriver.navigate_calls, 1) assert result.diagnostics.challenge is not None self.assertTrue(result.diagnostics.challenge.detected) - def test_soft_challenge_does_not_retry_when_work_budget_low(self) -> None: - """Below soft-retry floor, unclean assessment is challenge_block immediately.""" - _SoftChallengeDriver.reset() + def test_unreadable_hung_surface_returns_challenge_block(self) -> None: + _CleanTimeoutDriver.reset() payload = scrape_request( execution_mode=ExecutionMode.BROWSER, - navigation_mode=NavigationMode.AUTO, - max_retries=2, + navigation_mode=NavigationMode.GET, + max_retries=0, ) with tempfile.TemporaryDirectory() as tmp: engine = ScraperEngine(settings=get_settings(), runtime_root=Path(tmp)) with ( - patch("botasaurus.browser.Driver", _SoftChallengeDriver), + patch("botasaurus.browser.Driver", _CleanTimeoutDriver), patch( - "app.engine.browser_tier.remaining_work_seconds", - return_value=4, + "app.engine.browser_tier._inspect_assessment", + return_value=None, ), ): - result = engine.execute(payload, request_id="req-soft-low-budget") + result = engine.execute(payload, request_id="req-unreadable") self.assertIsInstance(result, ScrapeError) assert isinstance(result, ScrapeError) self.assertEqual(result.error_category, ErrorCategory.CHALLENGE_BLOCK) - self.assertEqual(result.diagnostics.attempts, 1) - self.assertEqual(_SoftChallengeDriver.navigate_calls, 1) + self.assertIsNone(result.diagnostics.timeout_phase) + assert result.diagnostics.challenge is not None + self.assertEqual(result.diagnostics.challenge.marker, "unreadable_surface") + self.assertEqual(_CleanTimeoutDriver.navigate_calls, 1) def test_mid_wait_challenge_returns_challenge_block(self) -> None: class _MidWaitChallengeDriver(_ScenarioDriver): diff --git a/tests/engine/test_scraper_engine.py b/tests/engine/test_scraper_engine.py index e3bcc3e..c3e1540 100644 --- a/tests/engine/test_scraper_engine.py +++ b/tests/engine/test_scraper_engine.py @@ -52,7 +52,6 @@ def get(self, *_args: object, **kwargs: Any) -> None: 1020.0, # browser ready after boot 1020.0, # remaining total (step budget) 1020.0, # remaining work (step budget) - 1020.0, # remaining work (soft-retry gate) 1020.0, # render_ms ] diff --git a/tests/infra/test_challenge_detector.py b/tests/infra/test_challenge_detector.py index 4d28817..feaf49e 100644 --- a/tests/infra/test_challenge_detector.py +++ b/tests/infra/test_challenge_detector.py @@ -25,21 +25,6 @@ def test_clean_response(self): self.assertFalse(res.blocked_detected) self.assertFalse(res.challenge_detected) - def test_soft_challenge_may_retry_strategies(self): - soft = ChallengeDetector.detect("Just a moment...", 200) - self.assertTrue(soft.may_retry_strategies(has_more=True, remaining_seconds=30)) - self.assertFalse( - soft.may_retry_strategies(has_more=False, remaining_seconds=30) - ) - self.assertFalse(soft.may_retry_strategies(has_more=True, remaining_seconds=4)) - - def test_hard_block_does_not_retry_strategies(self): - hard = ChallengeDetector.detect("Forbidden", 403) - self.assertFalse(hard.may_retry_strategies(has_more=True, remaining_seconds=30)) - self.assertFalse( - hard.may_retry_strategies(has_more=False, remaining_seconds=30) - ) - def test_driver_bot_detection_integration(self): mock_driver = MagicMock() mock_driver.is_bot_detected.return_value = True From dcb4ac7c2e9951574c5fd0932df06de151e72453 Mon Sep 17 00:00:00 2001 From: Gil Desmarais Date: Sun, 30 Aug 2026 10:30:11 +0200 Subject: [PATCH 3/5] test(challenge): share interstitial fixture corpus with gem Own Cloudflare/DataDome/Vercel HTML under tests/fixtures/challenge, trim one-sided detector markers, and document the shared corpus home. --- AGENTS.md | 25 ++++++----- app/infra/detector.py | 18 ++++---- tests/fixtures/challenge/README.md | 8 ++++ tests/fixtures/challenge/clean.html | 1 + .../challenge/cloudflare_interstitial.html | 9 ++++ .../challenge/datadome_interstitial.html | 6 +++ .../fixtures/challenge/vercel_checkpoint.html | 4 ++ tests/infra/test_challenge_corpus.py | 44 +++++++++++++++++++ 8 files changed, 95 insertions(+), 20 deletions(-) create mode 100644 tests/fixtures/challenge/README.md create mode 100644 tests/fixtures/challenge/clean.html create mode 100644 tests/fixtures/challenge/cloudflare_interstitial.html create mode 100644 tests/fixtures/challenge/datadome_interstitial.html create mode 100644 tests/fixtures/challenge/vercel_checkpoint.html create mode 100644 tests/infra/test_challenge_corpus.py diff --git a/AGENTS.md b/AGENTS.md index 7e6992d..d06c59e 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -42,19 +42,22 @@ app/ security/ # UrlGuard SSRF guardrails scripts/ bench_scrape.py # TestClient wall-time bench for POST /scrape (request tier) -tests/ - api/ # HTTP contract, request schema, 504 timeout envelope tests - domain/ # ScrapeService unit tests (timeout error mapping) - engine/ # ScraperEngine units, isolation regressions, timeout progress - infra/ # challenge, metadata, xhr, progress, sentry, telemetry, request-id, cleanup - security/ # UrlGuard tests - support/ - http.py # test_client() context manager + dependency_overrides helper - fakes.py # shared FakeDriver, FakeRequest, fake_request_cls, ... - factories.py # scrape_request(), example_url() - test_bench_regression.py # lightweight guard that bench script completes (root: guards scripts/) + tests/ + api/ # HTTP contract, request schema, 504 timeout envelope tests + domain/ # ScrapeService unit tests (timeout error mapping) + engine/ # ScraperEngine units, isolation regressions, timeout progress + fixtures/challenge/ # shared interstitial HTML corpus (gem BlockedSurface loads sibling path) + infra/ # challenge, metadata, xhr, sentry, telemetry, request-id, cleanup + security/ # UrlGuard tests + support/ + http.py # test_client() context manager + dependency_overrides helper + fakes.py # shared FakeDriver, FakeRequest, fake_request_cls, ... + factories.py # scrape_request(), example_url() + test_bench_regression.py # lightweight guard that bench script completes (root: guards scripts/) ``` +**Challenge corpus:** HTML under `tests/fixtures/challenge/` is the single fixture home for interstitial detection. Scrape-api `ChallengeDetector` markers and gem `BlockedSurface` signatures must both assert those files; do not reintroduce one-sided markers without a shared fixture. + Layer rules: | Layer | May import | Must not import | diff --git a/app/infra/detector.py b/app/infra/detector.py index f8bda28..248ebe7 100644 --- a/app/infra/detector.py +++ b/app/infra/detector.py @@ -8,20 +8,20 @@ from app.engine.driver_capabilities import DriverProtocol, call_if_available from app.schemas.response import ChallengeSignal +# Markers shared with html2rss BlockedSurface via tests/fixtures/challenge/. +# One-sided substrings without a shared fixture were deleted in the corpus trim. _CHALLENGE_MARKERS: tuple[str, ...] = ( - "challenge-error-text", - "Enable JavaScript and cookies to continue", "Just a moment...", + "checking your browser before accessing", + "Enable JavaScript and cookies to continue", "cf-challenge", - "cf-turnstile", + "cdn-cgi/challenge-platform", + "cloudflare ray id", "captcha-delivery.com", "datadome", - "DataDome CAPTCHA", - "/captcha/?", - "attention required! | cloudflare", - "cloudflare-ray-id", - "shield.recaptcha.net", - "geo.captcha-delivery.com", + "Vercel Security Checkpoint", + "checking the security", + "vercel.com/security", ) _CHALLENGE_MARKERS_PAIRS: tuple[tuple[str, str], ...] = tuple( diff --git a/tests/fixtures/challenge/README.md b/tests/fixtures/challenge/README.md new file mode 100644 index 0000000..87c6348 --- /dev/null +++ b/tests/fixtures/challenge/README.md @@ -0,0 +1,8 @@ +# Challenge corpus + +Shared HTML fixtures for Botasaurus scrape-api `ChallengeDetector` and the html2rss gem `BlockedSurface` module. + +**Owner:** `botasaurus-scrape-api/tests/fixtures/challenge/` +**Consumers:** scrape-api unit tests; gem `blocked_surface_spec` loads these files via the sibling path under the org workspace. + +Do not duplicate interstitial HTML in gem specs — assert against these files. diff --git a/tests/fixtures/challenge/clean.html b/tests/fixtures/challenge/clean.html new file mode 100644 index 0000000..41d078d --- /dev/null +++ b/tests/fixtures/challenge/clean.html @@ -0,0 +1 @@ +

Example Domain

This domain is for use in illustrative examples.

diff --git a/tests/fixtures/challenge/cloudflare_interstitial.html b/tests/fixtures/challenge/cloudflare_interstitial.html new file mode 100644 index 0000000..2af0bde --- /dev/null +++ b/tests/fixtures/challenge/cloudflare_interstitial.html @@ -0,0 +1,9 @@ + +Just a moment... + +Checking your browser before accessing example.com. +Please enable JavaScript and cookies to continue. +
+Cloudflare Ray ID: 9a1b2c3d4e5f + + diff --git a/tests/fixtures/challenge/datadome_interstitial.html b/tests/fixtures/challenge/datadome_interstitial.html new file mode 100644 index 0000000..d4314e2 --- /dev/null +++ b/tests/fixtures/challenge/datadome_interstitial.html @@ -0,0 +1,6 @@ + + + +
DataDome interstitial challenge
+ + diff --git a/tests/fixtures/challenge/vercel_checkpoint.html b/tests/fixtures/challenge/vercel_checkpoint.html new file mode 100644 index 0000000..f8df739 --- /dev/null +++ b/tests/fixtures/challenge/vercel_checkpoint.html @@ -0,0 +1,4 @@ + +Vercel Security Checkpoint +Checking the security of your connection before accessing example.com. + diff --git a/tests/infra/test_challenge_corpus.py b/tests/infra/test_challenge_corpus.py new file mode 100644 index 0000000..132e73c --- /dev/null +++ b/tests/infra/test_challenge_corpus.py @@ -0,0 +1,44 @@ +"""Shared challenge corpus: fixtures asserted by ChallengeDetector.""" + +from __future__ import annotations + +import unittest +from pathlib import Path + +from app.infra.detector import ChallengeDetector + +_FIXTURE_DIR = Path(__file__).resolve().parents[1] / "fixtures" / "challenge" + +_POSITIVE = ( + "cloudflare_interstitial.html", + "datadome_interstitial.html", + "vercel_checkpoint.html", +) + + +class ChallengeCorpusTests(unittest.TestCase): + def test_positive_fixtures_are_unclean(self) -> None: + for name in _POSITIVE: + with self.subTest(fixture=name): + html = (_FIXTURE_DIR / name).read_text(encoding="utf-8") + assessment = ChallengeDetector.detect(html, 200) + self.assertFalse(assessment.is_clean, msg=name) + self.assertTrue(assessment.blocked_detected, msg=name) + + def test_clean_fixture_is_clean(self) -> None: + html = (_FIXTURE_DIR / "clean.html").read_text(encoding="utf-8") + assessment = ChallengeDetector.detect(html, 200) + self.assertTrue(assessment.is_clean) + + def test_cloudflare_fixture_reports_marker(self) -> None: + html = (_FIXTURE_DIR / "cloudflare_interstitial.html").read_text( + encoding="utf-8" + ) + assessment = ChallengeDetector.detect(html, 200) + self.assertIsNotNone(assessment.detected_marker) + assert assessment.detected_marker is not None + self.assertIn("moment", assessment.detected_marker.lower()) + + +if __name__ == "__main__": + unittest.main() From 1c318d8c040beee0434adce279cbd65b81bc5785 Mon Sep 17 00:00:00 2001 From: Gil Desmarais Date: Sun, 30 Aug 2026 10:36:13 +0200 Subject: [PATCH 4/5] fix(session): keep reclaim close effective across cold-boot race Close idempotently by clearing the driver handle (not a pre-assign flag), retain a strong ref to shielded executor work, and cover early reclaim. --- app/engine/session.py | 20 +++++++++----- app/engine/work_lease.py | 8 ++++++ tests/engine/test_session_reclaim.py | 41 ++++++++++++++++++++++++++++ 3 files changed, 62 insertions(+), 7 deletions(-) create mode 100644 tests/engine/test_session_reclaim.py diff --git a/app/engine/session.py b/app/engine/session.py index b02e954..a5f169c 100644 --- a/app/engine/session.py +++ b/app/engine/session.py @@ -4,6 +4,7 @@ import errno import shutil +import threading from pathlib import Path from typing import TYPE_CHECKING, Any @@ -34,7 +35,7 @@ def __init__( self.adopted_profile_dir: Path | None = None self.warm_fingerprint: DriverFingerprint | None = None self.warm_hit: bool | None = None - self._closed = False + self._close_lock = threading.Lock() def __enter__(self) -> ScrapeSession: self.engine.register_request_id(self.request_id) @@ -48,12 +49,17 @@ def __enter__(self) -> ScrapeSession: return self def force_close(self) -> None: - """Lease-deadline reclaim: close Chromium once so the worker slot can free.""" - if self._closed: - return - self._closed = True - if self.driver is not None: - call_quietly(self.driver, "close") + """Lease-deadline reclaim: close Chromium if present (idempotent). + + Must not arm a one-shot flag while driver is still None — cold boot + assigns the driver after session enter; an early reclaim must leave + __exit__ able to close the later-assigned driver. + """ + with self._close_lock: + driver = self.driver + self.driver = None + if driver is not None: + call_quietly(driver, "close") def prepare_runtime_dir(self) -> None: """Create the request runtime dir only (warm-path adoption).""" diff --git a/app/engine/work_lease.py b/app/engine/work_lease.py index 84d52bb..4e9417e 100644 --- a/app/engine/work_lease.py +++ b/app/engine/work_lease.py @@ -82,6 +82,7 @@ def __init__( self._lock = threading.Lock() self._snap = LeaseSnapshot() self._reclaim_hooks: list[Callable[[], None]] = [] + self._inflight: asyncio.Future[object] | None = None @classmethod def tracking_only(cls, settings: Settings | None = None) -> WorkLease: @@ -180,6 +181,8 @@ async def run(self, *, host: str, work: Callable[[], T]) -> T: loop = asyncio.get_running_loop() awaitable = asyncio.ensure_future(loop.run_in_executor(self.executor, work)) + # Strong ref: shielded tasks are only weakly held by the loop. + self._inflight = awaitable try: return await asyncio.wait_for( asyncio.shield(awaitable), timeout=remaining @@ -187,5 +190,10 @@ async def run(self, *, host: str, work: Callable[[], T]) -> T: except TimeoutError: self.reclaim() raise + finally: + if awaitable.done(): + self._inflight = None + else: + awaitable.add_done_callback(lambda _f: None) finally: self.host_gate.release(host) diff --git a/tests/engine/test_session_reclaim.py b/tests/engine/test_session_reclaim.py new file mode 100644 index 0000000..cecd680 --- /dev/null +++ b/tests/engine/test_session_reclaim.py @@ -0,0 +1,41 @@ +"""ScrapeSession force_close remains effective after early reclaim.""" + +from __future__ import annotations + +import tempfile +import unittest +from pathlib import Path +from unittest.mock import MagicMock + +from app.config import get_settings +from app.engine import ScraperEngine +from app.engine.session import ScrapeSession +from app.engine.work_lease import WorkLease + + +class SessionReclaimTests(unittest.TestCase): + def test_early_reclaim_without_driver_still_closes_later_assign(self) -> None: + """Deadline reclaim during cold boot must not skip the later driver close.""" + with tempfile.TemporaryDirectory() as tmp: + engine = ScraperEngine(settings=get_settings(), runtime_root=Path(tmp)) + lease = WorkLease.tracking_only(engine.settings) + with ScrapeSession(engine, "req-reclaim-boot", lease=lease) as session: + lease.reclaim() # driver still None (cold boot race) + driver = MagicMock() + session.driver = driver + driver.close.assert_called_once() + + def test_reclaim_closes_assigned_driver_once(self) -> None: + with tempfile.TemporaryDirectory() as tmp: + engine = ScraperEngine(settings=get_settings(), runtime_root=Path(tmp)) + lease = WorkLease.tracking_only(engine.settings) + with ScrapeSession(engine, "req-reclaim-once", lease=lease) as session: + driver = MagicMock() + session.driver = driver + lease.reclaim() + lease.reclaim() + driver.close.assert_called_once() + + +if __name__ == "__main__": + unittest.main() From 837bf85ba03cd821f34ad9b886ce2e5181792d80 Mon Sep 17 00:00:00 2001 From: Gil Desmarais Date: Sun, 30 Aug 2026 10:40:58 +0200 Subject: [PATCH 5/5] fix(lease): abort before Chromium and fail-closed mid-wait unreadables Set lease.aborted on reclaim so workers exit without booting after host gate release; treat unreadable mid-wait probes as challenge_block. --- app/engine/browser_tier.py | 38 +++++++++++++++++------ app/engine/orchestrator.py | 12 ++++++++ app/engine/session.py | 5 ++-- app/engine/strategies.py | 20 +++++++++---- app/engine/work_lease.py | 13 ++++++-- app/infra/detector.py | 7 +++++ tests/domain/test_work_lease_reclaim.py | 15 ++++++++++ tests/engine/test_session_reclaim.py | 28 +++++++++++++++-- tests/engine/test_wait_for_readiness.py | 40 +++++++++++++++++++++++++ 9 files changed, 157 insertions(+), 21 deletions(-) create mode 100644 tests/engine/test_wait_for_readiness.py diff --git a/app/engine/browser_tier.py b/app/engine/browser_tier.py index d76030b..faf91a8 100644 --- a/app/engine/browser_tier.py +++ b/app/engine/browser_tier.py @@ -26,7 +26,11 @@ ) from app.engine.warm_pool import DriverFingerprint from app.engine.work_lease import WorkLease -from app.infra.detector import ChallengeAssessment, ChallengeDetector +from app.infra.detector import ( + UNREADABLE_SURFACE, + ChallengeAssessment, + ChallengeDetector, +) from app.infra.metadata import MetadataExtractor, MetadataResult from app.infra.xhr_collector import XhrCollector from app.logging_config import get_logger @@ -140,13 +144,6 @@ def _surface_unclean( ) -_UNREADABLE_SURFACE = ChallengeAssessment( - blocked_detected=True, - challenge_detected=True, - detected_marker="unreadable_surface", -) - - def _boot_storage_error( target_url: str, request_id: str, @@ -238,6 +235,17 @@ def run_browser_tier( TimeoutPhase.BOOT, execution_tier=ExecutionTier.BROWSER_DRIVER, ) + if lease.aborted: + return build_error( + target_url, + TIMEOUT_ERROR_BY_PHASE[TimeoutPhase.BOOT], + request_id=request_id, + error_category=ErrorCategory.TIMEOUT, + attempts=0, + render_ms=elapsed_ms(started_monotonic), + execution_tier=ExecutionTier.BROWSER_DRIVER, + timeout_phase=TimeoutPhase.BOOT, + ) fingerprint = DriverFingerprint.from_request(payload) session.warm_fingerprint = fingerprint warm_hit = False @@ -262,6 +270,18 @@ def run_browser_tier( except OSError as exc: return _boot_storage_error(target_url, request_id, started_monotonic, exc) + if lease.aborted: + return build_error( + target_url, + TIMEOUT_ERROR_BY_PHASE[TimeoutPhase.BOOT], + request_id=request_id, + error_category=ErrorCategory.TIMEOUT, + attempts=0, + render_ms=elapsed_ms(started_monotonic), + execution_tier=ExecutionTier.BROWSER_DRIVER, + timeout_phase=TimeoutPhase.BOOT, + ) + driver_window_size = ( [payload.window_size.width, payload.window_size.height] if payload.window_size @@ -410,7 +430,7 @@ def run_browser_tier( attempts=attempts, strategy=strategy, started_monotonic=started_monotonic, - assessment=assessment or _UNREADABLE_SURFACE, + assessment=assessment or UNREADABLE_SURFACE, ) if has_more: collector.reset() diff --git a/app/engine/orchestrator.py b/app/engine/orchestrator.py index bec0b1d..66d0f68 100644 --- a/app/engine/orchestrator.py +++ b/app/engine/orchestrator.py @@ -105,6 +105,18 @@ def execute( started_monotonic = now lease = lease or WorkLease.tracking_only(self.settings) + if lease.aborted: + # Outer deadline reclaimed before/without session enter — exit without + # Chromium so host-gate release stays honest vs live browsers. + phase = lease.snapshot().phase + return build_error( + target_url, + TIMEOUT_ERROR_BY_PHASE[phase], + request_id=resolved_request_id, + error_category=ErrorCategory.TIMEOUT, + timeout_phase=phase, + ) + if deadline_monotonic is not None and now >= deadline_monotonic: lease.mark(TimeoutPhase.QUEUE) return build_error( diff --git a/app/engine/session.py b/app/engine/session.py index a5f169c..f802c16 100644 --- a/app/engine/session.py +++ b/app/engine/session.py @@ -36,6 +36,9 @@ def __init__( self.warm_fingerprint: DriverFingerprint | None = None self.warm_hit: bool | None = None self._close_lock = threading.Lock() + # Register before __enter__ so reclaim during prepare_runtime still binds. + if self.lease is not None: + self.lease.register_reclaim(self.force_close) def __enter__(self) -> ScrapeSession: self.engine.register_request_id(self.request_id) @@ -44,8 +47,6 @@ def __enter__(self) -> ScrapeSession: except Exception: self.engine.unregister_request_id(self.request_id) raise - if self.lease is not None: - self.lease.register_reclaim(self.force_close) return self def force_close(self) -> None: diff --git a/app/engine/strategies.py b/app/engine/strategies.py index 8a9f4c7..60610da 100644 --- a/app/engine/strategies.py +++ b/app/engine/strategies.py @@ -12,7 +12,11 @@ resolve_callable, resolve_cdp_tab, ) -from app.infra.detector import ChallengeAssessment, ChallengeDetector +from app.infra.detector import ( + UNREADABLE_SURFACE, + ChallengeAssessment, + ChallengeDetector, +) from app.infra.xhr_collector import XhrCollector from app.logging_config import get_logger from app.schemas.enums import NavigationMode @@ -114,11 +118,15 @@ def configure_driver( def _mid_wait_challenge(driver: DriverProtocol) -> ChallengeAssessment | None: - """Best-effort challenge probe during readiness wait. None if clean or unreadable.""" + """Probe during readiness wait. None if clean; unclean/unreadable otherwise. + + Unreadable surfaces fail closed immediately so selector waits do not burn + remaining budget on a hostile/hung tab. + """ try: html = driver.page_html or "" except Exception: - return None + return UNREADABLE_SURFACE assessment = ChallengeDetector.detect(html, driver=driver) if assessment.is_clean: return None @@ -131,10 +139,10 @@ def wait_for_readiness( selector: str | None, timeout_seconds: int, ) -> ChallengeAssessment | None: - """Wait for selector / settle; return unclean assessment if challenge appears. + """Wait for selector / settle; return unclean/unreadable assessment to fail closed. - Selector waits run in ≤2s chunks so a challenge interstitial fails closed - before the full wait budget burns. Non-selector settles probe once. + Selector waits run in ≤2s chunks so a challenge or unreadable surface fails + closed before the full wait budget burns. Non-selector settles probe once. """ if not selector: if ( diff --git a/app/engine/work_lease.py b/app/engine/work_lease.py index 4e9417e..7ddd062 100644 --- a/app/engine/work_lease.py +++ b/app/engine/work_lease.py @@ -65,8 +65,9 @@ class WorkLease: """Owns admit → run → terminal timeout; reclaim kills the session driver once. Progress phase lives on the lease (replaces ScrapeProgress). Outer deadline - calls register_reclaim hooks (session force-close); Future.cancel is not the - Chromium reclaim story. + sets ``aborted`` and runs register_reclaim hooks (session force-close); + Future.cancel is not the Chromium reclaim story. Workers check ``aborted`` + before booting Chromium so host-gate release stays honest. """ def __init__( @@ -83,6 +84,7 @@ def __init__( self._snap = LeaseSnapshot() self._reclaim_hooks: list[Callable[[], None]] = [] self._inflight: asyncio.Future[object] | None = None + self._aborted = threading.Event() @classmethod def tracking_only(cls, settings: Settings | None = None) -> WorkLease: @@ -91,6 +93,11 @@ def tracking_only(cls, settings: Settings | None = None) -> WorkLease: return cls(settings=settings or get_settings()) + @property + def aborted(self) -> bool: + """True after outer deadline reclaim — workers must not boot Chromium.""" + return self._aborted.is_set() + def mark( self, phase: TimeoutPhase, @@ -124,6 +131,8 @@ def register_reclaim(self, hook: Callable[[], None]) -> None: self._reclaim_hooks.append(hook) def reclaim(self) -> None: + """Mark abort then run hooks so late-bound sessions still force-close.""" + self._aborted.set() with self._lock: hooks = list(self._reclaim_hooks) for hook in hooks: diff --git a/app/infra/detector.py b/app/infra/detector.py index 248ebe7..1341f79 100644 --- a/app/infra/detector.py +++ b/app/infra/detector.py @@ -54,6 +54,13 @@ def to_signal(self) -> ChallengeSignal: ) +UNREADABLE_SURFACE = ChallengeAssessment( + blocked_detected=True, + challenge_detected=True, + detected_marker="unreadable_surface", +) + + class ChallengeDetector: """Deep module encapsulating anti-bot challenge heuristics, HTTP status codes, and driver checks.""" diff --git a/tests/domain/test_work_lease_reclaim.py b/tests/domain/test_work_lease_reclaim.py index 1f0dd8d..47c70f9 100644 --- a/tests/domain/test_work_lease_reclaim.py +++ b/tests/domain/test_work_lease_reclaim.py @@ -75,8 +75,23 @@ def mark_reclaim() -> None: result = await lease.run(host="ok.example", work=lambda: "ok") self.assertEqual(result, "ok") self.assertEqual(reclaim_calls, 0) + self.assertFalse(lease.aborted) executor.shutdown(wait=False, cancel_futures=True) + async def test_reclaim_sets_aborted_before_hooks(self) -> None: + settings = _fast_settings(max_per_host=1, timeout_seconds=5) + lease = WorkLease.tracking_only(settings) + seen_aborted = False + + def hook() -> None: + nonlocal seen_aborted + seen_aborted = lease.aborted + + lease.register_reclaim(hook) + lease.reclaim() + self.assertTrue(lease.aborted) + self.assertTrue(seen_aborted) + if __name__ == "__main__": unittest.main() diff --git a/tests/engine/test_session_reclaim.py b/tests/engine/test_session_reclaim.py index cecd680..e783a28 100644 --- a/tests/engine/test_session_reclaim.py +++ b/tests/engine/test_session_reclaim.py @@ -1,16 +1,19 @@ -"""ScrapeSession force_close remains effective after early reclaim.""" +"""ScrapeSession force_close and abort-before-enter reclaim honesty.""" from __future__ import annotations import tempfile import unittest from pathlib import Path -from unittest.mock import MagicMock +from unittest.mock import MagicMock, patch from app.config import get_settings from app.engine import ScraperEngine from app.engine.session import ScrapeSession from app.engine.work_lease import WorkLease +from app.schemas.enums import ExecutionMode, NavigationMode +from app.schemas.response import ScrapeError +from tests.support.factories import scrape_request class SessionReclaimTests(unittest.TestCase): @@ -36,6 +39,27 @@ def test_reclaim_closes_assigned_driver_once(self) -> None: lease.reclaim() driver.close.assert_called_once() + def test_aborted_lease_skips_session_without_chromium(self) -> None: + """Deadline reclaim before worker enter must not start Chromium.""" + with tempfile.TemporaryDirectory() as tmp: + engine = ScraperEngine(settings=get_settings(), runtime_root=Path(tmp)) + lease = WorkLease.tracking_only(engine.settings) + lease.reclaim() + driver_ctor = MagicMock() + with patch("botasaurus.browser.Driver", driver_ctor): + result = engine.execute( + scrape_request( + execution_mode=ExecutionMode.BROWSER, + navigation_mode=NavigationMode.GET, + ), + request_id="req-aborted-early", + lease=lease, + ) + self.assertIsInstance(result, ScrapeError) + assert isinstance(result, ScrapeError) + self.assertEqual(result.error_category.value, "timeout") + driver_ctor.assert_not_called() + if __name__ == "__main__": unittest.main() diff --git a/tests/engine/test_wait_for_readiness.py b/tests/engine/test_wait_for_readiness.py new file mode 100644 index 0000000..a7510e9 --- /dev/null +++ b/tests/engine/test_wait_for_readiness.py @@ -0,0 +1,40 @@ +"""wait_for_readiness fail-closed on unreadable mid-wait surfaces.""" + +from __future__ import annotations + +import unittest +from typing import Any + +from app.engine.strategies import wait_for_readiness +from app.infra.detector import UNREADABLE_SURFACE + + +class _UnreadableWaitDriver: + wait_calls = 0 + + def wait_for_element(self, *_args: object, **_kwargs: Any) -> None: + type(self).wait_calls += 1 + raise TimeoutError("element not found") + + @property + def page_html(self) -> str: + raise RuntimeError("tab unreadable") + + +class WaitForReadinessTests(unittest.TestCase): + def test_unreadable_mid_wait_fails_closed_without_budget_burn(self) -> None: + _UnreadableWaitDriver.wait_calls = 0 + assessment = wait_for_readiness( + _UnreadableWaitDriver(), # type: ignore[arg-type] + selector="#content", + timeout_seconds=6, + ) + self.assertIsNotNone(assessment) + assert assessment is not None + self.assertEqual(assessment.detected_marker, UNREADABLE_SURFACE.detected_marker) + # One chunk then fail closed — not three 2s burns. + self.assertEqual(_UnreadableWaitDriver.wait_calls, 1) + + +if __name__ == "__main__": + unittest.main()