use std::collections::HashMap; use std::sync::Arc; use futures_util::StreamExt; use serde::Deserialize; use tokio::net::TcpStream; use tokio::sync::watch; use tokio::task::JoinHandle; use tokio_tungstenite::tungstenite::Message; use tokio_tungstenite::tungstenite::client::IntoClientRequest; use crate::AppState; use crate::db::{DatabaseBackend, adapt_sql, now_rfc3339}; use crate::event_log::{EventLog, Severity, log_event}; use crate::profile; // --------------------------------------------------------------------------- // Types // --------------------------------------------------------------------------- /// AT Protocol event stream frame header (DAG-CBOR). #[derive(Deserialize)] struct FrameHeader { op: i64, #[serde(default)] t: Option, } #[derive(Deserialize)] struct SubscribeLabelsMessage { seq: i64, labels: Vec