fix(orchestration): 审核基础设施故障不再让develop从头重做+failed不算依赖dead+dependency_blocked待办自动收口(2026-09-16 pbls双事故根治)

事故A(M1a QC超时→develop全量重做): qc/pm_review_run的LLM调用异常走mark_failed→failed_poller retry_task回submitted→develop把已交付完好的工作从头重做(10:42实测,交付件20分钟前刚提交)。审核方基础设施故障≠开发方交付失败,新增_requeue_after_review_infra_failure:回置qc_review/review重新认领(对齐无决策路径已有模式),retry达上限才mark_failed→fault冒泡,有限循环有出口。

事故B(dependency_blocked误报+永不关闭): DEP_DEAD含failed→上游任务瞬时failed(QC超时/重启打断)时给下游误发'永不可达'人工待办(10:11/10:42两起实锤,M1a failed后3秒给M1b发待办,9秒后failed_poller就重试成功)。failed是过渡态(failed_poller 60s内必接管:瞬时→retry,耗尽→fault→waiting+pause+fault_report通知人工),归wait语义;真dead只剩cancelled。配对新增_resolve_stale_dependency_blocks挂agent_poller:阻塞原因消失的pending待办自动done收口(result_data留痕),冒泡有去重+收口有配对,出口才完整。

verify_dep_gate.py同步:新增D5 failed/waiting→wait用例,20/20全绿
This commit is contained in:
ymq 2026-09-16 12:17:22 +08:00
parent 02c54264c8
commit f99b554e9e
2 changed files with 148 additions and 5 deletions

View File

@ -781,8 +781,9 @@ async def _task_deps_satisfied(sor, depends_on_raw, task_id='', dep_policy=None)
现改为先查回实际存在的依赖 id 集合**逐个比对**缺失的依赖视为未满足
D3 PM 可能编造 20 ID_pm_create_tasks len>=20 的引用不校验存在性
配合 D1 会让门控完全失效现在缺失即阻塞编造 ID 不再能绕过
D4 依赖处于 cancelled/failed 等永不可达终态时原实现让依赖方永久 submitted死锁
D4 依赖处于 cancelled 等永不可达终态时原实现让依赖方永久 submitted死锁
且无人知晓现在识别为 dead 依赖并在 reason 里标注由调用方冒泡人工处理
failed 不算 dead它是过渡态failed_poller 必然接管 DEP_DEAD 注释
自依赖depends_on 含自身 id 会永久阻塞识别为 dead 依赖
"""
deps = []
@ -823,10 +824,16 @@ async def _task_deps_satisfied(sor, depends_on_raw, task_id='', dep_policy=None)
# {"mode": "at_least", "n": k} 至少 k 个前置结束即启动
# 不变量(任何策略下都成立,代码守门):
# · 依赖任务不存在 → 阻塞(可能是编造 IDD1/D3
# · 依赖处于 cancelled/failed → 永不可达阻塞并冒泡D4
# · 依赖处于 cancelled → 永不可达阻塞并冒泡D4
# · 自依赖 / depends_on 解析失败 → 阻塞
# 🔴 failed 不是永不可达2026-09-16 pbls 误报根治failed 是过渡态——failed_poller
# 60s 内必处理瞬时→retry_task 回 submitted耗尽→fault→waiting+pause+fault_report
# 通知人工),永远有归属者。把 failed 当 dead 会在依赖任务瞬时故障(如 QC 的 LLM 超时)
# 时给下游任务误发 dependency_blocked 人工待办实测M1a QC 超时 failed 后 3 秒就给
# M1b 发了「永不可达」待办9 秒后 failed_poller 就重试成功——待办纯属误报且永不自动关闭)。
# failed/waiting 归入 wait 语义:继续等待,不冒泡;真死锁由 fault_report 通道负责通知。
DEP_DONE = ('completed', 'approved')
DEP_DEAD = ('cancelled', 'failed')
DEP_DEAD = ('cancelled',)
def _parse_dep_policy(policy):
@ -1083,6 +1090,67 @@ async def _bubble_blocked_dependency(sor, project_id, task_id, task_title, reaso
logger.warning(f"依赖阻塞冒泡: task={task_id} title={task_title} reason={reason}")
async def _resolve_stale_dependency_blocks(sor):
"""逃逸阀配对收口:依赖阻塞原因已消失的 pending 待办自动关闭2026-09-16
不变量任何自动门禁都必须配出口同样适用于出口本身dependency_blocked 待办
发出后若被依赖任务恢复流转failedretrysubmittedapproved阻塞原因
已不存在待办应自动收口否则用户界面永远挂着过时的依赖永不可达
还要人工逐条点掉今早实测design 任务瞬时 failed 触发的 M1a 待办design
approved 后待办仍 pending
判定复用 _task_deps_satisfied 单一语义来源重新求值后 kind 不再是 dead
依赖已恢复/已满足 关闭待办status=done + result_data 留痕并写审计
仍是 dead如依赖真被 cancelled 保留待办等人工
"""
recs = await sor.sqlExe(
"SELECT id, task_id, project_id FROM pipeline_human_tasks "
"WHERE task_type='dependency_blocked' AND status='pending' LIMIT 50", {})
await sor.sqlExe("COMMIT", {})
if not recs:
return 0
closed = 0
for r in recs:
hid = getattr(r, 'id', '')
tid = getattr(r, 'task_id', '') or ''
pid = getattr(r, 'project_id', '') or ''
if not tid:
continue
trecs = await sor.sqlExe(
"SELECT depends_on, params, state FROM pipeline_tasks WHERE id=${tid}$", {"tid": tid})
await sor.sqlExe("COMMIT", {})
if not trecs:
# 任务本身已被删除 → 待办失去载体,同样收口
still_dead, why = False, '任务已不存在'
else:
tstate = getattr(trecs[0], 'state', '') or ''
if tstate not in ('submitted', 'waiting'):
# 任务自己已恢复流转running/approved/...)→ 阻塞已解除
still_dead, why = False, f'任务已恢复state={tstate}'
else:
try:
_tp = json.loads(getattr(trecs[0], 'params', '{}') or '{}')
except (json.JSONDecodeError, TypeError):
_tp = {}
ok_dep, why = await _task_deps_satisfied(
sor, getattr(trecs[0], 'depends_on', '') or '', task_id=tid,
dep_policy=_tp.get('dep_policy'))
still_dead = (not ok_dep) and any(
k in why for k in ('永不可达', '不存在', '依赖自身', '解析失败'))
if still_dead:
continue
await sor.sqlExe(
"UPDATE pipeline_human_tasks SET status='done', qc_status='passed', "
"qc_comment='阻塞原因已消失,机制自动收口', submitted_by='system.orchestrator', "
"submitted_at=NOW(), result_data=${rd}$ WHERE id=${hid}$ AND status='pending'",
{"rd": json.dumps({"auto_resolved": True, "reason": why}, ensure_ascii=False),
"hid": hid})
await sor.sqlExe("COMMIT", {})
closed += 1
logger.info(f"dependency_blocked 自动收口: ht={hid} task={tid[:8]} why={why}")
return closed
async def _claim_task(sor, tenant_id, role, state='submitted', match_role=True, set_state='running'):
role = _normalize_role(role)
where_role = "AND (role=${role}$ OR role='')" if match_role else ""
@ -3329,8 +3397,16 @@ async def pm_review_run(project_id, agent_id=None, model_name=None):
raw = await llm_call_msgs(msgs, model=model_name, temperature=0.3, org_id=org_id, project_id=project_id, session_id='task:%s' % task_id)
except Exception as e:
err_msg = f"{type(e).__name__}: {str(e)[:400]}"
# PM 审核环节基础设施故障 → 回置 review 重审,不打 fail 让上游重做
# (与 qc_review_run 同款2026-09-16 pbls M1a 事故根治)。
ok_rq, msg_rq = await _requeue_after_review_infra_failure(
sor, task_id, project_id, TASK_REVIEW, 'agent.pm', agent_id, err_msg)
if ok_rq:
return {"status": "requeued", "task_id": task_id,
"error": err_msg, "requeue": msg_rq}
from .task_capability import mark_failed
await mark_failed(task_id, project_id, who="agent.pm", agent_id=agent_id, error=err_msg)
await mark_failed(task_id, project_id, who="agent.pm", agent_id=agent_id,
error=f"{err_msg}{msg_rq}")
return {"status": "failed", "task_id": task_id, "error": str(e)[:200]}
act = _parse_agent_action(raw)
@ -3699,8 +3775,18 @@ async def qc_review_run(project_id, agent_id=None, model_name=None):
raw = await llm_call_msgs(msgs, model=model_name, temperature=0.2, org_id=org_id, project_id=project_id, session_id='task:%s' % task_id)
except Exception as e:
err_msg = f"{type(e).__name__}: {str(e)[:400]}"
# 审核方基础设施故障 ≠ develop 交付失败2026-09-16 pbls M1a 事故根治):
# 原实现 mark_failed → failed_poller retry_task 回 submitted → develop 把
# 已交付的工作从头重做(交付件完好却被丢弃)。正确语义:回置 qc_review
# 重新认领重审;连续故障达上限才 mark_failed → fault 冒泡人工。
ok_rq, msg_rq = await _requeue_after_review_infra_failure(
sor, task_id, project_id, S_QC_REVIEW, 'agent.qc', agent_id, err_msg)
if ok_rq:
return {"status": "requeued", "task_id": task_id,
"error": err_msg, "requeue": msg_rq}
from .task_capability import mark_failed
await mark_failed(task_id, project_id, who="agent.qc", agent_id=agent_id, error=err_msg)
await mark_failed(task_id, project_id, who="agent.qc", agent_id=agent_id,
error=f"{err_msg}{msg_rq}")
return {"status": "failed", "task_id": task_id, "error": str(e)[:200]}
act = _parse_agent_action(raw)
@ -4039,6 +4125,54 @@ def _classify_failure(last_error: str) -> str:
return "fault"
async def _requeue_after_review_infra_failure(sor, task_id, project_id, review_state,
who, agent_id, err_msg):
"""审核阶段QC/PM基础设施故障 → 回置审核状态重新认领,不打 fail 让上游重做。
依据2026-09-16 pbls M1a 事故QC 环节上游 LLM 瞬时超时曾走 mark_failed
failed_poller retry_task 任务回 submitted develop 把整个开发工作**从头重做**
而它 20 分钟前提交的交付件完好无损审核方的基础设施故障 开发方的交付失败
正确语义是审核没做成 重新审核回置 qc_review/review + claimed_by
由对应 poller 重新认领
retry_count 上限保护连续基础设施故障上游持续超时窗 task_max_retry
返回 False由调用方 mark_failed failed_poller fault pause + fault_report
冒泡人工有限循环有出口不静默空转
Returns: (requeued: bool, msg: str)
"""
from .workspace import get_max_task_retry
_rc = await sor.sqlExe(
"SELECT retry_count, state FROM pipeline_tasks WHERE id=${tid}$", {"tid": task_id})
await sor.sqlExe("COMMIT", {})
rc, cur_state = 0, ''
if _rc:
try:
rc = int(getattr(_rc[0], 'retry_count', 0) or 0)
except (TypeError, ValueError):
rc = 0
cur_state = getattr(_rc[0], 'state', '') or ''
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}),需人工介入"
await sor.sqlExe(
"UPDATE pipeline_tasks SET claimed_by=NULL, retry_count=retry_count+1, "
"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", {})
from .audit import record_audit
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]}",
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}"
async def handle_failed_task(task_id: str, project_id: str) -> dict:
"""处理一个失败任务:判定重跑或报故障。由 failed poller 周期调用。

View File

@ -643,6 +643,15 @@ def load_pipeline_service():
"WHERE state='running' AND pipeline_id='role_task' "
"AND updated_at < (NOW() - INTERVAL 20 MINUTE)", {})
# 逃逸阀配对收口2026-09-16依赖阻塞原因已消失的 dependency_blocked
# 待办自动关闭。冒泡有去重、收口有配对,出口才是完整的(不变量:任何自动
# 门禁的冒泡待办都必须能随原因消失自动收口,不许永久挂着让人手动点)。
try:
from .agent_loop import _resolve_stale_dependency_blocks
await _resolve_stale_dependency_blocks(sor)
except Exception as _rse:
debug(f"dependency_blocked 收口失败: {_rse}")
# 全局并发 agent 数限制max_concurrent_agentsappbase params 表,默认 3
# 以 DB state='running' 计数为准(权威),超出上限时本轮回合不再派发。
from .workspace import get_max_concurrent_agents