# -*- coding: utf-8 -*- """pipeline_service/kb_ingest_capability.py — PM 建议文件入知识库能力。 用户需求(2026-09-04):PM agent 有权建议「将什么文件加入什么知识库」并说明理由。 设计(遵守平台铁律:建议不直接执行,人工确认后才动数据;待办内完成不跳转;只发一次): 1. PM 调 rag_suggest_ingest(kb_id, file_path, reason): - 校验项目工作空间内文件真实存在(不允许编造文件); - 落 pipeline_kb_suggestions(status=pending)+ 创建 owner 人工待办(一次性); 2. owner 在待办详情里看到文件(可预览/下载)+ 理由,批准/驳回: - 批准 → 后端用「项目 owner 身份的 rag API key」调 /rag/api/doc_upload.dspy 真实入库 (与检索同一权限模型),建议置 approved; - 驳回 → 建议置 rejected,记录意见。 状态机:pending → approved/rejected(终态,幂等防重复上传)。 """ import json import logging import os from appPublic.uniqueID import getID from sqlor.dbpools import DBPools logger = logging.getLogger("pipeline.kb_ingest") DBNAME = "pipeline" S_PENDING = "pending" S_APPROVED = "approved" S_REJECTED = "rejected" T_KB_INGEST = "kb_ingest_confirm" def _get_db(): db = DBPools() if not db.databases: from appPublic.jsonConfig import getConfig config = getConfig() if config.databases: db.databases = config.databases return db def _rec_to_dict(rec): if rec is None: return {} try: return dict(rec) except (TypeError, ValueError): return {} async def suggest_kb_ingest(project_id, kb_id, file_path, reason, who="", agent_id="", task_id=""): """PM 建议将文件加入知识库。返回 (True, suggestion_id) 或 (False, 错误)。""" if not project_id or not kb_id or not file_path: return False, "project_id/kb_id/file_path 均为必填" if not (reason or '').strip(): return False, "必须说明入库理由" db = _get_db() async with db.sqlorContext(DBNAME) as sor: proj = await sor.sqlExe( "SELECT id, workspace_dir, org_id FROM sd_projects WHERE id=${pid}$", {"pid": project_id}) await sor.sqlExe("COMMIT", {}) if not proj: return False, "项目不存在" workspace_dir = str(getattr(proj[0], "workspace_dir", "") or "") # 相对路径解析:兼容「项目目录相对路径」「机构工作空间相对路径 # (projects/{项目}/…)」「项目内子目录相对路径」三种写法 path = file_path.strip() if not os.path.isabs(path): candidates = [] if workspace_dir: candidates.append(os.path.join(workspace_dir, path)) projects_dir = os.path.dirname(workspace_dir) # projects/ 层 candidates.append(os.path.join(projects_dir, path)) space_dir = os.path.dirname(projects_dir) # 机构工作空间层 candidates.append(os.path.join(space_dir, path)) full = "" for c in candidates: if os.path.isfile(c): full = c break if not full: return False, ("文件不存在(已查项目工作空间): " + path + "。请先确认文件已产出再建议入库") else: full = path if not os.path.isfile(full): return False, "文件不存在: " + full # 幂等:同项目同文件同库未决建议不重复发待办(铁律:同事待办只发一次) dup = await sor.sqlExe( "SELECT id FROM pipeline_kb_suggestions WHERE project_id=${p}$ AND kb_id=${k}$ " "AND file_path=${f}$ AND status='pending' LIMIT 1", {"p": project_id, "k": kb_id, "f": full}) await sor.sqlExe("COMMIT", {}) if dup: return False, "该文件已有待确认的入库建议(" + str(getattr(dup[0], 'id', ''))[:8] + "…),勿重复提交" # 知识库可见性校验(用 owner 权限查,防止建议一个 owner 看不到的库) from .rag_client import rag_kb_list kbdata, kberr = await rag_kb_list(project_id) if kberr: return False, "校验知识库失败: " + kberr kbs = (kbdata or {}).get("kbs") if isinstance(kbdata, dict) else [] known = {str(k.get("id") or "") for k in kbs if isinstance(k, dict)} if known and str(kb_id) not in known: return False, "知识库不存在或项目 owner 无权检索该库: " + str(kb_id) sid = getID() await sor.C("pipeline_kb_suggestions", { "id": sid, "project_id": project_id, "kb_id": str(kb_id), "file_path": full, "file_name": os.path.basename(full), "reason": reason.strip(), "status": S_PENDING, "suggested_by": who or agent_id or "agent.pm", "task_id": task_id or "", "result_comment": "", }) await sor.sqlExe("COMMIT", {}) # 解析真人 owner 作为待办处理人 from .rag_client import resolve_project_owner owner, oerr = await resolve_project_owner(project_id) if oerr or not owner: owner = "" # 待办(一次一条;suggestion_id 塞进 form_schema 供待办详情读取, # fields 为空=无表单,只有批准/驳回按钮) from .human_task_capability import create_human_task desc = ("PM 建议将文件加入知识库,请审阅理由后决定。\n" "文件:" + os.path.basename(full) + "\n" "理由:" + reason.strip()) ok, ht = await create_human_task( project_id, "知识库入库建议:" + os.path.basename(full), description=desc, task_type=T_KB_INGEST, assignee_id=owner or None, assignee_role=None if owner else "owner", created_by="agent.pm", task_id=task_id or "", form_schema={"fields": [], "suggestion_id": sid}) if not ok: logger.warning("kb_ingest todo create failed: %s", ht) return True, sid async def decide_kb_ingest(suggestion_id, approve, operator_id, comment=""): """owner 决策入库建议。批准 → 真实上传(rag API,owner 身份)。 返回 (True, 消息) 或 (False, 错误)。 """ if not suggestion_id: return False, "缺少建议 id" if not operator_id: return False, "未登录" db = _get_db() async with db.sqlorContext(DBNAME) as sor: recs = await sor.sqlExe( "SELECT * FROM pipeline_kb_suggestions WHERE id=${i}$", {"i": suggestion_id}) await sor.sqlExe("COMMIT", {}) if not recs: return False, "建议不存在" sg = _rec_to_dict(recs[0]) if sg.get("status") != S_PENDING: return False, "建议已处理(当前状态: " + str(sg.get("status")) + ")" project_id = sg.get("project_id", "") # 操作者须为项目真人 owner(与待办处理人一致) from .rag_client import resolve_project_owner owner, oerr = await resolve_project_owner(project_id) if owner and str(operator_id) != str(owner): return False, "仅项目 owner 可批准/驳回入库建议" # 关联人工待办置为已处理(避免待办重复出现)——用 form_schema 里唯一的 # suggestion_id 精确匹配(description LIKE 文件名会误伤同名文件的待办) await sor.sqlExe( "UPDATE pipeline_human_tasks SET status='done', submitted_by=${o}$, submitted_at=NOW(), " "result_data=${rd}$ WHERE task_type='kb_ingest_confirm' AND status='pending' " "AND form_schema LIKE ${pat}$", {"o": operator_id, "rd": json.dumps({"suggestion_id": suggestion_id, "approve": bool(approve)}, ensure_ascii=False), "pat": "%\"suggestion_id\": \"" + str(suggestion_id) + "\"%"}) await sor.sqlExe("COMMIT", {}) if not approve: await sor.sqlExe( "UPDATE pipeline_kb_suggestions SET status='rejected', result_comment=${c}$, " "decided_by=${o}$, decided_at=NOW() WHERE id=${i}$", {"c": comment or "owner 驳回", "o": operator_id, "i": suggestion_id}) await sor.sqlExe("COMMIT", {}) return True, "已驳回,不会入库" # 批准 → 真实上传(API 模式,owner 身份) from .rag_client import rag_doc_upload data, err = await rag_doc_upload(project_id, sg.get("kb_id", ""), sg.get("file_path", ""), sg.get("file_name", "")) if err: await sor.sqlExe( "UPDATE pipeline_kb_suggestions SET result_comment=${c}$ WHERE id=${i}$", {"c": "上传失败: " + err, "i": suggestion_id}) await sor.sqlExe("COMMIT", {}) return False, "上传失败: " + err doc_id = str((data or {}).get("doc_id", "") or "") await sor.sqlExe( "UPDATE pipeline_kb_suggestions SET status='approved', result_comment=${c}$, " "decided_by=${o}$, decided_at=NOW() WHERE id=${i}$", {"c": "已入库,doc_id=" + doc_id, "o": operator_id, "i": suggestion_id}) await sor.sqlExe("COMMIT", {}) return True, "已批准并入库(文档号 " + (doc_id or "-") + ",后台解析中)"