feat(log): Seq 会话追踪、控制台降噪与推送质量增强

请求日志自动带 sid/request_id,支持 LogBind;挂 Seq 后控制台仅 Warn+,Seq 仍全量;并修毫秒时间戳、@m、失败重试与停机 flush。

Co-authored-by: Cursor <cursoragent@cursor.com>
This commit is contained in:
2026-07-23 05:26:44 +08:00
parent b254606003
commit 831a7f017a
6 changed files with 554 additions and 197 deletions
+193 -62
View File
@@ -47,18 +47,65 @@ type ErrorRecord struct {
Caller string
}
// Logger 日志核心结构体,封装 zerolog
type Logger struct {
zl zerolog.Logger
mw *multiWriter // 动态多输出,支持运行时追加 SeqWriter 等
logLevel int
// errorStore 错误历史环形缓冲(指针共享,供 WithFields 派生子 Logger 共用)
type errorStore struct {
mu sync.RWMutex
errors []ErrorRecord
errIdx int // 环形缓冲写入位置
errCount int // 实际写入总数(用于判断缓冲是否满)
errorsMu sync.RWMutex
errIdx int
errCount int
maxErrors int
}
// Logger 日志核心结构体,封装 zerolog
type Logger struct {
zl zerolog.Logger
mw *multiWriter // 动态多输出,支持运行时追加 SeqWriter 等
console *levelFilterWriter // 控制台出口,挂 Seq 后可降到 Warn+
logLevel int
store *errorStore
seqs []*SeqWriter // 本 Logger 挂载过的 SeqWriter,供 CloseSeq 冲刷
seqsMu sync.Mutex
}
// levelFilterWriter 按级别过滤的 Writer(实现 zerolog.LevelWriter
type levelFilterWriter struct {
inner io.Writer
minLevel atomic.Int32 // zerolog.Level 存为 int32
}
func newLevelFilterWriter(inner io.Writer, min zerolog.Level) *levelFilterWriter {
w := &levelFilterWriter{inner: inner}
w.minLevel.Store(int32(min))
return w
}
func (w *levelFilterWriter) setMinLevel(min zerolog.Level) {
if w == nil {
return
}
w.minLevel.Store(int32(min))
}
func (w *levelFilterWriter) Write(p []byte) (int, error) {
// 无级别信息时按 Info 处理:低于门槛则丢弃
if zerolog.Level(w.minLevel.Load()) > zerolog.InfoLevel {
return len(p), nil
}
return w.inner.Write(p)
}
func (w *levelFilterWriter) WriteLevel(l zerolog.Level, p []byte) (int, error) {
if l < zerolog.Level(w.minLevel.Load()) {
return len(p), nil
}
if lw, ok := w.inner.(zerolog.LevelWriter); ok {
return lw.WriteLevel(l, p)
}
return w.inner.Write(p)
}
const timeFormatMS = "2006-01-02 15:04:05.000"
// NewLogger 创建日志实例
// logLevel: 0=仅 error>=1=全部
// logFile: 文件路径模板(如 "logs/20060102.txt"),空则不写文件
@@ -66,7 +113,7 @@ type Logger struct {
func NewLogger(logLevel int, logFile string, maxErrors int) *Logger {
saveOrigStderr()
zerolog.CallerMarshalFunc = callerMarshalFunc
zerolog.TimeFieldFormat = "2006-01-02 15:04:05"
zerolog.TimeFieldFormat = timeFormatMS
var level zerolog.Level
if logLevel == 0 {
@@ -77,7 +124,7 @@ func NewLogger(logLevel int, logFile string, maxErrors int) *Logger {
consoleWriter := zerolog.ConsoleWriter{
Out: os.Stderr,
TimeFormat: "2006-01-02 15:04:05",
TimeFormat: timeFormatMS,
NoColor: false,
FormatLevel: formatLevelColor,
FormatCaller: func(i interface{}) string {
@@ -87,9 +134,10 @@ func NewLogger(logLevel int, logFile string, maxErrors int) *Logger {
return fmt.Sprintf("[%s]", i)
},
}
console := newLevelFilterWriter(consoleWriter, zerolog.DebugLevel)
var writers []io.Writer
writers = append(writers, consoleWriter)
writers = append(writers, console)
if logFile != "" {
fw := &TemplateFileWriter{pathTemplate: logFile}
@@ -105,13 +153,14 @@ func NewLogger(logLevel int, logFile string, maxErrors int) *Logger {
Logger()
l := &Logger{
zl: zl,
mw: mw,
logLevel: logLevel,
maxErrors: maxErrors,
zl: zl,
mw: mw,
console: console,
logLevel: logLevel,
store: &errorStore{maxErrors: maxErrors},
}
if maxErrors > 0 {
l.errors = make([]ErrorRecord, maxErrors)
l.store.errors = make([]ErrorRecord, maxErrors)
}
return l
@@ -121,7 +170,7 @@ func NewLogger(logLevel int, logFile string, maxErrors int) *Logger {
func NewLoggerNoCaller(logLevel int, logFile string, maxErrors int) *Logger {
saveOrigStderr()
zerolog.CallerMarshalFunc = callerMarshalFunc
zerolog.TimeFieldFormat = "2006-01-02 15:04:05"
zerolog.TimeFieldFormat = timeFormatMS
var level zerolog.Level
if logLevel == 0 {
@@ -132,13 +181,14 @@ func NewLoggerNoCaller(logLevel int, logFile string, maxErrors int) *Logger {
consoleWriter := zerolog.ConsoleWriter{
Out: os.Stderr,
TimeFormat: "2006-01-02 15:04:05",
TimeFormat: timeFormatMS,
NoColor: false,
FormatLevel: formatLevelColor,
}
console := newLevelFilterWriter(consoleWriter, zerolog.DebugLevel)
var writers []io.Writer
writers = append(writers, consoleWriter)
writers = append(writers, console)
if logFile != "" {
fw := &TemplateFileWriter{pathTemplate: logFile}
@@ -153,18 +203,40 @@ func NewLoggerNoCaller(logLevel int, logFile string, maxErrors int) *Logger {
Logger()
l := &Logger{
zl: zl,
mw: mw,
logLevel: logLevel,
maxErrors: maxErrors,
zl: zl,
mw: mw,
console: console,
logLevel: logLevel,
store: &errorStore{maxErrors: maxErrors},
}
if maxErrors > 0 {
l.errors = make([]ErrorRecord, maxErrors)
l.store.errors = make([]ErrorRecord, maxErrors)
}
return l
}
// WithFields 派生子 Logger,附加成对 key,value 字段;mw/错误历史与父共享,父不受污染。
func (l *Logger) WithFields(kvs ...string) *Logger {
if l == nil {
return nil
}
ctx := l.zl.With()
for i := 0; i+1 < len(kvs); i += 2 {
if kvs[i] == "" {
continue
}
ctx = ctx.Str(kvs[i], kvs[i+1])
}
return &Logger{
zl: ctx.Logger(),
mw: l.mw,
console: l.console,
logLevel: l.logLevel,
store: l.store,
}
}
// NewTestLogger 创建用于测试的静默 Logger(输出到 io.Discard
func NewTestLogger() *Logger {
zl := zerolog.New(io.Discard).Level(zerolog.Disabled)
@@ -180,7 +252,7 @@ func (l *Logger) SetOutput(w io.Writer) {
}
consoleWriter := zerolog.ConsoleWriter{
Out: w,
TimeFormat: "2006-01-02 15:04:05",
TimeFormat: timeFormatMS,
NoColor: true,
FormatLevel: formatLevelPlain,
FormatCaller: func(i interface{}) string {
@@ -295,14 +367,14 @@ func (l *Logger) Warn(args ...interface{}) *zerolog.Event {
func (l *Logger) Error(args ...interface{}) *zerolog.Event {
if len(args) > 0 {
msg := fmt.Sprint(args...)
if l.maxErrors > 0 {
if l.store != nil && l.store.maxErrors > 0 {
l.recordError(msg, nil)
}
l.zl.Error().Msg(msg)
nop := zerolog.Nop()
return nop.Error()
}
if l.maxErrors > 0 {
if l.store != nil && l.store.maxErrors > 0 {
l.recordError("", nil)
}
return l.zl.Error()
@@ -310,19 +382,20 @@ func (l *Logger) Error(args ...interface{}) *zerolog.Event {
// RecordError 手动记录一条错误到历史(用于需要指定详情的场景)
func (l *Logger) RecordError(err error, msg string, caller string) {
if l.maxErrors <= 0 {
if l == nil || l.store == nil || l.store.maxErrors <= 0 {
return
}
l.errorsMu.Lock()
defer l.errorsMu.Unlock()
l.errors[l.errIdx] = ErrorRecord{
s := l.store
s.mu.Lock()
defer s.mu.Unlock()
s.errors[s.errIdx] = ErrorRecord{
Err: err,
Msg: msg,
Time: time.Now(),
Caller: caller,
}
l.errIdx = (l.errIdx + 1) % l.maxErrors
l.errCount++
s.errIdx = (s.errIdx + 1) % s.maxErrors
s.errCount++
}
func (l *Logger) recordError(msg string, err error) {
@@ -336,15 +409,16 @@ func (l *Logger) recordError(msg string, err error) {
// GetRecentErrors 获取最近 N 条错误(不传则返回全部已存储的)
func (l *Logger) GetRecentErrors(n ...int) []ErrorRecord {
if l.maxErrors <= 0 {
if l == nil || l.store == nil || l.store.maxErrors <= 0 {
return nil
}
l.errorsMu.RLock()
defer l.errorsMu.RUnlock()
s := l.store
s.mu.RLock()
defer s.mu.RUnlock()
total := l.errCount
if total > l.maxErrors {
total = l.maxErrors
total := s.errCount
if total > s.maxErrors {
total = s.maxErrors
}
if total == 0 {
return nil
@@ -356,13 +430,12 @@ func (l *Logger) GetRecentErrors(n ...int) []ErrorRecord {
}
result := make([]ErrorRecord, 0, want)
// 从最新的往前读
for i := 0; i < want; i++ {
idx := (l.errIdx - 1 - i + l.maxErrors) % l.maxErrors
if l.errors[idx].Time.IsZero() {
idx := (s.errIdx - 1 - i + s.maxErrors) % s.maxErrors
if s.errors[idx].Time.IsZero() {
break
}
result = append(result, l.errors[idx])
result = append(result, s.errors[idx])
}
return result
}
@@ -379,7 +452,7 @@ func (l *Logger) Warnf(format string, v ...interface{}) {
l.zl.Warn().Msgf(format, v...)
}
func (l *Logger) Errorf(format string, v ...interface{}) {
if l.maxErrors > 0 {
if l.store != nil && l.store.maxErrors > 0 {
l.recordError(fmt.Sprintf(format, v...), nil)
}
l.zl.Error().Msgf(format, v...)
@@ -560,6 +633,8 @@ type SeqWriter struct {
queue chan []byte
dropped int64
client *http.Client
closed int32
done chan struct{}
}
func newSeqWriter(seqUrl, apiKey, instance string) *SeqWriter {
@@ -569,6 +644,7 @@ func newSeqWriter(seqUrl, apiKey, instance string) *SeqWriter {
instance: instance,
queue: make(chan []byte, 10000),
client: &http.Client{Timeout: 5 * time.Second},
done: make(chan struct{}),
}
go w.run()
return w
@@ -576,6 +652,9 @@ func newSeqWriter(seqUrl, apiKey, instance string) *SeqWriter {
// Write 将字节塞入 channel 即返回,由后台 goroutine 消费发送
func (w *SeqWriter) Write(p []byte) (int, error) {
if atomic.LoadInt32(&w.closed) != 0 {
return len(p), nil
}
clef := w.toClef(p)
select {
case w.queue <- clef:
@@ -588,15 +667,31 @@ func (w *SeqWriter) Write(p []byte) (int, error) {
return len(p), nil
}
// Close 停止接收新条目,并在约 2s 内冲刷残留批次(幂等)
func (w *SeqWriter) Close() {
if w == nil {
return
}
if !atomic.CompareAndSwapInt32(&w.closed, 0, 1) {
return
}
close(w.queue)
select {
case <-w.done:
case <-time.After(2 * time.Second):
seqWriteStderr("[seq] Close timeout, some events may be dropped\n")
}
}
// toClef 将 zerolog JSON 字段映射到 Seq CLEF 格式
// zerolog: time/level/message(msg) → CLEF: @t/@l/@mt
// zerolog: time/level/message(msg) → CLEF: @t/@l/@m
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))
return []byte(fmt.Sprintf(`{"@t":%q,"@l":"Information","@m":%q,"source":"stdout","instance":%q}`,
time.Now().Format("2006-01-02T15:04:05.000Z07:00"), strings.TrimSpace(string(p)), w.instance))
}
clef := make(map[string]interface{}, len(m)+3)
@@ -604,16 +699,18 @@ func (w *SeqWriter) toClef(p []byte) []byte {
clef[k] = v
}
// time → @tzerolog 格式 "2006-01-02 15:04:05" → ISO 8601
// time → @t优先毫秒格式,回退秒级
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)
if pt, err := time.ParseInLocation(timeFormatMS, t, time.Local); err == nil {
clef["@t"] = pt.Format("2006-01-02T15:04:05.000Z07:00")
} else if pt, err := time.ParseInLocation("2006-01-02 15:04:05", t, time.Local); err == nil {
clef["@t"] = pt.Format("2006-01-02T15:04:05.000Z07:00")
} else {
clef["@t"] = t
}
} else {
clef["@t"] = time.Now().Format(time.RFC3339)
clef["@t"] = time.Now().Format("2006-01-02T15:04:05.000Z07:00")
}
// level → @l
@@ -635,13 +732,13 @@ func (w *SeqWriter) toClef(p []byte) []byte {
}
}
// message/msg → @mt
// message/msg → @m(不用 @mt,避免 Seq 把 {xxx} 当消息模板)
if msg, ok := m["message"].(string); ok {
delete(clef, "message")
clef["@mt"] = msg
clef["@m"] = msg
} else if msg, ok := m["msg"].(string); ok {
delete(clef, "msg")
clef["@mt"] = msg
clef["@m"] = msg
}
if w.instance != "" {
@@ -652,14 +749,21 @@ func (w *SeqWriter) toClef(p []byte) []byte {
return b
}
// run 后台 goroutine:每 100 条或 500ms 批量 POST 到 Seq
// run 后台 goroutine:每 100 条或 500ms 批量 POST 到 Seqqueue 关闭后冲刷并退出
func (w *SeqWriter) run() {
defer close(w.done)
batch := make([][]byte, 0, 100)
ticker := time.NewTicker(500 * time.Millisecond)
defer ticker.Stop()
for {
select {
case entry := <-w.queue:
case entry, ok := <-w.queue:
if !ok {
if len(batch) > 0 {
w.flush(batch)
}
return
}
batch = append(batch, entry)
if len(batch) >= 100 {
w.flush(batch)
@@ -675,6 +779,15 @@ func (w *SeqWriter) run() {
}
func (w *SeqWriter) flush(batch [][]byte) {
if err := w.postBatch(batch); err != nil {
time.Sleep(1 * time.Second)
if err2 := w.postBatch(batch); err2 != nil {
seqWriteStderr("[seq] send failed after retry: %v (dropped=%d)\n", err2, atomic.LoadInt64(&w.dropped))
}
}
}
func (w *SeqWriter) postBatch(batch [][]byte) error {
var buf bytes.Buffer
for _, entry := range batch {
buf.Write(entry)
@@ -682,8 +795,7 @@ func (w *SeqWriter) flush(batch [][]byte) {
}
req, err := http.NewRequest("POST", w.seqUrl+"/api/events/raw?clef", &buf)
if err != nil {
seqWriteStderr("[seq] build request error: %v\n", err)
return
return err
}
req.Header.Set("Content-Type", "application/vnd.serilog.clef")
if w.apiKey != "" {
@@ -691,14 +803,14 @@ func (w *SeqWriter) flush(batch [][]byte) {
}
resp, err := w.client.Do(req)
if err != nil {
seqWriteStderr("[seq] send error: %v (dropped=%d)\n", err, atomic.LoadInt64(&w.dropped))
return
return err
}
defer resp.Body.Close()
if resp.StatusCode >= 400 {
body, _ := io.ReadAll(resp.Body)
seqWriteStderr("[seq] server returned %d: %s\n", resp.StatusCode, string(body))
return fmt.Errorf("status %d: %s", resp.StatusCode, string(body))
}
return nil
}
// SetSeqWriter 将 Seq HTTP 推送器附加到当前 Logger(不影响已有的控制台/文件输出)
@@ -706,11 +818,30 @@ func (w *SeqWriter) flush(batch [][]byte) {
// apiKey: Seq API Key,免费单用户版留空
// instance: 实例标识,建议用 config 的 port 字段区分同机多进程,如 "8085"
func (l *Logger) SetSeqWriter(seqUrl, apiKey, instance string) {
if seqUrl == "" || l.mw == nil {
if l == nil || seqUrl == "" || l.mw == nil {
return
}
sw := newSeqWriter(seqUrl, apiKey, instance)
l.mw.add(sw)
l.seqsMu.Lock()
l.seqs = append(l.seqs, sw)
l.seqsMu.Unlock()
// 挂上 Seq 后控制台只打 Warn+Info/Debug/SQL/访问日志仍进 Seq 与文件
l.console.setMinLevel(zerolog.WarnLevel)
}
// CloseSeq 冲刷并关闭本 Logger 挂载的全部 SeqWriter(幂等、nil 安全)
func (l *Logger) CloseSeq() {
if l == nil {
return
}
l.seqsMu.Lock()
seqs := l.seqs
l.seqs = nil
l.seqsMu.Unlock()
for _, sw := range seqs {
sw.Close()
}
}
// isInfrastructureFile 判断文件是否为基础设施(日志/DB/缓存/通用工具)管线,