Merge pull request #1 from kleo-dev/concept-arch
Working architecture concept
This commit is contained in:
Generated
+22
-2311
File diff suppressed because it is too large
Load Diff
+2
-3
@@ -5,13 +5,12 @@ edition = "2024"
|
|||||||
|
|
||||||
[dependencies]
|
[dependencies]
|
||||||
anyhow = "1.0.99"
|
anyhow = "1.0.99"
|
||||||
|
libloading = { version = "0.8.8", optional = true }
|
||||||
once_cell = "1.21.3"
|
once_cell = "1.21.3"
|
||||||
serde = { version = "1.0.219", features = ["serde_derive"] }
|
serde = { version = "1.0.219", features = ["serde_derive"] }
|
||||||
serde_json = "1.0.143"
|
serde_json = "1.0.143"
|
||||||
tungstenite = "0.27.0"
|
tungstenite = "0.27.0"
|
||||||
wasmtime = { version = "36.0.2", optional = true }
|
|
||||||
wasmtime-wasi = { version = "36.0.2", features = ["preview1"], optional = true }
|
|
||||||
|
|
||||||
[features]
|
[features]
|
||||||
default = ["loader"]
|
default = ["loader"]
|
||||||
loader = ["wasmtime", "wasmtime-wasi"]
|
loader = ["libloading"]
|
||||||
|
|||||||
+39
-24
@@ -1,6 +1,7 @@
|
|||||||
use std::{
|
use std::{
|
||||||
net::TcpListener,
|
net::{TcpListener, TcpStream},
|
||||||
path::{Path, PathBuf},
|
path::{Path, PathBuf},
|
||||||
|
sync::{Arc, Mutex},
|
||||||
};
|
};
|
||||||
|
|
||||||
#[cfg(feature = "loader")]
|
#[cfg(feature = "loader")]
|
||||||
@@ -13,7 +14,7 @@ pub mod vfs;
|
|||||||
pub use anyhow::Result;
|
pub use anyhow::Result;
|
||||||
use tungstenite::accept;
|
use tungstenite::accept;
|
||||||
|
|
||||||
use crate::plugin::{Plugin, PluginInstance};
|
use crate::plugin::DynPlugin;
|
||||||
pub use once_cell;
|
pub use once_cell;
|
||||||
|
|
||||||
#[derive(serde::Serialize, serde::Deserialize)]
|
#[derive(serde::Serialize, serde::Deserialize)]
|
||||||
@@ -22,7 +23,7 @@ pub struct ServerConfig {
|
|||||||
}
|
}
|
||||||
|
|
||||||
pub struct Server {
|
pub struct Server {
|
||||||
plugins: Vec<PluginInstance>,
|
plugins: Mutex<Vec<DynPlugin>>,
|
||||||
root: PathBuf,
|
root: PathBuf,
|
||||||
config: ServerConfig,
|
config: ServerConfig,
|
||||||
}
|
}
|
||||||
@@ -34,7 +35,7 @@ impl Default for ServerConfig {
|
|||||||
}
|
}
|
||||||
|
|
||||||
impl ServerConfig {
|
impl ServerConfig {
|
||||||
pub fn build(self, root: &Path) -> Server {
|
pub fn build(self, root: &Path) -> Arc<Server> {
|
||||||
Server::new_config(root, self)
|
Server::new_config(root, self)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
@@ -42,30 +43,33 @@ impl ServerConfig {
|
|||||||
impl Server {
|
impl Server {
|
||||||
logger!(LOGGER "Server");
|
logger!(LOGGER "Server");
|
||||||
|
|
||||||
pub fn new(root: &Path) -> Self {
|
pub fn new(root: &Path) -> Arc<Self> {
|
||||||
Self {
|
Arc::new(Self {
|
||||||
plugins: Vec::new(),
|
plugins: Mutex::new(Vec::new()),
|
||||||
root: root.to_path_buf(),
|
root: root.to_path_buf(),
|
||||||
config: ServerConfig::default(),
|
config: ServerConfig::default(),
|
||||||
}
|
})
|
||||||
}
|
}
|
||||||
|
|
||||||
pub fn new_config(root: &Path, config: ServerConfig) -> Self {
|
pub fn new_config(root: &Path, config: ServerConfig) -> Arc<Self> {
|
||||||
Self {
|
Arc::new(Self {
|
||||||
plugins: Vec::new(),
|
plugins: Mutex::new(Vec::new()),
|
||||||
root: root.to_path_buf(),
|
root: root.to_path_buf(),
|
||||||
config,
|
config,
|
||||||
}
|
})
|
||||||
}
|
}
|
||||||
|
|
||||||
pub fn run(&mut self) -> Result<()> {
|
pub fn run(self: &Arc<Self>) -> Result<()> {
|
||||||
// Load plugins
|
// Load plugins
|
||||||
#[cfg(feature = "loader")]
|
#[cfg(feature = "loader")]
|
||||||
loader::load_plugins(&mut self.plugins, &self.root.join("./plugins"))?;
|
loader::load_plugins(
|
||||||
|
&mut *self.plugins.lock().unwrap(),
|
||||||
|
&self.root.join("./plugins"),
|
||||||
|
)?;
|
||||||
|
|
||||||
// Initialize plugins
|
// Initialize plugins
|
||||||
for plugin in &mut self.plugins {
|
for plugin in self.plugins.lock().unwrap().iter_mut() {
|
||||||
plugin.init();
|
plugin.init(self);
|
||||||
}
|
}
|
||||||
|
|
||||||
// Start server
|
// Start server
|
||||||
@@ -75,14 +79,16 @@ impl Server {
|
|||||||
for stream in listener.incoming() {
|
for stream in listener.incoming() {
|
||||||
match stream {
|
match stream {
|
||||||
Ok(stream) => {
|
Ok(stream) => {
|
||||||
Self::LOGGER.info(format!("New connection: {}", stream.peer_addr()?));
|
std::thread::spawn({
|
||||||
let mut ws = accept(stream)?;
|
let srv = self.clone();
|
||||||
|
move || match srv.handle_client(stream) {
|
||||||
let msg = ws.read()?;
|
Ok(_) => {}
|
||||||
ws.send(msg)?;
|
Err(e) => Self::LOGGER.error(format!("Client handler failed: {e}")),
|
||||||
|
}
|
||||||
|
});
|
||||||
}
|
}
|
||||||
Err(e) => {
|
Err(e) => {
|
||||||
Self::LOGGER.error(format!("Connection failed: {}", e));
|
Self::LOGGER.error(format!("Connection failed: {e}"));
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
@@ -90,7 +96,16 @@ impl Server {
|
|||||||
Ok(())
|
Ok(())
|
||||||
}
|
}
|
||||||
|
|
||||||
pub fn add_plugin(&mut self, plugin: PluginInstance) {
|
pub fn add_plugin(self: &Arc<Self>, plugin: DynPlugin) {
|
||||||
self.plugins.push(plugin);
|
self.plugins.lock().unwrap().push(plugin);
|
||||||
|
}
|
||||||
|
|
||||||
|
fn handle_client(self: &Arc<Self>, stream: TcpStream) -> anyhow::Result<()> {
|
||||||
|
Self::LOGGER.info(format!("New connection: {}", stream.peer_addr()?));
|
||||||
|
let mut ws = accept(stream)?;
|
||||||
|
|
||||||
|
let msg = ws.read()?;
|
||||||
|
ws.send(msg)?;
|
||||||
|
Ok(())
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
+11
-22
@@ -1,41 +1,30 @@
|
|||||||
use std::path::Path;
|
use std::path::Path;
|
||||||
|
|
||||||
use wasmtime::{Engine, Linker, Module, Store};
|
use libloading::{Library, Symbol};
|
||||||
use wasmtime_wasi::WasiCtxBuilder;
|
|
||||||
|
|
||||||
use crate::{PluginInstance, logger, vfs};
|
use crate::{logger, plugin::DynPlugin, vfs};
|
||||||
|
|
||||||
logger! {
|
logger! {
|
||||||
const LOGGER "Loader"
|
const LOGGER "Loader"
|
||||||
}
|
}
|
||||||
|
|
||||||
pub fn load_plugin(wasm_path: &Path) -> anyhow::Result<PluginInstance> {
|
pub fn load_plugin(path: &Path) -> anyhow::Result<DynPlugin> {
|
||||||
let engine = Engine::default();
|
unsafe {
|
||||||
let mut store = Store::new(
|
LOGGER.info(format!("Loading plugin: {:?}", path));
|
||||||
&engine,
|
let lib = Library::new(path)?;
|
||||||
WasiCtxBuilder::new()
|
let func: Symbol<extern "C" fn() -> DynPlugin> = lib.get(b"load_plugin").unwrap();
|
||||||
.inherit_stdout()
|
Ok(func())
|
||||||
.inherit_stderr()
|
}
|
||||||
.build_p1(),
|
|
||||||
);
|
|
||||||
let module = Module::from_file(&engine, wasm_path)?;
|
|
||||||
let mut linker = Linker::new(&engine);
|
|
||||||
|
|
||||||
wasmtime_wasi::preview1::add_to_linker_sync(&mut linker, |s| s)?;
|
|
||||||
|
|
||||||
let instance = linker.instantiate(&mut store, &module)?;
|
|
||||||
Ok(PluginInstance::Wasm(instance, store))
|
|
||||||
}
|
}
|
||||||
|
|
||||||
pub fn load_plugins(arr: &mut Vec<PluginInstance>, path: &Path) -> crate::Result<()> {
|
pub fn load_plugins(arr: &mut Vec<DynPlugin>, path: &Path) -> crate::Result<()> {
|
||||||
LOGGER.info("Loading plugins");
|
LOGGER.info("Loading plugins");
|
||||||
vfs::dir(path)?;
|
vfs::dir(path)?;
|
||||||
if path.is_dir() {
|
if path.is_dir() {
|
||||||
for entry in std::fs::read_dir(path)? {
|
for entry in std::fs::read_dir(path)? {
|
||||||
let entry = entry?;
|
let entry = entry?;
|
||||||
let path = entry.path();
|
let path = entry.path();
|
||||||
if path.extension().and_then(|s| s.to_str()) == Some("wasm") {
|
if path.extension().and_then(|s| s.to_str()) == Some("dylib") {
|
||||||
LOGGER.info(format!("Loading plugin: {:?}", path));
|
|
||||||
match load_plugin(&path) {
|
match load_plugin(&path) {
|
||||||
Ok(plugin) => {
|
Ok(plugin) => {
|
||||||
arr.push(plugin);
|
arr.push(plugin);
|
||||||
|
|||||||
+3
-9
@@ -1,15 +1,9 @@
|
|||||||
#[macro_export]
|
#[macro_export]
|
||||||
macro_rules! export_plugin {
|
macro_rules! export_plugin {
|
||||||
($plugin_type:ty) => {
|
($p:expr) => {
|
||||||
static PLUGIN: std::sync::Mutex<Option<$plugin_type>> = std::sync::Mutex::new(None);
|
|
||||||
|
|
||||||
#[unsafe(no_mangle)]
|
#[unsafe(no_mangle)]
|
||||||
pub extern "C" fn init() {
|
pub extern "C" fn load_plugin() -> $crate::plugin::DynPlugin {
|
||||||
let mut plugin = PLUGIN.lock().unwrap();
|
$p
|
||||||
*plugin = Some(<$plugin_type>::default());
|
|
||||||
if let Some(plugin) = plugin.as_mut() {
|
|
||||||
plugin.init();
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
};
|
};
|
||||||
}
|
}
|
||||||
|
|||||||
+2
-3
@@ -3,9 +3,8 @@ use std::path::PathBuf;
|
|||||||
use voxa_server::{ServerConfig, vfs};
|
use voxa_server::{ServerConfig, vfs};
|
||||||
|
|
||||||
fn main() -> voxa_server::Result<()> {
|
fn main() -> voxa_server::Result<()> {
|
||||||
let root = PathBuf::from("./");
|
let root = PathBuf::from("");
|
||||||
let config: ServerConfig = vfs::read_config(&root.join("config.json"))?;
|
let config: ServerConfig = vfs::read_config(&root.join("config.json"))?;
|
||||||
let mut server = config.build(&root);
|
config.build(&root).run()?;
|
||||||
server.run()?;
|
|
||||||
Ok(())
|
Ok(())
|
||||||
}
|
}
|
||||||
|
|||||||
+6
-26
@@ -1,29 +1,9 @@
|
|||||||
#[cfg(feature = "loader")]
|
use std::sync::Arc;
|
||||||
use wasmtime::{Instance, Store};
|
|
||||||
#[cfg(feature = "loader")]
|
use crate::Server;
|
||||||
use wasmtime_wasi::preview1::WasiP1Ctx;
|
|
||||||
|
pub type DynPlugin = Box<dyn Plugin + Send + Sync>;
|
||||||
|
|
||||||
pub trait Plugin {
|
pub trait Plugin {
|
||||||
fn init(&mut self);
|
fn init(&mut self, server: &Arc<Server>);
|
||||||
}
|
|
||||||
|
|
||||||
pub enum PluginInstance {
|
|
||||||
#[cfg(feature = "loader")]
|
|
||||||
Wasm(Instance, Store<WasiP1Ctx>),
|
|
||||||
Static(Box<dyn Plugin>),
|
|
||||||
}
|
|
||||||
|
|
||||||
impl Plugin for PluginInstance {
|
|
||||||
fn init(&mut self) {
|
|
||||||
match self {
|
|
||||||
#[cfg(feature = "loader")]
|
|
||||||
PluginInstance::Wasm(instance, store) => {
|
|
||||||
let init = instance
|
|
||||||
.get_typed_func::<(), ()>(&mut *store, "init")
|
|
||||||
.unwrap();
|
|
||||||
init.call(store, ()).unwrap();
|
|
||||||
}
|
|
||||||
PluginInstance::Static(plugin) => plugin.init(),
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -1,4 +1,5 @@
|
|||||||
use voxa_server::{export_plugin, logger, plugin::Plugin};
|
use std::sync::Arc;
|
||||||
|
use voxa_server::{Server, export_plugin, logger, plugin::Plugin};
|
||||||
|
|
||||||
logger! {
|
logger! {
|
||||||
const LOGGER "My Plugin"
|
const LOGGER "My Plugin"
|
||||||
@@ -8,9 +9,9 @@ logger! {
|
|||||||
pub struct MyPlugin;
|
pub struct MyPlugin;
|
||||||
|
|
||||||
impl Plugin for MyPlugin {
|
impl Plugin for MyPlugin {
|
||||||
fn init(&mut self) {
|
fn init(&mut self, _server: &Arc<Server>) {
|
||||||
LOGGER.info("MyPlugin initialized!");
|
LOGGER.info("MyPlugin initialized!");
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
export_plugin!(MyPlugin);
|
export_plugin!(Box::new(MyPlugin));
|
||||||
|
|||||||
@@ -1,14 +1,12 @@
|
|||||||
use std::path::PathBuf;
|
use std::path::PathBuf;
|
||||||
|
|
||||||
use voxa_server::{ServerConfig, plugin::PluginInstance, vfs};
|
use voxa_server::{ServerConfig, vfs};
|
||||||
|
|
||||||
fn main() -> voxa_server::Result<()> {
|
fn main() -> voxa_server::Result<()> {
|
||||||
let root = PathBuf::from("./");
|
let root = PathBuf::from("./");
|
||||||
let config: ServerConfig = vfs::read_config(&root.join("config.json"))?;
|
let config: ServerConfig = vfs::read_config(&root.join("config.json"))?;
|
||||||
let mut server = config.build(&root);
|
let server = config.build(&root);
|
||||||
server.add_plugin(PluginInstance::Static(Box::new(
|
server.add_plugin(Box::new(test_plugin::MyPlugin::default()));
|
||||||
test_plugin::MyPlugin::default(),
|
|
||||||
)));
|
|
||||||
server.run()?;
|
server.run()?;
|
||||||
Ok(())
|
Ok(())
|
||||||
}
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user