package loggerx import ( "io" "log" "time" ) // 异步队列容量,跟随实例 const asyncQueueSize = 1000 // Close 等待异步日志落盘的最长时间 const closeDrainTimeout = 5 * time.Second // 写入,需要判断同步还是异步 func (l *Logger) write(event string, b []byte) (n int, err error) { if l.toAsync(event, b) { return len(b), nil } return l.store(event, b) } // 实际的存储 // 必须满足 io.Writer 契约:成功时 n == len(b),否则调用方(log / io.MultiWriter) // 会认为发生了短写并把日志吞掉 func (l *Logger) store(event string, b []byte) (n int, err error) { if l.option.isPrintFile { // 串行化写入:句柄缓冲区不是并发安全的, // 异步消费协程与同步调用可能同时写同一个句柄 l.writeMu.Lock() n, err = l.storeFile(event, b) l.writeMu.Unlock() if err != nil { return 0, err } } // 驱动(控制台等)的写入失败不改变对调用方的契约 _, _ = l.writeDrivers(b) return len(b), nil } // storeFile 写入日志文件 func (l *Logger) storeFile(event string, b []byte) (int, error) { f, err := l.getFile(event) if err != nil { return 0, err } n, err := f.Write(b) if err == nil && n < len(b) { err = io.ErrShortWrite } if err == nil { return n, nil } // 写入失败:丢弃这个句柄,落到磁盘后按最新文件名重开一次再写 // 只重试一次,避免原实现在短写时反复重开文件 l.discardFile(event, f) if nf, nerr := l.getFile(event); nerr == nil { if n2, err2 := nf.Write(b); err2 == nil && n2 == len(b) { return n2, nil } } return 0, err } // 写入额外的驱动(控制台 / 自定义 writer) func (l *Logger) writeDrivers(b []byte) (int, error) { if len(l.option.drivers) == 0 { return 0, nil } return io.MultiWriter(l.option.drivers...).Write(b) } // discardFile 关闭文件并从缓存中移除,使下次写入重新打开 func (l *Logger) discardFile(event string, f *logFile) { _ = f.Close() l.mu.Lock() defer l.mu.Unlock() key := fileKey{channel: l.channel, event: event} if cur, ok := l.filePath[key]; ok && cur == f { delete(l.filePath, key) } } // 异步队列的任务 type cacheData struct { Event string Data []byte } // toAsync 尝试异步写入 // 返回 true 表示已经交给异步队列,返回 false 表示需要同步写入 func (l *Logger) toAsync(event string, b []byte) bool { if l.writeType == writeTypeSync || // 指定同步模式 (l.writeType == writeTypeDefault && l.option.writeType != writeTypeAsync) { // 默认同步模式 return false } // 整段「检查开关 + 入队」都在同一把锁内:Close 也拿这把锁来停止投递, // 这样 Close 拿到锁时就能确定「要么这条还没入队、要么已经完整入队」。 // 否则消费协程可能先看到空队列就退出,把还在路上的这条日志整条丢掉。 // 队列满时这里会阻塞,但消费者是独立 goroutine 且不需要这把锁,不会死锁 l.async.mu.Lock() defer l.async.mu.Unlock() if l.async.closed { // 已开始关闭:退化为同步写入 return false } if l.async.ch == nil { l.async.ch = make(chan cacheData, asyncQueueSize) go l.asyncWorker(l.async.ch) } l.async.ch <- cacheData{Event: event, Data: b} return true } // asyncWorker 消费异步队列,直到队列被关闭 func (l *Logger) asyncWorker(q chan cacheData) { defer close(l.workerDone) defer func() { if r := recover(); r != nil { log.Println("loggerx: 异步写入协程异常:", r) } }() for val := range q { _, _ = l.store(val.Event, val.Data) } } // drainAsync 停止投递、关闭队列,并等消费协程把剩余任务全部写完 func (l *Logger) drainAsync() { // 持锁关闭投递:此后 toAsync 一律退化为同步写入 l.async.mu.Lock() l.async.closed = true q := l.async.ch l.async.mu.Unlock() if q == nil { // 从未启用过异步写入,没有后台 goroutine 需要等 return } // 关闭队列并等消费者把缓冲里的任务全部处理完。 // 这一步必须真的等到 workerDone:否则 Close 会在消费协程还在写缓冲时 // 就刷盘并关闭文件,最后几条日志会连着句柄一起丢掉 close(q) select { case <-l.workerDone: case <-time.After(closeDrainTimeout): log.Println("loggerx: 等待异步日志落盘超时") } }