use std::path::PathBuf;
use std::sync::{Arc, Mutex};
use async_trait::async_trait;
use tokio::io::AsyncWriteExt;
use crate::error::TelemetryError;
use crate::event::Event;
#[async_trait]
pub trait TelemetrySink: Send + Sync + 'static {
async fn emit(&self, event: &Event) -> Result<(), TelemetryError>;
async fn flush(&self) -> Result<(), TelemetryError> {
Ok(())
}
}
#[derive(Debug, Default, Clone, Copy)]
pub struct NoopSink;
#[async_trait]
impl TelemetrySink for NoopSink {
async fn emit(&self, _event: &Event) -> Result<(), TelemetryError> {
Ok(())
}
}
#[derive(Debug, Default, Clone)]
pub struct MemorySink {
inner: Arc<Mutex<Vec<Event>>>,
}
impl MemorySink {
#[must_use]
pub fn new() -> Self {
Self::default()
}
#[must_use]
pub fn snapshot(&self) -> Vec<Event> {
self.inner.lock().map(|v| v.clone()).unwrap_or_default()
}
#[must_use]
pub fn len(&self) -> usize {
self.inner.lock().map_or(0, |v| v.len())
}
#[must_use]
pub fn is_empty(&self) -> bool {
self.len() == 0
}
}
#[async_trait]
impl TelemetrySink for MemorySink {
async fn emit(&self, event: &Event) -> Result<(), TelemetryError> {
if let Ok(mut v) = self.inner.lock() {
v.push(event.clone());
}
Ok(())
}
}
#[derive(Debug, Clone)]
pub struct FileSink {
path: PathBuf,
gate: Arc<tokio::sync::Mutex<()>>,
}
impl FileSink {
#[must_use]
pub fn new(path: impl Into<PathBuf>) -> Self {
Self { path: path.into(), gate: Arc::new(tokio::sync::Mutex::new(())) }
}
}
#[async_trait]
impl TelemetrySink for FileSink {
async fn emit(&self, event: &Event) -> Result<(), TelemetryError> {
let redacted = event.redacted();
let mut line =
serde_json::to_string(&redacted).map_err(|e| TelemetryError::Serde(e.to_string()))?;
line.push('\n');
let _guard = self.gate.lock().await;
if let Some(parent) = self.path.parent() {
if !parent.as_os_str().is_empty() {
tokio::fs::create_dir_all(parent).await?;
}
}
let mut f =
tokio::fs::OpenOptions::new().create(true).append(true).open(&self.path).await?;
f.write_all(line.as_bytes()).await?;
f.flush().await?;
Ok(())
}
}