From b6ec1655587bb4fc299265c7218677ac4ecc115a Mon Sep 17 00:00:00 2001 From: Alex Emmet <111742636+Alex-Emmet@users.noreply.github.com> Date: Thu, 26 Mar 2026 02:01:26 +0100 Subject: [PATCH] [Fix] Hopefully omikron - iota comms work now ._. --- Cargo.lock | 4 +- src/main.rs | 3 +- src/omikron/omikron_connection.rs | 238 +++++++++++++++++++++++++----- src/omikron/ping_pong_task.rs | 16 +- 4 files changed, 216 insertions(+), 45 deletions(-) diff --git a/Cargo.lock b/Cargo.lock index 8143cc4..cea2fc0 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -3876,7 +3876,7 @@ checksum = "e421abadd41a4225275504ea4d6566923418b7f05506fbc9c0fe86ba7396114b" [[package]] name = "ttp-core" version = "0.1.0" -source = "git+https://github.com/Tensamin/TTP.git#7e46b440f847e8bf1ad1ce21932ce33a275dfcd1" +source = "git+https://github.com/Tensamin/TTP.git#c127d196082401f37b51ca5417b341ecbdb74318" dependencies = [ "base64", "byteorder", @@ -3888,7 +3888,7 @@ dependencies = [ [[package]] name = "ttp-native" version = "0.1.0" -source = "git+https://github.com/Tensamin/TTP.git#7e46b440f847e8bf1ad1ce21932ce33a275dfcd1" +source = "git+https://github.com/Tensamin/TTP.git#c127d196082401f37b51ca5417b341ecbdb74318" dependencies = [ "quinn", "rustls", diff --git a/src/main.rs b/src/main.rs index ae02090..767106f 100644 --- a/src/main.rs +++ b/src/main.rs @@ -155,8 +155,7 @@ async fn main() { ) .await; } - let omikron: Arc = Arc::new(OmikronConnection::new()); - omikron.connect().await; + let _ = omikron::omikron_connection::get_omikron_connection().await; log_t!("setup_completed"); loop { diff --git a/src/omikron/omikron_connection.rs b/src/omikron/omikron_connection.rs index af8ada4..6ecd05a 100755 --- a/src/omikron/omikron_connection.rs +++ b/src/omikron/omikron_connection.rs @@ -119,12 +119,22 @@ impl OmikronConnection { // ------------------------------------------------------------------------- pub async fn connect(self: &Arc) { + log!( + "OmikronConnection::connect [id={}, ptr={:p}]", + self.connection_id, + Arc::as_ptr(self) + ); if self.connection_loop_handle.lock().await.is_none() { self.clone().start().await; } } pub async fn start(self: Arc) { + log!( + "OmikronConnection::start [id={}, ptr={:p}]", + self.connection_id, + Arc::as_ptr(&self) + ); if let Some(handle) = self.connection_loop_handle.lock().await.take() { handle.abort(); } @@ -140,6 +150,11 @@ impl OmikronConnection { } pub async fn stop(&self) { + log!( + "OmikronConnection::stop [id={}, ptr={:p}]", + self.connection_id, + self + ); *self.reconnect_on_close.write().await = false; if let Some(tx) = self.shutdown_tx.lock().await.take() { @@ -159,6 +174,11 @@ impl OmikronConnection { } *self.state.write().await = ConnectionState::Disconnected; + log!( + "OmikronConnection::stop - setting sender to None [id={}, ptr={:p}]", + self.connection_id, + self + ); *self.sender.write().await = None; } @@ -221,6 +241,12 @@ impl OmikronConnection { let sender_arc = Arc::new(sender); *self.sender.write().await = Some(sender_arc.clone()); + log!( + "OmikronConnection::connect_once - sender set [id={}, ptr={:p}, sender_open={}]", + self.connection_id, + Arc::as_ptr(&self), + sender_arc.is_open() + ); *self.state.write().await = ConnectionState::Connected { identified: false }; // Handle registration/identification @@ -247,6 +273,11 @@ impl OmikronConnection { let result = read_handle.await; // Cleanup + log!( + "OmikronConnection::connect_once - setting sender to None (cleanup) [id={}, ptr={:p}]", + self.connection_id, + Arc::as_ptr(&self) + ); *self.sender.write().await = None; *self.state.write().await = ConnectionState::Disconnected; { @@ -345,6 +376,12 @@ impl OmikronConnection { // ------------------------------------------------------------------------- async fn read_loop(self: Arc, receiver: &mut Receiver) { + log!( + "OmikronConnection::read_loop [id={}, ptr={:p}, receiver_open={}]", + self.connection_id, + Arc::as_ptr(&self), + receiver.is_open() + ); loop { let result = receiver.receive().await; match result { @@ -353,22 +390,39 @@ impl OmikronConnection { } Err(e) => { log!( - "Receive error: {} connection_id={} receiver_open={}", + "Receive error: {} connection_id={} receiver_open={} ptr={:p}", e, self.connection_id, - receiver.is_open() + receiver.is_open(), + Arc::as_ptr(&self) ); + self.fail_all_waiting_tasks(format!( + "Connection receive error: {} (connection_id={})", + e, self.connection_id + )) + .await; break; } } if !receiver.is_open() { log!( - "Connection closed connection_id={} receiver_open=false", - self.connection_id + "Connection closed connection_id={} receiver_open=false ptr={:p}", + self.connection_id, + Arc::as_ptr(&self) ); + self.fail_all_waiting_tasks(format!( + "Connection closed (connection_id={}, receiver_open=false)", + self.connection_id + )) + .await; break; } } + log!( + "OmikronConnection::read_loop exit [id={}, ptr={:p}]", + self.connection_id, + Arc::as_ptr(&self) + ); } async fn heartbeat_loop(self: Arc) { @@ -396,23 +450,36 @@ impl OmikronConnection { // ------------------------------------------------------------------------- pub async fn handle_message(self: Arc, cv: CommunicationValue) { + log!( + "OmikronConnection::handle_message [id={}, ptr={:p}, type={:?}, msg_id={}]", + self.connection_id, + Arc::as_ptr(&self), + cv.get_type(), + cv.get_id() + ); if !cv.is_type(CommunicationType::ping) && !cv.is_type(CommunicationType::pong) { log_cv_in!(&cv); } let msg_id = cv.get_id(); - // Only check new waiting tasks - if let Some((_, task)) = WAITING_TASKS.remove(&msg_id) { - if (task.task)(self.clone(), cv.clone()) { - return; - } - } - - // Check new waiting tasks + // Dispatch waiting task for this message id if let Some((_, task)) = WAITING_TASKS.remove(&msg_id) { + log!( + "OmikronConnection::handle_message - found waiting task for msg_id={} [id={}, ptr={:p}]", + msg_id, + self.connection_id, + Arc::as_ptr(&self) + ); if (task.task)(self.clone(), cv.clone()) { return; + } else { + log!( + "Waiting task for msg_id={} returned false, continuing normal handling [id={}, ptr={:p}]", + msg_id, + self.connection_id, + Arc::as_ptr(&self) + ); } } @@ -426,30 +493,16 @@ impl OmikronConnection { return; } - if cv.is_type(CommunicationType::success) { - let iota_id = cv.get_data(DataTypes::iota_id).as_number().unwrap_or(0); - if iota_id != 0 { - let mut conf = CONFIG.write().await; - conf.change("iota_id", JsonValue::from(iota_id as i64)); - conf.update(); - log!("Iota registered with ID: {}", iota_id); - - let login_message = CommunicationValue::new(CommunicationType::identification) - .add_data(DataTypes::iota_id, DataValue::Number(iota_id)); - - let self_clone = self.clone(); - tokio::spawn(async move { - self_clone.send_message(&login_message).await; - }); - } - return; - } - if cv.is_type(CommunicationType::identification_response) { if let Some(accepted) = cv.get_data(DataTypes::accepted).as_bool() { let mut state = self.state.write().await; if let ConnectionState::Connected { identified: _ } = *state { *state = ConnectionState::Connected { identified: true }; + log!( + "OmikronConnection - identification accepted! [id={}, ptr={:p}]", + self.connection_id, + Arc::as_ptr(&self) + ); } } return; @@ -751,6 +804,14 @@ impl OmikronConnection { self.send_message(&response).await; return; } + + log!( + "OmikronConnection::handle_message - unhandled message type [id={}, ptr={:p}, type={:?}, msg_id={}]", + self.connection_id, + self, + cv.get_type(), + cv.get_id() + ); } async fn handle_challenge(&self, cv: &CommunicationValue) { @@ -795,28 +856,79 @@ impl OmikronConnection { // ------------------------------------------------------------------------- pub async fn send_message(&self, cv: &CommunicationValue) { + if let Err(err) = self.send_message_result(cv).await { + log_t!("send_message_failed", err); + } + } + + async fn send_message_result(&self, cv: &CommunicationValue) -> Result<(), String> { let sender_guard = self.sender.read().await; + log!( + "OmikronConnection::send_message_result [id={}, ptr={:p}, has_sender={}, sender_open={}]", + self.connection_id, + self, + sender_guard.is_some(), + sender_guard.as_ref().map(|s| s.is_open()).unwrap_or(false) + ); if let Some(sender) = sender_guard.as_ref() { if !sender.is_open() { - log_t!("send_message_failed", "connection closed".to_string()); drop(sender_guard); + log!( + "OmikronConnection::send_message_result - sender not open, setting to None [id={}, ptr={:p}]", + self.connection_id, + self + ); if let Some(sender) = self.sender.write().await.take() { sender.close(); } - return; + self.fail_all_waiting_tasks(format!( + "Send failed: connection closed (connection_id={})", + self.connection_id + )) + .await; + return Err("connection closed".to_string()); } let sender_clone = Arc::clone(sender); drop(sender_guard); + log!( + "OmikronConnection::send_message_result - sending [id={}, ptr={:p}, type={:?}, msg_id={}]", + self.connection_id, + self, + cv.get_type(), + cv.get_id() + ); + if !cv.is_type(CommunicationType::ping) && !cv.is_type(CommunicationType::pong) { log_cv_out!(&cv); } + if let Err(e) = sender_clone.send(cv).await { - log_t!("send_message_failed", e.to_string()); + self.fail_all_waiting_tasks(format!( + "Send failed: {} (connection_id={})", + e, self.connection_id + )) + .await; + return Err(e.to_string()); } + + Ok(()) } else { - log_t!("send_message_failed", "not connected".to_string()); + Err("not connected".to_string()) + } + } + + async fn fail_all_waiting_tasks(&self, reason: String) { + let keys: Vec = WAITING_TASKS.iter().map(|entry| *entry.key()).collect(); + + for key in keys { + if let Some((_, waiting_task)) = WAITING_TASKS.remove(&key) { + let response = CommunicationValue::new(CommunicationType::error) + .with_id(key) + .add_data(DataTypes::message, DataValue::Str(reason.clone())); + let _ = (waiting_task.task)(OMIKRON_CONNECTION.clone(), response); + } } } @@ -836,6 +948,12 @@ impl OmikronConnection { let (tx, mut rx) = mpsc::channel(1); let msg_id = cv.get_id(); + log!( + "OmikronConnection::await_response - inserting waiting task for msg_id={} [id={}, ptr={:p}]", + msg_id, + self.connection_id, + self + ); WAITING_TASKS.insert( msg_id, WaitingTask { @@ -850,16 +968,46 @@ impl OmikronConnection { }, ); - self.send_message(&cv).await; + if let Err(send_err) = self.send_message_result(cv).await { + WAITING_TASKS.remove(&msg_id); + return Err(format!( + "Request send failed (msg_id={}, reason={})", + msg_id, send_err + )); + } let timeout = timeout_duration.unwrap_or(Duration::from_secs(10)); match tokio::time::timeout(timeout, rx.recv()).await { - Ok(Some(response_cv)) => Ok(response_cv), - Ok(_) => Err("Channel closed".to_string()), - Err(_) => { + Ok(Some(response_cv)) => { + if response_cv.is_type(CommunicationType::error) { + let reason = response_cv + .get_data(DataTypes::message) + .as_str() + .unwrap_or("connection error") + .to_string(); + Err(format!( + "Request failed due to disconnect (msg_id={}, reason={})", + msg_id, reason + )) + } else { + Ok(response_cv) + } + } + Ok(_) => { WAITING_TASKS.remove(&msg_id); - Err("Request timed out".to_string()) + Err("Channel closed while awaiting response".to_string()) + } + Err(_) => { + let waiting_tasks_len = WAITING_TASKS.len(); + WAITING_TASKS.remove(&msg_id); + Err(format!( + "Request timed out (msg_id={}, timeout={}s, connected={}, waiting_tasks={})", + msg_id, + timeout.as_secs(), + self.is_connected().await, + waiting_tasks_len + )) } } } @@ -895,6 +1043,11 @@ impl OmikronConnection { pub static OMIKRON_CONNECTION: LazyLock> = LazyLock::new(|| { let conn = Arc::new(OmikronConnection::new()); + log!( + "OMIKRON_CONNECTION static initialized [id={}, ptr={:p}]", + conn.connection_id, + Arc::as_ptr(&conn) + ); start_task_cleanup_loop(); @@ -903,6 +1056,11 @@ pub static OMIKRON_CONNECTION: LazyLock> = LazyLock::new( pub async fn get_omikron_connection() -> Arc { let conn = OMIKRON_CONNECTION.clone(); + log!( + "get_omikron_connection() called [id={}, ptr={:p}]", + conn.connection_id, + Arc::as_ptr(&conn) + ); conn.connect().await; conn } diff --git a/src/omikron/ping_pong_task.rs b/src/omikron/ping_pong_task.rs index 4b757ba..446429b 100644 --- a/src/omikron/ping_pong_task.rs +++ b/src/omikron/ping_pong_task.rs @@ -1,5 +1,5 @@ -use crate::APP_STATE; use crate::omikron::omikron_connection::OmikronConnection; +use crate::{APP_STATE, log}; use dashmap::DashMap; use std::sync::LazyLock; use std::time::Instant; @@ -13,6 +13,13 @@ impl OmikronConnection { pub async fn send_ping(&self) { let id = rand_u32(); + log!( + "OmikronConnection::send_ping [id={}, ptr={:p}, ping_id={}]", + self.connection_id, + self, + id + ); + PING_TIMES.insert(id, Instant::now()); // Auto-cleanup old pings (optional) @@ -31,6 +38,13 @@ impl OmikronConnection { pub async fn handle_pong(&self, cv: &CommunicationValue) { let id = cv.get_id(); + log!( + "OmikronConnection::handle_pong [id={}, ptr={:p}, ping_id={}]", + self.connection_id, + self, + id + ); + if let Some((_, send_time)) = PING_TIMES.remove(&id) { let ping_ms = Instant::now().duration_since(send_time).as_millis() as i64; *self.last_ping.lock().await = ping_ms;