# ============================================================================= # 企微IT智能服务台 — 分诊业务逻辑服务 # ============================================================================= # 说明:分诊交互的核心业务逻辑,包括: # 1. start_triage — 发起分诊(含5秒超时自动转人工) # 2. submit_step — 提交步骤选择 # 3. skip_step — 跳过步骤 # 4. complete_triage — 分诊完成生成最终回复 # 5. transfer_to_human — 转人工 # 6. determine_urgency — 紧急度判断(关键词规则) # 坐席端:list_pending / get_detail / route_session / get_history / export / exclude_options # ============================================================================= import asyncio import io import logging from datetime import datetime from typing import Any, Dict, List, Optional from openpyxl import Workbook from sqlalchemy import func, select, and_, case from sqlalchemy.ext.asyncio import AsyncSession from app.config import settings from app.models.triage_session import TriageSession from app.services.dify_triage_service import get_dify_triage_service logger = logging.getLogger(__name__) # ============================================================================= # 紧急度判断关键词规则(决策 #3) # ============================================================================= # 扩展:同时用于"人工"按钮紧急直通判定 URGENCY_HIGH_KEYWORDS: List[str] = [ "紧急", "马上", "宕机", "无法工作", "崩溃", "死机", "蓝屏", # 新增:紧急直通人工关键词(电脑无法启动、网络无法连接等) "电脑无法启动", "网络无法连接", "多人不能上网", "无法上网", "开不了机", "连不上网", "全部断网", ] URGENCY_MEDIUM_KEYWORDS: List[str] = [ "报错", "失败", "连不上", "打不开", "不能用", ] # ============================================================================= # 信息锁定判定(决策 B2/B3) # ============================================================================= # 有效回答:不在以下集合中的回答。无效回答包括"人工""不知道"等。 INVALID_ANSWERS = frozenset({ "人工", "不知道", "不确定", "转人工", "跳过", "", }) # 有效回答占比阈值:≥70% 判定为信息锁定 INFO_LOCKED_THRESHOLD = 0.70 # ============================================================================= # 关闭关键词识别(决策 G2:AI解决确认支持关键词识别) # ============================================================================= RESOLVE_KEYWORDS: List[str] = [ "解决了", "谢谢", "没问题了", "可以了", "好了", "弄好了", "搞定了", "不需要了", "撤销", "关闭", ] class TriageService: """分诊业务逻辑服务。 管理 AI 分诊的完整生命周期,从发起分诊到最终路由。 """ def __init__(self): """初始化分诊服务。""" self.dify_service = get_dify_triage_service() # ========================================================================== # 紧急度判断(关键词规则) # ========================================================================== @staticmethod def determine_urgency(question: str, confidence: Optional[float] = None) -> str: """根据关键词 + 置信度判断紧急度。 规则: 1. 含高级关键词(紧急/宕机/崩溃等)→ high 2. 置信度 < 0.5 → high(低置信也视为紧急) 3. 含中级关键词(报错/失败/连不上等)→ medium 4. 其余 → low Args: question: 员工问题文本 confidence: AI 置信度(可选) Returns: str: 紧急度(high/medium/low) """ if any(kw in question for kw in URGENCY_HIGH_KEYWORDS): return "high" if confidence is not None and confidence < 0.5: return "high" if any(kw in question for kw in URGENCY_MEDIUM_KEYWORDS): return "medium" return "low" # ========================================================================== # H5 端方法 # ========================================================================== async def start_triage( self, db: AsyncSession, conversation_id: str, question: str, user_id: str, user_name: str = "", user_dept: str = "", device_info: str = "", ) -> Dict[str, Any]: """发起分诊(含5秒超时自动转人工)。 流程: 1. 创建 triage_sessions 记录(status=triaging) 2. 调用 Dify 分诊应用(5秒超时) 3. 超时则自动转人工(status=timeout) 4. 成功则更新分诊步骤和AI分析结果 Args: db: 数据库会话 conversation_id: 会话ID question: 员工问题文本 user_id: 员工ID user_name: 员工姓名 user_dept: 员工部门 device_info: 设备信息 Returns: Dict[str, Any]: 分诊结果或超时信息 """ # 1. 创建分诊会话记录 session = TriageSession( conversation_id=conversation_id, user_id=user_id, user_name=user_name, user_dept=user_dept, device_info=device_info, request_title=question[:200] if question else "", request_content=question, source="wecom_h5", status="triaging", urgency="medium", ) db.add(session) await db.commit() await db.refresh(session) triage_id = session.id logger.info("分诊会话已创建: triage_id=%s, user=%s", triage_id, user_id) # 2. 调用 Dify 分诊(5秒超时) try: result = await asyncio.wait_for( self.dify_service.analyze(question, context=[], step_index=0), timeout=float(settings.dify_triage_timeout), ) except asyncio.TimeoutError: # 超时自动转人工 logger.warning("分诊超时(>%s秒),自动转人工: triage_id=%s", settings.dify_triage_timeout, triage_id) await self._transfer_to_human_on_timeout(db, triage_id) return { "status": "timeout", "message": "分诊超时,已自动转人工", "triage_id": triage_id, } except RuntimeError as e: # Dify 不可用,降级转人工 logger.error("Dify 分诊不可用,降级转人工: triage_id=%s, error=%s", triage_id, e) await self._transfer_to_human_on_timeout(db, triage_id) return { "status": "timeout", "message": "分诊服务暂时不可用,已自动转人工", "triage_id": triage_id, } # 3. 更新分诊会话 confidence = result.get("confidence") urgency = self.determine_urgency(question, confidence) # 覆盖 Dify 返回的紧急度(以关键词规则为准) if result.get("urgency") and not any( kw in question for kw in URGENCY_HIGH_KEYWORDS + URGENCY_MEDIUM_KEYWORDS ): urgency = result.get("urgency", "medium") session.triage_steps = result.get("triage_steps", []) session.confidence = confidence session.urgency = urgency session.suggested_route = result.get("suggested_route") session.problem_type = result.get("problem_type") session.problem_category = result.get("problem_category") session.matched_knowledge = result.get("matched_knowledge") session.match_score = result.get("match_score") session.context_tags = result.get("context_tags", []) session.status = "triaging" session.updated_at = datetime.now() await db.commit() await db.refresh(session) logger.info( "分诊分析完成: triage_id=%s, problem_type=%s, urgency=%s, steps=%d", triage_id, result.get("problem_type"), urgency, len(result.get("triage_steps", [])), ) return { "triage_id": triage_id, "steps": result.get("triage_steps", []), "total": len(result.get("triage_steps", [])), "confidence": confidence, "urgency": urgency, "suggested_route": result.get("suggested_route"), } async def submit_step( self, db: AsyncSession, triage_id: str, step_index: int, selected_label: str, ) -> Dict[str, Any]: """提交步骤选择(含信息锁定判定)。 记录用户选择的上下文,并根据选择动态调整后续步骤。 当所有步骤完成时,判定信息是否锁定: - 有效回答占比 ≥ 70% → info_locked = true → WS推送 queue_segment_changed - 有效回答占比 < 70% → info_locked = false(员工需继续回答诊断题补充) Args: db: 数据库会话 triage_id: 信息梳理会话ID step_index: 当前步骤序号 selected_label: 选择的选项标签 Returns: Dict[str, Any]: 下一步骤数据、已收集上下文、信息锁定状态 """ session = await self._get_session(db, triage_id) if not session: return {"error": "信息梳理会话不存在"} # 记录已收集的上下文 collected = list(session.collected_context or []) if selected_label and selected_label not in collected: collected.append(selected_label) session.collected_context = collected session.updated_at = datetime.now() await db.commit() # 获取下一步骤(从预生成的步骤中取) steps = session.triage_steps or [] next_index = step_index + 1 all_steps_done = next_index >= len(steps) if not all_steps_done: next_step = steps[next_index] if next_index < len(steps) else None else: next_step = None # ================================================================== # 信息锁定判定:所有步骤完成时触发 # ================================================================== info_locked = False if all_steps_done: info_locked = self._check_info_locked(collected) if info_locked: # 更新关联的 Conversation 表 await self._update_conversation_info_locked( db, session.conversation_id, locked=True ) logger.info( "信息锁定成功: triage_id=%s, conversation_id=%s, " "有效回答=%d/%d (%.0f%%)", triage_id, session.conversation_id, sum(1 for a in collected if a.strip() not in INVALID_ANSWERS), len(collected), (sum(1 for a in collected if a.strip() not in INVALID_ANSWERS) / max(len(collected), 1)) * 100 ) # 标记信息梳理状态为完成 session.status = "routed" session.route_action = "info_locked" session.updated_at = datetime.now() await db.commit() # WS推送:队列段位变更 await self._push_queue_segment_changed( session.user_id, session.conversation_id, "incomplete", "completed", "信息梳理完成,已进入优先队列" ) return { "next_step": next_step, "collected_context": collected, "info_locked": info_locked, "all_steps_done": all_steps_done, } async def skip_step( self, db: AsyncSession, triage_id: str, step_index: int, ) -> Dict[str, Any]: """跳过步骤。 Args: db: 数据库会话 triage_id: 分诊会话ID step_index: 要跳过的步骤序号 Returns: Dict[str, Any]: 下一步骤数据 """ session = await self._get_session(db, triage_id) if not session: return {"error": "分诊会话不存在"} steps = session.triage_steps or [] next_index = step_index + 1 session.updated_at = datetime.now() await db.commit() if next_index < len(steps): next_step = steps[next_index] else: next_step = None return {"next_step": next_step} async def complete_triage( self, db: AsyncSession, triage_id: str, context: List[str], ) -> Dict[str, Any]: """分诊完成,生成最终回复。 Args: db: 数据库会话 triage_id: 分诊会话ID context: 已收集的上下文列表 Returns: Dict[str, Any]: AI 回复和置信度 """ session = await self._get_session(db, triage_id) if not session: return {"error": "分诊会话不存在"} # 更新收集的上下文 session.collected_context = context session.status = "routed" session.route_action = "ai_self" session.updated_at = datetime.now() try: # 调用 Dify 生成最终回复 result = await self.dify_service.generate_reply( session.request_content, context ) reply = result.get("reply", "根据您提供的信息,建议联系IT服务台获取进一步帮助。") confidence = result.get("confidence", 0.0) except RuntimeError as e: logger.warning("Dify 生成回复失败,使用降级回复: %s", e) reply = "根据您提供的信息,建议联系IT服务台获取进一步帮助。" confidence = 0.0 await db.commit() return {"reply": reply, "confidence": confidence} async def transfer_to_human( self, db: AsyncSession, triage_id: str, context: List[str], ) -> Dict[str, Any]: """转人工。 Args: db: 数据库会话 triage_id: 分诊会话ID context: 已收集的上下文列表 Returns: Dict[str, Any]: 转人工结果 """ session = await self._get_session(db, triage_id) if not session: return {"error": "分诊会话不存在"} session.collected_context = context session.status = "routed" session.route_action = "human" session.updated_at = datetime.now() await db.commit() return { "conversation_id": session.conversation_id, "status": "waiting_agent", } # ========================================================================== # 坐席端方法 # ========================================================================== async def list_pending( self, db: AsyncSession, urgency: Optional[str] = None, problem_type: Optional[str] = None, page: int = 1, page_size: int = 20, ) -> Dict[str, Any]: """获取待分诊列表(按紧急度排序)。 排序规则:high > medium > low,同紧急度按创建时间倒序。 Args: db: 数据库会话 urgency: 紧急度筛选 problem_type: 问题类型筛选 page: 页码 page_size: 每页数量 Returns: Dict[str, Any]: {total, items} """ # 构建查询条件 conditions = [TriageSession.status.in_(["pending", "triaging"])] if urgency: conditions.append(TriageSession.urgency == urgency) if problem_type: conditions.append(TriageSession.problem_type == problem_type) # 紧急度排序:用 CASE 表达式 urgency_order = case( (TriageSession.urgency == "high", 0), (TriageSession.urgency == "medium", 1), (TriageSession.urgency == "low", 2), else_=3, ) stmt = ( select(TriageSession) .where(and_(*conditions)) .order_by(urgency_order, TriageSession.created_at.desc()) ) # 统计总数 count_stmt = select(func.count()).select_from(TriageSession).where(and_(*conditions)) total_result = await db.execute(count_stmt) total = total_result.scalar() or 0 # 分页 offset = (page - 1) * page_size stmt = stmt.offset(offset).limit(page_size) result = await db.execute(stmt) items = result.scalars().all() return { "total": total, "items": [self._session_to_dict(s) for s in items], } async def get_stats(self, db: AsyncSession) -> Dict[str, Any]: """获取分诊看板统计概要。 Args: db: 数据库会话 Returns: Dict[str, Any]: 统计数据 """ now = datetime.now() today_start = now.replace(hour=0, minute=0, second=0, microsecond=0) # 待分诊总数 pending_result = await db.execute( select(func.count()).select_from(TriageSession).where( TriageSession.status.in_(["pending", "triaging"]) ) ) pending_total = pending_result.scalar() or 0 # 今日已分诊数 today_result = await db.execute( select(func.count()).select_from(TriageSession).where( and_( TriageSession.status == "routed", TriageSession.operated_at >= today_start, ) ) ) today_triaged = today_result.scalar() or 0 # AI 自答数 ai_self_result = await db.execute( select(func.count()).select_from(TriageSession).where( and_( TriageSession.route_action == "ai_self", TriageSession.operated_at >= today_start, ) ) ) ai_self_count = ai_self_result.scalar() or 0 # 转人工数 human_result = await db.execute( select(func.count()).select_from(TriageSession).where( and_( TriageSession.route_action == "human", TriageSession.operated_at >= today_start, ) ) ) human_count = human_result.scalar() or 0 # 自动审批数 auto_result = await db.execute( select(func.count()).select_from(TriageSession).where( and_( TriageSession.route_action == "auto_approval", TriageSession.operated_at >= today_start, ) ) ) auto_approval_count = auto_result.scalar() or 0 # 平均耗时(从创建到操作) avg_result = await db.execute( select( func.avg( func.extract("epoch", TriageSession.operated_at - TriageSession.created_at) ) ).where( and_( TriageSession.status == "routed", TriageSession.operated_at.isnot(None), TriageSession.operated_at >= today_start, ) ) ) avg_duration = avg_result.scalar() avg_duration_sec = float(avg_duration) if avg_duration else 0.0 return { "pending_total": pending_total, "today_triaged": today_triaged, "ai_self_count": ai_self_count, "human_count": human_count, "auto_approval_count": auto_approval_count, "avg_duration_sec": round(avg_duration_sec, 1), } async def get_detail(self, db: AsyncSession, triage_id: str) -> Optional[Dict[str, Any]]: """获取分诊详情。 Args: db: 数据库会话 triage_id: 分诊会话ID Returns: Optional[Dict[str, Any]]: 分诊详情字典,不存在返回 None """ session = await self._get_session(db, triage_id) if not session: return None return self._session_to_detail_dict(session) async def route_session( self, db: AsyncSession, triage_id: str, route_action: str, route_note: Optional[str], operator_id: str, ) -> Optional[Dict[str, Any]]: """坐席路由操作(覆盖 AI 建议)。 Args: db: 数据库会话 triage_id: 分诊会话ID route_action: 路由动作 route_note: 路由备注 operator_id: 操作坐席ID Returns: Optional[Dict[str, Any]]: 更新后的分诊会话字典 """ session = await self._get_session(db, triage_id) if not session: return None session.route_action = route_action session.route_note = route_note session.operator_id = operator_id session.operated_at = datetime.now() session.status = "routed" if route_action != "skip" else "skipped" session.updated_at = datetime.now() await db.commit() await db.refresh(session) return self._session_to_dict(session) async def get_history( self, db: AsyncSession, date_from: Optional[str] = None, date_to: Optional[str] = None, route_action: Optional[str] = None, page: int = 1, page_size: int = 20, ) -> Dict[str, Any]: """获取已分诊历史列表。 Args: db: 数据库会话 date_from: 开始日期 date_to: 结束日期 route_action: 路由动作筛选 page: 页码 page_size: 每页数量 Returns: Dict[str, Any]: {total, items} """ conditions = [TriageSession.status.in_(["routed", "skipped", "timeout"])] if date_from: try: dt_from = datetime.fromisoformat(date_from) conditions.append(TriageSession.created_at >= dt_from) except ValueError: pass if date_to: try: dt_to = datetime.fromisoformat(date_to) conditions.append(TriageSession.created_at <= dt_to) except ValueError: pass if route_action: conditions.append(TriageSession.route_action == route_action) stmt = ( select(TriageSession) .where(and_(*conditions)) .order_by(TriageSession.created_at.desc()) ) count_stmt = select(func.count()).select_from(TriageSession).where(and_(*conditions)) total_result = await db.execute(count_stmt) total = total_result.scalar() or 0 offset = (page - 1) * page_size stmt = stmt.offset(offset).limit(page_size) result = await db.execute(stmt) items = result.scalars().all() return { "total": total, "items": [self._session_to_dict(s) for s in items], } async def export_sessions( self, db: AsyncSession, date_from: Optional[str] = None, date_to: Optional[str] = None, ) -> bytes: """导出分诊记录为 xlsx。 导出基础字段 + 分诊步骤详情。 Args: db: 数据库会话 date_from: 开始日期 date_to: 结束日期 Returns: bytes: xlsx 文件内容 """ conditions = [] if date_from: try: dt_from = datetime.fromisoformat(date_from) conditions.append(TriageSession.created_at >= dt_from) except ValueError: pass if date_to: try: dt_to = datetime.fromisoformat(date_to) conditions.append(TriageSession.created_at <= dt_to) except ValueError: pass stmt = select(TriageSession).order_by(TriageSession.created_at.desc()) if conditions: stmt = stmt.where(and_(*conditions)) result = await db.execute(stmt) sessions = result.scalars().all() # 构建 Excel wb = Workbook() ws = wb.active ws.title = "分诊记录" # 表头 headers = [ "分诊ID", "会话ID", "员工ID", "员工姓名", "部门", "问题标题", "问题类型", "问题分类", "置信度", "紧急度", "AI建议路由", "最终路由", "路由备注", "操作坐席", "创建时间", "操作时间", "已收集上下文", "分诊步骤详情", ] ws.append(headers) # 数据行 for s in sessions: steps_detail = "" if s.triage_steps: for i, step in enumerate(s.triage_steps, 1): q = step.get("question", "") opts = " | ".join( f"{o.get('label', '')}({o.get('probability', 0):.0%})" for o in step.get("options", []) ) steps_detail += f"步骤{i}: {q} [{opts}]; " ws.append([ s.id, s.conversation_id, s.user_id, s.user_name or "", s.user_dept or "", s.request_title, s.problem_type or "", s.problem_category or "", round(s.confidence, 2) if s.confidence else "", s.urgency, s.suggested_route or "", s.route_action or "", s.route_note or "", s.operator_id or "", s.created_at.strftime("%Y-%m-%d %H:%M:%S") if s.created_at else "", s.operated_at.strftime("%Y-%m-%d %H:%M:%S") if s.operated_at else "", " / ".join(s.collected_context or []), steps_detail, ]) # 调整列宽 for col in ws.columns: max_length = max(len(str(cell.value or "")) for cell in col) ws.column_dimensions[col[0].column_letter].width = min(max_length + 2, 50) # 输出到内存 output = io.BytesIO() wb.save(output) output.seek(0) return output.getvalue() async def exclude_options( self, db: AsyncSession, triage_id: str, excluded_labels: List[str], recommended_label: Optional[str], ) -> Dict[str, Any]: """坐席排除/推荐分诊选项(通过 WS 推送到 H5)。 Args: db: 数据库会话 triage_id: 分诊会话ID excluded_labels: 要排除的选项标签列表 recommended_label: 推荐的选项标签 Returns: Dict[str, Any]: 排除结果 """ session = await self._get_session(db, triage_id) if not session: return {"error": "分诊会话不存在"} # 通过 WS 推送到 H5 端 from app.services.ws_manager import manager as ws_manager ws_data = { "type": "triage_exclude", "data": { "triage_id": triage_id, "excluded_labels": excluded_labels, "recommended_label": recommended_label, }, } await ws_manager.send_to_employee(session.user_id, ws_data) logger.info( "排除选项已推送: triage_id=%s, excluded=%s, recommended=%s", triage_id, excluded_labels, recommended_label, ) return {"excluded": True} # ========================================================================== # 内部辅助方法 # ========================================================================== @staticmethod def _check_info_locked(collected_context: List[str]) -> bool: """判定信息是否锁定(决策 B2/B3)。 条件:有效回答占比 ≥ 70%。 有效回答 = 不在 INVALID_ANSWERS 集合中的回答。 Args: collected_context: 已收集的上下文回答列表 Returns: bool: True=已锁定,False=未锁定 """ if not collected_context: return False total = len(collected_context) valid = sum(1 for ans in collected_context if ans.strip() not in INVALID_ANSWERS) return (valid / total) >= INFO_LOCKED_THRESHOLD async def _update_conversation_info_locked( self, db: AsyncSession, conversation_id: str, locked: bool ) -> None: """更新 Conversation 表的 info_locked 字段。 Args: db: 数据库会话 conversation_id: 会话ID locked: 是否锁定 """ from app.models.conversation import Conversation result = await db.execute( select(Conversation).where(Conversation.id == conversation_id) ) conv = result.scalar_one_or_none() if conv: conv.info_locked = locked conv.updated_at = datetime.now() await db.commit() logger.info("Conversation info_locked 更新: conv_id=%s, locked=%s", conversation_id, locked) async def _push_queue_segment_changed( self, employee_id: str, conversation_id: str, old_segment: str, new_segment: str, message: str, ) -> None: """推送队列段位变更 WS事件(queue_segment_changed)。 当 info_locked 变为 true 时,员工从"待梳理"段升级到"已梳理"段。 Args: employee_id: 员工ID conversation_id: 会话ID old_segment: 原段位(incomplete) new_segment: 新段位(completed) message: 提示消息 """ try: from app.services.ws_manager import manager as ws_manager ws_data = { "type": "queue_segment_changed", "data": { "conversation_id": conversation_id, "old_segment": old_segment, "new_segment": new_segment, "message": message, }, } await ws_manager.send_to_employee(employee_id, ws_data) except Exception as e: logger.warning("WS推送队列段位变更失败: %s", e) async def _get_session(self, db: AsyncSession, triage_id: str) -> Optional[TriageSession]: """获取分诊会话记录。""" result = await db.execute( select(TriageSession).where(TriageSession.id == triage_id) ) return result.scalar_one_or_none() async def _transfer_to_human_on_timeout( self, db: AsyncSession, triage_id: str ) -> None: """超时自动转人工。""" session = await self._get_session(db, triage_id) if session: session.status = "timeout" session.route_action = "human" session.route_note = "分诊超时,自动转人工" session.updated_at = datetime.now() await db.commit() @staticmethod def _session_to_dict(s: TriageSession) -> Dict[str, Any]: """将会话对象转为列表项字典。""" return { "id": s.id, "conversation_id": s.conversation_id, "user_id": s.user_id, "user_name": s.user_name, "user_dept": s.user_dept, "request_title": s.request_title, "problem_type": s.problem_type, "problem_category": s.problem_category, "confidence": s.confidence, "urgency": s.urgency, "suggested_route": s.suggested_route, "status": s.status, "route_action": s.route_action, "route_note": s.route_note, "operator_id": s.operator_id, "created_at": s.created_at.isoformat() if s.created_at else None, "operated_at": s.operated_at.isoformat() if s.operated_at else None, } @staticmethod def _session_to_detail_dict(s: TriageSession) -> Dict[str, Any]: """将会话对象转为详情字典。""" return { "id": s.id, "conversation_id": s.conversation_id, "user_id": s.user_id, "user_name": s.user_name, "user_dept": s.user_dept, "user_level": s.user_level, "device_info": s.device_info, "request_title": s.request_title, "request_content": s.request_content, "source": s.source, "problem_type": s.problem_type, "problem_category": s.problem_category, "confidence": s.confidence, "urgency": s.urgency, "suggested_route": s.suggested_route, "matched_knowledge": s.matched_knowledge, "match_score": s.match_score, "context_tags": s.context_tags or [], "triage_steps": s.triage_steps or [], "collected_context": s.collected_context or [], "status": s.status, "route_action": s.route_action, "route_note": s.route_note, "operator_id": s.operator_id, "created_at": s.created_at.isoformat() if s.created_at else None, "updated_at": s.updated_at.isoformat() if s.updated_at else None, "operated_at": s.operated_at.isoformat() if s.operated_at else None, } # 单例 _triage_service: Optional[TriageService] = None def get_triage_service() -> TriageService: """获取 TriageService 单例。 Returns: TriageService: 单例实例 """ global _triage_service if _triage_service is None: _triage_service = TriageService() return _triage_service