本文是“Phorge 现代化改造实战”系列第十五篇。上一篇为 Go 服务增加了统一的 Phorge API 入口;这一篇进一步接管有状态的后台任务,把队列存储、任务调度和仍留在 PHP 中的任务实现拆开迁移。

系列导航

  1. 从没有官方镜像到 Docker Compose 跑起来:建立最小可运行基线
  2. 改进容器化的七个细节:补齐权限、持久化、依赖和探活
  3. 接入 Stargate:把 Forward Auth 的信任边界做完整
  4. 拆分模块到 Gorge:无侵入改造不等于不碰文件
  5. 替换 diff 子进程:兼容不等于逐字一致
  6. 替换实时通知服务:为什么 HTTP 501 反而表示正常
  7. 迁移邮件服务:先分清哪些失败不该重试
  8. 迁移搜索服务:写进索引不等于搜得到
  9. 迁移文件存储:写得进去也要读得回来
  10. 迁移 Webhook 投递服务:先解决重复投递
  11. 升级 Gorge 的 HTTP 框架:接口没变,行为也不能变
  12. 兼容 Elasticsearch 5、6、7:版本配置决定整个索引结构
  13. 联调六项外部服务:容器在运行,不代表业务已经切换
  14. 为 Go 服务建立统一的 Phorge API 入口:网关可以换框架,Conduit 协议不能变
  15. 用 Go 接管 Phorge 工作队列:如何避免新旧消费者互相抢任务。

写在前面

前几篇迁走的高亮、搜索、发信和文件存储,大多可以概括成“Phorge 发一个 HTTP 请求,Gorge 返回结果”。工作队列还包含状态、租约、重试和消费权。两个 Go 服务解决代码实现,Phorge 新接口完成接线;接管还需要停掉旧的 taskmaster,并证明任务只由预期的消费者处理。

这次对应 Gorge 2026.09.08-r4Phorge 2026.09.08-r4。Gorge 一侧通过 PR #14 加入 gorge-taskqueuegorge-worker,Phorge 一侧则补上 HTTP 客户端、Conduit 执行入口、配置项、Setup Check 与 Compose 接线。

这次迁移需要继承 Phorge 队列的线协议:ID 从哪里分配、哪几列属于活跃表、租约如何抢占、yield 怎样表示,以及 Conduit 到底接受什么请求格式。

先拆清楚队列、消费者和任务实现

Phorge 原本由 PhabricatorTaskmasterDaemon 不断执行“租约 → 运行 → 归档”。任务主要落在三张表里:

  • worker_activetask:还在排队、等待重试或正在被租用的任务;
  • worker_taskdata:任务负载,与活跃任务通过 dataID 关联;
  • worker_archivetask:成功、永久失败或取消后的归档记录。

r4 把这条链拆成两个二进制:

flowchart LR
  P["Phorge PHP"] -->|"入队与队列操作"| Q["gorge-taskqueue :8090"]
  Q --> S[("MySQL 或 Redis")]
  W["gorge-worker :8170"] -->|"租约与回报"| Q
  W -->|"worker.execute"| C["gorge-conduit :8150"]
  C --> P

gorge-taskqueue 管队列状态,gorge-worker 管消费循环。两者保留 HTTP 边界,便于独立扩缩容和隔离故障:

组件 负责什么 不负责什么
gorge-taskqueue 入队、租约、重试、yield、取消、归档、队列统计 不执行具体业务任务
gorge-worker 领取任务、选择 handler、并发执行、回报结果 不直接读写 MySQL 或 Redis
Phorge PHP 产生任务;暂未 Go 化的任务仍在 PHP 中执行 配置接管后不再直接驱动队列

因此,一个 taskqueue 前可以挂多个 worker,也可以给重任务单独开一个带 TASK_CLASS_FILTER 的 worker 池。若把两者塞进同一进程,worker 要么绕过 HTTP 直接碰 Store,要么仍然兜一圈本机网络;两种做法都削弱了独立扩缩容和故障隔离。

这也解释了为什么 render 与 diff 可以合在一个二进制里,而 taskqueue 与 worker 不应合并:前者是无状态计算,后者之间存在明确的状态所有权和部署边界。

taskqueue 的难点集中在旧表语义

gorge-taskqueue/api/queue 下提供 enqueueleasecompletefailyieldcancelawakenstats 与任务查询接口。HTTP 层较薄,难点集中在底下三张旧表特殊的 ID、租约和归档语义。

worker_activetask.id 不能直接用 LastInsertId()

最初按常规思路插入 worker_activetask,MySQL 会直接报:

Error 1364 Field 'id' doesn't have a default value

原因是 PhabricatorWorkerActiveTask 使用 IDS_COUNTER。它的 id 由 Phorge 在应用层通过 lisk_counter 分配,没有使用 AUTO_INCREMENT。Go 侧如果另建一条自增序列,即使短期能写进去,也会和原生 phd 从两个来源取号,最终发生碰撞。

修复后的 enqueue 在同一个事务里复用 Phorge 的计数器写法:

func nextCounterValue(
    ctx context.Context,
    tx *sql.Tx,
    counterName string,
) (int64, error) {
    res, err := tx.ExecContext(ctx,
        `INSERT INTO lisk_counter (counterName, counterValue)
         VALUES (?, LAST_INSERT_ID(1))
         ON DUPLICATE KEY UPDATE
           counterValue = LAST_INSERT_ID(counterValue + 1)`,
        counterName)
    if err != nil {
        return 0, err
    }
    return res.LastInsertId()
}

这里虽然仍调用了 LastInsertId(),拿到的却是 LAST_INSERT_ID(...) 主动发布的计数器结果。worker_activetask 没有自增值;worker_taskdata 仍然使用 AUTO_INCREMENT,所以同一次入队事务里的两张表采用不同的 ID 机制。

这条约束最容易被“统一风格”的重构破坏。统一目标是与 Phorge 共用同一条序列,代码形状可以不同。

归档时要丢弃 failureTime

任务成功、永久失败或取消时,taskqueue 会在一个事务中把活跃行写入 worker_archivetask,再从 worker_activetask 删除。这样任务不会同时出现在两张表,也不会在两个操作之间凭空消失。

归档表与活跃表的字段并不对称。归档表没有 failureTime:这个字段只用于活跃任务的重试退避,归档时会被丢弃。归档行的结构可以概括成:

active columns - failureTime + result + duration + archivedEpoch

若把 failureTime 也塞进归档 SQL,MySQL 会报 Unknown column 'failureTime',已经执行完的任务留在活跃表,租约过期后又会被领走。

另一个容易写错的是 result。它使用 Phorge 定义的整数编码:0 表示成功,1 表示失败,2 表示取消。写成字符串后,Go 侧虽然能写入,Phorge 的管理界面却无法按原有语义读取。

没有 status 列,状态都藏在租约字段里

租约分两阶段完成:先取从未租过的任务,再用过期任务补足名额。

-- 第一阶段:新任务优先
SELECT id
FROM worker_activetask
WHERE leaseOwner IS NULL
  AND leaseExpires IS NULL
ORDER BY priority ASC, id ASC
LIMIT ?;

-- 第二阶段:租约已经过期的任务
SELECT id
FROM worker_activetask
WHERE leaseExpires < ?
ORDER BY priority ASC, id ASC
LIMIT ?;

priority 数字越小越紧急,同优先级按 id 保持 FIFO。选出候选行后,还要用带原条件的 UPDATE 取得租约,再只返回本次调用成功持有的行。这个二次条件非常重要:SELECT 只给出候选列表,不能证明所有权。

Phorge 的表里也没有单独的 status。yield 通过 leaseOwner="(yield)" 这个哨兵和未来的 leaseExpires 表示,awaken 再精确识别它并提前结束等待。字段名、空值语义和哨兵字符串都属于兼容协议,不能按 Go 结构体的内部习惯自由重命名。

MySQL 和 Redis 的回退能力不同

taskqueue 的 Store 同时有 MySQL 与 Redis 两个实现。Redis 用 hash、sorted set 和 Lua 脚本复现优先级、两阶段租约、yield 与归档移动,让并发操作保持原子性。worker 完全不知道后端是哪一个,它只通过 HTTP 领取和回报任务。

“配置缺失就回到原生 SQL”只描述了 PHP 的代码分支,两种后端的回滚能力仍有明显差异:

后端 数据放在哪里 清空 gorge.taskqueue.uri 后会怎样
MySQL 与 Phorge 共用 {namespace}_worker 三张表 原生 SQL 能继续看见原队列,适合直接切回
Redis Gorge 自己的 Redis key Phorge 只会回到 MySQL;Redis 中未完成的任务不会自动搬过去

所以 MySQL 模式可以理解为“同一份数据,换一个调度入口”;Redis 模式则涉及数据面迁移。切到 Redis 前要定义排空、迁移或回放方案,不能因为代码有 fallback 就假定数据也能无缝回退。

这里还有一个配置层的小口值得补上:r4 的 newStore()redis 之外的所有值都回落到 MySQL。把 GORGE_TASKQUEUE_BACKEND=redsi 拼错后,程序不会在启动时明确拒绝,而会尝试连接 MySQL。更稳妥的做法是只接受 mysqlredis,未知值直接启动失败,让错误尽量靠近配置现场。

worker 不碰数据库,未重写的任务经 Conduit 回到 PHP

gorge-worker 的核心是一个租约循环:

  1. 从 taskqueue 批量领取任务;
  2. taskClass 从 registry 找 handler;
  3. 在租约截止时间内执行;
  4. 成功走 complete,永久失败走 fail(permanent=true),yield 走 yield,其余错误进入临时失败与重试;
  5. 回报短暂失败时最多重试五次,仍失败则保留租约,等待正常的过期恢复。

本地 handler 目前只原生实现少量任务,例如 FeedPublisherHTTPWorker。配置 GORGE_WORKER_CONDUIT_URL 后,registry 会安装一个 fallback,把未在本地实现的任务统一交给 Phorge 新增的 worker.execute 方法。

于是执行路径形成一个回环:Phorge 产生任务,Go 负责排队与调度,尚未迁出的业务代码仍由 PHP 执行,结果再由 Go 回报给 taskqueue。这个过渡结构很实用,因为“替换队列运行时”和“重写所有任务业务逻辑”不必一次完成。

Conduit 不接受 JSON body

委派最初报的是:

invalid character '<' looking for beginning of value

这个错误通常说明客户端准备解析 JSON,服务器却返回了 HTML 错误页。Phorge 的 PhabricatorConduitAPIController 不接受 application/json 请求体;它要求表单编码,并把参数 JSON 放进 params 字段。

Go 客户端最终按 Phorge 自家客户端的线格式发送:

params["__conduit__"] = conduitMeta

paramsJSON, err := json.Marshal(params)
if err != nil {
    return nil, err
}

form := url.Values{}
form.Set("params", string(paramsJSON))
form.Set("output", "json")

req, err := http.NewRequestWithContext(
    ctx,
    http.MethodPost,
    baseURL+"/api/"+method,
    strings.NewReader(form.Encode()),
)
req.Header.Set("Content-Type", "application/x-www-form-urlencoded")
req.Header.Set("Accept", "application/json")

网关自己的认证仍放在 X-Service-Token。遇到非 JSON 响应时,客户端会返回 HTTP 状态和截断后的正文片段,不再只留下一个脱离上下文的字符解析错误。

worker.execute 只执行任务,不接管队列

Phorge r4 新增 PhabricatorWorkerExecuteConduitAPIMethod。它根据 taskClassdata 构造一个临时的 PhabricatorWorkerActiveTask,取得真实 worker 实例并调用 executeTask(),然后把结果分类返回:

PHP 执行结果 返回给 Go taskqueue 动作
正常返回 success complete
PhabricatorWorkerYieldException yield yield
PhabricatorWorkerPermanentFailureException permanent-failure 永久 fail 并归档
其它异常 failure 临时 fail,等待重试

临时对象只给 PHP worker 提供 task ID、类型和负载上下文,不写数据库。租约、失败次数与归档仍由 gorge-worker 通过 taskqueue 完成,避免两边同时做账。

PHP 侧保留原生路径,但“调用失败再写 SQL”要分后端看

Phorge 保留了原生 SQL 路径,并在几个关键入口检查 PhabricatorGorgeTaskQueueClient::isConfigured()

  • PhabricatorWorker::scheduleTask():入队;
  • PhabricatorWorker::awakenTaskIDs():唤醒 yield 的任务;
  • PhabricatorWorkerLeaseQuery::execute():租约或列表查询;
  • PhabricatorWorkerActiveTask:完成、失败与 yield 回报。

以租约为例,结构很清楚:

if (PhabricatorGorgeTaskQueueClient::isConfigured()) {
  if ($this->skipLease) {
    return $this->executeListViaGoService();
  }
  return $this->executeLeaseViaGoService();
}

// 原生 SQL 路径继续保留。

这让未配置 Gorge 的安装维持原行为,也让 MySQL 共享表模式具备直接回退能力。

不过完成、失败与 yield 的通知函数在 HTTP 调用失败时,会尝试回到本地 SQL。这个设计在 MySQL 后端有意义,因为服务与 PHP 看到的是同一行;在 Redis 后端,PHP 手里的任务对象只是临时对象,MySQL 里没有对应活跃行,本地 SQL 无法替 Redis 完成归档。Redis 模式应把 taskqueue 调用失败视为需要恢复的基础设施故障,不能依赖 SQL 兜底。

接管开关:为什么要把 phd.taskmasters 设为 0

GORGE_TASKQUEUE_URI 非空时,entrypoint 会写入 gorge.taskqueue.uri/token,并在默认开关下执行:

bin/config set phd.taskmasters 0

这里需要纠正一个容易流传的说法:**旧 taskmaster 和 gorge-worker 并存,并不意味着每个任务必然执行两次。**两边如果都正确遵循同一套条件更新和租约规则,通常会竞争并瓜分任务;同一行只能由一次带条件的更新成功租用。

仍然应该关掉旧消费者,理由是所有权必须唯一:

  • 两套运行时会争抢任务,无法预测某个 task class 最终在哪个环境执行;
  • 监控、容量规划和故障定位被拆成两套;
  • 灰度结果会被竞争比例污染,难以证明新 worker 覆盖了预期流量;
  • 租约过期、执行超时或结果回报失败时,任何“至少一次”队列都可能重试并产生重复副作用,这与是否同时部署两种消费者是另一层问题。

因此,phd.taskmasters=0 用来完成明确的消费权交接,解决的是所有权问题。r4 也提供 GORGE_TASKQUEUE_DISABLE_PHD_TASKMASTER=0 作为灰度逃生开关;使用它时需要接受混合消费,不宜当作常态扩容方案。

建议按下面顺序切换:

  1. 启动 taskqueue,并验证实际后端;
  2. 启动 worker 与 conduit,确认 worker 能领取并回报测试任务;
  3. 写入 Phorge 的 taskqueue 配置,让新入队和队列操作改走 HTTP;
  4. phd.taskmasters 设为 0,确认旧消费者退出;
  5. 观察活跃队列、归档队列与业务副作用,再宣布接管完成。

其中第 3、4 步应在同一个切换窗口内连续完成,尽量缩短新旧消费者同时竞争队列的时间。先把新链路验证好,再交接消费权,既不会在准备阶段留下无人处理的任务,也不会把混合消费长期当成正常状态。

这正是系列第十三篇反复强调的区别:服务“在运行”“可达”“就绪”“已被选择”和“业务生效”是五个不同状态。

启动顺序:避免循环依赖,也别高估健康探针

{namespace}_worker 数据库由 Phorge 的 bin/storage upgrade 创建。若让 Phorge 等 taskqueue 健康,而 taskqueue 又要先连到尚未创建的数据库,就会形成首次启动死锁。

r4 的 Compose 因此这样编排:

  • taskqueue 先等 MySQL 与授权初始化完成;
  • Phorge 对 taskqueue、worker 使用 service_started,不等它们健康;
  • worker 可以等 taskqueue 的 service_healthy
  • taskqueue 和 Redis 客户端都采用惰性连接,后端暂时不可达时进程仍可启动。

这里还要再收紧一次表述:taskqueue 的 /readyz 对 MySQL 只执行 PingContext(),不会查询 worker_activetask 是否存在。因此:

探针 能证明什么 不能证明什么
taskqueue /healthz HTTP 进程活着 后端可用、表已创建
taskqueue /readyz MySQL 或 Redis 可以连接 三张 worker 表一定存在、完整队列操作已通过
worker /readyz /healthz 等价,进程活着 能连 taskqueue、正在消费任务
worker /api/worker/stats 当前进程的处理计数 全队列状态;计数重启后会清零

所以首次启动后,至少要真实跑一次 enqueue → lease → complete,才能证明 schema、凭据、命名空间与队列协议同时正确。只看到两个绿色健康检查,还不足以宣告队列已接管。

Conduit 回调修了 Host 匹配,但信任边界仍需闭合

worker 把任务委派回 Phorge 时,gorge-conduit 默认访问 http://phorge:80。上游请求带 Host: phorge,若它不在 phabricator.base-uriproduction-uriallowed-uris 中,Phorge 会返回 Site Not Found。r4 的 entrypoint 会从 GORGE_CONDUIT_UPSTREAM_URL 解析主机并写入 phabricator.allowed-uris,解决这条内网回调链路。

这里有两个部署注意点。

第一,r4 生成的是一个只包含网关上游的数组,并通过 bin/config set --stdin phabricator.allowed-uris 写回。已有环境若本来维护了多个 allowed URI,应先合并旧值,避免启动时覆盖原列表。

第二点影响更大:worker.execute 当前设置了 shouldRequireAuthentication() = falseshouldAllowUnguardedWrites() = true。代码的安全前提是“调用只能经过带 X-Service-Token 的 gorge-conduit”,但这个前提并没有由该 Conduit method 自己验证。allowed-uris 只解决 Host 匹配,无法认证调用者。

因此,生产环境至少要阻止外部绕过网关,直接访问 Phorge 的这个内部入口。更稳妥的后续实现,是让 worker.execute 自己验证独立的服务凭据,或使用受限的 Conduit 身份,并对可执行的 taskClass 做白名单控制。机器间调用没有 CSRF 风险,仍然需要身份认证;这两件事不能混在一起。

验证闭环:不要只数容器,要数状态转移

这次验证分成三层。

第一层是 Go 的单元与契约测试。taskqueue 通过可注入的内存 Store 验证 HTTP 契约,MySQL 和 Redis 分别测试租约、归档与错误路径;worker 则验证类过滤、fallback、结果分类与回报重试。平台分层测试也加入 internal/taskqueueinternal/worker,防止公共层反向依赖业务域。

第二层是仓库里的 tests/e2e/taskqueue.sh,跑通真实 HTTP 的入队、租约、完成等状态转移。它比单纯访问健康探针更有价值,因为只有真实操作才会碰到库名、表结构与字段约束。

第三层是在完整 Phorge + Gorge 栈上,用 bin/worker flood 注入真实的 PhabricatorTestWorker

  • taskqueue 能持续领取任务;
  • gorge-worker 经 gorge-conduit 调用 worker.execute
  • 停止注入后,worker_activetask 下降、worker_archivetask 上升;
  • 成功归档行的 result=0,字段名是 failureCount,不是 failedCount
  • lisk_counterworker_activetask 的计数持续递增,证明 Go 与 PHP 共用同一 ID 来源;
  • 清空 taskqueue 配置后,MySQL 模式能重新走原生 SQL 入队。

最有说服力的证据来自状态转移与所有权:谁入队、谁领取、谁执行、谁归档,每一步都能对上。

其他

把 Phorge 的工作队列搬到 Go,核心在于讲清三种所有权:队列状态由谁保存,租约由谁授予,业务代码由谁执行。

r4 已经搭出了一个实用的过渡结构:Phorge 仍然产生任务,gorge-taskqueue 统一管理状态,gorge-worker 接管调度,尚未迁移的 task class 再经 Conduit 回到 PHP。它允许队列运行时先独立出来,而不要求一次性重写所有 worker。

增加两个容器还不足以代表服务化完成。完整交接还要满足:旧消费者明确退出、后端回退边界可解释、真实任务完成一次闭环、镜像发布件确实存在,以及内部执行入口没有可绕过的信任缺口。

最后

对这类“在旧系统旁边接入新实现”的改造而言,代码只提供新路径;租约、表结构、启动顺序、身份验证和回滚方式共同决定这条路径能否长期可信。

—EOF