1use std::{
4 fmt,
5 sync::{Arc, Mutex, OnceLock},
6 time::Duration,
7};
8
9use anyhow::anyhow;
10use futures_util::future::BoxFuture;
11use millipede_core::{
12 crawler::{
13 AttemptObservation, Crawler, CrawlerEnv, CrawlerHandle, CrawlerKind, RequestEnv,
14 RequestOutcome,
15 },
16 enqueue::EnqueueLinker,
17 errors::CrawlError,
18 events::CrawlerEvent,
19 http_client::{HttpClient, HttpClientError, HttpStatusError},
20 link_extraction::{ExtractedLink, LinkExtractor},
21 proxy::{ProxyConfiguration, ProxyInfo},
22 request::Request,
23 router::HasRequest,
24 session::{Session, SessionPool, SessionPoolOptions},
25 storage::StorageHandle,
26};
27
28use crate::{
29 BrowserError, BrowserHooks, BrowserPool, BrowserPoolOptions, BrowserPostHookCtx,
30 BrowserPostNavigationHook, BrowserPreHookCtx, BrowserPreNavigationHook, BrowserProvider,
31 BrowserResponse, GotoOptions, PageHandle, PageOptions, WaitUntil,
32};
33
34struct BrowserLinkExtractor {
35 page: PageHandle,
36}
37
38#[async_trait::async_trait]
39impl LinkExtractor for BrowserLinkExtractor {
40 async fn extract(&self, selector: Option<&str>) -> Result<Vec<ExtractedLink>, CrawlError> {
41 self.page
42 .evaluate_anchors(selector)
43 .await
44 .map_err(BrowserError::classify)
45 .map(|urls| {
46 urls.into_iter()
47 .map(|url| ExtractedLink {
48 url: url.to_string(),
49 base: None,
50 })
51 .collect()
52 })
53 }
54}
55
56#[derive(Clone)]
61#[non_exhaustive]
62pub struct BrowserContext {
63 pub request: Arc<Request>,
65 pub page: PageHandle,
70 pub response: Option<BrowserResponse>,
72 pub session: Option<Arc<Session>>,
74 pub proxy_info: Option<ProxyInfo>,
79 pub enqueue: EnqueueLinker,
81 pub storage: StorageHandle,
83 pub send_request: Arc<dyn HttpClient>,
85 pub crawler: CrawlerHandle,
87}
88
89impl HasRequest for BrowserContext {
90 fn request(&self) -> &Request {
91 &self.request
92 }
93}
94
95impl fmt::Debug for BrowserContext {
96 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
97 formatter
98 .debug_struct("BrowserContext")
99 .field("request", &self.request)
100 .field("page", &self.page)
101 .field("response", &self.response)
102 .field("session", &self.session)
103 .field("proxy_info", &self.proxy_info)
104 .field("enqueue", &self.enqueue)
105 .field("storage", &self.storage)
106 .field("send_request", &"<dyn HttpClient>")
107 .field("crawler", &self.crawler)
108 .finish()
109 }
110}
111
112enum SessionMode {
113 Disabled,
114 Owned(Arc<SessionPool>),
115 Shared(Arc<SessionPool>),
116}
117
118pub struct BrowserKind<P: BrowserProvider> {
126 pool: BrowserPool<P>,
127 sessions: SessionMode,
128 send_request: Arc<dyn HttpClient>,
129 retry_status_codes: Vec<u16>,
130 retry_server_errors: bool,
131 session_status_codes: Vec<u16>,
132 goto: GotoOptions,
133 pre_hooks: Vec<BrowserPreNavigationHook>,
134 post_hooks: Vec<BrowserPostNavigationHook>,
135 snapshot_errors: bool,
136 storage: OnceLock<StorageHandle>,
137 persist_task: Mutex<Option<tokio::task::JoinHandle<()>>>,
138}
139
140impl<P: BrowserProvider> BrowserKind<P> {
141 pub fn builder(provider: P) -> BrowserKindBuilder<P> {
143 BrowserKindBuilder::new(provider)
144 }
145
146 fn session_pool(&self) -> Option<&Arc<SessionPool>> {
147 match &self.sessions {
148 SessionMode::Disabled => None,
149 SessionMode::Owned(pool) | SessionMode::Shared(pool) => Some(pool),
150 }
151 }
152
153 async fn close_after_error(&self, page: &PageHandle, reason: &'static str) {
154 if let Err(close_error) = page.close().await {
155 tracing::warn!(%close_error, "failed to close page after {reason}");
156 }
157 }
158
159 async fn classify_status(
160 &self,
161 status: http::StatusCode,
162 session: Option<&Arc<Session>>,
163 ) -> Result<(), CrawlError> {
164 let code = status.as_u16();
165 if self.session_status_codes.contains(&code) {
166 if let Some(session) = session {
167 session.mark_bad().await;
168 }
169 return Err(CrawlError::session(HttpStatusError::new(status)));
170 }
171 if self.retry_status_codes.contains(&code)
172 || (self.retry_server_errors && status.is_server_error())
173 {
174 return Err(CrawlError::retry(HttpStatusError::new(status)));
175 }
176 if !status.is_success() && !status.is_redirection() {
177 return Err(CrawlError::non_retryable(HttpStatusError::new(status)));
178 }
179 if let Some(session) = session {
180 session.mark_good().await;
181 }
182 Ok(())
183 }
184
185 pub(crate) async fn execute_with_session(
186 &self,
187 env: RequestEnv<'_>,
188 session: Option<Arc<Session>>,
189 ) -> Result<BrowserContext, CrawlError> {
190 let mut page_opts = PageOptions::new();
191 if let Some(session) = &session {
192 page_opts = page_opts.with_session(Arc::clone(session));
193 }
194 let page = self
195 .pool
196 .new_page(page_opts)
197 .await
198 .map_err(BrowserError::classify)?;
199 for hook in &self.pre_hooks {
200 let context = BrowserPreHookCtx {
201 request: &env.request,
202 page: &page,
203 session: session.as_deref(),
204 proxy: page.proxy_info(),
205 };
206 if let Err(error) = hook(context).await {
207 self.close_after_error(&page, "pre-navigation hook error")
208 .await;
209 return Err(error);
210 }
211 }
212 let response = match page.goto(&env.request.url, self.goto.clone()).await {
213 Ok(response) => response,
214 Err(error) => {
215 self.close_after_error(&page, "navigation error").await;
216 return Err(error.classify());
217 }
218 };
219 if let Some(status) = response.as_ref().and_then(|response| response.status) {
220 if let Err(error) = self.classify_status(status, session.as_ref()).await {
221 self.close_after_error(&page, "HTTP status error").await;
222 return Err(error);
223 }
224 } else if let Some(session) = &session {
225 session.mark_good().await;
226 }
227 for hook in &self.post_hooks {
228 let context = BrowserPostHookCtx {
229 request: &env.request,
230 page: &page,
231 response: response.as_ref(),
232 session: session.as_deref(),
233 proxy: page.proxy_info(),
234 };
235 if let Err(error) = hook(context).await {
236 self.close_after_error(&page, "post-navigation hook error")
237 .await;
238 return Err(error);
239 }
240 }
241 let storage = match self.storage.get().cloned() {
242 Some(storage) => storage,
243 None => {
244 self.close_after_error(&page, "storage initialization error")
245 .await;
246 return Err(CrawlError::critical(anyhow!(
247 "BrowserKind::execute before start"
248 )));
249 }
250 };
251 let enqueue = EnqueueLinker::with_extractor(
252 env.crawler.clone(),
253 &env.request,
254 Arc::new(BrowserLinkExtractor { page: page.clone() }),
255 );
256 Ok(BrowserContext {
257 request: env.request.clone(),
258 proxy_info: page.proxy_info().cloned(),
259 page,
260 response,
261 session,
262 enqueue,
263 storage,
264 send_request: self.send_request.clone(),
265 crawler: env.crawler,
266 })
267 }
268}
269
270#[must_use = "builders do nothing unless consumed by build"]
272pub struct BrowserKindBuilder<P: BrowserProvider> {
273 provider: P,
274 pool_options: BrowserPoolOptions<P::LaunchOptions>,
275 session_pool: Option<SessionPoolOptions>,
276 shared_sessions: Option<Arc<SessionPool>>,
277 http_client: Option<Arc<dyn HttpClient>>,
278 retry_status_codes: Vec<u16>,
279 retry_server_errors: bool,
280 session_status_codes: Vec<u16>,
281 navigation_timeout: Duration,
282 wait_until: WaitUntil,
283 pre_hooks: Vec<BrowserPreNavigationHook>,
284 post_hooks: Vec<BrowserPostNavigationHook>,
285 snapshot_errors: bool,
286}
287
288impl<P: BrowserProvider> BrowserKindBuilder<P> {
289 fn new(provider: P) -> Self {
290 let pool_options = BrowserPoolOptions::default()
291 .with_hooks(BrowserHooks::default().with_session_cookie_sync());
292 Self {
293 provider,
294 pool_options,
295 session_pool: Some(SessionPoolOptions::default()),
296 shared_sessions: None,
297 http_client: None,
298 retry_status_codes: vec![408, 429],
299 retry_server_errors: true,
300 session_status_codes: vec![401, 403],
301 navigation_timeout: Duration::from_secs(30),
302 wait_until: WaitUntil::Load,
303 pre_hooks: Vec::new(),
304 post_hooks: Vec::new(),
305 snapshot_errors: false,
306 }
307 }
308
309 pub fn pool_options(mut self, options: BrowserPoolOptions<P::LaunchOptions>) -> Self {
311 self.pool_options = options;
312 self
313 }
314
315 pub fn launch_options(mut self, options: P::LaunchOptions) -> Self {
317 self.pool_options.launch_options = options;
318 self
319 }
320
321 pub fn max_open_pages_per_browser(mut self, value: usize) -> Self {
323 self.pool_options.max_open_pages_per_browser = value;
324 self
325 }
326
327 pub fn retire_browser_after_page_count(mut self, value: u64) -> Self {
329 self.pool_options.retire_browser_after_page_count = value;
330 self
331 }
332
333 pub fn max_browsers(mut self, value: usize) -> Self {
335 self.pool_options.max_browsers = Some(value);
336 self
337 }
338
339 pub fn proxy(mut self, proxy: ProxyConfiguration) -> Self {
341 self.pool_options.proxy = Some(proxy);
342 self
343 }
344
345 pub fn hooks(mut self, hooks: BrowserHooks) -> Self {
350 self.pool_options.hooks = hooks;
351 self
352 }
353
354 pub fn session_pool(mut self, options: SessionPoolOptions) -> Self {
356 self.session_pool = Some(options);
357 self.shared_sessions = None;
358 self
359 }
360
361 pub fn disable_sessions(mut self) -> Self {
363 self.session_pool = None;
364 self.shared_sessions = None;
365 self
366 }
367
368 pub fn shared_session_pool(mut self, pool: Arc<SessionPool>) -> Self {
370 self.shared_sessions = Some(pool);
371 self
372 }
373
374 pub fn http_client(mut self, client: Arc<dyn HttpClient>) -> Self {
376 self.http_client = Some(client);
377 self
378 }
379
380 pub fn retry_status_codes(mut self, codes: impl IntoIterator<Item = u16>) -> Self {
382 self.retry_status_codes = codes.into_iter().collect();
383 self
384 }
385
386 pub fn retry_server_errors(mut self, enabled: bool) -> Self {
388 self.retry_server_errors = enabled;
389 self
390 }
391
392 pub fn session_status_codes(mut self, codes: impl IntoIterator<Item = u16>) -> Self {
394 self.session_status_codes = codes.into_iter().collect();
395 self
396 }
397
398 pub fn navigation_timeout(mut self, timeout: Duration) -> Self {
400 self.navigation_timeout = timeout;
401 self
402 }
403
404 pub fn wait_until(mut self, wait_until: WaitUntil) -> Self {
406 self.wait_until = wait_until;
407 self
408 }
409
410 pub fn pre_navigation_hook<F>(mut self, hook: F) -> Self
412 where
413 F: for<'a> Fn(BrowserPreHookCtx<'a>) -> BoxFuture<'a, Result<(), CrawlError>>
414 + Send
415 + Sync
416 + 'static,
417 {
418 self.pre_hooks.push(Arc::new(hook));
419 self
420 }
421
422 pub fn post_navigation_hook<F>(mut self, hook: F) -> Self
424 where
425 F: for<'a> Fn(BrowserPostHookCtx<'a>) -> BoxFuture<'a, Result<(), CrawlError>>
426 + Send
427 + Sync
428 + 'static,
429 {
430 self.post_hooks.push(Arc::new(hook));
431 self
432 }
433
434 pub fn snapshot_errors_on_failure(mut self, enabled: bool) -> Self {
436 self.snapshot_errors = enabled;
437 self
438 }
439
440 pub fn build(self) -> Result<BrowserKind<P>, HttpClientError> {
442 let send_request = match self.http_client {
443 Some(client) => client,
444 None => Arc::new(millipede_http::ReqwestClient::new()?),
445 };
446 let sessions = if let Some(pool) = self.shared_sessions {
447 SessionMode::Shared(pool)
448 } else if let Some(options) = self.session_pool {
449 SessionMode::Owned(Arc::new(SessionPool::new(options)))
450 } else {
451 SessionMode::Disabled
452 };
453 Ok(BrowserKind {
454 pool: BrowserPool::new(self.provider, self.pool_options),
455 sessions,
456 send_request,
457 retry_status_codes: self.retry_status_codes,
458 retry_server_errors: self.retry_server_errors,
459 session_status_codes: self.session_status_codes,
460 goto: GotoOptions::default()
461 .with_timeout(self.navigation_timeout)
462 .with_wait_until(self.wait_until),
463 pre_hooks: self.pre_hooks,
464 post_hooks: self.post_hooks,
465 snapshot_errors: self.snapshot_errors,
466 storage: OnceLock::new(),
467 persist_task: Mutex::new(None),
468 })
469 }
470}
471
472pub type BrowserCrawler<P> = Crawler<BrowserKind<P>>;
474
475impl<P: BrowserProvider> CrawlerKind for BrowserKind<P> {
476 type Context = BrowserContext;
477
478 fn start<'a>(&'a self, env: &'a CrawlerEnv) -> BoxFuture<'a, Result<(), CrawlError>> {
479 Box::pin(async move {
480 let client = env.storage_client().cloned().ok_or_else(|| {
481 CrawlError::non_retryable(anyhow!("BrowserKind requires a storage client"))
482 })?;
483 let kvs = match env.kvs() {
484 Some(kvs) => kvs.clone(),
485 None => client
486 .open_key_value_store(Some(env.config().default_key_value_store_id()))
487 .await
488 .map_err(|error| CrawlError::retry(anyhow!(error)))?,
489 };
490 let dataset = client
491 .open_dataset(Some(env.config().default_dataset_id()))
492 .await
493 .map_err(|error| CrawlError::retry(anyhow!(error)))?;
494 let queue = env.request_queue().clone();
495 let _ = self
496 .storage
497 .set(StorageHandle::new(client, dataset, kvs.clone(), queue));
498
499 if let SessionMode::Owned(pool) = &self.sessions {
500 pool.attach_persistence(kvs);
501 pool.restore().await?;
502 let pool = Arc::clone(pool);
503 let mut events = env.events().subscribe();
504 let task = tokio::spawn(async move {
505 loop {
506 match events.recv().await {
507 Ok(CrawlerEvent::PersistState { .. }) => {
508 if let Err(error) = pool.persist().await {
509 tracing::warn!(%error, "session pool persistence failed");
510 }
511 }
512 Ok(CrawlerEvent::Exiting | CrawlerEvent::Aborting) => break,
513 Ok(_) | Err(tokio::sync::broadcast::error::RecvError::Lagged(_)) => {}
514 Err(tokio::sync::broadcast::error::RecvError::Closed) => break,
515 }
516 }
517 });
518 *self
519 .persist_task
520 .lock()
521 .unwrap_or_else(|error| error.into_inner()) = Some(task);
522 }
523 Ok(())
524 })
525 }
526
527 fn execute<'a>(
528 &'a self,
529 env: RequestEnv<'a>,
530 ) -> BoxFuture<'a, Result<Self::Context, CrawlError>> {
531 Box::pin(async move {
532 let session = if let Some(pool) = self.session_pool() {
533 Some(pool.session(None).await)
534 } else {
535 None
536 };
537 self.execute_with_session(env, session).await
538 })
539 }
540
541 fn observe(&self, ctx: &Self::Context) -> AttemptObservation {
542 let mut observation = AttemptObservation::default();
543 observation.status = ctx.response.as_ref().and_then(|response| response.status);
544 observation.loaded_url = ctx
545 .response
546 .as_ref()
547 .and_then(|response| response.url.clone())
548 .or_else(|| Some(ctx.request.url.clone()));
549 observation.session_id = ctx.session.as_ref().map(|session| session.id().clone());
550 observation.proxy_info = ctx.proxy_info.clone();
551 observation.response_bytes = None;
552 observation
553 }
554
555 fn cleanup(
556 &self,
557 outcome: RequestOutcome<Self::Context>,
558 ) -> BoxFuture<'_, Result<(), CrawlError>> {
559 Box::pin(async move {
560 match outcome {
561 RequestOutcome::Handled(ctx) => {
562 if let Err(error) = ctx.page.close().await {
563 tracing::warn!(%error, "failed to close handled browser page");
564 }
565 }
566 RequestOutcome::HandlerFailed { ctx, error } => {
567 if error.rotates_session() {
568 if let Some(session) = &ctx.session {
569 session.mark_bad().await;
570 }
571 }
572 if self.snapshot_errors {
573 let snapshotter = millipede_core::snapshot::ErrorSnapshotter::new(
574 ctx.storage.key_value_store().clone(),
575 );
576 match ctx.page.content().await {
577 Ok(html) => {
578 if let Err(error) = snapshotter
579 .capture(
580 &ctx.request,
581 "html",
582 bytes::Bytes::from(html),
583 "text/html",
584 )
585 .await
586 {
587 tracing::warn!(%error, "error HTML snapshot capture failed");
588 }
589 }
590 Err(error) => {
591 tracing::warn!(%error, "error snapshot content() failed");
592 }
593 }
594 match ctx
595 .page
596 .screenshot(crate::ScreenshotOptions::default())
597 .await
598 {
599 Ok(png) => {
600 if let Err(error) = snapshotter
601 .capture(&ctx.request, "png", png, "image/png")
602 .await
603 {
604 tracing::warn!(%error, "error screenshot capture failed");
605 }
606 }
607 Err(error) => {
608 tracing::warn!(%error, "error snapshot screenshot() failed");
609 }
610 }
611 }
612 if let Err(error) = ctx.page.close().await {
613 tracing::warn!(%error, "failed to close browser page after handler error");
614 }
615 }
616 RequestOutcome::ExecuteFailed { .. } => {
617 }
621 }
622 Ok(())
623 })
624 }
625
626 fn stop<'a>(&'a self, _env: &'a CrawlerEnv) -> BoxFuture<'a, Result<(), CrawlError>> {
627 Box::pin(async move {
628 if let Some(task) = self
629 .persist_task
630 .lock()
631 .unwrap_or_else(|error| error.into_inner())
632 .take()
633 {
634 task.abort();
635 }
636 if let SessionMode::Owned(pool) = &self.sessions {
637 if let Err(error) = pool.persist().await {
638 tracing::warn!(%error, "final session pool persistence failed");
639 }
640 }
641 if let Err(error) = self.pool.shutdown().await {
642 tracing::warn!(%error, "browser pool shutdown failed");
643 }
644 Ok(())
645 })
646 }
647}