taskgate

package module
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: 9 Imported by: 0

README

taskgate

CI PkgGoDev Go Report Card Go Version License

English | 简体中文

taskgate 是一个 go get 即用的 Go 任务排队限流库:排队、限流、重试、依赖、取消、优雅停止,一套接口五种后端。

它是库不是服务——没有 Web UI、不读环境变量和配置文件,只吃你传进来的 Config。典型场景:调用 LLM / OCR 这类有配额的外部网关时,把请求收进队列,按类型隔离限流,失败自动退避重试,进程崩了靠租约把任务捞回来。

特性

  • 类型级限流:每个队列独立的 {Workers, RPS, Burst},慢队列绝不拖累快队列;Routes 支持多类型共享同一队列(共享同一个网关配额)。
  • 周期配额(硬配额):{QuotaLimit, QuotaPeriod} 限"每个固定窗口最多启动 N 次 handler"——频率和配额是两码事(见下文);跨进程共享计数、绝不超发,任何故障只少放行不多放;后端不支持共享计数时 New() 直接报错,没有静默降级。
  • 五种后端一套契约:memorybroker(内存,零依赖)、sqlitebroker(单文件落盘,纯 Go 免 cgo)、redisbroker(多进程共享,Lua 原子流转)、pgbroker(PostgreSQL)、mysqlbroker(MySQL)——后两个是服务器型数据库后端,多进程共享,认领互斥靠 FOR UPDATE SKIP LOCKED(要求 MySQL 8.0+ / PostgreSQL 9.5+);同一套行为契约(brokertest 22 条)五后端验收。
  • 分布式限流(redis 后端):同一队列的 Workers/RPS 配额在所有接同一 Redis 的进程间共享,加机器不等于把网关打爆;进程崩溃占着的并发槽按租约自动回收。
  • 租约回收:认领即加租约,心跳自动续租;worker 崩溃后 reaper 把任务捞回重跑,毒任务封顶进死信;长任务可在 handler 里 RenewLease 手动续租,也可按队列关掉自动心跳改全手动(ManualHeartbeat)。代价是 at-least-once——崩溃回收会让 handler 重跑,所以 handler 必须幂等,见handler 的执行契约
  • 重试三计数分工:Attempts 管业务失败(指数退避 min(2^n×1s,10min)±20%,超 MaxRetry 进 failed)、Throttled 管网关限流(ErrThrottled 不占重试次数)、LeaseLost 管崩溃回收;ErrSkipRetry 直接死信。
  • 依赖流水线:DependsOn 串联/扇入,父任务终态与子任务唤醒在同一事务内完成,不丢唤醒;父失败默认连锁取消(可 IgnoreParentFailure)。
  • 取消:pending/blocked 直接取消并向下传播;running 任务的 handler ctx 被即时 cancel。
  • 优雅停止:Shutdown(ctx) 等在跑任务善终;超时则打断并把任务原样归还(不占任何计数),部署重启不消耗任务配额。
  • 可观测:Get / List / Stats / Overview / Wait 查询等待,OnStateChange 回调埋点;List 支持 Offset+Limit 稳定分页。

支持的后端

五种后端实现同一套 Broker 合同,由同一套 brokertest(22 条)验收,换后端只改构造一行、业务代码零改动。按部署形态挑一个:

后端 适用 依赖 持久化 多进程共享
memorybroker 单进程、测试、临时任务 否(进程退出即丢)
sqlitebroker 单机、要落盘、免 cgo modernc.org/sqlite(纯 Go) 单文件 同机多进程(文件锁)
redisbroker 多进程/多机、认领互斥 Redis + go-redis Redis 是(Lua 原子流转)
pgbroker 多进程/多机、已有 PG 库 github.com/jackc/pgx/v5(纯 Go) PostgreSQL 是(FOR UPDATE SKIP LOCKED)
mysqlbroker 多进程/多机、已有 MySQL 库 github.com/go-sql-driver/mysql(纯 Go) MySQL 是(FOR UPDATE SKIP LOCKED)

Redis 后端面向单实例 / 主从 / 哨兵拓扑,不支持 Redis Cluster(脚本内自行拼键、键不带 hash tag)。

安装

taskgate 需要 Go 1.25+ 并启用 modules。先初始化你的 module:

go mod init github.com/my/repo

再拉取 taskgate:

go get github.com/AmbroseX/taskgate@latest   # 最新正式版
go get github.com/AmbroseX/taskgate@v0.1.0   # 锁定某个版本

版本号遵循 语义化版本。当前处于 v0 阶段,API 尚未冻结,小版本(minor)升级可能带来不兼容改动,升级前请看发布说明。你的依赖锁在自己的 go.mod + go.sum 里,不主动 go get -u 就不会被新版本影响。

运行期依赖:modernc.org/sqlite(纯 Go)、golang.org/x/time/rategithub.com/oklog/ulid/v2github.com/redis/go-redis/v9github.com/go-redis/redis_rate/v10github.com/jackc/pgx/v5(纯 Go 免 cgo)、github.com/go-sql-driver/mysql(纯 Go 免 cgo);测试依赖 github.com/alicebob/miniredis/v2

快速开始

三级流水线:检索 → 生成 → 打分(完整可跑版本见 examples/llm,go run ./examples/llm)。

package main

import (
    "context"
    "encoding/json"
    "time"

    "github.com/AmbroseX/taskgate"
    "github.com/AmbroseX/taskgate/memorybroker"
)

func main() {
    g, err := taskgate.New(taskgate.Config{
        Broker: memorybroker.New(), // 或 sqlitebroker.Open("tasks.db") / redisbroker.New(...)
        Queues: map[string]taskgate.QueueConfig{
            "cpu": {Workers: 4},         // 本地轻活:4 并发不限速
            "llm": {Workers: 2, RPS: 3}, // 大模型网关:2 并发,每秒最多 3 个
        },
        Routes: map[string]string{ // 任务类型 → 队列
            "retrieve": "cpu", "generate": "llm", "score": "cpu",
        },
    })
    if err != nil {
        panic(err)
    }

    // 注册 handler:返回值写进 Result,返回错误决定重试路径。
    g.Handle("generate", func(ctx context.Context, t *taskgate.Task) ([]byte, error) {
        parent, _ := g.Get(ctx, t.DependsOn[0]) // 读上游的 Result
        _ = parent
        // ... 调大模型 ...
        return json.Marshal(map[string]string{"answer": "..."})
    })

    ctx := context.Background()

    // 起消费循环(阻塞,放 goroutine)。
    go g.Run(ctx)

    // 提交流水线:DependsOn 串联,父任务完成才唤醒子任务。
    rid, _ := g.Submit(ctx, "retrieve", nil)
    gid, _ := g.Submit(ctx, "generate", nil, taskgate.DependsOn(rid))
    sid, _ := g.Submit(ctx, "score", nil, taskgate.DependsOn(gid))
    _ = time.Second

    // 等最终结果;优雅停机。
    result, _ := g.Wait(ctx, sid)
    _ = result
    _ = g.Shutdown(ctx)
}

handler 的执行契约:必须幂等

taskgate 是 at-least-once(至少执行一次),不是 exactly-once。同一个任务的 handler 可能被执行多次,你的 handler 必须能容忍重跑。

这不是"出错了才会发生",而是正常运转的一部分。三条路径都会导致重跑:

路径 触发 计数
worker 进程崩溃(kill -9、OOM、断电) 心跳停 → 租约到期 → reaper 把任务捞回 LeaseLost
handler 卡死 + 队列开了 ManualHeartbeat 停止续租 → 租约到期被回收 LeaseLost
Shutdown(ctx) 超时打断 任务原样 Requeue 归还回 pending 三计数都不动(任务看起来"没跑过")

前两条崩溃发生时,第一跑可能已经跑到一半——LLM 已经调过了、钱已经花了、库可能已经写了一半。

一个会出事的例子
// ❌ 错的:重跑会重复扣额度、重复落库
g.Handle("summarize", func(ctx context.Context, t *taskgate.Task) ([]byte, error) {
    resp, err := callLLM(ctx, t.Payload) // ← 花钱了
    if err != nil {
        return nil, err
    }
    billing.Deduct(userID, resp.Tokens)  // ← 扣额度
    return db.Save(ctx, resp)            // ← 假设进程在这里崩了
})

进程崩在 db.Save 里 → 租约到期 → reaper 捞回重跑 → LLM 调了两次、额度扣了两次,而库里只有一份结果。

怎么防:三类东西要分开处理

重跑会伤到三类东西,防护手段各不相同,而且只有第一类能被本地事务彻底防住:

本地数据库副作用(扣额度、落库) 外部网关调用(调 LLM/OCR) 异步消息/通知(邮件、Webhook、短信)
能不能防住 ,用 Task.ID 做幂等键 取决于网关,taskgate 无能为力 只能防到"意图不丢不重",投递仍是 at-least-once
手段 数据库唯一约束 / 幂等表,和业务写入同一事务 网关支持 idempotency key 就传 Task.ID;不支持只能接受重复调用 事务性 outbox:把"要发这条通知"随业务事务落库;投递端/接收方仍要按幂等键去重

第三类容易被误当成第一类:outbox 只保证通知意图和业务数据同生共死,通知本身发出去这一步照样可能重复——邮件、Webhook 的接收方看到的还是 at-least-once,该带的幂等键(比如 Task.ID)一个都不能省。

本地副作用:用 Task.ID 原子保护

Task.ID 在重跑时不变,拿它当业务幂等键。关键是让"扣额度"和"落库"在同一个事务里、用同一个唯一键——分两步做就还是有窗口。

g.Handle("summarize", func(ctx context.Context, t *taskgate.Task) ([]byte, error) {
    resp, err := callLLM(ctx, t.Payload)
    if err != nil {
        return nil, err
    }
    // 扣额度 + 落库在同一事务内完成,以 t.ID 为唯一键;
    // 重跑时唯一约束冲突 → 返回首次的结果,不会二次扣费。
    return db.SaveOnce(ctx, t.ID, userID, resp)
})
外部调用:taskgate 消不掉这个窗口

上面这段不能阻止 LLM 被调用两次。进程如果崩在 callLLM 成功之后、db.SaveOnce 之前,重跑必然会再调一次网关——这一次的钱照花,只是不会重复扣用户额度、不会重复落库。

这个窗口是 at-least-once 的固有代价,处理方式只有两条:

  • 网关明确支持幂等键时,把 t.ID 作为幂等键传入,网关自己去重,重跑不产生第二次计费。这是唯一能真正消除重复调用的办法。 支付类网关普遍支持(如 Stripe 的 Idempotency-Key);LLM 网关是否支持、哪些端点支持,以网关官方文档为准,不要想当然。
  • 网关不支持:接受重复调用的成本,或在业务协议层做补偿(对账、按 t.ID 追溯去重)。taskgate 帮不上忙——"外部调用已成功、但本地还没落账"的不确定窗口,任何任务队列都消不掉。

WithBusinessKey 幂等的是"提交",不是"执行"。 同键重复 Submit 会返回 ErrTaskExists,保证同一件事只入队一次;但入队之后,这一个任务的 handler 仍然可能因为崩溃回收而跑多次。这是两件事,别混。(旧的 WithID 已弃用,现在是 WithBusinessKey 的别名——传入的值是业务键,不再是任务 ID。)

不要把关键副作用放进 OnStateChange 它是异步通知,不保证送达、不保证顺序、回调 panic 会被吞掉、也没有重试和"全局只消费一次"的保证。它只适合埋点和允许丢失/重复的观测逻辑,承载不了扣费、发通知、写业务终态。

为什么不做 exactly-once

这是分布式系统的老大难:handler 的执行和它的副作用落地不在同一个事务里,库层面给不了这个保证——任何"跑完就标记完成"的方案,都存在"跑完了、还没来得及标记就崩了"的窗口。能给的只有 at-least-once + 一个重跑时稳定的 Task.ID,让你在业务侧自己做幂等。

handler 的错误语义

handler 返回什么错,决定任务走哪条重试路径:

return nil, taskgate.ErrThrottled{RetryAfter: 30 * time.Second} // 网关限流:延后重排,不占重试次数
return nil, taskgate.ErrSkipRetry{Err: err}                     // 没救的错:直接进 failed 死信
return nil, err                                                 // 普通业务失败:指数退避重试

注意:ErrThrottled / ErrSkipRetry 必须按值返回(errors.As 按值匹配),不要返回其指针。

提交选项:WithBusinessKey(业务幂等键;WithID 已弃用为其别名)、Delay / RunAt(延迟执行)、MaxRetryDependsOnIgnoreParentFailure

长任务与手动续租

默认(自动档)scheduler 给每个在跑任务起自动心跳,每 LeaseTTL/3 续一次租,handler 什么都不用管。两种情况需要手动续租:

  • 自动档里想顺手多续一口:handler 在任务 ctx 上随时可调 taskgate.RenewLease(ctx),与自动心跳共用同一租约令牌,幂等延长,互不干扰。
  • 毒任务检测要更灵敏:自动心跳是调度器发的,handler 卡死心跳照跳,任务永远不会被回收。把队列的 ManualHeartbeat 设为 true 关掉自动心跳,handler 每处理完一个检查点自己续一次——卡死就停止续租,租约到期被 reaper 回收重跑。
Queues: map[string]taskgate.QueueConfig{
    "ocr": {Workers: 2, LeaseTTL: taskgate.Duration(60 * time.Second), ManualHeartbeat: true},
},

g.Handle("ocr", func(ctx context.Context, t *taskgate.Task) ([]byte, error) {
    for _, page := range pages {
        if err := taskgate.RenewLease(ctx); err != nil {
            return nil, err // ErrTaskCanceled / ErrLeaseLost 时 ctx 已被 cancel,尽快退出
        }
        // ... 处理一页 ...
    }
    return result, nil
})

RenewLease 的返回值语义:

返回 含义 handler 该做什么
nil 续租成功 继续干活
ErrTaskCanceled 任务被外部 Cancel(续租照做) ctx 已被 cancel,尽快退出
ErrLeaseLost 租约已丢(被 reaper 回收),结果注定作废 ctx 已被 cancel,立即放弃
ErrNoTask ctx 不是任务 ctx(handler 之外调用) 修代码
其他错误 网络抖动,租约没成也没丢 可稍后重试

手动档的跨进程取消语义:关掉自动心跳后,别的进程发起的 Cancel 只能由 handler 的下一次 RenewLease 发现(返回 ErrTaskCanceled);handler 一直不续租,则由租约过期兜底回收。同进程的 Cancel 不受影响,依然即时 cancel handler 的 ctx。手动档的语义就是"handler 对自己的租约负全责"。

List 分页

List 结果按 (CreatedAt, ID) 升序稳定排序(三后端一致),Filter 支持 Offset+Limit 翻页:先过滤 → 排序 → 跳过 Offset 条 → 取 Limit 条(0 = 不限);Offset 越界返回空列表不报错。

page2, _ := g.List(ctx, taskgate.Filter{Type: "ocr", Limit: 20, Offset: 20})

两条心里有数:

  • 弱一致翻页:翻页期间有任务入队/流转时不承诺快照一致,只承诺"未变动的任务不丢不重";要强一致游标等 M4 再议。
  • redis 后端的代价:List 走"索引集合求候选 → 逐个取回 → 内存排序切片",复杂度是 O(候选集) 而不是 O(页大小);大库存时先用 Filter 的 Type/Status/Queue 把候选集缩小再翻页。

Redis 后端(多进程)

多个 worker 进程接同一个 Redis 抢同一批任务时用它:Lua 原子流转保证每个任务在同一时刻只有一个有效租约,状态流转不会被并发撕裂;进程 kill -9 后在跑任务由租约回收重跑

注意这不等于"handler 绝不会并行重叠":租约失效后新 worker 可以重新认领,而旧 handler 可能因网络分区、进程暂停(STW/换页)或没及时响应 ctx 而还没退出,两份业务代码可能短暂并行。所以仍然是 at-least-once,handler 必须幂等——见handler 的执行契约

换后端只改构造这一行,其余代码零改动:

b, err := redisbroker.New(redisbroker.Options{
    Addr:      "127.0.0.1:6379",
    Password:  "",    // 空 = 不认证
    DB:        0,
    KeyPrefix: "tg:", // 默认 "tg:";多应用共用一个 Redis 时用它隔离
})
if err != nil { ... }
g, err := taskgate.New(taskgate.Config{Broker: b, Queues: ...})

所有"多步读写必须原子"的操作(认领、终态+依赖唤醒、连锁取消、计数维护)都在单段 Lua 脚本内完成,不存在"任务已离队却没有租约"的崩溃窗口;"父完成但子未唤醒"不可被观测到。

分布式限流:多进程共享配额

redis 后端额外实现了 LimiterProvider 能力接口,同一队列的 {Workers, RPS} 配额在所有接同一 Redis 的进程间共享——两个进程各配 {Workers: 2},全局同时在跑的也是 2 个,不是 4 个;RPS 走 GCRA(redis_rate),同样是全局速率。memory/sqlite 后端不受影响,维持进程内限流。

并发槽的自愈:每占一个槽记一个过期时刻(= 队列的 LeaseTTL),持有进程每 LeaseTTL/3 自动续期;进程崩溃 → 续期停 → 槽到期自动回收,最坏 2×LeaseTTL 内配额可再用,不会永久泄漏。

跨进程延迟(心里有数)
  • 新任务被别的进程发现:同进程提交有内部唤醒信号即时响应;别的进程写入靠 Dequeue 的兜底轮询发现,最坏 ≤100ms(与 sqlite 跨进程一致)。
  • 跨进程 Cancel:running 任务的取消标记由持有它的进程在下一次心跳发现,最坏 ≤ 一个心跳周期(LeaseTTL/3) 后其 handler ctx 被 cancel。
  • 崩溃任务被捞回:从最后一次成功续租算起,LeaseTTL(租约失效)+ 最坏 LeaseTTL/2(reaper 每 min(各队列 LeaseTTL)/2 扫一次)= LeaseTTL ~ 1.5×LeaseTTL(默认 60s → 60~90s)才会回到 pending;之后还要重新排队、占槽、等令牌才真正开跑。要更快发现崩溃就调小 LeaseTTL,代价是心跳更频繁(每 LeaseTTL/3 一次)。
Redis 键名速查(运维直查)

默认前缀 tg:(Options.KeyPrefix 可改)。积压、在跑数不用走应用,redis-cli 直接看:

类型 用途 直查示例
tg:task:{id} hash 任务全字段(时间存 unix 毫秒) HGETALL tg:task:01J...
tg:pending:{queue} list 就绪任务 ID 队列(FIFO) LLEN tg:pending:scoring
tg:delayed:{queue} zset 延迟/重试退避任务,score=run_at ZCARD tg:delayed:scoring
tg:inflight zset 在跑任务,score=租约到期时刻 ZCARD tg:inflight
tg:children:{id} set 反向依赖索引(依赖 {id} 的子任务) SMEMBERS tg:children:01J...
tg:idx:status:{status} set 状态索引(七态各一) SCARD tg:idx:status:failed
tg:idx:type:{type} set Type 索引(List 过滤用) SCARD tg:idx:type:ocr
tg:stats hash Type×Status 计数,字段 {type}:{status} HGETALL tg:stats
tg:types set 出现过的 Type SMEMBERS tg:types
tg:sem:{queue} zset 分布式并发槽(限流器私有) ZCARD tg:sem:scoring
rate:tg:{queue} string RPS 限速状态(redis_rate 的 GCRA 内部,限流器私有)。注意 rate: 前缀由 redis_rate 加在最外层,该键不在 KeyPrefix 命名空间内(实际键名 = rate: + KeyPrefix + 队列名),按前缀批量清理时别漏 GET rate:tg:scoring

Counts/Overview 就是读 tg:stats(每次流转时 Lua 顺手 HINCRBY 维护),QueueLen 就是 LLEN + ZCARD,都是计数器/长度读取,不扫全库。

测试与限制
  • 契约测试双档:miniredis 档离线进 CI;设 TASKGATE_REDIS_ADDR=127.0.0.1:6379 后同一套 22 条契约在真 Redis 上再跑一遍(随机 KeyPrefix 隔离、测后清理),验证 Lua 脚本兼容性。
  • 不支持 Redis Cluster:脚本内用前缀自行拼键、键不带 hash tag,面向单实例/主从/哨兵拓扑。
  • 限流键与任务键在同一个 Redis 实例:flushdb 级故障两者同生共死(诚实的取舍)。

SQL 后端:PostgreSQL / MySQL(多进程)

已经有一套 PostgreSQL 或 MySQL,不想再引一个 Redis 时用它:多个 worker 进程接同一个库抢同一批任务,认领互斥靠 FOR UPDATE SKIP LOCKED(要求 MySQL 8.0+ / PostgreSQL 9.5+),同样过 22 条契约。和 redis 一样,换后端只改构造函数一行,业务代码零改动:

import "github.com/AmbroseX/taskgate/pgbroker"
b, err := pgbroker.Open("postgres://user:pass@localhost:5432/db?sslmode=disable")
import "github.com/AmbroseX/taskgate/mysqlbroker"
b, err := mysqlbroker.Open("user:pass@tcp(localhost:3306)/db")

可选项:WithTablePrefix("myapp_")(默认 taskgate_,多应用共享一个库时用它隔离)、WithMaxOpenConns(n)(默认 10)、WithPollInterval(d)(默认 200ms)。首次 OpenInit 自动建表(CREATE TABLE IF NOT EXISTS,冷启动并发安全)。

已知限制
  • 跨进程新任务感知延迟 = 轮询间隔(默认 200ms,可调);未实现 PG LISTEN/NOTIFY 即时唤醒。
  • 不提供分布式限流:SQL 后端不实现 LimiterProvider,scheduler 自动退回进程内限流,多进程各限各的;需要精确跨进程限流请用 redis 后端。
  • 高并发依赖传播冲突时,事务会经历数据库死锁自动重试(有上限,默认 5 次),表现为个别调用延迟抬高;重试超限会把死锁错误原样抛出。
  • MySQL 专属:自定义 ID/type/queue 最长 255 字符(Enqueue 入口校验,超限清晰报错);payload/result 受服务器 max_allowed_packet 限制(默认 64M);表用 utf8mb4_bin 排序规则(DDL 内置,自定义 ID 大小写敏感)。
  • 契约测试需要真库(env 门控 TASKGATE_PG_DSN / TASKGATE_MYSQL_DSN),本地无库时 skip——本地全绿不代表跑过这两个后端,回归靠 CI。

Config 说明

库自己不读任何配置文件。Config 的字段带 yaml/json tag,应用自己 unmarshal 好再传进来:

# 应用自己的配置文件(taskgate 不读它,由应用 unmarshal 后注入)
queues:
  llm:
    workers: 2       # 并发上限(必填,>=1)
    rps: 3           # 每秒放行数,0 = 不限速
    burst: 3         # 突发额度,0 时取 max(1, int(rps))
    lease_ttl: 60s   # 租约时长,0 补默认 60s
    quota_limit: 5000   # 周期配额:每窗口最多启动 5000 次 handler,0 = 不启用
    quota_period: 24h   # 窗口时长(≥1s;固定窗口对齐 epoch,24h 即 UTC 零点,不是当地自然日)
    quota_key: my-gw    # 配额键,空 = 队列名;多队列同键 = 共享同一份窗口预算
  cpu:
    workers: 4
routes:              # 任务类型 → 队列;没配的类型用类型名当队列名
  generate: llm
default_queue:       # 兜底队列,可整个不配
  workers: 2
lease_lost_max: 3    # 租约丢失封顶(默认 3),超过进 failed
throttled_max: 100   # 限流重排封顶(默认 100),超过进 failed
var cfg taskgate.Config
_ = yaml.Unmarshal(raw, &cfg)      // 应用自己解;Duration 字段支持 "60s"、"10m" 写法
cfg.Broker = memorybroker.New()    // 运行期对象手动注入
cfg.OnStateChange = func(t taskgate.Task) { /* 埋点 */ }
g, err := taskgate.New(cfg)
频率 ≠ 配额:周期配额(硬配额)

RPS 管"每秒别把网关打爆",QuotaLimit 管"每个窗口累计别超预算"——这是两码事:网关允许你每天调 5000 次,不代表你必须匀速摊到每 17 秒一次,你可以一口气用完然后歇一天。三个维度正交并存:Workers 限同时在跑、RPS 限每秒新启动、QuotaLimit 限每窗口累计启动。

Queues: map[string]taskgate.QueueConfig{
    "llm": {
        Workers: 2, RPS: 3,                       // 别打爆:2 并发、每秒 3 个
        QuotaLimit: 5000,                          // 别超支:每窗口最多启动 5000 次 handler
        QuotaPeriod: taskgate.Duration(24 * time.Hour),
    },
}

合同(如实版):

  • 硬配额,绝不超发:窗口计数在共享介质里原子扣减(所有接同一介质的进程合计),任何故障——进程崩溃、介质断连、退还失败——方向都是只少放行,不多放行;
  • 窗口是对齐 epoch 的固定时长(24h = UTC 零点重置),不是当地自然日、不是自然月、也对不齐网关的账单周期;窗口时间用介质的服务端钟,应用机器的时钟偏差不影响;
  • 配额单位是 handler 启动次数,不是任务数(重试的再次认领同样计一次),也不是 token 用量;
  • 两个消耗漏口要心里有数:handler 启动后失败/被取消,以及类型没注册 handler 被判死信,这一次配额都已经消耗;
  • 额度耗尽不是错误:队列停止认领(不占 worker 槽),任务老实待在 pending 等下个窗口,不进 Throttled 计数、不进 failed;Stats(queue)QuotaExhausted 位可见;
  • 配额介质不可达时 fail-closed:暂停认领、零放行、退避重试(QuotaStalled 位可见)——绝不退回进程内计数假装还有保护;
  • 后端必须支持共享计数(五个内置后端都支持:memory=进程内、sqlite=库文件、redis/pg/mysql=服务器);不支持的后端配了配额,New() 直接报错。

边角用法

一些常用的小写法:

// 幂等提交:同一业务键重复提交返回 ErrTaskExists,不会重复入队。
// 任务 ID(ExecutionID)由系统生成并从 Submit 返回,业务键与任务 ID 是两个概念。
id, err := g.Submit(ctx, "generate", payload, taskgate.WithBusinessKey("order-42"))
var te *taskgate.TaskExistsError
if errors.As(err, &te) { /* 已排过队;te.ExecutionID/te.Status 是该键下最新执行 */ }

// 失败重跑(Replay):终态执行可以重放成一个新执行,旧记录永远不变。
// 典型 cron 配方:确定性业务键防双触发;失败后从拒绝错误里拿链尾,显式 Replay。
if te != nil && te.Status == taskgate.StatusFailed {
    newID, _ := g.Replay(ctx, te.ExecutionID)          // 新执行,ReplayOf 指回旧执行
    _ = newID
}
g.ReplayByKey(ctx, "order-42")                          // 按键重放,作用于该键最新执行
g.Replay(ctx, id, taskgate.AllowCompleted())            // 重放已成功的执行必须显式允许
g.Replay(ctx, id, taskgate.WithPayload(newPayload))     // 参数修正后重跑(默认复制旧 Payload)
history, _ := g.History(ctx, "order-42")                // 该键下的执行历史链(旧 → 新)

// 延迟执行:相对延迟 or 绝对时刻,二选一。
g.Submit(ctx, "reminder", payload, taskgate.Delay(30*time.Minute))
g.Submit(ctx, "reminder", payload, taskgate.RunAt(time.Now().Add(time.Hour)))

// 覆盖这条任务的重试上限(默认走 Config 的封顶)。
g.Submit(ctx, "flaky", payload, taskgate.MaxRetry(1))

// 扇入:一个任务依赖多个父任务,全部完成才唤醒。
g.Submit(ctx, "merge", nil, taskgate.DependsOn(idA, idB, idC))

// 父失败不连锁取消我(默认父 failed 会连锁取消子)。
g.Submit(ctx, "cleanup", nil, taskgate.DependsOn(job), taskgate.IgnoreParentFailure())

// 阻塞等最终结果 / 主动取消 / 查一条。
result, err := g.Wait(ctx, id)
err = g.Cancel(ctx, id)
task, err := g.Get(ctx, id)

// 看某队列积压与在跑数 / 全局各态计数。
stats, _ := g.Stats(ctx, "llm")
overview, _ := g.Overview(ctx)

错误类型速查

所有错误都导出,用 errors.Is / errors.As 判断:

// 哨兵错误(errors.Is)
taskgate.ErrTaskExists    // 业务键下已有执行(WithBusinessKey 幂等时会碰到;errors.As 可解构 *TaskExistsError 拿链尾)
taskgate.ErrTaskNotFound  // 任务不存在(Get/Cancel 找不到,或依赖的父任务缺失)
taskgate.ErrLeaseLost     // 租约令牌不匹配:任务已被回收或被别人重认领,结果作废
taskgate.ErrTaskCanceled  // 任务被打了取消标记,handler 该退出了
taskgate.ErrAlreadyFinal  // 对已进终态的任务再 Cancel
taskgate.ErrUnknownType   // Run 时遇到没注册 handler 的任务类型
taskgate.ErrShutdown      // Gate 已 Shutdown,拒绝新提交
taskgate.ErrNoTask        // 在 handler 之外的 ctx 上调 RenewLease
taskgate.ErrReplayNotFinal      // Replay 目标还没进终态
taskgate.ErrAlreadyReplayed     // Replay 目标已被重放过(历史链不分叉)
taskgate.ErrCompletedNotAllowed // 重放 completed 执行必须显式 AllowCompleted()

// 结构化错误(errors.As;handler 返回它们控制重试路径,必须按值返回)
taskgate.ErrThrottled{RetryAfter: d} // 网关限流:延后重排,不占重试次数
taskgate.ErrSkipRetry{Err: err}      // 没救的错:直接进 failed;Unwrap 可穿透到原错误

运行测试

全量离线可跑(L1 单元 → L2 brokertest 契约 → L3 集成 → L4 仿真 E2E):

go test ./... -race

想在真库上再跑一遍 22 条契约(可选;Redis 验证 Lua 脚本兼容性,PG/MySQL 本地无库时自动 skip):

TASKGATE_REDIS_ADDR=127.0.0.1:6379 go test ./redisbroker/... -race
TASKGATE_PG_DSN="postgres://postgres:pass@localhost:5432/postgres?sslmode=disable" go test ./pgbroker/... -race
TASKGATE_MYSQL_DSN="root:pass@tcp(localhost:3306)/taskgate" go test ./mysqlbroker/... -race

测试分层与 e2e

e2e/ 目录是 L4/L5 仿真:

  • e2e/mockgw/:可注入故障的 mock LLM/OCR 网关(测试组件,不属库 API)。把生产踩过的坑做成开关:Latency(延迟)、BusyAfterConcurrency(并发超限返 HTTP 200 但 body 里藏 busy 错误事件——复刻"状态码骗人"的真实网关)、FailRate(固定种子随机 500,CI 可复现)、CrashAfterConcurrency(并发超限直接断连)、BusyFirstN(前 N 个请求定向 busy);暴露 MaxConcurrency/BusyCount/CrashCount/Requests 原子观测口。
  • e2e/pipeline_test.go:五个核心用例——限流真的挡住 busy、busy 走 ErrThrottled 重排零 failed、断连走普通重试补完、三队列流水线 30/30 且结果逐级传递、中途取消连锁生效、SSE 藏错误重排后成功。
  • e2e/realgw_test.go:真实网关冒烟档,//go:build realgw 隔离,不进 CI(常规 go vet/go test 完全不编译它);读 LLM_GATEWAY_URL/LLM_GATEWAY_KEY(缺失自动 skip),手动执行:
LLM_GATEWAY_URL=https://网关地址 LLM_GATEWAY_KEY=密钥 \
  go test -tags realgw ./e2e/ -run RealGW -v

设计文档

里程碑

  • M1(已完成):核心排队、限流、重试、依赖、取消、Shutdown,memory / sqlite 双后端。
  • M2(已完成):redis 后端(Lua 原子流转、多进程认领互斥)、分布式限流(跨进程共享配额)、性能基线。
  • M3(已完成):L4 仿真 E2E(mockgw 故障注入五用例)、handler 手动续租(RenewLease/ManualHeartbeat)、List 分页、realgw 手动冒烟档。

明确不做(YAGNI):任务优先级、webhook 通知、游标分页、Web UI、cron 周期调度、DAG 工作流引擎、独立 server 模式。

Documentation

Overview

Package taskgate 是一个轻量的 Go 任务排队限流库:排队、限流、重试、依赖、取消。 它是库不是服务:不读环境变量、不读配置文件,只吃调用方传进来的 Config。

Index

Constants

This section is empty.

Variables

View Source
var (
	// ErrTaskExists 任务已存在:带 BusinessKey 提交撞键(此时错误可 errors.As 出
	// *TaskExistsError 拿链尾信息),或预置 ID 撞主键(测试/嵌入入口的存储防御)。
	ErrTaskExists = errors.New("taskgate: task already exists")
	// ErrTaskNotFound 任务不存在(Get/Cancel 找不到,或 Enqueue 时父任务缺失)。
	ErrTaskNotFound = errors.New("taskgate: task not found")
	// ErrLeaseLost 租约令牌不匹配:任务已被回收或被别人重新认领,结果作废。
	ErrLeaseLost = errors.New("taskgate: lease lost")
	// ErrTaskCanceled Heartbeat 发现任务被打了取消标记,scheduler 该 cancel handler 的 ctx 了。
	ErrTaskCanceled = errors.New("taskgate: task canceled")
	// ErrAlreadyFinal 对已经进终态的任务再 Cancel。
	ErrAlreadyFinal = errors.New("taskgate: task already in final state")
	// ErrUnknownType Run 时遇到没注册 handler 的任务类型(Submit 不校验)。
	ErrUnknownType = errors.New("taskgate: no handler registered for task type")
	// ErrShutdown Gate 已经 Shutdown,拒绝新提交。
	ErrShutdown = errors.New("taskgate: gate is shut down")
	// ErrNoTask 在任务 handler 之外的 ctx 上调 RenewLease:ctx 里没有续租闭包。
	ErrNoTask = errors.New("taskgate: no task associated with context")
	// ErrReplayNotFinal Replay 的目标执行还没进终态(只有终态执行可被重放)。
	ErrReplayNotFinal = errors.New("taskgate: replay target not in final state")
	// ErrAlreadyReplayed Replay 的目标已被重放过:历史链不分叉,每个执行至多被重放一次,
	// 重放只能打在链尾(最新执行)上。
	ErrAlreadyReplayed = errors.New("taskgate: execution already replayed (chain must not fork)")
	// ErrCompletedNotAllowed 重放 completed 的执行必须显式带 AllowCompleted() 选项,
	// 防止误触发重复计费。
	ErrCompletedNotAllowed = errors.New("taskgate: replaying a completed execution requires AllowCompleted")
)

哨兵错误:全部导出,调用方用 errors.Is 判断。

Functions

func CanTransition

func CanTransition(from, to Status) bool

CanTransition 导出状态机校验,给 memorybroker/sqlitebroker 等后端包在 每次写状态前做统一防御(合同要求:所有写入先过 canTransition 表)。 Phase 1 只定义了包内的 canTransition,后端在别的包里够不着,这里补一个只读出口。

func ParentFailureReason

func ParentFailureReason(parentID string, parentStatus Status) string

ParentFailureReason 连锁取消时写进子任务 LastError 的固定文案。 固定下来是为了 brokertest 能逐字断言,三个后端文案一致。

func RenewLease

func RenewLease(ctx context.Context) error

RenewLease 在 handler 里手动给当前任务续租(lease_until = now + LeaseTTL)。 自动档(默认)也可以调,与自动心跳互不干扰;手动档(QueueConfig.ManualHeartbeat=true) 必须靠它保活。ctx 必须是 handler 收到的那个任务 ctx(或它的子 ctx)。返回值:

  • nil:续租成功;
  • ErrTaskCanceled:任务已被外部 Cancel(续租照做),此时任务 ctx 已被 cancel, handler 应尽快退出;
  • ErrLeaseLost:租约已丢(任务被 reaper 回收),结果注定作废,任务 ctx 已被 cancel,handler 应立即放弃;
  • ErrNoTask:ctx 不是任务 ctx(handler 之外调用);
  • 其他错误(网络抖动等):续租没成也没丢,handler 可稍后重试。

Types

type Broker

type Broker interface {
	Init(opts BrokerOptions) error // New(cfg) 时调用一次,Dequeue 前必须先 Init
	Enqueue(ctx context.Context, t *Task) error
	Dequeue(ctx context.Context, queues []string) (*Task, error)
	Ack(ctx context.Context, id, leaseToken string, result []byte) error
	Fail(ctx context.Context, id, leaseToken, errMsg string, kind FailKind, retryAt time.Time) error
	Cancel(ctx context.Context, id string) error
	FinishCanceled(ctx context.Context, id, leaseToken string) error
	Requeue(ctx context.Context, id, leaseToken string) error
	Heartbeat(ctx context.Context, id, leaseToken string) error
	Get(ctx context.Context, id string) (*Task, error)
	Replay(ctx context.Context, req ReplayRequest) (*Task, error)
	List(ctx context.Context, f Filter) ([]*Task, error)
	QueueLen(ctx context.Context, queue string) (int, error)
	Counts(ctx context.Context) (map[string]map[Status]int64, error)
	ReapExpired(ctx context.Context) (int, error)
	Close() error
}

Broker 存储后端接口。只收 memory/sqlite/redis 三后端都能同语义实现的方法(宪法第 II 条); 各方法的行为合同见 contracts/broker-contract.md,由 brokertest 套件统一验收。

type BrokerOptions

type BrokerOptions struct {
	LeaseTTL        map[string]time.Duration // 队列→租约 TTL
	DefaultLeaseTTL time.Duration            // 缺省 60s
	LeaseLostMax    int                      // 缺省 3
	ThrottledMax    int                      // 缺省 100
	Notify          func(Task)               // 状态流转回调,可 nil;必须在锁/事务外异步调
	Clock           Clock                    // 可 nil=真时钟
}

BrokerOptions New(cfg) 装配时传给后端的运行参数,签名照 contracts/broker-contract.md。

type ChildAction

type ChildAction int

ChildAction 父任务到终态之后,对一个直接子任务要执行的动作。

const (
	// ChildNone 什么都不做:计数减了但还没减到 0,或者子任务已经在终态。
	ChildNone ChildAction = iota
	// ChildWake 唤醒:子任务 blocked → pending(父全部满足了)。
	ChildWake
	// ChildCancel 连锁取消:子任务应流转到 canceled(FailFast 且父失败/取消)。
	// 调用方仍需过 canTransition 校验;若子在 running 等特殊状态,由调用方自行防御处理。
	ChildCancel
)

func DecideOnParentFinal

func DecideOnParentFinal(parentStatus, childStatus Status, policy ParentFailurePolicy, pendingParents int) (int, ChildAction)

DecideOnParentFinal 判定"父任务进入终态 parentStatus 之后,对一个子任务怎么办"。 pendingParents 传入子任务当前剩余的未终态父计数(调用方在锁/事务内读出), 返回新的计数(递减不为负)与动作。约定:

  • 子任务已是终态(比如已被手动取消)→ 不动。
  • 父是 completed,或策略是 IgnoreParentFail(此时父任何终态都算满足)→ 计数减一; 减到 0 且子还在 blocked → 唤醒。
  • 父是 failed/canceled 且策略是 FailFast → 连锁取消(计数保持原样,反正任务要没了)。

type Clock

type Clock interface {
	// Now 当前时刻。
	Now() time.Time
	// After 到点后往返回的 channel 发一次当前时刻。
	After(d time.Duration) <-chan time.Time
	// Sleep 睡 d,ctx 先取消就提前返回 ctx.Err()。
	Sleep(ctx context.Context, d time.Duration) error
	// NewTicker 周期滴答,给 reaper/心跳循环用。
	NewTicker(d time.Duration) Ticker
}

Clock 可注入的时钟。租约、退避、限流全部通过它拿时间, 这样测试里用 fakeclock 手动推进,不用真 sleep(宪法第 V 条)。

func RealClock

func RealClock() Clock

RealClock 系统真时钟。BrokerOptions.Clock 传 nil 时后端应退回到它。

type Config

type Config struct {
	Broker        Broker                 `yaml:"-" json:"-"`
	Queues        map[string]QueueConfig `yaml:"queues" json:"queues"`
	Routes        map[string]string      `yaml:"routes" json:"routes"` // Type → Queue
	DefaultQueue  QueueConfig            `yaml:"default_queue" json:"default_queue"`
	OnStateChange func(Task)             `yaml:"-" json:"-"`
	LeaseLostMax  int                    `yaml:"lease_lost_max" json:"lease_lost_max"` // 0 补默认 3
	ThrottledMax  int                    `yaml:"throttled_max" json:"throttled_max"`   // 0 补默认 100
}

Config 全局配置。库不读 env/文件,应用自己 unmarshal 好再传进来; Broker 和 OnStateChange 是运行期对象,序列化时跳过。

type Duration

type Duration time.Duration

Duration 包一层 time.Duration,让 yaml/json 配置里能直接写 "10m"、"60s" 这种人话。

func (Duration) MarshalText

func (d Duration) MarshalText() ([]byte, error)

MarshalText 序列化回 "10m0s" 这种标准格式。

func (*Duration) UnmarshalText

func (d *Duration) UnmarshalText(b []byte) error

UnmarshalText 支持 "10m" 这类写法,yaml 和 json 解码都走这里。

type ErrSkipRetry

type ErrSkipRetry struct {
	Err error
}

ErrSkipRetry handler 返回它表示"这个错没救,别重试了",任务直接进 failed。 必须按值返回(errors.As 按值匹配),不要返回其指针。

func (ErrSkipRetry) Error

func (e ErrSkipRetry) Error() string

Error 实现 error 接口。

func (ErrSkipRetry) Unwrap

func (e ErrSkipRetry) Unwrap() error

Unwrap 让 errors.Is/As 能穿透到里面包的业务错误。

type ErrThrottled

type ErrThrottled struct {
	RetryAfter time.Duration
}

ErrThrottled handler 返回它表示"被网关限流了,过 RetryAfter 再来": 不占 Attempts,只涨 Throttled 计数,封顶(默认 100)才进 failed。 必须按值返回(errors.As 按值匹配),不要返回其指针。

func (ErrThrottled) Error

func (e ErrThrottled) Error() string

Error 实现 error 接口。

type FailKind

type FailKind int

FailKind Fail 的三种语义,决定动哪个计数、进 retrying 还是 failed。

const (
	// FailBusiness 业务失败:Attempts+1;Attempts>MaxRetry → failed,否则 retrying。
	FailBusiness FailKind = iota
	// FailThrottled 被网关限流:Throttled+1;≥ThrottledMax → failed,否则 retrying;Attempts 不动。
	FailThrottled
	// FailSkip 明确不重试:直接 failed。
	FailSkip
)

type Filter

type Filter struct {
	Type        string
	Queue       string
	Status      Status
	BusinessKey string // 非空时只返回该业务键下的执行;与其余条件是 AND 关系
	Limit       int    // 0=不限
	Offset      int    // 排序后跳过的条数,0=不跳过
}

Filter List 的过滤条件,零值字段表示不过滤。

排序与分页合同(M3 定型,三后端一致,见 contracts/broker-contract.md):

  • 结果一律按 (CreatedAt, ID) 升序:CreatedAt 由 broker 落库时统一写, 同一毫秒内再按 ID 定序,保证全序;
  • 执行顺序写死:先按 Type/Queue/Status 过滤 → 排序 → 跳过 Offset 条 → 取 Limit 条;
  • Offset ≥ 匹配总数 → 返回空列表(nil error);Offset < 0 按 0 处理;
  • 翻页弱一致:翻页期间数据变动不承诺快照一致,只承诺"未变动的任务不丢不重"。

type Gate

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

Gate 是 taskgate 的统一门面:提交、查询、等待、消费全从这里走。 一个 Gate 既可以只当生产者(New 后直接 Submit,不 Handle 不 Run), 也可以注册 handler 后 Run 起来当消费者,两者共用同一个 Broker。

func New

func New(cfg Config) (*Gate, error)

New 校验配置、补默认值、装配 BrokerOptions 并 Init 后端,返回可用的 Gate。 配置有问题直接返回 error,fail fast,绝不 panic。

func (*Gate) Cancel

func (g *Gate) Cancel(ctx context.Context, id string) error

Cancel 取消任务(US6):

  • blocked/pending/retrying:后端直接置 canceled 并向下传播(FailFast 子连锁取消);
  • running:后端只打取消标记,然后看任务是否正在本进程跑——在的话立即 cancel 它的 handler ctx(不在本进程跑的,由持有它的进程下一次 Heartbeat 发现标记); handler 退出后 scheduler 调 FinishCanceled 落库 canceled;
  • 终态:返回 ErrAlreadyFinal;不存在:返回 ErrTaskNotFound。

func (*Gate) Get

func (g *Gate) Get(ctx context.Context, id string) (*Task, error)

Get 查单个任务的当前快照。

func (*Gate) Handle

func (g *Gate) Handle(taskType string, h Handler)

Handle 注册任务类型对应的处理函数。必须在 Run 之前注册完。

func (*Gate) History

func (g *Gate) History(ctx context.Context, businessKey string) ([]*Task, error)

History 枚举该 BusinessKey 下的执行历史链,链序(旧 → 新),链尾即最新执行; 键不存在返回空切片。它是 List(Filter{BusinessKey: key}) 的便捷封装。

func (*Gate) List

func (g *Gate) List(ctx context.Context, f Filter) ([]*Task, error)

List 按过滤条件列任务。

func (*Gate) Overview

func (g *Gate) Overview(ctx context.Context) (map[string]map[Status]int64, error)

Overview 全局概览:Type × Status 的数量矩阵(就是 Broker.Counts)。

func (*Gate) Replay

func (g *Gate) Replay(ctx context.Context, executionID string, opts ...ReplayOption) (string, error)

Replay 按 ExecutionID 重放一次终态执行:创建**新执行**(新 ID、ReplayOf 指回目标、 三计数清零、默认复制目标 Payload)进入正常调度,目标记录逐字段不变。 目标必须是其链的链尾且已终态;completed 需显式 AllowCompleted()。返回新执行的 ID。

func (*Gate) ReplayByKey

func (g *Gate) ReplayByKey(ctx context.Context, businessKey string, opts ...ReplayOption) (string, error)

ReplayByKey 按 BusinessKey 重放,天然作用于链尾(该键下最新执行)。语义同 Replay。

func (*Gate) Run

func (g *Gate) Run(ctx context.Context) error

Run 启动消费:按注册过 handler 的类型对应的队列起认领循环,阻塞到 ctx 取消或 Shutdown, 然后停止认领、等在跑任务全部收尾后返回(因 Shutdown 退出同样返回 nil)。 生命周期细节在 scheduler.go。

func (*Gate) Shutdown

func (g *Gate) Shutdown(ctx context.Context) error

Shutdown 优雅停止(US7):

  1. 一进门置停机标记,此后 Submit 一律返回 ErrShutdown,认领循环停止拿新任务;
  2. 等所有在跑任务善终(Run 也随之退出并返回 nil);
  3. ctx 先到期:cancel 各在跑任务的 handler ctx,等 handler 退出后把这些任务 Requeue 回 pending(三计数与 RunAt 全不动),返回 ctx 的超时错误;
  4. 后台 goroutine(认领循环/心跳/reaper)全部同步收尾,返回后零泄漏;
  5. 重复调用幂等:第二次直接返回 nil。

注意 Shutdown 的打断不是用户取消:被打断的任务回 pending 等下次重跑,不会进 canceled。

func (*Gate) Stats

func (g *Gate) Stats(ctx context.Context, queue string) (QueueStats, error)

Stats 查单个队列的水位。队列没配置且没有 DefaultQueue 兜底时报错。

func (*Gate) Submit

func (g *Gate) Submit(ctx context.Context, taskType string, payload json.RawMessage, opts ...SubmitOption) (string, error)

Submit 提交任务:按 Routes 定队列、应用提交选项、入队,返回任务 ID。 Delay 和 RunAt 都设置时 RunAt 生效(RunAt 是绝对时刻,语义更强)。

func (*Gate) Wait

func (g *Gate) Wait(ctx context.Context, id string) (json.RawMessage, error)

Wait 阻塞等任务到终态:completed 返回 Result;failed/canceled 返回 *TaskFailedError; ctx 先取消返回 ctx.Err()(任务本身照常跑,Wait 只是不等了)。 实现是 50ms 轮询 Get,间隔走注入的 clock。

type Handler

type Handler func(ctx context.Context, t *Task) ([]byte, error)

Handler 业务处理函数:返回的 []byte 会作为 Result 写库;返回错误决定重试路径。

type LimiterProvider

type LimiterProvider interface {
	// QueueLimiter 按队列配置构造该队列的限流器;出错时 Gate.Run 直接返回该错误。
	QueueLimiter(queue string, qc QueueConfig) (QueueLimiter, error)
}

LimiterProvider 后端的**可选能力接口**:能为队列提供跨进程共享的限流器。

限流不是所有后端都能做(memory/sqlite 没有跨进程共享的介质),进不了 Broker 接口的"最小公倍数",所以单独拆成能力接口。scheduler 装配限流器时只做 `broker.(LimiterProvider)` 这一次**能力断言**——断言的是接口不是具体后端类型, 上层依然不 import 任何后端包,不违反宪法 II.2"上层不特判后端": 新后端想提供分布式限流,实现本接口即可;memory/sqlite 不实现, scheduler 自动退回进程内限流(localLimiter),行为与 M1 完全一致。

实现约束:QueueLimiter 的构造必须廉价、不得持有需要显式释放的资源—— 本接口没有 Close,某个队列构造失败时之前已建成的限流器不会被回收, 构造期占了资源就是泄漏(redisbroker 的实现只是复用 Broker 连接拼参数,零资源)。

type ParentFailurePolicy

type ParentFailurePolicy string

ParentFailurePolicy 父任务失败时子任务怎么办。

const (
	// FailFast 父任务失败/取消 → 子任务连锁取消(默认)。
	FailFast ParentFailurePolicy = "fail_fast"
	// IgnoreParentFail 父任务只要进了终态(哪怕失败)就照常唤醒子任务。
	// 注:同名选项函数是 IgnoreParentFailure(),常量名少个 ure 是为了避开重名。
	IgnoreParentFail ParentFailurePolicy = "ignore_parent_failure"
)

type ParentState

type ParentState struct {
	ID     string
	Status Status
}

ParentState 做决策需要的父任务快照:只要 ID(拼错误文案用)和状态。

type QueueConfig

type QueueConfig struct {
	Workers  int      `yaml:"workers" json:"workers"`
	RPS      float64  `yaml:"rps" json:"rps"`             // 0 = 不限速
	Burst    int      `yaml:"burst" json:"burst"`         // 0 时取 max(1, int(RPS))
	LeaseTTL Duration `yaml:"lease_ttl" json:"lease_ttl"` // 0 补默认 60s
	// ManualHeartbeat 手动续租开关。默认 false:scheduler 给每个在跑任务起自动心跳,
	// 每 LeaseTTL/3 续租一次。true:不起自动心跳,handler 必须自己定期调
	// taskgate.RenewLease 保活,否则租约到期会被 reaper 回收(LeaseLost+1)。
	// 注意手动档下跨进程 Cancel 只能靠 handler 下一次 RenewLease 发现
	// (返回 ErrTaskCanceled);handler 一直不续租,则由租约过期兜底回收。
	// 本地 Cancel(同进程)不受影响,依然即时打断 handler 的 ctx。
	ManualHeartbeat bool `yaml:"manual_heartbeat" json:"manual_heartbeat"`

	// 周期配额(spec 006,硬配额):每个固定时长窗口最多启动 QuotaLimit 次 handler。
	// 窗口对齐 epoch(windowStart = now/period×period),时间取共享介质的服务端钟,
	// 不是自然日/自然月;单位是"handler 启动次数",不是任务数(重试的再认领同样计数)。
	// QuotaLimit=0 完全不启用(零开销);启用时 QuotaPeriod 必须 >0,且后端必须实现
	// QuotaProvider 能力接口,否则 New() 直接报错——配额没有静默降级。
	QuotaLimit  int      `yaml:"quota_limit" json:"quota_limit"`
	QuotaPeriod Duration `yaml:"quota_period" json:"quota_period"`
	// QuotaKey 配额键,空 = 队列名。多个队列配同一个 key 时共享同一份窗口预算
	// (打同一个网关额度的场景),此时各队列的 (QuotaLimit, QuotaPeriod) 必须一致。
	QuotaKey string `yaml:"quota_key" json:"quota_key"`
}

QueueConfig 单个队列的限流参数。

type QueueLimiter

type QueueLimiter interface {
	// AcquireSlot 占一个并发槽,占不到就阻塞;ctx 取消返回 ctx.Err()。
	AcquireSlot(ctx context.Context) error
	// ReleaseSlot 归还并发槽。必须和 AcquireSlot 一一配对。
	ReleaseSlot()
	// WaitToken 等一个 RPS 令牌;不限速时立即放行;ctx 取消返回其错误。
	WaitToken(ctx context.Context) error
}

QueueLimiter 单个队列的限流器抽象,两层独立生效:

  • 并发槽(AcquireSlot/ReleaseSlot):限"同时在跑多少个";
  • RPS 令牌(WaitToken):限"每秒新启动多少个"。

scheduler 只依赖这个接口,不关心限流器是进程内的还是跨进程共享的: 后端实现了 LimiterProvider 就用后端给的,否则用进程内的 localLimiter。 两层怎么配合用见 scheduler.claimLoop:先占槽、再等令牌。

type QueueStats

type QueueStats struct {
	Workers  int     `json:"workers"`   // 配置的并发上限
	Running  int     `json:"running"`   // 本进程正在执行的任务数(纯生产者恒为 0)
	QueueLen int     `json:"queue_len"` // 积压:pending + retrying
	RPS      float64 `json:"rps"`       // 配置的限速,0 = 不限
	// 周期配额状态位(spec 006):"队列不动了"必须能靠这两个位区分原因。
	QuotaExhausted bool `json:"quota_exhausted"` // 本窗口额度已尽,认领暂停等下窗
	QuotaStalled   bool `json:"quota_stalled"`   // 配额介质不可达,fail-closed 暂停中
}

QueueStats 单个队列的水位:配置的并发/限速 + 当前在跑数 + 积压长度 + 配额状态位。

type QuotaGate

type QuotaGate interface {
	// Reserve 原子预留一份额度:在共享介质内一个原子单位完成
	// "取服务端时间 → 算窗口 → 检查余额 → 扣减",检查与扣减之间没有窗口。
	// 返回三态:
	//   - res≠nil:预留成功,res.Window 是本次预留落在的窗口起点;
	//   - res==nil 且 err==nil:本窗口额度耗尽(**不是错误**,等下个窗口);
	//   - err≠nil:介质故障,调用方必须 fail-closed(零放行,退避重试)。
	Reserve(ctx context.Context) (*QuotaReservation, error)
	// Release 尽力退还一份预留(认领扑空/出错的补偿),只作用于 r 的窗口;
	// 窗口已切走则落空无害。失败时调用方不重试——该份额度当 leaked(视同消耗),
	// 方向永远保守:任何故障只少放行、不多放行。
	Release(ctx context.Context, r *QuotaReservation) error
}

QuotaGate 单个队列(quota key)的配额闸。行为合同见 specs/006-periodic-quota/contracts/quota-capability-contract.md。

type QuotaProvider

type QuotaProvider interface {
	// QueueQuota 按队列配置构造配额闸;只在 qc.QuotaLimit > 0 时被调用。
	QueueQuota(queue string, qc QueueConfig) (QuotaGate, error)
}

QuotaProvider 后端的**可选能力接口**(spec 006):能为队列提供跨进程共享的周期配额。 与 LimiterProvider 平行,但合同相反——**没有静默降级**:配置了 QuotaLimit>0 而后端 未实现本接口,taskgate.New() 直接报错。硬配额的全部意义是"绝不超发",退回进程内 计数等于假保护,宁可不启动(模型裁决 #3)。

实现约束:构造必须廉价、不持有需显式释放的资源(同 LimiterProvider); quota key 相同的多个 QuotaGate 共享介质计数,实例之间不得有本地共享状态。

type QuotaReservation

type QuotaReservation struct {
	Window int64
}

QuotaReservation 一次额度预留。Window 是预留落在的窗口起点 (unix 秒,共享介质的服务端钟),Release 靠它定位退还目标。

type ReplayOption

type ReplayOption func(*replayOptions)

ReplayOption 重放时的函数式选项。

func AllowCompleted

func AllowCompleted() ReplayOption

AllowCompleted 显式允许重放 completed 的执行("重新生成报告"这类主动重跑)。 不带它重放 completed 会拿到 ErrCompletedNotAllowed——防止误触发重复计费。

func WithPayload

func WithPayload(p json.RawMessage) ReplayOption

WithPayload 用新 Payload 重放(参数修正后重跑)。不带它默认复制目标执行的 Payload; 要显式清空就传 json.RawMessage("null") 或 "{}"——nil 表示"没传",非 nil 即覆盖。

type ReplayRequest

type ReplayRequest struct {
	ExecutionID    string          // 目标执行,必须是其链的链尾
	BusinessKey    string          // 与 ExecutionID 二选一
	AllowCompleted bool            // 重放 completed 必须显式打开
	Payload        json.RawMessage // nil = 复制目标执行的 Payload;非 nil 即覆盖
}

ReplayRequest Replay 的入参:目标用 ExecutionID 或 BusinessKey 指定,恰好一个非空。 按键指定时天然作用于链尾(该键下最新执行)。整个校验+创建必须在后端一个原子单位 (同事务/同 Lua/同临界区)内完成,行为合同见 specs/005-identity-replay/contracts/。

type Status

type Status string

Status 任务状态,共七态。用字符串是为了落库和日志里直接可读。

const (
	StatusBlocked   Status = "blocked"   // 有父任务还没跑完,等唤醒
	StatusPending   Status = "pending"   // 排队中,可被认领
	StatusRunning   Status = "running"   // 已被 worker 认领,持有租约
	StatusRetrying  Status = "retrying"  // 失败后等退避时间到点重跑
	StatusCompleted Status = "completed" // 终态:成功
	StatusFailed    Status = "failed"    // 终态:失败(重试耗尽/跳过重试/计数封顶)
	StatusCanceled  Status = "canceled"  // 终态:被取消(主动取消或父失败传播)
)

func (Status) IsFinal

func (s Status) IsFinal() bool

IsFinal 是否终态。终态没有任何出边,不允许再流转。

type SubmitDecision

type SubmitDecision struct {
	// Status 只会是三者之一:pending(可直接排队)、blocked(等父)、canceled(父已失败且 FailFast)。
	Status Status
	// PendingParents 仅在 Status==blocked 时有意义:还没到终态的父任务数(同 ID 去重后)。
	PendingParents int
	// LastError 仅在 Status==canceled 时有意义:取消原因,如 "parent <id> failed"。
	LastError string
}

SubmitDecision 提交(Enqueue)时的初始状态判定结果。

func DecideOnSubmit

func DecideOnSubmit(parents []ParentState, policy ParentFailurePolicy) SubmitDecision

DecideOnSubmit 判定一个带依赖的任务在提交那一刻应该落成什么初始状态。 规则(照 broker-contract.md 的 Enqueue 合同):

  • FailFast 策略下,只要有任何一个父已经 failed/canceled → 直接 canceled。 哪怕其它父还没跑完也立即取消:这个子任务已经注定跑不成,等下去没有意义。
  • IgnoreParentFail 策略下,父只要进了终态(哪怕失败)就算"满足"。
  • 还有父没到终态 → blocked,并记下未完成父的数量(pending_parents)。
  • 父全部满足 → pending,可以直接排队。

同一个父 ID 写了多遍只算一个:否则 pending_parents 会多计,父完成一次只减一次, 子任务就永远唤不醒了。调用方(后端)拿到 PendingParents 后按这个数落库。

type SubmitOption

type SubmitOption func(*submitOptions)

SubmitOption 提交任务时的函数式选项。

func Delay

func Delay(d time.Duration) SubmitOption

Delay 延迟 d 之后才允许执行(提交时换算成 RunAt)。

func DependsOn

func DependsOn(ids ...string) SubmitOption

DependsOn 声明父任务,父全部完成才会被唤醒;父 ID 必须已存在,否则拒收。

func IgnoreParentFailure

func IgnoreParentFailure() SubmitOption

IgnoreParentFailure 父任务失败也照常执行(默认是 FailFast 连锁取消)。

func MaxRetry

func MaxRetry(n int) SubmitOption

MaxRetry 业务失败最多重试 n 次(Attempts > n 进 failed)。

func RunAt

func RunAt(t time.Time) SubmitOption

RunAt 指定最早执行时刻,和 Delay 二选一,后设置的生效。

func WithBusinessKey

func WithBusinessKey(key string) SubmitOption

WithBusinessKey 业务幂等键:同键下已存在任何执行(不论状态)时 Submit 拒绝, 错误满足 errors.Is(err, ErrTaskExists),且可 errors.As 出 *TaskExistsError 拿到 链尾执行的 ID 与状态。失败后想再跑同一件事,走 Replay,不走再次 Submit。

func WithID deprecated

func WithID(id string) SubmitOption

WithID 旧的"自定义任务 ID"选项。

Deprecated: 任务 ID(ExecutionID)已收紧为系统生成,用户不可指定;本选项现在 等同于 WithBusinessKey——传入的值成为业务幂等键,不再是任务 ID,**不能**拿去 Get/DependsOn(那两处只认 Submit 返回的 ID)。新代码请直接用 WithBusinessKey。

type Task

type Task struct {
	ID              string              `json:"id"`                     // ExecutionID:一次执行的永久身份,broker 生成 ulid,永不复用;公开 API 不提供写入口
	BusinessKey     string              `json:"business_key,omitempty"` // 业务幂等键:同键下存在任何执行则 Enqueue 拒绝;创建后不可变
	ReplayOf        string              `json:"replay_of,omitempty"`    // 本执行重放自哪个 ExecutionID;由 Replay 写入,创建后不可变
	Type            string              `json:"type"`                   // 决定 handler 和默认队列
	Queue           string              `json:"queue"`                  // 限流单元,入队那一刻定死
	Payload         json.RawMessage     `json:"payload,omitempty"`      // 入参
	Status          Status              `json:"status"`
	Result          json.RawMessage     `json:"result,omitempty"` // Ack 时写入
	LastError       string              `json:"last_error,omitempty"`
	Attempts        int                 `json:"attempts"`  // 业务失败次数,> MaxRetry → failed
	MaxRetry        int                 `json:"max_retry"` // 0 = 不重试
	LeaseLost       int                 `json:"lease_lost"`
	Throttled       int                 `json:"throttled"`
	RunAt           time.Time           `json:"run_at"` // 延迟执行和重试退避都靠它
	DependsOn       []string            `json:"depends_on,omitempty"`
	OnParentFailure ParentFailurePolicy `json:"on_parent_failure"`
	LeaseToken      string              `json:"lease_token,omitempty"` // Dequeue 时携带,对外只读
	CreatedAt       time.Time           `json:"created_at"`
	StartedAt       time.Time           `json:"started_at,omitzero"`
	FinishedAt      time.Time           `json:"finished_at,omitzero"`
}

Task 任务实体。字段语义照 data-model.md 第 1 节,Payload/Result 一律 json.RawMessage。

type TaskExistsError

type TaskExistsError struct {
	BusinessKey string // 撞的键
	ExecutionID string // 键下链尾执行的 ID
	Status      Status // 链尾执行的状态
}

TaskExistsError 带 BusinessKey 提交撞键时的错误:errors.Is(err, ErrTaskExists) 照常成立,同时携带键下链尾(最新执行)的身份与状态,调用方据此直接决定要不要 Replay,不必再按键查询绕一圈。errors.As 按 *TaskExistsError 匹配。

func (*TaskExistsError) Error

func (e *TaskExistsError) Error() string

Error 实现 error 接口。

func (*TaskExistsError) Unwrap

func (e *TaskExistsError) Unwrap() error

Unwrap 让 errors.Is(err, ErrTaskExists) 保持成立,存量判错代码零改动。

type TaskFailedError

type TaskFailedError struct {
	ID        string // 任务 ID
	Status    Status // failed 或 canceled
	LastError string // 最后一次失败/取消原因
}

TaskFailedError Wait 等到 failed/canceled 终态时返回的错误,带上任务现场方便定位。

func (*TaskFailedError) Error

func (e *TaskFailedError) Error() string

Error 实现 error 接口,带上任务 ID、终态与原因。

type Ticker

type Ticker interface {
	C() <-chan time.Time
	Stop()
}

Ticker 抽出接口是为了 fakeclock 能提供假的滴答。

Directories

Path Synopsis
Package brokertest 是 Broker 行为契约的统一验收套件。
Package brokertest 是 Broker 行为契约的统一验收套件。
e2e
mockgw
Package mockgw 是一个可注入故障的 mock LLM/OCR 网关,只给 e2e 测试用,不属于库的公开 API。
Package mockgw 是一个可注入故障的 mock LLM/OCR 网关,只给 e2e 测试用,不属于库的公开 API。
examples
llm command
examples/llm 三级 LLM 流水线示例:检索(retrieve)→ 生成(generate)→ 打分(score)。
examples/llm 三级 LLM 流水线示例:检索(retrieve)→ 生成(generate)→ 打分(score)。
internal
fakeclock
Package fakeclock 是测试专用的假时钟:时间只在调 Advance 时前进, 测试不真 sleep,时序完全确定(宪法第 V 条)。
Package fakeclock 是测试专用的假时钟:时间只在调 Advance 时前进, 测试不真 sleep,时序完全确定(宪法第 V 条)。
sqlbroker
Package sqlbroker 是 PostgreSQL / MySQL 两个服务器型后端的共享核心:基于标准库 database/sql,把两库"真正不同"的点收进 Dialect(见 dialect.go),其余标准 SQL 一份。
Package sqlbroker 是 PostgreSQL / MySQL 两个服务器型后端的共享核心:基于标准库 database/sql,把两库"真正不同"的点收进 Dialect(见 dialect.go),其余标准 SQL 一份。
Package memorybroker 是 Broker 的内存参考实现:单进程、单 sync.Mutex + sync.Cond。
Package memorybroker 是 Broker 的内存参考实现:单进程、单 sync.Mutex + sync.Cond。
Package mysqlbroker 是 taskgate 的 MySQL 后端:database/sql + go-sql-driver/mysql(纯 Go 免 cgo)。
Package mysqlbroker 是 taskgate 的 MySQL 后端:database/sql + go-sql-driver/mysql(纯 Go 免 cgo)。
Package pgbroker 是 taskgate 的 PostgreSQL 后端:database/sql + pgx(stdlib 模式,纯 Go 免 cgo)。
Package pgbroker 是 taskgate 的 PostgreSQL 后端:database/sql + pgx(stdlib 模式,纯 Go 免 cgo)。
prototype
identity
Package identity 是 Identity 领域模型的原型验证层(见 docs/plans/2026-07-16-Identity领域模型.md),不进正式代码。
Package identity 是 Identity 领域模型的原型验证层(见 docs/plans/2026-07-16-Identity领域模型.md),不进正式代码。
quota
Package quota 是 Quota 领域模型的原型验证(见 docs/plans/2026-07-16-Quota领域模型.md),不进正式代码。
Package quota 是 Quota 领域模型的原型验证(见 docs/plans/2026-07-16-Quota领域模型.md),不进正式代码。
Package redisbroker 是 Broker 的 Redis 后端:所有"多步读写必须原子"的操作 都收进单段 Lua 脚本执行(宪法 III:终态更新与子任务唤醒同一段脚本收敛), 语义以 memorybroker 为基准,由 brokertest 的 18 条契约统一验收。
Package redisbroker 是 Broker 的 Redis 后端:所有"多步读写必须原子"的操作 都收进单段 Lua 脚本执行(宪法 III:终态更新与子任务唤醒同一段脚本收敛), 语义以 memorybroker 为基准,由 brokertest 的 18 条契约统一验收。
Package sqlitebroker 是 Broker 的 sqlite 文件后端:纯 Go 驱动(modernc.org/sqlite,免 cgo), WAL 模式单文件落盘。
Package sqlitebroker 是 Broker 的 sqlite 文件后端:纯 Go 驱动(modernc.org/sqlite,免 cgo), WAL 模式单文件落盘。

Jump to

Keyboard shortcuts

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