Files
hotime/log/logger.go
T

756 lines
18 KiB
Go
Raw Normal View History

package log
import (
"bufio"
"bytes"
"encoding/json"
"fmt"
"io"
"net/http"
"os"
"path/filepath"
"runtime"
"strings"
"sync"
"sync/atomic"
"time"
"github.com/rs/zerolog"
)
const maxCaptureLine = 10 << 20 // 10MB,支持 GetReqMap 等大 body 单行日志
var (
origStderr io.Writer
origStderrOnce sync.Once
)
func saveOrigStderr() {
origStderrOnce.Do(func() {
origStderr = os.Stderr
})
}
func seqWriteStderr(format string, args ...interface{}) {
if origStderr != nil {
fmt.Fprintf(origStderr, format, args...)
} else {
fmt.Fprintf(os.Stderr, format, args...)
}
}
// ErrorRecord 错误历史记录条目
type ErrorRecord struct {
Err error
Msg string
Time time.Time
Caller string
}
// Logger 日志核心结构体,封装 zerolog
type Logger struct {
zl zerolog.Logger
mw *multiWriter // 动态多输出,支持运行时追加 SeqWriter 等
logLevel int
errors []ErrorRecord
errIdx int // 环形缓冲写入位置
errCount int // 实际写入总数(用于判断缓冲是否满)
errorsMu sync.RWMutex
maxErrors int
}
// NewLogger 创建日志实例
// logLevel: 0=仅 error>=1=全部
// logFile: 文件路径模板(如 "logs/20060102.txt"),空则不写文件
// maxErrors: 错误历史最大条数,0 则不记录
func NewLogger(logLevel int, logFile string, maxErrors int) *Logger {
saveOrigStderr()
zerolog.CallerMarshalFunc = callerMarshalFunc
zerolog.TimeFieldFormat = "2006-01-02 15:04:05"
var level zerolog.Level
if logLevel == 0 {
level = zerolog.ErrorLevel
} else {
level = zerolog.DebugLevel
}
consoleWriter := zerolog.ConsoleWriter{
2026-04-14 14:07:59 +08:00
Out: os.Stderr,
TimeFormat: "2006-01-02 15:04:05",
NoColor: false,
FormatLevel: formatLevelColor,
FormatCaller: func(i interface{}) string {
if i == nil {
return ""
}
return fmt.Sprintf("[%s]", i)
},
}
var writers []io.Writer
writers = append(writers, consoleWriter)
if logFile != "" {
fw := &TemplateFileWriter{pathTemplate: logFile}
writers = append(writers, fw)
}
mw := &multiWriter{writers: writers}
zl := zerolog.New(mw).
Level(level).
With().
Timestamp().
Caller().
Logger()
l := &Logger{
zl: zl,
mw: mw,
logLevel: logLevel,
maxErrors: maxErrors,
}
if maxErrors > 0 {
l.errors = make([]ErrorRecord, maxErrors)
}
return l
}
2026-04-14 14:07:59 +08:00
// NewLoggerNoCaller 创建不带自动 Caller 字段的日志实例(用于访问日志等 caller 无意义的场景)
func NewLoggerNoCaller(logLevel int, logFile string, maxErrors int) *Logger {
saveOrigStderr()
2026-04-14 14:07:59 +08:00
zerolog.CallerMarshalFunc = callerMarshalFunc
zerolog.TimeFieldFormat = "2006-01-02 15:04:05"
var level zerolog.Level
if logLevel == 0 {
level = zerolog.ErrorLevel
} else {
level = zerolog.DebugLevel
}
consoleWriter := zerolog.ConsoleWriter{
Out: os.Stderr,
TimeFormat: "2006-01-02 15:04:05",
NoColor: false,
FormatLevel: formatLevelColor,
}
var writers []io.Writer
writers = append(writers, consoleWriter)
if logFile != "" {
fw := &TemplateFileWriter{pathTemplate: logFile}
writers = append(writers, fw)
}
mw := &multiWriter{writers: writers}
zl := zerolog.New(mw).
2026-04-14 14:07:59 +08:00
Level(level).
With().
Timestamp().
Logger()
l := &Logger{
zl: zl,
mw: mw,
2026-04-14 14:07:59 +08:00
logLevel: logLevel,
maxErrors: maxErrors,
}
if maxErrors > 0 {
l.errors = make([]ErrorRecord, maxErrors)
}
return l
}
// NewTestLogger 创建用于测试的静默 Logger(输出到 io.Discard
func NewTestLogger() *Logger {
zl := zerolog.New(io.Discard).Level(zerolog.Disabled)
return &Logger{zl: zl, logLevel: 0}
}
// SetOutput 设置日志输出(用于测试等场景)
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{
2026-04-14 14:07:59 +08:00
Out: w,
TimeFormat: "2006-01-02 15:04:05",
NoColor: true,
FormatLevel: formatLevelPlain,
FormatCaller: func(i interface{}) string {
if i == nil {
return ""
}
return fmt.Sprintf("[%s]", i)
},
}
mw := &multiWriter{writers: []io.Writer{consoleWriter}}
l.mw = mw
l.zl = zerolog.New(mw).
With().
Timestamp().
Caller().
Logger()
}
// GetLevel 获取当前日志等级
func (l *Logger) GetLevel() int {
return l.logLevel
}
// EmitCaptured 将捕获的 stdout/stderr 等输出直写 multiWriter,绕过主 Logger 的 level 过滤。
// 保证 logLevel=0 时 fmt.Println / log.Println 仍能进 Seq。
func (l *Logger) EmitCaptured(level zerolog.Level, source, msg string) {
if l == nil || l.mw == nil || msg == "" {
return
}
zl := zerolog.New(l.mw)
zl.WithLevel(level).Timestamp().Str("source", source).Msg(msg)
}
// EmitRecoveredPanic 记录 recover 到的 panic:错误原因 + 调用栈,source=panic。
func (l *Logger) EmitRecoveredPanic(err interface{}) {
if l == nil || l.mw == nil || err == nil {
return
}
zl := zerolog.New(l.mw)
zl.WithLevel(zerolog.ErrorLevel).Timestamp().
Str("source", "panic").
Str("stack", string(debugStack())).
Msg(fmt.Sprint(err))
}
func debugStack() []byte {
buf := make([]byte, 64*1024)
n := runtime.Stack(buf, false)
return buf[:n]
}
// CaptureStream 从 reader 持续读取行并写入 Logger(用于 stdout/stderr 管道)。
// 单行超过 maxCaptureLine 时截断并标注,读取出错后重启读取,避免永久失效。
func CaptureStream(l *Logger, level zerolog.Level, source string, r io.Reader) {
rd := bufio.NewReader(r)
for {
line, err := rd.ReadString('\n')
if len(line) > 0 {
msg := strings.TrimRight(line, "\r\n")
if msg != "" {
if len(msg) > maxCaptureLine {
msg = msg[:maxCaptureLine] + "...(truncated)"
}
l.EmitCaptured(level, source, msg)
}
}
if err != nil {
if err == io.EOF {
if len(line) == 0 {
return
}
continue
}
l.EmitCaptured(zerolog.ErrorLevel, source+"-capture", err.Error())
return
}
}
}
// --- 链式调用 API ---
// 以下方法兼容两种调用风格:
// 无参数:返回 *zerolog.Event 用于链式调用,如 l.Error().Str("k","v").Msg("...")
// 有参数:直接拼接并打印日志(logrus 兼容),返回 no-op Event
func (l *Logger) Debug(args ...interface{}) *zerolog.Event {
if len(args) > 0 {
l.zl.Debug().Msg(fmt.Sprint(args...))
nop := zerolog.Nop()
return nop.Debug()
}
return l.zl.Debug()
}
func (l *Logger) Info(args ...interface{}) *zerolog.Event {
if len(args) > 0 {
l.zl.Info().Msg(fmt.Sprint(args...))
nop := zerolog.Nop()
return nop.Info()
}
return l.zl.Info()
}
func (l *Logger) Warn(args ...interface{}) *zerolog.Event {
if len(args) > 0 {
l.zl.Warn().Msg(fmt.Sprint(args...))
nop := zerolog.Nop()
return nop.Warn()
}
return l.zl.Warn()
}
func (l *Logger) Error(args ...interface{}) *zerolog.Event {
if len(args) > 0 {
msg := fmt.Sprint(args...)
if l.maxErrors > 0 {
l.recordError(msg, nil)
}
l.zl.Error().Msg(msg)
nop := zerolog.Nop()
return nop.Error()
}
if l.maxErrors > 0 {
l.recordError("", nil)
}
return l.zl.Error()
}
// RecordError 手动记录一条错误到历史(用于需要指定详情的场景)
func (l *Logger) RecordError(err error, msg string, caller string) {
if l.maxErrors <= 0 {
return
}
l.errorsMu.Lock()
defer l.errorsMu.Unlock()
l.errors[l.errIdx] = ErrorRecord{
Err: err,
Msg: msg,
Time: time.Now(),
Caller: caller,
}
l.errIdx = (l.errIdx + 1) % l.maxErrors
l.errCount++
}
func (l *Logger) recordError(msg string, err error) {
_, file, line, ok := runtime.Caller(2)
caller := ""
if ok {
caller = formatCaller(file, line)
}
l.RecordError(err, msg, caller)
}
// GetRecentErrors 获取最近 N 条错误(不传则返回全部已存储的)
func (l *Logger) GetRecentErrors(n ...int) []ErrorRecord {
if l.maxErrors <= 0 {
return nil
}
l.errorsMu.RLock()
defer l.errorsMu.RUnlock()
total := l.errCount
if total > l.maxErrors {
total = l.maxErrors
}
if total == 0 {
return nil
}
want := total
if len(n) > 0 && n[0] > 0 && n[0] < want {
want = n[0]
}
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() {
break
}
result = append(result, l.errors[idx])
}
return result
}
// --- 格式化快捷方法 ---
func (l *Logger) Debugf(format string, v ...interface{}) {
l.zl.Debug().Msgf(format, v...)
}
func (l *Logger) Infof(format string, v ...interface{}) {
l.zl.Info().Msgf(format, v...)
}
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 {
l.recordError(fmt.Sprintf(format, v...), nil)
}
l.zl.Error().Msgf(format, v...)
}
// --- 日志级别颜色 ---
func formatLevelColor(i interface{}) string {
level := strings.ToUpper(fmt.Sprintf("%s", i))
switch level {
case "DEBUG":
return fmt.Sprintf("\x1b[36m|%s|\x1b[0m", level) // cyan
case "INFO":
return fmt.Sprintf("\x1b[32m|%s|\x1b[0m", level) // green
case "WARN":
return fmt.Sprintf("\x1b[33m|%s|\x1b[0m", level) // yellow
case "ERROR":
return fmt.Sprintf("\x1b[31m|%s|\x1b[0m", level) // red
case "FATAL":
return fmt.Sprintf("\x1b[35m|%s|\x1b[0m", level) // magenta
default:
return fmt.Sprintf("|%s|", level)
}
}
func formatLevelPlain(i interface{}) string {
return fmt.Sprintf("|%s|", strings.ToUpper(fmt.Sprintf("%s", i)))
}
// --- TemplateFileWriter ---
// TemplateFileWriter 按时间模板切换文件路径的 Writer
type TemplateFileWriter struct {
pathTemplate string
mu sync.Mutex
currentPath string
file *os.File
writer *bufio.Writer
}
func (w *TemplateFileWriter) Write(p []byte) (n int, err error) {
w.mu.Lock()
defer w.mu.Unlock()
path := time.Now().Format(w.pathTemplate)
if path != w.currentPath {
if w.writer != nil {
_ = w.writer.Flush()
}
if w.file != nil {
_ = w.file.Close()
}
_ = os.MkdirAll(filepath.Dir(path), 0755)
w.file, err = os.OpenFile(path, os.O_APPEND|os.O_CREATE|os.O_WRONLY, 0644)
if err != nil {
return 0, err
}
w.writer = bufio.NewWriterSize(w.file, 4096)
w.currentPath = path
}
n, err = w.writer.Write(p)
_ = w.writer.Flush()
return
}
// --- 调用者智能过滤 ---
func callerMarshalFunc(_ uintptr, file string, line int) string {
return findCaller(file, line)
}
func findCaller(origFile string, origLine int) string {
var lastInfraFile string
var lastInfraLine int
2026-04-14 14:07:59 +08:00
var primaryCaller string
var appLayerCaller string
for i := 1; i < 20; i++ {
_, file, line, ok := runtime.Caller(i)
if !ok {
break
}
shortFile := shortenPath(file)
if isInfrastructureFile(shortFile) {
lastInfraFile = shortFile
lastInfraLine = line
continue
}
2026-04-14 14:07:59 +08:00
if primaryCaller == "" {
primaryCaller = fmt.Sprintf("%s:%d", shortFile, line)
}
if appLayerCaller == "" && !strings.Contains(strings.ToLower(file), "hotime") {
appLayerCaller = fmt.Sprintf("%s:%d", shortFile, line)
}
}
result := ""
if primaryCaller != "" {
result = primaryCaller
} else if lastInfraFile != "" {
result = fmt.Sprintf("%s:%d", lastInfraFile, lastInfraLine)
} else {
result = fmt.Sprintf("%s:%d", shortenPath(origFile), origLine)
}
2026-04-14 14:07:59 +08:00
if appLayerCaller != "" && result != appLayerCaller {
result += " <- " + appLayerCaller
}
2026-04-14 14:07:59 +08:00
return result
}
func formatCaller(file string, line int) string {
return fmt.Sprintf("%s:%d", shortenPath(file), line)
}
func shortenPath(file string) string {
n := 0
for i := len(file) - 1; i > 0; i-- {
if file[i] == '/' || file[i] == '\\' {
n++
if n >= 2 {
return file[i+1:]
}
}
}
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:
n := atomic.AddInt64(&w.dropped, 1)
if n == 1 || n%100 == 0 {
seqWriteStderr("[seq] queue full, dropped=%d\n", n)
}
}
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 {
seqWriteStderr("[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 {
seqWriteStderr("[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)
seqWriteStderr("[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/缓存/通用工具)管线,
// 这些文件在调用栈中会被跳过,以显示真正的业务调用者。
2026-04-14 14:07:59 +08:00
// context.go 是 Display/View 的纯中间层,也作为基础设施跳过。
// application.go、code/makecode.go 等保留,会作为有意义的调用者返回。
func isInfrastructureFile(file string) bool {
if strings.HasPrefix(file, "zerolog/") || strings.HasPrefix(file, "zerolog@") {
return true
}
if strings.HasPrefix(file, "logrus/") || strings.HasPrefix(file, "logrus@") {
return true
}
if strings.HasPrefix(file, "runtime/") {
return true
}
if strings.HasPrefix(file, "log/") {
return true
}
infraPrefixes := []string{"db/", "cache/", "common/", "dri/"}
for _, prefix := range infraPrefixes {
if strings.HasPrefix(file, prefix) {
return true
}
}
lowerFile := strings.ToLower(file)
if strings.Contains(lowerFile, "hotime") {
infraDirs := []string{"db/", "cache/", "common/", "log/", "dri/"}
for _, dir := range infraDirs {
if strings.Contains(file, dir) {
return true
}
}
2026-04-14 14:07:59 +08:00
if strings.HasSuffix(file, "/context.go") || strings.HasSuffix(file, "\\context.go") {
return true
}
}
return false
}