Skip to main content

cdp_server/
event.rs

1// @trace REQ-CDS-005 [entity:EventSubscription]
2// Event broadcaster: domain-based subscription filtering.
3
4use std::collections::HashMap;
5use std::sync::{Arc, Mutex};
6
7use serde_json::Value;
8
9use crate::protocol::{serialize_event, CdpEvent};
10use crate::session::{OutboxEvent, SessionHandle};
11use crate::EventSender;
12
13type SessionMap = Arc<Mutex<HashMap<String, Arc<SessionHandle>>>>;
14
15/// EventBroadcaster implements EventSender. It holds a reference to the
16/// session map and queues events into per-session outboxes; the server loop
17/// drains the outboxes into the WebSocket while holding the session lock.
18///
19/// Events are NEVER written to the socket directly here: a command dispatch
20/// running inside `CdpSession::process` may emit events for that very
21/// session — taking the session lock here would self-deadlock the server
22/// loop. Outbox + drain-at-send-time preserves the domain gating (applied
23/// when the drain holds the session).
24pub struct EventBroadcaster {
25    sessions: SessionMap,
26}
27
28impl EventBroadcaster {
29    pub fn new(sessions: SessionMap) -> Self {
30        EventBroadcaster { sessions }
31    }
32
33    /// Create a boxed clone-safe EventSender reference.
34    pub fn sender(&self) -> Box<dyn EventSender> {
35        Box::new(EventBroadcaster {
36            sessions: Arc::clone(&self.sessions),
37        })
38    }
39
40    fn enqueue(&self, entry: OutboxEvent) {
41        let sessions = match self.sessions.lock() {
42            Ok(s) => s,
43            Err(_) => return,
44        };
45        for handle in sessions.values() {
46            if let Ok(mut outbox) = handle.outbox.lock() {
47                outbox.push_back(entry.clone());
48            }
49        }
50    }
51}
52
53impl EventSender for EventBroadcaster {
54    fn send_event(&self, method: &str, params: Value) {
55        let domain = method.split('.').next().unwrap_or("").to_string();
56        let event = CdpEvent {
57            method: method.to_string(),
58            params: Some(params),
59        };
60        self.enqueue(OutboxEvent {
61            json: serialize_event(&event),
62            domain,
63            browser_only: false,
64        });
65    }
66
67    /// Session-scoped event (flattened CDP sessions): the event JSON carries
68    /// `sessionId`, so clients route it to the attached target session.
69    /// Delivered to browser-endpoint sessions (they own the flat sessions).
70    fn send_session_event(&self, session_id: &str, method: &str, params: Value) {
71        let domain = method.split('.').next().unwrap_or("").to_string();
72        let json = serde_json::json!({
73            "method": method,
74            "params": params,
75            "sessionId": session_id,
76        })
77        .to_string();
78        self.enqueue(OutboxEvent {
79            json,
80            domain,
81            browser_only: true,
82        });
83    }
84}
85
86// Clone: Arc-based shallow copy.
87impl Clone for EventBroadcaster {
88    fn clone(&self) -> Self {
89        EventBroadcaster {
90            sessions: Arc::clone(&self.sessions),
91        }
92    }
93}
94
95#[cfg(test)]
96mod tests {
97    use super::*;
98    use std::collections::HashMap;
99
100    fn empty_session_map() -> SessionMap {
101        Arc::new(Mutex::new(HashMap::new()))
102    }
103
104    // @trace TEST-CDS-005 [req:REQ-CDS-005] [level:unit]
105    #[test]
106    fn new_with_empty_sessions_no_panic() {
107        let _broadcaster = EventBroadcaster::new(empty_session_map());
108    }
109
110    #[test]
111    fn sender_returns_boxed_event_sender() {
112        let broadcaster = EventBroadcaster::new(empty_session_map());
113        let _sender: Box<dyn EventSender> = broadcaster.sender();
114    }
115
116    #[test]
117    fn send_event_empty_sessions_no_panic() {
118        let broadcaster = EventBroadcaster::new(empty_session_map());
119        broadcaster.send_event("Page.loadEventFired", serde_json::json!({}));
120    }
121
122    #[test]
123    fn send_session_event_empty_sessions_no_panic() {
124        let broadcaster = EventBroadcaster::new(empty_session_map());
125        broadcaster.send_session_event("sid", "Runtime.executionContextCreated", serde_json::json!({}));
126    }
127
128    #[test]
129    fn clone_shares_sessions_arc() {
130        let sessions = empty_session_map();
131        let a = EventBroadcaster::new(Arc::clone(&sessions));
132        let b = a.clone();
133        assert!(Arc::ptr_eq(&a.sessions, &b.sessions));
134    }
135
136    #[test]
137    fn send_event_method_domain_extraction_unit_test() {
138        assert_eq!("Page".split('.').next().unwrap_or(""), "Page");
139        assert_eq!(
140            "Runtime.consoleAPICalled".split('.').next().unwrap_or(""),
141            "Runtime"
142        );
143        assert_eq!(
144            "no_dot_method".split('.').next().unwrap_or(""),
145            "no_dot_method"
146        );
147        assert_eq!("".split('.').next().unwrap_or(""), "");
148    }
149}