jobstore

package
v0.1.4 Latest Latest
Warning

This package is not in the latest version of its module.

Go to latest
Published: Oct 4, 2026 License: Apache-2.0 Imports: 10 Imported by: 0

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

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 Lane

type Lane string

Lane 区分任务队列。快路径与重路径分开,是为了避免一个 LibreOffice 请求 把 Markdown/文本这类毫秒级任务一起堵住。

const (
	// LaneFast 是快路径:OFD、PDF、Markdown、文本与图片等纯 Go 转换。
	LaneFast Lane = "fast"
	// LaneHeavy 是重路径:Office 文档经 LibreOffice、HTML/MHTML 经 Chrome。
	LaneHeavy Lane = "heavy"
)

func Lanes

func Lanes() []Lane

Lanes 列出全部队列名,供管理接口展示积压情况。

type Order

type Order int

Order 是队列的出队顺序。同一优先级下按创建时间先进先出。

const (
	// OrderNormal 是普通任务。
	OrderNormal Order = iota
	// OrderHigh 让任务插到普通任务之前。预留给交互式转换。
	OrderHigh
)

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 State

type State string

State 是任务状态。

const (
	StateQueued    State = "queued"
	StateRunning   State = "running"
	StateSucceeded State = "succeeded"
	StateFailed    State = "failed"
	StateCancelled State = "cancelled"
)

func (State) Terminal

func (s State) Terminal() bool

Terminal 报告该状态是否已终结。终结的任务不再回到队列。

type StatsPair

type StatsPair struct {
	From string
	To   string
}

StatsPair 是字节计数的两个标签。

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

type StatsTriple struct {
	From  string
	To    string
	State string
}

StatsTriple 是转换计数的三个标签。

type Store

type Store struct {
	// contains filtered or unexported fields
}

Store 是任务存储。

func Open

func Open(path string) (*Store, error)

Open 打开或创建指定路径的存储。

func (*Store) Cancel

func (s *Store) Cancel(id string) error

Cancel 取消仍在排队的任务。

只允许取消 queued:已经 running 的任务正占着外部进程,中途杀掉它只会留下 半截输出,不如让它跑完再由调用方决定要不要丢弃结果。

func (*Store) Claim

func (s *Store) Claim(lane Lane) (*Job, error)

Claim 从指定通道取出队首任务并标记为 running。

返回 (nil, nil) 表示队列为空。取件与状态改写都在同一个写事务内完成, 因此多个 worker 并发调用不会取到同一个任务。 到期的延迟任务会先被搬进通道队列,因此重试任务到了 DueAt 就能被取走。

func (*Store) ClaimDueDeliveries

func (s *Store) ClaimDueDeliveries(now time.Time, limit int) ([]Delivery, error)

ClaimDueDeliveries 取出到期且待投递的通知,并就地标记为 inflight。

标记与取出在同一个写事务内完成,因此多个投递协程并发调用不会取到同一条。 取件不改变 DueAt,所以索引键保持不变,无需重排。 limit 为 0 表示取全部到期记录。

func (*Store) Close

func (s *Store) Close() error

Close 关闭数据库。

func (*Store) Compact

func (s *Store) Compact() error

Compact 把数据库收缩到实际使用量并原子替换原文件。

bbolt 删除键后只把页标记为空闲,文件不会自动变小;任务记录不断增删时需要定期 收缩,否则文件会一直涨。

Compact(dst, src *DB) 需要两个同时打开的句柄,而 bbolt 对可写文件加的是排他 锁,所以流程是:关掉当前句柄 → 用只读句柄重新打开源(共享锁)→ 压缩到临时 文件 → 替换原文件 → 重新以可写方式打开。任何一步失败都要把原库重新打开, 否则后续调用会因连接已关闭而全部失败。

func (*Store) CompleteDelivery

func (s *Store) CompleteDelivery(id string) error

CompleteDelivery 标记投递成功。

func (*Store) Count

func (s *Store) Count() (map[State]int, error)

Count 统计各状态的任务数量。

func (*Store) Enqueue

func (s *Store) Enqueue(job *Job) error

Enqueue 写入新任务并放入对应队列。与状态索引在同一事务内更新。

func (*Store) FailDelivery

func (s *Store) FailDelivery(id, reason string, next time.Duration) error

FailDelivery 记录一次失败。已达最大重试次数时转为 abandoned,否则按退避重新 排回 pending。

func (*Store) Finish

func (s *Store) Finish(id string, state State, output Output, failure string) error

Finish 写入终态。state 必须是终态,且任务当前处于 running。

不允许从 queued 直接置终态:那说明有人绕过 worker 直接改了结果, 与"只有 worker 能决定成败"的约定冲突。

func (*Store) Get

func (s *Store) Get(id string) (*Job, error)

func (*Store) ListByState

func (s *Store) ListByState(state State, limit int) ([]Job, error)

ListByState 按状态列出任务,limit 为 0 表示不限。用于管理接口与排障。

func (*Store) ListDeliveries

func (s *Store) ListDeliveries(state DeliveryState, limit int) ([]Delivery, error)

ListDeliveries 按状态列出投递记录,limit 为 0 表示不限。用于管理接口与排障。

func (*Store) Prune

func (s *Store) Prune(cutoff time.Time) (int, error)

Prune 删除早于 cutoff 的终态任务,返回删除数量。运行中与排队中的任务不受影响。 bbolt 的 mmap 只增不减:不断产生短命任务记录会让文件一直变大,因此需要定期 清理(配合 Store.Compact 做文件收缩)。

func (*Store) PruneDeliveries

func (s *Store) PruneDeliveries(cutoff time.Time) (int, error)

PruneDeliveries 删除早于 cutoff 的已完成投递,返回删除数量。

abandoned 不在此清理范围内:它们代表"通知没送出去",需要人工介入确认,抹掉 记录就看不出曾经漏过。

func (*Store) QueueDepth

func (s *Store) QueueDepth() (map[State]int, error)

QueueDepth 报告各非终态的任务数量。

用状态索引的前缀扫描实现,成本是 O(该状态的任务数)而不是 O(全部任务): bucketByState 的键首字节是状态码,同一状态的任务在 B+ 树里连续,所以 Seek 之后遇到前缀变化即可停止。终态任务即使积到百万条也不影响这里的耗时。

刻意不维护"计数器再自增自减":那需要在每个状态迁移点都记得同步, 漏一处就是永久性错误数字,而漏了不会报错。扫描慢一点,但读到的数 一定对。

func (*Store) RecordDirect

func (s *Store) RecordDirect(job *Job) error

RecordDirect 直接写入一条任务记录,不进任何队列。

供同步转换使用:这类转换由 HTTP 请求自己执行,没有 worker 参与,所以不能 经 Enqueue——那会把任务丢进通道队列,runner 就会把它领走再转换一遍, 等于同一个任务跑两次。

状态直接置为 running 而不是 queued:任务确实在执行中,而且这样它会出现在 队列深度的 running 计数里,运维能看到"有一个同步转换在跑"。

func (*Store) RecordUsage

func (s *Store) RecordUsage(id string, fromActual string, inputBytes int64) error

RecordUsage 回写任务的实际输入格式与输入字节数。

单独于 Finish 是因为这两个字段只有转换成功后才拿得到,而 Finish 会被 重试、取消、关停等路径调用。分开写让"终态"这个语义保持干净:一个方法 只管状态迁移,一个只管用量记账。

func (*Store) Recover

func (s *Store) Recover() (int, error)

Recover 把上次进程退出时停留在 running 的任务重新入队。

这类任务的转换进程已经随进程一起消失,不重排就会永远卡住。启动时调用一次。

func (*Store) RecoverInflight

func (s *Store) RecoverInflight() (int, error)

RecoverInflight 把上次进程退出时停在 inflight 的投递排回 pending。

与任务的 Recover 同一个理由:这些通知没有送达,但外部无法察觉。at-least-once 语义下重投是安全的,接收方按投递 ID 去重。

func (*Store) ScheduleDeliveries

func (s *Store) ScheduleDeliveries(items []Delivery) error

ScheduleDeliveries 写入待投递通知。同一批写入在一个事务内完成。

func (*Store) Snapshot

func (s *Store) Snapshot() (StatsSnapshot, error)

Snapshot 读取全部累计计数。

Jump to

Keyboard shortcuts

? : This menu
/ : Search site
f or F : Jump to
y or Y : Canonical URL