use std::sync::Arc;
use std::time::Duration;
use roder_api::events::{EventEnvelope, EventSinkFailed, RoderEvent};
use roder_api::extension::EventSink;
use time::OffsetDateTime;
use tokio::sync::mpsc;
use crate::bus::EventBus;
const QUEUE_DEPTH: usize = 256;
const HANDLE_TIMEOUT: Duration = Duration::from_secs(2);
const MAX_FAILURE_MESSAGE: usize = 300;
pub(crate) struct EventSinkDispatcher {
workers: Vec<Worker>,
}
struct Worker {
sink_id: String,
tx: mpsc::Sender<EventEnvelope>,
}
impl EventSinkDispatcher {
pub(crate) fn start(sinks: &[Arc<dyn EventSink>], bus: EventBus) -> Self {
let workers = sinks
.iter()
.map(|sink| {
let (tx, mut rx) = mpsc::channel::<EventEnvelope>(QUEUE_DEPTH);
let sink = sink.clone();
let sink_id = sink.id();
let worker_bus = bus.clone();
let worker_sink_id = sink_id.clone();
tokio::spawn(async move {
while let Some(envelope) = rx.recv().await {
let outcome =
tokio::time::timeout(HANDLE_TIMEOUT, sink.handle_event(&envelope))
.await;
let failure = match outcome {
Ok(Ok(())) => None,
Ok(Err(error)) => Some(truncate_message(&error.to_string())),
Err(_) => {
Some(format!("timed out after {}ms", HANDLE_TIMEOUT.as_millis()))
}
};
if let Some(message) = failure {
worker_bus.emit(RoderEvent::EventSinkFailed(EventSinkFailed {
sink_id: worker_sink_id.clone(),
event_kind: envelope.kind.clone(),
message,
timestamp: OffsetDateTime::now_utc(),
}));
}
}
});
Worker { sink_id, tx }
})
.collect();
Self { workers }
}
pub(crate) fn is_empty(&self) -> bool {
self.workers.is_empty()
}
pub(crate) fn dispatch(&self, envelope: &EventEnvelope, bus: &EventBus) {
if envelope.kind == "extension.event_sink_failed" {
return;
}
for worker in &self.workers {
if let Err(mpsc::error::TrySendError::Full(_)) = worker.tx.try_send(envelope.clone()) {
bus.emit(RoderEvent::EventSinkFailed(EventSinkFailed {
sink_id: worker.sink_id.clone(),
event_kind: envelope.kind.clone(),
message: format!("queue full ({QUEUE_DEPTH}); event dropped"),
timestamp: OffsetDateTime::now_utc(),
}));
}
}
}
}
fn truncate_message(message: &str) -> String {
if message.len() <= MAX_FAILURE_MESSAGE {
return message.to_string();
}
let mut end = MAX_FAILURE_MESSAGE;
while !message.is_char_boundary(end) {
end -= 1;
}
format!("{}…", &message[..end])
}