use dahua_camera_rtsp::CameraService;
use std::sync::Arc;
pub type AppState = Arc<AppStateInner>;
pub struct AppStateInner {
pub service: Arc<CameraService>,
#[cfg(feature = "alarms")]
pub alarms: Arc<std::sync::Mutex<dahua_camera_alarms::AlarmRegistry>>,
}
impl AppStateInner {
pub fn new(service: Arc<CameraService>) -> AppState {
#[cfg(feature = "alarms")]
let alarms = {
let registry = Arc::new(std::sync::Mutex::new(
dahua_camera_alarms::AlarmRegistry::new(Default::default()),
));
spawn_alarm_aggregator(service.clone(), registry.clone());
registry
};
Arc::new(Self {
service,
#[cfg(feature = "alarms")]
alarms,
})
}
}
#[cfg(feature = "alarms")]
fn spawn_alarm_aggregator(
service: Arc<CameraService>,
registry: Arc<std::sync::Mutex<dahua_camera_alarms::AlarmRegistry>>,
) {
let mut events = service.subscribe_events();
tokio::spawn(async move {
loop {
match events.recv().await {
Ok(event) => {
let now = now_unix_ms();
if let Ok(mut registry) = registry.lock() {
registry.observe(&event, now);
}
}
Err(tokio::sync::broadcast::error::RecvError::Lagged(skipped)) => {
tracing::warn!(
skipped,
"alarm aggregator lagged; any alarm whose STOP was dropped \
will auto-clear on its timeout"
);
}
Err(tokio::sync::broadcast::error::RecvError::Closed) => break,
}
}
});
}
#[cfg(feature = "alarms")]
pub(crate) fn now_unix_ms() -> u64 {
std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap_or_default()
.as_millis() as u64
}
impl std::ops::Deref for AppStateInner {
type Target = CameraService;
fn deref(&self) -> &CameraService {
&self.service
}
}