fix(flow): 上游QC过滤须无条件执行且建任务留在无编排任务分支内——上一版结构错位致编排任务在办时每15秒重复建(实测2分钟堆9个)

This commit is contained in:
yumoqing 2026-09-04 16:28:29 +08:00
parent 7733ff31cf
commit 442de85cc3

View File

@ -582,53 +582,56 @@ async def reconcile_project(sor, project_id, project_name=""):
# (评分项 →(资质 ∥ 要求+骨架)→ 成本收益),编排意图归 PM
# 依赖门禁由代码强制dispatch_analysis_dim 校验上游 QC 全过)。
# QC 退回重做仍走上面 qc_imp 分支,不受编排影响。
# 2026-09-04 修复(中电新项目实测):上游 QC 未过的维度是合法等待,
# 不是「掉链」——PM 即使接到编排任务也派发不了dispatch_analysis_dim
# 依赖门禁拒绝)。把这类维度留在 pending 会让编排任务空转(无可派发维度
# 即完成),空转轮次计入 orch_done → MAX_ORCH_ATTEMPT 轮后误抛
# 「PM 解析编排掉链」人工待办。只有真正可派发的维度才驱动编排/掉链计数。
# 注意:过滤必须无条件执行(放在「编排任务在办」判定之前)——上一版把
# 过滤塞进 else 分支、创建挪出分支外,导致编排任务在办时仍用未过滤列表
# 每 15 秒重复建任务(实测 2 分钟堆 9 个)。
if first_run_pending:
qc_passed = await _qc_passed_types(sor, project_id)
dispatchable, waiting_up = [], []
for d, lb in first_run_pending:
up_ok = all(t in qc_passed
for u in DIM_UPSTREAM.get(d, ())
for t in DIM_QC_TYPES[u])
(dispatchable if up_ok else waiting_up).append((d, lb))
if waiting_up:
acts.append("waiting: 维度「%s」等上游 QC 通过(合法等待,不计掉链)"
% "".join(lb for _, lb in waiting_up))
first_run_pending = dispatchable
if first_run_pending:
if await _open_orch_tasks(sor, project_id) > 0:
acts.append("waiting: PM 解析编排任务在办")
else:
# 2026-09-04 修复(中电新项目实测):上游 QC 未过的维度是合法等待,
# 不是「掉链」——PM 即使接到编排任务也派发不了dispatch_analysis_dim
# 依赖门禁拒绝)。把这类维度留在 pending 会让编排任务空转(无可派发维度
# 即完成),空转轮次计入 orch_done → MAX_ORCH_ATTEMPT 轮后误抛
# 「PM 解析编排掉链」人工待办。只有真正可派发的维度才驱动编排/掉链计数。
qc_passed = await _qc_passed_types(sor, project_id)
dispatchable, waiting_up = [], []
for d, lb in first_run_pending:
up_ok = all(t in qc_passed
for u in DIM_UPSTREAM.get(d, ())
for t in DIM_QC_TYPES[u])
(dispatchable if up_ok else waiting_up).append((d, lb))
if waiting_up:
acts.append("waiting: 维度「%s」等上游 QC 通过(合法等待,不计掉链)"
% "".join(lb for _, lb in waiting_up))
first_run_pending = dispatchable
if first_run_pending:
_ft2 = _ft or await _latest_tender_file_ts(sor, project_id)
orch_done = await _orch_done_count(sor, project_id, since=_ft2 or '')
if orch_done >= MAX_ORCH_ATTEMPT:
ht = await find_human_task(sor, project_id, "general", status="pending")
if not ht:
await create_human_task(
sor, project_id, "general",
"PM 解析编排掉链,需人工介入",
("解析阶段已完成 %d 次编排任务,仍有维度未启动/未过 QC%s\n"
"请人工核对编排任务卡点(任务列表,角色 agent.pm"
"参数 task_kind=analysis_orchestration"
"或人工派发缺失维度。"
% (orch_done, "".join(lb for _, lb in first_run_pending))),
assignee_role="owner.superuser")
_ft2 = _ft or await _latest_tender_file_ts(sor, project_id)
orch_done = await _orch_done_count(sor, project_id, since=_ft2 or '')
if orch_done >= MAX_ORCH_ATTEMPT:
ht = await find_human_task(sor, project_id, "general", status="pending")
if not ht:
await create_human_task(
sor, project_id, "general",
"PM 解析编排掉链,需人工介入",
("解析阶段已完成 %d 次编排任务,仍有维度未启动/未过 QC%s\n"
"请人工核对编排任务卡点(任务列表,角色 agent.pm"
"参数 task_kind=analysis_orchestration"
"或人工派发缺失维度。"
% (orch_done, "".join(lb for _, lb in first_run_pending))),
assignee_role="owner.superuser")
await sor.sqlExe("COMMIT", {})
acts.append("escalated: PM 解析编排掉链 → 人工介入")
else:
dims_brief = "".join("%s(%s)" % (lb, d) for d, lb in first_run_pending)
tid = await create_role_task(
sor, project_id, R_PM,
"%s 解析编排:按任务链派发分析维度" % pname,
{"stage": "analysis_orchestration",
"task_kind": "analysis_orchestration",
"pending_dims": [d for d, _ in first_run_pending]})
acts.append("created: PM 解析编排任务 %s(待派发维度:%s" % (tid, dims_brief))
await sor.sqlExe("COMMIT", {})
acts.append("escalated: PM 解析编排掉链 → 人工介入")
else:
dims_brief = "".join("%s(%s)" % (lb, d) for d, lb in first_run_pending)
tid = await create_role_task(
sor, project_id, R_PM,
"%s 解析编排:按任务链派发分析维度" % pname,
{"stage": "analysis_orchestration",
"task_kind": "analysis_orchestration",
"pending_dims": [d for d, _ in first_run_pending]})
acts.append("created: PM 解析编排任务 %s(待派发维度:%s" % (tid, dims_brief))
await sor.sqlExe("COMMIT", {})
# ── A2. 分析产出 QC 门禁(五个类型逐一审核,每类一个审核任务并行;阈值见 bid_qc_pass_score──
# 2026-09-03 章节级门禁:本段只负责「推进 QC」重做/审核/冒泡),不再提前 return