use crate::support::io::sink::file::FileSink;
use crate::support::io::sink::LogSink;
use crate::FileSinkConfig;
use crate::InklogError;
use async_trait::async_trait;
use std::collections::HashMap;
use std::sync::Arc;
use tracing::info;
#[async_trait]
pub trait SinkFactory: Send + Sync {
async fn create(&self) -> Result<Arc<dyn LogSink>, InklogError>;
fn sink_type(&self) -> &'static str;
fn metadata(&self) -> SinkMetadata;
}
#[derive(Debug, Clone)]
pub struct SinkMetadata {
pub name: String,
pub description: String,
pub features: Vec<String>,
pub config_schema: Option<serde_json::Value>,
}
pub struct SinkRegistry {
factories: HashMap<String, Box<dyn SinkFactory>>,
}
impl Default for SinkRegistry {
fn default() -> Self {
Self::new()
}
}
impl SinkRegistry {
pub fn new() -> Self {
Self {
factories: HashMap::new(),
}
}
pub fn register<F: SinkFactory + 'static>(&mut self, factory: F) {
let sink_type = factory.sink_type().to_string();
info!("Registering sink factory: {}", sink_type);
self.factories.insert(sink_type, Box::new(factory));
}
pub async fn create(&self, sink_type: &str) -> Result<Arc<dyn LogSink>, InklogError> {
let factory = self
.factories
.get(sink_type)
.ok_or_else(|| InklogError::ConfigError(format!("Unknown sink type: {}", sink_type)))?;
factory.create().await
}
pub fn list_sinks(&self) -> Vec<&str> {
self.factories.keys().map(|s| s.as_str()).collect()
}
pub fn get_metadata(&self, sink_type: &str) -> Option<SinkMetadata> {
self.factories.get(sink_type).map(|f| f.metadata())
}
pub fn has_sink(&self, sink_type: &str) -> bool {
self.factories.contains_key(sink_type)
}
pub fn unregister(&mut self, sink_type: &str) -> Option<Box<dyn SinkFactory>> {
self.factories.remove(sink_type)
}
pub fn clear(&mut self) {
self.factories.clear();
}
}
pub struct FileSinkFactory {
config: FileSinkConfig,
}
impl FileSinkFactory {
pub fn new(config: FileSinkConfig) -> Self {
Self { config }
}
}
#[async_trait]
impl SinkFactory for FileSinkFactory {
async fn create(&self) -> Result<Arc<dyn LogSink>, InklogError> {
let sink = FileSink::new(self.config.clone())?;
Ok(Arc::new(sink))
}
fn sink_type(&self) -> &'static str {
"file"
}
fn metadata(&self) -> SinkMetadata {
SinkMetadata {
name: "File Sink".to_string(),
description: "Writes logs to files with rotation, compression, and encryption support."
.to_string(),
features: vec![
"rotation".to_string(),
"compression".to_string(),
"encryption".to_string(),
"batching".to_string(),
],
config_schema: None,
}
}
}
#[cfg(test)]
mod tests {
use super::*;
use tempfile::tempdir;
#[test]
fn test_registry_registration() {
let mut registry = SinkRegistry::new();
let temp_dir = tempdir().unwrap();
let config = FileSinkConfig {
enabled: true,
path: temp_dir.path().join("test.log"),
..Default::default()
};
let factory = FileSinkFactory::new(config);
registry.register(factory);
assert!(registry.has_sink("file"));
assert!(!registry.has_sink("nonexistent"));
}
#[tokio::test]
async fn test_registry_create() {
let mut registry = SinkRegistry::new();
let temp_dir = tempdir().unwrap();
let config = FileSinkConfig {
enabled: true,
path: temp_dir.path().join("test.log"),
..Default::default()
};
let factory = FileSinkFactory::new(config);
registry.register(factory);
let sink = registry.create("file").await;
assert!(sink.is_ok());
let nonexistent = registry.create("nonexistent").await;
assert!(nonexistent.is_err());
}
#[test]
fn test_registry_list_sinks() {
let mut registry = SinkRegistry::new();
let temp_dir = tempdir().unwrap();
let config = FileSinkConfig {
enabled: true,
path: temp_dir.path().join("test.log"),
..Default::default()
};
let factory = FileSinkFactory::new(config);
registry.register(factory);
let sinks = registry.list_sinks();
assert_eq!(sinks.len(), 1);
assert!(sinks.contains(&"file"));
}
#[test]
fn test_registry_metadata() {
let mut registry = SinkRegistry::new();
let temp_dir = tempdir().unwrap();
let config = FileSinkConfig {
enabled: true,
path: temp_dir.path().join("test.log"),
..Default::default()
};
let factory = FileSinkFactory::new(config);
registry.register(factory);
let metadata = registry.get_metadata("file");
assert!(metadata.is_some());
let metadata = metadata.unwrap();
assert_eq!(metadata.name, "File Sink");
assert!(metadata.features.contains(&"rotation".to_string()));
}
#[test]
fn test_registry_unregister() {
let mut registry = SinkRegistry::new();
let temp_dir = tempdir().unwrap();
let config = FileSinkConfig {
enabled: true,
path: temp_dir.path().join("test.log"),
..Default::default()
};
let factory = FileSinkFactory::new(config);
registry.register(factory);
assert!(registry.has_sink("file"));
let removed = registry.unregister("file");
assert!(removed.is_some());
assert!(!registry.has_sink("file"));
}
#[test]
fn test_registry_default() {
let registry = SinkRegistry::default();
assert_eq!(registry.list_sinks().len(), 0);
assert!(!registry.has_sink("file"));
}
#[test]
fn test_registry_clear() {
let mut registry = SinkRegistry::new();
let temp_dir = tempdir().unwrap();
let config = FileSinkConfig {
enabled: true,
path: temp_dir.path().join("test1.log"),
..Default::default()
};
registry.register(FileSinkFactory::new(config));
assert_eq!(registry.list_sinks().len(), 1);
registry.clear();
assert_eq!(registry.list_sinks().len(), 0);
assert!(!registry.has_sink("file"));
}
#[test]
fn test_registry_unregister_nonexistent() {
let mut registry = SinkRegistry::new();
let removed = registry.unregister("nonexistent");
assert!(removed.is_none());
}
#[test]
fn test_registry_get_metadata_nonexistent() {
let registry = SinkRegistry::new();
assert!(registry.get_metadata("nonexistent").is_none());
}
#[tokio::test]
async fn test_registry_create_after_unregister() {
let mut registry = SinkRegistry::new();
let temp_dir = tempdir().unwrap();
let config = FileSinkConfig {
enabled: true,
path: temp_dir.path().join("test.log"),
..Default::default()
};
registry.register(FileSinkFactory::new(config));
let _ = registry.unregister("file");
let result = registry.create("file").await;
assert!(result.is_err());
}
#[test]
fn test_file_sink_factory_metadata() {
let temp_dir = tempdir().unwrap();
let config = FileSinkConfig {
enabled: true,
path: temp_dir.path().join("test.log"),
..Default::default()
};
let factory = FileSinkFactory::new(config);
let metadata = factory.metadata();
assert_eq!(metadata.name, "File Sink");
assert!(metadata.description.contains("rotation"));
assert!(metadata.features.contains(&"rotation".to_string()));
assert!(metadata.features.contains(&"compression".to_string()));
assert!(metadata.features.contains(&"encryption".to_string()));
assert!(metadata.features.contains(&"batching".to_string()));
assert!(metadata.config_schema.is_none());
}
}