Skip to main content

yggdrasil/api/logs/
capture.rs

1//! 进程内日志捕获 Layer(server-only)。
2//!
3//! 作为 `tracing_subscriber::Layer` 挂进 registry(见 main.rs),把全部 target
4//! 的日志事件复制一份送进两条无锁通道:
5//! - `mpsc`(容量 4096):[`crate::tasks::log_writer`] 攒批后批量 INSERT 进 logs 表;
6//! - `broadcast`(容量 1024):[`super::sse`] 的实时流订阅,按连接参数过滤后推送。
7//!
8//! 关键性质:
9//! - **绝不阻塞日志路径**:两条通道都用非阻塞发送。mpsc 满了只累加 dropped
10//!   计数;broadcast 满了最旧事件被覆盖,接收端走 `Lagged` → `gap` 事件。
11//! - **防递归**:`yggdrasil::api::logs` 与 `yggdrasil::tasks::log_writer`
12//!   前缀的 target 直接跳过——写库失败的 error 日志不能再进管道,否则
13//!   「写库失败 → 记日志 → 再写库 → 再失败」会自我放大。
14//! - **独立过滤**:Layer 在 main.rs 里以 per-layer `EnvFilter` 包裹
15//!   ([`log_viewer_filter`],读 `LOG_VIEWER_LEVEL`,默认 info,不吃 RUST_LOG)。
16//!   注意 `tracing` 依赖带 `release_max_level_info`:release 构建中
17//!   DEBUG/TRACE 在编译期即被剔除,env 调到 debug 也只有 info+。
18
19use 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
30/// 落库通道容量:writer 未启动(迁移窗口)或 DB 短暂不可达时的缓冲。
31/// 满了按条计入 dropped,绝不阻塞打日志的线程。
32const MPSC_CAP: usize = 4096;
33
34/// 实时流通道容量:SSE 客户端消费慢时最旧事件被覆盖(接收端收 Lagged)。
35const BROADCAST_CAP: usize = 1024;
36
37/// 单条消息上限(字节):防止异常大日志撑爆内存与表行。
38const MAX_MESSAGE_BYTES: usize = 4 * 1024;
39
40/// 防递归 target 前缀:日志管道自身的日志不再进管道。
41const EXCLUDED_TARGETS: [&str; 2] = ["yggdrasil::api::logs", "yggdrasil::tasks::log_writer"];
42
43/// 内部捕获记录(writer 落库与 SSE 实时流共用的最小表示,非 serde DTO)。
44#[derive(Debug, Clone)]
45pub struct LogRecord {
46    /// 事件捕获时刻(UTC)。
47    pub ts: DateTime<Utc>,
48    /// 级别大写静态串(ERROR/WARN/INFO/DEBUG/TRACE)转 owned。
49    pub level: String,
50    /// tracing target(模块路径)。
51    pub target: String,
52    /// 消息文本(含追加的结构化字段,已截断至 4KB)。
53    pub message: String,
54}
55
56/// 进程级日志通道枢纽。
57struct LogChannels {
58    /// 落库通道发送端(Layer 每条事件 try_send 一份)。
59    db_tx: mpsc::Sender<LogRecord>,
60    /// 落库通道接收端,启动时被 [`crate::tasks::log_writer`] 取走一次(take 语义)。
61    db_rx: Mutex<Option<mpsc::Receiver<LogRecord>>>,
62    /// 实时流广播发送端(SSE 每连接 subscribe 一个 receiver)。
63    live_tx: broadcast::Sender<LogRecord>,
64    /// 丢弃计数:mpsc 满 / 写库失败丢批,逐条累加(get_logs 响应透出)。
65    dropped: AtomicU64,
66}
67
68/// 全局通道实例(LazyLock:Layer 是 'static 的,无法从外部注入)。
69static 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
80/// 取走落库通道接收端(一次性;第二次调用返回 None)。
81pub fn take_db_receiver() -> Option<mpsc::Receiver<LogRecord>> {
82    CHANNELS.db_rx.lock().ok().and_then(|mut g| g.take())
83}
84
85/// 订阅实时流广播(每 SSE 连接一个 receiver)。
86pub fn subscribe_live() -> broadcast::Receiver<LogRecord> {
87    CHANNELS.live_tx.subscribe()
88}
89
90/// 进程启动以来累计丢弃的日志条数。
91pub fn dropped_count() -> u64 {
92    CHANNELS.dropped.load(Ordering::Relaxed)
93}
94
95/// 累加丢弃计数(writer 写库失败丢批时按批大小调用)。
96pub fn record_dropped(n: u64) {
97    CHANNELS.dropped.fetch_add(n, Ordering::Relaxed);
98}
99
100/// capture 层的独立 EnvFilter:读 `LOG_VIEWER_LEVEL`,非法/缺失时回退 "info"。
101/// 故意不吃 `RUST_LOG`——控制台级别与查看器级别解耦。
102pub 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
109/// 日志捕获 Layer。在 main.rs 里以 `.with_filter(log_viewer_filter())` 包裹后
110/// 挂进 registry;本体的 `enabled` 恒 true,级别过滤全部交给 per-layer filter。
111/// 泛型实现(非特化 Registry):registry().with(fmt).with(capture) 组合时,
112/// 外层 S 是 Layered<...> 而非裸 Registry,特化实现会导致 Layered 不再满足
113/// `Into<Dispatch>`。
114pub 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        // 防递归:日志管道自身(查询/导出/SSE/writer)的日志不进管道。
122        if EXCLUDED_TARGETS.iter().any(|p| target.starts_with(p)) {
123            return;
124        }
125
126        // 提取 message 字段;其余结构化字段以 ` key=value` 追加(task_id 等不丢)。
127        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        // mpsc 满 / 已关闭:只增 dropped 计数,绝不阻塞日志路径。
146        if CHANNELS.db_tx.try_send(record.clone()).is_err() {
147            CHANNELS.dropped.fetch_add(1, Ordering::Relaxed);
148        }
149
150        // 实时流:无订阅者时 broadcast::send 只会返回 Err,直接跳过省一次分发;
151        // 通道满时最旧事件被覆盖,慢客户端走 RecvError::Lagged → gap 事件。
152        if CHANNELS.live_tx.receiver_count() > 0 {
153            let _ = CHANNELS.live_tx.send(record);
154        }
155    }
156}
157
158/// 事件字段提取:message 单独存,其余字段按声明顺序收集。
159#[derive(Default)]
160struct MessageVisitor {
161    message: Option<String>,
162    extra: Vec<(String, String)>,
163}
164
165impl Visit for MessageVisitor {
166    /// `info!("...", k = v)` 的 message(format_args)与非字符串字段走这里。
167    /// Arguments 的 Debug 输出即格式化文本,不带引号。
168    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    /// 显式字符串字段(`k = "v"`)走这里,不带引号。
178    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
188/// 按字节上限截断,回退到 char 边界防止切坏 UTF-8。
189fn 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        // 4KB 上限内全是 ASCII 时原样保留
206        let mut s = "a".repeat(100);
207        truncate_message(&mut s);
208        assert_eq!(s.len(), 100);
209
210        // 超限时截断,且多字节字符不被切碎
211        let mut s = "日".repeat(MAX_MESSAGE_BYTES); // 每字 3 字节
212        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        // 注意:本测试消费进程级单例,只能有一个测试做 take 语义断言。
227        // take 之后再次 take 必须得到 None(writer 不会重复启动)。
228        let first = take_db_receiver();
229        let second = take_db_receiver();
230        assert!(second.is_none());
231        drop(first);
232    }
233}