Skip to main content

somatize_runtime/tracking/
jsonl_sink.rs

1//! Append-only JSONL event sink.
2
3use 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
12/// Writes every event as one JSON line to `events.jsonl`, teeing
13/// metric-bearing events (`TrialMetric`, `MetricReported`) into a flat
14/// `metrics.jsonl` for cheap time-series reads.
15///
16/// Writes are buffered and flushed every `flush_every` events (and on
17/// [`EventSink::flush`]/drop). I/O errors are logged and swallowed —
18/// tracking must never take down a training run.
19pub 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    /// Create fresh log files (truncating any existing ones).
28    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    /// Open existing logs in append mode, continuing at `start_seq`.
37    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    /// Next sequence number to be assigned (== events recorded so far
72    /// when starting from zero).
73    pub fn next_seq(&self) -> u64 {
74        self.seq.load(Ordering::SeqCst)
75    }
76
77    /// Flat metric line for `metrics.jsonl`, if the event carries one.
78    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}