use anyhow::{Result, anyhow};
use std::{collections::HashMap, sync::Arc};
use tokio::sync::{mpsc, oneshot};
use crate::{
ArcEvent, DroppedEventCallback, EventCallback, EventType, MatcherSet, SilencedEventCallback,
StatsSnapshot, StatsTracker,
};
pub type TabCallback = Arc<dyn Fn(ArcEvent) + Send + Sync>;
pub(crate) enum DispatcherCommand {
SetCallback(EventCallback, ResultSender),
SetSilencedCallback(SilencedEventCallback, ResultSender),
SetUncapturedCallback(DroppedEventCallback, ResultSender),
AddTab(String, MatcherSet, ResultSender),
UpdateTab(String, MatcherSet, ResultSender),
RemoveTab(String, ResultSender),
SetTabCallback(String, TabCallback, ResultSender),
RemoveTabCallback(String, ResultSender),
GetStats(oneshot::Sender<Option<StatsSnapshot>>),
ClearStats(ResultSender),
}
pub(crate) struct Dispatcher {
event_rx: mpsc::UnboundedReceiver<ArcEvent>,
command_rx: mpsc::UnboundedReceiver<DispatcherCommand>,
tabs: HashMap<String, MatcherSet>,
tab_callbacks: HashMap<String, TabCallback>,
callback: Option<EventCallback>,
silenced_callback: Option<SilencedEventCallback>,
dropped_callback: Option<DroppedEventCallback>,
pending_captured_events: Vec<(ArcEvent, Vec<String>)>,
stats: Option<Arc<StatsTracker>>,
}
impl Dispatcher {
pub fn new(
event_rx: mpsc::UnboundedReceiver<ArcEvent>,
command_rx: mpsc::UnboundedReceiver<DispatcherCommand>,
initial_tabs: HashMap<String, MatcherSet>,
stats_config: Option<crate::StatsConfig>,
) -> Self {
let stats = stats_config.map(|config| Arc::new(StatsTracker::new(config)));
Self {
event_rx,
command_rx,
tabs: initial_tabs,
tab_callbacks: HashMap::new(),
callback: None,
silenced_callback: None,
dropped_callback: None,
pending_captured_events: Vec::new(),
stats,
}
}
pub async fn run(mut self) {
loop {
tokio::select! {
Some(event) = self.event_rx.recv() => {
self.handle_event(event);
},
Some(cmd) = self.command_rx.recv() => {
self.handle_command(cmd).await;
},
else => break,
}
}
}
fn handle_event(&mut self, event: ArcEvent) {
let mut captured_by = Vec::new();
let mut silenced_by = Vec::new();
for (name, filter) in &self.tabs {
let mut is_silenced = false;
for matcher in filter.iter_matchers() {
if !matcher.include && matcher.matches(&event) {
silenced_by.push(name.clone());
is_silenced = true;
break;
}
}
if !is_silenced {
for matcher in filter.iter_matchers() {
if matcher.include && matcher.matches(&event) {
captured_by.push(name.clone());
break;
}
}
}
}
if let Some(stats) = &self.stats {
if !captured_by.is_empty() {
stats.record_event(EventType::Captured, &event);
} else if !silenced_by.is_empty() {
stats.record_event(EventType::Silenced, &event);
} else {
stats.record_event(EventType::Dropped, &event);
}
}
if !captured_by.is_empty() {
let mut handled_by_tab_callbacks = false;
for tab_name in &captured_by {
if let Some(tab_callback) = self.tab_callbacks.get(tab_name) {
tab_callback(Arc::clone(&event));
handled_by_tab_callbacks = true;
}
}
if !handled_by_tab_callbacks {
if let Some(cb) = &self.callback {
let tab_refs: Vec<&str> = captured_by.iter().map(String::as_str).collect();
cb(Arc::clone(&event), &tab_refs);
} else {
self.pending_captured_events
.push((Arc::clone(&event), captured_by));
}
}
} else if !silenced_by.is_empty() {
if let Some(silenced_cb) = &self.silenced_callback {
let silencer_refs: Vec<&str> = silenced_by.iter().map(String::as_str).collect();
silenced_cb(Arc::clone(&event), &silencer_refs);
}
} else if let Some(dropped_cb) = &self.dropped_callback {
dropped_cb(Arc::clone(&event));
}
}
fn handle_set_callback(&mut self, cb: EventCallback, response_tx: ResultSender) {
self.callback = Some(cb);
if let Some(callback) = &self.callback {
for (event, tabs) in self.pending_captured_events.drain(..) {
let tab_refs: Vec<&str> = tabs.iter().map(String::as_str).collect();
callback(event, &tab_refs);
}
}
response_tx.success();
}
fn handle_set_silenced_callback(
&mut self,
cb: SilencedEventCallback,
response_tx: ResultSender,
) {
self.silenced_callback = Some(cb);
response_tx.success();
}
fn handle_set_dropped_callback(&mut self, cb: DroppedEventCallback, response_tx: ResultSender) {
self.dropped_callback = Some(cb);
response_tx.success();
}
fn handle_get_stats(&self, response_tx: oneshot::Sender<Option<StatsSnapshot>>) {
let snapshot = self.stats.as_ref().map(|stats| stats.get_snapshot());
let _ = response_tx.send(snapshot);
}
fn handle_clear_stats(&self, response_tx: ResultSender) {
if let Some(stats) = &self.stats {
stats.clear();
response_tx.success();
} else {
response_tx.error("Stats not enabled");
}
}
async fn handle_command(&mut self, cmd: DispatcherCommand) {
match cmd {
DispatcherCommand::SetCallback(cb, response_tx) => {
self.handle_set_callback(cb, response_tx);
}
DispatcherCommand::SetSilencedCallback(cb, response_tx) => {
self.handle_set_silenced_callback(cb, response_tx);
}
DispatcherCommand::SetUncapturedCallback(cb, response_tx) => {
self.handle_set_dropped_callback(cb, response_tx);
}
DispatcherCommand::AddTab(name, filter_set, response_tx) => {
self.handle_add_tab(name, filter_set, response_tx);
}
DispatcherCommand::UpdateTab(name, filter_set, response_tx) => {
self.handle_update_tab(name, filter_set, response_tx);
}
DispatcherCommand::RemoveTab(name, response_tx) => {
self.handle_remove_tab(name, response_tx);
}
DispatcherCommand::SetTabCallback(name, callback, response_tx) => {
self.handle_set_tab_callback(name, callback, response_tx);
}
DispatcherCommand::RemoveTabCallback(name, response_tx) => {
self.handle_remove_tab_callback(name, response_tx);
}
DispatcherCommand::GetStats(response_tx) => {
self.handle_get_stats(response_tx);
}
DispatcherCommand::ClearStats(response_tx) => {
self.handle_clear_stats(response_tx);
}
}
}
fn handle_add_tab(
&mut self,
name: impl Into<String>,
filter_set: MatcherSet,
response_tx: ResultSender,
) {
self.tabs.insert(name.into(), filter_set);
response_tx.success();
}
fn handle_update_tab(
&mut self,
name: impl AsRef<str>,
filter_set: MatcherSet,
response_tx: ResultSender,
) {
let name = name.as_ref();
if !self.tabs.contains_key(name) {
response_tx.error(format!("Subscriber '{name:?}' not found"));
return;
}
self.tabs.insert(name.to_string(), filter_set);
response_tx.success();
}
fn handle_remove_tab(&mut self, name: impl AsRef<str>, response_tx: ResultSender) {
let name = name.as_ref();
if !self.tabs.contains_key(name) {
response_tx.error(format!("Subscriber '{name:?}' not found"));
return;
}
self.tabs.remove(name);
self.tab_callbacks.remove(name);
response_tx.success();
}
fn handle_set_tab_callback(
&mut self,
name: String,
callback: TabCallback,
response_tx: ResultSender,
) {
if !self.tabs.contains_key(&name) {
response_tx.error(format!("Tab '{name}' not found"));
return;
}
self.tab_callbacks.insert(name, callback);
response_tx.success();
}
fn handle_remove_tab_callback(&mut self, name: String, response_tx: ResultSender) {
if !self.tabs.contains_key(&name) {
response_tx.error(format!("Tab '{name}' not found"));
return;
}
self.tab_callbacks.remove(&name);
response_tx.success();
}
}
pub struct ResultSender(oneshot::Sender<Result<()>>);
impl ResultSender {
pub fn new(sender: oneshot::Sender<Result<()>>) -> Self {
Self(sender)
}
pub fn success(self) {
let _ = self.0.send(Ok(()));
}
pub fn error(self, msg: impl Into<String>) {
let _ = self.0.send(Err(anyhow!(msg.into())));
}
}