1use std::{
4 collections::HashMap,
5 fmt,
6 ops::Deref,
7 panic::AssertUnwindSafe,
8 sync::{
9 Arc, Mutex as StdMutex,
10 atomic::{AtomicU8, Ordering},
11 },
12 time::Duration,
13};
14
15use futures_util::FutureExt;
16use millipede_core::proxy::{ProxyConfiguration, ProxyInfo, ProxyResolveContext};
17use tokio::sync::{Mutex, Notify, mpsc, oneshot};
18
19use crate::{
20 BrowserError, BrowserHooks, BrowserPage, BrowserProvider, LaunchContext, PageId, PageOptions,
21};
22
23#[non_exhaustive]
25#[must_use = "browser pool options do nothing unless passed to BrowserPool::new"]
26pub struct BrowserPoolOptions<L> {
27 pub max_open_pages_per_browser: usize,
29 pub retire_browser_after_page_count: u64,
31 pub max_browsers: Option<usize>,
33 pub page_acquire_timeout: Duration,
35 pub launch_options: L,
37 pub proxy: Option<ProxyConfiguration>,
39 pub hooks: BrowserHooks,
41}
42
43impl<L: Default> Default for BrowserPoolOptions<L> {
44 fn default() -> Self {
45 Self {
46 max_open_pages_per_browser: 20,
47 retire_browser_after_page_count: 100,
48 max_browsers: None,
49 page_acquire_timeout: Duration::from_secs(60),
50 launch_options: L::default(),
51 proxy: None,
52 hooks: BrowserHooks::default(),
53 }
54 }
55}
56
57impl<L: Default> BrowserPoolOptions<L> {
58 pub fn new() -> Self {
60 Self::default()
61 }
62}
63
64impl<L> BrowserPoolOptions<L> {
65 pub fn with_max_open_pages_per_browser(mut self, value: usize) -> Self {
67 self.max_open_pages_per_browser = value;
68 self
69 }
70
71 pub fn with_retire_browser_after_page_count(mut self, value: u64) -> Self {
73 self.retire_browser_after_page_count = value;
74 self
75 }
76
77 pub fn with_max_browsers(mut self, value: Option<usize>) -> Self {
79 self.max_browsers = value;
80 self
81 }
82
83 pub fn with_page_acquire_timeout(mut self, value: Duration) -> Self {
85 self.page_acquire_timeout = value;
86 self
87 }
88
89 pub fn with_launch_options(mut self, value: L) -> Self {
91 self.launch_options = value;
92 self
93 }
94
95 pub fn with_proxy(mut self, value: Option<ProxyConfiguration>) -> Self {
97 self.proxy = value;
98 self
99 }
100
101 pub fn with_hooks(mut self, value: BrowserHooks) -> Self {
103 self.hooks = value;
104 self
105 }
106}
107
108impl<L> fmt::Debug for BrowserPoolOptions<L> {
109 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
110 formatter
111 .debug_struct("BrowserPoolOptions")
112 .field(
113 "max_open_pages_per_browser",
114 &self.max_open_pages_per_browser,
115 )
116 .field(
117 "retire_browser_after_page_count",
118 &self.retire_browser_after_page_count,
119 )
120 .field("max_browsers", &self.max_browsers)
121 .field("page_acquire_timeout", &self.page_acquire_timeout)
122 .field("launch_options", &"<provider-specific>")
123 .field("proxy", &self.proxy)
124 .field("hooks", &self.hooks)
125 .finish()
126 }
127}
128
129pub struct BrowserPool<P: BrowserProvider> {
131 inner: Arc<PoolInner<P>>,
132}
133
134impl<P: BrowserProvider> Clone for BrowserPool<P> {
135 fn clone(&self) -> Self {
136 Self {
137 inner: Arc::clone(&self.inner),
138 }
139 }
140}
141
142struct PoolInner<P: BrowserProvider> {
143 provider: P,
144 options: BrowserPoolOptions<P::LaunchOptions>,
145 state: Mutex<PoolState<P>>,
146 close_tx: mpsc::UnboundedSender<CloseCommand<P>>,
147 close_rx: StdMutex<Option<mpsc::UnboundedReceiver<CloseCommand<P>>>>,
148 capacity: Notify,
149 worker: StdMutex<Option<tokio::task::JoinHandle<()>>>,
150}
151
152struct PoolState<P: BrowserProvider> {
153 browsers: Vec<BrowserSlot<P>>,
154 pages: HashMap<PageId, PageEntry<P>>,
155 next_browser_id: u64,
156 shut_down: bool,
157}
158
159struct BrowserSlot<P: BrowserProvider> {
160 id: u64,
161 browser: Option<Arc<P::Browser>>,
162 open_pages: usize,
163 in_flight_pages: usize,
164 pages_created: u64,
165 launching: bool,
166 retired: bool,
167 proxy: Option<ProxyInfo>,
168}
169
170struct PageEntry<P: BrowserProvider> {
171 page: P::Page,
172 browser_id: u64,
173 opts: PageOptions,
174 closing: bool,
175 close_notify: Arc<Notify>,
176}
177
178enum CloseCommand<P: BrowserProvider> {
179 Page {
180 inner: Arc<PoolInner<P>>,
181 id: PageId,
182 },
183 Finalize {
184 inner: Arc<PoolInner<P>>,
185 },
186 Barrier(oneshot::Sender<()>),
187 Stop(oneshot::Sender<()>),
188}
189
190struct LaunchGuard<P: BrowserProvider> {
191 inner: Arc<PoolInner<P>>,
192 browser_id: u64,
193 armed: bool,
194}
195
196impl<P: BrowserProvider> LaunchGuard<P> {
197 fn new(inner: Arc<PoolInner<P>>, browser_id: u64) -> Self {
198 Self {
199 inner,
200 browser_id,
201 armed: true,
202 }
203 }
204
205 fn disarm(&mut self) {
206 self.armed = false;
207 }
208}
209
210impl<P: BrowserProvider> Drop for LaunchGuard<P> {
211 fn drop(&mut self) {
212 if !self.armed {
213 return;
214 }
215 let inner = Arc::clone(&self.inner);
216 let browser_id = self.browser_id;
217 if let Ok(mut state) = inner.state.try_lock() {
218 state.browsers.retain(|slot| slot.id != browser_id);
219 drop(state);
220 inner.capacity.notify_waiters();
221 } else if let Ok(runtime) = tokio::runtime::Handle::try_current() {
222 runtime.spawn(async move {
223 inner
224 .state
225 .lock()
226 .await
227 .browsers
228 .retain(|slot| slot.id != browser_id);
229 inner.capacity.notify_waiters();
230 });
231 } else {
232 tracing::warn!(
233 browser_id,
234 "cancelled browser launch dropped outside a Tokio runtime; launch placeholder cleanup could not be scheduled"
235 );
236 }
237 }
238}
239
240struct PageReservationGuard<P: BrowserProvider> {
241 inner: Arc<PoolInner<P>>,
242 browser_id: u64,
243 page: Option<P::Page>,
244 armed: bool,
245}
246
247impl<P: BrowserProvider> PageReservationGuard<P> {
248 fn new(inner: Arc<PoolInner<P>>, browser_id: u64) -> Self {
249 Self {
250 inner,
251 browser_id,
252 page: None,
253 armed: true,
254 }
255 }
256
257 fn set_page(&mut self, page: P::Page) {
258 self.page = Some(page);
259 }
260
261 fn clear_page(&mut self) {
262 self.page = None;
263 }
264
265 fn disarm(&mut self) {
266 self.armed = false;
267 self.page = None;
268 }
269}
270
271impl<P: BrowserProvider> Drop for PageReservationGuard<P> {
272 fn drop(&mut self) {
273 if !self.armed {
274 return;
275 }
276
277 let page = self.page.take();
278 let inner = Arc::clone(&self.inner);
279 let browser_id = self.browser_id;
280 if let Ok(runtime) = tokio::runtime::Handle::try_current() {
281 runtime.spawn(async move {
282 if let Some(page) = page {
283 if let Err(error) = inner.provider.close_page(page).await {
284 tracing::warn!(%error, "failed to close page after cancelled acquisition");
285 }
286 }
287 let browser = rollback_page_reservation(&inner, browser_id).await;
288 if let Some(browser) = browser {
289 close_browser_arc(inner.as_ref(), browser).await;
290 }
291 });
292 } else {
293 tracing::warn!(
294 "cancelled page acquisition dropped outside a Tokio runtime; provider cleanup was limited to dropping its handles"
295 );
296 }
297 }
298}
299
300async fn rollback_page_reservation<P: BrowserProvider>(
301 inner: &Arc<PoolInner<P>>,
302 browser_id: u64,
303) -> Option<Arc<P::Browser>> {
304 let mut state = inner.state.lock().await;
305 let shut_down = state.shut_down;
306 let mut browser = None;
307 if let Some(slot) = state.browsers.iter_mut().find(|slot| slot.id == browser_id) {
308 slot.open_pages = slot.open_pages.saturating_sub(1);
309 slot.in_flight_pages = slot.in_flight_pages.saturating_sub(1);
310 slot.pages_created = slot.pages_created.saturating_sub(1);
311 slot.retired = slot.pages_created >= inner.options.retire_browser_after_page_count;
312 if slot.open_pages == 0 && (slot.retired || shut_down) {
313 browser = slot.browser.take();
314 }
315 }
316 drop(state);
317 inner.capacity.notify_waiters();
318 browser
319}
320
321impl<P: BrowserProvider> BrowserPool<P> {
322 pub fn new(provider: P, options: BrowserPoolOptions<P::LaunchOptions>) -> Self {
328 let (close_tx, close_rx) = mpsc::unbounded_channel();
329 Self {
330 inner: Arc::new(PoolInner {
331 provider,
332 options,
333 state: Mutex::new(PoolState {
334 browsers: Vec::new(),
335 pages: HashMap::new(),
336 next_browser_id: 1,
337 shut_down: false,
338 }),
339 close_tx,
340 close_rx: StdMutex::new(Some(close_rx)),
341 capacity: Notify::new(),
342 worker: StdMutex::new(None),
343 }),
344 }
345 }
346
347 fn start_close_worker(&self) {
348 let mut worker = self
349 .inner
350 .worker
351 .lock()
352 .unwrap_or_else(|error| error.into_inner());
353 if worker.is_some() {
354 return;
355 }
356 let receiver = self
357 .inner
358 .close_rx
359 .lock()
360 .unwrap_or_else(|error| error.into_inner())
361 .take();
362 let Some(mut receiver) = receiver else {
363 return;
364 };
365 *worker = Some(tokio::spawn(async move {
366 while let Some(command) = receiver.recv().await {
367 match command {
368 CloseCommand::Page { inner, id } => {
369 match AssertUnwindSafe(close_page(inner, id)).catch_unwind().await {
370 Ok(Ok(())) => {}
371 Ok(Err(error)) => {
372 tracing::warn!(page_id = %id, %error, "background page close failed");
373 }
374 Err(_) => {
375 tracing::warn!(page_id = %id, "background page close panicked");
376 }
377 }
378 }
379 CloseCommand::Finalize { inner } => {
380 finalize_orphaned_pool(inner).await;
381 }
382 CloseCommand::Barrier(completion) => {
383 let _ = completion.send(());
384 }
385 CloseCommand::Stop(completion) => {
386 let _ = completion.send(());
387 break;
388 }
389 }
390 }
391 }));
392 }
393
394 pub async fn new_page(&self, opts: PageOptions) -> Result<PageHandle, BrowserError> {
396 self.start_close_worker();
397 let timeout = self.inner.options.page_acquire_timeout;
398 match tokio::time::timeout(timeout, self.acquire_page(opts)).await {
399 Ok(result) => result,
400 Err(_) => Err(BrowserError::PageCreate(anyhow::anyhow!(
401 "page acquisition timed out after {:?}",
402 timeout
403 ))),
404 }
405 }
406
407 async fn acquire_page(&self, opts: PageOptions) -> Result<PageHandle, BrowserError> {
408 loop {
409 let notified = self.inner.capacity.notified();
410 tokio::pin!(notified);
411 notified.as_mut().enable();
412
413 enum Choice<B> {
414 Existing {
415 id: u64,
416 browser: Arc<B>,
417 proxy: Option<ProxyInfo>,
418 },
419 Launch(u64),
420 Wait,
421 }
422
423 let choice = {
424 let mut state = self.inner.state.lock().await;
425 if state.shut_down {
426 return Err(BrowserError::Shutdown);
427 }
428 let available = state.browsers.iter().position(|slot| {
429 !slot.retired
430 && !slot.launching
431 && slot.browser.is_some()
432 && slot.open_pages < self.inner.options.max_open_pages_per_browser
433 && slot.pages_created < self.inner.options.retire_browser_after_page_count
434 });
435 if let Some(index) = available {
436 let slot = &mut state.browsers[index];
437 slot.open_pages += 1;
438 slot.in_flight_pages += 1;
439 slot.pages_created += 1;
440 Choice::Existing {
441 id: slot.id,
442 browser: Arc::clone(slot.browser.as_ref().expect("browser checked above")),
443 proxy: slot.proxy.clone(),
444 }
445 } else {
446 let live_count = state
447 .browsers
448 .iter()
449 .filter(|slot| slot.launching || slot.browser.is_some())
450 .count();
451 #[allow(clippy::unnecessary_map_or)]
452 let launch_allowed = self
453 .inner
454 .options
455 .max_browsers
456 .map_or(true, |maximum| live_count < maximum);
457 if launch_allowed {
458 let id = state.next_browser_id;
459 state.next_browser_id = state.next_browser_id.wrapping_add(1);
460 state.browsers.push(BrowserSlot {
461 id,
462 browser: None,
463 open_pages: 0,
464 in_flight_pages: 0,
465 pages_created: 0,
466 launching: true,
467 retired: false,
468 proxy: None,
469 });
470 Choice::Launch(id)
471 } else {
472 Choice::Wait
473 }
474 }
475 };
476
477 match choice {
478 Choice::Wait => notified.await,
479 Choice::Launch(id) => {
480 let mut guard = LaunchGuard::new(Arc::clone(&self.inner), id);
481 let proxy = if let Some(configuration) = &self.inner.options.proxy {
482 configuration
483 .new_proxy_info(ProxyResolveContext::new())
484 .await
485 .map_err(|error| BrowserError::Launch(anyhow::Error::new(error)))?
486 } else {
487 None
488 };
489 let mut context = LaunchContext::new();
490 context.proxy = proxy;
491 for hook in &self.inner.options.hooks.pre_launch {
492 hook(&mut context);
493 }
494 let launch_options = self.inner.options.launch_options.clone();
495 match self.inner.provider.launch(launch_options, &context).await {
496 Ok(browser) => {
497 let mut browser = Some(browser);
498 let install = {
499 let mut state = self.inner.state.lock().await;
500 if state.shut_down {
501 false
502 } else if let Some(slot) =
503 state.browsers.iter_mut().find(|slot| slot.id == id)
504 {
505 slot.browser = Some(Arc::new(
506 browser.take().expect("launched browser is available"),
507 ));
508 slot.launching = false;
509 slot.proxy = context.proxy;
510 true
511 } else {
512 false
513 }
514 };
515 if !install {
516 if let Err(error) = self
517 .inner
518 .provider
519 .close_browser(
520 browser.take().expect("uninstalled browser is available"),
521 )
522 .await
523 {
524 tracing::warn!(%error, "failed to close browser launched during shutdown");
525 }
526 self.inner
527 .state
528 .lock()
529 .await
530 .browsers
531 .retain(|slot| slot.id != id);
532 guard.disarm();
533 self.inner.capacity.notify_waiters();
534 return Err(BrowserError::Shutdown);
535 }
536 guard.disarm();
537 self.inner.capacity.notify_waiters();
538 }
539 Err(error) => return Err(error),
540 }
541 }
542 Choice::Existing { id, browser, proxy } => {
543 let mut guard = PageReservationGuard::new(Arc::clone(&self.inner), id);
544 let mut prepared_opts = opts.clone();
545 for hook in &self.inner.options.hooks.pre_page_create {
546 hook(&mut prepared_opts);
547 }
548 let page_result = self.inner.provider.new_page(browser.as_ref()).await;
549 drop(browser);
550 let page = match page_result {
551 Ok(page) => page,
552 Err(error) => return Err(error),
553 };
554 guard.set_page(page.clone());
555 for hook in &self.inner.options.hooks.post_page_create {
556 if let Err(error) = hook(&page, &prepared_opts).await {
557 if let Err(close_error) =
558 self.inner.provider.close_page(page.clone()).await
559 {
560 tracing::warn!(%close_error, "failed to close page after hook error");
561 }
562 guard.clear_page();
563 return Err(error);
564 }
565 }
566 let page_id = PageId::next();
567 let registered = {
568 let mut state = self.inner.state.lock().await;
569 if state.shut_down {
570 false
571 } else {
572 if let Some(slot) = state.browsers.iter_mut().find(|slot| slot.id == id)
573 {
574 slot.in_flight_pages = slot.in_flight_pages.saturating_sub(1);
575 }
576 state.pages.insert(
577 page_id,
578 PageEntry {
579 page: page.clone(),
580 browser_id: id,
581 opts: prepared_opts,
582 closing: false,
583 close_notify: Arc::new(Notify::new()),
584 },
585 );
586 true
587 }
588 };
589 if !registered {
590 if let Err(error) = self.inner.provider.close_page(page.clone()).await {
591 tracing::warn!(%error, "failed to close page created during shutdown");
592 }
593 guard.clear_page();
594 return Err(BrowserError::Shutdown);
595 }
596 guard.disarm();
597 let closer: Arc<dyn PoolCloser> = self.inner.clone();
598 return Ok(PageHandle {
599 page: Arc::new(page),
600 proxy,
601 shared: Arc::new(HandleShared {
602 id: page_id,
603 close_state: AtomicU8::new(HANDLE_OPEN),
604 closer: Some(closer),
605 }),
606 });
607 }
608 }
609 }
610 }
611
612 pub async fn shutdown(&self) -> Result<(), BrowserError> {
616 const SHUTDOWN_TIMEOUT: Duration = Duration::from_secs(5);
617
618 self.start_close_worker();
619 {
620 let mut state = self.inner.state.lock().await;
621 state.shut_down = true;
622 }
623 self.inner.capacity.notify_waiters();
624
625 let deadline = tokio::time::Instant::now() + SHUTDOWN_TIMEOUT;
626 let mut first_error = None;
627
628 let (barrier_tx, barrier_rx) = oneshot::channel();
629 if self
630 .inner
631 .close_tx
632 .send(CloseCommand::Barrier(barrier_tx))
633 .is_ok()
634 && tokio::time::timeout_at(deadline, barrier_rx).await.is_err()
635 {
636 first_error = Some(shutdown_timeout_error());
637 abort_close_worker(&self.inner).await;
638 }
639
640 loop {
641 let notified = self.inner.capacity.notified();
642 tokio::pin!(notified);
643 notified.as_mut().enable();
644 let has_in_flight = {
645 let state = self.inner.state.lock().await;
646 state
647 .browsers
648 .iter()
649 .any(|slot| slot.launching || slot.in_flight_pages != 0)
650 };
651 if !has_in_flight {
652 break;
653 }
654 if tokio::time::timeout_at(deadline, notified).await.is_err() {
655 if first_error.is_none() {
656 first_error = Some(shutdown_timeout_error());
657 }
658 break;
659 }
660 }
661
662 let page_ids = {
663 let state = self.inner.state.lock().await;
664 state.pages.keys().copied().collect::<Vec<_>>()
665 };
666
667 let mut close_tasks = Vec::with_capacity(page_ids.len());
668 for page_id in page_ids {
669 let inner = Arc::clone(&self.inner);
670 close_tasks.push(tokio::spawn(async move {
671 close_page_by_id(&inner, page_id).await
672 }));
673 }
674 for mut task in close_tasks {
675 match tokio::time::timeout_at(deadline, &mut task).await {
676 Ok(Ok(Ok(()))) => {}
677 Ok(Ok(Err(error))) => {
678 if first_error.is_none() {
679 first_error = Some(error);
680 }
681 }
682 Ok(Err(error)) => {
683 if first_error.is_none() {
684 first_error = Some(close_task_error(error));
685 }
686 }
687 Err(_) => {
688 if first_error.is_none() {
689 first_error = Some(shutdown_timeout_error());
690 }
691 }
692 }
693 }
694
695 let browsers = {
696 let mut state = self.inner.state.lock().await;
697 state
698 .browsers
699 .iter_mut()
700 .filter(|slot| slot.open_pages == 0 && slot.in_flight_pages == 0)
701 .filter_map(|slot| slot.browser.take())
702 .collect::<Vec<_>>()
703 };
704 let mut browser_tasks = Vec::with_capacity(browsers.len());
705 for browser in browsers {
706 let inner = Arc::clone(&self.inner);
707 browser_tasks.push(tokio::spawn(async move {
708 close_browser_arc(inner.as_ref(), browser).await;
709 }));
710 }
711 for mut task in browser_tasks {
712 if tokio::time::timeout_at(deadline, &mut task).await.is_err() && first_error.is_none()
713 {
714 first_error = Some(shutdown_timeout_error());
715 }
716 }
717
718 let (stop_tx, stop_rx) = oneshot::channel();
719 let stop_sent = self
720 .inner
721 .close_tx
722 .send(CloseCommand::Stop(stop_tx))
723 .is_ok();
724 if stop_sent
725 && tokio::time::timeout_at(deadline, stop_rx).await.is_err()
726 && first_error.is_none()
727 {
728 first_error = Some(shutdown_timeout_error());
729 }
730 let worker = self
731 .inner
732 .worker
733 .lock()
734 .unwrap_or_else(|error| error.into_inner())
735 .take();
736 if let Some(mut worker) = worker {
737 if tokio::time::timeout_at(deadline, &mut worker)
738 .await
739 .is_err()
740 {
741 worker.abort();
742 let _ = worker.await;
743 }
744 }
745
746 if let Some(error) = first_error {
747 Err(error)
748 } else {
749 Ok(())
750 }
751 }
752}
753
754async fn close_page_by_id<P: BrowserProvider>(
755 inner: &Arc<PoolInner<P>>,
756 id: PageId,
757) -> Result<(), BrowserError> {
758 close_page(Arc::clone(inner), id).await
759}
760
761async fn close_page<P: BrowserProvider>(
762 inner: Arc<PoolInner<P>>,
763 id: PageId,
764) -> Result<(), BrowserError> {
765 let entry: PageEntry<P> = loop {
766 let wait = {
767 let mut state = inner.state.lock().await;
768 let Some(entry) = state.pages.get_mut(&id) else {
769 return Ok(());
770 };
771 if entry.closing {
772 let mut notified = Box::pin(Arc::clone(&entry.close_notify).notified_owned());
773 notified.as_mut().enable();
774 Some(notified)
775 } else {
776 entry.closing = true;
777 break PageEntry {
778 page: entry.page.clone(),
779 browser_id: entry.browser_id,
780 opts: entry.opts.clone(),
781 closing: true,
782 close_notify: Arc::clone(&entry.close_notify),
783 };
784 }
785 };
786 if let Some(wait) = wait {
787 wait.await;
788 }
789 };
790 let mut close_guard = PageCloseGuard {
791 inner: Arc::clone(&inner),
792 id,
793 notify: Arc::clone(&entry.close_notify),
794 armed: true,
795 };
796
797 for hook in &inner.options.hooks.pre_page_close {
798 if let Err(error) = hook(&entry.page, &entry.opts).await {
799 tracing::warn!(page_id = %id, %error, "pre-page-close hook failed");
800 }
801 }
802 let close_result = inner.provider.close_page(entry.page).await;
803
804 let browser = {
805 let mut state = inner.state.lock().await;
806 state.pages.remove(&id);
807 let shut_down = state.shut_down;
808 let mut browser = None;
809 if let Some(slot) = state
810 .browsers
811 .iter_mut()
812 .find(|slot| slot.id == entry.browser_id)
813 {
814 slot.open_pages = slot.open_pages.saturating_sub(1);
815 if slot.pages_created >= inner.options.retire_browser_after_page_count {
816 slot.retired = true;
817 }
818 if (slot.retired || shut_down) && slot.open_pages == 0 {
819 browser = slot.browser.take();
820 }
821 }
822 browser
823 };
824 inner.capacity.notify_waiters();
825 close_guard.armed = false;
826 entry.close_notify.notify_waiters();
827 for hook in &inner.options.hooks.post_page_close {
828 hook(id);
829 }
830 if let Some(browser) = browser {
831 close_browser_arc(inner.as_ref(), browser).await;
832 }
833 close_result
834}
835
836struct PageCloseGuard<P: BrowserProvider> {
837 inner: Arc<PoolInner<P>>,
838 id: PageId,
839 notify: Arc<Notify>,
840 armed: bool,
841}
842
843impl<P: BrowserProvider> Drop for PageCloseGuard<P> {
844 fn drop(&mut self) {
845 if !self.armed {
846 return;
847 }
848 let inner = Arc::clone(&self.inner);
849 let id = self.id;
850 let notify = Arc::clone(&self.notify);
851 if let Ok(mut state) = inner.state.try_lock() {
852 if let Some(entry) = state.pages.get_mut(&id) {
853 entry.closing = false;
854 }
855 drop(state);
856 notify.notify_waiters();
857 } else if let Ok(runtime) = tokio::runtime::Handle::try_current() {
858 runtime.spawn(async move {
859 if let Some(entry) = inner.state.lock().await.pages.get_mut(&id) {
860 entry.closing = false;
861 }
862 notify.notify_waiters();
863 });
864 } else {
865 tracing::warn!(page_id = %id, "cancelled page close dropped outside a Tokio runtime");
866 }
867 }
868}
869
870fn shutdown_timeout_error() -> BrowserError {
871 BrowserError::PageCreate(anyhow::anyhow!(
872 "browser pool shutdown timed out after 5 seconds"
873 ))
874}
875
876fn close_task_error(error: tokio::task::JoinError) -> BrowserError {
877 BrowserError::PageCreate(anyhow::anyhow!("page close task failed: {error}"))
878}
879
880async fn abort_close_worker<P: BrowserProvider>(inner: &Arc<PoolInner<P>>) {
881 let worker = inner
882 .worker
883 .lock()
884 .unwrap_or_else(|error| error.into_inner())
885 .take();
886 if let Some(worker) = worker {
887 worker.abort();
888 let _ = worker.await;
889 }
890}
891
892async fn finalize_orphaned_pool<P: BrowserProvider>(inner: Arc<PoolInner<P>>) {
893 if Arc::strong_count(&inner) != 1 {
894 return;
895 }
896 let browsers = {
897 let mut state = inner.state.lock().await;
898 if !state.pages.is_empty() {
899 return;
900 }
901 state
902 .browsers
903 .iter_mut()
904 .filter_map(|slot| slot.browser.take())
905 .collect::<Vec<_>>()
906 };
907 for browser in browsers {
908 close_browser_arc(inner.as_ref(), browser).await;
909 }
910}
911
912async fn close_browser_arc<P: BrowserProvider>(inner: &PoolInner<P>, mut browser: Arc<P::Browser>) {
913 for _ in 0..16 {
914 match Arc::try_unwrap(browser) {
915 Ok(browser) => {
916 if let Err(error) = inner.provider.close_browser(browser).await {
917 tracing::warn!(%error, "failed to close retired browser");
918 }
919 return;
920 }
921 Err(still_shared) => {
922 browser = still_shared;
923 tokio::task::yield_now().await;
924 }
925 }
926 }
927 tracing::warn!(
928 outstanding_clones = Arc::strong_count(&browser) - 1,
929 "browser could not be explicitly closed because handles are still outstanding"
930 );
931}
932
933#[async_trait::async_trait]
934trait PoolCloser: Send + Sync {
935 async fn close_now(self: Arc<Self>, id: PageId) -> Result<(), BrowserError>;
936 fn enqueue_close(self: Arc<Self>, id: PageId);
937 fn enqueue_finalize(self: Arc<Self>);
938}
939
940#[async_trait::async_trait]
941impl<P: BrowserProvider> PoolCloser for PoolInner<P> {
942 async fn close_now(self: Arc<Self>, id: PageId) -> Result<(), BrowserError> {
943 match tokio::spawn(close_page(self, id)).await {
944 Ok(result) => result,
945 Err(error) => Err(close_task_error(error)),
946 }
947 }
948
949 fn enqueue_close(self: Arc<Self>, id: PageId) {
950 let _ = self.close_tx.send(CloseCommand::Page {
951 inner: Arc::clone(&self),
952 id,
953 });
954 }
955
956 fn enqueue_finalize(self: Arc<Self>) {
957 let _ = self.close_tx.send(CloseCommand::Finalize {
958 inner: Arc::clone(&self),
959 });
960 }
961}
962
963const HANDLE_OPEN: u8 = 0;
964const HANDLE_CLOSING: u8 = 1;
965const HANDLE_CLOSED: u8 = 2;
966
967struct HandleShared {
968 id: PageId,
969 close_state: AtomicU8,
970 closer: Option<Arc<dyn PoolCloser>>,
973}
974
975impl Drop for HandleShared {
976 fn drop(&mut self) {
977 let closer = self
978 .closer
979 .take()
980 .expect("handle closer is present until HandleShared::drop");
981 if self.close_state.load(Ordering::SeqCst) != HANDLE_CLOSED {
982 tracing::warn!(
983 page_id = %self.id,
984 "PageHandle dropped without close(); scheduling background close — prefer page.close().await"
985 );
986 Arc::clone(&closer).enqueue_close(self.id);
987 }
988 closer.enqueue_finalize();
989 }
990}
991
992#[derive(Clone)]
998pub struct PageHandle {
999 page: Arc<dyn BrowserPage>,
1000 proxy: Option<ProxyInfo>,
1001 shared: Arc<HandleShared>,
1002}
1003
1004impl PageHandle {
1005 pub fn id(&self) -> PageId {
1007 self.shared.id
1008 }
1009
1010 pub fn proxy_info(&self) -> Option<&ProxyInfo> {
1012 self.proxy.as_ref()
1013 }
1014
1015 pub async fn close(&self) -> Result<(), BrowserError> {
1019 if self
1020 .shared
1021 .close_state
1022 .compare_exchange(
1023 HANDLE_OPEN,
1024 HANDLE_CLOSING,
1025 Ordering::SeqCst,
1026 Ordering::SeqCst,
1027 )
1028 .is_err()
1029 {
1030 return Ok(());
1031 }
1032 let mut guard = HandleCloseGuard {
1033 shared: &self.shared,
1034 armed: true,
1035 };
1036 let result = Arc::clone(
1037 self.shared
1038 .closer
1039 .as_ref()
1040 .expect("handle closer is present while PageHandle exists"),
1041 )
1042 .close_now(self.shared.id)
1043 .await;
1044 if result.is_ok() {
1045 self.shared
1046 .close_state
1047 .store(HANDLE_CLOSED, Ordering::SeqCst);
1048 guard.armed = false;
1049 }
1050 result
1051 }
1052}
1053
1054impl Deref for PageHandle {
1055 type Target = dyn BrowserPage;
1056
1057 fn deref(&self) -> &Self::Target {
1058 self.page.as_ref()
1059 }
1060}
1061
1062impl fmt::Debug for PageHandle {
1063 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
1064 formatter
1065 .debug_struct("PageHandle")
1066 .field("id", &self.id())
1067 .field("proxy", &self.proxy)
1068 .field(
1069 "closed",
1070 &(self.shared.close_state.load(Ordering::SeqCst) == HANDLE_CLOSED),
1071 )
1072 .finish_non_exhaustive()
1073 }
1074}
1075
1076struct HandleCloseGuard<'a> {
1077 shared: &'a HandleShared,
1078 armed: bool,
1079}
1080
1081impl Drop for HandleCloseGuard<'_> {
1082 fn drop(&mut self) {
1083 if self.armed {
1084 self.shared.close_state.store(HANDLE_OPEN, Ordering::SeqCst);
1085 }
1086 }
1087}