You cannot select more than 25 topics Topics must start with a letter or number, can include dashes ('-') and can be up to 35 characters long.
nnet/internal/unpacker/delimiter.go

108 lines
3.1 KiB
Go

This file contains ambiguous Unicode characters!

This file contains ambiguous Unicode characters that may be confused with others in your current locale. If your use case is intentional and legitimate, you can safely ignore this warning. Use the Escape button to highlight these characters.

package unpacker
import (
"bytes"
unpackerpkg "git.noahlan.cn/noahlan/nnet/v2/pkg/unpacker"
)
// delimiterUnpacker 分隔符拆包器实现
type delimiterUnpacker struct {
delimiter []byte
buffer []byte
maxBufferSize int
}
// NewDelimiterUnpacker 创建分隔符拆包器
func NewDelimiterUnpacker(delimiter []byte) unpackerpkg.Unpacker {
return NewDelimiterUnpackerWithMaxBuffer(delimiter, unpackerpkg.DefaultMaxBufferSize)
}
// NewDelimiterUnpackerWithMaxBuffer 创建分隔符拆包器指定最大buffer大小
func NewDelimiterUnpackerWithMaxBuffer(delimiter []byte, maxBufferSize int) unpackerpkg.Unpacker {
if len(delimiter) == 0 {
delimiter = []byte{'\n'} // 默认换行符
}
if maxBufferSize <= 0 {
maxBufferSize = unpackerpkg.DefaultMaxBufferSize
}
return &delimiterUnpacker{
delimiter: delimiter,
buffer: make([]byte, 0, 4096), // 预分配初始容量
maxBufferSize: maxBufferSize,
}
}
// Unpack 拆包
func (u *delimiterUnpacker) Unpack(data []byte) ([][]byte, []byte, int, error) {
// 检查buffer大小限制
newSize := len(u.buffer) + len(data)
if newSize > u.maxBufferSize {
return nil, nil, 0, unpackerpkg.NewErrorf("unpacker buffer size exceeded: %d > %d", newSize, u.maxBufferSize)
}
// 优化:如果容量不足,预分配更大的容量(零拷贝优化)
if cap(u.buffer) < newSize {
newCap := cap(u.buffer) * 2
if newCap < newSize {
newCap = newSize
}
if newCap > u.maxBufferSize {
newCap = u.maxBufferSize
}
// 如果现有 buffer 为空,直接分配新 buffer避免不必要的复制
if len(u.buffer) == 0 {
u.buffer = make([]byte, 0, newCap)
} else {
newBuffer := make([]byte, len(u.buffer), newCap)
copy(newBuffer, u.buffer)
u.buffer = newBuffer
}
}
u.buffer = append(u.buffer, data...)
var messages [][]byte
for {
index := bytes.Index(u.buffer, u.delimiter)
if index == -1 {
// 没有找到分隔符,等待更多数据
break
}
// 提取消息(包含分隔符)
message := make([]byte, index+len(u.delimiter))
copy(message, u.buffer[:index+len(u.delimiter)])
messages = append(messages, message)
// 移除已处理的数据(优化:使用切片操作,避免复制)
u.buffer = u.buffer[index+len(u.delimiter):]
// 如果 buffer 太大但剩余数据很少,压缩 buffer减少内存占用
// 注意压缩不会改变buffer的长度只改变容量
if len(u.buffer) < cap(u.buffer)/4 && cap(u.buffer) > 4096 {
compressed := make([]byte, len(u.buffer), cap(u.buffer)/2)
copy(compressed, u.buffer)
u.buffer = compressed
}
}
// 输入字节已经被复制进连接级buffer调用方应从底层读缓冲中丢弃本次输入避免重复处理。
return messages, u.buffer, len(data), nil
}
// Pack 打包
func (u *delimiterUnpacker) Pack(data []byte) ([]byte, error) {
// 检查数据是否已包含分隔符
if bytes.HasSuffix(data, u.delimiter) {
return data, nil
}
// 添加分隔符
result := make([]byte, len(data)+len(u.delimiter))
copy(result, data)
copy(result[len(data):], u.delimiter)
return result, nil
}