Skip to main content

spectra_runtime/
async_writer.rs

1//! Off-thread console + NDJSON I/O so the Spectra hot path only enqueues.
2
3use std::sync::mpsc::{self, SyncSender, TrySendError};
4use std::sync::{Arc, OnceLock};
5use std::thread::{self, JoinHandle};
6
7use serde_json::Value;
8use spectra_core::{NdjsonFileSink, SchemaRegistry, SpectraSink};
9
10const CHANNEL_CAPACITY: usize = 65_536;
11
12enum EmitJob {
13    Counter {
14        name: String,
15        labels: Vec<(String, String)>,
16        delta: i64,
17    },
18    Gauge {
19        name: String,
20        labels: Vec<(String, String)>,
21        value: f64,
22    },
23    Event {
24        table: String,
25        fields: Value,
26    },
27}
28
29static WRITER_TX: OnceLock<SyncSender<EmitJob>> = OnceLock::new();
30
31/// Returns true when stderr console mirror is enabled (default: true).
32pub fn console_mirror_enabled() -> bool {
33    static ENABLED: OnceLock<bool> = OnceLock::new();
34    *ENABLED.get_or_init(|| {
35        !matches!(
36            std::env::var("SPECTRA_CONSOLE").as_deref(),
37            Ok("0") | Ok("false") | Ok("FALSE")
38        )
39    })
40}
41
42/// Returns true when hot-path emit should use the off-thread writer (default: true).
43pub fn off_thread_emit_enabled() -> bool {
44    static ENABLED: OnceLock<bool> = OnceLock::new();
45    *ENABLED.get_or_init(|| {
46        !matches!(
47            std::env::var("SPECTRA_SYNC_HOT_PATH").as_deref(),
48            Ok("1") | Ok("true") | Ok("TRUE")
49        )
50    })
51}
52
53/// Sink that enqueues NDJSON + optional console mirror work for a background thread.
54pub struct OffThreadSpectraSink {
55    ndjson: Arc<NdjsonFileSink>,
56    tx: SyncSender<EmitJob>,
57    _writer: Arc<JoinHandle<()>>,
58}
59
60impl OffThreadSpectraSink {
61    /// Spawn the background writer thread and wire the shared emit queue.
62    pub fn new(ndjson: NdjsonFileSink) -> Self {
63        let ndjson = Arc::new(ndjson);
64        let (tx, rx) = mpsc::sync_channel(CHANNEL_CAPACITY);
65        let ndjson_for_thread = Arc::clone(&ndjson);
66        let writer = thread::Builder::new()
67            .name("spectra-async-writer".into())
68            .spawn(move || writer_loop(rx, ndjson_for_thread))
69            .expect("spawn spectra-async-writer");
70        let _ = WRITER_TX.set(tx.clone());
71        Self {
72            ndjson,
73            tx,
74            _writer: Arc::new(writer),
75        }
76    }
77
78    /// Underlying NDJSON file sink (used by the worker thread).
79    pub fn ndjson(&self) -> &NdjsonFileSink {
80        &self.ndjson
81    }
82}
83
84impl SpectraSink for OffThreadSpectraSink {
85    fn record_counter(&self, name: &str, labels: &[(&str, &str)], delta: i64) {
86        if !off_thread_emit_enabled() {
87            self.ndjson.record_counter(name, labels, delta);
88            return;
89        }
90        let job = EmitJob::Counter {
91            name: name.to_string(),
92            labels: labels
93                .iter()
94                .map(|(k, v)| (k.to_string(), v.to_string()))
95                .collect(),
96            delta,
97        };
98        try_enqueue(&self.tx, job);
99    }
100
101    fn record_gauge(&self, name: &str, labels: &[(&str, &str)], value: f64) {
102        if !off_thread_emit_enabled() {
103            self.ndjson.record_gauge(name, labels, value);
104            return;
105        }
106        let job = EmitJob::Gauge {
107            name: name.to_string(),
108            labels: labels
109                .iter()
110                .map(|(k, v)| (k.to_string(), v.to_string()))
111                .collect(),
112            value,
113        };
114        try_enqueue(&self.tx, job);
115    }
116
117    fn log_event(&self, table: &str, fields: &Value) {
118        if !off_thread_emit_enabled() {
119            if console_mirror_enabled() {
120                if let Some(line) = format_console_line(table, fields) {
121                    eprintln!("{line}");
122                }
123            }
124            self.ndjson.log_event(table, fields);
125            return;
126        }
127        let job = EmitJob::Event {
128            table: table.to_string(),
129            fields: fields.clone(),
130        };
131        try_enqueue(&self.tx, job);
132    }
133}
134
135fn try_enqueue(tx: &SyncSender<EmitJob>, job: EmitJob) {
136    match tx.try_send(job) {
137        Ok(()) => {}
138        Err(TrySendError::Full(_)) => {
139            log::warn!("[spectra:async_writer] emit queue full; dropping telemetry job");
140        }
141        Err(TrySendError::Disconnected(_)) => {
142            log::warn!("[spectra:async_writer] emit queue disconnected");
143        }
144    }
145}
146
147fn writer_loop(rx: mpsc::Receiver<EmitJob>, ndjson: Arc<NdjsonFileSink>) {
148    while let Ok(job) = rx.recv() {
149        match job {
150            EmitJob::Counter {
151                name,
152                labels,
153                delta,
154            } => {
155                let label_refs: Vec<(&str, &str)> = labels
156                    .iter()
157                    .map(|(k, v)| (k.as_str(), v.as_str()))
158                    .collect();
159                ndjson.record_counter(&name, &label_refs, delta);
160            }
161            EmitJob::Gauge {
162                name,
163                labels,
164                value,
165            } => {
166                let label_refs: Vec<(&str, &str)> = labels
167                    .iter()
168                    .map(|(k, v)| (k.as_str(), v.as_str()))
169                    .collect();
170                ndjson.record_gauge(&name, &label_refs, value);
171            }
172            EmitJob::Event { table, fields } => {
173                if console_mirror_enabled() {
174                    if let Some(line) = format_console_line(&table, &fields) {
175                        eprintln!("{line}");
176                    }
177                }
178                ndjson.log_event(&table, &fields);
179            }
180        }
181    }
182}
183
184/// Build a single `[spectra:console]` line when schema fields are console-safe.
185pub fn format_console_line(table: &str, fields: &Value) -> Option<String> {
186    let parts = mirror_safe_console_parts(table, fields)?;
187    Some(format!("[spectra:console] {table} {}", parts.join(" ")))
188}
189
190fn mirror_safe_console_parts(table: &str, fields: &Value) -> Option<Vec<String>> {
191    let meta = SchemaRegistry::global().get_schema(table)?;
192    let obj = fields.as_object()?;
193
194    let mut parts = Vec::new();
195    for field in &meta.fields {
196        if !field.classification.safe_for_console {
197            continue;
198        }
199        if let Some(v) = obj.get(&field.name) {
200            if let Some(s) = value_as_str(v) {
201                if !s.is_empty() {
202                    parts.push(format!("{}={s}", field.name));
203                }
204            }
205        }
206    }
207
208    if parts.is_empty() {
209        None
210    } else {
211        Some(parts)
212    }
213}
214
215fn value_as_str(v: &Value) -> Option<String> {
216    match v {
217        Value::String(s) => Some(s.clone()),
218        Value::Bool(b) => Some(b.to_string()),
219        Value::Number(n) => Some(n.to_string()),
220        _ => None,
221    }
222}