本文是“Phorge 现代化改造实战”系列第十篇。前九篇从容器化、入口信任和模块边界一路写到六类外部能力;这一篇处理第一个主动读取 Phorge 数据库、消费共享队列的 Gorge 服务,并把重点放在任务所有权与重复投递上。

系列导航

  1. 从没有官方镜像到 Docker Compose 跑起来:建立最小可运行基线
  2. 改进容器化的七个细节:补齐权限、持久化、依赖和探活
  3. 接入 Stargate:把 Forward Auth 的信任边界做完整
  4. 拆分模块到 Gorge:无侵入改造不等于不碰文件
  5. 替换 diff 子进程:兼容不等于逐字一致
  6. 替换实时通知服务:为什么 HTTP 501 反而表示正常
  7. 迁移邮件服务:先分清哪些失败不该重试
  8. 迁移搜索服务:写进索引不等于搜得到
  9. 迁移文件存储:写得进去也要读得回来
  10. 迁移 Webhook 投递服务:先解决重复投递

写在前面

前面几次迁移都有一个共同形状:Phorge 收到请求,再通过 HTTP 把一项能力交给 Gorge。语法高亮、diff、实时通知、发信、搜索和文件存储虽然实现不同,主要调用方向始终是 PHP 到 Go。

Webhook 是第一个例外。

Phorge 仍然负责判断什么事件应该触发哪个 Hook,也仍然把待投递请求写进自己的数据库;Gorge 不提供“提交一个 Webhook”的 API,而是直接轮询那张表、抢占任务、发出请求,再把结果写回原来的记录。

Gorge 2026.09.07-r4 将原来的 Webhook 服务迁入单仓,新增 gorge-webhook 二进制、数据库 Store、投递调度器、只读诊断接口、5 份契约固件和一组单元测试,具体变化可以从 Gorge r3…r4 的代码差异 中查看;Phorge 2026.09.07-r4 则新增投递交接开关、配置检查、PHP worker 守卫和队列表索引,对应的宿主改动集中在 Phorge r3…r4 的代码差异 中。

独立服务的核心只有 6 个文件、784 行代码。搬进单仓后,难点却不是改 import 路径,而是回答几个共享队列才会出现的问题:一个任务被 SELECT 出来,是否已经属于当前进程;旧消费者什么时候必须停;全局静默配置怎样越过进程边界;数据库尚未初始化时,谁应该等谁;以及如何证明“没有重复投递”不是偶然。

这些问题做错之后,大多不会报错。接收方只是多收到一份请求,队列只是多停留几十秒,或者日志只是安静得像什么都没有发生。

先划清边界:数据库表就是交接协议

gorge-webhook 处理的是:

{namespace}_herald.herald_webhookrequest

Phorge 负责入队:

  • 判断一次事务是否命中 Herald 规则;
  • 找到对应的 Webhook;
  • 生成 request 行和 properties
  • 决定 retrynever 还是 forever

Gorge 负责出队:

  • 轮询 status = 'queued' 的候选行;
  • 抢占其中一行;
  • 读取 Hook 的 URI 和 HMAC key;
  • 生成与 Phorge 兼容的请求体并发送;
  • 把状态、HTTP 结果和错误信息写回同一行。

它不创建表、不修改表,也不需要 DDL 权限。herald_webhookrequest 的 schema、状态值、索引和垃圾回收仍由 Phorge 管理。

这也解释了为什么 Webhook 的 HTTP API 只有两个只读端点:

方法 路径 返回内容
GET /api/webhook/stats queued、sent、failed 和可用 Hook 数量
GET /api/webhook/hooks 包含禁用项在内的 Hook 总数

没有“立即投递”或“写入队列”接口。给外部再开一个入队入口,会让同一张队列表出现第二个写者,破坏 Phorge 对事件、权限和重试策略的控制。

这两个端点也不返回 Hook 列表,只返回数量。Hook 记录里包含目标 URI 和 HMAC key,完整管理面应该继续留在 Phorge 的权限体系内。诊断接口知道“有几个”,不应该知道“密钥是什么”。

Webhook 因此是单仓里第一个不由入站请求驱动工作的域。两个 API 都正常,只能证明诊断面可访问,不能证明轮询循环正在工作;相反,诊断接口暂时不可用,也不一定代表投递停止。真正的工作发生在后台 goroutine 和数据库之间。

第一个问题:SELECT 出来不等于已经拿到

独立服务原来的候选查询很直接:

SELECT ...
FROM herald_webhookrequest
WHERE status = 'queued'
ORDER BY id ASC
LIMIT ?

查询结果随后进入 goroutine 异步投递,只有请求完成后才更新状态。投递期间,这一行仍然是 queued

默认轮询间隔是 1 秒,单次 HTTP 投递最长允许 15 秒。只要接收方响应超过一个轮询周期,下一个 tick 就会再次取到同一行。第一次请求还在路上,第二次已经开始。

这个问题不需要多副本才能出现。单个进程、单个数据库连接池,只要采用“查询候选、异步处理、完成后改状态”的结构,就可能重复处理同一任务。

直觉里的“单实例不需要抢占”依赖一个没有说出口的前提:任务被取走后已经离开待办集合。轮询队列不具备这个性质。查询只得到候选,claim 成功才得到任务。

不能增加 claimed,也不能偷偷加列

常见修法是把状态改成 claimed,或者给表增加 lease_until。在 Phorge 的表上,两条路都不合适。

status 的合法值由 Phorge 定义:

const STATUS_QUEUED = 'queued';
const STATUS_FAILED = 'failed';
const STATUS_SENT = 'sent';

界面会根据这些值渲染图标,HeraldWebhookWorker::doWork() 还会在入口检查状态必须等于 queued。Go 侧增加第四个值,不只是扩展枚举,而是同时改变 PHP worker 与 UI 的协议。

给表加一列看起来影响更小,实际上更危险。Phorge 的 Lisk ORM 会根据 PHP 类声明维护 schema。没有出现在模型配置里的列,会被 bin/storage adjust 视为 surplus column,并在之后某次调整中删除。服务可能运行数周都正常,直到一次看似无关的维护把租约列删掉,重复投递才突然回来。

这类故障难查,不是因为删除没有日志,而是删除发生在另一个时间、另一条命令和另一位操作者手里。

dateModified 做一次乐观抢占

最终选择复用已有的 dateModified

  • Lisk 本来就维护它;
  • Phorge 页面不展示它;
  • Webhook 垃圾回收根据 dateCreated 判断保留时间;
  • 关闭 PHP 投递后,Go 服务是投递过程中的唯一更新者。

claim 使用一条 compare-and-set:

UPDATE herald_webhookrequest
SET dateModified = GREATEST(dateModified + 1, UNIX_TIMESTAMP())
WHERE id = ?
  AND status = 'queued'
  AND dateModified = ?

只有 RowsAffected() == 1 才表示当前进程抢到。返回 0 说明另一进程已经推进版本,或者这一行不再是 queued,当前 goroutine 必须什么都不发。

这条做法不需要长事务,也不依赖 SELECT ... FOR UPDATE SKIP LOCKED。后者要求 MySQL 8.0 或 MariaDB 10.6,而一个已有的 Phorge 部署使用什么数据库版本不能由新服务假设。

这里最容易被简化掉的是:

dateModified + 1

dateModified 只有秒级精度。如果一行在同一秒内被创建并 claim,直接写 UNIX_TIMESTAMP() 可能把它更新成原值。第二个抢占者手里的旧版本仍然匹配,于是两边都认为自己成功。

GREATEST(dateModified + 1, UNIX_TIMESTAMP()) 同时保证两个性质:数据库时间已经前进时跟随真实时间;多个动作落在同一秒时,版本仍然至少增加 1。它不是防御性写法,而是乐观锁成立的前提。

一个只在“同一秒内多次竞争”时触发、表现为多发一次 HTTP 请求的问题,很难靠普通功能测试碰到。r4 为它单独保留了测试,直接断言同秒 claim 也必须改变版本。

新任务为什么不需要先等 30 秒

候选查询不是简单加一个 dateModified <= cutoff

WHERE status = ?
  AND (dateModified = dateCreated OR dateModified <= ?)
  AND (lastRequestResult <> ? OR lastRequestEpoch <= ?)
ORDER BY id ASC
LIMIT ?

第一组时间条件是 claim lease。dateModified <= cutoff 表示上一次抢占已经过期,可以由别的进程接手;dateModified = dateCreated 则表示这是一行从未被投递进程碰过的新任务,应当立即进入候选集合。

Lisk 插入记录时会把两个时间戳写成同一秒,而 Gorge 的后续更新保证 dateModified 严格变大,所以两者相等可以表达“尚未处理”。

如果只保留 cutoff,新插入的任务也会被误认为正在租约内,平白等待 30 秒。系统不会报错,用户只会觉得 Webhook 比以前慢了。

三个时间窗,回答三个不同问题

r4 中有三套容易混淆的时间条件:

时间窗 默认值 判据 回答的问题
claim lease 30 秒 dateModified 抢到任务的进程死了,多久允许别人接手
请求重试退避 60 秒 lastRequestEpoch 这次投递失败后,多久再试同一请求
Hook 熔断 300 秒内 10 次失败 最近失败计数 一个接收端持续失败,多久暂停整个 Hook

这三者不能合并。

lease 必须覆盖一次完整的投递。默认 HTTP 超时是 15 秒,因此配置中的 lease 即使被设得更小,ClaimLease() 也会将实际值抬到不低于投递超时。否则 lease 在第一个 POST 结束前过期,第二个进程就能重新 claim,同一请求再次被发出去。

请求退避则是迁移时补上的另一条行为。独立服务在 retry = forever 的投递失败后仍把记录留在 queued,下一个 1 秒 tick 立即重试;默认熔断阈值是 10 次,于是十秒左右就能把一个 Hook 打进 300 秒熔断。迁入后单个失败请求默认等待 60 秒,与 Phorge worker 的重新租赁节奏保持一致。

熔断针对的是 Hook,不是某一行。一个 endpoint 在 300 秒窗口内累计 10 次失败后,属于它的请求暂时不再发出。窗口是滑动的,旧失败自然滑出后恢复,不需要额外的 reset 状态。

此外,同一个 Hook 的请求会被串行化,不同 Hook 才共享全局并发上限。Phorge 原来的 worker 也会按 Hook 加锁。原因不是顺序好看,而是避免一个积压的 Hook 瞬间发出大量并行请求,让熔断器还没来得及计数就已经把接收方压垮。

分开写的条件,不一定真的独立

lease 与请求退避在概念上是两套机制,SQL 里也写成了两个条件。但一次失败回写会同时刷新:

dateModified
lastRequestEpoch

候选查询又要求两条条件同时通过,所以实际重试间隔是:

max(claim lease, retry backoff)

默认值是 30 秒和 60 秒,最终 60 秒生效,实测两次失败重投间隔分别约为 60.01 秒和 59.98 秒,符合预期。

隐患出现在调参时。把 GORGE_WEBHOOK_RETRY_BACKOFF_SEC 调到 30 秒以下,不会得到更快的重试,lease 会把它静默抬回去;即使设置为 0,默认配置下也要等 30 秒。

这和 ClaimLease() 显式将 lease 抬到投递超时不同。后者能在读配置时算出实际值,前者是两个查询条件合取后的结果,目前没有配置校验或日志告诉操作者。

把两个概念拆成两个变量,只是完成了代码分离;只有它们各自依赖不被对方更新的状态,行为才真正独立。 r4 保留了当前实现,同时把有效值与根治方向记录进 findings,避免以后把“调小不生效”误判成环境问题。

接管必须是替换,不能新旧并跑

对无状态服务,先上线新实例、观察一段时间再切流量很常见。共享队列的消费者不能这样灰度。

PHP worker 与 Gorge 同时运行时,两边从同一张表读取同一行。最终结果不是吞吐翻倍,而是接收方收到两份内容和签名都相同的 POST。接收方无法判断这是消费重复,还是业务事件真的发生了两次。

因此 gorge.webhook.uri 与其他 Gorge 地址的语义不同:

配置项 实际作用
gorge.webhook.uri Webhook 投递的接管开关,同时供 setup check 探测服务
gorge.webhook.token 只保护 /stats/hooks 两个诊断端点

投递本身不通过这个 URI,也不使用 service token。Gorge 直接读数据库,发向接收方的签名使用每个 Hook 自己的 HMAC key。

Phorge 侧有两个守卫:

// 新请求入队后,不再创建 HeraldWebhookWorker 任务。
HeraldWebhookRequest::queueCall();

// 已经存在的 worker 任务运行时,也要再次退出。
HeraldWebhookWorker::doWork();

两处都调用同一个 PhabricatorGorgeWebhookClient::isDeliveryDelegated(),而不是分别读取配置。第一处挡住正常队列路径,第二处处理切换前已经排队的旧任务,以及 bin/webhook call 这类绕过标准调度的路径。

这意味着启用没有“第二步”:URI 一旦成功写进 Phorge,PHP 侧立即停止投递。反过来,容器已经运行而 URI 没写进去,是最危险的失配状态——PHP 与 Go 会同时排空同一队列。

entrypoint 因此会在写入 gorge.webhook.uri 失败时给出专门告警,明确说明重复投递风险。它仍然不让整个站点启动失败,因为“站点完全起不来”比“Webhook 可能重复”影响更大,但这个状态不能再只留下一句普通的配置失败。

全局静默模式必须留在交接条件里

接管谓词不只判断 URI:

public static function isDeliveryDelegated() {
  if (!self::isConfigured()) {
    return false;
  }

  if (PhabricatorEnv::getEnvConfig('phabricator.silent')) {
    return false;
  }

  return true;
}

phabricator.silent 是 Phorge 的全局配置,开启后邮件和 Webhook 都不应发出。Gorge 只读数据库里的请求记录,看不到这个配置;它能看到的 properties.silent 是单次事务的标记,不代表整个站点。

如果接管条件只判断 URI,一个处于全局静默状态的站点仍会被 Gorge 正常投递。

正确做法不是再给 Gorge 增加一个需要手工同步的环境变量,也不是让 Go 去读取 Phorge 的 local.json,而是让静默状态继续留在 PHP 路径。原 worker 会按照既有逻辑把请求标记为 failed / silent;Gorge 只查询 queued,自然不会碰到它。

守卫的位置同样重要。在 HeraldWebhookWorker::doWork() 中,全局静默检查必须先执行,委托守卫随后执行:

先处理 silent -> 请求变成 failed -> Gorge 不会读取
先判断委托   -> worker 直接退出    -> 请求仍是 queued -> Gorge 会发送

这是这次迁移里很有代表性的一条经验:新组件无法读取的旧系统全局状态,必须成为交接条件,而不能假设接管完成后再补。

HMAC 保护的是字节,不是 JSON 对象

迁移投递服务还要保持接收方看到的请求逐字节兼容。请求体继续使用 Phorge 的格式:

encoded, err := json.MarshalIndent(payload, "", "  ")
return string(encoded) + "\n", nil

两空格缩进、字段顺序和末尾换行都不是显示偏好。签名直接对这串字节计算:

mac := hmac.New(sha256.New, []byte(key))
mac.Write([]byte(payload))
signature := hex.EncodeToString(mac.Sum(nil))

结果放进:

X-Phabricator-Webhook-Signature

如果把格式化 JSON 改成紧凑 JSON,语义完全相同,HMAC 却完全不同。一个按照 Phorge 原实现验证签名的接收方会立即拒绝请求,而 Gorge 只能看到对方返回非 2xx,很难从表面判断是签名算法还是 JSON 格式发生变化。

几个看起来可以“清理”的细节也必须保留:

  • triggerstransactions 为空时是 [],不是 null
  • action.epoch 使用 request 的 dateCreated,重试仍然描述同一事件;
  • object.type 从 PHID 第二段提取,例如 PHID-TASK-... 得到 TASK
  • HMAC 输出是小写十六进制。

因此这里适合整值测试,而不是只断言 JSON 能解开。r4 的测试直接比较完整 payload 字符串,并验证签名确实基于包括末尾换行在内的原始字节。

回写结果也必须服从 Phorge 的值域

服务会把投递结果写回 Phorge 原来的字段:

字段 典型值 用途
status queued / sent / failed 请求生命周期与 UI 图标
lastRequestResult none / okay / fail 重试和熔断统计
lastRequestEpoch Unix 秒或 0 最近投递时间
properties.errorType hook / http / timeout UI 的错误分类
properties.errorCode HTTP 状态码或错误短码 请求详情

配置类失败与投递失败不能混在一起。Hook 被禁用、Hook 不存在或 request properties 无法解释时,请求会进入 failed,但 lastRequestResult 使用 nonelastRequestEpoch 使用 0。熔断只统计真正发出请求后得到的 fail;否则几条配置损坏的记录也会让一个健康 Hook 被判定成接收端故障。

成功投递也会写 errorType = httperrorCode = 200。这看起来不太自然,却是 Phorge 原 worker 的既有行为,界面会显示“HTTP Status Code / 200”。省略它不会影响功能,却会让同一页面上 PHP 与 Go 产生的记录长得不一样,给排查制造一条假线索。

properties 是整列覆盖写入,所以 Go 的 RequestProperties 还必须带回自己不使用的 transactionPHIDstriggerPHIDs。只反序列化关心的字段再写回,会把 Phorge 请求详情页依赖的内容清空。

key_status 不能只存在于数据库里

轮询查询按 status 过滤,再按 id 排序:

WHERE status = 'queued'
ORDER BY id ASC

原表没有覆盖 status 的索引,而 sent 记录会被垃圾回收保留 7 天。随着流量增加,每秒一次的轮询会不断扫描越来越大的保留窗口。

r4 使用:

key_status(status, id)

这个索引既支持状态过滤,也支持同一顺序读取最早请求。它不改变查询结果,只改变查询代价。

但 autopatch 创建索引只完成了一半。它还必须写入 HeraldWebhookRequest::CONFIG_KEY_SCHEMA。否则数据库里明明存在,Lisk 却认为它不属于模型;以后运行 bin/storage adjust 时,它会被当作 surplus key 删除。

这与不能偷偷增加租约列是同一类问题:在一个由 ORM 管理 schema 的旧系统里,数据库真实结构和框架声明必须同时更新,只做 DDL 不是完成。

就绪探针没有错,错的是谁在等待它

Webhook 服务没有本地或内存后端。数据库就是工作本身,所以 /readyz 只有一个条件:数据库能连接。

它不检查 herald_webhookrequest 表是否存在,也不会在 sql.Open() 后立即 ping,避免把数据库启动稍慢变成容器重启循环。看起来这样已经避开了首次初始化的循环依赖。

问题藏在 DSN 里:

user:pass@tcp(mysql:3306)/phabricator_herald?...
                          ^^^^^^^^^^^^^^^^^^^

go-sql-driver/mysql 在握手阶段就会选择数据库。即使 /readyz 只执行 ping、不查询任何表,只要 phabricator_herald 尚未创建,连接仍会失败:

GET /healthz -> 200
GET /readyz  -> 503 Error 1049: Unknown database 'phabricator_herald'

数据库服务器已经健康,不代表 DSN 指向的数据库已经存在。这个库要由 Phorge 容器里的 bin/storage upgrade 创建。

如果编排让 Phorge 等待 gorge-webhook: service_healthy,依赖就闭环了:

Phorge 等 Gorge ready
Gorge 等 phabricator_herald 存在
phabricator_herald 等 Phorge 执行 storage upgrade

修复不应该削弱 /readyz。库不存在时,服务确实无法投递,这个信号是正确的。真正需要调整的是依赖方:Phorge 只要求 gorge-webhook 已启动,不要求它已经 ready。

phorge:
  depends_on:
    gorge-webhook:
      condition: service_started

这样 Gorge 先启动,Phorge 随后创建数据库与表,Webhook 服务的 /readyz 再从 503 变成 200。首次启动会出现一段短暂 unhealthy 窗口,这是预期状态,不是需要重启的故障。

r4 还用同一个发现修正了 file-storage 的编排。它的 DSN 指向 phabricator_file,即使不检查 file_storageblob 表,ping 同样会因为库不存在而失败。之前“数据库服务器独立,所以 ping 安全”的判断只说对了一半:依赖不只由主机名决定,DSN 里的库名也是依赖图的一部分。

正确工作时,反而最难观察

claim 成功时,没有重复请求;claim 竞争失败时,当前 goroutine 什么都不做。这正是正确结果,也意味着最重要的行为表现为“什么都没有发生”。

代码里有一条专门日志:

slog.Debug("WEBHOOK_CLAIM_LOST", "request", req.PHID)

但仓库没有配置 slog 级别,Go 默认只输出 Info 及以上。这条能直接证明 claim 正在阻止重复投递的日志,在默认部署中看不到。

另一端则相反。Hook 进入熔断时,每个候选行、每个轮询 tick 都会记录一条 Warn:

slog.Warn("WEBHOOK_HOOK_IN_ERROR_BACKOFF", ...)

默认每秒轮询一次,并发上限 8;一个故障 Hook 带着积压时,大约会产生 8 条 Warn/秒,一天接近 69 万条。

需要的证据被 Debug 隐藏,不需要逐行重复的信息却用 Warn 刷屏。问题不在某一行该升还是该降,而在日志级别的判断标准:不能根据“这段代码重要不重要”分配,而应根据“线上排查需要怎样的证据”分配。

这两点在 r4 中仍是已知保留项。更合适的方向是给日志增加可配置级别,并对熔断跳过做聚合、采样或状态变化日志,例如只在进入和退出熔断时记录一次。

停止服务时,先停轮询再关数据库

Webhook 也是第一个在 srv.Run() 之外长期运行后台循环的域。启动时多了一个 goroutine:

ctx, stopPolling := context.WithCancel(context.Background())
polling := make(chan struct{})

go func() {
  defer close(polling)
  webhook.NewDispatcher(store, cfg).Run(ctx)
}()

收到 SIGINT 或 SIGTERM 后,HTTP 服务器先退出,随后取消轮询 context,等待在途投递完成,最后关闭数据库连接池:

runErr := srv.Run()

stopPolling()
<-polling
store.Close()

顺序不能倒过来。先关连接池会让已经发出的 POST 无法写回结果,request 仍停在 queued;lease 过期后它会再次投递,于是一次正常发布也可能制造重复请求。

监听端口失败时,整个进程仍然退出。技术上可以让投递 goroutine 继续工作,但那会留下一个仍在发请求、外部却无法查询健康状态的实例,编排也无法正确替换它。

怎么验证一个“没有发生”的结果

这个域的测试要分三层看。

HTTP 契约

5 份契约固件只覆盖诊断面:/stats/hooks、鉴权、query token 和数据库不可达。它们能保证字段名与 {data, error} 信封稳定,却碰不到一次真实投递。

单元测试

真正的兼容边界主要由单元测试保护:

  • store_test.go 固定候选查询与 claim SQL 的形状;
  • dispatcher_test.go 固定完整 JSON 字节、HMAC、结果状态与步骤顺序;
  • 多实例测试验证只有 claim 成功者会发出 POST;
  • 时间可注入,使 lease、退避和熔断无需等待真实分钟数。

直接断言 SQL 字符串在这里不是“测试实现细节”。status、两个时间条件、CAS 版本与 ORDER BY id 的组合就是并发协议,改写其中一个条件会改变谁有资格发出请求。

端到端与人工验证

现有 e2e 脚本可以验证服务启动、只读端点和统计不变量,但不会向 Phorge 表里伪造投递记录,因此不能证明出站 POST 的字节和次数。

完整验证需要从 Phorge 创建一个 Hook,让接收端保留原始 body 和请求头,再检查:

  1. 一个事件只收到一次;
  2. body 使用两空格缩进,并以换行结束;
  3. X-Phabricator-Webhook-Signature 等于对完整 body 计算的 HMAC-SHA256 小写十六进制;
  4. 接收端返回 500 时,同一 request 默认约 60 秒后重试;
  5. 达到阈值后,整个 Hook 进入约 300 秒熔断。

第 2、3 项必须基于原始字节。若接收器先把 body 解析成对象,再重新序列化,缩进与换行差异会被抹掉,测试会在签名已经不兼容时仍然通过。

最后

共享队列最麻烦的地方,是成功通常表现为“只发生了一次”,而重复、延迟和漏投递都可能被接收方当成业务本身。迁移这类后台消费者时,应该先设计所有权和可观察性,再搬代码。否则新服务看起来运行正常,只是世界的另一端悄悄多收到了一份请求。

–EOF