fix(agent_loop): 审核证据链四缺陷根治(2026-09-18 pbls M4a熔断实锤)——①infra回置独立预算params.infra_requeue(原与质量退回共享retry_count,M4a三轮=2次LLM故障+1次真实退回即触顶熔断;耗尽时抬retry_count到上限防failed_poller误判瞬时retry_task让develop从头重做;真实裁决/人工reset清零;派生任务pop)②qc_reject达上限路径不再丢弃questions明细,补写last_error随fault_report直达人工(原只留分数无清单,本次靠翻llm_call_trace才找回3项未过明细)③审核注入content[:6000]裸截断改_review_content_preview(limit=12000正常全量;超长保头60%+尾40%保引擎核验段+显式截断标记,QC不再把引擎截掉的证据判成证据缺失)④files_json构建:存在性过滤幽灵文件(实测_m4a_patch.py先写后删被QC判临时补丁混入)+并入_enforce_git_closure实际提交文件(实测sql/m4a_ddl.sql被add -A提交但清单漏收被QC判交付件缺失)

This commit is contained in:
ymq 2026-09-18 10:53:16 +08:00
parent c8b44fe423
commit 594d4f0501
2 changed files with 192 additions and 34 deletions

View File

@ -1406,6 +1406,32 @@ def _build_review_files_block(files_json):
+ "\n".join("- " + f for f in files)) + "\n".join("- " + f for f in files))
def _review_content_preview(content, limit=12000):
"""审核注入的交付件正文(2026-09-18 pbls M4a 事故根治,替代裸 [:6000])。
原实现 content[:6000] 硬切:引擎的机械核验/git 收口/import 闭包三段证据
**追加在正文尾部**,正文稍超 6000 字符就把 QC 最需要的证据段切掉——M4a 全长
6789 字符,QC 看到的机械核验段断在 `pbl_age`,据此判「验证证据不完整」退回;
判定依据是引擎自己截掉的证据(第 3 项未过实锤)。
策略:正常交付件(≤ limit)**完整注入,等效不截断**;病态超长才保头+保尾
(证据段在尾部必须保住)+ 显式截断标记——QC 明确知道中段被引擎略过、可
read_file 交付件文档补读,不再把「引擎截断」误判成「证据缺失」。
"""
if not content:
return ''
if len(content) <= limit:
return content
head = int(limit * 0.6)
tail = limit - head
omitted = len(content) - head - tail
return (content[:head]
+ f"\n\n……[中段 {omitted} 字符被引擎截断保护省略(交付件超长)。"
+ "注意:这不是证据缺失——如需中段内容请 read_file 交付件文档补读。"
+ "尾部引擎核验段完整保留如下]……\n\n"
+ content[-tail:])
async def _get_project_repos(sor, project_id): async def _get_project_repos(sor, project_id):
"""获取项目关联的所有仓库。""" """获取项目关联的所有仓库。"""
recs = await sor.sqlExe( recs = await sor.sqlExe(
@ -1882,6 +1908,8 @@ async def _create_next_task(sor, project_id, task, next_role, pm_comment='',
# 带 pm_assigned=True,design approved 后自动创建的 develop 任务继承了它,被 1785 行的 # 带 pm_assigned=True,design approved 后自动创建的 develop 任务继承了它,被 1785 行的
# deploy_test 触发判断误判为「模块级 develop」,导致应用脚手架 develop approved 后跳过 deploy_test。 # deploy_test 触发判断误判为「模块级 develop」,导致应用脚手架 develop approved 后跳过 deploy_test。
new_params.pop('pm_assigned', None) new_params.pop('pm_assigned', None)
# 2026-09-18:基础设施回置计数是「当前任务当前轮」的预算,派生任务从零开始
new_params.pop('infra_requeue', None)
# 幂等(2026-08-25 新增):并发场景下多个任务几乎同时 approved,每个都会走到这里尝试 # 幂等(2026-08-25 新增):并发场景下多个任务几乎同时 approved,每个都会走到这里尝试
# 创建下一阶段任务 → 同一迭代同一角色出现多个重复任务(历史上已出现「同一迭代两个 design # 创建下一阶段任务 → 同一迭代同一角色出现多个重复任务(历史上已出现「同一迭代两个 design
@ -1992,6 +2020,8 @@ async def _rollback_task_chain(sor, project_id, task_id, rollback_role, comment)
# pm_assigned 会污染——若 target 是应用级 develop(历史脏数据带 pm_assigned=True),approve 后被 # pm_assigned 会污染——若 target 是应用级 develop(历史脏数据带 pm_assigned=True),approve 后被
# 「模块级 develop → 跳过 deploy_test」误判,任务链断在 approved。与 _create_next_task 的 pop 对齐。 # 「模块级 develop → 跳过 deploy_test」误判,任务链断在 approved。与 _create_next_task 的 pop 对齐。
new_params.pop('pm_assigned', None) new_params.pop('pm_assigned', None)
# 2026-09-18:回退重做任务的基础设施回置预算从零开始(同 _create_next_task)
new_params.pop('infra_requeue', None)
if prev_task: if prev_task:
new_params['previous_task_id'] = getattr(prev_task, 'id', '') new_params['previous_task_id'] = getattr(prev_task, 'id', '')
new_params['previous_role'] = _normalize_role(getattr(prev_task, 'role', '')) new_params['previous_role'] = _normalize_role(getattr(prev_task, 'role', ''))
@ -2692,8 +2722,10 @@ async def _enforce_git_closure(space_dir, written_files, do_commit=True):
把「已提交」从 agent 的口头声明变成引擎保证的事实; 把「已提交」从 agent 的口头声明变成引擎保证的事实;
· 收口失败(git 命令 rc≠0)→ 返回 ok=False + 可行动原因,由调用方拒绝 deliver。 · 收口失败(git 命令 rc≠0)→ 返回 ok=False + 可行动原因,由调用方拒绝 deliver。
Returns: (ok: bool, report: str) report 为空=本任务未触及代码仓库;否则为核验记录 Returns: (ok: bool, report: str, committed_files: list[str]) report 为空=本任务未触及
(调用方回填交付件,给 QC/PM 当仪表盘:agent 声称的收口是否属实)。 代码仓库;否则为核验记录(调用方回填交付件,给 QC/PM 当仪表盘:agent 声称的
收口是否属实)。committed_files=本次引擎收口实际提交的文件(space_dir 相对
路径,2026-09-18 起供 files_json 对齐磁盘事实)。
""" """
touched = set() touched = set()
for f in written_files or []: for f in written_files or []:
@ -2707,9 +2739,10 @@ async def _enforce_git_closure(space_dir, written_files, do_commit=True):
if len(parts) >= 2 and parts[0] in ('apps', 'modules'): if len(parts) >= 2 and parts[0] in ('apps', 'modules'):
touched.add(os.path.join(space_dir, parts[0], parts[1])) touched.add(os.path.join(space_dir, parts[0], parts[1]))
if not touched: if not touched:
return True, '' return True, '', []
lines, ok = [], True lines, ok = [], True
committed_files = [] # 本次引擎收口实际提交的文件(space_dir 相对路径)
for rp in sorted(touched): for rp in sorted(touched):
label = os.path.relpath(rp, space_dir) label = os.path.relpath(rp, space_dir)
if not os.path.isdir(rp): if not os.path.isdir(rp):
@ -2746,11 +2779,24 @@ async def _enforce_git_closure(space_dir, written_files, do_commit=True):
ok = False ok = False
lines.append(f"{label}: 收口后仍有 {len(still.splitlines())} 项未提交变更(可能被并发写入)") lines.append(f"{label}: 收口后仍有 {len(still.splitlines())} 项未提交变更(可能被并发写入)")
continue continue
# 收口成功:从收口前 porcelain 提取实际提交的文件清单(2026-09-18 M4a 事故——
# agent 用 run_shell 写的 sql/m4a_ddl.sql 被 add -A 提交但不在 written_files
# 跟踪里,files_json 漏收 → QC 判「索引声明的交付件缺失」误退)
for ln in dirty.splitlines():
fp = ln[3:].strip() if len(ln) > 3 else ''
# 重命名格式 "old -> new"
if ' -> ' in fp:
fp = fp.split(' -> ', 1)[1].strip()
if not fp:
continue
rel = os.path.join(label, fp)
if rel not in committed_files:
committed_files.append(rel)
if n_dirty: if n_dirty:
lines.append(f"{label}: agent 未提交,引擎代为收口 {n_dirty} 项变更({str(r.get('message',''))[:80]})") lines.append(f"{label}: agent 未提交,引擎代为收口 {n_dirty} 项变更({str(r.get('message',''))[:80]})")
elif not was_git: elif not was_git:
lines.append(f"{label}: 非 git 目录,引擎已 git init 并首次提交") lines.append(f"{label}: 非 git 目录,引擎已 git init 并首次提交")
return ok, "\n".join(lines) return ok, "\n".join(lines), committed_files
def _validate_deliverable_type(role, deliverable_type, allowed_types): def _validate_deliverable_type(role, deliverable_type, allowed_types):
@ -3342,7 +3388,7 @@ async def role_agent_run(project_id, role, agent_id=None, model_name=None):
# 全部落盘之后执行(此前收口会漏提交 files 内容)。引擎代为 add+commit 本任务 # 全部落盘之后执行(此前收口会漏提交 files 内容)。引擎代为 add+commit 本任务
# 写入过的 apps/modules 仓库——「已提交」从 agent 口头声明变成引擎保证的事实; # 写入过的 apps/modules 仓库——「已提交」从 agent 口头声明变成引擎保证的事实;
# 核验记录(含「agent 未提交,引擎代为收口」)回填交付件正文,QC/PM 拿真实仪表盘。 # 核验记录(含「agent 未提交,引擎代为收口」)回填交付件正文,QC/PM 拿真实仪表盘。
_gc_ok, _gc_report = await _enforce_git_closure( _gc_ok, _gc_report, _gc_committed_files = await _enforce_git_closure(
space_dir, list(written_files) + files_written, do_commit=True) space_dir, list(written_files) + files_written, do_commit=True)
if _gc_report: if _gc_report:
result_text += ("\n\n---\n## git 收口核验(引擎自动执行,非 agent 声明)\n" result_text += ("\n\n---\n## git 收口核验(引擎自动执行,非 agent 声明)\n"
@ -3391,14 +3437,25 @@ async def role_agent_run(project_id, role, agent_id=None, model_name=None):
# 交付件本体文件清单(工作空间相对路径):write_file 实写 + deliver files + 摘要快照。 # 交付件本体文件清单(工作空间相对路径):write_file 实写 + deliver files + 摘要快照。
# PM/QC 审核注入时随清单强制 read_file 本体,不再只看 content 摘要(pbls 教训)。 # PM/QC 审核注入时随清单强制 read_file 本体,不再只看 content 摘要(pbls 教训)。
# 2026-09-18 M4a 事故修两处证据链失真:
# ① 存在性过滤——agent 先 write_file 后删除的临时文件(实测 _m4a_patch.py/
# _fix_budget.py)曾留在清单里,QC 判「临时补丁混入、范围不可追溯」退回;
# ② 并入引擎收口提交的文件——agent 用 run_shell 写的文件(实测 sql/m4a_ddl.sql)
# 不在 written_files 跟踪里,git 收口 add -A 提交了但清单漏收,QC 判
# 「索引声明的交付件缺失」退回。清单=引擎担保的磁盘事实,两侧对齐。
body_files = [] body_files = []
for f in list(written_files) + files_written + ([file_path] if file_path else []): for f in list(written_files) + files_written + ([file_path] if file_path else []):
try: try:
if not os.path.exists(f):
continue # 幽灵文件(写过又删)不进清单
rel = os.path.relpath(f, space_dir) rel = os.path.relpath(f, space_dir)
except Exception: except Exception:
continue continue
if rel and not rel.startswith('..') and rel not in body_files: if rel and not rel.startswith('..') and rel not in body_files:
body_files.append(rel) body_files.append(rel)
for rel in _gc_committed_files: # 引擎收口实际提交的文件补进清单
if rel and rel not in body_files:
body_files.append(rel)
files_json = json.dumps(body_files[:50], ensure_ascii=False) files_json = json.dumps(body_files[:50], ensure_ascii=False)
await sor.C("pipeline_deliverables", { await sor.C("pipeline_deliverables", {
@ -3798,7 +3855,8 @@ async def pm_review_run(project_id, agent_id=None, model_name=None):
repo_state = await _get_repo_state(space_dir, project_name) repo_state = await _get_repo_state(space_dir, project_name)
repos_str = ", ".join([r['name'] for r in repos]) if repos else "无" repos_str = ", ".join([r['name'] for r in repos]) if repos else "无"
content_preview = deliverable_content[:6000] # 2026-09-18:裸 [:6000] 曾把尾部引擎核验段切掉(M4a 误杀实锤),改走保尾预览
content_preview = _review_content_preview(deliverable_content)
role_skills = await _build_role_skills_block(sor, project_id, _normalize_role('pm'), org_id) role_skills = await _build_role_skills_block(sor, project_id, _normalize_role('pm'), org_id)
pm_system = PM_SYSTEM_PROMPT.replace('__TITLE__', title)\ pm_system = PM_SYSTEM_PROMPT.replace('__TITLE__', title)\
@ -3936,6 +3994,8 @@ async def pm_review_run(project_id, agent_id=None, model_name=None):
status = decision['status'] status = decision['status']
comment = decision.get('comment', '') comment = decision.get('comment', '')
# 真实裁决达成 → 基础设施回置计数清零(预算语义=连续故障,与 QC 循环对称)
await _reset_infra_requeue(sor, task_id)
# 防御:非终结角色被 review_complete 收尾 → 修正为 approved(避免断链)。 # 防御:非终结角色被 review_complete 收尾 → 修正为 approved(避免断链)。
# 只有无下一角色(deploy_prod)的 review_complete 才是合法终结。否则 requirement/design/develop/ # 只有无下一角色(deploy_prod)的 review_complete 才是合法终结。否则 requirement/design/develop/
@ -4199,7 +4259,8 @@ async def qc_review_run(project_id, agent_id=None, model_name=None):
await qc_reject_task(task_id, project_id, who="agent.qc", agent_id=agent_id, comment="没有交付件") await qc_reject_task(task_id, project_id, who="agent.qc", agent_id=agent_id, comment="没有交付件")
return {"status": "rejected", "task_id": task_id, "reason": "没有交付件"} return {"status": "rejected", "task_id": task_id, "reason": "没有交付件"}
content_preview = deliverable_content[:6000] # 2026-09-18:裸 [:6000] 曾把尾部引擎核验段切掉(M4a 误杀实锤),改走保尾预览
content_preview = _review_content_preview(deliverable_content)
role_skills = await _build_role_skills_block(sor, project_id, _normalize_role('qc'), org_id) role_skills = await _build_role_skills_block(sor, project_id, _normalize_role('qc'), org_id)
qc_system = QC_SYSTEM_PROMPT.replace('__TITLE__', title)\ qc_system = QC_SYSTEM_PROMPT.replace('__TITLE__', title)\
@ -4309,34 +4370,33 @@ async def qc_review_run(project_id, agent_id=None, model_name=None):
# 达上限 mark_failed → failed_poller 判 fault → pause + fault_report 冒泡人工。 # 达上限 mark_failed → failed_poller 判 fault → pause + fault_report 冒泡人工。
# ⚠️ 不走 retry_task(failed→submitted):QC 任务只被 qc poller 按 state='qc_review' # ⚠️ 不走 retry_task(failed→submitted):QC 任务只被 qc poller 按 state='qc_review'
# 认领,回 submitted 会搁浅(agent poller 不管 agent.qc 角色)。 # 认领,回 submitted 会搁浅(agent poller 不管 agent.qc 角色)。
# 2026-09-18 M4a 事故:无决策=审核方 LLM 故障,属基础设施回置,改走
# params.infra_requeue 独立预算(原 retry_count+1 会吃掉质量退回预算——
# M4a 三轮里两轮 LLM 故障+一次真实退回即触顶熔断的根因之一)。
from .workspace import get_max_task_retry from .workspace import get_max_task_retry
_rc = await sor.sqlExe(
"SELECT retry_count FROM pipeline_tasks WHERE id=${tid}$", {"tid": task_id})
await sor.sqlExe("COMMIT", {})
rc = 0
try:
rc = int(getattr(_rc[0], 'retry_count', 0) or 0) if _rc else 0
except (TypeError, ValueError):
rc = 0
max_retry = await get_max_task_retry(sor) max_retry = await get_max_task_retry(sor)
err = ("QC 审核循环耗尽 16 轮仍未输出有效决策(review_approve/review_reject)," err = ("QC 审核循环耗尽 16 轮仍未输出有效决策(review_approve/review_reject),"
"疑似上游模型输出异常") "疑似上游模型输出异常")
if rc < max_retry: ok_bump, infra_cnt = await _bump_infra_requeue(sor, task_id, max_retry)
if ok_bump:
await sor.sqlExe( await sor.sqlExe(
"UPDATE pipeline_tasks SET state='qc_review', claimed_by=NULL, " "UPDATE pipeline_tasks SET state='qc_review', claimed_by=NULL, "
"retry_count=retry_count+1, last_error=${e}$, updated_at=NOW() " "last_error=${e}$, updated_at=NOW() "
"WHERE id=${tid}$", {"e": err, "tid": task_id}) "WHERE id=${tid}$", {"e": err, "tid": task_id})
await sor.sqlExe("COMMIT", {}) await sor.sqlExe("COMMIT", {})
logger.warning(f"qc_review_run no decision: task={task_id} requeued " logger.warning(f"qc_review_run no decision: task={task_id} requeued "
f"(retry {rc + 1}/{max_retry})") f"(infra {infra_cnt}/{max_retry})")
return {"status": "requeued", "task_id": task_id, "error": err} return {"status": "requeued", "task_id": task_id, "error": err}
from .task_capability import mark_failed from .task_capability import mark_failed
await _floor_retry_count_at_max(sor, task_id, max_retry)
await mark_failed(task_id, project_id, who="agent.qc", agent_id=agent_id, await mark_failed(task_id, project_id, who="agent.qc", agent_id=agent_id,
error=err + f";重审 {rc} 次仍无有效决策,需人工介入") error=err + f";基础设施重审 {infra_cnt} 次仍无有效决策,需人工介入")
logger.error(f"qc_review_run no decision: task={task_id} -> failed (retry exhausted)") logger.error(f"qc_review_run no decision: task={task_id} -> failed (infra budget exhausted)")
return {"status": "failed", "task_id": task_id, "error": err} return {"status": "failed", "task_id": task_id, "error": err}
if decision['status'] == 'approved': if decision['status'] == 'approved':
# 真实裁决达成 → 基础设施回置计数清零(预算语义=连续故障)
await _reset_infra_requeue(sor, task_id)
from .task_capability import qc_approve_task from .task_capability import qc_approve_task
await qc_approve_task(task_id, project_id, who="agent.qc", agent_id=agent_id, comment=decision.get('comment', '')) await qc_approve_task(task_id, project_id, who="agent.qc", agent_id=agent_id, comment=decision.get('comment', ''))
if task_kind == 'human_task_qc' and human_task_id: if task_kind == 'human_task_qc' and human_task_id:
@ -4347,6 +4407,11 @@ async def qc_review_run(project_id, agent_id=None, model_name=None):
return {"status": "approved", "task_id": task_id, "comment": decision.get('comment', '')} return {"status": "approved", "task_id": task_id, "comment": decision.get('comment', '')}
else: else:
comment = decision.get('comment', '') comment = decision.get('comment', '')
# 未过项明细(QC 的 questions 字段)提前取出——2026-09-18 M4a 事故:
# 达上限路径原直接 return 丢弃明细,fault_report 只有「契合度 8.50、
# 未过项 3 项」没有清单,人工无法行动(本次靠翻 llm_call_trace 落盘才找回)。
rejection_q = decision.get("questions") or comment or "交付件不合规"
await _reset_infra_requeue(sor, task_id)
from .task_capability import qc_reject_task from .task_capability import qc_reject_task
ok, msg = await qc_reject_task(task_id, project_id, who="agent.qc", agent_id=agent_id, comment=comment) ok, msg = await qc_reject_task(task_id, project_id, who="agent.qc", agent_id=agent_id, comment=comment)
if task_kind == 'human_task_qc' and human_task_id: if task_kind == 'human_task_qc' and human_task_id:
@ -4356,11 +4421,18 @@ async def qc_review_run(project_id, agent_id=None, model_name=None):
logger.info(f"qc_review_run reject human_task: {human_task_id} 退回重做") logger.info(f"qc_review_run reject human_task: {human_task_id} 退回重做")
return {"status": "rejected", "task_id": task_id, "comment": comment} return {"status": "rejected", "task_id": task_id, "comment": comment}
if not ok and "failed" in msg: if not ok and "failed" in msg:
# 重复退回达上限 → 已转 failed,不再建 qc_reject 问题(failed poller 会报 fault_report 给人工) # 重复退回达上限 → 已转 failed。不另建 qc_reject 问题(failed poller
logger.info(f"qc_review_run reject->failed: task={task_id} {msg}") # 会报 fault_report,避免重复冒泡),但必须把未过项明细补写进
return {"status": "failed", "task_id": task_id, "comment": msg} # last_error——fault_report 正文带 last_error,明细随故障直达人工。
await sor.sqlExe(
"UPDATE pipeline_tasks SET last_error=CONCAT(IFNULL(last_error,''), "
"${q}$) WHERE id=${tid}$",
{"q": ("\n\n【QC 未过项明细(达上限转人工,引擎自动附上)】\n"
+ rejection_q)[:6000], "tid": task_id})
await sor.sqlExe("COMMIT", {})
logger.info(f"qc_review_run reject->failed: task={task_id} {msg} (questions persisted to last_error)")
return {"status": "failed", "task_id": task_id, "comment": msg, "question": rejection_q}
from .communication import raise_problem from .communication import raise_problem
rejection_q = decision.get("questions") or comment or "交付件不合规"
await raise_problem("qc_reject", rejection_q, "agent.qc", await raise_problem("qc_reject", rejection_q, "agent.qc",
tenant_id=project_id, task_id=task_id, tenant_id=project_id, task_id=task_id,
first_handler_role=task_role, first_handler_role=task_role,
@ -4608,6 +4680,76 @@ def _classify_failure(last_error: str) -> str:
return "fault" return "fault"
async def _bump_infra_requeue(sor, task_id, max_retry):
"""基础设施回置的独立预算计数(2026-09-18 pbls M4a 事故根治)。
事故形状:审核方 LLM 故障(16 轮无决策 / 上游超时)回置与开发方质量退回
**共享同一个 retry_count 预算(3)**——M4a 前两轮全是 QC 侧 LLM 基础设施
故障吃掉 2/3 预算,第三轮首次真实质量退回(8.50 分)即触顶 failed+项目熔断。
审核方的故障不该消耗开发方的退回预算。
计数器落 params.infra_requeue(不加列;reset_task_retry 人工恢复时清零,
QC 出真实裁决时清零——语义=「连续」基础设施故障预算)。
Returns: (ok: bool, cnt: int) ok=False 表示独立预算已耗尽。
"""
recs = await sor.sqlExe(
"SELECT params FROM pipeline_tasks WHERE id=${tid}$", {"tid": task_id})
await sor.sqlExe("COMMIT", {})
p, cnt = {}, 0
if recs:
try:
p = json.loads(getattr(recs[0], 'params', '') or '{}')
except Exception:
p = {}
if not isinstance(p, dict):
p = {}
try:
cnt = int(p.get('infra_requeue', 0) or 0)
except (TypeError, ValueError):
cnt = 0
if cnt >= max_retry:
return False, cnt
p['infra_requeue'] = cnt + 1
await sor.sqlExe(
"UPDATE pipeline_tasks SET params=${p}$, updated_at=NOW() WHERE id=${tid}$",
{"p": json.dumps(p, ensure_ascii=False), "tid": task_id})
await sor.sqlExe("COMMIT", {})
return True, cnt + 1
async def _reset_infra_requeue(sor, task_id):
"""QC/PM 出真实裁决后清零基础设施回置计数(预算语义=连续故障)。"""
try:
recs = await sor.sqlExe(
"SELECT params FROM pipeline_tasks WHERE id=${tid}$", {"tid": task_id})
await sor.sqlExe("COMMIT", {})
if not recs:
return
p = json.loads(getattr(recs[0], 'params', '') or '{}')
if isinstance(p, dict) and p.get('infra_requeue'):
p['infra_requeue'] = 0
await sor.sqlExe(
"UPDATE pipeline_tasks SET params=${p}$ WHERE id=${tid}$",
{"p": json.dumps(p, ensure_ascii=False), "tid": task_id})
await sor.sqlExe("COMMIT", {})
except Exception as e:
logger.warning(f"_reset_infra_requeue failed: task={task_id} {e}")
async def _floor_retry_count_at_max(sor, task_id, max_retry):
"""基础设施预算耗尽转 failed 前,把 retry_count 抬到上限。
否则 failed_poller 的 handle_failed_task 见 retry_count < max_retry 会走
_classify_failure——last_error 含「超时」被判瞬时 → retry_task 回 submitted
→ develop 把完好的交付件从头重做(正是 2026-09-16 M1a 事故的形状)。
抬到上限保证确定性走 fault → pause + fault_report 冒泡人工。
"""
await sor.sqlExe(
"UPDATE pipeline_tasks SET retry_count=GREATEST(retry_count, ${m}$) "
"WHERE id=${tid}$", {"m": max_retry, "tid": task_id})
await sor.sqlExe("COMMIT", {})
async def _requeue_after_review_infra_failure(sor, task_id, project_id, review_state, async def _requeue_after_review_infra_failure(sor, task_id, project_id, review_state,
who, agent_id, err_msg): who, agent_id, err_msg):
"""审核阶段(QC/PM)基础设施故障 → 回置审核状态重新认领,不打 fail 让上游重做。 """审核阶段(QC/PM)基础设施故障 → 回置审核状态重新认领,不打 fail 让上游重做。
@ -4618,8 +4760,11 @@ async def _requeue_after_review_infra_failure(sor, task_id, project_id, review_s
正确语义是「审核没做成 → 重新审核」:回置 qc_review/review + 清 claimed_by, 正确语义是「审核没做成 → 重新审核」:回置 qc_review/review + 清 claimed_by,
由对应 poller 重新认领。 由对应 poller 重新认领。
retry_count 上限保护:连续基础设施故障(上游持续超时窗)达 task_max_retry 时 独立预算(2026-09-18 pbls M4a 事故根治):基础设施回置计数走
返回 False,由调用方 mark_failed → failed_poller 判 fault → pause + fault_report params.infra_requeue,**不再吃 retry_count**(那是质量退回的预算,混用会让
首次真实退回即触顶熔断)。连续基础设施故障达 task_max_retry 时返回 False,
并先把 retry_count 抬到上限(防 failed_poller 误判瞬时→retry_task 让 develop
从头重做),由调用方 mark_failed → failed_poller 判 fault → pause + fault_report
冒泡人工——有限循环有出口,不静默空转。 冒泡人工——有限循环有出口,不静默空转。
Returns: (requeued: bool, msg: str) Returns: (requeued: bool, msg: str)
@ -4638,10 +4783,13 @@ async def _requeue_after_review_infra_failure(sor, task_id, project_id, review_s
if cur_state != review_state: if cur_state != review_state:
return False, f"任务状态 {cur_state or '不存在'} 非 {review_state},跳过回置" return False, f"任务状态 {cur_state or '不存在'} 非 {review_state},跳过回置"
max_retry = await get_max_task_retry(sor) max_retry = await get_max_task_retry(sor)
if rc >= max_retry: ok_bump, infra_cnt = await _bump_infra_requeue(sor, task_id, max_retry)
return False, f"审核回置 {rc} 次仍连续故障(达上限 {max_retry}),需人工介入" if not ok_bump:
await _floor_retry_count_at_max(sor, task_id, max_retry)
return False, (f"审核基础设施连续故障 {infra_cnt} 次达独立预算上限({max_retry}),"
f"需人工介入(质量退回预算 retry_count={rc} 未被消耗)")
await sor.sqlExe( await sor.sqlExe(
"UPDATE pipeline_tasks SET claimed_by=NULL, retry_count=retry_count+1, " "UPDATE pipeline_tasks SET claimed_by=NULL, "
"last_error=${e}$, updated_at=NOW() WHERE id=${tid}$ AND state=${st}$", "last_error=${e}$, updated_at=NOW() WHERE id=${tid}$ AND state=${st}$",
{"e": err_msg[:4000], "tid": task_id, "st": review_state}) {"e": err_msg[:4000], "tid": task_id, "st": review_state})
await sor.sqlExe("COMMIT", {}) await sor.sqlExe("COMMIT", {})
@ -4649,11 +4797,11 @@ async def _requeue_after_review_infra_failure(sor, task_id, project_id, review_s
await record_audit(project_id, 'pipeline_tasks', task_id, 'requeue', await record_audit(project_id, 'pipeline_tasks', task_id, 'requeue',
from_state=review_state, to_state=review_state, from_state=review_state, to_state=review_state,
who=who, agent_id=agent_id, who=who, agent_id=agent_id,
detail=f"审核基础设施故障回置重审({rc + 1}/{max_retry}):{err_msg[:300]}", detail=f"审核基础设施故障回置重审(infra {infra_cnt}/{max_retry}):{err_msg[:300]}",
sor=sor) sor=sor)
logger.warning(f"review infra failure requeued: task={task_id} state={review_state} " logger.warning(f"review infra failure requeued: task={task_id} state={review_state} "
f"({rc + 1}/{max_retry}) err={err_msg[:120]}") f"(infra {infra_cnt}/{max_retry}) err={err_msg[:120]}")
return True, f"requeued {rc + 1}/{max_retry}" return True, f"requeued infra {infra_cnt}/{max_retry}"
async def handle_failed_task(task_id: str, project_id: str) -> dict: async def handle_failed_task(task_id: str, project_id: str) -> dict:

View File

@ -311,11 +311,21 @@ async def reset_task_retry(task_id, tenant_id, who=None, agent_id=None):
from_state = getattr(recs[0], 'state', '') from_state = getattr(recs[0], 'state', '')
if from_state not in ('waiting', 'failed', 'submitted'): if from_state not in ('waiting', 'failed', 'submitted'):
return False, f"恢复失败(任务状态 {from_state} 不在 waiting/failed/submitted)" return False, f"恢复失败(任务状态 {from_state} 不在 waiting/failed/submitted)"
# 2026-09-18 M4a 事故配套:人工恢复时基础设施回置预算(params.infra_requeue)
# 一并清零——否则任务带着满 infra 计数重新认领,一次 LLM 瞬时故障即触顶。
_new_params = getattr(recs[0], 'params', '') or ''
try:
_p = json.loads(_new_params or '{}')
if isinstance(_p, dict) and _p.get('infra_requeue'):
_p['infra_requeue'] = 0
_new_params = json.dumps(_p, ensure_ascii=False)
except Exception:
pass
await sor.sqlExe( await sor.sqlExe(
"UPDATE pipeline_tasks SET retry_count=0, state='submitted', " "UPDATE pipeline_tasks SET retry_count=0, state='submitted', "
"claimed_by=NULL, last_error=NULL, updated_at=NOW() " "claimed_by=NULL, last_error=NULL, params=${p}$, updated_at=NOW() "
"WHERE id=${tid}$ AND tenant_id=${tn}$ AND state IN ('waiting','failed','submitted')", "WHERE id=${tid}$ AND tenant_id=${tn}$ AND state IN ('waiting','failed','submitted')",
{"tid": task_id, "tn": tenant_id}) {"p": _new_params, "tid": task_id, "tn": tenant_id})
chk = await sor.R('pipeline_tasks', {'id': task_id}) chk = await sor.R('pipeline_tasks', {'id': task_id})
await sor.sqlExe("COMMIT", {}) await sor.sqlExe("COMMIT", {})
# 成功判定=UPDATE 真实生效(state 与 retry_count 双核对),不只看终态—— # 成功判定=UPDATE 真实生效(state 与 retry_count 双核对),不只看终态——