Refactoring
This commit is contained in:
@@ -5,6 +5,9 @@ pub mod session;
|
|||||||
pub enum SessionFrame<T> {
|
pub enum SessionFrame<T> {
|
||||||
Typed(T),
|
Typed(T),
|
||||||
Binary(Vec<u8>),
|
Binary(Vec<u8>),
|
||||||
|
Ping,
|
||||||
|
Pong,
|
||||||
|
Close
|
||||||
}
|
}
|
||||||
|
|
||||||
pub type Result<T> = std::result::Result<T, Error>;
|
pub type Result<T> = std::result::Result<T, Error>;
|
||||||
|
|||||||
+39
-30
@@ -7,6 +7,8 @@ use tokio::{
|
|||||||
sync::Mutex,
|
sync::Mutex,
|
||||||
};
|
};
|
||||||
|
|
||||||
|
use crate::SessionFrame;
|
||||||
|
|
||||||
pub struct Session {
|
pub struct Session {
|
||||||
pub(crate) reader: Arc<Mutex<tokio::net::tcp::OwnedReadHalf>>,
|
pub(crate) reader: Arc<Mutex<tokio::net::tcp::OwnedReadHalf>>,
|
||||||
pub(crate) writer: Arc<Mutex<tokio::net::tcp::OwnedWriteHalf>>,
|
pub(crate) writer: Arc<Mutex<tokio::net::tcp::OwnedWriteHalf>>,
|
||||||
@@ -119,14 +121,14 @@ impl Session {
|
|||||||
impl Session {
|
impl Session {
|
||||||
/// Read a full WebSocket frame (handling masking and control frames)
|
/// Read a full WebSocket frame (handling masking and control frames)
|
||||||
/// Returns (opcode, payload)
|
/// Returns (opcode, payload)
|
||||||
pub async fn read_frame(&self) -> crate::Result<Option<(u8, Vec<u8>)>> {
|
pub async fn read_frame(&self) -> crate::Result<(bool, u8, Vec<u8>)> {
|
||||||
let mut reader = self.reader.lock().await;
|
let mut reader = self.reader.lock().await;
|
||||||
|
|
||||||
// --- 1. Read first 2-byte header ---
|
// --- 1. Read first 2-byte header ---
|
||||||
let mut header = [0u8; 2];
|
let mut header = [0u8; 2];
|
||||||
reader.read_exact(&mut header).await?;
|
reader.read_exact(&mut header).await?;
|
||||||
|
|
||||||
// let fin = header[0] & 0x80 != 0;
|
let fin = header[0] & 0x80 != 0;
|
||||||
let opcode = header[0] & 0x0F;
|
let opcode = header[0] & 0x0F;
|
||||||
let masked = header[1] & 0x80 != 0;
|
let masked = header[1] & 0x80 != 0;
|
||||||
let mut payload_len = (header[1] & 0x7F) as u64;
|
let mut payload_len = (header[1] & 0x7F) as u64;
|
||||||
@@ -163,35 +165,42 @@ impl Session {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
// --- 5. Handle control frames immediately ---
|
// --- 6. Return opcode + payload ---
|
||||||
match opcode {
|
Ok((fin, opcode, payload))
|
||||||
0x8 => {
|
|
||||||
// Close
|
|
||||||
self.close().await.ok();
|
|
||||||
return Err(crate::Error::ConnectionClosed);
|
|
||||||
}
|
|
||||||
0x9 => {
|
|
||||||
// Ping
|
|
||||||
self.send_pong().await.ok();
|
|
||||||
return Ok(None);
|
|
||||||
}
|
|
||||||
0xA => {
|
|
||||||
// Pong, ignore
|
|
||||||
return Ok(None);
|
|
||||||
}
|
|
||||||
0x0 | 0x1 | 0x2 => {
|
|
||||||
// Continuation / Text / Binary → valid payload
|
|
||||||
}
|
|
||||||
_ => {
|
|
||||||
self.close().await.ok();
|
|
||||||
return Err(crate::Error::InvalidFrame(format!(
|
|
||||||
"Unknown opcode: {}",
|
|
||||||
opcode
|
|
||||||
)));
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
|
|
||||||
// --- 6. Return opcode + payload ---
|
pub async fn read<T>(&self) -> crate::Result<SessionFrame<T>> {
|
||||||
Ok(Some((opcode, payload)))
|
let (fin, opcode, payload) = self.read_frame().await?;
|
||||||
|
|
||||||
|
match opcode {
|
||||||
|
// Close
|
||||||
|
0x8 => {
|
||||||
|
self.close().await.ok();
|
||||||
|
Ok(SessionFrame::Close)
|
||||||
|
}
|
||||||
|
|
||||||
|
// Ping
|
||||||
|
0x9 => {
|
||||||
|
self.send_pong().await.ok();
|
||||||
|
Ok(SessionFrame::Ping)
|
||||||
|
}
|
||||||
|
|
||||||
|
// Pong, ignore
|
||||||
|
0xA => Ok(SessionFrame::Pong),
|
||||||
|
|
||||||
|
// Continuation / Text / Binary → valid payload
|
||||||
|
0x0 => Ok(None),
|
||||||
|
|
||||||
|
0x1 => Ok(None),
|
||||||
|
|
||||||
|
0x2 => Ok(None),
|
||||||
|
|
||||||
|
_ => {
|
||||||
|
self.close().await.ok();
|
||||||
|
Err(crate::Error::InvalidFrame(format!(
|
||||||
|
"Unknown opcode: {opcode}"
|
||||||
|
)))
|
||||||
|
}
|
||||||
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user