diff --git a/pipeline_bidding/bid_ability.py b/pipeline_bidding/bid_ability.py index e96312d..920d842 100644 --- a/pipeline_bidding/bid_ability.py +++ b/pipeline_bidding/bid_ability.py @@ -76,6 +76,8 @@ BID_TOOLS = [ parameters={}, category="agent"), ToolDefinition(name="advance_bid", description="立即触发投标流程状态对账(补派应有的角色任务,通常无需手动调用)", parameters={}, category="agent"), + ToolDefinition(name="reset_task_retry", description="人工处理完失败任务后恢复重试:把失败/挂起的任务清零重试计数并重新执行", + parameters={"task_id": "任务ID(从diagnose_bid看到的失败任务)"}, category="agent"), ] # ══════════════════ 产线 prompt 片段 ══════════════════ @@ -380,6 +382,19 @@ async def _h_advance_bid(sor, p, ctx): return "对账结果:\n" + "\n".join("- " + a for a in (acts or ["无动作"])) +async def _h_reset_task_retry(sor, p, ctx): + pid, err = _need_project(ctx) + if err: + return err + tid = str(p.get("task_id") or "").strip() + if not tid: + return "ERROR: 缺少 task_id" + from pipeline_service.task_capability import reset_task_retry + ok, msg = await reset_task_retry(tid, pid, who="agent.main_agent", + agent_id=ctx.get("agent_id", "") or "") + return _fmt(ok, "任务已恢复,对账器将重新派发" if ok else msg) + + BID_HANDLERS = { "create_bid_project": _h_create_bid_project, "import_tender_file": _h_import_tender_file, @@ -400,6 +415,7 @@ BID_HANDLERS = { "bid_score_report": _h_bid_score_report, "diagnose_bid": _h_diagnose_bid, "advance_bid": _h_advance_bid, + "reset_task_retry": _h_reset_task_retry, } diff --git a/pipeline_bidding/bid_analysis_capability.py b/pipeline_bidding/bid_analysis_capability.py index 35de7d2..359da79 100644 --- a/pipeline_bidding/bid_analysis_capability.py +++ b/pipeline_bidding/bid_analysis_capability.py @@ -6,6 +6,7 @@ import json import logging +import os from .bid_common import ( get_db, new_id, rec_to_dict, rows_to_dicts, to_float, to_int, json_loads, @@ -78,7 +79,11 @@ async def list_tender_files(project_id, extract_status=""): async def read_tender_file(file_id, offset=0, length=20000): - """读招标文件正文(分段读,供 LLM 分析)。正文为空时尝试从磁盘提取。""" + """读招标文件正文(分段读,供 LLM 分析)。正文为空时尝试从磁盘提取。 + + 磁盘提取支持相对路径:裸文件名/相对路径按项目根目录解析 + (上传时可能只存了文件名,2026-09-01 修复)。 + """ db, dbname = get_db() async with db.sqlorContext(dbname) as sor: recs = await sor.sqlExe( @@ -90,11 +95,14 @@ async def read_tender_file(file_id, offset=0, length=20000): d = rec_to_dict(recs[0]) txt = d.get("content_text") or "" if not txt.strip() and d.get("file_path"): - txt = _extract_file_text(d["file_path"]) + pdir = await _project_dir_of(sor, d.get("project_id", "")) + real = _resolve_file_path(pdir, d["file_path"]) + txt = _extract_file_text(real) if real else "" if txt: await sor.sqlExe( - "UPDATE bid_tender_files SET content_text=${c}$, updated_at=NOW() " - "WHERE id=${fid}$", {"c": txt[:2000000], "fid": file_id}) + "UPDATE bid_tender_files SET content_text=${c}$, file_path=${p}$, " + "updated_at=NOW() WHERE id=${fid}$", + {"c": txt[:2000000], "p": real, "fid": file_id}) await sor.sqlExe("COMMIT", {}) if not txt.strip(): return ("招标文件正文为空(file_path=%s)。请确认文件已上传且格式为 docx/pdf/txt。" @@ -142,6 +150,93 @@ def _extract_file_text(path): return "" +async def _project_dir_of(sor, project_id): + """项目根目录(供相对路径解析)。解析口径与引擎上传落盘一致: + sd_projects.workspace_dir 优先;为空时按 {base}/{org}/{pipeline}/{项目名} 兜底。 + """ + if not project_id: + return "" + try: + recs = await sor.sqlExe( + "SELECT name, org_id, pipeline_id, workspace_dir FROM sd_projects " + "WHERE id=${p}$ LIMIT 1", {"p": project_id}) + if not recs: + return "" + r = recs[0] + ws = (getattr(r, 'workspace_dir', '') or '').strip() + if ws.startswith('/'): + return ws + from pipeline_service.workspace import get_workspace_base, build_workspace_path + base = await get_workspace_base(sor) + name = (getattr(r, 'name', '') or '').strip() + if not base or not name: + return "" + return build_workspace_path(base, getattr(r, 'org_id', '0') or '0', + getattr(r, 'pipeline_id', '') or 'general', name) + except Exception as e: + logger.warning("project dir resolve failed pid=%s: %s", project_id, str(e)[:120]) + return "" + + +def _resolve_file_path(project_dir, file_path): + """把存储的 file_path 解析成磁盘绝对路径。 + - 绝对路径且存在:原样返回; + - 相对/裸文件名:依次尝试 项目根目录、项目根目录下的 docs/、进程 cwd,存在即返回; + - 都找不到返回 ''。 + """ + fp = str(file_path or "").strip() + if not fp: + return "" + if os.path.isabs(fp): + return fp if os.path.isfile(fp) else "" + cands = [] + if project_dir: + cands.append(os.path.join(project_dir, fp)) + cands.append(os.path.join(project_dir, "docs", fp)) + cands.append(fp) + for c in cands: + if os.path.isfile(c): + return os.path.abspath(c) + return "" + + +async def _best_tender_record(sor, project_id): + """选最优招标文件记录(修复 2026-09-01「解析读旧截断稿」根因)。 + + 一个项目可能有多条招标文件记录(重复上传/导入)。选择规则: + **可得正文最多者胜**——正文取 len(content_text),磁盘文件取文件大小, + 两者取大;并列时取最新。返回 dict 附带 _resolved_path(磁盘绝对路径提示)。 + + 旧实现 ORDER BY created_at LIMIT 1 固定取最老记录——最老的往往是早期 + 截断摘要,完整文件后来才传,导致解析永远读不到完整文件。 + 只按「有正文优先」也不够:截断稿有正文、完整稿只在磁盘上,会误选截断稿。 + """ + recs = await sor.sqlExe( + "SELECT id, project_id, file_name, file_path, content_text, created_at " + "FROM bid_tender_files WHERE project_id=${pid}$ AND file_type='tender_doc' " + "ORDER BY created_at DESC", {"pid": project_id}) + if not recs: + return {} + rows = rows_to_dicts(recs, limit=50) + pdir = await _project_dir_of(sor, project_id) + best, best_score, best_real = {}, -1, "" + for r in rows: + score = len((r.get("content_text") or "").strip()) + real = "" + if (r.get("file_path") or "").strip(): + real = _resolve_file_path(pdir, r["file_path"]) + if real: + try: + score = max(score, os.path.getsize(real)) + except OSError: + pass + if score > best_score: + best, best_score, best_real = r, score, real + if best and best_real: + best["_resolved_path"] = best_real + return best + + async def mark_file_extracted(file_id, note="", who=None, agent_id=None): """标记招标文件已抽取完成。""" db, dbname = get_db() @@ -338,25 +433,31 @@ async def extract_from_file(project_id, file_id="", model_name="", offset=0, recs = await sor.sqlExe( "SELECT id, content_text, file_path, file_name FROM bid_tender_files " "WHERE id=${fid}$", {"fid": file_id}) + await sor.sqlExe("COMMIT", {}) + if not recs: + return False, "项目下没有招标文件(bid_tender_files),无法抽取" + d = rec_to_dict(recs[0]) else: - recs = await sor.sqlExe( - "SELECT id, content_text, file_path, file_name FROM bid_tender_files " - "WHERE project_id=${pid}$ AND file_type='tender_doc' ORDER BY created_at LIMIT 1", - {"pid": project_id}) - await sor.sqlExe("COMMIT", {}) - if not recs: - return False, "项目下没有招标文件(bid_tender_files),无法抽取" - d = rec_to_dict(recs[0]) + # 选最优记录:有正文优先、其次磁盘有文件,同类取最新 + # (2026-09-01 修复:旧实现固定取最老记录,完整文件后来上传时永远读不到) + d = await _best_tender_record(sor, project_id) + if not d: + return False, "项目下没有招标文件(bid_tender_files),无法抽取" txt = d.get("content_text") or "" if not txt.strip() and d.get("file_path"): - txt = _extract_file_text(d["file_path"]) + pdir = await _project_dir_of(sor, project_id) + real = d.get("_resolved_path") or _resolve_file_path(pdir, d["file_path"]) + txt = _extract_file_text(real) if real else "" if txt: await sor.sqlExe( - "UPDATE bid_tender_files SET content_text=${c}$ WHERE id=${fid}$", - {"c": txt[:2000000], "fid": d["id"]}) + "UPDATE bid_tender_files SET content_text=${c}$, file_path=${p}$ " + "WHERE id=${fid}$", + {"c": txt[:2000000], "p": real, "fid": d["id"]}) await sor.sqlExe("COMMIT", {}) if not txt.strip(): - return False, "招标文件正文为空,无法抽取(file=%s)" % d.get("file_name", "") + return False, ("招标文件正文为空,无法抽取(file=%s,file_path=%s);" + "请确认文件已上传到项目目录" % + (d.get("file_name", ""), d.get("file_path", ""))) o = to_int(offset, 0) seg = txt[o:o + to_int(length, 30000)] data, err = await llm_json(EXTRACT_PROMPT.replace("__CONTENT__", seg), diff --git a/pipeline_bidding/bid_flow.py b/pipeline_bidding/bid_flow.py index 521aa3b..bd40262 100644 --- a/pipeline_bidding/bid_flow.py +++ b/pipeline_bidding/bid_flow.py @@ -53,12 +53,85 @@ async def _count(sor, sql, p): return to_int(getattr(recs[0], "c", 0) if recs else 0) -async def _done_task_count(sor, project_id, role): - return await _count(sor, - "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": project_id, "r": role}) +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) + + +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"): + 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): @@ -131,15 +204,6 @@ async def _qc_collect_improvements(sor, project_id): return items -async def _redo_done_count(sor, project_id): - """已完成的 QC 退回重做任务数(解析重做轮次计数,防无限重派)。""" - return await _count( - sor, "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": project_id, "r": R_ANALYST, "kw": "%qc_redo%"}) - - async def reconcile_project(sor, project_id, project_name=""): """对单个投标项目做一次状态对账。返回动作说明列表(幂等,可反复调用)。""" acts = [] @@ -229,6 +293,17 @@ async def reconcile_project(sor, project_id, project_name=""): {"p": project_id, "r": R_ANALYST, "kw": "%qc_redo%"}) attempts = await _done_task_count(sor, project_id, R_ANALYST) n_redo_done = await _redo_done_count(sor, project_id) + # 自动恢复(2026-09-01):解析升级后来了**新的**招标文件(上传时间晚于 + # 最后一次解析完成)→ 清掉旧文件残留产出与旧 QC 意见,重试计数从新文件 + # 时刻重起。否则项目永久卡死:旧计数 ≥ 上限,永远不再派解析 + # (中传项目「走不下去」的另一半根因)。 + _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: 检测到新招标文件,已清空旧解析残留,重新解析") + attempts = await _done_task_count(sor, project_id, R_ANALYST, since=_ft) + n_redo_done = await _redo_done_count(sor, project_id, since=_ft) first_attempts = max(0, attempts - n_redo_done) qc_imp = await _qc_collect_improvements(sor, project_id) if qc_imp or n_redo_open > 0: