dahua-camera-server 0.3.1

axum HTTP, SSE and WebSocket API for the dahua-camera stack
Documentation
//! `GET /ws/events` — the merged analytics-event feed, as JSON text frames.

use crate::AppState;
use axum::extract::ws::{Message, WebSocket, WebSocketUpgrade};
use axum::extract::State;
use axum::response::IntoResponse;
use tokio::sync::broadcast::error::RecvError;

/// Upgrade a request to the event socket.
pub async fn events_handler(
    State(state): State<AppState>,
    upgrade: WebSocketUpgrade,
) -> impl IntoResponse {
    upgrade.on_upgrade(move |socket| pump(socket, state))
}

async fn pump(mut socket: WebSocket, state: AppState) {
    let mut events = state.subscribe_events();

    loop {
        tokio::select! {
            incoming = socket.recv() => {
                match incoming {
                    None | Some(Err(_)) | Some(Ok(Message::Close(_))) => break,
                    _ => continue,
                }
            }
            received = events.recv() => {
                match received {
                    Ok(event) => {
                        let Ok(json) = serde_json::to_string(&event) else { continue };
                        if socket.send(Message::Text(json.into())).await.is_err() {
                            break;
                        }
                    }
                    Err(RecvError::Closed) => break,
                    // Events are advisory; a client that fell behind during a
                    // burst simply misses those and carries on.
                    Err(RecvError::Lagged(dropped)) => {
                        tracing::debug!(dropped, "event client lagged");
                        continue;
                    }
                }
            }
        }
    }
}