1use 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
75const WORKER_QUEUE_CAPACITY: usize = 10_000;
77
78static UNIX_SOCKET_SINK: OnceLock<Arc<UnixSocketSink>> = OnceLock::new();
80
81pub struct UnixSocketSink {
83 inner: Mutex<UnixSocketSinkInner>,
86 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 sink: Arc<UnixSocketSink>,
126 target_filter: Targets,
127}
128
129impl UnixSocketSink {
130 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 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 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 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 pub fn dropped_frames(&self) -> u64 {
199 self.dropped.load(Ordering::Relaxed)
200 }
201
202 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 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 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
395pub(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
405pub 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
413pub 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
421pub 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 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
517fn 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 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 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 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}