Documentation
¶
Overview ¶
Package sqlitebroker 是 Broker 的 sqlite 文件后端:纯 Go 驱动(modernc.org/sqlite,免 cgo), WAL 模式单文件落盘。所有"终态更新 + 子任务唤醒/连锁取消"都在同一个事务里完成(宪法 III), 语义以 memorybroker 为基准,由 brokertest 的 18 条契约统一验收。
并发模型:连接池收紧到 1 个连接(单进程内所有读写串行),配合 WAL + busy_timeout, 避免 SQLITE_BUSY;跨进程的写入靠 Dequeue 的 100ms 兜底轮询发现。
Index ¶
- type Broker
- func (b *Broker) Ack(ctx context.Context, id, leaseToken string, result []byte) error
- func (b *Broker) Cancel(ctx context.Context, id string) error
- func (b *Broker) Close() error
- func (b *Broker) Counts(ctx context.Context) (map[string]map[taskgate.Status]int64, error)
- func (b *Broker) Dequeue(ctx context.Context, queues []string) (*taskgate.Task, error)
- func (b *Broker) Enqueue(ctx context.Context, t *taskgate.Task) error
- func (b *Broker) Fail(ctx context.Context, id, leaseToken, errMsg string, kind taskgate.FailKind, ...) error
- func (b *Broker) FinishCanceled(ctx context.Context, id, leaseToken string) error
- func (b *Broker) Get(ctx context.Context, id string) (*taskgate.Task, error)
- func (b *Broker) Heartbeat(ctx context.Context, id, leaseToken string) error
- func (b *Broker) Init(opts taskgate.BrokerOptions) error
- func (b *Broker) List(ctx context.Context, f taskgate.Filter) ([]*taskgate.Task, error)
- func (b *Broker) QueueLen(ctx context.Context, queue string) (int, error)
- func (b *Broker) QueueQuota(queue string, qc taskgate.QueueConfig) (taskgate.QuotaGate, error)
- func (b *Broker) ReapExpired(ctx context.Context) (int, error)
- func (b *Broker) Replay(ctx context.Context, req taskgate.ReplayRequest) (*taskgate.Task, error)
- func (b *Broker) Requeue(ctx context.Context, id, leaseToken string) error
Constants ¶
This section is empty.
Variables ¶
This section is empty.
Functions ¶
This section is empty.
Types ¶
type Broker ¶
type Broker struct {
// contains filtered or unexported fields
}
Broker sqlite 后端。实现 taskgate.Broker。
func Open ¶
Open 打开(不存在则创建)path 指向的 sqlite 库文件并建表。 连接参数:WAL 日志、busy_timeout 5s、synchronous NORMAL、事务一律 BEGIN IMMEDIATE。 返回的 Broker 用之前必须先 Init(由 taskgate.New(cfg) 统一调用)。
func (*Broker) Ack ¶
Ack 成功完结:completed + Result + FinishedAt,并在同一事务内唤醒子任务。 testHookBeforeAckCommit 在全部写入之后、提交之前触发(崩溃专项测试用,默认 nil)。
func (*Broker) Cancel ¶
Cancel 取消:排队类状态(blocked/pending/retrying)直接 canceled 并同事务传播; running 只打 cancel_requested 标记(终态由 FinishCanceled 落库); 终态报 ErrAlreadyFinal,不存在报 ErrTaskNotFound。
func (*Broker) Dequeue ¶
Dequeue 阻塞认领:直到某队列出现"status∈{pending,retrying} 且 run_at≤now"的任务, 或 ctx 取消(返回 ctx.Err())。循环结构:试认领一次 → 无果则挂起等待,三个唤醒源:
- 同进程写入踢的内部信号(Enqueue/Ack/Fail/... 后 wakeAll);
- 注入 clock 的到点信号:等 min(100ms, 最近的 run_at - now),延迟任务到点自动醒, 100ms 同时兜底跨进程写入(fakeclock 下不推时间就纯挂起,不空转);
- ctx 取消。
func (*Broker) Enqueue ¶
Enqueue 入队。同一个事务内完成:ID 查重、父任务存在性校验、初始状态判定 (DecideOnSubmit)、tasks 与 task_deps 落库;生成的 ID 与判定结果回填到 *t。
func (*Broker) Fail ¶
func (b *Broker) Fail(ctx context.Context, id, leaseToken, errMsg string, kind taskgate.FailKind, retryAt time.Time) error
Fail 失败路径:按 FailKind 动对应计数,封顶或耗尽进 failed(同事务连锁传播), 否则进 retrying、RunAt=retryAt 到点重跑。
func (*Broker) FinishCanceled ¶
FinishCanceled worker 响应取消后收尾:running → canceled 落库并同事务传播。
func (*Broker) Heartbeat ¶
Heartbeat 续租:lease_until = now + TTL。发现取消标记时**续租照做**(先提交), 再返回 ErrTaskCanceled 提醒 scheduler 去 cancel handler 的 ctx。
func (*Broker) Init ¶
func (b *Broker) Init(opts taskgate.BrokerOptions) error
Init 装配运行参数,零值补默认(TTL 60s / LeaseLostMax 3 / ThrottledMax 100 / 真时钟)。
func (*Broker) List ¶
List 按 Filter 过滤,零值字段不过滤;先过滤 → ORDER BY (created_at, id) 升序 → LIMIT/OFFSET 分页(排序分页合同见 broker-contract.md,Offset 越界返回空)。
func (*Broker) QueueQuota ¶
QueueQuota 构造该队列的配额闸(taskgate.QuotaProvider 能力实现)。
func (*Broker) ReapExpired ¶
ReapExpired 回收过期租约(lease_until < now 严格小于,压线不算过期): 带 cancel_requested=1 的过期任务直接落 canceled(不占 LeaseLost,触发传播), 其余一条 UPDATE ... RETURNING 原子完成 LeaseLost+1、封顶进 failed(固定文案)或回 pending、 清令牌;翻 failed 的行在同一事务内触发连锁传播。顺带做防御性修复: blocked 但父实际全部终态的任务,按提交时同一套决策函数补唤醒/补取消 (这不是正常路径,是给"唤醒中途崩"这类事故兜底)。返回值只算租约回收条数。