persistence + proper disconnect

This commit is contained in:
2025-11-28 17:29:33 -05:00
parent 029d62cb53
commit 7557507f27
16 changed files with 413 additions and 67 deletions

View File

@@ -1,10 +1,15 @@
// mod data;
mod db;
mod net;
use crate::db::{DbUtil, open_db};
use clap::Parser;
use net::{ClientSender, ConAccepter, listen};
use openworm::{
net::{ClientMsg, DisconnectReason, Msg, RecvHandler, ServerMsg, install_crypto_provider},
net::{ClientMsg, DisconnectReason, RecvHandler, ServerMsg, install_crypto_provider},
rsc::DataDir,
};
use sled::{Db, Tree};
use std::{
collections::HashMap,
sync::{
@@ -14,27 +19,39 @@ use std::{
};
use tokio::sync::RwLock;
#[derive(Parser, Debug)]
#[command(version, about, long_about = None)]
struct Args {
/// port to listen on
#[arg(short, long)]
port: u16,
}
fn main() {
let args = Args::parse();
install_crypto_provider();
run_server();
run_server(args.port);
}
#[tokio::main]
pub async fn run_server() {
pub async fn run_server(port: u16) {
let dir = DataDir::default();
let path = dir.get();
let db: Db = open_db(path.join("server.db"));
let handler = ServerListener {
msgs: Default::default(),
msgs: db.open_tree("msgs").unwrap(),
senders: Default::default(),
count: 0.into(),
db,
};
listen(path, handler).await;
listen(port, path, handler).await;
}
type ClientId = u64;
struct ServerListener {
msgs: Arc<RwLock<Vec<Msg>>>,
db: Db,
msgs: Tree,
senders: Arc<RwLock<HashMap<ClientId, ClientSender>>>,
count: AtomicU64,
}
@@ -44,6 +61,7 @@ impl ConAccepter for ServerListener {
let id = self.count.fetch_add(1, Ordering::Release);
self.senders.write().await.insert(id, send.clone());
ClientHandler {
db: self.db.clone(),
msgs: self.msgs.clone(),
senders: self.senders.clone(),
send,
@@ -53,17 +71,22 @@ impl ConAccepter for ServerListener {
}
struct ClientHandler {
msgs: Arc<RwLock<Vec<Msg>>>,
db: Db,
msgs: Tree,
send: ClientSender,
senders: Arc<RwLock<HashMap<ClientId, ClientSender>>>,
id: ClientId,
}
impl RecvHandler<ClientMsg> for ClientHandler {
async fn connect(&self) -> () {
println!("connected: {:?}", self.send.remote());
}
async fn msg(&self, msg: ClientMsg) {
match msg {
ClientMsg::SendMsg(msg) => {
self.msgs.write().await.push(msg.clone());
let id = self.db.generate_id().unwrap();
self.msgs.insert_(id.to_be_bytes(), &msg);
let mut handles = Vec::new();
for (&id, send) in self.senders.read().await.iter() {
if id == self.id {
@@ -81,13 +104,14 @@ impl RecvHandler<ClientMsg> for ClientHandler {
}
}
ClientMsg::RequestMsgs => {
let msgs = self.msgs.read().await.clone();
let msgs = self.msgs.iter_all().collect();
let _ = self.send.send(ServerMsg::LoadMsgs(msgs)).await;
}
}
}
async fn disconnect(&self, reason: DisconnectReason) -> () {
println!("disconnected: {:?}", self.send.remote());
match reason {
DisconnectReason::Closed | DisconnectReason::Timeout => (),
DisconnectReason::Other(e) => println!("connection issue: {e}"),