use chrono::{DateTime, Utc};
use serde::{Deserialize, Serialize};
use tokio::sync::broadcast;
#[derive(Debug, Clone, Serialize, Deserialize)]
#[serde(tag = "type", rename_all = "snake_case")]
pub enum UiEvent {
SessionCreated {
session_key: String,
agent_id: String,
channel: String,
peer_id: String,
},
SessionUpdated {
session_key: String,
update: SessionUpdate,
},
MessageReceived {
session_key: String,
content: String,
peer_id: String,
},
MessageSent {
session_key: String,
content: String,
},
ToolExecuted {
session_key: String,
tool: String,
result: serde_json::Value,
success: bool,
},
ChannelStatusChanged {
channel_id: String,
connected: bool,
error: Option<String>,
},
Heartbeat {
timestamp: DateTime<Utc>,
},
}
#[derive(Debug, Clone, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum SessionUpdate {
StateChanged {
new_state: String,
},
MessageCount {
count: u64,
},
Ended {
reason: String,
},
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct UiEventEnvelope {
pub id: String,
pub timestamp: DateTime<Utc>,
pub event: UiEvent,
}
impl UiEventEnvelope {
#[must_use]
pub fn new(event: UiEvent) -> Self {
use rand::RngCore;
let mut rng = rand::thread_rng();
let mut bytes = [0u8; 8];
rng.fill_bytes(&mut bytes);
Self {
id: hex::encode(bytes),
timestamp: Utc::now(),
event,
}
}
}
const DEFAULT_CHANNEL_CAPACITY: usize = 256;
pub struct EventBroadcaster {
sender: broadcast::Sender<UiEventEnvelope>,
}
impl EventBroadcaster {
#[must_use]
pub fn new() -> Self {
let (sender, _) = broadcast::channel(DEFAULT_CHANNEL_CAPACITY);
Self { sender }
}
#[must_use]
pub fn with_capacity(capacity: usize) -> Self {
let (sender, _) = broadcast::channel(capacity);
Self { sender }
}
#[must_use]
pub fn broadcast(&self, event: UiEvent) -> usize {
let envelope = UiEventEnvelope::new(event);
self.sender.send(envelope).unwrap_or(0)
}
#[must_use]
pub fn subscribe(&self) -> broadcast::Receiver<UiEventEnvelope> {
self.sender.subscribe()
}
#[must_use]
pub fn subscriber_count(&self) -> usize {
self.sender.receiver_count()
}
}
impl Default for EventBroadcaster {
fn default() -> Self {
Self::new()
}
}
impl Clone for EventBroadcaster {
fn clone(&self) -> Self {
Self {
sender: self.sender.clone(),
}
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn test_event_envelope() {
let event = UiEvent::Heartbeat {
timestamp: Utc::now(),
};
let envelope = UiEventEnvelope::new(event);
assert!(!envelope.id.is_empty());
assert_eq!(envelope.id.len(), 16); }
#[tokio::test]
async fn test_broadcaster() {
let broadcaster = EventBroadcaster::new();
let mut rx = broadcaster.subscribe();
assert_eq!(broadcaster.subscriber_count(), 1);
let event = UiEvent::Heartbeat {
timestamp: Utc::now(),
};
let count = broadcaster.broadcast(event);
assert_eq!(count, 1);
let received = rx.recv().await.unwrap();
match received.event {
UiEvent::Heartbeat { .. } => {}
_ => panic!("Wrong event type"),
}
}
#[test]
fn test_no_subscribers() {
let broadcaster = EventBroadcaster::new();
let count = broadcaster.broadcast(UiEvent::Heartbeat {
timestamp: Utc::now(),
});
assert_eq!(count, 0);
}
}