chore: buffer opt
This commit is contained in:
@@ -67,6 +67,9 @@ func NewExecutorService(
|
|||||||
settingsService SettingsService,
|
settingsService SettingsService,
|
||||||
envService EnvService,
|
envService EnvService,
|
||||||
) *ExecutorService {
|
) *ExecutorService {
|
||||||
|
// 0. 清理旧临时日志
|
||||||
|
CleanupOrphanedTinyLogs()
|
||||||
|
|
||||||
es := &ExecutorService{
|
es := &ExecutorService{
|
||||||
taskService: taskService,
|
taskService: taskService,
|
||||||
taskLogService: taskLogService,
|
taskLogService: taskLogService,
|
||||||
|
|||||||
@@ -7,10 +7,11 @@ import (
|
|||||||
"fmt"
|
"fmt"
|
||||||
"io"
|
"io"
|
||||||
"os"
|
"os"
|
||||||
|
"path/filepath"
|
||||||
"sync"
|
"sync"
|
||||||
"unicode/utf8"
|
|
||||||
|
|
||||||
"github.com/engigu/baihu-panel/internal/constant"
|
"github.com/engigu/baihu-panel/internal/constant"
|
||||||
|
"github.com/engigu/baihu-panel/internal/logger"
|
||||||
"github.com/engigu/baihu-panel/internal/utils"
|
"github.com/engigu/baihu-panel/internal/utils"
|
||||||
)
|
)
|
||||||
|
|
||||||
@@ -57,7 +58,7 @@ type TinyLog struct {
|
|||||||
path string
|
path string
|
||||||
writer *bufio.Writer
|
writer *bufio.Writer
|
||||||
subscribers []chan []byte
|
subscribers []chan []byte
|
||||||
remainder []byte // Leftover bytes from previous write (partial multi-byte characters)
|
remainder []byte // Leftover bytes from previous write (partial lines)
|
||||||
masks []string // Secrets to mask
|
masks []string // Secrets to mask
|
||||||
closed bool
|
closed bool
|
||||||
}
|
}
|
||||||
@@ -90,7 +91,6 @@ func (l *TinyLog) Write(p []byte) (n int, err error) {
|
|||||||
return 0, os.ErrClosed
|
return 0, os.ErrClosed
|
||||||
}
|
}
|
||||||
|
|
||||||
// 1. 合并上次调用剩余的字节(可能是半个 UTF-8 字符)
|
|
||||||
originalInputLen := len(p)
|
originalInputLen := len(p)
|
||||||
payload := p
|
payload := p
|
||||||
if len(l.remainder) > 0 {
|
if len(l.remainder) > 0 {
|
||||||
@@ -98,42 +98,43 @@ func (l *TinyLog) Write(p []byte) (n int, err error) {
|
|||||||
l.remainder = nil
|
l.remainder = nil
|
||||||
}
|
}
|
||||||
|
|
||||||
// 2. 识别结尾不完整的 UTF-8 序列
|
// 1. 寻找最后一个换行符
|
||||||
lastSafe := len(payload)
|
lastNewline := bytes.LastIndexByte(payload, '\n')
|
||||||
// UTF-8 字符最多 4 字节,检查最后几个字节
|
if lastNewline == -1 {
|
||||||
for i := len(payload) - 1; i >= 0 && i >= len(payload)-4; i-- {
|
// 没有换行符,且如果长度超过 4KB,强制截断并输出,防止内存无限制增长
|
||||||
if utf8.RuneStart(payload[i]) {
|
if len(payload) > 4096 {
|
||||||
if !utf8.FullRune(payload[i:]) {
|
lastNewline = len(payload) - 1
|
||||||
// 发现末尾存在不完整字符
|
} else {
|
||||||
lastSafe = i
|
// 保留当前所有内容到下一轮
|
||||||
l.remainder = make([]byte, len(payload)-i)
|
l.remainder = payload
|
||||||
copy(l.remainder, payload[i:])
|
return originalInputLen, nil
|
||||||
}
|
|
||||||
break
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
// 如果整个负载都不完整且不超过一个 UTF-8 字符的最大长度,
|
// 2. 提取出完整的行
|
||||||
// 则全部保留到下次调用
|
completeBytes := payload[:lastNewline+1]
|
||||||
if lastSafe == 0 && len(l.remainder) > 0 {
|
|
||||||
return originalInputLen, nil
|
// 3. 剥离并保存剩余的部分
|
||||||
|
if lastNewline+1 < len(payload) {
|
||||||
|
l.remainder = make([]byte, len(payload)-(lastNewline+1))
|
||||||
|
copy(l.remainder, payload[lastNewline+1:])
|
||||||
}
|
}
|
||||||
|
|
||||||
// 3. 仅将完整的部分转换为 UTF-8,并调用封装的函数进行脱敏处理
|
// 4. 将完整行转换为 UTF-8 并脱敏
|
||||||
text := utils.MaskSecrets(utils.ToUTF8(payload[:lastSafe]), l.masks)
|
text := utils.MaskSecrets(utils.ToUTF8(completeBytes), l.masks)
|
||||||
data := []byte(text)
|
outData := []byte(text)
|
||||||
|
|
||||||
// 4. 写入文件缓冲区
|
// 5. 输出安全部分
|
||||||
_, err = l.writer.Write(data)
|
_, err = l.writer.Write(outData)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return 0, err
|
return 0, err
|
||||||
}
|
}
|
||||||
|
|
||||||
// 5. 广播给所有订阅者
|
// 6. 广播给所有订阅者
|
||||||
if len(l.subscribers) > 0 {
|
if len(l.subscribers) > 0 {
|
||||||
for _, ch := range l.subscribers {
|
for _, ch := range l.subscribers {
|
||||||
select {
|
select {
|
||||||
case ch <- data:
|
case ch <- outData:
|
||||||
default:
|
default:
|
||||||
// 如果订阅者处理太慢,丢弃消息以避免阻塞写入
|
// 如果订阅者处理太慢,丢弃消息以避免阻塞写入
|
||||||
}
|
}
|
||||||
@@ -324,3 +325,23 @@ func (l *TinyLog) ReadLastLines(n int) ([]byte, error) {
|
|||||||
func (l *TinyLog) GetPath() string {
|
func (l *TinyLog) GetPath() string {
|
||||||
return l.path
|
return l.path
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// CleanupOrphanedTinyLogs 启动时清理残留的临时日志文件
|
||||||
|
func CleanupOrphanedTinyLogs() {
|
||||||
|
tmpDir := os.TempDir()
|
||||||
|
files, err := os.ReadDir(tmpDir)
|
||||||
|
if err != nil {
|
||||||
|
return
|
||||||
|
}
|
||||||
|
|
||||||
|
count := 0
|
||||||
|
for _, file := range files {
|
||||||
|
if !file.IsDir() && len(file.Name()) > 9 && file.Name()[:9] == "task_log_" && filepath.Ext(file.Name()) == ".log" {
|
||||||
|
os.Remove(filepath.Join(tmpDir, file.Name()))
|
||||||
|
count++
|
||||||
|
}
|
||||||
|
}
|
||||||
|
if count > 0 {
|
||||||
|
logger.Infof("[System] 清理了 %d 个残留的任务日志临时文件", count)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|||||||
@@ -2,10 +2,12 @@ package utils
|
|||||||
|
|
||||||
import (
|
import (
|
||||||
"fmt"
|
"fmt"
|
||||||
|
"sync"
|
||||||
)
|
)
|
||||||
|
|
||||||
// TailBuffer 是一个只保留最后 N 字节数据的缓冲区
|
// TailBuffer 是一个只保留最后 N 字节数据的缓冲区
|
||||||
type TailBuffer struct {
|
type TailBuffer struct {
|
||||||
|
mu sync.Mutex
|
||||||
limit int
|
limit int
|
||||||
data []byte
|
data []byte
|
||||||
size int // 当前实际存储的大小
|
size int // 当前实际存储的大小
|
||||||
@@ -22,6 +24,9 @@ func NewTailBuffer(limit int) *TailBuffer {
|
|||||||
|
|
||||||
// Write 实现 io.Writer 接口
|
// Write 实现 io.Writer 接口
|
||||||
func (b *TailBuffer) Write(p []byte) (n int, err error) {
|
func (b *TailBuffer) Write(p []byte) (n int, err error) {
|
||||||
|
b.mu.Lock()
|
||||||
|
defer b.mu.Unlock()
|
||||||
|
|
||||||
n = len(p)
|
n = len(p)
|
||||||
if n >= b.limit {
|
if n >= b.limit {
|
||||||
// 如果单次写入就超过了限制,直接取最后 limit 字节
|
// 如果单次写入就超过了限制,直接取最后 limit 字节
|
||||||
@@ -43,16 +48,24 @@ func (b *TailBuffer) Write(p []byte) (n int, err error) {
|
|||||||
|
|
||||||
// Bytes 返回缓冲区内的所有数据
|
// Bytes 返回缓冲区内的所有数据
|
||||||
func (b *TailBuffer) Bytes() []byte {
|
func (b *TailBuffer) Bytes() []byte {
|
||||||
return b.data
|
b.mu.Lock()
|
||||||
|
defer b.mu.Unlock()
|
||||||
|
res := make([]byte, len(b.data))
|
||||||
|
copy(res, b.data)
|
||||||
|
return res
|
||||||
}
|
}
|
||||||
|
|
||||||
// String 返回缓冲区内的字符串表示
|
// String 返回缓冲区内的字符串表示
|
||||||
func (b *TailBuffer) String() string {
|
func (b *TailBuffer) String() string {
|
||||||
|
b.mu.Lock()
|
||||||
|
defer b.mu.Unlock()
|
||||||
return string(b.data)
|
return string(b.data)
|
||||||
}
|
}
|
||||||
|
|
||||||
// Len 返回当前存储的数据长度
|
// Len 返回当前存储的数据长度
|
||||||
func (b *TailBuffer) Len() int {
|
func (b *TailBuffer) Len() int {
|
||||||
|
b.mu.Lock()
|
||||||
|
defer b.mu.Unlock()
|
||||||
return len(b.data)
|
return len(b.data)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
Reference in New Issue
Block a user