Plugin actions
This commit is contained in:
+15
-1
@@ -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(),
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
+47
-9
@@ -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<TcpStream>);
|
||||
pub struct Plugin {
|
||||
stream: TcpStream,
|
||||
reader: BufReader<TcpStream>,
|
||||
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<PluginMessage> {
|
||||
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<Server>) -> 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");
|
||||
});
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
+6
-1
@@ -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,
|
||||
},
|
||||
}
|
||||
|
||||
@@ -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(())
|
||||
}
|
||||
|
||||
+10
-6
@@ -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<Self>, stream: TcpStream) -> anyhow::Result<Client> {
|
||||
fn init_client(self: &Arc<Self>, stream: TcpStream) -> crate::Result<Client> {
|
||||
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<Self>, client: &Client) -> anyhow::Result<()> {
|
||||
fn handle_client(self: &Arc<Self>, 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(())
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user