use crate::actor::{Actor, ActorContext};
use crate::message::Message;
use futures_util::SinkExt;
use futures_util::stream::{SplitSink, SplitStream};
use async_trait::async_trait;
use log::{debug, error, info};
use web_time::Duration;
use tokio_tungstenite::tungstenite::Message as WsMessage;
use tokio_tungstenite::{MaybeTlsStream, WebSocketStream};
type WsStream = WebSocketStream<MaybeTlsStream<tokio::net::TcpStream>>;
type WsSender = SplitSink<WsStream, WsMessage>;
type WsReceiver = SplitStream<WsStream>;
pub struct WsConn {
sender: WsSender,
receiver: Option<WsReceiver>,
allow_public_space: bool,
}
impl WsConn {
pub fn new(sender: WsSender, receiver: WsReceiver, allow_public_space: bool) -> Self {
Self {
sender,
receiver: Some(receiver),
allow_public_space,
}
}
}
#[async_trait]
impl Actor for WsConn {
async fn handle(&mut self, mut msg: Message, _ctx: &ActorContext) {
if let Message::Put(ref mut put) = msg {
put.json_str = None;
}
let wire = msg.to_string();
debug!("[WS→] SENDING {} bytes: {}", wire.len(), &wire[..wire.len().min(300)]);
let _ = self
.sender
.send(WsMessage::Text(wire.into()))
.await;
}
async fn pre_start(&mut self, ctx: &ActorContext) {
info!("WsConn starting");
let hi = Message::Hi {
from: ctx.addr.clone(),
peer_id: ctx.peer_id.read().clone(),
};
let _ = self
.sender
.send(WsMessage::Text(hi.to_string().into()))
.await;
let receiver = self.receiver.take().unwrap();
let mut ctx2 = ctx.clone();
let allow_public_space = self.allow_public_space;
ctx.child_task(async move {
use futures_util::StreamExt;
let mut receiver = receiver;
while let Some(result) = receiver.next().await {
let ws_msg = match result {
Ok(m) => m,
Err(e) => {
debug!("[WS] recv error: {}", e);
break;
}
};
let text = match ws_msg {
WsMessage::Text(t) => t,
WsMessage::Binary(_) => { debug!("[WS] binary frame (ignored)"); continue; }
WsMessage::Ping(_) => { debug!("[WS] ping frame (ignored)"); continue; }
WsMessage::Pong(_) => { debug!("[WS] pong frame (ignored)"); continue; }
WsMessage::Close(_) => { debug!("[WS] close frame received from peer"); break; }
WsMessage::Frame(_) => { debug!("[WS] raw frame (ignored)"); continue; }
};
if text.is_empty() {
debug!("[WS] empty text frame (ignored)");
continue;
}
debug!("[WS←] RECV {} bytes: {}", text.len(), &text[..text.len().min(300)]);
match Message::try_from(&text, ctx2.addr.clone(), allow_public_space) {
Ok(msgs) => {
for msg in msgs {
if ctx2.router.send(msg).is_err() {
error!("failed to forward incoming message to router");
}
}
}
Err(e) => {
debug!("[WS] parse error: {} (len={})", e, text.len());
}
}
}
debug!("[WS] receive loop ended — stopping actor");
ctx2.stop();
});
}
async fn stopping(&mut self, _context: &ActorContext) {
info!("WsConn stopping — sending WebSocket Close frame");
let close_result = crate::tokio_time::timeout(Duration::from_secs(2), self.sender.close()).await;
match close_result {
Ok(Ok(())) => debug!("WsConn Close frame acknowledged"),
Ok(Err(e)) => debug!("WsConn Close error (non-fatal): {}", e),
Err(_) => debug!("WsConn Close timed out — connection dropped"),
}
}
}