队列藏在我们自己的接口后面

更新于 · 在 sijie.xyz 查看原条目 ↗

状态: 已在 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 可以 import github.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。