use std::time::SystemTime;
use tokio::sync::broadcast;
#[derive(Debug, Clone)]
pub struct EventSource {
sender: broadcast::Sender<BrokerEvent>,
}
#[derive(Debug, Clone)]
pub struct BrokerEvent {
pub at: SystemTime,
pub kind: BrokerEventKind,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum BrokerEventKind {
KeyRotated {
key_id: String,
new_version: u32,
},
BundleChanged {
trust_domain: String,
},
Revoked {
trust_domain: String,
id: String,
},
}
impl Default for EventSource {
fn default() -> Self {
Self::new()
}
}
impl EventSource {
#[must_use]
pub fn new() -> Self {
let (sender, _) = broadcast::channel(1024);
Self { sender }
}
#[must_use]
pub fn subscribe(&self) -> broadcast::Receiver<BrokerEvent> {
self.sender.subscribe()
}
pub fn key_rotated(&self, key_id: impl Into<String>, new_version: u32) {
self.publish(BrokerEventKind::KeyRotated {
key_id: key_id.into(),
new_version,
});
}
pub fn bundle_changed(&self, trust_domain: impl Into<String>) {
self.publish(BrokerEventKind::BundleChanged {
trust_domain: trust_domain.into(),
});
}
pub fn revoked(&self, trust_domain: impl Into<String>, id: impl Into<String>) {
self.publish(BrokerEventKind::Revoked {
trust_domain: trust_domain.into(),
id: id.into(),
});
}
fn publish(&self, kind: BrokerEventKind) {
let _receivers = self.sender.send(BrokerEvent {
at: SystemTime::now(),
kind,
});
}
}
#[cfg(test)]
mod tests {
use super::*;
#[tokio::test]
async fn subscribers_receive_events() {
let source = EventSource::new();
let mut rx = source.subscribe();
source.key_rotated("app.key", 2);
let event = rx.recv().await.expect("event received");
assert_eq!(
event.kind,
BrokerEventKind::KeyRotated {
key_id: "app.key".to_string(),
new_version: 2,
}
);
}
}