feat(logging): 增加 Seq 日志集成与标准输出重定向

- 新增 SeqWriter 支持,将日志异步推送到 Seq 平台,支持多进程实例区分
- 实现 redirectStdout 函数,重定向 os.Stdout 和标准 log 包输出,确保 fmt.Println 等输出被捕获并记录
- 更新 README 文档,增加 Seq 日志集成的说明与配置链接
- 扩展 Logger 结构,支持动态添加输出目标,提升日志记录灵活性
This commit is contained in:
2026-05-18 13:12:51 +08:00
parent ef10437a5a
commit 94a251ec44
5 changed files with 388 additions and 5 deletions
+214 -5
View File
@@ -2,13 +2,17 @@ package log
import (
"bufio"
"bytes"
"encoding/json"
"fmt"
"io"
"net/http"
"os"
"path/filepath"
"runtime"
"strings"
"sync"
"sync/atomic"
"time"
"github.com/rs/zerolog"
@@ -25,6 +29,7 @@ type ErrorRecord struct {
// Logger 日志核心结构体,封装 zerolog
type Logger struct {
zl zerolog.Logger
mw *multiWriter // 动态多输出,支持运行时追加 SeqWriter 等
logLevel int
errors []ErrorRecord
errIdx int // 环形缓冲写入位置
@@ -69,8 +74,8 @@ func NewLogger(logLevel int, logFile string, maxErrors int) *Logger {
writers = append(writers, fw)
}
multi := zerolog.MultiLevelWriter(writers...)
zl := zerolog.New(multi).
mw := &multiWriter{writers: writers}
zl := zerolog.New(mw).
Level(level).
With().
Timestamp().
@@ -79,6 +84,7 @@ func NewLogger(logLevel int, logFile string, maxErrors int) *Logger {
l := &Logger{
zl: zl,
mw: mw,
logLevel: logLevel,
maxErrors: maxErrors,
}
@@ -116,8 +122,8 @@ func NewLoggerNoCaller(logLevel int, logFile string, maxErrors int) *Logger {
writers = append(writers, fw)
}
multi := zerolog.MultiLevelWriter(writers...)
zl := zerolog.New(multi).
mw := &multiWriter{writers: writers}
zl := zerolog.New(mw).
Level(level).
With().
Timestamp().
@@ -125,6 +131,7 @@ func NewLoggerNoCaller(logLevel int, logFile string, maxErrors int) *Logger {
l := &Logger{
zl: zl,
mw: mw,
logLevel: logLevel,
maxErrors: maxErrors,
}
@@ -145,6 +152,7 @@ func NewTestLogger() *Logger {
func (l *Logger) SetOutput(w io.Writer) {
if w == io.Discard {
l.zl = zerolog.New(w).Level(zerolog.Disabled)
l.mw = nil
return
}
consoleWriter := zerolog.ConsoleWriter{
@@ -159,7 +167,9 @@ func (l *Logger) SetOutput(w io.Writer) {
return fmt.Sprintf("[%s]", i)
},
}
l.zl = zerolog.New(consoleWriter).
mw := &multiWriter{writers: []io.Writer{consoleWriter}}
l.mw = mw
l.zl = zerolog.New(mw).
With().
Timestamp().
Caller().
@@ -422,6 +432,205 @@ func shortenPath(file string) string {
return file
}
// --- multiWriter: 支持运行时追加输出目标 ---
// multiWriter 实现 zerolog.LevelWriter,允许在 Logger 创建后动态追加新的输出(如 SeqWriter)
type multiWriter struct {
mu sync.RWMutex
writers []io.Writer
}
func (mw *multiWriter) Write(p []byte) (int, error) {
mw.mu.RLock()
defer mw.mu.RUnlock()
for _, w := range mw.writers {
_, _ = w.Write(p)
}
return len(p), nil
}
// WriteLevel 实现 zerolog.LevelWriter,使 zerolog 的级别过滤正确传递给各子 Writer
func (mw *multiWriter) WriteLevel(l zerolog.Level, p []byte) (int, error) {
mw.mu.RLock()
defer mw.mu.RUnlock()
for _, w := range mw.writers {
if lw, ok := w.(zerolog.LevelWriter); ok {
_, _ = lw.WriteLevel(l, p)
} else {
_, _ = w.Write(p)
}
}
return len(p), nil
}
func (mw *multiWriter) add(w io.Writer) {
mw.mu.Lock()
mw.writers = append(mw.writers, w)
mw.mu.Unlock()
}
// --- SeqWriter: 异步推送到 SeqCLEF 格式)---
// SeqWriter 将 zerolog JSON 日志通过 channel 队列异步发送到 Seq。
// Write() 仅做 channel <- bytesO(1) 非阻塞,绝不阻塞 web 请求处理 goroutine。
// channel 容量 10000,满时丢弃新条目(dropped 计数),不会影响主服务。
type SeqWriter struct {
seqUrl string
apiKey string
instance string
queue chan []byte
dropped int64
client *http.Client
}
func newSeqWriter(seqUrl, apiKey, instance string) *SeqWriter {
w := &SeqWriter{
seqUrl: strings.TrimRight(seqUrl, "/"),
apiKey: apiKey,
instance: instance,
queue: make(chan []byte, 10000),
client: &http.Client{Timeout: 5 * time.Second},
}
go w.run()
return w
}
// Write 将字节塞入 channel 即返回,由后台 goroutine 消费发送
func (w *SeqWriter) Write(p []byte) (int, error) {
clef := w.toClef(p)
select {
case w.queue <- clef:
default:
atomic.AddInt64(&w.dropped, 1)
}
return len(p), nil
}
// toClef 将 zerolog JSON 字段映射到 Seq CLEF 格式
// zerolog: time/level/message(msg) → CLEF: @t/@l/@mt
func (w *SeqWriter) toClef(p []byte) []byte {
trimmed := bytes.TrimSpace(p)
var m map[string]interface{}
if err := json.Unmarshal(trimmed, &m); err != nil {
// 非 JSON(如 fmt.Println 直接输出),包装为纯文本日志
return []byte(fmt.Sprintf(`{"@t":%q,"@l":"Information","@mt":%q,"source":"stdout","instance":%q}`,
time.Now().Format(time.RFC3339), strings.TrimSpace(string(p)), w.instance))
}
clef := make(map[string]interface{}, len(m)+3)
for k, v := range m {
clef[k] = v
}
// time → @tzerolog 格式 "2006-01-02 15:04:05" → ISO 8601
if t, ok := m["time"].(string); ok {
delete(clef, "time")
if pt, err := time.ParseInLocation("2006-01-02 15:04:05", t, time.Local); err == nil {
clef["@t"] = pt.Format(time.RFC3339)
} else {
clef["@t"] = t
}
} else {
clef["@t"] = time.Now().Format(time.RFC3339)
}
// level → @l
if lv, ok := m["level"].(string); ok {
delete(clef, "level")
switch lv {
case "debug":
clef["@l"] = "Debug"
case "info":
clef["@l"] = "Information"
case "warn":
clef["@l"] = "Warning"
case "error":
clef["@l"] = "Error"
case "fatal":
clef["@l"] = "Fatal"
default:
clef["@l"] = lv
}
}
// message/msg → @mt
if msg, ok := m["message"].(string); ok {
delete(clef, "message")
clef["@mt"] = msg
} else if msg, ok := m["msg"].(string); ok {
delete(clef, "msg")
clef["@mt"] = msg
}
if w.instance != "" {
clef["instance"] = w.instance
}
b, _ := json.Marshal(clef)
return b
}
// run 后台 goroutine:每 100 条或 500ms 批量 POST 到 Seq
func (w *SeqWriter) run() {
batch := make([][]byte, 0, 100)
ticker := time.NewTicker(500 * time.Millisecond)
defer ticker.Stop()
for {
select {
case entry := <-w.queue:
batch = append(batch, entry)
if len(batch) >= 100 {
w.flush(batch)
batch = batch[:0]
}
case <-ticker.C:
if len(batch) > 0 {
w.flush(batch)
batch = batch[:0]
}
}
}
}
func (w *SeqWriter) flush(batch [][]byte) {
var buf bytes.Buffer
for _, entry := range batch {
buf.Write(entry)
buf.WriteByte('\n')
}
req, err := http.NewRequest("POST", w.seqUrl+"/api/events/raw?clef", &buf)
if err != nil {
fmt.Fprintf(os.Stderr, "[seq] build request error: %v\n", err)
return
}
req.Header.Set("Content-Type", "application/vnd.serilog.clef")
if w.apiKey != "" {
req.Header.Set("X-Seq-ApiKey", w.apiKey)
}
resp, err := w.client.Do(req)
if err != nil {
fmt.Fprintf(os.Stderr, "[seq] send error: %v (dropped=%d)\n", err, atomic.LoadInt64(&w.dropped))
return
}
defer resp.Body.Close()
if resp.StatusCode >= 400 {
body, _ := io.ReadAll(resp.Body)
fmt.Fprintf(os.Stderr, "[seq] server returned %d: %s\n", resp.StatusCode, string(body))
}
}
// SetSeqWriter 将 Seq HTTP 推送器附加到当前 Logger(不影响已有的控制台/文件输出)
// seqUrl: Seq 服务地址,如 "http://127.0.0.1:5341"
// apiKey: Seq API Key,免费单用户版留空
// instance: 实例标识,建议用 config 的 port 字段区分同机多进程,如 "8085"
func (l *Logger) SetSeqWriter(seqUrl, apiKey, instance string) {
if seqUrl == "" || l.mw == nil {
return
}
sw := newSeqWriter(seqUrl, apiKey, instance)
l.mw.add(sw)
}
// isInfrastructureFile 判断文件是否为基础设施(日志/DB/缓存/通用工具)管线,
// 这些文件在调用栈中会被跳过,以显示真正的业务调用者。
// context.go 是 Display/View 的纯中间层,也作为基础设施跳过。