mod batching_http_transport;
mod http_transport;
mod transport;
mod websocket_transport;
use chrono::{DateTime, Utc};
use serde::{Deserialize, Serialize};
use serde_json::Value;
use tokio::sync::{mpsc, oneshot};
use tracing::{field::Visit, span, Event, Id, Subscriber};
use tracing_subscriber::{layer::Context, registry::LookupSpan, Layer};
use url::Url;
use uuid::Uuid;
pub use batching_http_transport::{BatchConfig, BatchingHttpTransport};
pub use http_transport::HttpTransport;
pub use transport::TransportError;
pub use websocket_transport::WebSocketTransport;
use transport::Transport;
#[derive(Debug, Clone, Serialize, Deserialize)]
pub(crate) struct EventData {
event_type: String,
event_data: Value,
event_timestamp: DateTime<Utc>,
}
#[derive(Debug, Clone)]
pub struct EyesLayer {
sender: mpsc::UnboundedSender<EventData>,
}
#[derive(Debug)]
pub struct EyesShutdownHandle {
shutdown_tx: oneshot::Sender<()>,
completion_rx: oneshot::Receiver<()>,
}
impl EyesShutdownHandle {
pub async fn shutdown(self) -> Result<(), Box<dyn std::error::Error>> {
let _ = self.shutdown_tx.send(());
self.completion_rx.await?;
Ok(())
}
}
#[derive(Debug, Clone)]
pub struct EyesSubscriberBuilder {
base_url: Url,
org_id: Uuid,
app_id: Uuid,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum TransportType {
Http,
BatchingHttp,
WebSocket,
}
impl EyesSubscriberBuilder {
pub fn new(
base_url: impl Into<String>,
org_id: Uuid,
app_id: Uuid,
) -> Result<Self, url::ParseError> {
Ok(Self {
base_url: Url::parse(&base_url.into())?,
org_id,
app_id,
})
}
pub fn new_with_default(org_id: Uuid, app_id: Uuid) -> Result<Self, url::ParseError> {
Self::new("https://eyes.coreyja.com", org_id, app_id)
}
pub fn from_env_with_transport(
org_id: Uuid,
app_id: Uuid,
) -> Result<(Self, TransportType), url::ParseError> {
let base_url =
std::env::var("EYES_URL").unwrap_or_else(|_| "https://eyes.coreyja.com".to_string());
let transport = match std::env::var("EYES_TRANSPORT")
.unwrap_or_else(|_| "http".to_string())
.to_lowercase()
.as_str()
{
"websocket" | "ws" => TransportType::WebSocket,
"batching" | "batch" | "batching_http" => TransportType::BatchingHttp,
_ => TransportType::Http,
};
Ok((Self::new(base_url, org_id, app_id)?, transport))
}
pub fn from_env(org_id: Uuid, app_id: Uuid) -> Result<Self, url::ParseError> {
let base_url =
std::env::var("EYES_URL").unwrap_or_else(|_| "https://eyes.coreyja.com".to_string());
Self::new(base_url, org_id, app_id)
}
pub fn build(self) -> (EyesLayer, EyesShutdownHandle) {
self.build_with_transport(TransportType::Http)
}
pub fn build_from_env(
org_id: Uuid,
app_id: Uuid,
) -> Result<(EyesLayer, EyesShutdownHandle), url::ParseError> {
let (builder, transport) = Self::from_env_with_transport(org_id, app_id)?;
Ok(builder.build_with_transport(transport))
}
pub fn build_with_transport(
self,
transport_type: TransportType,
) -> (EyesLayer, EyesShutdownHandle) {
self.build_with_transport_and_config(transport_type, BatchConfig::default())
}
pub fn build_with_transport_and_config(
self,
transport_type: TransportType,
batch_config: BatchConfig,
) -> (EyesLayer, EyesShutdownHandle) {
let (sender, receiver) = mpsc::unbounded_channel::<EventData>();
let (shutdown_tx, shutdown_rx) = oneshot::channel();
let (completion_tx, completion_rx) = oneshot::channel();
let transport: Box<dyn Transport> = match transport_type {
TransportType::Http => Box::new(
HttpTransport::new(self.base_url.clone(), self.org_id, self.app_id)
.expect("Failed to create HTTP transport"),
),
TransportType::BatchingHttp => Box::new(
BatchingHttpTransport::new(
self.base_url.clone(),
self.org_id,
self.app_id,
batch_config,
)
.expect("Failed to create batching HTTP transport"),
),
TransportType::WebSocket => Box::new(
WebSocketTransport::new(self.base_url.clone(), self.org_id, self.app_id)
.expect("Failed to create WebSocket transport"),
),
};
tokio::spawn(transport::run_transport_loop(
transport,
receiver,
shutdown_rx,
completion_tx,
));
let layer = EyesLayer { sender };
let handle = EyesShutdownHandle {
shutdown_tx,
completion_rx,
};
(layer, handle)
}
}
impl<S> Layer<S> for EyesLayer
where
S: Subscriber + for<'a> LookupSpan<'a>,
{
fn on_new_span(&self, attrs: &span::Attributes<'_>, id: &Id, ctx: Context<'_, S>) {
let span = ctx.span(id).expect("Span not found");
let mut visitor = JsonVisitor::default();
attrs.record(&mut visitor);
let mut event_data = serde_json::json!({
"span_id": format!("{:?}", id),
"name": span.metadata().name(),
"target": span.metadata().target(),
"level": format!("{:?}", span.metadata().level()),
"fields": visitor.fields,
});
let parent_id = span
.parent()
.map(|p| p.id())
.or_else(|| ctx.current_span().id().cloned());
if let Some(parent_id) = parent_id {
event_data["parent_id"] = serde_json::json!(format!("{:?}", parent_id));
}
let event = EventData {
event_type: "span_new".to_string(),
event_data,
event_timestamp: Utc::now(),
};
let _ = self.sender.send(event);
}
fn on_event(&self, event: &Event<'_>, ctx: Context<'_, S>) {
let mut visitor = JsonVisitor::default();
event.record(&mut visitor);
let mut event_data = serde_json::json!({
"level": format!("{:?}", event.metadata().level()),
"target": event.metadata().target(),
"fields": visitor.fields,
});
if let Some(span) = ctx.event_span(event) {
event_data["span_id"] = serde_json::json!(format!("{:?}", span.id()));
}
let event_msg = EventData {
event_type: "event".to_string(),
event_data,
event_timestamp: Utc::now(),
};
let _ = self.sender.send(event_msg);
}
fn on_enter(&self, id: &Id, _ctx: Context<'_, S>) {
let event = EventData {
event_type: "span_enter".to_string(),
event_data: serde_json::json!({
"span_id": format!("{:?}", id),
}),
event_timestamp: Utc::now(),
};
let _ = self.sender.send(event);
}
fn on_exit(&self, id: &Id, _ctx: Context<'_, S>) {
let event = EventData {
event_type: "span_exit".to_string(),
event_data: serde_json::json!({
"span_id": format!("{:?}", id),
}),
event_timestamp: Utc::now(),
};
let _ = self.sender.send(event);
}
fn on_close(&self, id: Id, ctx: Context<'_, S>) {
let span = ctx.span(&id).expect("Span not found");
let event = EventData {
event_type: "span_close".to_string(),
event_data: serde_json::json!({
"span_id": format!("{:?}", id),
"name": span.metadata().name(),
}),
event_timestamp: Utc::now(),
};
let _ = self.sender.send(event);
}
}
#[derive(Default)]
struct JsonVisitor {
fields: serde_json::Map<String, Value>,
}
impl Visit for JsonVisitor {
fn record_debug(&mut self, field: &tracing::field::Field, value: &dyn std::fmt::Debug) {
self.fields.insert(
field.name().to_string(),
Value::String(format!("{:?}", value)),
);
}
fn record_str(&mut self, field: &tracing::field::Field, value: &str) {
self.fields
.insert(field.name().to_string(), Value::String(value.to_string()));
}
fn record_i64(&mut self, field: &tracing::field::Field, value: i64) {
self.fields
.insert(field.name().to_string(), Value::Number(value.into()));
}
fn record_u64(&mut self, field: &tracing::field::Field, value: u64) {
self.fields
.insert(field.name().to_string(), Value::Number(value.into()));
}
fn record_bool(&mut self, field: &tracing::field::Field, value: bool) {
self.fields
.insert(field.name().to_string(), Value::Bool(value));
}
fn record_f64(&mut self, field: &tracing::field::Field, value: f64) {
self.fields.insert(
field.name().to_string(),
serde_json::Number::from_f64(value)
.map(Value::Number)
.unwrap_or(Value::Null),
);
}
}
#[cfg(test)]
mod tests {
use super::*;
use tracing::{info, span, Level};
use tracing_subscriber::layer::SubscriberExt;
#[test]
fn test_builder_creation() {
let org_id = Uuid::new_v4();
let app_id = Uuid::new_v4();
let builder = EyesSubscriberBuilder::new("http://localhost:4318", org_id, app_id).unwrap();
assert_eq!(builder.app_id, app_id);
assert_eq!(builder.org_id, org_id);
}
#[test]
fn test_builder_invalid_url() {
let org_id = Uuid::new_v4();
let app_id = Uuid::new_v4();
let result = EyesSubscriberBuilder::new("invalid-url", org_id, app_id);
assert!(result.is_err());
}
#[test]
fn test_builder_new_with_default() {
let org_id = Uuid::new_v4();
let app_id = Uuid::new_v4();
let builder = EyesSubscriberBuilder::new_with_default(org_id, app_id).unwrap();
assert_eq!(builder.app_id, app_id);
assert_eq!(builder.org_id, org_id);
assert_eq!(builder.base_url.as_str(), "https://eyes.coreyja.com/");
}
#[test]
fn test_http_transport_creation() {
let org_id = Uuid::new_v4();
let app_id = Uuid::new_v4();
let base_url = Url::parse("http://localhost:4318").unwrap();
let transport = HttpTransport::new(base_url, org_id, app_id);
assert!(transport.is_ok());
}
#[test]
fn test_websocket_transport_creation() {
let org_id = Uuid::new_v4();
let app_id = Uuid::new_v4();
let base_url = Url::parse("http://localhost:4318").unwrap();
let transport = WebSocketTransport::new(base_url, org_id, app_id);
assert!(transport.is_ok());
}
#[test]
fn test_websocket_url_conversion() {
let org_id = Uuid::new_v4();
let app_id = Uuid::new_v4();
let https_url = Url::parse("https://example.com").unwrap();
let _transport = WebSocketTransport::new(https_url, org_id, app_id).unwrap();
}
#[test]
fn test_event_data_serialization() {
let event = EventData {
event_type: "test_event".to_string(),
event_data: serde_json::json!({"key": "value", "number": 42}),
event_timestamp: Utc::now(),
};
let serialized = serde_json::to_string(&event).unwrap();
let deserialized: EventData = serde_json::from_str(&serialized).unwrap();
assert_eq!(event.event_type, deserialized.event_type);
assert_eq!(event.event_data, deserialized.event_data);
}
#[test]
fn test_json_visitor_basic_functionality() {
let mut visitor = JsonVisitor::default();
assert_eq!(visitor.fields.len(), 0);
visitor.fields.insert(
"test_key".to_string(),
Value::String("test_value".to_string()),
);
assert_eq!(visitor.fields.len(), 1);
assert_eq!(
visitor.fields.get("test_key"),
Some(&Value::String("test_value".to_string()))
);
}
#[test]
fn test_transport_type_debug() {
let http = TransportType::Http;
let ws = TransportType::WebSocket;
assert_eq!(format!("{:?}", http), "Http");
assert_eq!(format!("{:?}", ws), "WebSocket");
}
#[test]
fn test_transport_type_equality() {
assert_eq!(TransportType::Http, TransportType::Http);
assert_eq!(TransportType::WebSocket, TransportType::WebSocket);
assert_ne!(TransportType::Http, TransportType::WebSocket);
}
#[tokio::test]
async fn test_layer_integration() {
let org_id = Uuid::new_v4();
let app_id = Uuid::new_v4();
let (layer, shutdown_handle) =
EyesSubscriberBuilder::new("http://localhost:4318", org_id, app_id)
.unwrap()
.build();
let subscriber = tracing_subscriber::registry().with(layer);
tracing::subscriber::with_default(subscriber, || {
let span = span!(Level::INFO, "test_span", user_id = 123);
let _enter = span.enter();
info!("Test message in span");
});
shutdown_handle.shutdown().await.unwrap();
}
#[tokio::test]
async fn test_layer_with_websocket_transport() {
let org_id = Uuid::new_v4();
let app_id = Uuid::new_v4();
let (layer, shutdown_handle) =
EyesSubscriberBuilder::new("http://localhost:4318", org_id, app_id)
.unwrap()
.build_with_transport(TransportType::WebSocket);
let subscriber = tracing_subscriber::registry().with(layer);
tracing::subscriber::with_default(subscriber, || {
info!("Test WebSocket transport");
});
shutdown_handle.shutdown().await.unwrap();
}
#[tokio::test]
async fn test_build_from_env_with_defaults() {
std::env::remove_var("EYES_URL");
std::env::remove_var("EYES_TRANSPORT");
let org_id = Uuid::new_v4();
let app_id = Uuid::new_v4();
let result = EyesSubscriberBuilder::build_from_env(org_id, app_id);
assert!(result.is_ok());
if let Ok((_, shutdown_handle)) = result {
shutdown_handle.shutdown().await.unwrap();
}
}
#[tokio::test]
async fn test_build_from_env_with_custom_values() {
std::env::set_var("EYES_URL", "http://custom.example.com");
std::env::set_var("EYES_TRANSPORT", "websocket");
let org_id = Uuid::new_v4();
let app_id = Uuid::new_v4();
let result = EyesSubscriberBuilder::build_from_env(org_id, app_id);
assert!(result.is_ok());
if let Ok((_, shutdown_handle)) = result {
shutdown_handle.shutdown().await.unwrap();
}
std::env::remove_var("EYES_URL");
std::env::remove_var("EYES_TRANSPORT");
}
#[tokio::test]
async fn test_layer_with_batching_transport() {
let org_id = Uuid::new_v4();
let app_id = Uuid::new_v4();
let (layer, shutdown_handle) =
EyesSubscriberBuilder::new("http://localhost:4318", org_id, app_id)
.unwrap()
.build_with_transport(TransportType::BatchingHttp);
let subscriber = tracing_subscriber::registry().with(layer);
tracing::subscriber::with_default(subscriber, || {
info!("Test batching HTTP transport");
});
shutdown_handle.shutdown().await.unwrap();
}
#[tokio::test]
async fn test_layer_with_batching_transport_custom_config() {
use std::time::Duration;
let org_id = Uuid::new_v4();
let app_id = Uuid::new_v4();
let custom_config = BatchConfig::new(50, Duration::from_millis(100));
let (layer, shutdown_handle) =
EyesSubscriberBuilder::new("http://localhost:4318", org_id, app_id)
.unwrap()
.build_with_transport_and_config(TransportType::BatchingHttp, custom_config);
let subscriber = tracing_subscriber::registry().with(layer);
tracing::subscriber::with_default(subscriber, || {
info!("Test batching HTTP transport with custom config");
});
shutdown_handle.shutdown().await.unwrap();
}
#[test]
fn test_batching_transport_creation() {
let org_id = Uuid::new_v4();
let app_id = Uuid::new_v4();
let base_url = Url::parse("http://localhost:4318").unwrap();
let transport =
BatchingHttpTransport::with_default_config(base_url.clone(), org_id, app_id);
assert!(transport.is_ok());
use std::time::Duration;
let custom_config = BatchConfig::new(50, Duration::from_millis(100));
let transport = BatchingHttpTransport::new(base_url, org_id, app_id, custom_config);
assert!(transport.is_ok());
}
#[test]
fn test_transport_type_batching_http() {
assert_eq!(TransportType::BatchingHttp, TransportType::BatchingHttp);
assert_ne!(TransportType::BatchingHttp, TransportType::Http);
assert_ne!(TransportType::BatchingHttp, TransportType::WebSocket);
}
}