fix: 修正任务归维度组(轮次以QC审核记录为准含第4轮);落库核验不通过标verify_failed不再撒谎approved

This commit is contained in:
yumoqing 2026-09-03 18:34:57 +08:00
parent 2223112e0e
commit 3b718a0869
3 changed files with 110 additions and 67 deletions

View File

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

View File

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

View File

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