From 594d4f0501a39dcb03e9ce5d0983471f0ccaa2c6 Mon Sep 17 00:00:00 2001 From: ymq Date: Fri, 18 Sep 2026 10:53:16 +0800 Subject: [PATCH] =?UTF-8?q?fix(agent=5Floop):=20=E5=AE=A1=E6=A0=B8?= =?UTF-8?q?=E8=AF=81=E6=8D=AE=E9=93=BE=E5=9B=9B=E7=BC=BA=E9=99=B7=E6=A0=B9?= =?UTF-8?q?=E6=B2=BB(2026-09-18=20pbls=20M4a=E7=86=94=E6=96=AD=E5=AE=9E?= =?UTF-8?q?=E9=94=A4)=E2=80=94=E2=80=94=E2=91=A0infra=E5=9B=9E=E7=BD=AE?= =?UTF-8?q?=E7=8B=AC=E7=AB=8B=E9=A2=84=E7=AE=97params.infra=5Frequeue(?= =?UTF-8?q?=E5=8E=9F=E4=B8=8E=E8=B4=A8=E9=87=8F=E9=80=80=E5=9B=9E=E5=85=B1?= =?UTF-8?q?=E4=BA=ABretry=5Fcount,M4a=E4=B8=89=E8=BD=AE=3D2=E6=AC=A1LLM?= =?UTF-8?q?=E6=95=85=E9=9A=9C+1=E6=AC=A1=E7=9C=9F=E5=AE=9E=E9=80=80?= =?UTF-8?q?=E5=9B=9E=E5=8D=B3=E8=A7=A6=E9=A1=B6=E7=86=94=E6=96=AD;?= =?UTF-8?q?=E8=80=97=E5=B0=BD=E6=97=B6=E6=8A=ACretry=5Fcount=E5=88=B0?= =?UTF-8?q?=E4=B8=8A=E9=99=90=E9=98=B2failed=5Fpoller=E8=AF=AF=E5=88=A4?= =?UTF-8?q?=E7=9E=AC=E6=97=B6retry=5Ftask=E8=AE=A9develop=E4=BB=8E?= =?UTF-8?q?=E5=A4=B4=E9=87=8D=E5=81=9A;=E7=9C=9F=E5=AE=9E=E8=A3=81?= =?UTF-8?q?=E5=86=B3/=E4=BA=BA=E5=B7=A5reset=E6=B8=85=E9=9B=B6;=E6=B4=BE?= =?UTF-8?q?=E7=94=9F=E4=BB=BB=E5=8A=A1pop)=E2=91=A1qc=5Freject=E8=BE=BE?= =?UTF-8?q?=E4=B8=8A=E9=99=90=E8=B7=AF=E5=BE=84=E4=B8=8D=E5=86=8D=E4=B8=A2?= =?UTF-8?q?=E5=BC=83questions=E6=98=8E=E7=BB=86,=E8=A1=A5=E5=86=99last=5Fe?= =?UTF-8?q?rror=E9=9A=8Ffault=5Freport=E7=9B=B4=E8=BE=BE=E4=BA=BA=E5=B7=A5?= =?UTF-8?q?(=E5=8E=9F=E5=8F=AA=E7=95=99=E5=88=86=E6=95=B0=E6=97=A0?= =?UTF-8?q?=E6=B8=85=E5=8D=95,=E6=9C=AC=E6=AC=A1=E9=9D=A0=E7=BF=BBllm=5Fca?= =?UTF-8?q?ll=5Ftrace=E6=89=8D=E6=89=BE=E5=9B=9E3=E9=A1=B9=E6=9C=AA?= =?UTF-8?q?=E8=BF=87=E6=98=8E=E7=BB=86)=E2=91=A2=E5=AE=A1=E6=A0=B8?= =?UTF-8?q?=E6=B3=A8=E5=85=A5content[:6000]=E8=A3=B8=E6=88=AA=E6=96=AD?= =?UTF-8?q?=E6=94=B9=5Freview=5Fcontent=5Fpreview(limit=3D12000=E6=AD=A3?= =?UTF-8?q?=E5=B8=B8=E5=85=A8=E9=87=8F;=E8=B6=85=E9=95=BF=E4=BF=9D?= =?UTF-8?q?=E5=A4=B460%+=E5=B0=BE40%=E4=BF=9D=E5=BC=95=E6=93=8E=E6=A0=B8?= =?UTF-8?q?=E9=AA=8C=E6=AE=B5+=E6=98=BE=E5=BC=8F=E6=88=AA=E6=96=AD?= =?UTF-8?q?=E6=A0=87=E8=AE=B0,QC=E4=B8=8D=E5=86=8D=E6=8A=8A=E5=BC=95?= =?UTF-8?q?=E6=93=8E=E6=88=AA=E6=8E=89=E7=9A=84=E8=AF=81=E6=8D=AE=E5=88=A4?= =?UTF-8?q?=E6=88=90=E8=AF=81=E6=8D=AE=E7=BC=BA=E5=A4=B1)=E2=91=A3files=5F?= =?UTF-8?q?json=E6=9E=84=E5=BB=BA:=E5=AD=98=E5=9C=A8=E6=80=A7=E8=BF=87?= =?UTF-8?q?=E6=BB=A4=E5=B9=BD=E7=81=B5=E6=96=87=E4=BB=B6(=E5=AE=9E?= =?UTF-8?q?=E6=B5=8B=5Fm4a=5Fpatch.py=E5=85=88=E5=86=99=E5=90=8E=E5=88=A0?= =?UTF-8?q?=E8=A2=ABQC=E5=88=A4=E4=B8=B4=E6=97=B6=E8=A1=A5=E4=B8=81?= =?UTF-8?q?=E6=B7=B7=E5=85=A5)+=E5=B9=B6=E5=85=A5=5Fenforce=5Fgit=5Fclosur?= =?UTF-8?q?e=E5=AE=9E=E9=99=85=E6=8F=90=E4=BA=A4=E6=96=87=E4=BB=B6(?= =?UTF-8?q?=E5=AE=9E=E6=B5=8Bsql/m4a=5Fddl.sql=E8=A2=ABadd=20-A=E6=8F=90?= =?UTF-8?q?=E4=BA=A4=E4=BD=86=E6=B8=85=E5=8D=95=E6=BC=8F=E6=94=B6=E8=A2=AB?= =?UTF-8?q?QC=E5=88=A4=E4=BA=A4=E4=BB=98=E4=BB=B6=E7=BC=BA=E5=A4=B1)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- pipeline_service/agent_loop.py | 212 +++++++++++++++++++++++----- pipeline_service/task_capability.py | 14 +- 2 files changed, 192 insertions(+), 34 deletions(-) diff --git a/pipeline_service/agent_loop.py b/pipeline_service/agent_loop.py index 1080a54..a35fe7e 100644 --- a/pipeline_service/agent_loop.py +++ b/pipeline_service/agent_loop.py @@ -1406,6 +1406,32 @@ def _build_review_files_block(files_json): + "\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): """获取项目关联的所有仓库。""" 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 行的 # deploy_test 触发判断误判为「模块级 develop」,导致应用脚手架 develop approved 后跳过 deploy_test。 new_params.pop('pm_assigned', None) + # 2026-09-18:基础设施回置计数是「当前任务当前轮」的预算,派生任务从零开始 + new_params.pop('infra_requeue', None) # 幂等(2026-08-25 新增):并发场景下多个任务几乎同时 approved,每个都会走到这里尝试 # 创建下一阶段任务 → 同一迭代同一角色出现多个重复任务(历史上已出现「同一迭代两个 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 后被 # 「模块级 develop → 跳过 deploy_test」误判,任务链断在 approved。与 _create_next_task 的 pop 对齐。 new_params.pop('pm_assigned', None) + # 2026-09-18:回退重做任务的基础设施回置预算从零开始(同 _create_next_task) + new_params.pop('infra_requeue', None) if prev_task: new_params['previous_task_id'] = getattr(prev_task, 'id', '') 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 的口头声明变成引擎保证的事实; · 收口失败(git 命令 rc≠0)→ 返回 ok=False + 可行动原因,由调用方拒绝 deliver。 - Returns: (ok: bool, report: str) report 为空=本任务未触及代码仓库;否则为核验记录 - (调用方回填交付件,给 QC/PM 当仪表盘:agent 声称的收口是否属实)。 + Returns: (ok: bool, report: str, committed_files: list[str]) report 为空=本任务未触及 + 代码仓库;否则为核验记录(调用方回填交付件,给 QC/PM 当仪表盘:agent 声称的 + 收口是否属实)。committed_files=本次引擎收口实际提交的文件(space_dir 相对 + 路径,2026-09-18 起供 files_json 对齐磁盘事实)。 """ touched = set() 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'): touched.add(os.path.join(space_dir, parts[0], parts[1])) if not touched: - return True, '' + return True, '', [] lines, ok = [], True + committed_files = [] # 本次引擎收口实际提交的文件(space_dir 相对路径) for rp in sorted(touched): label = os.path.relpath(rp, space_dir) if not os.path.isdir(rp): @@ -2746,11 +2779,24 @@ async def _enforce_git_closure(space_dir, written_files, do_commit=True): ok = False lines.append(f"{label}: 收口后仍有 {len(still.splitlines())} 项未提交变更(可能被并发写入)") 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: lines.append(f"{label}: agent 未提交,引擎代为收口 {n_dirty} 项变更({str(r.get('message',''))[:80]})") elif not was_git: 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): @@ -3342,7 +3388,7 @@ async def role_agent_run(project_id, role, agent_id=None, model_name=None): # 全部落盘之后执行(此前收口会漏提交 files 内容)。引擎代为 add+commit 本任务 # 写入过的 apps/modules 仓库——「已提交」从 agent 口头声明变成引擎保证的事实; # 核验记录(含「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) if _gc_report: 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 + 摘要快照。 # 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 = [] for f in list(written_files) + files_written + ([file_path] if file_path else []): try: + if not os.path.exists(f): + continue # 幽灵文件(写过又删)不进清单 rel = os.path.relpath(f, space_dir) except Exception: continue if rel and not rel.startswith('..') and rel not in body_files: 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) 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) 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) 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'] comment = decision.get('comment', '') + # 真实裁决达成 → 基础设施回置计数清零(预算语义=连续故障,与 QC 循环对称) + await _reset_infra_requeue(sor, task_id) # 防御:非终结角色被 review_complete 收尾 → 修正为 approved(避免断链)。 # 只有无下一角色(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="没有交付件") 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) 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 冒泡人工。 # ⚠️ 不走 retry_task(failed→submitted):QC 任务只被 qc poller 按 state='qc_review' # 认领,回 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 - _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) 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( "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}) await sor.sqlExe("COMMIT", {}) 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} 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, - error=err + f";重审 {rc} 次仍无有效决策,需人工介入") - logger.error(f"qc_review_run no decision: task={task_id} -> failed (retry exhausted)") + error=err + f";基础设施重审 {infra_cnt} 次仍无有效决策,需人工介入") + logger.error(f"qc_review_run no decision: task={task_id} -> failed (infra budget exhausted)") return {"status": "failed", "task_id": task_id, "error": err} if decision['status'] == 'approved': + # 真实裁决达成 → 基础设施回置计数清零(预算语义=连续故障) + await _reset_infra_requeue(sor, task_id) 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', '')) 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', '')} else: 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 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: @@ -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} 退回重做") return {"status": "rejected", "task_id": task_id, "comment": comment} if not ok and "failed" in msg: - # 重复退回达上限 → 已转 failed,不再建 qc_reject 问题(failed poller 会报 fault_report 给人工) - logger.info(f"qc_review_run reject->failed: task={task_id} {msg}") - return {"status": "failed", "task_id": task_id, "comment": msg} + # 重复退回达上限 → 已转 failed。不另建 qc_reject 问题(failed poller + # 会报 fault_report,避免重复冒泡),但必须把未过项明细补写进 + # 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 - rejection_q = decision.get("questions") or comment or "交付件不合规" await raise_problem("qc_reject", rejection_q, "agent.qc", tenant_id=project_id, task_id=task_id, first_handler_role=task_role, @@ -4608,6 +4680,76 @@ def _classify_failure(last_error: str) -> str: 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, who, agent_id, err_msg): """审核阶段(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, 由对应 poller 重新认领。 - retry_count 上限保护:连续基础设施故障(上游持续超时窗)达 task_max_retry 时 - 返回 False,由调用方 mark_failed → failed_poller 判 fault → pause + fault_report + 独立预算(2026-09-18 pbls M4a 事故根治):基础设施回置计数走 + 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) @@ -4638,10 +4783,13 @@ async def _requeue_after_review_infra_failure(sor, task_id, project_id, review_s if cur_state != review_state: return False, f"任务状态 {cur_state or '不存在'} 非 {review_state},跳过回置" max_retry = await get_max_task_retry(sor) - if rc >= max_retry: - return False, f"审核回置 {rc} 次仍连续故障(达上限 {max_retry}),需人工介入" + ok_bump, infra_cnt = await _bump_infra_requeue(sor, task_id, 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( - "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}$", {"e": err_msg[:4000], "tid": task_id, "st": review_state}) 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', from_state=review_state, to_state=review_state, 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) logger.warning(f"review infra failure requeued: task={task_id} state={review_state} " - f"({rc + 1}/{max_retry}) err={err_msg[:120]}") - return True, f"requeued {rc + 1}/{max_retry}" + f"(infra {infra_cnt}/{max_retry}) err={err_msg[:120]}") + return True, f"requeued infra {infra_cnt}/{max_retry}" async def handle_failed_task(task_id: str, project_id: str) -> dict: diff --git a/pipeline_service/task_capability.py b/pipeline_service/task_capability.py index 64b4b0d..83fbdde 100644 --- a/pipeline_service/task_capability.py +++ b/pipeline_service/task_capability.py @@ -311,11 +311,21 @@ async def reset_task_retry(task_id, tenant_id, who=None, agent_id=None): from_state = getattr(recs[0], 'state', '') if from_state not in ('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( "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')", - {"tid": task_id, "tn": tenant_id}) + {"p": _new_params, "tid": task_id, "tn": tenant_id}) chk = await sor.R('pipeline_tasks', {'id': task_id}) await sor.sqlExe("COMMIT", {}) # 成功判定=UPDATE 真实生效(state 与 retry_count 双核对),不只看终态——