pipeline-bidding/pipeline_bidding/bid_orchestration_capability.py

144 lines
6.7 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

# -*- 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 "无"))