Session server
This commit is contained in:
@@ -66,16 +66,12 @@ impl Method for Data {
|
|||||||
|
|
||||||
let server = SessionServer::bind("127.0.0.1:8080").await?;
|
let server = SessionServer::bind("127.0.0.1:8080").await?;
|
||||||
|
|
||||||
loop {
|
server
|
||||||
let session = server.accept().await;
|
.session_loop(async |session, addr| {
|
||||||
|
// This will run on every new client
|
||||||
|
|
||||||
session
|
Ok(())
|
||||||
.on::<Data, _>(async |_, req| {
|
}).await;
|
||||||
// Response is required
|
|
||||||
Ok("Hello from server".to_string())
|
|
||||||
})
|
|
||||||
.await;
|
|
||||||
}
|
|
||||||
```
|
```
|
||||||
|
|
||||||
## Protocol
|
## Protocol
|
||||||
|
|||||||
+8
-21
@@ -1,7 +1,5 @@
|
|||||||
use serde::{Deserialize, Serialize};
|
use serde::{Deserialize, Serialize};
|
||||||
use tokio::net::TcpListener;
|
use session_rs::{Method, server::SessionServer};
|
||||||
|
|
||||||
use session_rs::{Method, session::Session, ws::WebSocket};
|
|
||||||
|
|
||||||
#[derive(Debug, Serialize, Deserialize)]
|
#[derive(Debug, Serialize, Deserialize)]
|
||||||
struct Data;
|
struct Data;
|
||||||
@@ -15,23 +13,10 @@ impl Method for Data {
|
|||||||
|
|
||||||
#[tokio::main(flavor = "current_thread")]
|
#[tokio::main(flavor = "current_thread")]
|
||||||
async fn main() -> session_rs::Result<()> {
|
async fn main() -> session_rs::Result<()> {
|
||||||
let listener = TcpListener::bind("127.0.0.1:8080").await?;
|
let server = SessionServer::bind("127.0.0.1:8080").await?;
|
||||||
println!("Server listening on ws://127.0.0.1:8080");
|
|
||||||
|
|
||||||
loop {
|
|
||||||
let (stream, addr) = listener.accept().await?;
|
|
||||||
println!("New connection: {}", addr);
|
|
||||||
|
|
||||||
tokio::spawn(async move {
|
|
||||||
// Wrap session in Arc so tasks can share it
|
|
||||||
let session = Session::from_ws(
|
|
||||||
WebSocket::handshake(stream)
|
|
||||||
.await
|
|
||||||
.expect("Failed to initialize websocket"),
|
|
||||||
);
|
|
||||||
|
|
||||||
session.start_receiver();
|
|
||||||
|
|
||||||
|
server
|
||||||
|
.session_loop(async |session, _| {
|
||||||
session
|
session
|
||||||
.on::<Data, _>(async |_, req| {
|
.on::<Data, _>(async |_, req| {
|
||||||
println!("Msg from client: {req}");
|
println!("Msg from client: {req}");
|
||||||
@@ -43,6 +28,8 @@ async fn main() -> session_rs::Result<()> {
|
|||||||
Ok("Hello from server".to_string())
|
Ok("Hello from server".to_string())
|
||||||
})
|
})
|
||||||
.await;
|
.await;
|
||||||
});
|
|
||||||
}
|
Ok(())
|
||||||
|
})
|
||||||
|
.await
|
||||||
}
|
}
|
||||||
|
|||||||
+41
-1
@@ -1 +1,41 @@
|
|||||||
pub struct SessionServer();
|
use std::{net::SocketAddr, sync::Arc};
|
||||||
|
|
||||||
|
use tokio::net::TcpListener;
|
||||||
|
|
||||||
|
use crate::{session::Session, ws::WebSocket};
|
||||||
|
|
||||||
|
pub struct SessionServer {
|
||||||
|
listener: TcpListener,
|
||||||
|
}
|
||||||
|
|
||||||
|
impl SessionServer {
|
||||||
|
pub async fn bind(addr: &str) -> crate::Result<Self> {
|
||||||
|
Ok(Self {
|
||||||
|
listener: TcpListener::bind(addr).await?,
|
||||||
|
})
|
||||||
|
}
|
||||||
|
|
||||||
|
pub async fn accept(&self) -> crate::Result<(Session, SocketAddr)> {
|
||||||
|
let (stream, addr) = self.listener.accept().await?;
|
||||||
|
|
||||||
|
let ws = WebSocket::handshake(stream).await?;
|
||||||
|
|
||||||
|
Ok((Session::from_ws(ws), addr))
|
||||||
|
}
|
||||||
|
|
||||||
|
pub async fn session_loop<Fut: Future<Output = crate::Result<()>> + Send + 'static>(
|
||||||
|
&self,
|
||||||
|
on_conn: impl Fn(Session, SocketAddr) -> Fut + 'static,
|
||||||
|
) -> crate::Result<()> {
|
||||||
|
let conn_handler = Arc::new(on_conn);
|
||||||
|
|
||||||
|
loop {
|
||||||
|
let (session, addr) = self.accept().await?;
|
||||||
|
let conn_handler = conn_handler.clone();
|
||||||
|
|
||||||
|
session.start_receiver();
|
||||||
|
|
||||||
|
tokio::spawn(conn_handler(session, addr));
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user