Skip to main content

eyes_subscriber/
lib.rs

1//! # Eyes Subscriber
2//!
3//! A tracing subscriber for sending structured trace data to Eyes (eyes.coreyja.com).
4//!
5//! ## Quick Start
6//!
7//! ```no_run
8//! use eyes_subscriber::EyesSubscriberBuilder;
9//! use tracing_subscriber::prelude::*;
10//! use uuid::Uuid;
11//!
12//! # async fn example() -> Result<(), Box<dyn std::error::Error>> {
13//! let org_id = Uuid::parse_str("your-org-id")?;
14//! let app_id = Uuid::parse_str("your-app-id")?;
15//!
16//! // Simplest: auto-configure from environment variables
17//! let (eyes_layer, shutdown_handle) = EyesSubscriberBuilder::build_from_env(org_id, app_id)?;
18//!
19//! tracing_subscriber::registry()
20//!     .with(eyes_layer)
21//!     .init();
22//!
23//! // Your application code here
24//! tracing::info!("Application started");
25//!
26//! // Graceful shutdown
27//! shutdown_handle.shutdown().await?;
28//! # Ok(())
29//! # }
30//! ```
31//!
32//! ## Configuration
33//!
34//! The subscriber can be configured in several ways:
35//!
36//! 1. **Environment variables** (recommended):
37//!    - `EYES_URL`: Override the default URL (defaults to https://eyes.coreyja.com)
38//!    - `EYES_TRANSPORT`: Set to "websocket" or "ws" for WebSocket, defaults to HTTP
39//!    - `EYES_QUEUE_CAPACITY`: Capacity of the bounded event queue (defaults to 65536)
40//!    - `EYES_TOKEN`: Bearer token sent on every request (HTTP, batch, WebSocket
41//!      upgrade, manifest and heartbeat). Required once the server runs with
42//!      `EYES_API_AUTH=enforce`.
43//!    - `EYES_EMIT_ENTER_EXIT`: Set to "1" or "true" (case-insensitive) to emit
44//!      `span_enter`/`span_exit` events (disabled by default; see
45//!      [`EyesSubscriberBuilder::with_emit_enter_exit`])
46//! 2. **Default production**: Use `new_with_default()` for https://eyes.coreyja.com
47//! 3. **Custom URL**: Use `new()` with any URL for self-hosted instances
48//!
49//! ## Transports
50//!
51//! Three transport methods are available:
52//! - **HTTP** (default): Reliable, request/response based
53//! - **BatchingHttp**: HTTP with client-side batching for high-volume use cases
54//! - **WebSocket**: Lower latency, persistent connection
55//!
56//! ## Emitted measurements
57//!
58//! A measurement is an ordinary event with `event_type = "measurement"` and a
59//! **versioned** payload:
60//!
61//! ```json
62//! {"version": 1, "metric_name": "cpu", "metric_kind": "gauge", "value": 42.5,
63//!  "unit": "percent", "description": "CPU utilisation",
64//!  "fields": {"host": "web-1"}, "level": "INFO", "target": "eyes::measurement"}
65//! ```
66//!
67//! Emit one with [`emit_gauge`], [`emit_counter`], [`emit_sample`], or the
68//! [`measurement!`] macro when you have dimensions:
69//!
70//! ```no_run
71//! eyes_subscriber::emit_gauge("cpu", 42.5);
72//! eyes_subscriber::emit_counter("requests", 5);
73//! eyes_subscriber::measurement!(
74//!     "gauge", "cpu", 42.5_f64,
75//!     unit = "percent", description = "CPU utilisation", host = "web-1"
76//! );
77//! ```
78//!
79//! ### The v1 instrument matrix
80//!
81//! | kind | meaning | accepted values |
82//! | --- | --- | --- |
83//! | `gauge` | instantaneous value | any finite number |
84//! | `counter` | cumulative, monotonically non-decreasing total | any finite number >= 0 |
85//! | `sample` | one discrete observation | any finite number |
86//!
87//! Int and float are both valid for every kind. Counter **rate and reset
88//! semantics are deferred**: counters are stored and aggregable, but nothing
89//! computes a rate or detects a reset.
90//!
91//! ### Required filter directive
92//!
93//! Measurements travel as `tracing` events on the reserved
94//! [`MEASUREMENT_TARGET`], so a global `EnvFilter` that does not enable INFO for
95//! `eyes::measurement` drops them before this layer ever runs. An app with a
96//! target-scoped filter must include `eyes::measurement=info`.
97
98mod batching_http_transport;
99mod http_transport;
100mod manifest;
101mod transport;
102mod websocket_transport;
103
104use std::sync::atomic::{AtomicU64, Ordering};
105use std::sync::Arc;
106use std::time::{Duration, Instant};
107
108use chrono::{DateTime, Utc};
109use serde::{Deserialize, Serialize};
110use serde_json::Value;
111use tokio::sync::{mpsc, oneshot};
112use tracing::{field::Visit, span, Event, Id, Subscriber};
113use tracing_subscriber::{layer::Context, registry::LookupSpan, Layer};
114use url::Url;
115use uuid::Uuid;
116
117pub use batching_http_transport::{BatchConfig, BatchingHttpTransport};
118pub use http_transport::HttpTransport;
119pub use manifest::{
120    is_forbidden_monitor_ip, monitor_origin, resolve_monitor_target, send_manifest,
121    send_manifest_from_env, send_process_heartbeat, send_process_shutdown, AppManifest, CronEntry,
122    ExpectedProcessRole, HttpMethod, HttpMonitor, ManifestError, MonitorTargetError,
123    ProcessHeartbeat, ProcessHeartbeatConfig, ProcessHeartbeatHandle, ProcessIdentity,
124    ProcessSignal, ProcessSignalError, ProcessSignalPayload, MANIFEST_VERSION,
125};
126pub use transport::TransportError;
127pub use websocket_transport::WebSocketTransport;
128
129/// Re-exported so [`measurement!`] can name `tracing` hygienically in a crate
130/// that does not depend on it directly.
131#[doc(hidden)]
132pub use tracing;
133
134use transport::Transport;
135
136/// The reserved `tracing` target that routes an event through the measurement
137/// contract instead of the ordinary log-event path.
138///
139/// Emission goes through `tracing::event!` because the layer has no
140/// back-reference to the registry: only `on_event` receives a `Context` and can
141/// therefore resolve the containing span's eyes id. The cost is that a global
142/// `EnvFilter` which does not enable INFO for `eyes::measurement` drops
143/// measurements *before* the layer runs — any app installing a target-scoped
144/// filter must include `eyes::measurement=info`.
145pub const MEASUREMENT_TARGET: &str = "eyes::measurement";
146
147/// The measurement wire-contract version this subscriber emits.
148pub const MEASUREMENT_VERSION: u64 = 1;
149
150/// Field names lifted out of the dimension bag into the measurement contract.
151///
152/// A field with one of these names is part of the instrument, not a dimension.
153/// `version` is deliberately absent: the layer writes the discriminator itself
154/// and never lifts a caller-supplied value, so an app is free to use `version`
155/// as an ordinary dimension.
156pub const MEASUREMENT_RESERVED_FIELDS: [&str; 5] =
157    ["metric_name", "metric_kind", "value", "unit", "description"];
158
159/// Emits a gauge: an instantaneous value.
160///
161/// A non-finite value (NaN, ±infinity) cannot be represented in JSON and is
162/// dropped by the layer, counted in its drop accounting.
163pub fn emit_gauge(metric_name: &str, value: f64) {
164    tracing::event!(
165        target: "eyes::measurement",
166        tracing::Level::INFO,
167        metric_name = metric_name,
168        metric_kind = "gauge",
169        value = value,
170    );
171}
172
173/// Emits a counter: a cumulative, monotonically non-decreasing total.
174///
175/// Takes `u64` so a negative cumulative total is unrepresentable rather than a
176/// runtime 400. A value above `i64::MAX` is carried as a float by the server's
177/// value tagger and is therefore lossy above 2^53.
178pub fn emit_counter(metric_name: &str, value: u64) {
179    tracing::event!(
180        target: "eyes::measurement",
181        tracing::Level::INFO,
182        metric_name = metric_name,
183        metric_kind = "counter",
184        value = value,
185    );
186}
187
188/// Emits a sample: one discrete observation.
189///
190/// A non-finite value (NaN, ±infinity) cannot be represented in JSON and is
191/// dropped by the layer, counted in its drop accounting.
192pub fn emit_sample(metric_name: &str, value: f64) {
193    tracing::event!(
194        target: "eyes::measurement",
195        tracing::Level::INFO,
196        metric_name = metric_name,
197        metric_kind = "sample",
198        value = value,
199    );
200}
201
202/// Emits a typed measurement with dimensions.
203///
204/// ```no_run
205/// eyes_subscriber::measurement!(
206///     "gauge", "cpu", 42.5_f64,
207///     unit = "percent", description = "CPU utilisation", host = "web-1"
208/// );
209/// ```
210///
211/// Dimension tokens pass straight through to `tracing`'s field grammar, so
212/// dotted names and sigils (`?err`, `%value`) work. Names in
213/// [`MEASUREMENT_RESERVED_FIELDS`] are lifted into the measurement contract and
214/// cannot be used as dimensions; `version` is not reserved.
215#[macro_export]
216macro_rules! measurement {
217    ($kind:expr, $name:expr, $value:expr $(,)?) => {
218        $crate::tracing::event!(
219            target: "eyes::measurement",
220            $crate::tracing::Level::INFO,
221            metric_name = $name,
222            metric_kind = $kind,
223            value = $value,
224        )
225    };
226    ($kind:expr, $name:expr, $value:expr, $($dimensions:tt)+) => {
227        $crate::tracing::event!(
228            target: "eyes::measurement",
229            $crate::tracing::Level::INFO,
230            metric_name = $name,
231            metric_kind = $kind,
232            value = $value,
233            $($dimensions)+
234        )
235    };
236}
237
238#[derive(Debug, Clone, Serialize, Deserialize)]
239pub(crate) struct EventData {
240    event_type: String,
241    event_data: Value,
242    event_timestamp: DateTime<Utc>,
243    #[serde(default)]
244    process_instance_id: Option<Uuid>,
245}
246
247/// Default capacity of the bounded event queue between the tracing layer and
248/// the transport loop. Overridable via `EYES_QUEUE_CAPACITY` or
249/// [`EyesSubscriberBuilder::with_queue_capacity`].
250pub const DEFAULT_QUEUE_CAPACITY: usize = 65_536;
251
252#[derive(Debug, Clone)]
253pub struct EyesLayer {
254    sender: mpsc::Sender<EventData>,
255    dropped: Arc<AtomicU64>,
256    emit_enter_exit: bool,
257    process_instance_id: Option<Uuid>,
258}
259
260impl EyesLayer {
261    /// Enqueue an event without ever blocking the tracing hot path. If the
262    /// queue is full the event is dropped and counted; the transport loop
263    /// reports drops as a synthetic WARN event once the pipeline recovers.
264    fn dispatch(&self, event: EventData) {
265        match self.sender.try_send(event) {
266            Ok(()) => {}
267            Err(mpsc::error::TrySendError::Full(_)) => {
268                self.dropped.fetch_add(1, Ordering::Relaxed);
269            }
270            Err(mpsc::error::TrySendError::Closed(_)) => {}
271        }
272    }
273}
274
275#[derive(Debug)]
276pub struct EyesShutdownHandle {
277    shutdown_tx: oneshot::Sender<()>,
278    completion_rx: oneshot::Receiver<()>,
279}
280
281impl EyesShutdownHandle {
282    pub async fn shutdown(self) -> Result<(), Box<dyn std::error::Error>> {
283        let _ = self.shutdown_tx.send(());
284        self.completion_rx.await?;
285        Ok(())
286    }
287}
288
289/// `Debug` rendering for a bearer token that never prints the secret.
290///
291/// The types carrying a token are public API in a published crate, and
292/// downstream apps debug-print their configuration freely. One
293/// `tracing::debug!(?config)` in a cja app would otherwise land a live token
294/// in that app's telemetry — which for these apps means Eyes' own event store.
295pub(crate) struct RedactedToken<'a>(pub(crate) Option<&'a str>);
296
297impl std::fmt::Debug for RedactedToken<'_> {
298    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
299        match self.0 {
300            None => f.write_str("None"),
301            Some(token) => {
302                let prefix: String = token.chars().take(12).collect();
303                write!(f, "Some(\"{prefix}\u{2026}\")")
304            }
305        }
306    }
307}
308
309#[derive(Clone)]
310pub struct EyesSubscriberBuilder {
311    base_url: Url,
312    org_id: Uuid,
313    app_id: Uuid,
314    queue_capacity: Option<usize>,
315    emit_enter_exit: Option<bool>,
316    process_instance_id: Option<Uuid>,
317    auth_token: Option<String>,
318}
319
320impl std::fmt::Debug for EyesSubscriberBuilder {
321    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
322        f.debug_struct("EyesSubscriberBuilder")
323            .field("base_url", &self.base_url)
324            .field("org_id", &self.org_id)
325            .field("app_id", &self.app_id)
326            .field("queue_capacity", &self.queue_capacity)
327            .field("emit_enter_exit", &self.emit_enter_exit)
328            .field("process_instance_id", &self.process_instance_id)
329            .field("auth_token", &RedactedToken(self.auth_token.as_deref()))
330            .finish()
331    }
332}
333
334#[derive(Debug, Clone, Copy, PartialEq, Eq)]
335pub enum TransportType {
336    /// Standard HTTP transport - one request per event
337    Http,
338    /// Batching HTTP transport - buffers events and sends in batches
339    BatchingHttp,
340    /// WebSocket transport - persistent connection
341    WebSocket,
342}
343
344impl EyesSubscriberBuilder {
345    pub fn new(
346        base_url: impl Into<String>,
347        org_id: Uuid,
348        app_id: Uuid,
349    ) -> Result<Self, url::ParseError> {
350        Ok(Self {
351            base_url: Url::parse(&base_url.into())?,
352            org_id,
353            app_id,
354            queue_capacity: None,
355            emit_enter_exit: None,
356            process_instance_id: None,
357            auth_token: None,
358        })
359    }
360
361    /// Associate emitted telemetry with the boot-stable [`ProcessIdentity`].
362    pub fn with_process_instance_id(mut self, instance_id: Uuid) -> Self {
363        self.process_instance_id = Some(instance_id);
364        self
365    }
366
367    /// Set the capacity of the bounded event queue.
368    ///
369    /// Defaults to [`DEFAULT_QUEUE_CAPACITY`], overridable via the
370    /// `EYES_QUEUE_CAPACITY` environment variable. When the queue is full,
371    /// new events are dropped rather than blocking the host application.
372    pub fn with_queue_capacity(mut self, capacity: usize) -> Self {
373        self.queue_capacity = Some(capacity);
374        self
375    }
376
377    /// Enable or disable emission of `span_enter`/`span_exit` events.
378    ///
379    /// Disabled by default: tokio's tracing enters and exits a span on every
380    /// poll of an instrumented future, so long-lived instrumented loops emit a
381    /// steady stream of enter/exit noise. Span durations are computed from
382    /// `span_new`/`span_close`, which are always emitted, and enter/exit
383    /// pairs are aggregated locally into `busy_ms`/`poll_count` fields on
384    /// `span_close` regardless of this setting.
385    ///
386    /// Can also be enabled via the `EYES_EMIT_ENTER_EXIT` environment variable
387    /// ("1" or "true", case-insensitive), read once when the layer is built.
388    /// An explicit call to this method wins over the environment variable.
389    pub fn with_emit_enter_exit(mut self, emit: bool) -> Self {
390        self.emit_enter_exit = Some(emit);
391        self
392    }
393
394    fn resolve_emit_enter_exit(&self) -> bool {
395        self.emit_enter_exit.unwrap_or_else(|| {
396            std::env::var("EYES_EMIT_ENTER_EXIT")
397                .map(|v| {
398                    let v = v.to_lowercase();
399                    v == "1" || v == "true"
400                })
401                .unwrap_or(false)
402        })
403    }
404
405    /// Set the bearer token sent with every event request.
406    ///
407    /// Defaults to the `EYES_TOKEN` environment variable, read once when the
408    /// layer is built. An explicit call wins over the environment.
409    pub fn with_auth_token(mut self, token: impl Into<String>) -> Self {
410        self.auth_token = Some(token.into());
411        self
412    }
413
414    /// Resolved at build time rather than in `from_env`, so every construction
415    /// path — `build_from_env`, `from_env_with_transport`, and a hand-built
416    /// builder — picks up `EYES_TOKEN` without its own env read.
417    fn resolve_auth_token(&self) -> Option<String> {
418        self.auth_token.clone().or_else(|| {
419            std::env::var("EYES_TOKEN")
420                .ok()
421                .map(|s| s.trim().to_string())
422                .filter(|s| !s.is_empty())
423        })
424    }
425
426    fn resolve_queue_capacity(&self) -> usize {
427        self.queue_capacity
428            .or_else(|| {
429                std::env::var("EYES_QUEUE_CAPACITY")
430                    .ok()
431                    .and_then(|v| v.parse().ok())
432            })
433            .unwrap_or(DEFAULT_QUEUE_CAPACITY)
434            .max(1)
435    }
436
437    /// Create a new builder with the default production URL (eyes.coreyja.com)
438    pub fn new_with_default(org_id: Uuid, app_id: Uuid) -> Result<Self, url::ParseError> {
439        Self::new("https://eyes.coreyja.com", org_id, app_id)
440    }
441
442    /// Create a new builder, checking environment variables for configuration
443    ///
444    /// Checks the following environment variables:
445    /// - `EYES_URL`: Base URL for the Eyes server (defaults to https://eyes.coreyja.com)
446    /// - `EYES_TRANSPORT`: Transport type - "http", "batching" or "websocket" (defaults to "http")
447    ///
448    /// Returns a tuple of (builder, transport_type) to allow customization
449    pub fn from_env_with_transport(
450        org_id: Uuid,
451        app_id: Uuid,
452    ) -> Result<(Self, TransportType), url::ParseError> {
453        let base_url =
454            std::env::var("EYES_URL").unwrap_or_else(|_| "https://eyes.coreyja.com".to_string());
455
456        let transport = match std::env::var("EYES_TRANSPORT")
457            .unwrap_or_else(|_| "http".to_string())
458            .to_lowercase()
459            .as_str()
460        {
461            "websocket" | "ws" => TransportType::WebSocket,
462            "batching" | "batch" | "batching_http" => TransportType::BatchingHttp,
463            _ => TransportType::Http,
464        };
465
466        Ok((Self::new(base_url, org_id, app_id)?, transport))
467    }
468
469    /// Create a new builder, checking environment variables for configuration
470    ///
471    /// Checks the following in order:
472    /// 1. EYES_URL environment variable
473    /// 2. Falls back to https://eyes.coreyja.com
474    ///
475    /// Uses HTTP transport by default. For transport configuration, use `from_env_with_transport`
476    pub fn from_env(org_id: Uuid, app_id: Uuid) -> Result<Self, url::ParseError> {
477        let base_url =
478            std::env::var("EYES_URL").unwrap_or_else(|_| "https://eyes.coreyja.com".to_string());
479        Self::new(base_url, org_id, app_id)
480    }
481
482    pub fn build(self) -> (EyesLayer, EyesShutdownHandle) {
483        self.build_with_transport(TransportType::Http)
484    }
485
486    /// Build directly from environment variables in one step
487    ///
488    /// This is a convenience method that combines `from_env_with_transport` and `build_with_transport`.
489    ///
490    /// Environment variables:
491    /// - `EYES_URL`: Base URL (defaults to https://eyes.coreyja.com)
492    /// - `EYES_TRANSPORT`: Transport type - "http" or "websocket" (defaults to "http")
493    pub fn build_from_env(
494        org_id: Uuid,
495        app_id: Uuid,
496    ) -> Result<(EyesLayer, EyesShutdownHandle), url::ParseError> {
497        let (builder, transport) = Self::from_env_with_transport(org_id, app_id)?;
498        Ok(builder.build_with_transport(transport))
499    }
500
501    pub fn build_with_transport(
502        self,
503        transport_type: TransportType,
504    ) -> (EyesLayer, EyesShutdownHandle) {
505        self.build_with_transport_and_config(transport_type, BatchConfig::default())
506    }
507
508    /// Build with a specific transport type and batch configuration
509    ///
510    /// The batch config is only used when `transport_type` is `BatchingHttp`.
511    pub fn build_with_transport_and_config(
512        self,
513        transport_type: TransportType,
514        batch_config: BatchConfig,
515    ) -> (EyesLayer, EyesShutdownHandle) {
516        let emit_enter_exit = self.resolve_emit_enter_exit();
517        let auth_token = self.resolve_auth_token();
518        let (sender, receiver) = mpsc::channel::<EventData>(self.resolve_queue_capacity());
519        let (shutdown_tx, shutdown_rx) = oneshot::channel();
520        let (completion_tx, completion_rx) = oneshot::channel();
521        let dropped = Arc::new(AtomicU64::new(0));
522
523        let transport: Box<dyn Transport> = match transport_type {
524            TransportType::Http => Box::new(
525                HttpTransport::new(self.base_url.clone(), self.org_id, self.app_id, auth_token)
526                    .expect("Failed to create HTTP transport"),
527            ),
528            TransportType::BatchingHttp => Box::new(
529                BatchingHttpTransport::new(
530                    self.base_url.clone(),
531                    self.org_id,
532                    self.app_id,
533                    batch_config,
534                    auth_token,
535                )
536                .expect("Failed to create batching HTTP transport"),
537            ),
538            TransportType::WebSocket => Box::new(
539                WebSocketTransport::new(
540                    self.base_url.clone(),
541                    self.org_id,
542                    self.app_id,
543                    auth_token,
544                )
545                .expect("Failed to create WebSocket transport"),
546            ),
547        };
548
549        // Spawn background task to send events
550        tokio::spawn(transport::run_transport_loop(
551            transport,
552            receiver,
553            shutdown_rx,
554            completion_tx,
555            dropped.clone(),
556            transport::TransportLoopConfig::default(),
557        ));
558
559        let layer = EyesLayer {
560            sender,
561            dropped,
562            emit_enter_exit,
563            process_instance_id: self.process_instance_id,
564        };
565        let handle = EyesShutdownHandle {
566            shutdown_tx,
567            completion_rx,
568        };
569
570        (layer, handle)
571    }
572}
573
574/// Globally unique span id stored in the span's extensions.
575///
576/// The tracing registry reuses its numeric span ids aggressively (every
577/// process restart begins again at Id(1)), so registry ids collide across
578/// restarts and processes. Each span instead gets a random 32-hex-char id at
579/// creation; all lifecycle events and parent references resolve through it.
580struct EyesSpanId(String);
581
582/// Per-span busy-time aggregation stored in the span's extensions.
583///
584/// tokio's tracing enters and exits a span on every poll of an instrumented
585/// future, so instead of emitting an event per enter/exit we accumulate the
586/// time spent inside the span locally (the way tracing-subscriber's timing
587/// layer and tokio-console do) and report it once on `span_close` as
588/// `busy_ms` alongside `poll_count`. This runs regardless of whether
589/// `span_enter`/`span_exit` event emission is enabled.
590#[derive(Default)]
591struct BusyTimings {
592    last_enter: Option<Instant>,
593    busy: Duration,
594    poll_count: u64,
595}
596
597/// Fields recorded on a span after creation (`span.record(..)`), stored in the
598/// span's extensions until `span_close`.
599///
600/// tracing only delivers a span's creation-time attributes to `on_new_span`;
601/// anything recorded later (tower-http's `http.response.status_code` at
602/// response time, a job worker's `job.id` once a job is claimed) arrives via
603/// `on_record`. We accumulate those here — later records overwrite earlier
604/// values for the same field name — and merge them into the `span_close`
605/// event's `fields` object, so `span_new` carries creation-time attributes
606/// and `span_close` carries the final recorded state. The extension is only
607/// inserted once a span actually records something, so spans that never call
608/// `record` pay nothing.
609struct RecordedFields(serde_json::Map<String, Value>);
610
611fn generate_span_id() -> String {
612    Uuid::new_v4().simple().to_string()
613}
614
615/// Read a span's unique id from its extensions, falling back to the debug
616/// format of the registry id for spans created before this layer was attached.
617fn eyes_span_id<S>(span: &tracing_subscriber::registry::SpanRef<'_, S>) -> String
618where
619    S: Subscriber + for<'a> LookupSpan<'a>,
620{
621    span.extensions()
622        .get::<EyesSpanId>()
623        .map(|eyes_id| eyes_id.0.clone())
624        .unwrap_or_else(|| format!("{:?}", span.id()))
625}
626
627impl EyesLayer {
628    /// Builds the versioned measurement payload for an event on the reserved
629    /// target.
630    ///
631    /// The five reserved names are lifted out of the recorded fields and onto
632    /// the top level of `event_data`; whatever is left is the dimension bag.
633    /// `JsonVisitor` stores integers as JSON integers and floats as JSON
634    /// floats, so an integer counter stays an integer on the wire (the server's
635    /// tagger reads it as `int`) and an `f64` gauge stays a float — which is
636    /// the whole reason the contract carries a bare number.
637    /// Returns `None` when the lifted `value` is not a JSON number — a NaN or
638    /// infinite `f64` reaches here as `Value::Null` (`record_f64` cannot
639    /// represent it in JSON), and the server would reject the payload with a
640    /// 400 the transports cannot report. Dropping it here keeps the loss inside
641    /// the subscriber's own drop accounting instead of nowhere.
642    fn measurement_event<S>(
643        &self,
644        event: &Event<'_>,
645        ctx: &Context<'_, S>,
646        mut fields: serde_json::Map<String, Value>,
647    ) -> Option<EventData>
648    where
649        S: Subscriber + for<'a> LookupSpan<'a>,
650    {
651        let lifted: Vec<(&str, Option<Value>)> = MEASUREMENT_RESERVED_FIELDS
652            .iter()
653            .map(|name| (*name, fields.remove(*name)))
654            .collect();
655        if !lifted
656            .iter()
657            .any(|(name, value)| *name == "value" && matches!(value, Some(Value::Number(_))))
658        {
659            return None;
660        }
661
662        let mut event_data = serde_json::json!({
663            "version": MEASUREMENT_VERSION,
664            // Display form ("INFO"), never Debug ("Level(Info)") — see the span
665            // serialization above.
666            "level": event.metadata().level().to_string(),
667            "target": event.metadata().target(),
668            "fields": fields,
669        });
670        for (name, value) in lifted {
671            if let Some(value) = value {
672                event_data[name] = value;
673            }
674        }
675
676        // `eyes_span_id` reads a random UUID stashed in the span's extensions.
677        // `span.id()` is a per-process counter that resets to `Id(1)` on
678        // restart and must never reach the wire.
679        if let Some(span) = ctx.event_span(event) {
680            event_data["span_id"] = serde_json::json!(eyes_span_id(&span));
681        }
682
683        Some(EventData {
684            event_type: "measurement".to_string(),
685            event_data,
686            event_timestamp: Utc::now(),
687            process_instance_id: self.process_instance_id,
688        })
689    }
690}
691
692impl<S> Layer<S> for EyesLayer
693where
694    S: Subscriber + for<'a> LookupSpan<'a>,
695{
696    fn on_new_span(&self, attrs: &span::Attributes<'_>, id: &Id, ctx: Context<'_, S>) {
697        let span = ctx.span(id).expect("Span not found");
698
699        let unique_id = generate_span_id();
700        span.extensions_mut().insert(EyesSpanId(unique_id.clone()));
701
702        let mut visitor = JsonVisitor::default();
703        attrs.record(&mut visitor);
704
705        let mut event_data = serde_json::json!({
706            "span_id": unique_id,
707            "name": span.metadata().name(),
708            "target": span.metadata().target(),
709            // Display form ("INFO"), not Debug ("Level(Info)"): the eyes
710            // server filters/aggregates on these exact strings (`level =
711            // 'INFO'`, `level = 'ERROR'`), so the Debug wrapper silently
712            // broke cron fire detection, log level filters, and job-error
713            // detection.
714            "level": span.metadata().level().to_string(),
715            "fields": visitor.fields,
716        });
717
718        // Get parent from either explicit parent OR contextual parent (current span).
719        // span.parent() only returns explicitly-set parents, but #[tracing::instrument]
720        // uses contextual parents via the current span stack.
721        let parent = span
722            .parent()
723            .or_else(|| ctx.current_span().id().and_then(|pid| ctx.span(pid)));
724
725        if let Some(parent) = parent {
726            event_data["parent_id"] = serde_json::json!(eyes_span_id(&parent));
727        }
728
729        let event = EventData {
730            event_type: "span_new".to_string(),
731            event_data,
732            event_timestamp: Utc::now(),
733            process_instance_id: self.process_instance_id,
734        };
735
736        self.dispatch(event);
737    }
738
739    fn on_record(&self, id: &Id, values: &span::Record<'_>, ctx: Context<'_, S>) {
740        let Some(span) = ctx.span(id) else {
741            return;
742        };
743
744        let mut visitor = JsonVisitor::default();
745        values.record(&mut visitor);
746        if visitor.fields.is_empty() {
747            return;
748        }
749
750        let mut extensions = span.extensions_mut();
751        if let Some(recorded) = extensions.get_mut::<RecordedFields>() {
752            // Later records overwrite earlier values for the same field name.
753            recorded.0.extend(visitor.fields);
754        } else {
755            extensions.insert(RecordedFields(visitor.fields));
756        }
757    }
758
759    fn on_event(&self, event: &Event<'_>, ctx: Context<'_, S>) {
760        let mut visitor = JsonVisitor::default();
761        event.record(&mut visitor);
762
763        if event.metadata().target() == MEASUREMENT_TARGET {
764            match self.measurement_event(event, &ctx, visitor.fields) {
765                Some(measurement) => self.dispatch(measurement),
766                // Non-numeric value (NaN/infinity lands here as JSON null):
767                // count it with the queue-full drops so the loss is visible.
768                None => {
769                    self.dropped.fetch_add(1, Ordering::Relaxed);
770                }
771            }
772            return;
773        }
774
775        let mut event_data = serde_json::json!({
776            // Display form ("INFO"), not Debug ("Level(Info)") — see the span
777            // serialization above.
778            "level": event.metadata().level().to_string(),
779            "target": event.metadata().target(),
780            "fields": visitor.fields,
781        });
782
783        if let Some(span) = ctx.event_span(event) {
784            event_data["span_id"] = serde_json::json!(eyes_span_id(&span));
785        }
786
787        let event_msg = EventData {
788            event_type: "event".to_string(),
789            event_data,
790            event_timestamp: Utc::now(),
791            process_instance_id: self.process_instance_id,
792        };
793
794        self.dispatch(event_msg);
795    }
796
797    fn on_enter(&self, id: &Id, ctx: Context<'_, S>) {
798        // Busy-time aggregation runs unconditionally: the emit_enter_exit flag
799        // only gates event emission, not timing.
800        if let Some(span) = ctx.span(id) {
801            let mut extensions = span.extensions_mut();
802            if let Some(timings) = extensions.get_mut::<BusyTimings>() {
803                timings.last_enter = Some(Instant::now());
804            } else {
805                extensions.insert(BusyTimings {
806                    last_enter: Some(Instant::now()),
807                    ..BusyTimings::default()
808                });
809            }
810        }
811
812        if !self.emit_enter_exit {
813            return;
814        }
815        let span_id = ctx
816            .span(id)
817            .map(|span| eyes_span_id(&span))
818            .unwrap_or_else(|| format!("{:?}", id));
819        let event = EventData {
820            event_type: "span_enter".to_string(),
821            event_data: serde_json::json!({
822                "span_id": span_id,
823            }),
824            event_timestamp: Utc::now(),
825            process_instance_id: self.process_instance_id,
826        };
827
828        self.dispatch(event);
829    }
830
831    fn on_exit(&self, id: &Id, ctx: Context<'_, S>) {
832        // Busy-time aggregation runs unconditionally: the emit_enter_exit flag
833        // only gates event emission, not timing.
834        if let Some(span) = ctx.span(id) {
835            let mut extensions = span.extensions_mut();
836            if let Some(timings) = extensions.get_mut::<BusyTimings>() {
837                if let Some(entered_at) = timings.last_enter.take() {
838                    timings.busy += entered_at.elapsed();
839                    timings.poll_count += 1;
840                }
841            }
842        }
843
844        if !self.emit_enter_exit {
845            return;
846        }
847        let span_id = ctx
848            .span(id)
849            .map(|span| eyes_span_id(&span))
850            .unwrap_or_else(|| format!("{:?}", id));
851        let event = EventData {
852            event_type: "span_exit".to_string(),
853            event_data: serde_json::json!({
854                "span_id": span_id,
855            }),
856            event_timestamp: Utc::now(),
857            process_instance_id: self.process_instance_id,
858        };
859
860        self.dispatch(event);
861    }
862
863    fn on_close(&self, id: Id, ctx: Context<'_, S>) {
864        let span = ctx.span(&id).expect("Span not found");
865
866        // Report the locally aggregated busy time. A span that was never
867        // entered has no BusyTimings extension; report zeros so the fields
868        // are always present and queryable.
869        let (busy_ms, poll_count) = span
870            .extensions()
871            .get::<BusyTimings>()
872            .map(|timings| {
873                (
874                    u64::try_from(timings.busy.as_millis()).unwrap_or(u64::MAX),
875                    timings.poll_count,
876                )
877            })
878            .unwrap_or((0, 0));
879
880        // Start from any fields recorded after creation via `span.record(..)`
881        // (see RecordedFields); the extension is absent for spans that never
882        // recorded anything, in which case this allocates nothing.
883        let mut fields = span
884            .extensions_mut()
885            .remove::<RecordedFields>()
886            .map(|recorded| recorded.0)
887            .unwrap_or_default();
888
889        // Inserted after the recorded fields on purpose: busy_ms/poll_count
890        // are eyes-owned aggregation keys on span_close, so if a recorded
891        // field happens to share one of these names the aggregation values
892        // win.
893        fields.insert("busy_ms".to_string(), Value::from(busy_ms));
894        fields.insert("poll_count".to_string(), Value::from(poll_count));
895
896        let event = EventData {
897            event_type: "span_close".to_string(),
898            event_data: serde_json::json!({
899                "span_id": eyes_span_id(&span),
900                "name": span.metadata().name(),
901                "fields": fields,
902            }),
903            event_timestamp: Utc::now(),
904            process_instance_id: self.process_instance_id,
905        };
906
907        self.dispatch(event);
908    }
909}
910
911#[derive(Default)]
912struct JsonVisitor {
913    fields: serde_json::Map<String, Value>,
914}
915
916impl Visit for JsonVisitor {
917    fn record_debug(&mut self, field: &tracing::field::Field, value: &dyn std::fmt::Debug) {
918        self.fields.insert(
919            field.name().to_string(),
920            Value::String(format!("{:?}", value)),
921        );
922    }
923
924    fn record_str(&mut self, field: &tracing::field::Field, value: &str) {
925        self.fields
926            .insert(field.name().to_string(), Value::String(value.to_string()));
927    }
928
929    fn record_i64(&mut self, field: &tracing::field::Field, value: i64) {
930        self.fields
931            .insert(field.name().to_string(), Value::Number(value.into()));
932    }
933
934    fn record_u64(&mut self, field: &tracing::field::Field, value: u64) {
935        self.fields
936            .insert(field.name().to_string(), Value::Number(value.into()));
937    }
938
939    fn record_bool(&mut self, field: &tracing::field::Field, value: bool) {
940        self.fields
941            .insert(field.name().to_string(), Value::Bool(value));
942    }
943
944    fn record_f64(&mut self, field: &tracing::field::Field, value: f64) {
945        self.fields.insert(
946            field.name().to_string(),
947            serde_json::Number::from_f64(value)
948                .map(Value::Number)
949                .unwrap_or(Value::Null),
950        );
951    }
952}
953
954#[cfg(test)]
955mod tests {
956    use super::*;
957    use tracing::{info, span, Level};
958    use tracing_subscriber::layer::SubscriberExt;
959
960    #[test]
961    fn test_builder_creation() {
962        let org_id = Uuid::new_v4();
963        let app_id = Uuid::new_v4();
964        let builder = EyesSubscriberBuilder::new("http://localhost:4318", org_id, app_id).unwrap();
965        assert_eq!(builder.app_id, app_id);
966        assert_eq!(builder.org_id, org_id);
967    }
968
969    #[test]
970    fn test_builder_invalid_url() {
971        let org_id = Uuid::new_v4();
972        let app_id = Uuid::new_v4();
973        let result = EyesSubscriberBuilder::new("invalid-url", org_id, app_id);
974        assert!(result.is_err());
975    }
976
977    #[test]
978    fn test_builder_new_with_default() {
979        let org_id = Uuid::new_v4();
980        let app_id = Uuid::new_v4();
981        let builder = EyesSubscriberBuilder::new_with_default(org_id, app_id).unwrap();
982        assert_eq!(builder.app_id, app_id);
983        assert_eq!(builder.org_id, org_id);
984        // URL should be set to production default
985        assert_eq!(builder.base_url.as_str(), "https://eyes.coreyja.com/");
986    }
987
988    #[test]
989    fn test_http_transport_creation() {
990        let org_id = Uuid::new_v4();
991        let app_id = Uuid::new_v4();
992        let base_url = Url::parse("http://localhost:4318").unwrap();
993        let transport = HttpTransport::new(base_url, org_id, app_id, None);
994        assert!(transport.is_ok());
995    }
996
997    #[test]
998    fn test_websocket_transport_creation() {
999        let org_id = Uuid::new_v4();
1000        let app_id = Uuid::new_v4();
1001        let base_url = Url::parse("http://localhost:4318").unwrap();
1002        let transport = WebSocketTransport::new(base_url, org_id, app_id, None);
1003        assert!(transport.is_ok());
1004    }
1005
1006    #[test]
1007    fn test_websocket_url_conversion() {
1008        let org_id = Uuid::new_v4();
1009        let app_id = Uuid::new_v4();
1010        let https_url = Url::parse("https://example.com").unwrap();
1011        let _transport = WebSocketTransport::new(https_url, org_id, app_id, None).unwrap();
1012        // The URL should be converted internally to wss://
1013    }
1014
1015    #[test]
1016    fn test_event_data_serialization() {
1017        let event = EventData {
1018            event_type: "test_event".to_string(),
1019            event_data: serde_json::json!({"key": "value", "number": 42}),
1020            event_timestamp: Utc::now(),
1021            process_instance_id: None,
1022        };
1023
1024        let serialized = serde_json::to_string(&event).unwrap();
1025        let deserialized: EventData = serde_json::from_str(&serialized).unwrap();
1026
1027        assert_eq!(event.event_type, deserialized.event_type);
1028        assert_eq!(event.event_data, deserialized.event_data);
1029    }
1030
1031    #[test]
1032    fn test_json_visitor_basic_functionality() {
1033        let mut visitor = JsonVisitor::default();
1034
1035        // Test that visitor starts empty
1036        assert_eq!(visitor.fields.len(), 0);
1037
1038        // Test that we can add fields
1039        visitor.fields.insert(
1040            "test_key".to_string(),
1041            Value::String("test_value".to_string()),
1042        );
1043        assert_eq!(visitor.fields.len(), 1);
1044        assert_eq!(
1045            visitor.fields.get("test_key"),
1046            Some(&Value::String("test_value".to_string()))
1047        );
1048    }
1049
1050    #[test]
1051    fn test_transport_type_debug() {
1052        let http = TransportType::Http;
1053        let ws = TransportType::WebSocket;
1054
1055        assert_eq!(format!("{:?}", http), "Http");
1056        assert_eq!(format!("{:?}", ws), "WebSocket");
1057    }
1058
1059    #[test]
1060    fn test_transport_type_equality() {
1061        assert_eq!(TransportType::Http, TransportType::Http);
1062        assert_eq!(TransportType::WebSocket, TransportType::WebSocket);
1063        assert_ne!(TransportType::Http, TransportType::WebSocket);
1064    }
1065
1066    #[tokio::test]
1067    async fn test_span_ids_unique_across_registry_reuse() {
1068        let (sender, mut receiver) = mpsc::channel::<EventData>(64);
1069        let layer = EyesLayer {
1070            sender,
1071            dropped: std::sync::Arc::new(std::sync::atomic::AtomicU64::new(0)),
1072            // Opt in so the enter/exit lifecycle assertions below stay covered.
1073            emit_enter_exit: true,
1074            process_instance_id: None,
1075        };
1076        let subscriber = tracing_subscriber::registry().with(layer);
1077
1078        tracing::subscriber::with_default(subscriber, || {
1079            {
1080                let parent = span!(Level::INFO, "parent_span");
1081                let _parent_guard = parent.enter();
1082                let child = span!(Level::INFO, "child_span");
1083                let _child_guard = child.enter();
1084                info!("inside child");
1085            }
1086            // Both spans closed: the registry may now reuse their numeric ids.
1087            {
1088                let reused = span!(Level::INFO, "reused_slot_span");
1089                let _guard = reused.enter();
1090            }
1091        });
1092
1093        let mut events = Vec::new();
1094        while let Ok(event) = receiver.try_recv() {
1095            events.push(event);
1096        }
1097
1098        let span_id_of = |name: &str| -> String {
1099            events
1100                .iter()
1101                .find(|e| e.event_type == "span_new" && e.event_data["name"] == name)
1102                .unwrap_or_else(|| panic!("no span_new for {name}"))
1103                .event_data["span_id"]
1104                .as_str()
1105                .unwrap()
1106                .to_string()
1107        };
1108
1109        let parent_id = span_id_of("parent_span");
1110        let child_id = span_id_of("child_span");
1111        let reused_id = span_id_of("reused_slot_span");
1112
1113        // Unique ids are 32 hex chars, never the registry debug format.
1114        for id in [&parent_id, &child_id, &reused_id] {
1115            assert_eq!(id.len(), 32, "unexpected id shape: {id}");
1116            assert!(!id.starts_with("Id("), "registry id leaked: {id}");
1117        }
1118        assert_ne!(parent_id, child_id);
1119        // The regression: a reused registry slot must still get a fresh id.
1120        assert_ne!(reused_id, parent_id);
1121        assert_ne!(reused_id, child_id);
1122
1123        // The child's parent reference uses the parent's unique id.
1124        let child_new = events
1125            .iter()
1126            .find(|e| e.event_type == "span_new" && e.event_data["name"] == "child_span")
1127            .unwrap();
1128        assert_eq!(child_new.event_data["parent_id"], parent_id.as_str());
1129
1130        // Every lifecycle event for the child carries the same unique id.
1131        let child_lifecycle: Vec<_> = events
1132            .iter()
1133            .filter(|e| {
1134                matches!(
1135                    e.event_type.as_str(),
1136                    "span_enter" | "span_exit" | "span_close"
1137                ) && e.event_data["span_id"] == child_id.as_str()
1138            })
1139            .collect();
1140        assert!(
1141            child_lifecycle.len() >= 3,
1142            "expected enter/exit/close with the child's unique id, got {}",
1143            child_lifecycle.len()
1144        );
1145
1146        // The in-span event is attributed to the child's unique id.
1147        let in_span_event = events
1148            .iter()
1149            .find(|e| {
1150                e.event_type == "event" && e.event_data["fields"]["message"] == "inside child"
1151            })
1152            .unwrap();
1153        assert_eq!(in_span_event.event_data["span_id"], child_id.as_str());
1154    }
1155
1156    #[tokio::test]
1157    async fn test_layer_integration() {
1158        let org_id = Uuid::new_v4();
1159        let app_id = Uuid::new_v4();
1160
1161        // Create a builder and build the layer
1162        let (layer, shutdown_handle) =
1163            EyesSubscriberBuilder::new("http://localhost:4318", org_id, app_id)
1164                .unwrap()
1165                .build();
1166
1167        // Set up tracing with our layer
1168        let subscriber = tracing_subscriber::registry().with(layer);
1169
1170        // Use the subscriber in a limited scope
1171        tracing::subscriber::with_default(subscriber, || {
1172            let span = span!(Level::INFO, "test_span", user_id = 123);
1173            let _enter = span.enter();
1174            info!("Test message in span");
1175        });
1176
1177        // Shutdown gracefully
1178        shutdown_handle.shutdown().await.unwrap();
1179    }
1180
1181    #[tokio::test]
1182    async fn test_layer_with_websocket_transport() {
1183        let org_id = Uuid::new_v4();
1184        let app_id = Uuid::new_v4();
1185
1186        let (layer, shutdown_handle) =
1187            EyesSubscriberBuilder::new("http://localhost:4318", org_id, app_id)
1188                .unwrap()
1189                .build_with_transport(TransportType::WebSocket);
1190
1191        // Test that the layer was created successfully
1192        let subscriber = tracing_subscriber::registry().with(layer);
1193
1194        tracing::subscriber::with_default(subscriber, || {
1195            info!("Test WebSocket transport");
1196        });
1197
1198        shutdown_handle.shutdown().await.unwrap();
1199    }
1200
1201    #[tokio::test]
1202    async fn test_build_from_env_with_defaults() {
1203        // Clear environment
1204        std::env::remove_var("EYES_URL");
1205        std::env::remove_var("EYES_TRANSPORT");
1206
1207        let org_id = Uuid::new_v4();
1208        let app_id = Uuid::new_v4();
1209
1210        let result = EyesSubscriberBuilder::build_from_env(org_id, app_id);
1211        assert!(result.is_ok());
1212
1213        // Test shutdown
1214        if let Ok((_, shutdown_handle)) = result {
1215            shutdown_handle.shutdown().await.unwrap();
1216        }
1217    }
1218
1219    #[tokio::test]
1220    async fn test_build_from_env_with_custom_values() {
1221        std::env::set_var("EYES_URL", "http://custom.example.com");
1222        std::env::set_var("EYES_TRANSPORT", "websocket");
1223
1224        let org_id = Uuid::new_v4();
1225        let app_id = Uuid::new_v4();
1226
1227        let result = EyesSubscriberBuilder::build_from_env(org_id, app_id);
1228        assert!(result.is_ok());
1229
1230        // Test shutdown
1231        if let Ok((_, shutdown_handle)) = result {
1232            shutdown_handle.shutdown().await.unwrap();
1233        }
1234
1235        // Clean up
1236        std::env::remove_var("EYES_URL");
1237        std::env::remove_var("EYES_TRANSPORT");
1238    }
1239
1240    #[tokio::test]
1241    async fn test_layer_with_batching_transport() {
1242        let org_id = Uuid::new_v4();
1243        let app_id = Uuid::new_v4();
1244
1245        let (layer, shutdown_handle) =
1246            EyesSubscriberBuilder::new("http://localhost:4318", org_id, app_id)
1247                .unwrap()
1248                .build_with_transport(TransportType::BatchingHttp);
1249
1250        // Test that the layer was created successfully
1251        let subscriber = tracing_subscriber::registry().with(layer);
1252
1253        tracing::subscriber::with_default(subscriber, || {
1254            info!("Test batching HTTP transport");
1255        });
1256
1257        shutdown_handle.shutdown().await.unwrap();
1258    }
1259
1260    #[tokio::test]
1261    async fn test_layer_with_batching_transport_custom_config() {
1262        use std::time::Duration;
1263
1264        let org_id = Uuid::new_v4();
1265        let app_id = Uuid::new_v4();
1266
1267        let custom_config = BatchConfig::new(50, Duration::from_millis(100));
1268
1269        let (layer, shutdown_handle) =
1270            EyesSubscriberBuilder::new("http://localhost:4318", org_id, app_id)
1271                .unwrap()
1272                .build_with_transport_and_config(TransportType::BatchingHttp, custom_config);
1273
1274        let subscriber = tracing_subscriber::registry().with(layer);
1275
1276        tracing::subscriber::with_default(subscriber, || {
1277            info!("Test batching HTTP transport with custom config");
1278        });
1279
1280        shutdown_handle.shutdown().await.unwrap();
1281    }
1282
1283    #[test]
1284    fn test_batching_transport_creation() {
1285        let org_id = Uuid::new_v4();
1286        let app_id = Uuid::new_v4();
1287        let base_url = Url::parse("http://localhost:4318").unwrap();
1288
1289        let transport =
1290            BatchingHttpTransport::with_default_config(base_url.clone(), org_id, app_id, None);
1291        assert!(transport.is_ok());
1292
1293        use std::time::Duration;
1294        let custom_config = BatchConfig::new(50, Duration::from_millis(100));
1295        let transport = BatchingHttpTransport::new(base_url, org_id, app_id, custom_config, None);
1296        assert!(transport.is_ok());
1297    }
1298
1299    #[test]
1300    fn test_transport_type_batching_http() {
1301        assert_eq!(TransportType::BatchingHttp, TransportType::BatchingHttp);
1302        assert_ne!(TransportType::BatchingHttp, TransportType::Http);
1303        assert_ne!(TransportType::BatchingHttp, TransportType::WebSocket);
1304    }
1305
1306    #[test]
1307    fn test_dispatch_drops_and_counts_when_queue_full() {
1308        let (sender, mut receiver) = mpsc::channel(1);
1309        let layer = EyesLayer {
1310            sender,
1311            dropped: Arc::new(AtomicU64::new(0)),
1312            emit_enter_exit: false,
1313            process_instance_id: None,
1314        };
1315
1316        let event = |event_type: &str| EventData {
1317            event_type: event_type.to_string(),
1318            event_data: serde_json::json!({}),
1319            event_timestamp: Utc::now(),
1320            process_instance_id: None,
1321        };
1322
1323        layer.dispatch(event("first"));
1324        layer.dispatch(event("second"));
1325
1326        // The queue held the first event; the second was dropped and counted.
1327        assert_eq!(layer.dropped.load(Ordering::Relaxed), 1);
1328        assert_eq!(receiver.try_recv().unwrap().event_type, "first");
1329        assert!(receiver.try_recv().is_err());
1330    }
1331
1332    #[test]
1333    #[serial_test::serial]
1334    fn test_queue_capacity_default_and_builder() {
1335        std::env::remove_var("EYES_QUEUE_CAPACITY");
1336
1337        let org_id = Uuid::new_v4();
1338        let app_id = Uuid::new_v4();
1339        let builder = EyesSubscriberBuilder::new("http://localhost:4318", org_id, app_id).unwrap();
1340
1341        assert_eq!(builder.resolve_queue_capacity(), DEFAULT_QUEUE_CAPACITY);
1342        assert_eq!(
1343            builder
1344                .clone()
1345                .with_queue_capacity(123)
1346                .resolve_queue_capacity(),
1347            123
1348        );
1349        // Zero is clamped so try_send always has somewhere to go.
1350        assert_eq!(builder.with_queue_capacity(0).resolve_queue_capacity(), 1);
1351    }
1352
1353    #[test]
1354    #[serial_test::serial]
1355    fn test_queue_capacity_from_env() {
1356        let org_id = Uuid::new_v4();
1357        let app_id = Uuid::new_v4();
1358        let builder = EyesSubscriberBuilder::new("http://localhost:4318", org_id, app_id).unwrap();
1359
1360        std::env::set_var("EYES_QUEUE_CAPACITY", "1024");
1361        assert_eq!(builder.resolve_queue_capacity(), 1024);
1362        // An explicit builder setting wins over the environment.
1363        assert_eq!(
1364            builder
1365                .clone()
1366                .with_queue_capacity(123)
1367                .resolve_queue_capacity(),
1368            123
1369        );
1370
1371        // Unparseable values fall back to the default.
1372        std::env::set_var("EYES_QUEUE_CAPACITY", "not-a-number");
1373        assert_eq!(builder.resolve_queue_capacity(), DEFAULT_QUEUE_CAPACITY);
1374
1375        std::env::remove_var("EYES_QUEUE_CAPACITY");
1376    }
1377
1378    #[test]
1379    #[serial_test::serial]
1380    fn test_emit_enter_exit_default_off() {
1381        std::env::remove_var("EYES_EMIT_ENTER_EXIT");
1382
1383        let org_id = Uuid::new_v4();
1384        let app_id = Uuid::new_v4();
1385        let builder = EyesSubscriberBuilder::new("http://localhost:4318", org_id, app_id).unwrap();
1386
1387        assert!(!builder.resolve_emit_enter_exit());
1388        // The builder can opt in without the environment variable.
1389        assert!(builder.with_emit_enter_exit(true).resolve_emit_enter_exit());
1390    }
1391
1392    #[test]
1393    #[serial_test::serial]
1394    fn test_emit_enter_exit_from_env() {
1395        let org_id = Uuid::new_v4();
1396        let app_id = Uuid::new_v4();
1397        let builder = EyesSubscriberBuilder::new("http://localhost:4318", org_id, app_id).unwrap();
1398
1399        for enabled in ["1", "true", "TRUE", "True"] {
1400            std::env::set_var("EYES_EMIT_ENTER_EXIT", enabled);
1401            assert!(
1402                builder.resolve_emit_enter_exit(),
1403                "{enabled:?} should enable enter/exit emission"
1404            );
1405        }
1406
1407        for disabled in ["0", "false", "yes", ""] {
1408            std::env::set_var("EYES_EMIT_ENTER_EXIT", disabled);
1409            assert!(
1410                !builder.resolve_emit_enter_exit(),
1411                "{disabled:?} should not enable enter/exit emission"
1412            );
1413        }
1414
1415        std::env::remove_var("EYES_EMIT_ENTER_EXIT");
1416    }
1417
1418    #[test]
1419    #[serial_test::serial]
1420    fn test_emit_enter_exit_builder_overrides_env() {
1421        let org_id = Uuid::new_v4();
1422        let app_id = Uuid::new_v4();
1423        let builder = EyesSubscriberBuilder::new("http://localhost:4318", org_id, app_id).unwrap();
1424
1425        // An explicit builder setting wins over the environment.
1426        std::env::set_var("EYES_EMIT_ENTER_EXIT", "1");
1427        assert!(!builder
1428            .clone()
1429            .with_emit_enter_exit(false)
1430            .resolve_emit_enter_exit());
1431
1432        std::env::set_var("EYES_EMIT_ENTER_EXIT", "false");
1433        assert!(builder.with_emit_enter_exit(true).resolve_emit_enter_exit());
1434
1435        std::env::remove_var("EYES_EMIT_ENTER_EXIT");
1436    }
1437
1438    /// Build a bare layer + receiver for direct event inspection.
1439    fn test_layer(emit_enter_exit: bool) -> (EyesLayer, mpsc::Receiver<EventData>) {
1440        let (sender, receiver) = mpsc::channel::<EventData>(64);
1441        let layer = EyesLayer {
1442            sender,
1443            dropped: Arc::new(AtomicU64::new(0)),
1444            emit_enter_exit,
1445            process_instance_id: None,
1446        };
1447        (layer, receiver)
1448    }
1449
1450    fn drain_events(receiver: &mut mpsc::Receiver<EventData>) -> Vec<EventData> {
1451        let mut events = Vec::new();
1452        while let Ok(event) = receiver.try_recv() {
1453            events.push(event);
1454        }
1455        events
1456    }
1457
1458    fn span_close_for<'a>(events: &'a [EventData], name: &str) -> &'a EventData {
1459        events
1460            .iter()
1461            .find(|e| e.event_type == "span_close" && e.event_data["name"] == name)
1462            .unwrap_or_else(|| panic!("no span_close for {name}"))
1463    }
1464
1465    /// Enter/exit a span N times with a sleep per poll, regardless of the
1466    /// emission flag, and return the captured events.
1467    fn poll_span_n_times(
1468        layer: EyesLayer,
1469        receiver: &mut mpsc::Receiver<EventData>,
1470        polls: u32,
1471        sleep_per_poll: std::time::Duration,
1472    ) -> Vec<EventData> {
1473        let subscriber = tracing_subscriber::registry().with(layer);
1474        tracing::subscriber::with_default(subscriber, || {
1475            let span = span!(Level::INFO, "polled_span");
1476            for _ in 0..polls {
1477                let guard = span.enter();
1478                std::thread::sleep(sleep_per_poll);
1479                drop(guard);
1480            }
1481            drop(span);
1482        });
1483        drain_events(receiver)
1484    }
1485
1486    #[tokio::test]
1487    async fn test_level_serialized_as_display_not_debug() {
1488        // The eyes server filters/aggregates on the exact level string
1489        // (`level = 'INFO'`). Levels must serialize as the Display form
1490        // ("WARN"), never the Debug wrapper ("Level(Warn)").
1491        let (layer, mut receiver) = test_layer(false);
1492        let subscriber = tracing_subscriber::registry().with(layer);
1493        tracing::subscriber::with_default(subscriber, || {
1494            let _span = span!(Level::WARN, "warn_span");
1495            tracing::error!("boom");
1496        });
1497        let events = drain_events(&mut receiver);
1498
1499        let span_new = events
1500            .iter()
1501            .find(|e| e.event_type == "span_new" && e.event_data["name"] == "warn_span")
1502            .expect("span_new for warn_span");
1503        assert_eq!(span_new.event_data["level"], "WARN");
1504
1505        let log = events
1506            .iter()
1507            .find(|e| e.event_type == "event")
1508            .expect("log event");
1509        assert_eq!(log.event_data["level"], "ERROR");
1510    }
1511
1512    #[tokio::test]
1513    async fn test_span_close_reports_busy_ms_and_poll_count() {
1514        let (layer, mut receiver) = test_layer(false);
1515        let events =
1516            poll_span_n_times(layer, &mut receiver, 3, std::time::Duration::from_millis(5));
1517
1518        let close = span_close_for(&events, "polled_span");
1519        let fields = &close.event_data["fields"];
1520        assert_eq!(
1521            fields["poll_count"].as_u64(),
1522            Some(3),
1523            "poll_count should match the number of enter/exit cycles"
1524        );
1525        // 3 polls x >=5ms each. Lower bound only: CI timing is unreliable.
1526        let busy_ms = fields["busy_ms"].as_u64().expect("busy_ms should be a u64");
1527        assert!(
1528            busy_ms >= 10,
1529            "busy_ms should reflect time in span, got {busy_ms}"
1530        );
1531    }
1532
1533    #[tokio::test]
1534    async fn test_span_close_busy_fields_present_when_never_entered() {
1535        let (layer, mut receiver) = test_layer(false);
1536        let subscriber = tracing_subscriber::registry().with(layer);
1537
1538        tracing::subscriber::with_default(subscriber, || {
1539            // Created and dropped without ever being entered.
1540            let _span = span!(Level::INFO, "never_entered_span");
1541        });
1542
1543        let events = drain_events(&mut receiver);
1544        let close = span_close_for(&events, "never_entered_span");
1545        let fields = &close.event_data["fields"];
1546        assert_eq!(fields["busy_ms"].as_u64(), Some(0));
1547        assert_eq!(fields["poll_count"].as_u64(), Some(0));
1548    }
1549
1550    #[tokio::test]
1551    async fn test_busy_aggregation_works_with_emit_enter_exit_enabled() {
1552        let (layer, mut receiver) = test_layer(true);
1553        let events =
1554            poll_span_n_times(layer, &mut receiver, 2, std::time::Duration::from_millis(5));
1555
1556        // Enter/exit events are still emitted when opted in...
1557        assert_eq!(
1558            events
1559                .iter()
1560                .filter(|e| e.event_type == "span_enter")
1561                .count(),
1562            2
1563        );
1564        assert_eq!(
1565            events
1566                .iter()
1567                .filter(|e| e.event_type == "span_exit")
1568                .count(),
1569            2
1570        );
1571
1572        // ...and the aggregation still lands on span_close.
1573        let close = span_close_for(&events, "polled_span");
1574        let fields = &close.event_data["fields"];
1575        assert_eq!(fields["poll_count"].as_u64(), Some(2));
1576        let busy_ms = fields["busy_ms"].as_u64().expect("busy_ms should be a u64");
1577        assert!(
1578            busy_ms >= 5,
1579            "busy_ms should reflect time in span, got {busy_ms}"
1580        );
1581    }
1582
1583    #[tokio::test]
1584    async fn test_span_close_includes_recorded_fields() {
1585        let (layer, mut receiver) = test_layer(false);
1586        let subscriber = tracing_subscriber::registry().with(layer);
1587
1588        tracing::subscriber::with_default(subscriber, || {
1589            // Fields must be declared at creation (as Empty) to be recordable
1590            // later — exactly how tower-http's MakeSpan/OnResponse work.
1591            let span = span!(
1592                Level::INFO,
1593                "recording_span",
1594                status_code = tracing::field::Empty,
1595                content_type = tracing::field::Empty
1596            );
1597            let _guard = span.enter();
1598            span.record("status_code", 200_u64);
1599            span.record("content_type", "text/html");
1600        });
1601
1602        let events = drain_events(&mut receiver);
1603
1604        // span_new carries no values for the Empty fields...
1605        let span_new = events
1606            .iter()
1607            .find(|e| e.event_type == "span_new" && e.event_data["name"] == "recording_span")
1608            .expect("no span_new for recording_span");
1609        assert!(span_new.event_data["fields"]
1610            .get("status_code")
1611            .is_none_or(|v| v.is_null()));
1612
1613        // ...and span_close carries the final recorded state.
1614        let close = span_close_for(&events, "recording_span");
1615        let fields = &close.event_data["fields"];
1616        assert_eq!(fields["status_code"].as_u64(), Some(200));
1617        assert_eq!(fields["content_type"].as_str(), Some("text/html"));
1618        // The aggregation fields are still present alongside them.
1619        assert!(fields["busy_ms"].is_u64());
1620        assert!(fields["poll_count"].is_u64());
1621    }
1622
1623    #[tokio::test]
1624    async fn test_recorded_field_last_write_wins() {
1625        let (layer, mut receiver) = test_layer(false);
1626        let subscriber = tracing_subscriber::registry().with(layer);
1627
1628        tracing::subscriber::with_default(subscriber, || {
1629            let span = span!(
1630                Level::INFO,
1631                "rerecord_span",
1632                attempt = tracing::field::Empty
1633            );
1634            let _guard = span.enter();
1635            span.record("attempt", 1_u64);
1636            span.record("attempt", 2_u64);
1637            span.record("attempt", 3_u64);
1638        });
1639
1640        let events = drain_events(&mut receiver);
1641        let close = span_close_for(&events, "rerecord_span");
1642        assert_eq!(
1643            close.event_data["fields"]["attempt"].as_u64(),
1644            Some(3),
1645            "the last recorded value for a field should win"
1646        );
1647    }
1648
1649    #[tokio::test]
1650    async fn test_recorded_fields_cannot_clobber_busy_aggregation() {
1651        let (layer, mut receiver) = test_layer(false);
1652        let subscriber = tracing_subscriber::registry().with(layer);
1653
1654        tracing::subscriber::with_default(subscriber, || {
1655            // Hostile field names colliding with the eyes-owned aggregation
1656            // keys on span_close.
1657            let span = span!(
1658                Level::INFO,
1659                "hostile_span",
1660                busy_ms = tracing::field::Empty,
1661                poll_count = tracing::field::Empty
1662            );
1663            let guard = span.enter();
1664            span.record("busy_ms", "not-a-duration");
1665            span.record("poll_count", "lots");
1666            drop(guard);
1667        });
1668
1669        let events = drain_events(&mut receiver);
1670        let close = span_close_for(&events, "hostile_span");
1671        let fields = &close.event_data["fields"];
1672        // The aggregation values win over the recorded strings.
1673        assert!(
1674            fields["busy_ms"].is_u64(),
1675            "busy_ms must remain the aggregated u64, got {:?}",
1676            fields["busy_ms"]
1677        );
1678        assert_eq!(
1679            fields["poll_count"].as_u64(),
1680            Some(1),
1681            "poll_count must remain the aggregated value, got {:?}",
1682            fields["poll_count"]
1683        );
1684    }
1685
1686    #[tokio::test]
1687    async fn test_span_close_shape_unchanged_without_records() {
1688        let (layer, mut receiver) = test_layer(false);
1689        let subscriber = tracing_subscriber::registry().with(layer);
1690
1691        tracing::subscriber::with_default(subscriber, || {
1692            let span = span!(Level::INFO, "no_record_span", user_id = 7);
1693            let _guard = span.enter();
1694        });
1695
1696        let events = drain_events(&mut receiver);
1697        let close = span_close_for(&events, "no_record_span");
1698        let fields = close.event_data["fields"]
1699            .as_object()
1700            .expect("span_close fields should be an object");
1701        // Exactly the aggregation fields, nothing else: creation-time
1702        // attributes stay on span_new only.
1703        assert_eq!(fields.len(), 2, "unexpected span_close fields: {fields:?}");
1704        assert!(fields["busy_ms"].is_u64());
1705        assert_eq!(fields["poll_count"].as_u64(), Some(1));
1706    }
1707
1708    #[tokio::test]
1709    async fn test_enter_exit_not_emitted_by_default_layer() {
1710        let (sender, mut receiver) = mpsc::channel::<EventData>(64);
1711        let layer = EyesLayer {
1712            sender,
1713            dropped: Arc::new(AtomicU64::new(0)),
1714            emit_enter_exit: false,
1715            process_instance_id: None,
1716        };
1717        let subscriber = tracing_subscriber::registry().with(layer);
1718
1719        tracing::subscriber::with_default(subscriber, || {
1720            let span = span!(Level::INFO, "quiet_span");
1721            let _guard = span.enter();
1722            info!("inside quiet span");
1723        });
1724
1725        let mut events = Vec::new();
1726        while let Ok(event) = receiver.try_recv() {
1727            events.push(event);
1728        }
1729
1730        assert!(
1731            events
1732                .iter()
1733                .all(|e| e.event_type != "span_enter" && e.event_type != "span_exit"),
1734            "span_enter/span_exit must not be emitted by default"
1735        );
1736        // The lifecycle events the server derives durations from still flow.
1737        for expected in ["span_new", "event", "span_close"] {
1738            assert!(
1739                events.iter().any(|e| e.event_type == expected),
1740                "missing {expected} event"
1741            );
1742        }
1743    }
1744
1745    // ---------------------------------------------------------- measurements
1746
1747    /// Drains every event a closure emits through a fresh layer.
1748    fn drain_measurements(build: impl FnOnce(EyesLayer)) -> Vec<EventData> {
1749        let (sender, mut receiver) = mpsc::channel::<EventData>(64);
1750        let layer = EyesLayer {
1751            sender,
1752            dropped: std::sync::Arc::new(std::sync::atomic::AtomicU64::new(0)),
1753            emit_enter_exit: false,
1754            process_instance_id: None,
1755        };
1756        build(layer);
1757        let mut events = Vec::new();
1758        while let Ok(event) = receiver.try_recv() {
1759            events.push(event);
1760        }
1761        events
1762    }
1763
1764    fn emit_through_layer(body: impl FnOnce()) -> Vec<EventData> {
1765        drain_measurements(|layer| {
1766            let subscriber = tracing_subscriber::registry().with(layer);
1767            tracing::subscriber::with_default(subscriber, body);
1768        })
1769    }
1770
1771    /// `record_f64` maps NaN/±infinity to JSON null, which the server would
1772    /// reject with a 400 the transports cannot surface. The layer must drop the
1773    /// measurement itself and count it, not ship a doomed payload.
1774    #[test]
1775    fn a_non_finite_measurement_value_is_dropped_and_counted() {
1776        let (sender, mut receiver) = mpsc::channel::<EventData>(64);
1777        let dropped = std::sync::Arc::new(std::sync::atomic::AtomicU64::new(0));
1778        let layer = EyesLayer {
1779            sender,
1780            dropped: dropped.clone(),
1781            emit_enter_exit: false,
1782            process_instance_id: None,
1783        };
1784        let subscriber = tracing_subscriber::registry().with(layer);
1785        tracing::subscriber::with_default(subscriber, || {
1786            emit_gauge("rate", f64::NAN);
1787            emit_sample("latency", f64::INFINITY);
1788            emit_gauge("ok", 1.5);
1789        });
1790        let mut events = Vec::new();
1791        while let Ok(event) = receiver.try_recv() {
1792            events.push(event);
1793        }
1794        assert_eq!(events.len(), 1, "only the finite gauge survives");
1795        assert_eq!(events[0].event_data["metric_name"], "ok");
1796        assert_eq!(dropped.load(Ordering::Relaxed), 2);
1797    }
1798
1799    #[test]
1800    fn emit_gauge_dispatches_the_versioned_measurement_contract() {
1801        let events = emit_through_layer(|| emit_gauge("cpu", 42.5));
1802        assert_eq!(events.len(), 1);
1803        let data = &events[0].event_data;
1804        assert_eq!(events[0].event_type, "measurement");
1805        assert_eq!(data["version"], serde_json::json!(MEASUREMENT_VERSION));
1806        assert_eq!(data["metric_name"], "cpu");
1807        assert_eq!(data["metric_kind"], "gauge");
1808        assert_eq!(data["value"], serde_json::json!(42.5));
1809        // Display form, never Debug: `Level(Info)` silently broke every
1810        // server-side level predicate once.
1811        assert_eq!(data["level"], "INFO");
1812        assert_eq!(data["target"], MEASUREMENT_TARGET);
1813        assert_eq!(data["fields"], serde_json::json!({}));
1814    }
1815
1816    #[test]
1817    fn emit_counter_keeps_its_integer_identity_on_the_wire() {
1818        let events = emit_through_layer(|| emit_counter("requests", 5));
1819        assert_eq!(events.len(), 1);
1820        let value = &events[0].event_data["value"];
1821        assert!(value.is_i64() || value.is_u64(), "not an integer: {value}");
1822        assert_eq!(value, &serde_json::json!(5));
1823    }
1824
1825    #[test]
1826    fn emit_sample_names_its_kind() {
1827        let events = emit_through_layer(|| emit_sample("latency", 1.5));
1828        assert_eq!(events[0].event_data["metric_kind"], "sample");
1829    }
1830
1831    #[test]
1832    fn the_macro_lifts_reserved_names_and_leaves_version_a_dimension() {
1833        let events = emit_through_layer(|| {
1834            measurement!(
1835                "gauge",
1836                "cpu",
1837                42.5_f64,
1838                unit = "percent",
1839                description = "d",
1840                host = "web-1",
1841                version = "app-2.1"
1842            );
1843        });
1844        let data = &events[0].event_data;
1845        assert_eq!(data["unit"], "percent");
1846        assert_eq!(data["description"], "d");
1847        // `version` is NOT reserved: the layer writes the discriminator itself,
1848        // so an app keeps `version` as an ordinary dimension name.
1849        assert_eq!(data["version"], serde_json::json!(MEASUREMENT_VERSION));
1850        assert_eq!(data["fields"]["host"], "web-1");
1851        assert_eq!(data["fields"]["version"], "app-2.1");
1852        for reserved in MEASUREMENT_RESERVED_FIELDS {
1853            assert!(
1854                data["fields"].get(reserved).is_none(),
1855                "{reserved} left in the dimension bag"
1856            );
1857        }
1858    }
1859
1860    #[test]
1861    fn a_measurement_inside_a_span_carries_that_span_s_eyes_id() {
1862        let events = emit_through_layer(|| {
1863            let span = span!(Level::INFO, "outer");
1864            let _guard = span.enter();
1865            emit_gauge("cpu", 1.0);
1866        });
1867        let opened = events
1868            .iter()
1869            .find(|e| e.event_type == "span_new")
1870            .expect("span_new");
1871        let measured = events
1872            .iter()
1873            .find(|e| e.event_type == "measurement")
1874            .expect("measurement");
1875        let span_id = opened.event_data["span_id"].as_str().unwrap();
1876        assert_eq!(measured.event_data["span_id"], span_id);
1877        // Never the registry debug form, which resets to Id(1) on restart.
1878        assert_eq!(span_id.len(), 32);
1879        assert!(!span_id.starts_with("Id("));
1880    }
1881
1882    #[test]
1883    fn a_filter_that_does_not_enable_the_measurement_target_drops_measurements() {
1884        use tracing_subscriber::EnvFilter;
1885
1886        let dropped = drain_measurements(|layer| {
1887            let subscriber = tracing_subscriber::registry()
1888                .with(EnvFilter::new("warn"))
1889                .with(layer);
1890            tracing::subscriber::with_default(subscriber, || emit_gauge("cpu", 1.0));
1891        });
1892        assert!(dropped.is_empty(), "{dropped:?}");
1893
1894        let kept = drain_measurements(|layer| {
1895            let subscriber = tracing_subscriber::registry()
1896                .with(EnvFilter::new("warn,eyes::measurement=info"))
1897                .with(layer);
1898            tracing::subscriber::with_default(subscriber, || emit_gauge("cpu", 1.0));
1899        });
1900        assert_eq!(kept.len(), 1);
1901    }
1902
1903    #[test]
1904    fn every_emitter_writes_the_same_literal_target() {
1905        // `tracing::event!` needs a literal `target:` at the callsite, so the
1906        // constant and the literals can only be kept in step by asserting it.
1907        for events in [
1908            emit_through_layer(|| emit_gauge("g", 1.0)),
1909            emit_through_layer(|| emit_counter("c", 1)),
1910            emit_through_layer(|| emit_sample("s", 1.0)),
1911            emit_through_layer(|| measurement!("gauge", "m", 1.0_f64)),
1912            emit_through_layer(|| measurement!("gauge", "m", 1.0_f64, host = "h")),
1913        ] {
1914            assert_eq!(events.len(), 1);
1915            assert_eq!(events[0].event_data["target"], MEASUREMENT_TARGET);
1916            assert_eq!(events[0].event_type, "measurement");
1917        }
1918    }
1919}