- D1: Registry 3 prompt(策划/对抗/fix)+eval 骨架+registry 注册 - D2: agent-loop 2 schema + 4 模板 schema(对齐 gen_spike 真源, additionalProperties:false) - D3: 编排器+裁决引擎(judge D1-D10 决策表/账本幂等重放/预算闸熔断/eval 回流/批报告), 59 单测绿(mini-desktop) - D5: 玩家 CDP 取证 player_cdp.py(真点击/蛇形→驼峰映射/demo 三信号/三角合成判定), 9001 自测 PASS - (D4 后端件已在 9bf3d54 先行入库, 三条 mvn 两轮独立实跑全绿) - spec: 评审版(已拍板)+execution 版(已审定+建设期补充裁决 D4-a~d/D5-a~c/D2-a/D3-a) - W2 试点/种子 9003-9008/B4 四链路烟测 产物与总账/作战清单同步 - eval 种子: C1 spike 52 条转入 config.clicker-designer(inputs 52/labels 13) Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
376 lines
19 KiB
Python
376 lines
19 KiB
Python
#!/usr/bin/env python3
|
||
# -*- coding: utf-8 -*-
|
||
"""
|
||
JSONL 账本 + designId 幂等重放 + 预算闸 + 熔断(D3 件)—— spec:HJ-AGENT-LOOP-EXEC-001 §7.2/§7.3/§7.4
|
||
|
||
账本契约(§7.2,append-only,路径 runs/<batchId>/ledger.jsonl):
|
||
- 阶段事件行:{type:"stage", ts, batchId, designId, round, stage, status, taskId?, versionId?, traceId?, data, llmCalls}
|
||
- 终判行:{type:"verdict", ...Verdict 全字段(§6.2), ts, batchId}
|
||
|
||
幂等重放键 = designId(§7.2 + §16 D2-a 裁决口径):
|
||
- designId = sha256(idea+templateId+round) hex64(评审版 §3.2 冻结公式,**哈希因子含 round**):
|
||
round=0 首轮与 round=1 回炉轮的 designId 必不同;ledger/重放仍以**各轮 designId** 为幂等键;
|
||
fix 前后对照以 rootDesignId(=round0 的 designId)为父子关联键(§7.7 labels.jsonl 落该字段);
|
||
- resume 时按账本扫描每个(轮级)designId 已完成的最后阶段,跳过已 ok 阶段从断点续跑;
|
||
- 任一轮 designId 已存在**终判** verdict(decision ∈ accept/kill)→ 该创意整体跳过;
|
||
注:decision=fix 的 round=0 中间 Verdict(D5 行语义=「全流程重走 E1」)不构成终判,
|
||
resume 时该创意以 round=1 的 designId 从回炉轮续跑——否则 fix 中断的创意会被错误跳过;
|
||
- 含 infra_fail 的 designId 从 submit 阶段重走(新 taskId/traceId,旧 task 留档不复用);
|
||
落地语义=「infra_fail 行使该 designId 此前积累的 ok 阶段产物失效」(resume_state 即清空重积累,
|
||
infra 后若有新 ok 行则以新断点为准——支持多次中断/重放的时序正确性)。
|
||
"""
|
||
|
||
import hashlib
|
||
import json
|
||
import math
|
||
import os
|
||
import time
|
||
from datetime import datetime, timezone
|
||
|
||
# 账本受控枚举(§7.2 stage / status 取值,越界写入直接拒绝)
|
||
STAGES = (
|
||
"design", "schema_gate", "dedup_gate", "adversary", "fix", "submit",
|
||
"callback", "verify_package", "play", "judge", "publish", "postcheck",
|
||
)
|
||
STATUSES = ("ok", "fail", "infra_fail", "skip")
|
||
|
||
# 终判 decision 集合(resume 整体跳过的判据;fix 为中间 Verdict 不在内)
|
||
_TERMINAL_DECISIONS = ("accept", "kill")
|
||
|
||
|
||
class LedgerError(Exception):
|
||
"""账本契约违规(stage/status 越界、文件损坏等),必须立即暴露。"""
|
||
|
||
|
||
def now_iso():
|
||
"""统一 ISO8601 UTC 时间戳(账本 ts 字段)。"""
|
||
return datetime.now(timezone.utc).isoformat()
|
||
|
||
|
||
def mint_design_id(idea, template_id, round_):
|
||
"""
|
||
铸造 designId:sha256(idea+templateId+round) hex64(§6.2 公式原样,§16 D2-a 裁决确认
|
||
哈希因子含 round——round=0 与 round=1 的 designId 必不同,幂等键按轮独立)。
|
||
纯拼接、UTF-8 编码,保证可复现;rootDesignId = mint_design_id(idea, templateId, 0)。
|
||
"""
|
||
if round_ not in (0, 1):
|
||
# 防御:round 越出 0/1(回炉额度合计 1 轮,§10.1)即调用方缺陷
|
||
raise LedgerError("mint_design_id round 只能取 0/1,实际=%r" % (round_,))
|
||
raw = "%s%s%d" % (idea, template_id, round_)
|
||
return hashlib.sha256(raw.encode("utf-8")).hexdigest()
|
||
|
||
|
||
class Ledger(object):
|
||
"""append-only JSONL 账本写入器(每行落盘即 flush+fsync,崩溃最多丢最后一行)。"""
|
||
|
||
def __init__(self, path, batch_id):
|
||
self.path = path
|
||
self.batch_id = batch_id
|
||
# 确保目录存在(runs/<batchId>/)
|
||
os.makedirs(os.path.dirname(os.path.abspath(path)), exist_ok=True)
|
||
|
||
def _append_line(self, obj):
|
||
"""底层落一行 JSON(中文原样 ensure_ascii=False;append-only 不回写历史行)。"""
|
||
line = json.dumps(obj, ensure_ascii=False)
|
||
with open(self.path, "a", encoding="utf-8") as fh:
|
||
fh.write(line + "\n")
|
||
fh.flush()
|
||
os.fsync(fh.fileno())
|
||
|
||
def append_stage(self, design_id, round_, stage, status, data=None,
|
||
task_id=None, version_id=None, trace_id=None, llm_calls=0):
|
||
"""写阶段事件行(§7.2);stage/status 越出受控枚举立即拒绝。"""
|
||
if stage not in STAGES:
|
||
raise LedgerError("stage=%r 越出受控枚举 %s" % (stage, list(STAGES)))
|
||
if status not in STATUSES:
|
||
raise LedgerError("status=%r 越出受控枚举 %s" % (status, list(STATUSES)))
|
||
row = {
|
||
"type": "stage", "ts": now_iso(), "batchId": self.batch_id,
|
||
"designId": design_id, "round": round_, "stage": stage, "status": status,
|
||
"data": data if data is not None else {}, "llmCalls": llm_calls,
|
||
}
|
||
# 可选关联键:拿到即记,全链路可追溯
|
||
if task_id is not None:
|
||
row["taskId"] = task_id
|
||
if version_id is not None:
|
||
row["versionId"] = version_id
|
||
if trace_id is not None:
|
||
row["traceId"] = trace_id
|
||
self._append_line(row)
|
||
return row
|
||
|
||
def append_verdict(self, verdict):
|
||
"""写终判/回炉 Verdict 行(§7.2:Verdict 全字段 + ts + batchId)。"""
|
||
row = dict(verdict)
|
||
row["type"] = "verdict"
|
||
row["ts"] = now_iso()
|
||
row["batchId"] = self.batch_id
|
||
self._append_line(row)
|
||
return row
|
||
|
||
# ---------------------------------------------------------------- 读取与重放
|
||
|
||
@staticmethod
|
||
def load_rows(path):
|
||
"""读取账本全行;坏行(半截 JSON,崩溃尾行)跳过并告警计数,不让单行损坏废掉整本账。"""
|
||
rows = []
|
||
bad = 0
|
||
if not os.path.exists(path):
|
||
return rows, bad
|
||
with open(path, "r", encoding="utf-8") as fh:
|
||
for ln, line in enumerate(fh, 1):
|
||
line = line.strip()
|
||
if not line:
|
||
continue
|
||
try:
|
||
rows.append(json.loads(line))
|
||
except ValueError:
|
||
bad += 1 # 崩溃残行:跳过(append-only 语义下仅可能是最后一行)
|
||
return rows, bad
|
||
|
||
@staticmethod
|
||
def resume_state(path):
|
||
"""
|
||
扫描账本 → 每个(轮级)designId 的重放档案(§7.2 幂等重放语义 + §16 D2-a 口径):
|
||
{
|
||
designId: {
|
||
"terminal": bool, # 已有终判 verdict(accept/kill)→ 该创意整体跳过
|
||
"fix_issued": bool, # 已落 decision=fix 的中间 Verdict(round=0 designId 上)
|
||
# → 创意以 round=1 的 designId 从回炉轮续跑
|
||
"replay_from_submit": bool, # 最近一段以 infra_fail 收尾 → 从 submit 阶段重走(新 taskId)
|
||
"ok_stages": {(round, stage): data}, # 已 ok 阶段产物(断点续跑的状态恢复源)
|
||
"published": bool, # publish 阶段已 ok(金丝雀计数跨 resume 累计)
|
||
}
|
||
}
|
||
时序语义:infra_fail 行使该 designId 此前积累的 ok 产物失效(清空 ok_stages、标记重走);
|
||
其后若出现新 ok 行(上次 resume 已重放过一段)则以新断点为准、重走标记清回——
|
||
保证多次中断/重放序列下「最后一段连续记录」即有效断点。
|
||
"""
|
||
rows, bad = Ledger.load_rows(path)
|
||
state = {}
|
||
|
||
def slot(design_id):
|
||
return state.setdefault(design_id, {
|
||
"terminal": False, "fix_issued": False,
|
||
"replay_from_submit": False, "ok_stages": {}, "published": False,
|
||
})
|
||
|
||
for row in rows:
|
||
design_id = row.get("designId")
|
||
if not design_id:
|
||
continue
|
||
info = slot(design_id)
|
||
if row.get("type") == "verdict":
|
||
if row.get("decision") in _TERMINAL_DECISIONS:
|
||
info["terminal"] = True
|
||
elif row.get("decision") == "fix":
|
||
# D5 中间 Verdict:落在 round=0 的 designId 上,标记创意须续跑回炉轮
|
||
info["fix_issued"] = True
|
||
elif row.get("type") == "stage":
|
||
data = row.get("data") or {}
|
||
if row.get("status") == "infra_fail":
|
||
# infra_fail → 该 designId 重放时从 submit 重走(§7.2);此前 ok 产物失效
|
||
info["replay_from_submit"] = True
|
||
info["ok_stages"] = {}
|
||
elif row.get("status") == "ok":
|
||
# 批级 QA 行(换模型抽检 audit / 稳定性重测 retest)不作断点产物:
|
||
# 防止其 findings(可能来自不同模型)在 resume 时顶替原始对抗产物
|
||
if data.get("audit") or data.get("retest"):
|
||
continue
|
||
# infra 后出现新 ok 行 = 重放已推进,重走标记清回(最后一段为准)
|
||
info["replay_from_submit"] = False
|
||
key = (row.get("round"), row.get("stage"))
|
||
info["ok_stages"][key] = data
|
||
if row.get("stage") == "publish":
|
||
info["published"] = True
|
||
return state, bad
|
||
|
||
|
||
# ============================== 预算闸(§7.3:八项数字硬编码常量,改值=改代码走 PR) ==============================
|
||
|
||
class BudgetExceeded(Exception):
|
||
"""预算闸触发(携带受控记账类别,编排器按 D7 infra 路径处置)。"""
|
||
|
||
def __init__(self, message, category):
|
||
super(BudgetExceeded, self).__init__(message)
|
||
self.category = category # 受控词后缀,如 "llm_budget"
|
||
|
||
|
||
class BudgetGuard(object):
|
||
"""
|
||
预算闸(§7.3 八项;全部硬编码常量+中文注释,改值=改代码走 PR)。
|
||
计数状态按 designId 维护;时钟可注入便于单测。
|
||
"""
|
||
|
||
# ① LLM 调用 ≤8 次/创意:design(1)+adversary(1)+fix 轮 design+adversary(2)+schema 重出(1)+重试余量(3);
|
||
# 超限即该 designId 判 infra_fail 停止
|
||
MAX_LLM_CALLS_PER_IDEA = 8
|
||
# ② 实玩 ≤3 次/创意:首玩 + runnable 重试 1 + infra 重试 1
|
||
MAX_PLAYS_PER_IDEA = 3
|
||
# ③ 批时长 30 分钟强制收口:超时未完成 designId 全记 infra_fail(不计分母),出报告
|
||
BATCH_TIMEOUT_SECONDS = 30 * 60
|
||
# ④ 浏览器池 ≤2 实例,实玩串行/池(防互扰,评审版 §3.4)
|
||
BROWSER_POOL_MAX = 2
|
||
# ⑤ 策划并发 ≤3,间隔 0.3s(沿 spike 通道纪律;v1 编排器顺序执行天然满足 ≤3,间隔由 llm_client 节流强制)
|
||
DESIGNER_CONCURRENCY_MAX = 3
|
||
DESIGNER_CALL_INTERVAL_SECONDS = 0.3
|
||
# ⑥ 金丝雀:每批入 feed ≤10;accept 超出部分账本记 accept_unpublished,不发布
|
||
CANARY_FEED_MAX_PER_BATCH = 10
|
||
# ⑦ 换模型抽检:accept 的 20%(向上取整)重裁;分歧率 >10% → 冻结本批发布开关 + 告警行落账本
|
||
AUDIT_RECHECK_RATIO = 0.20
|
||
AUDIT_DIVERGENCE_FREEZE_RATIO = 0.10
|
||
# ⑧ 对抗稳定性重测:每批随机 ≥3 条 GameDesign 同 prompt 重测,P0/P1 判定翻转即不一致,一致率入批报告
|
||
ADVERSARY_RETEST_MIN = 3
|
||
|
||
def __init__(self, clock=time.monotonic):
|
||
self._clock = clock # 可注入时钟(单测用桩)
|
||
self._batch_started_at = None # 批开始时刻(③ 闸基准)
|
||
self._llm_calls = {} # designId → 已发起 LLM 调用次数
|
||
self._plays = {} # designId → 已发起实玩次数
|
||
|
||
# ---------- ③ 批时长 ----------
|
||
def start_batch(self):
|
||
"""记录批开始时刻(30 分钟闸基准)。"""
|
||
self._batch_started_at = self._clock()
|
||
|
||
def batch_timed_out(self):
|
||
"""批是否超过 30 分钟强制收口线。"""
|
||
if self._batch_started_at is None:
|
||
raise LedgerError("BudgetGuard 未 start_batch 即查询超时——编排器调用缺陷")
|
||
return (self._clock() - self._batch_started_at) > self.BATCH_TIMEOUT_SECONDS
|
||
|
||
# ---------- ① LLM 调用 ----------
|
||
def note_llm_call(self, design_id):
|
||
"""登记一次 LLM 调用(每次 HTTP 尝试都计数);超 ①闸 抛 BudgetExceeded。"""
|
||
used = self._llm_calls.get(design_id, 0)
|
||
if used >= self.MAX_LLM_CALLS_PER_IDEA:
|
||
raise BudgetExceeded(
|
||
"designId=%s LLM 调用已达上限 %d 次/创意(§7.3-①),判 infra_fail 停止"
|
||
% (design_id, self.MAX_LLM_CALLS_PER_IDEA), category="llm_budget")
|
||
self._llm_calls[design_id] = used + 1
|
||
|
||
def llm_calls_used(self, design_id):
|
||
"""该创意已消耗的 LLM 调用数(账本 llmCalls 字段与批报告成本统计用)。"""
|
||
return self._llm_calls.get(design_id, 0)
|
||
|
||
# ---------- ② 实玩 ----------
|
||
def note_play(self, design_id):
|
||
"""登记一次实玩;超 ②闸(3 次=首玩+runnable 重试+infra 重试)抛 BudgetExceeded。"""
|
||
used = self._plays.get(design_id, 0)
|
||
if used >= self.MAX_PLAYS_PER_IDEA:
|
||
raise BudgetExceeded(
|
||
"designId=%s 实玩已达上限 %d 次/创意(§7.3-②),判 infra_fail 停止"
|
||
% (design_id, self.MAX_PLAYS_PER_IDEA), category="play_budget")
|
||
self._plays[design_id] = used + 1
|
||
|
||
def plays_used(self, design_id):
|
||
"""该创意已消耗的实玩次数(playAttempt 推算与报告用)。"""
|
||
return self._plays.get(design_id, 0)
|
||
|
||
# ---------- ⑥ 金丝雀 ----------
|
||
def canary_slot_available(self, published_count):
|
||
"""本批是否还有金丝雀发布额度(入参=本批已发布数,跨 resume 从账本累计)。"""
|
||
return published_count < self.CANARY_FEED_MAX_PER_BATCH
|
||
|
||
# ---------- ⑦ 换模型抽检 ----------
|
||
def audit_sample_size(self, accept_count):
|
||
"""换模型抽检条数 = accept 的 20% 向上取整(accept=0 → 0)。"""
|
||
return int(math.ceil(accept_count * self.AUDIT_RECHECK_RATIO)) if accept_count > 0 else 0
|
||
|
||
|
||
# ============================== 熔断(§7.4 四条:触发即停批/冻结,写报告,退出码非 0) ==============================
|
||
|
||
class CircuitBreaker(object):
|
||
"""
|
||
熔断器(§7.4)。计数口径:
|
||
- 分母 = 已处理创意数(到达终判或 infra_fail 的 designId 数,本次运行内);
|
||
- 条款1:批内 infra_fail 占比 >30% → 停批告警;
|
||
- 条款2:连续 5 条同因 kill(reasons 首项相同)→ 停批告警(prompt/环境系统性问题);
|
||
- 条款3:试玩环境不可用(浏览器/前端 serve/staging 探活失败)→ 整批暂停,禁降级为「跳过试玩」;
|
||
- 条款4:换模型抽检分歧 >10% → 仅冻结发布段(裁决与账本照常),报告显著标注。
|
||
另含 F10 发布链连败护栏:连续 2 条 accept_publish_fail → 停发布段(§11 F10)。
|
||
"""
|
||
|
||
INFRA_RATIO_TRIP = 0.30 # 条款1 阈值(严格大于才触发)
|
||
CONSECUTIVE_SAME_KILL_TRIP = 5 # 条款2 阈值
|
||
PUBLISH_FAIL_CONSECUTIVE_STOP = 2 # F10:发布链连续失败 2 条 → 停发布段
|
||
|
||
def __init__(self):
|
||
self.processed = 0 # 分母:已处理创意数
|
||
self.infra_failed = 0 # 条款1 分子
|
||
self._kill_streak_reason = None
|
||
self._kill_streak = 0 # 条款2 连续同因 kill 计数
|
||
self.halted_reason = None # 条款3:整批暂停原因(非空即停批)
|
||
self.publish_frozen_reason = None # 条款4 / F10:发布段冻结原因
|
||
self._publish_fail_streak = 0 # F10 连续发布失败计数
|
||
|
||
# ---------- 计数入口 ----------
|
||
def note_outcome(self, kind, first_reason=None):
|
||
"""
|
||
登记一个创意的最终处理结果。
|
||
kind ∈ {"accept","kill","infra_fail"};kill 须携 reasons 首项(条款2 同因判定)。
|
||
"""
|
||
if kind not in ("accept", "kill", "infra_fail"):
|
||
raise LedgerError("熔断计数 kind=%r 越界" % (kind,))
|
||
self.processed += 1
|
||
if kind == "infra_fail":
|
||
self.infra_failed += 1
|
||
# infra 不打断 kill 连击计数?——条款2 口径为「连续 5 条同因 kill」,
|
||
# 以创意处理序列为准:非 kill 的结果(accept/infra)打断连续性。
|
||
self._kill_streak_reason, self._kill_streak = None, 0
|
||
elif kind == "kill":
|
||
if first_reason is not None and first_reason == self._kill_streak_reason:
|
||
self._kill_streak += 1
|
||
else:
|
||
self._kill_streak_reason, self._kill_streak = first_reason, 1
|
||
else: # accept
|
||
self._kill_streak_reason, self._kill_streak = None, 0
|
||
|
||
# ---------- 条款判定 ----------
|
||
def infra_ratio_tripped(self):
|
||
"""条款1:infra_fail 占比 >30%(分母=已处理创意数;分母为 0 不触发)。"""
|
||
if self.processed == 0:
|
||
return False
|
||
return (self.infra_failed / float(self.processed)) > self.INFRA_RATIO_TRIP
|
||
|
||
def kill_streak_tripped(self):
|
||
"""条款2:连续 5 条同因 kill(reasons 首项相同)。"""
|
||
return self._kill_streak >= self.CONSECUTIVE_SAME_KILL_TRIP
|
||
|
||
def halt_env(self, reason):
|
||
"""条款3:试玩环境不可用 → 整批暂停(runnable 证据不可豁免,禁降级为跳过试玩)。"""
|
||
self.halted_reason = reason
|
||
|
||
def is_halted(self):
|
||
return self.halted_reason is not None
|
||
|
||
def freeze_publish(self, reason):
|
||
"""条款4 / F10 / F12:冻结发布段(裁决与账本照常;已发布金丝雀保留——现无下架接口 R7)。"""
|
||
if self.publish_frozen_reason is None:
|
||
self.publish_frozen_reason = reason
|
||
|
||
def publish_frozen(self):
|
||
return self.publish_frozen_reason is not None
|
||
|
||
def note_publish_result(self, ok):
|
||
"""F10:登记发布链结果;连续 2 条失败 → 停发布段(告警冻结)。"""
|
||
if ok:
|
||
self._publish_fail_streak = 0
|
||
else:
|
||
self._publish_fail_streak += 1
|
||
if self._publish_fail_streak >= self.PUBLISH_FAIL_CONSECUTIVE_STOP:
|
||
self.freeze_publish("F10:发布链连续 %d 条失败,停发布段" % self._publish_fail_streak)
|
||
|
||
def tripped_summary(self):
|
||
"""汇总当前熔断态(批报告/退出码用);返回 [中文描述],空列表=未熔断。"""
|
||
out = []
|
||
if self.infra_ratio_tripped():
|
||
out.append("条款1:infra_fail 占比 %d/%d 超 30%%,停批" % (self.infra_failed, self.processed))
|
||
if self.kill_streak_tripped():
|
||
out.append("条款2:连续 %d 条同因 kill(%s),停批" % (self._kill_streak, self._kill_streak_reason))
|
||
if self.is_halted():
|
||
out.append("条款3:试玩环境不可用(%s),整批暂停" % self.halted_reason)
|
||
if self.publish_frozen():
|
||
out.append("发布段冻结:%s" % self.publish_frozen_reason)
|
||
return out
|