[feat][evaluation] Centralized experiment scheduling: OSS seams, quota accounting, and reclaim paths - #646
Open
xueyizheng wants to merge 58 commits into
Open
[feat][evaluation] Centralized experiment scheduling: OSS seams, quota accounting, and reclaim paths#646xueyizheng wants to merge 58 commits into
xueyizheng wants to merge 58 commits into
Conversation
xueyizheng
force-pushed
the
feat/evalx-evaluation
branch
from
August 31, 2026 13:12
0f191a7 to
a476a67
Compare
为中心化实验 item 调度铺地基,本提交只加字段与校验,不改任何派发行为:
legacy 路径读到的仍是 legacy/1,行为与改动前一致。
IDL:
- domain/expt.thrift: 加 ExptTriggerType Evalx 常量(该 typedef 是 string 而非
enum,追加常量无兼容风险);新增 ExpectedResourceConsumption /
ExpectedQuotaConsumption 结构;Experiment 读视图开 116~119 段位回显
priority_level / scheduler_mode
- coze.loop.evaluation.expt.thrift: Create/Submit 两个 Request 各加
92 priority_level、93 expected_quota_consumption、94 scheduler_mode。
94 故意不加 api.body —— 与既有 trigger_type 同为服务端内部覆写字段,
公网调用方不得自行指定 enforce,只允许 commercial wrapper 按灰度白名单填入。
DDL(四个文件,非 spec 原述的三处):
- docker-compose init-sql + patch-sql、helm init-sql + init-sql/xxx_alter.sql
各加 priority_level / scheduler_mode 两列与 idx_scheduler_queue
- 索引不以 space_id 打头(跨空间调度队列扫描不应带 space_id),与现有 12 个
以 space_id 开头的索引形态不同,已在 SQL 注释写明理由
- priority_level DESC 需 MySQL 8.0+,低版本会静默退化为升序,注释中标注上线前
须确认实例版本并 EXPLAIN 验证
- 新列同时写入 CREATE TABLE 与 alter 两处;注意仓库既有漂移(eval_set_space_id /
target_space_id / eval_set_access_level 仅在 alter、不在 CREATE TABLE),
本次不顺带修复以免混淆 diff
Go:
- entity/expt_dispatch_mode.go: ExptDispatchMode 常量与 Normalize/IsValid/
IsCentralDispatch 收敛函数。未知模式一律按 legacy 处理 —— 安全侧是走旧链路,
而非让实验既跳过 legacy 闸又拿不到 reservation
- entity/expt_quota_consumption.go: Validate 校验非空/amount>0/键唯一/禁 wildcard,
Normalize 去空白(带空白的 key 在调度期拼 constraint key 时会匹配不上上限配置,
静默降级成「未登记资源」而被放行)
- Experiment DO 加 PriorityLevel / ExptDispatchMode。字段名刻意避开
SchedulerMode —— entity.ExptSchedulerMode 已被「实验跑法调度器」占用,
同名会让两个无关概念在阅读时混淆
- eval_conf 走 json.Marshal(非 thrift binary),故 ExpectedQuotaConsumption
挂在 EvaluationConfiguration 上零 DDL、零代码生成,老数据反序列化后为 nil
顺带修一处会咬到后续开发的隐患:
- QuotaSpaceExpt 新增 Clone() 并让 repo 改用它。原 repo 手工构造
&QuotaSpaceExpt{ExptID2RunTime: maps.Clone(...)} 只覆盖当时唯一的字段,
一旦该 struct 加第二个字段就会在每次 CreateOrUpdate 被零值静默写回 Redis。
GORM model/query 手工补两列(gen 程序 UseDB 从真实库反查,需先应用 DDL 才能跑);
query gen 本就缺三个跨空间列、与 model 已不同步,本次只保证新列两侧对齐。
测试: entity 新增两个测试文件覆盖全部分支;
go build ./modules/evaluation/... 通过,相关包 go test 全绿,gofmt 干净。
Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
中心化调度需要按全局优先级跨空间取候选实验,而既有 IExperimentRepo.List 全部要求 space_id。若按空间分别扫描,低优空间的实验会先于高优空间被处理,全局优先级语义即失效, 因此新增一个不带 space_id 的扫描方法。 - entity/expt_scheduler_queue.go:扫描参数与 keyset 游标 用 keyset 而非 offset 分页:候选集合在扫描期间会被并发写入(新实验提交、实验进终态), offset 在这种场景下会漏掉或重复元素 - mysql/expt.go:DAO 实现,走 idx_scheduler_queue 用裸 gorm 而非 gen DSL —— keyset 的三元组比较是带括号 OR 的复合条件,gen 链式 API 要嵌套多层 Or(),可读性差得多 未加 FORCE INDEX:status IN (...) 是 range 条件,MySQL 可能因此放弃用索引满足 ORDER BY 而走 filesort;灰度期数据量小可接受,留给 EXPLAIN 实测再定 —— 过早 FORCE 会在数据分布 变化后选到更差的计划 条件含 latest_run_id > 0,排除"只 Create 尚未 Run"的实验(没有 run 可派 item,扫进来白跑) - expt_repo_impl.go:PO2DO 转换。单条实验 eval_conf 损坏时跳过并告警,不让整拍失败 —— 一条脏数据不应永久阻塞所有实验的调度。evaluator refs 传 nil,调度不需要,避免 N+1 查询 mock 手工补 ScanSchedulerQueue 而非重跑 go generate:本地 mockgen 版本较新,全量重生成会 带来 ~1000 行无关 diff(参数名重命名 + isgomock 字段),淹没本次真实改动。 go build ./modules/evaluation/... 通过,相关包 go test 全绿,gofmt 干净。 Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
两处最小接缝,legacy 实验行为完全不变(读到的 mode 是 legacy,两条分支都短路)。 1. domain/component/central_reservation_guard.go:窄 port ICentralReservationGuard 只有 ConfirmRunning / Release 两个方法。完整账本、调度算法、Adapter 都在商业版, 由 Wire 注入;开源部署注入 noop。 noop 的 ConfirmRunning 刻意 fail-closed 返回 false:本方法只会被 enforce 实验触达, 而 enforce 意味着"额度由中心账本管控"。没有账本却放行等于零约束跑 item —— 这种失败是 静默的(资源打爆才发现);拒绝执行则可见(实验不动会被察觉)。 Release 永不返回错误:终态收口不应因额度模块缺席而失败。 2. item consumer 新增 HandleCentralReservation 中间件 插在 HandleEventCheck 之后、HandleEventLock 之前:Check 已排除终态 run(无需额度校验), 放在 Lock 之前可避免为一条注定丢弃的消息去抢 item 锁。 **模式判定回查 experiment.scheduler_mode DB 列,不看 event 上的任何标记** —— 若模式随 event 传递,字段丢失或取零值时 central 消息会被当作 legacy 处理,跳过校验、静默绕过额度; 这个方向的失败无声,比多查一次 DB 危险得多。 Guard 用 setter 注入而非构造参数:该构造函数已有一个 variadic 参数(Go 不允许第二个), 且它是可选依赖。 reservation 缺失 → 丢弃消息;账本报错 → 返回 error 让 MQ 重试(item 已被预占, 丢弃会让它停在 Queueing 白等一轮 reservation 超时)。 3. 旧 per-experiment daemon 加 enforce 薄分支(防双驱动) enforce 实验在此丢弃 toSubmit,但**保留其余全部职责**:完成 item 归档、zombie / sandbox terminated 处理、run/实验终态收口、NextTick 续跳。若直接 return,实验会因为 没人收口而永远停在 Processing。 必须丢弃而非"让它也派":旧链路按实验自身配置并发补 item,既不看全局优先级也不经额度账本, 两个驱动并存会直接超发。 go build ./modules/evaluation/... 通过;domain/component 与 domain/service 全部测试通过 (后者 227s,覆盖既有全量用例,确认接缝未破坏 legacy 行为);gofmt 干净。 Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
上一版把 centralGuard 做成 setter 注入,理由是"构造函数已有 variadic 参数"。 但这让 Wire 无法自动装配:Guard 被构造出来却没人调 setter,闭环实际是断的 (wire_gen 里能看到 NewCentralReservationGuardAdapter 被调用,但 consumer 的 centralGuard 字段始终为 nil,enforce 消息会走"无 guard → fail-closed"分支被全部丢弃)。 改为放在 variadic 之前的常规参数:Go 允许(variadic 只需是最后一个),Wire 能自动填, 调用点只有 wire_gen(生成物)与一个单测。同时删掉 WithCentralReservationGuard setter, 避免两条注入路径并存造成"到底哪条生效"的歧义。 domain/service/wire.go 注册 component.NewNoopCentralReservationGuard: 该 set 被商业版复用,商业版在自己的 set 里注入真实适配器覆盖它 —— 与既有 NewSandboxAgentNotifier(开源 no-op 桩 / 商业版真实 Lark send)同一模式。 go build ./modules/evaluation/... 通过,domain/component 测试通过,gofmt 干净。 Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Step 1 只改了 thrift 未跑 codegen,导致 kitex_gen 里没有 ExpectedResourceConsumption /
ExpectedQuotaConsumption / Evalx 等符号。下游 cozeloop-gen-commercial 的 expt-ref.go 会
ref_expt 引用这些类型,于是 commercial 编译报 undefined —— 根因在本仓漏了这一步。
生成命令:GOPATH=$(go env GOPATH) bash script/cloudwego/code_gen.sh
(脚本用 ${GOPATH}/bin 作为 loopgen 安装目标;环境未导出 GOPATH 时会解析成 /bin 而权限失败,
故需显式带上。脚本内部 NO_PUSH_REMOTE=true,纯本地生成不推远端。)
仅 6 个生成文件变更,全部对应本次改动的两个 thrift。
Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
修一个我在实现中心调度 Adapter 时引入的真实缺陷。 原实现的 LoadRuntimeState 只把 status=Processing 算作活跃 item,还写了段注释自我说服: 「reservation 侧的活跃 item 由账本幂等保证,不必查 Redis 求并集,收益仅是 deficit 精确一点点」。 这个推理是错的:已在 Redis 预占但尚未被 consumer 消费的 item,在 MySQL 里仍是纯 Queueing, 于是它既不计入并发占用、又会再次进入待授予队列。后果不是「精确一点点」—— - deficit 被持续高估:并发 5 的实验即使 5 个 item 都已预占待消费,deficit 仍算成 5 - 每拍都为同一批 item 重复申请额度(ReserveBatch 幂等挡住了重复扣额,但没挡住重复申请) - 高优实验反复占用授予机会,挤掉本该拿到额度的低优实验 解法(按 spec 7ae163c):在 MySQL 侧显式表达「已预占」。 DDL(四处:docker init/patch + helm init/alter,两两 diff 一致) - quota_reservation_state TINYINT UNSIGNED NOT NULL DEFAULT 0 - idx_expt_run_dispatch(space_id, expt_id, expt_run_id, status, quota_reservation_state, id) 末位带 id 让 keyset 分页与稳定排序都走索引,避免 filesort entity - QuotaReservationState 只有 none/reserved 两值,**不是** Redis 状态机的复制品 (Redis 侧有 Reserved/Dispatched/DispatchUncertain/Running 四个非终态) Redis reservation 仍是账本真值;本列只是调度投影,不进 IDL/OpenAPI/Stats - 未知值按未预占处理:安全侧是让它进候选被重新预占,而非当成已占用而永不派发 IExptItemDispatchRepo(新建窄接口,不往 IExptItemResultRepo 加方法) - ClaimQuotaReserved:Queueing/none → Queueing/reserved,Redis 预占与 MQ 发布之间的必经关卡。 **逐个 CAS 而非批量 UPDATE**:批量只能拿到总 RowsAffected,无法知道哪些成功; 而调用方必须精确知道才能只发布成功项、并释放失败项的 Redis reservation - ResetQuotaReserved:MQ 明确失败时退回。条件带 status=Queueing —— 若已被 consumer 推进到 Processing 说明消息其实投递成功了,退回投影会让它被重复授予 - LoadDispatchRuntime:**一次查询**取回占用与候选。分两次查会在两次之间漏掉刚从 Queueing 变 Processing 的 item(第一次查它还是 Queueing、第二次只查 Queueing 已查不到), 导致占用少算、超发 - StartReservedItem:Queueing/reserved → Processing/none。清掉 reserved 标记, 因为 Processing 本身已代表占用,留着会让对账看到自相矛盾的状态 - MGetDispatchObservations:供分钟级对账识别四类漂移 分类逻辑抽成纯函数 classifyDispatchRuntime 以便直测 —— 它是本次正确性核心。 14 个用例覆盖:Processing 计占用、Queueing/reserved 计占用且不进候选、Queueing/none 唯一候选、 混合场景、limit 只约束候选不约束占用、Processing 带脏 reserved 不重复计数、nil 跳过、 chunk 分批(含 size<=0 防死循环)。 go build ./modules/evaluation/... 通过,新增测试全绿,gofmt 干净。 Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
…ng/none Guard 此前只做了 Redis 侧 ConfirmRunning,但 spec 要求 consumer 还要把 run log 投影推进到 Processing/none。缺这一步,item 已开始执行而投影仍停在 Queueing/reserved:下一拍把它当 「已预占未消费」继续计入占用(这本身没错),但一旦 reservation 因超时被清理,它就变成 「既不 Processing、也无 reservation」的孤儿,对账要多绕一圈才能修。就地兑现让两侧同步收敛。 几处刻意的错误处理: - 投影写失败返回 error 让 MQ 重试:额度已预占且 reservation 已转 Running, 丢弃消息会让这份额度占着直到超时清理 - CAS 未命中(started=false)不阻断执行:可能是重复投递(已 Processing)或投影已被 repair 修正;此时 reservation 校验已通过说明额度是真的,继续执行安全 - dispatchRepo 为 nil 时跳过:保持 legacy 路径与开源部署不受影响 repo wire set 注册 NewExptItemDispatchRepo + NewExptItemDispatchDAO,wire 自动接入 consumer。 mock 手写而非跑 mockgen:本地 mockgen 版本较新,全量重生成会带来大量无关 diff (参数名重命名 + isgomock 字段),淹没真实改动 —— 与前几次同样处理。 go build ./modules/evaluation/... 通过,consumer 构造测试通过,gofmt 干净。 Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
线上与所有 PPE 泳道共用同一个 MySQL 库,而中心调度是"跨空间扫全局 → 抢租约 → 扣额度"的后台任务。缺少所有权边界时,泳道实例会扫出线上的 enforce 实验、为其 预占额度并把 item 发进泳道 topic,由泳道 consumer 执行、结果写回共享库 —— 线上侧对此毫无感知。这是数据污染,不是资源浪费。 引入 scheduler_scope 作为调度所有权与 Priority 排序边界: - entity.Experiment.SchedulerScope:创建时冻结的不透明稳定 ID,Retry 继承, legacy 为空串;业务代码不得解析该字符串(泳道/空间/App/Region 只是生成规则的输入) - SchedulerQueueScanParam.SchedulerScope 必填,DAO 下推为 WHERE 等值条件。 空值直接报错而非退化成扫全表 —— 这是挡住越界的物理闸门,宁可可见地报错, 不可静默地污染 - 四份部署 SQL 同步加列,索引改为 (scheduler_mode, scheduler_scope, status, deleted_at, priority_level DESC, created_at, id)。 scope 插在第 2 位而非追加末尾:它是等值条件,必须排在 range 条件 status IN(...) 之前,否则 MySQL 无法用索引满足 ORDER BY - GORM model 同步加字段,priority_level 的索引 priority 从 4 改为 5 (scope 占用第 2 位,不改会与 status 撞位、生成的列序与 SQL 不一致) Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
BOE 真机 EXPLAIN(MySQL 8.0.27-18-ndb)证明原索引后 3 列不产生任何收益: status IN (2,3) → type=range, Extra="Using index condition; Using where; Using filesort" status=3 单值 → type=ref, Extra 无 filesort 真实查询用 status IN (Pending,Processing),这是 range 条件。MySQL 用索引满足 ORDER BY 的前提是排序列不位于 range 列之后,所以 priority_level/created_at/id 放进索引消不掉 filesort,只让每个索引条目更宽、写放大更高。单值那组证明索引本身 没问题 —— filesort 来自查询里的 IN,改索引解决不了。 索引收敛为 (scheduler_mode, scheduler_scope, status, deleted_at): - deleted_at 保留:查询含 deleted_at IS NULL,索引内判掉可省回表 - id 不写:InnoDB 二级索引条目天然带主键,显式写是冗余 排序开销可接受:候选集是当前活跃的 enforce 实验,量级几十到几百而非全表。 彻底消除 filesort 需把 IN 拆成两条等值查询再归并,但那要引入双流归并游标, spec §1.4 已否决该方案。 同步修正 GORM model 的索引 tag(此前 status/deleted_at 未登记进 idx_scheduler_queue,与 SQL 不一致)。 Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
调度侧的 WHERE scheduler_scope 只防住"泳道去调度线上实验",防不住反方向: item MQ 的泳道路由依赖 producer 的 x_tt_env tag,而该 tag 会因环境变量缺失、 消息重投、broker 配置差异而失效。届时一条 PPE 的 item 消息可能被线上 consumer 取到(或反之),而两个环境共用同一个库 —— 不校验归属就会用一个环境的进程去跑 另一个环境的 item,结果直接写进对方的数据。 新增窄 port ICentralSchedulerScopeOwner 回答"这个 item 该不该由我来跑": - 与 ICentralReservationGuard 分开:Guard 回答"有没有额度",本 port 回答归属。 即使额度充足也可能必须拒绝(消息投错环境) - noop 实现取**放行**(与 Guard 的 fail-closed 相反):单环境部署不存在"别的环境", 拒绝执行只会让所有 enforce item 永久卡住,那是自造故障而非防护 - 解析失败返回 error 而非 false:false 表示"确定不属于我,丢弃",error 表示 "无法判定,请重试"。混淆二者会让一次环境探测抖动静默丢弃本该执行的 item, 而它已预占额度、要等 reservation 超时才回收 HandleCentralReservation 增加两道闸: - enforce 实验 scheduler_scope 为空 → fail-closed 丢弃(没有 Scope 就无法确定 去哪本账查 reservation,猜一本账等于用别人的额度跑这个 item) - Scope 不属于本进程 → 丢弃而非报错重试(Scope 不匹配是路由问题,重试只会在同一个 错误进程上再失败;正确的进程会从自己队列拿到消息,或由下一拍重新派发) ConfirmRunning/Release 增加 schedulerScope 参数,由调用方从 DB 读出后传入, 不由实现方按当前运行环境推断 —— 推断会让"实验属于哪本账"取决于谁在处理消息, 而它本该只取决于数据本身。 新增参数放在 variadic 之前而非 setter 注入:setter 会让 wire 构造出实例但无人 调用 setter、字段恒为 nil,而 nil 在本文件里被解释为"跳过校验",等于静默关闭防护 (centralGuard 此前已踩过一次)。 Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
三个 P0,任一不修中心调度都无法在泳道跑通。 全仓零个 Release 生产调用点:CompleteItemRun / completeItemRunOnUnretriableErr 只写 MySQL,zombie 与沙箱提前终态所在的 ExptSchedulerImpl 压根没有 guard 字段。结果 Redis used 单调递增,第一批额度跑满后全线卡死。症状还很像"额度配小了",第一反应会去调大 TCC, 调大后再次跑满,很难联想到根因。 释放接在两处,不逐个终态分支加: - consumer 侧收在 HandleCentralReservation 出口一处。终态路径有四条(success / fail / 不可重试前置失败 / indebt 终止),逐条手动调释放意味着以后任何人新增分支都可能漏, 而漏掉是静默的额度泄漏。 **不能无条件释放**:MQ 重试也从这里返回,此刻释放会让重投消息在 ConfirmRunning 处 因 reservation 不存在被丢弃、item 永久卡 Processing。因此回查 run log 投影, 只在 IsItemRunFinished 时释放 —— 不靠 execErr 判断(execErr==nil 未必终态, asyncAbort 下 item 仍在跑;execErr!=nil 未必非终态,已被兜底落 Fail)。 - zombie / 沙箱提前终态由 daemon 直接判终态,对应 consumer 消息可能永不返回, 必须单独接(releaseCentralQuotaForItems)。 原 Lua 对 state 不存在的账本置 ready,注释写"首次启用没有 active item"。该前提只在 真·首次启用成立;Redis 重启 / key 淘汰 / 误删同样走这个分支,而此时 MySQL 里一批 Processing item 仍占着资源 → 按 used=0 再授予一遍即超卖。 改为置 rebuilding + 拒绝授予(LedgerNotReady)。代码分不清"首次启用"与"账本丢失", 就取安全侧:首次启用要等一次 rebuild 才能调度(可见可等),超卖是静默的。 GetState 的同款默认值一并修正 —— 它是运维查"账本健不健康"的入口,报 ready 会让人 以为一切正常、也会让对账跳过重建。 IDL 有 priority_level / expected_quota_consumption / scheduler_mode,但业务转换全程 不使用,正常入口落成 legacy / 1 / "" / nil —— 中心调度没有任何候选。 按"trigger=evalx 即可信"接通: - entity.ShouldEnforceByTrigger:只有 EvalX 进 enforce。不按空间白名单 —— 白名单会把 同空间里人手点的实验也拽进 enforce,而那些不申报消耗向量,调度器只能跳过, 表现为"实验建好了但一个 item 都不跑"。按来源区分天然对齐"谁申报、谁被管控"。 - CreateExptParam 加 priority / scope / 消耗向量;ConvertCreateReq 映射 priority 与 向量,**刻意忽略请求里的 scheduler_mode/scheduler_scope** —— 否则任何内部 RPC 调用方 都能自行声明 enforce 并伪造 scope,绕过额度、甚至去动别的环境的账本。 - CreateExpt 冻结:mode 由 trigger 派生,scope 经新增窄 port ICentralSchedulerScopeProvider 解析(解析失败/空值报错,不降级成 legacy —— 降级会让 EvalX 以为受管控、实际走旧链路绕过额度),消耗向量校验后冻结进 eval_conf (调度侧 frozenConsumptionOf 正是从 EvalConf 读)。 - Submit → Create 透传 priority 与向量。 Retry 无需改动:它走 AllowExptRun + 直接发 schedule event,不经 Create 转换, 已冻结的 mode/scope/向量天然被继承。 新增 port 的参数一律放在 variadic 之前而非 setter 注入:setter 会让 wire 构造出实例 但无人调用、字段恒 nil,而 nil 在这些位置被解释为"跳过校验/跳过释放",等于静默关闭防护 (centralGuard 已踩过一次)。 测试:trigger admission 11 例(含大小写/空白容忍、前缀后缀不得误命中)+ IDL 常量一致性; 账本 fail-closed 两例(拒绝授予且 used 保持 0、状态落 rebuilding)。 evaluation 全量 55 个包测试通过。 Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
ScanSchedulerQueue 此前并入宽接口 IExperimentRepo(15+ 方法、十余处依赖), 导致多处手写 fake 因缺少新方法而编译失败 —— 而那个编译失败又掩盖了同包内既有的 测试失败(22 个 _BitsUTGen / region 路由用例),排查时先看到的是编译错误, 真实失败要等编译修好才浮现。 拆成单方法窄接口后,"新增调度能力"不再波及无关调用方:本次一并删掉了上一个 commit 里为 5 个手写 fake 补的桩方法 —— 它们现在完全不需要存在,这正是拆分的收益。 实现方仍是同一个 exptRepoImpl(同表、同 DAO),拆的只是消费侧契约; 中心调度 adapter 只依赖窄接口。 Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
内场实测后判定该索引收益接近零、代价明确,故取消;experiment 表的 idx_scheduler_queue 保留不动(那张表确实需要跨实验按优先级扫描)。 三条依据: 1. 所有 dispatch 查询的 WHERE 恒以 (space_id, expt_id, expt_run_id) 打头 —— ClaimQuotaReserved / ResetQuotaReserved / StartReservedItem 三个再加 item_id, LoadDispatchRuntime / MGetDispatchObservations 走三列前缀。没有任何查询需要 status 或 quota_reservation_state 作索引前缀:中心调度只扫「当前 run 的 run log」, 跨实验扫描发生在 experiment 表、不在本表。 2. 该前缀已被既有索引完全覆盖:uk_expt_run_item_turn(space_id,expt_id,expt_run_id,item_id) UNIQUE 让带 item_id 的精确 CAS 直接定位单行;idx_expt_run_result_state (space_id,expt_id,expt_run_id,result_state) 覆盖不带 item_id 的那两个。 3. 单个 run 的 run log 实测仅 ~900 行(内场最大 914),三列前缀定位后按 status/quota_reservation_state 过滤是内存操作。 代价对比:本表内场 7800 万行(experiment 才 23 万)。为几百行的内存过滤给 7800 万行表 建 6 列复合索引,换来的是在线 DDL 与长期写放大。只加列不加索引后,带默认值的 tinyint 在 MySQL 8.0 走 instant add column,秒级完成、无需 gh-ost(内场库 8.0.27-18-ndb)。 改动 6 处:4 份部署 SQL(docker/helm 双路径已 diff 验证一致,CI mysql-schema-check 会校验)、GORM model 的 index tag、IExptItemDispatchRepo 的接口注释 —— 注释里补上「新增方法请保持该前缀形状」与跨 run 对账的 caveat,避免以后有人按旧假设 写出需要索引的查询却发现索引已不存在。 Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
## 问题(实测确认,不是假想)
exptDAO.Update 用 struct 做 Updates,GORM 只跳过**零值**字段。而 DO2PO 会把未设置的
调度字段 Normalize 成非零值(mode ""→"legacy"、priority 0→1,见 convert/expt.go:61-63)。
实测 DO2PO(&Experiment{ID:123, LatestRunID:456}) 产出
SchedulerMode="legacy" / PriorityLevel=1 —— 两者非零,必然被写进 UPDATE。
全仓有 8 处「只带 ID + 一两个业务字段」的部分更新(LogRun 写 latest_run_id、
ScheduleStart 改 status 等),每一处都会顺手把 enforce 实验改回 legacy、
把 EvalX 申报的优先级重置为 1。
后果静默且严重,三条链路同时失效:
- 中心调度扫不到它了(扫描条件 scheduler_mode='enforce');
- 旧 daemon 的抑制判断 IsCentralDispatch 读到 legacy → **恢复自主派发**,
同一 run 出现两个派发驱动、绕过全局额度账本,正是设计明令禁止的情形;
- consumer 侧 guard 同样读到 legacy,item 执行时不再校验 reservation。
且 scheduler_scope 是零值会被跳过,最终留下 mode=legacy + scope 非空 的不可能组合。
注意代码注释一直声称这些值"创建时冻结",但写路径并未兑现 ——
struct Updates 不构成冻结。
## 修法
DAO 层显式 Omit(schedulingFrozenColumns...):这三列的唯一合法写入点是 Create。
选 Omit 而非"让 DO2PO 不 Normalize":后者依赖零值跳过,而零值跳过正是本次踩的坑,
再靠它一次仍然脆弱(且 Create 会因此需要额外保证)。
## 测试
expt_update_frozen_columns_test.go 用 dry-run gorm 渲染真实 UPDATE 语句做断言 ——
断言 SQL 文本而非转换结果,因为本 bug 的成因正是"转换结果看起来合理、
但 GORM 会把它写进 UPDATE",只测转换层测不出来。
已验证:临时移除 Omit 后该测试失败并点名 priority_level 与 scheduler_mode,
即它确实守住了不变量而非恒真。
另加一条测试把冻结列名单钉死,将来新增冻结列漏加会失败。
(harness 需 SkipDefaultTransaction:GORM 的 Updates 默认包事务,sqlmock 会拒绝 Begin。)
## 数据影响
无需修数据:线上 fornax_evaluation 目前尚无这三列(DDL 未执行),
本修复先于任何 enforce 实验落地,不存在已损坏的行。
evaluation 全量 55 个包测试通过。
Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
trigger 判据(trigger_type=evalx)之上的第二道闸,语义是 AND 而非 OR: 只有 trigger 已判定为 EvalX(因而调用方一定申报了 expected_quota_consumption)时 才咨询本 policy,由它决定"这个空间/这类评测对象是否纳入本轮灰度"。 为什么必须是收窄而不能扩大:enforce 实验强制要求资源消耗向量(缺向量在创建期即报错), 而非 EvalX 入口(控制台手动 / OpenAPI / 定时)目前没有传向量的字段。若 policy 能把这些 入口的实验也拽进 enforce,结果要么创建报错、要么被调度器永远跳过 —— 后者表现为 "实验建好了但一个 item 都不跑",且那条分支是静默 return。等 OpenAPI 具备申报能力后 才可考虑放宽成 OR。 CentralAdmissionSubject 只带"这个实验是什么"(space/target 类型与 ID),不带 "该不该 enforce"—— 后者是 policy 的职责。字段用基础类型而非 entity.EvalTarget: policy 实现在 commercial,让它依赖 OSS 领域实体会把整个 target 模型拖进配置层。 开源部署注入 noop(恒定放行):开源侧不产生 EvalX trigger,本 policy 不会被咨询到。 与 Guard 的 noop 取 fail-closed 相反 —— 那里放行会绕过额度,这里拒绝只会让 单环境部署的所有 enforce 实验退回 legacy,是自造故障。 Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
trigger_type=evalx 之上的第二道闸:只有 trigger 已判定为 EvalX 时才咨询 policy, 由它按空间 / 评测对象类型 / 评测对象 ID 决定"这个实验是否纳入本轮灰度"。 语义是 AND(收窄)而非 OR —— 理由见 ICentralAdmissionPolicy 的 port 注释。 CreateExpt 的判定链变成: trigger==evalx → policy 允许 → 校验并冻结消耗向量 → 解析并冻结 scope → enforce 任一环不成立即落 legacy(policy 不匹配是灰度常态,记 Info 不记 Warn)。 policy 返回 error 时**不降级放行**:那意味着配置读取本身坏了,放行会让实验在无灰度 控制的情况下进 enforce。与"不匹配"(正常结果,落 legacy)严格区分。 新增 EvalTargetType.ConfigName() 提供 snake_case 稳定名,供人工维护的 TCC 白名单使用。 与 String() 分开是受众不同:String() 是驼峰名给日志/错误用;ConfigName() 与 OpenAPI 公开枚举("sandbox_agent")和 CLI --target-type 取值一致,配置里写名字比写枚举数字 可读且不会因枚举调整而失配。未登记类型返回空串 —— 调用方应视为"不匹配任何配置项", 不得回落默认类型,否则新增类型会意外命中既有灰度规则。仅记录型(*Online)不给配置名: 它们不执行评测对象、不会进调度。 evaluation 全量 55 个包测试通过。 Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
IDL 早就预留了 domain/expt.thrift 的 116/117 并注明"中心化调度读视图", 但 ToExptDTO 从未赋值 —— 字段定义好了、管子没接,调用方只能查库才知道自己的实验 有没有被中心调度纳管、申报了多少额度。OpenAPI 面连字段都没有。 本次接通两套读模型: - domain/expt.thrift 补 118 expected_quota_consumption - domain_openapi/experiment.thrift 补 116~118 + 对等结构 ExpectedQuotaConsumption (另立一份而非 include domain,与 RunModeConfig/ExptEvalSetSourceType 同套模式,避免符号冲突) - ToExptDTO / DomainExperimentDTO2OpenAPI / entity 直转路径三处赋值 两个刻意的取舍: 1. legacy 实验也回显 priority_level / scheduler_mode,不做"仅 enforce 才给"的裁剪。 放量期最高频的疑问是"我这个实验为什么没进中心调度",那时最需要看的恰恰是一个 legacy 实验的 scheduler_mode —— 回显 legacy 就一眼确认没纳管,不必去猜是 trigger 没带 evalx、灰度范围没命中还是代码没生效。若只有 enforce 才给,调用方还无法分辨 "这实验是 legacy"与"这接口版本不支持该字段"。二者都有 DB 默认值,天然非空。 2. scheduler_scope 一律不进读模型。它是不透明调度域 ID(形如 fornax_cn_prod), 对调用方没有可用语义却泄露部署拓扑;业务代码本就不允许解析 Scope 字符串, 回显只会诱使调用方依赖这个不稳定契约。内部运维需要时直接查表。 TestToExptDTO_SchedulerScopeNeverExposed 用整个 DTO 的字符串形式做断言, 而不是只查某个已知字段 —— 将来新增字段若误带 scope 也会被挡住。 expected_quota_consumption 取"有则回显、无则省略":legacy 确实没申报, 省略比返回空结构更如实;调用方据此区分"没申报"与"申报了空向量"(数据异常)。 entity→openapi 走 entity→domain→openapi 两跳而非直转,复用 Normalize* 的 历史/异常取值收敛(0→1、越界夹取、非法模式→legacy),避免规则复制两份后漂移。
一、额度泄漏(必然发生,非罕见组合)
额度闸(HandleCentralReservation)在中间件链内层按 run log 状态决定是否归还额度,
而有两条终态路径是在它**外层**写的,届时它已经返回:
① 不可重试的前置失败:BuildExptRecordEvalCtx 等阶段失败时,item 由
HandleEventErr 里的 completeItemRunOnUnretriableErr 兜底落 Fail。
额度闸判定时 item 还是 Processing,正确地保留了 reservation,之后无人释放
② 欠费终止:把整个实验落 Terminated 却完全不动 item run log,
因此额度闸按状态判定**永远**不会释放
既无 reservation TTL 清理也还没有对账,所以这两条路径上的额度永久泄漏。现象是
"额度慢慢跑满后整个 Scope 再也调度不动",看起来像上限配小了,极难反推。
修法:额度闸取得执行权后往 ctx 挂预占凭据,HandleEventErr 在**落 Fail 之后**据凭据补一次
释放。顺序不能反——先释放会出现"额度已还但 item 仍算 Processing"的窗口,下一拍按虚高的
占用少派 item。与额度闸内的释放重复调用是安全的(Redis HDEL + used 有下限保护)。
不把 Scope/guard 泄进 HandleEventErr 签名,是因为那一层不该知道中心调度的存在。
二、取执行权移进 item 锁内
原顺序是 ConfirmRunning → StartReservedItem → 抢 item 锁 → 执行。正确性此前由幂等兜住
(Lua 对已 running 返回 1、StartReservedItem 是 CAS、Release 是 HDEL),不会重复扣减,
但两条并发消息会各自发一次无谓的 Redis 写和一次注定失败的 CAS,且"谁在执行"这个事实被
拆到了锁的两边。
拆成两层而不是简单调换顺序:
- HandleCentralAdmission(锁外):纯读,"这条消息该不该由我处理"——模式/Scope 非空/
Scope 归属/guard 已注入。不该处理的消息不必去抢锁
- HandleCentralReservation(锁内):写,取执行权 + 兑现投影
为什么不把 Lock 直接挪到最外层:那会让归属校验也进锁,一个路由错误的进程要先抢到
item 锁才发现不归自己。两层之间靠 ctx 传已查到的实验,省掉锁内第二次 GetByID。
测试 14 个。RunsInsideLock 断言的是调用序列本身
(lock → confirm → start_reserved → exec → release → unlock),
以后有人把顺序改回去会失败而不是静默通过;另有反向用例钉住"重试路径不得释放"
——那个方向的错误更隐蔽:会让重投消息因 reservation 不存在被丢弃,item 永久卡 Processing。
Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
产品要求「只有 Fornax 管理员才能在发起实验时选择 priority,其他人一律走 default」。
此前**完全没有门控**:priority_level 在 IDL 里带 json/form 绑定标签(field 92,Submit 与
Create 两个入口都有),任何有 createLoopEvaluationExperiment 权限的调用方都能填 99 插队,
且不违反任何校验、不报错。中心调度按严格优先级排序,这意味着一个人把自己所有实验设成 99
就能让别人的实验饿死。
两道闸,都在 application 层:
expt_priority_white_list 谁能指定 priority(user_emails / space_ids / caller_psms,OR)
expt_trigger_trust_conf 谁能自称 evalx 从而进 enforce(按 caller PSM)
为什么不放 Authorization 层:商业版 allowlist_decorator 对特定 caller+method 直接
跳过整个 Authorization,门控放那里会被整条绕过。CreateExperiment 是四个入口
(EvalX / 控制台 / OpenAPI / 定时)的唯一汇聚点,Submit 也转成本请求后调进来。
★ 两个缺省方向刻意相反,因为失败代价不对称:
priority 配不到 → 拒绝。大家退回缺省优先级,可见且无损
trigger 配不到 → 放行。若一律拒绝会让全部 EvalX 实验静默退回 legacy,
中心调度突然没有任何候选,现象是"实验都在跑但一个都不受额度管控"
trigger 闸另有独立 enabled 开关,默认关闭:先在灰度确认 PSM 名单无误再打开,
不打断当前靠自报 evalx 的测试路径。上线前必须打开(已记入 spec 待办)。
★ 用 user_email 而非 user_id:名单靠人维护、靠人 review,zhangsan@bytedance.com
一眼知道是谁而 7123456789012345678 要另查一次 —— 加错人是"给了插队权",
最不该靠肉眼比对 19 位数字来防。邮箱取自已验证的 ByteTIM ticket claim
(商业版 CtxUser 中间件),不是请求体字段,调用方无法伪造,故可作授权键。
同理 trigger 判据用 kitex caller 而非请求体里的 trigger_type —— 后者谁都能自称 evalx。
★ space_ids 的正确用法是「只有管理员在的私有空间」,此时它是受控人员名单的代理。
类型注释里写明了红线:绝不填普通业务空间(谁都能建实验、成员随时增减)。
★ 新增 ConfigIDList:19 位雪花 ID 的配置读写两侧要求恰好相反 —— 写入侧
bytedcli tcc 把 JSON number 当 double 会静默截断,读取侧 encoding/json 解到
[]int64 时字符串形态直接报错、整份配置回落缺省值。两种失败都不报错给运维,
现象都是"配了却不生效"。故两种写法都接受(ParseInt 不经 float64)。
附带修复:加接口方法后三个包的**手写** fake configer(非 mockgen 生成)编译失败
—— infra/storage、infra/repo/target、infra/repo/evaluator,已补齐新方法。
Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
原先为兼容"数字与字符串两种写法"造了个 ConfigIDList 自定义类型(含 UnmarshalJSON /
MarshalJSON 与一组精度回归测试)。按约定简化:TCC 里的 Space ID 本来就该写字符串,
不需要兼容层。
保留的必要部分:
- 比对时把入参 int64 格式化成字符串再比
- spaceID=0 一律不匹配(防运维在名单里误填 "0" 就把所有无空间上下文的请求放行)
- 字段注释保留"为什么用字符串"(雪花 ID 超 float64 安全范围、bytedcli 写入会截断),
避免后来人"顺手改成 []int64 更规范"
端到端解析测试仍在,它现在证明的是 19 位 ID 从字符串配置解出后精确匹配。
净减 96 行。
Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
并发也是一种「单 item 占一份、终态归还一份」的资源,因此不新造计数器,而是登记进额度 体系当一个维度,直接复用现成的预占 / 释放 / 木桶 / 优先级排序。 为什么必须有这一维:资源向量(sandbox seat / model token)管的是"耗多少外部资源", 管不住"同时有多少 item 在跑"。若某类实验只申报 model token 不申报 sandbox,并发就完全 失控 —— 中心调度此前只受单实验 deficit + 资源额度两道约束,全 Scope 在跑的 item 总数 没有任何上限。有了这一维,TCC 里给 concurrency|item 配 global_quota=N 就等于 "全 Scope 最多 N 个 item 同时在跑",且高优先级实验先拿并发名额。 WithConcurrencyDimension 幂等:若已显式申报 concurrency|item(未来支持"重型 item 占 2 份并发")则保留申报值不覆盖。不原地改 receiver —— ExpectedQuotaConsumption 是创建期 冻结进 eval_conf 的快照,原地修改会让"冻结"语义失效。 注:本提交是接手另一会话的未提交工作(用户要求一起提交)。实现与 4 个测试均已核对, build 与测试全绿。commercial 侧的调用点随后一并提交。 Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
TCC 里的 central_expt_scheduler_space_config.default_priority 此前**完全无效且静默**:
GetDefaultPriority() 在生产代码里零个调用点,实际生效的是 entity.NormalizeExptPriorityLevel,
它把未申报的优先级**硬编码**收敛成 1。配成 5 看着生效,实际所有实验还是落 1,不报错不告警
—— 性质同雪花 ID 截断那类坑。
为什么现在必须修:priority 授权白名单落地后,非白名单调用方的申报值会被丢弃、强制走
default,这条路成了**主路径** —— 大部分实验的优先级都由这个值决定,而它锁死在 1。
改法(跨仓的 OSS 侧):
- ICentralAdmissionPolicy.AllowCentralScheduling 返回 bool → CentralAdmissionDecision
{Admitted, DefaultPriority}
- 新增 entity.NormalizeExptPriorityLevelWithDefault;旧入口保留为 WithDefault(p, 0)
的薄包装,读路径与 DB 转换那些既有调用点一个都不用改
★ 为什么回结构体而不加第二个 getter:两次独立读 TCC 之间配置可能热变更,
那样同一个实验会出现"按 A 配置准入、按 B 配置定优先级"。一次读取一起返回没有这个窗口。
★ DefaultPriority=0 统一表示"没有意见",回落到 DefaultExptPriorityLevel(=1)。
一个约定同时覆盖三种情况:TCC 没配该字段、noop policy(开源部署)、policy 未注入。
★ 越界的 default_priority(如 999)回落到 1,**不截断到 99**。它来自人工配置,
999 更可能是笔误而非"想要最高优";截断会让一次笔误静默变成"该空间所有实验都最高优"
—— 最难发现的一类事故(每个实验单看都正常,只有整体排序全乱)。
注意这与**申报值**越界的处理故意不同:调用方显式传 150 则截断到 99,那是明确意图。
★ 只在 enforce 分支采纳缺省值:legacy 实验不参与优先级排序,给它套非 1 的值
只会让 DB 里多一堆无意义数据、干扰排查。
顺带:go generate 补齐了 central_scheduler 四个此前从未生成的 mock(本仓 mocks 入 git)。
Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
产品口径:priority / 资源消耗向量 / trigger_type=evalx 都可以让调用方申报, 但**必须是白名单里的身份**。于是把上一轮单独做的 expt_trigger_trust_conf 删掉, 判据统一收进一份名单 expt_scheduling_privilege_white_list: priority_level 未授权 → 清空,走缺省优先级 expected_quota_consumption 未授权 → 清空 trigger_type = evalx 未授权 → 降级 manual,走 legacy ★ 为什么必须同一判据:这三样共同决定"这个实验拿多少资源、排在谁前面"。 只挡其中一两样等于没挡 —— 只挡 priority 却放开 trigger,任何人仍能自称 evalx 把实验塞进中心调度;只挡 trigger 却放开 quota,纳管范围内的实验仍能虚报消耗。 顺带这也消掉了上一轮"两份配置缺省方向刻意相反"的别扭设计:现在只有一份名单、一个方向。 ★ 判定条件覆盖"只申报向量"这一形态(有专门用例钉住):调用方不设 priority、 trigger 也不是 evalx、只塞一个资源向量时同样要过闸。向量在 legacy 下不生效, 但会被冻结进 eval_conf —— 该实验后续一旦被纳管就直接按虚报值扣额度。 ★ 只降级 evalx,不动 openapi/schedule/manual:只有 evalx 带来特权, 伪造其它 trigger 没有好处,而一并降级会把"调用来源"这个排查依据抹掉。 OpenAPI 侧新增三个字段(60/61/62),并把 SubmitExperimentOApi 的 trigger_type 从硬编码 openapi 改为可申报(未传仍回落 openapi,与改动前一致)。透传后由下游 唯一汇聚点 CreateExperiment 裁决 —— 门控不在每个入口各判一次。 ExpectedQuotaConsumption 结构此前做读视图时已在 domain_openapi 建好,直接复用, 本次只补反向转换 ExpectedQuotaConsumptionOpenAPI2Domain。 key 改名(priority → scheduling_privilege)是因为它现在管三样,叫 priority 会误导; 该 key 尚未在任何环境创建过,改名零成本。 附带:三个包的手写 fake configer 跟着接口改名同步(非 mockgen 生成,不会自动更新)。 Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
BOE 泳道实测到 panic: CheckBenefit → session.UserID → invalid memory address or nil pointer dereference 栈上是 HandleCentralAdmission → HandleEventLock → HandleCentralReservation → HandleEventExec,即**中心调度这条新派发链路**。legacy 不受影响。 根因在派发侧(中心调度构造 item 事件时漏填 Session,另提交修复),但这一层也不该 因为调用方漏填而 panic —— 该 panic 发生在 item 执行链里、被 HandleEventErr 的 recover 转成 error,现象是"每个派发出去的 item 都失败"而非进程崩溃,从现象极难反推到 "某个字段没填"。而且额度此时已预占,item 到不了终态,配额会被持续占住。 取"退化成匿名"而不是提前报错:权益校验拿到空 UserID 会返回明确的业务错误,那是可读的 失败;在此自造 error 会掩盖真实原因(调用方漏填),下次再有新链路漏填时同样难查。 同时打 Error 日志点出 nil session,让根因直接可见。 回归测试已反向验证:把 nil 保护去掉后,TestCheckBenefit_NilSessionDoesNotPanic 会复现日志里那个一模一样的 panic。 Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
releaseQuotaIfItemTerminal 是 item 正常跑完/正常失败时的**主释放点**,此前没有任何 测试。既有 expt_central_quota_release_test.go 覆盖的是 HandleEventErr 那层的兜底释放, 两者是不同的分支。 这一层同时是测试矩阵三格的共同判据: R2 终态即释放,不分成败 R3 MQ 重试不得释放(代码里刻意的反直觉设计) R8 不重复释放 为什么用单测而不是泳道 E2E 覆盖:本层全部分支都由"回查 run log 投影拿到什么状态" 决定,单测能精确摆出每种状态(含"查不到"与"仍在 Processing");而在泳道上构造 "失败但可重试"与"失败且已终态"的区别要靠让评测对象按特定方式报错,既慢又不稳定。 真机负责验证链路连通(已完成:22:05 实测 used 1808→708),分支穷举交给这里。 R3 那条尤其值得钉死:判据刻意是"回查投影的真实状态"而不是"execErr 是否为空" —— execErr != nil 既可能是可重试的瞬时错、也可能是已落终态的失败,只看 err 无法区分。 错误地在此释放会让重投消息在 ConfirmRunning 处被丢弃,item 永久停在 Processing, 比"额度多占一会"严重得多。用例特意带上 execErr + status=Processing 这个最危险的组合。 变异验证(三个守卫逐个确认承重): - 去掉 IsItemRunFinished 判定(无条件释放)→ KeepsReservationWhileRetriable FAIL - 改成只按 execErr 判(失败就不释放) → ReleasesOnEveryTerminalState FAIL - 去掉 len(obs)==0 守卫 → SkipsWhenProjectionMissing FAIL(并 panic: index out of range —— 该守卫同时在防一个真实的越界崩溃) Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
## 释放 reason 说谎(可观测性失真,非额度算错)
BOE 实测:129 条 `[CentralReservation] quota released on item terminal` **全部**写着
`reason: item success`,而其中 104 个 item 的 DB 状态是 `status=3`(Fail)——
失败被伪装成成功,靠日志根本发现不了实验在大面积失败。我自己就据此对外报过
"107/107 全部成功",实际 105 个在失败。
根因:`CompleteItemRun` 写完 `status=Fail + err_msg` 之后,只有 in-debt 错误
(`evalErrNeedTerminateExpt`)才 `return evalErr`,普通 item 失败一律 `return nil`;
`Eval` 又只透传其返回值。于是普通失败对释放点**完全不可见**,`execErr` 恒为 nil,
旧实现 `reason := "item success"; if execErr != nil {...}` 永远走第一行。
修法:抽出 `terminalReleaseReason(status, execErr)`,由**已经回查到的** run log 终态
推导(Success/Fail/Terminal 各自成文),execErr 非 nil 时附错误文本(in-debt 路径仍有
信息量)。未识别的终态带出原始状态码而非静默标 success —— 同样的错不该换个形式再犯。
**只修标签,不改行为**:终态判定仍是 `IsItemRunFinished`,释放时机与幂等性不变。
也没有让 `HandleEventErr` 的错误分支重新执行 —— 那需要改 `CompleteItemRun` 的返回契约,
影响面大得多(每个失败 item 都会进重试判定),单独评估。
⚠️ 由此推论一条容易误读的事实:既有观测「129 条全是 retry:false、无重投」**不能**读作
"重试机制健康",而是重试判定压根没运行过(nextErr 恒 nil 直接 return)。
变异验证:把 reason 改回 execErr 推导 → ReleasesOnEveryTerminalState +
FailedItemGetsFailedReason 两个用例 FAIL。
## 两处注释与实现不符(审计发现,本仓部分)
1. `expt_run_item_event_impl.go`「额度对账(spec §3.11)是最终防线」→ **对账不存在**。
目前只有调度器每拍在 Reserve 前跑的一小段(只清"投影 none + 账本 reserved"、
且只覆盖当前 LatestRunID)。本函数释放失败后那条 reservation **无任何兜底**,
Warn 日志是唯一线索。
2. `central_reservation_guard.go`「取得**一次性**执行权」→ `ConfirmRunning` 实际是
**幂等**的(对已 Running 返回 true、不做 CAS)。防重复执行靠的是 consumer 侧的
item 锁 `expt_item_eval_run_lock`,不是它。
这两处注释让代码读起来比实际更安全,我自己已被误导过(写悬挂预占对账时就建立在
"对账是最终防线"这个假设上)。
Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
审计发现的确定泄漏:`expt_manage_execution_impl.go` 里 `centralGuard` 出现 **0 次** ——
`Kill()` / `terminateItemTurns()` 都不释放中心额度。
为什么 consumer 侧那个释放点兜不住:它在 item 执行链的出口,靠"item 到终态"触发。
而实验被 Kill / Cancel 或落 Failed 之后,这些 item **不会再被执行** —— 消息可能已被丢弃、
可能压根还没投递、也可能 consumer 早已放弃,于是那个出口永远走不到。
更糟的是 running 态的悬挂连调度器每拍的对账都不碰
(`reapDanglingReservations` 显式 `if !view.IsReserved() { continue }`)——
**这条路径没有任何兜底**。而"取消实验"是常规操作,线上必然被走到:
用户取消一个正在跑 100 个 item 的 enforce 实验 = 100 份额度永久泄漏。
## 实现
新增 `releaseCentralQuotaForIncompleteItems`,挂在 `CompleteExpt` 的
`incompleteTurnIDs` 之后、终态 switch **之前**:
- 放在 switch 外只写一次,Terminated 与 Failed(default) 两条分支都覆盖,
且将来新增终态分支自动生效 —— "新增分支漏掉清理"正是本文件 default 分支注释里
记录过的历史教训(沙箱曾因此漏回收,一次漏两个)
- 按 **item 去重**:reservation 是 item 粒度而 incompleteTurnIDs 是 turn 粒度,
不去重会对同一条 reservation 发多次释放(Release 幂等,但会放大往返、且日志计数失真)
- `exptRunID` 为 nil 时回落 `LatestRunID`:账本 key 是 (run_id, item_id),run 号错了
释放就变成**静默 no-op**(比报错更糟)。`Kill` 的签名里 exptRunID 是 *int64
- legacy 实验直接返回:它们从不预占,而 CompleteExpt 是所有实验的公共收口,
白打 Redis 往返会按实验数放大
- enforce 却无 Scope 时**跳过而非猜**:猜错会归还别人的额度 → 超发,比不归还更糟
- best-effort:单个失败只告警、继续释放其余(一个失败就放弃剩下 99 个,
等于把"泄漏 1 份"放大成"泄漏 100 份")
构造函数新增 `centralGuard` 参数,wire 重新生成(未手改 wire_gen.go,diff 仅
6 行:把 guard 的构造提到 NewExptManager 之前)。开源部署注入既有的 noop 实现。
## 变异验证
- 去掉 legacy 判断 → SkipsLegacy FAIL
- 把去重的 key 从 ItemID 改成 TurnID → DedupesItemIDs FAIL
- **删掉 CompleteExpt 里的整行调用 → CompleteExpt_ReleasesCentralQuota FAIL**
★ 最后那条是补上的一道防线:前 7 个用例都直接调那个函数,所以删掉调用点它们**依然全绿**
(变异时实测过)。没有调用点测试,"函数写对了但没接上"这种失败模式毫无防线 ——
而这恰恰就是本次修复要补的缺口本身的形态。
Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
PPE 实测:5 个 item 已 Processing 14 小时,而 results 接口返回 run_state=queueing ×30。 现象是「实验看着没动」,实际跑得很正常——这个误导今天带偏了一整天的排查。 根因是平行实现漏写。legacy 的 handleToSubmits 一直成对写两张表 (UpdateItemRunLog + UpdateItemsResult,expt_run_scheduler_event_impl.go:673/678), 而中心调度这条新派发路径的 StartReservedItem 只写 run log 一张 (expt_item_dispatch.go:162)。run log 是执行真值,但用户看到的是主表 ——MGetExperimentResult 走 expt_item_result 构造 run_state。 修法:取得执行权后同步把主表推进到 Processing。失败只告警不阻断——主表是展示投影, 写不进去不影响执行与额度正确性,而返回错误会让已拿到执行权的 item 被 MQ 重投一遍。 配一条用例,断言主表被推进且 ufields 里 status 确为 Processing(只调方法不算修好)。 反向变异验证:关掉该修复 → 用例 FAIL;恢复 → 通过。全包测试绿。 与团队记忆 mem-20260820-new-dispatch-path-missing-enum-field 同族: 老路径对、新路径漏,且单测不显式断言就发现不了。 Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
## 为什么要传这个 打错 `category` 是唯一一条「打错字导致真实资源失控」的路径: - 打错 `resource_key`(`claude-opus-5` → `claude-opus5`)时**类级 wildcard 兜得住** —— 展开会同时记 `model|claude-opus5` 与 `model|*` 两笔,总量约束照样生效。 - 打错 `category`(`sandbox` → `sanbox`)时**连 wildcard 都兜不住** —— 记的是 `sanbox|*` 而不是 `sandbox|*`,两个 key 在账本里毫不相干。 该 item 在沙箱这一维上完全不受限,却真的会去占沙箱。 所以 category 需要在创建期就拒掉,而 resource_key 刻意不做名称校验 (真源是平台侧资源目录,迭代远快于发版节奏;做成白名单等于"平台上了新模型 要等评测侧发版",而卡住的表现是创建报参数错误、排查方向指不到"某个常量没加")。 ## 为什么校验放在 commercial 而不是这里 category 的登记表是内部资源目录(含内部模型/机型标识),不进开源仓。 本仓只把「申报了哪些 category」这个事实交出去,判定留给 policy 实现。 ## 改动 - `CentralAdmissionSubject` 加 `QuotaCategories`(去重、已 TrimSpace、顺序不保证) - `ExpectedQuotaConsumption.Categories()`:nil 安全、TrimSpace 后去重、跳过空值 去重是必需的 —— 同 category 下申报多个具体资源是正常形态(model|A + model|B), 不去重会让 policy 对同一 category 反复判定、错误信息里也重复列出 - `allowCentralScheduling` 多收一个向量参数并填进 subject policy 未注入 / noop 时行为不变(恒定放行、不看这个字段)。 ## 验证 `go build ./modules/evaluation/...` 零输出;`domain/entity` + `domain/component` + `domain/service` 全 PASS(service 包 225s 全量跑过)。 `Categories()` 三条用例:nil 安全、去重保序、trim 后同名去重。 Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
kill 一个 enforce 实验后额度一条不还,PPE 实测 9 个实验 45 条 reservation
卡了 11~19 小时。不是释放逻辑缺失(它存在且被调用),是**判据用错了维度**。
reservation 是 item 粒度(账本 key = res:<run_id>:<item_id>),而待释放集合此前来自
GetIncompleteTurns —— 它只收 turn_status ∈ {Queueing, Processing}。于是
**turn 已终态、item 仍 Processing** 的那些 item 拿不进列表,一条都不释放。
那不是边缘情形,而是沙箱执行进程死亡的典型形态:turn 被判完(超时/失败落终态),
item 的 run log 却没人去改。而这一格没有任何兜底 —— 此时 Redis 侧 state 已是 running:
reap 只处理 reserved(显式 `if !view.IsReserved() { continue }`);
对账的 isReleasableWithoutEvidence = {reserved, dispatched},刻意排除 running
(判据只能用预占时刻,分不出"卡死很久"与"刚被接管",按它释放会真超发);
zombie 清理只扫 Processing,但实验已终态、daemon 不再跳。
三条路都不接 ⇒ 永久泄漏。
改动:
1. 待释放集合改由 incompleteItemIDsForRelease 按 **item run log status** 自查
(Queueing ∪ Processing,即 !IsItemRunFinished),不再依赖调用方传 turn 列表。
扫描失败时**跳过释放而不猜 item** —— 释放不属于本实验的 item 会归还别人的额度,
那是超发,比泄漏严重。
2. 释放调用移到 `if !opt.NoCompleteItemTurn` **之外**。额度释放与"要不要改写 turn 状态"
是两件无关的事,没理由被同一开关控制;该选项当前零 caller,但一旦有人加 caller,
把释放留在块内就会连带静默失效。
穷举审计确认这一处同时是 4 个场景的根因:沙箱进程死亡、用户 kill、
实验落 Failed(尤其 run 级 36h 僵尸)、consumer SIGKILL 后实验已终态。
测试:既有 8 条全部改为经 repo 注入 run log(谁改回 turn 判据,mock 期望就不被满足),
并新增两条:
- ReleasesWhenTurnTerminalButItemNot —— 核心回归,turn 侧刻意返回空、只有 item 有数据
- ScanFailDoesNotGuessItems —— 扫描失败不得瞎释放
变异验证(三个都被抓到):
待释放集合恒为空 → ReleasesOnTerminal ×3 + ReleasesWhenTurnTerminalButItemNot
+ DedupesItemIDs 红
判据改回 turn 反推 → ReleasesOnTerminal 红
删掉 CompleteExpt 调用 → TestCompleteExpt_Releases 红(它专守"有没有被接上")
55 个包全绿。
Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
删除是一条**独立于 CompleteExpt 的泄漏路径**,此前两个删除入口(Delete / MDelete)
grep centralGuard 均为 0。三条独立原因让它比其它泄漏更彻底:
① 实验软删后 consumer 侧 GetByID 拿不到实验、直接退出 ⇒ 那些 item 永不执行,
"item 到终态才释放"的出口永远不会被走到;
② 删除**压根不经过** CompleteExpt,那里的释放不在这条路径上;而 CompleteExpt
对已删实验还有一条 early return,所以"先删再 kill"同样救不回来;
③ 软删后 ScanSchedulerQueue 带 deleted_at IS NULL ⇒ **连 full recovery 都扫不到**,
这些 reservation 会永久留在账本里 —— 其它泄漏至少还能靠运维 token 触发整本重建抹掉。
改动:
1. Delete 与 MDelete 都在**软删之前**调 releaseCentralQuotaForIncompleteItems。
顺序不是风格问题:删完再释放的话,中途任何失败都让额度彻底失去归还机会
(实验已不可查,SchedulerScope / LatestRunID 都拿不到)。反过来"先释放但删除失败"
只是让一个仍存在的实验少占额度,下一拍调度会重新预占,无损。
best-effort:释放失败只告警不阻断删除 —— 让额度问题挡住用户删实验是把后台问题
升级成前台故障。
2. CompleteExpt 对已删实验的 early return **刻意不补释放**(拿不到 scope/run 只能瞎猜,
猜错会归还别人的额度=超发),改为写明它不再是泄漏点的前提是"删除自己会释放",
并点出谁删掉那处会让这条路径重新变成永久泄漏。
测试(4 条,两个入口各自独立守):
- TestDelete_ReleasesCentralQuota 单实验入口
- TestMDelete_ReleasesCentralQuotaForEachExpt 批量入口,且断言 run 号各归各的
(账本 key 是 (run_id,item_id),串了就找不到)
- TestMDelete_ReleasesBeforeSoftDelete ★ 用 gomock .After() 钉死顺序
- TestMDelete_LegacyExptSkipsRelease legacy 连 run log 都不该扫(零期望 mock 保证)
变异验证(三个,且互不掩盖):
删 MDelete 的释放 → ReleasesCentralQuotaForEachExpt + ReleasesBeforeSoftDelete 红
删 Delete 的释放 → TestDelete_ReleasesCentralQuota 红(批量那条不会替它兜底)
释放挪到软删之后 → **只有** ReleasesBeforeSoftDelete 红
—— 坐实只测"释放了几条"的用例测不出顺序
55 个包全绿。
Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
## 为什么 同一个 `resource_key` 经不同来源拿到的可能是**不同的池子**:同一个模型走 LiteLLM 与走业务方自备通道,配额各自独立。混在一个账本条目里记账会让两边互相挤占。 字段对**所有 category** 一视同仁,不只模型 —— sandbox 也可以区分自建集群与平台池。 ## ★ 空 source 与该字段引入之前**完全等价** 这是它能安全加进已有账本的前提,也是本次改动的核心约束: - 存量 reservation 的 field 名不变 → 释放照样找得到 - 存量 TCC 上限的登记 key 不变 → 照样查得到 - 不申报 source 的调用方零改动 对应的 key 拼法(commercial 侧):空 source 不追加任何分隔符,产出与旧的**字节级 相同**的 key。若图省事写成无条件 `key + "@" + source`,后果是两条静默故障 —— 释放按新 key 找而账本存旧 key(HDEL 扑空 → 额度永久泄漏)、上限按新 key 查而 TCC 登记旧 key(查不到 → 按不受限放行)。 ## 改动 - `ExpectedResourceConsumption.Source`(`omitempty`:不申报时不写空串, 让"没有 source"在数据里可区分、也让存量快照字节形态不变) - `Validate()`:source 不得为 `*`(同 resource_key,通配只允许出现在上限配置里); **去重键带上 source** —— 不带会把"同一模型分来源申报"这个合法形态误判成重复键而拒绝创建, 那恰恰是本字段要支持的用法 - `Normalize()`:source 一并 TrimSpace(它进账本 key,带空白会与上限配置里的同名来源匹配不上) - 两份 thrift 加 `4: optional string source` ## 验证 `go build ./modules/evaluation/...` 零输出;`domain/entity` 全 PASS。 新增 6 条用例:同资源分来源合法、有无 source 算不同键、同来源重复仍判重、 通配拒绝、空 source 合法、Normalize 去空白。 **变异验证**:把去重键改回不带 source → `TestValidate_SourceRules` 三个子用例 FAIL。 convertor 与 `kitex_gen` 生成物在下一个 commit(本 commit 只含手写且可独立编译的部分)。 Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
承接上一个 commit(IDL 与 entity 已改,convertor 因缺生成物无法编译)。
## 生成物范围:恰好 4 个文件、273 行**纯新增**、零删除
跑 `backend/script/cloudwego/kitex_tool.sh`(`NO_PUSH_REMOTE=true`)。该脚本会
`rm -rf kitex_gen` 后全量重生成,所以先确认了生成物基线干净、再逐项核对 diff:
kitex_gen/coze/loop/evaluation/domain/expt/expt.go +81
kitex_gen/coze/loop/evaluation/domain/expt/k-expt.go +56
kitex_gen/coze/loop/evaluation/domain_openapi/experiment/*.go +136
全部是 field 4 的读写脚手架(ReadField4 / writeField4 / GetSource / IsSetSource /
fastRead / fastWrite)。**没有一行删除、没有触及其它服务** —— 工具链版本与仓库
约定一致(thriftgo v0.4.1 / kitex v0.13.1 / hz v0.9.7 / validator v0.2.6),
不存在版本差异导致的重排。
## 转换层:四个搬运点都补上 source
`expectedQuotaConsumptionDTO2DO` / `DO2DTO`、`ExpectedQuotaConsumptionDomain2OpenAPI` /
`OpenAPI2Domain`。四处都是逐字段手写赋值,漏一个的后果极隐蔽 —— 申报方传了来源、
落库被吞掉,于是不同来源的用量记进同一账本条目互相挤占,而接口回显看起来
"没申报过来源"、与调用方真的没传完全一样。
## 验证
`go build ./...`(**全仓**)零输出;`domain/entity` + `convertor/experiment` 全 PASS。
新增端到端往返用例,覆盖内部 DTO↔DO 与 OpenAPI↔domain 两条链,且刻意包含
「同资源有来源 + 同资源无来源」并存的形态(那正是本字段要支持的用法)。
**变异验证**:把 DTO2DO 那处的 `Source` 赋值删掉 → 用例 FAIL。
用例里用 sandbox 而非 model 验 OpenAPI 那条链,顺带钉住"本字段不是模型专用"。
Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
PPE 长跑压测实测:enforce 实验的 processing_turn_count **走负数**
(一个 14 题实验 fail 累到 4 时 processing = -4),pending_turn_count 恒不下降。
legacy 对照实验计数正常,所以只影响 enforce。
## 根因
完成侧(expt_result_impl.go:175-177)做的是「从 item 原状态减 1、往新状态加 1」,
减的是 Processing 桶。两条链路谁往那个桶里加过:
legacy handleToSubmits 派发时 ArithOperateCount{Processing:+n, Queueing:-n} ← 加了
enforce 派发只 ClaimQuotaReserved(Queueing/none → Queueing/**reserved**) ← 没加
enforce 的 status 全程是 Queueing,真正翻成 Processing 的是 consumer 侧的
StartReservedItem,而那里只写了 run log 与主表,没记 stats。
于是完成时**减一个从未被加过的计数** —— 单向下溢,而不是"少加了一次"。
## 修在哪
补在 StartReservedItem 成功之后,与已有的主表推进并列 —— 那里才是
「item 真正进入 Processing」的时刻。派发侧不是修复点:那时 status 还是 Queueing,
在派发侧加会让"已预占未消费"的 item 被计成 Processing,与投影口径不一致。
★ 必须绑定 started(CAS 真的翻了状态)。started=false 是重复投递或已被 repair 修正,
此时 item 早已计入 Processing,再加一次就从"少计"变成"多计" ——
CAS 结果是这条路径上唯一的"恰好一次"信号。
失败只告警不阻断,与相邻的主表推进同策:stats 是展示投影,
返回错误会让已取得执行权的 item 被 MQ 重投一遍。
## 测试
两个用例互为反面,缺一个就有漏网的改法:
- AdvancesStatsToProcessing —— 断言 op 的**内容**(Processing:+1 且 Queueing:-1)。
只断言"调过"的话,只加不减也能通过,而那样 pending 依旧不降。
- SkipsStatsOnDuplicateDelivery —— started=false 时 Times(0)。
没有它,把记账写在 started 判断之外也能全绿。
反向变异 3 个全部被检出:
① 去掉 `started &&` → 重复投递用例 FAIL
② 去掉 `Queueing: -1` → 内容断言 FAIL(expected -1, actual 0)
③ 方向写反 → 内容断言 FAIL(两项都反)
`go test ./modules/evaluation/domain/service/` 全绿(226s)。
## 备注
与团队记忆 mem-20260820-new-dispatch-path-missing-enum-field 同族:
新增派发路径时漏掉旧路径的副作用。那次漏 ExptRunMode 字段,这次漏 stats 记账,
同一文件里还有人修过同族的"漏写主表"。**平行实现漏副作用**已出现三次,
共性是老路径把多件事做在一处,新路径只搬了其中一件。
Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
handleZombies / sweepTerminatedSandboxItems 判定 item 终态后,除了写 run log 还抢先把主表 expt_item_result.status 改成 Fail。而随后 RecordItemRunLogs 算 expt_stats 增量用的是「主表旧值 -1 / run log 新值 +1」的差分 —— 两边都读到 Fail 就算成净零,该 item 从计数上凭空蒸发(各桶总和 < 总行数)。 去掉这两处的 status 键,err_msg 保留原地写入(#602 要的用户可见超时原因不受影响)。 主表 status 由 RecordItemRunLogs 统一落,run log 已是 Fail,最终态不变, 只是晚一个 tick 内的间隔(约 2.5s)可见 —— 这个代价已确认可接受。 回归来源:a2f11f2a(#602) 为带 err_msg 新增主表写入时顺手带上了 status; fc6c6ff(#606) 照抄了这个形状,把同一个 bug 复制到 sandbox sweep 路径。 顺带修一处失真注释:原注释称顺序是为「避免额度已放但 item 仍显示 Processing 被下一拍读到而重复授予」,但调度侧 LoadDispatchRuntime 读的是 run log 且候选 必须 status=Queueing,Processing 的 item 永远进不了候选 —— 该机制不可达, 方向也相反(少派而非超发)。改为真实理由:先落 run log 的 Fail 再放额度, 才不会留下「额度已归还、run log 仍算占用」的窗口。 测试侧两处 UpdateItemsResult 的 gomock.Any() 换成捕获实参、断言 status 键不存在 —— 原来用 Any() 意味着删/加这个键测试都不会红,正是 #602 能悄悄溜进来的原因。 两个反向变异(分别把 status 加回两条路径)均被检出。 Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
中心调度里 item 的最后一道兜底是「实验级 zombie 超时 → 实验判失败 → 归还未跑完 item 的额度」,判据 now - event.CreatedAt >= ZombieIntervalSecond(默认 36h)。 而在线实验的 daemon(ExptAppendExec.NextTick)每转一圈就把 event.CreatedAt 刷成当前 时间 —— 那是在线实验该有的行为(长期存活、不断追加 item),但副作用是这个绝对时钟 永不到期。离线四种 mode 都不刷,所以只有在线这一支没有兜底。 后果不止额度多占:投递结果不明的 item 投影停在 Queueing/reserved,而那个形态被调度 候选(只取 Queueing/none)、item zombie(只扫 Processing)、reap(只处理 reserved) 同时漏掉。对账器的退回队列是它唯一出路,若连实验级兜底也没有,任何未被对账覆盖的 窗口都会变成永久泄漏 + item 丢失。 用 ExptType 而非 ExptRunMode:run mode 是 Run 时才定的,而这道闸必须在创建期生效 (dispatch_mode 一次性冻结进 DB 列,之后 Run/Retry 一律回查该列)。ExptType 在 CreateExptParam 里就有,是创建期唯一可用且权威的判据。 本期在线实验不走 evalx trigger,所以这道闸当前不改变任何行为 —— 它挡的是"哪天在线 实验开始发 evalx"那个未来。在此之前那条约定只存在于口头,代码里没有任何拦点。 刻意不动另两个 ShouldEnforceByTrigger 调用点(特权字段读取、飞书通知抑制): 它们回答的是"是不是 evalx 来的",与"要不要进 enforce"无关。 零值 ExptType 不视为 Online:否则没填 expt_type 的调用方会静默失去中心调度。 变异验证:去掉在线判定 → 3 格用例 FAIL。
nil 分支:两个生产调用点(商业版 toRequirements / frozenConstraintsOf)都在它之前挡了 nil,所以生产不可达。但**不能删** —— 实测 nil receiver 读 c.Resources 直接 panic, 这个分支是本方法唯一的 nil 安全保障,删掉等于把"nil 安全"变成"nil panic"。 补注释说明为何不可达、为何仍要留。 幂等分支:原注释写"例如未来支持重型 item 占 2 份并发",把它说成未来能力。实际 concurrency 已在商业版 category 白名单内、Validate 也不拦,**现在申报即生效**。 TestWithConcurrencyDimension_NilReceiver 的注释断言「若返回 nil,调度器会把该实验当 无向量跳过,enforce 实验永远不跑」—— 后半句不成立:调用点在此之前就已经跳过了 (scheduler.go 的 len(requirements)==0 → skipReasonNoRequirements),这个函数返不返回 非 nil 都改变不了结果。改成如实描述它守的是方法自身的 nil 安全契约,并记上变异实测: 把该分支改成 return nil,全 evaluation 模块只有本用例 FAIL、商业版 51 包零失败。
两条都是 2026-08-26 PPE 长跑实测出来的,共同点是**没有任何兜底会兜住**。
① 实验终态时 item run log 漏写终态
kill 只写了主表 expt_item_result(terminateItemTurns),run log 留在 Processing。
而调度侧三条判据全读 run log:释放判据 IsItemRunFinished 永远 false、对账见
Processing 判「有执行证据、绝不可释放」、zombie 只扫 Processing 但实验已终态
daemon 不再跳 —— 三条路都不接 ⇒ 飞行中的 item 预占永久泄漏。
实测 6 条 48h+ 僵尸吃掉 sandbox|default 8 里的 6,把 prio 99 实验有效并发压到 2,
造成结果层优先级倒挂。不是对账有 bug,是我们喂了它一个假事实。
修法:CompleteExpt 释放额度之后,把 run log 里 Queueing/Processing 的行一并置 Terminal。
置终态后对账会按 terminal_leftover 正常归还,于是它也成了那次 best-effort 释放
万一失败时唯一的兜底。顺序不可调换(先置终态会让释放侧一条都查不到)。
只对中心调度实验生效:run log 的 status 不是用户可见字段,legacy 不走那三条链路。
② 重试不释放被顶替的旧 run 预占
FailRetry / RetryAll / RetryItems 三处 reset 都是 UpdateItemsResult{expt_run_id: 新值}
→ BatchCreateNXRunLogs,中间零 Release。改完 DB 里不再存在携带旧 runID 的主表行,
而增量对账只按 LatestRunID 建投影 ⇒ obs == nil ⇒ ActionReportOnly,只有全量 recovery
才可能碰到 —— 增量对账永远够不着。
修法:落在 LogRun / LogRetryItemsRun,即新 run 号写进 LatestRunID **之前**(账本 key 是
(run_id,item_id),旧 run 号一丢 field 就再也拼不出来)。retried 分支刻意不释放:
那说明锁被活着的 run 持有,释放等于把正在跑的额度还回去(真超发)。
测试:8 个变异逐一验证能被抓住(删调用行 / 挪顺序 / 写成 Fail / 用新 run 号释放 /
去掉模式闸 / 去掉 legacy 闸 / 去掉活跑闸 / 去掉同号短路),不靠「测试绿就算过」。
Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
供中心调度每拍授予点使用:某维度 AmountPerItem > Limit 即结构性不可能(等多久 都白等),调度器据此把该实验待跑 item 置 Fail 并写入此 err_msg,而非继续静默排队 ——此前它与"暂时排队"在日志和页面上完全无法区分(实测第三方实验静默 38h+)。 结果层同步识别并暴露到 ItemSystemInfo.Error,用户在前端能看到具体维度与数值。 NoAffectStability=true:这是用户申报配置错误,非系统稳定性问题。 Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
ConfirmRunning 判定 reservation 不存在时,原实现只打一条 Info 就 return,不动 run log 投影。而此时 consumer 往往已用 StartReservedItem 把投影兑现成 Processing/none,于是 item 成了「既无额度、也无执行者」的孤儿 —— 却仍以 Processing 计入 item_concur_num: ScanEvalItems 与 LoadDispatchRuntime 都把它当"在跑",既不派新 item 也等不到它完成, 只能靠异步僵尸阈值(默认 3h)兜底判 Fail。 2026-08-28 线上实测:两个实验各 20 个 item 占满 20 的并发槽位 77 分钟, turn 表零记录、账本零 reservation,实验整体停摆;日志里连一条 Warn 都没有。 对账(reconcile)覆盖不到这一类:它遍历的是账本里的 reservation(groupForObservation 的入参就是 reservation map),而这类漂移的特征恰恰是账本侧无记录 —— 压根进不了输入集。 所以只能在 consumer 侧收口。 改动: - 该分支改为先读投影再分三路处理,日志从 Info 升为 Warn 并带上判定结果 · 终态 —— 迟到消息,一个字段都不动(退回会让已跑完的 item 重跑) · Processing —— 孤儿态,退回 Queueing 让它重新被授予,连带回滚 stats 与主表投影 · Queueing —— reserved 的走 ResetQuotaReserved 清回 none;none 的本就在候选里 - 新增 IExptItemDispatchRepo.RequeueProcessingItem:Processing/none → Queueing/none 的 CAS, 条件同时钉住 status 与 quota_reservation_state,只作用于孤儿这一种形状 - 全程只告警不返回 error:这条消息注定不执行,返回 error 只会让 MQ 无休止重投 选择退回 Queueing 而不是落 Fail:这类 item 是在创建 turn 之前被丢弃的,从未真正执行, 判失败会凭空吃掉一道题;退回队列它能被重新授予并真正跑完。 Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
44f52a9 的两处漏洞,均由独立 review 指出: 1. 主表 UpdateItemsResult 没绑 requeued,而 stats 绑了。 主表 status 不是纯展示字段,是 stats 的锚点 —— 完成侧 statsCntOp 读 items_result.Status 做「-1」。CAS 未命中时(并发 handleZombies / sweepTerminatedSandboxItems 跑在实验锁下、不持 item 锁,可在读投影与 CAS 之间 抢先落终态)stats 仍挂着 Processing +1 而主表已被改成 Queueing, 完成侧就去减 Queueing 桶 ⇒ processing_turn_count 永远归不了零。 同仓 expt_run_scheduler_event_impl.go 的 zombie 路径已就此写过警示, c4a6a95 也已就此修过一次。 2. RequeueProcessingItem 用 Update 会刷新 updated_at,而僵尸判定正是 `Processing 且 time.Since(updated_at) > zombieSecond`。若「reservation 消失」 是持续性成因(账本损坏、reap 竞态反复触发),item 会在 Queueing ↔ Processing 之间无限往返、每次把僵尸时钟拨回零,实验永不收敛 —— 等于把「3 小时后必定 收敛」的有界故障换成无界故障。改用 UpdateColumn 保留原始 updated_at: 一次性成因下 item 正常重跑,持续性成因下仍在原定 3 小时被判 Fail, 最坏情况不劣于修复前。 补一条 CAS 未命中的用例 —— 原有四个用例挡不住第 1 条。 Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
此前这条错误码零测试。补两层,分别钉住两类不同的失效: errno 层 —— 序列化 round-trip 与「不与僵尸超时串味」。round-trip 是这条码的真实 契约(调度侧写进 err_msg 落库、结果层反解展示),任一侧编解码不一致,用户看到的 就是"失败但无原因"。互不命中同样关键:两者共用 err_msg 字段,标错会给出相反的 处置建议(僵尸该等/重跑,额度不可满足该改配置)。 结果层 —— 错误码有测试 ≠ 有人真去读它。分支漏接的表现不是报错,而是 item 显示 失败却没有原因;另配一条反向用例,确认僵尸超时不会被新分支抢走。 四条反向变异逐条实跑确认检出:文案漏掉 amount/limit、Parse 不校验 code、 结果层漏接分支、僵尸分支错标成额度码。modules/evaluation 55 包全绿。 Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
staticcheck 报 `could remove embedded field "Dialector" from selector`: gorm.DB 内嵌了 Dialector,Explain 直接从 gormDB 上就能调。行为完全等价。 Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
codecov 的 patch 门槛是 80%(threshold 0%、informational:false,会 gate),本分支 实测只有 79.18% —— 差 0.82 个点。补的两个文件此前都是 0 覆盖。 派发投影 repo 是纯透传层,所以用例只钉两件事:入参**按位置**原样转发、返回值与 错误原样上抛。每个方法用精确实参匹配而不是 gomock.Any()——透传层唯一真实的失效 模式就是把 spaceID/exptID/exptRunID 这几个同类型 int64 传错位置,用 Any() 就恰好 放过了唯一要防的那类 bug。CAS 未命中必须回 (false, nil) 也单独钉住:改写成 true 会让重复投递的 item 执行两次。 ScanSchedulerQueue 三条分支各对应一个真实后果:入参原样下推(Scope 丢了会扫出别的 Scope 的实验)、单条 eval_conf 损坏跳过而非整批失败(否则一条脏数据永久阻塞该 Scope 下所有实验的调度)、DAO 出错必须上抛(静默回空会被调度器读成"队列里没候选", 表现是集体不动)。 顺带生成缺失的 IExptItemDispatchDAO mock(按文件里既有的 go:generate 指令)。 四条反向变异逐条实跑确认检出:坏 payload 改为整批失败、DAO 出错静默回空、 exptID/exptRunID 传反、CAS 未命中改写成 true。 55 包全绿;CI 同款 golangci-lint v2.2.1 三档(新增/模块全量/全 backend)均 0 issues。 Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
…ota_consumption 给已发起实验开出运行中调整调度优先级与单 item 预期资源消耗的入口,普通面与 OpenAPI 面同时开放。字段号与各自的创建接口对齐(92/93、60/61),便于两处对照。 只改 IDL 与生成物;授权闸门与写入逻辑随后提交。
两面(普通 / OpenAPI)共用一份白名单闸与校验:命中才生效,未命中丢弃 + WARN, 与创建期 enforceSchedulingPrivilege 同口径。判据抽成包级函数,两处各写一遍 迟早漂移成"某一面能绕过白名单"。 几处刻意的取舍: - 校验放在闸门之后。与创建期一致,也避免用"报不报错"把「你不在白名单里」 泄漏给未授权调用方。 - priority 不接受 0 当"不修改",nil 才是。0 落到下游 Normalize 会被收敛成缺省 优先级,于是"我传了 0"和"我想设成缺省"无法区分,静默改掉一个本不该动的值。 - 只允许改 enforce 实验。legacy 既不参与优先级排序也没有额度账本,照写只会得到 一个没人读的值,而调用方收到成功。 - priority 走 UpdateFields 的显式列名 map,不走 exptRepo.Update。后者是 struct Updates + Omit(schedulingFrozenColumns),正是为防部分更新把 enforce 打回 legacy 而存在的;显式列名只可能改到写进 map 的那一列。 - 向量改动后必须给该 run 在飞的预占重新计价,失败上抛。全量恢复刻意只读 MySQL 重建 used,隐含"每条活预占按当前向量计价"这条前提;只改库会让释放按旧价、 恢复按新价,雷埋在修账本的唯一手段里。 新增 ICentralReservationGuard.RepriceRunConsumption 承载改价,开源部署 noop。
main 的 item run count 给 NewPayloadBuilder 插了两个参数
(第 8 位 exptItemResultRepo、第 10 位 evalTargetRepo),
而本分支新增的 TestNewPayloadBuilder_ItemQuotaImpossibleErrParsing 是按旧签名写的。
两处改动在不同位置,git 文本合并成功、不报冲突,但测试包编译不过:
expt_result_impl_test.go:8444: not enough arguments in call to NewPayloadBuilder
本用例只验"额度不可满足 / 僵尸超时"两种 err_msg 的反解分支,构造期不触碰这两个依赖,
故传 nil 并标注参数名,避免下次再有人数不清位置。
验证:go build ./... 通过;modules/evaluation 55 个包全绿。
调度拍是「读一页候选(连冻结向量)→ 逐个 Reserve」,改价可能插在某个候选的 已读向量与未 Reserve 之间,那一拍新建的 reservation 会按旧价计价。 写清三件事,因为缺任何一件这个窗口就不可接受:影响面被并发上限封住且下一拍 不再发生;释放按每条自存 amount 归还所以不漏不负;随 item 终态自然消退。 并写明残留影响只在个位数 limit 的维度上不可忽略,以及根治方向是向量修订号 而不是在用户面接口里抢调度租约。
跨空间共享评测对象时,传给 operator 的 spaceID 已被 resolveLoadSpaceID 换成 **评测对象来源空间** —— 那是「去哪个空间读这个对象」的口径。而「按哪个空间取模型凭据 (api_key / base_url)」应当跟着发起实验的空间走:evaluator 侧早已如此 (evaluatorSpaceID = expt.SpaceID),callTarget 的埋点也一直用 Event.SpaceID, 只有 target 的模型凭据还在按来源空间取。 两个口径必须并存,所以另开一个字段,**不改那个 spaceID 形参** —— 它同时决定沙箱 workspace、eval_target_record 的落库/读取空间、空间 AK/SK 与 TCC 空间级配置, 换掉会撕裂销毁链(四条回收链全读 record.SpaceID)与回传链 (ReportInvokeRecords 按 (space, record_id) 定位)。 - ExecuteTargetCtx / ExecuteEvalTargetParam 各加 ExptSpaceID,纯加字段 - expt_run_item_turn_impl.go 从 etec.Event.SpaceID 填入;target_impl.go 透传 - 调试链路 (DebugTarget / AsyncDebugTarget) 不填 → 0,消费侧据此回落原行为 ★ 恒取 Event.SpaceID,不从 ItemConfig 派生:多集执行恒用顶层 target,而 ItemConfig 的来源空间是 per-set 的,按 per-set 派生会重演「拿 B 的空间去加载 A 的 target」那类错配。 测试:新增跨空间守卫用例,专门造「Event.SpaceID=42 / spaceID 形参=99」的 fixture。 同空间 fixture 上写这条断言是没有牙齿的(两值相等,分辨不出实现取了哪个)—— 变异验证过:把实现改成取形参,该用例因断言失败而红。 验证:go build ./... 通过;modules/evaluation 55 个包全绿。
线上实测:PPE 一个 900 题 enforce 实验,expt_stats 的 pending 虚高 281、 success 少 249,而 item 表与 turn 表分毫不差;同泳道的 legacy 实验完全一致。 缺口还在持续增长(一小时内 21 个完成只被记了 17 个)。 两处修改: ① 派发侧记账换判据。原先绑在 run log 的 CAS(started)上,而主表推进是无条件的 —— 两个判据一分叉,计数行就与主表对不上。而完成侧 statsCntOp 恰恰是按 items_result.Status 做「-1」的,于是它去减一个计数行里从没加过的桶。 现在以「主表真的发生了状态迁移」为唯一判据,两边同源;减的桶取主表实际所在的 状态而非写死 Queueing(repair/重投会让它停在别的状态)。 ② 调度器每拍对账。终态路径有多条(正常完成 / zombie / 沙箱 sweep),每条都得自己 记得记一笔增量账,漏一笔就永久偏一笔,而那些调用点全是 warn-only 或干脆没有 —— 没有任何机制会发现。运行期此前也没有重算:CalculateStats 只在 CompleteExpt 与在线实验 daemon 里调。现在每拍用一条 GROUP BY 拿主表真实分布,不一致才写回, 并打 Warn —— 那条日志是「还在漏」的唯一信号。 以主表而非 turn 表为准:完成侧的「-1」读的就是主表,要能修正它必须同源; 拿 turn 表对账会在多轮实验上得出另一套数字(所以没复用 CalculateStats)。 派发侧仍保持失败只告警:返回错误会让已取得执行权的 item 被 MQ 重投, 重投一次 Agent 执行的代价远大于一次计数偏差 —— 残留偏差由②收敛。
上一版是「每实验每拍一条 GROUP BY」,一个空间几百个在跑的实验就是每分钟几百条 count。 两道闸把成本降下来,判据都不是拍脑袋: - 只对 enforce 实验。legacy 的派发把 run log、主表、turn 表、stats 四张无条件一起写, 失败 return err 重投整批 —— 计数行天然镜像主表。2026-09-01 PPE 实测两个正在跑的 legacy 实验(30 题 / 3 题)分毫不差,而同泳道 enforce 差 249/281。给 legacy 对账是纯开销。 - 同一实验 5 分钟内只对一次(走 idem SetNX)。偏差只在有人正看时才有意义,且实验进终态 时 CompleteExpt 一定会全量重算兜底 —— 对账只需要"最终收敛",不需要"实时精确"。 SetNX 出错时跳过而不是继续:Redis 抖动本就该少做事,"出错就对账"会让 Redis 一挂 就退化成每拍全量 count,正好是最不该加压的时候。
main 引入 ProvideExptItemRefRepos 后 wire 的匿名 slice 变量编号整体后移, 逐 commit 解冲突无法保证生成物一致,改由 wire 重新生成; 并补上解冲突时丢掉的函数间空行(gofumpt)。
xueyizheng
force-pushed
the
feat/evalx-evaluation
branch
from
September 1, 2026 11:35
c8e594b to
fcb523b
Compare
legacy 的 handleToSubmits 派发时五写(run log / 主表 / CK 加速表 / turn 主表 / stats), 中心调度这条平行路径只写了前两项与 stats。 漏 CK 的表现:item_run_state 筛选打的是 etrf.status,开启加速器时结果里的 run_state 也从 CK 读,于是 enforce 实验执行期间按「运行中」筛恒为空、按「排队中」筛反而捞出 正在跑的 item 并显示成排队中;完成时的 upsert 会纠回来,只有 in-flight 窗口失真。 判据沿用主表是否真的发生迁移,两写彼此独立且只告警不阻断。
派发侧推进五项(run log / 主表 / stats / turn 主表 / CK),此前回退只做前三项, 留下「item 已回队列、turn 与 CK 仍显示执行中」的错位。它不自愈:实验若以 Failed 收口,CompleteExpt 的 default 分支只销毁沙箱、不动 turn 状态。 判据与 stats、主表一致(绑定 requeued),失败只告警。
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
Summary
Today each experiment drives its own item dispatch loop, filling up to its configured concurrency. That works in isolation but has no notion of global ordering or of a shared resource budget: when many experiments run in one deployment, whoever ticks first wins, and nothing prevents the sum of all experiments from exceeding the resources actually available.
This PR adds the seams and bookkeeping needed for a centralized scheduler — one that decides across experiments, by priority and by declared resource consumption — while keeping the existing per-experiment path byte-for-byte unchanged for experiments that do not opt in.
The scheduler implementation itself is pluggable and is not part of this PR; what lands here is the open-source side: schema, entities, narrow ports, projections, quota accounting, and the reclaim paths that keep the ledger honest.
What's in it
Opt-in, frozen at creation. An experiment carries a dispatch mode that is decided once at creation and never re-derived. Retries and re-runs read the frozen column, so changing configuration after the fact cannot silently move a running experiment between modes. Only experiments whose trigger source opts in are enrolled; everything else behaves exactly as before.
Declared consumption, expanded into constraints. An experiment declares what a single item costs (sandbox:default=1, model:some-model=1000, …). Each declaration expands into layered ledger keys — the concrete resource, its cross-source total, and the category-level total — so a per-resource cap and a per-category cap can both be enforced, and adding a new source cannot be used to slip past a total.
Two drivers cannot both dispatch. For enrolled experiments the legacy per-experiment tick keeps every other duty (archiving finished items, timeout handling, terminal wrap-up, next tick) but drops item submission. Letting both dispatch would overrun the budget immediately, since the legacy path neither consults global priority nor touches the ledger.
A reservation projection on the run log. Item run logs carry a reservation state so "reserved but not yet consumed" is distinguishable from "queued" and from "running". The consumer redeems the reservation when it takes execution rights, which is what makes at-most-once dispatch verifiable rather than hopeful.
Reclaim paths — the bulk of the fixes here. Quota that is taken must come back on every exit, and most of this PR's bug fixes are exactly that:
terminal wrap-up releases quota for items that never finished, judged by item state rather
than turn state (the two diverge when an external execution process dies);
deleting an experiment releases its quota (previously leaked, irreversibly);
retry releases the superseded run's reservations before the run id is rewritten — afterwards
the ledger keys can no longer be reconstructed;
terminal wrap-up also drives the run log to a terminal state, not just the main table: the
release predicate, the reconciliation view and the timeout sweep all read the run log, so
leaving it behind hands all three a false fact;
items whose reservation has gone missing are returned to the queue instead of being dropped,
which previously left them stuck in Processing forever.
A distinguishable failure for impossible declarations. If a single item declares more of a dimension than the deployment's registered ceiling, no amount of waiting will ever admit it. Previously this was indistinguishable from ordinary contention — the same "nothing granted" line every tick, with the experiment simply never progressing. There is now a dedicated error code carrying the dimension, the declared amount and the ceiling, surfaced on the item so the user can see which dimension to change and to what. A registered ceiling of zero is deliberately not treated this way: that is an operator throttle, and experiments are expected to queue until it is lifted.
Statistics stay consistent. Two daemon terminal paths were racing the main-table status write and a counter could drift; enrolled experiments now record statistics when execution rights are taken, and items failed for an impossible declaration adjust the queued/failed buckets directly (they never entered execution, so no consumer will do it for them).
Compatibility
Default behaviour is unchanged: without opt-in, dispatch, statistics and error surfaces are
identical to before.
The quota gate ships with a no-op implementation registered in the OSS wiring, so a build with
no scheduler configured behaves exactly as it does today.
Schema changes are additive columns plus one narrowed index; no column is dropped or retyped.