兴趣笔记
我觉得有意思的学习内容 ![]()
以前我一直认为自己对知识的渴望是纯粹的,认知、建模是我所有的意义。应试的学习让我感到不屑,却又因这份借口而无需过多进行实践
现在看来,对学习进步的需求无非来自三种追求:
- 获得知识本身
- 优越感
- 逃避劣等的恐惧
也许大部份人都有这些追求,只是不同的人和学习的内容,三种占比不同吧
我曾经自以为的纯粹知识渴望,也不过很大程度是对劣等感的逃避,为了没有胆量承认的虚荣
然而也许内心并不纯粹,追求俗不可耐,这样的自己才是真实的,而非想象中某种完美的幻影
无论出发点如何,如果还有兴趣能让我感到幸福,那便足够了吧 🙂
SplitArgs
- 把字符串根据空格/引号分割成多个 token
- 无引号 token 按照字面字符处理
- 单引号 token 仅 unescape
''->' - 双引号 token 按照 YAML 双引号字符串 unescape
FSM

Go 实现
这样用闭包按照我的理解应该不会有性能问题?可惜因为需要 unescape 无法做到 zero-allocation
const (
splitArgsLookForToken int = iota
splitArgsInToken
splitArgsInSingleQuotes
splitArgsInDoubleQuotes
splitArgsSingleQuoteEscape
splitArgsDoubleQuoteEscape
splitArgsEndSuccess
splitArgsEndFail
)
func SplitArgs(s string) (result []string, err error) {
var sb strings.Builder
var current rune
var currentSize int
var errMsg string
peeked := false
eoi := false
idx := 0
state := splitArgsLookForToken
peek := func() {
if peeked {
return
}
peeked = true
if idx >= len(s) {
eoi = true
} else {
current, currentSize = utf8.DecodeRuneInString(s[idx:])
}
}
peekHex := func(size int) error {
if idx+size > len(s) {
return errors.New("early EOI at index " + strconv.Itoa(idx))
}
r := rune(0)
for i := range size {
c := s[idx+i]
var v rune
switch {
case c >= '0' && c <= '9':
v = rune(c - '0')
case c >= 'a' && c <= 'f':
v = rune(c - 'a' + 10)
case c >= 'A' && c <= 'F':
v = rune(c - 'A' + 10)
default:
return errors.New("bad hex digit at index " + strconv.Itoa(idx+i))
}
r = (r << 4) | v
}
current = r
currentSize = size
return nil
}
// Must only be called after calling peek()
pop := func() {
idx += currentSize
peeked = false
}
begin := func() {
sb.Reset()
}
end := func() {
result = append(result, sb.String())
}
// Must only be called after calling peek() and before calling pop()
write := func() {
sb.WriteRune(current)
}
for {
switch state {
case splitArgsEndFail:
err = errors.New("SplitArgs: " + errMsg)
fallthrough
case splitArgsEndSuccess:
return
case splitArgsLookForToken:
peek()
if eoi {
state = splitArgsEndSuccess
continue
}
switch current {
case '\t', '\n', '\v', '\f', '\r', ' ':
pop()
continue
case '\'':
pop()
begin()
state = splitArgsInSingleQuotes
continue
case '"':
pop()
begin()
state = splitArgsInDoubleQuotes
continue
default:
begin()
state = splitArgsInToken
continue
}
case splitArgsInToken:
peek()
if eoi {
end()
state = splitArgsEndSuccess
continue
}
switch current {
case '\t', '\n', '\v', '\f', '\r', ' ':
pop()
end()
state = splitArgsLookForToken
continue
default:
write()
pop()
continue
}
case splitArgsInSingleQuotes:
peek()
if eoi {
errMsg = "early EOI at index " + strconv.Itoa(idx)
state = splitArgsEndFail
continue
}
switch current {
case '\'':
pop()
state = splitArgsSingleQuoteEscape
continue
default:
write()
pop()
continue
}
case splitArgsInDoubleQuotes:
peek()
if eoi {
errMsg = "early EOI at index " + strconv.Itoa(idx)
state = splitArgsEndFail
continue
}
switch current {
case '"':
end()
pop()
state = splitArgsLookForToken
continue
case '\\':
pop()
state = splitArgsDoubleQuoteEscape
continue
default:
write()
pop()
continue
}
case splitArgsSingleQuoteEscape:
peek()
if !eoi && current == '\'' {
write()
pop()
state = splitArgsInSingleQuotes
continue
} else {
end()
state = splitArgsLookForToken
continue
}
case splitArgsDoubleQuoteEscape:
peek()
if eoi {
errMsg = "early EOI at index " + strconv.Itoa(idx)
state = splitArgsEndFail
continue
}
switch current {
case '"', '\\', '/', ' ', '\x09':
write()
pop()
state = splitArgsInDoubleQuotes
continue
case '0', 'a', 'b', 't', 'n', 'v', 'f', 'r', 'e', 'N', '_', 'L', 'P':
switch current {
case '0':
current = '\x00'
case 'a':
current = '\a'
case 'b':
current = '\b'
case 't':
current = '\t'
case 'n':
current = '\n'
case 'v':
current = '\v'
case 'f':
current = '\f'
case 'r':
current = '\r'
case 'e':
current = '\x1b'
case 'N':
current = '\x85'
case '_':
current = '\xa0'
case 'L':
current = '\u2028'
case 'P':
current = '\u2029'
}
write()
pop()
state = splitArgsInDoubleQuotes
continue
case 'x', 'u', 'U':
var size int
switch current {
case 'x':
size = 2
case 'u':
size = 4
case 'U':
size = 8
}
pop()
if err := peekHex(size); err != nil {
errMsg = err.Error()
state = splitArgsEndFail
continue
}
write()
pop()
state = splitArgsInDoubleQuotes
continue
default:
errMsg = "bad escape sequence at index " + strconv.Itoa(idx)
state = splitArgsEndFail
continue
}
default:
errMsg = "internal error"
state = splitArgsEndFail
continue
}
}
}
感觉接口设计有点问题,似乎应该把 unescape 和 split 分开? 🤔
才发现状态名写错了应该是 unescape
方法论
写这种复杂状态机之前,先用表格大致设计出状态和迁移会顺利很多

我第一次写 IP Prefix trie 的时候完全梦到哪写到哪,哪里有问题就多加一个分支,最后变得非常莫名其妙 😅
每过一两个月就觉得以前写的东西依托构使,恨不得狠狠 refactor 但是又没精力有没有懂的
Ring Buffer Read/Write
行为:
- 写入/读取 异步安全
- 读取时,如果有效数据不足,仅读取有效数据并报错
- 写入时,如果有效空间不足,仅写入有效空间并报错
- 保留 1 字节区分 完全空/完全满
结构体
type Ring struct {
buffer []byte
wp atomic.Int64 // write pointer
_ [64 - 24 - 8]byte
rp atomic.Int64 // read pointer
_ [64 - 8]byte
}
读取
func (ring *Ring) Read(p []byte) (n int, err error) {
wp := int(ring.wp.Load()) // Might be stale but it's ok
rp := int(ring.rp.Load())
defer func() {
ring.rp.Store(int64(rp))
}()
// data size
n = wp - rp
if n < 0 {
n += len(ring.buffer)
}
if n < len(p) {
err = io.EOF // B, D
} else {
n = len(p) // A, C
}
if rem := n - copy(p[:n], ring.buffer[rp:]); rem > 0 {
copy(p[n-rem:], ring.buffer[:rem]) // C, D
rp = rem
} else {
rp += n // A, B
}
return
}
状态变量:
- EOF:
len(p)大于有效数据 - Wrap: 读取跨越切片结尾
| Not EOF | EOF | |
|---|---|---|
| Not Wrap | A | B |
| Wrap | C | D |
A
0 1 2 3 4 5 6 7
x x o o #[x x]x
^ ^
wp rp
Read(p), len(p): 2
n = 2 - 5 + 8 = 5 > 2: n = 2
rem = 2 - copy(p[:2], buf[5:]) = 0
rp += 2 = 7
B
0 1 2 3 4 5 6 7
o #[x x x]o o o
^ ^
rp wp
Read(p), len(p): 4
n = 5 - 2 = 3 < 4: EOF
rem = 3 - copy(p[:3], buf[2:]) = 0
rp += 3 = 5
C
0 1 2 3 4 5 6 7
x]x o o #[x x x
^ ^
wp rp
Read(p), len(p): 4
n = 2 - 5 + 8 = 5 > 4: n = 4
rem = 4 - copy(p[:4], buf[5:]) = 1
copy(p[4-1:], buf[:1])
rp = 1
D
0 1 2 3 4 5 6 7
x]x o o #[x x x
^ ^
wp rp
Read(p), len(p): 6
n = 2 - 5 + 8 = 5 < 6: EOF
rem = 5 - copy(p[:5], buf[5:]) = 2
copy(p[5-2:], buf[:2])
rp = 2
写入
func (ring *Ring) Write(p []byte) (n int, err error) {
rp := int(ring.rp.Load()) // Might be stale but it's ok
wp := int(ring.wp.Load())
defer func() {
ring.wp.Store(int64(wp))
}()
// free space
n = rp - wp - 1
if n < 0 {
n += len(ring.buffer)
}
if n < len(p) {
err = ErrBufferFull // B, D
} else {
n = len(p) // A, C
}
if rem := n - copy(ring.buffer[wp:], p[:n]); rem > 0 {
copy(ring.buffer[:rem], p[n-rem:]) // C, D
wp = rem
} else {
wp += n // A, B
}
return
}
状态变量:
- Full:
len(p)大于有效空间 - Wrap: 写入跨越切片结尾
| Not Full | Full | |
|---|---|---|
| Not Wrap | A | B |
| Wrap | C | D |
A
0 1 2 3 4 5 6 7
x x[o o]# x x x
^ ^
wp rp
Write(p), len(p): 2
n = 5 - 2 - 1 = 2 >= 2: n = 2
rem = 2 - copy(buf[2:], p[:2]) = 0
wp += 2 = 4
B
0 1 2 3 4 5 6 7
x x[o o]# x x x
^ ^
wp rp
Write(p), len(p): 4
n = 5 - 2 - 1 = 2 < 4: Full
rem = 2 - copy(buf[2:], p[:2]) = 0
wp += 2 = 4
C
0 1 2 3 4 5 6 7
o #[x x x]o o o
^ ^
rp wp
Write(p), len(p): 4
n = 2 - 5 - 1 + 8 = 4 >= 4: n = 4
rem = 4 - copy(p[:4], buf[5:]) = 1
copy(buf[:1], p[4-1:])
wp = 1
D
0 1 2 3 4 5 6 7
o #[x x x]o o o
^ ^
rp wp
Write(p), len(p): 6
n = 2 - 5 - 1 + 8 = 4 < 6: Full
rem = 4 - copy(p[:4], buf[5:]) = 1
copy(buf[:1], p[4-1:])
wp = 1
ID Counter
type IDCounter interface {
GetID() uint64
}
行为:
- 异步安全
- 有限时间内 ID 唯一
SimpleIDCounter
type SimpleIDCounter struct{ atomic.Uint64 }
func (c *SimpleIDCounter) GetID() uint64 {
return c.Add(1)
}
异步安全,uint64 可生成 1 << 64 个唯一 ID,高并发下由于锁竞争延迟线性提高
RandomShardedIDCounter
type shard struct {
value atomic.Uint64
_ [64 - 8]byte
}
type RandomShardedIDCounter struct {
shards []shard
mask int
shifts int
}
func NewRandomShardedIDCounter(n int) *RandomShardedIDCounter {
if n <= 0 {
n = 1
}
numShards := 1
shifts := 1
for numShards < n {
numShards <<= 1
shifts++
}
return &RandomShardedIDCounter{
shards: make([]shard, numShards),
mask: numShards - 1,
shifts: shifts,
}
}
func (c *RandomShardedIDCounter) GetID() uint64 {
idx := rand.Int() & c.mask
val := c.shards[idx].value.Add(1)
return (val << uint64(c.shifts)) | uint64(idx)
}
异步安全,用 shardIdx 作为 ID 低位,每个 shard 可连续生成 1 << (64 - bitWidth(shardIdx)) 个唯一 ID
用随机 shardIdx 均分请求到各 shard 缓解锁竞争
这里 rand 是 math/rand/v2,因为 v1 的全局函数是同步的
跑分

Sharded 2 <= GOMAXPROCS <= 4 突然很慢是怎么回事有没有懂的? 🫠
pprof 我还不太能用明白
无状态垃圾桶 Total: 2.23s 2.33s (flat, cum) 98.73%
60 . . mask: numShards - 1,
61 . . shifts: shifts,
62 . . }
63 . . }
64 . .
65 10ms 10ms func (c *RandomShardedIDCounter) GetID() uint64 {
66 60ms 160ms idx := rand.Int() & c.mask
67 . .
68 60ms 60ms val := c.shards[idx].value.Add(1)
69 . .
70 2.10s 2.10s return (val << uint64(c.shifts)) | uint64(idx)
71 . . }
何意味
Btw 上传 webp 似乎有问题,转换后的 avif 图片无法正确显示
无状态垃圾桶 我这边可以正常显示 😍
年月擦身过,暂且问,会更好吗?
Log file change watcher
接口
type Change uint8
const (
ChangeCreated Change = iota
ChangeDeleted
ChangeMoved
ChangeModified
ChangeTruncated
)
type Watcher interface {
BlockUntilChange(ctx context.Context) (Change, error)
}
Polling 实现
package tail
import (
"context"
"os"
"time"
)
var _ Watcher = (*PollWatcher)(nil)
type PollWatcher struct {
interval time.Duration
filepath string
fi os.FileInfo
size int64
modified time.Time
}
func NewPollWatcher(path string, interval time.Duration) (*PollWatcher, error) {
fi, err := os.Stat(path)
if err != nil {
if os.IsNotExist(err) {
fi = nil
} else {
return nil, err
}
}
pw := &PollWatcher{
filepath: path,
interval: interval,
fi: fi,
}
if fi != nil {
pw.size = fi.Size()
pw.modified = fi.ModTime()
}
return pw, nil
}
func (pw *PollWatcher) BlockUntilChange(ctx context.Context) (Change, error) {
ticker := time.NewTicker(pw.interval)
defer ticker.Stop()
for {
select {
case <-ctx.Done():
return 0, ctx.Err()
case <-ticker.C:
fi, err := os.Stat(pw.filepath)
if err != nil {
if os.IsNotExist(err) {
if pw.fi != nil {
pw.fi = nil
return ChangeDeleted, nil
}
continue
}
return 0, err
}
lastSize := pw.size
lastModified := pw.modified
lastFi := pw.fi
pw.size = fi.Size()
pw.modified = fi.ModTime()
pw.fi = fi
if lastFi == nil {
return ChangeCreated, nil
}
if !os.SameFile(lastFi, fi) {
return ChangeMoved, nil
}
if pw.size < lastSize {
return ChangeTruncated, nil
}
if pw.size > lastSize || !pw.modified.Equal(lastModified) {
return ChangeModified, nil
}
}
}
}
| 情况 | 判定 |
|---|---|
| 正常 append | pw.size > lastSize -> ChangeModified |
| Log rotate (重命名 + 新建) | !os.SameFile(lastFi, fi) -> ChangeMoved |
| Log rotate (复制 + 截短) | pw.size < lastSize -> ChangeTruncated |
| 新建 | lastFi == nil -> ChangeCreated |
| 删除 | pw.fi != nil + os.IsNotExist -> ChangeDeleted |

刚才那是 png,这个能显示吗?
无状态垃圾桶 FSNotify 实现
package tail
import (
"context"
"errors"
"fmt"
"os"
"path/filepath"
"github.com/fsnotify/fsnotify"
)
var _ Watcher = (*FSNotifyWatcher)(nil)
type FSNotifyWatcher struct {
watcher *fsnotify.Watcher
filepath string
size int64
wasPresent bool
}
func NewFSNotifyWatcher(path string) (*FSNotifyWatcher, error) {
watcher, err := fsnotify.NewWatcher()
if err != nil {
return nil, err
}
cleanPath := filepath.Clean(path)
parentDir := filepath.Dir(cleanPath)
// Watch parent and filter by Event.Name
err = watcher.Add(parentDir)
if err != nil {
watcher.Close()
return nil, err
}
// Get initial size and presence
var size int64
wasPresent := false
fi, err := os.Stat(path)
if err == nil {
size = fi.Size()
wasPresent = true
}
return &FSNotifyWatcher{
watcher: watcher,
filepath: cleanPath,
size: size,
wasPresent: wasPresent,
}, nil
}
func (fw *FSNotifyWatcher) Close() error {
return fw.watcher.Close()
}
func (fw *FSNotifyWatcher) BlockUntilChange(ctx context.Context) (Change, error) {
for {
select {
case <-ctx.Done():
return 0, ctx.Err()
case err, ok := <-fw.watcher.Errors:
if !ok {
return 0, errors.New("watcher channel closed")
}
return 0, err
case event, ok := <-fw.watcher.Events:
if !ok {
return 0, errors.New("watcher channel closed")
}
if filepath.Clean(event.Name) != fw.filepath {
continue
}
switch {
case event.Has(fsnotify.Create):
if fw.wasPresent {
if fi, err := os.Stat(fw.filepath); err == nil {
fw.size = fi.Size()
}
continue
}
fw.wasPresent = true
return ChangeCreated, nil
case event.Has(fsnotify.Write):
if fi, err := os.Stat(fw.filepath); err == nil {
lastSize := fw.size
fw.size = fi.Size()
if fw.size < lastSize {
return ChangeTruncated, nil
} else {
return ChangeModified, nil
}
} else {
return 0, fmt.Errorf("os.Stat failed: %w", err)
}
case event.Has(fsnotify.Remove):
fw.wasPresent = false
return ChangeDeleted, nil
case event.Has(fsnotify.Rename):
return ChangeMoved, nil
}
}
}
}
无状态垃圾桶 能正常显示。所有的图片除了svg之外,不管是什么格式,都会被ImageMagick转换为AVIF格式。如果没法正常显示,可能是浏览器不支持: https://caniuse.com/avif
其实普遍的做法是根据不同的浏览器UserAgent,服务器发送支持的格式。
另外AVIF虽然压缩率比较高,浏览器广泛支持,但也存在一些问题:
- 压缩/转码成avif较慢。
- 不支持progrssive decoding,导致没法在网络较慢时像jpg一样一行一行地加载,只能全部下载后才能加载。
比较理想的是jpeg-xl,支持progressive decoding, encoding较快。
但是目前仅支持最新版Chrome/Firefox并需要启用experimental setting:
- Chrome:
访问 chrome://flags/
enable-jxl-image-format - Firefox:
访问 about:config
image.jxl.enabled
IO 现在可以了,之前疑似没加载出来
不过之前确实上传 webp 无法立即显示,但是 png 可以
现在无法复现了
Edit: 编辑器内上传图片无法立即加载,手动在新标签页访问图片路径后就可以在编辑器内显示了,不知道是不是网络问题
Tail
package tail
import (
"bufio"
"bytes"
"context"
"errors"
"fmt"
"io"
"os"
"sync"
"time"
)
const (
pollInterval = 250 * time.Millisecond
)
type fakeReader struct{}
func (*fakeReader) Read(p []byte) (n int, err error) {
return 0, nil
}
type Config struct {
// Seek to before tailing.
Seek SeekInfo
// Provide appended data as the file grows.
Follow bool
// Keep trying to open a file if it is inaccessible.
// Ignored when Follow is false.
Retry bool
// Use PollWatcher for file changes.
// FSNotifyWatcher is used if false.
Poll bool
// Disable seeking.
// Enable this for named pipes (mkfifo).
NoSeek bool
// Line delimiter.
Delimiter byte
}
func DefaultConfig() Config {
return Config{
Delimiter: '\n',
}
}
// SeekInfo
type SeekInfo struct {
Offset int64
Whence int
}
type Tail struct {
cfg Config
filename string
file *os.File
watcher Watcher
reader *bufio.Reader
firstOpen bool
mu sync.Mutex
ctx context.Context
cancel context.CancelFunc
}
func New(ctx context.Context, filename string, cfg Config) (*Tail, error) {
var watcher Watcher
var err error
if cfg.Poll {
watcher, err = NewPollWatcher(filename, pollInterval)
if err != nil {
return nil, fmt.Errorf("failed to create PollWatcher: %w", err)
}
} else {
watcher, err = NewFSNotifyWatcher(filename)
if err != nil {
return nil, fmt.Errorf("failed to create FSNotifyWatcher: %w", err)
}
}
return newTailWithWatcher(ctx, filename, cfg, watcher), nil
}
func newTailWithWatcher(ctx context.Context, filename string, cfg Config, watcher Watcher) *Tail {
ctx, cancel := context.WithCancel(ctx)
return &Tail{
cfg: cfg,
filename: filename,
watcher: watcher,
reader: bufio.NewReader(&fakeReader{}),
firstOpen: true,
ctx: ctx,
cancel: cancel,
}
}
var (
ErrBufferTooShort = errors.New("buffer is too short to read complete line")
)
func (t *Tail) ReadLine(b []byte) (n int, err error) {
t.mu.Lock()
defer t.mu.Unlock()
for {
// Try to read a complete line already in the buffer
if t.file != nil {
n, found, err := t.tryReadLine(b)
if found || err != nil {
return n, err
}
// No complete line in buffer
if !t.cfg.Follow {
return 0, io.EOF
}
}
// Eagerly open file
if t.file == nil && !t.cfg.Follow {
if err := t.openFile(); err == nil {
continue
}
}
// Wait for next file change
ch, bErr := t.watcher.BlockUntilChange(t.ctx)
if bErr != nil {
return 0, bErr
}
fileWasOpen := t.file != nil
if t.file == nil {
if ch != ChangeDeleted {
err = t.openFile()
if err != nil {
if t.cfg.Retry && t.cfg.Follow {
continue
}
return
}
} else {
continue
}
}
switch ch {
case ChangeCreated:
if fileWasOpen {
err = t.reopenFile()
if err != nil {
return
}
}
case ChangeDeleted:
err = t.closeFile()
if err != nil {
return
}
case ChangeModified:
t.reader.Reset(t.file)
case ChangeTruncated:
if !t.cfg.NoSeek {
_, err = t.file.Seek(0, io.SeekStart)
if err != nil {
return
}
}
}
}
}
func (t *Tail) tryReadLine(b []byte) (n int, found bool, err error) {
nBuf := t.reader.Buffered()
if nBuf == 0 {
if _, err := t.reader.Peek(1); err != nil {
if errors.Is(err, io.EOF) {
return 0, false, nil
}
return 0, false, err
}
nBuf = t.reader.Buffered()
}
peeked, err := t.reader.Peek(nBuf)
if err != nil && !errors.Is(err, io.EOF) {
return 0, false, err
}
idx := bytes.IndexByte(peeked, t.cfg.Delimiter)
if idx < 0 {
return 0, false, nil
}
if len(b) < idx {
return 0, false, ErrBufferTooShort
}
copied := copy(b, peeked[:idx])
t.reader.Discard(idx + 1)
return copied, true, nil
}
func (t *Tail) Close() error {
t.cancel()
return t.closeFile()
}
func (t *Tail) closeFile() (err error) {
if t.file != nil {
err = t.file.Close()
t.file = nil
}
return
}
func (t *Tail) openFile() error {
file, err := os.Open(t.filename)
if err != nil {
return err
}
if !t.cfg.NoSeek && t.firstOpen {
if _, err := file.Seek(t.cfg.Seek.Offset, t.cfg.Seek.Whence); err != nil {
file.Close()
return err
}
}
t.firstOpen = false
t.file = file
t.reader.Reset(t.file)
return nil
}
func (t *Tail) reopenFile() (err error) {
err = t.closeFile()
if err != nil {
return
}
err = t.openFile()
return
}
合并了 ChangeMoved 和 ChangeDeleted(全部当作 ChangeDeleted)
FSM
好复杂,能 syslog over UDP 就别搞这种(
IO Minor issue

无状态垃圾桶
已配置,刷新一下
无状态垃圾桶 重构
package tail
import (
"bufio"
"bytes"
"context"
"errors"
"fmt"
"io"
"os"
"sync"
"time"
)
const (
pollInterval = 250 * time.Millisecond
)
type fakeReader struct{}
func (*fakeReader) Read(p []byte) (n int, err error) {
return 0, nil
}
type Config struct {
// Seek to before tailing.
Seek SeekInfo
// Provide appended data as the file grows.
Follow bool
// Keep trying to open a file if it is inaccessible.
// Requires Follow to be set.
Retry bool
// Use PollWatcher for file changes.
// FSNotifyWatcher is used if false.
Poll bool
// Disable seeking.
// Enable this for named pipes (mkfifo).
NoSeek bool
// MaxLineSize sets the maximum line length in bytes.
// Lines exceeding this are returned in chunks of MaxLineSize with error ErrMaxLineSizeExceeded.
// Defaults to 4096 when <= 0.
MaxLineSize int
// Line delimiter.
Delimiter byte
}
func DefaultConfig() Config {
return Config{
Delimiter: '\n',
MaxLineSize: 4096,
}
}
// SeekInfo
type SeekInfo struct {
Offset int64
Whence int
}
type tailState int8
const (
stateInit tailState = iota
stateReady
stateEOF
stateNotOpen
stateFailed
)
type Tail struct {
cfg Config
filename string
file *os.File
watcher Watcher
reader *bufio.Reader
state tailState
failErr error
mu sync.Mutex
ctx context.Context
cancel context.CancelFunc
}
func New(ctx context.Context, filename string, cfg Config) (*Tail, error) {
if !cfg.Follow && cfg.Retry {
return nil, errors.New("Retry set without Follow")
}
var watcher Watcher
var err error
if cfg.Poll {
watcher, err = NewPollWatcher(filename, pollInterval)
if err != nil {
return nil, fmt.Errorf("failed to create PollWatcher: %w", err)
}
} else {
watcher, err = NewFSNotifyWatcher(filename)
if err != nil {
return nil, fmt.Errorf("failed to create FSNotifyWatcher: %w", err)
}
}
return newTailWithWatcher(ctx, filename, cfg, watcher), nil
}
func newTailWithWatcher(ctx context.Context, filename string, cfg Config, watcher Watcher) *Tail {
ctx, cancel := context.WithCancel(ctx)
if cfg.MaxLineSize <= 0 {
cfg.MaxLineSize = 4096
}
return &Tail{
cfg: cfg,
filename: filename,
watcher: watcher,
reader: bufio.NewReaderSize(&fakeReader{}, cfg.MaxLineSize),
state: stateInit,
ctx: ctx,
cancel: cancel,
}
}
var (
ErrMaxLineSizeExceeded = errors.New("line exceeds MaxLineSize")
ErrShortBuffer = io.ErrShortBuffer
ErrClosed = errors.New("closed")
)
func (t *Tail) ReadLine(b []byte) (int, error) {
t.mu.Lock()
defer t.mu.Unlock()
for {
switch t.state {
case stateInit:
// Try to open file
file, err := os.Open(t.filename)
if err != nil {
t.failErr = err
t.state = stateNotOpen
continue
} else {
if !t.cfg.NoSeek {
_, err = file.Seek(t.cfg.Seek.Offset, t.cfg.Seek.Whence)
if err != nil {
t.failErr = err
t.state = stateFailed
return 0, err
}
}
t.file = file
t.reader.Reset(file)
t.state = stateReady
continue
}
case stateReady:
// Try to read line
n := t.reader.Buffered()
if n == 0 {
if _, err := t.reader.Peek(1); err != nil {
if errors.Is(err, io.EOF) {
t.state = stateEOF
continue
} else {
t.failErr = err
t.state = stateFailed
return 0, err
}
}
n = t.reader.Buffered()
}
peeked, err := t.reader.Peek(n)
if err != nil && !errors.Is(err, io.EOF) {
t.failErr = err
t.state = stateFailed
return 0, err
}
idx := bytes.IndexByte(peeked, t.cfg.Delimiter)
if idx < 0 {
// No delimiter found
if n >= t.cfg.MaxLineSize {
// Line exceeds MaxLineSize
// Return chunk
n = t.cfg.MaxLineSize
if len(b) < n {
return 0, ErrShortBuffer
}
peeked, _ = t.reader.Peek(n)
copied := copy(b, peeked[:n])
t.reader.Discard(n)
return copied, ErrMaxLineSizeExceeded
}
t.state = stateEOF
continue
}
// Delimiter found
// Chunk if line exceeds MaxLineSize
if idx >= t.cfg.MaxLineSize {
n = t.cfg.MaxLineSize
if len(b) < n {
return 0, ErrShortBuffer
}
peeked, _ = t.reader.Peek(n)
copied := copy(b, peeked[:n])
t.reader.Discard(n)
return copied, ErrMaxLineSizeExceeded
}
if len(b) < idx {
return 0, ErrShortBuffer
}
copied := copy(b, peeked[:idx])
t.reader.Discard(idx + 1)
return copied, nil
case stateEOF:
if !t.cfg.Follow {
n := t.reader.Buffered()
if n > 0 {
// Probe for more data
_, peekErr := t.reader.Peek(n + 1)
if errors.Is(peekErr, io.EOF) {
// True EOF
// Return as final line
peeked, _ := t.reader.Peek(n)
copied := copy(b, peeked[:n])
t.reader.Discard(n)
return copied, nil
}
if peekErr != nil {
t.failErr = peekErr
t.state = stateFailed
return 0, peekErr
}
// Has more
t.state = stateReady
continue
}
t.failErr = io.EOF
t.state = stateFailed
return 0, io.EOF
}
ch, err := t.watcher.BlockUntilChange(t.ctx)
if err != nil {
t.failErr = err
t.state = stateFailed
return 0, err
}
switch ch {
case ChangeModified, ChangeCreated:
// Drain EOF error in bufio.Reader
if n := t.reader.Buffered(); n > 0 && n < t.cfg.MaxLineSize {
t.reader.Peek(n + 1) // trigger file read
}
t.state = stateReady
continue
case ChangeDeleted:
// Drain buffer from old file
if n := t.reader.Buffered(); n > 0 {
if len(b) < n {
return 0, ErrShortBuffer
}
peeked, _ := t.reader.Peek(n)
copied := copy(b, peeked[:n])
t.reader.Discard(n)
t.file.Close()
t.file = nil
t.state = stateNotOpen
return copied, nil
}
err = t.file.Close()
if err != nil {
t.failErr = err
t.state = stateFailed
return 0, err
}
t.file = nil
t.state = stateNotOpen
continue
case ChangeTruncated:
if !t.cfg.NoSeek {
_, err := t.file.Seek(0, io.SeekStart)
if err != nil {
t.failErr = err
t.state = stateFailed
return 0, err
}
}
t.reader.Reset(t.file)
t.state = stateReady
continue
}
case stateNotOpen:
if !t.cfg.Retry {
t.state = stateFailed
return 0, errors.New("cannot open file")
}
ch, err := t.watcher.BlockUntilChange(t.ctx)
if err != nil {
t.failErr = err
t.state = stateFailed
return 0, err
}
switch ch {
case ChangeCreated, ChangeModified, ChangeTruncated:
// Try to open file
file, err := os.Open(t.filename)
if err != nil {
continue
} else {
t.file = file
t.reader.Reset(file)
t.state = stateReady
continue
}
case ChangeDeleted:
t.failErr = errors.New("got ChangeDeleted in stateNotOpen")
t.state = stateFailed
return 0, t.failErr
}
case stateFailed:
return 0, t.failErr
}
}
}
func (t *Tail) Close() error {
t.cancel()
t.mu.Lock()
defer t.mu.Unlock()
t.failErr = ErrClosed
t.state = stateFailed
if t.file != nil {
if err := t.file.Close(); err != nil {
return err
}
}
return nil
}
