From d8462efff612d216e3cb34eabb3aa7f4b516419b Mon Sep 17 00:00:00 2001 From: Leo dev Date: Mon, 15 Dec 2025 23:57:20 +0100 Subject: [PATCH] Plugin actions --- src/plugin/loader.rs | 16 +++++++++++- src/plugin/mod.rs | 56 ++++++++++++++++++++++++++++++++++------- src/plugin/types.rs | 7 +++++- src/requests/message.rs | 4 +-- src/server.rs | 16 +++++++----- 5 files changed, 79 insertions(+), 20 deletions(-) diff --git a/src/plugin/loader.rs b/src/plugin/loader.rs index a919852..d43a9da 100644 --- a/src/plugin/loader.rs +++ b/src/plugin/loader.rs @@ -55,7 +55,11 @@ impl PluginLoader { .unwrap() .try_clone() .unwrap(); - Plugin(a.try_clone().unwrap(), BufReader::new(a)) + Plugin { + stream: a.try_clone().unwrap(), + reader: BufReader::new(a), + id: plugin_json.id, + } } pub fn start_server(&self) { @@ -77,3 +81,13 @@ impl PluginLoader { }); } } + +impl Clone for Plugin { + fn clone(&self) -> Self { + Plugin { + stream: self.stream.try_clone().unwrap(), + reader: BufReader::new(self.stream.try_clone().unwrap()), + id: self.id.clone(), + } + } +} diff --git a/src/plugin/mod.rs b/src/plugin/mod.rs index 0047378..399e57c 100644 --- a/src/plugin/mod.rs +++ b/src/plugin/mod.rs @@ -4,22 +4,60 @@ pub mod types; use std::{ io::{BufRead, BufReader, Write}, net::TcpStream, + sync::Arc, }; -use crate::plugin::types::{LoaderMessage, PluginMessage}; +use crate::{ + plugin::types::{LoaderMessage, PluginMessage}, + server::Server, + types::message::ServerMessage, +}; -pub struct Plugin(TcpStream, BufReader); +pub struct Plugin { + stream: TcpStream, + reader: BufReader, + id: String, +} impl Plugin { - pub fn send(&mut self, m: &LoaderMessage) { - self.0 - .write(serde_json::to_string(&m).unwrap().as_bytes()) - .unwrap(); + pub fn send(&mut self, m: &LoaderMessage) -> crate::Result<()> { + self.stream + .write(serde_json::to_string(&m).unwrap().as_bytes())?; + Ok(()) } - pub fn read(&mut self) -> PluginMessage { + pub fn read(&mut self) -> crate::Result { let mut buf = String::new(); - self.1.read_line(&mut buf).unwrap(); - serde_json::from_str(&buf).unwrap() + self.reader.read_line(&mut buf)?; + Ok(serde_json::from_str(&buf)?) + } + + pub fn run(&mut self, server: &Arc) -> crate::Result<()> { + loop { + match self.read()? { + PluginMessage::SendMessage { + channel_id, + contents, + } => { + let msg = server.db.insert_message( + &channel_id, + &self.id, + &contents, + chrono::Utc::now().timestamp(), + )?; + + for c in server.clients.lock().unwrap().iter() { + let c = c.clone(); + let server = server.clone(); + let msg = msg.clone(); + std::thread::spawn(move || { + server + .wrap_err(&c, c.send(ServerMessage::MessageCreate(msg))) + .expect("Failed to broadcast"); + }); + } + } + } + } } } diff --git a/src/plugin/types.rs b/src/plugin/types.rs index 1d094ff..169ce4b 100644 --- a/src/plugin/types.rs +++ b/src/plugin/types.rs @@ -35,4 +35,9 @@ pub enum LoaderMessage { #[derive(Debug, Clone, Serialize, Deserialize)] #[serde(tag = "type", content = "params", rename_all = "snake_case")] -pub enum PluginMessage {} +pub enum PluginMessage { + SendMessage { + channel_id: String, + contents: String, + }, +} diff --git a/src/requests/message.rs b/src/requests/message.rs index f58941b..929b547 100644 --- a/src/requests/message.rs +++ b/src/requests/message.rs @@ -1,7 +1,5 @@ use std::sync::Arc; -use serde_json::ser; - use crate::{plugin::types::LoaderMessage, server::Server, types, utils::client::Client}; crate::logger!(LOGGER "Message Manager"); @@ -48,7 +46,7 @@ pub fn send( server.send_plugin_message(&LoaderMessage::MessageSent { user_id: client.get_uuid().unwrap_or_default(), msg: msg, - }); + })?; Ok(()) } diff --git a/src/server.rs b/src/server.rs index 1580dad..0385b5a 100644 --- a/src/server.rs +++ b/src/server.rs @@ -79,7 +79,10 @@ impl Server { let path = entry.path(); if path.extension().and_then(|s| s.to_str()) == Some("json") { - self.plugins.lock().unwrap().push(plugin_loader.load(&path)); + let mut p = plugin_loader.load(&path); + let s = self.clone(); + self.plugins.lock().unwrap().push(p.clone()); + std::thread::spawn(move || p.run(&s)); } } Self::LOGGER.info("Plugins loaded"); @@ -120,7 +123,7 @@ impl Server { Ok(()) } - fn init_client(self: &Arc, stream: TcpStream) -> anyhow::Result { + fn init_client(self: &Arc, stream: TcpStream) -> crate::Result { Self::LOGGER.info(format!("New connection: {}", stream.peer_addr()?)); // Initialize client let mut client = Client::new(stream)?; @@ -173,7 +176,7 @@ impl Server { Ok(client) } - fn handle_client(self: &Arc, client: &Client) -> anyhow::Result<()> { + fn handle_client(self: &Arc, client: &Client) -> crate::Result<()> { // The main req/res loop loop { let req = client.read()?; @@ -181,7 +184,7 @@ impl Server { self.send_plugin_message(&LoaderMessage::Request { user_id: client.get_uuid().unwrap_or_default(), msg: r.clone(), - }); + })?; self.wrap_err(&client, self.call_request(r, &client))?; } } @@ -204,9 +207,10 @@ impl Server { res } - pub fn send_plugin_message(&self, msg: &LoaderMessage) { + pub fn send_plugin_message(&self, msg: &LoaderMessage) -> crate::Result<()> { for p in self.plugins.lock().unwrap().iter_mut() { - p.send(msg); + p.send(msg)?; } + Ok(()) } }