work (new network + db initial working state)
This commit is contained in:
@@ -1,14 +1,13 @@
|
||||
// mod data;
|
||||
mod db;
|
||||
mod net;
|
||||
|
||||
use crate::db::{Db, Msg, User, open_db};
|
||||
use crate::db::{Db, Msg, User};
|
||||
use clap::Parser;
|
||||
use net::{ClientSender, ConAccepter, listen};
|
||||
use openworm::{
|
||||
net::{
|
||||
ClientMsg, ClientMsgInst, CreateAccount, DisconnectReason, LoadMsg, RecvHandler,
|
||||
ServerError, ServerMsg, install_crypto_provider,
|
||||
AccountCreated, ClientMsg, ClientRequestMsg, CreateAccount, DisconnectReason, LoadMsg,
|
||||
RecvHandler, ServerError, ServerMsg, install_crypto_provider,
|
||||
},
|
||||
rsc::DataDir,
|
||||
};
|
||||
@@ -43,7 +42,7 @@ fn main() {
|
||||
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 db = Db::open(path.join("server_db"));
|
||||
let handler = ServerListener {
|
||||
senders: Default::default(),
|
||||
count: 0.into(),
|
||||
@@ -56,9 +55,6 @@ pub async fn run_server(port: u16) {
|
||||
endpoint.close(0u32.into(), &[]);
|
||||
let _ = handle.await;
|
||||
endpoint.wait_idle().await;
|
||||
println!("saving...");
|
||||
db.flush();
|
||||
println!("saved");
|
||||
}
|
||||
|
||||
type ClientId = u64;
|
||||
@@ -76,7 +72,7 @@ pub enum ClientState {
|
||||
}
|
||||
|
||||
impl ConAccepter for ServerListener {
|
||||
async fn accept(&self, send: ClientSender) -> impl RecvHandler<ClientMsgInst> {
|
||||
async fn accept(&self, send: ClientSender) -> impl RecvHandler<ClientRequestMsg> {
|
||||
let id = self.count.fetch_add(1, Ordering::Release);
|
||||
self.senders.write().await.insert(id, send.clone());
|
||||
ClientHandler {
|
||||
@@ -97,29 +93,33 @@ struct ClientHandler {
|
||||
state: Arc<RwLock<ClientState>>,
|
||||
}
|
||||
|
||||
impl RecvHandler<ClientMsgInst> for ClientHandler {
|
||||
impl RecvHandler<ClientRequestMsg> for ClientHandler {
|
||||
async fn connect(&self) -> () {
|
||||
println!("connected: {:?}", self.send.remote().ip());
|
||||
}
|
||||
async fn msg(&self, msg: ClientMsgInst) {
|
||||
let msg = ClientMsg::from(msg);
|
||||
async fn msg(&self, req: ClientRequestMsg) {
|
||||
let msg = ClientMsg::from(req.msg);
|
||||
let replier = self.send.replier(req.id);
|
||||
match msg {
|
||||
ClientMsg::SendMsg(msg) => {
|
||||
let ClientState::Authed(uid) = &*self.state.read().await else {
|
||||
let _ = self.send.send(ServerError::NotLoggedIn).await;
|
||||
let _ = replier.send(ServerError::NotLoggedIn).await;
|
||||
return;
|
||||
};
|
||||
let msg = Msg {
|
||||
user: *uid,
|
||||
author: *uid,
|
||||
content: msg.content,
|
||||
};
|
||||
let id = self.db.generate_id().unwrap();
|
||||
self.db.msgs.insert(&id, &msg);
|
||||
// TODO: it is technically possible to send 2 messages at the exact same time...
|
||||
// should probably append a number if one already exists at that time,
|
||||
// but also I can't see this ever happening...?
|
||||
// should be an easy fix later (write tx)
|
||||
let timestamp = time::OffsetDateTime::now_utc().unix_timestamp_nanos();
|
||||
self.db.msgs.insert(×tamp, &msg);
|
||||
let mut handles = Vec::new();
|
||||
let user: User = self.db.users.get(uid).unwrap();
|
||||
let msg = LoadMsg {
|
||||
content: msg.content,
|
||||
user: user.username,
|
||||
author: *uid,
|
||||
};
|
||||
for (&id, send) in self.senders.read().await.iter() {
|
||||
if id == self.id {
|
||||
@@ -128,7 +128,12 @@ impl RecvHandler<ClientMsgInst> for ClientHandler {
|
||||
let send = send.clone();
|
||||
let msg = msg.clone();
|
||||
let fut = async move {
|
||||
let _ = send.send(msg).await;
|
||||
let _ = send
|
||||
.send(LoadMsg {
|
||||
content: msg.content,
|
||||
author: msg.author,
|
||||
})
|
||||
.await;
|
||||
};
|
||||
handles.push(tokio::spawn(fut));
|
||||
}
|
||||
@@ -138,27 +143,19 @@ impl RecvHandler<ClientMsgInst> for ClientHandler {
|
||||
}
|
||||
ClientMsg::RequestMsgs => {
|
||||
let ClientState::Authed(_uid) = &*self.state.read().await else {
|
||||
let _ = self.send.send(ServerError::NotLoggedIn).await;
|
||||
let _ = replier.send(ServerError::NotLoggedIn).await;
|
||||
return;
|
||||
};
|
||||
let msgs = self
|
||||
.db
|
||||
.msgs
|
||||
.iter_all()
|
||||
.map(|msg| {
|
||||
let user = self
|
||||
.db
|
||||
.users
|
||||
.get(&msg.user)
|
||||
.map(|user| user.username.to_string())
|
||||
.unwrap_or("deleted user".to_string());
|
||||
LoadMsg {
|
||||
content: msg.content,
|
||||
user,
|
||||
}
|
||||
.values()
|
||||
.map(|msg| LoadMsg {
|
||||
content: msg.content,
|
||||
author: msg.author,
|
||||
})
|
||||
.collect();
|
||||
let _ = self.send.send(ServerMsg::LoadMsgs(msgs)).await;
|
||||
let _ = replier.send(ServerMsg::LoadMsgs(msgs)).await;
|
||||
}
|
||||
ClientMsg::CreateAccount(info) => {
|
||||
let CreateAccount {
|
||||
@@ -167,28 +164,41 @@ impl RecvHandler<ClientMsgInst> for ClientHandler {
|
||||
password,
|
||||
login_key,
|
||||
} = &info;
|
||||
if !self.db.usernames.init_unique(username) {
|
||||
let _ = self.send.send(ServerError::UsernameTaken).await;
|
||||
return;
|
||||
}
|
||||
let id = self.db.generate_id().unwrap();
|
||||
let salt = SaltString::generate(&mut OsRng);
|
||||
let params = scrypt::Params::new(11, 8, 1, 32).unwrap();
|
||||
let hash = Scrypt
|
||||
.hash_password_customized(password.as_bytes(), None, None, params, &salt)
|
||||
.unwrap()
|
||||
.to_string();
|
||||
self.db.users.insert(
|
||||
&id,
|
||||
&User {
|
||||
username: username.clone(),
|
||||
password_hash: hash,
|
||||
},
|
||||
);
|
||||
let mut id;
|
||||
loop {
|
||||
let mut tx = self.db.write_tx();
|
||||
if tx.has_key(&self.db.usernames, username.clone()) {
|
||||
let _ = replier.send(ServerError::UsernameTaken).await;
|
||||
return;
|
||||
}
|
||||
id = rand::random();
|
||||
while tx.has_key(&self.db.users, id) {
|
||||
id = rand::random();
|
||||
}
|
||||
tx.insert(
|
||||
&self.db.users,
|
||||
&id,
|
||||
&User {
|
||||
username: username.clone(),
|
||||
password_hash: hash.clone(),
|
||||
bio: String::new(),
|
||||
pfp: None,
|
||||
},
|
||||
);
|
||||
tx.insert(&self.db.usernames, username, &id);
|
||||
if tx.commit() {
|
||||
break;
|
||||
}
|
||||
}
|
||||
println!("account created: \"{username}\"");
|
||||
self.db.usernames.insert(&username, &id);
|
||||
*self.state.write().await = ClientState::Authed(id);
|
||||
// let _ = self.send.send(ServerMsg::Login()).await;
|
||||
let _ = replier.send(AccountCreated {}).await;
|
||||
} // ClientMsgType::Login { username, password } => {
|
||||
// let Some(id) = self.db.usernames.get(&username) else {
|
||||
// let _ = self.send.send(ServerError::UnknownUsername).await;
|
||||
|
||||
Reference in New Issue
Block a user