mod batching_http_transport;
mod http_transport;
mod manifest;
mod transport;
mod websocket_transport;
use std::sync::atomic::{AtomicU64, Ordering};
use std::sync::Arc;
use std::time::{Duration, Instant};
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 manifest::{
is_forbidden_monitor_ip, monitor_origin, resolve_monitor_target, send_manifest,
send_manifest_from_env, send_process_heartbeat, send_process_shutdown, AppManifest, CronEntry,
ExpectedProcessRole, HttpMethod, HttpMonitor, ManifestError, MonitorTargetError,
ProcessHeartbeat, ProcessHeartbeatConfig, ProcessHeartbeatHandle, ProcessIdentity,
ProcessSignal, ProcessSignalError, ProcessSignalPayload, MANIFEST_VERSION,
};
pub use transport::TransportError;
pub use websocket_transport::WebSocketTransport;
#[doc(hidden)]
pub use tracing;
use transport::Transport;
pub const MEASUREMENT_TARGET: &str = "eyes::measurement";
pub const MEASUREMENT_VERSION: u64 = 1;
pub const MEASUREMENT_RESERVED_FIELDS: [&str; 5] =
["metric_name", "metric_kind", "value", "unit", "description"];
pub fn emit_gauge(metric_name: &str, value: f64) {
tracing::event!(
target: "eyes::measurement",
tracing::Level::INFO,
metric_name = metric_name,
metric_kind = "gauge",
value = value,
);
}
pub fn emit_counter(metric_name: &str, value: u64) {
tracing::event!(
target: "eyes::measurement",
tracing::Level::INFO,
metric_name = metric_name,
metric_kind = "counter",
value = value,
);
}
pub fn emit_sample(metric_name: &str, value: f64) {
tracing::event!(
target: "eyes::measurement",
tracing::Level::INFO,
metric_name = metric_name,
metric_kind = "sample",
value = value,
);
}
#[macro_export]
macro_rules! measurement {
($kind:expr, $name:expr, $value:expr $(,)?) => {
$crate::tracing::event!(
target: "eyes::measurement",
$crate::tracing::Level::INFO,
metric_name = $name,
metric_kind = $kind,
value = $value,
)
};
($kind:expr, $name:expr, $value:expr, $($dimensions:tt)+) => {
$crate::tracing::event!(
target: "eyes::measurement",
$crate::tracing::Level::INFO,
metric_name = $name,
metric_kind = $kind,
value = $value,
$($dimensions)+
)
};
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub(crate) struct EventData {
event_type: String,
event_data: Value,
event_timestamp: DateTime<Utc>,
#[serde(default)]
process_instance_id: Option<Uuid>,
}
pub const DEFAULT_QUEUE_CAPACITY: usize = 65_536;
#[derive(Debug, Clone)]
pub struct EyesLayer {
sender: mpsc::Sender<EventData>,
dropped: Arc<AtomicU64>,
emit_enter_exit: bool,
process_instance_id: Option<Uuid>,
}
impl EyesLayer {
fn dispatch(&self, event: EventData) {
match self.sender.try_send(event) {
Ok(()) => {}
Err(mpsc::error::TrySendError::Full(_)) => {
self.dropped.fetch_add(1, Ordering::Relaxed);
}
Err(mpsc::error::TrySendError::Closed(_)) => {}
}
}
}
#[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(())
}
}
pub(crate) struct RedactedToken<'a>(pub(crate) Option<&'a str>);
impl std::fmt::Debug for RedactedToken<'_> {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match self.0 {
None => f.write_str("None"),
Some(token) => {
let prefix: String = token.chars().take(12).collect();
write!(f, "Some(\"{prefix}\u{2026}\")")
}
}
}
}
#[derive(Clone)]
pub struct EyesSubscriberBuilder {
base_url: Url,
org_id: Uuid,
app_id: Uuid,
queue_capacity: Option<usize>,
emit_enter_exit: Option<bool>,
process_instance_id: Option<Uuid>,
auth_token: Option<String>,
}
impl std::fmt::Debug for EyesSubscriberBuilder {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("EyesSubscriberBuilder")
.field("base_url", &self.base_url)
.field("org_id", &self.org_id)
.field("app_id", &self.app_id)
.field("queue_capacity", &self.queue_capacity)
.field("emit_enter_exit", &self.emit_enter_exit)
.field("process_instance_id", &self.process_instance_id)
.field("auth_token", &RedactedToken(self.auth_token.as_deref()))
.finish()
}
}
#[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,
queue_capacity: None,
emit_enter_exit: None,
process_instance_id: None,
auth_token: None,
})
}
pub fn with_process_instance_id(mut self, instance_id: Uuid) -> Self {
self.process_instance_id = Some(instance_id);
self
}
pub fn with_queue_capacity(mut self, capacity: usize) -> Self {
self.queue_capacity = Some(capacity);
self
}
pub fn with_emit_enter_exit(mut self, emit: bool) -> Self {
self.emit_enter_exit = Some(emit);
self
}
fn resolve_emit_enter_exit(&self) -> bool {
self.emit_enter_exit.unwrap_or_else(|| {
std::env::var("EYES_EMIT_ENTER_EXIT")
.map(|v| {
let v = v.to_lowercase();
v == "1" || v == "true"
})
.unwrap_or(false)
})
}
pub fn with_auth_token(mut self, token: impl Into<String>) -> Self {
self.auth_token = Some(token.into());
self
}
fn resolve_auth_token(&self) -> Option<String> {
self.auth_token.clone().or_else(|| {
std::env::var("EYES_TOKEN")
.ok()
.map(|s| s.trim().to_string())
.filter(|s| !s.is_empty())
})
}
fn resolve_queue_capacity(&self) -> usize {
self.queue_capacity
.or_else(|| {
std::env::var("EYES_QUEUE_CAPACITY")
.ok()
.and_then(|v| v.parse().ok())
})
.unwrap_or(DEFAULT_QUEUE_CAPACITY)
.max(1)
}
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 emit_enter_exit = self.resolve_emit_enter_exit();
let auth_token = self.resolve_auth_token();
let (sender, receiver) = mpsc::channel::<EventData>(self.resolve_queue_capacity());
let (shutdown_tx, shutdown_rx) = oneshot::channel();
let (completion_tx, completion_rx) = oneshot::channel();
let dropped = Arc::new(AtomicU64::new(0));
let transport: Box<dyn Transport> = match transport_type {
TransportType::Http => Box::new(
HttpTransport::new(self.base_url.clone(), self.org_id, self.app_id, auth_token)
.expect("Failed to create HTTP transport"),
),
TransportType::BatchingHttp => Box::new(
BatchingHttpTransport::new(
self.base_url.clone(),
self.org_id,
self.app_id,
batch_config,
auth_token,
)
.expect("Failed to create batching HTTP transport"),
),
TransportType::WebSocket => Box::new(
WebSocketTransport::new(
self.base_url.clone(),
self.org_id,
self.app_id,
auth_token,
)
.expect("Failed to create WebSocket transport"),
),
};
tokio::spawn(transport::run_transport_loop(
transport,
receiver,
shutdown_rx,
completion_tx,
dropped.clone(),
transport::TransportLoopConfig::default(),
));
let layer = EyesLayer {
sender,
dropped,
emit_enter_exit,
process_instance_id: self.process_instance_id,
};
let handle = EyesShutdownHandle {
shutdown_tx,
completion_rx,
};
(layer, handle)
}
}
struct EyesSpanId(String);
#[derive(Default)]
struct BusyTimings {
last_enter: Option<Instant>,
busy: Duration,
poll_count: u64,
}
struct RecordedFields(serde_json::Map<String, Value>);
fn generate_span_id() -> String {
Uuid::new_v4().simple().to_string()
}
fn eyes_span_id<S>(span: &tracing_subscriber::registry::SpanRef<'_, S>) -> String
where
S: Subscriber + for<'a> LookupSpan<'a>,
{
span.extensions()
.get::<EyesSpanId>()
.map(|eyes_id| eyes_id.0.clone())
.unwrap_or_else(|| format!("{:?}", span.id()))
}
impl EyesLayer {
fn measurement_event<S>(
&self,
event: &Event<'_>,
ctx: &Context<'_, S>,
mut fields: serde_json::Map<String, Value>,
) -> Option<EventData>
where
S: Subscriber + for<'a> LookupSpan<'a>,
{
let lifted: Vec<(&str, Option<Value>)> = MEASUREMENT_RESERVED_FIELDS
.iter()
.map(|name| (*name, fields.remove(*name)))
.collect();
if !lifted
.iter()
.any(|(name, value)| *name == "value" && matches!(value, Some(Value::Number(_))))
{
return None;
}
let mut event_data = serde_json::json!({
"version": MEASUREMENT_VERSION,
"level": event.metadata().level().to_string(),
"target": event.metadata().target(),
"fields": fields,
});
for (name, value) in lifted {
if let Some(value) = value {
event_data[name] = value;
}
}
if let Some(span) = ctx.event_span(event) {
event_data["span_id"] = serde_json::json!(eyes_span_id(&span));
}
Some(EventData {
event_type: "measurement".to_string(),
event_data,
event_timestamp: Utc::now(),
process_instance_id: self.process_instance_id,
})
}
}
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 unique_id = generate_span_id();
span.extensions_mut().insert(EyesSpanId(unique_id.clone()));
let mut visitor = JsonVisitor::default();
attrs.record(&mut visitor);
let mut event_data = serde_json::json!({
"span_id": unique_id,
"name": span.metadata().name(),
"target": span.metadata().target(),
"level": span.metadata().level().to_string(),
"fields": visitor.fields,
});
let parent = span
.parent()
.or_else(|| ctx.current_span().id().and_then(|pid| ctx.span(pid)));
if let Some(parent) = parent {
event_data["parent_id"] = serde_json::json!(eyes_span_id(&parent));
}
let event = EventData {
event_type: "span_new".to_string(),
event_data,
event_timestamp: Utc::now(),
process_instance_id: self.process_instance_id,
};
self.dispatch(event);
}
fn on_record(&self, id: &Id, values: &span::Record<'_>, ctx: Context<'_, S>) {
let Some(span) = ctx.span(id) else {
return;
};
let mut visitor = JsonVisitor::default();
values.record(&mut visitor);
if visitor.fields.is_empty() {
return;
}
let mut extensions = span.extensions_mut();
if let Some(recorded) = extensions.get_mut::<RecordedFields>() {
recorded.0.extend(visitor.fields);
} else {
extensions.insert(RecordedFields(visitor.fields));
}
}
fn on_event(&self, event: &Event<'_>, ctx: Context<'_, S>) {
let mut visitor = JsonVisitor::default();
event.record(&mut visitor);
if event.metadata().target() == MEASUREMENT_TARGET {
match self.measurement_event(event, &ctx, visitor.fields) {
Some(measurement) => self.dispatch(measurement),
None => {
self.dropped.fetch_add(1, Ordering::Relaxed);
}
}
return;
}
let mut event_data = serde_json::json!({
"level": event.metadata().level().to_string(),
"target": event.metadata().target(),
"fields": visitor.fields,
});
if let Some(span) = ctx.event_span(event) {
event_data["span_id"] = serde_json::json!(eyes_span_id(&span));
}
let event_msg = EventData {
event_type: "event".to_string(),
event_data,
event_timestamp: Utc::now(),
process_instance_id: self.process_instance_id,
};
self.dispatch(event_msg);
}
fn on_enter(&self, id: &Id, ctx: Context<'_, S>) {
if let Some(span) = ctx.span(id) {
let mut extensions = span.extensions_mut();
if let Some(timings) = extensions.get_mut::<BusyTimings>() {
timings.last_enter = Some(Instant::now());
} else {
extensions.insert(BusyTimings {
last_enter: Some(Instant::now()),
..BusyTimings::default()
});
}
}
if !self.emit_enter_exit {
return;
}
let span_id = ctx
.span(id)
.map(|span| eyes_span_id(&span))
.unwrap_or_else(|| format!("{:?}", id));
let event = EventData {
event_type: "span_enter".to_string(),
event_data: serde_json::json!({
"span_id": span_id,
}),
event_timestamp: Utc::now(),
process_instance_id: self.process_instance_id,
};
self.dispatch(event);
}
fn on_exit(&self, id: &Id, ctx: Context<'_, S>) {
if let Some(span) = ctx.span(id) {
let mut extensions = span.extensions_mut();
if let Some(timings) = extensions.get_mut::<BusyTimings>() {
if let Some(entered_at) = timings.last_enter.take() {
timings.busy += entered_at.elapsed();
timings.poll_count += 1;
}
}
}
if !self.emit_enter_exit {
return;
}
let span_id = ctx
.span(id)
.map(|span| eyes_span_id(&span))
.unwrap_or_else(|| format!("{:?}", id));
let event = EventData {
event_type: "span_exit".to_string(),
event_data: serde_json::json!({
"span_id": span_id,
}),
event_timestamp: Utc::now(),
process_instance_id: self.process_instance_id,
};
self.dispatch(event);
}
fn on_close(&self, id: Id, ctx: Context<'_, S>) {
let span = ctx.span(&id).expect("Span not found");
let (busy_ms, poll_count) = span
.extensions()
.get::<BusyTimings>()
.map(|timings| {
(
u64::try_from(timings.busy.as_millis()).unwrap_or(u64::MAX),
timings.poll_count,
)
})
.unwrap_or((0, 0));
let mut fields = span
.extensions_mut()
.remove::<RecordedFields>()
.map(|recorded| recorded.0)
.unwrap_or_default();
fields.insert("busy_ms".to_string(), Value::from(busy_ms));
fields.insert("poll_count".to_string(), Value::from(poll_count));
let event = EventData {
event_type: "span_close".to_string(),
event_data: serde_json::json!({
"span_id": eyes_span_id(&span),
"name": span.metadata().name(),
"fields": fields,
}),
event_timestamp: Utc::now(),
process_instance_id: self.process_instance_id,
};
self.dispatch(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, None);
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, None);
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, None).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(),
process_instance_id: None,
};
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_span_ids_unique_across_registry_reuse() {
let (sender, mut receiver) = mpsc::channel::<EventData>(64);
let layer = EyesLayer {
sender,
dropped: std::sync::Arc::new(std::sync::atomic::AtomicU64::new(0)),
emit_enter_exit: true,
process_instance_id: None,
};
let subscriber = tracing_subscriber::registry().with(layer);
tracing::subscriber::with_default(subscriber, || {
{
let parent = span!(Level::INFO, "parent_span");
let _parent_guard = parent.enter();
let child = span!(Level::INFO, "child_span");
let _child_guard = child.enter();
info!("inside child");
}
{
let reused = span!(Level::INFO, "reused_slot_span");
let _guard = reused.enter();
}
});
let mut events = Vec::new();
while let Ok(event) = receiver.try_recv() {
events.push(event);
}
let span_id_of = |name: &str| -> String {
events
.iter()
.find(|e| e.event_type == "span_new" && e.event_data["name"] == name)
.unwrap_or_else(|| panic!("no span_new for {name}"))
.event_data["span_id"]
.as_str()
.unwrap()
.to_string()
};
let parent_id = span_id_of("parent_span");
let child_id = span_id_of("child_span");
let reused_id = span_id_of("reused_slot_span");
for id in [&parent_id, &child_id, &reused_id] {
assert_eq!(id.len(), 32, "unexpected id shape: {id}");
assert!(!id.starts_with("Id("), "registry id leaked: {id}");
}
assert_ne!(parent_id, child_id);
assert_ne!(reused_id, parent_id);
assert_ne!(reused_id, child_id);
let child_new = events
.iter()
.find(|e| e.event_type == "span_new" && e.event_data["name"] == "child_span")
.unwrap();
assert_eq!(child_new.event_data["parent_id"], parent_id.as_str());
let child_lifecycle: Vec<_> = events
.iter()
.filter(|e| {
matches!(
e.event_type.as_str(),
"span_enter" | "span_exit" | "span_close"
) && e.event_data["span_id"] == child_id.as_str()
})
.collect();
assert!(
child_lifecycle.len() >= 3,
"expected enter/exit/close with the child's unique id, got {}",
child_lifecycle.len()
);
let in_span_event = events
.iter()
.find(|e| {
e.event_type == "event" && e.event_data["fields"]["message"] == "inside child"
})
.unwrap();
assert_eq!(in_span_event.event_data["span_id"], child_id.as_str());
}
#[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, None);
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, None);
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);
}
#[test]
fn test_dispatch_drops_and_counts_when_queue_full() {
let (sender, mut receiver) = mpsc::channel(1);
let layer = EyesLayer {
sender,
dropped: Arc::new(AtomicU64::new(0)),
emit_enter_exit: false,
process_instance_id: None,
};
let event = |event_type: &str| EventData {
event_type: event_type.to_string(),
event_data: serde_json::json!({}),
event_timestamp: Utc::now(),
process_instance_id: None,
};
layer.dispatch(event("first"));
layer.dispatch(event("second"));
assert_eq!(layer.dropped.load(Ordering::Relaxed), 1);
assert_eq!(receiver.try_recv().unwrap().event_type, "first");
assert!(receiver.try_recv().is_err());
}
#[test]
#[serial_test::serial]
fn test_queue_capacity_default_and_builder() {
std::env::remove_var("EYES_QUEUE_CAPACITY");
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.resolve_queue_capacity(), DEFAULT_QUEUE_CAPACITY);
assert_eq!(
builder
.clone()
.with_queue_capacity(123)
.resolve_queue_capacity(),
123
);
assert_eq!(builder.with_queue_capacity(0).resolve_queue_capacity(), 1);
}
#[test]
#[serial_test::serial]
fn test_queue_capacity_from_env() {
let org_id = Uuid::new_v4();
let app_id = Uuid::new_v4();
let builder = EyesSubscriberBuilder::new("http://localhost:4318", org_id, app_id).unwrap();
std::env::set_var("EYES_QUEUE_CAPACITY", "1024");
assert_eq!(builder.resolve_queue_capacity(), 1024);
assert_eq!(
builder
.clone()
.with_queue_capacity(123)
.resolve_queue_capacity(),
123
);
std::env::set_var("EYES_QUEUE_CAPACITY", "not-a-number");
assert_eq!(builder.resolve_queue_capacity(), DEFAULT_QUEUE_CAPACITY);
std::env::remove_var("EYES_QUEUE_CAPACITY");
}
#[test]
#[serial_test::serial]
fn test_emit_enter_exit_default_off() {
std::env::remove_var("EYES_EMIT_ENTER_EXIT");
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!(!builder.resolve_emit_enter_exit());
assert!(builder.with_emit_enter_exit(true).resolve_emit_enter_exit());
}
#[test]
#[serial_test::serial]
fn test_emit_enter_exit_from_env() {
let org_id = Uuid::new_v4();
let app_id = Uuid::new_v4();
let builder = EyesSubscriberBuilder::new("http://localhost:4318", org_id, app_id).unwrap();
for enabled in ["1", "true", "TRUE", "True"] {
std::env::set_var("EYES_EMIT_ENTER_EXIT", enabled);
assert!(
builder.resolve_emit_enter_exit(),
"{enabled:?} should enable enter/exit emission"
);
}
for disabled in ["0", "false", "yes", ""] {
std::env::set_var("EYES_EMIT_ENTER_EXIT", disabled);
assert!(
!builder.resolve_emit_enter_exit(),
"{disabled:?} should not enable enter/exit emission"
);
}
std::env::remove_var("EYES_EMIT_ENTER_EXIT");
}
#[test]
#[serial_test::serial]
fn test_emit_enter_exit_builder_overrides_env() {
let org_id = Uuid::new_v4();
let app_id = Uuid::new_v4();
let builder = EyesSubscriberBuilder::new("http://localhost:4318", org_id, app_id).unwrap();
std::env::set_var("EYES_EMIT_ENTER_EXIT", "1");
assert!(!builder
.clone()
.with_emit_enter_exit(false)
.resolve_emit_enter_exit());
std::env::set_var("EYES_EMIT_ENTER_EXIT", "false");
assert!(builder.with_emit_enter_exit(true).resolve_emit_enter_exit());
std::env::remove_var("EYES_EMIT_ENTER_EXIT");
}
fn test_layer(emit_enter_exit: bool) -> (EyesLayer, mpsc::Receiver<EventData>) {
let (sender, receiver) = mpsc::channel::<EventData>(64);
let layer = EyesLayer {
sender,
dropped: Arc::new(AtomicU64::new(0)),
emit_enter_exit,
process_instance_id: None,
};
(layer, receiver)
}
fn drain_events(receiver: &mut mpsc::Receiver<EventData>) -> Vec<EventData> {
let mut events = Vec::new();
while let Ok(event) = receiver.try_recv() {
events.push(event);
}
events
}
fn span_close_for<'a>(events: &'a [EventData], name: &str) -> &'a EventData {
events
.iter()
.find(|e| e.event_type == "span_close" && e.event_data["name"] == name)
.unwrap_or_else(|| panic!("no span_close for {name}"))
}
fn poll_span_n_times(
layer: EyesLayer,
receiver: &mut mpsc::Receiver<EventData>,
polls: u32,
sleep_per_poll: std::time::Duration,
) -> Vec<EventData> {
let subscriber = tracing_subscriber::registry().with(layer);
tracing::subscriber::with_default(subscriber, || {
let span = span!(Level::INFO, "polled_span");
for _ in 0..polls {
let guard = span.enter();
std::thread::sleep(sleep_per_poll);
drop(guard);
}
drop(span);
});
drain_events(receiver)
}
#[tokio::test]
async fn test_level_serialized_as_display_not_debug() {
let (layer, mut receiver) = test_layer(false);
let subscriber = tracing_subscriber::registry().with(layer);
tracing::subscriber::with_default(subscriber, || {
let _span = span!(Level::WARN, "warn_span");
tracing::error!("boom");
});
let events = drain_events(&mut receiver);
let span_new = events
.iter()
.find(|e| e.event_type == "span_new" && e.event_data["name"] == "warn_span")
.expect("span_new for warn_span");
assert_eq!(span_new.event_data["level"], "WARN");
let log = events
.iter()
.find(|e| e.event_type == "event")
.expect("log event");
assert_eq!(log.event_data["level"], "ERROR");
}
#[tokio::test]
async fn test_span_close_reports_busy_ms_and_poll_count() {
let (layer, mut receiver) = test_layer(false);
let events =
poll_span_n_times(layer, &mut receiver, 3, std::time::Duration::from_millis(5));
let close = span_close_for(&events, "polled_span");
let fields = &close.event_data["fields"];
assert_eq!(
fields["poll_count"].as_u64(),
Some(3),
"poll_count should match the number of enter/exit cycles"
);
let busy_ms = fields["busy_ms"].as_u64().expect("busy_ms should be a u64");
assert!(
busy_ms >= 10,
"busy_ms should reflect time in span, got {busy_ms}"
);
}
#[tokio::test]
async fn test_span_close_busy_fields_present_when_never_entered() {
let (layer, mut receiver) = test_layer(false);
let subscriber = tracing_subscriber::registry().with(layer);
tracing::subscriber::with_default(subscriber, || {
let _span = span!(Level::INFO, "never_entered_span");
});
let events = drain_events(&mut receiver);
let close = span_close_for(&events, "never_entered_span");
let fields = &close.event_data["fields"];
assert_eq!(fields["busy_ms"].as_u64(), Some(0));
assert_eq!(fields["poll_count"].as_u64(), Some(0));
}
#[tokio::test]
async fn test_busy_aggregation_works_with_emit_enter_exit_enabled() {
let (layer, mut receiver) = test_layer(true);
let events =
poll_span_n_times(layer, &mut receiver, 2, std::time::Duration::from_millis(5));
assert_eq!(
events
.iter()
.filter(|e| e.event_type == "span_enter")
.count(),
2
);
assert_eq!(
events
.iter()
.filter(|e| e.event_type == "span_exit")
.count(),
2
);
let close = span_close_for(&events, "polled_span");
let fields = &close.event_data["fields"];
assert_eq!(fields["poll_count"].as_u64(), Some(2));
let busy_ms = fields["busy_ms"].as_u64().expect("busy_ms should be a u64");
assert!(
busy_ms >= 5,
"busy_ms should reflect time in span, got {busy_ms}"
);
}
#[tokio::test]
async fn test_span_close_includes_recorded_fields() {
let (layer, mut receiver) = test_layer(false);
let subscriber = tracing_subscriber::registry().with(layer);
tracing::subscriber::with_default(subscriber, || {
let span = span!(
Level::INFO,
"recording_span",
status_code = tracing::field::Empty,
content_type = tracing::field::Empty
);
let _guard = span.enter();
span.record("status_code", 200_u64);
span.record("content_type", "text/html");
});
let events = drain_events(&mut receiver);
let span_new = events
.iter()
.find(|e| e.event_type == "span_new" && e.event_data["name"] == "recording_span")
.expect("no span_new for recording_span");
assert!(span_new.event_data["fields"]
.get("status_code")
.is_none_or(|v| v.is_null()));
let close = span_close_for(&events, "recording_span");
let fields = &close.event_data["fields"];
assert_eq!(fields["status_code"].as_u64(), Some(200));
assert_eq!(fields["content_type"].as_str(), Some("text/html"));
assert!(fields["busy_ms"].is_u64());
assert!(fields["poll_count"].is_u64());
}
#[tokio::test]
async fn test_recorded_field_last_write_wins() {
let (layer, mut receiver) = test_layer(false);
let subscriber = tracing_subscriber::registry().with(layer);
tracing::subscriber::with_default(subscriber, || {
let span = span!(
Level::INFO,
"rerecord_span",
attempt = tracing::field::Empty
);
let _guard = span.enter();
span.record("attempt", 1_u64);
span.record("attempt", 2_u64);
span.record("attempt", 3_u64);
});
let events = drain_events(&mut receiver);
let close = span_close_for(&events, "rerecord_span");
assert_eq!(
close.event_data["fields"]["attempt"].as_u64(),
Some(3),
"the last recorded value for a field should win"
);
}
#[tokio::test]
async fn test_recorded_fields_cannot_clobber_busy_aggregation() {
let (layer, mut receiver) = test_layer(false);
let subscriber = tracing_subscriber::registry().with(layer);
tracing::subscriber::with_default(subscriber, || {
let span = span!(
Level::INFO,
"hostile_span",
busy_ms = tracing::field::Empty,
poll_count = tracing::field::Empty
);
let guard = span.enter();
span.record("busy_ms", "not-a-duration");
span.record("poll_count", "lots");
drop(guard);
});
let events = drain_events(&mut receiver);
let close = span_close_for(&events, "hostile_span");
let fields = &close.event_data["fields"];
assert!(
fields["busy_ms"].is_u64(),
"busy_ms must remain the aggregated u64, got {:?}",
fields["busy_ms"]
);
assert_eq!(
fields["poll_count"].as_u64(),
Some(1),
"poll_count must remain the aggregated value, got {:?}",
fields["poll_count"]
);
}
#[tokio::test]
async fn test_span_close_shape_unchanged_without_records() {
let (layer, mut receiver) = test_layer(false);
let subscriber = tracing_subscriber::registry().with(layer);
tracing::subscriber::with_default(subscriber, || {
let span = span!(Level::INFO, "no_record_span", user_id = 7);
let _guard = span.enter();
});
let events = drain_events(&mut receiver);
let close = span_close_for(&events, "no_record_span");
let fields = close.event_data["fields"]
.as_object()
.expect("span_close fields should be an object");
assert_eq!(fields.len(), 2, "unexpected span_close fields: {fields:?}");
assert!(fields["busy_ms"].is_u64());
assert_eq!(fields["poll_count"].as_u64(), Some(1));
}
#[tokio::test]
async fn test_enter_exit_not_emitted_by_default_layer() {
let (sender, mut receiver) = mpsc::channel::<EventData>(64);
let layer = EyesLayer {
sender,
dropped: Arc::new(AtomicU64::new(0)),
emit_enter_exit: false,
process_instance_id: None,
};
let subscriber = tracing_subscriber::registry().with(layer);
tracing::subscriber::with_default(subscriber, || {
let span = span!(Level::INFO, "quiet_span");
let _guard = span.enter();
info!("inside quiet span");
});
let mut events = Vec::new();
while let Ok(event) = receiver.try_recv() {
events.push(event);
}
assert!(
events
.iter()
.all(|e| e.event_type != "span_enter" && e.event_type != "span_exit"),
"span_enter/span_exit must not be emitted by default"
);
for expected in ["span_new", "event", "span_close"] {
assert!(
events.iter().any(|e| e.event_type == expected),
"missing {expected} event"
);
}
}
fn drain_measurements(build: impl FnOnce(EyesLayer)) -> Vec<EventData> {
let (sender, mut receiver) = mpsc::channel::<EventData>(64);
let layer = EyesLayer {
sender,
dropped: std::sync::Arc::new(std::sync::atomic::AtomicU64::new(0)),
emit_enter_exit: false,
process_instance_id: None,
};
build(layer);
let mut events = Vec::new();
while let Ok(event) = receiver.try_recv() {
events.push(event);
}
events
}
fn emit_through_layer(body: impl FnOnce()) -> Vec<EventData> {
drain_measurements(|layer| {
let subscriber = tracing_subscriber::registry().with(layer);
tracing::subscriber::with_default(subscriber, body);
})
}
#[test]
fn a_non_finite_measurement_value_is_dropped_and_counted() {
let (sender, mut receiver) = mpsc::channel::<EventData>(64);
let dropped = std::sync::Arc::new(std::sync::atomic::AtomicU64::new(0));
let layer = EyesLayer {
sender,
dropped: dropped.clone(),
emit_enter_exit: false,
process_instance_id: None,
};
let subscriber = tracing_subscriber::registry().with(layer);
tracing::subscriber::with_default(subscriber, || {
emit_gauge("rate", f64::NAN);
emit_sample("latency", f64::INFINITY);
emit_gauge("ok", 1.5);
});
let mut events = Vec::new();
while let Ok(event) = receiver.try_recv() {
events.push(event);
}
assert_eq!(events.len(), 1, "only the finite gauge survives");
assert_eq!(events[0].event_data["metric_name"], "ok");
assert_eq!(dropped.load(Ordering::Relaxed), 2);
}
#[test]
fn emit_gauge_dispatches_the_versioned_measurement_contract() {
let events = emit_through_layer(|| emit_gauge("cpu", 42.5));
assert_eq!(events.len(), 1);
let data = &events[0].event_data;
assert_eq!(events[0].event_type, "measurement");
assert_eq!(data["version"], serde_json::json!(MEASUREMENT_VERSION));
assert_eq!(data["metric_name"], "cpu");
assert_eq!(data["metric_kind"], "gauge");
assert_eq!(data["value"], serde_json::json!(42.5));
assert_eq!(data["level"], "INFO");
assert_eq!(data["target"], MEASUREMENT_TARGET);
assert_eq!(data["fields"], serde_json::json!({}));
}
#[test]
fn emit_counter_keeps_its_integer_identity_on_the_wire() {
let events = emit_through_layer(|| emit_counter("requests", 5));
assert_eq!(events.len(), 1);
let value = &events[0].event_data["value"];
assert!(value.is_i64() || value.is_u64(), "not an integer: {value}");
assert_eq!(value, &serde_json::json!(5));
}
#[test]
fn emit_sample_names_its_kind() {
let events = emit_through_layer(|| emit_sample("latency", 1.5));
assert_eq!(events[0].event_data["metric_kind"], "sample");
}
#[test]
fn the_macro_lifts_reserved_names_and_leaves_version_a_dimension() {
let events = emit_through_layer(|| {
measurement!(
"gauge",
"cpu",
42.5_f64,
unit = "percent",
description = "d",
host = "web-1",
version = "app-2.1"
);
});
let data = &events[0].event_data;
assert_eq!(data["unit"], "percent");
assert_eq!(data["description"], "d");
assert_eq!(data["version"], serde_json::json!(MEASUREMENT_VERSION));
assert_eq!(data["fields"]["host"], "web-1");
assert_eq!(data["fields"]["version"], "app-2.1");
for reserved in MEASUREMENT_RESERVED_FIELDS {
assert!(
data["fields"].get(reserved).is_none(),
"{reserved} left in the dimension bag"
);
}
}
#[test]
fn a_measurement_inside_a_span_carries_that_span_s_eyes_id() {
let events = emit_through_layer(|| {
let span = span!(Level::INFO, "outer");
let _guard = span.enter();
emit_gauge("cpu", 1.0);
});
let opened = events
.iter()
.find(|e| e.event_type == "span_new")
.expect("span_new");
let measured = events
.iter()
.find(|e| e.event_type == "measurement")
.expect("measurement");
let span_id = opened.event_data["span_id"].as_str().unwrap();
assert_eq!(measured.event_data["span_id"], span_id);
assert_eq!(span_id.len(), 32);
assert!(!span_id.starts_with("Id("));
}
#[test]
fn a_filter_that_does_not_enable_the_measurement_target_drops_measurements() {
use tracing_subscriber::EnvFilter;
let dropped = drain_measurements(|layer| {
let subscriber = tracing_subscriber::registry()
.with(EnvFilter::new("warn"))
.with(layer);
tracing::subscriber::with_default(subscriber, || emit_gauge("cpu", 1.0));
});
assert!(dropped.is_empty(), "{dropped:?}");
let kept = drain_measurements(|layer| {
let subscriber = tracing_subscriber::registry()
.with(EnvFilter::new("warn,eyes::measurement=info"))
.with(layer);
tracing::subscriber::with_default(subscriber, || emit_gauge("cpu", 1.0));
});
assert_eq!(kept.len(), 1);
}
#[test]
fn every_emitter_writes_the_same_literal_target() {
for events in [
emit_through_layer(|| emit_gauge("g", 1.0)),
emit_through_layer(|| emit_counter("c", 1)),
emit_through_layer(|| emit_sample("s", 1.0)),
emit_through_layer(|| measurement!("gauge", "m", 1.0_f64)),
emit_through_layer(|| measurement!("gauge", "m", 1.0_f64, host = "h")),
] {
assert_eq!(events.len(), 1);
assert_eq!(events[0].event_data["target"], MEASUREMENT_TARGET);
assert_eq!(events[0].event_type, "measurement");
}
}
}