兴趣笔记

深入交流 490 views 20 replies
#1 ·

我觉得有意思的学习内容 :innocent_cat:

以前我一直认为自己对知识的渴望是纯粹的,认知、建模是我所有的意义。应试的学习让我感到不屑,却又因这份借口而无需过多进行实践

现在看来,对学习进步的需求无非来自三种追求:

  • 获得知识本身
  • 优越感
  • 逃避劣等的恐惧

也许大部份人都有这些追求,只是不同的人和学习的内容,三种占比不同吧

我曾经自以为的纯粹知识渴望,也不过很大程度是对劣等感的逃避,为了没有胆量承认的虚荣

然而也许内心并不纯粹,追求俗不可耐,这样的自己才是真实的,而非想象中某种完美的幻影

无论出发点如何,如果还有兴趣能让我感到幸福,那便足够了吧 🙂

❤️wave_gif👍
7
#2 ·

SplitArgs

  • 把字符串根据空格/引号分割成多个 token
  • 无引号 token 按照字面字符处理
  • 单引号 token 仅 unescape '' -> '
  • 双引号 token 按照 YAML 双引号字符串 unescape

FSM

splitargs|690x343

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

方法论

写这种复杂状态机之前,先用表格大致设计出状态和迁移会顺利很多

image|689x102

我第一次写 IP Prefix trie 的时候完全梦到哪写到哪,哪里有问题就多加一个分支,最后变得非常莫名其妙 😅

❤️
3
#3 ·

每过一两个月就觉得以前写的东西依托构使,恨不得狠狠 refactor 但是又没精力有没有懂的

❤️plus1
2
#4 ·

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
#5 ·

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 缓解锁竞争

这里 randmath/rand/v2,因为 v1 的全局函数是同步的


跑分

benchplot|690x431

Sharded 2 <= GOMAXPROCS <= 4 突然很慢是怎么回事有没有懂的? 🫠
pprof 我还不太能用明白

#6 ·
stateless_trash_can 无状态垃圾桶
  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            .          .           } 

何意味

#7 · (edited)
签名的 Anon 酱老是被 GC 掉,在这里放一个也许就不会了吧

Btw 上传 webp 似乎有问题,转换后的 avif 图片无法正确显示

cat_rabbit❤️
1
mod
#8 ·
stateless_trash_can 无状态垃圾桶

我这边可以正常显示 😍

年月擦身过,暂且问,会更好吗?

plus1wave_gif
2
#9 ·

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
#10 ·
一只玉米人 玉米🌽

anon.webp|375x500
刚才那是 png,这个能显示吗?

#11 ·
stateless_trash_can 无状态垃圾桶
stateless_trash_can:

Log rotate (复制 + 截短)
pw.size < lastSize -> ChangeTruncated

这里如果在一次迭代期间,截短并写入超过原来文件大小的内容,会误判为 ChangeModified,暂时想不到怎么办 🫠 (不过这种情况对于日志来说不太可能?)

后果是丢失一部分行

#12 ·
stateless_trash_can 无状态垃圾桶

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
			}
		}
	}
}
admin
#13 · (edited)
stateless_trash_can 无状态垃圾桶

能正常显示。所有的图片除了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
👍
1
#14 ·
IO IO

现在可以了,之前疑似没加载出来

不过之前确实上传 webp 无法立即显示,但是 png 可以

现在无法复现了

Edit: 编辑器内上传图片无法立即加载,手动在新标签页访问图片路径后就可以在编辑器内显示了,不知道是不是网络问题

#15 ·

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
}

合并了 ChangeMovedChangeDeleted(全部当作 ChangeDeleted

FSM

tail.svg|690x496

ref: https://github.com/nxadm/tail

好复杂,能 syslog over UDP 就别搞这种(

#16 ·
IO IO

Minor issue
Screenshot From 2026-07-14 12-24-23.png|690x77

admin
#17 ·
stateless_trash_can 无状态垃圾桶

github.svg|16x16
已配置,刷新一下

👍
1
#18 · (edited)
stateless_trash_can 无状态垃圾桶
stateless_trash_can:

Tail

我是垃圾桶,这写的就是垃圾

❤️
1
#20 ·
stateless_trash_can 无状态垃圾桶
stateless_trash_can:

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.
	// 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
}