Skip to main content

millipede_browser/
pool.rs

1//! Concurrent browser and page lifecycle management.
2
3use 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/// Configuration for a [`BrowserPool`].
24#[non_exhaustive]
25#[must_use = "browser pool options do nothing unless passed to BrowserPool::new"]
26pub struct BrowserPoolOptions<L> {
27    /// Maximum number of simultaneously open pages in one browser.
28    pub max_open_pages_per_browser: usize,
29    /// Number of created pages after which a browser is retired.
30    pub retire_browser_after_page_count: u64,
31    /// Maximum number of live or launching browsers, or `None` for no limit.
32    pub max_browsers: Option<usize>,
33    /// Maximum time to wait for a page, including browser launch and hook work.
34    pub page_acquire_timeout: Duration,
35    /// Provider-specific options cloned for every browser launch.
36    pub launch_options: L,
37    /// Optional proxy configuration resolved once per browser launch.
38    pub proxy: Option<ProxyConfiguration>,
39    /// Browser lifecycle hooks.
40    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    /// Creates options with the standard pool limits and default provider launch options.
59    pub fn new() -> Self {
60        Self::default()
61    }
62}
63
64impl<L> BrowserPoolOptions<L> {
65    /// Sets the maximum number of simultaneously open pages in one browser.
66    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    /// Sets the page count after which a browser is retired.
72    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    /// Sets the maximum number of live or launching browsers.
78    pub fn with_max_browsers(mut self, value: Option<usize>) -> Self {
79        self.max_browsers = value;
80        self
81    }
82
83    /// Sets the maximum duration allowed for page acquisition.
84    pub fn with_page_acquire_timeout(mut self, value: Duration) -> Self {
85        self.page_acquire_timeout = value;
86        self
87    }
88
89    /// Replaces the provider-specific browser launch options.
90    pub fn with_launch_options(mut self, value: L) -> Self {
91        self.launch_options = value;
92        self
93    }
94
95    /// Sets the proxy configuration used for browser launches.
96    pub fn with_proxy(mut self, value: Option<ProxyConfiguration>) -> Self {
97        self.proxy = value;
98        self
99    }
100
101    /// Replaces the browser lifecycle hooks.
102    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
129/// A lazily launched pool of provider browsers and their pages.
130pub 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    /// Creates an empty pool without requiring a Tokio runtime.
323    ///
324    /// Browser launch and the fallback close worker are both lazy: the first call to
325    /// [`Self::new_page`] starts them from within the caller's runtime. This intentionally differs
326    /// from the asynchronous constructor sketched in INTERFACE §12.
327    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    /// Acquires a page, launching or waiting for browser capacity as needed.
395    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    /// Closes all pages and browsers and rejects future acquisitions.
613    ///
614    /// Calling shutdown more than once is safe.
615    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    /// A strong closer is cloned into every queued command, so the worker can finish cleanup even
971    /// after the owning `BrowserPool` and the last page handle have both been dropped.
972    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/// Provider-erased RAII handle for a page checked out from a [`BrowserPool`].
993///
994/// `Clone` is required by the engine's `Context: Clone` contract. All clones share an atomic
995/// close guard, guaranteeing at-most-once cleanup: an explicit [`Self::close`] wins, while dropping
996/// the last clone schedules background cleanup only when explicit close was never called.
997#[derive(Clone)]
998pub struct PageHandle {
999    page: Arc<dyn BrowserPage>,
1000    proxy: Option<ProxyInfo>,
1001    shared: Arc<HandleShared>,
1002}
1003
1004impl PageHandle {
1005    /// Returns this page's stable process-local identifier.
1006    pub fn id(&self) -> PageId {
1007        self.shared.id
1008    }
1009
1010    /// Returns the proxy resolved for the browser that owns this page.
1011    pub fn proxy_info(&self) -> Option<&ProxyInfo> {
1012        self.proxy.as_ref()
1013    }
1014
1015    /// Explicitly closes the page.
1016    ///
1017    /// This is idempotent across every clone of the handle.
1018    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}