1812 lines
99 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.

"""投标产线任务流转引擎:状态收敛器(reconciler)+ 后台 poller。
设计原则(与开发产线一致):**流转归机制,不归 LLM**。
不依赖 agent 主动"调下一步工具",而是每 15 秒对每个投标项目做一次幂等状态对账:
按「领域状态」(招标文件/评分项/章节/标书)计算应存在哪些角色任务,缺则补建,已有则跳过。
因此:
- 章节评审不达标 → 章节置 rejected → 下一轮对账自动派重写任务(章节级循环)
- 整书评分不达标 → 命中章节置 rejected → 自动重写 → 重新评审 → 重新合成 → 重新评分(整书级循环)
- 人类任务 pending(等招标文件 / 等人类资料)→ 对账直接跳过该项目(阻塞门禁)
RoleSpec.next_role 全部留空:投标产线不用「任务链跳下一角色」,一律由本对账器决定。
"""
import asyncio
import json
import logging
import os
from .bid_common import (
get_db, rows_to_dicts, rec_to_dict, to_int, to_float, json_loads, record,
create_role_task, count_open_tasks, has_done_task,
has_pending_blocking_human_task, find_human_task, get_thresholds,
create_human_task, get_project, PIPELINE_ID,
CH_PENDING, CH_WRITING, CH_WRITTEN, CH_REJECTED, CH_APPROVED,
DOC_DRAFT, DOC_REVIEWING, DOC_PASSED, DOC_REJECTED, DOC_DELIVERED,
HT_TENDER_FILE, HT_HUMAN_DOCS, HT_QC_ESCALATION, HT_DELIVERY_CONFIRM,
HT_TEMPLATE_CONFIRM,
QC_TYPE_SCORING, QC_TYPE_QUALS, QC_TYPE_REQS, QC_TYPE_OUTLINE, QC_TYPE_COST,
QC_TYPE_TECH, QC_TYPES, TECH_QC_TYPES,
QC_OUTPUT_TABLE, DIM_TO_STAGE,
STAGE_ANALYSIS, STAGE_QUALS, STAGE_REQS_OUTLINE, STAGE_COST_BENEFIT,
STAGE_QC, STAGE_PREP, STAGE_CH_REVIEW, STAGE_WHOLE_SCORE, STAGE_DELIVERY,
FLOW_KEY_TECH_PROPOSAL, STAGE_TECH_ANALYSIS, STAGE_TECH_QC,
STAGE_TEMPLATE_CONFIRM, STAGE_TECH_CH_WRITE, STAGE_TECH_CH_QC,
STAGE_TECH_COMPOSE, STAGE_TECH_DELIVERY,
)
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"
R_REVIEWER = "agent.bid_reviewer"
R_COMPOSITOR = "agent.bid_compositor"
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
# poller 心跳(供 /pipeline-bidding/api/bid_probe.dspy 诊断;服务内可见)
POLLER_STATE = {"last_run": "", "last_error": "", "rounds": 0}
async def _count(sor, sql, p):
recs = await sor.sqlExe(sql, p)
await sor.sqlExe("COMMIT", {})
return to_int(getattr(recs[0], "c", 0) if recs else 0)
async def _done_task_count(sor, project_id, role, since=''):
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')")
p = {"p": project_id, "r": role}
if since:
sql += " AND updated_at >= ${since}$"
p["since"] = since
return await _count(sor, sql, p)
async def _redo_done_count(sor, project_id, since=''):
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 params LIKE ${kw}$")
p = {"p": project_id, "r": R_ANALYST, "kw": "%qc_redo%"}
if since:
sql += " AND updated_at >= ${since}$"
p["since"] = since
return await _count(sor, sql, p)
# ── 四维度分析辅助(2026-09-02:解析拆成四个并行维度)──
# 维度定义:dim key → (任务标签, 对应 QC 类型)
ANALYSIS_DIMS = (
("scoring", "评分项与得分规则", (QC_TYPE_SCORING,)),
("quals", "资质清单与要求", (QC_TYPE_QUALS,)),
("reqs_outline", "投标文件要求与章节骨架", (QC_TYPE_REQS, QC_TYPE_OUTLINE)),
("cost_benefit", "成本收益分析", (QC_TYPE_COST,)),
)
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",),
# 2026-09-15 用户裁定:资质有无/资质QC不阻塞标书推进。cost_benefit 测算
# (保证金/预算/付款条件)实际依据评分项与投标文件要求,不依赖资质清单;
# 且 BID_FLOW_STAGES 声明层 analysis_cost_benefit.deps 本就无 analysis_quals
# ——机制层与声明层对齐(原不一致:quals QC 卡住连报价章全链冻结)。
"cost_benefit": ("scoring", "reqs_outline"),
}
# 维度 → 落库表(落库核验用;reqs_outline 维度同时写两张表)
DIM_TABLES = {
"scoring": (QC_OUTPUT_TABLE[QC_TYPE_SCORING],),
"quals": (QC_OUTPUT_TABLE[QC_TYPE_QUALS],),
"reqs_outline": (QC_OUTPUT_TABLE[QC_TYPE_REQS], QC_OUTPUT_TABLE[QC_TYPE_OUTLINE]),
"cost_benefit": (QC_OUTPUT_TABLE[QC_TYPE_COST],),
}
# 章节级上游依赖矩阵(2026-09-03):每种章节 section 只依赖自己必需的分析产出
# 通过 QC——某一维度 QC 卡住只阻塞依赖它的章节,其余章节照常流转。
# 公共事实:章节骨架(bid_chapters)本身就是章节,OUTLINE 过了章节才可信;
# REQS(投标文件格式要求)决定章节怎么写,所有章节都要。
# 2026-09-15 用户裁定:资质有无/QC过否不阻塞标书编写——business 移除 QUALS 依赖
# (资质缺口走 human_docs_request 待办通知人类,属信息流不是门禁流)。
CH_SECTION_DEPS = {
"technical": (QC_TYPE_SCORING, QC_TYPE_REQS, QC_TYPE_OUTLINE),
"price": (QC_TYPE_SCORING, QC_TYPE_COST, QC_TYPE_OUTLINE),
"business": (QC_TYPE_REQS, QC_TYPE_OUTLINE),
}
async def _qc_passed_types(sor, project_id):
"""最新一轮审核已通过(或人工强制放行)的 QC 类型集合。章节级依赖判定用。
逃逸阀闭环:QC 轮次用尽冒泡人工后,未通过类型不在集合里 → 依赖它的章节被挡;
人工在记录页把 passed 改为 1 强制放行后,该类型进集合 → 依赖章节自动解锁,无需改码。
"""
passed = set()
for t in QC_TYPES:
recs = await sor.sqlExe(
"SELECT 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 and str(rec_to_dict(recs[0]).get("passed")) == "1":
passed.add(t)
return passed
async def _missing_dims(n_items, n_struct, n_quals, n_cost, chapters):
"""返回缺失的分析维度 [(dim, label)]。各维度独立判定,缺几个并行派几个。"""
missing = []
if n_items == 0:
missing.append(("scoring", "评分项与得分规则"))
if n_quals == 0:
missing.append(("quals", "资质清单与要求"))
if n_struct == 0 or not chapters:
missing.append(("reqs_outline", "投标文件要求与章节骨架"))
if n_cost == 0:
missing.append(("cost_benefit", "成本收益分析"))
return missing
async def _open_dim_tasks(sor, project_id, dim):
"""某维度在办分析任务数(按 params.analysis_dim 匹配)。"""
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,'$.analysis_dim'))=${d}$",
{"p": project_id, "r": R_ANALYST, "d": dim})
await sor.sqlExe("COMMIT", {})
return to_int(getattr(recs[0], "c", 0) if recs else 0)
async def _tender_text_len(sor, project_id):
"""项目最新招标文件正文字符数(无文件/无正文返回 0)。拆解阈值判定用。"""
recs = await sor.sqlExe(
"SELECT LENGTH(content_text) AS n FROM bid_tender_files "
"WHERE project_id=${p}$ ORDER BY created_at DESC LIMIT 1", {"p": project_id})
await sor.sqlExe("COMMIT", {})
if not recs:
return 0
return to_int(getattr(recs[0], "n", 0) or 0, 0)
async def _dim_round_count(sor, project_id, dim, since='', redo=None):
"""某维度已消耗的「轮次」数(轮 = 去重后的 params.round 个数,非任务条数)。
2026-09-15 拆解机制配套:一个拆解轮会派 N 个分段子任务(同 round),
按任务条数计轮会把一轮误算成 N 轮 → 轮次上限误判。首跑与重做统一口径:
任务 params 带 round(拆解轮)则按 round 去重,不带 round 的老任务每条算一轮。
redo=None 全部 / True 只重做(qc_redo) / False 只首跑(无 qc_redo)。
"""
sql = ("SELECT params FROM pipeline_tasks WHERE tenant_id=${p}$ "
"AND pipeline_id='role_task' AND role=${r}$ "
"AND state IN ('approved','completed','qc_rejected') "
"AND JSON_UNQUOTE(JSON_EXTRACT(params,'$.analysis_dim'))=${d}$")
p = {"p": project_id, "r": R_ANALYST, "d": dim}
if redo is True:
sql += " AND params LIKE ${kw}$"
p["kw"] = "%qc_redo%"
elif redo is False:
sql += " AND params NOT LIKE ${kw}$"
p["kw"] = "%qc_redo%"
if since:
sql += " AND updated_at >= ${since}$"
p["since"] = since
recs = await sor.sqlExe(sql, p)
await sor.sqlExe("COMMIT", {})
rounds = set()
n_plain = 0
for r in recs or []:
try:
pp = json.loads(getattr(r, "params", "") or "{}")
except (json.JSONDecodeError, TypeError):
pp = {}
rd = pp.get("round") if isinstance(pp, dict) else None
if rd:
rounds.add(str(rd))
else:
n_plain += 1
return len(rounds) + n_plain
async def _create_split_tasks(sor, project_id, dim, label, pname, base_params,
text_len, split_chars, round_tag, qc_imp=None):
"""大文件分析维度机械拆解:按 split_chars 把正文切成 N 段,每段一个子任务。
子任务 params 带 char_range=[start,end)/seg=i/seg_n=N/round=round_tag/is_subtask=1,
分析员技能硬规则:有 char_range 只读本区间(抽取工具 offset/length 原生支持)。
返回创建的子任务 id 列表。N 上限 6(防极端文件拆爆并发)。
"""
n = min(6, max(2, -(-text_len // split_chars))) # 向上取整,[2,6]
step = -(-text_len // n)
ids = []
for i in range(n):
start = i * step
end = min(text_len, start + step)
if start >= text_len:
break
params = dict(base_params)
params.update({"char_range": [start, end], "seg": i + 1, "seg_n": n,
"round": round_tag, "is_subtask": 1})
if qc_imp:
params["qc_improvements"] = qc_imp
tid = await create_role_task(
sor, project_id, R_ANALYST,
"%s 分析%s:%s 第%d/%d段(字符 %d-%d)"
% (pname, "重做" if qc_imp else "", label, i + 1, n, start, end),
params)
ids.append(tid)
return ids
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}$ "
"AND pipeline_id='role_task' AND role=${r}$ AND state IN ('approved','completed') "
"AND JSON_UNQUOTE(JSON_EXTRACT(params,'$.analysis_dim'))=${d}$")
p = {"p": project_id, "r": R_ANALYST, "d": dim}
if since:
sql += " AND updated_at >= ${since}$"
p["since"] = since
return await _count(sor, sql, p)
async def _dim_redo_done_count(sor, project_id, dim, since=''):
"""某维度重做(qc_redo=1)已终结任务数(approved/completed/qc_rejected,不计入首跑重试上限)。
qc_rejected(QC 判定不通过)计入:防止对账器对同一轮反复派重做任务。
"""
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','qc_rejected') "
"AND JSON_UNQUOTE(JSON_EXTRACT(params,'$.analysis_dim'))=${d}$ AND params LIKE ${kw}$")
p = {"p": project_id, "r": R_ANALYST, "d": dim, "kw": "%qc_redo%"}
if since:
sql += " AND updated_at >= ${since}$"
p["since"] = since
return await _count(sor, sql, p)
async def _dim_qc_improvements(sor, project_id, dim):
"""某维度各 QC 类型最新一轮未通过的改进意见(全通过/未审核返回空列表)。"""
imps = []
for t in DIM_QC_TYPES.get(dim, ()):
recs = await sor.sqlExe(
"SELECT round, passed, improvement 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 not recs:
continue
r = rec_to_dict(recs[0])
if str(r.get("passed")) == "1" or not (r.get("improvement") or "").strip():
continue
imps.append("[%s 第%s轮] %s" % (t, r.get("round"), r["improvement"]))
return imps
async def _open_qc_type_tasks(sor, project_id, qc_type):
"""某类型在办 QC 审核任务数(params.qc_types 含该类型)。"""
kw = '%"qc_types": ["' + qc_type + '"]%' # 拼接构造:SQL 里 % 是 LIKE 通配符,不能用 % 格式化
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 params LIKE ${kw}$",
{"p": project_id, "r": R_QC, "kw": kw})
await sor.sqlExe("COMMIT", {})
return to_int(getattr(recs[0], "c", 0) if recs else 0)
async def _latest_tender_file_ts(sor, project_id):
"""项目最新招标文件的 created_at(无文件返回 '')。"""
recs = await sor.sqlExe(
"SELECT created_at FROM bid_tender_files WHERE project_id=${p}$ "
"ORDER BY created_at DESC LIMIT 1", {"p": project_id})
await sor.sqlExe("COMMIT", {})
if not recs:
return ''
return str(getattr(recs[0], "created_at", "") or "")
async def _latest_done_analysis_ts(sor, project_id):
"""项目最后一次已完成解析任务的 updated_at(无则返回 '')。"""
recs = await sor.sqlExe(
"SELECT updated_at FROM pipeline_tasks WHERE tenant_id=${p}$ "
"AND pipeline_id='role_task' AND role=${r}$ AND state IN "
"('approved','completed') ORDER BY updated_at DESC LIMIT 1",
{"p": project_id, "r": R_ANALYST})
await sor.sqlExe("COMMIT", {})
if not recs:
return ''
return str(getattr(recs[0], "updated_at", "") or "")
async def _reset_analysis_for_new_file(sor, project_id, since_ts):
"""新招标文件到来 → 清空旧文件时代的解析残留,让新文件干净重跑(幂等)。
清空范围:五类产出物(评分项/资质/投标文件要求/章节骨架/成本收益,旧文件抽的
对新文件无意义)+ 早于新文件的 QC 审核记录(旧改进意见不能带入新文件解析)。
顺手关闭旧「解析未产出完整结果」人工介入任务(自动恢复即处理完毕)。
"""
# 仅当存在早于新文件的 QC 记录时执行(说明此前有过旧解析,属恢复场景)
recs = await sor.sqlExe(
"SELECT id FROM bid_qc_reviews WHERE project_id=${p}$ "
"AND created_at < ${ts}$ LIMIT 1", {"p": project_id, "ts": since_ts})
await sor.sqlExe("COMMIT", {})
if not recs:
return False
for tbl in ("bid_scoring_items", "bid_qualifications", "bid_doc_requirements",
"bid_chapters", "bid_cost_benefit"):
await sor.sqlExe(
"DELETE FROM %s WHERE project_id=${p}$" % tbl, {"p": project_id})
await sor.sqlExe("COMMIT", {})
await sor.sqlExe(
"DELETE FROM bid_qc_reviews WHERE project_id=${p}$ AND created_at < ${ts}$",
{"p": project_id, "ts": since_ts})
await sor.sqlExe("COMMIT", {})
ht = await find_human_task(sor, project_id, "general", status="pending")
if ht and "解析" in (ht.get("title") or ""):
await sor.sqlExe(
"UPDATE pipeline_human_tasks SET status='done', qc_status='passed', "
"qc_comment='已上传新招标文件,系统自动重新解析(旧文件残留产出已清空)', "
"submitted_at=NOW() WHERE id=${h}$", {"h": ht.get("id")})
await sor.sqlExe("COMMIT", {})
logger.info("bid_flow: reset analysis outputs for new tender file project=%s",
project_id)
return True
async def _auto_close_file_upload_task(sor, project_id):
"""招标文件已上传 → 自动完成「等待上传招标文件」人类任务(上传即完成,不必人工再点一次)。"""
ht = await find_human_task(sor, project_id, HT_TENDER_FILE, status="pending")
if not ht:
return False
n = await _count(sor, "SELECT COUNT(*) AS c FROM bid_tender_files WHERE project_id=${p}$",
{"p": project_id})
if n <= 0:
return False
await sor.sqlExe(
"UPDATE pipeline_human_tasks SET status='done', qc_status='passed', "
"qc_comment='招标文件已上传,系统自动完成', submitted_at=NOW() WHERE id=${h}$",
{"h": ht.get("id")})
await sor.sqlExe("COMMIT", {})
logger.info("bid_flow: auto-closed tender_file_upload task project=%s", project_id)
return True
# ══════════════ QC 门禁辅助 ══════════════
async def _qc_pending_types(sor, project_id):
"""返回未放行的产出类型清单(按依赖序),每项 (qc_type, state, improvement, round, out_newer):
state=unaudited:产出已存在但从未审核;state=rework:最近一轮审核未通过。
已通过(或人工在记录页强制放行)的类型不再出现。
out_newer(2026-09-03):产出表最新更新时间晚于最近一轮审核时间——
说明审核后产出被人工/修正任务更新过,即使轮次用尽也应再派一轮复审
(闭环「修正落库→自动复审」;只有产出真更新才开轮,不会无限烧轮)。
"""
pending = []
for t in QC_TYPES:
tbl = QC_OUTPUT_TABLE[t]
recs = await sor.sqlExe(
"SELECT round, passed, improvement, updated_at 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 not recs:
n_out = await _count(sor, "SELECT COUNT(*) AS c FROM " + tbl +
" WHERE project_id=${p}$", {"p": project_id})
if n_out > 0:
pending.append((t, "unaudited", "", 0, False))
continue
r = rec_to_dict(recs[0])
if str(r.get("passed")) == "1":
continue
rnd = to_int(r.get("round"), 1)
rev_t = str(r.get("updated_at") or r.get("created_at") or "")
out_rec = await sor.sqlExe(
"SELECT MAX(GREATEST(created_at, IFNULL(updated_at, created_at))) AS mt "
"FROM " + tbl + " WHERE project_id=${p}$", {"p": project_id})
await sor.sqlExe("COMMIT", {})
out_t = str(getattr(out_rec[0], 'mt', '') or '') if out_rec else ""
out_newer = bool(out_t and rev_t and out_t > rev_t)
n_out = await _count(sor, "SELECT COUNT(*) AS c FROM " + tbl +
" WHERE project_id=${p}$", {"p": project_id})
if n_out > 0:
pending.append((t, "rework", r.get("improvement") or "", rnd, out_newer))
else:
# 产出已被清空待重做:round 信息带回,供对账器判断是否已超轮次上限
pending.append((t, "awaiting_redo", r.get("improvement") or "", rnd, False))
return pending
async def _qc_collect_improvements(sor, project_id):
"""汇总各类产出最新一轮的改进意见(解析重做任务的 params.qc_improvements)。"""
items = []
for t in QC_TYPES:
recs = await sor.sqlExe(
"SELECT round, passed, improvement 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 not recs:
continue
r = rec_to_dict(recs[0])
if str(r.get("passed")) == "1" or not (r.get("improvement") or "").strip():
continue
items.append("[%s] 第%s轮:%s" % (t, r.get("round"), r["improvement"]))
return items
async def _auto_import_workspace_tender(sor, project_id):
"""工作空间已有招标书文件但 bid_tender_files 无记录 → 自动注册并启动解析。
根治(2026-09-01):用户把招标书直接放进项目工作空间(或 agent 落盘)时,
旧逻辑只认 bid_tender_files 表记录,对账器每轮都发「上传招标文件」人类待办,
用户怎么答「文件就在目录里」都没用——死循环。现在对账器主动扫描工作空间,
发现招标书即 import_tender_file 注册(落库+自动关闭上传待办+收敛器派解析),
与上传入口完全等价,不再催人。
"""
try:
from pipeline_service.workspace import get_project_dir_by_id
pdir, _base = await get_project_dir_by_id(sor, project_id)
if not pdir or not os.path.isdir(pdir):
return ""
cands = []
for fn in sorted(os.listdir(pdir)):
ext = os.path.splitext(fn)[1].lower()
if ext not in (".docx", ".doc", ".pdf", ".txt", ".md"):
continue
fp = os.path.join(pdir, fn)
if os.path.isfile(fp) and os.path.getsize(fp) > 0:
cands.append((os.path.getsize(fp), fp, fn))
if not cands:
return ""
cands.sort(reverse=True) # 最大文件优先(招标书通常是项目里最大的文档)
_sz, fp, fn = cands[0]
from .bid_tender_capability import import_tender_file
ok, msg = await import_tender_file(
project_id, file_path=fp, file_name=fn,
who="system.reconciler", agent_id="bid_flow")
if ok:
logger.info("bid_flow: auto-imported workspace tender %s project=%s",
fn, project_id)
return "auto_imported: %s(%s)" % (fn, msg)
logger.warning("bid_flow: auto-import tender failed project=%s: %s",
project_id, str(msg)[:160])
except Exception as e:
logger.warning("bid_flow: auto-import workspace tender error project=%s: %s",
project_id, str(e)[:160])
return ""
async def _flow_stage_table(project_id):
"""取项目裁剪计划的阶段启用表(流程裁剪机制,pipeline_flow_plans)。
返回 None = 无已确认计划(标准全流程,存量项目零影响);
dict = {stage_key: bool},False 即该阶段被裁(用户已在待办确认)。
引擎不可用时诚实降级 None(不阻塞流转),并记 warning。
"""
try:
from pipeline_service.flow_plan_capability import plan_stage_enabled
return await plan_stage_enabled(project_id)
except Exception as e:
logger.warning("bid_flow: flow plan unavailable project=%s: %s",
project_id, str(e)[:160])
return None
async def _latest_flow_plan_status(project_id):
"""项目最新流程计划状态('' = 从未提案)。驳回后暂停派发用。"""
try:
from pipeline_service.flow_plan_capability import get_latest_plan
plan = await get_latest_plan(project_id)
return (plan or {}).get("status", "") or "", (plan or {}).get("reject_comment", "") or ""
except Exception:
return "", ""
async def _plan_flow_key(project_id):
"""已确认计划的 base_flow_key('' = 无计划/单流程标准链路)。"""
try:
from pipeline_service.flow_plan_capability import get_plan_flow_key
return await get_plan_flow_key(project_id)
except Exception:
return ""
async def _auto_detect_flow(sor, project_id, pname):
"""多流程产线:首个文件导入后 LLM 自动判流 → propose 草案 → 确认待办(阻塞等用户)。
触发条件:产线声明 ≥2 流程模板 且 项目从未提案(无 pipeline_flow_plans 记录)。
判流失败不默认硬跑(语义判断禁静默降级):发人工介入待办说明原因。
单流程产线(模板 <2)本函数直接返回,行为与存量完全一致。
"""
try:
from pipeline_service.flow_plan_capability import (
detect_flow_template, propose_flow_plan, get_latest_plan)
except Exception as e:
logger.warning("bid_flow: flow detect unavailable project=%s: %s",
project_id, str(e)[:160])
return ""
proj = await get_project(sor, project_id)
pid = proj.get("pipeline_id") or PIPELINE_ID
if await get_latest_plan(project_id):
return "" # 已提案过(任何状态),不重复判流
# 读最新文件正文供判流
recs = await sor.sqlExe(
"SELECT file_name, LEFT(content_text, 6000) AS head FROM bid_tender_files "
"WHERE project_id=${p}$ ORDER BY created_at DESC LIMIT 1", {"p": project_id})
await sor.sqlExe("COMMIT", {})
if not recs:
return ""
fname = getattr(recs[0], "file_name", "") or ""
fhead = getattr(recs[0], "head", "") or ""
fk, note = await detect_flow_template(
pid, file_name=fname, file_text=fhead,
project_id=project_id, org_id=proj.get("org_id") or "")
if not fk:
# 判流失败:发人工介入待办(不默认选流程硬跑)
ht = await find_human_task(sor, project_id, "general", status="pending")
if not ht and note:
await create_human_task(
sor, project_id, "general",
"流程自动识别失败,需人工指定流程:%s" % (pname or project_id),
("系统无法自动判断该文件应走哪条流程:%s\n\n"
"请在会话里用 propose_flow_plan(flow_key=...) 指定流程,"
"或联系管理员核对文件内容。" % note),
assignee_role="owner.superuser")
await sor.sqlExe("COMMIT", {})
return ""
ok, msg = await propose_flow_plan(
project_id, trim_keys=[], user_requirements="(系统自动判流,见识别结论)",
propose_note="", pipeline_id=pid, flow_key=fk, detect_note=note,
who="system.reconciler", agent_id="bid_flow")
if ok:
logger.info("bid_flow: auto-detected flow %s for project=%s", fk, project_id)
return "created: 自动判流 %s → 流程确认待办(等用户确认)" % fk
logger.warning("bid_flow: auto propose failed project=%s: %s", project_id, str(msg)[:160])
return ""
async def reconcile_project(sor, project_id, project_name=""):
"""对单个投标项目做一次状态对账。返回动作说明列表(幂等,可反复调用)。"""
acts = []
pname = project_name or project_id
# ── 0a. 失败任务优先处理:角色 agent 执行失败(如机构未配置 LLM)→
# 抛人工介入,且不再造新任务(防每轮对账重复造)。
n_failed = await _count(
sor, "SELECT COUNT(*) AS c FROM pipeline_tasks WHERE tenant_id=${p}$ "
"AND pipeline_id='role_task' AND state='failed'", {"p": project_id})
if n_failed > 0:
ht = await find_human_task(sor, project_id, "general", status="pending")
if not ht:
tasks = await sor.sqlExe(
"SELECT role, title FROM pipeline_tasks WHERE tenant_id=${p}$ "
"AND pipeline_id='role_task' AND state='failed' LIMIT 5", {"p": project_id})
await sor.sqlExe("COMMIT", {})
tlist = "\n".join("- %s: %s" % (getattr(r, "role", ""), getattr(r, "title", ""))
for r in (tasks or []))
await create_human_task(
sor, project_id, "general",
"投标任务执行失败,需人工处理(%d 个)" % n_failed,
("以下角色任务执行失败,对账器已暂停自动派发:\n" + tlist +
"\n\n常见原因:机构未配置 LLM 模型(项目缺省模型/角色模型未设置)、"
"模型调用超时、任务内容异常。\n请处理后用 reset_task_retry 恢复,或取消任务。"),
assignee_role="owner.superuser")
await sor.sqlExe("COMMIT", {})
return ["blocked: %d 个任务执行失败,已抛人工介入,暂停派发" % n_failed]
# ── 0b. 招标文件到位则自动关闭上传任务,再判阻塞门禁 ──
await _auto_close_file_upload_task(sor, project_id)
# 产线起点兜底:项目无招标文件记录时——
# ① 先扫工作空间:招标书已在项目目录里就直接自动注册并启动解析(根治
# 「文件在目录里还反复催上传」死循环,2026-09-01),不发待办;
# ② 真没有且无 pending 上传待办才发一次(同一事项 pending 期间不重发)。
_n_files_pre = await _count(sor, "SELECT COUNT(*) AS c FROM bid_tender_files "
"WHERE project_id=${p}$", {"p": project_id})
if _n_files_pre <= 0:
_imp = await _auto_import_workspace_tender(sor, project_id)
if _imp:
acts.append(_imp)
else:
_up_ht = await find_human_task(sor, project_id, HT_TENDER_FILE, status="pending")
if not _up_ht:
p = await get_project(sor, project_id)
hid = await create_human_task(
sor, project_id, HT_TENDER_FILE,
"上传招标文件:%s" % (p.get("name", "") or project_id),
("请上传招标文件(含答疑/补遗文件)。\n"
"上传后系统自动进入招标文件分析(四个维度并行:评分项与得分规则、资质清单、"
"投标文件要求与章节骨架、成本收益分析)。"))
await sor.sqlExe("COMMIT", {})
return acts + ["created: 上传招标文件任务 %s" % hid]
if await has_pending_blocking_human_task(sor, project_id):
return ["blocked: 有待办人类任务(等招标文件/等流程裁剪确认),本轮不派发"]
# ── 0c. 流程裁剪门禁(2026-09-11)──
# 最新计划被驳回 → 暂停派发,等会话 agent 按驳回意见修订重提(用户已表达要改流程,
# 不得按旧流程抢跑);pending_confirm 由 0b 的阻塞待办挡住;confirmed 计划在下方
# 各段按 stage_table 生效;从未提案('')→ 标准全流程(存量项目零影响)。
stage_table = await _flow_stage_table(project_id)
_fp_status, _fp_comment = await _latest_flow_plan_status(project_id)
if _fp_status == "rejected" and stage_table is None:
return acts + ["blocked: 流程裁剪方案被用户驳回(意见:%s),等会话 agent 修订重提,本轮不派发"
% (_fp_comment[:200] or "未填")]
# ── 0d. 多流程分支(2026-09-14 批2)──
# 已确认计划带 base_flow_key=bid_tech_proposal → 走纯技术方案流程(独立对账分支,
# 以需求书为唯一基准,不碰评分项/资质/成本收益)。
# 多流程产线 + 从未提案 + 已有文件 → LLM 自动判流出草案(确认待办阻塞,等用户拍板)。
_fk = await _plan_flow_key(project_id)
if _fk == FLOW_KEY_TECH_PROPOSAL:
return acts + await _reconcile_tech_proposal(
sor, project_id, pname, stage_table or {})
if not _fk and not _fp_status:
_n_files_now = await _count(sor, "SELECT COUNT(*) AS c FROM bid_tender_files "
"WHERE project_id=${p}$", {"p": project_id})
if _n_files_now > 0:
_det = await _auto_detect_flow(sor, project_id, pname)
if _det:
return acts + [_det]
# 判流失败/待人工指定:人工待办已在 _auto_detect_flow 里发(general 非阻塞),
# 为避免标准流程抢跑,检测到该待办即本轮不派发;无待办 = 单流程产线(现状不变)。
_dj = await find_human_task(sor, project_id, "general", status="pending")
if _dj and "流程自动识别失败" in (_dj.get("title") or ""):
return acts + ["blocked: 流程自动识别失败,等人工指定流程(会话里 propose_flow_plan)"]
def _stage_on(key):
"""阶段是否执行:无确认计划(None)= 标准全流程都执行。"""
if stage_table is None:
return True
return bool(stage_table.get(key, True))
async def _qc_passed_eff():
"""章节依赖判定用的 QC 通过集(裁剪语义,2026-09-11):
- qc_analysis 被裁 → 全部类型视为通过(用户已确认放弃解析QC门禁);
- 某分析维度被裁 → 该维度对应 QC 类型视为通过(不会再有产出与审核);
- 其余按 bid_qc_reviews 最新一轮实际判定。
"""
if not _stage_on(STAGE_QC):
return set(QC_TYPES)
passed = await _qc_passed_types(sor, project_id)
for dim, stage in DIM_TO_STAGE.items():
if not _stage_on(stage):
for t in DIM_QC_TYPES.get(dim, ()):
passed.add(t)
return passed
th = await get_thresholds(sor)
n_files = await _count(sor, "SELECT COUNT(*) AS c FROM bid_tender_files "
"WHERE project_id=${p}$", {"p": project_id})
n_items = await _count(sor, "SELECT COUNT(*) AS c FROM bid_scoring_items "
"WHERE project_id=${p}$", {"p": project_id})
n_struct = await _count(sor, "SELECT COUNT(*) AS c FROM bid_doc_requirements "
"WHERE project_id=${p}$ AND req_type='structure'",
{"p": project_id})
n_quals = await _count(sor, "SELECT COUNT(*) AS c FROM bid_qualifications "
"WHERE project_id=${p}$", {"p": project_id})
n_quals_pending = await _count(sor, "SELECT COUNT(*) AS c FROM bid_qualifications "
"WHERE project_id=${p}$ AND match_status='pending'",
{"p": project_id})
n_cost = await _count(sor, "SELECT COUNT(*) AS c FROM bid_cost_benefit "
"WHERE project_id=${p}$ AND kind='item'", {"p": project_id})
chs = await sor.sqlExe(
"SELECT id, chapter_no, title, section, status, revise_count FROM bid_chapters "
"WHERE project_id=${p}$ ORDER BY order_no, chapter_no", {"p": project_id})
await sor.sqlExe("COMMIT", {})
chapters = rows_to_dicts(chs, limit=500)
docs = await sor.sqlExe(
"SELECT id, version, status FROM bid_documents WHERE project_id=${p}$ "
"ORDER BY version DESC LIMIT 1", {"p": project_id})
await sor.sqlExe("COMMIT", {})
doc = rec_to_dict(docs[0]) if docs else {}
# ── A. 招标文件分析(2026-09-02 起拆成四个并行维度,各派一个任务)──
# 1) 评分项与得分规则 2) 资质清单与要求 3) 投标文件要求与章节骨架 4) 成本收益分析
# 各维度独立判定缺失 → 独立派发,引擎每项目并发上限(默认 4)保证真并行。
# 新招标文件自动恢复(清旧产出重跑)在派发前统一执行。
if n_files > 0:
_ft = await _latest_tender_file_ts(sor, project_id)
_at = await _latest_done_analysis_ts(sor, project_id)
if _ft and _at and _ft > _at:
if await _reset_analysis_for_new_file(sor, project_id, _ft):
acts.append("recovered: 检测到新招标文件,已清空旧分析残留,重新四维度分析")
missing = await _missing_dims(n_items, n_struct, n_quals, n_cost, chapters)
# 流程裁剪:被裁分析维度不派(用户已在确认待办里确认)
_trimmed_dims = [d for d, _ in missing if not _stage_on(DIM_TO_STAGE.get(d, ""))]
if _trimmed_dims:
missing = [(d, lb) for d, lb in missing
if _stage_on(DIM_TO_STAGE.get(d, ""))]
acts.append("trimmed: 分析维度「%s」已被用户裁剪,跳过" % "、".join(_trimmed_dims))
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)
continue
_ft2 = _ft or await _latest_tender_file_ts(sor, project_id)
attempts = await _dim_round_count(sor, project_id, dim, since=_ft2, redo=False)
n_redo_done = await _dim_round_count(sor, project_id, dim, since=_ft2, redo=True)
# QC 退回重做(携带改进意见)优先于首跑判定,不计入首跑重试上限
qc_imp = await _dim_qc_improvements(sor, project_id, dim)
if qc_imp:
if n_redo_done >= th["qc_max_round"]:
ht = await find_human_task(sor, project_id, "general", status="pending")
if not ht:
await create_human_task(
sor, project_id, "general",
"分析维度「%s」重做轮次已用尽,需人工介入" % label,
("QC 退回后该维度已重做 %d 轮,产出仍未通过契合度审核:\n%s\n\n"
"请人工核对招标文件并修正产出物(对应管理页),"
"或在 QC 审核记录页把对应记录 passed 改为 1 强制放行。"
% (n_redo_done, "\n".join(qc_imp)[:3000])),
assignee_role="owner.superuser")
await sor.sqlExe("COMMIT", {})
acts.append("escalated: 维度「%s」QC 重做轮次用尽 → 人工介入" % label)
continue
# 大文件机械拆解(2026-09-15 用户要求):正文超阈值 → 重做也拆分段子任务,
# 避免单任务长读全文几十分钟;round 标记保证轮次按轮计不按任务条数计。
_tlen = await _tender_text_len(sor, project_id)
if _tlen > th["split_chars"]:
_base = {"stage": "analysis_redo", "task_kind": "bid_analysis",
"analysis_dim": dim, "qc_redo": 1}
_ids = await _create_split_tasks(
sor, project_id, dim, label, pname, _base, _tlen,
th["split_chars"], "redo%d" % (n_redo_done + 1), qc_imp=qc_imp)
acts.append("created: 分析重做拆解 %d 个分段子任务(%s,正文 %d 字)"
% (len(_ids), label, _tlen))
created_dims.append(dim)
continue
tid = await create_role_task(
sor, project_id, R_ANALYST,
"%s 分析重做:%s(QC 退回,按改进意见修正)" % (pname, label),
{"stage": "analysis_redo", "task_kind": "bid_analysis",
"analysis_dim": dim, "qc_redo": 1, "qc_improvements": qc_imp})
acts.append("created: 分析重做任务 %s(%s,改进意见 %d 条)" % (tid, label, len(qc_imp)))
created_dims.append(dim)
continue
first_attempts = max(0, attempts - n_redo_done)
if first_attempts >= MAX_ANALYST_ATTEMPT:
ht = await find_human_task(sor, project_id, "general", status="pending")
if not ht:
await create_human_task(
sor, project_id, "general",
"分析维度「%s」未产出结果,需人工介入" % label,
("该维度分析任务已完成 %d 次,产出仍为空。\n"
"请人工检查招标文件是否可读(bid_tender_files.content_text),"
"或人工补录该维度产出后完成本任务。" % first_attempts),
assignee_role="owner.superuser")
await sor.sqlExe("COMMIT", {})
acts.append("escalated: 维度「%s」重试超限 → 人工介入" % label)
continue
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 分支,不受编排影响。
# 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_eff()
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:
_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)逐章判定——
# 某一维度卡住只阻塞依赖它的章节,其余章节照常流转。
# 2026-09-11 流程裁剪:qc_analysis 阶段被裁 → 整段跳过(章节依赖由 _qc_passed_eff 视为全过)。
qc_pending = [] if not _stage_on(STAGE_QC) else await _qc_pending_types(sor, project_id)
# 2026-09-15 拆解配套门禁:维度有在办分析任务(含拆解分段子任务)时不派该类型
# QC——否则 QC 审到半成品产出必退(拆解机制成立的前提,见 _create_split_tasks)。
if qc_pending:
_kept = []
for _qp in qc_pending:
_d = QC_TYPE_TO_DIM.get(_qp[0], "")
if _d and await _open_dim_tasks(sor, project_id, _d) > 0:
acts.append("waiting: 维度「%s」分析任务在办,%s 暂缓 QC 审核(防审半成品)"
% (_d, _qp[0]))
continue
_kept.append(_qp)
qc_pending = _kept
over_types = set()
_over_active = False # 2026-09-16:本轮是否仍存在「轮次用尽等人工」的类型
if qc_pending:
# 轮次用尽守卫:最新审核已达上限仍未通过且产出未更新 → 不再空转派审核任务,
# 确保有阻塞人工任务(逃逸阀:人工修产出物或在记录页强制放行后流程自恢复)。
# 2026-09-03:产出在审核后更新过(out_newer,人工/修正任务落库)→ 不算用尽,
# 闭环「修正落库→自动复审」;只有产出真更新才开轮,不会无限烧轮。
over = [(t, s, imp, rnd, on) for t, s, imp, rnd, on in qc_pending
if s == "rework" and rnd >= th["qc_max_round"] and not on]
if over:
over_types = set(t for t, _, _, _, _ in over)
_over_active = True
ht = await find_human_task(sor, project_id, HT_QC_ESCALATION, status="pending")
if not ht:
await create_human_task(
sor, project_id, HT_QC_ESCALATION,
"解析产出 QC 连续不达标需人工介入",
("以下产出物审核已达 %d 轮上限仍未达标:\n%s\n\n"
"请人工核对招标文件并修正产出物,或在 QC 审核记录页(bid_qc_reviews)"
"把对应记录 passed 改为 1 强制放行。\n处理完成本任务后流程自动恢复。"
% (th["qc_max_round"],
"\n".join("- %s(第%s轮):%s" % (t, rnd, (imp or "")[:300])
for t, s, imp, rnd, _on in over))),
assignee_role="owner.superuser")
await sor.sqlExe("COMMIT", {})
acts.append("escalated: QC 审核轮次用尽(%s),等人工介入;不依赖它的章节照常流转"
% "/".join(sorted(over_types)))
# awaiting_redo:产出已被 QC 清空 → 重派对应维度的重做任务,不能对空产出审核
needs_redo = [(t, s, imp, rnd) for t, s, imp, rnd, _on in qc_pending if s == "awaiting_redo"]
if needs_redo:
redo_dims = []
for t, _, _, rnd in needs_redo:
dim = QC_TYPE_TO_DIM.get(t, "scoring")
if dim in redo_dims:
continue
if not _stage_on(DIM_TO_STAGE.get(dim, "")):
continue # 流程裁剪:被裁维度不派重做
if await _open_dim_tasks(sor, project_id, dim) > 0:
continue
n_redo_done = await _dim_redo_done_count(sor, project_id, dim)
if n_redo_done >= th["qc_max_round"]:
ht = await find_human_task(sor, project_id, "general", status="pending")
if not ht:
await create_human_task(
sor, project_id, "general",
"分析维度「%s」重做轮次已用尽,需人工介入" % dim,
("QC 退回后该维度已重做 %d 轮,产出仍未通过契合度审核。\n"
"请人工核对招标文件并修正产出物,或在 QC 审核记录页强制放行。"
% n_redo_done),
assignee_role="owner.superuser")
await sor.sqlExe("COMMIT", {})
continue
qc_imp = await _dim_qc_improvements(sor, project_id, dim)
# 大文件机械拆解(同 A 块口径):超阈值拆分段子任务重做
_tlen2 = await _tender_text_len(sor, project_id)
if _tlen2 > th["split_chars"]:
_base2 = {"stage": "analysis_redo", "task_kind": "bid_analysis",
"analysis_dim": dim, "qc_redo": 1}
_ids2 = await _create_split_tasks(
sor, project_id, dim, dim, pname, _base2, _tlen2,
th["split_chars"], "redo%d" % (int(rnd) + 1), qc_imp=qc_imp)
redo_dims.append(dim)
acts.append("created: 分析重做拆解 %d 段(%s,产出被 QC 清空)"
% (len(_ids2), dim))
continue
await create_role_task(
sor, project_id, R_ANALYST,
"%s 分析重做:%s(QC 退回产出已清空,按改进意见重做)" % (pname, dim),
{"stage": "analysis_redo", "task_kind": "bid_analysis",
"analysis_dim": dim, "qc_redo": 1, "qc_improvements": qc_imp})
redo_dims.append(dim)
acts.append("created: 分析重做任务(%s,产出被 QC 清空)" % dim)
if redo_dims:
await sor.sqlExe("COMMIT", {})
else:
acts.append("waiting: 分析重做任务在办(QC 退回产出已清空)")
# 逐类型并行审核:每类一个审核任务(跳过:在办的类型 / 空产出待重做 / 轮次用尽等人工)
needs_qc = [(t, s, imp, rnd) for t, s, imp, rnd, _on in qc_pending
if s == "unaudited" or (s == "rework" and t not in over_types)]
created_qc = []
for t, _, _, rnd in needs_qc:
if await _open_qc_type_tasks(sor, project_id, t) > 0:
continue
# 竞态守卫(2026-09-02):该维度分析/重做任务在办 → 不派审核。
# 根治:产出被清空后重做任务还在跑,QC 审核就先派出去对着重写中
# 的半成品打分(实测 cost_benefit 对空表打 0、qualifications 被连烧两轮)。
dim = QC_TYPE_TO_DIM.get(t, "")
if dim and await _open_dim_tasks(sor, project_id, dim) > 0:
acts.append("waiting: 维度「%s」分析任务在办,暂不派 %s 审核" % (dim, t))
continue
# 回归+增量结构(2026-09-03):复审任务带上「上一轮意见+本轮修正记录」,
# 审核输出必须分两段——上轮意见逐条回归验证、本轮新发现问题。
# 根治:审核方每轮全量重读无回归意识,旧问题修了没人验收、
# 新问题源源不断,产出在 8.0 分横住冲不过 8.5(实测 chapter_outline r4-r7)。
prev_imp, fixes = "", ""
if rnd and int(rnd) >= 2:
prev = await sor.sqlExe(
"SELECT improvement FROM bid_qc_reviews WHERE project_id=${p}$ "
"AND qc_type=${t}$ AND round=${r}$",
{"p": project_id, "t": t, "r": int(rnd) - 1})
await sor.sqlExe("COMMIT", {})
if prev:
prev_imp = rec_to_dict(prev[0]).get("improvement") or ""
# ⚠ audit_log 表只有 created_at(models/audit_log.json 与库一致),
# 曾写 updated_at → 1054 使整个 reconcile_project 崩溃、全项目派发停摆
# (2026-09-15 甘肃项目实测:第3轮复审派发路径首次触雷,累计崩 175 次)。
fx = await sor.sqlExe(
"SELECT action, entity, detail, created_at FROM audit_log "
"WHERE tenant_id=${p}$ AND entity IN ('bid_chapters','bid_doc_requirements',"
"'bid_scoring_items','bid_qualifications','bid_cost_benefit') "
"ORDER BY created_at DESC LIMIT 30",
{"p": project_id})
await sor.sqlExe("COMMIT", {})
fixes = "\n".join(
"- [%s] %s/%s: %s" % (rec_to_dict(f).get("created_at", "")[:16],
rec_to_dict(f).get("entity", ""),
rec_to_dict(f).get("action", ""),
(rec_to_dict(f).get("detail") or "")[:80])
for f in fx[:20])
qparams = {"stage": "qc_analysis", "task_kind": "bid_qc", "qc_types": [t]}
if prev_imp:
qparams["prev_round"] = int(rnd) - 1
qparams["prev_improvement"] = prev_imp
qparams["recent_fixes"] = fixes
tid = await create_role_task(
sor, project_id, R_QC,
"%s 解析产出契合度审核(%s)" % (pname, t),
qparams)
acts.append("created: QC 审核任务 %s(%s)" % (tid, t))
created_qc.append(t)
if created_qc:
await sor.sqlExe("COMMIT", {})
# 不提前 return:继续往下走章节级派发(依赖已满足的章节立即启动)
# ── A2b. qc_escalation 待办自动关闭(2026-09-16 甘肃项目实测 bug)──
# 升级待办发出后,人工修产出/记录页强制放行/重做通过 → 轮次用尽状态消失
# (_over_active=False),但 pending 的 qc_escalation 待办无人关闭就永远挂着,
# 角标常驻误导「还有事没办」。逃逸阀闭环的另一半:问题解决 → 待办自动收。
if not _over_active:
_ht_esc = await find_human_task(sor, project_id, HT_QC_ESCALATION, status="pending")
if _ht_esc:
await sor.sqlExe(
"UPDATE pipeline_human_tasks SET status='done', qc_status='passed', "
"qc_comment='对账器自动关闭:QC 轮次用尽的类型已全部通过/放行,流程已恢复', "
"submitted_at=NOW() WHERE id=${h}$ AND status='pending'",
{"h": _ht_esc.get("id")})
await sor.sqlExe("COMMIT", {})
acts.append("auto-closed: qc_escalation 待办(轮次用尽类型已全部通过/放行)")
# ── B. 资料准备(知识库资质匹配 + 人类文件清单)──
# 2026-09-03 章节级依赖:资料准备是商务章(business)的上游依赖,
# 不再阻塞全局流转——技术/报价章在上游 QC 通过后立即开写。
# 2026-09-11 流程裁剪:prep 被裁 → 不派。
# 2026-09-15 起资料准备不再阻塞商务章(资质缺口=信息流,见下方 C 段注释)。
prep_done = (not _stage_on(STAGE_PREP)) or await has_done_task(sor, project_id, R_PREP)
if _stage_on(STAGE_PREP) and n_quals > 0 and not prep_done:
if await count_open_tasks(sor, project_id, role=R_PREP) > 0:
acts.append("waiting: 资料准备任务在办(只阻塞商务章)")
else:
tid = await create_role_task(
sor, project_id, R_PREP,
"%s 投标资料准备(资质匹配 + 人类文件清单)" % pname,
{"stage": "prep", "task_kind": "bid_prep", "qual_pending": n_quals_pending})
await sor.sqlExe("COMMIT", {})
acts.append("created: 资料准备任务 %s" % tid)
if not chapters:
return acts + ["idle: 无章节,等待解析产出章节骨架"]
# ── C0. writing 孤儿章节回收(2026-09-16 甘肃项目实测 bug)──
# 写者任务 write_chapter 落 writing 后进程被打断/超时终态(任务 approved/failed/
# cancelled),章节停在 writing:C 段只扫 pending/rejected、D 段只扫 written,
# writing+无任务 = 永久死态(1.9/2.1.4 两章卡死 11 小时无人回收)。
# 回收规则:无任何在办任务时——有正文 → written(D 段自动派评审续走质量循环);
# 无正文 → rejected(C 段重派编写)。revise_count 不动(评审才是计轮主体)。
for c in chapters:
if c.get("status") != CH_WRITING:
continue
if await count_open_tasks(sor, project_id, chapter_id=c["id"]) > 0:
continue # 写者/评审任务还在办,不是孤儿
_has_txt = await sor.sqlExe(
"SELECT LENGTH(content) AS n FROM bid_chapters WHERE id=${i}$", {"i": c["id"]})
await sor.sqlExe("COMMIT", {})
_n = to_int(getattr(_has_txt[0], "n", 0) if _has_txt else 0, 0)
_new_st = CH_WRITTEN if _n > 0 else CH_REJECTED
await sor.sqlExe(
"UPDATE bid_chapters SET status=${st}$, updated_at=NOW() "
"WHERE id=${i}$ AND status=${old}$",
{"st": _new_st, "i": c["id"], "old": CH_WRITING})
await sor.sqlExe("COMMIT", {})
await record(sor, project_id, "bid_chapters", c["id"], "orphan_recover",
from_state=CH_WRITING, to_state=_new_st, who="bid_flow",
detail="writing 孤儿回收:无在办任务,正文 %d 字符 → %s" % (_n, _new_st))
c["status"] = _new_st # 本轮内存态同步,后续段直接可用
acts.append("recovered: 章[%s] writing 孤儿回收 → %s(无在办任务)"
% (c.get("chapter_no"), _new_st))
if any("orphan_recover" in (a or "") or "recovered: 章" in (a or "") for a in acts):
await sor.sqlExe("COMMIT", {})
# ── C1. 超重做上限章节的人工闭环(2026-09-16 甘肃项目实测 bug)──
# 旧行为缺陷:升级待办只在 review_chapter 调用瞬间发,且 general 类型级去重
# 会把第 2~N 章的信号吞掉;人工完成待办后章节永久冻结(revise_count 已超限,
# 无任何复位路径)。三态闭环:
# ① 有 pending 升级待办(标题精确匹配本章)→ 真在等人工,跳过;
# ② 无 pending 但有 done → 人工已处理 → revise_count 复位 0,C 段自动重派
# 写作(人工确认换一轮写作机会,留审计);
# ③ 无 pending 也无 done(信号曾丢失/从未发过)→ 按标题补发升级待办。
_esc_types = ("material_supply", "general")
for c in chapters:
if c.get("status") != CH_REJECTED:
continue
if to_int(c.get("revise_count"), 0) <= th["max_revise"]:
continue
_titles = ("章节反复不达标需人工介入:%s %s" % (c.get("chapter_no"), c.get("title")),
"章节实体材料缺口(信息流):%s %s" % (c.get("chapter_no"), c.get("title")))
# ⚠ sqlor 的 ${key}$ 替换自带引号——动态 IN 列表占位符必须裸写,
# 手工包单引号会变成 ''value'' 1064(2026-09-16 实测触雷)
_in = ",".join("${t%d}$" % i for i in range(len(_titles)))
_p = {"pid": project_id}
for i, t in enumerate(_titles):
_p["t%d" % i] = t[:480]
_recs = await sor.sqlExe(
"SELECT id, status FROM pipeline_human_tasks WHERE project_id=${pid}$ "
"AND task_type IN ('" + "','".join(_esc_types) + "') AND title IN (" + _in + ")",
_p)
await sor.sqlExe("COMMIT", {})
_rows = rows_to_dicts(_recs)
if any(r.get("status") == "pending" for r in _rows):
continue # ① 等人工中
if any(r.get("status") == "done" for r in _rows):
# ② 人工已处理 → 复位重派
await sor.sqlExe(
"UPDATE bid_chapters SET revise_count=0, updated_at=NOW() "
"WHERE id=${i}$ AND status=${old}$",
{"i": c["id"], "old": CH_REJECTED})
await sor.sqlExe("COMMIT", {})
await record(sor, project_id, "bid_chapters", c["id"], "human_reset",
to_state=CH_REJECTED, who="bid_flow",
detail="人工已完成升级待办 → revise_count 复位,自动重派写作")
c["revise_count"] = 0
acts.append("recovered: 章[%s] 人工处理完成,revise_count 复位自动重派"
% c.get("chapter_no"))
else:
# ③ 信号丢失/从未发过 → 补发(材料缺失挂起的章按评审意见判类型)
_cmr = await sor.sqlExe(
"SELECT LEFT(review_comment, 1600) AS rc FROM bid_chapters WHERE id=${i}$",
{"i": c["id"]})
await sor.sqlExe("COMMIT", {})
_cm = (getattr(_cmr[0], "rc", "") if _cmr else "") or ""
if "材料缺失" in _cm and "非写作缺陷" in _cm:
_tt, _ti = "material_supply", _titles[1]
_desc = ("章节「%s %s」评审判定实体材料缺失(补发信息流待办)。\n评审意见:\n%s\n\n"
"如需补材料请上传后走增补轮重写该章;不补则完成本待办即可。"
% (c.get("chapter_no"), c.get("title"), _cm[:1500]))
else:
_tt, _ti = "general", _titles[0]
_desc = ("章节「%s %s」已重做 %d 次仍未达门限(补发升级待办,原信号丢失)。\n"
"评审意见:\n%s\n\n"
"请人工判断:补充素材/调整章节范围/降低目标分,处理完成后完成本任务。"
% (c.get("chapter_no"), c.get("title"),
to_int(c.get("revise_count"), 0), _cm[:1500]))
await create_human_task(
sor, project_id, _tt, _ti, _desc,
assignee_role="owner.superuser", dedup_by_title=True)
await sor.sqlExe("COMMIT", {})
acts.append("escalated: 章[%s] 超重做上限且无待办记录,补发人工介入待办"
% c.get("chapter_no"))
# ── C. 章节编写(一章一任务、并发派发;pending/rejected 且无在办任务)──
# 商务/技术分角色:section=business → 商务标写者;其余(technical/price)→ 技术标写者。
# 并发:单轮最多派 write_concurrency 个编写任务(引擎另有每项目并发上限兜底)。
# 章节级上游依赖门禁(2026-09-03):每章只等自己 section 依赖的上游产出通过 QC
# (CH_SECTION_DEPS),依赖满足立即派发——某一维度卡住只阻塞依赖它的章节,
# 其余章节照常流转(根治:此前五类 QC 全过才放行,cost_benefit 一个维度慢
# 导致全部章节冻结)。商务章额外等资料准备完成。
qc_passed = await _qc_passed_eff()
# 2026-09-15 用户裁定:资质文件待办(human_docs_request)不再阻塞任何章节编写。
# 原 biz_ok=(prep_done or n_quals==0) and not hd_pending 会把「人类没上传资质文件」
# 变成商务章硬门禁——用户明确说忽略也不管用(答复还会被 agent QC 退回,见
# human_task_capability 的人类权威规则)。资质缺口是信息流(待办通知),不是门禁流。
# prep(资料准备)也不再挡商务章:与声明层对齐(chapter_write.deps 只有
# analysis_reqs_outline),写者需要资质信息时自查知识库/待办清单。
created_w = 0
waiting_deps = {}
for c in chapters:
if created_w >= th["write_concurrency"]:
break
if c.get("status") not in (CH_PENDING, CH_REJECTED):
continue
if to_int(c.get("revise_count"), 0) > th["max_revise"]:
continue # 已超重做上限,等人工介入(review 阶段已抛任务)
section = c.get("section") or "technical"
role = R_BIZ_WRITER if section == "business" else R_WRITER
if await count_open_tasks(sor, project_id, role=role, chapter_id=c["id"]) > 0:
continue
# 上游依赖门禁:该 section 依赖的所有产出类型最新一轮必须已通过契合度审核
deps = CH_SECTION_DEPS.get(section, CH_SECTION_DEPS["technical"])
unmet = [t for t in deps if t not in qc_passed]
if unmet:
waiting_deps.setdefault("/".join(sorted(unmet)), []).append(str(c.get("chapter_no")))
continue
stage = "revise" if c.get("status") == CH_REJECTED else "write"
tid = await create_role_task(
sor, project_id, role,
"%s 第%s章《%s》%s" % (pname, c.get("chapter_no"), c.get("title"),
"修改重写" if stage == "revise" else
("编写(商务标)" if role == R_BIZ_WRITER else "编写(技术标)")),
{"stage": stage, "task_kind": "bid_write", "chapter_id": c["id"],
"chapter_no": c.get("chapter_no"), "section": section,
"revise_count": c.get("revise_count")})
created_w += 1
acts.append("created: 章节编写任务 %s (%s/%s)" % (tid, c.get("chapter_no"), role))
if waiting_deps:
acts.append("waiting: 章节上游依赖未满足:" + "; ".join(
"章[%s]等(%s)" % ("/".join(nums), deps)
for deps, nums in sorted(waiting_deps.items())))
if created_w:
await sor.sqlExe("COMMIT", {})
# ── D. 章节评审(written 且无在办评审任务)──
# 2026-09-11 流程裁剪:chapter_review 被裁 → written 章节直接置 approved(用户已确认
# 放弃章节级质量把关,确认待办里有风险提示),不派评审任务。
if not _stage_on(STAGE_CH_REVIEW):
_auto_appr = [c for c in chapters if c.get("status") == CH_WRITTEN]
for c in _auto_appr:
await sor.sqlExe(
"UPDATE bid_chapters SET status=${st}$, updated_at=NOW() WHERE id=${i}$ "
"AND status=${old}$",
{"st": CH_APPROVED, "i": c["id"], "old": CH_WRITTEN})
if _auto_appr:
await sor.sqlExe("COMMIT", {})
acts.append("trimmed: chapter_review 已裁剪,%d 个章节 written→approved 直通"
% len(_auto_appr))
# 状态变了,本轮重读章节供 E 段判定
chs2 = await sor.sqlExe(
"SELECT id, chapter_no, title, section, status, revise_count FROM bid_chapters "
"WHERE project_id=${p}$ ORDER BY order_no, chapter_no", {"p": project_id})
await sor.sqlExe("COMMIT", {})
chapters = rows_to_dicts(chs2, limit=500)
created_r = 0
for c in chapters:
if not _stage_on(STAGE_CH_REVIEW):
break
if c.get("status") != CH_WRITTEN:
continue
if await count_open_tasks(sor, project_id, role=R_REVIEWER, chapter_id=c["id"]) > 0:
continue
tid = await create_role_task(
sor, project_id, R_REVIEWER,
"%s 第%s章《%s》评审打分" % (pname, c.get("chapter_no"), c.get("title")),
{"stage": "review", "task_kind": "bid_review", "chapter_id": c["id"],
"chapter_no": c.get("chapter_no")})
created_r += 1
acts.append("created: 章节评审任务 %s (%s)" % (tid, c.get("chapter_no")))
if created_r:
await sor.sqlExe("COMMIT", {})
# ── D2. 落库核验(2026-09-03,机制硬门禁)──
# 实测根因:人工修正任务 approved 但产出只写文件没落库(bid_chapters 停在旧时间),
# 机制对旧数据不复审、流程静默卡死,且任务树显示 approved 撒谎。
# 核验范围:仅人工修正/补录任务(系统 analyst 任务走工具原语落库、由 QC 对库审核兜底)。
# 不通过 → 状态改 verify_failed + 抛人工补落库/重跑。
verify_bad = []
_vtasks = await sor.sqlExe(
"SELECT id, title, params, created_at FROM pipeline_tasks WHERE tenant_id=${p}$ "
"AND role=${r}$ AND state IN ('approved','completed')",
{"p": project_id, "r": R_ANALYST}) or []
for t in _vtasks:
try:
p = json.loads(getattr(t, 'params', '') or '{}')
except Exception:
p = {}
dim = (p.get('analysis_dim') or '').strip()
title = getattr(t, 'title', '') or ''
if dim or p.get('qc_redo'):
continue # 系统分析/重做任务不在核验范围
if not ('修正' in title or '补录' in title):
continue
# 人工修正任务按标题关键词定目标表,避免跨维度误判
if 'chapter_outline' in title or '骨架' in title:
tbls = (QC_OUTPUT_TABLE[QC_TYPE_OUTLINE],)
elif 'doc_requirements' in title or '文件要求' in title or '补录' in title:
tbls = (QC_OUTPUT_TABLE[QC_TYPE_REQS],)
else:
tbls = (QC_OUTPUT_TABLE[QC_TYPE_REQS], QC_OUTPUT_TABLE[QC_TYPE_OUTLINE])
t_created = str(getattr(t, 'created_at', '') or '')
fresh = False
for tbl in tbls:
rec = await sor.sqlExe(
"SELECT MAX(GREATEST(created_at, IFNULL(updated_at, created_at))) AS mt FROM " + tbl +
" WHERE project_id=${p}$", {"p": project_id})
mt = str(getattr(rec[0], 'mt', '') or '') if rec else ''
if mt and mt >= t_created:
fresh = True
break
if not fresh:
tid = getattr(t, 'id', '')
await sor.sqlExe(
"UPDATE pipeline_tasks SET state='verify_failed', updated_at=NOW() WHERE id=${t}$",
{"t": tid})
verify_bad.append((tid, title[:50]))
if verify_bad:
await sor.sqlExe("COMMIT", {})
ht = await find_human_task(sor, project_id, "output_verify", status="pending")
if not ht:
await create_human_task(
sor, project_id, "output_verify",
"分析/修正任务声称完成但产出未落库",
("以下任务已 approved,但对应产出表的最大更新时间早于任务创建时间"
"(只写了文件或根本没写库),已标 verify_failed:\n%s\n\n"
"请处理:在对应管理页补录落库,或重跑该任务;落库后流程自动复审。"
"处理完完成本任务。"
% "\n".join("- %s %s" % (tid, ti) for tid, ti in verify_bad)),
assignee_role="owner.superuser")
await sor.sqlExe("COMMIT", {})
acts.append("escalated: 产出落库核验不通过(%d 个任务标 verify_failed),抛人工" % len(verify_bad))
else:
# 2026-09-03:核验全通过且有待处理的 output_verify → 自动关闭(问题已解决);
# 之前标 verify_failed 的人工修正任务,落库后恢复 approved(状态与事实一致)
await sor.sqlExe(
"UPDATE pipeline_tasks SET state='approved', updated_at=NOW() "
"WHERE tenant_id=${p}$ AND role=${r}$ AND state='verify_failed'",
{"p": project_id, "r": R_ANALYST})
ht = await find_human_task(sor, project_id, "output_verify", status="pending")
if ht:
await sor.sqlExe(
"UPDATE pipeline_human_tasks SET status='done', submitted_at=NOW() "
"WHERE id=${h}$", {"h": ht.get("id")})
acts.append("auto-closed: output_verify 待办(落库核验已通过)")
await sor.sqlExe("COMMIT", {})
if created_w or created_r:
return acts
# ── E. 合成标书(全部章节 approved;无标书或最新标书已被打回)──
all_approved = all(c.get("status") == CH_APPROVED for c in chapters)
if all_approved:
if not doc or doc.get("status") == DOC_REJECTED:
if await count_open_tasks(sor, project_id, role=R_COMPOSITOR) > 0:
return acts + ["waiting: 标书合成任务在办"]
tid = await create_role_task(
sor, project_id, R_COMPOSITOR, "%s 合成投标文件(标书)" % pname,
{"stage": "compose", "task_kind": "bid_compose",
"prev_version": doc.get("version") if doc else 0})
await sor.sqlExe("COMMIT", {})
return acts + ["created: 标书合成任务 %s" % tid]
# ── F. 整书评分 ──
# 2026-09-11 流程裁剪:whole_score 被裁 → 合成后直接置 passed(用户已确认
# 放弃模拟评标门禁),进入交付环节。
if doc.get("status") in (DOC_DRAFT, DOC_REVIEWING):
if not _stage_on(STAGE_WHOLE_SCORE):
await sor.sqlExe(
"UPDATE bid_documents SET status=${st}$, updated_at=NOW() WHERE id=${d}$ "
"AND status IN ('" + DOC_DRAFT + "','" + DOC_REVIEWING + "')",
{"st": DOC_PASSED, "d": doc.get("id")})
await sor.sqlExe("COMMIT", {})
return acts + ["trimmed: whole_score 已裁剪,标书 v%s 直接置 passed"
% doc.get("version")]
if await count_open_tasks(sor, project_id, role=R_SCORER) > 0:
return acts + ["waiting: 整书评分任务在办"]
tid = await create_role_task(
sor, project_id, R_SCORER,
"%s 标书整体评审与评分(v%s)" % (pname, doc.get("version")),
{"stage": "score", "task_kind": "bid_score",
"document_id": doc.get("id"), "version": doc.get("version")})
await sor.sqlExe("COMMIT", {})
return acts + ["created: 整书评分任务 %s" % tid]
# ── G. 通过 → 交付确认 → 项目完成 ──
if doc.get("status") == DOC_PASSED:
# 2026-09-11 流程裁剪:delivery_confirm 被裁 → 跳过人工交付确认直接完结
if not _stage_on(STAGE_DELIVERY):
await sor.sqlExe(
"UPDATE bid_documents SET status=${st}$, updated_at=NOW() WHERE id=${d}$",
{"st": DOC_DELIVERED, "d": doc.get("id")})
await sor.sqlExe(
"UPDATE sd_projects SET status='completed', updated_at=NOW() "
"WHERE id=${p}$ AND status<>'completed'", {"p": project_id})
await sor.sqlExe("COMMIT", {})
return acts + ["trimmed: delivery_confirm 已裁剪,标书直接置 delivered,项目完成"]
ht = await find_human_task(sor, project_id, HT_DELIVERY_CONFIRM)
# cancelled/rejected 的旧待办视为不存在(2026-09-16:章节退回重写时联动作废
# 过旧待办,重合成后必须重发新版确认待办,否则卡 waiting 死锁)
if not ht or ht.get("status") in ("cancelled", "rejected"):
hid = await create_human_task(
sor, project_id, HT_DELIVERY_CONFIRM, "标书评分通过,请确认交付",
"标书已通过评分门限,请人工复核用印/密封/份数/递交后确认交付。")
await sor.sqlExe("COMMIT", {})
return acts + ["created: 交付确认任务 %s" % hid]
if ht.get("status") == "done":
await sor.sqlExe(
"UPDATE bid_documents SET status=${st}$, updated_at=NOW() WHERE id=${d}$",
{"st": DOC_DELIVERED, "d": doc.get("id")})
await sor.sqlExe(
"UPDATE sd_projects SET status='completed', updated_at=NOW() "
"WHERE id=${p}$ AND status<>'completed'", {"p": project_id})
await sor.sqlExe("COMMIT", {})
return acts + ["completed: 标书已交付,项目完成"]
return acts + ["waiting: 等待人工交付确认"]
return acts or ["idle: 无待办动作"]
async def reconcile_all(sor=None):
"""对账全部活跃投标项目。返回 {project_id: [actions]}。"""
if sor is not None:
return await _reconcile_all_with(sor)
db, dbname = get_db()
async with db.sqlorContext(dbname) as s:
return await _reconcile_all_with(s)
# ══════════════ 纯技术方案流程分支(bid_tech_proposal,2026-09-14 批2) ══════════════
# 阶段链:tech_analysis(六类评估项)→ tech_qc(对标需求书原文)→ template_confirm
# (模版检索+用户确认骨架)→ tech_chapter_write → tech_chapter_qc(对标需求书覆盖度)
# → tech_compose(合成技术方案书)→ tech_delivery_confirm。
# 与标准流程的根本差异:**用户提供的技术需求书是唯一基准**——QC/评审全部对标需求书,
# 不碰评分项/资质/成本收益(那是招标文件的产物,本流程的输入不是招标文件)。
TECH_KIND = "tech" # 分析维度值(任务 params.analysis_dim,与标准四维度并列)
async def _tech_qc_latest(sor, project_id):
"""tech_items 最新一轮 QC(无记录返回 {})。"""
recs = await sor.sqlExe(
"SELECT id, round, fit_score, passed, improvement, task_id FROM bid_qc_reviews "
"WHERE project_id=${p}$ AND qc_type=${t}$ ORDER BY round DESC LIMIT 1",
{"p": project_id, "t": QC_TYPE_TECH})
await sor.sqlExe("COMMIT", {})
return rec_to_dict(recs[0]) if recs else {}
async def _open_tech_tasks(sor, project_id, task_kind):
"""某类技术方案任务在办数(按 params.task_kind 匹配)。"""
kw = '%"task_kind": "' + task_kind + '"%'
recs = await sor.sqlExe(
"SELECT COUNT(*) AS c FROM pipeline_tasks WHERE tenant_id=${p}$ "
"AND pipeline_id='role_task' AND state IN "
"('submitted','running','review','qc_review','waiting') AND params LIKE ${kw}$",
{"p": project_id, "kw": kw})
await sor.sqlExe("COMMIT", {})
return to_int(getattr(recs[0], "c", 0) if recs else 0)
async def _reconcile_tech_proposal(sor, project_id, pname, stage_table):
"""纯技术方案流程对账(幂等)。stage_table=已确认计划的阶段启用表。"""
acts = []
def _on(key):
return bool(stage_table.get(key, True))
th = await get_thresholds(sor)
n_items = await _count(sor, "SELECT COUNT(*) AS c FROM bid_tech_items "
"WHERE project_id=${p}$ AND kind='item'", {"p": project_id})
chs = await sor.sqlExe(
"SELECT id, chapter_no, title, section, status, revise_count FROM bid_chapters "
"WHERE project_id=${p}$ ORDER BY order_no, chapter_no", {"p": project_id})
await sor.sqlExe("COMMIT", {})
chapters = rows_to_dicts(chs, limit=500)
docs = await sor.sqlExe(
"SELECT id, version, status FROM bid_documents WHERE project_id=${p}$ "
"ORDER BY version DESC LIMIT 1", {"p": project_id})
await sor.sqlExe("COMMIT", {})
doc = rec_to_dict(docs[0]) if docs else {}
latest_qc = await _tech_qc_latest(sor, project_id)
qc_passed = (str(latest_qc.get("passed")) == "1") if latest_qc else False
# tech_qc 被裁 → T2 段的 _on(STAGE_TECH_QC) 判定即视为通过(用户已在确认待办里放弃该门禁),
# 不再单独算 qc_eff 死变量(2026-09-15 清理:原变量算了未被引用)。
# ── T1. 技术需求解析(六类评估项)──
if n_items == 0 and _on(STAGE_TECH_ANALYSIS):
if await _open_tech_tasks(sor, project_id, "bid_tech_analysis") > 0:
acts.append("waiting: 技术需求解析任务在办")
else:
qc_imp = (latest_qc.get("improvement") or "") if (latest_qc and str(latest_qc.get("passed")) != "1") else ""
attempts = await _count(
sor, "SELECT COUNT(*) AS c FROM pipeline_tasks WHERE tenant_id=${p}$ "
"AND pipeline_id='role_task' AND role=${r}$ AND state IN "
"('approved','completed','qc_rejected') AND params LIKE ${kw}$",
{"p": project_id, "r": R_ANALYST, "kw": '%bid_tech_analysis%'})
if attempts >= MAX_ANALYST_ATTEMPT + th["qc_max_round"]:
ht = await find_human_task(sor, project_id, "general", status="pending")
if not ht:
await create_human_task(
sor, project_id, "general",
"技术需求解析多次未产出,需人工介入:%s" % pname,
("技术评估项解析任务已执行 %d 次,产出仍为空。\n"
"请检查需求书是否可读(bid_tender_files.content_text),"
"或用 list_tech_items/add_tech_item 人工补录评估项。" % attempts),
assignee_role="owner.superuser")
await sor.sqlExe("COMMIT", {})
acts.append("escalated: 技术需求解析重试超限 → 人工介入")
else:
tparams = {"stage": "tech_analysis", "task_kind": "bid_tech_analysis",
"analysis_dim": TECH_KIND, "flow_key": FLOW_KEY_TECH_PROPOSAL}
if qc_imp:
tparams["qc_redo"] = 1
tparams["qc_improvements"] = ["[tech_items] %s" % qc_imp[:2000]]
tid = await create_role_task(
sor, project_id, R_ANALYST,
"%s 技术需求解析(六类评估项%s)" % (pname, ",QC退回重做" if qc_imp else ""),
tparams)
await sor.sqlExe("COMMIT", {})
acts.append("created: 技术需求解析任务 %s" % tid)
return acts or ["idle: 技术需求解析中"]
# ── T2. 解析产出 QC(对标需求书原文,qc_type=tech_items)──
if _on(STAGE_TECH_QC) and n_items > 0 and not qc_passed:
rnd = to_int(latest_qc.get("round"), 0) if latest_qc else 0
if latest_qc and rnd >= th["qc_max_round"] and str(latest_qc.get("passed")) != "1":
# 轮次用尽:产出未更新则不空转,抛人工(同标准流程逃逸阀)
ht = await find_human_task(sor, project_id, HT_QC_ESCALATION, status="pending")
if not ht:
await create_human_task(
sor, project_id, HT_QC_ESCALATION,
"技术评估项 QC 连续不达标需人工介入",
("「tech_items」已审核 %d 轮仍未达标(最近改进意见:%s)。\n\n"
"请人工核对技术需求书修正评估项(list_tech_items/add_tech_item),"
"或在 QC 审核记录页把该记录 passed 改为 1 强制放行。"
% (rnd, (latest_qc.get("improvement") or "")[:1500])),
assignee_role="owner.superuser")
await sor.sqlExe("COMMIT", {})
acts.append("escalated: 技术评估项 QC 轮次用尽 → 人工介入")
return acts
if await _open_tech_tasks(sor, project_id, "bid_qc") > 0:
acts.append("waiting: 技术评估项 QC 任务在办")
return acts
# 竞态守卫:解析/重做任务在办 → 不派审核(防对半成品打分)
if await _open_tech_tasks(sor, project_id, "bid_tech_analysis") > 0:
acts.append("waiting: 解析任务在办,暂不派 tech QC")
return acts
tid = await create_role_task(
sor, project_id, R_QC,
"%s 技术评估项契合度审核(对标需求书原文)" % pname,
{"stage": "tech_qc", "task_kind": "bid_qc", "qc_types": [QC_TYPE_TECH],
"flow_key": FLOW_KEY_TECH_PROPOSAL})
await sor.sqlExe("COMMIT", {})
return acts + ["created: 技术评估项 QC 任务 %s" % tid]
# ── T3. 模版确定 + 用户确认骨架 ──
if not chapters:
if not _on(STAGE_TEMPLATE_CONFIRM):
# 被裁(trim=no 理论上裁不掉,防御):无人确认不得开写
return acts + ["blocked: 模版确认阶段被裁但骨架缺失,无法开写章节"]
# 阻塞门禁:pending 的模版确认待办在 0b 已全局阻塞,走到这里只处理「未发/被驳回」
ht_pend = await find_human_task(sor, project_id, HT_TEMPLATE_CONFIRM, status="pending")
if ht_pend:
return acts + ["waiting: 技术方案骨架待用户确认"]
if await _open_tech_tasks(sor, project_id, "tech_template") > 0:
return acts + ["waiting: 模版/骨架拟定任务在办"]
from .bid_tech_capability import get_template_reject_comment
rej = await get_template_reject_comment(sor, project_id)
tid = await create_role_task(
sor, project_id, R_PM,
"%s 拟定技术方案模版骨架%s" % (pname, "(按驳回意见修订)" if rej else ""),
{"stage": "template_confirm", "task_kind": "tech_template",
"flow_key": FLOW_KEY_TECH_PROPOSAL,
**({"reject_comment": rej[:2000]} if rej else {})})
await sor.sqlExe("COMMIT", {})
return acts + ["created: 模版骨架任务 %s(PM 检索模版→拟骨架→发用户确认)" % tid]
# ── T4. 分章编写(骨架已确认;一章一任务并发派发)──
created_w = 0
for c in chapters:
if created_w >= th["write_concurrency"]:
break
if c.get("status") not in (CH_PENDING, CH_REJECTED):
continue
if to_int(c.get("revise_count"), 0) > th["max_revise"]:
continue # 超重做上限,评审阶段已抛人工
if await count_open_tasks(sor, project_id, role=R_WRITER, chapter_id=c["id"]) > 0:
continue
stage = "revise" if c.get("status") == CH_REJECTED else "write"
tid = await create_role_task(
sor, project_id, R_WRITER,
"%s 技术方案 第%s章《%s》%s" % (pname, c.get("chapter_no"), c.get("title"),
"修改重写" if stage == "revise" else "编写"),
{"stage": stage, "task_kind": "bid_write", "chapter_id": c["id"],
"chapter_no": c.get("chapter_no"), "section": "technical",
"flow_key": FLOW_KEY_TECH_PROPOSAL,
"revise_count": c.get("revise_count")})
created_w += 1
acts.append("created: 技术方案章节编写任务 %s (%s)" % (tid, c.get("chapter_no")))
if created_w:
await sor.sqlExe("COMMIT", {})
# ── T5. 章节 QC 评审(对标需求书覆盖度)──
if not _on(STAGE_TECH_CH_QC):
_auto_appr = [c for c in chapters if c.get("status") == CH_WRITTEN]
for c in _auto_appr:
await sor.sqlExe(
"UPDATE bid_chapters SET status=${st}$, updated_at=NOW() WHERE id=${i}$ "
"AND status=${old}$",
{"st": CH_APPROVED, "i": c["id"], "old": CH_WRITTEN})
if _auto_appr:
await sor.sqlExe("COMMIT", {})
acts.append("trimmed: tech_chapter_qc 已裁剪,%d 个章节 written→approved 直通"
% len(_auto_appr))
chs2 = await sor.sqlExe(
"SELECT id, chapter_no, title, section, status, revise_count FROM bid_chapters "
"WHERE project_id=${p}$ ORDER BY order_no, chapter_no", {"p": project_id})
await sor.sqlExe("COMMIT", {})
chapters = rows_to_dicts(chs2, limit=500)
created_r = 0
for c in chapters:
if not _on(STAGE_TECH_CH_QC):
break
if c.get("status") != CH_WRITTEN:
continue
if await count_open_tasks(sor, project_id, role=R_REVIEWER, chapter_id=c["id"]) > 0:
continue
tid = await create_role_task(
sor, project_id, R_REVIEWER,
"%s 技术方案 第%s章《%s》评审(对标需求书)" % (pname, c.get("chapter_no"), c.get("title")),
{"stage": "review", "task_kind": "bid_review", "chapter_id": c["id"],
"chapter_no": c.get("chapter_no"), "flow_key": FLOW_KEY_TECH_PROPOSAL})
created_r += 1
acts.append("created: 技术方案章节评审任务 %s (%s)" % (tid, c.get("chapter_no")))
if created_r:
await sor.sqlExe("COMMIT", {})
if created_w or created_r:
return acts
# ── T6. 合成技术方案书 ──
all_approved = all(c.get("status") == CH_APPROVED for c in chapters)
if all_approved:
if not doc or doc.get("status") == DOC_REJECTED:
if not _on(STAGE_TECH_COMPOSE):
return acts + ["blocked: 合成阶段被裁但流程要求产出技术方案书,无法完成"]
if await count_open_tasks(sor, project_id, role=R_COMPOSITOR) > 0:
return acts + ["waiting: 技术方案书合成任务在办"]
tid = await create_role_task(
sor, project_id, R_COMPOSITOR, "%s 合成技术方案书" % pname,
{"stage": "compose", "task_kind": "bid_compose",
"flow_key": FLOW_KEY_TECH_PROPOSAL,
"prev_version": doc.get("version") if doc else 0})
await sor.sqlExe("COMMIT", {})
return acts + ["created: 技术方案书合成任务 %s" % tid]
# ── T7. 交付确认 ──
if doc.get("status") in (DOC_DRAFT, DOC_REVIEWING, DOC_PASSED):
if not _on(STAGE_TECH_DELIVERY):
await sor.sqlExe(
"UPDATE bid_documents SET status=${st}$, updated_at=NOW() WHERE id=${d}$",
{"st": DOC_DELIVERED, "d": doc.get("id")})
await sor.sqlExe(
"UPDATE sd_projects SET status='completed', updated_at=NOW() "
"WHERE id=${p}$ AND status<>'completed'", {"p": project_id})
await sor.sqlExe("COMMIT", {})
return acts + ["trimmed: tech_delivery_confirm 已裁剪,技术方案书直接置 delivered,项目完成"]
ht = await find_human_task(sor, project_id, HT_DELIVERY_CONFIRM)
if not ht:
hid = await create_human_task(
sor, project_id, HT_DELIVERY_CONFIRM,
"技术方案书已完成,请确认交付",
"技术方案书已合成并通过章节评审,请下载核对后确认交付。")
await sor.sqlExe("COMMIT", {})
return acts + ["created: 交付确认任务 %s" % hid]
if ht.get("status") == "done":
await sor.sqlExe(
"UPDATE bid_documents SET status=${st}$, updated_at=NOW() WHERE id=${d}$",
{"st": DOC_DELIVERED, "d": doc.get("id")})
await sor.sqlExe(
"UPDATE sd_projects SET status='completed', updated_at=NOW() "
"WHERE id=${p}$ AND status<>'completed'", {"p": project_id})
await sor.sqlExe("COMMIT", {})
return acts + ["completed: 技术方案书已交付,项目完成"]
return acts + ["waiting: 等待人工交付确认"]
return acts or ["idle: 技术方案流程无待办动作"]
async def _reconcile_all_with(sor):
recs = await sor.sqlExe(
"SELECT id, name FROM sd_projects WHERE pipeline_id=${pl}$ "
"AND status NOT IN ('paused','archived','completed','cancelled') LIMIT 200",
{"pl": PIPELINE_ID})
await sor.sqlExe("COMMIT", {})
out = {}
for r in (recs or []):
pid = getattr(r, "id", "")
pname = getattr(r, "name", "")
if not pid:
continue
try:
out[pid] = await reconcile_project(sor, pid, pname)
except Exception as e:
logger.exception("reconcile project %s failed", pid)
out[pid] = ["error: %s" % str(e)[:200]]
return out
async def diagnose_bid_project(project_id):
"""诊断投标项目:阶段/门禁/章节分布/在办任务/最近评分,一次看清卡在哪。"""
db, dbname = get_db()
async with db.sqlorContext(dbname) as sor:
prec = await sor.sqlExe(
"SELECT id, name, status, pipeline_id FROM sd_projects WHERE id=${p}$",
{"p": project_id})
await sor.sqlExe("COMMIT", {})
if not prec:
return "项目不存在: %s" % project_id
proj = rec_to_dict(prec[0])
# 门禁:非 in_progress 的项目只诊断不派发(2026-09-02:诊断入口曾绕过
# paused 排除派发任务,违反"项目 paused 零消耗"铁律)
_proj_paused = proj.get("status") not in ("in_progress", "active", "running")
hts = await sor.sqlExe(
"SELECT id, task_type, title, status FROM pipeline_human_tasks "
"WHERE project_id=${p}$ ORDER BY created_at DESC LIMIT 20", {"p": project_id})
await sor.sqlExe("COMMIT", {})
tasks = await sor.sqlExe(
"SELECT id, role, title, state, retry_count FROM pipeline_tasks "
"WHERE tenant_id=${p}$ AND pipeline_id='role_task' "
"ORDER BY created_at DESC LIMIT 30", {"p": project_id})
await sor.sqlExe("COMMIT", {})
chs = await sor.sqlExe(
"SELECT status, COUNT(*) AS c FROM bid_chapters WHERE project_id=${p}$ "
"GROUP BY status", {"p": project_id})
await sor.sqlExe("COMMIT", {})
docs = await sor.sqlExe(
"SELECT version, status, total_score, max_score FROM bid_documents "
"WHERE project_id=${p}$ ORDER BY version DESC LIMIT 5", {"p": project_id})
await sor.sqlExe("COMMIT", {})
counts = {}
for tbl in ("bid_tender_files", "bid_scoring_items", "bid_qualifications",
"bid_doc_requirements", "bid_qc_reviews", "bid_chapters"):
counts[tbl] = await _count(
sor, "SELECT COUNT(*) AS c FROM " + tbl + " WHERE project_id=${p}$",
{"p": project_id})
acts = (["skipped: 项目状态 %s(非 in_progress),只诊断不派发" % proj.get("status")]
if _proj_paused
else await reconcile_project(sor, project_id, proj.get("name", "")))
return json.dumps({
"project": proj, "counts": counts,
"chapter_status": {getattr(r, "status", ""): to_int(getattr(r, "c", 0))
for r in (chs or [])},
"documents": rows_to_dicts(docs),
"human_tasks": rows_to_dicts(hts),
"role_tasks": rows_to_dicts(tasks),
"reconcile_actions": acts,
}, ensure_ascii=False, default=str)
# ══════════════ 后台 poller ══════════════
def start_poller():
"""注册投标产线对账 poller(ahserver 启动钩子)。
on_startup handler 只 spawn 后台任务,绝不内联无限循环(否则阻塞 aiohttp 启动)。
"""
try:
from ahserver.configuredServer import add_startup
except ImportError:
logger.warning("ahserver.add_startup 不可用,投标 poller 未启动")
return False
async def _bid_poller(app):
async def _loop():
db, dbname = get_db()
POLLER_STATE["last_error"] = ""
while True:
try:
async def _once():
async with db.sqlorContext(dbname) as sor:
res = await _reconcile_all_with(sor)
return res
import datetime
# watchdog:每轮最多 60 秒,防连接池/MDL 锁卡死导致 poller 永久停摆
res = await asyncio.wait_for(_once(), timeout=60)
POLLER_STATE["rounds"] += 1
POLLER_STATE["last_run"] = datetime.datetime.now().strftime("%H:%M:%S")
POLLER_STATE["last_error"] = ""
for pid, acts in (res or {}).items():
for a in acts:
# waiting 也输出(2026-09-15):拆解门禁「维度在办暂缓 QC」
# 等合法等待态是可观测性关键,不输出会误判门禁没生效。
if a.startswith(("created", "escalated", "completed",
"error", "trimmed", "waiting")):
try:
from appPublic.log import debug as _dbg
_dbg("bid_flow %s: %s" % (pid, a))
except Exception:
pass
logger.info("bid_flow %s: %s", pid, a)
except Exception as e:
POLLER_STATE["last_error"] = "%s: %s" % (type(e).__name__, str(e)[:160])
logger.warning("bid_poller error/timeout: %s", str(e)[:200])
await asyncio.sleep(POLL_INTERVAL)
try:
from appPublic.log import debug as _dbg0
_dbg0("bid_flow poller handler invoked, creating loop task")
except Exception:
pass
asyncio.create_task(_loop())
try:
from appPublic.log import debug as _dbg
_dbg("bid_flow poller started (interval=%ss)" % POLL_INTERVAL)
except Exception:
pass
logger.info("bid_flow poller started (interval=%ss)", POLL_INTERVAL)
add_startup(_bid_poller)
return True