From 1f878b93142437704b95d4448360afcbc5ed55e7 Mon Sep 17 00:00:00 2001 From: Alex Emmet Date: Thu, 18 Dec 2025 17:00:48 +0100 Subject: [PATCH] Fixes Reboot, I64 instead of Option --- Cargo.lock | 51 +- src/data/communication.rs | 8 +- src/gui/input_handler.rs | 54 +- src/gui/log_panel.rs | 10 +- src/gui/tui.rs | 1 + src/langu/language_creator.rs | 5 +- src/main.rs | 5 +- src/omikron/omikron_connection.rs | 853 +++++++++++++++--------------- src/omikron/ping_pong_task.rs | 1 + src/users/contact.rs | 16 +- src/users/user_manager.rs | 5 +- src/util/chats_util.rs | 4 +- src/util/config_util.rs | 1 - 13 files changed, 537 insertions(+), 477 deletions(-) diff --git a/Cargo.lock b/Cargo.lock index 5bcd2ce..b3d2fd9 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -222,9 +222,9 @@ checksum = "72b3254f16251a8381aa12e40e3c4d2f0199f8c6508fbecb9d91f575e0fbb8c6" [[package]] name = "base64ct" -version = "1.8.0" +version = "1.8.1" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "55248b47b0caf0546f7988906588779981c43bb1bc9d0c44087278f80cdb44ba" +checksum = "0e050f626429857a27ddccb31e0aca21356bfa709c04041aefddac081a8f068a" [[package]] name = "bitflags" @@ -285,9 +285,9 @@ dependencies = [ [[package]] name = "cc" -version = "1.2.48" +version = "1.2.49" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "c481bdbf0ed3b892f6f806287d72acd515b352a4ec27a208489b8c1bc839633a" +checksum = "90583009037521a116abf44494efecd645ba48b6622457080f080b85544e2215" dependencies = [ "find-msvc-tools", "jobserver", @@ -348,9 +348,9 @@ dependencies = [ [[package]] name = "cmake" -version = "0.1.54" +version = "0.1.56" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "e7caa3f9de89ddbe2c607f4101924c5abec803763ae9534e4f4d7d8f84aa81f0" +checksum = "b042e5d8a74ae91bb0961acd039822472ec99f8ab0948cbf6d1369588f8be586" dependencies = [ "cc", ] @@ -800,13 +800,13 @@ checksum = "3a3076410a55c90011c298b04d0cfa770b00fa04e1e3c97d3f6c9de105a03844" [[package]] name = "flate2" -version = "1.1.7" +version = "1.1.5" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "a2152dbcb980c05735e2a651d96011320a949eb31a0c8b38b72645ce97dec676" +checksum = "bfe33edd8e85a12a67454e37f8c75e730830d83e313556ab9ebf9ee7fbeb3bfb" dependencies = [ "crc32fast", + "libz-rs-sys", "miniz_oxide", - "zlib-rs", ] [[package]] @@ -1283,9 +1283,9 @@ checksum = "7aedcccd01fc5fe81e6b489c15b247b8b0690feb23304303a9e560f37efc560a" [[package]] name = "icu_properties" -version = "2.1.1" +version = "2.1.2" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "e93fcd3157766c0c8da2f8cff6ce651a31f0810eaa1c51ec363ef790bbb5fb99" +checksum = "020bfc02fe870ec3a66d93e677ccca0562506e5872c650f893269e08615d74ec" dependencies = [ "icu_collections", "icu_locale_core", @@ -1297,9 +1297,9 @@ dependencies = [ [[package]] name = "icu_properties_data" -version = "2.1.1" +version = "2.1.2" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "02845b3647bb045f1100ecd6480ff52f34c35f82d9880e029d329c21d1054899" +checksum = "616c294cf8d725c6afcd8f55abc17c56464ef6211f9ed59cccffe534129c77af" [[package]] name = "icu_provider" @@ -1540,6 +1540,15 @@ version = "0.2.178" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "37c93d8daa9d8a012fd8ab92f088405fb202ea0b6ab73ee2482ae66af4f42091" +[[package]] +name = "libz-rs-sys" +version = "0.5.4" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "15413ef615ad868d4d65dce091cb233b229419c7c0c4bcaa746c0901c49ff39c" +dependencies = [ + "zlib-rs", +] + [[package]] name = "linux-raw-sys" version = "0.4.15" @@ -2271,9 +2280,9 @@ checksum = "7a2d987857b319362043e95f5353c0535c1f58eec5336fdfcf626430af7def58" [[package]] name = "reqwest" -version = "0.12.24" +version = "0.12.25" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "9d0946410b9f7b082a427e4ef5c8ff541a88b357bc6c637c40db3a68ac70a36f" +checksum = "b6eff9328d40131d43bd911d42d79eb6a47312002a4daefc9e37f17e74a7701a" dependencies = [ "base64", "bytes", @@ -2616,9 +2625,9 @@ dependencies = [ [[package]] name = "simd-adler32" -version = "0.3.7" +version = "0.3.8" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "d66dc143e6b11c1eddc06d5c423cfc97062865baf299914ab64caa38182078fe" +checksum = "e320a6c5ad31d271ad523dcf3ad13e2767ad8b1cb8f047f75a8aeaf8da139da2" [[package]] name = "slab" @@ -2966,9 +2975,9 @@ dependencies = [ [[package]] name = "tower-http" -version = "0.6.7" +version = "0.6.8" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "9cf146f99d442e8e68e585f5d798ccd3cad9a7835b917e09728880a862706456" +checksum = "d4e6559d53cc268e5031cd8429d05415bc4cb4aefc4aa5d6cc35fbf5b924a1f8" dependencies = [ "bitflags", "bytes", @@ -3881,9 +3890,9 @@ dependencies = [ [[package]] name = "zlib-rs" -version = "0.5.3" +version = "0.5.4" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "36134c44663532e6519d7a6dfdbbe06f6f8192bde8ae9ed076e9b213f0e31df7" +checksum = "51f936044d677be1a1168fae1d03b583a285a5dd9d8cbf7b24c23aa1fc775235" [[package]] name = "zopfli" diff --git a/src/data/communication.rs b/src/data/communication.rs index 0620c36..2f8c689 100644 --- a/src/data/communication.rs +++ b/src/data/communication.rs @@ -355,22 +355,22 @@ impl CommunicationValue { object! { id: self.id.to_string(), type: format!("{:?}", self.comm_type), - sender: self.sender.to_string(), - receiver: self.receiver.to_string(), + sender: self.sender, + receiver: self.receiver, data: jdata } } else if self.sender > 0 { object! { id: self.id.to_string(), type: format!("{:?}", self.comm_type), - sender: self.sender.to_string(), + sender: self.sender, data: jdata } } else if self.receiver > 0 { object! { id: self.id.to_string(), type: format!("{:?}", self.comm_type), - receiver: self.receiver.to_string(), + receiver: self.receiver, data: jdata } } else { diff --git a/src/gui/input_handler.rs b/src/gui/input_handler.rs index 6528bdd..5a59837 100644 --- a/src/gui/input_handler.rs +++ b/src/gui/input_handler.rs @@ -1,33 +1,47 @@ -use crate::{ACTIVE_TASKS, RELOAD, SHUTDOWN, gui::tui::UNIQUE, util::config_util::CONFIG}; -use crossterm::event::{Event, KeyCode, read}; -use crossterm::event::{KeyEvent, KeyModifiers}; +use std::time::Duration; + +use crate::ACTIVE_TASKS; +use crate::{RELOAD, SHUTDOWN, gui::tui::UNIQUE, util::config_util::CONFIG}; +// Switched to poll/read which are available by default in crossterm +use crossterm::event::{Event, KeyCode, KeyEvent, KeyModifiers, poll, read}; + use json::JsonValue; -use tokio::{self}; pub fn setup_input_handler() { tokio::spawn(async move { { - ACTIVE_TASKS - .lock() - .unwrap() - .push("Input Handler".to_string()); + let mut tasks = ACTIVE_TASKS.lock().unwrap(); + tasks.push("Input Handler".to_string()); } - while let Ok(event) = read() { - if *SHUTDOWN.read().await { - break; + + loop { + { + let should_shutdown = *SHUTDOWN.read().await; + if should_shutdown { + break; + } } - if let Event::Key(key) = event { - handle_input(key).await; - } - if *SHUTDOWN.read().await { - break; + + let has_event = match poll(Duration::from_millis(100)) { + Ok(true) => true, + Ok(false) => false, + Err(_) => false, + }; + + if has_event { + match read() { + Ok(event) => { + if let Event::Key(key_event) = event { + handle_input(key_event).await; + } + } + Err(_) => (), + } } } { - ACTIVE_TASKS - .lock() - .unwrap() - .retain(|t| !t.eq(&"Input Handler".to_string())); + let mut tasks = ACTIVE_TASKS.lock().unwrap(); + tasks.retain(|t| t != "Input Handler"); } }); } diff --git a/src/gui/log_panel.rs b/src/gui/log_panel.rs index 6c7280e..85a9b68 100644 --- a/src/gui/log_panel.rs +++ b/src/gui/log_panel.rs @@ -33,9 +33,17 @@ pub fn log_message(msg: impl Into) { *UNIQUE.write().await = true; }); } +pub fn log_message_format(msg: impl Into, args: &[&str]) { + APP_STATE + .lock() + .unwrap() + .push_log(format(&msg.into(), args)); + tokio::spawn(async move { + *UNIQUE.write().await = true; + }); +} pub fn setup() { - // Start a background thread to sample metrics tokio::spawn(async move { { ACTIVE_TASKS.lock().unwrap().push("metrics".to_string()); diff --git a/src/gui/tui.rs b/src/gui/tui.rs index 064b63b..05c2e57 100644 --- a/src/gui/tui.rs +++ b/src/gui/tui.rs @@ -17,6 +17,7 @@ use crossterm::{ use once_cell::sync::Lazy; use ratatui::{ Frame, Terminal, + crossterm::event::poll, layout::{Constraint, Direction, Layout, Rect}, prelude::CrosstermBackend, style::Color, diff --git a/src/langu/language_creator.rs b/src/langu/language_creator.rs index 1d9ce29..fc6d1c5 100644 --- a/src/langu/language_creator.rs +++ b/src/langu/language_creator.rs @@ -22,7 +22,10 @@ pub fn create_languages() -> Result<(), JsonError> { "identification_response", "IOTA identified on Omikron, {} users!", )?; - omikron_messages.insert("send_message_failed", "Failed to send message to Omikron")?; + omikron_messages.insert( + "send_message_failed", + "Failed to send message to Omikron: {}", + )?; omikron_messages.insert("connection_failed", "Failed to connect to Omikron: {}")?; // BUTTONS diff --git a/src/main.rs b/src/main.rs index fe76679..815932d 100644 --- a/src/main.rs +++ b/src/main.rs @@ -44,7 +44,6 @@ pub static APP_STATE: LazyLock>> = pub static SHUTDOWN: Lazy> = Lazy::new(|| RwLock::new(false)); pub static RELOAD: Lazy> = Lazy::new(|| RwLock::new(true)); pub static ACTIVE_TASKS: Lazy>> = Lazy::new(|| Mutex::new(Vec::new())); -pub static RECONNECT: Lazy> = Lazy::new(|| RwLock::new(false)); #[tokio::main(flavor = "multi_thread", worker_threads = 8)] #[allow(unused_must_use, dead_code)] @@ -76,10 +75,10 @@ async fn main() { CONFIG.write().await.change( "iota_id", JsonValue::Number(Number::from( - (SystemTime::now() + SystemTime::now() .duration_since(UNIX_EPOCH) .unwrap() - .as_millis() as i64), + .as_millis() as i64, )), ); CONFIG.write().await.update(); diff --git a/src/omikron/omikron_connection.rs b/src/omikron/omikron_connection.rs index fd527be..4529155 100644 --- a/src/omikron/omikron_connection.rs +++ b/src/omikron/omikron_connection.rs @@ -1,12 +1,12 @@ use crate::auth::local_auth; use crate::data::communication::{CommunicationType, CommunicationValue, DataTypes}; -use crate::gui::log_panel::{log_cv, log_message, log_message_trans}; +use crate::gui::log_panel::{log_cv, log_message, log_message_format}; use crate::users::contact::Contact; use crate::users::user_community_util::UserCommunityUtil; use crate::util::chat_files; use crate::util::chats_util::{get_user, get_users, mod_user}; use crate::util::file_util::{get_children, load_file, save_file}; -use crate::{ACTIVE_TASKS, RECONNECT, SHUTDOWN}; +use crate::{ACTIVE_TASKS, SHUTDOWN}; use futures::Stream; use futures::stream::{SplitSink, SplitStream}; use futures_util::sink::Sink; @@ -15,6 +15,7 @@ use hyper::upgrade::Upgraded; use hyper_util::rt::TokioIo; use json::JsonValue; use json::number::Number; +use ratatui::crossterm::event::poll; use std::collections::HashMap; use std::sync::{Arc, LazyLock}; use std::time::{SystemTime, UNIX_EPOCH}; @@ -90,9 +91,6 @@ impl OmikronConnection { if *SHUTDOWN.read().await { break; } - if *RECONNECT.read().await { - break; - } match connect_async("wss://app.tensamin.net/ws/iota/").await { Ok((ws_stream, _)) => { let (write_half, read_half) = ws_stream.split(); @@ -111,7 +109,7 @@ impl OmikronConnection { if *SHUTDOWN.read().await { break; } - if *RECONNECT.read().await { + if *cloned_self.is_connected.lock().await == false { break; } cloned_self.send_ping().await; @@ -137,7 +135,7 @@ impl OmikronConnection { } } pub async fn send_message(&self, msg: String) { - Self::send_message_static(&self.writer, msg).await + Self::send_message_static(&self.writer, Arc::clone(&self.is_connected), msg).await; } pub async fn set_variant(self: &Arc, variant: ConnectionVariant) { @@ -163,410 +161,24 @@ impl OmikronConnection { ACTIVE_TASKS.lock().unwrap().push("Listener".to_string()); } tokio::spawn(async move { - while let Some(msg) = read_half.next().await { - if *SHUTDOWN.read().await { - break; - } - if *RECONNECT.read().await { - break; - } - let waiting = waiting_out.clone(); - let writer = writer_out.clone(); - let is_connected = is_connected_out.clone(); - let sel = sel_out.clone(); - let variant = variant.clone(); - let sel_arc = sel_arc_out.clone(); - tokio::spawn(async move { - match msg { - Ok(Message::Close(Some(frame))) => { - log_message(format!("[Omikron] Closed: {:?}", frame)); - *is_connected.lock().await = false; - return; + while !*SHUTDOWN.read().await { + if poll(Duration::from_millis(100)).unwrap() { + if let Some(msg) = read_half.next().await { + if *is_connected_out.lock().await == false { + log_message("Disconnected, not handeling incomming"); + break; } - Ok(Message::Text(text)) => { - let mut cv = CommunicationValue::from_json(&text); - if cv.is_type(CommunicationType::pong) { - sel.handle_pong(&cv, true).await; - return; - } - let com = variant.read().await.clone(); - if com == ConnectionVariant::ClientUnauthenticated { - if cv.is_type(CommunicationType::identification) { - // Extract user ID - let user_id: i64 = cv - .get_data(DataTypes::user_id) - .unwrap_or(&JsonValue::Null) - .as_i64() - .unwrap_or(0); - if user_id == 0 { - sel_arc - .send_message( - CommunicationValue::new( - CommunicationType::error_invalid_user_id, - ) - .with_id(cv.get_id()) - .to_json() - .to_string(), - ) - .await; - return; - } - - // Validate private key - if let Some(private_key_hash) = - cv.get_data(DataTypes::private_key_hash) - { - log_message(format!( - "private_key_hash: {}", - private_key_hash - )); - let is_valid = local_auth::is_private_key_valid( - &user_id, - &private_key_hash.to_string(), - ); - - if !is_valid { - log_message("Invalid private key"); - sel_arc.send_message( - CommunicationValue::new( - CommunicationType::error_invalid_private_key, - ) - .with_id(cv.get_id()) - .to_json() - .to_string(), - ) - .await; - return; - } - } else { - log_message("Missing private key"); - sel_arc - .send_message( - CommunicationValue::new( - CommunicationType::error_invalid_private_key, - ) - .with_id(cv.get_id()) - .to_json() - .to_string(), - ) - .await; - return; - } - - // Set identification data - - sel_arc.set_user_id(user_id).await; - sel_arc - .set_variant(ConnectionVariant::ClientAuthenticated) - .await; - - let response = CommunicationValue::new( - CommunicationType::identification_response, - ) - .with_id(cv.get_id()); - sel_arc.send_message(response.to_json().to_string()).await; - } - } - // ************************************************ // - // Direct messages // - // ************************************************ // - log_cv(&cv); - if let Some(x) = waiting.lock().await.remove(&cv.get_id()) { - x(cv); - return; - } - if cv.is_type(CommunicationType::message_other_iota) { - let sender_id = &cv.get_sender(); - let receiver_id = &cv.get_receiver(); - - chat_files::add_message( - cv.get_data(DataTypes::send_time) - .unwrap_or(&JsonValue::new_object()) - .as_i64() - .unwrap_or( - SystemTime::now() - .duration_since(UNIX_EPOCH) - .unwrap() - .as_millis() - as i64, - ) as u128, - false, - *receiver_id, - *sender_id, - cv.get_data(DataTypes::content).unwrap().as_str().unwrap(), - ); - let user_forward = - CommunicationValue::new(CommunicationType::message_live) - .with_id(cv.get_id()) - .with_receiver(*receiver_id) - .add_data( - DataTypes::send_time, - cv.get_data(DataTypes::send_time).unwrap().clone(), - ) - .add_data( - DataTypes::message, - cv.get_data(DataTypes::content).unwrap().clone(), - ) - .add_data( - DataTypes::sender_id, - JsonValue::Number(Number::from(cv.get_sender())), - ); - Self::send_message_static( - &writer.clone(), - user_forward.to_json().to_string(), - ) - .await; - return; - } - - if cv.is_type(CommunicationType::message_send) { - let my_id = cv.get_sender(); - let other_id = cv - .get_data(DataTypes::receiver_id) - .unwrap_or(&JsonValue::Null) - .as_i64() - .unwrap_or(0); - chat_files::add_message( - SystemTime::now() - .duration_since(UNIX_EPOCH) - .unwrap() - .as_millis() as u128, - true, - my_id, - other_id, - &*cv.get_data(DataTypes::content).unwrap().to_string(), - ); - let ack = CommunicationValue::new(CommunicationType::message) - .with_id(cv.get_id()) - .with_receiver(my_id); - Self::send_message_static( - &writer.clone(), - ack.to_json().to_string(), - ) - .await; - let forward = CommunicationValue::forward_to_other_iota(&mut cv); - Self::send_message_static( - &writer.clone(), - forward.to_json().to_string(), - ) - .await; - return; - } - - if cv.is_type(CommunicationType::messages_get) { - let my_id = cv.get_sender(); - let partner_id = cv - .get_data(DataTypes::user_id) - .unwrap_or(&JsonValue::Null) - .as_i64() - .unwrap_or(0); - let offset = cv - .get_data(DataTypes::offset) - .unwrap_or(&JsonValue::Null) - .to_string() - .parse::() - .unwrap_or(0); - let amount = cv - .get_data(DataTypes::amount) - .unwrap_or(&JsonValue::Null) - .to_string() - .parse::() - .unwrap_or(0); - let messages = - chat_files::get_messages(my_id, partner_id, offset, amount); - let resp = CommunicationValue::new(CommunicationType::messages_get) - .with_id(cv.get_id()) - .with_receiver(my_id) - .add_data(DataTypes::messages, messages); - - Self::send_message_static( - &writer.clone(), - resp.to_json().to_string(), - ) - .await; - return; - } - - if cv.is_type(CommunicationType::get_chats) { - let user_id = cv.get_sender(); - let users = get_users(user_id); - let resp = CommunicationValue::new(CommunicationType::get_chats) - .with_id(cv.get_id()) - .with_receiver(user_id) - .add_data(DataTypes::user_ids, users); - Self::send_message_static( - &writer.clone(), - resp.to_json().to_string(), - ) - .await; - return; - } - - if cv.is_type(CommunicationType::add_chat) { - let user_id = cv.get_sender(); - let other_id = cv - .get_data(DataTypes::user_id) - .unwrap_or(&JsonValue::Null) - .as_i64() - .unwrap_or(0); - let mut contact = - get_user(user_id, other_id).unwrap_or(Contact::new(other_id)); // needs ChatsUtil + Contact - contact.set_last_message_at( - SystemTime::now() - .duration_since(UNIX_EPOCH) - .unwrap() - .as_millis() as i64, - ); - mod_user(user_id, &contact); - let resp = CommunicationValue::new(CommunicationType::add_chat) - .with_id(cv.get_id()) - .with_receiver(user_id); - Self::send_message_static( - &writer.clone(), - resp.to_json().to_string(), - ) - .await; - return; - } - - if cv.is_type(CommunicationType::add_community) { - UserCommunityUtil::add_community( - cv.get_sender(), - cv.get_data(DataTypes::community_address) - .unwrap() - .to_string(), - cv.get_data(DataTypes::community_title).unwrap().to_string(), - cv.get_data(DataTypes::position).unwrap().to_string(), - ); - let resp = - CommunicationValue::new(CommunicationType::add_community) - .with_id(cv.get_id()) - .with_receiver(cv.get_sender()); - Self::send_message_static( - &writer.clone(), - resp.to_json().to_string(), - ) - .await; - return; - } - - if cv.is_type(CommunicationType::get_communities) { - let resp = - CommunicationValue::new(CommunicationType::get_communities) - .with_id(cv.get_id()) - .with_receiver(cv.get_sender()) - .add_array( - DataTypes::communities, - UserCommunityUtil::get_communities(cv.get_sender()), - ); - Self::send_message_static( - &writer.clone(), - resp.to_json().to_string(), - ) - .await; - return; - } - - if cv.is_type(CommunicationType::remove_community) { - UserCommunityUtil::remove_community( - cv.get_sender(), - cv.get_data(DataTypes::community_address) - .unwrap() - .to_string(), - ); // needs UserCommunityUtil - let resp = - CommunicationValue::new(CommunicationType::remove_community) - .with_id(cv.get_id()) - .with_receiver(cv.get_sender()); - Self::send_message_static( - &writer.clone(), - resp.to_json().to_string(), - ) - .await; - return; - } - - if cv.is_type(CommunicationType::settings_save) { - let my_id = cv.get_sender(); - let settings_name = - cv.get_data(DataTypes::settings_name).unwrap().to_string(); - let settings_value = - cv.get_data(DataTypes::payload).unwrap().to_string(); - - save_file( - &format!("users/{}/settings/", my_id), - &format!("{}.settings", settings_name), - &settings_value, - ); - - let response = - CommunicationValue::new(CommunicationType::settings_save) - .with_receiver(my_id) - .with_id(cv.get_id()); - - Self::send_message_static( - &writer.clone(), - response.to_json().to_string(), - ) - .await; - return; - } - if cv.is_type(CommunicationType::settings_load) { - let my_id = cv.get_sender(); - let settings_name = - cv.get_data(DataTypes::settings_name).unwrap().to_string(); - let settings_value_str = load_file( - &format!("users/{}/settings/", my_id), - &format!("{}.settings", settings_name), - ); - let settings_value_json = JsonValue::from(settings_value_str); - let response = - CommunicationValue::new(CommunicationType::settings_load) - .with_id(cv.get_id()) - .with_receiver(my_id) - .add_data(DataTypes::payload, settings_value_json) - .add_data_str(DataTypes::settings_name, settings_name); - - Self::send_message_static( - &writer.clone(), - response.to_json().to_string(), - ) - .await; - return; - } - if cv.is_type(CommunicationType::settings_list) { - let my_id = cv.get_sender(); - let settings = get_children(&format!("users/{}/settings/", my_id)); - let mut settings_json = JsonValue::new_array(); - for s in settings { - let s = s.replace(".settings", ""); - if s.is_empty() { - continue; - } - let _ = settings_json.push(JsonValue::String(s)); - } - let response = - CommunicationValue::new(CommunicationType::settings_list) - .with_id(cv.get_id()) - .with_receiver(my_id) - .add_data(DataTypes::settings, settings_json); - - Self::send_message_static( - &writer.clone(), - response.to_json().to_string(), - ) - .await; - return; - } - } - Err(e) => { - log_message(format!("[Omikron] Error: {}", e)); - *is_connected.lock().await = false; - return; - } - _ => {} + Self::handle_message( + msg, + waiting_out.clone(), + writer_out.clone(), + is_connected_out.clone(), + sel_out.clone(), + variant.clone(), + sel_arc_out.clone(), + ); } - }); + } } }); { @@ -576,20 +188,435 @@ impl OmikronConnection { .retain(|t| !t.eq(&"Listener".to_string())); } } + pub fn handle_message( + msg: Result, + waiting: Arc>>>, + writer: Arc< + Mutex< + Option + Send + Unpin + 'static>>, + >, + >, + is_connected: Arc>, + sel: Arc, + variant: Arc>, + sel_arc: Arc, + ) { + tokio::spawn(async move { + match msg { + Ok(Message::Close(Some(frame))) => { + log_message(format!("[Omikron] Closed: {:?}", frame)); + *is_connected.lock().await = false; + return; + } + Ok(Message::Text(text)) => { + let mut cv = CommunicationValue::from_json(&text); + if cv.is_type(CommunicationType::pong) { + sel.handle_pong(&cv, true).await; + return; + } + let com = variant.read().await.clone(); + if com == ConnectionVariant::ClientUnauthenticated { + if cv.is_type(CommunicationType::identification) { + // Extract user ID + let user_id: i64 = cv + .get_data(DataTypes::user_id) + .unwrap_or(&JsonValue::Null) + .as_i64() + .unwrap_or(0); + if user_id == 0 { + sel_arc + .send_message( + CommunicationValue::new( + CommunicationType::error_invalid_user_id, + ) + .with_id(cv.get_id()) + .to_json() + .to_string(), + ) + .await; + return; + } + + // Validate private key + if let Some(private_key_hash) = cv.get_data(DataTypes::private_key_hash) + { + log_message(format!("private_key_hash: {}", private_key_hash)); + let is_valid = local_auth::is_private_key_valid( + &user_id, + &private_key_hash.to_string(), + ); + + if !is_valid { + log_message("Invalid private key"); + sel_arc + .send_message( + CommunicationValue::new( + CommunicationType::error_invalid_private_key, + ) + .with_id(cv.get_id()) + .to_json() + .to_string(), + ) + .await; + return; + } + } else { + log_message("Missing private key"); + sel_arc + .send_message( + CommunicationValue::new( + CommunicationType::error_invalid_private_key, + ) + .with_id(cv.get_id()) + .to_json() + .to_string(), + ) + .await; + return; + } + + // Set identification data + + sel_arc.set_user_id(user_id).await; + sel_arc + .set_variant(ConnectionVariant::ClientAuthenticated) + .await; + + let response = + CommunicationValue::new(CommunicationType::identification_response) + .with_id(cv.get_id()); + sel_arc.send_message(response.to_json().to_string()).await; + } + } + // ************************************************ // + // Direct messages // + // ************************************************ // + log_cv(&cv); + if let Some(x) = waiting.lock().await.remove(&cv.get_id()) { + x(cv); + return; + } + if cv.is_type(CommunicationType::message_other_iota) { + let sender_id = &cv.get_sender(); + let receiver_id = &cv.get_receiver(); + + chat_files::add_message( + cv.get_data(DataTypes::send_time) + .unwrap_or(&JsonValue::new_object()) + .as_i64() + .unwrap_or( + SystemTime::now() + .duration_since(UNIX_EPOCH) + .unwrap() + .as_millis() as i64, + ) as u128, + false, + *receiver_id, + *sender_id, + cv.get_data(DataTypes::content).unwrap().as_str().unwrap(), + ); + let user_forward = CommunicationValue::new(CommunicationType::message_live) + .with_id(cv.get_id()) + .with_receiver(*receiver_id) + .add_data( + DataTypes::send_time, + cv.get_data(DataTypes::send_time).unwrap().clone(), + ) + .add_data( + DataTypes::message, + cv.get_data(DataTypes::content).unwrap().clone(), + ) + .add_data( + DataTypes::sender_id, + JsonValue::Number(Number::from(cv.get_sender())), + ); + Self::send_message_static( + &writer.clone(), + is_connected, + user_forward.to_json().to_string(), + ) + .await; + return; + } + + if cv.is_type(CommunicationType::message_send) { + let my_id = cv.get_sender(); + let other_id = cv + .get_data(DataTypes::receiver_id) + .unwrap_or(&JsonValue::Null) + .as_i64() + .unwrap_or(0); + chat_files::add_message( + SystemTime::now() + .duration_since(UNIX_EPOCH) + .unwrap() + .as_millis() as u128, + true, + my_id, + other_id, + &*cv.get_data(DataTypes::content).unwrap().to_string(), + ); + let ack = CommunicationValue::new(CommunicationType::message) + .with_id(cv.get_id()) + .with_receiver(my_id); + Self::send_message_static( + &writer.clone(), + Arc::clone(&is_connected), + ack.to_json().to_string(), + ) + .await; + let forward = CommunicationValue::forward_to_other_iota(&mut cv); + Self::send_message_static( + &writer.clone(), + is_connected, + forward.to_json().to_string(), + ) + .await; + return; + } + + if cv.is_type(CommunicationType::messages_get) { + let my_id = cv.get_sender(); + let partner_id = cv + .get_data(DataTypes::user_id) + .unwrap_or(&JsonValue::Null) + .as_i64() + .unwrap_or(0); + let offset = cv + .get_data(DataTypes::offset) + .unwrap_or(&JsonValue::Null) + .to_string() + .parse::() + .unwrap_or(0); + let amount = cv + .get_data(DataTypes::amount) + .unwrap_or(&JsonValue::Null) + .to_string() + .parse::() + .unwrap_or(0); + let messages = chat_files::get_messages(my_id, partner_id, offset, amount); + let resp = CommunicationValue::new(CommunicationType::messages_get) + .with_id(cv.get_id()) + .with_receiver(my_id) + .add_data(DataTypes::messages, messages); + + Self::send_message_static( + &writer.clone(), + is_connected, + resp.to_json().to_string(), + ) + .await; + return; + } + + if cv.is_type(CommunicationType::get_chats) { + let user_id = cv.get_sender(); + let users = get_users(user_id); + let resp = CommunicationValue::new(CommunicationType::get_chats) + .with_id(cv.get_id()) + .with_receiver(user_id) + .add_data(DataTypes::user_ids, users); + Self::send_message_static( + &writer.clone(), + is_connected, + resp.to_json().to_string(), + ) + .await; + return; + } + + if cv.is_type(CommunicationType::add_chat) { + let user_id = cv.get_sender(); + let other_id = cv + .get_data(DataTypes::user_id) + .unwrap_or(&JsonValue::Null) + .as_i64() + .unwrap_or(0); + let mut contact = + get_user(user_id, other_id).unwrap_or(Contact::new(other_id)); + contact.set_last_message_at( + SystemTime::now() + .duration_since(UNIX_EPOCH) + .unwrap() + .as_millis() as i64, + ); + mod_user(user_id, &contact); + let resp = CommunicationValue::new(CommunicationType::add_chat) + .with_id(cv.get_id()) + .with_receiver(user_id); + Self::send_message_static( + &writer.clone(), + is_connected, + resp.to_json().to_string(), + ) + .await; + return; + } + + if cv.is_type(CommunicationType::add_community) { + UserCommunityUtil::add_community( + cv.get_sender(), + cv.get_data(DataTypes::community_address) + .unwrap() + .to_string(), + cv.get_data(DataTypes::community_title).unwrap().to_string(), + cv.get_data(DataTypes::position).unwrap().to_string(), + ); + let resp = CommunicationValue::new(CommunicationType::add_community) + .with_id(cv.get_id()) + .with_receiver(cv.get_sender()); + Self::send_message_static( + &writer.clone(), + is_connected, + resp.to_json().to_string(), + ) + .await; + return; + } + + if cv.is_type(CommunicationType::get_communities) { + let resp = CommunicationValue::new(CommunicationType::get_communities) + .with_id(cv.get_id()) + .with_receiver(cv.get_sender()) + .add_array( + DataTypes::communities, + UserCommunityUtil::get_communities(cv.get_sender()), + ); + Self::send_message_static( + &writer.clone(), + is_connected, + resp.to_json().to_string(), + ) + .await; + return; + } + + if cv.is_type(CommunicationType::remove_community) { + UserCommunityUtil::remove_community( + cv.get_sender(), + cv.get_data(DataTypes::community_address) + .unwrap() + .to_string(), + ); // needs UserCommunityUtil + let resp = CommunicationValue::new(CommunicationType::remove_community) + .with_id(cv.get_id()) + .with_receiver(cv.get_sender()); + Self::send_message_static( + &writer.clone(), + is_connected, + resp.to_json().to_string(), + ) + .await; + return; + } + + if cv.is_type(CommunicationType::settings_save) { + let my_id = cv.get_sender(); + let settings_name = + cv.get_data(DataTypes::settings_name).unwrap().to_string(); + let settings_value = cv.get_data(DataTypes::payload).unwrap().to_string(); + + save_file( + &format!("users/{}/settings/", my_id), + &format!("{}.settings", settings_name), + &settings_value, + ); + + let response = CommunicationValue::new(CommunicationType::settings_save) + .with_receiver(my_id) + .with_id(cv.get_id()); + + Self::send_message_static( + &writer.clone(), + is_connected, + response.to_json().to_string(), + ) + .await; + return; + } + if cv.is_type(CommunicationType::settings_load) { + let my_id = cv.get_sender(); + let settings_name = + cv.get_data(DataTypes::settings_name).unwrap().to_string(); + let settings_value_str = load_file( + &format!("users/{}/settings/", my_id), + &format!("{}.settings", settings_name), + ); + let settings_value_json = JsonValue::from(settings_value_str); + let response = CommunicationValue::new(CommunicationType::settings_load) + .with_id(cv.get_id()) + .with_receiver(my_id) + .add_data(DataTypes::payload, settings_value_json) + .add_data_str(DataTypes::settings_name, settings_name); + + Self::send_message_static( + &writer.clone(), + is_connected, + response.to_json().to_string(), + ) + .await; + return; + } + if cv.is_type(CommunicationType::settings_list) { + let my_id = cv.get_sender(); + let settings = get_children(&format!("users/{}/settings/", my_id)); + let mut settings_json = JsonValue::new_array(); + for s in settings { + let s = s.replace(".settings", ""); + if s.is_empty() { + continue; + } + let _ = settings_json.push(JsonValue::String(s)); + } + let response = CommunicationValue::new(CommunicationType::settings_list) + .with_id(cv.get_id()) + .with_receiver(my_id) + .add_data(DataTypes::settings, settings_json); + + Self::send_message_static( + &writer.clone(), + is_connected, + response.to_json().to_string(), + ) + .await; + return; + } + } + Err(e) => { + log_message(format!("[Omikron] Error: {}", e)); + *is_connected.lock().await = false; + return; + } + _ => {} + } + }); + } + pub async fn send_message_static( writer: &Arc< Mutex + Send + Unpin>>>, >, + connected: Arc>, msg: String, ) { let mut guard = writer.lock().await; if let Some(writer) = guard.as_mut() { - if let Ok(_) = writer.send(Message::Text(Utf8Bytes::from(msg))).await { - if let Ok(_) = writer.flush().await { - return; + match writer.send(Message::Text(Utf8Bytes::from(msg))).await { + Ok(_) => match writer.flush().await { + Ok(_) => return, + Err(e) => { + log_message_format("send_message_failed", &[&e.to_string()]); + *connected.lock().await = false; + } + }, + Err(e) => { + log_message_format("send_message_failed", &[&e.to_string()]); + *connected.lock().await = false; } } + } else { + log_message_format("send_message_failed", &["Immutable Writer"]); + *connected.lock().await = false; } - log_message_trans("send_message_failed"); } } diff --git a/src/omikron/ping_pong_task.rs b/src/omikron/ping_pong_task.rs index 5a96a5a..12e8d42 100644 --- a/src/omikron/ping_pong_task.rs +++ b/src/omikron/ping_pong_task.rs @@ -1,5 +1,6 @@ use crate::APP_STATE; use crate::data::communication::{CommunicationType, CommunicationValue, DataTypes}; +use crate::gui::log_panel::log_message; use crate::omikron::omikron_connection::OmikronConnection; use json::number::Number; use tokio::time::Instant; diff --git a/src/users/contact.rs b/src/users/contact.rs index bdfa6b2..632cea3 100644 --- a/src/users/contact.rs +++ b/src/users/contact.rs @@ -1,9 +1,9 @@ -use json::{self, JsonValue}; +use json::{self, JsonValue, number::Number}; use std::time::{SystemTime, UNIX_EPOCH}; #[derive(Debug, Clone)] pub struct Contact { - pub user_id: Option, + pub user_id: i64, pub user_name: Option, pub last_message_at: Option, } @@ -15,7 +15,7 @@ impl Default for Contact { .unwrap() .as_millis() as i64; Contact { - user_id: None, + user_id: 0, user_name: None, last_message_at: Some(now), } @@ -25,7 +25,7 @@ impl Default for Contact { impl Contact { pub fn new(user_id: i64) -> Self { Contact { - user_id: Some(user_id), + user_id: user_id, user_name: None, last_message_at: None, } @@ -36,19 +36,17 @@ impl Contact { pub fn to_json(&self) -> JsonValue { let mut obj = JsonValue::new_object(); - if let Some(id) = &self.user_id { - obj["user_id"] = JsonValue::from(id.to_string()); - } + obj["user_id"] = JsonValue::Number(Number::from(self.user_id)); if let Some(name) = &self.user_name { obj["user_name"] = JsonValue::from(name.as_str()); } if let Some(ts) = &self.last_message_at { - obj["last_message_at"] = JsonValue::from(ts.to_string()); + obj["last_message_at"] = JsonValue::Number(Number::from(*ts)); } obj } pub fn from_json(o: &JsonValue) -> Contact { - let user_id = o["user_id"].as_i64(); + let user_id = o["user_id"].as_i64().unwrap_or(0); let user_name = o["user_name"].as_str().map(|s| s.to_string()); diff --git a/src/users/user_manager.rs b/src/users/user_manager.rs index 90fcd45..b714638 100644 --- a/src/users/user_manager.rs +++ b/src/users/user_manager.rs @@ -1,9 +1,9 @@ -use crate::RECONNECT; use crate::auth::{auth_connector, crypto_helper}; use crate::gui::log_panel::log_message; use crate::users::user_profile::UserProfile; use crate::util::config_util::CONFIG; use crate::util::file_util::{load_file, save_file}; +use crate::{RELOAD, SHUTDOWN}; use base64::{Engine as _, engine::general_purpose::STANDARD}; use hex::{self}; use json::JsonValue; @@ -71,7 +71,8 @@ pub async fn create_user(username: &str) -> (Option, Option ); auth_connector::complete_register(&up, &CONFIG.read().await.get_iota_id().to_string()).await; - *RECONNECT.write().await = true; + *SHUTDOWN.write().await = true; + *RELOAD.write().await = true; log_message("Created User"); save_file( "", diff --git a/src/util/chats_util.rs b/src/util/chats_util.rs index a1b1833..ed71bf5 100644 --- a/src/util/chats_util.rs +++ b/src/util/chats_util.rs @@ -1,5 +1,6 @@ use json::{self, JsonValue, array}; +use crate::gui::log_panel::log_message; use crate::users::contact::Contact; use crate::util::file_util::{load_file, save_file}; @@ -14,7 +15,7 @@ pub fn mod_user(storage_owner: i64, contact: &Contact) { }; for i in 0..contacts.len() { - if contacts[i]["user_id"].as_str() == Some(&contact.user_id.unwrap().to_string()) { + if contacts[i]["user_id"] == contact.user_id { contacts.array_remove(i); break; } @@ -56,6 +57,5 @@ pub fn get_users(storage_owner: i64) -> JsonValue { } } } - contacts_out } diff --git a/src/util/config_util.rs b/src/util/config_util.rs index 68cad76..4b09921 100644 --- a/src/util/config_util.rs +++ b/src/util/config_util.rs @@ -2,7 +2,6 @@ use crate::util::file_util::{load_file, save_file}; use json::JsonValue; use once_cell::sync::Lazy; use tokio::sync::RwLock; -use uuid::Uuid; pub static CONFIG: Lazy> = Lazy::new(|| RwLock::new(ConfigUtil::new()));