This commit is contained in:
2026-02-26 19:18:33 -05:00
parent 17eaf3e324
commit 54eb7d4261
11 changed files with 169 additions and 22 deletions

View File

@@ -1,4 +1,4 @@
use crate::ClientEvent;
use crate::ClientSender;
use dashmap::DashMap;
use openworm::net::{
ClientMsg, ClientMsgInst, RecvHandler, RequestId, RequestMsg, SERVER_NAME, ServerMsg,
@@ -26,6 +26,7 @@ pub struct ConnectInfo {
#[derive(Clone)]
pub struct NetHandle {
send: UnboundedSender<NetCtrlMsg>,
event_sender: ClientSender,
}
type NetResult<T> = Result<T, String>;
@@ -33,9 +34,12 @@ type NetResult<T> = Result<T, String>;
pub enum NetCtrlMsg {
Send(ClientMsg),
Request(ClientMsg, oneshot::Sender<ServerMsg>),
RequestSync(ClientMsg, Box<dyn FnOnce(ServerMsg) + Send + Sync>),
Exit,
}
type Resp<R: RequestMsg> = Result<R::Result, ()>;
impl NetHandle {
fn send_(&self, msg: NetCtrlMsg) {
let _ = self.send.send(msg);
@@ -45,7 +49,7 @@ impl NetHandle {
self.send_(NetCtrlMsg::Send(msg.into()));
}
pub async fn request<R: RequestMsg>(&self, msg: R) -> Result<R::Result, ()> {
pub async fn request<R: RequestMsg>(&self, msg: R) -> Resp<R> {
let (send, recv) = oneshot::channel();
self.send_(NetCtrlMsg::Request(msg.into(), send));
let Ok(recv) = recv.await else { return Err(()) };
@@ -56,11 +60,42 @@ impl NetHandle {
}
}
pub fn request_sync<R: RequestMsg>(&self, msg: R) -> SyncRecv<R> {
let (send, recv) = oneshot::channel();
let sender = self.event_sender.clone();
self.send_(NetCtrlMsg::RequestSync(
msg.into(),
Box::new(move |msg| {
let _ = send.send(if let Some(res) = R::result(msg) {
Ok(res)
} else {
Err(())
});
sender.run();
}),
));
SyncRecv::<R> { recv }
}
pub fn exit(self) {
self.send_(NetCtrlMsg::Exit);
}
}
pub struct SyncRecv<R: RequestMsg> {
recv: oneshot::Receiver<Resp<R>>,
}
impl<R: RequestMsg> SyncRecv<R> {
pub fn try_recv(&mut self) -> Option<Resp<R>> {
match self.recv.try_recv() {
Ok(res) => Some(res),
Err(oneshot::error::TryRecvError::Empty) => None,
Err(oneshot::error::TryRecvError::Closed) => Some(Err(())),
}
}
}
async fn connection_cert(
addr: SocketAddr,
cert: CertificateDer<'_>,
@@ -117,7 +152,11 @@ async fn connection_no_cert(addr: SocketAddr) -> NetResult<(Endpoint, Connection
}
impl NetHandle {
pub async fn connect(msg: impl MsgHandler, info: ConnectInfo) -> Result<Self, String> {
pub async fn connect(
msg: impl MsgHandler,
info: ConnectInfo,
event_sender: ClientSender,
) -> Result<Self, String> {
let (send, mut ui_recv) = tokio::sync::mpsc::unbounded_channel::<NetCtrlMsg>();
let cert = CertificateDer::from_slice(&info.cert);
@@ -133,6 +172,7 @@ impl NetHandle {
let mut req_id = RequestId::first();
let recv = Arc::new(ServerRecv {
msg,
requests_sync: DashMap::default(),
requests: DashMap::default(),
});
tokio::spawn(recv_uni(conn_, recv.clone()));
@@ -161,6 +201,17 @@ impl NetHandle {
break;
}
}
NetCtrlMsg::RequestSync(msg, f) => {
let msg = ClientMsgInst {
id: request_id,
msg,
};
recv.requests_sync.insert(request_id, f);
if send_uni(&conn, msg).await.is_err() {
println!("disconnected from server");
break;
}
}
NetCtrlMsg::Exit => {
conn.close(0u32.into(), &[]);
endpoint.wait_idle().await;
@@ -170,7 +221,7 @@ impl NetHandle {
}
});
Ok(NetHandle { send })
Ok(NetHandle { send, event_sender })
}
}
@@ -188,15 +239,18 @@ where
struct ServerRecv<F: MsgHandler> {
requests: DashMap<RequestId, oneshot::Sender<ServerMsg>>,
requests_sync: DashMap<RequestId, Box<dyn FnOnce(ServerMsg) + Send + Sync>>,
msg: F,
}
impl<F: MsgHandler> RecvHandler<ServerMsgInst> for ServerRecv<F> {
async fn msg(&self, resp: ServerMsgInst) {
if let Some(id) = resp.id
&& let Some((_, send)) = self.requests.remove(&id)
{
let _ = send.send(resp.msg);
if let Some(id) = resp.id {
if let Some((_, send)) = self.requests.remove(&id) {
let _ = send.send(resp.msg);
} else if let Some((_, f)) = self.requests_sync.remove(&id) {
f(resp.msg)
}
} else {
self.msg.run(resp.msg).await;
}