From 9904e54914d45a376f4fc7a3f1dcc04d6ddd57e2 Mon Sep 17 00:00:00 2001 From: Alex Emmet <111742636+Alex-Emmet@users.noreply.github.com> Date: Sun, 17 May 2026 21:34:55 +0200 Subject: [PATCH] [Add] basic Tauri --- Cargo.lock | 2 - Cargo.toml | 4 +- src/main.rs | 9 ++ src/notifications/mod.rs | 1 + src/notifications/tauri.rs | 160 ++++++++++++++++++++++++++++ src/transport/omikron_connection.rs | 93 ++++++++++------ src/transport/omikron_manager.rs | 10 ++ src/util/logger.rs | 37 +++---- 8 files changed, 261 insertions(+), 55 deletions(-) create mode 100644 src/notifications/mod.rs create mode 100644 src/notifications/tauri.rs diff --git a/Cargo.lock b/Cargo.lock index 21d6ccc..b34b40b 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -3599,7 +3599,6 @@ checksum = "e421abadd41a4225275504ea4d6566923418b7f05506fbc9c0fe86ba7396114b" [[package]] name = "ttp-core" version = "0.1.0" -source = "git+https://git.methanium.net/Tensamin/TTP.git#929cb9d3a6aebebe6365973f13062ac2a8e03af6" dependencies = [ "base64", "byteorder", @@ -3612,7 +3611,6 @@ dependencies = [ [[package]] name = "ttp-native" version = "0.1.0" -source = "git+https://git.methanium.net/Tensamin/TTP.git#929cb9d3a6aebebe6365973f13062ac2a8e03af6" dependencies = [ "quinn", "rustls", diff --git a/Cargo.toml b/Cargo.toml index 22ec471..c6c6bf4 100755 --- a/Cargo.toml +++ b/Cargo.toml @@ -4,8 +4,8 @@ version = "0.1.0" edition = "2024" [dependencies] -ttp-core = { git = "https://git.methanium.net/Tensamin/TTP.git", package = "ttp-core" } -ttp-native = { git = "https://git.methanium.net/Tensamin/TTP.git", package = "ttp-native" } +ttp-core = { path = "../ttp/core" } +ttp-native = { path = "../ttp/native" } actix-web = { version = "4.12.1", features = ["rustls-0_23"] } aes-gcm = "*" diff --git a/src/main.rs b/src/main.rs index 8143c9f..a6b127a 100644 --- a/src/main.rs +++ b/src/main.rs @@ -1,8 +1,10 @@ +mod notifications; mod server; mod sql; mod transport; mod util; +use crate::notifications::tauri; use crate::sql::sql::initialize_db; use crate::sql::sql::print_users; use crate::transport::omikron_connection; @@ -42,6 +44,13 @@ async fn main() { } }); + tokio::spawn(async move { + match tauri::start(9189).await { + Err(e) => log_err!(0, PrintType::General, "{:?}", e), + _ => {} + } + }); + log!("Started"); log!(" .env"); if let Err(e) = initialize_db().await { diff --git a/src/notifications/mod.rs b/src/notifications/mod.rs new file mode 100644 index 0000000..62411f7 --- /dev/null +++ b/src/notifications/mod.rs @@ -0,0 +1 @@ +pub mod tauri; diff --git a/src/notifications/tauri.rs b/src/notifications/tauri.rs new file mode 100644 index 0000000..c57c26d --- /dev/null +++ b/src/notifications/tauri.rs @@ -0,0 +1,160 @@ +use crate::log; +use crate::util::file_util::load_file_vec; +use crate::util::logger::PrintType; +use dashmap::DashMap; +use once_cell::sync::Lazy; +use std::sync::Arc; +use std::time::Duration; +use tokio::sync::futures; +use ttp_core::{CommunicationType, CommunicationValue, DataTypes, DataValue}; +use ttp_native::{Host, Policy, Receiver, SendMode, Sender}; + +pub struct TauriConnection { + pub user_id: i64, + pub sender: Arc, +} + +static TAURI_CONNECTIONS: Lazy>>> = Lazy::new(DashMap::new); + +pub async fn start(port: u16) -> Result<(), Box> { + let cert_pem = load_file_vec("certs", "transport_cert.pem").expect("Error loading Pemfile"); + let key_pem = load_file_vec("certs", "transport_key.pem").expect("Error loading Keyfile"); + + let mut host: Host = ttp_native::host( + port, + cert_pem, + key_pem, + Policy { + send_mode: SendMode::SingleStreamPerMessage, + max_message_size: 1_000_000_000, + close_frame_len: u32::MAX, + application_close_code: 0, + open_stream_timeout: Duration::from_millis(2_000), + write_timeout: Duration::from_millis(2_000), + accept_stream_timeout: Duration::from_millis(10_000), + read_timeout: Duration::from_millis(30_000), + keep_alive_interval: Some(Duration::from_secs(6)), + max_idle_timeout: Some(Duration::from_secs(30)), + force_close_delay: Duration::from_millis(300), + max_transient_recv_errors: 20, + transient_recv_backoff: Duration::from_millis(100), + receiver_queue_capacity: 1000, + }, + ) + .await?; + + log!("TauriServer listening on port {}", port); + + while let Some((sender, mut receiver)) = host.next().await { + tokio::spawn(async move { + handle_connection(sender, &mut receiver).await; + }); + } + + Ok(()) +} + +async fn handle_connection(sender: Sender, receiver: &mut Receiver) { + let mut current_user_id = 0; + let sender = Arc::new(sender); + + while let Ok(cv) = receiver.receive().await { + match cv.get_type() { + CommunicationType::tauri_identification => { + let user_id = cv.get_data(DataTypes::user_id).as_number().unwrap_or(0); + if user_id != 0 { + current_user_id = user_id; + let conn = Arc::new(TauriConnection { + user_id, + sender: sender.clone(), + }); + + TAURI_CONNECTIONS + .entry(user_id) + .and_modify(|conns| { + conns.retain(|c| !Arc::ptr_eq(&c.sender.handle(), &sender.handle())); + conns.push(conn.clone()); + }) + .or_insert_with(|| vec![conn]); + + let response = + CommunicationValue::new(CommunicationType::success).with_id(cv.get_id()); + if let Err(e) = sender.send(&response).await { + log!( + PrintType::General, + "Failed to send tauri success response: {}", + e + ); + } else { + log!(user_id, PrintType::Client, "Tauri device registered"); + } + } + } + CommunicationType::ping => { + let response = + CommunicationValue::new(CommunicationType::pong).with_id(cv.get_id()); + if let Err(_) = sender.send(&response).await { + break; + } + } + _ => {} + } + } + + // Cleanup + if current_user_id != 0 { + if let Some(mut conns) = TAURI_CONNECTIONS.get_mut(¤t_user_id) { + conns.retain(|c| !Arc::ptr_eq(c.sender.handle(), sender.handle())); + } + } +} + +pub async fn send_notification(user_id: i64, sender_id: i64) { + let conns_opt = { + let entry = TAURI_CONNECTIONS.get(&user_id); + entry.map(|e| e.clone()) + }; + + if let Some(conns) = conns_opt { + let cv = CommunicationValue::new(CommunicationType::push_notification) + .add_data(DataTypes::sender_id, DataValue::Number(sender_id)); + + let mut remove_needed = false; + for conn in conns.iter() { + if conn.sender.send(&cv).await.is_err() { + remove_needed = true; + } + } + + if remove_needed { + if let Some(mut conns_mut) = TAURI_CONNECTIONS.get_mut(&user_id) { + conns_mut.retain(|conn| conn.sender.is_open()); + } + } + } +} + +pub async fn remove_notification(user_id: i64, sender_id: i64) { + let conns_opt = { + let entry = TAURI_CONNECTIONS.get(&user_id); + entry.map(|e| e.clone()) + }; + + if let Some(conns) = conns_opt { + let cv = CommunicationValue::new(CommunicationType::read_notification) + .add_data(DataTypes::sender_id, DataValue::Number(sender_id)); + + let mut remove_needed = false; + for conn in conns.iter() { + if conn.sender.send(&cv).await.is_err() { + remove_needed = true; + } + } + + if remove_needed { + if let Some(mut conns_mut) = TAURI_CONNECTIONS.get_mut(&user_id) { + conns_mut.retain(|conn| conn.sender.is_open()); + } + } + } +} diff --git a/src/transport/omikron_connection.rs b/src/transport/omikron_connection.rs index d3453e0..8205775 100644 --- a/src/transport/omikron_connection.rs +++ b/src/transport/omikron_connection.rs @@ -966,8 +966,8 @@ impl OmikronConnection { cv: CommunicationValue, ) -> OmikronResult<()> { let user_id = cv.get_sender() as i64; - if let Ok(notifications) = sql::get_notifications(user_id).await { - let json_array: Vec = notifications + let response_array = match sql::get_notifications(user_id).await { + Ok(notifications) => notifications .into_iter() .map(|(sender, amount)| { DataValue::Container(vec![ @@ -975,59 +975,90 @@ impl OmikronConnection { (DataTypes::amount, DataValue::Number(amount)), ]) }) - .collect(); + .collect(), + Err(e) => { + log!(PrintType::General, "SQL get_notifications error: {}", e); + vec![] + } + }; - let response = CommunicationValue::new(CommunicationType::get_notifications) - .with_id(cv.get_id()) - .add_data(DataTypes::notifications, DataValue::Array(json_array)); - self.send(&response).await - } else { - Ok(()) - } + let response = CommunicationValue::new(CommunicationType::get_notifications) + .with_id(cv.get_id()) + .add_data(DataTypes::notifications, DataValue::Array(response_array)); + self.send(&response).await } async fn handle_read_notification( self: Arc, cv: CommunicationValue, ) -> OmikronResult<()> { - let user_id = cv.get_sender() as i64; + let receiver_id = match cv.get_sender() { + s if s > 0 => s as i64, + _ => match cv.get_data(DataTypes::receiver_id).as_number() { + Some(id) => id as i64, + None => return Ok(()), + }, + }; + if let Some(other_id) = cv .get_data(DataTypes::sender_id) .as_number() .map(|n| n as i64) { - if sql::read_notification(user_id, other_id).await.is_ok() { + if let Err(e) = sql::read_notification(receiver_id, other_id).await { + log!(PrintType::General, "SQL read_notification error: {}", e); + } else { let response = CommunicationValue::new(CommunicationType::read_notification) .with_id(cv.get_id()); - self.send(&response).await - } else { - Ok(()) + let _ = self.send(&response).await; + + // Sync with Tauri + crate::notifications::tauri::remove_notification(receiver_id, other_id).await; + + // Sync with other Omikron clients + let sync_cv = CommunicationValue::new(CommunicationType::read_notification) + .with_receiver(receiver_id as u64) + .add_data(DataTypes::sender_id, DataValue::Number(other_id)); + crate::transport::omikron_manager::send_to_user(receiver_id, &sync_cv).await; } - } else { - Ok(()) } + Ok(()) } async fn handle_push_notification( self: Arc, cv: CommunicationValue, ) -> OmikronResult<()> { - let user_id = cv.get_sender() as i64; - if let Some(other_id) = cv - .get_data(DataTypes::sender_id) - .as_number() - .map(|n| n as i64) - { - if sql::add_notification(user_id, other_id).await.is_ok() { - let response = CommunicationValue::new(CommunicationType::push_notification) - .with_id(cv.get_id()); - self.send(&response).await - } else { - Ok(()) - } + let receiver_id = match cv.get_receiver() { + r if r > 0 => r as i64, + _ => match cv.get_data(DataTypes::receiver_id).as_number() { + Some(id) => id as i64, + None => return Ok(()), + }, + }; + + let sender_id = match cv.get_data(DataTypes::sender_id).as_number() { + Some(id) => id as i64, + None => cv.get_sender() as i64, + }; + + if let Err(e) = sql::add_notification(receiver_id, sender_id).await { + log!(PrintType::General, "SQL add_notification error: {}", e); } else { - Ok(()) + let response = + CommunicationValue::new(CommunicationType::push_notification).with_id(cv.get_id()); + let _ = self.send(&response).await; + + // Sync with Tauri + crate::notifications::tauri::send_notification(receiver_id, sender_id).await; + + // Sync with other Omikron clients + let push_cv = CommunicationValue::new(CommunicationType::push_notification) + .with_receiver(receiver_id as u64) + .add_data(DataTypes::sender_id, DataValue::Number(sender_id)); + crate::transport::omikron_manager::send_to_user(receiver_id, &push_cv).await; } + Ok(()) } async fn handle_ping(self: Arc, cv: CommunicationValue) -> OmikronResult<()> { diff --git a/src/transport/omikron_manager.rs b/src/transport/omikron_manager.rs index fbee4a9..70c73cd 100644 --- a/src/transport/omikron_manager.rs +++ b/src/transport/omikron_manager.rs @@ -3,6 +3,7 @@ use dashmap::DashMap; use once_cell::sync::Lazy; use rand::prelude::IteratorRandom; use std::sync::Arc; +use ttp_core::CommunicationValue; pub static OMIKRON_CONNECTIONS: Lazy>> = Lazy::new(|| DashMap::new()); @@ -24,6 +25,7 @@ pub async fn add_omikron(conn: Arc) { pub async fn remove_omikron(omikron_id: i64) { OMIKRON_CONNECTIONS.remove(&omikron_id); } + pub async fn get_random_omikron() -> Result, ()> { let mut rng = rand::thread_rng(); @@ -37,3 +39,11 @@ pub async fn get_random_omikron() -> Result, ()> { Err(()) } + +pub async fn send_to_user(user_id: i64, cv: &CommunicationValue) { + if let Some(user_conn) = crate::sql::user_online_tracker::get_user_status(user_id) { + if let Some(omikron_conn) = OMIKRON_CONNECTIONS.get(&user_conn.omikron_id) { + let _ = omikron_conn.value().clone().send_message(cv).await; + } + } +} diff --git a/src/util/logger.rs b/src/util/logger.rs index e2e6c92..7fa3947 100644 --- a/src/util/logger.rs +++ b/src/util/logger.rs @@ -1,5 +1,5 @@ use std::{ - collections::{BTreeMap, HashMap}, + collections::BTreeMap, fs::{self, OpenOptions}, io::Write, path::Path, @@ -9,7 +9,6 @@ use std::{ }; use ansi_term::Color; -use json::JsonValue; use ttp_core::{CommunicationValue, DataTypes, DataValue}; static LOGGER: OnceLock> = OnceLock::new(); @@ -119,6 +118,16 @@ pub fn log_internal( #[macro_export] macro_rules! log { + // sender + actor + ($sender:expr, $kind:path, $($arg:tt)*) => { + $crate::util::logger::log_internal(Some($sender), $kind, "", false, format!($($arg)*)) + }; + + // actor only + ($kind:path, $($arg:tt)*) => { + $crate::util::logger::log_internal(None, $kind, "", false, format!($($arg)*)) + }; + // plain ($($arg:tt)*) => { $crate::util::logger::log_internal( @@ -129,27 +138,17 @@ macro_rules! log { format!($($arg)*) ) }; - - // sender + actor - ($sender:expr, $kind:expr, $($arg:tt)*) => { - $crate::util::logger::log_internal(Some($sender), $kind, "", false, format!($($arg)*)) - }; - - // actor only - ($kind:expr, $($arg:tt)*) => { - $crate::util::logger::log_internal(None, $kind, "", false, format!($($arg)*)) - }; } #[macro_export] macro_rules! log_in { // sender + actor - ($sender:expr, $kind:expr, $($arg:tt)*) => { + ($sender:expr, $kind:path, $($arg:tt)*) => { $crate::util::logger::log_internal(Some($sender), $kind, ">", false, format!($($arg)*)) }; // actor only - ($kind:expr, $($arg:tt)*) => { + ($kind:path, $($arg:tt)*) => { $crate::util::logger::log_internal(None, $kind, ">", false, format!($($arg)*)) }; @@ -167,14 +166,13 @@ macro_rules! log_in { #[macro_export] macro_rules! log_out { - // sender + actor - ($sender:expr, $kind:expr, $($arg:tt)*) => { + ($sender:expr, $kind:path, $($arg:tt)*) => { $crate::util::logger::log_internal(Some($sender), $kind, "<", false, format!($($arg)*)) }; // actor only - ($kind:expr, $($arg:tt)*) => { + ($kind:path, $($arg:tt)*) => { $crate::util::logger::log_internal(None, $kind, "<", false, format!($($arg)*)) }; @@ -192,14 +190,13 @@ macro_rules! log_out { #[macro_export] macro_rules! log_err { - // sender + actor - ($sender:expr, $kind:expr, $($arg:tt)*) => { + ($sender:expr, $kind:path, $($arg:tt)*) => { $crate::util::logger::log_internal(Some($sender), $kind, ">>", true, format!($($arg)*)) }; // actor only - ($kind:expr, $($arg:tt)*) => { + ($kind:path, $($arg:tt)*) => { $crate::util::logger::log_internal(None, $kind, ">>", true, format!($($arg)*)) };