scirs2_core/structured_logging/
logger.rs1use std::io::{self, Write};
18use std::sync::{Arc, Mutex, RwLock};
19use std::time::{SystemTime, UNIX_EPOCH};
20
21use super::types::{FieldValue, LogLevel, LogRecord, TraceConfig};
22
23pub trait LogSink: Send + Sync {
31 fn write(&self, record: &LogRecord);
33}
34
35pub struct ConsoleSink {
41 use_stdout: bool,
42}
43
44impl ConsoleSink {
45 pub fn new() -> Self {
47 Self { use_stdout: false }
48 }
49
50 pub fn stdout() -> Self {
52 Self { use_stdout: true }
53 }
54
55 pub fn format_json(record: &LogRecord) -> String {
59 let mut fields_json = String::new();
60 for (i, (k, v)) in record.fields.iter().enumerate() {
61 if i > 0 {
62 fields_json.push(',');
63 }
64 let escaped_key = k.replace('\\', "\\\\").replace('"', "\\\"");
65 fields_json.push('"');
66 fields_json.push_str(&escaped_key);
67 fields_json.push_str("\":");
68 fields_json.push_str(&v.to_json_value());
69 }
70
71 let mut out = format!(
72 "{{\"level\":\"{}\",\"msg\":{},\"ts\":{}",
73 record.level.as_str(),
74 FieldValue::Str(record.message.clone()).to_json_value(),
75 record.timestamp_ns,
76 );
77
78 if let Some(tid) = record.trace_id {
79 out.push_str(&format!(",\"trace_id\":{}", tid));
80 }
81 if let Some(sid) = record.span_id {
82 out.push_str(&format!(",\"span_id\":{}", sid));
83 }
84 if !fields_json.is_empty() {
85 out.push_str(&format!(",\"fields\":{{{}}}", fields_json));
86 }
87 out.push('}');
88 out
89 }
90}
91
92impl Default for ConsoleSink {
93 fn default() -> Self {
94 Self::new()
95 }
96}
97
98impl LogSink for ConsoleSink {
99 fn write(&self, record: &LogRecord) {
100 let line = Self::format_json(record);
101 if self.use_stdout {
102 let stdout = io::stdout();
103 let mut handle = stdout.lock();
104 let _ = writeln!(handle, "{}", line);
105 } else {
106 let stderr = io::stderr();
107 let mut handle = stderr.lock();
108 let _ = writeln!(handle, "{}", line);
109 }
110 }
111}
112
113pub struct MemorySink {
119 records: Mutex<Vec<LogRecord>>,
120}
121
122impl MemorySink {
123 pub fn new() -> Self {
125 Self {
126 records: Mutex::new(Vec::new()),
127 }
128 }
129
130 pub fn records(&self) -> Vec<LogRecord> {
132 self.records.lock().map(|g| g.clone()).unwrap_or_default()
133 }
134
135 pub fn clear(&self) {
137 if let Ok(mut g) = self.records.lock() {
138 g.clear();
139 }
140 }
141
142 pub fn len(&self) -> usize {
144 self.records.lock().map(|g| g.len()).unwrap_or(0)
145 }
146
147 pub fn is_empty(&self) -> bool {
149 self.len() == 0
150 }
151}
152
153impl Default for MemorySink {
154 fn default() -> Self {
155 Self::new()
156 }
157}
158
159impl LogSink for MemorySink {
160 fn write(&self, record: &LogRecord) {
161 if let Ok(mut g) = self.records.lock() {
162 g.push(record.clone());
163 }
164 }
165}
166
167#[derive(Debug, Clone)]
175pub struct OtlpLogRecord {
176 pub time_unix_nano: u64,
178 pub severity_number: u32,
180 pub severity_text: String,
182 pub body: String,
184 pub attributes: Vec<(String, String)>,
186 pub trace_id: Option<u64>,
188 pub span_id: Option<u64>,
190}
191
192pub struct OtelLogSink {
194 records: Mutex<Vec<OtlpLogRecord>>,
195}
196
197impl OtelLogSink {
198 pub fn new() -> Self {
200 Self {
201 records: Mutex::new(Vec::new()),
202 }
203 }
204
205 pub fn to_otlp(record: &LogRecord) -> OtlpLogRecord {
207 let attributes: Vec<(String, String)> = record
208 .fields
209 .iter()
210 .map(|(k, v)| (k.clone(), v.to_json_value()))
211 .collect();
212
213 OtlpLogRecord {
214 time_unix_nano: record.timestamp_ns,
215 severity_number: record.level.to_otlp_severity(),
216 severity_text: record.level.as_str().to_owned(),
217 body: record.message.clone(),
218 attributes,
219 trace_id: record.trace_id,
220 span_id: record.span_id,
221 }
222 }
223
224 pub fn otlp_records(&self) -> Vec<OtlpLogRecord> {
226 self.records.lock().map(|g| g.clone()).unwrap_or_default()
227 }
228}
229
230impl Default for OtelLogSink {
231 fn default() -> Self {
232 Self::new()
233 }
234}
235
236impl LogSink for OtelLogSink {
237 fn write(&self, record: &LogRecord) {
238 let otlp = Self::to_otlp(record);
239 if let Ok(mut g) = self.records.lock() {
240 g.push(otlp);
241 }
242 }
243}
244
245pub struct LogBuilder {
251 fields: Vec<(String, FieldValue)>,
252 trace_id: Option<u64>,
253 span_id: Option<u64>,
254 logger: Arc<StructuredLogger>,
255}
256
257impl LogBuilder {
258 fn new(logger: Arc<StructuredLogger>) -> Self {
259 Self {
260 fields: Vec::new(),
261 trace_id: None,
262 span_id: None,
263 logger,
264 }
265 }
266
267 pub fn field(mut self, key: impl Into<String>, value: impl Into<FieldValue>) -> Self {
269 self.fields.push((key.into(), value.into()));
270 self
271 }
272
273 pub fn span_context(mut self, trace_id: u64, span_id: u64) -> Self {
275 self.trace_id = Some(trace_id);
276 self.span_id = Some(span_id);
277 self
278 }
279
280 pub fn emit(self, level: LogLevel, message: impl Into<String>) {
282 let mut record = LogRecord::new(level, message);
283 record.fields = self.fields;
284 record.trace_id = self.trace_id;
285 record.span_id = self.span_id;
286 self.logger.log(&record);
287 }
288
289 pub fn info(self, message: impl Into<String>) {
291 self.emit(LogLevel::Info, message);
292 }
293
294 pub fn warn(self, message: impl Into<String>) {
296 self.emit(LogLevel::Warn, message);
297 }
298
299 pub fn error(self, message: impl Into<String>) {
301 self.emit(LogLevel::Error, message);
302 }
303
304 pub fn debug(self, message: impl Into<String>) {
306 self.emit(LogLevel::Debug, message);
307 }
308}
309
310pub struct StructuredLogger {
330 min_level: LogLevel,
331 sinks: RwLock<Vec<Box<dyn LogSink>>>,
332}
333
334impl StructuredLogger {
335 pub fn new(min_level: LogLevel, sinks: Vec<Box<dyn LogSink>>) -> Self {
337 Self {
338 min_level,
339 sinks: RwLock::new(sinks),
340 }
341 }
342
343 pub fn log(&self, record: &LogRecord) {
345 if record.level < self.min_level {
347 return;
348 }
349 if let Ok(sinks) = self.sinks.read() {
350 for sink in sinks.iter() {
351 sink.write(record);
352 }
353 }
354 }
355
356 pub fn add_sink(&self, sink: Box<dyn LogSink>) {
358 if let Ok(mut sinks) = self.sinks.write() {
359 sinks.push(sink);
360 }
361 }
362
363 pub fn with_fields(self: &Arc<Self>, fields: Vec<(String, FieldValue)>) -> LogBuilder {
365 let mut builder = LogBuilder::new(Arc::clone(self));
366 builder.fields = fields;
367 builder
368 }
369
370 pub fn builder(self: &Arc<Self>) -> LogBuilder {
372 LogBuilder::new(Arc::clone(self))
373 }
374
375 pub fn info(&self, msg: &str) {
377 let record = LogRecord::new(LogLevel::Info, msg);
378 self.log(&record);
379 }
380
381 pub fn warn(&self, msg: &str) {
383 let record = LogRecord::new(LogLevel::Warn, msg);
384 self.log(&record);
385 }
386
387 pub fn error(&self, msg: &str) {
389 let record = LogRecord::new(LogLevel::Error, msg);
390 self.log(&record);
391 }
392
393 pub fn debug(&self, msg: &str) {
395 let record = LogRecord::new(LogLevel::Debug, msg);
396 self.log(&record);
397 }
398
399 pub fn trace(&self, msg: &str) {
401 let record = LogRecord::new(LogLevel::Trace, msg);
402 self.log(&record);
403 }
404}
405
406pub fn init_logger(_config: TraceConfig) -> Arc<StructuredLogger> {
414 Arc::new(StructuredLogger::new(
415 LogLevel::Info,
416 vec![Box::new(ConsoleSink::new())],
417 ))
418}
419
420pub(crate) fn now_ns() -> u64 {
422 SystemTime::now()
423 .duration_since(UNIX_EPOCH)
424 .map(|d| d.as_nanos() as u64)
425 .unwrap_or(0)
426}
427
428#[cfg(test)]
433mod tests {
434 use super::*;
435
436 fn make_logger(level: LogLevel) -> (Arc<StructuredLogger>, Arc<MemorySink>) {
437 let sink = Arc::new(MemorySink::new());
438 struct SharedSink(Arc<MemorySink>);
440 impl LogSink for SharedSink {
441 fn write(&self, record: &LogRecord) {
442 self.0.write(record);
443 }
444 }
445 let logger = Arc::new(StructuredLogger::new(
446 level,
447 vec![Box::new(SharedSink(Arc::clone(&sink)))],
448 ));
449 (logger, sink)
450 }
451
452 #[test]
453 fn test_log_record_fields() {
454 let record = LogRecord::new(LogLevel::Info, "test").with_field("key", FieldValue::Int(42));
455 assert_eq!(record.message, "test");
456 assert_eq!(record.fields.len(), 1);
457 assert_eq!(record.fields[0].0, "key");
458 }
459
460 #[test]
461 fn test_memory_sink_captures() {
462 let (logger, sink) = make_logger(LogLevel::Debug);
463 logger.info("hello");
464 logger.warn("world");
465 let recs = sink.records();
466 assert_eq!(recs.len(), 2);
467 assert_eq!(recs[0].message, "hello");
468 assert_eq!(recs[1].level, LogLevel::Warn);
469 }
470
471 #[test]
472 fn test_log_level_filter() {
473 let (logger, sink) = make_logger(LogLevel::Warn);
474 logger.debug("should be filtered");
475 logger.info("also filtered");
476 logger.warn("passes");
477 let recs = sink.records();
478 assert_eq!(recs.len(), 1);
479 assert_eq!(recs[0].message, "passes");
480 }
481
482 #[test]
483 fn test_log_builder_fluent() {
484 let (logger, sink) = make_logger(LogLevel::Debug);
485 logger
486 .builder()
487 .field("user", "alice")
488 .field("count", 7i64)
489 .info("user logged in");
490 let recs = sink.records();
491 assert_eq!(recs.len(), 1);
492 assert_eq!(recs[0].fields.len(), 2);
493 }
494
495 #[test]
496 fn test_console_sink_json_format() {
497 let record = LogRecord::new(LogLevel::Info, "msg").with_field("k", FieldValue::Int(1));
498 let json = ConsoleSink::format_json(&record);
499 assert!(json.contains("\"level\":\"INFO\""));
500 assert!(json.contains("\"msg\":\"msg\""));
501 assert!(json.contains("\"k\":1"));
502 }
503
504 #[test]
505 fn test_otlp_log_conversion() {
506 let record =
507 LogRecord::new(LogLevel::Error, "boom").with_field("code", FieldValue::Int(500));
508 let otlp = OtelLogSink::to_otlp(&record);
509 assert_eq!(otlp.severity_text, "ERROR");
510 assert_eq!(otlp.severity_number, 17);
511 assert_eq!(otlp.body, "boom");
512 assert!(!otlp.attributes.is_empty());
513 }
514}