use std::sync::mpsc::{self, Receiver, RecvTimeoutError, Sender};
use std::sync::Arc;
use std::time::Duration;
use serde_json::Value;
use crate::bridge::event_translator::{translate, ServoEvent};
use crate::error::{CdpError, Result};
use super::r#trait::{CdpEvent, Transport};
use super::TransportKind;
#[derive(Debug, Clone)]
pub enum InMemoryBridgeResponse {
Ok(Value),
Err(String),
}
pub trait InMemoryBridge: Send + Sync {
fn dispatch_command(
&self,
method: &str,
params: Value,
session_id: Option<&str>,
) -> InMemoryBridgeResponse;
}
pub struct InMemoryTransport {
bridge: Arc<dyn InMemoryBridge>,
event_tx: Sender<CdpEvent>,
event_rx: Receiver<CdpEvent>,
servo_event_rx: Option<Receiver<ServoEvent>>,
pending_cdp_events: std::collections::VecDeque<CdpEvent>,
closed: bool,
command_timeout: Duration,
event_timeout: Duration,
}
impl std::fmt::Debug for InMemoryTransport {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("InMemoryTransport")
.field("closed", &self.closed)
.field("command_timeout", &self.command_timeout)
.field("event_timeout", &self.event_timeout)
.finish()
}
}
impl InMemoryTransport {
pub fn new(bridge: Arc<dyn InMemoryBridge>) -> Self {
let (event_tx, event_rx) = mpsc::channel::<CdpEvent>();
Self {
bridge,
event_tx,
event_rx,
servo_event_rx: None,
pending_cdp_events: std::collections::VecDeque::new(),
closed: false,
command_timeout: Duration::from_secs(30),
event_timeout: Duration::from_millis(100),
}
}
pub fn event_sender(&self) -> Sender<CdpEvent> {
self.event_tx.clone()
}
pub fn attach_servo_event_receiver(&mut self, rx: Receiver<ServoEvent>) {
self.servo_event_rx = Some(rx);
}
pub fn is_closed(&self) -> bool {
self.closed
}
}
impl Transport for InMemoryTransport {
fn kind(&self) -> TransportKind {
TransportKind::InMemory
}
fn send_command(
&mut self,
method: &str,
params: Value,
session_id: Option<&str>,
) -> Result<Value> {
if self.closed {
return Err(CdpError::ConnectionClosed);
}
match self.bridge.dispatch_command(method, params, session_id) {
InMemoryBridgeResponse::Ok(v) => Ok(v),
InMemoryBridgeResponse::Err(msg) => Err(CdpError::ProtocolError(msg)),
}
}
fn recv_event(&mut self) -> Result<Option<CdpEvent>> {
if self.closed {
return Err(CdpError::ConnectionClosed);
}
if let Some(ev) = self.pending_cdp_events.pop_front() {
return Ok(Some(ev));
}
if let Some(servo_rx) = &self.servo_event_rx {
match servo_rx.recv_timeout(self.event_timeout) {
Ok(se) => {
let mut cdp_events = translate(se);
if let Some(first) = cdp_events.pop() {
for ev in cdp_events.into_iter().rev() {
self.pending_cdp_events.push_front(ev);
}
return Ok(Some(first));
}
return self.recv_event();
}
Err(RecvTimeoutError::Timeout) => {
}
Err(RecvTimeoutError::Disconnected) => {
}
}
}
match self.event_rx.recv_timeout(self.event_timeout) {
Ok(ev) => Ok(Some(ev)),
Err(RecvTimeoutError::Timeout) => Ok(None),
Err(RecvTimeoutError::Disconnected) => Err(CdpError::ConnectionClosed),
}
}
fn close(&mut self) -> Result<()> {
if !self.closed {
self.closed = true;
}
Ok(())
}
fn set_command_timeout(&mut self, timeout: Duration) {
self.command_timeout = timeout;
}
fn set_event_timeout(&mut self, timeout: Duration) {
self.event_timeout = timeout;
}
}
#[cfg(test)]
mod mock_bridge {
use super::*;
use std::sync::Mutex;
pub struct MockInMemoryBridge {
pub history: Mutex<Vec<(String, Value, Option<String>)>>,
pub responder:
Box<dyn Fn(&str, &Value, Option<&str>) -> InMemoryBridgeResponse + Send + Sync>,
}
impl MockInMemoryBridge {
pub fn new<F>(responder: F) -> Self
where
F: Fn(&str, &Value, Option<&str>) -> InMemoryBridgeResponse + Send + Sync + 'static,
{
Self {
history: Mutex::new(Vec::new()),
responder: Box::new(responder),
}
}
pub fn ok_null() -> Self {
Self::new(|_, _, _| InMemoryBridgeResponse::Ok(Value::Null))
}
pub fn always_err(msg: impl Into<String>) -> Self {
let msg = msg.into();
Self::new(move |_, _, _| InMemoryBridgeResponse::Err(msg.clone()))
}
}
impl InMemoryBridge for MockInMemoryBridge {
fn dispatch_command(
&self,
method: &str,
params: Value,
session_id: Option<&str>,
) -> InMemoryBridgeResponse {
self.history.lock().unwrap().push((
method.to_string(),
params.clone(),
session_id.map(|s| s.to_string()),
));
(self.responder)(method, ¶ms, session_id)
}
}
#[test]
fn in_memory_kind_is_in_memory() {
let bridge = Arc::new(MockInMemoryBridge::ok_null());
let mut t = InMemoryTransport::new(bridge);
assert_eq!(t.kind(), TransportKind::InMemory);
assert!(!t.is_closed());
let _ = t.close();
assert!(t.is_closed());
}
#[test]
fn in_memory_send_command_returns_ok_response() {
let bridge = Arc::new(MockInMemoryBridge::new(|_m, _p, _s| {
InMemoryBridgeResponse::Ok(serde_json::json!({"title": "Test Page"}))
}));
let mut t = InMemoryTransport::new(bridge);
let r = t
.send_command("Page.getTitle", serde_json::json!({}), None)
.unwrap();
assert_eq!(r["title"], "Test Page");
}
#[test]
fn in_memory_send_command_propagates_error() {
let bridge = Arc::new(MockInMemoryBridge::always_err("method not found"));
let mut t = InMemoryTransport::new(bridge);
let err = t
.send_command("Unknown.method", serde_json::json!({}), None)
.unwrap_err();
assert!(matches!(err, CdpError::ProtocolError(_)));
assert!(err.to_string().contains("method not found"));
}
#[test]
fn in_memory_send_command_after_close_returns_connection_closed() {
let bridge = Arc::new(MockInMemoryBridge::ok_null());
let mut t = InMemoryTransport::new(bridge);
t.close().unwrap();
let err = t
.send_command("X", serde_json::json!({}), None)
.unwrap_err();
assert!(matches!(err, CdpError::ConnectionClosed));
}
#[test]
fn in_memory_recv_event_gets_pushed_event() {
let bridge = Arc::new(MockInMemoryBridge::ok_null());
let mut t = InMemoryTransport::new(bridge);
let sender = t.event_sender();
sender
.send(CdpEvent::new(
"Page.frameNavigated",
serde_json::json!({"url": "x"}),
))
.unwrap();
let ev = t.recv_event().unwrap().expect("expected an event");
assert_eq!(ev.method, "Page.frameNavigated");
assert_eq!(ev.params["url"], "x");
}
#[test]
fn in_memory_recv_event_returns_none_on_timeout() {
let bridge = Arc::new(MockInMemoryBridge::ok_null());
let mut t = InMemoryTransport::new(bridge);
let ev = t.recv_event().unwrap();
assert!(ev.is_none());
}
#[test]
fn in_memory_recv_event_after_close_returns_connection_closed() {
let bridge = Arc::new(MockInMemoryBridge::ok_null());
let mut t = InMemoryTransport::new(bridge);
t.close().unwrap();
let err = t.recv_event().unwrap_err();
assert!(matches!(err, CdpError::ConnectionClosed));
}
#[test]
fn in_memory_close_is_idempotent() {
let bridge = Arc::new(MockInMemoryBridge::ok_null());
let mut t = InMemoryTransport::new(bridge);
t.close().unwrap();
t.close().unwrap();
}
#[test]
fn in_memory_session_id_passed_to_bridge() {
let bridge = Arc::new(MockInMemoryBridge::new(|_m, _p, s| {
let sid = s.unwrap_or("default");
InMemoryBridgeResponse::Ok(serde_json::json!({"echo": sid}))
}));
let mut t = InMemoryTransport::new(bridge);
let r = t
.send_command("X.y", serde_json::json!({}), Some("TARGET-99"))
.unwrap();
assert_eq!(r["echo"], "TARGET-99");
}
#[test]
fn in_memory_event_timeout_overridable() {
let bridge = Arc::new(MockInMemoryBridge::ok_null());
let mut t = InMemoryTransport::new(bridge);
t.set_event_timeout(Duration::from_millis(1));
let start = std::time::Instant::now();
let ev = t.recv_event().unwrap();
let elapsed = start.elapsed();
assert!(ev.is_none());
assert!(elapsed.as_millis() < 200, "elapsed: {:?}", elapsed);
}
#[test]
fn in_memory_bridge_response_debug() {
let r1 = InMemoryBridgeResponse::Ok(Value::Null);
let r2 = InMemoryBridgeResponse::Err("boom".into());
let s1 = format!("{:?}", r1);
let s2 = format!("{:?}", r2);
assert!(s1.contains("Ok"));
assert!(s2.contains("Err"));
}
#[test]
fn in_memory_transport_debug_format() {
let bridge = Arc::new(MockInMemoryBridge::ok_null());
let t = InMemoryTransport::new(bridge);
let s = format!("{:?}", t);
assert!(s.contains("InMemoryTransport"));
assert!(s.contains("closed"));
}
}