From 495de7e87d542681e0b6b96dc7eeed7892fbf2bb Mon Sep 17 00:00:00 2001 From: Klesti Selimaj Date: Wed, 18 Feb 2026 21:54:16 +0100 Subject: [PATCH] Server client example --- examples/client.rs | 37 ++++++++++++++++++++++++++- examples/server.rs | 62 +++++++++++++++++++++++++++++++++++++++++++++- src/lib.rs | 1 + 3 files changed, 98 insertions(+), 2 deletions(-) diff --git a/examples/client.rs b/examples/client.rs index ab9323e..5d09db5 100644 --- a/examples/client.rs +++ b/examples/client.rs @@ -1,2 +1,37 @@ +use std::sync::Arc; +use tokio::net::TcpStream; + +use session_rs::session::Session; + #[tokio::main(flavor = "current_thread")] -async fn main() {} +async fn main() -> session_rs::Result<()> { + let stream = TcpStream::connect("127.0.0.1:8080").await?; + let session = Arc::new(Session::new(stream).await?); + + // Spawn read loop + let read_session = Arc::clone(&session); + tokio::spawn(async move { + loop { + match read_session.read_frame().await { + Ok(Some((opcode, payload))) => { + if opcode == 0x1 { + let text = String::from_utf8(payload).unwrap_or_default(); + println!("Server says: {}", text); + } + } + Ok(None) => {} + Err(_) => break, + } + } + }); + + // Send a few messages + for i in 0..5 { + let msg = serde_json::json!({ "hello": i }); + session.send(&msg).await?; + tokio::time::sleep(std::time::Duration::from_secs(1)).await; + } + + session.close().await?; + Ok(()) +} diff --git a/examples/server.rs b/examples/server.rs index ab9323e..35b294c 100644 --- a/examples/server.rs +++ b/examples/server.rs @@ -1,2 +1,62 @@ +use std::sync::Arc; +use tokio::net::TcpListener; + +use session_rs::session::Session; + #[tokio::main(flavor = "current_thread")] -async fn main() {} +async fn main() -> session_rs::Result<()> { + let listener = TcpListener::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 = match Session::new(stream).await { + Ok(s) => Arc::new(s), + Err(e) => { + eprintln!("Handshake failed: {:?}", e); + return; + } + }; + + // Simple ping loop + let ping_session = Arc::clone(&session); + tokio::spawn(async move { + let mut interval = tokio::time::interval(std::time::Duration::from_secs(15)); + loop { + interval.tick().await; + if ping_session.send_ping().await.is_err() { + break; + } + } + }); + + // Read loop + loop { + match session.read_frame().await { + Ok(Some((opcode, payload))) => { + if opcode == 0x1 { + // Text frame → parse JSON if possible + let text = String::from_utf8(payload).unwrap_or_default(); + println!("Received text: {}", text); + + // Echo back + if let Err(e) = session.send(&serde_json::json!({"echo": text})).await { + eprintln!("Send error: {:?}", e); + break; + } + } + } + Ok(None) => {} + Err(_) => break, // connection closed + } + } + + let _ = session.close().await; + println!("Connection {} closed", addr); + }); + } +} diff --git a/src/lib.rs b/src/lib.rs index 935aa31..64c4fe1 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -9,6 +9,7 @@ pub enum SessionFrame { pub type Result = std::result::Result; +#[derive(Debug)] pub enum Error { Io(std::io::Error), Json(serde_json::Error),