From 01792915289e787786f6f79c3b1e9355ca47140b Mon Sep 17 00:00:00 2001 From: yumoqing Date: Fri, 4 Sep 2026 13:15:30 +0800 Subject: [PATCH] =?UTF-8?q?feat:=20=E6=96=B9=E6=A1=88=E4=B9=99=E2=80=94?= =?UTF-8?q?=E2=80=94=E8=A7=A3=E6=9E=90=E9=A6=96=E8=B7=91=E6=94=B9PM?= =?UTF-8?q?=E7=BC=96=E6=8E=92(dispatch=5Fanalysis=5Fdim=E4=BE=9D=E8=B5=96?= =?UTF-8?q?=E9=97=A8=E7=A6=81+analysis=5Fprogress+=E4=B8=8A=E6=B8=B8?= =?UTF-8?q?=E6=B3=A8=E5=85=A5)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- pipeline_bidding/bid_flow.py | 81 +++++++++- .../bid_orchestration_capability.py | 143 ++++++++++++++++++ pipeline_bidding/role_tool_schemas.py | 15 ++ 3 files changed, 233 insertions(+), 6 deletions(-) create mode 100644 pipeline_bidding/bid_orchestration_capability.py diff --git a/pipeline_bidding/bid_flow.py b/pipeline_bidding/bid_flow.py index fcf2f66..fa6de46 100644 --- a/pipeline_bidding/bid_flow.py +++ b/pipeline_bidding/bid_flow.py @@ -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)逐章判定—— diff --git a/pipeline_bidding/bid_orchestration_capability.py b/pipeline_bidding/bid_orchestration_capability.py new file mode 100644 index 0000000..1590ee5 --- /dev/null +++ b/pipeline_bidding/bid_orchestration_capability.py @@ -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 "无")) diff --git a/pipeline_bidding/role_tool_schemas.py b/pipeline_bidding/role_tool_schemas.py index 0e57a13..d3ce881 100644 --- a/pipeline_bidding/role_tool_schemas.py +++ b/pipeline_bidding/role_tool_schemas.py @@ -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"], + }, }