yggdrasil/api/logs/
capture.rs1use std::sync::atomic::{AtomicU64, Ordering};
20use std::sync::{LazyLock, Mutex};
21
22use chrono::{DateTime, Utc};
23use tokio::sync::{broadcast, mpsc};
24use tracing::field::{Field, Visit};
25use tracing::Event;
26use tracing::Subscriber;
27use tracing_subscriber::layer::{Context, Layer};
28use tracing_subscriber::EnvFilter;
29
30const MPSC_CAP: usize = 4096;
33
34const BROADCAST_CAP: usize = 1024;
36
37const MAX_MESSAGE_BYTES: usize = 4 * 1024;
39
40const EXCLUDED_TARGETS: [&str; 2] = ["yggdrasil::api::logs", "yggdrasil::tasks::log_writer"];
42
43#[derive(Debug, Clone)]
45pub struct LogRecord {
46 pub ts: DateTime<Utc>,
48 pub level: String,
50 pub target: String,
52 pub message: String,
54}
55
56struct LogChannels {
58 db_tx: mpsc::Sender<LogRecord>,
60 db_rx: Mutex<Option<mpsc::Receiver<LogRecord>>>,
62 live_tx: broadcast::Sender<LogRecord>,
64 dropped: AtomicU64,
66}
67
68static CHANNELS: LazyLock<LogChannels> = LazyLock::new(|| {
70 let (db_tx, db_rx) = mpsc::channel(MPSC_CAP);
71 let (live_tx, _) = broadcast::channel(BROADCAST_CAP);
72 LogChannels {
73 db_tx,
74 db_rx: Mutex::new(Some(db_rx)),
75 live_tx,
76 dropped: AtomicU64::new(0),
77 }
78});
79
80pub fn take_db_receiver() -> Option<mpsc::Receiver<LogRecord>> {
82 CHANNELS.db_rx.lock().ok().and_then(|mut g| g.take())
83}
84
85pub fn subscribe_live() -> broadcast::Receiver<LogRecord> {
87 CHANNELS.live_tx.subscribe()
88}
89
90pub fn dropped_count() -> u64 {
92 CHANNELS.dropped.load(Ordering::Relaxed)
93}
94
95pub fn record_dropped(n: u64) {
97 CHANNELS.dropped.fetch_add(n, Ordering::Relaxed);
98}
99
100pub fn log_viewer_filter() -> EnvFilter {
103 std::env::var("LOG_VIEWER_LEVEL")
104 .ok()
105 .and_then(|v| EnvFilter::try_new(v).ok())
106 .unwrap_or_else(|| EnvFilter::new("info"))
107}
108
109pub struct CaptureLayer;
115
116impl<S: Subscriber> Layer<S> for CaptureLayer {
117 fn on_event(&self, event: &Event<'_>, _ctx: Context<'_, S>) {
118 let meta = event.metadata();
119 let target = meta.target();
120
121 if EXCLUDED_TARGETS.iter().any(|p| target.starts_with(p)) {
123 return;
124 }
125
126 let mut visitor = MessageVisitor::default();
128 event.record(&mut visitor);
129 let mut message = visitor.message.unwrap_or_default();
130 for (key, value) in visitor.extra {
131 message.push(' ');
132 message.push_str(&key);
133 message.push('=');
134 message.push_str(&value);
135 }
136 truncate_message(&mut message);
137
138 let record = LogRecord {
139 ts: Utc::now(),
140 level: meta.level().as_str().to_string(),
141 target: target.to_string(),
142 message,
143 };
144
145 if CHANNELS.db_tx.try_send(record.clone()).is_err() {
147 CHANNELS.dropped.fetch_add(1, Ordering::Relaxed);
148 }
149
150 if CHANNELS.live_tx.receiver_count() > 0 {
153 let _ = CHANNELS.live_tx.send(record);
154 }
155 }
156}
157
158#[derive(Default)]
160struct MessageVisitor {
161 message: Option<String>,
162 extra: Vec<(String, String)>,
163}
164
165impl Visit for MessageVisitor {
166 fn record_debug(&mut self, field: &Field, value: &dyn std::fmt::Debug) {
169 if field.name() == "message" {
170 self.message = Some(format!("{value:?}"));
171 } else {
172 self.extra
173 .push((field.name().to_string(), format!("{value:?}")));
174 }
175 }
176
177 fn record_str(&mut self, field: &Field, value: &str) {
179 if field.name() == "message" {
180 self.message = Some(value.to_string());
181 } else {
182 self.extra
183 .push((field.name().to_string(), value.to_string()));
184 }
185 }
186}
187
188fn truncate_message(s: &mut String) {
190 if s.len() > MAX_MESSAGE_BYTES {
191 let mut end = MAX_MESSAGE_BYTES;
192 while !s.is_char_boundary(end) {
193 end -= 1;
194 }
195 s.truncate(end);
196 }
197}
198
199#[cfg(all(test, feature = "server"))]
200mod tests {
201 use super::*;
202
203 #[test]
204 fn truncate_respects_char_boundary() {
205 let mut s = "a".repeat(100);
207 truncate_message(&mut s);
208 assert_eq!(s.len(), 100);
209
210 let mut s = "日".repeat(MAX_MESSAGE_BYTES); truncate_message(&mut s);
213 assert!(s.len() <= MAX_MESSAGE_BYTES);
214 assert!(s.is_char_boundary(s.len()));
215 }
216
217 #[test]
218 fn dropped_counter_accumulates() {
219 let before = dropped_count();
220 record_dropped(7);
221 assert_eq!(dropped_count(), before + 7);
222 }
223
224 #[test]
225 fn db_receiver_take_once() {
226 let first = take_db_receiver();
229 let second = take_db_receiver();
230 assert!(second.is_none());
231 drop(first);
232 }
233}