use crate::data::communication::{CommunicationType, CommunicationValue, DataTypes}; use crate::sql::iota_omikron_tracker::{ track_iota_omikron, untrack_by_omikron as untrack_iota_by_omikron, untrack_iota, }; use crate::sql::sql::get_omikron_by_id; use crate::sql::user_online_tracker::{ track_user_omikron, untrack_by_omikron as untrack_user_by_omikron, untrack_user, }; use crate::util::crypto_helper::encrypt; use crate::util::logger::PrintType; use crate::{get_private_key, log_out}; use crate::{get_public_key, log_in}; use base64::{Engine as _, engine::general_purpose::STANDARD}; use dashmap::DashMap; use futures::SinkExt; use futures::stream::SplitSink; use futures::stream::SplitStream; use hyper::upgrade::Upgraded; use hyper_util::rt::TokioIo; use json::JsonValue; use rand::Rng; use rand::distributions::Alphanumeric; use std::sync::Arc; use tokio::sync::RwLock; use tokio_tungstenite::WebSocketStream; use tungstenite::Message; use tungstenite::Utf8Bytes; use uuid::Uuid; use x448::PublicKey; pub struct OmikronConnection { pub sender: Arc>, Message>>>, pub receiver: Arc>>>>, pub omikron_id: Arc>, pub pub_key: Arc>>>, identified: Arc>, challenged: Arc>, challenge: Arc>, pub ping: Arc>, waiting_tasks: DashMap< Uuid, Box, CommunicationValue) -> bool + Send + Sync>, >, } impl OmikronConnection { pub fn new( sender: SplitSink>, Message>, receiver: SplitStream>>, ) -> Arc { Arc::new(Self { sender: Arc::new(RwLock::new(sender)), receiver: Arc::new(RwLock::new(receiver)), omikron_id: Arc::new(RwLock::new(0)), pub_key: Arc::new(RwLock::new(None)), identified: Arc::new(RwLock::new(false)), challenged: Arc::new(RwLock::new(false)), challenge: Arc::new(RwLock::new(String::new())), ping: Arc::new(RwLock::new(-1)), waiting_tasks: DashMap::new(), }) } pub async fn send_message(&self, cv: &CommunicationValue) { let mut sender = self.sender.write().await; let message_text = Message::Text(Utf8Bytes::from(cv.to_json().to_string())); if !cv.is_type(CommunicationType::pong) { log_out!( *self.omikron_id.read().await, PrintType::Omikron, "{}", message_text ); } sender.send(message_text).await.unwrap(); } pub async fn get_user_id(&self) -> i64 { *self.omikron_id.read().await } pub async fn is_identified(&self) -> bool { *self.identified.read().await && *self.challenged.read().await } pub async fn get_public_key(&self) -> PublicKey { PublicKey::from_bytes(self.pub_key.read().await.as_ref().unwrap()).unwrap() } pub async fn handle_message(self: Arc, message: String) { let cv = CommunicationValue::from_json(&message); if cv.is_type(CommunicationType::ping) { self.handle_ping(cv).await; return; } log_in!( *self.omikron_id.read().await, PrintType::Omikron, "{}", message ); if let Some((_, task)) = self.waiting_tasks.remove(&cv.get_id()) { let _ = task(self.clone(), cv.clone()); return; } // Handle identification if !*self.identified.read().await && cv.is_type(CommunicationType::identification) { let omikron_id = cv .get_data(DataTypes::omikron) .unwrap_or(&JsonValue::Null) .as_i64() .unwrap_or(0); match get_omikron_by_id(omikron_id).await { Ok((public_key, _)) => { // Generate Challenge, encrypt it and send it to the omikron *self.omikron_id.write().await = omikron_id; let challenge_str: String = rand::thread_rng() .sample_iter(&Alphanumeric) .take(32) .map(char::from) .collect(); *self.challenge.write().await = challenge_str.clone(); let user_public_key_bytes = match STANDARD.decode(&public_key) { Ok(bytes) => bytes, Err(_) => { self.send_error_response( &cv.get_id(), CommunicationType::error_invalid_omikron_id, ) .await; return; } }; *self.pub_key.write().await = Some(user_public_key_bytes.clone()); let omikron_pub_key: PublicKey = match PublicKey::from_bytes(&user_public_key_bytes) { Some(key) => key, None => { self.send_error_response( &cv.get_id(), CommunicationType::error_invalid_public_key, ) .await; return; } }; let encrypted_challenge = encrypt(get_private_key(), omikron_pub_key, &challenge_str) .unwrap_or("".to_string()); let response = CommunicationValue::new(CommunicationType::challenge) .add_data_str( DataTypes::public_key, STANDARD.encode(get_public_key().as_bytes()), ) .add_data_str(DataTypes::challenge, encrypted_challenge) .with_id(cv.get_id()); self.send_message(&response).await; *self.identified.write().await = true; return; } Err(e) => { self.send_message( &CommunicationValue::new(CommunicationType::error_not_authenticated) .with_id(cv.get_id()) .add_data_str(DataTypes::error_type, e.to_string()), ) .await; return; } } } // Handle challenge response if *self.identified.read().await && !*self.challenged.read().await && cv.is_type(CommunicationType::challenge_response) { let client_response = cv .get_data(DataTypes::challenge) .unwrap_or(&JsonValue::Null) .as_str() .unwrap_or(""); let expected_challenge = self.challenge.read().await.clone(); if client_response == expected_challenge { *self.challenged.write().await = true; let response = CommunicationValue::new(CommunicationType::identification_response) .with_id(cv.get_id()); self.send_message(&response).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; } let omikron_id = self.get_user_id().await; if cv.is_type(CommunicationType::user_connected) { if let Some(user_id) = cv.get_data(DataTypes::user_id).and_then(|v| v.as_i64()) { track_user_omikron(user_id, omikron_id).await; } return; } if cv.is_type(CommunicationType::user_disconnected) { if let Some(user_id) = cv.get_data(DataTypes::user_id).and_then(|v| v.as_i64()) { untrack_user(user_id).await; } return; } if cv.is_type(CommunicationType::iota_connected) { if let Some(iota_id) = cv.get_data(DataTypes::iota_id).and_then(|v| v.as_i64()) { track_iota_omikron(iota_id, omikron_id).await; } return; } if cv.is_type(CommunicationType::iota_disconnected) { if let Some(iota_id) = cv.get_data(DataTypes::iota_id).and_then(|v| v.as_i64()) { untrack_iota(iota_id).await; } return; } if cv.is_type(CommunicationType::sync_client_iota_status) { if let Some(json::JsonValue::Array(user_ids)) = cv.get_data(DataTypes::user_ids).cloned() { for user_id_json in user_ids { if let Some(user_id) = user_id_json.as_i64() { track_user_omikron(user_id, omikron_id).await; } } } if let Some(json::JsonValue::Array(iota_ids)) = cv.get_data(DataTypes::iota_ids).cloned() { for iota_id_json in iota_ids { if let Some(iota_id) = iota_id_json.as_i64() { track_iota_omikron(iota_id, omikron_id).await; } } } return; } } 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; } pub async fn close(&self) { let mut sender = self.sender.write().await; let _ = sender.close().await; } pub async fn handle_close(self: Arc) { if self.is_identified().await { let omikron_id = self.get_user_id().await; if omikron_id != 0 { untrack_iota_by_omikron(omikron_id).await; untrack_user_by_omikron(omikron_id).await; } } } async fn handle_ping(&self, cv: CommunicationValue) { if let Some(last_ping) = cv.get_data(DataTypes::last_ping) { if let Ok(ping_val) = last_ping.to_string().parse::() { let mut ping_guard = self.ping.write().await; *ping_guard = ping_val; } } let response = CommunicationValue::new(CommunicationType::pong).with_id(cv.get_id()); self.send_message(&response).await; } }