Skip to main content

bobby_browser_client/
http.rs

1use chrono::{Duration as ChronoDuration, Utc};
2use reqwest::{Client, Method};
3use serde::de::DeserializeOwned;
4use serde::Serialize;
5use std::time::Duration;
6use thiserror::Error;
7use uuid::Uuid;
8
9use crate::{
10    CommandEnvelope, CommandOutcome, CreateSessionRequest, OpenPageRequest, PageState, RuntimeInfo,
11    SessionId, SessionState, CURRENT_INTERFACE_VERSION,
12};
13
14/// Errors returned by [`BrowserRuntimeClient`].
15///
16/// Messages are redacted so the bearer token never appears in transport,
17/// HTTP, or protocol error text.
18#[derive(Debug, Error)]
19pub enum ClientError {
20    /// Network or HTTP-client failure before a valid response body.
21    #[error("transport error: {0}")]
22    Transport(String),
23    /// Non-success HTTP status with response body text.
24    #[error("HTTP {status}: {message}")]
25    Http { status: u16, message: String },
26    /// Response body could not be interpreted, or client preconditions failed.
27    #[error("protocol error: {0}")]
28    Protocol(String),
29}
30
31impl ClientError {
32    fn redact(self, bearer: &str) -> Self {
33        match self {
34            Self::Transport(message) => Self::Transport(message.replace(bearer, "")),
35            Self::Http { status, message } => Self::Http {
36                status,
37                message: message.replace(bearer, ""),
38            },
39            Self::Protocol(message) => Self::Protocol(message.replace(bearer, "")),
40        }
41    }
42}
43
44/// Options for a single request.
45#[derive(Debug, Clone, Default)]
46pub struct RequestOptions {
47    /// Request timeout (defaults to 30 seconds).
48    pub timeout: Option<Duration>,
49    /// Value for `x-correlation-id` (generated UUID when omitted).
50    pub correlation_id: Option<String>,
51    /// Value for `idempotency-key` when the operation supports it.
52    pub idempotency_key: Option<String>,
53}
54
55/// Authenticated HTTP client for the Bobby Browser `/v1` runtime interface.
56///
57/// Construct with [`BrowserRuntimeClient::new`]. `base_url` is the runtime
58/// origin; a trailing slash and a trailing `/v1` are stripped.
59#[derive(Debug, Clone)]
60pub struct BrowserRuntimeClient {
61    base_url: String,
62    bearer_token: String,
63    http: Client,
64    default_timeout: Duration,
65}
66
67impl BrowserRuntimeClient {
68    /// Create a client.
69    ///
70    /// `base_url` is the runtime origin (`http://127.0.0.1:7777` or `…/v1`).
71    /// Returns [`ClientError::Protocol`] when `base_url` or `bearer_token` is empty.
72    pub fn new(
73        base_url: impl Into<String>,
74        bearer_token: impl Into<String>,
75    ) -> Result<Self, ClientError> {
76        let bearer_token = bearer_token.into();
77        if bearer_token.is_empty() {
78            return Err(ClientError::Protocol(
79                "bearerToken must not be empty".into(),
80            ));
81        }
82        let base_url = normalize_base_url(base_url.into());
83        if base_url.is_empty() {
84            return Err(ClientError::Protocol("baseUrl must not be empty".into()));
85        }
86        let http = Client::builder()
87            .timeout(Duration::from_secs(30))
88            .build()
89            .map_err(|error| ClientError::Transport(error.to_string()))?;
90        Ok(Self {
91            base_url,
92            bearer_token,
93            http,
94            default_timeout: Duration::from_secs(30),
95        })
96    }
97
98    /// `GET /v1/runtime` — version, capabilities, and load counters.
99    pub async fn runtime_info(
100        &self,
101        options: Option<RequestOptions>,
102    ) -> Result<RuntimeInfo, ClientError> {
103        self.json(Method::GET, "/v1/runtime", None::<()>, options)
104            .await
105    }
106
107    /// `POST /v1/sessions` — create a browser session.
108    pub async fn create_session(
109        &self,
110        input: &CreateSessionRequest,
111        options: Option<RequestOptions>,
112    ) -> Result<SessionState, ClientError> {
113        self.json(Method::POST, "/v1/sessions", Some(input), options)
114            .await
115    }
116
117    /// `GET /v1/sessions` — list active sessions.
118    pub async fn list_sessions(
119        &self,
120        options: Option<RequestOptions>,
121    ) -> Result<Vec<SessionState>, ClientError> {
122        self.json(Method::GET, "/v1/sessions", None::<()>, options)
123            .await
124    }
125
126    /// `DELETE /v1/sessions/{id}` — tear down a session (`204` on success).
127    pub async fn delete_session(
128        &self,
129        session_id: &SessionId,
130        options: Option<RequestOptions>,
131    ) -> Result<(), ClientError> {
132        self.empty(
133            Method::DELETE,
134            &format!("/v1/sessions/{}", session_id.0),
135            options,
136        )
137        .await
138    }
139
140    /// `POST /v1/pages` — open a page in a session.
141    pub async fn open_page(
142        &self,
143        input: &OpenPageRequest,
144        options: Option<RequestOptions>,
145    ) -> Result<PageState, ClientError> {
146        self.json(Method::POST, "/v1/pages", Some(input), options)
147            .await
148    }
149
150    /// `POST /v1/commands` — submit a command envelope (primitive or intent).
151    pub async fn submit(
152        &self,
153        input: &CommandEnvelope,
154        options: Option<RequestOptions>,
155    ) -> Result<CommandOutcome, ClientError> {
156        self.json(Method::POST, "/v1/commands", Some(input), options)
157            .await
158    }
159
160    async fn empty(
161        &self,
162        method: Method,
163        path: &str,
164        options: Option<RequestOptions>,
165    ) -> Result<(), ClientError> {
166        let options = options.unwrap_or_default();
167        let timeout = options.timeout.unwrap_or(self.default_timeout);
168        let deadline =
169            Utc::now() + ChronoDuration::from_std(timeout).unwrap_or(ChronoDuration::seconds(30));
170        let correlation = options
171            .correlation_id
172            .unwrap_or_else(|| Uuid::new_v4().to_string());
173
174        let mut request = self
175            .http
176            .request(method, format!("{}{path}", self.base_url))
177            .timeout(timeout)
178            .header("authorization", format!("Bearer {}", self.bearer_token))
179            .header("x-interface-version", CURRENT_INTERFACE_VERSION)
180            .header("x-correlation-id", correlation)
181            .header("x-deadline", deadline.to_rfc3339());
182        if let Some(key) = options.idempotency_key {
183            request = request.header("idempotency-key", key);
184        }
185
186        let response = request.send().await.map_err(|error| {
187            ClientError::Transport(error.to_string()).redact(&self.bearer_token)
188        })?;
189        let status = response.status();
190        if status == reqwest::StatusCode::NO_CONTENT || status.is_success() {
191            return Ok(());
192        }
193        let text = response.text().await.unwrap_or_default();
194        Err(ClientError::Http {
195            status: status.as_u16(),
196            message: text,
197        }
198        .redact(&self.bearer_token))
199    }
200
201    async fn json<B, T>(
202        &self,
203        method: Method,
204        path: &str,
205        body: Option<B>,
206        options: Option<RequestOptions>,
207    ) -> Result<T, ClientError>
208    where
209        B: Serialize,
210        T: DeserializeOwned,
211    {
212        let options = options.unwrap_or_default();
213        let timeout = options.timeout.unwrap_or(self.default_timeout);
214        let deadline =
215            Utc::now() + ChronoDuration::from_std(timeout).unwrap_or(ChronoDuration::seconds(30));
216        let correlation = options
217            .correlation_id
218            .unwrap_or_else(|| Uuid::new_v4().to_string());
219
220        let mut request = self
221            .http
222            .request(method, format!("{}{path}", self.base_url))
223            .timeout(timeout)
224            .header("authorization", format!("Bearer {}", self.bearer_token))
225            .header("x-interface-version", CURRENT_INTERFACE_VERSION)
226            .header("x-correlation-id", correlation)
227            .header("x-deadline", deadline.to_rfc3339());
228        if let Some(key) = options.idempotency_key {
229            request = request.header("idempotency-key", key);
230        }
231        if let Some(body) = body {
232            request = request
233                .header("content-type", "application/json")
234                .json(&body);
235        }
236
237        let response = request.send().await.map_err(|error| {
238            ClientError::Transport(error.to_string()).redact(&self.bearer_token)
239        })?;
240        let status = response.status();
241        let text = response.text().await.map_err(|error| {
242            ClientError::Transport(error.to_string()).redact(&self.bearer_token)
243        })?;
244        if !status.is_success() {
245            return Err(ClientError::Http {
246                status: status.as_u16(),
247                message: text,
248            }
249            .redact(&self.bearer_token));
250        }
251        serde_json::from_str(&text).map_err(|error| {
252            ClientError::Protocol(format!("invalid JSON body: {error}")).redact(&self.bearer_token)
253        })
254    }
255}
256
257fn normalize_base_url(value: String) -> String {
258    let trimmed = value.trim_end_matches('/').to_string();
259    if let Some(stripped) = trimmed.strip_suffix("/v1") {
260        stripped.trim_end_matches('/').to_string()
261    } else {
262        trimmed
263    }
264}
265
266#[cfg(test)]
267mod tests {
268    use super::*;
269    use axum::http::HeaderMap;
270    use axum::http::StatusCode;
271    use axum::response::IntoResponse;
272    use axum::routing::get;
273    use axum::Router;
274    use serde_json::json;
275    use tokio::net::TcpListener;
276
277    async fn runtime_handler(headers: HeaderMap) -> impl IntoResponse {
278        assert!(headers
279            .get("authorization")
280            .and_then(|v| v.to_str().ok())
281            .is_some_and(|v| v == "Bearer test-token"));
282        assert_eq!(
283            headers
284                .get("x-interface-version")
285                .and_then(|v| v.to_str().ok()),
286            Some(CURRENT_INTERFACE_VERSION)
287        );
288        assert!(headers.get("x-correlation-id").is_some());
289        assert!(headers.get("x-deadline").is_some());
290        (
291            StatusCode::OK,
292            [(axum::http::header::CONTENT_TYPE, "application/json")],
293            json!({
294                "version": env!("CARGO_PKG_VERSION"),
295                "capabilities": ["session:read"],
296                "active_sessions": 0,
297                "queued_jobs": 0,
298                "uptime_ms": 1,
299            })
300            .to_string(),
301        )
302    }
303
304    #[tokio::test]
305    async fn runtime_info_sends_required_headers() {
306        let app = Router::new().route("/v1/runtime", get(runtime_handler));
307        let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
308        let addr = listener.local_addr().unwrap();
309        tokio::spawn(async move {
310            axum::serve(listener, app).await.unwrap();
311        });
312
313        let client = BrowserRuntimeClient::new(format!("http://{addr}/v1"), "test-token").unwrap();
314        let info = client.runtime_info(None).await.unwrap();
315        assert_eq!(info.version, env!("CARGO_PKG_VERSION"));
316        assert_eq!(info.active_sessions, 0);
317    }
318
319    #[tokio::test]
320    async fn normalize_strips_v1_suffix() {
321        assert_eq!(
322            normalize_base_url("http://127.0.0.1:7777/v1/".into()),
323            "http://127.0.0.1:7777"
324        );
325    }
326
327    #[tokio::test]
328    async fn rejects_empty_bearer() {
329        let err = BrowserRuntimeClient::new("http://127.0.0.1:7777", "").unwrap_err();
330        assert!(matches!(err, ClientError::Protocol(_)));
331    }
332}