somatize_runtime/tracking/
jsonl_sink.rs1use chrono::Utc;
4use somatize_core::event::Event;
5use somatize_core::tracking::{EventEnvelope, EventSink};
6use std::fs::{File, OpenOptions};
7use std::io::{BufWriter, Write};
8use std::path::Path;
9use std::sync::Mutex;
10use std::sync::atomic::{AtomicU64, Ordering};
11
12pub struct JsonlEventSink {
20 events: Mutex<BufWriter<File>>,
21 metrics: Option<Mutex<BufWriter<File>>>,
22 seq: AtomicU64,
23 flush_every: u64,
24}
25
26impl JsonlEventSink {
27 pub fn create(
29 events_path: &Path,
30 metrics_path: Option<&Path>,
31 flush_every: usize,
32 ) -> std::io::Result<Self> {
33 Self::open(events_path, metrics_path, flush_every, 0, false)
34 }
35
36 pub fn append(
38 events_path: &Path,
39 metrics_path: Option<&Path>,
40 flush_every: usize,
41 start_seq: u64,
42 ) -> std::io::Result<Self> {
43 Self::open(events_path, metrics_path, flush_every, start_seq, true)
44 }
45
46 fn open(
47 events_path: &Path,
48 metrics_path: Option<&Path>,
49 flush_every: usize,
50 start_seq: u64,
51 append: bool,
52 ) -> std::io::Result<Self> {
53 let open = |path: &Path| -> std::io::Result<BufWriter<File>> {
54 let mut opts = OpenOptions::new();
55 opts.create(true).write(true);
56 if append {
57 opts.append(true);
58 } else {
59 opts.truncate(true);
60 }
61 Ok(BufWriter::new(opts.open(path)?))
62 };
63 Ok(Self {
64 events: Mutex::new(open(events_path)?),
65 metrics: metrics_path.map(|p| open(p).map(Mutex::new)).transpose()?,
66 seq: AtomicU64::new(start_seq),
67 flush_every: flush_every.max(1) as u64,
68 })
69 }
70
71 pub fn next_seq(&self) -> u64 {
74 self.seq.load(Ordering::SeqCst)
75 }
76
77 fn metric_line(event: &Event) -> Option<serde_json::Value> {
79 match event {
80 Event::TrialMetric {
81 study_id: _,
82 trial_id,
83 metric,
84 } => Some(serde_json::json!({
85 "ts": metric.timestamp,
86 "name": metric.name,
87 "value": metric.value,
88 "step": metric.step,
89 "trial_id": trial_id,
90 "node_id": null,
91 })),
92 Event::MetricReported {
93 run_id: _,
94 metric,
95 node_id,
96 trial_id,
97 } => Some(serde_json::json!({
98 "ts": metric.timestamp,
99 "name": metric.name,
100 "value": metric.value,
101 "step": metric.step,
102 "trial_id": trial_id,
103 "node_id": node_id,
104 })),
105 _ => None,
106 }
107 }
108
109 fn write_line(writer: &Mutex<BufWriter<File>>, line: &str, do_flush: bool) {
110 let mut guard = match writer.lock() {
111 Ok(g) => g,
112 Err(poisoned) => poisoned.into_inner(),
113 };
114 if let Err(e) = writeln!(guard, "{line}") {
115 tracing::warn!("tracking: failed to write event line: {e}");
116 return;
117 }
118 if do_flush && let Err(e) = guard.flush() {
119 tracing::warn!("tracking: failed to flush event log: {e}");
120 }
121 }
122}
123
124impl EventSink for JsonlEventSink {
125 fn record(&self, event: &Event) {
126 let seq = self.seq.fetch_add(1, Ordering::SeqCst);
127 let envelope = EventEnvelope {
128 seq,
129 ts: Utc::now(),
130 event: event.clone(),
131 };
132 let do_flush = seq % self.flush_every == self.flush_every - 1;
133 match serde_json::to_string(&envelope) {
134 Ok(line) => Self::write_line(&self.events, &line, do_flush),
135 Err(e) => tracing::warn!("tracking: failed to serialize event: {e}"),
136 }
137 if let Some(metrics) = &self.metrics
138 && let Some(line) = Self::metric_line(event)
139 {
140 Self::write_line(metrics, &line.to_string(), do_flush);
141 }
142 }
143
144 fn flush(&self) {
145 for writer in std::iter::once(&self.events).chain(self.metrics.iter()) {
146 let mut guard = match writer.lock() {
147 Ok(g) => g,
148 Err(poisoned) => poisoned.into_inner(),
149 };
150 if let Err(e) = guard.flush() {
151 tracing::warn!("tracking: failed to flush event log: {e}");
152 }
153 }
154 }
155}
156
157impl Drop for JsonlEventSink {
158 fn drop(&mut self) {
159 self.flush();
160 }
161}