Skip to main content

yggdrasil/tasks/
log_writer.rs

1//! 日志批量落库后台任务(server-only)。
2//!
3//! 从 [`crate::api::logs::capture`] 的 mpsc 通道接收捕获的日志事件,
4//! 攒批(200 条或 500ms 窗口,先到先发)后用一条 `INSERT ... UNNEST`
5//! 批量写入 logs 表。
6//!
7//! 失败语义:
8//! - INSERT 失败:记 error 日志(target 在 capture 的防递归排除名单内,
9//!   不会回流进管道)、整批丢弃、按批大小累加 dropped 计数——
10//!   日志查看器允许丢数据,绝不允许反过来拖垮主服务;
11//! - `get_conn()` 失败:保留本批,sleep 后重试(不丢批也不重试风暴),
12//!   重试期间通道继续缓冲,缓冲溢出走 capture 的 dropped 计数。
13
14use std::time::Duration;
15
16use crate::api::logs::capture::{self, LogRecord};
17use crate::db::pool::get_conn;
18
19/// 单批最大条数。
20const BATCH_SIZE: usize = 200;
21/// 攒批窗口:首条到达后最长等待这么久就冲刷(不足一批也发)。
22const FLUSH_WINDOW: Duration = Duration::from_millis(500);
23/// `get_conn()` 失败后的重试间隔(保留本批,不丢)。
24const CONN_RETRY_INTERVAL: Duration = Duration::from_secs(2);
25
26/// 启动日志落库循环(serve() 内 spawn 一次)。
27pub async fn run_writer() {
28    let mut rx = match capture::take_db_receiver() {
29        Some(rx) => rx,
30        None => {
31            // target 在 capture 排除名单内,此日志只到控制台,不会回流。
32            tracing::error!("log writer: capture receiver already taken; task exiting");
33            return;
34        }
35    };
36    tracing::info!(
37        batch_size = BATCH_SIZE,
38        flush_window_ms = FLUSH_WINDOW.as_millis() as u64,
39        "log writer started"
40    );
41
42    let mut batch: Vec<LogRecord> = Vec::with_capacity(BATCH_SIZE);
43    loop {
44        // 阻塞等首条;通道关闭(进程退出)时冲刷余量后退出。
45        let first = match rx.recv().await {
46            Some(r) => r,
47            None => {
48                flush(&mut batch).await;
49                return;
50            }
51        };
52        batch.push(first);
53
54        // 窗口内尽力拉满一批;窗口到点或通道关闭即冲刷。
55        let deadline = tokio::time::Instant::now() + FLUSH_WINDOW;
56        while batch.len() < BATCH_SIZE {
57            match tokio::time::timeout_at(deadline, rx.recv()).await {
58                Ok(Some(r)) => batch.push(r),
59                Ok(None) | Err(_) => break,
60            }
61        }
62        flush(&mut batch).await;
63    }
64}
65
66/// 冲刷一批:UNNEST 数组批量 INSERT。空批直接返回。
67async fn flush(batch: &mut Vec<LogRecord>) {
68    if batch.is_empty() {
69        return;
70    }
71
72    // 连接失败:保留本批,sleep 后重试——启动窗口内 DB 可能尚未就绪,
73    // 运行期 DB 短暂抖动也不该丢批。
74    let client = loop {
75        match get_conn().await {
76            Ok(c) => break c,
77            Err(e) => {
78                tracing::error!(error = %e, "log writer: failed to get DB connection; retrying");
79                tokio::time::sleep(CONN_RETRY_INTERVAL).await;
80            }
81        }
82    };
83
84    let ts: Vec<chrono::DateTime<chrono::Utc>> = batch.iter().map(|r| r.ts).collect();
85    let levels: Vec<&str> = batch.iter().map(|r| r.level.as_str()).collect();
86    let targets: Vec<&str> = batch.iter().map(|r| r.target.as_str()).collect();
87    let messages: Vec<&str> = batch.iter().map(|r| r.message.as_str()).collect();
88
89    let result = client
90        .execute(
91            "INSERT INTO logs (ts, level, target, message) \
92             SELECT * FROM UNNEST($1::timestamptz[], $2::text[], $3::text[], $4::text[])",
93            &[&ts, &levels, &targets, &messages],
94        )
95        .await;
96
97    if let Err(e) = result {
98        // target 在 capture 排除名单内,此 error 不会回流进管道。
99        tracing::error!(
100            error = %e,
101            dropped = batch.len() as u64,
102            "log writer: batch insert failed; dropping batch"
103        );
104        capture::record_dropped(batch.len() as u64);
105    }
106    batch.clear();
107}