Server client example
This commit is contained in:
+36
-1
@@ -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(())
|
||||
}
|
||||
|
||||
+61
-1
@@ -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);
|
||||
});
|
||||
}
|
||||
}
|
||||
|
||||
@@ -9,6 +9,7 @@ pub enum SessionFrame<T> {
|
||||
|
||||
pub type Result<T> = std::result::Result<T, Error>;
|
||||
|
||||
#[derive(Debug)]
|
||||
pub enum Error {
|
||||
Io(std::io::Error),
|
||||
Json(serde_json::Error),
|
||||
|
||||
Reference in New Issue
Block a user