Skip to main content

allwright/
client_tab.rs

1use crate::proto::context_session_command::Command as ContextCommand;
2use crate::proto::context_session_event::Event as ContextEvent;
3use crate::proto::{
4    CloseContextSessionCommand, ContextSessionCommand, ContextSessionPingCommand,
5    NavigatePageCommand,
6};
7use tokio::sync::mpsc;
8use tokio_stream::wrappers::ReceiverStream;
9
10use super::command::command_retry_options;
11use super::selectors::normalize_selector_for_transport;
12use super::types::{CommandOptions, Error, NavigateResult, Result, Tab, TabHandle, TabState};
13
14impl Tab {
15    pub fn locator(&self, css_selector: impl Into<String>) -> super::types::Locator {
16        super::types::Locator {
17            page: self.clone(),
18            selector: normalize_selector_for_transport(&css_selector.into()),
19        }
20    }
21
22    pub fn session_id(&self) -> &str {
23        &self.inner.session_id
24    }
25
26    pub async fn goto(&self, url: impl Into<String>) -> Result<NavigateResult> {
27        self.navigate(url).await
28    }
29
30    pub async fn ping(&self, message: impl Into<String>) -> Result<String> {
31        let mut state = self.inner.state.lock().await;
32        let handle = self.ensure_handle(&mut state).await?;
33        ensure_tab_open(handle, &self.inner.session_id)?;
34
35        handle
36            .command_tx
37            .send(ContextSessionCommand {
38                surface_session_id: self.inner.surface_session_id.clone(),
39                context_session_id: self.inner.session_id.clone(),
40                command: Some(ContextCommand::Ping(ContextSessionPingCommand {
41                    message: message.into(),
42                })),
43            })
44            .await
45            .map_err(|_| Error::new("failed to send ContextSessionPingCommand"))?;
46
47        loop {
48            let event = handle
49                .events
50                .message()
51                .await?
52                .ok_or_else(|| Error::new("tab session closed while waiting for pong"))?;
53
54            match event.event {
55                Some(ContextEvent::Attached(_)) => {}
56                Some(ContextEvent::Pong(pong)) => return Ok(pong.message),
57                Some(ContextEvent::Error(error)) => {
58                    return Err(Error::new(format!(
59                        "tab session error while pinging: {}",
60                        error.message
61                    )));
62                }
63                Some(ContextEvent::Closed(_)) => {
64                    handle.closed = true;
65                    return Err(Error::new(format!(
66                        "tab session {} closed while waiting for pong",
67                        self.inner.session_id
68                    )));
69                }
70                _ => {}
71            }
72        }
73    }
74
75    pub async fn navigate(&self, url: impl Into<String>) -> Result<NavigateResult> {
76        self.navigate_with_options(url, CommandOptions::default())
77            .await
78    }
79
80    pub async fn navigate_with_options(
81        &self,
82        url: impl Into<String>,
83        options: CommandOptions,
84    ) -> Result<NavigateResult> {
85        let mut state = self.inner.state.lock().await;
86        let handle = self.ensure_handle(&mut state).await?;
87        ensure_tab_open(handle, &self.inner.session_id)?;
88
89        handle
90            .command_tx
91            .send(ContextSessionCommand {
92                surface_session_id: self.inner.surface_session_id.clone(),
93                context_session_id: self.inner.session_id.clone(),
94                command: Some(ContextCommand::Navigate(NavigatePageCommand {
95                    url: url.into(),
96                    retry_options: command_retry_options(options.timeout_ms),
97                })),
98            })
99            .await
100            .map_err(|_| Error::new("failed to send NavigatePageCommand"))?;
101
102        let mut navigated = None;
103        let mut injection = None;
104
105        loop {
106            let event = handle
107                .events
108                .message()
109                .await?
110                .ok_or_else(|| Error::new("tab session closed while waiting for navigation"))?;
111
112            match event.event {
113                Some(ContextEvent::Attached(_)) => {}
114                Some(ContextEvent::Navigated(navigated_event)) => {
115                    navigated = Some(navigated_event);
116                }
117                Some(ContextEvent::ChromiumBidiInjection(injection_event)) => {
118                    injection = Some(injection_event);
119                }
120                Some(ContextEvent::Error(error)) => {
121                    return Err(Error::new(format!(
122                        "tab session error while navigating: {}",
123                        error.message
124                    )));
125                }
126                Some(ContextEvent::Closed(_)) => {
127                    handle.closed = true;
128                    return Err(Error::new(format!(
129                        "tab session {} closed while navigating",
130                        self.inner.session_id
131                    )));
132                }
133                _ => {}
134            }
135
136            if navigated.is_some() && injection.is_some() {
137                let navigated_event = navigated
138                    .take()
139                    .ok_or_else(|| Error::new("navigation event disappeared unexpectedly"))?;
140                let injection_event = injection
141                    .take()
142                    .ok_or_else(|| Error::new("bidi injection event disappeared unexpectedly"))?;
143                return Ok(NavigateResult {
144                    url: navigated_event.url,
145                    note: navigated_event.note,
146                    bidi_session_id: injection_event.bidi_session_id,
147                    mapper_target_id: injection_event.mapper_target_id,
148                    mapper_session_id: injection_event.mapper_session_id,
149                    package_version: injection_event.package_version,
150                });
151            }
152        }
153    }
154
155    pub async fn close(&self) -> Result<()> {
156        let mut state = self.inner.state.lock().await;
157        let handle = self.ensure_handle(&mut state).await?;
158        if handle.closed {
159            return Ok(());
160        }
161
162        handle
163            .command_tx
164            .send(ContextSessionCommand {
165                surface_session_id: self.inner.surface_session_id.clone(),
166                context_session_id: self.inner.session_id.clone(),
167                command: Some(ContextCommand::Close(CloseContextSessionCommand {})),
168            })
169            .await
170            .map_err(|_| Error::new("failed to send CloseContextSessionCommand"))?;
171
172        loop {
173            let event = handle
174                .events
175                .message()
176                .await?
177                .ok_or_else(|| Error::new("tab session closed before close confirmation"))?;
178
179            match event.event {
180                Some(ContextEvent::Attached(_)) => {}
181                Some(ContextEvent::Closed(_)) => {
182                    handle.closed = true;
183                    return Ok(());
184                }
185                Some(ContextEvent::Error(error)) => {
186                    return Err(Error::new(format!(
187                        "tab session error while closing: {}",
188                        error.message
189                    )));
190                }
191                _ => {}
192            }
193        }
194    }
195
196    pub(crate) async fn ensure_handle<'a>(
197        &self,
198        state: &'a mut TabState,
199    ) -> Result<&'a mut TabHandle> {
200        if state.handle.is_none() {
201            let mut engine = self.inner.runtime.engine.clone();
202            let (command_tx, command_rx) = mpsc::channel(16);
203            let response = engine
204                .context_session(tonic::Request::new(ReceiverStream::new(command_rx)))
205                .await?;
206            state.handle = Some(TabHandle {
207                command_tx,
208                events: response.into_inner(),
209                closed: false,
210            });
211        }
212
213        state
214            .handle
215            .as_mut()
216            .ok_or_else(|| Error::new("tab session handle was not initialized"))
217    }
218}
219
220pub(crate) fn ensure_tab_open(handle: &TabHandle, session_id: &str) -> Result<()> {
221    if handle.closed {
222        return Err(Error::new(format!("tab session {} is closed", session_id)));
223    }
224    Ok(())
225}