Skip to main content

browser_commander/puppeteer/
bridge.rs

1//! The client side of `browser-commander serve --stdio` (issue #108).
2//!
3//! The JavaScript CLI serves Puppeteer's live objects as remote handles over
4//! JSON-RPC 2.0, one message per line (docs/cli-and-bridge.md). This module
5//! starts that server through command-stream, matches responses to requests,
6//! routes `events.emit` notifications to subscriptions and converts values
7//! between Rust and the bridge's value encoding. The generated wrappers in
8//! [`super::api`] are thin typed calls on top of [`RemoteHandle`].
9
10use std::collections::{HashMap, VecDeque};
11use std::fmt;
12use std::path::{Path, PathBuf};
13use std::sync::atomic::{AtomicU64, Ordering};
14use std::sync::{Arc, Mutex as StdMutex};
15use std::time::Duration;
16
17use base64::Engine as _;
18use command_stream::{quote::quote, ProcessRunner, RunOptions, StdinOption};
19use serde_json::{json, Map, Value};
20use thiserror::Error;
21use tokio::io::{AsyncBufReadExt, AsyncRead, AsyncWrite, AsyncWriteExt, BufReader};
22use tokio::sync::{mpsc, oneshot, Mutex};
23
24use crate::utilities::subprocess::kill_owned_process_tree;
25
26/// Path of the JavaScript CLI (`js/bin/browser-commander.js`), when it is not
27/// found next to the crate or in `node_modules`.
28pub const JS_CLI_ENV: &str = "BROWSER_COMMANDER_JS_CLI";
29
30const STDERR_LINES: usize = 50;
31const EXIT_GRACE: Duration = Duration::from_secs(5);
32
33/// Errors from the bridge.
34#[derive(Debug, Error)]
35pub enum BridgeError {
36    /// The server, or Puppeteer behind it, rejected a call.
37    #[error("{name}: {message}")]
38    Remote {
39        /// JSON-RPC error code (`-32000` for errors thrown by Puppeteer).
40        code: i64,
41        /// Error class (`TimeoutError`, `TypeError`, …) or `RpcError`.
42        name: String,
43        /// Error message.
44        message: String,
45        /// Server-side stack, when it sent one.
46        stack: Option<String>,
47    },
48    /// The server went away, or the bridge was closed.
49    #[error("serve --stdio closed: {0}")]
50    Closed(String),
51    /// Reading or writing the pipe failed.
52    #[error("serve --stdio I/O error: {0}")]
53    Io(#[from] std::io::Error),
54    /// A message was not valid JSON.
55    #[error("serve --stdio sent invalid JSON: {0}")]
56    Json(#[from] serde_json::Error),
57    /// A value did not have the type the wrapper declares.
58    #[error("expected {expected} from the bridge, got {value}")]
59    Decode {
60        /// The declared type.
61        expected: &'static str,
62        /// What arrived.
63        value: Value,
64    },
65    /// Node.js or the JavaScript CLI could not be found or started.
66    #[error("serve --stdio unavailable: {0}")]
67    Unavailable(String),
68}
69
70impl BridgeError {
71    /// Whether Puppeteer reported a timeout.
72    pub fn is_timeout(&self) -> bool {
73        matches!(self, BridgeError::Remote { name, .. } if name == "TimeoutError")
74    }
75
76    fn decode(expected: &'static str, value: Value) -> Self {
77        BridgeError::Decode { expected, value }
78    }
79}
80
81fn remote_error(error: &Value) -> BridgeError {
82    let data = error.get("data");
83    let text = |value: Option<&Value>| value.and_then(Value::as_str).map(str::to_string);
84    BridgeError::Remote {
85        code: error.get("code").and_then(Value::as_i64).unwrap_or(-32000),
86        name: text(data.and_then(|d| d.get("name"))).unwrap_or_else(|| "RpcError".into()),
87        message: text(error.get("message")).unwrap_or_default(),
88        stack: text(data.and_then(|d| d.get("stack"))),
89    }
90}
91
92/// A caller waiting for its response; `subscribe` marks `events.subscribe`.
93struct Waiter {
94    sender: oneshot::Sender<Result<Value, BridgeError>>,
95    subscribe: bool,
96}
97
98type Pending = HashMap<u64, Waiter>;
99type Subscribers = HashMap<String, mpsc::UnboundedSender<Vec<Value>>>;
100type Receivers = HashMap<String, mpsc::UnboundedReceiver<Vec<Value>>>;
101
102struct Inner {
103    writer: Mutex<Box<dyn AsyncWrite + Send + Unpin>>,
104    next_id: AtomicU64,
105    pending: StdMutex<Pending>,
106    subscribers: StdMutex<Subscribers>,
107    /// Channels of subscriptions whose `events.subscribe` has been answered
108    /// but not yet picked up by [`RemoteHandle::subscribe`].
109    receivers: StdMutex<Receivers>,
110    closed: StdMutex<Option<String>>,
111    reader: StdMutex<Option<tokio::task::JoinHandle<()>>>,
112}
113
114impl Drop for Inner {
115    fn drop(&mut self) {
116        if let Some(task) = self.reader.get_mut().ok().and_then(Option::take) {
117            task.abort();
118        }
119    }
120}
121
122/// A JSON-RPC conversation with one `serve --stdio` server.
123///
124/// Requests are pipelined: any number can be in flight, and each call waits
125/// for the response with its own id.
126#[derive(Clone)]
127pub struct BridgeClient {
128    inner: Arc<Inner>,
129}
130
131impl fmt::Debug for BridgeClient {
132    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
133        f.debug_struct("BridgeClient")
134            .field("closed", &self.close_reason())
135            .finish()
136    }
137}
138
139impl BridgeClient {
140    /// Talk to a server over its stdout (`reader`) and stdin (`writer`).
141    /// Lines are read on a background task.
142    pub fn new<R, W>(reader: R, writer: W) -> Self
143    where
144        R: AsyncRead + Send + Unpin + 'static,
145        W: AsyncWrite + Send + Unpin + 'static,
146    {
147        let inner = Arc::new(Inner {
148            writer: Mutex::new(Box::new(writer)),
149            next_id: AtomicU64::new(1),
150            pending: StdMutex::new(HashMap::new()),
151            subscribers: StdMutex::new(HashMap::new()),
152            receivers: StdMutex::new(HashMap::new()),
153            closed: StdMutex::new(None),
154            reader: StdMutex::new(None),
155        });
156        let weak = Arc::downgrade(&inner);
157        let task = tokio::spawn(async move {
158            let mut lines = BufReader::new(reader).lines();
159            let reason = loop {
160                let line = match lines.next_line().await {
161                    Ok(Some(line)) => line,
162                    Ok(None) => break "the server closed its output".to_string(),
163                    Err(err) => break format!("reading from the server failed: {err}"),
164                };
165                let Some(inner) = weak.upgrade() else {
166                    return;
167                };
168                let client = BridgeClient { inner };
169                match serde_json::from_str::<Value>(&line) {
170                    Ok(message) => client.dispatch(message),
171                    Err(err) => {
172                        tracing::debug!(target: "browser_commander::puppeteer", "ignored line {line:?}: {err}");
173                    }
174                }
175            };
176            if let Some(inner) = weak.upgrade() {
177                BridgeClient { inner }.mark_closed(reason);
178            }
179        });
180        if let Ok(mut slot) = inner.reader.lock() {
181            *slot = Some(task);
182        }
183        Self { inner }
184    }
185
186    /// Send a request and wait for its result.
187    pub async fn request(&self, method: &str, params: Value) -> Result<Value, BridgeError> {
188        self.send_request(method, params, false).await
189    }
190
191    async fn send_request(
192        &self,
193        method: &str,
194        params: Value,
195        subscribe: bool,
196    ) -> Result<Value, BridgeError> {
197        if let Some(reason) = self.close_reason() {
198            return Err(BridgeError::Closed(reason));
199        }
200        let id = self.inner.next_id.fetch_add(1, Ordering::Relaxed);
201        let (sender, receiver) = oneshot::channel();
202        self.lock_pending().insert(id, Waiter { sender, subscribe });
203        let mut line = serde_json::to_vec(&json!({
204            "jsonrpc": "2.0",
205            "id": id,
206            "method": method,
207            "params": params,
208        }))?;
209        line.push(b'\n');
210        let written = {
211            let mut writer = self.inner.writer.lock().await;
212            match writer.write_all(&line).await {
213                Ok(()) => writer.flush().await,
214                Err(err) => Err(err),
215            }
216        };
217        if let Err(err) = written {
218            self.lock_pending().remove(&id);
219            return Err(err.into());
220        }
221        match receiver.await {
222            Ok(result) => result,
223            Err(_) => Err(BridgeError::Closed(
224                self.close_reason()
225                    .unwrap_or_else(|| "the response was dropped".to_string()),
226            )),
227        }
228    }
229
230    /// `handle.root`: the engine's entry object (`puppeteer`, `playwright`).
231    pub async fn root(&self, name: &str) -> Result<RemoteHandle, BridgeError> {
232        let value = self.request("handle.root", json!({ "name": name })).await?;
233        RemoteHandle::from_wire(self, value)
234    }
235
236    /// Why the bridge closed, if it has.
237    pub fn close_reason(&self) -> Option<String> {
238        self.inner
239            .closed
240            .lock()
241            .ok()
242            .and_then(|reason| reason.clone())
243    }
244
245    /// End the server's input. `serve --stdio` finishes in-flight requests,
246    /// closes its sessions and exits; later calls fail with
247    /// [`BridgeError::Closed`].
248    pub async fn close_input(&self) {
249        let mut writer = self.inner.writer.lock().await;
250        let _ = writer.shutdown().await;
251        // Shutting a child's stdin down only flushes it; the pipe closes when
252        // the handle is dropped.
253        *writer = Box::new(tokio::io::sink());
254        drop(writer);
255        if let Ok(mut closed) = self.inner.closed.lock() {
256            closed.get_or_insert_with(|| "the bridge was closed".to_string());
257        }
258    }
259
260    fn lock_pending(&self) -> std::sync::MutexGuard<'_, Pending> {
261        self.inner
262            .pending
263            .lock()
264            .unwrap_or_else(std::sync::PoisonError::into_inner)
265    }
266
267    fn lock_subscribers(&self) -> std::sync::MutexGuard<'_, Subscribers> {
268        self.inner
269            .subscribers
270            .lock()
271            .unwrap_or_else(std::sync::PoisonError::into_inner)
272    }
273
274    fn lock_receivers(&self) -> std::sync::MutexGuard<'_, Receivers> {
275        self.inner
276            .receivers
277            .lock()
278            .unwrap_or_else(std::sync::PoisonError::into_inner)
279    }
280
281    fn mark_closed(&self, reason: String) {
282        if let Ok(mut closed) = self.inner.closed.lock() {
283            closed.get_or_insert(reason.clone());
284        }
285        let pending: Vec<_> = self.lock_pending().drain().collect();
286        for (_, waiter) in pending {
287            let _ = waiter.sender.send(Err(BridgeError::Closed(reason.clone())));
288        }
289        self.lock_subscribers().clear();
290        self.lock_receivers().clear();
291    }
292
293    /// Route one message: a response to its caller, an `events.emit`
294    /// notification to its subscription.
295    fn dispatch(&self, message: Value) {
296        if let Some(id) = message.get("id").and_then(Value::as_u64) {
297            let Some(waiter) = self.lock_pending().remove(&id) else {
298                return;
299            };
300            let outcome = match message.get("error") {
301                Some(error) => Err(remote_error(error)),
302                None => Ok(message.get("result").cloned().unwrap_or(Value::Null)),
303            };
304            // The server may emit right after answering, so the channel must
305            // exist before the reader handles the next line.
306            if let (true, Ok(result)) = (waiter.subscribe, &outcome) {
307                if let Some(subscription) = result.get("subscription").and_then(Value::as_str) {
308                    let (sender, receiver) = mpsc::unbounded_channel();
309                    self.lock_subscribers()
310                        .insert(subscription.to_string(), sender);
311                    self.lock_receivers()
312                        .insert(subscription.to_string(), receiver);
313                }
314            }
315            let _ = waiter.sender.send(outcome);
316            return;
317        }
318        if message.get("method").and_then(Value::as_str) != Some("events.emit") {
319            return;
320        }
321        let params = message.get("params").cloned().unwrap_or(Value::Null);
322        let Some(subscription) = params.get("subscription").and_then(Value::as_str) else {
323            return;
324        };
325        let args = match params.get("args") {
326            Some(Value::Array(args)) => args.clone(),
327            _ => Vec::new(),
328        };
329        let mut subscribers = self.lock_subscribers();
330        if let Some(sender) = subscribers.get(subscription) {
331            if sender.send(args).is_err() {
332                subscribers.remove(subscription);
333            }
334        }
335    }
336}
337
338/// A Puppeteer object that lives in the server, such as a `Page`.
339///
340/// Handles are cheap to clone; every clone names the same object. The
341/// generated wrappers hold one each.
342#[derive(Clone)]
343pub struct RemoteHandle {
344    client: BridgeClient,
345    id: Arc<str>,
346    type_name: Arc<str>,
347}
348
349impl fmt::Debug for RemoteHandle {
350    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
351        f.debug_struct("RemoteHandle")
352            .field("id", &self.id)
353            .field("type", &self.type_name)
354            .finish()
355    }
356}
357
358impl PartialEq for RemoteHandle {
359    fn eq(&self, other: &Self) -> bool {
360        Arc::ptr_eq(&self.client.inner, &other.client.inner) && self.id == other.id
361    }
362}
363
364impl RemoteHandle {
365    /// The server's id for the object (`h1`, `h2`, …).
366    pub fn id(&self) -> &str {
367        &self.id
368    }
369
370    /// The object's runtime class as the server reports it (`CdpPage`).
371    pub fn type_name(&self) -> &str {
372        &self.type_name
373    }
374
375    /// The bridge the object lives in.
376    pub fn client(&self) -> &BridgeClient {
377        &self.client
378    }
379
380    /// The handle as an argument: `{"$handle": "h1"}`.
381    pub fn to_wire(&self) -> Value {
382        json!({ "$handle": &*self.id })
383    }
384
385    /// `handle.call`: call a method. `None` arguments are sent as
386    /// `undefined`; trailing ones are left out, as JavaScript does.
387    pub async fn call<T: FromWire>(
388        &self,
389        method: &str,
390        args: Vec<Option<Value>>,
391    ) -> Result<T, BridgeError> {
392        let value = self
393            .client
394            .request(
395                "handle.call",
396                json!({ "handle": &*self.id, "method": method, "args": encode_args(args) }),
397            )
398            .await?;
399        T::from_wire(&self.client, value)
400    }
401
402    /// `handle.get`: read a property (awaited when it is a promise).
403    pub async fn get<T: FromWire>(&self, property: &str) -> Result<T, BridgeError> {
404        let value = self
405            .client
406            .request(
407                "handle.get",
408                json!({ "handle": &*self.id, "property": property }),
409            )
410            .await?;
411        T::from_wire(&self.client, value)
412    }
413
414    /// `handle.describe`: the object's type and member names.
415    pub async fn describe(&self) -> Result<Value, BridgeError> {
416        self.client
417            .request("handle.describe", json!({ "handle": &*self.id }))
418            .await
419    }
420
421    /// `handle.dispose`: forget the handle on the server. The object itself
422    /// stays alive; dispose a JS handle with its own `dispose()` method.
423    pub async fn release(&self) -> Result<(), BridgeError> {
424        self.client
425            .request("handle.dispose", json!({ "handle": &*self.id }))
426            .await
427            .map(|_| ())
428    }
429
430    /// `events.subscribe`: receive the arguments of every `event` the object
431    /// emits until the [`Subscription`] is closed or dropped.
432    pub async fn subscribe(&self, event: &str) -> Result<Subscription, BridgeError> {
433        let result = self
434            .client
435            .send_request(
436                "events.subscribe",
437                json!({ "handle": &*self.id, "event": event }),
438                true,
439            )
440            .await?;
441        let id = result
442            .get("subscription")
443            .and_then(Value::as_str)
444            .ok_or_else(|| BridgeError::decode("a subscription id", result.clone()))?
445            .to_string();
446        let receiver = self.client.lock_receivers().remove(&id).ok_or_else(|| {
447            BridgeError::Closed(
448                self.client
449                    .close_reason()
450                    .unwrap_or_else(|| "the subscription was dropped".to_string()),
451            )
452        })?;
453        Ok(Subscription {
454            client: self.client.clone(),
455            id,
456            receiver,
457        })
458    }
459}
460
461/// Arguments for `handle.call`: `None` is `undefined`, trailing ones dropped.
462fn encode_args(mut args: Vec<Option<Value>>) -> Value {
463    while matches!(args.last(), Some(None)) {
464        args.pop();
465    }
466    Value::Array(
467        args.into_iter()
468            .map(|arg| arg.unwrap_or_else(|| json!({ "$undefined": true })))
469            .collect(),
470    )
471}
472
473/// Events from one `events.subscribe`.
474#[derive(Debug)]
475pub struct Subscription {
476    client: BridgeClient,
477    id: String,
478    receiver: mpsc::UnboundedReceiver<Vec<Value>>,
479}
480
481impl Subscription {
482    /// The server's subscription id.
483    pub fn id(&self) -> &str {
484        &self.id
485    }
486
487    /// The arguments of the next event, in the bridge's value encoding
488    /// ([`decode_handle`] turns `{"$handle": …}` into a typed wrapper).
489    /// `None` once the subscription or the bridge is closed.
490    pub async fn next(&mut self) -> Option<Vec<Value>> {
491        self.receiver.recv().await
492    }
493
494    /// `events.unsubscribe`.
495    pub async fn close(self) -> Result<(), BridgeError> {
496        self.client.lock_subscribers().remove(&self.id);
497        self.client
498            .request("events.unsubscribe", json!({ "subscription": &self.id }))
499            .await
500            .map(|_| ())
501    }
502}
503
504impl Drop for Subscription {
505    fn drop(&mut self) {
506        // Events stop being delivered; the server-side listener stays until
507        // `close` or the end of the bridge.
508        self.client.lock_subscribers().remove(&self.id);
509    }
510}
511
512/// A function argument: `$function` source compiled on the server, or a
513/// string, which Puppeteer evaluates as an expression (or treats as a
514/// selector where one is expected).
515#[derive(Debug, Clone, PartialEq)]
516pub struct JsFunction(Value);
517
518impl JsFunction {
519    /// A function from its source, such as `"(a, b) => a + b"`.
520    pub fn source(source: impl Into<String>) -> Self {
521        Self(json!({ "$function": source.into() }))
522    }
523
524    /// A string passed as is: an expression for `evaluate`, a selector for
525    /// `locator`.
526    pub fn text(text: impl Into<String>) -> Self {
527        Self(Value::String(text.into()))
528    }
529
530    /// The encoded argument.
531    pub fn to_wire(&self) -> Value {
532        self.0.clone()
533    }
534}
535
536impl From<&str> for JsFunction {
537    fn from(text: &str) -> Self {
538        Self::text(text)
539    }
540}
541
542impl From<String> for JsFunction {
543    fn from(text: String) -> Self {
544        Self::text(text)
545    }
546}
547
548/// Bytes as an argument (`$binary`).
549pub fn binary(bytes: &[u8]) -> Value {
550    json!({ "$binary": base64::engine::general_purpose::STANDARD.encode(bytes) })
551}
552
553/// A value as it arrives from the bridge, converted to a declared type.
554pub trait FromWire: Sized {
555    /// Convert, failing with [`BridgeError::Decode`] on a type mismatch.
556    fn from_wire(client: &BridgeClient, value: Value) -> Result<Self, BridgeError>;
557}
558
559fn is_undefined(value: &Value) -> bool {
560    value.get("$undefined").is_some()
561}
562
563impl FromWire for () {
564    fn from_wire(_: &BridgeClient, _: Value) -> Result<Self, BridgeError> {
565        Ok(())
566    }
567}
568
569impl FromWire for Value {
570    fn from_wire(_: &BridgeClient, value: Value) -> Result<Self, BridgeError> {
571        Ok(value)
572    }
573}
574
575impl FromWire for String {
576    fn from_wire(_: &BridgeClient, value: Value) -> Result<Self, BridgeError> {
577        match value {
578            Value::String(text) => Ok(text),
579            other => Err(BridgeError::decode("a string", other)),
580        }
581    }
582}
583
584impl FromWire for f64 {
585    fn from_wire(_: &BridgeClient, value: Value) -> Result<Self, BridgeError> {
586        match &value {
587            Value::Number(number) => number
588                .as_f64()
589                .ok_or_else(|| BridgeError::decode("a number", value.clone())),
590            // The bridge sends NaN and the infinities as null.
591            Value::Null => Ok(f64::NAN),
592            _ => Err(BridgeError::decode("a number", value)),
593        }
594    }
595}
596
597impl FromWire for bool {
598    fn from_wire(_: &BridgeClient, value: Value) -> Result<Self, BridgeError> {
599        value
600            .as_bool()
601            .ok_or_else(|| BridgeError::decode("a boolean", value))
602    }
603}
604
605impl FromWire for Vec<u8> {
606    fn from_wire(_: &BridgeClient, value: Value) -> Result<Self, BridgeError> {
607        value
608            .get("$binary")
609            .and_then(Value::as_str)
610            .and_then(|data| base64::engine::general_purpose::STANDARD.decode(data).ok())
611            .ok_or_else(|| BridgeError::decode("bytes", value))
612    }
613}
614
615impl<T: FromWire> FromWire for Option<T> {
616    fn from_wire(client: &BridgeClient, value: Value) -> Result<Self, BridgeError> {
617        if value.is_null() || is_undefined(&value) {
618            return Ok(None);
619        }
620        T::from_wire(client, value).map(Some)
621    }
622}
623
624impl<T: FromWire> FromWire for Vec<T> {
625    fn from_wire(client: &BridgeClient, value: Value) -> Result<Self, BridgeError> {
626        match value {
627            Value::Array(items) => items
628                .into_iter()
629                .map(|item| T::from_wire(client, item))
630                .collect(),
631            other => Err(BridgeError::decode("a list", other)),
632        }
633    }
634}
635
636impl FromWire for RemoteHandle {
637    fn from_wire(client: &BridgeClient, value: Value) -> Result<Self, BridgeError> {
638        let Some(id) = value.get("$handle").and_then(Value::as_str) else {
639            return Err(BridgeError::decode("a remote object", value));
640        };
641        Ok(RemoteHandle {
642            client: client.clone(),
643            id: Arc::from(id),
644            type_name: Arc::from(value.get("type").and_then(Value::as_str).unwrap_or("")),
645        })
646    }
647}
648
649/// A typed wrapper of a remote Puppeteer object.
650pub trait Remote: Sized {
651    /// The Puppeteer type it wraps (`Page`).
652    const TYPE: &'static str;
653
654    /// Wrap a handle without checking its type.
655    fn from_remote(remote: RemoteHandle) -> Self;
656
657    /// The handle behind the wrapper.
658    fn remote(&self) -> &RemoteHandle;
659
660    /// View the same object as another type, for example an `ElementHandle`
661    /// as the `JSHandle` it extends.
662    fn cast<T: Remote>(&self) -> T {
663        T::from_remote(self.remote().clone())
664    }
665}
666
667/// Decode `{"$handle": …}`, such as an event argument, as a typed wrapper.
668pub fn decode_handle<T: Remote>(client: &BridgeClient, value: Value) -> Result<T, BridgeError> {
669    RemoteHandle::from_wire(client, value).map(T::from_remote)
670}
671
672/// Declare a wrapper struct around a [`RemoteHandle`].
673macro_rules! remote_type {
674    ($(#[$meta:meta])* $name:ident) => {
675        $(#[$meta])*
676        #[derive(Clone, Debug, PartialEq)]
677        pub struct $name {
678            pub(crate) remote: $crate::puppeteer::bridge::RemoteHandle,
679        }
680
681        impl $crate::puppeteer::bridge::Remote for $name {
682            const TYPE: &'static str = stringify!($name);
683
684            fn from_remote(remote: $crate::puppeteer::bridge::RemoteHandle) -> Self {
685                Self { remote }
686            }
687
688            fn remote(&self) -> &$crate::puppeteer::bridge::RemoteHandle {
689                &self.remote
690            }
691        }
692
693        impl $crate::puppeteer::bridge::FromWire for $name {
694            fn from_wire(
695                client: &$crate::puppeteer::bridge::BridgeClient,
696                value: serde_json::Value,
697            ) -> Result<Self, $crate::puppeteer::bridge::BridgeError> {
698                $crate::puppeteer::bridge::decode_handle(client, value)
699            }
700        }
701    };
702}
703pub(crate) use remote_type;
704
705/// Where the JavaScript CLI is: [`JS_CLI_ENV`], then `js/bin` next to the
706/// crate sources, then `node_modules/browser-commander` under `working_dir`
707/// (or the current directory).
708pub fn js_cli_path(working_dir: Option<&Path>) -> Result<PathBuf, BridgeError> {
709    let candidates = if let Some(configured) = std::env::var_os(JS_CLI_ENV) {
710        vec![PathBuf::from(configured)]
711    } else {
712        let base = working_dir
713            .map(Path::to_path_buf)
714            .or_else(|| std::env::current_dir().ok())
715            .unwrap_or_default();
716        vec![
717            PathBuf::from(env!("CARGO_MANIFEST_DIR")).join("../js/bin/browser-commander.js"),
718            base.join("node_modules/browser-commander/bin/browser-commander.js"),
719        ]
720    };
721    candidates
722        .into_iter()
723        .find(|path| path.is_file())
724        .ok_or_else(|| {
725            BridgeError::Unavailable(format!(
726                "the JavaScript CLI was not found; install the browser-commander npm package or set {JS_CLI_ENV}"
727            ))
728        })
729}
730
731/// Options for [`PuppeteerBridge::launch`].
732#[derive(Debug, Clone, Default)]
733pub struct BridgeOptions {
734    /// Node.js executable. Defaults to `BROWSER_COMMANDER_NODE`, then `node`.
735    pub node: Option<PathBuf>,
736    /// The JavaScript CLI. Defaults to [`js_cli_path`].
737    pub cli: Option<PathBuf>,
738    /// Working directory of the server.
739    pub working_dir: Option<PathBuf>,
740    /// Echo the server's stderr as it arrives.
741    pub verbose: bool,
742}
743
744/// A running `browser-commander serve --stdio` and its client.
745pub struct PuppeteerBridge {
746    client: BridgeClient,
747    runner: Mutex<Option<ProcessRunner>>,
748    pid: Option<u32>,
749    stderr: Arc<StdMutex<VecDeque<String>>>,
750}
751
752impl fmt::Debug for PuppeteerBridge {
753    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
754        f.debug_struct("PuppeteerBridge")
755            .field("pid", &self.pid)
756            .finish()
757    }
758}
759
760impl PuppeteerBridge {
761    /// Start `node <cli> serve --stdio` through command-stream.
762    pub async fn launch(options: BridgeOptions) -> Result<Self, BridgeError> {
763        let node = options
764            .node
765            .clone()
766            .or_else(|| std::env::var_os("BROWSER_COMMANDER_NODE").map(PathBuf::from))
767            .unwrap_or_else(|| PathBuf::from("node"));
768        let node = node.to_string_lossy().into_owned();
769        let cli = match options.cli.clone() {
770            Some(cli) => cli,
771            None => js_cli_path(options.working_dir.as_deref())?,
772        };
773        let cli = cli.to_string_lossy().into_owned();
774        let command = [node.as_str(), cli.as_str(), "serve", "--stdio"]
775            .into_iter()
776            .map(quote)
777            .collect::<Vec<_>>()
778            .join(" ");
779
780        let mut runner = ProcessRunner::new(
781            command,
782            RunOptions {
783                mirror: false,
784                capture: true,
785                stdin: StdinOption::Pipe,
786                cwd: options.working_dir.clone(),
787                shell_operators: false,
788                trace: false,
789                ..RunOptions::default()
790            },
791        );
792        runner
793            .start()
794            .await
795            .map_err(|err| BridgeError::Unavailable(format!("failed to start {node}: {err}")))?;
796        let pid = runner.pid();
797        let (stdin, stdout, stderr) = {
798            let mut child = runner.child().ok_or_else(|| {
799                BridgeError::Unavailable("the server process did not start".to_string())
800            })?;
801            let native = child.native_mut();
802            (
803                native.stdin.take(),
804                native.stdout.take(),
805                native.stderr.take(),
806            )
807        };
808        let (Some(stdin), Some(stdout)) = (stdin, stdout) else {
809            if let Some(pid) = pid {
810                kill_owned_process_tree(pid);
811            }
812            return Err(BridgeError::Unavailable(
813                "the server's stdin and stdout were not piped".to_string(),
814            ));
815        };
816
817        let stderr_lines = Arc::new(StdMutex::new(VecDeque::new()));
818        if let Some(stderr) = stderr {
819            let lines = Arc::clone(&stderr_lines);
820            let verbose = options.verbose;
821            tokio::spawn(async move {
822                let mut reader = BufReader::new(stderr).lines();
823                while let Ok(Some(line)) = reader.next_line().await {
824                    if verbose {
825                        eprintln!("[serve --stdio] {line}");
826                    }
827                    tracing::debug!(target: "browser_commander::puppeteer", "{line}");
828                    if let Ok(mut lines) = lines.lock() {
829                        if lines.len() == STDERR_LINES {
830                            lines.pop_front();
831                        }
832                        lines.push_back(line);
833                    }
834                }
835            });
836        }
837
838        Ok(Self {
839            client: BridgeClient::new(stdout, stdin),
840            runner: Mutex::new(Some(runner)),
841            pid,
842            stderr: stderr_lines,
843        })
844    }
845
846    /// The JSON-RPC client.
847    pub fn client(&self) -> &BridgeClient {
848        &self.client
849    }
850
851    /// The `puppeteer` module's default export, a `PuppeteerNode`.
852    pub async fn puppeteer(&self) -> Result<super::api::PuppeteerNode, BridgeError> {
853        let remote = self.client.root("puppeteer").await?;
854        Ok(<super::api::PuppeteerNode as Remote>::from_remote(remote))
855    }
856
857    /// Process id of the server.
858    pub fn pid(&self) -> Option<u32> {
859        self.pid
860    }
861
862    /// The server's most recent stderr lines.
863    pub fn stderr_tail(&self) -> Vec<String> {
864        self.stderr
865            .lock()
866            .map(|lines| lines.iter().cloned().collect())
867            .unwrap_or_default()
868    }
869
870    /// Stop the server: end its input, give it time to finish, then kill
871    /// whatever is left of its process tree. Browsers started with
872    /// `PuppeteerNode::launch` should be closed first.
873    pub async fn close(&self) {
874        self.client.close_input().await;
875        let Some(mut runner) = self.runner.lock().await.take() else {
876            return;
877        };
878        let deadline = tokio::time::Instant::now() + EXIT_GRACE;
879        loop {
880            let exited = runner
881                .child()
882                .map(|mut child| matches!(child.native_mut().try_wait(), Ok(Some(_))))
883                .unwrap_or(true);
884            if exited || tokio::time::Instant::now() >= deadline {
885                break;
886            }
887            tokio::time::sleep(Duration::from_millis(50)).await;
888        }
889        if let Some(pid) = self.pid {
890            kill_owned_process_tree(pid);
891        }
892    }
893}
894
895impl Drop for PuppeteerBridge {
896    fn drop(&mut self) {
897        let still_owned = self
898            .runner
899            .try_lock()
900            .map(|runner| runner.is_some())
901            .unwrap_or(true);
902        if still_owned {
903            if let Some(pid) = self.pid {
904                kill_owned_process_tree(pid);
905            }
906        }
907    }
908}
909
910/// An object literal from optional fields, leaving out the `None`s.
911pub fn options(fields: impl IntoIterator<Item = (&'static str, Option<Value>)>) -> Value {
912    Value::Object(
913        fields
914            .into_iter()
915            .filter_map(|(key, value)| value.map(|value| (key.to_string(), value)))
916            .collect::<Map<String, Value>>(),
917    )
918}
919
920#[cfg(test)]
921#[path = "bridge_tests.rs"]
922mod tests;