tramex_tools/interface/websocket/
ws_connection.rs1use core::fmt::{Debug, Formatter};
3use ewebsock::{WsEvent, WsMessage, WsReceiver, WsSender};
4use std::vec;
5
6use crate::interface::interface_types::InterfaceTrait;
7use crate::interface::types::BaseMessage;
8use crate::tramex_error;
9use crate::{data::Data, errors::TramexError};
10
11use crate::interface::{layer::Layer, layer::Layers, log_get::LogGet, types::WebSocketLog};
12pub struct WsConnection {
14 pub ws_sender: WsSender,
16
17 pub ws_receiver: WsReceiver,
19
20 pub msg_id: u64,
22
23 pub connecting: bool,
25
26 pub asking_size_max: u64,
28
29 pub available: bool,
31
32 pub name: String,
34
35 pub auto_loading: bool,
37
38 pub waiting_for_response: bool,
40}
41
42impl WsConnection {
43 pub fn new(ws_sender: WsSender, ws_receiver: WsReceiver) -> Self {
45 Self {
46 ws_sender,
47 ws_receiver,
48 msg_id: 1,
49 connecting: true,
50 asking_size_max: 1024,
51 available: true,
52 name: "".to_string(),
53 auto_loading: true, waiting_for_response: false,
55 }
56 }
57
58 pub fn connect(url: &str, wakeup: impl Fn() + Send + Sync + 'static) -> Result<(WsSender, WsReceiver), String> {
62 let options = ewebsock::Options {
63 #[cfg(not(target_arch = "wasm32"))]
64 additional_headers: vec![("Origin".to_string(), "tramex".to_string())],
65 ..Default::default()
66 };
67 ewebsock::connect_with_wakeup(url, options, wakeup)
68 }
69
70 pub fn close_impl(&mut self) -> Result<(), TramexError> {
74 self.ws_sender.close();
75 Ok(())
76 }
77}
78
79impl InterfaceTrait for WsConnection {
80 fn get_more_data(&mut self, layer_list: Layers, _data: &mut Data) -> Result<(), Vec<TramexError>> {
81 if self.waiting_for_response {
83 return Ok(());
84 }
85
86 let msg = LogGet::new(self.msg_id, layer_list, self.asking_size_max);
87 log::debug!("📤 Sending log_get request #{}", self.msg_id);
88 match serde_json::to_string(&msg) {
89 Ok(msg_stringed) => {
90 log::debug!("📤 JSON payload: {}", msg_stringed);
91 self.ws_sender.send(WsMessage::Text(msg_stringed));
92 self.msg_id += 1;
93 self.waiting_for_response = true;
94 log::debug!("⏳ Waiting for response...");
95 }
96 Err(err) => {
97 log::error!("Error encoding message: {err:?}");
98 return Err(vec![tramex_error!(
99 err.to_string(),
100 crate::errors::ErrorCode::WebSocketErrorEncodingMessage
101 )]);
102 }
103 }
104 Ok(())
105 }
106
107 fn close(&mut self) -> Result<(), TramexError> {
108 self.close_impl()
109 }
110
111 fn supports_preloading(&self) -> bool {
112 false }
114
115 fn get_total_event_count(&self) -> Option<usize> {
116 None }
118
119 fn is_fully_read(&self) -> bool {
120 !self.available }
122}
123
124impl WsConnection {
125 pub fn should_request_more(&self) -> bool {
129 self.auto_loading && !self.waiting_for_response
130 }
131
132 pub fn try_recv(&mut self, data: &mut Data) -> Result<(), Vec<TramexError>> {
136 while let Some(event) = self.ws_receiver.try_recv() {
137 log::debug!("🔵 WebSocket event received");
139 self.connecting = false;
140 match event {
141 WsEvent::Message(msg) => {
142 self.available = true;
143 match msg {
144 WsMessage::Text(event_text) => {
145 let decoded: Result<WebSocketLog, serde_json::Error> = serde_json::from_str(&event_text);
147 match decoded {
148 Ok(decoded_data) => {
149 log::debug!("✅ Received log_get response with {} logs", decoded_data.logs.len());
150
151 self.waiting_for_response = false;
153
154 let mut errors = vec![];
155 for one_log in decoded_data.logs {
156 if Layer::RRC == one_log.layer
158 && one_log.data.iter().any(|line| line.contains("NR band combinations"))
159 {
160 log::debug!("Skipping NR band combinations log");
161 continue;
162 }
163 match one_log.extract_data() {
164 Ok(trace) => {
165 data.events.push(trace);
166 }
167 Err(err) => {
168 log::error!("Error while extracting data: {err:?}");
169 errors.push(err);
170 }
171 }
172 }
173 if !errors.is_empty() {
174 return Err(errors);
175 }
176 }
177 Err(_err) => {
178 let decoded_base: Result<BaseMessage, serde_json::Error> =
179 serde_json::from_str(&event_text);
180 match decoded_base {
181 Ok(decoded_data) => {
182 if decoded_data.message == "ready" {
183 log::debug!(
184 "✅ Received 'ready' message from server: {}",
185 decoded_data.name
186 );
187 log::debug!(
188 "💡 Server is ready. You need to click 'Load More' or enable auto-loading to request logs."
189 );
190 }
191 log::debug!("📥 Received BaseMessage: {decoded_data:?}");
192 self.name = decoded_data.name;
193 }
194 Err(err) => {
195 log::error!("Error decoding message: {err:?}");
196 log::error!("Message: {event_text:?}");
197 return Err(vec![tramex_error!(
198 err.to_string(),
199 crate::errors::ErrorCode::WebSocketErrorDecodingMessage
200 )]);
201 }
202 }
203 }
204 }
205 }
206 WsMessage::Unknown(str_error) => {
207 log::error!("Unknown message: {str_error:?}");
208 return Err(vec![tramex_error!(
209 str_error,
210 crate::errors::ErrorCode::WebSocketUnknownMessageReceived
211 )]);
212 }
213 WsMessage::Binary(bin) => {
214 log::error!("Unknown binary message: {bin:?}");
215 return Err(vec![tramex_error!(
216 format!("Unknown binary message: {bin:?}"),
217 crate::errors::ErrorCode::WebSocketUnknownBinaryMessageReceived
218 )]);
219 }
220 _ => {
221 log::debug!("Received Ping-Pong")
222 }
223 }
224 }
225 WsEvent::Opened => {
226 self.available = true;
227 log::debug!("✅ WebSocket connection opened successfully");
228 }
229 WsEvent::Closed => {
230 self.available = false;
231 log::debug!("WebSocket closed");
232 return Err(vec![tramex_error!(
233 "WebSocket closed".to_string(),
234 crate::errors::ErrorCode::WebSocketClosed
235 )]);
236 }
237 WsEvent::Error(str_err) => {
238 self.available = false;
239 log::error!("WebSocket error: {str_err:?}");
240 return Err(vec![tramex_error!(str_err, crate::errors::ErrorCode::WebSocketError)]);
241 }
242 }
243 }
244 Ok(())
245 }
246}
247
248impl Debug for WsConnection {
249 fn fmt(&self, formatter: &mut Formatter<'_>) -> core::fmt::Result {
250 formatter
251 .debug_struct("Interface")
252 .field("ws_sender", &"Box<WsSender>")
253 .field("ws_receiver", &"Box<WsReceiver>")
254 .field("connecting", &self.connecting)
255 .finish()
256 }
257}
258
259impl Drop for WsConnection {
260 fn drop(&mut self) {
261 log::debug!("Cleaning WsConnection");
262 self.ws_sender.close()
263 }
264}