use std::fs;
use std::path::Path;
use std::sync::Arc;
use std::time::Duration;
use serde_json::Value;
use tap_http::event::{EventLoggerConfig, HttpEvent, LogDestination};
use tap_http::{TapHttpConfig, TapHttpServer};
use tap_node::{NodeConfig, TapNode};
use tempfile::tempdir;
use tokio::sync::Mutex;
use tokio::time::sleep;
#[tokio::test(flavor = "multi_thread")]
async fn test_event_logging_config() {
let temp_dir = tempdir().unwrap();
let log_path = temp_dir.path().join("test-events.log");
let log_path_str = log_path.to_str().unwrap().to_string();
let config = TapHttpConfig {
event_logger: Some(EventLoggerConfig {
destination: LogDestination::File {
path: log_path_str.clone(),
max_size: None,
rotate: false,
},
structured: true,
log_level: tracing::Level::INFO,
}),
..Default::default()
};
let node_config = NodeConfig {
storage_path: None,
..Default::default()
};
let node = TapNode::new(node_config);
let server = TapHttpServer::new(config, node);
assert!(server.event_bus().subscriber_count() > 0);
}
#[tokio::test(flavor = "multi_thread")]
async fn test_server_events() {
struct TestSubscriber {
events: Arc<Mutex<Vec<HttpEvent>>>,
}
#[async_trait::async_trait]
impl tap_http::event::HandleEvent<'_> for TestSubscriber {
async fn handle_event_async(&self, event: HttpEvent) {
self.events.lock().await.push(event);
}
}
let node_config = NodeConfig {
storage_path: None,
..Default::default()
};
let node = TapNode::new(node_config);
let config = TapHttpConfig::default();
let mut server = TapHttpServer::new(config, node);
let events = Arc::new(Mutex::new(Vec::new()));
let subscriber = TestSubscriber {
events: events.clone(),
};
server.event_bus().subscribe(subscriber);
server.start().await.unwrap();
sleep(Duration::from_millis(100)).await;
server.stop().await.unwrap();
sleep(Duration::from_millis(100)).await;
let captured_events = events.lock().await;
assert!(captured_events.len() >= 2);
match &captured_events[0] {
HttpEvent::ServerStarted { address } => {
assert!(address.starts_with("127.0.0.1:"));
}
_ => panic!("First event should be ServerStarted"),
}
match &captured_events[captured_events.len() - 1] {
HttpEvent::ServerStopped => {}
_ => panic!("Last event should be ServerStopped"),
}
}
#[tokio::test(flavor = "multi_thread")]
async fn test_json_event_logging() {
let temp_dir = tempdir().unwrap();
let log_path = temp_dir.path().join("json-events.log");
let log_path_str = log_path.to_str().unwrap().to_string();
let config = TapHttpConfig {
port: 8001,
event_logger: Some(EventLoggerConfig {
destination: LogDestination::File {
path: log_path_str.clone(),
max_size: None,
rotate: false,
},
structured: true,
log_level: tracing::Level::INFO,
}),
..Default::default()
};
let node_config = NodeConfig {
storage_path: None,
..Default::default()
};
let node = TapNode::new(node_config);
let mut server = TapHttpServer::new(config, node);
server.start().await.unwrap();
sleep(Duration::from_millis(100)).await;
server.stop().await.unwrap();
sleep(Duration::from_millis(100)).await;
assert!(Path::new(&log_path_str).exists());
let log_content = fs::read_to_string(&log_path_str).unwrap();
let log_lines: Vec<&str> = log_content.trim().split('\n').collect();
assert!(!log_lines.is_empty());
for line in log_lines {
let json: Value = serde_json::from_str(line).unwrap();
assert!(json["timestamp"].is_string());
assert!(json["event_type"].is_string());
assert!(json["data"].is_object());
}
}