Better binary message handling
This commit is contained in:
+9
-1
@@ -24,7 +24,7 @@ $$ | $$ |$$ / $$ |\\$$\\ $$ |$$ / $$ |
|
|||||||
use crate::{
|
use crate::{
|
||||||
cli, logger,
|
cli, logger,
|
||||||
plugin::{Plugin, loader::PluginLoader, types::LoaderMessage},
|
plugin::{Plugin, loader::PluginLoader, types::LoaderMessage},
|
||||||
types,
|
types::{self, message::WsMessage},
|
||||||
utils::{self, auth, client::Client},
|
utils::{self, auth, client::Client},
|
||||||
};
|
};
|
||||||
|
|
||||||
@@ -198,10 +198,18 @@ impl Server {
|
|||||||
while !self.shutting_down.load(Ordering::SeqCst) {
|
while !self.shutting_down.load(Ordering::SeqCst) {
|
||||||
let req = client.read()?;
|
let req = client.read()?;
|
||||||
if let Some(r) = &req {
|
if let Some(r) = &req {
|
||||||
|
match r {
|
||||||
|
WsMessage::Binary(_) => {
|
||||||
|
// ignore binary
|
||||||
|
}
|
||||||
|
_ => {
|
||||||
self.send_plugin_message(&LoaderMessage::Request {
|
self.send_plugin_message(&LoaderMessage::Request {
|
||||||
user_id: client.get_uuid().unwrap_or_default(),
|
user_id: client.get_uuid().unwrap_or_default(),
|
||||||
msg: r.clone(),
|
msg: r.clone(),
|
||||||
})?;
|
})?;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
self.wrap_err(&client, self.call_request(r, &client))?;
|
self.wrap_err(&client, self.call_request(r, &client))?;
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
+87
-17
@@ -199,6 +199,34 @@ impl Client {
|
|||||||
Ok(())
|
Ok(())
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// Send a binary WebSocket frame (server -> client)
|
||||||
|
pub fn send_bin(&self, payload: &[u8]) -> crate::Result<()> {
|
||||||
|
let mut stream = self.0.try_clone()?;
|
||||||
|
|
||||||
|
let mut header = Vec::with_capacity(10);
|
||||||
|
|
||||||
|
// FIN=1, opcode=2 (binary)
|
||||||
|
header.push(0x82);
|
||||||
|
|
||||||
|
let len = payload.len();
|
||||||
|
|
||||||
|
if len < 126 {
|
||||||
|
header.push(len as u8); // mask bit = 0
|
||||||
|
} else if len <= 0xFFFF {
|
||||||
|
header.push(126);
|
||||||
|
header.extend_from_slice(&(len as u16).to_be_bytes());
|
||||||
|
} else {
|
||||||
|
header.push(127);
|
||||||
|
header.extend_from_slice(&(len as u64).to_be_bytes());
|
||||||
|
}
|
||||||
|
|
||||||
|
stream.write_all(&header)?;
|
||||||
|
stream.write_all(payload)?;
|
||||||
|
stream.flush()?;
|
||||||
|
|
||||||
|
Ok(())
|
||||||
|
}
|
||||||
|
|
||||||
/// Read a full WebSocket message, handling fragmentation and control frames.
|
/// Read a full WebSocket message, handling fragmentation and control frames.
|
||||||
///
|
///
|
||||||
/// Returns:
|
/// Returns:
|
||||||
@@ -211,21 +239,22 @@ impl Client {
|
|||||||
let mut stream = self.0.try_clone()?;
|
let mut stream = self.0.try_clone()?;
|
||||||
|
|
||||||
let mut message_payload = Vec::new();
|
let mut message_payload = Vec::new();
|
||||||
|
let mut expecting_continuation = false;
|
||||||
|
let mut message_type: Option<u8> = None; // 0x1 for text, 0x2 for binary
|
||||||
|
|
||||||
loop {
|
loop {
|
||||||
// read the 2-byte header
|
// Read 2-byte header
|
||||||
let mut header = [0u8; 2];
|
let mut header = [0u8; 2];
|
||||||
if let Err(e) = stream.read_exact(&mut header) {
|
match stream.read_exact(&mut header) {
|
||||||
if e.kind() == io::ErrorKind::WouldBlock || e.kind() == io::ErrorKind::TimedOut {
|
Ok(_) => {}
|
||||||
|
Err(e) => match e.kind() {
|
||||||
|
io::ErrorKind::WouldBlock | io::ErrorKind::TimedOut => {
|
||||||
self.send_ping()?;
|
self.send_ping()?;
|
||||||
continue;
|
continue;
|
||||||
}
|
}
|
||||||
|
io::ErrorKind::UnexpectedEof | io::ErrorKind::BrokenPipe => return Ok(None),
|
||||||
if e.kind() == io::ErrorKind::UnexpectedEof || e.kind() == io::ErrorKind::BrokenPipe
|
_ => return Err(e.into()),
|
||||||
{
|
},
|
||||||
return Ok(None);
|
|
||||||
}
|
|
||||||
return Err(e.into());
|
|
||||||
}
|
}
|
||||||
|
|
||||||
let fin = header[0] & 0x80 != 0;
|
let fin = header[0] & 0x80 != 0;
|
||||||
@@ -233,7 +262,7 @@ impl Client {
|
|||||||
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;
|
||||||
|
|
||||||
// Extended payload lengths
|
// Extended payload length
|
||||||
if payload_len == 126 {
|
if payload_len == 126 {
|
||||||
let mut ext_len = [0u8; 2];
|
let mut ext_len = [0u8; 2];
|
||||||
stream.read_exact(&mut ext_len)?;
|
stream.read_exact(&mut ext_len)?;
|
||||||
@@ -244,7 +273,7 @@ impl Client {
|
|||||||
payload_len = u64::from_be_bytes(ext_len);
|
payload_len = u64::from_be_bytes(ext_len);
|
||||||
}
|
}
|
||||||
|
|
||||||
// Mask key (client→server MUST be masked)
|
// Mask key
|
||||||
let mut mask = [0u8; 4];
|
let mut mask = [0u8; 4];
|
||||||
if masked {
|
if masked {
|
||||||
stream.read_exact(&mut mask)?;
|
stream.read_exact(&mut mask)?;
|
||||||
@@ -265,7 +294,7 @@ impl Client {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
// Read payload + unmask
|
// Read payload
|
||||||
let mut payload = vec![0u8; payload_len as usize];
|
let mut payload = vec![0u8; payload_len as usize];
|
||||||
if payload_len > 0 {
|
if payload_len > 0 {
|
||||||
stream.read_exact(&mut payload)?;
|
stream.read_exact(&mut payload)?;
|
||||||
@@ -275,13 +304,45 @@ impl Client {
|
|||||||
}
|
}
|
||||||
|
|
||||||
match opcode {
|
match opcode {
|
||||||
0x0 | 0x1 | 0x2 => {
|
0x0 => {
|
||||||
// Continuation / Text / Binary
|
// Continuation
|
||||||
|
if !expecting_continuation {
|
||||||
|
let _ = self.send_close(1002, "Unexpected continuation frame");
|
||||||
|
return Ok(None);
|
||||||
|
}
|
||||||
message_payload.extend(payload);
|
message_payload.extend(payload);
|
||||||
|
if fin {
|
||||||
|
break;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
0x1 => {
|
||||||
|
// Text
|
||||||
|
if expecting_continuation {
|
||||||
|
let _ =
|
||||||
|
self.send_close(1002, "New data frame while expecting continuation");
|
||||||
|
return Ok(None);
|
||||||
|
}
|
||||||
|
message_payload.extend(payload);
|
||||||
|
message_type = Some(0x1);
|
||||||
if fin {
|
if fin {
|
||||||
break;
|
break;
|
||||||
} else {
|
} else {
|
||||||
continue;
|
expecting_continuation = true;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
0x2 => {
|
||||||
|
// Binary
|
||||||
|
if expecting_continuation {
|
||||||
|
let _ =
|
||||||
|
self.send_close(1002, "New data frame while expecting continuation");
|
||||||
|
return Ok(None);
|
||||||
|
}
|
||||||
|
message_payload.extend(payload);
|
||||||
|
message_type = Some(0x2);
|
||||||
|
if fin {
|
||||||
|
break;
|
||||||
|
} else {
|
||||||
|
expecting_continuation = true;
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
0x8 => {
|
0x8 => {
|
||||||
@@ -301,10 +362,12 @@ impl Client {
|
|||||||
return Ok(None);
|
return Ok(None);
|
||||||
}
|
}
|
||||||
0x9 => {
|
0x9 => {
|
||||||
|
// Ping
|
||||||
self.send_pong()?;
|
self.send_pong()?;
|
||||||
continue;
|
continue;
|
||||||
}
|
}
|
||||||
0xA => {
|
0xA => {
|
||||||
|
// Pong
|
||||||
continue;
|
continue;
|
||||||
}
|
}
|
||||||
_ => {
|
_ => {
|
||||||
@@ -314,13 +377,20 @@ impl Client {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
// Try parsing JSON into ClientMessage
|
// Convert payload into proper message type
|
||||||
let message = match String::from_utf8(message_payload.clone()) {
|
let message = match message_type {
|
||||||
|
Some(0x1) => {
|
||||||
|
// Text frame → try JSON, otherwise keep text
|
||||||
|
match String::from_utf8(message_payload.clone()) {
|
||||||
Ok(text) => match serde_json::from_str(&text) {
|
Ok(text) => match serde_json::from_str(&text) {
|
||||||
Ok(msg) => WsMessage::Message(msg),
|
Ok(msg) => WsMessage::Message(msg),
|
||||||
Err(_) => WsMessage::String(text),
|
Err(_) => WsMessage::String(text),
|
||||||
},
|
},
|
||||||
Err(_) => WsMessage::Binary(message_payload),
|
Err(_) => WsMessage::Binary(message_payload),
|
||||||
|
}
|
||||||
|
}
|
||||||
|
Some(0x2) => WsMessage::Binary(message_payload),
|
||||||
|
_ => return Ok(None), // Should not happen
|
||||||
};
|
};
|
||||||
|
|
||||||
Ok(Some(message))
|
Ok(Some(message))
|
||||||
|
|||||||
Reference in New Issue
Block a user