41 lines
1.4 KiB
Go
41 lines
1.4 KiB
Go
// Package checkpoint 实现 Agent 断点续传所需的检查点存储。
|
||
//
|
||
// 日志文件检查点记录「文件路径 → 已消费 offset」,指标上报检查点记录
|
||
// 「host_id → 已确认 seq」。默认提供 memory / file / redis 三种驱动,
|
||
// 生产环境使用 Redis 记录 offset,实现进程重启后从已确认点继续。
|
||
package checkpoint
|
||
|
||
import "fmt"
|
||
|
||
// Store 是检查点存储的抽象接口。
|
||
type Store interface {
|
||
// GetLogOffset 返回指定日志文件已消费到的 offset,不存在时返回 0。
|
||
GetLogOffset(file string) (int64, error)
|
||
// SetLogOffset 推进日志文件 offset。
|
||
SetLogOffset(file string, offset int64) error
|
||
|
||
// GetAckedSeq 返回指定主机已确认的上报序号,不存在时返回 0。
|
||
GetAckedSeq(hostID string) (int64, error)
|
||
// SetAckedSeq 推进主机已确认的上报序号。
|
||
SetAckedSeq(hostID string, seq int64) error
|
||
|
||
Close() error
|
||
}
|
||
|
||
// New 根据驱动名创建检查点存储。
|
||
//
|
||
// 支持:memory(进程内,测试/零依赖)、file(本地 JSON 持久化)、
|
||
// redis(Redis 记录 offset,生产推荐)。
|
||
func New(driver, file, redisURL string) (Store, error) {
|
||
switch driver {
|
||
case "", "memory":
|
||
return NewMemoryStore(), nil
|
||
case "file":
|
||
return NewFileStore(file)
|
||
case "redis":
|
||
return NewRedisStore(redisURL)
|
||
default:
|
||
return nil, fmt.Errorf("checkpoint: unknown driver %q", driver)
|
||
}
|
||
}
|