Documentation
¶
Overview ¶
Package jobstore 为 ofd-server 提供基于 bbolt 的任务持久化与队列。
选 bbolt 的原因:单文件、事务是 ACID 的、纯 Go 不需要 cgo,适合与本项目现有 构建矩阵共存。代价是它对数据库文件加排他锁,**同一文件只能由一个进程打开**, 因此本服务按单节点设计;需要多副本时必须换成 SQLite 或外部队列。
四类数据分四个 bucket:
meta schema_version jobs jobID → JSON 任务记录 queue_fast 快速路径的待办,按创建时间排序 queue_heavy 重路径的待办(Office/HTML 走 LibreOffice/Chrome) jobs_by_state <state><创建时间><jobID>,按状态查与清理
队列键与状态索引都用大端时间戳,保证 bbolt 的字节序比较等价于时间序。 所有会同时改动多个 bucket 的操作都在单个写事务内完成:bbolt 串行化写事务, 因此多个 worker 并发取件是安全的。
Index ¶
- type Delivery
- type DeliveryState
- type Job
- type Lane
- type Order
- type Output
- type State
- type StatsPair
- type StatsSnapshot
- type StatsTriple
- type Store
- func (s *Store) Cancel(id string) error
- func (s *Store) Claim(lane Lane) (*Job, error)
- func (s *Store) ClaimDueDeliveries(now time.Time, limit int) ([]Delivery, error)
- func (s *Store) Close() error
- func (s *Store) Compact() error
- func (s *Store) CompleteDelivery(id string) error
- func (s *Store) Count() (map[State]int, error)
- func (s *Store) Enqueue(job *Job) error
- func (s *Store) FailDelivery(id, reason string, next time.Duration) error
- func (s *Store) Finish(id string, state State, output Output, failure string) error
- func (s *Store) Get(id string) (*Job, error)
- func (s *Store) ListByState(state State, limit int) ([]Job, error)
- func (s *Store) ListDeliveries(state DeliveryState, limit int) ([]Delivery, error)
- func (s *Store) Prune(cutoff time.Time) (int, error)
- func (s *Store) PruneDeliveries(cutoff time.Time) (int, error)
- func (s *Store) QueueDepth() (map[State]int, error)
- func (s *Store) RecordDirect(job *Job) error
- func (s *Store) RecordUsage(id string, fromActual string, inputBytes int64) error
- func (s *Store) Recover() (int, error)
- func (s *Store) RecoverInflight() (int, error)
- func (s *Store) ScheduleDeliveries(items []Delivery) error
- func (s *Store) Snapshot() (StatsSnapshot, error)
Constants ¶
This section is empty.
Variables ¶
This section is empty.
Functions ¶
This section is empty.
Types ¶
type Delivery ¶
type Delivery struct {
// ID 全局唯一,同时作为回调里的 X-OFD-Delivery,供接收方去重。
ID string `json:"id"`
// JobID 关联的转换任务,仅供排障时串联。
JobID string `json:"job_id,omitempty"`
// Target 是预注册的通知目标名,不是 URL——URL 与密钥都在服务端侧配置,
// 调用方碰不到,避免被当作伪造回调的凭据。
Target string `json:"target"`
// Event 为 "succeeded" 或 "failed"。
Event string `json:"event"`
// Payload 是回调请求体,已序列化好的 JSON。
Payload json.RawMessage `json:"payload"`
// Attempts 已经尝试的次数。
Attempts int `json:"attempts"`
// MaxAttempts 超过后转为 abandoned。
MaxAttempts int `json:"max_attempts"`
// DueAt 下次尝试时间;重试时按退避往后推。
DueAt time.Time `json:"due_at"`
// State 见 DeliveryState。
State DeliveryState `json:"state"`
// LastError 最近一次失败原因。
LastError string `json:"last_error,omitempty"`
}
Delivery 是一条待投递的通知。
type DeliveryState ¶
type DeliveryState string
DeliveryState 是通知投递的状态。
const ( // DeliveryPending 等待到期投递。 DeliveryPending DeliveryState = "pending" // DeliveryInflight 已取出、正在投递。进程崩溃时该记录会被 RecoverInflight // 重新排回 pending——投递语义是 at-least-once,重复投递由接收方按 // X-OFD-Delivery 去重。 DeliveryInflight DeliveryState = "inflight" // DeliveryDone 已确认(收到 2xx),不再重试。 DeliveryDone DeliveryState = "done" // DeliveryAbandoned 超过最大重试次数,放弃投递等待人工处理。 DeliveryAbandoned DeliveryState = "abandoned" )
type Job ¶
type Job struct {
ID string `json:"id"`
State State `json:"state"`
Lane Lane `json:"lane"`
Order Order `json:"order"`
// From 是提交时调用方声明的输入格式,可能为空——大多数请求靠服务端
// 按魔数与文件名推断,声明值只是覆盖手段。
From string `json:"from"`
// FromActual 是服务端实际判定的输入格式,成功后回写。
//
// 不能拿 From 当统计口径:README 里写明输入格式由服务端解析、不采信
// 调用方的文件名,而 From 恰好就是那个声明值。一个 OFD 内容配了
// .html 名字、按 html 导入器转换的任务,在 From 里会记成 ofd→pdf。
// 空表示尚未判定(任务失败或仍在排队)。
FromActual string `json:"from_actual,omitempty"`
To string `json:"to"`
// InputBytes 是输入字节数。
//
// 提交时先记已收到的部分(上传的字节,URL 输入为 0),转换成功后用
// 实际落盘的尺寸覆盖,所以 URL 输入失败时这一项是 0——不是漏记,
// 是那时确实没有拿到完整内容。
InputBytes int64 `json:"input_bytes,omitempty"`
// Request 是提交时的原始请求 JSON。
Request json.RawMessage `json:"request,omitempty"`
// Output 是终态时的结果位置。
Output Output `json:"output,omitempty"`
// Error 是失败原因,终态时对调用方可见。
Error string `json:"error,omitempty"`
CreatedAt time.Time `json:"created_at"`
// DueAt 是任务最早可被取走的时刻。为零表示立即可取;重试时用它实现退避,
// 靠把它塞进 delayed 队列实现,而不是靠 sleep——sleep 会占住一个 worker。
DueAt time.Time `json:"due_at,omitempty"`
StartedAt time.Time `json:"started_at,omitempty"`
FinishedAt time.Time `json:"finished_at,omitempty"`
Attempt int `json:"attempt"`
NotifyTarget string `json:"notify_target,omitempty"`
// NotifyEvents 限定要通知的事件,空表示用目标的默认集合。
NotifyEvents []string `json:"notify_events,omitempty"`
}
Job 是一次转换任务的持久化记录。Request 保存原始请求体,worker 取出后据此 重建 Source/Sink——重试时需要能完全重放。
type Output ¶
type Output struct {
Kind string `json:"kind"`
Path string `json:"path,omitempty"`
Bucket string `json:"bucket,omitempty"`
Key string `json:"key,omitempty"`
URL string `json:"url,omitempty"`
Size int64 `json:"size"`
// Files 是产物的文件名清单,逐页输出会有多项。
//
// 单有这个字段是不够的:调用方能指定 output.filename,但产物名由格式扩展名
// 决定,逐页输出还会带页号(thumb-0001.png)。只报目录等于让调用方还得去
// 猜文件名或者列目录——那正是"不然用户如何知道"要解决的问题。
Files []string `json:"files,omitempty"`
}
Output 描述结果位置,与 transfer.Location 字段对应,但不直接复用——任务记录要 保持结构稳定,不随传输层实现变化。
type StatsSnapshot ¶
type StatsSnapshot struct {
// Since 是第一次记账的时间。没有转换发生过时为零值。
//
// 没有它,调用方只看到一个总数却不知道该按什么时间跨度算速率——
// "1234 次转换"在一周和一年里是完全不同的两件事。
Since time.Time `json:"since"`
// Jobs 按 输入格式 -> 输出格式 -> 状态 索引的转换次数。
Jobs map[StatsTriple]uint64
// InputBytes 按输入格式索引的输入字节总量。
InputBytes map[string]uint64
// OutputBytes 按 输入格式 -> 输出格式 索引的输出字节总量。
OutputBytes map[StatsPair]uint64
}
StatsSnapshot 是某一时刻的全部计数。
type StatsTriple ¶
StatsTriple 是转换计数的三个标签。
type Store ¶
type Store struct {
// contains filtered or unexported fields
}
Store 是任务存储。
func (*Store) Cancel ¶
Cancel 取消仍在排队的任务。
只允许取消 queued:已经 running 的任务正占着外部进程,中途杀掉它只会留下 半截输出,不如让它跑完再由调用方决定要不要丢弃结果。
func (*Store) Claim ¶
Claim 从指定通道取出队首任务并标记为 running。
返回 (nil, nil) 表示队列为空。取件与状态改写都在同一个写事务内完成, 因此多个 worker 并发调用不会取到同一个任务。 到期的延迟任务会先被搬进通道队列,因此重试任务到了 DueAt 就能被取走。
func (*Store) ClaimDueDeliveries ¶
ClaimDueDeliveries 取出到期且待投递的通知,并就地标记为 inflight。
标记与取出在同一个写事务内完成,因此多个投递协程并发调用不会取到同一条。 取件不改变 DueAt,所以索引键保持不变,无需重排。 limit 为 0 表示取全部到期记录。
func (*Store) Compact ¶
Compact 把数据库收缩到实际使用量并原子替换原文件。
bbolt 删除键后只把页标记为空闲,文件不会自动变小;任务记录不断增删时需要定期 收缩,否则文件会一直涨。
Compact(dst, src *DB) 需要两个同时打开的句柄,而 bbolt 对可写文件加的是排他 锁,所以流程是:关掉当前句柄 → 用只读句柄重新打开源(共享锁)→ 压缩到临时 文件 → 替换原文件 → 重新以可写方式打开。任何一步失败都要把原库重新打开, 否则后续调用会因连接已关闭而全部失败。
func (*Store) CompleteDelivery ¶
CompleteDelivery 标记投递成功。
func (*Store) FailDelivery ¶
FailDelivery 记录一次失败。已达最大重试次数时转为 abandoned,否则按退避重新 排回 pending。
func (*Store) Finish ¶
Finish 写入终态。state 必须是终态,且任务当前处于 running。
不允许从 queued 直接置终态:那说明有人绕过 worker 直接改了结果, 与"只有 worker 能决定成败"的约定冲突。
func (*Store) ListByState ¶
ListByState 按状态列出任务,limit 为 0 表示不限。用于管理接口与排障。
func (*Store) ListDeliveries ¶
func (s *Store) ListDeliveries(state DeliveryState, limit int) ([]Delivery, error)
ListDeliveries 按状态列出投递记录,limit 为 0 表示不限。用于管理接口与排障。
func (*Store) Prune ¶
Prune 删除早于 cutoff 的终态任务,返回删除数量。运行中与排队中的任务不受影响。 bbolt 的 mmap 只增不减:不断产生短命任务记录会让文件一直变大,因此需要定期 清理(配合 Store.Compact 做文件收缩)。
func (*Store) PruneDeliveries ¶
PruneDeliveries 删除早于 cutoff 的已完成投递,返回删除数量。
abandoned 不在此清理范围内:它们代表"通知没送出去",需要人工介入确认,抹掉 记录就看不出曾经漏过。
func (*Store) QueueDepth ¶
QueueDepth 报告各非终态的任务数量。
用状态索引的前缀扫描实现,成本是 O(该状态的任务数)而不是 O(全部任务): bucketByState 的键首字节是状态码,同一状态的任务在 B+ 树里连续,所以 Seek 之后遇到前缀变化即可停止。终态任务即使积到百万条也不影响这里的耗时。
刻意不维护"计数器再自增自减":那需要在每个状态迁移点都记得同步, 漏一处就是永久性错误数字,而漏了不会报错。扫描慢一点,但读到的数 一定对。
func (*Store) RecordDirect ¶
RecordDirect 直接写入一条任务记录,不进任何队列。
供同步转换使用:这类转换由 HTTP 请求自己执行,没有 worker 参与,所以不能 经 Enqueue——那会把任务丢进通道队列,runner 就会把它领走再转换一遍, 等于同一个任务跑两次。
状态直接置为 running 而不是 queued:任务确实在执行中,而且这样它会出现在 队列深度的 running 计数里,运维能看到"有一个同步转换在跑"。
func (*Store) RecordUsage ¶
RecordUsage 回写任务的实际输入格式与输入字节数。
单独于 Finish 是因为这两个字段只有转换成功后才拿得到,而 Finish 会被 重试、取消、关停等路径调用。分开写让"终态"这个语义保持干净:一个方法 只管状态迁移,一个只管用量记账。
func (*Store) RecoverInflight ¶
RecoverInflight 把上次进程退出时停在 inflight 的投递排回 pending。
与任务的 Recover 同一个理由:这些通知没有送达,但外部无法察觉。at-least-once 语义下重投是安全的,接收方按投递 ID 去重。
func (*Store) ScheduleDeliveries ¶
ScheduleDeliveries 写入待投递通知。同一批写入在一个事务内完成。