1use 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
31pub 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
42pub 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
53pub struct OffThreadSpectraSink {
55 ndjson: Arc<NdjsonFileSink>,
56 tx: SyncSender<EmitJob>,
57 _writer: Arc<JoinHandle<()>>,
58}
59
60impl OffThreadSpectraSink {
61 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 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
184pub 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}