队列藏在我们自己的接口后面
状态: 已在 v0.1.76 发布(2026-09-27)—— 设计与落地记录见 StandMeet 仓库的 docs/design/event-bus-outbox-webhooks.md。
上层(各域、admin 面板、MCP)只依赖我们自己的接口:events.Recorder,以及 internal/infra/jobs 里的 jobs.Jobs、jobs.Inspector、jobs.Runtime。River 是任务接口现在的实现;以后换实现,上层代码不动。换得掉的是传输;outbox 换不掉——不管底下是什么,事件都先和业务改动同事务写进 outbox。
| 接口 | 方法 | 实现 |
|---|---|---|
events.Recorder(固定) |
With(tx) Recorder、Record(ctx, ownerID, type, subject, data) error |
outbox 写入器(Postgres),唯一实现 |
jobs.Jobs |
With(tx) Jobs、Enqueue(ctx, kind, args, EnqueueOpts) (JobID, error) |
River |
jobs.Inspector |
Overview(kind)、List(filter)、Get、Retry、Cancel、Periodic、RunPeriodic |
River |
jobs.Runtime |
Jobs + Inspector + Wait(ctx, id, max)、Start、Stop |
River |
jobs.Handler |
func(ctx, args json.RawMessage) error,幂等 |
任务种类与订阅 |
- 没有单独的
Transport接口。 订阅是数据(events.Subscription),总线把它变成任务种类;relay 通过jobs.Jobs入队。换传输只换internal/infra/events里的 relay,域代码不动。 - 任务种类(
jobs.Kind)和周期任务(jobs.Periodic{Name, Every, Run})都是数据,构造 runtime 时交进去。没有Periodic(...)注册方法。 EnqueueOpts有RunAt(延后执行)和UniqueByArgs(同一 args 只一个任务,webhook 扇出用)。- 状态:
pending、running、retryable、completed、discarded、cancelled。
classDiagram
direction TB
class Recorder {
«interface · never changes»
+With(tx) Recorder
+Record(ctx, ownerID, type, subject, data) error
}
class Jobs {
«interface · swappable»
+With(tx) Jobs
+Enqueue(ctx, kind, args, opts) JobID
}
class Inspector {
«interface · swappable»
+Overview(kind) / List(filter) / Get(id)
+Retry(id) / Cancel(id) / Periodic() / RunPeriodic(name)
}
class Runtime {
«interface · swappable»
+Wait(ctx, id, max) State
+Start(ctx) / Stop(ctx)
}
class Bus {
«infra/events · relay loop»
+Kinds() one per Subscription
+Periodics() retention and sweep
}
class OutboxRecorder {
«Postgres · only implementation»
INSERT events in caller tx
}
class RiverAdapter {
«infra/jobs/river · now»
}
class Domains {
«upper layer»
corpus / access / owner …
}
class TasksPanel {
«upper layer»
admin · MCP tasks.*
}
Recorder <|.. OutboxRecorder
Runtime --|> Jobs
Runtime --|> Inspector
Runtime <|.. RiverAdapter
Domains ..> Recorder
Domains ..> Jobs : Enqueue
Domains ..> Bus : Subscription data
Bus ..> Jobs : relay enqueues
TasksPanel ..> Inspector
保持可换的规则
- 接口只认字符串 + JSON:
type、kind是字符串,data、args是 JSON。在 River 上,kind 通过一个原始 JSON 参数类型保持为字符串,它的KindAliases()列出所有声明过的 kind(internal/infra/jobs/river/worker.go)。River 的泛型参数类型不漏到上层。 - 契约不随实现变:至少一次投递、handler 必须幂等、不保证跨事件顺序。
- 门禁保证不漏出去:
check-queue-behind-port.sh——只有internal/infra/jobs/river可以 importgithub.com/riverqueue/**(测试也算),没有排除清单。 - 一致性套件写明契约:约 27 条 UT 覆盖
Jobs/Inspector/Runtime的契约,现在对 River 跑,将来换实现原样跑。套件 + 全套 e2e 验收仍绿,就是“上层不感知”的证据。
放在哪
internal/infra/jobs 放端口、队列表、失败分类和 DefaultBackoff。internal/infra/jobs/river 是适配器:River 自己的迁移器(jobsriver.Migrate,启动时紧跟 pgstore.Migrate 执行)、worker、周期任务、inspector 和 Wait。两者都不认识任何域。internal/infra/events 放 Recorder、relay、保留期清理和 webhook 签名(event-model)。
River 简介
River(riverqueue.com)v0.47 是跑在 Postgres 上的 Go 任务队列,原生 pgx/v5,许可证 MPL-2.0,与 AGPL 兼容。
| River 能力 | 我们拿来做什么 |
|---|---|
InsertTx(事务内入队) |
任务和引起它的变更一起提交 |
SKIP LOCKED 抢占 |
多进程的 worker 不会拿到同一个任务 |
| LISTEN/NOTIFY 唤醒 | 任务不必等下一次轮询 |
| 退避重试 | 唯一的重试主人(retry-has-one-owner) |
| 按参数唯一入队 | 每个(端点,事件)只一个 webhook.deliver;不用来做合并(relay-claims-rows-not-cursor) |
| 选主的周期任务 | 周期任务只跑一份,不是每进程一份(relay 不在其中) |
| rescuer | 卡在 running 超过阈值的任务被救回重试 |
| 清理服务 | completed 24 小时、discarded 7 天后删除(River 默认值) |
| 任务行本身就是可查日志 | tasks-panel 读它 |
- 不加服务:River 只在现有 Postgres 里多几张表(
river_job、river_leader等)。 - River 的 DDL 归 River:启动时由它自己的迁移器建表,
schema.sql不抄一份(抄了下次升级 River 就会漂移)。 - 不引入 River UI:它是独立服务 + 自己的认证。tasks-panel 参考它的信息结构,用 dispatcher ops 实现。
为什么队列放在 Postgres 上:why-not-a-broker。