1use 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#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
35#[serde(rename_all = "snake_case")]
36pub enum AcquisitionMode {
37 Fast,
39 Resilient,
41 Hostile,
43 Investigate,
45}
46
47#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
49#[serde(rename_all = "snake_case")]
50pub enum StrategyUsed {
51 DirectHttp,
53 TlsProfiledHttp,
55 BrowserLightStealth,
57 StickyProxyBrowserSession,
59 #[cfg(feature = "browserbase")]
61 BrowserbaseManagedSession,
62 InvestigateEntry,
64}
65
66#[derive(Debug, Clone)]
79pub struct ReplayDefenseContext {
80 pub policy: ReplayDefensePolicy,
82 pub state: ReplayDefenseState,
84}
85
86impl ReplayDefenseContext {
87 #[must_use]
89 pub fn new(state: ReplayDefenseState) -> Self {
90 Self {
91 policy: ReplayDefensePolicy::default(),
92 state,
93 }
94 }
95
96 #[must_use]
98 pub const fn with_policy(policy: ReplayDefensePolicy, state: ReplayDefenseState) -> Self {
99 Self { policy, state }
100 }
101}
102
103#[derive(Debug, Clone)]
116pub struct TransportRealismContext {
117 pub profile: TransportProfile,
119 pub observation: Option<TransportObservation>,
124}
125
126impl TransportRealismContext {
127 #[must_use]
129 pub const fn new(profile: TransportProfile) -> Self {
130 Self {
131 profile,
132 observation: None,
133 }
134 }
135
136 #[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 #[must_use]
150 pub fn with_observation_opt(mut self, observation: Option<TransportObservation>) -> Self {
151 self.observation = observation;
152 self
153 }
154
155 #[must_use]
157 pub fn with_profile(mut self, profile: TransportProfile) -> Self {
158 self.profile = profile;
159 self
160 }
161}
162
163#[derive(Debug, Clone)]
189pub struct InterstitialContext {
190 pub signature: PageSignature,
192 pub policy: InterstitialPolicy,
196}
197
198impl InterstitialContext {
199 #[must_use]
201 pub fn new(signature: PageSignature) -> Self {
202 Self {
203 signature,
204 policy: InterstitialPolicy::default(),
205 }
206 }
207
208 #[must_use]
211 pub const fn with_policy(policy: InterstitialPolicy, signature: PageSignature) -> Self {
212 Self { signature, policy }
213 }
214
215 #[must_use]
217 pub const fn with_policy_opt(mut self, policy: InterstitialPolicy) -> Self {
218 self.policy = policy;
219 self
220 }
221}
222
223#[derive(Debug, Clone)]
225pub struct AcquisitionRequest {
226 pub url: String,
228 pub mode: AcquisitionMode,
230 pub wait_for_selector: Option<String>,
232 pub extraction_js: Option<String>,
234 pub total_timeout: Duration,
236 pub navigation_timeout: Duration,
238 pub request_timeout: Duration,
240 pub html_excerpt_bytes: usize,
242 pub investigate_start: Option<StrategyUsed>,
244 pub browserbase_enabled: bool,
246 pub freshness_contract: Option<FreshnessContract>,
255 pub replay_defense: Option<ReplayDefenseContext>,
266 pub transport_realism: Option<TransportRealismContext>,
277 pub interstitial: Option<InterstitialContext>,
294 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#[derive(Debug, Clone, Serialize, Deserialize)]
334pub struct BrowserbaseSessionConfig {
335 pub session_id: Option<String>,
343 pub warmup: bool,
351 pub warmup_stabilize_ms: u64,
353 pub max_retries: u8,
356 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#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
377#[serde(rename_all = "snake_case")]
378pub enum StageFailureKind {
379 Setup,
381 Timeout,
383 Blocked,
385 Transport,
387 Extraction,
389 ReplayDefenseTriggered,
398 InterstitialRouted,
413 RateLimited,
416}
417
418#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
420pub struct StageFailure {
421 pub strategy: StrategyUsed,
423 pub kind: StageFailureKind,
425 pub message: String,
427}
428
429#[derive(Debug, Clone, Serialize, Deserialize)]
431pub struct AcquisitionResult {
432 pub success: bool,
434 pub strategy_used: Option<StrategyUsed>,
436 pub attempted: Vec<StrategyUsed>,
438 pub final_url: Option<String>,
440 pub status_code: Option<u16>,
442 pub html_excerpt: Option<String>,
444 pub extracted: Option<Value>,
446 pub failures: Vec<StageFailure>,
448 pub timed_out: bool,
450 pub freshness: Option<FreshnessReport>,
456 pub replay_defense: Option<ReplayDefenseReport>,
467 pub transport_realism: Option<TransportRealismReport>,
475 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#[derive(Clone)]
532pub struct AcquisitionRunner {
533 pool: Arc<BrowserPool>,
534}
535
536impl AcquisitionRunner {
537 #[must_use]
551 pub const fn new(pool: Arc<BrowserPool>) -> Self {
552 Self { pool }
553 }
554
555 #[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 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 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 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 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 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 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 if self.evaluate_replay_defense(request, &mut result).await {
825 return result;
826 }
827
828 if Self::evaluate_interstitial(request, &mut result) {
838 return result;
839 }
840
841 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#[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#[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#[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#[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#[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#[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 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 #[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 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 let policy = ReplayDefensePolicy {
2099 force_reset_on_drift: true,
2100 ..ReplayDefensePolicy::default()
2101 };
2102 let context = ReplayDefenseContext::with_policy(policy, state);
2103
2104 let request = AcquisitionRequest {
2113 url: "https://example.com/path".to_string(),
2114 mode: AcquisitionMode::Fast,
2115 total_timeout: Duration::from_secs(15),
2116 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 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 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 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 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 #[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 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}