diff --git a/pipeline_service/agent_loop.py b/pipeline_service/agent_loop.py index 226ba3a..d2d1591 100644 --- a/pipeline_service/agent_loop.py +++ b/pipeline_service/agent_loop.py @@ -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 个前置结束即启动 # 不变量(任何策略下都成立,代码守门): # · 依赖任务不存在 → 阻塞(可能是编造 ID,D1/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 待办 + 发出后,若被依赖任务恢复流转(failed→retry→submitted→…→approved),阻塞原因 + 已不存在,待办应自动收口——否则用户界面永远挂着过时的「依赖永不可达」, + 还要人工逐条点掉(今早实测: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 周期调用。 diff --git a/pipeline_service/init.py b/pipeline_service/init.py index cca668b..e21c902 100644 --- a/pipeline_service/init.py +++ b/pipeline_service/init.py @@ -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_agents,appbase params 表,默认 3)。 # 以 DB state='running' 计数为准(权威),超出上限时本轮回合不再派发。 from .workspace import get_max_concurrent_agents