From 831a7f017ad962f9edb675e1f80c64a9a8d26b31 Mon Sep 17 00:00:00 2001 From: hoteas <925970985@qq.com> Date: Thu, 23 Jul 2026 05:26:44 +0800 Subject: [PATCH] =?UTF-8?q?feat(log):=20Seq=20=E4=BC=9A=E8=AF=9D=E8=BF=BD?= =?UTF-8?q?=E8=B8=AA=E3=80=81=E6=8E=A7=E5=88=B6=E5=8F=B0=E9=99=8D=E5=99=AA?= =?UTF-8?q?=E4=B8=8E=E6=8E=A8=E9=80=81=E8=B4=A8=E9=87=8F=E5=A2=9E=E5=BC=BA?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 请求日志自动带 sid/request_id,支持 LogBind;挂 Seq 后控制台仅 Warn+,Seq 仍全量;并修毫秒时间戳、@m、失败重试与停机 flush。 Co-authored-by: Cursor --- application.go | 82 ++++++++++---- context.go | 48 +++++--- docs/Seq_日志集成.md | 184 +++++++++++++++---------------- log/capture_test.go | 180 +++++++++++++++++++++++++++++- log/logger.go | 255 ++++++++++++++++++++++++++++++++----------- var.go | 2 +- 6 files changed, 554 insertions(+), 197 deletions(-) diff --git a/application.go b/application.go index 59589da..d137a5f 100644 --- a/application.go +++ b/application.go @@ -2,7 +2,9 @@ package hotime import ( "context" + "crypto/rand" "database/sql" + "encoding/hex" "fmt" "io/ioutil" stdlog "log" @@ -25,8 +27,8 @@ import ( . "code.hoteas.com/golang/hotime/db" "code.hoteas.com/golang/hotime/log" mysql "github.com/go-sql-driver/mysql" - logrus "github.com/sirupsen/logrus" "github.com/rs/zerolog" + logrus "github.com/sirupsen/logrus" ) type Application struct { @@ -97,10 +99,45 @@ func (that *Application) initiateGracefulShutdown(drain, shutdown time.Duration) } else { that.Log.Infof("服务已安全关闭") } + // 冲刷 Seq 队列,避免停机前后关键日志丢失 + that.Log.CloseSeq() + if that.WebConnectLog != nil { + that.WebConnectLog.CloseSeq() + } os.Exit(0) }) } +// clientIP 取客户端真实 IP:X-Real-IP → X-Forwarded-For 第一段 → RemoteAddr +func clientIP(req *http.Request) string { + if req == nil { + return "" + } + if ip := strings.TrimSpace(req.Header.Get("X-Real-IP")); ip != "" { + return ip + } + if xff := req.Header.Get("X-Forwarded-For"); xff != "" { + if i := strings.Index(xff, ","); i >= 0 { + return strings.TrimSpace(xff[:i]) + } + return strings.TrimSpace(xff) + } + host, _, err := net.SplitHostPort(req.RemoteAddr) + if err == nil && host != "" { + return host + } + return req.RemoteAddr +} + +// shortID 生成 n 字节随机 hex(n=6 → 12 位) +func shortID(n int) string { + b := make([]byte, n) + if _, err := rand.Read(b); err != nil { + return Md5(strconv.Itoa(Rand(10)))[:n*2] + } + return hex.EncodeToString(b) +} + // Run 启动实例 func (that *Application) Run(router Router) { //如果没有设置配置自动生成配置 @@ -438,10 +475,22 @@ func (that *Application) handler(w http.ResponseWriter, req *http.Request) { if err != nil { unescapeUrl = req.RequestURI } + + // 请求级追踪:sid=sessionId 前 12 位(脱敏),request_id 回写响应头 + requestId := shortID(6) + sid := sessionId + if len(sid) > 12 { + sid = sid[:12] + } + w.Header().Set("X-Request-Id", requestId) + reqLog := that.Log.WithFields("sid", sid, "request_id", requestId) + dbCopy := that.Db + dbCopy.Log = reqLog + //访问实例 context := Context{SessionIns: SessionIns{SessionId: sessionId, HoTimeCache: that.HoTimeCache}, - Resp: w, Req: req, Application: that, RouterString: s, Config: that.Config, Db: &that.Db, - HandlerStr: unescapeUrl} + Resp: w, Req: req, Application: that, RouterString: s, Config: that.Config, Db: &dbCopy, + HandlerStr: unescapeUrl, Logger: reqLog} //header默认设置 header := w.Header() header.Set("Content-Type", "text/html; charset=utf-8") @@ -455,26 +504,17 @@ func (that *Application) handler(w http.ResponseWriter, req *http.Request) { defer func() { //是否展示日志 if that.WebConnectLog != nil { - - //负载均衡优化 - ipStr := "" - if req.Header.Get("X-Forwarded-For") != "" { - ipStr = req.Header.Get("X-Forwarded-For") - } else if req.Header.Get("X-Real-IP") != "" { - ipStr = req.Header.Get("X-Real-IP") - } - //负载均衡优化 - if ipStr == "" { - //RemoteAddr := that.Req.RemoteAddr - ipStr = Substr(context.Req.RemoteAddr, 0, strings.Index(context.Req.RemoteAddr, ":")) - } - - that.WebConnectLog.Info(). - Str("ip", ipStr). + evt := that.WebConnectLog.Info(). + Str("ip", clientIP(req)). Str("method", context.Req.Method). + Str("sid", sid). + Str("request_id", requestId). Float64("cost_ms", ObjToFloat64(time.Now().UnixNano()-nowUnixTime.UnixNano())/1000000.00). - Float64("size_kb", ObjToFloat64(context.DataSize)/1000.00). - Msg(context.HandlerStr) + Float64("size_kb", ObjToFloat64(context.DataSize)/1000.00) + if country := req.Header.Get("EO-Client-IPCountry"); country != "" { + evt = evt.Str("ip_country", country) + } + evt.Msg(context.HandlerStr) } }() diff --git a/context.go b/context.go index 368b33b..69c3630 100644 --- a/context.go +++ b/context.go @@ -6,19 +6,20 @@ import ( "io" "mime/multipart" "net/http" - "strings" "sync" "time" . "code.hoteas.com/golang/hotime/common" . "code.hoteas.com/golang/hotime/db" + htlog "code.hoteas.com/golang/hotime/log" ) type Context struct { *Application Resp http.ResponseWriter Req *http.Request - Log Map //日志有则创建 + Log Map //业务 logs 表字段(有则写入),与 Logger 不同 + Logger *htlog.Logger // 请求级日志(带 sid/request_id),业务与 SQL 共用 RouterString []string Config Map Db *HoTimeDB @@ -60,7 +61,7 @@ func (that *Context) Display(statu int, data interface{}) { //兼容android等需要json转对象的服务 resp["error"] = temp - that.Application.Log.Warn().Int("status", statu).Msg(resp.ToJsonString()) + that.reqLogger().Warn().Int("status", statu).Msg(resp.ToJsonString()) } else { resp["result"] = data @@ -71,6 +72,33 @@ func (that *Context) Display(statu int, data interface{}) { //that.Data=d; } +// reqLogger 返回请求级 Logger,nil 时回退 Application.Log +func (that *Context) reqLogger() *htlog.Logger { + if that.Logger != nil { + return that.Logger + } + if that.Application != nil { + return that.Application.Log + } + return nil +} + +// LogBind 为本请求追加自定义日志字段(如 user_id),后续业务日志与 SQL 日志均携带。 +// 供 main.go 的 connectListener 调用;nil 安全。 +func (that *Context) LogBind(key string, value interface{}) { + if that == nil || key == "" { + return + } + base := that.reqLogger() + if base == nil { + return + } + that.Logger = base.WithFields(key, ObjToStr(value)) + if that.Db != nil { + that.Db.Log = that.Logger + } +} + func (that *Context) View() { if that.RespFunc != nil { that.RespFunc() @@ -87,19 +115,7 @@ func (that *Context) View() { if that.Session("user_id").Data != nil { that.Log["user_id"] = that.Session("user_id").ToCeilInt() } - //负载均衡优化 - ipStr := "" - if that.Req.Header.Get("X-Forwarded-For") != "" { - ipStr = that.Req.Header.Get("X-Forwarded-For") - } else if that.Req.Header.Get("X-Real-IP") != "" { - ipStr = that.Req.Header.Get("X-Real-IP") - } - //负载均衡优化 - if ipStr == "" { - //RemoteAddr := that.Req.RemoteAddr - ipStr = Substr(that.Req.RemoteAddr, 0, strings.Index(that.Req.RemoteAddr, ":")) - } - that.Log["ip"] = ipStr + that.Log["ip"] = clientIP(that.Req) that.Db.Insert("logs", that.Log) } diff --git a/docs/Seq_日志集成.md b/docs/Seq_日志集成.md index e8a1278..1195889 100644 --- a/docs/Seq_日志集成.md +++ b/docs/Seq_日志集成.md @@ -19,7 +19,61 @@ HoTime 框架内置 Seq 日志推送支持。通过在 `config.json` 填写 `seq | `seqUrl` | 空(不激活) | Seq 服务地址,空值时功能静默不生效 | | `seqApiKey` | 空 | API Key,免费单用户版留空 | -`instance` 字段由框架自动拼接为 `ip:port`(如 `192.168.1.10:8085`),无需手动填写,同机多进程靠端口区分,跨服务器靠 IP 区分,支持未来集群扩展。 +`instance` 字段由框架自动拼接为 `ip:port`(如 `192.168.1.10:8085`),无需手动填写。 + +激活后每个 HTTP 请求会自动: + +- 生成 `request_id`(12 位 hex),回写响应头 `X-Request-Id` +- 将 `sessionId` 前 12 位作为 `sid` 写入请求级日志(脱敏,避免把登录凭据写进 Seq) +- 业务日志、SQL 日志、访问日志均携带 `sid` / `request_id` +- **控制台自动降噪为 Warn+**(Info/Debug/SQL/访问日志仍进 Seq 与文件,零额外配置) + +--- + +## 会话追踪与 LogBind + +框架在 `handler` 入口派生请求级 Logger,并浅拷贝 `Db` 将其 `Log` 指向同一 Logger,因此 **SQL 日志自动带会话字段**。字段会出现在同条日志的控制台/文件/Seq 出口上。 + +业务在 `SetConnectListener` 鉴权通过后绑定(xbc `main.go`): + +```go +// app 鉴权通过后 +context.LogBind("user_id", context.Session("user_id").ToCeilInt64()) + +// admin 鉴权通过后(或 session 已有 admin_id) +context.LogBind("admin_id", context.Session("admin_id").ToCeilInt64()) +``` + +`LogBind` 后,本请求后续业务日志与 SQL 日志都会带上该字段。 + +**不带会话字段的边界:** + +- `fmt.Println` / stdout 捕获、panic、MySQL driver、定时任务等无请求上下文的日志 +- 直接写 `that.Application.Log` 的旧代码(应改用 `that.Logger` 或 `that.LogBind` 后的请求级 Logger) + +--- + +## 控制台降噪(配了 seqUrl) + +| 出口 | 行为 | +|---|---| +| 控制台 | 仅 Warn / Error(含 `Display` 非 0 的 Warn) | +| Seq | 按 `logLevel` 全量(Info/Debug/SQL/访问日志等) | +| 本地文件 | 与原先一致,不受控制台过滤影响 | + +未配置 `seqUrl` 时控制台仍按 `logLevel` 全打。 + +--- + +## 客户端 IP 与地域 + +零配置,优先级: + +1. `X-Real-IP`(EdgeOne / 反代透传的真实 IP) +2. `X-Forwarded-For` **第一段** +3. `RemoteAddr` + +请求头有 `EO-Client-IPCountry` 时,访问日志追加 `ip_country`(两位国家码);没有则不记。 --- @@ -27,26 +81,24 @@ HoTime 框架内置 Seq 日志推送支持。通过在 `config.json` 填写 `seq ``` 业务代码 - │ l.Info().Msg("...") ← HoTime Logger 正常调用路径 - │ fmt.Println("...") ← 被 redirectStdout 捕获后转入同一路径 + │ that.Logger.Info().Msg("...") ← 请求级(含 sid/request_id) + │ Db.Query → SQL 日志 ← 同一请求级 Logger + │ fmt.Println("...") ← 捕获后无 sid ↓ multiWriter(hotimev1.5/log/logger.go) - ├─ ConsoleWriter → 彩色终端输出(不变) - ├─ FileWriter → 本地日志文件(按需,logFile 配置) - └─ SeqWriter + ├─ Console(挂 Seq 后 Warn+)→ 终端 + ├─ FileWriter → 本地文件(按需,全量) + └─ SeqWriter → Seq(全量) │ Write() 只做 channel <- bytes,O(1) 非阻塞 ↓ channel(容量 10000) ↓ 后台 goroutine 批量打包(100 条 或 500ms) - ↓ HTTP POST + ↓ HTTP POST(失败重试 1 次) Seq 服务(CLEF 格式) ``` -**关键特性:** -- `SeqWriter.Write()` 仅向 channel 投递字节即返回,**绝不阻塞** web 请求处理 goroutine -- channel 满时(Seq 宕机/网络故障)新日志被丢弃并计数,主服务完全不受影响 -- HTTP POST 设 5s 超时,失败仅打印到 stderr +优雅停机时调用 `CloseSeq()` 冲刷残留批次,避免停机前后日志丢失。 --- @@ -54,82 +106,16 @@ multiWriter(hotimev1.5/log/logger.go) | zerolog 字段 | Seq CLEF 字段 | 说明 | |---|---|---| -| `time` | `@t` | 时间,自动转 ISO 8601 格式 | -| `level` | `@l` | 级别,映射为 Debug/Information/Warning/Error/Fatal | -| `message` / `msg` | `@mt` | 消息正文 | -| `caller` | `caller` | 调用位置,原样保留 | -| 其余自定义字段 | 原字段名 | 原样保留,可在 Seq 中直接查询 | -| — | `instance` | 框架自动注入,值为 `port` 配置(如 `"8085"`) | -| — | `source` | fmt.Println 等捕获的输出标记为 `stdout` | - ---- - -## 单机多进程实例区分 - -框架启动时自动获取本机出口 IP,拼接为 `ip:port` 格式作为 `instance`: - -``` -单机多进程: - 192.168.1.10:8085 ─┐ - 192.168.1.10:8086 ─┼─ HTTP CLEF ──→ Seq - 192.168.1.10:8087 ─┘ - -多服务器集群: - 192.168.1.10:8085 ─┐ - 192.168.1.11:8085 ─┼─ HTTP CLEF ──→ Seq(中央日志服务器) - 192.168.1.12:8085 ─┘ -``` - -Seq 中按实例筛选: -- 单台机器所有进程:`instance like '192.168.1.10%'` -- 精确到某个进程:`instance = '192.168.1.10:8085'` - ---- - -## stdout / stderr 全量捕获 - -`SetConfig()` 中自动 `redirectStdout` + `redirectStderr`,并桥接 **logrus**(微信 SDK)到同一管道。 - -| 写法 | Seq 字段 | 说明 | -|------|----------|------| -| `fmt.Println` / `log.Println` | `source=stdout` | **不受 `logLevel` 影响**,`logLevel=0` 也会进 Seq | -| 写 `os.Stderr` 的包 | `source=stderr` | 同上,Error 级别 | -| `logrus.Info`(wechat) | `source=stdout` | redirect 后 `logrus.SetOutput` 桥接 | -| `that.Log.Info/Error...` | 结构化字段 | 经 `multiWriter` → SeqWriter | -| MySQL driver 内部错误 | `source=mysql-driver` | `SetLogger` 适配器 | -| 框架 `recover` 到的 panic | `source=panic` + `stack` | 代码层原因与调用栈 | - -单行日志上限约 **10MB**(适配 `GetReqMap` 打整包 body);超长会截断并标注 `...(truncated)`。 - -**已知不进 Seq(文档边界):** - -- 未 `recover`、进程直接崩溃的 runtime 栈(写 fd2,未做 Dup2) -- 达梦驱动自有文件日志(vendor 独立写盘) -- `seqUrl` 挂上之前的极早期引导日志 - ---- - -## 捕获矩阵速查 - -| 来源 | 进 Seq? | 检索示例 | -|------|----------|----------| -| `that.Log.*` | 是 | `@mt like '%关键词%'` | -| `log.Println` / `fmt.Println` | 是 | `source = 'stdout'` | -| logrus / 标准 `log` | 是 | `source = 'stdout'` | -| stderr 重定向 | 是 | `source = 'stderr'` | -| recover panic | 是 | `source = 'panic'` | -| MySQL driver | 是 | `source = 'mysql-driver'` | -| 队列满丢弃 | 否(计数) | 终端可见 `[seq] queue full` | - ---- - -## Seq 安装 - -Seq 提供 Windows MSI 安装包和 Docker 镜像,单机免费,无外部数据库依赖: - -- Windows:[https://datalust.co/download/seq](https://datalust.co/download/seq),安装后自动注册为 Windows 服务 -- Docker:`docker run -d --restart always --name seq -p 5341:80 -e ACCEPT_EULA=Y datalust/seq` -- 访问 `http://localhost:5341` 使用 Web UI +| `time` | `@t` | 毫秒精度 ISO 8601 | +| `level` | `@l` | Debug/Information/Warning/Error/Fatal | +| `message` / `msg` | `@m` | 消息正文(不用 `@mt`,避免 `{xxx}` 被当模板) | +| `caller` | `caller` | 调用位置 | +| `sid` | `sid` | sessionId 前 12 位 | +| `request_id` | `request_id` | 单次请求 id | +| `ip_country` | `ip_country` | 有 `EO-Client-IPCountry` 时 | +| 其余自定义字段 | 原字段名 | `LogBind` 追加的字段原样保留 | +| — | `instance` | `ip:port` | +| — | `source` | stdout / stderr / panic 等 | --- @@ -137,15 +123,23 @@ Seq 提供 Windows MSI 安装包和 Docker 镜像,单机免费,无外部数 | 目标 | 查询语句 | |---|---| -| 关键词搜索 | 直接输入,如 `支付失败` | +| 按会话追踪 | `sid = 'abc123def456'` | +| 按单次请求 | `request_id = 'fedcba987654'` | +| 关键词 | 直接输入,如 `支付失败` | | 日志级别 | `@l = 'Error'` | -| 特定实例 | `instance = '8085'` | -| stdout 来源 | `source = 'stdout'` | -| panic 恢复 | `source = 'panic'` | -| 请求体调试 | `请求参数GetReqMap` 或 `source = 'stdout'` | -| 调用位置 | `caller like '%order.go%'` | -| 组合查询 | `@l = 'Error' and instance = '8086' and @mt like '%超时%'` | -| 日期范围 | 右上角时间选择器,支持精确到秒 | +| 特定实例 | `instance = '192.168.1.10:8085'` | +| 地域 | `ip_country = 'CN'` | +| stdout | `source = 'stdout'` | +| panic | `source = 'panic'` | +| 组合 | `@l = 'Error' and sid = 'abc123def456'` | + +--- + +## Seq 安装 + +- Windows:[https://datalust.co/download/seq](https://datalust.co/download/seq) +- Docker:`docker run -d --restart always --name seq -p 5341:80 -e ACCEPT_EULA=Y datalust/seq` +- 访问 `http://localhost:5341` --- diff --git a/log/capture_test.go b/log/capture_test.go index 04f0e59..8e75ae0 100644 --- a/log/capture_test.go +++ b/log/capture_test.go @@ -96,8 +96,8 @@ func TestToClef_PlainText(t *testing.T) { if err := json.Unmarshal(clef, &m); err != nil { t.Fatal(err) } - if m["@mt"] != "plain log line" { - t.Fatalf("@mt=%v", m["@mt"]) + if m["@m"] != "plain log line" { + t.Fatalf("@m=%v", m["@m"]) } if m["source"] != "stdout" { t.Fatalf("source=%v", m["source"]) @@ -107,6 +107,182 @@ func TestToClef_PlainText(t *testing.T) { } } +func TestWithFields_ChildHasFieldsParentClean(t *testing.T) { + parent, buf := newLoggerWithCapture(1) + child := parent.WithFields("sid", "abc123def456", "request_id", "fedcba987654") + child.Info().Msg("child-msg") + m := buf.lastJSON() + if m == nil { + t.Fatal("no output") + } + if m["sid"] != "abc123def456" || m["request_id"] != "fedcba987654" { + t.Fatalf("child fields missing: %v", m) + } + parent.Info().Msg("parent-msg") + m2 := buf.lastJSON() + if m2["sid"] != nil || m2["request_id"] != nil { + t.Fatalf("parent polluted: %v", m2) + } + // 错误历史父子共享 + parent2 := NewLogger(1, "", 10) + child2 := parent2.WithFields("sid", "x") + child2.Error("shared-err") + if n := len(parent2.GetRecentErrors()); n != 1 { + t.Fatalf("shared error store want 1 got %d", n) + } +} + +func TestToClef_MillisecondAndAtM(t *testing.T) { + sw := &SeqWriter{instance: "test:1"} + in := []byte(`{"level":"info","time":"2026-07-23 04:50:01.123","message":"hello {name}"}`) + clef := sw.toClef(in) + var m map[string]interface{} + if err := json.Unmarshal(clef, &m); err != nil { + t.Fatal(err) + } + if m["@m"] != "hello {name}" { + t.Fatalf("@m=%v", m["@m"]) + } + if m["@mt"] != nil { + t.Fatalf("should not have @mt: %v", m["@mt"]) + } + at, ok := m["@t"].(string) + if !ok || !strings.Contains(at, ".123") { + t.Fatalf("@t missing ms: %v", m["@t"]) + } +} + +func TestSeqWriter_RetryOnce(t *testing.T) { + var mu sync.Mutex + calls := 0 + var gotBody string + srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + body, _ := io.ReadAll(r.Body) + mu.Lock() + calls++ + n := calls + if n == 1 { + mu.Unlock() + w.WriteHeader(http.StatusInternalServerError) + return + } + gotBody = string(body) + mu.Unlock() + w.WriteHeader(http.StatusCreated) + })) + defer srv.Close() + + l := NewLogger(1, "", 0) + l.SetSeqWriter(srv.URL, "", "test:retry") + l.Info().Msg("retry-me") + + deadline := time.Now().Add(4 * time.Second) + for time.Now().Before(deadline) { + mu.Lock() + ok := calls >= 2 && strings.Contains(gotBody, "retry-me") + mu.Unlock() + if ok { + return + } + time.Sleep(50 * time.Millisecond) + } + mu.Lock() + defer mu.Unlock() + t.Fatalf("retry failed: calls=%d body=%q", calls, gotBody) +} + +func TestLevelFilterWriter_WarnOnly(t *testing.T) { + buf := &captureBuf{} + w := newLevelFilterWriter(buf, zerolog.WarnLevel) + mw := &multiWriter{writers: []io.Writer{w}} + zl := zerolog.New(mw).Level(zerolog.DebugLevel) + zl.Info().Msg("info-hidden") + zl.Warn().Msg("warn-shown") + buf.mu.Lock() + n := len(buf.lines) + buf.mu.Unlock() + if n != 1 { + t.Fatalf("want 1 console line got %d: %v", n, buf.lines) + } + if !strings.Contains(buf.lines[0], "warn-shown") { + t.Fatalf("unexpected: %v", buf.lines[0]) + } +} + +func TestSetSeqWriter_QuietsConsoleKeepsSeqFull(t *testing.T) { + var mu sync.Mutex + var gotBody string + srv := startFakeSeqServer(t, func(body string) { + mu.Lock() + gotBody += body + mu.Unlock() + }) + defer srv.Close() + + l := NewLogger(1, "", 0) + consoleBuf := &captureBuf{} + l.console.inner = consoleBuf + l.SetSeqWriter(srv.URL, "", "test:quiet") + l.Info().Msg("info-to-seq-only") + l.Warn().Msg("warn-both") + + deadline := time.Now().Add(3 * time.Second) + for time.Now().Before(deadline) { + mu.Lock() + ok := strings.Contains(gotBody, "info-to-seq-only") && strings.Contains(gotBody, "warn-both") + mu.Unlock() + if ok { + break + } + time.Sleep(50 * time.Millisecond) + } + mu.Lock() + body := gotBody + mu.Unlock() + if !strings.Contains(body, "info-to-seq-only") || !strings.Contains(body, "warn-both") { + t.Fatalf("Seq incomplete: %q", body) + } + consoleBuf.mu.Lock() + defer consoleBuf.mu.Unlock() + for _, line := range consoleBuf.lines { + if strings.Contains(line, "info-to-seq-only") { + t.Fatalf("Info leaked to console: %v", consoleBuf.lines) + } + } + foundWarn := false + for _, line := range consoleBuf.lines { + if strings.Contains(line, "warn-both") { + foundWarn = true + } + } + if !foundWarn { + t.Fatalf("Warn missing on console: %v", consoleBuf.lines) + } + l.CloseSeq() +} + +func TestSeqWriter_CloseFlushes(t *testing.T) { + var mu sync.Mutex + var gotBody string + srv := startFakeSeqServer(t, func(body string) { + mu.Lock() + gotBody += body + mu.Unlock() + }) + defer srv.Close() + + l := NewLogger(1, "", 0) + l.SetSeqWriter(srv.URL, "", "test:close") + l.Info().Msg("flush-on-close") + l.CloseSeq() + + mu.Lock() + defer mu.Unlock() + if !strings.Contains(gotBody, "flush-on-close") { + t.Fatalf("Close did not flush: %q", gotBody) + } +} + func TestStdoutCapture_LongLine(t *testing.T) { l, buf := newLoggerWithCapture(1) pr, pw := io.Pipe() diff --git a/log/logger.go b/log/logger.go index 318e579..ba465cb 100644 --- a/log/logger.go +++ b/log/logger.go @@ -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 → @t(zerolog 格式 "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 到 Seq;queue 关闭后冲刷并退出 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/缓存/通用工具)管线, diff --git a/var.go b/var.go index 3d7f3ca..d73207f 100644 --- a/var.go +++ b/var.go @@ -61,7 +61,7 @@ var ConfigNote = Map{ "logHistory": "默认100,非必须,内存中保留最近N条错误日志,用于调试调阅,通过 Log.GetRecentErrors() 获取", "webConnectLogShow": "默认true,非必须,访问日志如果需要web访问链接、访问ip、访问时间打印,false为关闭true开启此功能", "webConnectLogFile": "无默认,非必须,webConnectLogShow开启之后才能使用,如果需要存储日志文件时使用,保存格式为:a/b/c/20060102150405.txt,将生成:a/b/c/年月日时分秒.txt,按需设置", - "seqUrl": "无默认,非必须,Seq 日志平台地址,如 http://127.0.0.1:5341,填写后自动将所有日志通过异步 channel 队列推送到 Seq,空值则不激活;同时自动将 fmt.Println/标准log 包的输出也一并捕获推送;instance 字段自动拼接为 ip:port(如 192.168.1.10:8085),支持单机多进程和多服务器集群", + "seqUrl": "无默认,非必须,Seq 日志平台地址,如 http://127.0.0.1:5341,填写后自动将所有日志通过异步 channel 队列推送到 Seq,空值则不激活;同时自动将 fmt.Println/标准log 包的输出也一并捕获推送;instance 字段自动拼接为 ip:port;请求日志自动带 sid(sessionId 前12位脱敏)与 request_id(并回写 X-Request-Id),可用 that.LogBind 追加 user_id 等;访问日志优先取 X-Real-IP,有 EO-Client-IPCountry 时记 ip_country;配置后控制台自动仅打 Warn/Error,Seq 与文件仍按 logLevel 全量", "seqApiKey": "无默认,非必须,Seq API Key,单机免费版留空即可,多用户或有认证要求时填写", //"codeConfig": Map{ // "注释": "配置即启用,非必须,默认无",