yumoqing e387c7c5ce fix(bidding): 解析读不到完整招标书根治——最优文件选择+路径解析+升级自动恢复
- extract/read 选最优记录: 磁盘文件大小与正文长度取大, 替代固定取最老
- 相对/裸文件名按项目根目录解析, 解析后回写绝对路径
- bid_flow: 新招标文件到来自动重置解析计数+清旧产出残留+关旧人工任务(防永久卡死)
- 主 agent 补 reset_task_retry 工具(待办文案承诺但一直未实现)
2026-09-01 09:26:45 +08:00

662 lines
34 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
from .bid_common import (
get_db, rows_to_dicts, rec_to_dict, to_int, to_float, json_loads,
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,
QC_TYPE_SCORING, QC_TYPE_QUALS, QC_TYPE_REQS, QC_TYPE_OUTLINE, QC_TYPES,
QC_OUTPUT_TABLE,
)
logger = logging.getLogger("pipeline.bidding")
R_ANALYST = "agent.tender_analyst"
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_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)
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"):
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):
state=unaudited:产出已存在但从未审核;state=rework:最近一轮审核未通过。
已通过(或人工在记录页强制放行)的类型不再出现。
"""
pending = []
for t in QC_TYPES:
tbl = QC_OUTPUT_TABLE[t]
n_out = await _count(sor, "SELECT COUNT(*) AS c FROM " + tbl +
" WHERE project_id=${p}$", {"p": project_id})
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:
if n_out > 0:
pending.append((t, "unaudited", "", 0))
continue
r = rec_to_dict(recs[0])
if str(r.get("passed")) == "1":
continue
rnd = to_int(r.get("round"), 1)
if n_out > 0:
pending.append((t, "rework", r.get("improvement") or "", rnd))
else:
# 产出已被清空待重做:round 信息带回,供对账器判断是否已超轮次上限
pending.append((t, "awaiting_redo", r.get("improvement") or "", rnd))
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 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)
# 产线起点兜底:项目无任何招标文件且无待办上传任务 → 自动发「上传招标文件」任务
# (无论项目经 create_bid_project 还是通用 create_project 创建,都从上传文件起步)
_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:
_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: 有待办人类任务(等招标文件/等人类资料/等QC人工介入),本轮不派发"]
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})
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. 招标文件解析(评分项 / 资质 / 投标文件要求 / 章节骨架)──
need_analysis = n_files > 0 and (n_items == 0 or n_struct == 0 or not chapters)
if need_analysis:
if await count_open_tasks(sor, project_id, role=R_ANALYST) > 0:
return acts + ["waiting: 招标文件解析任务在办"]
# 携带 QC 改进意见的重做任务(qc_redo=1)不计入首跑重试上限
n_redo_open = 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 "
"('submitted','running','review','qc_review') AND params LIKE ${kw}$",
{"p": project_id, "r": R_ANALYST, "kw": "%qc_redo%"})
attempts = await _done_task_count(sor, project_id, R_ANALYST)
n_redo_done = await _redo_done_count(sor, project_id)
# 自动恢复(2026-09-01):解析升级后来了**新的**招标文件(上传时间晚于
# 最后一次解析完成)→ 清掉旧文件残留产出与旧 QC 意见,重试计数从新文件
# 时刻重起。否则项目永久卡死:旧计数 ≥ 上限,永远不再派解析
# (中传项目「走不下去」的另一半根因)。
_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: 检测到新招标文件,已清空旧解析残留,重新解析")
attempts = await _done_task_count(sor, project_id, R_ANALYST, since=_ft)
n_redo_done = await _redo_done_count(sor, project_id, since=_ft)
first_attempts = max(0, attempts - n_redo_done)
qc_imp = await _qc_collect_improvements(sor, project_id)
if qc_imp or n_redo_open > 0:
# QC 退回重做:携带改进意见,不受首跑重试上限约束;但重做轮次有上限(防无限循环)
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",
"解析重做轮次已用尽,需人工介入",
("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", {})
return acts + ["escalated: QC 重做轮次用尽 → 人工介入"]
tid = await create_role_task(
sor, project_id, R_ANALYST,
"%s 招标文件解析重做(QC 退回,按改进意见修正)" % pname,
{"stage": "analysis_redo", "task_kind": "bid_analysis", "qc_redo": 1,
"qc_improvements": qc_imp})
await sor.sqlExe("COMMIT", {})
return acts + ["created: 解析重做任务 %s(QC 改进意见 %d 条)" % (tid, len(qc_imp))]
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", "招标文件解析未产出完整结果,需人工介入",
("解析任务已完成 %d 次,但仍缺少:%s。\n"
"请人工检查招标文件是否可读(bid_tender_files.content_text),"
"或人工补录评分项/投标文件要求后完成本任务。"
% (first_attempts,
"、".join(x for x in [
"评分项" if n_items == 0 else "",
"投标文件章节结构" if n_struct == 0 else "",
"章节骨架" if not chapters else ""] if x))),
assignee_role="owner.superuser")
await sor.sqlExe("COMMIT", {})
acts.append("escalated: 解析重试超限 → 人工介入")
return acts
tid = await create_role_task(
sor, project_id, R_ANALYST,
"%s 招标文件解析(评分项/资质/投标文件要求/章节骨架)" % pname,
{"stage": "analysis", "task_kind": "bid_analysis"})
await sor.sqlExe("COMMIT", {})
return acts + ["created: 解析任务 %s" % tid]
# ── A2. 解析产出 QC 门禁(评分标准/得分规则/资质/投标文件要求/章节骨架的契合度审核)──
# 每类产出须经 agent.qc 对照招标文件原文按 10 分制打分,得分高于通过分(默认 9.5)才放行。
qc_pending = await _qc_pending_types(sor, project_id)
if qc_pending:
# 轮次用尽守卫:最新审核已达上限仍未通过 → 不再空转派审核任务,
# 确保有阻塞人工任务(逃逸阀:人工修产出物或在记录页强制放行后流程自恢复)。
over = [(t, s, imp, rnd) for t, s, imp, rnd in qc_pending
if s == "rework" and rnd >= th["qc_max_round"]]
if over:
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 in over))),
assignee_role="owner.superuser")
await sor.sqlExe("COMMIT", {})
return acts + ["blocked: QC 审核轮次用尽(%s),等人工介入"
% "/".join(t for t, _, _, _ in over)]
# awaiting_redo:产出已被 QC 清空 → 先重派解析重做,不能对空产出审核
needs_redo = [(t, s, imp, rnd) for t, s, imp, rnd in qc_pending if s == "awaiting_redo"]
if needs_redo:
if await count_open_tasks(sor, project_id, role=R_ANALYST) > 0:
return acts + ["waiting: 解析重做任务在办(QC 退回产出已清空)"]
n_redo_done = await _redo_done_count(sor, project_id)
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",
"解析重做轮次已用尽,需人工介入",
("QC 退回后解析已重做 %d 轮,产出仍未通过契合度审核。\n"
"请人工核对招标文件并修正产出物,或在 QC 审核记录页强制放行。"
% n_redo_done),
assignee_role="owner.superuser")
await sor.sqlExe("COMMIT", {})
return acts + ["escalated: QC 重做轮次用尽 → 人工介入"]
qc_imp = await _qc_collect_improvements(sor, project_id)
tid = await create_role_task(
sor, project_id, R_ANALYST,
"%s 招标文件解析重做(QC 退回产出已清空,按改进意见重抽)" % pname,
{"stage": "analysis_redo", "task_kind": "bid_analysis", "qc_redo": 1,
"qc_improvements": qc_imp})
await sor.sqlExe("COMMIT", {})
return acts + ["created: 解析重做任务 %s(产出被 QC 清空:%s)"
% (tid, "/".join(t for t, _, _, _ in needs_redo))]
needs_qc = [(t, s, imp, rnd) for t, s, imp, rnd in qc_pending]
if await count_open_tasks(sor, project_id, role=R_QC) > 0:
return acts + ["waiting: 解析产出 QC 审核任务在办"]
types_str = "/".join(t for t, _, _, _ in needs_qc)
tid = await create_role_task(
sor, project_id, R_QC,
"%s 解析产出契合度审核(%s)" % (pname, types_str),
{"stage": "qc_analysis", "task_kind": "bid_qc",
"qc_types": [t for t, _, _, _ in needs_qc]})
await sor.sqlExe("COMMIT", {})
return acts + ["created: QC 审核任务 %s(%s)" % (tid, types_str)]
# ── B. 资料准备(知识库资质匹配 + 人类文件清单)──
prep_done = await has_done_task(sor, project_id, R_PREP)
if n_quals > 0 and not prep_done:
if await count_open_tasks(sor, project_id, role=R_PREP) > 0:
return acts + ["waiting: 资料准备任务在办"]
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", {})
return acts + ["created: 资料准备任务 %s" % tid]
if not chapters:
return acts + ["idle: 无章节,等待解析产出章节骨架"]
# ── C. 章节编写(一章一任务、并发派发;pending/rejected 且无在办任务)──
# 商务/技术分角色:section=business → 商务标写者;其余(technical/price)→ 技术标写者。
# 并发:单轮最多派 write_concurrency 个编写任务(引擎另有每项目并发上限兜底)。
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 # 已超重做上限,等人工介入(review 阶段已抛任务)
role = R_BIZ_WRITER if (c.get("section") or "technical") == "business" else R_WRITER
if await count_open_tasks(sor, project_id, role=role, chapter_id=c["id"]) > 0:
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": c.get("section") or "technical",
"revise_count": c.get("revise_count")})
created_w += 1
acts.append("created: 章节编写任务 %s (%s/%s)" % (tid, c.get("chapter_no"), role))
if created_w:
await sor.sqlExe("COMMIT", {})
# ── D. 章节评审(written 且无在办评审任务)──
created_r = 0
for c in chapters:
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", {})
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. 整书评分 ──
if doc.get("status") in (DOC_DRAFT, DOC_REVIEWING):
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:
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(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)
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])
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 = 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:
if a.startswith(("created", "escalated", "completed", "error")):
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