use std::sync::Arc;
use anyhow::Result;
use async_trait::async_trait;
use tokio::sync::RwLock;
use zelos_trace_types::ipc::{IpcMessageWithId, Receiver, Sender};
use crate::{filter::Filter, router::DEFAULT_CHANNEL_SIZE};
#[async_trait]
pub(crate) trait TraceSinkHandle: Send + Sync {
async fn send_async(&self, msg: &IpcMessageWithId) -> Result<()>;
}
pub(crate) struct TraceSinkHandleFiltered {
pub sender: Sender,
pub filters: Arc<RwLock<Vec<Filter>>>,
}
#[async_trait]
impl TraceSinkHandle for TraceSinkHandleFiltered {
async fn send_async(&self, msg: &IpcMessageWithId) -> Result<()> {
for filter in self.filters.read().await.iter() {
if filter.matches(msg) {
self.sender.try_send(msg.clone())?;
continue;
}
}
Ok(())
}
}
pub(crate) struct TraceSinkHandleAllBlocking {
pub sender: Sender,
}
impl TraceSinkHandleAllBlocking {
pub(crate) fn new() -> (Self, Receiver) {
let (sender, receiver) = flume::bounded::<IpcMessageWithId>(DEFAULT_CHANNEL_SIZE);
(Self { sender }, receiver)
}
}
#[async_trait]
impl TraceSinkHandle for TraceSinkHandleAllBlocking {
async fn send_async(&self, msg: &IpcMessageWithId) -> Result<()> {
self.sender.send_async(msg.clone()).await?;
Ok(())
}
}
#[derive(Debug)]
pub struct TraceSink {
filters: Arc<RwLock<Vec<Filter>>>,
}
impl TraceSink {
pub(crate) fn new() -> (Self, Receiver, TraceSinkHandleFiltered) {
let (sender, receiver) = flume::bounded::<IpcMessageWithId>(1024);
let filters = Arc::new(RwLock::new(Vec::new()));
(
Self {
filters: filters.clone(),
},
receiver,
TraceSinkHandleFiltered { sender, filters },
)
}
pub async fn subscribe(&self, filter: Filter) {
let mut filters = self.filters.write().await;
filters.push(filter);
}
pub async fn unsubscribe(&self, filter: Filter) {
let mut filters = self.filters.write().await;
filters.retain(|f| f != &filter);
}
}