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, parse_config::FileMetadata, 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 pub headers_received: bool,
43}
44
45impl WsConnection {
46 pub fn new(ws_sender: WsSender, ws_receiver: WsReceiver) -> Self {
48 log::info!("🔌 WsConnection::new() - creating connection (connecting=true, available=false)");
49 Self {
50 ws_sender,
51 ws_receiver,
52 msg_id: 1,
53 connecting: true,
54 asking_size_max: 1024,
55 available: false, name: "".to_string(),
57 auto_loading: true, waiting_for_response: false,
59 headers_received: false,
60 }
61 }
62
63 pub fn connect(url: &str, wakeup: impl Fn() + Send + Sync + 'static) -> Result<(WsSender, WsReceiver), String> {
67 let options = ewebsock::Options {
68 #[cfg(not(target_arch = "wasm32"))]
69 additional_headers: vec![("Origin".to_string(), "tramex".to_string())],
70 ..Default::default()
71 };
72 ewebsock::connect_with_wakeup(url, options, wakeup)
73 }
74
75 pub fn close_impl(&mut self) -> Result<(), TramexError> {
79 self.ws_sender.close();
80 Ok(())
81 }
82}
83
84impl InterfaceTrait for WsConnection {
85 fn get_more_data(&mut self, layer_list: Layers, _data: &mut Data) -> Result<(), Vec<TramexError>> {
86 if self.waiting_for_response {
88 return Ok(());
89 }
90
91 let msg = LogGet::new(self.msg_id, layer_list, self.asking_size_max);
92 let msg = if !self.headers_received { msg.with_headers() } else { msg };
94 log::debug!("📤 Sending log_get request #{}", self.msg_id);
95 match serde_json::to_string(&msg) {
96 Ok(msg_stringed) => {
97 log::debug!("📤 JSON payload: {}", msg_stringed);
98 self.ws_sender.send(WsMessage::Text(msg_stringed));
99 self.msg_id += 1;
100 self.waiting_for_response = true;
101 log::debug!("⏳ Waiting for response...");
102 }
103 Err(err) => {
104 log::error!("Error encoding message: {err:?}");
105 return Err(vec![tramex_error!(
106 err.to_string(),
107 crate::errors::ErrorCode::WebSocketErrorEncodingMessage
108 )]);
109 }
110 }
111 Ok(())
112 }
113
114 fn close(&mut self) -> Result<(), TramexError> {
115 self.close_impl()
116 }
117
118 fn supports_preloading(&self) -> bool {
119 false }
121
122 fn get_total_event_count(&self) -> Option<usize> {
123 None }
125
126 fn is_fully_read(&self) -> bool {
127 !self.available }
129}
130
131impl WsConnection {
132 pub fn should_request_more(&self) -> bool {
136 self.auto_loading && !self.waiting_for_response
137 }
138
139 pub fn try_recv(&mut self, data: &mut Data) -> Result<(), Vec<TramexError>> {
143 while let Some(event) = self.ws_receiver.try_recv() {
144 log::debug!(
145 "🔵 WsConnection::try_recv() - event received, current state: connecting={}, available={}",
146 self.connecting,
147 self.available
148 );
149 match event {
150 WsEvent::Message(msg) => {
151 if !self.connecting {
154 self.available = true;
155 }
156 match msg {
157 WsMessage::Text(event_text) => {
158 let decoded: Result<WebSocketLog, serde_json::Error> = serde_json::from_str(&event_text);
160 match decoded {
161 Ok(decoded_data) => {
162 log::debug!("✅ Received log_get response with {} logs", decoded_data.logs.len());
163
164 self.waiting_for_response = false;
166
167 if !self.headers_received
169 && let Some(ref headers) = decoded_data.headers
170 {
171 data.metadata = FileMetadata::parse_from_lines(headers);
172 log::info!(
173 "📋 Parsed WebSocket metadata: technology={:?}, pci={:?}, mode={:?}",
174 data.metadata.technology,
175 data.metadata.pci,
176 data.metadata.mode
177 );
178 self.headers_received = true;
179 }
180
181 let mut errors = vec![];
182 for one_log in decoded_data.logs {
183 if Layer::RRC == one_log.layer
185 && one_log.data.iter().any(|line| line.contains("NR band combinations"))
186 {
187 log::debug!("Skipping NR band combinations log");
188 continue;
189 }
190 match one_log.extract_data() {
191 Ok(trace) => {
192 data.events.push(trace);
193 }
194 Err(err) => {
195 log::error!("Error while extracting data: {err:?}");
196 errors.push(err);
197 }
198 }
199 }
200 if !errors.is_empty() {
201 return Err(errors);
202 }
203 }
204 Err(_err) => {
205 let decoded_base: Result<BaseMessage, serde_json::Error> =
206 serde_json::from_str(&event_text);
207 match decoded_base {
208 Ok(decoded_data) => {
209 if decoded_data.message == "ready" {
210 log::debug!(
211 "✅ Received 'ready' message from server: {}",
212 decoded_data.name
213 );
214 log::debug!(
215 "💡 Server is ready. You need to click 'Load More' or enable auto-loading to request logs."
216 );
217 }
218 log::debug!("📥 Received BaseMessage: {decoded_data:?}");
219 self.name = decoded_data.name;
220 }
221 Err(err) => {
222 log::error!("Error decoding message: {err:?}");
223 log::error!("Message: {event_text:?}");
224 return Err(vec![tramex_error!(
225 err.to_string(),
226 crate::errors::ErrorCode::WebSocketErrorDecodingMessage
227 )]);
228 }
229 }
230 }
231 }
232 }
233 WsMessage::Unknown(str_error) => {
234 log::error!("Unknown message: {str_error:?}");
235 return Err(vec![tramex_error!(
236 str_error,
237 crate::errors::ErrorCode::WebSocketUnknownMessageReceived
238 )]);
239 }
240 WsMessage::Binary(bin) => {
241 log::error!("Unknown binary message: {bin:?}");
242 return Err(vec![tramex_error!(
243 format!("Unknown binary message: {bin:?}"),
244 crate::errors::ErrorCode::WebSocketUnknownBinaryMessageReceived
245 )]);
246 }
247 _ => {
248 log::debug!("Received Ping-Pong")
249 }
250 }
251 }
252 WsEvent::Opened => {
253 log::info!("✅ WsEvent::Opened - connection established! Setting connecting=false, available=true");
254 self.connecting = false;
255 self.available = true;
256 }
257 WsEvent::Closed => {
258 log::warn!("⚠️ WsEvent::Closed - connection closed. Setting connecting=false, available=false");
259 self.available = false;
260 self.connecting = false;
261 return Err(vec![tramex_error!(
262 "WebSocket closed".to_string(),
263 crate::errors::ErrorCode::WebSocketClosed
264 )]);
265 }
266 WsEvent::Error(str_err) => {
267 log::error!(
268 "❌ WsEvent::Error - connection failed: {}. Setting connecting=false, available=false",
269 str_err
270 );
271 self.available = false;
272 self.connecting = false;
273 return Err(vec![tramex_error!(str_err, crate::errors::ErrorCode::WebSocketError)]);
274 }
275 }
276 }
277 Ok(())
278 }
279}
280
281impl Debug for WsConnection {
282 fn fmt(&self, formatter: &mut Formatter<'_>) -> core::fmt::Result {
283 formatter
284 .debug_struct("Interface")
285 .field("ws_sender", &"Box<WsSender>")
286 .field("ws_receiver", &"Box<WsReceiver>")
287 .field("connecting", &self.connecting)
288 .finish()
289 }
290}
291
292impl Drop for WsConnection {
293 fn drop(&mut self) {
294 log::debug!("Cleaning WsConnection");
295 self.ws_sender.close()
296 }
297}