feat: QC不通过新增qc_rejected状态+失败轮产出快照bid_analysis_history(不清空历史)+存量回填

This commit is contained in:
yumoqing 2026-09-03 17:31:36 +08:00
parent 7b85943ffa
commit 2d9c1264e3
5 changed files with 124 additions and 5 deletions

View File

@ -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_rejectedQC 判定不通过计入防止对账器对同一轮反复派重做任务
"""
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:

View File

@ -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))

View File

@ -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_historyQC 不通过清空前的产出快照表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
# 末位任务看该维度最新一轮 QCpassed=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()

View File

@ -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:

View File

@ -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不通过',
}
# 投标产线角色分组(按业务流转顺序)