更新
This commit is contained in:
+115
-20
@@ -1,8 +1,9 @@
|
||||
package loggerx
|
||||
|
||||
import (
|
||||
"errors"
|
||||
"fmt"
|
||||
"io"
|
||||
"log"
|
||||
"time"
|
||||
)
|
||||
|
||||
@@ -31,17 +32,23 @@ func (l *Logger) store(event string, b []byte) (n int, err error) {
|
||||
func (l *Logger) storeTo(channel, event string, b []byte) (n int, err error) {
|
||||
if l.option.isPrintFile {
|
||||
// 串行化写入:句柄缓冲区不是并发安全的,
|
||||
// 异步消费协程与同步调用可能同时写同一个句柄
|
||||
l.writeMu.Lock()
|
||||
n, err = l.storeFileTo(channel, event, b)
|
||||
l.writeMu.Unlock()
|
||||
// 异步消费协程与同步调用可能同时写同一个句柄。
|
||||
// 这里必须用 defer 解锁:一旦未来写入路径里出现 panic,
|
||||
// 非 defer 的 Unlock 会被跳过,锁永久不释放,整个进程的日志全卡死
|
||||
n, err = func() (int, error) {
|
||||
l.writeMu.Lock()
|
||||
defer l.writeMu.Unlock()
|
||||
return l.storeFileTo(channel, event, b)
|
||||
}()
|
||||
if err != nil {
|
||||
return 0, err
|
||||
}
|
||||
}
|
||||
|
||||
// 驱动(控制台等)的写入失败不改变对调用方的契约
|
||||
_, _ = l.writeDrivers(b)
|
||||
// 驱动(控制台等)的写入失败不改变对调用方的契约,但要上报
|
||||
if _, derr := l.writeDrivers(b); derr != nil {
|
||||
l.reportError(fmt.Errorf("loggerx: 写入额外驱动失败: %w", derr))
|
||||
}
|
||||
|
||||
return len(b), nil
|
||||
}
|
||||
@@ -88,11 +95,16 @@ func (l *Logger) storeFileTo(channel, event string, b []byte) (int, error) {
|
||||
}
|
||||
|
||||
// rollFileTo 归档当前文件并返回新文件句柄
|
||||
//
|
||||
// 任何失败路径都必须保证:filePath[key] 不会留下一个「已关闭」的句柄。
|
||||
// 否则后续写入只会进那个句柄的内存缓冲,而它已不在句柄表里,
|
||||
// 刷新和关闭都遍历不到 —— 数据静默丢失且接口返回成功
|
||||
func (l *Logger) rollFileTo(channel, event string, f *logFile) (*logFile, error) {
|
||||
key := fileKey{channel: channel, event: event}
|
||||
|
||||
// 先把缓冲清空再关句柄,保证归档内容完整
|
||||
if err := f.Close(); err != nil {
|
||||
l.repairHandle(key, f)
|
||||
return nil, err
|
||||
}
|
||||
|
||||
@@ -100,6 +112,7 @@ func (l *Logger) rollFileTo(channel, event string, f *logFile) (*logFile, error)
|
||||
// (按小时切割出来的名字本身就长这样:2026/09/13/06_info.log)
|
||||
path, err := l.archive(f.baseName, f.fileName)
|
||||
if err != nil {
|
||||
l.repairHandle(key, f)
|
||||
return nil, err
|
||||
}
|
||||
|
||||
@@ -112,6 +125,7 @@ func (l *Logger) rollFileTo(channel, event string, f *logFile) (*logFile, error)
|
||||
// 新文件带走递增序号,避免覆盖刚归档出去的同名文件
|
||||
nf, err := l.openNumberedFile(key, l.nextIndex())
|
||||
if err != nil {
|
||||
l.repairHandle(key, nil)
|
||||
return nil, err
|
||||
}
|
||||
|
||||
@@ -124,12 +138,60 @@ func (l *Logger) rollFileTo(channel, event string, f *logFile) (*logFile, error)
|
||||
return nf, nil
|
||||
}
|
||||
|
||||
// repairHandle 给某个 key 装回一个「可写」的句柄
|
||||
//
|
||||
// 滚动失败时现场可能残留一个已关闭的句柄(还在表里或已被摘掉),
|
||||
// 这里统一替换成一个新打开的同名文件句柄,保证后续写入有地方落盘;
|
||||
// 实在打不开就把表项摘掉,让下一次写入重新走完整的打开流程
|
||||
func (l *Logger) repairHandle(key fileKey, stale *logFile) {
|
||||
l.mu.Lock()
|
||||
if stale != nil {
|
||||
if cur, ok := l.filePath[key]; ok && cur == stale {
|
||||
delete(l.filePath, key)
|
||||
}
|
||||
} else {
|
||||
delete(l.filePath, key)
|
||||
}
|
||||
l.mu.Unlock()
|
||||
|
||||
nf, err := l.openNewFile(key)
|
||||
if err != nil {
|
||||
l.reportError(fmt.Errorf("loggerx: 滚动失败后无法重新打开日志文件 %s/%s: %w", key.channel, key.event, err))
|
||||
return
|
||||
}
|
||||
l.mu.Lock()
|
||||
// 期间可能有别的 goroutine 已经装好了句柄,别覆盖
|
||||
if _, ok := l.filePath[key]; !ok {
|
||||
l.filePath[key] = nf
|
||||
} else {
|
||||
_ = nf.Close()
|
||||
}
|
||||
l.mu.Unlock()
|
||||
}
|
||||
|
||||
// 写入额外的驱动(控制台 / 自定义 writer)
|
||||
//
|
||||
// 用户传进来的 driver 往往是 *bytes.Buffer / *bufio.Writer 这类非并发安全的对象,
|
||||
// 所以这里必须和文件写入一样串行化;同时逐个写、单个失败不影响其它 driver
|
||||
// (io.MultiWriter 会在第一个错误处短路,把后面的输出一起吞掉)
|
||||
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)
|
||||
|
||||
l.driverMu.Lock()
|
||||
defer l.driverMu.Unlock()
|
||||
|
||||
var errs []error
|
||||
for _, d := range l.option.drivers {
|
||||
if d == nil {
|
||||
continue
|
||||
}
|
||||
if _, err := d.Write(b); err != nil {
|
||||
errs = append(errs, err)
|
||||
}
|
||||
}
|
||||
return len(b), joinErrors(errs)
|
||||
}
|
||||
|
||||
// discardFile 关闭文件并从缓存中移除,使下次写入重新打开
|
||||
@@ -190,22 +252,40 @@ func (l *Logger) toAsync(event string, b []byte) bool {
|
||||
|
||||
defer l.async.end() // 必须在解锁之后 Done,保证 begin/end 覆盖整段投递
|
||||
|
||||
ch <- cacheData{Channel: l.channel, Event: event, Data: b}
|
||||
// 必须复制一份再入队:b 可能是标准库 log 从 sync.Pool 借来的行缓冲,
|
||||
// io.Writer 契约明确禁止保留传入的切片 —— log 在 Write 返回后立刻把
|
||||
// 缓冲还池并复用,异步消费协程读到的就会是被改写的内存
|
||||
// (曾经导致 1374/2000 条日志内容错乱)
|
||||
data := make([]byte, len(b))
|
||||
copy(data, b)
|
||||
|
||||
ch <- cacheData{Channel: l.channel, Event: event, Data: data}
|
||||
return true
|
||||
}
|
||||
|
||||
// asyncWorker 消费异步队列,直到队列被关闭
|
||||
//
|
||||
// 关键:单条任务 panic 绝不能让消费循环退出。
|
||||
// 否则队列没人消费、channel 很快写满,之后所有写入方都会永久阻塞在
|
||||
// `ch <- ...` 上(应用整体卡死),Close 也会永远等不到 workerDone。
|
||||
func (l *Logger) asyncWorker(q chan cacheData) {
|
||||
defer close(l.workerDone)
|
||||
|
||||
for val := range q {
|
||||
l.consumeOne(val)
|
||||
}
|
||||
}
|
||||
|
||||
// consumeOne 处理一条异步任务,把 panic 限制在这一条之内
|
||||
func (l *Logger) consumeOne(val cacheData) {
|
||||
defer func() {
|
||||
if r := recover(); r != nil {
|
||||
log.Println("loggerx: 异步写入协程异常:", r)
|
||||
// 回调要包一层 recover:用户的错误处理函数自己也可能 panic
|
||||
l.reportError(fmt.Errorf("loggerx: 异步写入单条日志时 panic: %v", r))
|
||||
}
|
||||
}()
|
||||
for val := range q {
|
||||
// 按任务里记的 channel 落盘,避免多个 channel 的日志混到根目录
|
||||
_, _ = l.storeTo(val.Channel, val.Event, val.Data)
|
||||
}
|
||||
// 按任务里记的 channel 落盘,避免多个 channel 的日志混到根目录
|
||||
_, _ = l.storeTo(val.Channel, val.Event, val.Data)
|
||||
}
|
||||
|
||||
// drainAsync 停止投递、关闭队列,并等消费协程把剩余任务全部写完
|
||||
@@ -221,18 +301,33 @@ func (l *Logger) drainAsync() {
|
||||
return
|
||||
}
|
||||
|
||||
// 关键:必须等「已经进入投递临界区」的写入全部入队之后才能关队列。
|
||||
// 这里能看到 wg 已经是 0,就说明没有写入还卡在投递路径上,
|
||||
// 否则 close(q) 之后它们再发就会 panic: send on closed channel
|
||||
// 等「已进入投递临界区」的写入全部入队,之后才能关队列。
|
||||
//
|
||||
// 这里刻意【不加超时】:wg 的非零计数正说明有 goroutine 卡在
|
||||
// 临界区里(通常是队列满导致 `ch <-` 阻塞)。若超时后就 close(q),
|
||||
// 那些还停在 `ch <-` 上的生产者会被唤醒并 panic: send on closed channel,
|
||||
// 而它们跑在应用自己的 goroutine 上(Info/Write 的调用方),
|
||||
// 没有 recover,直接把进程打挂。
|
||||
//
|
||||
// 不设超时也不会死等:只有消费者倒下才会让队列永久满,
|
||||
// 而消费者现在对每条任务单独 recover(见 consumeOne),不会死。
|
||||
l.async.wg.Wait()
|
||||
|
||||
// 关闭队列并等消费者把缓冲里的任务全部处理完。
|
||||
// 这一步必须真的等到 workerDone:否则 Close 会在消费协程还在写缓冲时
|
||||
// 就刷盘并关闭文件,最后几条日志会连着句柄一起丢掉
|
||||
close(q)
|
||||
if !waitChanTimeout(l.workerDone, closeDrainTimeout) {
|
||||
l.reportError(errors.New("loggerx: 等待异步日志落盘超时,队列中剩余日志可能丢失"))
|
||||
}
|
||||
}
|
||||
|
||||
// waitChanTimeout 等通道关闭,超时返回 false
|
||||
func waitChanTimeout(ch <-chan struct{}, d time.Duration) bool {
|
||||
select {
|
||||
case <-l.workerDone:
|
||||
case <-time.After(closeDrainTimeout):
|
||||
log.Println("loggerx: 等待异步日志落盘超时")
|
||||
case <-ch:
|
||||
return true
|
||||
case <-time.After(d):
|
||||
return false
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user