Working WebSocket
This commit is contained in:
+3
-6
@@ -1,6 +1,6 @@
|
|||||||
use std::sync::Arc;
|
use std::sync::Arc;
|
||||||
|
|
||||||
use session_rs::ws::WebSocket;
|
use session_rs::{SessionFrame, ws::WebSocket};
|
||||||
|
|
||||||
#[tokio::main(flavor = "current_thread")]
|
#[tokio::main(flavor = "current_thread")]
|
||||||
async fn main() -> session_rs::Result<()> {
|
async fn main() -> session_rs::Result<()> {
|
||||||
@@ -11,13 +11,10 @@ async fn main() -> session_rs::Result<()> {
|
|||||||
tokio::spawn(async move {
|
tokio::spawn(async move {
|
||||||
loop {
|
loop {
|
||||||
match read_session.read().await {
|
match read_session.read().await {
|
||||||
Ok(Some((opcode, payload))) => {
|
Ok(SessionFrame::Text(text)) => {
|
||||||
if opcode == 0x1 {
|
|
||||||
let text = String::from_utf8(payload).unwrap_or_default();
|
|
||||||
println!("Server says: {}", text);
|
println!("Server says: {}", text);
|
||||||
}
|
}
|
||||||
}
|
Ok(_) => {}
|
||||||
Ok(None) => {}
|
|
||||||
Err(_) => break,
|
Err(_) => break,
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
+4
-8
@@ -1,7 +1,7 @@
|
|||||||
use std::sync::Arc;
|
use std::sync::Arc;
|
||||||
use tokio::net::TcpListener;
|
use tokio::net::TcpListener;
|
||||||
|
|
||||||
use session_rs::session::WebSocket;
|
use session_rs::{SessionFrame, ws::WebSocket};
|
||||||
|
|
||||||
#[tokio::main(flavor = "current_thread")]
|
#[tokio::main(flavor = "current_thread")]
|
||||||
async fn main() -> session_rs::Result<()> {
|
async fn main() -> session_rs::Result<()> {
|
||||||
@@ -26,11 +26,8 @@ async fn main() -> session_rs::Result<()> {
|
|||||||
|
|
||||||
// Read loop
|
// Read loop
|
||||||
loop {
|
loop {
|
||||||
match session.read_frame().await {
|
match session.read().await {
|
||||||
Ok(Some((opcode, payload))) => {
|
Ok(SessionFrame::Text(text)) => {
|
||||||
if opcode == 0x1 {
|
|
||||||
// Text frame → parse JSON if possible
|
|
||||||
let text = String::from_utf8(payload).unwrap_or_default();
|
|
||||||
println!("Received text: {}", text);
|
println!("Received text: {}", text);
|
||||||
|
|
||||||
// Echo back
|
// Echo back
|
||||||
@@ -39,8 +36,7 @@ async fn main() -> session_rs::Result<()> {
|
|||||||
break;
|
break;
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
Ok(_) => {}
|
||||||
Ok(None) => {}
|
|
||||||
Err(e) => {
|
Err(e) => {
|
||||||
eprintln!("{e:?}");
|
eprintln!("{e:?}");
|
||||||
break;
|
break;
|
||||||
|
|||||||
+12
-3
@@ -1,13 +1,15 @@
|
|||||||
|
use std::string::FromUtf8Error;
|
||||||
|
|
||||||
pub mod server;
|
pub mod server;
|
||||||
pub mod session;
|
pub mod session;
|
||||||
pub mod ws;
|
pub mod ws;
|
||||||
|
|
||||||
pub enum SessionFrame<T> {
|
pub enum SessionFrame {
|
||||||
Typed(T),
|
Text(String),
|
||||||
Binary(Vec<u8>),
|
Binary(Vec<u8>),
|
||||||
Ping,
|
Ping,
|
||||||
Pong,
|
Pong,
|
||||||
Close
|
Close,
|
||||||
}
|
}
|
||||||
|
|
||||||
pub type Result<T> = std::result::Result<T, Error>;
|
pub type Result<T> = std::result::Result<T, Error>;
|
||||||
@@ -19,6 +21,7 @@ pub enum Error {
|
|||||||
InvalidFrame(String),
|
InvalidFrame(String),
|
||||||
HandshakeFailed(String),
|
HandshakeFailed(String),
|
||||||
ConnectionClosed,
|
ConnectionClosed,
|
||||||
|
Utf8(FromUtf8Error),
|
||||||
}
|
}
|
||||||
|
|
||||||
impl From<std::io::Error> for Error {
|
impl From<std::io::Error> for Error {
|
||||||
@@ -32,3 +35,9 @@ impl From<serde_json::Error> for Error {
|
|||||||
Self::Json(value)
|
Self::Json(value)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
impl From<FromUtf8Error> for Error {
|
||||||
|
fn from(value: FromUtf8Error) -> Self {
|
||||||
|
Self::Utf8(value)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|||||||
+32
-9
@@ -171,8 +171,36 @@ impl WebSocket {
|
|||||||
Ok((fin, opcode, payload))
|
Ok((fin, opcode, payload))
|
||||||
}
|
}
|
||||||
|
|
||||||
pub async fn read<T>(&self) -> crate::Result<SessionFrame<T>> {
|
pub async fn read(&self) -> crate::Result<SessionFrame> {
|
||||||
let (fin, opcode, payload) = self.read_frame().await?;
|
let (fin, opcode, mut payload) = self.read_frame().await?;
|
||||||
|
|
||||||
|
if !fin {
|
||||||
|
// Continuation loop
|
||||||
|
while let (fin, o, mut p) = self.read_frame().await?
|
||||||
|
&& !fin
|
||||||
|
{
|
||||||
|
match o {
|
||||||
|
// Continuation
|
||||||
|
0x0 => payload.append(&mut p),
|
||||||
|
// Close
|
||||||
|
0x8 => {
|
||||||
|
self.close().await.ok();
|
||||||
|
}
|
||||||
|
// Ping
|
||||||
|
0x9 => {
|
||||||
|
self.send_pong().await.ok();
|
||||||
|
}
|
||||||
|
// Pong
|
||||||
|
0xA => {}
|
||||||
|
_ => {
|
||||||
|
self.close().await.ok();
|
||||||
|
return Err(crate::Error::InvalidFrame(format!(
|
||||||
|
"Unknown opcode: {opcode}"
|
||||||
|
)));
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
match opcode {
|
match opcode {
|
||||||
// Close
|
// Close
|
||||||
@@ -190,16 +218,11 @@ impl WebSocket {
|
|||||||
// Pong
|
// Pong
|
||||||
0xA => Ok(SessionFrame::Pong),
|
0xA => Ok(SessionFrame::Pong),
|
||||||
|
|
||||||
// Continuation
|
|
||||||
0x0 => Ok(SessionFrame::Pong),
|
|
||||||
|
|
||||||
// Text
|
// Text
|
||||||
// 0x1 => {
|
0x1 => Ok(SessionFrame::Text(String::from_utf8(payload)?)),
|
||||||
|
|
||||||
// },
|
|
||||||
|
|
||||||
// Binary
|
// Binary
|
||||||
0x2 => Ok(SessionFrame::Pong),
|
0x2 => Ok(SessionFrame::Binary(payload)),
|
||||||
|
|
||||||
_ => {
|
_ => {
|
||||||
self.close().await.ok();
|
self.close().await.ok();
|
||||||
|
|||||||
Reference in New Issue
Block a user