Axum WebSocket Server Architecture for Real-Time CSI Streaming in RuView
RuView implements a dual-port Axum WebSocket server that uses a Tokio broadcast channel to push Channel State Information (CSI) updates to clients in real time, supporting both raw sensing data and pose-oriented streams on separate HTTP and dedicated WebSocket ports.
The RuView project (ruvnet/RuView) processes Wi-Fi Channel State Information to infer human pose and vital signs. Its sensing server uses Axum to stream this data via WebSockets, enabling browsers and external services to consume live CSI updates with minimal latency.
Shared Application State with Broadcast Channel
The server centers on a thread-safe application state that manages a Tokio broadcast channel. This pattern ensures every connected client receives identical JSON payloads without blocking the data producers.
In rust-port/wifi-densepose-rs/crates/wifi-densepose-sensing-server/src/main.rs at line 3562, the state is initialized:
struct AppStateInner {
// … many fields …
tx: broadcast::Sender<String>, // ← central broadcast channel
}
type SharedState = Arc<RwLock<AppStateInner>>;
The tx sender is cloned into a broadcast::Receiver for each new WebSocket client. Because broadcast::Sender is non-blocking, background CSI collection tasks continue running at full speed regardless of client count or connection quality.
Dual-Port Server Architecture
RuView splits its networking into two logical layers that share the same SharedState:
- HTTP API & UI – Serves static UI files, REST endpoints, and WebSocket connections on the HTTP port (configurable, default typically 8080)
- Dedicated WebSocket Server – Listens on a separate port (
--ws-port, default 8765) for raw sensing streams
Both routers mount the same ws_sensing_handler, allowing UI developers to connect via the same origin to avoid CORS issues, while external services can use the dedicated port.
The setup occurs at lines 3631-3642:
// WebSocket app (dedicated port)
let ws_state = state.clone();
let ws_app = Router::new()
.route("/ws/sensing", get(ws_sensing_handler))
.route("/health", get(health))
.with_state(ws_state);
let ws_addr = SocketAddr::from((bind_ip, args.ws_port));
let ws_listener = tokio::net::TcpListener::bind(ws_addr).await.unwrap();
tokio::spawn(async move { axum::serve(ws_listener, ws_app).await.unwrap(); });
The HTTP application router (lines 3634-3670) mirrors the WebSocket route:
let http_app = Router::new()
.route("/", get(info_page))
.route("/ws/sensing", get(ws_sensing_handler)) // same WS on HTTP port
// … many REST endpoints …
.with_state(state.clone());
let http_addr = SocketAddr::from((bind_ip, args.http_port));
axum::serve(tokio::net::TcpListener::bind(http_addr).await.unwrap(), http_app).await.unwrap();
WebSocket Handlers
Generic Sensing Stream
The primary handler ws_sensing_handler (lines 1492-1497) upgrades HTTP connections to WebSockets and delegates to handle_ws_client (lines 1499-1529):
async fn ws_sensing_handler(
ws: WebSocketUpgrade,
State(state): State<SharedState>,
) -> impl IntoResponse {
ws.on_upgrade(|socket| handle_ws_client(socket, state))
}
The client handler subscribes to the broadcast channel and uses tokio::select! to forward messages while monitoring for client disconnects:
async fn handle_ws_client(mut socket: WebSocket, state: SharedState) {
let mut rx = {
let s = state.read().await;
s.tx.subscribe()
};
loop {
tokio::select! {
msg = rx.recv() => {
if let Ok(json) = msg {
if socket.send(Message::Text(json.into())).await.is_err() { break; }
}
}
msg = socket.recv() => {
if matches!(msg, Some(Ok(Message::Close(_))) | None) { break; }
}
}
}
}
Pose-Oriented Stream
For UI compatibility with DensePose visualizations, the ws_pose_handler (lines 1533-1537) provides a specialized endpoint:
async fn ws_pose_handler(
ws: WebSocketUpgrade,
State(state): State<SharedState>,
) -> impl IntoResponse {
ws.on_upgrade(|socket| handle_ws_pose_client(socket, state))
}
The handle_ws_pose_client function (lines 1560-1645) transforms each SensingUpdate into pose-data messages, selecting between model-inferred keypoints or signal-derived keypoints based on availability.
Real-Time Data Production Flow
Background tasks—such as udp_receiver_task, windows_wifi_task, or simulated_data_task—generate SensingUpdate structs and push them to the broadcast channel. At line 1297, the serialization and broadcast occur:
let json = serde_json::to_string(&sensing_update).unwrap();
let _ = s.tx.send(json);
This non-blocking send means CSI packets arriving from ESP32 devices via UDP port 5005 (or other sources) are immediately fanned out to all WebSocket clients without backpressure.
Implementation Examples
Browser Client Connection
Connect to the dedicated WebSocket port from JavaScript:
const ws = new WebSocket('ws://localhost:8765/ws/sensing');
ws.onmessage = ev => {
const data = JSON.parse(ev.data);
console.log('CSI update:', data);
};
ws.onopen = () => console.log('Connected to RuView CSI stream');
ws.onclose = () => console.log('Disconnected');
Server Launch Configuration
Start the server with specific ports and an ESP32 CSI source:
cargo run --release -- \
--http-port 8080 \
--ws-port 8765 \
--bind-addr 0.0.0.0 \
--source esp32 \
--udp-port 5005
Adding Custom WebSocket Endpoints
Extend the router with additional handlers that reuse the shared broadcast channel:
async fn custom_ws_handler(
ws: WebSocketUpgrade,
State(state): State<SharedState>,
) -> impl IntoResponse {
ws.on_upgrade(|socket| async move {
let mut rx = { let s = state.read().await; s.tx.subscribe() };
while let Ok(msg) = rx.recv().await {
// filter or transform here
let _ = socket.send(Message::Text(msg)).await;
}
})
}
Mount it to the router:
ws_app = ws_app.route("/ws/custom", get(custom_ws_handler));
Summary
- Dual-port architecture separates UI traffic from high-volume CSI streams while sharing state via
SharedState - Broadcast channel pattern in
AppStateInnerenables non-blocking, multi-consumer real-time data distribution - Two WebSocket endpoints provide raw sensing data (
/ws/sensing) and UI-optimized pose streams via dedicated handlers - Axum WebSocket upgrades use
tokio::select!for concurrent message forwarding and connection lifecycle management - Source-agnostic design supports ESP32 UDP packets, Windows Wi-Fi scans, or simulated data feeding the same streaming infrastructure
Frequently Asked Questions
What ports does the RuView WebSocket server use?
By default, the HTTP server listens on port 8080 and the dedicated WebSocket server listens on port 8765. Both are configurable via --http-port and --ws-port command-line arguments. The WebSocket handler is mounted on /ws/sensing on both ports.
How does the broadcast channel prevent blocking?
The server uses tokio::sync::broadcast::Sender, which provides a non-blocking send() method. When background tasks push CSI updates, the sender immediately returns even if clients are slow or disconnected, ensuring the UDP receiver or Wi-Fi scanner never stalls.
What is the difference between the sensing and pose WebSocket endpoints?
The generic /ws/sensing endpoint streams raw SensingUpdate JSON as-is, suitable for data logging or custom processing. The pose-specific endpoint transforms these updates into DensePose-compatible keypoint formats, making it ideal for the bundled visualization UI.
Can the WebSocket be accessed from the browser UI?
Yes. Because the ws_sensing_handler is mounted on both the HTTP and dedicated WebSocket routers, browser-based UIs served from the HTTP port can connect to ws://localhost:8080/ws/sensing without triggering cross-origin restrictions, while external services use ws://localhost:8765/ws/sensing.
Have a question about this repo?
These articles cover the highlights, but your codebase questions are specific. Give your agent direct access to the source. Share this with your agent to get started:
curl -s "https://instagit.com/install.md" Maintain an open-source project? Get it listed too →