sqlitebroker

package
v0.2.0 Latest Latest
Warning

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

Go to latest
Published: Jul 16, 2026 License: Apache-2.0 Imports: 13 Imported by: 0

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

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

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

Open 打开(不存在则创建)path 指向的 sqlite 库文件并建表。 连接参数:WAL 日志、busy_timeout 5s、synchronous NORMAL、事务一律 BEGIN IMMEDIATE。 返回的 Broker 用之前必须先 Init(由 taskgate.New(cfg) 统一调用)。

func (*Broker) Ack

func (b *Broker) Ack(ctx context.Context, id, leaseToken string, result []byte) error

Ack 成功完结:completed + Result + FinishedAt,并在同一事务内唤醒子任务。 testHookBeforeAckCommit 在全部写入之后、提交之前触发(崩溃专项测试用,默认 nil)。

func (*Broker) Cancel

func (b *Broker) Cancel(ctx context.Context, id string) error

Cancel 取消:排队类状态(blocked/pending/retrying)直接 canceled 并同事务传播; running 只打 cancel_requested 标记(终态由 FinishCanceled 落库); 终态报 ErrAlreadyFinal,不存在报 ErrTaskNotFound。

func (*Broker) Close

func (b *Broker) Close() error

Close 关库:标记 closed 并广播,让阻塞中的 Dequeue 尽快退出,再关连接池。

func (*Broker) Counts

func (b *Broker) Counts(ctx context.Context) (map[string]map[taskgate.Status]int64, error)

Counts 出现过的 Type×Status 稀疏矩阵,和逐个 Get 汇总必须一致(brokertest 验证)。

func (*Broker) Dequeue

func (b *Broker) Dequeue(ctx context.Context, queues []string) (*taskgate.Task, error)

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

func (b *Broker) Enqueue(ctx context.Context, t *taskgate.Task) error

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

func (b *Broker) FinishCanceled(ctx context.Context, id, leaseToken string) error

FinishCanceled worker 响应取消后收尾:running → canceled 落库并同事务传播。

func (*Broker) Get

func (b *Broker) Get(ctx context.Context, id string) (*taskgate.Task, error)

Get 取单个任务。scanRec 扫出来的就是全新副本,调用方改了不影响存储。

func (*Broker) Heartbeat

func (b *Broker) Heartbeat(ctx context.Context, id, leaseToken string) error

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

func (b *Broker) List(ctx context.Context, f taskgate.Filter) ([]*taskgate.Task, error)

List 按 Filter 过滤,零值字段不过滤;先过滤 → ORDER BY (created_at, id) 升序 → LIMIT/OFFSET 分页(排序分页合同见 broker-contract.md,Offset 越界返回空)。

func (*Broker) QueueLen

func (b *Broker) QueueLen(ctx context.Context, queue string) (int, error)

QueueLen 队列积压:status∈{pending,retrying} 的数量(不看 RunAt 到没到点)。

func (*Broker) QueueQuota

func (b *Broker) QueueQuota(queue string, qc taskgate.QueueConfig) (taskgate.QuotaGate, error)

QueueQuota 构造该队列的配额闸(taskgate.QuotaProvider 能力实现)。

func (*Broker) ReapExpired

func (b *Broker) ReapExpired(ctx context.Context) (int, error)

ReapExpired 回收过期租约(lease_until < now 严格小于,压线不算过期): 带 cancel_requested=1 的过期任务直接落 canceled(不占 LeaseLost,触发传播), 其余一条 UPDATE ... RETURNING 原子完成 LeaseLost+1、封顶进 failed(固定文案)或回 pending、 清令牌;翻 failed 的行在同一事务内触发连锁传播。顺带做防御性修复: blocked 但父实际全部终态的任务,按提交时同一套决策函数补唤醒/补取消 (这不是正常路径,是给"唤醒中途崩"这类事故兜底)。返回值只算租约回收条数。

func (*Broker) Replay

func (b *Broker) Replay(ctx context.Context, req taskgate.ReplayRequest) (*taskgate.Task, error)

Replay 重放一次终态执行(spec 005):定位目标、校验前置条件(终态/未被重放/ completed 需显式允许)、创建新执行,全部在同一个事务内原子完成——并发同目标 重放恰好一个成功(BEGIN IMMEDIATE 串行化写者,uq_replay_of 唯一索引兜底)。 目标行零改写;新执行沿用目标的 Type/Queue/MaxRetry/OnParentFailure 与 BusinessKey, 三计数清零、无依赖、pending 落库。

func (*Broker) Requeue

func (b *Broker) Requeue(ctx context.Context, id, leaseToken string) error

Requeue 优雅停机时归还任务:running → pending,三计数与 RunAt 全不动, 清租约和取消标记。这不算失败,一个计数都不许占(合同要求)。

Jump to

Keyboard shortcuts

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