587 lines
30 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):
return 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')",
{"p": project_id, "r": role})
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 _redo_done_count(sor, project_id):
"""已完成的 QC 退回重做任务数(解析重做轮次计数,防无限重派)。"""
return 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') AND params LIKE ${kw}$",
{"p": project_id, "r": R_ANALYST, "kw": "%qc_redo%"})
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)
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