From c0652516c6ba33d90fd80d1ff9b22a9c4523869b Mon Sep 17 00:00:00 2001 From: Alex Emmet Date: Sat, 10 Jan 2026 16:54:41 +0000 Subject: [PATCH] mplemented commtypes, change user data --- src/calls/call_manager.rs | 30 ++++++++++++++++++++++++++---- src/data/communication.rs | 18 ++++++++++++++++++ src/omega/omega_connection.rs | 2 +- src/rho/client_connection.rs | 21 +++++++++++++++++++++ 4 files changed, 66 insertions(+), 5 deletions(-) diff --git a/src/calls/call_manager.rs b/src/calls/call_manager.rs index 622c108..d02413d 100644 --- a/src/calls/call_manager.rs +++ b/src/calls/call_manager.rs @@ -1,12 +1,13 @@ use livekit_api::services::room::RoomClient; use once_cell::sync::Lazy; -use std::{str::FromStr, sync::Arc, time::Duration}; +use std::{env, str::FromStr, sync::Arc, time::Duration}; use tokio::sync::RwLock; + use uuid::Uuid; use crate::{ calls::{call_group::CallGroup, caller::Caller}, - log, + log, log_err, util::logger::PrintType, }; @@ -93,8 +94,29 @@ pub fn garbage_collect_calls() { }); } pub async fn clean_calls() { - let room_service = RoomClient::new("https://call.tensamin.net").unwrap(); - let rooms = room_service.list_rooms(Vec::new()).await.unwrap(); + let api_key = match env::var("LIVEKIT_API_KEY") { + Ok(key) => key, + Err(_) => { + log_err!(PrintType::General, "LIVEKIT_API_KEY not set!"); + return; + } + }; + let api_secret = match env::var("LIVEKIT_API_SECRET") { + Ok(secret) => secret, + Err(_) => { + log_err!(PrintType::General, "LIVEKIT_API_SECRET not set!"); + return; + } + }; + 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, + Err(e) => { + log_err!(PrintType::General, "Could not get rooms! {}", e); + return; + } + }; let mut call_ids: Vec = Vec::new(); let mut no_users: Vec = Vec::new(); for room in rooms { diff --git a/src/data/communication.rs b/src/data/communication.rs index e21dcae..4dd6d43 100644 --- a/src/data/communication.rs +++ b/src/data/communication.rs @@ -81,6 +81,8 @@ pub enum DataTypes { online_status, omikron_id, omikron_connections, + reset_token, + new_token, } impl DataTypes { @@ -161,6 +163,8 @@ impl DataTypes { "onlinestatus" => DataTypes::online_status, "omikronid" => DataTypes::omikron_id, "omikronconnections" => DataTypes::omikron_connections, + "resettoken" => DataTypes::reset_token, + "newtoken" => DataTypes::new_token, _ => DataTypes::error_type, // fallback if unknown } } @@ -170,6 +174,7 @@ impl DataTypes { #[allow(non_camel_case_types, dead_code)] pub enum CommunicationType { error, + error_invalid_data, error_invalid_user_id, error_invalid_omikron_id, error_not_found, @@ -244,6 +249,12 @@ pub enum CommunicationType { change_user_data, change_iota_data, + get_register, + complete_register_user, + complete_register_iota, + delete_user, + delete_iota, + start_register, complete_register, } @@ -262,6 +273,7 @@ impl CommunicationType { "function" => CommunicationType::function, "update" => CommunicationType::update, "createuser" => CommunicationType::create_user, + "errorinvaliddata" => CommunicationType::error_invalid_data, "errorinvaliduserid" => CommunicationType::error_invalid_user_id, "errorinvalidomikronid" => CommunicationType::error_invalid_omikron_id, "errornotfound" => CommunicationType::error_not_found, @@ -326,6 +338,12 @@ impl CommunicationType { "changeuserdata" => CommunicationType::change_user_data, "changeiotadata" => CommunicationType::change_iota_data, + "getregister" => CommunicationType::get_register, + "completeregisteruser" => CommunicationType::complete_register_user, + "completeregisteriota" => CommunicationType::complete_register_iota, + "deleteuser" => CommunicationType::delete_user, + "deleteiota" => CommunicationType::delete_iota, + "startregister" => CommunicationType::start_register, "completeregister" => CommunicationType::complete_register, diff --git a/src/omega/omega_connection.rs b/src/omega/omega_connection.rs index 164349c..49d73d3 100644 --- a/src/omega/omega_connection.rs +++ b/src/omega/omega_connection.rs @@ -29,7 +29,7 @@ use crate::{ }; use crate::{auth::crypto_helper::secret_key_to_base64, log_err}; -static WAITING_TASKS: Lazy< +pub static WAITING_TASKS: Lazy< DashMap, CommunicationValue) -> bool + Send + Sync>>, > = Lazy::new(DashMap::new); diff --git a/src/rho/client_connection.rs b/src/rho/client_connection.rs index eb840cf..ac1cb65 100644 --- a/src/rho/client_connection.rs +++ b/src/rho/client_connection.rs @@ -1,6 +1,7 @@ use async_tungstenite::tungstenite::Message; use async_tungstenite::{WebSocketReceiver, WebSocketSender}; use json::number::Number; +use rustls::sign::SingleCertAndKey; use std::sync::{Arc, Weak}; use tokio::sync::RwLock; use tokio_util::compat::Compat; @@ -9,6 +10,7 @@ use uuid::Uuid; use super::{rho_connection::RhoConnection, rho_manager}; 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, @@ -144,6 +146,25 @@ impl ClientConnection { return; } + if cv.is_type(CommunicationType::change_user_data) { + let client_for_closure = self.clone(); + WAITING_TASKS.insert( + cv.get_id(), + Box::new(move |_, response_cv| { + let client = client_for_closure.clone(); + tokio::spawn(async move { + client.send_message(&response_cv).await; + }); + true + }), + ); + + get_omega_connection() + .send_message(&cv.with_sender(*self.user_id.read().await)) + .await; + return; + } + // Forward other messages to Iota self.forward_to_iota(cv).await; });