use monoloop_contracts::InterpreterOutputEvent;
use monoloop_loop::{CanonicalEventSubscription, SubscriptionPublisher, SubscriptionStatus};
use std::sync::Arc;
use tokio::sync::mpsc;
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub enum SubscriberPolicy {
Lossless,
BestEffort,
}
struct Sub {
name: String,
policy: SubscriberPolicy,
pub_: SubscriptionPublisher,
}
pub struct EventDistributor {
subs: Vec<Sub>,
}
impl EventDistributor {
pub fn new() -> Self {
Self { subs: Vec::new() }
}
pub fn subscribe(
&mut self,
name: impl Into<String>,
policy: SubscriberPolicy,
capacity: usize,
) -> CanonicalEventSubscription {
let name = name.into();
let (pub_, sub) = SubscriptionPublisher::channel(name.clone(), capacity);
self.subs.push(Sub { name, policy, pub_ });
sub
}
pub async fn publish(&self, event: InterpreterOutputEvent) {
for sub in &self.subs {
match sub.policy {
SubscriberPolicy::Lossless => {
let _ = sub.pub_.publish(event.clone()).await;
}
SubscriberPolicy::BestEffort => {
let _ = sub.pub_.publish(event.clone()).await;
}
}
}
}
pub fn close(self) {
drop(self.subs);
}
}
impl Default for EventDistributor {
fn default() -> Self {
Self::new()
}
}
pub async fn pump_interpreter_to_distributor(
events: Arc<monoloop_interpreter::CanonicalEventStream>,
distributor: EventDistributor,
) {
loop {
match events.recv().await {
Some(ev) => {
let done = matches!(ev, InterpreterOutputEvent::Ended(_));
distributor.publish(ev).await;
if done {
break;
}
}
None => break,
}
}
distributor.close();
}
pub type StatusTx = mpsc::Sender<SubscriptionStatus>;