kerf 0.1.2

Simple tokio-based trace event collector
Documentation
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,
        }
    }

    // Main run loop that consumes self
    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();

        // Check each tab
        for (name, filter) in &self.tabs {
            // First check if any exclusion filters apply - these take precedence
            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 not silenced, check if any inclusion filters match
            if !is_silenced {
                for matcher in filter.iter_matchers() {
                    if matcher.include && matcher.matches(&event) {
                        captured_by.push(name.clone());
                        break;
                    }
                }
            }
        }

        // Record stats
        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);
            }
        }

        // Determine status and call callbacks
        if !captured_by.is_empty() {
            // First, try per-tab callbacks
            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 no per-tab callbacks handled it, use global callback
            if !handled_by_tab_callbacks {
                if let Some(cb) = &self.callback {
                    // Collect references to tab names
                    let tab_refs: Vec<&str> = captured_by.iter().map(String::as_str).collect();
                    cb(Arc::clone(&event), &tab_refs);
                } else {
                    // Store event for later processing when callback is set
                    self.pending_captured_events
                        .push((Arc::clone(&event), captured_by));
                }
            }
        } else if !silenced_by.is_empty() {
            if let Some(silenced_cb) = &self.silenced_callback {
                // Collect references to silencer names
                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) {
        // Set the new callback
        self.callback = Some(cb);

        // Drain pending captured events
        if let Some(callback) = &self.callback {
            for (event, tabs) in self.pending_captured_events.drain(..) {
                // Convert to references for the callback
                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,
    ) {
        // Set the new callback
        self.silenced_callback = Some(cb);
        response_tx.success();
    }

    fn handle_set_dropped_callback(&mut self, cb: DroppedEventCallback, response_tx: ResultSender) {
        // Set the new callback
        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");
        }
    }

    // Handle a dispatcher command
    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);
            }
        }
    }

    // Simplified add_tab handler
    fn handle_add_tab(
        &mut self,
        name: impl Into<String>,
        filter_set: MatcherSet,
        response_tx: ResultSender,
    ) {
        // Add the tab to the map
        self.tabs.insert(name.into(), filter_set);
        response_tx.success();
    }

    // Simplified update_tab handler
    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;
        }

        // Update the filter set
        self.tabs.insert(name.to_string(), filter_set);
        response_tx.success();
    }

    // Simplified remove_tab handler
    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;
        }

        // Remove the tab and its callback
        self.tabs.remove(name);
        self.tab_callbacks.remove(name);
        response_tx.success();
    }

    // Set callback for a specific tab
    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();
    }

    // Remove callback for a specific tab
    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();
    }
}

// Result sender for operation responses
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())));
    }
}