1use 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
15pub struct EventBroadcaster {
25 sessions: SessionMap,
26}
27
28impl EventBroadcaster {
29 pub fn new(sessions: SessionMap) -> Self {
30 EventBroadcaster { sessions }
31 }
32
33 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 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
86impl 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 #[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}