diff --git a/src/data/communication.rs b/src/data/communication.rs index 7c1660d..f6b2fe3 100644 --- a/src/data/communication.rs +++ b/src/data/communication.rs @@ -53,6 +53,7 @@ pub enum DataTypes { signature, signed, message, + message_state, last_ping, ping_iota, ping_clients, @@ -112,6 +113,7 @@ impl DataTypes { #[allow(non_camel_case_types, dead_code)] pub enum CommunicationType { error, + error_anonymous, error_internal, error_invalid_data, error_invalid_user_id, @@ -134,6 +136,7 @@ pub enum CommunicationType { settings_load, settings_list, message, + message_state, message_send, message_live, message_other_iota, diff --git a/src/eula/terms_checker.rs b/src/eula/terms_checker.rs index 3a9fd47..c2ca9c6 100644 --- a/src/eula/terms_checker.rs +++ b/src/eula/terms_checker.rs @@ -93,15 +93,15 @@ impl ConsentUiState { .as_secs(); format!( "\"EULA=true\" indicates that you read and accepted the End User Licence agreement. You can find our EULA at https://docs.tensamin.net/legal/eula/\ - \nEULA={}\ - \n\"PrivacyPolicy=true\" indicates that you read and accepted the Privacy Policy. You can find our Privacy Policy at https://docs.tensamin.net/legal/privacy-policy/\ - \nPrivacyPolicy={}\ - \n\"ToS=true\" indicates that you read and accepted the Terms of Service. You can find our Terms of Service at https://docs.tensamin.net/legal/terms-of-service/\ - \nToS={}\ - \nThis file reflects the current consent state used by the application.\ - \nIt may be regenerated or overwritten by the application.\ - \nThis file was last edited by Tensamin at:\ - \nUNIX-SECOND={}", + \nEULA={}\ + \n\"PrivacyPolicy=true\" indicates that you read and accepted the Privacy Policy. You can find our Privacy Policy at https://docs.tensamin.net/legal/privacy-policy/\ + \nPrivacyPolicy={}\ + \n\"ToS=true\" indicates that you read and accepted the Terms of Service. You can find our Terms of Service at https://docs.tensamin.net/legal/terms-of-service/\ + \nToS={}\ + \nThis file reflects the current consent state used by the application.\ + \nIt may be regenerated or overwritten by the application.\ + \nThis file was last edited by Tensamin at:\ + \nUNIX-SECOND={}", self.eula, self.pp, self.tos, current_secs ) } diff --git a/src/omikron/omikron_connection.rs b/src/omikron/omikron_connection.rs index 0adc9d0..40656b4 100644 --- a/src/omikron/omikron_connection.rs +++ b/src/omikron/omikron_connection.rs @@ -2,6 +2,7 @@ use crate::auth::local_auth; use crate::gui::log_panel::{log_cv, log_message_format}; use crate::users::contact::Contact; use crate::users::user_community_util::UserCommunityUtil; +use crate::util::chat_files::MessageState; use crate::util::chats_util::{get_user, mod_user}; use crate::util::crypto_util::{DataFormat, SecurePayload}; use crate::util::file_util::{get_children, load_file, save_file}; @@ -13,14 +14,15 @@ use crate::{ util::{config_util::CONFIG, crypto_helper}, }; use dashmap::DashMap; -use futures::Stream; use futures::stream::{SplitSink, SplitStream}; +use futures::{FutureExt, Stream}; use futures_util::sink::Sink; use futures_util::{SinkExt, StreamExt}; use hyper::upgrade::Upgraded; use hyper_util::rt::TokioIo; use json::JsonValue; use json::number::Number; +use pkcs8::DecodePrivateKey; use std::collections::HashMap; use std::sync::{Arc, LazyLock}; use std::time::{SystemTime, UNIX_EPOCH}; @@ -29,6 +31,7 @@ use tokio::time::{Duration, Instant, sleep}; use tokio_tungstenite::{connect_async, tungstenite::protocol::Message}; use tungstenite::Utf8Bytes; use uuid::Uuid; +use warp::reply::Json; pub static OMIKRON_CONNECTION: LazyLock>>>> = LazyLock::new(|| Arc::new(RwLock::new(None))); @@ -148,13 +151,10 @@ impl OmikronConnection { // Login flow if self.connect_internal().await { self.send_message( - CommunicationValue::new(CommunicationType::identification) - .add_data( - DataTypes::iota_id, - JsonValue::Number(json::number::Number::from(iota_id)), - ) - .to_json() - .to_string(), + &CommunicationValue::new(CommunicationType::identification).add_data( + DataTypes::iota_id, + JsonValue::Number(json::number::Number::from(iota_id)), + ), ) .await; } @@ -207,8 +207,13 @@ impl OmikronConnection { true } - pub async fn send_message(&self, msg: String) { - Self::send_message_static(&self.writer, Arc::clone(&self.is_connected), msg).await; + pub async fn send_message(&self, cv: &CommunicationValue) { + Self::send_message_static( + &self.writer, + Arc::clone(&self.is_connected), + cv.to_json().to_string(), + ) + .await; } pub async fn set_variant(self: &Arc, variant: ConnectionVariant) { @@ -228,7 +233,6 @@ impl OmikronConnection { let is_connected_out = self.is_connected.clone(); let sel_out = self.clone(); let variant = self.variant.clone(); - let sel_arc_out = self.clone(); { ACTIVE_TASKS.lock().unwrap().push("Listener".to_string()); @@ -238,14 +242,12 @@ impl OmikronConnection { if *SHUTDOWN.read().await { break; } - Self::handle_message( + sel_out.clone().handle_message( msg, waiting_out.clone(), writer_out.clone(), is_connected_out.clone(), - sel_out.clone(), variant.clone(), - sel_arc_out.clone(), ); } *is_connected_out.lock().await = false; @@ -259,6 +261,7 @@ impl OmikronConnection { } } pub fn handle_message( + self: Arc, msg: Result, waiting: Arc>>, writer: Arc< @@ -267,9 +270,7 @@ impl OmikronConnection { >, >, is_connected: Arc>, - sel: Arc, variant: Arc>, - sel_arc: Arc, ) { tokio::spawn(async move { match msg { @@ -281,7 +282,7 @@ impl OmikronConnection { Ok(Message::Text(text)) => { let cv = CommunicationValue::from_json(&text); if cv.is_type(CommunicationType::pong) { - sel.handle_pong(&cv, true).await; + self.handle_pong(&cv, true).await; return; } if cv.is_type(CommunicationType::challenge) { @@ -324,7 +325,7 @@ impl OmikronConnection { JsonValue::String(decrypted.export(DataFormat::Base64)), ); - sel_arc.send_message(response.to_json().to_string()).await; + self.send_message(&response).await; } else { log_message("Failed to decrypt challenge"); } @@ -349,13 +350,11 @@ impl OmikronConnection { .add_data( DataTypes::iota_id, JsonValue::Number(json::number::Number::from(iota_id)), - ) - .to_json() - .to_string(); + ); - let sel_arc_clone = sel_arc.clone(); + let self_clone = self.clone(); tokio::spawn(async move { - sel_arc_clone.send_message(login_message).await; + self_clone.send_message(&login_message).await; }); } else { log_message("Iota registration failed."); @@ -379,16 +378,13 @@ impl OmikronConnection { .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(), + self.send_message( + &CommunicationValue::new( + CommunicationType::error_invalid_user_id, ) - .await; + .with_id(cv.get_id()), + ) + .await; return; } @@ -403,44 +399,37 @@ impl OmikronConnection { 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(), + self.send_message( + &CommunicationValue::new( + CommunicationType::error_invalid_private_key, ) - .await; + .with_id(cv.get_id()), + ) + .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(), + self.send_message( + &CommunicationValue::new( + CommunicationType::error_invalid_private_key, ) - .await; + .with_id(cv.get_id()), + ) + .await; return; } // Set identification data - sel_arc.set_user_id(user_id).await; - sel_arc - .set_variant(ConnectionVariant::ClientAuthenticated) + self.set_user_id(user_id).await; + self.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; + self.send_message(&response).await; } } // ************************************************ // @@ -451,6 +440,25 @@ impl OmikronConnection { y(cv); return; } + if cv.is_type(CommunicationType::message_state) { + let sender_id = &cv.get_sender(); + let receiver_id = &cv.get_receiver(); + + chat_files::change_message_state( + cv.get_data(DataTypes::send_time) + .unwrap_or(&JsonValue::new_object()) + .as_i64() + .unwrap_or(0) as i64, + *receiver_id, + *sender_id, + MessageState::from_str( + cv.get_data(DataTypes::message_state) + .unwrap_or(&JsonValue::Null) + .as_str() + .unwrap_or(""), + ), + ); + } if cv.is_type(CommunicationType::message_other_iota) { let sender_id = &cv.get_sender(); let receiver_id = &cv.get_receiver(); @@ -485,12 +493,50 @@ impl OmikronConnection { 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; + let user_resp = self + .clone() + .await_response(&user_forward, Some(Duration::from_secs(10))) + .await; + + if let Ok(user_resp) = user_resp { + let ms: MessageState = MessageState::from_str( + user_resp + .get_data(DataTypes::message_state) + .unwrap_or(&JsonValue::Null) + .as_str() + .unwrap_or(""), + ) + .upgrade(MessageState::Send); + self.send_message( + &CommunicationValue::new(CommunicationType::message_state) + .with_id(cv.get_id()) + .with_receiver(*sender_id) + .with_sender(*receiver_id) + .add_data( + DataTypes::send_time, + cv.get_data(DataTypes::send_time).unwrap().clone(), + ) + .add_data( + DataTypes::message_state, + JsonValue::from(ms.as_str()), + ), + ); + } else { + self.send_message( + &CommunicationValue::new(CommunicationType::message_state) + .with_id(cv.get_id()) + .with_receiver(*sender_id) + .with_sender(*receiver_id) + .add_data( + DataTypes::send_time, + cv.get_data(DataTypes::send_time).unwrap().clone(), + ) + .add_data( + DataTypes::message_state, + JsonValue::from(MessageState::Send.as_str()), + ), + ); + } return; } @@ -821,13 +867,13 @@ impl OmikronConnection { }), ); - self.send_message(cv.to_json().to_string()).await; + self.send_message(&cv).await; 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(None) => Err("Failed to receive response, channel was closed.".to_string()), + Ok(_) => Err("Failed to receive response, channel was closed.".to_string()), Err(_) => { self.waiting.remove(&msg_id); Err(format!( diff --git a/src/omikron/ping_pong_task.rs b/src/omikron/ping_pong_task.rs index 5a96a5a..c851c44 100644 --- a/src/omikron/ping_pong_task.rs +++ b/src/omikron/ping_pong_task.rs @@ -20,11 +20,9 @@ impl OmikronConnection { .add_data_num( DataTypes::last_ping, Number::from(*self.last_ping.lock().await), - ) - .to_json() - .to_string(); + ); - self.send_message(ping_message).await; + self.send_message(&ping_message).await; } /// Handles incoming pong and calculates latency diff --git a/src/util/chat_files.rs b/src/util/chat_files.rs index ed4e51f..0eb2640 100644 --- a/src/util/chat_files.rs +++ b/src/util/chat_files.rs @@ -5,20 +5,40 @@ use std::path::Path; use crate::gui::log_panel::log_message; +#[derive(PartialEq, Debug, Clone)] pub enum MessageState { Read, Received, + Send, Sending, - Error, } impl MessageState { - fn as_str(&self) -> &'static str { + pub fn as_str(&self) -> &'static str { match self { MessageState::Read => "READ", MessageState::Received => "RECEIVED", + MessageState::Send => "SEND", MessageState::Sending => "SENDING", - MessageState::Error => "ERROR", + } + } + pub fn from_str(str: &str) -> Self { + match str.to_uppercase().as_str() { + "READ" => MessageState::Read, + "RECEIVED" => MessageState::Received, + "SEND" => MessageState::Send, + _ => MessageState::Sending, + } + } + pub fn upgrade(self, other: Self) -> Self { + if other == Self::Read || self == Self::Read { + Self::Read + } else if other == Self::Received || self == Self::Received { + Self::Received + } else if other == Self::Send || self == Self::Send { + Self::Send + } else { + Self::Sending } } } @@ -86,9 +106,9 @@ pub fn add_message( save_file(&user_dir, &file_name, &message_chunk.dump()); } pub fn change_message_state( + timestamp: i64, storage_owner: i64, external_user: i64, - timestamp: i64, new_state: MessageState, ) -> std::io::Result<()> { let user_dir = format!("users/{}/chats/{}", storage_owner, external_user); @@ -114,7 +134,13 @@ pub fn change_message_state( let mut modified = false; for i in 0..chunk.len() { if chunk[i]["message_time"].as_i64() == Some(timestamp) { - chunk[i]["message_state"] = JsonValue::from(new_state.as_str()); + chunk[i]["message_state"] = JsonValue::from( + MessageState::from_str( + chunk[i]["message_state"].as_str().unwrap_or("SENDING"), + ) + .upgrade(new_state.clone()) + .as_str(), + ); modified = true; break; } diff --git a/src/util/crypto_util.rs b/src/util/crypto_util.rs index 04e967c..df010e0 100644 --- a/src/util/crypto_util.rs +++ b/src/util/crypto_util.rs @@ -109,11 +109,6 @@ impl SecurePayload { let peer_pub = public_key.into(); let shared_secret = self.private_key.as_diffie_hellman(&peer_pub).unwrap(); - println!( - "Encryption Shared Secret (Hex): {}", - hex::encode(shared_secret.as_bytes()) - ); - // 3. Key & Nonce Derivation (HKDF) // We derive 32 bytes for the key and 12 bytes for a deterministic nonce. let hkdf = Hkdf::::new(None, shared_secret.as_bytes()); @@ -168,12 +163,6 @@ impl SecurePayload { let peer_pub = peer_public_key_bytes.into(); let shared_secret = self.private_key.as_diffie_hellman(&peer_pub).unwrap(); - // LOGGING: Shared Secret - println!( - "Decryption Shared Secret (Hex): {}", - hex::encode(shared_secret.as_bytes()) - ); - // 2. Key & Nonce Derivation (Must match encryption exactly) let hkdf = Hkdf::::new(None, shared_secret.as_bytes()); let mut okm = [0u8; 44];