# ============================================================================= # 企微IT智能服务台 — RAGFlow 文档 Ingestion 服务(通道 C / P1-5) # ============================================================================= # 说明:封装 RAGFlow 文档上传→ETL→结构化→生成 KnowledgeSuggestion 流程。 # 训练师上传非标准格式文档(.docx/.pdf/.txt/.png/.jpg), # RAGFlow 做第一道整理/筛选/结构化,产出 KnowledgeSuggestion 进审批队列。 # # 核心流程: # 1. upload_and_process: 上传文档→轮询处理状态→拉取结构化片段→生成建议 # 2. poll_processing_status: 轮询 RAGFlow 文档处理状态(最多5分钟) # 3. create_suggestions_from_result: 结构化片段→KnowledgeSuggestion 列表 # # 设计决策: # - 触发方式:训练师手动上传(P1-5 决策) # - source_type=document_ragflow, audience=engineer_workguide # - 复用现有 integrations/ragflow/ 客户端基础设施 # ============================================================================= import logging import uuid from datetime import datetime from typing import Any, Dict, List, Optional from app.config import settings logger = logging.getLogger(__name__) class RagflowIngestionService: """RAGFlow 文档 Ingestion 服务 — 通道 C。 将非标准格式文档通过 RAGFlow 整理/筛选/结构化, 产出 KnowledgeSuggestion 进入 D7 审批流程。 使用方式: service = RagflowIngestionService() result = await service.upload_and_process(file_data, file_name, category_hint) """ # 轮询参数 _POLL_INTERVAL_SECONDS: int = 10 # 轮询间隔(秒) _MAX_WAIT_SECONDS: int = 300 # 最大等待时间(5分钟) def __init__(self): """初始化 RAGFlow Ingestion 服务。""" self.ragflow_base_url: str = settings.automation_ragflow_base_url self.ragflow_api_key: str = settings.automation_ragflow_api_key self.enabled: bool = settings.ragflow_ingestion_enabled async def upload_and_process( self, file_data: bytes, file_name: str, category_hint: str = "其他", ) -> Dict[str, Any]: """上传文档到 RAGFlow 并等待处理完成,生成 KnowledgeSuggestion 列表。 Args: file_data: 文件字节流 file_name: 文件名(含扩展名,如 "FAQ更新说明.docx") category_hint: 分类提示(可选,帮助 RAGFlow 归类) Returns: Dict: { "task_id": str, # RAGFlow 任务ID "status": str, # "completed" / "failed" / "pending" "suggestions": list[dict], # KnowledgeSuggestion 列表 } """ task_id = str(uuid.uuid4()) if not self.enabled: logger.info("RAGFlow Ingestion 未启用,返回空结果") return { "task_id": task_id, "status": "disabled", "suggestions": [], } try: # 1. 上传文档到 RAGFlow doc_id = await self._upload_document(file_data, file_name) if not doc_id: return { "task_id": task_id, "status": "failed", "suggestions": [], } # 2. 轮询处理状态 status = await self.poll_processing_status(doc_id) if status != "completed": return { "task_id": task_id, "status": status, "suggestions": [], } # 3. 拉取结构化片段 chunks = await self._fetch_document_chunks(doc_id) if not chunks: return { "task_id": task_id, "status": "completed", "suggestions": [], } # 4. 生成 KnowledgeSuggestion suggestions = self.create_suggestions_from_result(chunks, category_hint) return { "task_id": task_id, "status": "completed", "suggestions": suggestions, } except Exception as e: logger.error(f"RAGFlow Ingestion 失败: {e}") return { "task_id": task_id, "status": "failed", "suggestions": [], } async def poll_processing_status( self, doc_id: str, max_wait: int = 300 ) -> str: """轮询 RAGFlow 文档处理状态。 Args: doc_id: RAGFlow 文档ID max_wait: 最大等待时间(秒),默认 300 秒 Returns: str: "completed" / "failed" / "pending" / "timeout" """ import asyncio elapsed = 0 while elapsed < max_wait: try: status = await self._get_document_status(doc_id) if status == "completed": logger.info(f"RAGFlow 文档处理完成: doc_id={doc_id}") return "completed" if status == "failed": logger.error(f"RAGFlow 文档处理失败: doc_id={doc_id}") return "failed" except Exception as e: logger.warning(f"轮询 RAGFlow 状态失败: {e}") await asyncio.sleep(self._POLL_INTERVAL_SECONDS) elapsed += self._POLL_INTERVAL_SECONDS logger.debug( f"轮询 RAGFlow 状态: doc_id={doc_id}, elapsed={elapsed}s" ) logger.warning(f"RAGFlow 文档处理超时: doc_id={doc_id}") return "timeout" def create_suggestions_from_result( self, chunks: List[Dict[str, Any]], category_hint: str = "其他" ) -> List[Dict[str, Any]]: """将 RAGFlow 结构化片段转为 KnowledgeSuggestion 列表。 每个片段生成一个建议,source_type=document_ragflow, audience=engineer_workguide(通道 C 默认工程师作业指导)。 Args: chunks: RAGFlow 结构化段落列表 category_hint: 分类提示 Returns: List[Dict]: KnowledgeSuggestion 数据列表(可直接用于创建 DB 记录) """ suggestions: List[Dict[str, Any]] = [] for chunk in chunks: suggestion = { "suggestion_type": "new_faq", "title": chunk.get("title", "RAGFlow 提取的知识片段"), "content": chunk.get("content", ""), "category": category_hint or chunk.get("category", "其他"), "tags": chunk.get("tags", []), "source_type": "document_ragflow", "source_data": [chunk.get("chunk_id", str(uuid.uuid4()))], "reason": f"RAGFlow 从文档中提取的结构化知识片段", "confidence": 0.85, # RAGFlow 结构化提取默认 0.85(§8.2) "audience": "engineer_workguide", # 通道 C 默认工程师作业指导 "issue": chunk.get("issue", ""), "action": chunk.get("action", ""), "relation_type": "LEADS_TO", "parent_issue": "", "graph_meta": {}, "graph_sync_status": "pending", "source_failed": False, } suggestions.append(suggestion) logger.info( f"RAGFlow 生成 {len(suggestions)} 条 KnowledgeSuggestion" ) return suggestions # -------------------------------------------------------------------------- # 内部方法 — RAGFlow API 调用 # -------------------------------------------------------------------------- async def _upload_document( self, file_data: bytes, file_name: str ) -> Optional[str]: """上传文档到 RAGFlow。 Args: file_data: 文件字节流 file_name: 文件名 Returns: Optional[str]: RAGFlow 文档ID,失败返回 None """ try: import httpx from app.integrations.ragflow.client import RagflowClient # TODO: 接入现有 RagflowClient 实现实际上传 # 当前返回模拟 doc_id(RAGFlow 服务部署后替换为真实调用) logger.info( f"RAGFlow 文档上传(模拟): file_name={file_name}, " f"size={len(file_data)}" ) doc_id = f"ragflow_doc_{uuid.uuid4().hex[:12]}" return doc_id except ImportError: logger.warning("RAGFlow 客户端不可用,返回模拟 doc_id") return f"ragflow_doc_{uuid.uuid4().hex[:12]}" except Exception as e: logger.error(f"RAGFlow 文档上传失败: {e}") return None async def _get_document_status(self, doc_id: str) -> str: """查询 RAGFlow 文档处理状态。 Args: doc_id: RAGFlow 文档ID Returns: str: "processing" / "completed" / "failed" """ # TODO: 接入现有 RagflowClient 实现状态查询 # 当前返回 completed(占位,RAGFlow 服务部署后替换为真实调用) logger.debug(f"RAGFlow 文档状态查询(模拟): doc_id={doc_id}") return "completed" async def _fetch_document_chunks( self, doc_id: str ) -> List[Dict[str, Any]]: """拉取 RAGFlow 处理后的结构化段落。 Args: doc_id: RAGFlow 文档ID Returns: List[Dict]: 结构化段落列表 """ # TODO: 接入现有 RagflowClient 实现段落拉取 # 当前返回空列表(占位,RAGFlow 服务部署后替换为真实调用) logger.info(f"RAGFlow 文档段落拉取(模拟): doc_id={doc_id}") return []