use msgtrans::{
event::ServerEvent, packet::Packet, protocol::WebSocketServerConfig, tokio,
transport::TransportServerBuilder,
};
use std::env;
const BIZ_DROP: u8 = 250;
const DEFAULT_PORT: u16 = 18080;
#[tokio::main]
async fn main() -> Result<(), Box<dyn std::error::Error>> {
let port: u16 = env::var("MSGTRANS_E2E_PORT")
.ok()
.and_then(|s| s.parse().ok())
.unwrap_or(DEFAULT_PORT);
let bind = format!("127.0.0.1:{}", port);
let ws_config = WebSocketServerConfig::new(&bind)?;
let transport = TransportServerBuilder::new()
.max_connections(64)
.with_protocol(ws_config)
.build()
.await?;
let mut events = transport.subscribe_events();
let transport_for_echo = transport.clone();
let transport_for_serve = transport.clone();
let event_task = tokio::spawn(async move {
let transport = transport_for_echo;
while let Ok(event) = events.recv().await {
if let ServerEvent::MessageReceived {
session_id,
context,
} = event
{
let biz_type = context.biz_type;
let payload = context.data.clone();
if biz_type == BIZ_DROP {
continue;
}
if context.is_request() {
context.respond(payload);
} else {
let transport_clone = transport.clone();
tokio::spawn(async move {
let mut packet = Packet::one_way(0, payload);
packet.set_biz_type(biz_type);
let _ = transport_clone.send_to_session(session_id, packet).await;
});
}
}
}
});
println!("MSGTRANS_E2E_READY ws://{}", bind);
let serve_result = transport_for_serve.serve().await;
let _ = event_task.await;
serve_result?;
Ok(())
}