feat: 方案乙——解析首跑改PM编排(dispatch_analysis_dim依赖门禁+analysis_progress+上游注入)

This commit is contained in:
yumoqing 2026-09-04 13:15:30 +08:00
parent 32f1a386b4
commit 0179291528
3 changed files with 233 additions and 6 deletions

View File

@ -32,6 +32,7 @@ from .bid_common import (
logger = logging.getLogger("pipeline.bidding")
R_ANALYST = "agent.tender_analyst"
R_PM = "agent.pm"
R_PREP = "agent.bid_prep"
R_WRITER = "agent.bid_writer"
R_BIZ_WRITER = "agent.bid_biz_writer"
@ -41,6 +42,7 @@ R_SCORER = "agent.bid_scorer"
R_QC = "agent.qc"
MAX_ANALYST_ATTEMPT = 2 # 解析阶段首跑最多重派次数(防死循环;QC 退回的重做不计入)
MAX_ORCH_ATTEMPT = 3 # PM 分析编排任务最多重派次数(防掉链死循环)
MAX_PREP_ATTEMPT = 2
POLL_INTERVAL = 15
@ -88,6 +90,16 @@ ANALYSIS_DIMS = (
DIM_QC_TYPES = {d: ts for d, _, ts in ANALYSIS_DIMS}
QC_TYPE_TO_DIM = {q: d for d, _, ts in ANALYSIS_DIMS for q in ts}
# 分析维度依赖链(2026-09-04 方案乙:PM 编排,代码守依赖门禁):
# 评分项 →(资质 ∥ 要求+骨架)→ 成本收益。上游维度的全部 QC 类型通过
# 才准派发下游(依赖是不变量;编排意图归 PM,dispatch_analysis_dim 强制校验)。
DIM_UPSTREAM = {
"scoring": (),
"quals": ("scoring",),
"reqs_outline": ("scoring",),
"cost_benefit": ("scoring", "quals", "reqs_outline"),
}
# 维度 → 落库表(落库核验用;reqs_outline 维度同时写两张表)
DIM_TABLES = {
"scoring": (QC_OUTPUT_TABLE[QC_TYPE_SCORING],),
@ -150,6 +162,31 @@ async def _open_dim_tasks(sor, project_id, dim):
return to_int(getattr(recs[0], "c", 0) if recs else 0)
async def _open_orch_tasks(sor, project_id):
"""PM 分析编排任务在办数(2026-09-04 方案乙)。"""
recs = await sor.sqlExe(
"SELECT COUNT(*) AS c FROM pipeline_tasks WHERE tenant_id=${p}$ "
"AND pipeline_id='role_task' AND role=${r}$ AND state IN "
"('submitted','running','review','qc_review') "
"AND JSON_UNQUOTE(JSON_EXTRACT(params,'$.task_kind'))='analysis_orchestration'",
{"p": project_id, "r": R_PM})
await sor.sqlExe("COMMIT", {})
return to_int(getattr(recs[0], "c", 0) if recs else 0)
async def _orch_done_count(sor, project_id, since=''):
"""PM 分析编排任务已完成数(掉链检测用)。"""
sql = ("SELECT COUNT(*) AS c FROM pipeline_tasks WHERE tenant_id=${p}$ "
"AND pipeline_id='role_task' AND role=${r}$ AND state IN "
"('approved','completed') "
"AND JSON_UNQUOTE(JSON_EXTRACT(params,'$.task_kind'))='analysis_orchestration'")
p = {"p": project_id, "r": R_PM}
if since:
sql += " AND updated_at >= ${since}$"
p["since"] = since
return await _count(sor, sql, p)
async def _dim_done_count(sor, project_id, dim, since=''):
"""某维度已完成(首跑+重做)任务数。"""
sql = ("SELECT COUNT(*) AS c FROM pipeline_tasks WHERE tenant_id=${p}$ "
@ -488,6 +525,7 @@ async def reconcile_project(sor, project_id, project_name=""):
acts.append("recovered: 检测到新招标文件,已清空旧分析残留,重新四维度分析")
missing = await _missing_dims(n_items, n_struct, n_quals, n_cost, chapters)
created_dims = []
first_run_pending = []
for dim, label in missing:
if await _open_dim_tasks(sor, project_id, dim) > 0:
acts.append("waiting: 分析维度「%s」任务在办" % label)
@ -534,16 +572,47 @@ async def reconcile_project(sor, project_id, project_name=""):
await sor.sqlExe("COMMIT", {})
acts.append("escalated: 维度「%s」重试超限 → 人工介入" % label)
continue
tid = await create_role_task(
sor, project_id, R_ANALYST,
"%s 招标文件分析:%s" % (pname, label),
{"stage": "analysis", "task_kind": "bid_analysis", "analysis_dim": dim})
acts.append("created: 分析任务 %s(%s)" % (tid, label))
created_dims.append(dim)
first_run_pending.append((dim, label))
if created_dims:
await sor.sqlExe("COMMIT", {})
# 2026-09-03:不再提前 return——同轮继续章节级派发(依赖已满足的章节立即启动)
# ── A1b. PM 解析编排(2026-09-04 方案乙):首跑不再由对账器并行硬派,
# 而是派一个「解析编排」任务给 PM:PM 按任务链顺序派发
# (评分项 →(资质 ∥ 要求+骨架)→ 成本收益),编排意图归 PM;
# 依赖门禁由代码强制(dispatch_analysis_dim 校验上游 QC 全过)。
# QC 退回重做仍走上面 qc_imp 分支,不受编排影响。
if first_run_pending:
if await _open_orch_tasks(sor, project_id) > 0:
acts.append("waiting: PM 解析编排任务在办")
else:
_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", {})
# ── A2. 分析产出 QC 门禁(五个类型逐一审核,每类一个审核任务并行;阈值见 bid_qc_pass_score)──
# 2026-09-03 章节级门禁:本段只负责「推进 QC」(重做/审核/冒泡),不再提前 return
# 阻塞全局流转。哪些章节可以开写,由 C 段按章节依赖矩阵(CH_SECTION_DEPS)逐章判定——

View File

@ -0,0 +1,143 @@
# -*- coding: utf-8 -*-
"""PM 分析编排能力(方案乙,2026-09-04 用户选定)。
解析首跑任务不再由对账器并行硬派:对账器派「分析编排」任务给 PM,
PM 按任务链编排派发:评分项 →(资质 ∥ 要求+骨架)→ 成本收益。
分工:编排意图(顺序判断/异常招标文档的指令调整/交叉核对要点)归 PM;
依赖门禁(上游维度 QC 全过才准派下游)由代码强制,PM 跳步会被工具拒绝。
QC 退回重做(携带改进意见)仍由对账器自动派发,不归 PM。
"""
import json
import logging
from .bid_common import get_db, resolve_project_id, rec_to_dict, create_role_task, record
logger = logging.getLogger("pipeline.bidding")
MAX_PM_INSTRUCTIONS = 2000
async def _progress(sor, project_id):
"""各维度进度:产出计数 / 最新 QC / 在办任务 / 依赖门禁状态。"""
from .bid_flow import (ANALYSIS_DIMS, DIM_QC_TYPES, DIM_UPSTREAM,
QC_OUTPUT_TABLE, _qc_passed_types, _open_dim_tasks)
passed = await _qc_passed_types(sor, project_id)
rows = []
for dim, label, qc_ts in ANALYSIS_DIMS:
info = {"dim": dim, "label": label}
cnt = 0
for t in qc_ts:
tbl = QC_OUTPUT_TABLE[t]
recs = await sor.sqlExe(
"SELECT COUNT(*) AS c FROM " + tbl + " WHERE project_id=${p}$",
{"p": project_id})
await sor.sqlExe("COMMIT", {})
cnt += int(getattr(recs[0], "c", 0) if recs else 0)
info["产出条数"] = cnt
qc_brief = []
for t in qc_ts:
recs = await sor.sqlExe(
"SELECT round, fit_score, passed FROM bid_qc_reviews "
"WHERE project_id=${p}$ AND qc_type=${t}$ ORDER BY round DESC LIMIT 1",
{"p": project_id, "t": t})
await sor.sqlExe("COMMIT", {})
if recs:
r = rec_to_dict(recs[0])
qc_brief.append("%s 第%s轮 %s分 %s" % (
t, r.get("round"), r.get("fit_score"),
"通过" if str(r.get("passed")) == "1" else "未过"))
else:
qc_brief.append("%s 未审核" % t)
info["QC"] = ";".join(qc_brief)
info["在办任务"] = await _open_dim_tasks(sor, project_id, dim) > 0
up = DIM_UPSTREAM.get(dim, ())
miss = []
for u in up:
for t in DIM_QC_TYPES[u]:
if t not in passed:
miss.append("%s(%s)" % (u, t))
info["依赖门禁"] = "就绪" if not miss else "阻塞:%s 未通过" % "、".join(miss)
rows.append(info)
return rows
async def analysis_progress(project_id="", who=None, agent_id=None):
"""解析编排决策依据:各维度产出计数、最新 QC 结果、在办任务、依赖门禁状态(就绪/阻塞及缺什么)。"""
try:
project_id = await resolve_project_id(project_id)
except ValueError as e:
return "ERROR: %s" % str(e)[:200]
db, dbname = get_db()
async with db.sqlorContext(dbname) as sor:
rows = await _progress(sor, project_id)
return rows
async def dispatch_analysis_dim(dim, instructions="", project_id="",
who=None, agent_id=None):
"""派发一个分析维度的首跑任务(PM 编排专用)。
代码门禁:该维度的全部上游维度 QC 通过才准派发,否则拒绝(依赖是不变量)。
自动把已通过的上游维度(含 QC 轮次)注入任务参数 upstream_dims,
供分析员先读上游产出再抽取(交叉核对/矛盾照录)。
instructions:PM 对该维度的补充指令(异常招标文档的特殊交代/交叉核对要点)。
"""
from .bid_flow import (ANALYSIS_DIMS, DIM_QC_TYPES, DIM_UPSTREAM,
R_ANALYST, _qc_passed_types, _open_dim_tasks)
dims = {d: lb for d, lb, _ts in ANALYSIS_DIMS}
if dim not in dims:
return False, "未知维度 %s(可选:%s)" % (dim, "、".join(dims))
try:
project_id = await resolve_project_id(project_id)
except ValueError as e:
return False, str(e)[:200]
label = dims[dim]
db, dbname = get_db()
async with db.sqlorContext(dbname) as sor:
if await _open_dim_tasks(sor, project_id, dim) > 0:
return False, "维度「%s」已有在办分析任务,不得重复派发" % label
passed = await _qc_passed_types(sor, project_id)
miss = []
for u in DIM_UPSTREAM.get(dim, ()):
for t in DIM_QC_TYPES[u]:
if t not in passed:
miss.append("%s(%s)" % (u, t))
if miss:
return False, ("依赖门禁:维度「%s」的上游未全部通过 QC——%s 未通过。"
"请先推进上游,不得跳步。" % (label, "、".join(miss)))
# 上游参照(已通过的维度 + QC 轮次),供分析员先读上游再抽取
up_refs = []
for u in DIM_UPSTREAM.get(dim, ()):
rnd = 0
for t in DIM_QC_TYPES[u]:
recs = await sor.sqlExe(
"SELECT round FROM bid_qc_reviews WHERE project_id=${p}$ "
"AND qc_type=${t}$ AND passed='1' ORDER BY round DESC LIMIT 1",
{"p": project_id, "t": t})
await sor.sqlExe("COMMIT", {})
if recs:
rnd = max(rnd, int(rec_to_dict(recs[0]).get("round") or 0))
up_refs.append({"dim": u, "label": dims[u], "qc_round": rnd})
pn = await sor.sqlExe(
"SELECT name FROM sd_projects WHERE id=${p}$", {"p": project_id})
await sor.sqlExe("COMMIT", {})
pname = rec_to_dict(pn[0]).get("name", "") if pn else project_id
params = {"stage": "analysis", "task_kind": "bid_analysis",
"analysis_dim": dim, "upstream_dims": up_refs,
"orchestrated_by": "agent.pm"}
if instructions:
params["pm_instructions"] = str(instructions)[:MAX_PM_INSTRUCTIONS]
tid = await create_role_task(
sor, project_id, R_ANALYST,
"%s 招标文件分析:%s(PM编排,先读上游再抽取)" % (pname, label),
params)
await sor.sqlExe("COMMIT", {})
await record(sor, project_id, "pipeline_tasks", tid, "pm_dispatch",
who=who, agent_id=agent_id,
detail="PM 编排派发分析维度 %s(上游 %s)"
% (dim, ",".join(u["dim"] for u in up_refs) or "无"))
await sor.sqlExe("COMMIT", {})
return True, ("已派发维度「%s」分析任务 %s(上游参照 %d 个:%s)"
% (label, tid, len(up_refs),
"、".join(u["label"] for u in up_refs) or "无"))

View File

@ -334,6 +334,21 @@ BID_ROLE_TOOL_SCHEMAS = {
"params": {"project_id": "项目ID(可选,默认当前)"},
"required": [],
},
# ── PM 分析编排(2026-09-04 方案乙)──
"analysis_progress": {
"module": f"{_M}.bid_orchestration_capability",
"description": "解析编排决策依据:各维度产出计数/最新QC/在办任务/依赖门禁状态(就绪或阻塞缺什么)",
"params": {"project_id": "项目ID(可选,默认当前)"},
"required": [],
},
"dispatch_analysis_dim": {
"module": f"{_M}.bid_orchestration_capability",
"description": "派发一个分析维度首跑任务。代码门禁:上游维度QC未全过会拒绝(不得跳步)。已过的上游自动注入任务参数供分析员先读上游再抽取。可选维度:scoring 评分项 / quals 资质 / reqs_outline 要求与骨架 / cost_benefit 成本收益",
"params": {"dim": "维度名 scoring|quals|reqs_outline|cost_benefit",
"instructions": "对该维度的补充指令:异常招标文档的特殊交代/交叉核对要点/上游产出怎么用(可选)",
"project_id": "项目ID(可选,默认当前)"},
"required": ["dim"],
},
}