1use 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
26pub 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#[derive(Debug, Error)]
35pub enum BridgeError {
36 #[error("{name}: {message}")]
38 Remote {
39 code: i64,
41 name: String,
43 message: String,
45 stack: Option<String>,
47 },
48 #[error("serve --stdio closed: {0}")]
50 Closed(String),
51 #[error("serve --stdio I/O error: {0}")]
53 Io(#[from] std::io::Error),
54 #[error("serve --stdio sent invalid JSON: {0}")]
56 Json(#[from] serde_json::Error),
57 #[error("expected {expected} from the bridge, got {value}")]
59 Decode {
60 expected: &'static str,
62 value: Value,
64 },
65 #[error("serve --stdio unavailable: {0}")]
67 Unavailable(String),
68}
69
70impl BridgeError {
71 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
92struct 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 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#[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 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 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 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 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 pub async fn close_input(&self) {
249 let mut writer = self.inner.writer.lock().await;
250 let _ = writer.shutdown().await;
251 *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 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 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#[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 pub fn id(&self) -> &str {
367 &self.id
368 }
369
370 pub fn type_name(&self) -> &str {
372 &self.type_name
373 }
374
375 pub fn client(&self) -> &BridgeClient {
377 &self.client
378 }
379
380 pub fn to_wire(&self) -> Value {
382 json!({ "$handle": &*self.id })
383 }
384
385 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 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 pub async fn describe(&self) -> Result<Value, BridgeError> {
416 self.client
417 .request("handle.describe", json!({ "handle": &*self.id }))
418 .await
419 }
420
421 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 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
461fn 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#[derive(Debug)]
475pub struct Subscription {
476 client: BridgeClient,
477 id: String,
478 receiver: mpsc::UnboundedReceiver<Vec<Value>>,
479}
480
481impl Subscription {
482 pub fn id(&self) -> &str {
484 &self.id
485 }
486
487 pub async fn next(&mut self) -> Option<Vec<Value>> {
491 self.receiver.recv().await
492 }
493
494 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 self.client.lock_subscribers().remove(&self.id);
509 }
510}
511
512#[derive(Debug, Clone, PartialEq)]
516pub struct JsFunction(Value);
517
518impl JsFunction {
519 pub fn source(source: impl Into<String>) -> Self {
521 Self(json!({ "$function": source.into() }))
522 }
523
524 pub fn text(text: impl Into<String>) -> Self {
527 Self(Value::String(text.into()))
528 }
529
530 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
548pub fn binary(bytes: &[u8]) -> Value {
550 json!({ "$binary": base64::engine::general_purpose::STANDARD.encode(bytes) })
551}
552
553pub trait FromWire: Sized {
555 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 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
649pub trait Remote: Sized {
651 const TYPE: &'static str;
653
654 fn from_remote(remote: RemoteHandle) -> Self;
656
657 fn remote(&self) -> &RemoteHandle;
659
660 fn cast<T: Remote>(&self) -> T {
663 T::from_remote(self.remote().clone())
664 }
665}
666
667pub fn decode_handle<T: Remote>(client: &BridgeClient, value: Value) -> Result<T, BridgeError> {
669 RemoteHandle::from_wire(client, value).map(T::from_remote)
670}
671
672macro_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
705pub 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#[derive(Debug, Clone, Default)]
733pub struct BridgeOptions {
734 pub node: Option<PathBuf>,
736 pub cli: Option<PathBuf>,
738 pub working_dir: Option<PathBuf>,
740 pub verbose: bool,
742}
743
744pub 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 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 pub fn client(&self) -> &BridgeClient {
848 &self.client
849 }
850
851 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 pub fn pid(&self) -> Option<u32> {
859 self.pid
860 }
861
862 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 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
910pub 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;