use helix_core::effect::DomainEventBytes;
use helix_core::ports::EventSink;
use helix_driver_host::{
AsyncMetricSink, LabelKey, MetricEvent, MetricId, MetricLabels, NoopMetricSink,
};
use std::sync::Arc;
use tokio::sync::broadcast;
const BROADCAST_CAPACITY: usize = 1024;
pub const BUS_CHANNEL: &str = "im:__bus__";
pub fn to_bus_envelope(ev: &DomainEventBytes) -> serde_json::Value {
let payload: serde_json::Value =
serde_json::from_slice(ev.0.as_ref()).unwrap_or(serde_json::Value::Null);
let channel = payload
.get("event")
.and_then(|e| e.as_str())
.unwrap_or("")
.to_string();
serde_json::json!({
"channel": channel,
"payload": payload,
})
}
pub fn bus_event_name(ev: &DomainEventBytes) -> Option<String> {
let payload: serde_json::Value = serde_json::from_slice(ev.0.as_ref()).ok()?;
payload
.get("event")
.and_then(|e| e.as_str())
.map(str::to_string)
}
#[derive(Debug)]
pub enum RecvOutcome {
Event(DomainEventBytes),
Lagged(u64),
Closed,
}
pub struct EventReceiver {
inner: broadcast::Receiver<DomainEventBytes>,
lagged_total: u64,
metrics: Arc<dyn AsyncMetricSink>,
consumer: &'static str,
}
impl EventReceiver {
pub fn new(inner: broadcast::Receiver<DomainEventBytes>) -> Self {
Self {
inner,
lagged_total: 0,
metrics: Arc::new(NoopMetricSink),
consumer: "unobserved",
}
}
fn with_metrics(mut self, metrics: Arc<dyn AsyncMetricSink>, consumer: &'static str) -> Self {
self.metrics = metrics;
self.consumer = consumer;
self
}
pub async fn recv(&mut self) -> RecvOutcome {
match self.inner.recv().await {
Ok(ev) => RecvOutcome::Event(ev),
Err(broadcast::error::RecvError::Lagged(n)) => {
self.lagged_total = self.lagged_total.saturating_add(n);
if self.metrics.is_enabled() {
let _ = self.metrics.try_record(MetricEvent::counter(
MetricId::EventLaggedTotal,
n as f64,
MetricLabels::one(LabelKey::Stage, "event")
.with(LabelKey::Consumer, self.consumer),
));
}
RecvOutcome::Lagged(n)
}
Err(broadcast::error::RecvError::Closed) => {
if self.metrics.is_enabled() {
let _ = self.metrics.try_record(MetricEvent::counter(
MetricId::EventConsumerClosedTotal,
1.0,
MetricLabels::one(LabelKey::Stage, "event")
.with(LabelKey::Consumer, self.consumer),
));
}
RecvOutcome::Closed
}
}
}
pub fn lagged_total(&self) -> u64 {
self.lagged_total
}
}
#[derive(Clone)]
pub struct NativeEventSink {
sender: broadcast::Sender<DomainEventBytes>,
metrics: Arc<dyn AsyncMetricSink>,
}
impl NativeEventSink {
pub fn new() -> (Self, broadcast::Receiver<DomainEventBytes>) {
Self::new_with_capacity(BROADCAST_CAPACITY)
}
pub fn new_with_capacity(capacity: usize) -> (Self, broadcast::Receiver<DomainEventBytes>) {
let (sender, receiver) = broadcast::channel(capacity);
(
Self {
sender,
metrics: Arc::new(NoopMetricSink),
},
receiver,
)
}
pub fn with_metric_sink(mut self, metrics: Arc<dyn AsyncMetricSink>) -> Self {
self.metrics = metrics;
self
}
pub fn subscribe(&self) -> broadcast::Receiver<DomainEventBytes> {
self.sender.subscribe()
}
pub fn subscribe_observed(&self) -> EventReceiver {
self.subscribe_observed_as("native")
}
pub fn subscribe_observed_as(&self, consumer: &'static str) -> EventReceiver {
let receiver = EventReceiver::new(self.sender.subscribe())
.with_metrics(Arc::clone(&self.metrics), consumer);
if self.metrics.is_enabled() {
let _ = self.metrics.try_record(MetricEvent::gauge(
MetricId::EventReceiverCount,
self.sender.receiver_count() as f64,
MetricLabels::one(LabelKey::Stage, "event").with(LabelKey::Consumer, consumer),
));
}
receiver
}
pub fn receiver_count(&self) -> usize {
self.sender.receiver_count()
}
}
impl EventSink for NativeEventSink {
fn emit(&self, event: DomainEventBytes) {
if self.sender.send(event).is_err() && self.metrics.is_enabled() {
let _ = self.metrics.try_record(MetricEvent::counter(
MetricId::EventNoReceiverTotal,
1.0,
MetricLabels::one(LabelKey::Stage, "event"),
));
}
}
}
impl helix_driver_host::BatchSink for NativeEventSink {}
#[cfg(test)]
#[path = "event_sink_tests.rs"]
mod tests;