diff --git a/pipeline_bidding/bid_flow.py b/pipeline_bidding/bid_flow.py index 5e04140..3b493d7 100644 --- a/pipeline_bidding/bid_flow.py +++ b/pipeline_bidding/bid_flow.py @@ -155,9 +155,13 @@ async def _dim_done_count(sor, project_id, dim, since=''): async def _dim_redo_done_count(sor, project_id, dim, since=''): - """某维度重做(qc_redo=1)已完成任务数(不计入首跑重试上限)。""" + """某维度重做(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') " + "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: diff --git a/pipeline_bidding/bid_qc_capability.py b/pipeline_bidding/bid_qc_capability.py index c3844ee..baba630 100644 --- a/pipeline_bidding/bid_qc_capability.py +++ b/pipeline_bidding/bid_qc_capability.py @@ -204,6 +204,38 @@ async def start_qc(qc_type, project_id="", who=None, agent_id=None, task_id=""): }, ensure_ascii=False, default=str) +async def _snapshot_outputs(sor, pid, qc_type, round_no, task_id, fit_score, improvement): + """QC 不通过清空产出前,把本轮产出整表快照到 bid_analysis_history(幂等:同轮重复快照跳过)。 + + 保留失败轮次的输出,供任务树/人工追溯;不清空历史。 + """ + tbl = QC_OUTPUT_TABLE.get(qc_type or "", "") + if not tbl: + return 0 + exist = await sor.sqlExe( + "SELECT COUNT(*) AS c FROM bid_analysis_history WHERE project_id=${p}$ " + "AND qc_type=${t}$ AND round=${r}$", {"p": pid, "t": qc_type, "r": round_no}) + if exist and to_int(getattr(exist[0], 'c', 0), 0) > 0: + return 0 + rows = await sor.sqlExe( + "SELECT * FROM " + tbl + " WHERE project_id=${p}$", {"p": pid}) or [] + n = 0 + for r in rows: + d = rec_to_dict(r) + d.pop('id', None) + await sor.sqlExe( + "INSERT INTO bid_analysis_history " + "(id, project_id, qc_type, round, task_id, fit_score, improvement, payload, created_at) " + "VALUES (${id}$, ${p}$, ${t}$, ${r}$, ${tid}$, ${sc}$, ${imp}$, ${pl}$, NOW())", + {"id": new_id(), "p": pid, "t": qc_type, "r": round_no, "tid": task_id or "", + "sc": fit_score, "imp": (improvement or "")[:20000], + "pl": json.dumps(d, ensure_ascii=False, default=str)[:200000]}) + n += 1 + await sor.sqlExe("COMMIT", {}) + logger.info("bid_qc: snapshot %s round=%d rows=%d", qc_type, round_no, n) + return n + + async def finish_qc(review_id, fit_score, improvement="", comments="", who=None, agent_id=None, task_id=""): """落分并判定:得分 > 通过分 → 该类产出放行;否则清空产出物 + 触发解析重做。 @@ -269,14 +301,24 @@ async def finish_qc(review_id, fit_score, improvement="", comments="", return True, ("QC 不通过且已达审核轮次上限 %d → 已抛人工介入(阻塞门禁)。" % th["qc_max_round"]) - # ── 清空该类产出物(章节骨架退回时章节表一并清空),触发解析重做 ── + # ── 不通过且未达轮次上限:先快照本轮产出(保留失败轮次的输出), + # 再把该轮分析任务标 qc_rejected(任务状态与QC判定一致,不再显示误导的"已批准"), + # 最后清空产出触发重做。失败轮次的输入(qc_improvements 在重做任务 params) + # 与输出(快照表)全程可查。── tbl = QC_OUTPUT_TABLE.get(qc_type or "", "") if not tbl: - return False, "QC 记录类型非法: %s" % qc_type + return False, "QC 记录类型非法: %s" + await _snapshot_outputs(sor, pid, qc_type, to_int(review.get("round"), 1), + review.get("task_id") or "", sc, improvement or "") + if review.get("task_id"): + await sor.sqlExe( + "UPDATE pipeline_tasks SET state='qc_rejected', updated_at=NOW() " + "WHERE id=${t}$", {"t": review.get("task_id")}) await sor.sqlExe("DELETE FROM " + tbl + " WHERE project_id=${p}$", {"p": pid}) await sor.sqlExe("COMMIT", {}) cleared = tbl if qc_type != QC_TYPE_OUTLINE else "bid_chapters(章节骨架)" - return True, ("QC 不通过:%s 契合度 %.2f 未高于 %.2f。已清空该类产出物(%s)," + return True, ("QC 不通过:%s 契合度 %.2f 未高于 %.2f。本轮产出已快照保留(bid_analysis_history)," + "任务标记 qc_rejected,产出物(%s)已清空," "系统将在下一轮对账重派解析任务(携带改进意见)重做后重审。" % (qc_type, sc, th["qc_pass_score"], cleared)) diff --git a/scripts/import_init_bidding.py b/scripts/import_init_bidding.py index 3b8e456..d697416 100644 --- a/scripts/import_init_bidding.py +++ b/scripts/import_init_bidding.py @@ -15,6 +15,13 @@ sys.path.insert(0, ROOT_DIR) from appPublic.jsonConfig import getConfig from appPublic.aes import aes_decode_b64 +import pymysql + + +def pymysql_connect(kw, pwd): + return pymysql.connect(host=str(kw.host), port=int(kw.port), + user=str(kw.user), password=pwd, + database=str(kw.db), charset='utf8mb4') def mysql_exec(kw, pwd, sql): @@ -70,6 +77,69 @@ def main(): "UPDATE pipelines SET description='%s' WHERE id='bidding_general'" % desc) print("pipelines.bidding_general: description updated") + # ── 5. bid_analysis_history:QC 不通过清空前的产出快照表(2026-09-03,保留失败轮次输出)── + mysql_exec(kw, pwd, """ + CREATE TABLE IF NOT EXISTS bid_analysis_history ( + id VARCHAR(32) NOT NULL, + project_id VARCHAR(32) NOT NULL, + qc_type VARCHAR(32) NOT NULL, + round INT NOT NULL DEFAULT 0, + task_id VARCHAR(32) NOT NULL DEFAULT '', + fit_score DECIMAL(5,2) NOT NULL DEFAULT 0, + improvement TEXT, + payload LONGTEXT, + created_at DATETIME NOT NULL, + PRIMARY KEY (id), + KEY idx_bah_project (project_id, qc_type, round) + ) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4""") + print("bid_analysis_history: ensured") + + # ── 6. 存量任务状态回填:approved 但实际触发过重做的分析任务 → qc_rejected(幂等)── + # 规则:同维度内其后还有任务 ⇒ 该轮失败触发过重做 ⇒ qc_rejected; + # 末位任务看该维度最新一轮 QC:passed=0 ⇒ qc_rejected。只改 approved/completed。 + import json as _json + conn = pymysql_connect(kw, pwd) + cur = conn.cursor() + cur.execute("SELECT DISTINCT tenant_id FROM pipeline_tasks WHERE role='agent.tender_analyst'") + pids = [r[0] for r in cur.fetchall()] + DIM_QC = {'scoring': ('scoring_items',), 'quals': ('qualifications',), + 'reqs_outline': ('doc_requirements', 'chapter_outline'), + 'cost_benefit': ('cost_benefit',)} + n_fix = 0 + for pid in pids: + cur.execute("SELECT id, state, params, created_at FROM pipeline_tasks " + "WHERE tenant_id=%s AND role='agent.tender_analyst' ORDER BY created_at", (pid,)) + tasks = cur.fetchall() + groups = {} + for tid, st, params, cat in tasks: + try: + dim = (_json.loads(params if isinstance(params, str) else (params or '{}')) or {}).get('analysis_dim') or '' + except Exception: + dim = '' + if dim: + groups.setdefault(dim, []).append((tid, st, cat)) + for dim, ts in groups.items(): + for i, (tid, st, cat) in enumerate(ts): + new_state = None + if i < len(ts) - 1 and st in ('approved', 'completed'): + new_state = 'qc_rejected' + elif i == len(ts) - 1 and st in ('approved', 'completed'): + qcs = DIM_QC.get(dim, ()) + if qcs: + fmt = ','.join(['%s'] * len(qcs)) + cur.execute("SELECT passed FROM bid_qc_reviews WHERE project_id=%s " + "AND qc_type IN (" + fmt + ") ORDER BY round DESC LIMIT 1", + (pid,) + qcs) + rr = cur.fetchone() + if rr and str(rr[0]) != '1': + new_state = 'qc_rejected' + if new_state: + cur.execute("UPDATE pipeline_tasks SET state=%s WHERE id=%s", (new_state, tid)) + n_fix += 1 + conn.commit() + conn.close() + print("backfill qc_rejected: %d tasks" % n_fix) + if __name__ == "__main__": main() diff --git a/wwwroot/api/bid_task_io.dspy b/wwwroot/api/bid_task_io.dspy index 727e4c0..1775504 100644 --- a/wwwroot/api/bid_task_io.dspy +++ b/wwwroot/api/bid_task_io.dspy @@ -17,6 +17,7 @@ _state_colors = { 'waiting': '#f0a040', 'submitted': '#64748b', 'review': '#2196f3', 'rejected': '#dc2626', 'paused': '#f0a040', 'cancelled': '#64748b', 'qc_review': '#8b5cf6', + 'qc_rejected': '#dc2626', } if not raw_id: diff --git a/wwwroot/api/bid_task_tree.dspy b/wwwroot/api/bid_task_tree.dspy index 06d9799..042a1e6 100644 --- a/wwwroot/api/bid_task_tree.dspy +++ b/wwwroot/api/bid_task_tree.dspy @@ -17,12 +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', } _state_zh = { 'completed': '已完成', 'approved': '已批准', 'running': '运行中', 'failed': '失败', 'waiting': '等待中', 'submitted': '已提交', 'review': '审核中', 'rejected': '已驳回', 'paused': '已暂停', 'cancelled': '已取消', 'qc_review': '审核中', + 'qc_rejected': 'QC不通过', } # 投标产线角色分组(按业务流转顺序)