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{ 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 } // NewLoggerNoCaller 创建不带自动 Caller 字段的日志实例(用于访问日志等 caller 无意义的场景) func NewLoggerNoCaller(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{ 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). Level(level). With(). Timestamp(). Logger() l := &Logger{ zl: zl, mw: mw, 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{ 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 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 } 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) } if appLayerCaller != "" && result != appLayerCaller { result += " <- " + appLayerCaller } 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: 异步推送到 Seq(CLEF 格式)--- // SeqWriter 将 zerolog JSON 日志通过 channel 队列异步发送到 Seq。 // Write() 仅做 channel <- bytes,O(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 → @t(zerolog 格式 "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/缓存/通用工具)管线, // 这些文件在调用栈中会被跳过,以显示真正的业务调用者。 // 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 } } if strings.HasSuffix(file, "/context.go") || strings.HasSuffix(file, "\\context.go") { return true } } return false }