Working client
This commit is contained in:
+1
-3
@@ -1,12 +1,10 @@
|
|||||||
use std::sync::Arc;
|
use std::sync::Arc;
|
||||||
use tokio::net::TcpStream;
|
|
||||||
|
|
||||||
use session_rs::session::Session;
|
use session_rs::session::Session;
|
||||||
|
|
||||||
#[tokio::main(flavor = "current_thread")]
|
#[tokio::main(flavor = "current_thread")]
|
||||||
async fn main() -> session_rs::Result<()> {
|
async fn main() -> session_rs::Result<()> {
|
||||||
let stream = TcpStream::connect("127.0.0.1:8080").await?;
|
let session = Arc::new(Session::new_server("127.0.0.1:8080", "/").await?);
|
||||||
let session = Arc::new(Session::new(stream).await?);
|
|
||||||
|
|
||||||
// Spawn read loop
|
// Spawn read loop
|
||||||
let read_session = Arc::clone(&session);
|
let read_session = Arc::clone(&session);
|
||||||
|
|||||||
+1
-1
@@ -14,7 +14,7 @@ async fn main() -> session_rs::Result<()> {
|
|||||||
|
|
||||||
tokio::spawn(async move {
|
tokio::spawn(async move {
|
||||||
// Wrap session in Arc so tasks can share it
|
// Wrap session in Arc so tasks can share it
|
||||||
let session = match Session::new(stream).await {
|
let session = match Session::new_client(stream).await {
|
||||||
Ok(s) => Arc::new(s),
|
Ok(s) => Arc::new(s),
|
||||||
Err(e) => {
|
Err(e) => {
|
||||||
eprintln!("Handshake failed: {:?}", e);
|
eprintln!("Handshake failed: {:?}", e);
|
||||||
|
|||||||
@@ -1,8 +1,13 @@
|
|||||||
|
use sha1::{Digest, Sha1};
|
||||||
|
use std::sync::Arc;
|
||||||
use tokio::{
|
use tokio::{
|
||||||
io::{AsyncBufReadExt, AsyncWriteExt, BufReader},
|
io::{AsyncBufReadExt, AsyncWriteExt, BufReader},
|
||||||
net::TcpStream,
|
net::TcpStream,
|
||||||
|
sync::Mutex,
|
||||||
};
|
};
|
||||||
|
|
||||||
|
use crate::session::Session;
|
||||||
|
|
||||||
pub async fn handle_websocket_handshake(stream: &mut TcpStream) -> std::io::Result<()> {
|
pub async fn handle_websocket_handshake(stream: &mut TcpStream) -> std::io::Result<()> {
|
||||||
let (read_half, mut write_half) = stream.split();
|
let (read_half, mut write_half) = stream.split();
|
||||||
let mut reader = BufReader::new(read_half);
|
let mut reader = BufReader::new(read_half);
|
||||||
@@ -75,3 +80,90 @@ pub async fn handle_websocket_handshake(stream: &mut TcpStream) -> std::io::Resu
|
|||||||
write_half.write_all(response.as_bytes()).await?;
|
write_half.write_all(response.as_bytes()).await?;
|
||||||
Ok(())
|
Ok(())
|
||||||
}
|
}
|
||||||
|
|
||||||
|
impl Session {
|
||||||
|
pub async fn new_client(mut stream: TcpStream) -> crate::Result<Self> {
|
||||||
|
crate::handshake::handle_websocket_handshake(&mut stream).await?;
|
||||||
|
|
||||||
|
let (read, write) = stream.into_split();
|
||||||
|
|
||||||
|
Ok(Self {
|
||||||
|
reader: Arc::new(Mutex::new(read)),
|
||||||
|
writer: Arc::new(Mutex::new(write)),
|
||||||
|
id: rand::random(),
|
||||||
|
})
|
||||||
|
}
|
||||||
|
|
||||||
|
/// Connect to a WebSocket server and perform the handshake
|
||||||
|
pub async fn new_server(addr: &str, path: &str) -> crate::Result<Self> {
|
||||||
|
// 1. TCP connect
|
||||||
|
let mut stream = TcpStream::connect(addr).await?;
|
||||||
|
|
||||||
|
// 2. Generate Sec-WebSocket-Key
|
||||||
|
let key_bytes: [u8; 16] = rand::random();
|
||||||
|
let key = base64::encode(&key_bytes);
|
||||||
|
|
||||||
|
// 3. Send HTTP Upgrade request
|
||||||
|
let request = format!(
|
||||||
|
"GET {} HTTP/1.1\r\n\
|
||||||
|
Host: {}\r\n\
|
||||||
|
Upgrade: websocket\r\n\
|
||||||
|
Connection: Upgrade\r\n\
|
||||||
|
Sec-WebSocket-Key: {}\r\n\
|
||||||
|
Sec-WebSocket-Version: 13\r\n\
|
||||||
|
\r\n",
|
||||||
|
path, addr, key
|
||||||
|
);
|
||||||
|
stream.write_all(request.as_bytes()).await?;
|
||||||
|
stream.flush().await?;
|
||||||
|
|
||||||
|
// 4. Read HTTP response
|
||||||
|
let mut reader = BufReader::new(&mut stream);
|
||||||
|
let mut status_line = String::new();
|
||||||
|
reader.read_line(&mut status_line).await?;
|
||||||
|
if !status_line.starts_with("HTTP/1.1 101") {
|
||||||
|
return Err(crate::Error::HandshakeFailed(format!(
|
||||||
|
"Expected 101 Switching Protocols, got: {}",
|
||||||
|
status_line.trim_end()
|
||||||
|
)));
|
||||||
|
}
|
||||||
|
|
||||||
|
// Read headers
|
||||||
|
let mut sec_accept = None;
|
||||||
|
loop {
|
||||||
|
let mut line = String::new();
|
||||||
|
reader.read_line(&mut line).await?;
|
||||||
|
let line = line.trim_end();
|
||||||
|
if line.is_empty() {
|
||||||
|
break; // end of headers
|
||||||
|
}
|
||||||
|
if let Some((k, v)) = line.split_once(':') {
|
||||||
|
if k.eq_ignore_ascii_case("sec-websocket-accept") {
|
||||||
|
sec_accept = Some(v.trim().to_string());
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// 5. Verify Sec-WebSocket-Accept
|
||||||
|
let expected = {
|
||||||
|
let mut sha1 = Sha1::new();
|
||||||
|
sha1.update(key.as_bytes());
|
||||||
|
sha1.update(b"258EAFA5-E914-47DA-95CA-C5AB0DC85B11");
|
||||||
|
base64::encode(sha1.finalize())
|
||||||
|
};
|
||||||
|
if sec_accept.as_deref() != Some(expected.as_str()) {
|
||||||
|
return Err(crate::Error::HandshakeFailed(
|
||||||
|
"Sec-WebSocket-Accept mismatch".into(),
|
||||||
|
));
|
||||||
|
}
|
||||||
|
|
||||||
|
// 6. Upgrade succeeded, split stream
|
||||||
|
let (read, write) = stream.into_split();
|
||||||
|
|
||||||
|
Ok(Self {
|
||||||
|
reader: Arc::new(Mutex::new(read)),
|
||||||
|
writer: Arc::new(Mutex::new(write)),
|
||||||
|
id: rand::random(),
|
||||||
|
})
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|||||||
@@ -14,6 +14,7 @@ pub enum Error {
|
|||||||
Io(std::io::Error),
|
Io(std::io::Error),
|
||||||
Json(serde_json::Error),
|
Json(serde_json::Error),
|
||||||
InvalidFrame(String),
|
InvalidFrame(String),
|
||||||
|
HandshakeFailed(String),
|
||||||
ConnectionClosed,
|
ConnectionClosed,
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
+3
-17
@@ -9,23 +9,9 @@ use tokio::{
|
|||||||
};
|
};
|
||||||
|
|
||||||
pub struct Session {
|
pub struct Session {
|
||||||
reader: Arc<Mutex<tokio::net::tcp::OwnedReadHalf>>,
|
pub(crate) reader: Arc<Mutex<tokio::net::tcp::OwnedReadHalf>>,
|
||||||
writer: Arc<Mutex<tokio::net::tcp::OwnedWriteHalf>>,
|
pub(crate) writer: Arc<Mutex<tokio::net::tcp::OwnedWriteHalf>>,
|
||||||
id: u64,
|
pub(crate) id: u64,
|
||||||
}
|
|
||||||
|
|
||||||
impl Session {
|
|
||||||
pub async fn new(mut stream: TcpStream) -> crate::Result<Self> {
|
|
||||||
crate::handshake::handle_websocket_handshake(&mut stream).await?;
|
|
||||||
|
|
||||||
let (read, write) = stream.into_split();
|
|
||||||
|
|
||||||
Ok(Self {
|
|
||||||
reader: Arc::new(Mutex::new(read)),
|
|
||||||
writer: Arc::new(Mutex::new(write)),
|
|
||||||
id: rand::random(),
|
|
||||||
})
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
|
|
||||||
impl Clone for Session {
|
impl Clone for Session {
|
||||||
|
|||||||
Reference in New Issue
Block a user