diff --git a/src/auth/auth_connector.rs b/src/auth/auth_connector.rs deleted file mode 100644 index 485d952..0000000 --- a/src/auth/auth_connector.rs +++ /dev/null @@ -1,184 +0,0 @@ -use crate::data::communication::{CommunicationType, CommunicationValue, DataTypes}; -use crate::log; -use crate::util::config_util::CONFIG; -use crate::util::logger::PrintType; -use json::number::Number; -use reqwest::{Client, Response}; -use std::time::Duration; -use uuid::Uuid; - -#[derive(Debug, Clone)] -pub struct AuthUser { - pub created_at: i64, - pub username: String, - pub display: String, - pub avatar: String, - pub about: String, - pub status: String, - pub public_key: String, - pub sub_level: i32, - pub sub_end: i32, -} - -fn client() -> Client { - Client::builder() - .connect_timeout(Duration::from_secs(100)) - .timeout(Duration::from_secs(150)) - .build() - .unwrap() -} -pub async fn get_auth_public_key() -> Option { - let url = format!("https://auth.tensamin.net/api/get/public_key"); - let client = client(); - let res = client.get(&url).send().await.ok()?; - let json = res.text().await.ok()?; - - let cv = CommunicationValue::from_json(&json); - if cv.comm_type != CommunicationType::success { - return None; - } - - Some(cv.get_data(DataTypes::public_key).unwrap().to_string()) -} - -pub async fn get_user(user_id: Uuid) -> Option { - let url = format!("https://auth.tensamin.net/api/get/{}", user_id); - let client = client(); - let res = client.get(&url).send().await.ok()?; - let json = res.text().await.ok()?; - - let cv = CommunicationValue::from_json(&json); - if cv.comm_type != CommunicationType::success { - return None; - } - - Some(AuthUser { - created_at: cv - .get_data(DataTypes::created_at) - .unwrap() - .to_string() - .parse::() - .unwrap_or(-1), - username: cv.get_data(DataTypes::username).unwrap().to_string(), - display: cv.get_data(DataTypes::display).unwrap().to_string(), - avatar: cv.get_data(DataTypes::avatar).unwrap().to_string(), - about: cv.get_data(DataTypes::about).unwrap().to_string(), - status: cv.get_data(DataTypes::status).unwrap().to_string(), - public_key: cv.get_data(DataTypes::public_key).unwrap().to_string(), - sub_level: cv - .get_data(DataTypes::sub_level) - .unwrap() - .to_string() - .parse::() - .unwrap_or(-1), - sub_end: cv - .get_data(DataTypes::sub_end) - .unwrap() - .to_string() - .parse::() - .unwrap_or(-1), - }) -} - -pub async fn get_iota_by_user_id(user_id: i64) -> Option { - let url = format!("https://auth.tensamin.net/api/get/iota-id/{}", user_id); - - let client = client(); - let res = client - .get(&url) - .header("Authorization", CONFIG.read().await.omikron_id.to_string()) - .header("Content-Type", "application/json") - .send() - .await - .ok()?; - - let json = res.text().await.ok()?; - let json = json.replace("iota_uuid", "iota_id"); - - let cv = CommunicationValue::from_json(&json); - if cv.comm_type != CommunicationType::success { - log!(PrintType::Iota, "{}", &cv.to_json().to_string()); - return None; - } - - let iota_id = cv.get_data(DataTypes::iota_id)?.as_i64().unwrap_or(0); - if iota_id == 0 { - return None; - } - Some(iota_id) -} -pub async fn is_private_key_valid(user_id: i64, pk_hash: &str) -> bool { - let url = format!( - "https://auth.tensamin.net/api/get/private-key-hash/{}/", - user_id - ); - - let client = client(); - let res = client - .get(&url) - .header("Authorization", CONFIG.read().await.omikron_id.to_string()) - .header("PrivateKeyHash", pk_hash) - .header("Accept", "application/json") - .send() - .await; - - let Ok(response) = res else { - return false; - }; - - let Ok(body) = response.text().await else { - return false; - }; - - let cv = CommunicationValue::from_json(&body); - if cv.comm_type != CommunicationType::success { - return false; - } - - match cv.get_data(DataTypes::matches) { - Some(val) => val.as_bool().unwrap_or(false), - None => false, - } -} -pub async fn get_public_key(user_id: i64) -> Option { - let url = format!("https://auth.tensamin.net/api/{}/public-key", user_id); - - let client = client(); - let res = client - .get(&url) - .header("Accept", "application/json") - .send() - .await - .ok()?; - - let body = res.text().await.ok()?; - let cv = CommunicationValue::from_json(&body); - - if cv.comm_type != CommunicationType::message_send { - return None; - } - - Some(cv.get_data(DataTypes::ping_clients)?.to_string()) -} - -pub async fn get_register() -> Option { - let url = "https://auth.tensamin.net/api/register/init".to_string(); - let client = client(); - let res = client.get(&url).send().await.ok()?; - let json = res.text().await.ok()?; - - let cv = CommunicationValue::from_json(&json); - cv.get_data(DataTypes::user_id) - .unwrap_or(&json::JsonValue::Number(Number::from(0))) - .as_i64() -} - -async fn handle_response(resp: Response) -> bool { - match resp.text().await { - Ok(text) => { - let cv = CommunicationValue::from_json(&text.to_string()); - cv.comm_type == CommunicationType::success - } - Err(_) => false, - } -} diff --git a/src/auth/mod.rs b/src/auth/mod.rs deleted file mode 100644 index 8eef5c3..0000000 --- a/src/auth/mod.rs +++ /dev/null @@ -1,2 +0,0 @@ -pub mod auth_connector; -pub mod crypto_helper; diff --git a/src/calls/call_manager.rs b/src/calls/call_manager.rs index d02413d..e0a06ae 100644 --- a/src/calls/call_manager.rs +++ b/src/calls/call_manager.rs @@ -108,7 +108,8 @@ pub async fn clean_calls() { return; } }; - let room_service = RoomClient::with_api_key("https://call.tensamin.net", &api_key, &api_secret); + let room_service = + RoomClient::with_api_key("https://call.tensamin.net/", &api_key, &api_secret); let rooms = match room_service.list_rooms(Vec::new()).await { Ok(rooms) => rooms, diff --git a/src/data/communication.rs b/src/data/communication.rs index 4dd6d43..de68e6b 100644 --- a/src/data/communication.rs +++ b/src/data/communication.rs @@ -10,6 +10,7 @@ pub enum DataTypes { error_type, accepted_ids, uuid, + register_id, settings, settings_name, chat_partner_id, @@ -92,6 +93,7 @@ impl DataTypes { match normalized.as_str() { "errortype" => DataTypes::error_type, "chatpartnerid" => DataTypes::chat_partner_id, + "registerid" => DataTypes::register_id, "uuid" => DataTypes::uuid, "settings" => DataTypes::settings, "settingsname" => DataTypes::settings_name, @@ -174,6 +176,7 @@ impl DataTypes { #[allow(non_camel_case_types, dead_code)] pub enum CommunicationType { error, + error_internal, error_invalid_data, error_invalid_user_id, error_invalid_omikron_id, @@ -211,6 +214,8 @@ pub enum CommunicationType { register_response, identification, identification_response, + register_iota, + register_iota_success, ping, pong, add_chat, @@ -273,6 +278,7 @@ impl CommunicationType { "function" => CommunicationType::function, "update" => CommunicationType::update, "createuser" => CommunicationType::create_user, + "errorinternal" => CommunicationType::error_internal, "errorinvaliddata" => CommunicationType::error_invalid_data, "errorinvaliduserid" => CommunicationType::error_invalid_user_id, "errorinvalidomikronid" => CommunicationType::error_invalid_omikron_id, @@ -310,6 +316,8 @@ impl CommunicationType { "registerresponse" => CommunicationType::register_response, "identification" => CommunicationType::identification, "identificationresponse" => CommunicationType::identification_response, + "registeriota" => CommunicationType::register_iota, + "registeriotasuccess" => CommunicationType::register_iota_success, "ping" => CommunicationType::ping, "pong" => CommunicationType::pong, "addchat" => CommunicationType::add_chat, diff --git a/src/main.rs b/src/main.rs index d641ef4..d7fa6ae 100644 --- a/src/main.rs +++ b/src/main.rs @@ -1,4 +1,3 @@ -mod auth; mod calls; mod data; mod omega; @@ -15,12 +14,12 @@ use tokio_util::compat::TokioAsyncReadCompatExt; use tungstenite::handshake::server::{Request, Response}; use crate::{ - auth::crypto_helper::{load_public_key, load_secret_key}, calls::call_manager::garbage_collect_calls, omega::omega_connection::OmegaConnection, rho::{client_connection::ClientConnection, iota_connection::IotaConnection}, util::{ config_util::CONFIG, + crypto_helper::{load_public_key, load_secret_key}, logger::{PrintType, startup}, }, }; diff --git a/src/omega/omega_connection.rs b/src/omega/omega_connection.rs old mode 100644 new mode 100755 index 49d73d3..1b902fb --- a/src/omega/omega_connection.rs +++ b/src/omega/omega_connection.rs @@ -11,23 +11,23 @@ use once_cell::sync::Lazy; use std::{collections::HashMap, env, sync::Arc, time::Duration}; use tokio::{ net::TcpStream, - sync::{Mutex, RwLock}, + sync::{Mutex, RwLock, mpsc}, time::{Instant, sleep}, }; use tokio_native_tls::TlsStream; use uuid::Uuid; +use crate::log_err; use crate::{ - auth::crypto_helper::decrypt, data::{ communication::{CommunicationType, CommunicationValue, DataTypes}, user::UserStatus, }, get_private_key, log, log_in, log_out, rho::rho_manager::{self, RHO_CONNECTIONS}, + util::crypto_helper::{decrypt, load_public_key}, util::logger::PrintType, }; -use crate::{auth::crypto_helper::secret_key_to_base64, log_err}; pub static WAITING_TASKS: Lazy< DashMap, CommunicationValue) -> bool + Send + Sync>>, @@ -163,9 +163,11 @@ impl OmegaConnection { .to_string() })?; + let server_pub_key_obj = load_public_key(server_pub_key).ok_or("Failed to load public key".to_string())?; + let decrypted_challenge = decrypt( - &secret_key_to_base64(&get_private_key()), - server_pub_key, + get_private_key(), + server_pub_key_obj, challenge, ) .map_err(|e| { @@ -195,6 +197,16 @@ impl OmegaConnection { return false; } + if let Some(accepted) = final_cv.get_data(DataTypes::accepted).and_then(|v| v.as_bool()) { + if !accepted { + log_err!(PrintType::Omega, "Omega did not accept identification."); + return false; + } + } else { + log_err!(PrintType::Omega, "Omega response did not contain 'accepted' field."); + return false; + } + tokio::spawn(async move { let mut connected_iota_ids: Vec = Vec::new(); let mut connected_user_ids: Vec = Vec::new(); @@ -399,4 +411,47 @@ impl OmegaConnection { async fn send_global(cv: CommunicationValue) { OMEGA_CONNECTION.send_message(&cv).await; } + + pub async fn await_response( + &self, + cv: &CommunicationValue, + timeout_duration: Option, + ) -> Result { + let (tx, mut rx) = mpsc::channel(1); + let msg_id = cv.get_id(); + + let task_tx = tx.clone(); + WAITING_TASKS.insert( + msg_id, + Box::new(move |_, response_cv| { + let inner_tx = task_tx.clone(); + tokio::spawn(async move { + if let Err(e) = inner_tx.send(response_cv).await { + log_err!( + PrintType::Omega, + "Failed to send response back to awaiter: {}", + e + ); + } + }); + true + }), + ); + + 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()), + Err(_) => { + WAITING_TASKS.remove(&msg_id); + Err(format!( + "Request timed out after {} seconds.", + timeout.as_secs() + )) + } + } + } } diff --git a/src/rho/client_connection.rs b/src/rho/client_connection.rs index 838d579..bf6cfe9 100644 --- a/src/rho/client_connection.rs +++ b/src/rho/client_connection.rs @@ -12,7 +12,6 @@ use crate::calls::call_manager; use crate::omega::omega_connection::{WAITING_TASKS, get_omega_connection}; use crate::util::logger::PrintType; use crate::{ - auth::auth_connector, // calls::call_manager::CallManager, data::{ communication::{CommunicationType, CommunicationValue, DataTypes}, @@ -190,8 +189,9 @@ impl ClientConnection { // Validate private key if let Some(private_key_hash) = cv.get_data(DataTypes::private_key_hash) { println!("private_key_hash: {}", private_key_hash); - let is_valid = - auth_connector::is_private_key_valid(user_id, &private_key_hash.to_string()).await; + let is_valid = true; // NO VALIDATION, + // SWAP TO AUTH VIA CHALLENGE + // auth_connector::is_private_key_valid(user_id, &private_key_hash.to_string()).await; if !is_valid { println!("Invalid private key"); diff --git a/src/rho/iota_connection.rs b/src/rho/iota_connection.rs index 11d87f6..0a56366 100644 --- a/src/rho/iota_connection.rs +++ b/src/rho/iota_connection.rs @@ -1,6 +1,3 @@ -use crate::auth::crypto_helper::encrypt; -use crate::auth::crypto_helper::load_public_key; -use crate::auth::crypto_helper::public_key_to_base64; use crate::calls::call_group::CallGroup; use crate::calls::call_manager; use crate::get_private_key; @@ -10,14 +7,15 @@ use crate::log_in; use crate::log_out; use crate::omega::omega_connection::WAITING_TASKS; use crate::omega::omega_connection::get_omega_connection; +use crate::util::crypto_helper::encrypt; +use crate::util::crypto_helper::load_public_key; +use crate::util::crypto_helper::public_key_to_base64; use crate::util::logger::PrintType; use async_tungstenite::WebSocketReceiver; use async_tungstenite::WebSocketSender; use async_tungstenite::tungstenite::Message; -use base64::alphabet::STANDARD; use dashmap::DashMap; use json::JsonValue; -use json::number::Number; use rand::Rng; use rand::distributions::Alphanumeric; use std::{ @@ -32,7 +30,6 @@ use x448::PublicKey; use super::{rho_connection::RhoConnection, rho_manager}; use crate::{ - auth::auth_connector, // calls::call_manager::CallManager, data::communication::{CommunicationType, CommunicationValue, DataTypes}, omega::omega_connection::OmegaConnection, @@ -140,111 +137,241 @@ impl IotaConnection { /// Handle incoming message from Iota pub async fn handle_message(self: Arc, message: Utf8Bytes) { let cv = CommunicationValue::from_json(&message); + let identified = *self.identified.read().await; + let challenged = *self.challenged.read().await; - // Handle identification - if !self.is_identified().await { - let identified = *self.identified.read().await; - let challenged = *self.challenged.read().await; - - if !identified && cv.is_type(CommunicationType::identification) { - let iota_id = cv - .get_data(DataTypes::iota_id) - .and_then(|v| v.as_i64()) - .unwrap_or(0); - - let iota_for_closure: Arc = self.clone(); - WAITING_TASKS.insert( - cv.get_id(), - Box::new(|omega_conn: Arc, cv: CommunicationValue| { - let base64_pub = cv - .get_data(DataTypes::public_key) - .unwrap_or(&JsonValue::Null) - .as_str() - .unwrap_or(""); - let pub_key: PublicKey = match load_public_key(base64_pub) { - Some(b) => { - let pub_key_bytes: Vec = b.as_bytes().to_vec(); - let iota: Arc = iota_for_closure.clone(); - tokio::spawn(async move { - *iota.pub_key.write().await = Some(pub_key_bytes) - }); - b - } - _ => { - let iota: Arc = iota_for_closure.clone(); - tokio::spawn(async move { - iota.send_message( - &CommunicationValue::new( - CommunicationType::error_invalid_omikron_id, - ) - .with_id(cv.get_id()), - ) - .await; - }); - return true; - } - }; - tokio::spawn(async move { - let challenge: String = rand::thread_rng() - .sample_iter(&Alphanumeric) - .take(32) - .map(char::from) - .collect(); - - *self.iota_id.write().await = iota_id; - *self.challenge.write().await = challenge.clone(); - *self.identified.write().await = true; - - let encrypted = - encrypt(get_private_key(), pub_key, &challenge).unwrap_or_default(); - - let response = CommunicationValue::new(CommunicationType::challenge) - .with_id(cv.get_id()) - .add_data_str( - DataTypes::public_key, - public_key_to_base64(&get_public_key()), - ) - .add_data_str(DataTypes::challenge, encrypted); - - self.send_message(&response).await; - }); - true - }), - ); - } - - // ────────────────────────────── - // Challenge response - // ────────────────────────────── - if identified && !challenged && cv.is_type(CommunicationType::challenge_response) { - let client_response = cv - .get_data(DataTypes::challenge) - .and_then(|v| v.as_str()) - .unwrap_or(""); - - if client_response == *self.challenge.read().await { - *self.challenged.write().await = true; - let _ = sql::set_omikron_active(self.iota_id.await, true); - - self.send_message( - &CommunicationValue::new(CommunicationType::identification_response) - .with_id(cv.get_id()), - ) + if !identified && cv.is_type(CommunicationType::identification) { + let iota_id = cv + .get_data(DataTypes::iota_id) + .and_then(|v| v.as_i64()) + .unwrap_or(0); + if iota_id == 0 { + self.send_error_response(&cv.get_id(), CommunicationType::error_invalid_data) .await; - } else { - self.send_error_response( - &cv.get_id(), - CommunicationType::error_invalid_challenge, - ) - .await; - self.close().await; - } + self.close().await; return; } + let user_ids_json = cv + .get_data(DataTypes::user_ids) + .unwrap_or(&JsonValue::Null) + .clone(); + let mut user_ids = Vec::new(); + if let JsonValue::Array(ids) = user_ids_json { + for id_val in ids { + if let Some(id) = id_val.as_i64() { + user_ids.push(id); + } + } + } + + *self.iota_id.write().await = iota_id; + *self.user_ids.write().await = user_ids; + + let get_pub_key_msg = CommunicationValue::new(CommunicationType::get_iota_data) + .add_data(DataTypes::iota_id, JsonValue::from(iota_id)); + + let msg_id = get_pub_key_msg.get_id(); + let iota_conn_clone = self.clone(); + let original_cv_id = cv.get_id(); + + WAITING_TASKS.insert( + msg_id, + Box::new(move |_, response_cv: CommunicationValue| { + let iota_conn_for_task = iota_conn_clone.clone(); + tokio::spawn(async move { + if !response_cv.is_type(CommunicationType::get_iota_data) { + iota_conn_for_task + .send_error_response( + &original_cv_id, + CommunicationType::error_internal, + ) + .await; + iota_conn_for_task.close().await; + return; + } + + let base64_pub = response_cv + .get_data(DataTypes::public_key) + .and_then(|v| v.as_str()) + .unwrap_or(""); + + let pub_key = match load_public_key(base64_pub) { + Some(pk) => pk, + None => { + iota_conn_for_task + .send_error_response( + &original_cv_id, + CommunicationType::error_invalid_public_key, + ) + .await; + iota_conn_for_task.close().await; + return; + } + }; + + *iota_conn_for_task.pub_key.write().await = + Some(pub_key.as_bytes().to_vec()); + + let challenge: String = rand::thread_rng() + .sample_iter(&Alphanumeric) + .take(32) + .map(char::from) + .collect(); + + *iota_conn_for_task.challenge.write().await = challenge.clone(); + + let encrypted_challenge = + encrypt(get_private_key(), pub_key, &challenge).unwrap_or_default(); + + *iota_conn_for_task.identified.write().await = true; + + let challenge_msg = CommunicationValue::new(CommunicationType::challenge) + .with_id(original_cv_id) + .add_data_str( + DataTypes::public_key, + public_key_to_base64(&get_public_key()), + ) + .add_data_str(DataTypes::challenge, encrypted_challenge); + + iota_conn_for_task.send_message(&challenge_msg).await; + }); + true + }), + ); + + get_omega_connection().send_message(&get_pub_key_msg).await; + return; + } else if !identified && cv.is_type(CommunicationType::complete_register_iota) { + let base64_pub = cv + .get_data(DataTypes::public_key) + .and_then(|v| v.as_str()) + .unwrap_or(""); + + if base64_pub.is_empty() { + self.send_error_response(&cv.get_id(), CommunicationType::error_invalid_public_key) + .await; + self.close().await; + return; + } + + let register_msg = CommunicationValue::new(CommunicationType::complete_register_iota) + .add_data( + DataTypes::public_key, + JsonValue::String(base64_pub.to_string()), + ); + + let msg_id = register_msg.get_id(); + let iota_conn_clone = self.clone(); + let original_cv_id = cv.get_id(); + + WAITING_TASKS.insert( + msg_id, + Box::new(move |_, response_cv: CommunicationValue| { + let iota_conn_for_task = iota_conn_clone.clone(); + tokio::spawn(async move { + if !response_cv.is_type(CommunicationType::complete_register_iota) { + iota_conn_for_task + .send_error_response( + &original_cv_id, + CommunicationType::error_internal, + ) + .await; + iota_conn_for_task.close().await; + return; + } + + let new_iota_id = response_cv + .get_data(DataTypes::iota_id) + .and_then(|v| v.as_i64()) + .unwrap_or(0); + + if new_iota_id == 0 { + iota_conn_for_task + .send_error_response( + &original_cv_id, + CommunicationType::error_internal, + ) + .await; + iota_conn_for_task.close().await; + return; + } + + *iota_conn_for_task.iota_id.write().await = new_iota_id; + *iota_conn_for_task.identified.write().await = true; + + let success_msg = CommunicationValue::new(CommunicationType::success) + .with_id(original_cv_id) + .add_data(DataTypes::iota_id, JsonValue::from(new_iota_id)); + + iota_conn_for_task.send_message(&success_msg).await; + }); + true + }), + ); + get_omega_connection().send_message(®ister_msg).await; + return; + } + + if identified && !challenged && cv.is_type(CommunicationType::challenge_response) { + let client_response = cv + .get_data(DataTypes::challenge) + .and_then(|v| v.as_str()) + .unwrap_or(""); + + if client_response == *self.challenge.read().await { + *self.challenged.write().await = true; + + let iota_id = self.get_iota_id().await; + let user_ids = self.get_user_ids().await; + + let mut validated_user_ids: Vec = Vec::new(); + for user_id in user_ids { + validated_user_ids.push(user_id); + } + + if rho_manager::contains_iota(iota_id).await { + if let Some(existing_rho) = rho_manager::get_rho_by_iota(iota_id).await { + existing_rho.close_iota_connection().await; + } + } + + let rho_connection = + Arc::new(RhoConnection::new(self.clone(), validated_user_ids.clone()).await); + + self.set_rho_connection(Arc::downgrade(&rho_connection)) + .await; + + rho_manager::add_rho(rho_connection).await; + + let mut str = String::new(); + for id in &validated_user_ids { + str.push_str(&format!(",{}", id)); + } + if !str.is_empty() { + str.remove(0); + } + + self.send_message( + &CommunicationValue::new(CommunicationType::identification_response) + .with_id(cv.get_id()) + .add_data_str(DataTypes::accepted_ids, str) + .add_data_str(DataTypes::accepted, validated_user_ids.len().to_string()), + ) + .await; + } else { + self.send_error_response(&cv.get_id(), CommunicationType::error_invalid_challenge) + .await; + self.close().await; + } + return; + } + + if !self.is_identified().await { self.send_error_response(&cv.get_id(), CommunicationType::error_not_authenticated) .await; self.close().await; + return; } // Handle ping @@ -282,6 +409,17 @@ impl IotaConnection { // Forward to client self.forward_to_client(cv).await; } + + async fn send_error_response(&self, message_id: &Uuid, error_type: CommunicationType) { + let error = CommunicationValue::new(error_type).with_id(*message_id); + self.send_message(&error).await; + } + + async fn close(&self) { + let mut sender = self.sender.write().await; + let _ = sender.close(None).await; + } + async fn handle_omega_forward(self: Arc, cv: CommunicationValue) { let iota_for_closure = self.clone(); WAITING_TASKS.insert( @@ -298,98 +436,6 @@ impl IotaConnection { .send_message(&cv.with_sender(*self.iota_id.read().await)) .await; } - /// Handle identification message - async fn handle_identification(self: Arc, cv: CommunicationValue) { - let iota_id: i64 = cv - .get_data(DataTypes::iota_id) - .unwrap_or(&JsonValue::Number(Number::from(0))) - .as_i64() - .unwrap_or(0); - - if iota_id == 0 { - let error = CommunicationValue::new(CommunicationType::error).with_id(cv.get_id()); - self.send_message(&error).await; - return; - } - - // Parse user IDs - let mut validated_user_ids: Vec = Vec::new(); - if let Some(user_ids_str) = cv.get_data(DataTypes::user_ids) { - for id_str in user_ids_str.to_string().split(',') { - match id_str.parse::() { - Ok(user_id) => { - if let Some(auth_iota_id) = auth_connector::get_iota_by_id(user_id).await { - log_in!( - PrintType::Iota, - "auth for {} should be {} is {}", - user_id, - iota_id, - auth_iota_id - ); - if auth_iota_id == iota_id { - validated_user_ids.push(user_id); - } - } else { - log_in!(PrintType::Iota, "User ID {} not parsed", user_id); - } - } - Err(e) => { - log_in!( - PrintType::Iota, - "Failed to parse '{}' as i64: {:?}", - id_str, - e, - ); - } - } - } - } - - // Set identification data - { - let mut iota_id_guard = self.iota_id.write().await; - *iota_id_guard = iota_id; - } - { - let mut user_ids_guard = self.user_ids.write().await; - *user_ids_guard = validated_user_ids.clone(); - } - { - let mut identified_guard = self.identified.write().await; - *identified_guard = true; - } - - // Check for existing connection and close it - if rho_manager::contains_iota(iota_id).await { - if let Some(existing_rho) = rho_manager::get_rho_by_iota(iota_id).await { - existing_rho.close_iota_connection().await; - } - } - - // Create RhoConnection - let rho_connection = - Arc::new(RhoConnection::new(self.clone(), validated_user_ids.clone()).await); - - // Set up bidirectional reference - self.set_rho_connection(Arc::downgrade(&rho_connection)) - .await; - - // Add to manager - rho_manager::add_rho(rho_connection).await; - - // Send response - let mut str = String::new(); - for id in &validated_user_ids { - str.push_str(&format!(",{}", id)); - } - let response = CommunicationValue::new(CommunicationType::identification_response) - .with_id(cv.get_id()) - .add_data_str(DataTypes::accepted_ids, str) - .add_data_str(DataTypes::accepted, validated_user_ids.len().to_string()); - - self.send_message(&response).await; - } - /// Handle ping message async fn handle_ping(&self, cv: CommunicationValue) { if let Some(last_ping) = cv.get_data(DataTypes::last_ping) { diff --git a/src/auth/crypto_helper.rs b/src/util/crypto_helper.rs similarity index 100% rename from src/auth/crypto_helper.rs rename to src/util/crypto_helper.rs diff --git a/src/util/mod.rs b/src/util/mod.rs index 58403c4..0ea6aef 100644 --- a/src/util/mod.rs +++ b/src/util/mod.rs @@ -1,3 +1,4 @@ pub mod config_util; +pub mod crypto_helper; pub mod file_util; pub mod logger;