"""投标产线任务流转引擎:状态收敛器(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, 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_TYPE_COST, 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) # ── 四维度分析辅助(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-03):每种章节 section 只依赖自己必需的分析产出 # 通过 QC——某一维度 QC 卡住只阻塞依赖它的章节,其余章节照常流转。 # 公共事实:章节骨架(bid_chapters)本身就是章节,OUTLINE 过了章节才可信; # REQS(投标文件格式要求)决定章节怎么写,所有章节都要。 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_QUALS, 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 _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)已完成任务数(不计入首跑重试上限)。""" 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}$ 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): 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 _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 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: 有待办人类任务(等招标文件/等人类资料/等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}) 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) created_dims = [] 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_done_count(sor, project_id, dim, since=_ft2) n_redo_done = await _dim_redo_done_count(sor, project_id, dim, since=_ft2) # 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 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 tid = await create_role_task( sor, project_id, R_ANALYST, "%s 招标文件分析:%s" % (pname, label), {"stage": "analysis", "task_kind": "bid_analysis", "analysis_dim": dim}) acts.append("created: 分析任务 %s(%s)" % (tid, label)) created_dims.append(dim) if created_dims: await sor.sqlExe("COMMIT", {}) # 2026-09-03:不再提前 return——同轮继续章节级派发(依赖已满足的章节立即启动) # ── A2. 分析产出 QC 门禁(五个类型逐一审核,每类一个审核任务并行;阈值见 bid_qc_pass_score)── # 2026-09-03 章节级门禁:本段只负责「推进 QC」(重做/审核/冒泡),不再提前 return # 阻塞全局流转。哪些章节可以开写,由 C 段按章节依赖矩阵(CH_SECTION_DEPS)逐章判定—— # 某一维度卡住只阻塞依赖它的章节,其余章节照常流转。 qc_pending = await _qc_pending_types(sor, project_id) over_types = set() 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: over_types = set(t for t, _, _, _ in 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", {}) acts.append("escalated: QC 审核轮次用尽(%s),等人工介入;不依赖它的章节照常流转" % "/".join(sorted(over_types))) # awaiting_redo:产出已被 QC 清空 → 重派对应维度的重做任务,不能对空产出审核 needs_redo = [(t, s, imp, rnd) for t, s, imp, rnd in qc_pending if s == "awaiting_redo"] if needs_redo: redo_dims = [] for t, _, _, _ in needs_redo: dim = QC_TYPE_TO_DIM.get(t, "scoring") if dim in redo_dims: 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) 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 in qc_pending if s == "unaudited" or (s == "rework" and t not in over_types)] created_qc = [] for t, _, _, _ 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 tid = await create_role_task( sor, project_id, R_QC, "%s 解析产出契合度审核(%s)" % (pname, t), {"stage": "qc_analysis", "task_kind": "bid_qc", "qc_types": [t]}) acts.append("created: QC 审核任务 %s(%s)" % (tid, t)) created_qc.append(t) if created_qc: await sor.sqlExe("COMMIT", {}) # 不提前 return:继续往下走章节级派发(依赖已满足的章节立即启动) # ── B. 资料准备(知识库资质匹配 + 人类文件清单)── # 2026-09-03 章节级依赖:资料准备是商务章(business)的上游依赖, # 不再阻塞全局流转——技术/报价章在上游 QC 通过后立即开写。 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: 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: 无章节,等待解析产出章节骨架"] # ── 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_types(sor, project_id) hd_pending = await find_human_task(sor, project_id, HT_HUMAN_DOCS, status="pending") biz_ok = (prep_done or n_quals == 0) and not hd_pending 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 if section == "business" and not biz_ok: waiting_deps.setdefault("bid_prep(资料准备)", []).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 且无在办评审任务)── 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]) # 门禁:非 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: 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