use anyhow::Result;
use serde_json::{json, Value};
use std::collections::HashMap;
use std::sync::Arc;
use tokio::io::{AsyncBufReadExt, AsyncWriteExt, BufReader};
use tokio::sync::{Mutex, RwLock};
use crate::cli::env_resolver::ResolvedBrowser;
use crate::session::backend::TabBackend;
pub type BidiCache = Arc<Mutex<Option<Arc<crate::bidi::BidiClient>>>>;
#[derive(Clone)]
pub struct ServerState {
pub browser: Arc<RwLock<ResolvedBrowser>>,
pub bidi: BidiCache,
pub bidi_lock: Arc<Mutex<BidiLockState>>,
pub backend: Arc<Mutex<Option<TabBackend>>>,
pub active_target_id: Arc<Mutex<Option<String>>>,
pub origin_target_ids: Arc<Mutex<HashMap<String, String>>>,
pub sidecar: Arc<Mutex<Option<crate::sidecar::Sidecar>>>,
pub sidecar_config: crate::sidecar::SidecarConfig,
pub op_barrier: Arc<RwLock<()>>,
}
#[derive(Default)]
pub enum BidiLockState {
#[default]
Pending,
Acquired(crate::registry::BidiLockGuard),
NotApplicable,
}
impl ServerState {
pub fn new(browser: ResolvedBrowser) -> Self {
Self::with_sidecar_config(browser, crate::sidecar::SidecarConfig::default())
}
pub fn with_sidecar_config(
browser: ResolvedBrowser,
sidecar_config: crate::sidecar::SidecarConfig,
) -> Self {
Self {
browser: Arc::new(RwLock::new(browser)),
bidi: Arc::new(Mutex::new(None)),
bidi_lock: Arc::new(Mutex::new(BidiLockState::Pending)),
backend: Arc::new(Mutex::new(None)),
active_target_id: Arc::new(Mutex::new(None)),
origin_target_ids: Arc::new(Mutex::new(HashMap::new())),
sidecar: Arc::new(Mutex::new(None)),
sidecar_config,
op_barrier: Arc::new(RwLock::new(())),
}
}
pub async fn ensure_sidecar(&self, tool_name: &str) -> Result<crate::sidecar::Sidecar> {
let resolved = self.browser_snapshot().await;
if resolved.engine != crate::detect::Engine::Cdp {
return Err(crate::errors::SessionError::EngineUnsupported {
tool: tool_name.to_string(),
required_engine: "Chromium (CDP)".into(),
current_engine: format!("{:?}", resolved.engine),
hint: "use engine-agnostic tools such as browser_get_html, browser_fetch, browser_take_screenshot, or switch to a Chromium browser via browser_select",
}
.into());
}
let mut guard = self.sidecar.lock().await;
if let Some(sc) = guard.as_ref() {
return Ok(sc.clone());
}
let sc = crate::sidecar::Sidecar::start(self.sidecar_config.clone()).await?;
sc.connect(&resolved.endpoint).await?;
*guard = Some(sc.clone());
Ok(sc)
}
pub async fn browser_snapshot(&self) -> ResolvedBrowser {
self.browser.read().await.clone()
}
pub async fn ensure_bidi_lock(&self) -> Result<()> {
use crate::cli::env_resolver::Source;
use crate::cli::mcp::acquire_bidi_lock_if_needed;
use crate::detect::Engine;
use crate::registry::Registry;
let mut guard = self.bidi_lock.lock().await;
if matches!(*guard, BidiLockState::Pending) {
let resolved = self.browser_snapshot().await;
if resolved.engine != Engine::Bidi || matches!(resolved.source, Source::External) {
*guard = BidiLockState::NotApplicable;
return Ok(());
}
let acquired = tokio::task::spawn_blocking(move || {
let registry = Registry::open()?;
acquire_bidi_lock_if_needed(®istry, &resolved)
})
.await??;
*guard = match acquired {
Some(lock) => BidiLockState::Acquired(lock),
None => BidiLockState::NotApplicable,
};
}
Ok(())
}
pub async fn ensure_backend(&self) -> Result<TabBackend> {
self.ensure_bidi_lock().await?;
let mut guard = self.backend.lock().await;
if let Some(b) = guard.as_ref() {
return Ok(b.clone());
}
let resolved = self.browser_snapshot().await;
let b = crate::session::backend::open_backend(&resolved.endpoint, resolved.engine).await?;
*guard = Some(b.clone());
Ok(b)
}
pub async fn current_tab(&self) -> Result<(TabBackend, String)> {
let backend = self.ensure_backend().await?;
let mut pointer = self.active_target_id.lock().await;
if let Some(tid) = pointer.as_ref() {
let live = backend.live_target_ids().await?;
if live.contains(tid) {
return Ok((backend, tid.clone()));
}
}
let new_tid = backend.create_tab("about:blank").await?;
*pointer = Some(new_tid.clone());
Ok((backend, new_tid))
}
pub async fn resolve_or_create_for_origin(&self, url: &str) -> Result<(TabBackend, String)> {
let want =
url::Url::parse(url).map_err(|e| anyhow::anyhow!("invalid fetch URL `{url}`: {e}"))?;
let origin_root = crate::session::attach::origin_root_url(&want);
let backend = self.ensure_backend().await?;
let mut origin_targets = self.origin_target_ids.lock().await;
let live_targets = backend.live_targets().await?;
let live_ids: std::collections::HashSet<&str> =
live_targets.iter().map(|t| t.id.as_str()).collect();
if let Some(cached) = origin_targets.get(&origin_root) {
if live_ids.contains(cached.as_str()) {
return Ok((backend, cached.clone()));
}
origin_targets.remove(&origin_root);
}
if let Some(existing) = live_targets.iter().find(|t| {
url::Url::parse(&t.url)
.map(|parsed| crate::session::attach::same_origin(&parsed, &want))
.unwrap_or(false)
}) {
origin_targets.insert(origin_root, existing.id.clone());
return Ok((backend, existing.id.clone()));
}
let new_tid = backend.create_tab(&origin_root).await?;
origin_targets.insert(origin_root, new_tid.clone());
Ok((backend, new_tid))
}
pub async fn resolve_target_for_args(&self, args: &Value) -> Result<(TabBackend, String)> {
let tab = args.get("tab").and_then(|v| v.as_str()).map(String::from);
let target = args
.get("target")
.and_then(|v| v.as_str())
.map(String::from);
match (tab, target) {
(Some(_), Some(_)) => Err(anyhow::anyhow!("`tab` and `target` are mutually exclusive")),
(Some(name), None) => {
let backend = self.ensure_backend().await?;
let browser_name = self.registered_browser_name().await?;
let bn = browser_name.clone();
let n = name.clone();
let row = sync_registry_op(move |reg| reg.tab_get(&bn, &n))
.await?
.ok_or_else(|| crate::errors::SessionError::TabNotFound {
browser: browser_name.clone(),
name: name.clone(),
})?;
let live = backend.live_target_ids().await?;
if !live.contains(&row.target_id) {
let bn = browser_name.clone();
let n = name.clone();
sync_registry_op(move |reg| reg.tab_delete(&bn, &n)).await?;
return Err(crate::errors::SessionError::TabNotFound {
browser: browser_name,
name,
}
.into());
}
let bn = browser_name.clone();
let n = name.clone();
sync_registry_op(move |reg| reg.tab_touch(&bn, &n)).await?;
Ok((backend, row.target_id))
}
(None, Some(regex)) => {
let backend = self.ensure_backend().await?;
let target_id = resolve_target_by_regex(&backend, ®ex).await?;
Ok((backend, target_id))
}
(None, None) => self.current_tab().await,
}
}
pub async fn registered_browser_name(&self) -> Result<String> {
use crate::cli::env_resolver::Source;
let resolved = self.browser_snapshot().await;
match resolved.source {
Source::Registered { name } => Ok(name),
Source::External => Err(anyhow::anyhow!(
"operation requires a registered browser; external URL endpoints \
don't have a stable identity"
)),
}
}
pub async fn switch_browser(&self, new_browser: ResolvedBrowser) -> Result<()> {
{
let mut bidi = self.bidi.lock().await;
if let Some(client) = bidi.take() {
let _ = client.session_end().await;
}
}
{
let mut backend = self.backend.lock().await;
*backend = None;
}
{
let mut lock = self.bidi_lock.lock().await;
*lock = BidiLockState::Pending;
}
{
let mut pointer = self.active_target_id.lock().await;
*pointer = None;
}
{
let mut origins = self.origin_target_ids.lock().await;
origins.clear();
}
{
let mut sidecar = self.sidecar.lock().await;
if let Some(sc) = sidecar.take() {
let _ = sc.call("dispose", serde_json::json!({})).await;
drop(sc);
}
}
{
let mut br = self.browser.write().await;
*br = new_browser;
}
self.ensure_bidi_lock().await?;
Ok(())
}
}
pub(crate) async fn sync_registry_op<T, F>(f: F) -> Result<T>
where
F: FnOnce(&crate::registry::Registry) -> Result<T> + Send + 'static,
T: Send + 'static,
{
tokio::task::spawn_blocking(move || {
let reg = crate::registry::Registry::open()?;
f(®)
})
.await?
}
pub(crate) async fn resolve_browser_send(
selector: crate::cli::env_resolver::BrowserSelector,
) -> Result<ResolvedBrowser> {
use crate::cli::env_resolver::{BrowserSelector, DefaultResolver, Resolver};
match selector {
BrowserSelector::Url(u) => match u.scheme() {
"ws" | "wss" => Ok(ResolvedBrowser {
engine: if u.path().contains("/session") {
crate::detect::Engine::Bidi
} else {
crate::detect::Engine::Cdp
},
endpoint: u.to_string(),
source: crate::cli::env_resolver::Source::External,
}),
"http" | "https" => {
let base = u.as_str().trim_end_matches('/').to_string();
let ws = DefaultResolver.fetch_version(&base).await?;
let ws_url = url::Url::parse(&ws)?;
Ok(ResolvedBrowser {
engine: if ws_url.path().contains("/session") {
crate::detect::Engine::Bidi
} else {
crate::detect::Engine::Cdp
},
endpoint: ws,
source: crate::cli::env_resolver::Source::External,
})
}
other => anyhow::bail!("unsupported URL scheme: {other}"),
},
other => {
tokio::task::spawn_blocking(move || {
let reg = crate::registry::Registry::open()?;
let rt = tokio::runtime::Builder::new_current_thread().build()?;
rt.block_on(crate::cli::env_resolver::resolve_with(
other,
®,
&DefaultResolver,
))
})
.await?
}
}
}
async fn resolve_target_by_regex(backend: &TabBackend, regex: &str) -> Result<String> {
use crate::errors::SessionError;
use regex::Regex;
use std::time::Duration;
const PROBE: Duration = Duration::from_millis(500);
let re = Regex::new(regex)?;
let targets = backend.live_targets().await?;
let matches: Vec<_> = targets.iter().filter(|t| re.is_match(&t.url)).collect();
if matches.is_empty() {
return Err(anyhow::anyhow!("no target matched URL regex `{regex}`"));
}
let mut last_id: Option<String> = None;
let mut last_url: Option<String> = None;
for t in &matches {
last_id = Some(t.id.clone());
last_url = Some(t.url.clone());
let ok = matches!(
tokio::time::timeout(PROBE, backend.evaluate(&t.id, "1", false, PROBE)).await,
Ok(Ok(_))
);
if ok {
return Ok(t.id.clone());
}
}
Err(SessionError::TabHung {
target_id: last_id,
url: last_url,
timeout_ms: PROBE.as_millis() as u64,
hint: "all-matches-hung",
}
.into())
}
impl std::fmt::Debug for ServerState {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("ServerState").finish()
}
}
pub type ToolHandler = std::sync::Arc<
dyn Fn(ServerState, Value) -> futures_util::future::BoxFuture<'static, Result<Value>>
+ Send
+ Sync,
>;
pub struct RegisteredTool {
pub name: String,
pub description: String,
pub input_schema: Value,
pub handler: ToolHandler,
}
#[derive(Clone, Default)]
pub struct ToolRegistry {
inner: std::sync::Arc<std::sync::Mutex<Vec<RegisteredTool>>>,
}
impl ToolRegistry {
pub fn new() -> Self {
Self::default()
}
pub fn register(&self, t: RegisteredTool) {
self.inner.lock().unwrap().push(t);
}
pub fn list(&self) -> Vec<Value> {
self.inner
.lock()
.unwrap()
.iter()
.map(|t| {
json!({
"name": t.name,
"description": t.description,
"inputSchema": t.input_schema,
})
})
.collect()
}
pub fn handler(&self, name: &str) -> Option<ToolHandler> {
self.inner
.lock()
.unwrap()
.iter()
.find(|t| t.name == name)
.map(|t| t.handler.clone())
}
}
pub async fn run(state: ServerState, tools: ToolRegistry) -> Result<()> {
run_with_streams(state, tools, tokio::io::stdin(), tokio::io::stdout()).await
}
pub async fn run_with_streams<I, O>(
state: ServerState,
tools: ToolRegistry,
stdin: I,
mut stdout: O,
) -> Result<()>
where
I: tokio::io::AsyncRead + Unpin,
O: tokio::io::AsyncWrite + Unpin + Send + 'static,
{
let (tx, mut rx) = tokio::sync::mpsc::unbounded_channel::<Vec<u8>>();
let writer = tokio::spawn(async move {
while let Some(frame) = rx.recv().await {
if stdout.write_all(&frame).await.is_err() {
break;
}
if stdout.flush().await.is_err() {
break;
}
}
});
let mut lines = BufReader::new(stdin).lines();
while let Some(line) = lines.next_line().await? {
if line.trim().is_empty() {
continue;
}
let req: Value = match serde_json::from_str(&line) {
Ok(v) => v,
Err(e) => {
let _ = tx.send(error_frame(
Value::Null,
-32700,
&format!("parse error: {e}"),
));
continue;
}
};
let id = req.get("id").cloned().unwrap_or(Value::Null);
let method = req.get("method").and_then(|m| m.as_str()).unwrap_or("");
let params = req.get("params").cloned().unwrap_or(Value::Null);
if id.is_null() && method == "notifications/initialized" {
continue;
}
match method {
"initialize" => {
let _ = tx.send(handle_initialize(id));
}
"ping" => {
let _ = tx.send(handle_ping(id));
}
"tools/list" => {
let _ = tx.send(handle_tools_list(id, &tools));
}
"tools/call" => {
handle_tools_call(id, ¶ms, &state, &tools, &tx);
}
_ => {
let _ = tx.send(error_frame(
id,
-32601,
&format!("method not found: {method}"),
));
}
}
}
drop(tx);
let _ = writer.await;
Ok(())
}
fn handle_initialize(id: Value) -> Vec<u8> {
result_frame(
id,
json!({
"protocolVersion": "2024-11-05",
"capabilities": {"tools": {}},
"serverInfo": {
"name": "browser-control",
"version": env!("CARGO_PKG_VERSION"),
},
}),
)
}
fn handle_ping(id: Value) -> Vec<u8> {
result_frame(id, json!({}))
}
fn handle_tools_list(id: Value, tools: &ToolRegistry) -> Vec<u8> {
result_frame(id, json!({"tools": tools.list()}))
}
fn handle_tools_call(
id: Value,
params: &Value,
state: &ServerState,
tools: &ToolRegistry,
tx: &tokio::sync::mpsc::UnboundedSender<Vec<u8>>,
) {
let name = params
.get("name")
.and_then(|v| v.as_str())
.unwrap_or("")
.to_string();
let args = params.get("arguments").cloned().unwrap_or(Value::Null);
let handler = tools.handler(&name);
let state = state.clone();
let tx = tx.clone();
let barrier = state.op_barrier.clone();
let is_exclusive = name == "browser_select";
tokio::spawn(async move {
let frame = match handler {
None => error_frame(id, -32602, &format!("tool not found: {name}")),
Some(h) => {
if is_exclusive {
let _guard = barrier.write().await;
match h(state, args).await {
Ok(v) => result_frame(id, v),
Err(e) => tool_error_frame(id, &e),
}
} else {
let _guard = barrier.read().await;
match h(state, args).await {
Ok(v) => result_frame(id, v),
Err(e) => tool_error_frame(id, &e),
}
}
}
};
let _ = tx.send(frame);
});
}
fn result_frame(id: Value, result: Value) -> Vec<u8> {
let resp = json!({"jsonrpc": "2.0", "id": id, "result": result});
let mut s = serde_json::to_vec(&resp).unwrap_or_else(|_| b"{}".to_vec());
s.push(b'\n');
s
}
fn tool_error_frame(id: Value, err: &anyhow::Error) -> Vec<u8> {
let message = tool_error_message(err);
result_frame(
id,
json!({
"content": [{ "type": "text", "text": message }],
"isError": true,
}),
)
}
fn tool_error_message(err: &anyhow::Error) -> String {
use crate::errors::SessionError;
use crate::registry::bidi_lock::BidiLockBusy;
if let Some(se) = err.downcast_ref::<SessionError>() {
return se.to_string();
}
if let Some(busy) = err.downcast_ref::<BidiLockBusy>() {
return busy.to_string();
}
format!("{err:#}")
}
fn error_frame(id: Value, code: i64, message: &str) -> Vec<u8> {
let resp = json!({
"jsonrpc": "2.0",
"id": id,
"error": {"code": code, "message": message},
});
let mut s = serde_json::to_vec(&resp).unwrap_or_else(|_| b"{}".to_vec());
s.push(b'\n');
s
}
#[cfg(test)]
mod tests {
use super::*;
use crate::cli::env_resolver::Source;
use crate::detect::Engine;
use tokio::io::{AsyncBufReadExt, AsyncWriteExt, BufReader};
fn dummy_resolved() -> ResolvedBrowser {
ResolvedBrowser {
endpoint: "ws://localhost:9999".into(),
engine: Engine::Cdp,
source: Source::External,
}
}
fn dummy_state() -> ServerState {
ServerState::new(dummy_resolved())
}
async fn send_recv(tools: ToolRegistry, requests: &[Value]) -> Vec<Value> {
let (mut client_w, server_r) = tokio::io::duplex(8192);
let (server_w, client_r) = tokio::io::duplex(8192);
let state = dummy_state();
let join = tokio::spawn(async move {
let _ = run_with_streams(state, tools, server_r, server_w).await;
});
for req in requests {
let mut s = serde_json::to_vec(req).unwrap();
s.push(b'\n');
client_w.write_all(&s).await.unwrap();
}
drop(client_w);
let mut reader = BufReader::new(client_r);
let mut responses = Vec::new();
loop {
let mut line = String::new();
let n = reader.read_line(&mut line).await.unwrap();
if n == 0 {
break;
}
responses.push(serde_json::from_str(line.trim()).unwrap());
}
let _ = join.await;
responses
}
fn echo_tool() -> RegisteredTool {
RegisteredTool {
name: "echo".to_string(),
description: "Echo arguments back".to_string(),
input_schema: json!({"type": "object"}),
handler: std::sync::Arc::new(|_state, args| {
Box::pin(async move { Ok(json!({"echoed": args})) })
}),
}
}
fn failing_tool() -> RegisteredTool {
RegisteredTool {
name: "boom".to_string(),
description: "Always fails with a typed SessionError".to_string(),
input_schema: json!({"type": "object"}),
handler: std::sync::Arc::new(|_state, _args| {
Box::pin(async move {
Err(anyhow::Error::new(crate::errors::SessionError::TabHung {
target_id: Some("T1".into()),
url: Some("https://example.test".into()),
timeout_ms: 20_000,
hint: "renderer wedged",
}))
})
}),
}
}
#[tokio::test]
async fn initialize_round_trip() {
let resp = send_recv(
ToolRegistry::new(),
&[json!({"jsonrpc":"2.0","id":1,"method":"initialize","params":{}})],
)
.await;
assert_eq!(resp.len(), 1);
assert_eq!(resp[0]["id"], 1);
assert_eq!(resp[0]["result"]["protocolVersion"], "2024-11-05");
assert_eq!(resp[0]["result"]["serverInfo"]["name"], "browser-control");
}
#[tokio::test]
async fn tools_list_empty() {
let resp = send_recv(
ToolRegistry::new(),
&[json!({"jsonrpc":"2.0","id":2,"method":"tools/list"})],
)
.await;
assert_eq!(resp[0]["result"]["tools"], json!([]));
}
#[tokio::test]
async fn tools_list_after_register() {
let tools = ToolRegistry::new();
tools.register(echo_tool());
let resp = send_recv(
tools,
&[json!({"jsonrpc":"2.0","id":3,"method":"tools/list"})],
)
.await;
let list = resp[0]["result"]["tools"].as_array().unwrap();
assert_eq!(list.len(), 1);
assert_eq!(list[0]["name"], "echo");
}
#[tokio::test]
async fn tools_call_unknown_errors() {
let resp = send_recv(
ToolRegistry::new(),
&[json!({
"jsonrpc":"2.0","id":4,"method":"tools/call",
"params":{"name":"nope","arguments":{}}
})],
)
.await;
assert_eq!(resp[0]["error"]["code"], -32602);
assert!(resp[0]["error"]["message"]
.as_str()
.unwrap()
.contains("nope"));
}
#[tokio::test]
async fn tool_failure_returns_iserror_result_not_protocol_error() {
let tools = ToolRegistry::new();
tools.register(failing_tool());
let resp = send_recv(
tools,
&[json!({
"jsonrpc":"2.0","id":7,"method":"tools/call",
"params":{"name":"boom","arguments":{}}
})],
)
.await;
assert!(resp[0]["error"].is_null());
assert_eq!(resp[0]["result"]["isError"], true);
let text = resp[0]["result"]["content"][0]["text"].as_str().unwrap();
assert!(text.contains("tab hung"), "got: {text}");
assert!(text.contains("renderer wedged"), "got: {text}");
assert!(text.contains("T1"), "got: {text}");
}
#[tokio::test]
async fn tools_call_registered_returns_result() {
let tools = ToolRegistry::new();
tools.register(echo_tool());
let resp = send_recv(
tools,
&[json!({
"jsonrpc":"2.0","id":5,"method":"tools/call",
"params":{"name":"echo","arguments":{"hello":"world"}}
})],
)
.await;
assert_eq!(resp[0]["result"]["echoed"], json!({"hello":"world"}));
}
#[tokio::test]
async fn unknown_method_returns_minus_32601() {
let resp = send_recv(
ToolRegistry::new(),
&[json!({"jsonrpc":"2.0","id":6,"method":"bogus"})],
)
.await;
assert_eq!(resp[0]["error"]["code"], -32601);
}
#[tokio::test]
async fn ping_returns_empty_object() {
let resp = send_recv(
ToolRegistry::new(),
&[json!({"jsonrpc":"2.0","id":7,"method":"ping"})],
)
.await;
assert_eq!(resp[0]["result"], json!({}));
}
use futures_util::{SinkExt, StreamExt};
use std::sync::atomic::{AtomicUsize, Ordering};
use tokio_tungstenite::tungstenite::Message;
async fn spawn_counting_cdp_mock(
live: Vec<String>,
) -> (String, Arc<AtomicUsize>, tokio::sync::oneshot::Sender<()>) {
let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
let addr = listener.local_addr().unwrap();
let (stop_tx, mut stop_rx) = tokio::sync::oneshot::channel::<()>();
let conns = Arc::new(AtomicUsize::new(0));
let conns_srv = conns.clone();
tokio::spawn(async move {
loop {
let accept = tokio::select! {
_ = &mut stop_rx => break,
a = listener.accept() => a,
};
let (stream, _) = match accept {
Ok(s) => s,
Err(_) => break,
};
conns_srv.fetch_add(1, Ordering::SeqCst);
let live = live.clone();
tokio::spawn(async move {
let mut ws = match tokio_tungstenite::accept_async(stream).await {
Ok(w) => w,
Err(_) => return,
};
while let Some(Ok(msg)) = ws.next().await {
if let Message::Text(t) = msg {
let req: Value = match serde_json::from_str(&t) {
Ok(v) => v,
Err(_) => continue,
};
let id = req["id"].as_u64().unwrap_or(0);
let method = req["method"].as_str().unwrap_or("");
let result = match method {
"Target.getTargets" => {
let infos: Vec<Value> = live
.iter()
.map(|tid| {
json!({"targetId": tid, "type": "page", "url": ""})
})
.collect();
json!({"targetInfos": infos})
}
"Target.attachToTarget" => json!({"sessionId": "S1"}),
"Runtime.evaluate" => json!({"result": {"value": 1}}),
_ => json!({}),
};
let resp = json!({"id": id, "result": result});
if ws.send(Message::Text(resp.to_string())).await.is_err() {
break;
}
}
}
});
}
});
(format!("ws://{addr}"), conns, stop_tx)
}
fn registered_state(name: &str, endpoint: &str) -> ServerState {
ServerState::new(ResolvedBrowser {
endpoint: endpoint.to_string(),
engine: Engine::Cdp,
source: Source::Registered { name: name.into() },
})
}
#[allow(clippy::await_holding_lock)]
#[tokio::test]
async fn resolve_named_tab_live_resolves_and_touches() {
let _g = crate::test_support::ENV_LOCK
.lock()
.unwrap_or_else(|e| e.into_inner());
let tmp = tempfile::TempDir::new().unwrap();
std::env::set_var("BROWSER_CONTROL_DATA_DIR", tmp.path());
{
let reg = crate::registry::Registry::open().unwrap();
reg.tab_upsert("bx", "work", "T1", "about:blank", true)
.unwrap();
}
let (url, _conns, _stop) = spawn_counting_cdp_mock(vec!["T1".into()]).await;
let state = registered_state("bx", &url);
let target_id = match state.resolve_target_for_args(&json!({"tab": "work"})).await {
Ok((_backend, tid)) => tid,
Err(e) => panic!("named tab should resolve: {e:#}"),
};
assert_eq!(target_id, "T1");
std::env::remove_var("BROWSER_CONTROL_DATA_DIR");
}
#[allow(clippy::await_holding_lock)]
#[tokio::test]
async fn resolve_named_tab_stale_sweeps_and_errors() {
let _g = crate::test_support::ENV_LOCK
.lock()
.unwrap_or_else(|e| e.into_inner());
let tmp = tempfile::TempDir::new().unwrap();
std::env::set_var("BROWSER_CONTROL_DATA_DIR", tmp.path());
{
let reg = crate::registry::Registry::open().unwrap();
reg.tab_upsert("bx", "gone", "T_DEAD", "about:blank", true)
.unwrap();
}
let (url, _conns, _stop) = spawn_counting_cdp_mock(vec!["T1".into()]).await;
let state = registered_state("bx", &url);
let err = match state.resolve_target_for_args(&json!({"tab": "gone"})).await {
Ok(_) => panic!("stale named tab must error"),
Err(e) => e,
};
let typed = err
.downcast_ref::<crate::errors::SessionError>()
.expect("typed SessionError");
assert!(
matches!(typed, crate::errors::SessionError::TabNotFound { .. }),
"expected TabNotFound, got {typed:?}"
);
let reg = crate::registry::Registry::open().unwrap();
assert!(
reg.tab_get("bx", "gone").unwrap().is_none(),
"row not swept"
);
std::env::remove_var("BROWSER_CONTROL_DATA_DIR");
}
#[allow(clippy::await_holding_lock)]
#[tokio::test]
async fn resolve_named_tab_missing_row_errors() {
let _g = crate::test_support::ENV_LOCK
.lock()
.unwrap_or_else(|e| e.into_inner());
let tmp = tempfile::TempDir::new().unwrap();
std::env::set_var("BROWSER_CONTROL_DATA_DIR", tmp.path());
let (url, _conns, _stop) = spawn_counting_cdp_mock(vec!["T1".into()]).await;
let state = registered_state("bx", &url);
let err = match state.resolve_target_for_args(&json!({"tab": "nope"})).await {
Ok(_) => panic!("unknown named tab must error"),
Err(e) => e,
};
let typed = err
.downcast_ref::<crate::errors::SessionError>()
.expect("typed SessionError");
assert!(
matches!(typed, crate::errors::SessionError::TabNotFound { .. }),
"expected TabNotFound, got {typed:?}"
);
std::env::remove_var("BROWSER_CONTROL_DATA_DIR");
}
#[tokio::test]
async fn resolve_tab_and_target_mutually_exclusive() {
let state = dummy_state();
let err = match state
.resolve_target_for_args(&json!({"tab": "a", "target": "b"}))
.await
{
Ok(_) => panic!("tab+target must error"),
Err(e) => e,
};
assert!(
err.to_string().contains("mutually exclusive"),
"got: {err:#}"
);
}
#[tokio::test]
async fn concurrent_ensure_backend_opens_one_backend() {
let (url, conns, _stop) = spawn_counting_cdp_mock(vec!["T1".into()]).await;
let state = ServerState::new(ResolvedBrowser {
endpoint: url,
engine: Engine::Cdp,
source: Source::External,
});
let mut handles = Vec::new();
for _ in 0..8 {
let s = state.clone();
handles.push(tokio::spawn(async move { s.ensure_backend().await }));
}
for h in handles {
h.await.unwrap().expect("ensure_backend should succeed");
}
assert_eq!(
conns.load(Ordering::SeqCst),
1,
"expected exactly one backend (one WS connection) under concurrency"
);
}
#[allow(clippy::await_holding_lock)]
#[tokio::test]
async fn external_cdp_backend_does_not_open_registry() {
let _g = crate::test_support::ENV_LOCK
.lock()
.unwrap_or_else(|e| e.into_inner());
let stale_path = {
let tmp = tempfile::TempDir::new().unwrap();
tmp.path().to_path_buf()
};
std::env::set_var("BROWSER_CONTROL_DATA_DIR", &stale_path);
let (url, _conns, _stop) = spawn_counting_cdp_mock(vec!["T1".into()]).await;
let state = ServerState::new(ResolvedBrowser {
endpoint: url,
engine: Engine::Cdp,
source: Source::External,
});
let result = state.ensure_backend().await;
std::env::remove_var("BROWSER_CONTROL_DATA_DIR");
result.expect("external CDP backend must not depend on registry env");
}
#[tokio::test]
async fn initialized_notification_is_silently_ignored() {
let resp = send_recv(
ToolRegistry::new(),
&[
json!({"jsonrpc":"2.0","method":"notifications/initialized"}),
json!({"jsonrpc":"2.0","id":8,"method":"ping"}),
],
)
.await;
assert_eq!(resp.len(), 1);
assert_eq!(resp[0]["id"], 8);
}
}