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#[derive(Debug, Error)]
19pub enum ClientError {
20 #[error("transport error: {0}")]
22 Transport(String),
23 #[error("HTTP {status}: {message}")]
25 Http { status: u16, message: String },
26 #[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#[derive(Debug, Clone, Default)]
46pub struct RequestOptions {
47 pub timeout: Option<Duration>,
49 pub correlation_id: Option<String>,
51 pub idempotency_key: Option<String>,
53}
54
55#[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 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 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 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 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 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 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 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}