diff --git a/pipeline_bidding/bid_flow.py b/pipeline_bidding/bid_flow.py index ddc9ec5..40fe9b0 100644 --- a/pipeline_bidding/bid_flow.py +++ b/pipeline_bidding/bid_flow.py @@ -719,13 +719,15 @@ async def reconcile_project(sor, project_id, project_name=""): await sor.sqlExe("COMMIT", {}) # ── D2. 落库核验(2026-09-03,机制硬门禁)── - # 实测根因:修正任务 approved 但产出只写文件没落库(bid_chapters 停在旧时间), - # 机制对旧数据不复审、流程静默卡死。核验:approved 的分析/修正任务, - # 对应产出表 max 更新时间必须 ≥ 任务创建时间,否则抛人工补落库/重跑。 + # 实测根因:人工修正任务 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='approved'", {"p": project_id, "r": R_ANALYST}) or [] + "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 '{}') @@ -733,19 +735,17 @@ async def reconcile_project(sor, project_id, project_name=""): p = {} dim = (p.get('analysis_dim') or '').strip() title = getattr(t, 'title', '') or '' - is_manual = ('修正' in title or '补录' in title) and not p.get('qc_redo') - if not dim and not is_manual: + if dim or p.get('qc_redo'): + continue # 系统分析/重做任务不在核验范围 + if not ('修正' in title or '补录' in title): continue - if dim: - tbls = DIM_TABLES.get(dim) or (QC_OUTPUT_TABLE[QC_TYPE_REQS], QC_OUTPUT_TABLE[QC_TYPE_OUTLINE]) + # 人工修正任务按标题关键词定目标表,避免跨维度误判 + 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: - # 人工修正任务按标题关键词定目标表,避免跨维度误判 - if 'chapter_outline' in title or '骨架' in title: - tbls = (QC_OUTPUT_TABLE[QC_TYPE_OUTLINE],) - elif 'doc_requirements' 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]) + 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: @@ -757,21 +757,26 @@ async def reconcile_project(sor, project_id, project_name=""): fresh = True break if not fresh: - verify_bad.append((getattr(t, 'id', ''), title[:50], dim or 'reqs_outline')) + 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,但对应产出表的最大更新时间早于任务创建时间" - "(只写了文件或根本没写库):\n%s\n\n" + "(只写了文件或根本没写库),已标 verify_failed:\n%s\n\n" "请处理:在对应管理页补录落库,或重跑该任务;落库后流程自动复审。" "处理完完成本任务。" - % "\n".join("- %s [%s] %s" % (tid, dm, ti) for tid, ti, dm in verify_bad)), + % "\n".join("- %s %s" % (tid, ti) for tid, ti in verify_bad)), assignee_role="owner.superuser") await sor.sqlExe("COMMIT", {}) - acts.append("escalated: 产出落库核验不通过(%d 个任务),抛人工" % len(verify_bad)) + acts.append("escalated: 产出落库核验不通过(%d 个任务标 verify_failed),抛人工" % len(verify_bad)) if created_w or created_r: return acts diff --git a/wwwroot/api/bid_task_io.dspy b/wwwroot/api/bid_task_io.dspy index 1775504..1bea2e5 100644 --- a/wwwroot/api/bid_task_io.dspy +++ b/wwwroot/api/bid_task_io.dspy @@ -18,6 +18,7 @@ _state_colors = { 'review': '#2196f3', 'rejected': '#dc2626', 'paused': '#f0a040', 'cancelled': '#64748b', 'qc_review': '#8b5cf6', 'qc_rejected': '#dc2626', + 'verify_failed': '#f0a040', } if not raw_id: diff --git a/wwwroot/api/bid_task_tree.dspy b/wwwroot/api/bid_task_tree.dspy index 3a44e0d..f1dede8 100644 --- a/wwwroot/api/bid_task_tree.dspy +++ b/wwwroot/api/bid_task_tree.dspy @@ -17,14 +17,14 @@ _state_icons = { 'failed': '\u274c', 'waiting': '\u26a0\ufe0f', 'submitted': '\u2b1c', 'review': '\U0001f440', 'rejected': '\U0001f6ab', 'paused': '\u23f8\ufe0f', 'cancelled': '\U0001f6d1', 'qc_review': '\U0001f50d', - 'qc_rejected': '\U0001f6ab', + 'qc_rejected': '\U0001f6ab', 'verify_failed': '\u26a0\ufe0f', } _state_zh = { 'completed': '已完成', 'approved': '已批准', 'running': '运行中', 'failed': '失败', 'waiting': '等待中', 'submitted': '已提交', 'review': '审核中', 'rejected': '已驳回', 'paused': '已暂停', 'cancelled': '已取消', 'qc_review': '审核中', - 'qc_rejected': 'QC不通过', + 'qc_rejected': 'QC不通过', 'verify_failed': '落库核验不通过', } # 投标产线角色分组(按业务流转顺序) @@ -153,26 +153,37 @@ async with DBPools().sqlorContext(dbname) as sor: # ── 维度/类型分组 ── if role == 'agent.tender_analyst': - # 各维度最新一轮 QC 结果(组节点标签用,消除"任务已批准≠QC通过"误导) - DIM_QC = {'scoring': 'scoring_items', 'quals': 'qualifications', - 'reqs_outline': 'doc_requirements', 'cost_benefit': 'cost_benefit'} - latest_qc = {} + # 各维度 QC 审核:最大轮次 + 最新一轮结果(组节点标签用, + # 轮次以审核记录为准——人工修正任务产出的第4轮也计入) + DIM_QCS = {'scoring': ('scoring_items',), 'quals': ('qualifications',), + 'reqs_outline': ('doc_requirements', 'chapter_outline'), + 'cost_benefit': ('cost_benefit',)} + qc_by_type = {} for rv in await sor.sqlExe( "SELECT qc_type, fit_score, passed, round FROM bid_qc_reviews " "WHERE project_id=${p}$ ORDER BY qc_type, round DESC", {"p": pid}) or []: qt = getattr(rv, 'qc_type', '') or '' - if qt not in latest_qc: - latest_qc[qt] = rv + if qt not in qc_by_type: + qc_by_type[qt] = rv groups = {} for t in role_tasks: key = (_params_of(t).get('analysis_dim') or '').strip() if not key: p = _params_of(t) - # 无维度:人工定向修正/补录任务单独成组(排除系统 qc_redo 重做任务), - # 其余归遗留全量解析 - if not p.get('qc_redo') and ('修正' in (getattr(t, 'title', '') or '') - or '补录' in (getattr(t, 'title', '') or '')): - key = '__manual__' + title = getattr(t, 'title', '') or '' + if p.get('qc_redo'): + key = '__legacy__' + elif '修正' in title or '补录' in title: + # 人工修正任务归入它修正的维度组(标题关键词定维度), + # 不再单独成组——第4轮产出在维度组里可见 + if 'scoring' in title or '评分项' in title: + key = 'scoring' + elif 'qual' in title or '资质' in title: + key = 'quals' + elif 'cost' in title or '成本收益' in title: + key = 'cost_benefit' + else: + key = 'reqs_outline' else: key = '__legacy__' groups.setdefault(key, []).append(t) @@ -181,10 +192,9 @@ async with DBPools().sqlorContext(dbname) as sor: dim_labels = { 'scoring': '评分项与得分规则', 'quals': '资质清单与要求', 'reqs_outline': '投标文件要求与章节骨架', 'cost_benefit': '成本收益分析', - '__legacy__': '全量解析(维度拆分前)', '__manual__': '人工定向修正', + '__legacy__': '全量解析(维度拆分前)', } - order = ['scoring', 'quals', 'reqs_outline', 'cost_benefit', - '__manual__', '__legacy__'] + order = ['scoring', 'quals', 'reqs_outline', 'cost_benefit', '__legacy__'] nodes = [] for key in order + [k for k in groups if k not in order]: if key not in groups: @@ -202,18 +212,27 @@ async with DBPools().sqlorContext(dbname) as sor: del_cnt.get(getattr(t, 'id', ''), 0) > 0), }) else: - final_state = getattr(ts[-1], 'state', '') or '' - lq = latest_qc.get(DIM_QC.get(key, ''), None) + # 轮次 = 该维度 QC 审核记录最大轮次(含人工修正触发的复审) + max_round = 0 + lq = None + for qt in DIM_QCS.get(key, ()): + rv = qc_by_type.get(qt) + if rv: + rr = int(getattr(rv, 'round', 0) or 0) + if rr > max_round: + max_round = rr + lq = rv + if not max_round: + max_round = len(ts) if lq: lq_passed = str(getattr(lq, 'passed', '0')) == '1' tail = "QC %.1f 分 %s" % (float(getattr(lq, 'fit_score', 0) or 0), '通过' if lq_passed else '未通过') else: - tail = "任务" + _state_zh.get(final_state, final_state) + tail = "任务" + _state_zh.get(getattr(ts[-1], 'state', '') or '', '') nodes.append({ "id": "grp:" + role + ":" + key, - "label": "{}({} {})最新[{}]".format( - label, len(ts), '个任务' if key == '__manual__' else '轮', tail), + "label": "{}({} 轮)最新[{}]".format(label, max_round, tail), "is_leaf": False, }) return json.dumps(nodes, ensure_ascii=False) @@ -388,8 +407,8 @@ async with DBPools().sqlorContext(dbname) as sor: }) return json.dumps(nodes, ensure_ascii=False) - # analyst 组:按任务创建序展开,轮次号与QC结果取自审核记录 - # (匹配规则:审核记录的 created_at 落在 [本任务updated_at, 下一任务created_at) 区间) + # analyst 组展开:与分组分支同逻辑过滤(人工修正任务按标题归维度组)。 + # 轮次:系统任务按时间区间匹配审核记录;人工修正任务取该维度最大轮次。 ts = [] for t in all_tasks: if (getattr(t, 'role', '') or '').strip() != role: @@ -399,57 +418,75 @@ async with DBPools().sqlorContext(dbname) as sor: k = (_params_of2(t).get('analysis_dim') or '').strip() if not k: p = _params_of2(t) - if not p.get('qc_redo') and ('修正' in (getattr(t, 'title', '') or '') - or '补录' in (getattr(t, 'title', '') or '')): - k = '__manual__' + title = getattr(t, 'title', '') or '' + if p.get('qc_redo'): + k = '__legacy__' + elif '修正' in title or '补录' in title: + if 'scoring' in title or '评分项' in title: + k = 'scoring' + elif 'qual' in title or '资质' in title: + k = 'quals' + elif 'cost' in title or '成本收益' in title: + k = 'cost_benefit' + else: + k = 'reqs_outline' else: k = '__legacy__' if k == key: ts.append(t) ts.sort(key=lambda x: str(getattr(x, 'created_at', '') or '')) - qc_type = DIM_QC.get(key, '') + # reqs_outline 维度含两类审核记录,合并取每轮最大 + DIM_QCS = {'scoring': ('scoring_items',), 'quals': ('qualifications',), + 'reqs_outline': ('doc_requirements', 'chapter_outline'), + 'cost_benefit': ('cost_benefit',)} rev_by_round = {} - if qc_type: + for qt in DIM_QCS.get(key, ()): for rv in await sor.sqlExe( "SELECT round, fit_score, passed, created_at FROM bid_qc_reviews " "WHERE project_id=${p}$ AND qc_type=${k}$ ORDER BY round ASC", - {"p": pid, "k": qc_type}) or []: - rev_by_round[int(getattr(rv, 'round', 0) or 0)] = rv + {"p": pid, "k": qt}) or []: + rr = int(getattr(rv, 'round', 0) or 0) + if rr not in rev_by_round: + rev_by_round[rr] = rv nodes = [] used_rounds = set() + max_round = max(rev_by_round) if rev_by_round else 0 for i, t in enumerate(ts): tid = getattr(t, 'id', '') st = getattr(t, 'state', '') or '' t_upd = str(getattr(t, 'updated_at', '') or '') + upd = t_upd[5:16] nxt_created = str(getattr(ts[i + 1], 'created_at', '') or '9999') if i + 1 < len(ts) else '9999' + is_manual = not (_params_of2(t).get('analysis_dim') or '').strip() and \ + not _params_of2(t).get('qc_redo') and \ + ('修正' in (getattr(t, 'title', '') or '') or '补录' in (getattr(t, 'title', '') or '')) matched = None - for rnd, rv in rev_by_round.items(): - if rnd in used_rounds: - continue - rv_t = str(getattr(rv, 'created_at', '') or '') - if t_upd and rv_t >= t_upd and rv_t < nxt_created: - matched = (rnd, rv) - break + if not is_manual: + for rnd, rv in rev_by_round.items(): + if rnd in used_rounds: + continue + rv_t = str(getattr(rv, 'created_at', '') or '') + if t_upd and rv_t >= t_upd and rv_t < nxt_created: + matched = (rnd, rv) + break if matched: rnd, rv = matched used_rounds.add(rnd) passed = str(getattr(rv, 'passed', '0')) == '1' sc = float(getattr(rv, 'fit_score', 0) or 0) qc_label = "→ QC %.1f 分 %s" % (sc, '通过' if passed else '未通过') - else: - rnd = i + 1 - qc_label = '' - upd = t_upd[5:16] - if key == '__manual__': - lbl = "%s [任务%s] %s" % (_state_icons.get(st, '\u2b1c'), - _state_zh.get(st, st), - (getattr(t, 'title', '') or '')[:40]) - else: lbl = "第%d轮 %s [任务%s] %s %s" % ( - rnd, _state_icons.get(st, '\u2b1c'), _state_zh.get(st, st), - qc_label or '→ 未审核/审核记录未关联', upd) + rnd, _state_icons.get(st, '\u2b1c'), _state_zh.get(st, st), qc_label, upd) + elif is_manual: + # 修正任务不挂轮次号(轮次以 QC 审核记录为准,见 QC 组) + lbl = "修正 %s [任务%s] %s" % (_state_icons.get(st, '\u2b1c'), + _state_zh.get(st, st), + (getattr(t, 'title', '') or '')[:40]) + else: + lbl = "第%d轮 %s [任务%s] → 未审核/审核记录未关联 %s" % ( + i + 1, _state_icons.get(st, '\u2b1c'), _state_zh.get(st, st), upd) nodes.append({ "id": tid, "label": lbl,