diff --git a/src/arbiter/models.rs b/src/arbiter/models.rs index 15a3387..c156221 100644 --- a/src/arbiter/models.rs +++ b/src/arbiter/models.rs @@ -18,6 +18,8 @@ pub struct User { pub sessions: Arc>>, pub user_type: String, pub pretty_name: String, + pub parent_id: Option, + pub children: Vec, } impl User { @@ -28,6 +30,8 @@ impl User { sessions: self.sessions, user_type: self.user_type, pretty_name: self.pretty_name, + parent_id: self.parent_id, + children: self.children } } } @@ -39,6 +43,8 @@ pub struct UserWithId { pub sessions: Arc>>, pub user_type: String, pub pretty_name: String, + pub parent_id: Option, + pub children: Vec, } impl Into for UserWithId { @@ -48,6 +54,8 @@ impl Into for UserWithId { sessions: self.sessions, user_type: self.user_type, pretty_name: self.pretty_name, + parent_id: self.parent_id, + children: self.children } } } @@ -62,12 +70,12 @@ pub struct ApiKey { } impl ApiKey { - pub fn to_api_key_with_key(self, key: String) -> ApiKeyWithKey { + pub fn to_api_key_with_key(&self, key: &String) -> ApiKeyWithKey { ApiKeyWithKey { - key, - allowed_events_to: self.allowed_events_to, - allowed_events_from: self.allowed_events_from, - user_id: self.user_id, + key: key.clone(), + allowed_events_to: self.allowed_events_to.clone(), + allowed_events_from: self.allowed_events_from.clone(), + user_id: self.user_id.clone(), echo: self.echo, proxy: self.proxy } diff --git a/src/client/mod.rs b/src/client/mod.rs index a47d472..75a9714 100644 --- a/src/client/mod.rs +++ b/src/client/mod.rs @@ -1,7 +1,4 @@ -use crate::server::models::{ - IPCMessageWithId, - UserConfig, -}; +use crate::{arbiter::models::ApiKey, server::{models::IPCMessageWithId, websockets::WsIn, AUTH_HEADER}, user::NexusUser}; use fastwebsockets::{ FragmentCollector, Frame, @@ -21,22 +18,38 @@ use log::info; use std::future::Future; use std::sync::Arc; use std::time::Duration; -use tokio::net::TcpStream; +use std::collections::HashMap; +use tokio::{net::TcpStream, sync::broadcast}; use tokio::{ - sync::{ - Mutex, - mpsc::{ - UnboundedReceiver, - UnboundedSender, - unbounded_channel, - }, - }, + sync::Mutex, task::{ JoinHandle, spawn, }, }; use tokio_util::sync::CancellationToken; +use log::error; + +#[derive(Debug, Clone)] +pub struct ClientStatus { + pub connected: bool +} + +impl ClientStatus { + pub fn new(connected: bool) -> Self { + ClientStatus { + connected + } + } + + pub fn set(&mut self, connected: bool) { + self.connected = connected + } + + pub fn get(&self) -> bool { + self.connected + } +} /// Use when connecting a program to nexus *remotely*. Use the nexus::server::listener's /// nexus_listener() function to run an embedded server. @@ -45,11 +58,14 @@ pub struct NexusClient { url: String, // TODO: Make optional as we can connect via OS socket as well. port: u16, - config: UserConfig, - connected: bool, + api_key: String, + status: Arc>, keep_trying: bool, - to_server: Option>, - from_server: Option>, + to_server_tx: broadcast::Sender, + from_server: broadcast::Sender, + handles: Vec>, + api_keys: Arc>>, + cancellation_token: CancellationToken // TODO: Add user registry to see if a user is connected via this client to route back instead of sending to server. } @@ -75,18 +91,24 @@ impl NexusClient { secure: bool, url: &String, port: &u16, - config: &UserConfig, + api_key: &String, keep_trying: bool, ) -> Self { + let (from_server, _) = broadcast::channel::(usize::MAX / 2); + let (to_server_tx, _) = broadcast::channel::(usize::MAX / 2); + NexusClient { secure, url: url.clone(), port: port.clone(), - config: config.clone(), - connected: false, + api_key: api_key.clone(), + status: Arc::new(Mutex::new(ClientStatus::new(false))), keep_trying, - to_server: None, - from_server: None + to_server_tx, + from_server, + handles: Vec::new(), + api_keys: Arc::new(Mutex::new(HashMap::new())), + cancellation_token: CancellationToken::new() } } @@ -96,35 +118,33 @@ impl NexusClient { pub async fn connect_to_ws( &mut self, ) -> Result< - ( - JoinHandle<()>, - (CancellationToken, CancellationToken) - ), + JoinHandle<()>, anyhow::Error, > { let mut keep_trying = true; - let cancellation_tokens = (CancellationToken::new(), CancellationToken::new()); let mut ws_opt = None; let mut error: Option = None; + let host = format!( + "{}:{}", + self.url.clone(), + self.port.clone().to_string() + ); + let uri = format!( + "{}://{}:{}/ws", + (if self.secure { "https" } else { "http" }), + self.url.clone(), + self.port.clone().to_string() + ); while keep_trying { - match TcpStream::connect(format!( - "{}:{}", - self.url.clone(), - self.port.clone().to_string() - )).await { + match TcpStream::connect(host.clone()).await { Ok(stream) => { match Request::builder() .method("GET") - .uri(format!( - "{}://{}:{}/", - (if self.secure { "https" } else { "http" }), - self.url.clone(), - self.port.clone().to_string() - )) + .uri(uri.clone()) .header( "Host", - format!("{}:{}", self.url.clone(), self.port.clone().to_string()), + host.clone(), ) .header(UPGRADE, "websocket") .header(CONNECTION, "upgrade") @@ -132,7 +152,7 @@ impl NexusClient { "Sec-WebSocket-Key", fastwebsockets::handshake::generate_key(), ) - // TODO: Add API key! + .header(AUTH_HEADER, self.api_key.clone()) .header("Sec-WebSocket-Version", "13") .body(Empty::::new()) { @@ -141,7 +161,9 @@ impl NexusClient { ws_opt = Some(the_socket); } Err(e) => { + error!("Failed to perform WebSocket Handshake with \"{}\":\n{}", uri.clone(), e); if self.keep_trying { + error!("Retrying WebSocket Handshake with \"{}\"...", uri.clone()); tokio::time::sleep(Duration::from_secs(1)).await; } else { error = Some(e.into()); @@ -150,7 +172,9 @@ impl NexusClient { } }, Err(e) => { + error!("Failed to send WebSocket Upgrade request to \"{}\":\n{}", uri.clone(), e); if self.keep_trying { + info!("Retrying WebSocket Upgrade request for \"{}\"...", uri.clone()); tokio::time::sleep(Duration::from_secs(1)).await; } else { error = Some(e.into()); @@ -160,7 +184,9 @@ impl NexusClient { } } Err(e) => { + error!("Failed to open TCP connection to \"{}\":\n{}", host.clone(), e); if self.keep_trying { + info!("Retrying TCP connection to \"{}\"...", host.clone()); tokio::time::sleep(Duration::from_secs(1)).await; } else { error = Some(e.into()); @@ -172,12 +198,11 @@ impl NexusClient { match ws_opt { Some(ws) => { - self.connected = true; + self.status.lock().await.set(true); let (mut reader, og_writer) = ws.split(tokio::io::split); let writer = Arc::new(Mutex::new(AddSend(og_writer))); - let (handle_from_tx, handle_from_rx) = unbounded_channel::(); - let (handle_to_tx, mut handle_to_rx) = unbounded_channel::(); + let (from_tx, _) = broadcast::channel::(usize::MAX / 2); let senders_writer = writer.clone(); let mut sender = move |frame| { @@ -185,12 +210,13 @@ impl NexusClient { async move { senders_writer_two.lock().await.0.write_frame(frame).await } }; - let handle_tokens = cancellation_tokens.clone(); + let handle_from_tx = from_tx.clone(); + let handle_token = self.cancellation_token.clone(); + let mut handle_to_rx = self.to_server_tx.subscribe(); let handle = spawn(async move { - let recv_tokens = Arc::new(handle_tokens.clone()); + let recv_token = Arc::new(handle_token.clone()); let recv_handle = spawn(async move { - recv_tokens - .0 + recv_token .run_until_cancelled(async move { // TODO: Use FragmentCollector! while let Ok(mut message) = reader.read_frame(&mut sender).await { @@ -221,13 +247,12 @@ impl NexusClient { .await; }); - let send_tokens = Arc::new(handle_tokens.clone()); + let send_token = Arc::new(handle_token.clone()); let send_handle = spawn(async move { let sender_writer = writer.clone(); - send_tokens - .0 + send_token .run_until_cancelled(async move { - while let Some(message) = handle_to_rx.recv().await { + while let Ok(message) = handle_to_rx.recv().await { match sender_writer .lock() .await @@ -237,7 +262,9 @@ impl NexusClient { )) .await { - Ok(_) => {} + Ok(_) => { + // TODO: + } Err(e) => { // TODO: } @@ -269,18 +296,43 @@ impl NexusClient { send_handle ]) => { info!("WS threads have exited!"); - handle_tokens.1.cancel(); }} }); - self.to_server = Some(handle_to_tx); - self.from_server = Some(handle_from_rx); - - Ok((handle, cancellation_tokens)) + Ok(handle) } None => { return Err(error.unwrap()); } } } + + pub async fn new_user(&self, api_key_str: &String) -> Result { + if self.status.lock().await.get() { + match self.api_keys.lock().await.get(&api_key_str.clone()) { + Some(api_key) => { + let status = self.status.clone(); + let cancellation_token = self.cancellation_token.clone(); + let to_server = self.to_server_tx.clone(); + let from_server_tx = self.from_server.clone(); + + Ok(NexusUser::new( + false, + status, + cancellation_token, + api_key.to_api_key_with_key(&api_key_str.clone()), + to_server, + from_server_tx + )) + }, + None => { + Err(anyhow::anyhow!("Client's API key does not exist in the mini-store... sure ya have the right one?")) + } + } + } else { + Err(anyhow::anyhow!("Client is not connected yet!")) + } + } + + // TODO: Get API keys that this proxy user registered. } diff --git a/src/server/listener.rs b/src/server/listener.rs index 831c3e2..c7c4ec7 100644 --- a/src/server/listener.rs +++ b/src/server/listener.rs @@ -5,7 +5,6 @@ use crate::arbiter::models::{ use crate::server::models::{ Client, ClientWithId, - CoreUserConfig, IPCMessageWithId, NexusStore, Session, @@ -13,9 +12,7 @@ use crate::server::models::{ use crate::server::websockets::handle_ws_client; use crate::server::{AUTH_HEADER, DEAUTH_EVENT}; use crate::utils::{ - gen_cid_with_check, - gen_ipc_message, - iso8601, + gen_cid_with_check, gen_message_id_with_check, iso8601 }; use log::{ debug, @@ -50,7 +47,6 @@ use warp::{ Filter, http::StatusCode, }; - use super::models::UserConfig; // example error response @@ -59,6 +55,13 @@ struct ApiErrorResult { detail: String, } +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct IPCWithKey { + kind: String, + message: String, + api_key: String +} + // errors thrown by handlers and custom filters, // such as `ensure_authentication` filter #[derive(Error, Debug)] @@ -163,7 +166,7 @@ async fn ensure_authentication( .get(&api_key_str.clone()) .unwrap() .clone() - .to_api_key_with_key(api_key_str.clone()); + .to_api_key_with_key(&api_key_str.clone()); let user_id = api_key.clone().user_id; debug!("Registering as client: {}", client_id.clone()); @@ -261,80 +264,14 @@ pub struct ServerHealth { up_since: String, } -async fn handle_ipc_send( - sender: Arc>, - msg: IPCMessageWithId, - user_config: &Arc, - store: &NexusStore, -) { - let users_mutex = &store.users.to_owned(); - let users = users_mutex.lock().await; - let user_conf = users.get(&user_config.id.clone()).expect(&format!( - "ERROR: Core user not found: {}", - user_config.id.clone() - )); - let keys_mutex = &store.api_keys.to_owned(); - let keys = keys_mutex.lock().await; - - let api_key_conf = keys.get(&user_config.api_keys[0].key.clone()).expect(&format!( - "ERROR: Core user api_key not found: {}", - user_config.api_keys[0].key.clone() - )); - let mut event_sent = false; - - for allowed_event_regex in api_key_conf.allowed_events_from.clone() { - match Regex::new(&allowed_event_regex.clone()) { - Ok(regex) => { - if regex.is_match(&msg.kind.clone()) { - match sender.send(msg.clone()) { - Ok(_) => { - event_sent = true; - } - Err(e) => { - error!( - "Core user: {}, IPC channel: {}, Failed to send message: {{ \"author\": \"{}\", \"kind\": \"{}\", \"message\": \"{}\" }}, due to:\n{}", - user_config.id.clone(), - user_conf.user_type.clone(), - msg.author.clone(), - msg.kind.clone(), - msg.message.clone(), - e - ); - } - }; - } - } - Err(e) => { - error!( - "Core user: {}, api key's \"allowed events from\", regex: {}, is invalid! Regex Error: {}", - user_config.id.clone(), - allowed_event_regex.clone(), - e - ); - } - } - } - - if !event_sent { - debug!( - "Core user: {}, event \"{}\" not sent.", - user_config.id, - msg.kind.clone() - ); - } -} - pub async fn nexus_listener( port: u16, store: Arc, - mut internal_client_senders: Vec<( + internal_client_senders: Vec<( &UserConfig, UnboundedSender, )>, - internal_client_recivers: Vec<( - &UserConfig, - UnboundedReceiver, - )>, + internal_client_recivers: Vec>, cancellation_tokens: (CancellationToken, CancellationToken), ) { info!("Starting nexus on port: {}...", port); @@ -424,7 +361,6 @@ pub async fn nexus_listener( let ipc_dispatch_store = Arc::new(store.clone()); let ipc_dispatch_clients_tx = Arc::new(clients_tx.clone()); - let ipc_dispatch_user_config = Arc::new(nexus_user_config.clone()); let ipc_dispatch_token = cancellation_tokens.0.clone(); let ipc_dispatch_handle = tokio::task::spawn(async move { tokio::select! { @@ -452,7 +388,7 @@ pub async fn nexus_listener( match Regex::new(&allowed_event_regex) { Ok(regex) => { // TODO: Generate an internal CID in CoreUserConfig! - if regex.is_match(&allowed_event_regex) { // && !(message.author.clone().split("?client=").collect::>()[1] == *client_id.clone()) { + if regex.is_match(&message.kind.clone()) { // && !(message.author.clone().split("?client=").collect::>()[1] == *client_id.clone()) { debug!("Sending event: \"{}\", to client: {}...", message.kind.clone(), client_id.clone()); match client_sender.send(message.clone()) { Ok(_) => { @@ -464,6 +400,8 @@ pub async fn nexus_listener( }; break; + } else { + // TODO: Not permitted to send message } }, Err(e) => { @@ -499,12 +437,14 @@ pub async fn nexus_listener( .finish() .to_string(); - let generated_message = gen_ipc_message( - &ipc_dispatch_store.clone(), - &ipc_dispatch_user_config.clone(), + let generated_message = IPCMessageWithId { + // TODO: fix this!!!!! + author: "ipc://com.reboot-codes.nexus.listener".to_string(), kind, - "api key removed from store".to_string() - ).await; + message: "api key removed from store".to_string(), + id: gen_message_id_with_check(&ipc_dispatch_store).await + }; + ipc_dispatch_store.messages.lock().await.insert(generated_message.id.clone(), generated_message.clone().into()); let _ = client_sender.send(generated_message.clone()); @@ -531,20 +471,67 @@ pub async fn nexus_listener( let mut internal_ipc_handles = vec![]; - for client in internal_client_recivers { + for mut client in internal_client_recivers { // Internal IPC Handles - let client_cfg = Arc::new(client.0.clone()); let internal_ipc_store = Arc::new(store.clone()); let internal_ipc_tx = Arc::new(from_client_tx.clone()); let internal_ipc_token = cancellation_tokens.0.clone(); internal_ipc_handles.push(tokio::task::spawn(async move { tokio::select! { _ = internal_ipc_token.cancelled() => { - debug!("Internal IPC handle for client: \"{}\", exited!", client_cfg.id.clone()); + debug!("Internal IPC Handle Exited!"); }, _ = async move { - while let Some(msg) = client.1.recv().await { - handle_ipc_send(internal_ipc_tx.clone(), msg, &client_cfg.clone(), &internal_ipc_store.clone()).await; + while let Some(msg) = client.recv().await { + match internal_ipc_store.api_keys.lock().await.get(&msg.api_key.to_string()) { + Some(api_key) => { + match internal_ipc_store.users.lock().await.get(&api_key.user_id) { + Some(_user) => { + for allowed_event_regex in api_key.allowed_events_from.clone() { + match Regex::new(&allowed_event_regex.clone()) { + Ok(regex) => { + if regex.is_match(&msg.kind.clone()) { + let author = format!("ipc://{}", api_key.user_id.clone()); + match internal_ipc_tx.send(IPCMessageWithId { + author: author.clone(), + kind: msg.kind.clone(), + message: msg.message.clone(), + id: gen_message_id_with_check(&internal_ipc_store).await + }) { + Ok(_) => {} + Err(e) => { + error!( + "IPC user: {}, Failed to send message: {{ \"author\": \"{}\", \"kind\": \"{}\", \"message\": \"{}\" }}, due to:\n{}", + api_key.user_id.clone(), + author.clone(), + msg.kind.clone(), + msg.message.clone(), + e + ); + } + }; + } + } + Err(e) => { + error!( + "IPC user: {}, api key's \"allowed events from\", regex: {}, is invalid! Regex Error: {}", + api_key.user_id.clone(), + allowed_event_regex.clone(), + e + ); + } + } + } + }, + None => { + // TODO: + } + } + }, + None => { + // TODO: + } + } } } => {} } diff --git a/src/server/models.rs b/src/server/models.rs index a57c3b4..9097f1e 100644 --- a/src/server/models.rs +++ b/src/server/models.rs @@ -2,7 +2,7 @@ use crate::{ arbiter::models::{ ApiKey, ApiKeyWithKeyWithoutUID, - User, + User, UserWithId, }, utils::{ gen_api_key_with_check, @@ -86,13 +86,20 @@ pub struct CoreUserConfig { } #[derive(Debug, Clone, Serialize, Deserialize)] -pub struct UserConfig { +pub struct UserConfigWithId { pub user_type: String, pub pretty_name: String, pub id: String, pub api_keys: Vec, } +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct UserConfig { + pub user_type: String, + pub pretty_name: String, + pub api_keys: Vec, +} + // TODO: Add serialization/deserialization functions... // TODO: Add options for making certain models ephemeral or persistent. #[derive(Debug, Clone)] @@ -107,7 +114,7 @@ impl NexusStore { const MASTER_USER_TYPE: &str = "com.reboot-codes.nexus.master"; /// Create a new store with a set master user. - pub async fn new(master_user_pretty_name: &String) -> (NexusStore, UserConfig) { + pub async fn new(master_user_pretty_name: &String) -> (NexusStore, UserWithId) { let mut ret = NexusStore { users: Arc::new(Mutex::new(HashMap::new())), api_keys: Arc::new(Mutex::new(HashMap::new())), @@ -115,43 +122,78 @@ impl NexusStore { messages: Arc::new(Mutex::new(HashMap::new())), }; - let master_user_config = ret.add_master_user(&master_user_pretty_name).await; + let master_user = ret.add_master_user(&master_user_pretty_name).await; - (ret, master_user_config) + (ret, master_user) } - pub async fn add_user(&mut self, user_config: UserConfig) { - let mut key_ids: Vec = vec![]; - for key_config in user_config.api_keys.iter() { - key_ids.push(key_config.key.clone()); + pub async fn add_user(&mut self, user_config: UserConfig, parent: Option) -> Result { + let mut parent_id = None; + let mut error = None; + + match parent.clone() { + Some(target_parent_id) => { + match self.users.lock().await.get(&target_parent_id) { + Some(_parent) => { + parent_id = parent.clone(); + }, + None => { + error = Some(anyhow::anyhow!("Parent ID does not exist in store!")); + } + } + }, + None => {} } - self.users.lock().await.insert( - user_config.id.clone(), - User { - pretty_name: user_config.pretty_name, - user_type: user_config.user_type, - api_keys: key_ids, - sessions: Arc::new(Mutex::new(HashMap::new())), - }, - ); - - for key_config in user_config.api_keys.iter() { - self.api_keys.lock().await.insert( - key_config.key.clone(), - ApiKey { - allowed_events_to: key_config.allowed_events_to.clone(), - allowed_events_from: key_config.allowed_events_from.clone(), - user_id: user_config.id.clone(), - echo: key_config.echo.clone(), - }, - ); + match error { + Some(e) => { Err(e) }, + None => { + let mut key_ids: Vec = vec![]; + for key_config in user_config.api_keys.iter() { + key_ids.push(key_config.key.clone()); + } + let id = gen_uid_with_check(self).await; + + self.users.lock().await.insert( + id.clone(), + User { + pretty_name: user_config.pretty_name.clone(), + user_type: user_config.user_type.clone(), + api_keys: key_ids.clone(), + sessions: Arc::new(Mutex::new(HashMap::new())), + parent_id: parent.clone(), + children: Vec::new() + }, + ); + + for key_config in user_config.api_keys.iter() { + self.api_keys.lock().await.insert( + key_config.key.clone(), + ApiKey { + allowed_events_to: key_config.allowed_events_to.clone(), + allowed_events_from: key_config.allowed_events_from.clone(), + user_id: id.clone(), + echo: key_config.echo, + proxy: key_config.proxy + }, + ); + } + + Ok(UserWithId { + pretty_name: user_config.pretty_name, + user_type: user_config.user_type, + api_keys: key_ids, + sessions: Arc::new(Mutex::new(HashMap::new())), + parent_id: parent, + children: Vec::new(), + id: id.clone() + }) + } } } - pub async fn add_master_user(&mut self, pretty_name: &String) -> UserConfig { + pub async fn add_master_user(&mut self, pretty_name: &String) -> UserWithId { let ret = UserConfig { - id: gen_uid_with_check(self).await, pretty_name: pretty_name.clone(), user_type: NexusStore::MASTER_USER_TYPE.to_string(), api_keys: vec![ApiKeyWithKeyWithoutUID { @@ -159,11 +201,10 @@ impl NexusStore { allowed_events_from: vec![".*".to_string()], key: gen_api_key_with_check(self).await, echo: true, + proxy: true }], }; - self.add_user(ret.clone()).await; - - ret + self.add_user(ret.clone(), None).await.unwrap() } } diff --git a/src/user/mod.rs b/src/user/mod.rs index 113e6c1..49c8167 100644 --- a/src/user/mod.rs +++ b/src/user/mod.rs @@ -1,6 +1,92 @@ +use log::error; +use regex::Regex; +use tokio::{sync::{broadcast, mpsc::{unbounded_channel, UnboundedReceiver, UnboundedSender}, Mutex}, task::JoinHandle}; +use tokio_util::sync::CancellationToken; +use std::sync::Arc; +use crate::{arbiter::models::ApiKeyWithKey, client::ClientStatus, server::{models::IPCMessageWithId, websockets::WsIn}}; + /// Use in a thread to connect to send messages to/recieve messages from nexus with. /// Requires a [NexusClient](nexus::client::NexusClient) or a tokio IPC channel /// connected to an embedded Nexus server to actually send messages. /// /// This struct basically just ensures that messages are formatted properly. -pub struct NexusUser {} +#[derive(Debug, Clone)] +pub struct NexusUser { + is_child: bool, + connected: Arc>, + cancellation_token: CancellationToken, + api_key: ApiKeyWithKey, + to_client: broadcast::Sender, + from_server_tx: broadcast::Sender +} + +impl NexusUser { + pub fn new( + is_child: bool, + connected: Arc>, + cancellation_token: CancellationToken, + api_key: ApiKeyWithKey, + to_client: broadcast::Sender, + from_server_tx: broadcast::Sender + ) -> Self { + NexusUser { + is_child, + connected, + cancellation_token, + api_key, + to_client, + from_server_tx + } + } + + pub fn send_msg(&self, kind: &String, message: &String) -> Result<(), anyhow::Error> { + match self.to_client.send(WsIn { + kind: kind.clone(), + message: message.clone(), + api_key: Some(self.api_key.key.clone()) + }) { + Ok(_) => { + Ok(()) + }, + Err(e) => { + Err(e.into()) + } + } + } + + pub fn subscribe(&self) -> (UnboundedReceiver, JoinHandle<()>) { + let (tx, rx) = unbounded_channel::(); + let mut from_client = self.from_server_tx.subscribe(); + let cancellation_token = self.cancellation_token.clone(); + let this = self.clone(); + + ( + rx, + tokio::task::spawn(async move { + cancellation_token.run_until_cancelled(async move { + while let Ok(message) = from_client.recv().await { + for allowed_event_regex in &this.api_key.allowed_events_to.clone() { + match Regex::new(&allowed_event_regex) { + Ok(regex) => { + if regex.is_match(&message.kind.clone()) { + match tx.send(message.clone()) { + Ok(_) => {}, + Err(e) => { + error!("Failed to send message to user: {}, due to:\n{}", this.api_key.user_id.clone(), e); + } + }; + + break; + } + }, + Err(e) => { + error!("Allowed event regular expression: \"{}\", for user id: {}, errored with: {}", allowed_event_regex.clone(), this.api_key.user_id.clone(), e); + } + } + } + } + }).await; + }) + ) + } +} diff --git a/src/utils.rs b/src/utils.rs index d8fc81a..0813bb3 100644 --- a/src/utils.rs +++ b/src/utils.rs @@ -99,21 +99,6 @@ pub async fn gen_message_id_with_check(store: &NexusStore) -> String { } } -pub async fn gen_ipc_message( - store: &NexusStore, - user_config: &CoreUserConfig, - kind: String, - message: String, -) -> IPCMessageWithId { - let message_id = gen_message_id_with_check(&store.clone()).await; - IPCMessageWithId { - id: message_id.clone(), - author: user_config.id.to_string(), - kind, - message, - } -} - pub async fn gen_cid_with_check(store: &NexusStore) -> String { loop { let client_id = Uuid::new_v4().to_string();