Skip to main content

stygian_browser/
acquisition.rs

1//! Opinionated acquisition runner with deterministic escalation.
2//!
3//! The runner executes a mode-specific strategy ladder and returns a terminal
4//! [`AcquisitionResult`] for every request, including setup-failure and timeout
5//! paths.
6
7use std::sync::Arc;
8use std::time::{Duration, Instant};
9
10#[cfg(feature = "browserbase")]
11use chromiumoxide::Browser;
12#[cfg(feature = "browserbase")]
13use futures::StreamExt;
14use serde::{Deserialize, Serialize};
15use serde_json::Value;
16#[cfg(feature = "browserbase")]
17use tokio::time::timeout;
18
19use crate::BrowserPool;
20use crate::error::BrowserError;
21use crate::freshness::{FreshnessCheckInput, FreshnessContract, FreshnessReport};
22use crate::interstitial_router::{
23    InterstitialPolicy, InterstitialRouter, PageSignature, RouterDecision,
24};
25use crate::page::WaitUntil;
26use crate::replay_defense::ReplayDefensePolicy;
27use crate::replay_defense::{ReplayDefenseCheckInput, ReplayDefenseReport, ReplayDefenseState};
28use crate::transport_realism::{
29    TransportObservation, TransportProfile, TransportRealismReport,
30    score as score_transport_realism,
31};
32
33/// Opinionated acquisition mode for the escalation ladder.
34#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
35#[serde(rename_all = "snake_case")]
36pub enum AcquisitionMode {
37    /// Prioritize lowest-latency paths.
38    Fast,
39    /// Favor reliability with broader escalation.
40    Resilient,
41    /// Start from stronger anti-bot paths.
42    Hostile,
43    /// Enter from a policy-guided start point.
44    Investigate,
45}
46
47/// Strategy stage attempted by the acquisition runner.
48#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
49#[serde(rename_all = "snake_case")]
50pub enum StrategyUsed {
51    /// Plain HTTP fetch.
52    DirectHttp,
53    /// HTTP fetch using a TLS-profiled client.
54    TlsProfiledHttp,
55    /// Browser session with opinionated light-stealth defaults.
56    BrowserLightStealth,
57    /// Browser session scoped to a sticky context id.
58    StickyProxyBrowserSession,
59    /// Managed remote browser session routed through Browserbase.
60    #[cfg(feature = "browserbase")]
61    BrowserbaseManagedSession,
62    /// Policy-guided entry marker for investigation mode.
63    InvestigateEntry,
64}
65
66/// Replay-defense context supplied to an [`AcquisitionRequest`].
67///
68/// Carries the [`ReplayDefensePolicy`] (which determines the
69/// rotation / nonce / drift levers) and the live
70/// [`ReplayDefenseState`] (the per-session record) into the runner.
71/// When the context is set, the runner evaluates the policy
72/// before any stage executes and, if the decision requires a
73/// forced refresh, calls
74/// [`BrowserPool::release_context`][crate::pool::BrowserPool::release_context]
75/// to invalidate the sticky session for the target host before
76/// short-circuiting with a structured
77/// [`StageFailureKind::Setup`] failure.
78#[derive(Debug, Clone)]
79pub struct ReplayDefenseContext {
80    /// Policy to apply to the supplied state.
81    pub policy: ReplayDefensePolicy,
82    /// Per-session record to evaluate.
83    pub state: ReplayDefenseState,
84}
85
86impl ReplayDefenseContext {
87    /// Build a context with the default policy.
88    #[must_use]
89    pub fn new(state: ReplayDefenseState) -> Self {
90        Self {
91            policy: ReplayDefensePolicy::default(),
92            state,
93        }
94    }
95
96    /// Build a context with the supplied policy and state.
97    #[must_use]
98    pub const fn with_policy(policy: ReplayDefensePolicy, state: ReplayDefenseState) -> Self {
99        Self { policy, state }
100    }
101}
102
103/// Transport-realism strategy hint supplied to an
104/// [`AcquisitionRequest`].
105///
106/// The context carries the [`TransportProfile`] (the per-target
107/// expected fingerprints, e.g. Chrome 136) and an optional
108/// [`TransportObservation`] (live capture data). When supplied, the
109/// runner evaluates the observation against the profile via
110/// [`score_transport_realism`][crate::transport_realism::score] and
111/// attaches the resulting [`TransportRealismReport`] to the
112/// [`AcquisitionResult::transport_realism`] field so downstream
113/// policy mapping (T83 / T85 / T89 / T93) can consume it as a
114/// strategy hint.
115#[derive(Debug, Clone)]
116pub struct TransportRealismContext {
117    /// Per-target transport profile the runner should score against.
118    pub profile: TransportProfile,
119    /// Optional live observation. When `None`, the score collapses
120    /// to the documented "no signal" defaults — the runner still
121    /// attaches the report so callers can detect the missing-data
122    /// path deterministically.
123    pub observation: Option<TransportObservation>,
124}
125
126impl TransportRealismContext {
127    /// Build a context with the default profile and no observation.
128    #[must_use]
129    pub const fn new(profile: TransportProfile) -> Self {
130        Self {
131            profile,
132            observation: None,
133        }
134    }
135
136    /// Build a context with the supplied profile and observation.
137    #[must_use]
138    pub const fn with_observation(
139        profile: TransportProfile,
140        observation: TransportObservation,
141    ) -> Self {
142        Self {
143            profile,
144            observation: Some(observation),
145        }
146    }
147
148    /// Replace the observation on an existing context.
149    #[must_use]
150    pub fn with_observation_opt(mut self, observation: Option<TransportObservation>) -> Self {
151        self.observation = observation;
152        self
153    }
154
155    /// Replace the profile on an existing context.
156    #[must_use]
157    pub fn with_profile(mut self, profile: TransportProfile) -> Self {
158        self.profile = profile;
159        self
160    }
161}
162
163/// Interstitial routing context supplied to an
164/// [`AcquisitionRequest`].
165///
166/// Carries the [`PageSignature`] observed on a previous
167/// attempt plus the [`InterstitialPolicy`] that controls
168/// the router's behaviour. When the context is set, the
169/// runner evaluates the signature via the
170/// [`InterstitialRouter`]
171/// **before** any stage executes:
172///
173/// 1. The resulting [`RouterDecision`] is attached to
174///    [`AcquisitionResult::interstitial`] regardless of
175///    the decision's kind.
176/// 2. When the decision is non-`Transient` **and**
177///    [`InterstitialPolicy::short_circuit_on_classified`]
178///    is `true` (the default), the runner short-circuits
179///    with a structured
180///    [`StageFailureKind::InterstitialRouted`]
181///    failure tagged with the decision so the calling
182///    layer can dispatch the dedicated
183///    [`InterstitialRoute`][crate::interstitial_router::InterstitialRoute]
184///    without burning through the generic ladder.
185///
186/// Default-on (no new feature gate). Purely additive on
187/// [`AcquisitionRequest`] and [`AcquisitionResult`].
188#[derive(Debug, Clone)]
189pub struct InterstitialContext {
190    /// Page signature observed on a previous attempt.
191    pub signature: PageSignature,
192    /// Routing policy (queue interval, challenge solve
193    /// budget, hard-block escalation, short-circuit
194    /// toggle).
195    pub policy: InterstitialPolicy,
196}
197
198impl InterstitialContext {
199    /// Build a context with the default policy.
200    #[must_use]
201    pub fn new(signature: PageSignature) -> Self {
202        Self {
203            signature,
204            policy: InterstitialPolicy::default(),
205        }
206    }
207
208    /// Build a context with the supplied policy and
209    /// signature.
210    #[must_use]
211    pub const fn with_policy(policy: InterstitialPolicy, signature: PageSignature) -> Self {
212        Self { signature, policy }
213    }
214
215    /// Replace the policy on an existing context.
216    #[must_use]
217    pub const fn with_policy_opt(mut self, policy: InterstitialPolicy) -> Self {
218        self.policy = policy;
219        self
220    }
221}
222
223/// One acquisition request.
224#[derive(Debug, Clone)]
225pub struct AcquisitionRequest {
226    /// Target URL.
227    pub url: String,
228    /// Acquisition mode.
229    pub mode: AcquisitionMode,
230    /// Optional selector that must be present for browser-stage success.
231    pub wait_for_selector: Option<String>,
232    /// Optional JavaScript extraction expression evaluated in browser stages.
233    pub extraction_js: Option<String>,
234    /// Hard wall-clock timeout for the whole acquisition attempt.
235    pub total_timeout: Duration,
236    /// Per-navigation timeout for browser stages.
237    pub navigation_timeout: Duration,
238    /// Per-request timeout for HTTP stages.
239    pub request_timeout: Duration,
240    /// Maximum HTML bytes captured into `html_excerpt`.
241    pub html_excerpt_bytes: usize,
242    /// Optional policy-guided stage that `Investigate` mode starts from.
243    pub investigate_start: Option<StrategyUsed>,
244    /// Opt into the optional Browserbase-managed stage when available.
245    pub browserbase_enabled: bool,
246    /// Optional previously-captured [`FreshnessContract`] for the
247    /// sticky identity being reused. When set, the runner evaluates
248    /// freshness against this contract before any stage executes.
249    /// If the contract is invalid (stale TTL, signature mismatch,
250    /// or domain mismatch), the runner short-circuits with a
251    /// structured rejection and the
252    /// [`AcquisitionResult::freshness`] field is populated with the
253    /// [`FreshnessReport`] describing why.
254    pub freshness_contract: Option<FreshnessContract>,
255    /// Optional [`ReplayDefenseContext`] (T81). When set, the runner
256    /// evaluates the policy against the supplied state before any
257    /// stage executes. If the decision requires a forced refresh
258    /// (rotation due, nonce expired/rotated, or signature drift
259    /// with `force_reset_on_drift = true`), the runner calls
260    /// [`BrowserPool::release_context`][crate::pool::BrowserPool::release_context]
261    /// to invalidate the sticky session for the target host and
262    /// short-circuits with a structured rejection. The full
263    /// [`ReplayDefenseReport`] is attached to
264    /// [`AcquisitionResult::replay_defense`].
265    pub replay_defense: Option<ReplayDefenseContext>,
266    /// Optional [`TransportRealismContext`] (T82) — typed
267    /// `AcquisitionRunner` strategy hint. When set, the runner
268    /// evaluates the supplied [`TransportObservation`]
269    /// against the supplied [`TransportProfile`] via the
270    /// transport-realism scorer and attaches the resulting
271    /// [`TransportRealismReport`] to
272    /// [`AcquisitionResult::transport_realism`]. The runner does
273    /// not short-circuit on low scores — strategy hints are
274    /// observed by downstream policy mapping (T83 / T85 / T89 /
275    /// T93), not enforced by the runner itself.
276    pub transport_realism: Option<TransportRealismContext>,
277    /// Optional [`InterstitialContext`] (T94) — typed
278    /// `AcquisitionRunner` failure-recovery hint. When set,
279    /// the runner classifies the supplied
280    /// [`PageSignature`]
281    /// via the [`InterstitialRouter`]
282    /// before any stage executes. The resulting
283    /// [`RouterDecision`]
284    /// is attached to [`AcquisitionResult::interstitial`].
285    /// When the decision is non-`Transient` **and** the
286    /// policy's
287    /// [`short_circuit_on_classified`][InterstitialPolicy::short_circuit_on_classified]
288    /// is `true` (the default), the runner short-circuits
289    /// with a structured
290    /// [`StageFailureKind::InterstitialRouted`] failure
291    /// so the calling layer can dispatch the dedicated
292    /// route without burning through the generic ladder.
293    pub interstitial: Option<InterstitialContext>,
294    /// Optional [`BrowserbaseSessionConfig`] (T114) tuning the
295    /// Browserbase-managed stage's session reuse, warmup, and
296    /// rate-limit retry behavior. Only consulted when
297    /// `browserbase_enabled` is `true`; `None` is equivalent to
298    /// `Some(BrowserbaseSessionConfig::default())`.
299    pub browserbase_session: Option<BrowserbaseSessionConfig>,
300}
301
302impl Default for AcquisitionRequest {
303    fn default() -> Self {
304        Self {
305            url: String::new(),
306            mode: AcquisitionMode::Resilient,
307            wait_for_selector: None,
308            extraction_js: None,
309            total_timeout: Duration::from_secs(45),
310            navigation_timeout: Duration::from_secs(30),
311            request_timeout: Duration::from_secs(15),
312            html_excerpt_bytes: 4_096,
313            investigate_start: None,
314            browserbase_enabled: false,
315            freshness_contract: None,
316            replay_defense: None,
317            transport_realism: None,
318            interstitial: None,
319            browserbase_session: None,
320        }
321    }
322}
323
324/// Session reuse, warmup, and rate-limit retry configuration for the
325/// Browserbase-managed acquisition stage (T114).
326///
327/// Closes the gap reported against real-world Browserbase usage:
328/// the stage previously minted a brand-new session on every call (no
329/// reuse), never gave the remote session a chance to settle before
330/// the real navigation (no warmup), and treated a `429` from
331/// Browserbase's own session-management API the same as any other
332/// transport failure (no backoff).
333#[derive(Debug, Clone, Serialize, Deserialize)]
334pub struct BrowserbaseSessionConfig {
335    /// Reuse an existing Browserbase session instead of creating a
336    /// new one. When `None` (or empty), the `BROWSERBASE_SESSION_ID`
337    /// environment variable is consulted next; only when both are
338    /// unset does the stage fall back to creating — and, on
339    /// release, deleting — a fresh session. A session supplied here
340    /// or via the environment variable is never deleted by the
341    /// stage; the caller owns its lifecycle.
342    pub session_id: Option<String>,
343    /// Run a warmup navigation to the target URL, and let it settle,
344    /// before the real extraction navigation. Gives anti-bot
345    /// challenge JS a chance to execute and cookies/fingerprint
346    /// state a chance to stabilize on a freshly connected session
347    /// before the page that actually gets scraped loads. Best
348    /// effort: a warmup failure is logged and does not fail the
349    /// stage.
350    pub warmup: bool,
351    /// Settle delay after the warmup navigation, in milliseconds.
352    pub warmup_stabilize_ms: u64,
353    /// Maximum retry attempts for session creation when Browserbase
354    /// responds `429` (rate limited). `0` disables retries.
355    pub max_retries: u8,
356    /// Base delay for exponential backoff between retries, in
357    /// milliseconds. Attempt `n` (0-indexed) waits
358    /// `backoff_base_ms * 2^n`, unless Browserbase's `Retry-After`
359    /// response header supplies an explicit delay.
360    pub backoff_base_ms: u64,
361}
362
363impl Default for BrowserbaseSessionConfig {
364    fn default() -> Self {
365        Self {
366            session_id: None,
367            warmup: true,
368            warmup_stabilize_ms: 500,
369            max_retries: 3,
370            backoff_base_ms: 500,
371        }
372    }
373}
374
375/// Failure class recorded per strategy stage.
376#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
377#[serde(rename_all = "snake_case")]
378pub enum StageFailureKind {
379    /// Stage initialization/setup failed.
380    Setup,
381    /// Stage hit a timeout.
382    Timeout,
383    /// Stage reached a known anti-bot block class.
384    Blocked,
385    /// Transport/runtime failure.
386    Transport,
387    /// Extraction/validation failure.
388    Extraction,
389    /// Replay-defense policy forced a refresh of the sticky session.
390    ///
391    /// Emitted by [`AcquisitionRunner::run`] when the supplied
392    /// [`ReplayDefenseContext`][crate::replay_defense::ReplayDefenseState]
393    /// decision (`RotationDue` / `NonceExpired` / `NonceRotated` /
394    /// `SignatureDrift` with `force_reset_on_drift = true`)
395    /// instructs the runner to invalidate the sticky session and
396    /// short-circuit. Callers should retry with a fresh session.
397    ReplayDefenseTriggered,
398    /// Interstitial router short-circuited the run with a
399    /// classified decision (`Queue` / `Challenge` / `HardBlock`).
400    ///
401    /// Emitted by [`AcquisitionRunner::run`] when the supplied
402    /// [`InterstitialContext`]
403    /// classifies a previously-observed
404    /// [`PageSignature`]
405    /// as a queue / challenge / hard block and the configured
406    /// [`InterstitialPolicy::short_circuit_on_classified`][crate::interstitial_router::InterstitialPolicy::short_circuit_on_classified]
407    /// is `true` (the default). The full
408    /// [`RouterDecision`]
409    /// is attached to [`AcquisitionResult::interstitial`] so
410    /// downstream tooling can dispatch the dedicated route
411    /// without burning through the generic ladder.
412    InterstitialRouted,
413    /// The remote session provider (e.g. Browserbase) rate-limited
414    /// session management requests and retries were exhausted.
415    RateLimited,
416}
417
418/// Captured failure record for one stage.
419#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
420pub struct StageFailure {
421    /// Stage where the failure happened.
422    pub strategy: StrategyUsed,
423    /// Coarse failure kind.
424    pub kind: StageFailureKind,
425    /// Compact diagnostic message.
426    pub message: String,
427}
428
429/// Terminal acquisition result.
430#[derive(Debug, Clone, Serialize, Deserialize)]
431pub struct AcquisitionResult {
432    /// `true` when any stage satisfied success criteria.
433    pub success: bool,
434    /// Stage that produced the terminal success, if any.
435    pub strategy_used: Option<StrategyUsed>,
436    /// Ordered stage attempts.
437    pub attempted: Vec<StrategyUsed>,
438    /// Final URL observed from the successful stage.
439    pub final_url: Option<String>,
440    /// HTTP status code observed from the successful stage.
441    pub status_code: Option<u16>,
442    /// Best-effort HTML excerpt from the successful stage.
443    pub html_excerpt: Option<String>,
444    /// Optional extraction payload.
445    pub extracted: Option<Value>,
446    /// Failure bundle collected across stages.
447    pub failures: Vec<StageFailure>,
448    /// `true` when the wall-clock timeout fired before completion.
449    pub timed_out: bool,
450    /// Freshness report for the contract (if any) supplied via
451    /// [`AcquisitionRequest::freshness_contract`]. `None` when no
452    /// contract was supplied. Always populated when a contract was
453    /// supplied — `Valid` if the contract held, an invalid
454    /// `FreshnessDecision` variant if it was rejected.
455    pub freshness: Option<FreshnessReport>,
456    /// Replay-defense report for the context (if any) supplied via
457    /// [`AcquisitionRequest::replay_defense`]. `None` when no
458    /// context was supplied. Always populated when a context was
459    /// supplied — `Valid` if the policy held, an invalid
460    /// [`ReplayDefenseDecision`][crate::replay_defense::ReplayDefenseDecision]
461    /// variant otherwise. When `forced_refresh = true` the runner
462    /// has already invalidated the sticky session for the target
463    /// host via
464    /// [`BrowserPool::release_context`][crate::pool::BrowserPool::release_context]
465    /// and short-circuited the run.
466    pub replay_defense: Option<ReplayDefenseReport>,
467    /// Transport-realism report for the context (if any) supplied via
468    /// [`AcquisitionRequest::transport_realism`]. `None` when no
469    /// context was supplied. Always populated when a context was
470    /// supplied — carries the per-target compatibility score,
471    /// confidence/coverage markers, and structured mismatch list.
472    /// Consumed by downstream policy mapping (T83 / T85 / T89 /
473    /// T93) as a strategy hint.
474    pub transport_realism: Option<TransportRealismReport>,
475    /// Interstitial routing decision for the context (if
476    /// any) supplied via
477    /// [`AcquisitionRequest::interstitial`]. `None` when no
478    /// context was supplied. Always populated when a
479    /// context was supplied — carries the classified
480    /// [`InterstitialKind`][crate::interstitial_router::InterstitialKind],
481    /// the dedicated
482    /// [`InterstitialSeverity`][crate::interstitial_router::InterstitialSeverity]
483    /// tier (retryable / requires-solve / terminal), the
484    /// dedicated
485    /// [`InterstitialRoute`][crate::interstitial_router::InterstitialRoute],
486    /// and the per-signature evidence. When the decision
487    /// is non-`Transient` and the policy's
488    /// `short_circuit_on_classified` is `true`, the runner
489    /// has already short-circuited the run with a
490    /// [`StageFailureKind::InterstitialRouted`] failure
491    /// and the decision is the authoritative answer.
492    pub interstitial: Option<RouterDecision>,
493}
494
495impl AcquisitionResult {
496    const fn empty() -> Self {
497        Self {
498            success: false,
499            strategy_used: None,
500            attempted: Vec::new(),
501            final_url: None,
502            status_code: None,
503            html_excerpt: None,
504            extracted: None,
505            failures: Vec::new(),
506            timed_out: false,
507            freshness: None,
508            replay_defense: None,
509            transport_realism: None,
510            interstitial: None,
511        }
512    }
513}
514
515#[derive(Debug, Clone)]
516struct StageSuccess {
517    final_url: Option<String>,
518    status_code: Option<u16>,
519    html_excerpt: Option<String>,
520    extracted: Option<Value>,
521}
522
523#[derive(Debug, Clone)]
524enum StageOutcome {
525    Marker,
526    Success(StageSuccess),
527    Failure(StageFailure),
528}
529
530/// Runner facade for opinionated acquisition.
531#[derive(Clone)]
532pub struct AcquisitionRunner {
533    pool: Arc<BrowserPool>,
534}
535
536impl AcquisitionRunner {
537    /// Create a new acquisition runner.
538    ///
539    /// # Example
540    ///
541    /// ```no_run
542    /// use stygian_browser::{AcquisitionRunner, BrowserConfig, BrowserPool};
543    ///
544    /// # async fn run() -> stygian_browser::Result<()> {
545    /// let pool = BrowserPool::new(BrowserConfig::default()).await?;
546    /// let _runner = AcquisitionRunner::new(pool);
547    /// # Ok(())
548    /// # }
549    /// ```
550    #[must_use]
551    pub const fn new(pool: Arc<BrowserPool>) -> Self {
552        Self { pool }
553    }
554
555    /// Return the deterministic stage ladder for a mode.
556    ///
557    /// Investigation mode starts at `investigate_start` when provided.
558    #[must_use]
559    pub fn strategy_ladder(
560        mode: AcquisitionMode,
561        investigate_start: Option<StrategyUsed>,
562    ) -> Vec<StrategyUsed> {
563        let mut stages = match mode {
564            AcquisitionMode::Fast => vec![
565                StrategyUsed::DirectHttp,
566                StrategyUsed::TlsProfiledHttp,
567                StrategyUsed::BrowserLightStealth,
568            ],
569            AcquisitionMode::Resilient => vec![
570                StrategyUsed::DirectHttp,
571                StrategyUsed::TlsProfiledHttp,
572                StrategyUsed::BrowserLightStealth,
573                StrategyUsed::StickyProxyBrowserSession,
574            ],
575            AcquisitionMode::Hostile => vec![
576                StrategyUsed::BrowserLightStealth,
577                StrategyUsed::StickyProxyBrowserSession,
578                StrategyUsed::TlsProfiledHttp,
579                StrategyUsed::DirectHttp,
580            ],
581            AcquisitionMode::Investigate => {
582                let start = investigate_start.unwrap_or(StrategyUsed::BrowserLightStealth);
583                vec![
584                    StrategyUsed::InvestigateEntry,
585                    start,
586                    StrategyUsed::StickyProxyBrowserSession,
587                    StrategyUsed::TlsProfiledHttp,
588                ]
589            }
590        };
591
592        dedupe_preserve_order(&mut stages);
593        stages
594    }
595
596    /// Execute the acquisition ladder and return a terminal result.
597    ///
598    /// This method never panics and always returns an [`AcquisitionResult`],
599    /// including timeout and setup-failure paths.
600    ///
601    /// # Example
602    ///
603    /// ```no_run
604    /// use stygian_browser::{AcquisitionMode, AcquisitionRequest, AcquisitionRunner, BrowserConfig, BrowserPool};
605    ///
606    /// # async fn run() -> stygian_browser::Result<()> {
607    /// let pool = BrowserPool::new(BrowserConfig::default()).await?;
608    /// let runner = AcquisitionRunner::new(pool);
609    /// let request = AcquisitionRequest {
610    ///     url: "https://example.com".to_string(),
611    ///     mode: AcquisitionMode::Resilient,
612    ///     ..AcquisitionRequest::default()
613    /// };
614    /// let _result = runner.run(request).await;
615    /// # Ok(())
616    /// # }
617    /// ```
618    pub async fn run(&self, request: AcquisitionRequest) -> AcquisitionResult {
619        let timeout = request.total_timeout;
620        let timeout_strategy = Self::strategy_ladder(request.mode, request.investigate_start)
621            .into_iter()
622            .find(|strategy| *strategy != StrategyUsed::InvestigateEntry)
623            .unwrap_or(StrategyUsed::DirectHttp);
624        let mut result = tokio::time::timeout(timeout, self.run_inner(&request))
625            .await
626            .unwrap_or_else(|_| {
627                let mut timed_out = AcquisitionResult::empty();
628                timed_out.timed_out = true;
629                timed_out.failures.push(StageFailure {
630                    strategy: timeout_strategy,
631                    kind: StageFailureKind::Timeout,
632                    message: format!("acquisition timed out after {}ms", timeout.as_millis()),
633                });
634                timed_out
635            });
636
637        if !result.success {
638            // Guarantee deterministic terminal output for all unsuccessful runs.
639            if result.failures.is_empty() {
640                result.failures.push(StageFailure {
641                    strategy: timeout_strategy,
642                    kind: StageFailureKind::Transport,
643                    message: "acquisition ended without stage output".to_string(),
644                });
645            }
646        }
647
648        result
649    }
650
651    /// Evaluate the supplied `replay_defense` context against the
652    /// request URL, attach the [`ReplayDefenseReport`] to `result`,
653    /// and — when the decision mandates a forced refresh — release
654    /// the sticky pool slots for the target host and push a
655    /// structured [`StageFailureKind::ReplayDefenseTriggered`]
656    /// failure onto `result.failures`. Returns `true` when the
657    /// runner should short-circuit.
658    async fn evaluate_replay_defense(
659        &self,
660        request: &AcquisitionRequest,
661        result: &mut AcquisitionResult,
662    ) -> bool {
663        let Some(context) = request.replay_defense.as_ref() else {
664            return false;
665        };
666        let observed_host = host_hint(&request.url).unwrap_or_else(|| context.state.domain.clone());
667        let observed_signature = context.state.signature.clone();
668        let observed_nonce = context.state.nonce.clone();
669        let input = ReplayDefenseCheckInput::new(
670            &observed_host,
671            observed_signature.as_deref(),
672            observed_nonce.as_deref(),
673            crate::replay_defense::unix_epoch_ms(),
674        );
675        let report = ReplayDefenseReport::evaluate(&context.policy, &context.state, &input);
676        report.log();
677        let forced_refresh = report.forced_refresh;
678        result.replay_defense = Some(report);
679        if !forced_refresh {
680            return false;
681        }
682        let decision_label = result
683            .replay_defense
684            .as_ref()
685            .map_or("replay_defense", |r| r.decision.label());
686        let reason = result.replay_defense.as_ref().and_then(|r| r.decision.reason()).map_or_else(
687            || "replay defense forced refresh".to_string(),
688            |r| {
689                format!(
690                    "replay defense forced refresh ({reason}, contract_domain={cd}, observed_domain={od}, elapsed_ms={e})",
691                    reason = r.kind,
692                    cd = r.contract_domain,
693                    od = r.observed_domain,
694                    e = r.elapsed_ms,
695                )
696            },
697        );
698        // Invalidate the sticky session for the observed host so
699        // the next acquisition starts from a clean pool slot.
700        let released = self.pool.release_context(&observed_host).await;
701        tracing::info!(
702            target: "stygian::replay_defense",
703            host = %observed_host,
704            released_idle_browsers = released,
705            decision = decision_label,
706            "replay defense forced refresh released sticky pool slots",
707        );
708        result.failures.push(StageFailure {
709            strategy: StrategyUsed::InvestigateEntry,
710            kind: StageFailureKind::ReplayDefenseTriggered,
711            message: reason,
712        });
713        true
714    }
715
716    /// Evaluate the supplied `interstitial` context, attach
717    /// the resulting [`RouterDecision`] to `result`, and —
718    /// when the decision is classified (non-`Transient`) and
719    /// the policy mandates a short-circuit — push a
720    /// structured [`StageFailureKind::InterstitialRouted`]
721    /// failure onto `result.failures`. Returns `true` when
722    /// the runner should short-circuit.
723    fn evaluate_interstitial(request: &AcquisitionRequest, result: &mut AcquisitionResult) -> bool {
724        let Some(context) = request.interstitial.as_ref() else {
725            return false;
726        };
727        let router = InterstitialRouter::new(context.policy.clone());
728        let decision = router.classify_and_route(&context.signature);
729        decision.log();
730        let should_short_circuit = router.should_short_circuit(decision.kind());
731        result.interstitial = Some(decision);
732        if !should_short_circuit {
733            return false;
734        }
735        let kind_label = result
736            .interstitial
737            .as_ref()
738            .map_or("interstitial", |d| d.kind().label());
739        let severity_label = result
740            .interstitial
741            .as_ref()
742            .map_or("terminal", |d| d.severity().label());
743        let reason = result
744            .interstitial
745            .as_ref()
746            .map_or_else(
747                || "interstitial routed".to_string(),
748                |d| {
749                    format!(
750                        "interstitial routed ({kind}, severity={sev}, host={host}, status_code={status:?}, route={route})",
751                        kind = d.kind().label(),
752                        sev = d.severity().label(),
753                        host = d.evidence().host.as_deref().unwrap_or(""),
754                        status = d.evidence().status_code,
755                        route = d.route().label(),
756                    )
757                },
758            );
759        result.failures.push(StageFailure {
760            strategy: StrategyUsed::InvestigateEntry,
761            kind: StageFailureKind::InterstitialRouted,
762            message: reason,
763        });
764        tracing::info!(
765            target: "stygian::interstitial_router",
766            kind = kind_label,
767            severity = severity_label,
768            "interstitial routing short-circuited the runner",
769        );
770        true
771    }
772
773    async fn run_inner(&self, request: &AcquisitionRequest) -> AcquisitionResult {
774        let mut result = AcquisitionResult::empty();
775
776        // Freshness short-circuit: when a contract is supplied with the
777        // request, evaluate it against the request URL before any stage
778        // executes. An invalid contract is a deterministic, structured
779        // rejection — no I/O is performed and the runner returns early.
780        if let Some(contract) = request.freshness_contract.as_ref() {
781            let observed_host = host_hint(&request.url).unwrap_or_else(|| contract.domain.clone());
782            let observed_signature: Option<String> = None;
783            let input = FreshnessCheckInput::new(
784                &observed_host,
785                observed_signature.as_deref(),
786                crate::freshness::unix_epoch_ms(),
787            );
788            let report = FreshnessReport::evaluate(contract, &input);
789            report.log();
790            let rejected = report.decision.is_invalid();
791            result.freshness = Some(report);
792            if rejected {
793                let reason = result
794                    .freshness
795                    .as_ref()
796                    .and_then(|r| r.decision.reason())
797                    .map_or_else(
798                        || "freshness contract invalidated".to_string(),
799                        |r| {
800                            format!(
801                                "freshness contract invalidated ({reason}, contract_domain={cd}, observed_domain={od}, elapsed_ms={e}, max_age_ms={m})",
802                                reason = r.kind,
803                                cd = r.contract_domain,
804                                od = r.observed_domain,
805                                e = r.elapsed_ms,
806                                m = r.max_age_ms,
807                            )
808                        },
809                    );
810                result.failures.push(StageFailure {
811                    strategy: StrategyUsed::InvestigateEntry,
812                    kind: StageFailureKind::Setup,
813                    message: reason,
814                });
815                return result;
816            }
817        }
818
819        // Replay-defense short-circuit (T81): when a context is supplied,
820        // evaluate the policy against the request URL + the supplied state.
821        // A decision that mandates a forced refresh invalidates the sticky
822        // session via `BrowserPool::release_context` and short-circuits the
823        // run with a structured `ReplayDefenseTriggered` failure.
824        if self.evaluate_replay_defense(request, &mut result).await {
825            return result;
826        }
827
828        // Interstitial routing short-circuit (T94): when a context is
829        // supplied, classify the previously-observed page signature via
830        // the `InterstitialRouter` and attach the resulting
831        // `RouterDecision` to the result. A classified (non-`Transient`)
832        // decision with the policy's `short_circuit_on_classified` flag
833        // enabled short-circuits the run with a structured
834        // `InterstitialRouted` failure so the calling layer can dispatch
835        // the dedicated route (queue wait / challenge solve / hard-block
836        // escalation) without burning through the generic ladder.
837        if Self::evaluate_interstitial(request, &mut result) {
838            return result;
839        }
840
841        // Transport-realism strategy hint (T82): when a context is supplied,
842        // score the observation against the per-target profile and attach
843        // the resulting `TransportRealismReport` to the result. The runner
844        // never short-circuits on low scores — strategy hints are observed
845        // by downstream policy mapping (T83 / T85 / T89 / T93), not
846        // enforced by the runner itself.
847        if let Some(context) = request.transport_realism.as_ref() {
848            let observation = context.observation.clone().unwrap_or_default();
849            let report = score_transport_realism(&context.profile, &observation);
850            tracing::debug!(
851                target: "stygian::transport_realism",
852                profile = %report.profile_name,
853                score = report.compatibility.score,
854                confidence = report.compatibility.confidence,
855                coverage = report.compatibility.coverage,
856                matched = report.compatibility.matched_count,
857                total = report.compatibility.total_checks,
858                mismatches = report.compatibility.mismatches.len(),
859                "transport realism scored",
860            );
861            result.transport_realism = Some(report);
862        }
863
864        #[cfg(feature = "browserbase")]
865        let mut ladder = Self::strategy_ladder(request.mode, request.investigate_start);
866
867        #[cfg(not(feature = "browserbase"))]
868        let ladder = Self::strategy_ladder(request.mode, request.investigate_start);
869
870        #[cfg(feature = "browserbase")]
871        {
872            maybe_insert_browserbase_stage(&mut ladder, request.browserbase_enabled);
873        }
874        let started = Instant::now();
875
876        for strategy in ladder {
877            if started.elapsed() >= request.total_timeout {
878                result.timed_out = true;
879                result.failures.push(StageFailure {
880                    strategy,
881                    kind: StageFailureKind::Timeout,
882                    message: "wall-clock timeout reached before stage execution".to_string(),
883                });
884                break;
885            }
886
887            result.attempted.push(strategy);
888            match self.execute_stage(strategy, request).await {
889                StageOutcome::Marker => {}
890                StageOutcome::Success(success) => {
891                    result.success = true;
892                    result.strategy_used = Some(strategy);
893                    result.final_url = success.final_url;
894                    result.status_code = success.status_code;
895                    result.html_excerpt = success.html_excerpt;
896                    result.extracted = success.extracted;
897                    break;
898                }
899                StageOutcome::Failure(failure) => result.failures.push(failure),
900            }
901        }
902
903        result
904    }
905
906    async fn execute_stage(
907        &self,
908        strategy: StrategyUsed,
909        request: &AcquisitionRequest,
910    ) -> StageOutcome {
911        match strategy {
912            StrategyUsed::DirectHttp => {
913                #[cfg(feature = "tls-config")]
914                {
915                    self.run_http_stage(request, false).await
916                }
917
918                #[cfg(not(feature = "tls-config"))]
919                {
920                    self.run_http_stage(request, false)
921                }
922            }
923            StrategyUsed::TlsProfiledHttp => {
924                #[cfg(feature = "tls-config")]
925                {
926                    self.run_http_stage(request, true).await
927                }
928
929                #[cfg(not(feature = "tls-config"))]
930                {
931                    self.run_http_stage(request, true)
932                }
933            }
934            StrategyUsed::BrowserLightStealth => self.run_browser_stage(request, false).await,
935            StrategyUsed::StickyProxyBrowserSession => self.run_browser_stage(request, true).await,
936            #[cfg(feature = "browserbase")]
937            StrategyUsed::BrowserbaseManagedSession => Self::run_browserbase_stage(request).await,
938            StrategyUsed::InvestigateEntry => StageOutcome::Marker,
939        }
940    }
941
942    #[cfg(feature = "browserbase")]
943    #[allow(clippy::too_many_lines)]
944    async fn run_browserbase_stage(request: &AcquisitionRequest) -> StageOutcome {
945        if !request.browserbase_enabled {
946            return StageOutcome::Failure(StageFailure {
947                strategy: StrategyUsed::BrowserbaseManagedSession,
948                kind: StageFailureKind::Setup,
949                message: "browserbase stage disabled for this request".to_string(),
950            });
951        }
952
953        let api_key = match std::env::var("BROWSERBASE_API_KEY") {
954            Ok(value) if !value.trim().is_empty() => value,
955            _ => {
956                return StageOutcome::Failure(StageFailure {
957                    strategy: StrategyUsed::BrowserbaseManagedSession,
958                    kind: StageFailureKind::Setup,
959                    message: "browserbase requires BROWSERBASE_API_KEY".to_string(),
960                });
961            }
962        };
963
964        let project_id = match std::env::var("BROWSERBASE_PROJECT_ID") {
965            Ok(value) if !value.trim().is_empty() => value,
966            _ => {
967                return StageOutcome::Failure(StageFailure {
968                    strategy: StrategyUsed::BrowserbaseManagedSession,
969                    kind: StageFailureKind::Setup,
970                    message: "browserbase requires BROWSERBASE_PROJECT_ID".to_string(),
971                });
972            }
973        };
974
975        let config = request.browserbase_session.clone().unwrap_or_default();
976        let env_session_id = std::env::var("BROWSERBASE_SESSION_ID").ok();
977        let reused_session_id =
978            resolve_browserbase_session_id(config.session_id.as_deref(), env_session_id.as_deref());
979
980        let (session, owns_session) = if let Some(session_id) = reused_session_id {
981            match retrieve_browserbase_session(request, &api_key, &session_id).await {
982                Ok(session) => (session, false),
983                Err(err) => {
984                    return StageOutcome::Failure(StageFailure {
985                        strategy: StrategyUsed::BrowserbaseManagedSession,
986                        kind: classify_browser_error(&err),
987                        message: format!("browserbase session reuse failed: {err}"),
988                    });
989                }
990            }
991        } else {
992            match create_browserbase_session_with_retry(request, &api_key, &project_id, &config)
993                .await
994            {
995                Ok(session) => (session, true),
996                Err(err) => {
997                    return StageOutcome::Failure(StageFailure {
998                        strategy: StrategyUsed::BrowserbaseManagedSession,
999                        kind: classify_browser_error(&err),
1000                        message: err.to_string(),
1001                    });
1002                }
1003            }
1004        };
1005
1006        let connect_timeout = request.request_timeout.min(request.total_timeout);
1007        let (mut browser, mut handler) = match timeout(
1008            connect_timeout,
1009            Browser::connect(session.connect_url.clone()),
1010        )
1011        .await
1012        {
1013            Ok(Ok(pair)) => pair,
1014            Ok(Err(err)) => {
1015                if owns_session {
1016                    let _ = delete_browserbase_session(request, &api_key, &session.id).await;
1017                }
1018                return StageOutcome::Failure(StageFailure {
1019                    strategy: StrategyUsed::BrowserbaseManagedSession,
1020                    kind: StageFailureKind::Transport,
1021                    message: format!("browserbase connect failed: {err}"),
1022                });
1023            }
1024            Err(_) => {
1025                if owns_session {
1026                    let _ = delete_browserbase_session(request, &api_key, &session.id).await;
1027                }
1028                return StageOutcome::Failure(StageFailure {
1029                    strategy: StrategyUsed::BrowserbaseManagedSession,
1030                    kind: StageFailureKind::Timeout,
1031                    message: format!(
1032                        "browserbase connect timed out after {}ms",
1033                        connect_timeout.as_millis()
1034                    ),
1035                });
1036            }
1037        };
1038
1039        let handler_task = tokio::spawn(async move {
1040            while let Some(event) = handler.next().await {
1041                if let Err(error) = event {
1042                    tracing::warn!(%error, "browserbase handler error");
1043                    break;
1044                }
1045            }
1046        });
1047
1048        let run_result = async {
1049            let raw_page =
1050                browser
1051                    .new_page("about:blank")
1052                    .await
1053                    .map_err(|err| BrowserError::CdpError {
1054                        operation: "Browser.newPage".to_string(),
1055                        message: err.to_string(),
1056                    })?;
1057
1058            let mut page = crate::page::PageHandle::new(raw_page, request.navigation_timeout);
1059
1060            if config.warmup {
1061                let warmup_timeout_ms =
1062                    u64::try_from(request.navigation_timeout.as_millis()).unwrap_or(u64::MAX);
1063                let warmup_result = page
1064                    .warmup(crate::page::WarmupOptions {
1065                        url: request.url.clone(),
1066                        wait: crate::page::WarmupWait::DomContentLoaded,
1067                        timeout_ms: warmup_timeout_ms,
1068                        stabilize_ms: config.warmup_stabilize_ms,
1069                    })
1070                    .await;
1071                if let Err(error) = warmup_result {
1072                    tracing::warn!(
1073                        %error,
1074                        "browserbase warmup navigation failed; continuing with primary navigation"
1075                    );
1076                }
1077            }
1078
1079            page.navigate(
1080                &request.url,
1081                WaitUntil::DomContentLoaded,
1082                request.navigation_timeout,
1083            )
1084            .await?;
1085
1086            if let Some(selector) = &request.wait_for_selector {
1087                page.wait_for_selector(selector, request.navigation_timeout)
1088                    .await?;
1089            }
1090
1091            let extracted = match request.extraction_js.as_deref() {
1092                Some(script) => Some(page.eval::<Value>(script).await.map_err(|err| {
1093                    BrowserError::ScriptExecutionFailed {
1094                        script: script.to_string(),
1095                        reason: err.to_string(),
1096                    }
1097                })?),
1098                None => None,
1099            };
1100
1101            let html = page.content().await?;
1102            let final_url = page.url().await.ok();
1103            let status_code = page.status_code().ok().flatten();
1104
1105            Ok::<StageSuccess, BrowserError>(StageSuccess {
1106                final_url,
1107                status_code,
1108                html_excerpt: Some(truncate_html(&html, request.html_excerpt_bytes)),
1109                extracted,
1110            })
1111        }
1112        .await;
1113
1114        let _ = timeout(Duration::from_secs(5), browser.close()).await;
1115        handler_task.abort();
1116        if owns_session {
1117            let _ = delete_browserbase_session(request, &api_key, &session.id).await;
1118        }
1119
1120        match run_result {
1121            Ok(success) => {
1122                if is_block_status(success.status_code) {
1123                    StageOutcome::Failure(StageFailure {
1124                        strategy: StrategyUsed::BrowserbaseManagedSession,
1125                        kind: StageFailureKind::Blocked,
1126                        message: format!(
1127                            "blocked status during browserbase stage: {:?}",
1128                            success.status_code
1129                        ),
1130                    })
1131                } else {
1132                    StageOutcome::Success(success)
1133                }
1134            }
1135            Err(err) => StageOutcome::Failure(StageFailure {
1136                strategy: StrategyUsed::BrowserbaseManagedSession,
1137                kind: classify_browser_error(&err),
1138                message: err.to_string(),
1139            }),
1140        }
1141    }
1142
1143    async fn run_browser_stage(&self, request: &AcquisitionRequest, sticky: bool) -> StageOutcome {
1144        let strategy = if sticky {
1145            StrategyUsed::StickyProxyBrowserSession
1146        } else {
1147            StrategyUsed::BrowserLightStealth
1148        };
1149
1150        let handle_result = if sticky {
1151            let context = host_hint(&request.url).unwrap_or_else(|| "default".to_string());
1152            self.pool.acquire_for(&context).await
1153        } else {
1154            self.pool.acquire().await
1155        };
1156
1157        let handle = match handle_result {
1158            Ok(handle) => handle,
1159            Err(err) => {
1160                return StageOutcome::Failure(StageFailure {
1161                    strategy,
1162                    kind: StageFailureKind::Setup,
1163                    message: format!("browser acquire failed: {err}"),
1164                });
1165            }
1166        };
1167
1168        let page_result = async {
1169            let browser = handle.browser().ok_or_else(|| {
1170                BrowserError::ConfigError("browser handle already released".to_string())
1171            })?;
1172            let mut page = browser.new_page().await?;
1173            page.navigate(
1174                &request.url,
1175                WaitUntil::DomContentLoaded,
1176                request.navigation_timeout,
1177            )
1178            .await?;
1179
1180            if let Some(selector) = &request.wait_for_selector {
1181                page.wait_for_selector(selector, request.navigation_timeout)
1182                    .await?;
1183            }
1184
1185            let extracted = match request.extraction_js.as_deref() {
1186                Some(script) => Some(page.eval::<Value>(script).await.map_err(|err| {
1187                    BrowserError::ScriptExecutionFailed {
1188                        script: script.to_string(),
1189                        reason: err.to_string(),
1190                    }
1191                })?),
1192                None => None,
1193            };
1194
1195            let html = page.content().await?;
1196            let final_url = page.url().await.ok();
1197            let status_code = page.status_code().ok().flatten();
1198            let html_excerpt = truncate_html(&html, request.html_excerpt_bytes);
1199
1200            drop(page);
1201
1202            Ok::<StageSuccess, BrowserError>(StageSuccess {
1203                final_url,
1204                status_code,
1205                html_excerpt: Some(html_excerpt),
1206                extracted,
1207            })
1208        }
1209        .await;
1210
1211        handle.release().await;
1212
1213        match page_result {
1214            Ok(success) => {
1215                if is_block_status(success.status_code) {
1216                    StageOutcome::Failure(StageFailure {
1217                        strategy,
1218                        kind: StageFailureKind::Blocked,
1219                        message: format!(
1220                            "blocked status during browser stage: {:?}",
1221                            success.status_code
1222                        ),
1223                    })
1224                } else {
1225                    StageOutcome::Success(success)
1226                }
1227            }
1228            Err(err) => StageOutcome::Failure(StageFailure {
1229                strategy,
1230                kind: classify_browser_error(&err),
1231                message: err.to_string(),
1232            }),
1233        }
1234    }
1235
1236    #[cfg(feature = "tls-config")]
1237    async fn run_http_stage(
1238        &self,
1239        request: &AcquisitionRequest,
1240        tls_profiled: bool,
1241    ) -> StageOutcome {
1242        if request.wait_for_selector.is_some() || request.extraction_js.is_some() {
1243            return StageOutcome::Failure(StageFailure {
1244                strategy: if tls_profiled {
1245                    StrategyUsed::TlsProfiledHttp
1246                } else {
1247                    StrategyUsed::DirectHttp
1248                },
1249                kind: StageFailureKind::Extraction,
1250                message: "HTTP stages cannot satisfy selector/extraction requirements".to_string(),
1251            });
1252        }
1253
1254        self.run_http_stage_impl(request, tls_profiled).await
1255    }
1256
1257    #[cfg(not(feature = "tls-config"))]
1258    fn run_http_stage(&self, request: &AcquisitionRequest, tls_profiled: bool) -> StageOutcome {
1259        if request.wait_for_selector.is_some() || request.extraction_js.is_some() {
1260            return StageOutcome::Failure(StageFailure {
1261                strategy: if tls_profiled {
1262                    StrategyUsed::TlsProfiledHttp
1263                } else {
1264                    StrategyUsed::DirectHttp
1265                },
1266                kind: StageFailureKind::Extraction,
1267                message: "HTTP stages cannot satisfy selector/extraction requirements".to_string(),
1268            });
1269        }
1270
1271        self.run_http_stage_impl(request, tls_profiled)
1272    }
1273
1274    #[cfg(feature = "tls-config")]
1275    async fn run_http_stage_impl(
1276        &self,
1277        request: &AcquisitionRequest,
1278        tls_profiled: bool,
1279    ) -> StageOutcome {
1280        use crate::tls::{CHROME_131, build_profiled_client_preset};
1281
1282        let strategy = if tls_profiled {
1283            StrategyUsed::TlsProfiledHttp
1284        } else {
1285            StrategyUsed::DirectHttp
1286        };
1287
1288        let client = if tls_profiled {
1289            match build_profiled_client_preset(&CHROME_131, None) {
1290                Ok(client) => client,
1291                Err(err) => {
1292                    return StageOutcome::Failure(StageFailure {
1293                        strategy,
1294                        kind: StageFailureKind::Setup,
1295                        message: format!("tls-profiled client setup failed: {err}"),
1296                    });
1297                }
1298            }
1299        } else {
1300            match reqwest::Client::builder()
1301                .timeout(request.request_timeout)
1302                .cookie_store(true)
1303                .build()
1304            {
1305                Ok(client) => client,
1306                Err(err) => {
1307                    return StageOutcome::Failure(StageFailure {
1308                        strategy,
1309                        kind: StageFailureKind::Setup,
1310                        message: format!("http client setup failed: {err}"),
1311                    });
1312                }
1313            }
1314        };
1315
1316        let response = match client
1317            .get(&request.url)
1318            .timeout(request.request_timeout)
1319            .send()
1320            .await
1321        {
1322            Ok(response) => response,
1323            Err(err) => {
1324                return StageOutcome::Failure(StageFailure {
1325                    strategy,
1326                    kind: if err.is_timeout() {
1327                        StageFailureKind::Timeout
1328                    } else {
1329                        StageFailureKind::Transport
1330                    },
1331                    message: err.to_string(),
1332                });
1333            }
1334        };
1335
1336        let status_code = Some(response.status().as_u16());
1337        let final_url = Some(response.url().to_string());
1338        let html = match response.text().await {
1339            Ok(text) => text,
1340            Err(err) => {
1341                return StageOutcome::Failure(StageFailure {
1342                    strategy,
1343                    kind: StageFailureKind::Transport,
1344                    message: format!("response body read failed: {err}"),
1345                });
1346            }
1347        };
1348
1349        if is_block_status(status_code) {
1350            return StageOutcome::Failure(StageFailure {
1351                strategy,
1352                kind: StageFailureKind::Blocked,
1353                message: format!("blocked status from HTTP stage: {status_code:?}"),
1354            });
1355        }
1356
1357        StageOutcome::Success(StageSuccess {
1358            final_url,
1359            status_code,
1360            html_excerpt: Some(truncate_html(&html, request.html_excerpt_bytes)),
1361            extracted: None,
1362        })
1363    }
1364
1365    #[cfg(not(feature = "tls-config"))]
1366    #[expect(
1367        clippy::unused_self,
1368        reason = "signature must match the tls-config variant for uniform call sites"
1369    )]
1370    fn run_http_stage_impl(
1371        &self,
1372        _request: &AcquisitionRequest,
1373        tls_profiled: bool,
1374    ) -> StageOutcome {
1375        let strategy = if tls_profiled {
1376            StrategyUsed::TlsProfiledHttp
1377        } else {
1378            StrategyUsed::DirectHttp
1379        };
1380        StageOutcome::Failure(StageFailure {
1381            strategy,
1382            kind: StageFailureKind::Setup,
1383            message: "HTTP acquisition requires the `tls-config` feature".to_string(),
1384        })
1385    }
1386}
1387
1388#[cfg(feature = "browserbase")]
1389#[derive(Debug, Clone)]
1390struct BrowserbaseSession {
1391    id: String,
1392    connect_url: String,
1393}
1394
1395/// Shared non-success handling for Browserbase session-management
1396/// responses: classifies `429` into [`BrowserError::RateLimited`]
1397/// (carrying the `Retry-After` delay when present) ahead of the
1398/// generic non-2xx path, then parses the JSON body.
1399#[cfg(feature = "browserbase")]
1400async fn browserbase_response_payload(
1401    response: reqwest::Response,
1402    url: &str,
1403    action: &str,
1404) -> Result<Value, BrowserError> {
1405    if response.status() == reqwest::StatusCode::TOO_MANY_REQUESTS {
1406        return Err(BrowserError::RateLimited {
1407            retry_after_ms: parse_retry_after_ms(response.headers()),
1408        });
1409    }
1410
1411    if !response.status().is_success() {
1412        let status = response.status();
1413        let body = response.text().await.unwrap_or_default();
1414        return Err(BrowserError::ConnectionError {
1415            url: url.to_string(),
1416            reason: format!("session {action} failed ({status}): {body}"),
1417        });
1418    }
1419
1420    response
1421        .json()
1422        .await
1423        .map_err(|err| BrowserError::ConnectionError {
1424            url: url.to_string(),
1425            reason: format!("session {action} response parse failed: {err}"),
1426        })
1427}
1428
1429/// Parses an integer-seconds `Retry-After` header into milliseconds.
1430///
1431/// Only the delay-seconds form is supported (the HTTP-date form is
1432/// rare on rate-limit responses); callers fall back to their own
1433/// backoff schedule when this returns `None`.
1434#[cfg(feature = "browserbase")]
1435fn parse_retry_after_ms(headers: &reqwest::header::HeaderMap) -> Option<u64> {
1436    headers
1437        .get(reqwest::header::RETRY_AFTER)
1438        .and_then(|value| value.to_str().ok())
1439        .and_then(|value| value.trim().parse::<u64>().ok())
1440        .map(|secs| secs.saturating_mul(1000))
1441}
1442
1443/// Resolves which Browserbase session id (if any) should be reused
1444/// instead of minting a fresh session, preferring an explicitly
1445/// configured id over the `BROWSERBASE_SESSION_ID` environment
1446/// variable. Blank values are treated as unset.
1447#[cfg(feature = "browserbase")]
1448fn resolve_browserbase_session_id(
1449    configured: Option<&str>,
1450    env_value: Option<&str>,
1451) -> Option<String> {
1452    configured
1453        .filter(|value| !value.trim().is_empty())
1454        .or_else(|| env_value.filter(|value| !value.trim().is_empty()))
1455        .map(ToString::to_string)
1456}
1457
1458/// Exponential backoff delay for retry attempt `attempt` (0-indexed):
1459/// `base_ms * 2^attempt`, saturating rather than overflowing.
1460#[cfg(feature = "browserbase")]
1461const fn backoff_delay_ms(base_ms: u64, attempt: u8) -> u64 {
1462    let shift = if attempt as u32 > 62 {
1463        62
1464    } else {
1465        attempt as u32
1466    };
1467    base_ms.saturating_mul(1u64 << shift)
1468}
1469
1470#[cfg(feature = "browserbase")]
1471async fn create_browserbase_session(
1472    request: &AcquisitionRequest,
1473    api_key: &str,
1474    project_id: &str,
1475) -> Result<BrowserbaseSession, BrowserError> {
1476    let client = reqwest::Client::builder()
1477        .timeout(request.request_timeout)
1478        .build()
1479        .map_err(|err| {
1480            BrowserError::ConfigError(format!("browserbase client setup failed: {err}"))
1481        })?;
1482
1483    let create_url = format!("{}/sessions", browserbase_api_base());
1484    let response = client
1485        .post(create_url.clone())
1486        .bearer_auth(api_key)
1487        .header("x-bb-api-key", api_key)
1488        .json(&serde_json::json!({ "projectId": project_id }))
1489        .send()
1490        .await
1491        .map_err(|err| BrowserError::ConnectionError {
1492            url: create_url.clone(),
1493            reason: err.to_string(),
1494        })?;
1495
1496    let payload = browserbase_response_payload(response, &create_url, "create").await?;
1497
1498    let connect_url = browserbase_connect_url(&payload).ok_or_else(|| {
1499        BrowserError::ConfigError("browserbase response missing connect URL".to_string())
1500    })?;
1501    let session_id = browserbase_session_id(&payload).ok_or_else(|| {
1502        BrowserError::ConfigError("browserbase response missing session id".to_string())
1503    })?;
1504
1505    Ok(BrowserbaseSession {
1506        id: session_id,
1507        connect_url,
1508    })
1509}
1510
1511/// Retries [`create_browserbase_session`] with exponential backoff
1512/// when Browserbase responds `429`, honoring its `Retry-After`
1513/// header when present. Any other error is returned immediately.
1514#[cfg(feature = "browserbase")]
1515async fn create_browserbase_session_with_retry(
1516    request: &AcquisitionRequest,
1517    api_key: &str,
1518    project_id: &str,
1519    config: &BrowserbaseSessionConfig,
1520) -> Result<BrowserbaseSession, BrowserError> {
1521    let mut attempt = 0u8;
1522    loop {
1523        match create_browserbase_session(request, api_key, project_id).await {
1524            Ok(session) => return Ok(session),
1525            Err(BrowserError::RateLimited { retry_after_ms }) if attempt < config.max_retries => {
1526                let delay_ms = retry_after_ms
1527                    .unwrap_or_else(|| backoff_delay_ms(config.backoff_base_ms, attempt));
1528                tokio::time::sleep(Duration::from_millis(delay_ms)).await;
1529                attempt += 1;
1530            }
1531            Err(err) => return Err(err),
1532        }
1533    }
1534}
1535
1536/// Fetches connection details for an existing Browserbase session by
1537/// id, for the session-reuse path — no new session is created.
1538#[cfg(feature = "browserbase")]
1539async fn retrieve_browserbase_session(
1540    request: &AcquisitionRequest,
1541    api_key: &str,
1542    session_id: &str,
1543) -> Result<BrowserbaseSession, BrowserError> {
1544    let client = reqwest::Client::builder()
1545        .timeout(request.request_timeout)
1546        .build()
1547        .map_err(|err| {
1548            BrowserError::ConfigError(format!("browserbase client setup failed: {err}"))
1549        })?;
1550
1551    let get_url = format!("{}/sessions/{session_id}", browserbase_api_base());
1552    let response = client
1553        .get(get_url.clone())
1554        .bearer_auth(api_key)
1555        .header("x-bb-api-key", api_key)
1556        .send()
1557        .await
1558        .map_err(|err| BrowserError::ConnectionError {
1559            url: get_url.clone(),
1560            reason: err.to_string(),
1561        })?;
1562
1563    let payload = browserbase_response_payload(response, &get_url, "retrieve").await?;
1564
1565    let connect_url = browserbase_connect_url(&payload).ok_or_else(|| {
1566        BrowserError::ConfigError("browserbase retrieve response missing connect URL".to_string())
1567    })?;
1568
1569    Ok(BrowserbaseSession {
1570        id: session_id.to_string(),
1571        connect_url,
1572    })
1573}
1574
1575#[cfg(feature = "browserbase")]
1576async fn delete_browserbase_session(
1577    request: &AcquisitionRequest,
1578    api_key: &str,
1579    session_id: &str,
1580) -> Result<(), BrowserError> {
1581    let client = reqwest::Client::builder()
1582        .timeout(request.request_timeout)
1583        .build()
1584        .map_err(|err| {
1585            BrowserError::ConfigError(format!("browserbase client setup failed: {err}"))
1586        })?;
1587
1588    let delete_url = format!("{}/sessions/{session_id}", browserbase_api_base());
1589    let response = client
1590        .delete(delete_url.clone())
1591        .bearer_auth(api_key)
1592        .header("x-bb-api-key", api_key)
1593        .send()
1594        .await
1595        .map_err(|err| BrowserError::ConnectionError {
1596            url: delete_url.clone(),
1597            reason: err.to_string(),
1598        })?;
1599
1600    if response.status().is_success() {
1601        Ok(())
1602    } else {
1603        Err(BrowserError::ConnectionError {
1604            url: delete_url,
1605            reason: format!("session delete failed with status {}", response.status()),
1606        })
1607    }
1608}
1609
1610#[cfg(feature = "browserbase")]
1611fn browserbase_api_base() -> String {
1612    std::env::var("BROWSERBASE_API_BASE")
1613        .unwrap_or_else(|_| "https://api.browserbase.com/v1".to_string())
1614        .trim_end_matches('/')
1615        .to_string()
1616}
1617
1618#[cfg(feature = "browserbase")]
1619fn browserbase_session_id(payload: &Value) -> Option<String> {
1620    payload
1621        .get("id")
1622        .or_else(|| payload.get("sessionId"))
1623        .or_else(|| payload.get("session_id"))
1624        .or_else(|| payload.get("data").and_then(|v| v.get("id")))
1625        .or_else(|| payload.get("data").and_then(|v| v.get("sessionId")))
1626        .or_else(|| payload.get("data").and_then(|v| v.get("session_id")))
1627        .and_then(Value::as_str)
1628        .map(ToString::to_string)
1629}
1630
1631#[cfg(feature = "browserbase")]
1632fn browserbase_connect_url(payload: &Value) -> Option<String> {
1633    [
1634        "connectUrl",
1635        "connect_url",
1636        "wsUrl",
1637        "ws_url",
1638        "websocketUrl",
1639        "websocket_url",
1640        "browserWSEndpoint",
1641        "wsEndpoint",
1642        "ws_endpoint",
1643    ]
1644    .iter()
1645    .find_map(|key| payload.get(*key).and_then(Value::as_str))
1646    .or_else(|| {
1647        payload.get("data").and_then(|data| {
1648            [
1649                "connectUrl",
1650                "connect_url",
1651                "wsUrl",
1652                "ws_url",
1653                "websocketUrl",
1654                "websocket_url",
1655                "browserWSEndpoint",
1656                "wsEndpoint",
1657                "ws_endpoint",
1658            ]
1659            .iter()
1660            .find_map(|key| data.get(*key).and_then(Value::as_str))
1661        })
1662    })
1663    .map(ToString::to_string)
1664}
1665
1666fn dedupe_preserve_order(stages: &mut Vec<StrategyUsed>) {
1667    let mut seen = Vec::new();
1668    stages.retain(|stage| {
1669        if seen.contains(stage) {
1670            false
1671        } else {
1672            seen.push(*stage);
1673            true
1674        }
1675    });
1676}
1677
1678#[cfg(feature = "browserbase")]
1679fn maybe_insert_browserbase_stage(stages: &mut Vec<StrategyUsed>, enabled: bool) {
1680    if !enabled || stages.contains(&StrategyUsed::BrowserbaseManagedSession) {
1681        return;
1682    }
1683
1684    if let Some(pos) = stages
1685        .iter()
1686        .position(|stage| *stage == StrategyUsed::StickyProxyBrowserSession)
1687    {
1688        stages.insert(pos, StrategyUsed::BrowserbaseManagedSession);
1689    } else {
1690        stages.push(StrategyUsed::BrowserbaseManagedSession);
1691    }
1692}
1693
1694fn classify_browser_error(error: &BrowserError) -> StageFailureKind {
1695    match error {
1696        BrowserError::Timeout { .. } => StageFailureKind::Timeout,
1697        BrowserError::NavigationFailed { reason, .. } if reason.contains("selector") => {
1698            StageFailureKind::Blocked
1699        }
1700        BrowserError::ScriptExecutionFailed { .. } => StageFailureKind::Extraction,
1701        BrowserError::ConfigError(_) | BrowserError::PoolExhausted { .. } => {
1702            StageFailureKind::Setup
1703        }
1704        BrowserError::RateLimited { .. } => StageFailureKind::RateLimited,
1705        BrowserError::ProxyUnavailable { .. }
1706        | BrowserError::ConnectionError { .. }
1707        | BrowserError::CdpError { .. }
1708        | BrowserError::LaunchFailed { .. }
1709        | BrowserError::NavigationFailed { .. }
1710        | BrowserError::Io(_)
1711        | BrowserError::StaleNode { .. } => StageFailureKind::Transport,
1712        #[cfg(feature = "extract")]
1713        BrowserError::ExtractionFailed(_) => StageFailureKind::Extraction,
1714    }
1715}
1716
1717const fn is_block_status(status: Option<u16>) -> bool {
1718    matches!(status, Some(401 | 403 | 407 | 429 | 503))
1719}
1720
1721fn truncate_html(html: &str, max_bytes: usize) -> String {
1722    if html.len() <= max_bytes {
1723        return html.to_string();
1724    }
1725
1726    let mut out = String::new();
1727    for ch in html.chars() {
1728        if out.len() + ch.len_utf8() > max_bytes {
1729            break;
1730        }
1731        out.push(ch);
1732    }
1733    out
1734}
1735
1736fn host_hint(url: &str) -> Option<String> {
1737    let without_scheme = url.split_once("://")?.1;
1738    let authority = without_scheme.split('/').next()?;
1739    let host = authority.rsplit('@').next()?.split(':').next()?;
1740    if host.is_empty() {
1741        None
1742    } else {
1743        Some(host.to_ascii_lowercase())
1744    }
1745}
1746
1747#[cfg(test)]
1748#[allow(
1749    clippy::unwrap_used,
1750    clippy::expect_used,
1751    clippy::panic,
1752    clippy::indexing_slicing
1753)]
1754mod tests {
1755    use super::*;
1756
1757    #[test]
1758    fn ladder_is_deterministic_for_modes() {
1759        assert_eq!(
1760            AcquisitionRunner::strategy_ladder(AcquisitionMode::Fast, None),
1761            vec![
1762                StrategyUsed::DirectHttp,
1763                StrategyUsed::TlsProfiledHttp,
1764                StrategyUsed::BrowserLightStealth,
1765            ]
1766        );
1767
1768        assert_eq!(
1769            AcquisitionRunner::strategy_ladder(
1770                AcquisitionMode::Investigate,
1771                Some(StrategyUsed::StickyProxyBrowserSession)
1772            ),
1773            vec![
1774                StrategyUsed::InvestigateEntry,
1775                StrategyUsed::StickyProxyBrowserSession,
1776                StrategyUsed::TlsProfiledHttp,
1777            ]
1778        );
1779    }
1780
1781    #[test]
1782    fn block_statuses_are_classified() {
1783        assert!(is_block_status(Some(403)));
1784        assert!(is_block_status(Some(429)));
1785        assert!(!is_block_status(Some(200)));
1786        assert!(!is_block_status(None));
1787    }
1788
1789    #[test]
1790    fn host_hint_extracts_authority() {
1791        assert_eq!(
1792            host_hint("https://user:pass@example.com:8443/path"),
1793            Some("example.com".to_string())
1794        );
1795    }
1796
1797    #[test]
1798    fn truncate_html_respects_utf8_boundaries() {
1799        let src = "abc😀def";
1800        let out = truncate_html(src, 5);
1801        assert_eq!(out, "abc");
1802    }
1803
1804    #[cfg(feature = "browserbase")]
1805    #[test]
1806    fn browserbase_connect_url_is_extracted_from_nested_data() {
1807        let payload = serde_json::json!({
1808            "data": {
1809                "connectUrl": "wss://connect.browserbase.example/devtools/browser/abc"
1810            }
1811        });
1812
1813        assert_eq!(
1814            browserbase_connect_url(&payload),
1815            Some("wss://connect.browserbase.example/devtools/browser/abc".to_string())
1816        );
1817    }
1818
1819    #[cfg(feature = "browserbase")]
1820    #[test]
1821    fn browserbase_stage_is_inserted_before_sticky_stage() {
1822        let mut ladder = vec![
1823            StrategyUsed::DirectHttp,
1824            StrategyUsed::StickyProxyBrowserSession,
1825            StrategyUsed::TlsProfiledHttp,
1826        ];
1827
1828        maybe_insert_browserbase_stage(&mut ladder, true);
1829
1830        assert_eq!(
1831            ladder,
1832            vec![
1833                StrategyUsed::DirectHttp,
1834                StrategyUsed::BrowserbaseManagedSession,
1835                StrategyUsed::StickyProxyBrowserSession,
1836                StrategyUsed::TlsProfiledHttp,
1837            ]
1838        );
1839    }
1840
1841    #[cfg(feature = "browserbase")]
1842    #[test]
1843    fn browserbase_session_config_defaults_enable_warmup_and_retry() {
1844        let config = BrowserbaseSessionConfig::default();
1845        assert!(config.session_id.is_none());
1846        assert!(config.warmup);
1847        assert_eq!(config.warmup_stabilize_ms, 500);
1848        assert_eq!(config.max_retries, 3);
1849        assert_eq!(config.backoff_base_ms, 500);
1850    }
1851
1852    #[cfg(feature = "browserbase")]
1853    #[test]
1854    fn resolve_session_id_prefers_configured_over_env() {
1855        assert_eq!(
1856            resolve_browserbase_session_id(Some("configured"), Some("from-env")),
1857            Some("configured".to_string())
1858        );
1859    }
1860
1861    #[cfg(feature = "browserbase")]
1862    #[test]
1863    fn resolve_session_id_falls_back_to_env_when_unconfigured() {
1864        assert_eq!(
1865            resolve_browserbase_session_id(None, Some("from-env")),
1866            Some("from-env".to_string())
1867        );
1868    }
1869
1870    #[cfg(feature = "browserbase")]
1871    #[test]
1872    fn resolve_session_id_treats_blank_values_as_unset() {
1873        assert_eq!(
1874            resolve_browserbase_session_id(Some("  "), Some("from-env")),
1875            Some("from-env".to_string())
1876        );
1877        assert_eq!(resolve_browserbase_session_id(Some(""), None), None);
1878        assert_eq!(resolve_browserbase_session_id(None, None), None);
1879    }
1880
1881    #[cfg(feature = "browserbase")]
1882    #[test]
1883    fn backoff_delay_doubles_per_attempt() {
1884        assert_eq!(backoff_delay_ms(500, 0), 500);
1885        assert_eq!(backoff_delay_ms(500, 1), 1_000);
1886        assert_eq!(backoff_delay_ms(500, 2), 2_000);
1887        assert_eq!(backoff_delay_ms(500, 3), 4_000);
1888    }
1889
1890    #[cfg(feature = "browserbase")]
1891    #[test]
1892    fn backoff_delay_saturates_instead_of_overflowing() {
1893        assert_eq!(backoff_delay_ms(u64::MAX, 10), u64::MAX);
1894        assert_eq!(backoff_delay_ms(1, 100), 1u64 << 62);
1895    }
1896
1897    #[cfg(feature = "browserbase")]
1898    #[test]
1899    fn retry_after_header_parses_seconds_to_millis() {
1900        let mut headers = reqwest::header::HeaderMap::new();
1901        headers.insert(
1902            reqwest::header::RETRY_AFTER,
1903            reqwest::header::HeaderValue::from_static("2"),
1904        );
1905        assert_eq!(parse_retry_after_ms(&headers), Some(2_000));
1906    }
1907
1908    #[cfg(feature = "browserbase")]
1909    #[test]
1910    fn retry_after_header_missing_returns_none() {
1911        let headers = reqwest::header::HeaderMap::new();
1912        assert_eq!(parse_retry_after_ms(&headers), None);
1913    }
1914
1915    #[cfg(feature = "browserbase")]
1916    #[test]
1917    fn retry_after_header_non_numeric_returns_none() {
1918        let mut headers = reqwest::header::HeaderMap::new();
1919        headers.insert(
1920            reqwest::header::RETRY_AFTER,
1921            reqwest::header::HeaderValue::from_static("Wed, 21 Oct 2026 07:28:00 GMT"),
1922        );
1923        assert_eq!(parse_retry_after_ms(&headers), None);
1924    }
1925
1926    #[tokio::test]
1927    async fn stale_freshness_contract_short_circuits_runner() {
1928        use crate::freshness::{FreshnessContract, FreshnessPolicyKind};
1929        use std::time::Duration;
1930
1931        let past_ms = crate::freshness::unix_epoch_ms().saturating_sub(60_000);
1932        let stale = FreshnessContract::with_signature(
1933            "example.com",
1934            "sha256:abc",
1935            past_ms,
1936            Duration::from_secs(1),
1937            FreshnessPolicyKind::Standard,
1938        )
1939        .expect("contract");
1940
1941        let request = AcquisitionRequest {
1942            url: "https://example.com/path".to_string(),
1943            mode: AcquisitionMode::Fast,
1944            total_timeout: Duration::from_secs(5),
1945            freshness_contract: Some(stale),
1946            ..AcquisitionRequest::default()
1947        };
1948
1949        // Synchronous contract check: we don't actually need a pool
1950        // because the runner should short-circuit before acquiring.
1951        let runner = AcquisitionRunner::new(crate::BrowserPool::placeholder());
1952        let result = runner.run(request).await;
1953
1954        assert!(!result.success, "stale contract must not succeed");
1955        assert!(
1956            result.freshness.is_some(),
1957            "freshness report must be attached"
1958        );
1959        let report = result.freshness.as_ref().expect("report");
1960        assert!(
1961            report.decision.is_invalid(),
1962            "decision should be invalid for stale contract, got {report:?}"
1963        );
1964        assert_eq!(
1965            report.decision.label(),
1966            "stale_ttl",
1967            "expected stale_ttl, got {}",
1968            report.decision.label()
1969        );
1970        assert_eq!(result.attempted.len(), 0, "no stages should be attempted");
1971        assert_eq!(result.failures.len(), 1, "exactly one structured failure");
1972        assert_eq!(
1973            result.failures.first().map(|f| f.kind),
1974            Some(StageFailureKind::Setup)
1975        );
1976    }
1977
1978    #[tokio::test]
1979    async fn domain_mismatch_freshness_short_circuits_runner() {
1980        use crate::freshness::{FreshnessContract, FreshnessPolicyKind};
1981        use std::time::Duration;
1982
1983        let captured = crate::freshness::unix_epoch_ms();
1984        let contract = FreshnessContract::with_signature(
1985            "example.com",
1986            "sha256:abc",
1987            captured,
1988            Duration::from_mins(1),
1989            FreshnessPolicyKind::Standard,
1990        )
1991        .expect("contract");
1992
1993        let request = AcquisitionRequest {
1994            url: "https://other.example/path".to_string(),
1995            mode: AcquisitionMode::Fast,
1996            total_timeout: Duration::from_secs(5),
1997            freshness_contract: Some(contract),
1998            ..AcquisitionRequest::default()
1999        };
2000
2001        let runner = AcquisitionRunner::new(crate::BrowserPool::placeholder());
2002        let result = runner.run(request).await;
2003
2004        assert!(!result.success);
2005        let report = result.freshness.as_ref().expect("report");
2006        assert_eq!(report.decision.label(), "domain_mismatch");
2007        assert_eq!(result.attempted.len(), 0);
2008    }
2009
2010    // ─── Replay defense (T81) ───────────────────────────────────────────────
2011
2012    #[tokio::test]
2013    async fn rotation_due_replay_defense_short_circuits_runner() {
2014        use crate::ReplayDefenseContext;
2015        use crate::replay_defense::{ReplayDefensePolicy, ReplayDefenseState};
2016        use std::time::Duration;
2017
2018        let past_ms = crate::replay_defense::unix_epoch_ms().saturating_sub(120_000);
2019        let state = ReplayDefenseState::new("example.com", None, None, past_ms);
2020        // 1 second rotation interval — anything older is "rotation due".
2021        let policy = ReplayDefensePolicy {
2022            rotation_interval: Duration::from_secs(1),
2023            ..ReplayDefensePolicy::default()
2024        };
2025        let context = ReplayDefenseContext::with_policy(policy, state);
2026
2027        let request = AcquisitionRequest {
2028            url: "https://example.com/path".to_string(),
2029            mode: AcquisitionMode::Fast,
2030            total_timeout: Duration::from_secs(5),
2031            replay_defense: Some(context),
2032            ..AcquisitionRequest::default()
2033        };
2034
2035        let runner = AcquisitionRunner::new(crate::BrowserPool::placeholder());
2036        let result = runner.run(request).await;
2037
2038        assert!(!result.success);
2039        let report = result
2040            .replay_defense
2041            .as_ref()
2042            .expect("replay defense report attached");
2043        assert_eq!(report.decision.label(), "rotation_due");
2044        assert!(report.forced_refresh);
2045        assert_eq!(result.attempted.len(), 0, "no stages attempted");
2046        assert_eq!(result.failures.len(), 1);
2047        assert_eq!(
2048            result.failures.first().map(|f| f.kind),
2049            Some(StageFailureKind::ReplayDefenseTriggered)
2050        );
2051    }
2052
2053    #[tokio::test]
2054    async fn nonce_expired_replay_defense_short_circuits_runner() {
2055        use crate::ReplayDefenseContext;
2056        use crate::replay_defense::{ReplayDefensePolicy, ReplayDefenseState};
2057        use std::time::Duration;
2058
2059        let past_ms = crate::replay_defense::unix_epoch_ms().saturating_sub(120_000);
2060        let state = ReplayDefenseState::new("example.com", None, Some("nonce-001"), past_ms);
2061        let policy = ReplayDefensePolicy {
2062            nonce_validity_window: Duration::from_secs(1),
2063            ..ReplayDefensePolicy::default()
2064        };
2065        let context = ReplayDefenseContext::with_policy(policy, state);
2066
2067        let request = AcquisitionRequest {
2068            url: "https://example.com/path".to_string(),
2069            mode: AcquisitionMode::Fast,
2070            total_timeout: Duration::from_secs(5),
2071            replay_defense: Some(context),
2072            ..AcquisitionRequest::default()
2073        };
2074
2075        let runner = AcquisitionRunner::new(crate::BrowserPool::placeholder());
2076        let result = runner.run(request).await;
2077
2078        assert!(!result.success);
2079        let report = result
2080            .replay_defense
2081            .as_ref()
2082            .expect("replay defense report attached");
2083        assert_eq!(report.decision.label(), "nonce_expired");
2084        assert!(report.forced_refresh);
2085        assert_eq!(result.attempted.len(), 0);
2086    }
2087
2088    #[tokio::test]
2089    async fn signature_drift_replay_defense_short_circuits_runner() {
2090        use crate::ReplayDefenseContext;
2091        use crate::replay_defense::{ReplayDefensePolicy, ReplayDefenseState};
2092        use std::time::Duration;
2093
2094        let captured = crate::replay_defense::unix_epoch_ms();
2095        let state =
2096            ReplayDefenseState::with_fingerprint("example.com", "sha256:abc", None, captured);
2097        // force_reset_on_drift = true (default)
2098        let policy = ReplayDefensePolicy {
2099            force_reset_on_drift: true,
2100            ..ReplayDefensePolicy::default()
2101        };
2102        let context = ReplayDefenseContext::with_policy(policy, state);
2103
2104        // URL has a #fragment that the state doesn't, but the host
2105        // matches. The runner reads the host out via host_hint, so
2106        // the observed signature in the input is the state signature
2107        // — and a forced refresh is triggered by the **policy**
2108        // check (state.signature != input.observed_signature) only
2109        // when they actually differ. The integration test below
2110        // covers that path on a real browser; here we just confirm
2111        // the runner accepts the context and emits a valid report.
2112        let request = AcquisitionRequest {
2113            url: "https://example.com/path".to_string(),
2114            mode: AcquisitionMode::Fast,
2115            total_timeout: Duration::from_secs(15),
2116            // Short request timeout so the HTTP stages fail fast on
2117            // the placeholder pool instead of being cut off by the
2118            // outer total_timeout.
2119            request_timeout: Duration::from_millis(100),
2120            replay_defense: Some(context),
2121            ..AcquisitionRequest::default()
2122        };
2123
2124        let runner = AcquisitionRunner::new(crate::BrowserPool::placeholder());
2125        let result = runner.run(request).await;
2126
2127        // Observed signature comes from the state itself (mirrors
2128        // the freshness check), so the decision is Valid here.
2129        let report = result
2130            .replay_defense
2131            .as_ref()
2132            .expect("replay defense report attached");
2133        assert_eq!(report.decision.label(), "valid");
2134        assert!(!report.forced_refresh);
2135    }
2136
2137    #[tokio::test]
2138    async fn valid_replay_defense_state_does_not_short_circuit() {
2139        use crate::ReplayDefenseContext;
2140        use crate::replay_defense::{ReplayDefensePolicy, ReplayDefenseState};
2141        use std::time::Duration;
2142
2143        let captured = crate::replay_defense::unix_epoch_ms();
2144        let state = ReplayDefenseState::new("example.com", None, None, captured);
2145        let policy = ReplayDefensePolicy {
2146            rotation_interval: Duration::from_mins(30),
2147            ..ReplayDefensePolicy::default()
2148        };
2149        let context = ReplayDefenseContext::with_policy(policy, state);
2150
2151        let request = AcquisitionRequest {
2152            url: "https://example.com/path".to_string(),
2153            mode: AcquisitionMode::Fast,
2154            total_timeout: Duration::from_secs(15),
2155            // Short request timeout so the HTTP stages fail fast on
2156            // the placeholder pool instead of being cut off by the
2157            // outer total_timeout.
2158            request_timeout: Duration::from_millis(100),
2159            replay_defense: Some(context),
2160            ..AcquisitionRequest::default()
2161        };
2162
2163        let runner = AcquisitionRunner::new(crate::BrowserPool::placeholder());
2164        let result = runner.run(request).await;
2165
2166        // No forced refresh — but the placeholder pool will still
2167        // fail the run with PoolExhausted, so the run is reported as
2168        // unsuccessful (success = false) but the replay defense
2169        // report itself must be Valid.
2170        let report = result
2171            .replay_defense
2172            .as_ref()
2173            .expect("replay defense report attached");
2174        assert_eq!(report.decision.label(), "valid");
2175        assert!(!report.forced_refresh);
2176    }
2177
2178    #[test]
2179    fn replay_defense_context_with_default_policy_uses_baseline() {
2180        use crate::ReplayDefenseContext;
2181        use crate::replay_defense::ReplayDefenseState;
2182
2183        let state = ReplayDefenseState::new("example.com", None, None, 0);
2184        let context = ReplayDefenseContext::new(state);
2185        // Default policy: 30 min rotation, 5 min nonce, force_reset_on_drift
2186        assert_eq!(context.policy.rotation_interval, Duration::from_mins(30));
2187        assert_eq!(context.policy.nonce_validity_window, Duration::from_mins(5));
2188        assert!(context.policy.force_reset_on_drift);
2189    }
2190
2191    // ─── Interstitial routing (T94) ──────────────────────────────────────────
2192
2193    #[tokio::test]
2194    async fn interstitial_hard_block_short_circuits_runner() {
2195        use crate::InterstitialContext;
2196        use crate::interstitial_router::PageSignature;
2197        use std::time::Duration;
2198
2199        let signature = PageSignature::new("https://example.com/blocked", Some(403))
2200            .with_body_marker("access denied");
2201        let context = InterstitialContext::new(signature);
2202
2203        let request = AcquisitionRequest {
2204            url: "https://example.com/path".to_string(),
2205            mode: AcquisitionMode::Fast,
2206            total_timeout: Duration::from_secs(15),
2207            request_timeout: Duration::from_millis(100),
2208            interstitial: Some(context),
2209            ..AcquisitionRequest::default()
2210        };
2211
2212        let runner = AcquisitionRunner::new(crate::BrowserPool::placeholder());
2213        let result = runner.run(request).await;
2214
2215        let decision = result
2216            .interstitial
2217            .as_ref()
2218            .expect("interstitial decision attached");
2219        assert_eq!(decision.kind().label(), "hard_block");
2220        assert!(decision.is_terminal());
2221        assert_eq!(result.attempted.len(), 0, "no stages attempted");
2222        assert_eq!(result.failures.len(), 1);
2223        assert_eq!(
2224            result.failures.first().map(|f| f.kind),
2225            Some(StageFailureKind::InterstitialRouted)
2226        );
2227    }
2228
2229    #[tokio::test]
2230    async fn interstitial_queue_short_circuits_runner() {
2231        use crate::InterstitialContext;
2232        use crate::interstitial_router::PageSignature;
2233        use std::time::Duration;
2234
2235        let signature = PageSignature::new("https://example.com/queue", Some(200))
2236            .with_body_marker("please wait")
2237            .with_queue_position(3);
2238        let context = InterstitialContext::new(signature);
2239
2240        let request = AcquisitionRequest {
2241            url: "https://example.com/path".to_string(),
2242            mode: AcquisitionMode::Fast,
2243            total_timeout: Duration::from_secs(15),
2244            request_timeout: Duration::from_millis(100),
2245            interstitial: Some(context),
2246            ..AcquisitionRequest::default()
2247        };
2248
2249        let runner = AcquisitionRunner::new(crate::BrowserPool::placeholder());
2250        let result = runner.run(request).await;
2251
2252        let decision = result
2253            .interstitial
2254            .as_ref()
2255            .expect("interstitial decision attached");
2256        assert_eq!(decision.kind().label(), "queue");
2257        assert!(decision.is_retryable());
2258        assert_eq!(result.attempted.len(), 0, "no stages attempted");
2259        assert_eq!(
2260            result.failures.first().map(|f| f.kind),
2261            Some(StageFailureKind::InterstitialRouted)
2262        );
2263    }
2264
2265    #[tokio::test]
2266    async fn interstitial_challenge_short_circuits_runner() {
2267        use crate::InterstitialContext;
2268        use crate::interstitial_router::PageSignature;
2269        use std::time::Duration;
2270
2271        let signature = PageSignature::new(
2272            "https://example.com/cdn-cgi/challenge-platform/h/b",
2273            Some(403),
2274        )
2275        .with_body_marker("cf-chl-bypass")
2276        .with_vendor_hint("cloudflare");
2277        let context = InterstitialContext::new(signature);
2278
2279        let request = AcquisitionRequest {
2280            url: "https://example.com/path".to_string(),
2281            mode: AcquisitionMode::Fast,
2282            total_timeout: Duration::from_secs(15),
2283            request_timeout: Duration::from_millis(100),
2284            interstitial: Some(context),
2285            ..AcquisitionRequest::default()
2286        };
2287
2288        let runner = AcquisitionRunner::new(crate::BrowserPool::placeholder());
2289        let result = runner.run(request).await;
2290
2291        let decision = result
2292            .interstitial
2293            .as_ref()
2294            .expect("interstitial decision attached");
2295        assert_eq!(decision.kind().label(), "challenge");
2296        assert!(decision.requires_solve());
2297        assert_eq!(result.attempted.len(), 0, "no stages attempted");
2298        assert_eq!(
2299            result.failures.first().map(|f| f.kind),
2300            Some(StageFailureKind::InterstitialRouted)
2301        );
2302    }
2303
2304    #[tokio::test]
2305    async fn interstitial_transient_does_not_short_circuit() {
2306        use crate::InterstitialContext;
2307        use crate::interstitial_router::PageSignature;
2308        use std::time::Duration;
2309
2310        let signature = PageSignature::new("https://example.com/redirect", Some(302));
2311        let context = InterstitialContext::new(signature);
2312
2313        let request = AcquisitionRequest {
2314            url: "https://example.com/path".to_string(),
2315            mode: AcquisitionMode::Fast,
2316            total_timeout: Duration::from_secs(15),
2317            request_timeout: Duration::from_millis(100),
2318            interstitial: Some(context),
2319            ..AcquisitionRequest::default()
2320        };
2321
2322        let runner = AcquisitionRunner::new(crate::BrowserPool::placeholder());
2323        let result = runner.run(request).await;
2324
2325        // The transient decision is still attached but
2326        // the runner does NOT short-circuit, so the
2327        // `InterstitialRouted` failure must be absent.
2328        let decision = result
2329            .interstitial
2330            .as_ref()
2331            .expect("interstitial decision attached");
2332        assert_eq!(decision.kind().label(), "transient");
2333        assert!(!decision.is_classified());
2334        assert!(
2335            result
2336                .failures
2337                .iter()
2338                .all(|f| f.kind != StageFailureKind::InterstitialRouted),
2339            "transient interstitial must not short-circuit"
2340        );
2341    }
2342
2343    #[test]
2344    fn interstitial_context_with_default_policy_uses_baseline() {
2345        use crate::interstitial_router::PageSignature;
2346
2347        let signature = PageSignature::new("https://example.com", None);
2348        let context = InterstitialContext::new(signature);
2349        assert_eq!(
2350            context.policy.queue_max_retries,
2351            crate::interstitial_router::DEFAULT_QUEUE_MAX_RETRIES
2352        );
2353        assert!(context.policy.short_circuit_on_classified);
2354    }
2355}