Skip to main content

hyperactor_telemetry/
unix_sink.rs

1/*
2 * Copyright (c) Meta Platforms, Inc. and affiliates.
3 * All rights reserved.
4 *
5 * This source code is licensed under the BSD-style license found in the
6 * LICENSE file in the root directory of this source tree.
7 */
8
9//! Producer-side Unix-socket sink for telemetry sidecars.
10//!
11//! This module is the producer-side transport for telemetry. It installs
12//! an inactive process-global `UnixSocketSink`, buffers trace and entity events into
13//! schema-specific `RecordBatchBuffer`s, and, once activated with a socket
14//! path, forwards flushed batches to the sidecar over a Unix socket.
15//!
16//! The dispatcher thread only does cheap row buffering and bounded channel
17//! sends. Arrow IPC serialization and socket I/O run on a dedicated writer
18//! thread so slow or missing sidecars do not block application tracing.
19//!
20//! Frame layout matches `socket_ingest.rs`: table-name length, UTF-8 table
21//! name, Arrow IPC payload length, and one non-empty `RecordBatch` payload.
22//! Frames are lossy by design: if the queue is full, serialization fails, or a
23//! socket write fails, we increment `dropped` and keep the tracing path moving.
24
25use std::io::Write;
26use std::os::unix::net::UnixStream;
27use std::path::PathBuf;
28use std::sync::Arc;
29use std::sync::Mutex;
30use std::sync::OnceLock;
31use std::sync::atomic::AtomicU64;
32use std::sync::atomic::Ordering;
33use std::sync::mpsc;
34use std::time::SystemTime;
35
36use monarch_record_batch::RecordBatchBuffer;
37use monarch_telemetry_schema::MAX_FRAME_LEN;
38use monarch_telemetry_schema::entity_tables::ACTOR_STATUS_EVENTS;
39use monarch_telemetry_schema::entity_tables::ACTORS;
40use monarch_telemetry_schema::entity_tables::Actor;
41use monarch_telemetry_schema::entity_tables::ActorBuffer;
42use monarch_telemetry_schema::entity_tables::ActorStatusEvent as ActorStatusEventRow;
43use monarch_telemetry_schema::entity_tables::ActorStatusEventBuffer;
44use monarch_telemetry_schema::entity_tables::MESHES;
45use monarch_telemetry_schema::entity_tables::MESSAGE_STATUS_EVENTS;
46use monarch_telemetry_schema::entity_tables::MESSAGES;
47use monarch_telemetry_schema::entity_tables::Mesh;
48use monarch_telemetry_schema::entity_tables::MeshBuffer;
49use monarch_telemetry_schema::entity_tables::Message;
50use monarch_telemetry_schema::entity_tables::MessageBuffer;
51use monarch_telemetry_schema::entity_tables::MessageStatusEvent as MessageStatusEventRow;
52use monarch_telemetry_schema::entity_tables::MessageStatusEventBuffer;
53use monarch_telemetry_schema::entity_tables::SENT_MESSAGES;
54use monarch_telemetry_schema::entity_tables::SentMessage;
55use monarch_telemetry_schema::entity_tables::SentMessageBuffer;
56use monarch_telemetry_schema::trace_tables::EVENTS;
57use monarch_telemetry_schema::trace_tables::Event;
58use monarch_telemetry_schema::trace_tables::EventBuffer;
59use monarch_telemetry_schema::trace_tables::SPAN_EVENTS;
60use monarch_telemetry_schema::trace_tables::SPANS;
61use monarch_telemetry_schema::trace_tables::Span;
62use monarch_telemetry_schema::trace_tables::SpanBuffer;
63use monarch_telemetry_schema::trace_tables::SpanEvent;
64use monarch_telemetry_schema::trace_tables::SpanEventBuffer;
65use monarch_telemetry_schema::write_frame_header;
66use tracing_subscriber::filter::Targets;
67
68use crate::EntityEvent;
69use crate::FieldValue;
70use crate::TraceEvent;
71use crate::TraceEventSink;
72use crate::config::get_tracing_targets;
73use crate::generate_sent_message_id;
74
75/// Maximum queued table batches before the tracing path starts dropping frames.
76const WORKER_QUEUE_CAPACITY: usize = 10_000;
77
78/// Process-global sink installed with the tracing subscriber and activated later.
79static UNIX_SOCKET_SINK: OnceLock<Arc<UnixSocketSink>> = OnceLock::new();
80
81/// Producer-side sink that forwards telemetry batches to a sidecar socket.
82pub struct UnixSocketSink {
83    // Shared between the dispatcher adapter and the global activation handle.
84    // The mutex covers buffered rows plus the path/worker activation state.
85    inner: Mutex<UnixSocketSinkInner>,
86    // The writer thread increments this for asynchronous serialization and
87    // socket failures, so callers can observe producer-side frame loss.
88    dropped: Arc<AtomicU64>,
89}
90
91struct UnixSocketSinkInner {
92    spans_buffer: SpanBuffer,
93    span_events_buffer: SpanEventBuffer,
94    events_buffer: EventBuffer,
95    actors_buffer: ActorBuffer,
96    meshes_buffer: MeshBuffer,
97    actor_status_events_buffer: ActorStatusEventBuffer,
98    sent_messages_buffer: SentMessageBuffer,
99    messages_buffer: MessageBuffer,
100    message_status_events_buffer: MessageStatusEventBuffer,
101    path: Option<PathBuf>,
102    worker: Option<WorkerHandle>,
103}
104
105struct WorkerHandle {
106    sender: mpsc::SyncSender<TelemetryTableBuffer>,
107    _join_handle: std::thread::JoinHandle<()>,
108}
109
110enum TelemetryTableBuffer {
111    Spans(SpanBuffer),
112    SpanEvents(SpanEventBuffer),
113    Events(EventBuffer),
114    Actors(ActorBuffer),
115    Meshes(MeshBuffer),
116    ActorStatusEvents(ActorStatusEventBuffer),
117    SentMessages(SentMessageBuffer),
118    Messages(MessageBuffer),
119    MessageStatusEvents(MessageStatusEventBuffer),
120}
121
122struct UnixSocketSinkAdapter {
123    // The dispatcher owns this adapter, while `UNIX_SOCKET_SINK` keeps another
124    // handle so callers can supply the socket path after logging initialization.
125    sink: Arc<UnixSocketSink>,
126    target_filter: Targets,
127}
128
129impl UnixSocketSink {
130    /// Create a sink with no socket path and no writer thread.
131    pub fn new() -> Self {
132        Self {
133            inner: Mutex::new(UnixSocketSinkInner {
134                spans_buffer: SpanBuffer::default(),
135                span_events_buffer: SpanEventBuffer::default(),
136                events_buffer: EventBuffer::default(),
137                actors_buffer: ActorBuffer::default(),
138                meshes_buffer: MeshBuffer::default(),
139                actor_status_events_buffer: ActorStatusEventBuffer::default(),
140                sent_messages_buffer: SentMessageBuffer::default(),
141                messages_buffer: MessageBuffer::default(),
142                message_status_events_buffer: MessageStatusEventBuffer::default(),
143                path: None,
144                worker: None,
145            }),
146            dropped: Arc::new(AtomicU64::new(0)),
147        }
148    }
149
150    /// Activate the sink and lazily spawn a writer thread.
151    ///
152    /// Reapplying the same path is a no-op; changing paths is rejected because
153    /// the process-global sink should not be retargeted after activation.
154    pub fn set_path(&self, path: PathBuf) -> anyhow::Result<()> {
155        let mut inner = self
156            .inner
157            .lock()
158            .map_err(|_| anyhow::anyhow!("lock poisoned"))?;
159
160        if inner.path.as_ref() == Some(&path) {
161            return Ok(());
162        }
163        if let Some(existing) = &inner.path {
164            anyhow::bail!(
165                "unix socket sink already activated with {}",
166                existing.display()
167            );
168        }
169
170        // The tracing subscriber is installed before the sidecar path is
171        // available. Start the writer only after activation so the inactive
172        // sink is cheap and does not repeatedly try to connect to nowhere.
173        let (sender, receiver) = mpsc::sync_channel(WORKER_QUEUE_CAPACITY);
174        let dropped = Arc::clone(&self.dropped);
175        let writer_path = path.clone();
176        let join_handle = std::thread::Builder::new()
177            .name("monarch-telemetry-unix-sink".into())
178            .spawn(move || writer_loop(writer_path, receiver, dropped))
179            .map_err(|error| anyhow::anyhow!("failed to spawn unix socket sink: {error}"))?;
180
181        inner.path = Some(path);
182        inner.worker = Some(WorkerHandle {
183            sender,
184            _join_handle: join_handle,
185        });
186        Ok(())
187    }
188
189    /// Return whether this sink has an active socket path.
190    pub fn is_active(&self) -> bool {
191        self.inner
192            .lock()
193            .map(|inner| inner.worker.is_some())
194            .unwrap_or(false)
195    }
196
197    /// Return the cumulative number of dropped socket frames.
198    pub fn dropped_frames(&self) -> u64 {
199        self.dropped.load(Ordering::Relaxed)
200    }
201
202    /// Convert one dispatcher event into its table row and buffer it locally.
203    fn consume_shared(&self, event: &TraceEvent) -> anyhow::Result<()> {
204        let mut inner = self
205            .inner
206            .lock()
207            .map_err(|_| anyhow::anyhow!("lock poisoned"))?;
208
209        match event {
210            TraceEvent::NewSpan {
211                id,
212                name,
213                target,
214                level,
215                fields,
216                timestamp,
217                parent_id,
218                thread_name,
219                file,
220                line,
221            } => {
222                inner.spans_buffer.insert(Span {
223                    id: *id,
224                    name: name.to_string(),
225                    target: target.to_string(),
226                    level: level.to_string(),
227                    fields_json: fields_to_json(fields),
228                    timestamp_us: timestamp_to_micros(timestamp),
229                    parent_id: *parent_id,
230                    thread_name: thread_name.to_string(),
231                    file: file.map(|s| s.to_string()),
232                    line: *line,
233                });
234            }
235            TraceEvent::SpanEnter { id, timestamp, .. } => {
236                inner.span_events_buffer.insert(SpanEvent {
237                    id: *id,
238                    timestamp_us: timestamp_to_micros(timestamp),
239                    event_type: "enter".to_string(),
240                });
241            }
242            TraceEvent::SpanExit { id, timestamp, .. } => {
243                inner.span_events_buffer.insert(SpanEvent {
244                    id: *id,
245                    timestamp_us: timestamp_to_micros(timestamp),
246                    event_type: "exit".to_string(),
247                });
248            }
249            TraceEvent::SpanClose { id, timestamp } => {
250                inner.span_events_buffer.insert(SpanEvent {
251                    id: *id,
252                    timestamp_us: timestamp_to_micros(timestamp),
253                    event_type: "close".to_string(),
254                });
255            }
256            TraceEvent::Event {
257                name,
258                target,
259                level,
260                fields,
261                timestamp,
262                parent_span,
263                thread_id,
264                thread_name,
265                module_path,
266                file,
267                line,
268            } => {
269                inner.events_buffer.insert(Event {
270                    name: name.to_string(),
271                    target: target.to_string(),
272                    level: level.to_string(),
273                    fields_json: fields_to_json(fields),
274                    timestamp_us: timestamp_to_micros(timestamp),
275                    parent_span: *parent_span,
276                    thread_id: thread_id.to_string(),
277                    thread_name: thread_name.to_string(),
278                    module_path: module_path.map(|s| s.to_string()),
279                    file: file.map(|s| s.to_string()),
280                    line: *line,
281                });
282            }
283            TraceEvent::Entity(event) => buffer_entity_event(&mut inner, event),
284        }
285        Ok(())
286    }
287
288    /// Move non-empty table buffers to the writer queue, or drop them while inactive.
289    fn flush_shared(&self) -> anyhow::Result<()> {
290        let mut inner = self
291            .inner
292            .lock()
293            .map_err(|_| anyhow::anyhow!("lock poisoned"))?;
294
295        if inner.worker.is_none() {
296            // The buffers provide only a small pre-activation window: rows
297            // consumed before `set_path` can be sent if activation wins the
298            // race with the next dispatcher flush. Otherwise, drop them here
299            // so startup telemetry does not accumulate without a socket path.
300            inner.spans_buffer = SpanBuffer::default();
301            inner.span_events_buffer = SpanEventBuffer::default();
302            inner.events_buffer = EventBuffer::default();
303            inner.actors_buffer = ActorBuffer::default();
304            inner.meshes_buffer = MeshBuffer::default();
305            inner.actor_status_events_buffer = ActorStatusEventBuffer::default();
306            inner.sent_messages_buffer = SentMessageBuffer::default();
307            inner.messages_buffer = MessageBuffer::default();
308            inner.message_status_events_buffer = MessageStatusEventBuffer::default();
309            return Ok(());
310        }
311
312        let sender = inner
313            .worker
314            .as_ref()
315            .expect("worker should exist")
316            .sender
317            .clone();
318
319        flush_buffer(
320            &mut inner.spans_buffer,
321            TelemetryTableBuffer::Spans,
322            &sender,
323            &self.dropped,
324        );
325        flush_buffer(
326            &mut inner.span_events_buffer,
327            TelemetryTableBuffer::SpanEvents,
328            &sender,
329            &self.dropped,
330        );
331        flush_buffer(
332            &mut inner.events_buffer,
333            TelemetryTableBuffer::Events,
334            &sender,
335            &self.dropped,
336        );
337        flush_buffer(
338            &mut inner.actors_buffer,
339            TelemetryTableBuffer::Actors,
340            &sender,
341            &self.dropped,
342        );
343        flush_buffer(
344            &mut inner.meshes_buffer,
345            TelemetryTableBuffer::Meshes,
346            &sender,
347            &self.dropped,
348        );
349        flush_buffer(
350            &mut inner.actor_status_events_buffer,
351            TelemetryTableBuffer::ActorStatusEvents,
352            &sender,
353            &self.dropped,
354        );
355        flush_buffer(
356            &mut inner.sent_messages_buffer,
357            TelemetryTableBuffer::SentMessages,
358            &sender,
359            &self.dropped,
360        );
361        flush_buffer(
362            &mut inner.messages_buffer,
363            TelemetryTableBuffer::Messages,
364            &sender,
365            &self.dropped,
366        );
367        flush_buffer(
368            &mut inner.message_status_events_buffer,
369            TelemetryTableBuffer::MessageStatusEvents,
370            &sender,
371            &self.dropped,
372        );
373        Ok(())
374    }
375}
376
377impl TraceEventSink for UnixSocketSinkAdapter {
378    fn consume(&mut self, event: &TraceEvent) -> Result<(), anyhow::Error> {
379        self.sink.consume_shared(event)
380    }
381
382    fn target_filter(&self) -> Option<&Targets> {
383        Some(&self.target_filter)
384    }
385
386    fn flush(&mut self) -> Result<(), anyhow::Error> {
387        self.sink.flush_shared()
388    }
389
390    fn name(&self) -> &str {
391        "UnixSocketSink"
392    }
393}
394
395/// Install the inactive process-global Unix socket sink.
396pub(crate) fn install_unix_socket_sink_inactive() -> Box<dyn TraceEventSink> {
397    let sink = Arc::new(UnixSocketSink::new());
398    let _ = UNIX_SOCKET_SINK.set(Arc::clone(&sink));
399    Box::new(UnixSocketSinkAdapter {
400        sink,
401        target_filter: get_tracing_targets(),
402    })
403}
404
405/// Activate the process-global Unix socket sink against a socket path.
406pub fn set_unix_socket_sink_path(path: impl Into<PathBuf>) -> anyhow::Result<()> {
407    let sink = UNIX_SOCKET_SINK
408        .get()
409        .ok_or_else(|| anyhow::anyhow!("unix socket sink is not installed"))?;
410    sink.set_path(path.into())
411}
412
413/// Return whether the process-global Unix socket sink is active.
414pub fn unix_socket_sink_is_active() -> bool {
415    UNIX_SOCKET_SINK
416        .get()
417        .map(|sink| sink.is_active())
418        .unwrap_or(false)
419}
420
421/// Return the process-global Unix socket sink's cumulative dropped frames.
422pub fn unix_socket_sink_dropped_frames() -> Option<u64> {
423    UNIX_SOCKET_SINK.get().map(|sink| sink.dropped_frames())
424}
425
426fn flush_buffer<B>(
427    buffer: &mut B,
428    table_buffer: impl FnOnce(B) -> TelemetryTableBuffer,
429    sender: &mpsc::SyncSender<TelemetryTableBuffer>,
430    dropped: &AtomicU64,
431) where
432    B: RecordBatchBuffer + Default,
433{
434    if buffer.is_empty() {
435        return;
436    }
437
438    let batch = table_buffer(std::mem::take(buffer));
439    match sender.try_send(batch) {
440        Ok(()) => {}
441        Err(mpsc::TrySendError::Full(_)) | Err(mpsc::TrySendError::Disconnected(_)) => {
442            // Socket delivery must never backpressure the tracing dispatcher.
443            // Drop the whole table batch when the worker cannot accept it.
444            dropped.fetch_add(1, Ordering::Relaxed);
445        }
446    }
447}
448
449fn buffer_entity_event(inner: &mut UnixSocketSinkInner, event: &EntityEvent) {
450    match event {
451        EntityEvent::Actor(event) => {
452            inner.actors_buffer.insert(Actor {
453                id: event.id,
454                timestamp_us: timestamp_to_micros(&event.timestamp),
455                mesh_id: event.mesh_id,
456                rank: event.rank,
457                full_name: event.full_name.clone(),
458                display_name: event.display_name.clone(),
459            });
460        }
461        EntityEvent::Mesh(event) => {
462            inner.meshes_buffer.insert(Mesh {
463                id: event.id,
464                timestamp_us: timestamp_to_micros(&event.timestamp),
465                class: event.class.clone(),
466                given_name: event.given_name.clone(),
467                full_name: event.full_name.clone(),
468                shape_json: event.shape_json.clone(),
469                parent_mesh_id: event.parent_mesh_id,
470                parent_view_json: event.parent_view_json.clone(),
471            });
472        }
473        EntityEvent::ActorStatus(event) => {
474            inner
475                .actor_status_events_buffer
476                .insert(ActorStatusEventRow {
477                    id: event.id,
478                    timestamp_us: timestamp_to_micros(&event.timestamp),
479                    actor_id: event.actor_id,
480                    new_status: event.new_status.clone(),
481                    reason: event.reason.clone(),
482                });
483        }
484        EntityEvent::SentMessage(event) => {
485            inner.sent_messages_buffer.insert(SentMessage {
486                id: generate_sent_message_id(event.sender_actor_id),
487                timestamp_us: timestamp_to_micros(&event.timestamp),
488                sender_actor_id: event.sender_actor_id,
489                actor_mesh_id: event.actor_mesh_id,
490                view_json: event.view_json.clone(),
491                shape_json: event.shape_json.clone(),
492            });
493        }
494        EntityEvent::Message(event) => {
495            inner.messages_buffer.insert(Message {
496                id: event.id,
497                timestamp_us: timestamp_to_micros(&event.timestamp),
498                from_actor_id: event.from_actor_id,
499                to_actor_id: event.to_actor_id,
500                endpoint: event.endpoint.clone(),
501                port_index: event.port_index,
502            });
503        }
504        EntityEvent::MessageStatus(event) => {
505            inner
506                .message_status_events_buffer
507                .insert(MessageStatusEventRow {
508                    id: event.id,
509                    timestamp_us: timestamp_to_micros(&event.timestamp),
510                    message_id: event.message_id,
511                    status: event.status.clone(),
512                });
513        }
514    }
515}
516
517/// Serialize drained table buffers and write framed Arrow IPC payloads to the sidecar.
518fn writer_loop(
519    path: PathBuf,
520    receiver: mpsc::Receiver<TelemetryTableBuffer>,
521    dropped: Arc<AtomicU64>,
522) {
523    let mut stream = None;
524
525    while let Ok(mut buffer) = receiver.recv() {
526        // Convert buffered rows into the same one-table, one-batch frame shape
527        // that the ingest server validates. Failures here are producer-side
528        // frame drops; the writer keeps processing later batches.
529        //
530        // Keep the table-name selection next to the concrete buffer variant so
531        // adding a socket table requires an explicit writer-loop update.
532        let (table_name, batch) = match &mut buffer {
533            TelemetryTableBuffer::Spans(buffer) => (SPANS, buffer.drain_to_record_batch()),
534            TelemetryTableBuffer::SpanEvents(buffer) => {
535                (SPAN_EVENTS, buffer.drain_to_record_batch())
536            }
537            TelemetryTableBuffer::Events(buffer) => (EVENTS, buffer.drain_to_record_batch()),
538            TelemetryTableBuffer::Actors(buffer) => (ACTORS, buffer.drain_to_record_batch()),
539            TelemetryTableBuffer::Meshes(buffer) => (MESHES, buffer.drain_to_record_batch()),
540            TelemetryTableBuffer::ActorStatusEvents(buffer) => {
541                (ACTOR_STATUS_EVENTS, buffer.drain_to_record_batch())
542            }
543            TelemetryTableBuffer::SentMessages(buffer) => {
544                (SENT_MESSAGES, buffer.drain_to_record_batch())
545            }
546            TelemetryTableBuffer::Messages(buffer) => (MESSAGES, buffer.drain_to_record_batch()),
547            TelemetryTableBuffer::MessageStatusEvents(buffer) => {
548                (MESSAGE_STATUS_EVENTS, buffer.drain_to_record_batch())
549            }
550        };
551        let Ok(batch) = batch else {
552            dropped.fetch_add(1, Ordering::Relaxed);
553            continue;
554        };
555        if batch.num_rows() == 0 {
556            continue;
557        }
558
559        let Ok(payload) = monarch_telemetry_schema::serialize_batch(&batch) else {
560            dropped.fetch_add(1, Ordering::Relaxed);
561            continue;
562        };
563        if payload.len() > MAX_FRAME_LEN {
564            dropped.fetch_add(1, Ordering::Relaxed);
565            continue;
566        }
567
568        // Exact frame size: u16 table-name length, table-name bytes, u32
569        // payload length, then the Arrow IPC payload bytes.
570        let mut frame = Vec::with_capacity(2 + table_name.len() + 4 + payload.len());
571        if write_frame_header(&mut frame, table_name, payload.len()).is_err() {
572            dropped.fetch_add(1, Ordering::Relaxed);
573            continue;
574        }
575        frame.extend_from_slice(&payload);
576
577        if write_frame(&path, &mut stream, &frame).is_err() {
578            dropped.fetch_add(1, Ordering::Relaxed);
579            stream = None;
580        }
581    }
582}
583
584fn write_frame(
585    path: &PathBuf,
586    stream: &mut Option<UnixStream>,
587    frame: &[u8],
588) -> std::io::Result<()> {
589    if stream.is_none() {
590        *stream = Some(UnixStream::connect(path)?);
591    }
592
593    // Reuse a connected stream across frames. On any write error, discard the
594    // stream so the next frame attempts a fresh connect to a restarted sidecar.
595    match stream
596        .as_mut()
597        .expect("stream should exist")
598        .write_all(frame)
599    {
600        Ok(()) => Ok(()),
601        Err(error) => {
602            *stream = None;
603            Err(error)
604        }
605    }
606}
607
608fn fields_to_json(fields: &[(&str, FieldValue)]) -> String {
609    monarch_telemetry_schema::fields_to_json(fields.iter().map(|(key, value)| {
610        let json_value = match value {
611            FieldValue::Bool(b) => serde_json::Value::Bool(*b),
612            FieldValue::I64(i) => serde_json::Value::Number((*i).into()),
613            FieldValue::U64(u) => serde_json::Value::Number((*u).into()),
614            FieldValue::F64(f) => serde_json::Number::from_f64(*f)
615                .map(serde_json::Value::Number)
616                .unwrap_or(serde_json::Value::Null),
617            FieldValue::Str(s) => serde_json::Value::String(s.clone()),
618            FieldValue::Debug(d) => serde_json::Value::String(d.clone()),
619        };
620        (*key, json_value)
621    }))
622}
623
624fn timestamp_to_micros(timestamp: &SystemTime) -> i64 {
625    timestamp
626        .duration_since(std::time::UNIX_EPOCH)
627        .unwrap_or_default()
628        .as_micros() as i64
629}
630
631#[cfg(test)]
632pub(crate) fn adapter_for_test(sink: Arc<UnixSocketSink>) -> Box<dyn TraceEventSink> {
633    Box::new(UnixSocketSinkAdapter {
634        sink,
635        target_filter: get_tracing_targets(),
636    })
637}
638
639#[cfg(test)]
640mod tests {
641    use std::io::ErrorKind;
642    use std::io::Read;
643    use std::os::unix::net::UnixListener;
644    use std::os::unix::net::UnixStream;
645    use std::sync::atomic::AtomicU64;
646    use std::sync::atomic::Ordering;
647    use std::sync::mpsc;
648    use std::time::Duration;
649    use std::time::Instant;
650    use std::time::SystemTime;
651
652    use monarch_record_batch::RecordBatchBuffer;
653
654    use super::*;
655    use crate::ActorEvent;
656    use crate::ActorStatusEvent;
657    use crate::MeshEvent;
658    use crate::MessageEvent;
659    use crate::MessageStatusEvent;
660    use crate::SentMessageEvent;
661
662    static TEST_SEQ: AtomicU64 = AtomicU64::new(0);
663
664    struct Frame {
665        table_name: String,
666        payload: Vec<u8>,
667    }
668
669    fn socket_path(name: &str) -> PathBuf {
670        let seq = TEST_SEQ.fetch_add(1, Ordering::Relaxed);
671        let dir = std::env::temp_dir().join(format!(
672            "monarch_unix_socket_sink_{}_{}",
673            std::process::id(),
674            seq
675        ));
676        std::fs::create_dir_all(&dir).unwrap();
677        dir.join(name)
678    }
679
680    fn timestamp() -> SystemTime {
681        SystemTime::UNIX_EPOCH + Duration::from_micros(123)
682    }
683
684    fn fields() -> crate::trace_dispatcher::TraceFields {
685        let mut fields = crate::trace_dispatcher::TraceFields::new();
686        fields.push(("count", FieldValue::U64(3)));
687        fields
688    }
689
690    fn span() -> TraceEvent {
691        TraceEvent::NewSpan {
692            id: 7,
693            name: "test_span",
694            target: "test_target",
695            level: tracing::Level::INFO,
696            fields: fields(),
697            timestamp: timestamp(),
698            parent_id: None,
699            thread_name: "test_thread",
700            file: Some("test.rs"),
701            line: Some(42),
702        }
703    }
704
705    fn span_enter() -> TraceEvent {
706        TraceEvent::SpanEnter {
707            id: 7,
708            timestamp: timestamp(),
709            thread_name: "test_thread",
710        }
711    }
712
713    fn event() -> TraceEvent {
714        TraceEvent::Event {
715            name: "test_event",
716            target: "test_target",
717            level: tracing::Level::INFO,
718            fields: fields(),
719            timestamp: timestamp(),
720            parent_span: Some(7),
721            thread_id: "99",
722            thread_name: "test_thread",
723            module_path: Some("test_module"),
724            file: Some("test.rs"),
725            line: Some(43),
726        }
727    }
728
729    fn entity_events() -> Vec<TraceEvent> {
730        vec![
731            TraceEvent::Entity(EntityEvent::Actor(ActorEvent {
732                id: 1,
733                timestamp: timestamp(),
734                mesh_id: 2,
735                rank: 3,
736                full_name: "actor/full".to_string(),
737                display_name: Some("actor".to_string()),
738            })),
739            TraceEvent::Entity(EntityEvent::Mesh(MeshEvent {
740                id: 2,
741                timestamp: timestamp(),
742                class: "Host".to_string(),
743                given_name: "hosts".to_string(),
744                full_name: "hosts/full".to_string(),
745                shape_json: r#"{"dims":[1]}"#.to_string(),
746                parent_mesh_id: None,
747                parent_view_json: None,
748            })),
749            TraceEvent::Entity(EntityEvent::ActorStatus(ActorStatusEvent {
750                id: 3,
751                timestamp: timestamp(),
752                actor_id: 1,
753                new_status: "Running".to_string(),
754                reason: Some("test".to_string()),
755            })),
756            TraceEvent::Entity(EntityEvent::SentMessage(SentMessageEvent {
757                timestamp: timestamp(),
758                sender_actor_id: 1,
759                actor_mesh_id: 2,
760                view_json: r#"{"rank":0}"#.to_string(),
761                shape_json: r#"{"dims":[1]}"#.to_string(),
762            })),
763            TraceEvent::Entity(EntityEvent::Message(MessageEvent {
764                timestamp: timestamp(),
765                id: 4,
766                from_actor_id: 1,
767                to_actor_id: 5,
768                endpoint: Some("endpoint".to_string()),
769                port_index: Some(6),
770            })),
771            TraceEvent::Entity(EntityEvent::MessageStatus(MessageStatusEvent {
772                timestamp: timestamp(),
773                id: 7,
774                message_id: 4,
775                status: "complete".to_string(),
776            })),
777        ]
778    }
779
780    fn read_frame(stream: &mut UnixStream) -> Frame {
781        let mut name_len_bytes = [0; 2];
782        stream.read_exact(&mut name_len_bytes).unwrap();
783        let name_len = u16::from_be_bytes(name_len_bytes) as usize;
784
785        let mut name_bytes = vec![0; name_len];
786        stream.read_exact(&mut name_bytes).unwrap();
787        let table_name = String::from_utf8(name_bytes).unwrap();
788
789        let mut payload_len_bytes = [0; 4];
790        stream.read_exact(&mut payload_len_bytes).unwrap();
791        let payload_len = u32::from_be_bytes(payload_len_bytes) as usize;
792        assert!((1..=MAX_FRAME_LEN).contains(&payload_len));
793
794        let mut payload = vec![0; payload_len];
795        stream.read_exact(&mut payload).unwrap();
796        Frame {
797            table_name,
798            payload,
799        }
800    }
801
802    fn read_frame_tables(stream: UnixStream, count: usize) -> Vec<String> {
803        read_frames(stream, count)
804            .into_iter()
805            .map(|frame| frame.table_name)
806            .collect()
807    }
808
809    fn read_frames(mut stream: UnixStream, count: usize) -> Vec<Frame> {
810        stream
811            .set_read_timeout(Some(Duration::from_secs(5)))
812            .unwrap();
813
814        (0..count).map(|_| read_frame(&mut stream)).collect()
815    }
816
817    fn assert_no_frame_available(stream: &mut UnixStream) {
818        stream
819            .set_read_timeout(Some(Duration::from_millis(500)))
820            .unwrap();
821        let mut name_len_bytes = [0; 2];
822        match stream.read_exact(&mut name_len_bytes) {
823            Ok(()) => panic!("unexpected extra telemetry frame"),
824            Err(error) if matches!(error.kind(), ErrorKind::WouldBlock | ErrorKind::TimedOut) => {}
825            Err(error) => panic!("unexpected read error: {error}"),
826        }
827    }
828
829    fn wait_for_dropped_frames(sink: &UnixSocketSink, expected: u64) {
830        let deadline = Instant::now() + Duration::from_secs(5);
831        while Instant::now() < deadline {
832            if sink.dropped_frames() == expected {
833                return;
834            }
835            std::thread::sleep(Duration::from_millis(10));
836        }
837        assert_eq!(sink.dropped_frames(), expected);
838    }
839
840    fn assert_frame_batch<B>(frame: &Frame)
841    where
842        B: RecordBatchBuffer + Default,
843    {
844        let batch = monarch_telemetry_schema::deserialize_one_batch(&frame.payload).unwrap();
845        let expected = B::default().drain_to_record_batch().unwrap();
846
847        assert_eq!(batch.num_rows(), 1);
848        assert_eq!(batch.schema(), expected.schema());
849    }
850
851    fn assert_entity_frame_schema(frame: &Frame) {
852        match frame.table_name.as_str() {
853            ACTORS => assert_frame_batch::<ActorBuffer>(frame),
854            MESHES => assert_frame_batch::<MeshBuffer>(frame),
855            ACTOR_STATUS_EVENTS => assert_frame_batch::<ActorStatusEventBuffer>(frame),
856            SENT_MESSAGES => assert_frame_batch::<SentMessageBuffer>(frame),
857            MESSAGES => assert_frame_batch::<MessageBuffer>(frame),
858            MESSAGE_STATUS_EVENTS => assert_frame_batch::<MessageStatusEventBuffer>(frame),
859            other => panic!("unexpected entity table frame: {other}"),
860        }
861    }
862
863    #[test]
864    fn inactive_flush_discards_buffered_events_without_counting_drops() {
865        let sink = UnixSocketSink::new();
866
867        sink.consume_shared(&span()).unwrap();
868        sink.consume_shared(&span_enter()).unwrap();
869        sink.consume_shared(&event()).unwrap();
870        for event in entity_events() {
871            sink.consume_shared(&event).unwrap();
872        }
873
874        {
875            let inner = sink.inner.lock().unwrap();
876            assert_eq!(inner.spans_buffer.len(), 1);
877            assert_eq!(inner.span_events_buffer.len(), 1);
878            assert_eq!(inner.events_buffer.len(), 1);
879            assert_eq!(inner.actors_buffer.len(), 1);
880            assert_eq!(inner.meshes_buffer.len(), 1);
881            assert_eq!(inner.actor_status_events_buffer.len(), 1);
882            assert_eq!(inner.sent_messages_buffer.len(), 1);
883            assert_eq!(inner.messages_buffer.len(), 1);
884            assert_eq!(inner.message_status_events_buffer.len(), 1);
885        }
886
887        sink.flush_shared().unwrap();
888
889        {
890            let inner = sink.inner.lock().unwrap();
891            assert_eq!(inner.spans_buffer.len(), 0);
892            assert_eq!(inner.span_events_buffer.len(), 0);
893            assert_eq!(inner.events_buffer.len(), 0);
894            assert_eq!(inner.actors_buffer.len(), 0);
895            assert_eq!(inner.meshes_buffer.len(), 0);
896            assert_eq!(inner.actor_status_events_buffer.len(), 0);
897            assert_eq!(inner.sent_messages_buffer.len(), 0);
898            assert_eq!(inner.messages_buffer.len(), 0);
899            assert_eq!(inner.message_status_events_buffer.len(), 0);
900        }
901        assert!(!sink.is_active());
902        assert_eq!(sink.dropped_frames(), 0);
903    }
904
905    #[test]
906    fn set_path_is_idempotent_for_same_path() {
907        let path = socket_path("same.sock");
908        let sink = UnixSocketSink::new();
909
910        sink.set_path(path.clone()).unwrap();
911        sink.set_path(path).unwrap();
912    }
913
914    #[test]
915    fn set_path_rejects_different_path_after_activation() {
916        let sink = UnixSocketSink::new();
917
918        sink.set_path(socket_path("first.sock")).unwrap();
919
920        assert!(sink.set_path(socket_path("second.sock")).is_err());
921    }
922
923    #[test]
924    fn active_flush_sends_one_frame_per_non_empty_trace_table() {
925        let path = socket_path("telemetry.sock");
926        let listener = UnixListener::bind(&path).unwrap();
927        let (sender, receiver) = mpsc::channel();
928        let read_handle = std::thread::spawn(move || {
929            let (stream, _addr) = listener.accept().unwrap();
930            sender.send(read_frame_tables(stream, 3)).unwrap();
931        });
932
933        let sink = UnixSocketSink::new();
934        sink.set_path(path).unwrap();
935        sink.consume_shared(&span()).unwrap();
936        sink.consume_shared(&span_enter()).unwrap();
937        sink.consume_shared(&event()).unwrap();
938
939        sink.flush_shared().unwrap();
940
941        let tables = receiver.recv_timeout(Duration::from_secs(5)).unwrap();
942        read_handle.join().unwrap();
943        assert_eq!(
944            tables,
945            vec![
946                SPANS.to_string(),
947                SPAN_EVENTS.to_string(),
948                EVENTS.to_string()
949            ]
950        );
951        assert_eq!(sink.dropped_frames(), 0);
952    }
953
954    #[test]
955    fn active_flush_sends_one_frame_per_non_empty_entity_table() {
956        let path = socket_path("entities.sock");
957        let listener = UnixListener::bind(&path).unwrap();
958        let (sender, receiver) = mpsc::channel();
959        let read_handle = std::thread::spawn(move || {
960            let (stream, _addr) = listener.accept().unwrap();
961            sender.send(read_frames(stream, 6)).unwrap();
962        });
963
964        let sink = UnixSocketSink::new();
965        sink.set_path(path).unwrap();
966        for event in entity_events() {
967            sink.consume_shared(&event).unwrap();
968        }
969
970        sink.flush_shared().unwrap();
971
972        let frames = receiver.recv_timeout(Duration::from_secs(5)).unwrap();
973        read_handle.join().unwrap();
974        assert_eq!(
975            frames
976                .iter()
977                .map(|frame| frame.table_name.as_str())
978                .collect::<Vec<_>>(),
979            vec![
980                ACTORS,
981                MESHES,
982                ACTOR_STATUS_EVENTS,
983                SENT_MESSAGES,
984                MESSAGES,
985                MESSAGE_STATUS_EVENTS,
986            ]
987        );
988        frames.iter().for_each(assert_entity_frame_schema);
989        assert_eq!(sink.dropped_frames(), 0);
990    }
991
992    #[test]
993    fn active_flush_omits_empty_trace_tables() {
994        let path = socket_path("events_only.sock");
995        let listener = UnixListener::bind(&path).unwrap();
996        let (sender, receiver) = mpsc::channel();
997        let read_handle = std::thread::spawn(move || {
998            let (mut stream, _addr) = listener.accept().unwrap();
999            stream
1000                .set_read_timeout(Some(Duration::from_secs(5)))
1001                .unwrap();
1002            let table_name = read_frame(&mut stream).table_name;
1003            assert_no_frame_available(&mut stream);
1004            sender.send(table_name).unwrap();
1005        });
1006
1007        let sink = UnixSocketSink::new();
1008        sink.set_path(path).unwrap();
1009        sink.consume_shared(&event()).unwrap();
1010
1011        sink.flush_shared().unwrap();
1012
1013        let table = receiver.recv_timeout(Duration::from_secs(5)).unwrap();
1014        read_handle.join().unwrap();
1015        assert_eq!(table, EVENTS);
1016        assert_eq!(sink.dropped_frames(), 0);
1017    }
1018
1019    #[test]
1020    fn flush_counts_missing_sidecar_connection_as_dropped_frame() {
1021        let sink = UnixSocketSink::new();
1022        sink.set_path(socket_path("missing.sock")).unwrap();
1023        sink.consume_shared(&event()).unwrap();
1024
1025        sink.flush_shared().unwrap();
1026
1027        wait_for_dropped_frames(&sink, 1);
1028    }
1029}