Files
wecom_it_smart_desk/backend/app/tasks/h5_ai_task.py
T

1329 lines
53 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
# =============================================================================
# 企微IT智能服务台 — H5 员工端 AI 回复后台任务
# =============================================================================
# 背景:原 h5_send_message 在同步 HTTP 请求内 await AI 推理(Dify 3~15s),
# 整条请求被阻塞,前端表现为"发送中"长时间卡顿。
# 本模块将 AI 推理移出请求,改为 asyncio 后台任务,结果经 WebSocket
# 流式推回(ai_reply_chunk / ai_reply),发送瞬时完成。
#
# 关键约束(详见 docs/02-需求分析/技术架构演进/员工端消息发送延时改造方案.md):
# 1. 必须单 worker 运行(docker-compose --workers 1):
# ws_manager 是进程内单例,多 worker 时后台任务与员工 WS 连接可能不在
# 同进程,broadcast 会静默丢失(约 50%)。
# 2. 使用独立 DB session_get_session_factory),不可复用请求的 db
# (请求返回后该 session 会被关闭)。
# =============================================================================
import asyncio
import logging
import os
from datetime import datetime, timedelta
from pathlib import Path
from sqlalchemy import select
from app.api.byod import _byod_keyword_prefilter
from app.database import _get_session_factory
from app.dependencies import get_shared_ai_handler
from app.models.conversation import Conversation
from app.models.message import Message
from app.services.routing_service import (
routing_keyword_prefilter,
detect_routing_intent,
get_contact_by_category,
send_contact_card,
record_routing_event,
_keyword_fallback_category,
)
from app.services.vision_service import VisionService
from app.services.ws_manager import manager as ws_manager
from app.services.asset_recommend_service import get_asset_recommend_service
from app.services.employee_profile_service import get_employee_profile_service
from app.api.approval import APPROVAL_TEMPLATES
logger = logging.getLogger(__name__)
# =============================================================================
# Phase 4A: VisionService 接入 — 图片消息视觉理解
# =============================================================================
# 图片文件本地存储根目录(与 upload.py 中 UPLOAD_DIR 一致)
_UPLOAD_DIR = Path(os.getenv("UPLOAD_DIR", "./uploads"))
# 视觉理解置信度阈值:低于此值不注入描述(避免错误描述误导 AI)
_VISION_CONFIDENCE_THRESHOLD = 0.6
def _media_url_to_local_path(media_url: str) -> Path:
"""将媒体 URL 路径转换为本地文件系统路径。
做什么:把 "/api/media/2026/07/13/abc.png" 转换为
"./uploads/2026/07/13/abc.png"
为什么:VisionService 需要读取原始图片字节流,
而媒体 URL 是 HTTP 访问路径,不是文件系统路径。
Args:
media_url: 媒体文件 URL(如 /api/media/2026/07/13/abc.png
Returns:
Path: 本地文件路径对象
"""
# 去掉 URL 前缀 /api/media/,拼接到 UPLOAD_DIR
# 例: "/api/media/2026/07/13/abc.png" → "2026/07/13/abc.png"
relative = media_url.replace("/api/media/", "", 1)
return _UPLOAD_DIR / relative
async def _fetch_recent_employee_text(
db, conversation_id: str, employee_id: str, within_seconds: int = 5
) -> str:
"""获取最近 N 秒内员工的文字消息(Phase 4B 消息融合)。
做什么:查询同一会话中,当前图片消息之前 within_seconds 秒内,
员工发送的文本消息内容。
为什么:用户经常先打字描述问题再发截图,或先发截图再补充文字。
将文字与图片视觉描述融合后一次性传给 Dify,
避免 AI 分别处理两条消息导致上下文割裂。
Args:
db: 异步 DB session
conversation_id: 会话 ID
employee_id: 员工企微 UserID
within_seconds: 时间窗口(秒),默认 5 秒
Returns:
str: 最近的员工文字消息内容(多条用换行拼接),无则返回空字符串
"""
cutoff = datetime.now() - timedelta(seconds=within_seconds)
stmt = (
select(Message)
.where(
Message.conversation_id == conversation_id,
Message.sender_type == "employee",
Message.sender_id == employee_id,
Message.msg_type == "text",
Message.created_at >= cutoff,
)
.order_by(Message.created_at.desc())
.limit(3) # 最多取 3 条,避免内容过长
)
result = await db.execute(stmt)
messages = result.scalars().all()
if not messages:
return ""
# 按时间正序拼接(先发的在前)
texts = [m.content for m in reversed(messages) if m.content]
return "\n".join(texts)
async def _enrich_image_content(
db,
media_url: str,
original_content: str,
conversation_id: str,
employee_id: str,
) -> str:
"""用 VisionService 分析图片,生成增强后的消息内容。
做什么:
1. 从本地文件系统读取图片
2. 调用 VisionService.analyze_screenshot() 获取视觉描述
3. 查询最近 5 秒内的员工文字消息(消息融合)
4. 拼接视觉描述 + 用户文字 → 传给 Dify
为什么:Dify 文本模型无法直接""图片,需要先将图片转为
文字描述,再与用户输入融合后传给 Dify 推理。
降级策略:
- 图片文件不存在 → 返回原始 content
- VisionService 调用失败 → 返回 "我收到了您的截图,但暂时无法识别内容"
- 置信度 < 0.6 → 不注入视觉描述,仅使用用户文字
Args:
db: 异步 DB session
media_url: 图片 URL(如 /api/media/2026/07/13/abc.png
original_content: 原始消息内容(如 "[图片] 截图"
conversation_id: 会话 ID
employee_id: 员工企微 UserID
Returns:
str: 增强后的消息内容(视觉描述 + 用户文字)
"""
# 1. 读取本地图片文件
local_path = _media_url_to_local_path(media_url)
if not local_path.exists():
logger.warning(f"图片文件不存在: {local_path} (media_url={media_url})")
return original_content
try:
image_bytes = local_path.read_bytes()
except Exception as e:
logger.error(f"读取图片文件失败: {local_path} - {e}")
return original_content
# 2. 调用 VisionService 分析截图
vision_service = VisionService()
try:
result = await vision_service.analyze_screenshot(
image_bytes, conversation_id
)
description = result.get("description", "")
confidence = result.get("confidence", 0.0)
logger.info(
f"VisionService 分析完成: conversation={conversation_id}, "
f"confidence={confidence:.2f}, desc_len={len(description)}"
)
# 3. 注入视觉描述到会话上下文(供后续多轮对话使用)
if description and confidence >= _VISION_CONFIDENCE_THRESHOLD:
await vision_service.inject_to_conversation_context(
description, conversation_id
)
except Exception as e:
logger.error(f"VisionService 调用异常: {e}")
description = ""
confidence = 0.0
finally:
await vision_service.close()
# 4. 消息融合:查询最近 5 秒内员工的文字消息
recent_text = await _fetch_recent_employee_text(
db, conversation_id, employee_id, within_seconds=5
)
# 5. 拼接增强内容
# 格式:[视觉描述] + [用户最近文字] + [原始消息内容]
parts = []
if description and confidence >= _VISION_CONFIDENCE_THRESHOLD:
parts.append(f"[用户发送了截图,视觉理解结果] {description}")
if recent_text:
parts.append(f"[用户最近的文字描述] {recent_text}")
# 原始内容如果不是纯占位符(如"[图片] 截图"),也加入
if original_content and not original_content.startswith("[图片]"):
parts.append(original_content)
if not parts:
# 降级:视觉分析失败且无文字补充
return "我收到了您的截图,但暂时无法识别内容,请描述一下您遇到的问题。"
return "\n".join(parts)
async def _persist_and_push_solution(
db,
conversation,
employee_id: str,
solution,
):
"""处理图谱命中的解决方案(图谱查询结果)
做什么:
1. 创建AI回复消息记录(图谱命中的解决方案)
2. 通过WS推送给员工端
为什么:图谱命中的解决方案直接返回,不需要调用Dify
Args:
db: 数据库会话
conversation: 会话对象
employee_id: 员工ID
solution: SolutionResult图谱查询结果
"""
from app.models.message import Message
from app.services.ws_manager import manager as ws_manager
# 1. 创建AI回复消息记录
message = Message(
conversation_id=conversation.id,
sender_type="ai",
sender_id="graph",
content=solution.solution,
msg_type="text",
)
db.add(message)
# 更新会话计数
conversation.ai_substantive_reply_count += 1
conversation.updated_at = datetime.now()
await db.commit()
# 2. 构建推送数据
msg_data = {
"id": str(message.id),
"conversation_id": str(conversation.id),
"sender_type": "ai",
"sender_id": "graph",
"sender_name": "智能助手",
"content": solution.solution,
"msg_type": "text",
"created_at": message.created_at.isoformat(),
"reply_source": "graph_hit", # 标记来源为图谱命中
}
# 3. 通过WS推送给员工端
try:
await ws_manager.broadcast_to_employees([employee_id], {
"type": "ai_reply",
"data": msg_data,
})
logger.info(
f"图谱命中推送成功: employee={employee_id}, "
f"solution={solution.action_name}"
)
except Exception as push_err:
logger.error(f"图谱命中推送失败: {push_err}")
async def _persist_and_push(
db,
conversation: Conversation,
employee_id: str,
content: str,
is_guidance: bool,
should_count: bool,
should_transfer: bool,
dify_conversation_id,
):
"""持久化 AI 回复并推送给员工端 + 广播坐席端。
做什么:
1. 存 AI 消息到 DB
2. 更新会话状态(dify 上下文 / 计数 / 转人工)
3. 经 WS 向员工推 ai_reply 终态(前端据此替换打字机气泡)
4. 经 WS 向坐席端广播 new_message + conversation_updated
为什么:把"落库 + 推送"封装为单点,供同步路径与流式路径复用。
"""
# 1. 存 AI 消息
ai_message = Message(
conversation_id=conversation.id,
sender_type="ai",
sender_id="ai_bot",
sender_name="Duckula(达寇拉)",
content=content,
msg_type="text",
is_read=True,
)
db.add(ai_message)
await db.flush()
# 2. 更新会话状态
if dify_conversation_id:
conversation.dify_conversation_id = dify_conversation_id
if should_count:
conversation.ai_substantive_reply_count += 1
if should_transfer:
conversation.status = "queued"
conversation.updated_at = datetime.now()
db.add(conversation)
await db.flush()
await db.commit()
# 3. 推 ai_reply 终态给员工(前端替换打字机气泡)
await ws_manager.broadcast_to_employees([employee_id], {
"type": "ai_reply",
"data": {
"message_id": str(ai_message.id),
"conversation_id": str(conversation.id),
"sender_type": "ai",
"sender_id": "ai_bot",
"sender_name": "Duckula(达寇拉)",
"content": content,
"msg_type": "text",
"is_guidance": is_guidance,
"ai_reply_count": conversation.ai_substantive_reply_count,
"can_call_agent": conversation.ai_substantive_reply_count >= 3,
"conversation_status": conversation.status,
},
})
# 4. 广播坐席端(new_message + conversation_updated
try:
await ws_manager.broadcast({
"type": "new_message",
"data": {
"conversation_id": str(conversation.id),
"message_id": str(ai_message.id),
"sender_type": "ai",
"sender_id": "ai_bot",
"sender_name": "Duckula(达寇拉)",
"content": content,
"msg_type": "text",
},
})
await ws_manager.broadcast({
"type": "conversation_updated",
"data": {
"conversation_id": str(conversation.id),
"status": conversation.status,
"assigned_agent_id": str(conversation.assigned_agent_id) if conversation.assigned_agent_id else None,
},
})
except Exception as ws_err:
# WS 广播失败不阻塞消息存储,只记录 warning
logger.warning(f"WS 广播 AI 回复给坐席失败(消息已存储): {ws_err}")
async def _persist_and_push_structured(
db,
conversation: Conversation,
employee_id: str,
result: dict,
):
"""持久化结构化 AI 回复并推送给员工端 + 广播坐席端(v2.0 双 WS 通道)。
改造后的核心变化(2026-07-13):
- Dify 返回 JSON {text, action, options},后端解析后同时发两条 WS:
① ai_reply → 聊天气泡(text + options
② dynamic_recommend → 侧边栏推荐(action 卡片)
- 两条消息同一时刻发出,零时间差到达
- 文字明确引用侧边栏内容(如"右侧已为您准备好入口"),语义强关联
命中判断规则:
- 结构化回复且有 action 或 options → 视为命中(AI 在主动引导)
- 纯文本回复 → 走原有 _check_knowledge_hit 判断
Args:
db: 异步 DB session
conversation: 当前会话对象
employee_id: 员工企微 UserID
result: get_structured_reply() 返回的结构化结果
"""
text = result.get("text", "")
# ★ 防御性类型保护:确保 content 始终是 String
# 如果 Dify 返回的 text 是 dict/listWS 推送后前端会显示 [object Object]
if not isinstance(text, str):
import json as _json
text = _json.dumps(text, ensure_ascii=False) if text else ""
logger.warning(f"_persist_and_push_structured: text 非 String 类型,已转换: {text[:80]}...")
action = result.get("action")
options = result.get("options")
hit = result.get("hit", False)
is_structured = result.get("is_structured", False)
dify_conv_id = result.get("conversation_id")
# Phase 6A: 提取诊断阶段
diagnosis_stage = result.get("diagnosis_stage")
# 结构化回复且有 action 或 options → 视为命中(AI 在主动引导/推荐)
if is_structured and (action or options):
hit = True
# Phase 6A: 基于 diagnosis_stage 调整会话状态
# escalating → AI 建议转人工
# resolved → AI 认为问题已解决
if diagnosis_stage == "escalating":
hit = False # 不计为有效回复,触发转人工
elif diagnosis_stage == "resolved":
hit = True # 计为有效回复
should_count = hit
should_transfer = not hit
# 确定消息类型
# v2.4 修复:有 action(审批卡片)时,强制设置为 ai_structured
# 确保前端能渲染审批卡片入口,不依赖 Dify 返回的 is_structured 字段
if action or options or is_structured:
msg_type = "ai_structured"
else:
msg_type = "text"
# v2.5 调试日志
logger.info(f"[DEBUG] msg_type = {msg_type}, action = {bool(action)}, options = {bool(options)}, is_structured = {is_structured}")
# 构建 extra_data(存储 options 和 action 供前端渲染)
extra_data = {}
if options:
extra_data["options"] = options
# 为审批卡片注入标准化 card_data(替换原有的分散匹配逻辑)
if action:
# v3.2 修复:降级路径已构建好 card_data 时,跳过重复匹配(避免 approval_type=None 误报"匹配失败"
if action.get("card_data"):
matched_card = action["card_data"]
options = matched_card.get("options", [])
if options and not action.get("url"):
action["url"] = options[0].get("url", "")
logger.info(f"[ApprovalMatcher] 使用预构建 card_data: {matched_card.get('title')}")
else:
approval_type = action.get("approval_type")
title = action.get("title")
# v3.0 重构:委托 ApprovalMatcher 统一完成模板匹配 + 卡片构建
from app.services.approval_matcher import get_approval_matcher
matcher = get_approval_matcher()
matched_card = matcher.match_and_build_card(approval_type, title)
if matched_card:
# 注入 URL 供 approve-direct-card 兼容路径使用
options = matched_card.get("options", [])
if options:
action["url"] = options[0].get("url", "")
action["card_data"] = matched_card
logger.info(f"[ApprovalMatcher] 匹配成功: {approval_type} -> card_type={matched_card.get('card_type')}")
else:
# v4.0 P1-4 后:matcher 仅在 approval_type 为空时返回 None(其余情况兜底全量卡片)
logger.info(f"[ApprovalMatcher] approval_type 为空,无卡片: title={title}")
extra_data["action"] = action
# ★ 调试日志:打印推送到前端的 extra_data 内容
logger.info(f"[DEBUG] 推送到前端的 extra_data: {extra_data}")
if extra_data.get("action"):
logger.info(f"[DEBUG] extra_data.action.approval_type = {extra_data['action'].get('approval_type')}")
# 1. 存 AI 消息
ai_message = Message(
conversation_id=conversation.id,
sender_type="ai",
sender_id="ai_bot",
sender_name="Duckula(达寇拉)",
content=text,
msg_type=msg_type,
extra_data=extra_data if extra_data else None,
is_read=True,
)
db.add(ai_message)
await db.flush()
# 2. 更新会话状态
if dify_conv_id:
conversation.dify_conversation_id = dify_conv_id
if should_count:
conversation.ai_substantive_reply_count += 1
if should_transfer:
conversation.status = "queued"
# Phase 6A: 将 diagnosis_stage 存入 tags(无需迁移,利用现有 JSON 字段)
if diagnosis_stage:
tags = conversation.tags or {}
tags["diagnosis_stage"] = diagnosis_stage
tags["diagnosis_updated_at"] = datetime.now().isoformat()
conversation.tags = tags
conversation.updated_at = datetime.now()
db.add(conversation)
await db.flush()
await db.commit()
# 3. 推 ai_reply 给员工端(聊天气泡:text + options
# 添加异常处理,避免整体失败
try:
await ws_manager.broadcast_to_employees([employee_id], {
"type": "ai_reply",
"data": {
"message_id": str(ai_message.id),
"conversation_id": str(conversation.id),
"sender_type": "ai",
"sender_id": "ai_bot",
"sender_name": "Duckula(达寇拉)",
"content": text,
"msg_type": msg_type,
"extra_data": extra_data if extra_data else None,
"is_guidance": False,
"ai_reply_count": conversation.ai_substantive_reply_count,
"can_call_agent": conversation.ai_substantive_reply_count >= 3,
"conversation_status": conversation.status,
# Phase 6A: 诊断阶段(前端可据此调整 UI/提示)
"diagnosis_stage": diagnosis_stage,
},
})
except Exception as emp_err:
# 员工端推送失败不应该导致整个任务失败
logger.warning(f"员工端 AI 回复推送失败(不影响坐席端): {emp_err}")
# 4. 推 dynamic_recommend 给员工端侧边栏(仅当 action 非空且不是审批类型时)
# 与 ai_reply 同一时刻发出 → 零时间差到达
# v2.3 修改:审批类型只推送到消息气泡(左边),不推送到侧边栏(右边)
if action and not action.get("approval_type"):
# 为审批卡片注入运维平台跳转URL(实现免登录跳转)
approval_type = action.get("approval_type")
action_url = ""
location = "运维平台"
# 优先精确匹配:直接用 approval_type 查找模板
if approval_type and approval_type in APPROVAL_TEMPLATES:
template = APPROVAL_TEMPLATES[approval_type]
action_url = template.get("url", "")
location = template.get("location", "运维平台")
# 关键字匹配:当精确匹配失败时,通过关键字查找模板
# Dify返回的 approval_type 可能是中文分类名(如"账号权限申请"、"VPN账号申请"
elif approval_type:
for template_id, template in APPROVAL_TEMPLATES.items():
keywords = template.get("keywords", [])
# 检查 approval_type 是否包含任意一个关键字
if any(kw.lower() in approval_type.lower() for kw in keywords):
action_url = template.get("url", "")
location = template.get("location", "运维平台")
break
# v2.2 新增:根据 Dify 返回的 approval_type 设置 filtered_options
filtered_options = []
if approval_type:
filtered_options = [approval_type]
recommend_data = {
"recommend_id": f"rec_{ai_message.id}",
"card_type": action.get("type", "approval_card"),
"title": action.get("title", ""),
"description": action.get("description", ""),
"approval_type": action.get("approval_type"),
"filtered_options": filtered_options, # v2.2: 精确匹配的选项
"action_url": action_url, # 运维平台跳转URL
"action_label": f"打开{location}" if action_url else "打开审批表单", # 按钮文字
"location": location, # 平台名称
"confidence": action.get("confidence", 0.85),
"message_id": str(ai_message.id),
"conversation_id": str(conversation.id),
}
try:
await ws_manager.broadcast_to_employees([employee_id], {
"type": "dynamic_recommend",
"data": recommend_data,
})
logger.info(
f"动态推荐已推送: employee={employee_id}, "
f"card_type={recommend_data['card_type']}, "
f"title={recommend_data['title']}"
)
except Exception as rec_err:
logger.warning(f"员工端动态推荐推送失败: {rec_err}")
# 5. 广播坐席端(new_message + conversation_updated
try:
await ws_manager.broadcast({
"type": "new_message",
"data": {
"conversation_id": str(conversation.id),
"message_id": str(ai_message.id),
"sender_type": "ai",
"sender_id": "ai_bot",
"sender_name": "Duckula(达寇拉)",
"content": text,
"msg_type": msg_type,
"extra_data": extra_data if extra_data else None,
},
})
await ws_manager.broadcast({
"type": "conversation_updated",
"data": {
"conversation_id": str(conversation.id),
"status": conversation.status,
"assigned_agent_id": (
str(conversation.assigned_agent_id)
if conversation.assigned_agent_id else None
),
},
})
except Exception as ws_err:
logger.warning(f"WS 广播结构化 AI 回复给坐席失败(消息已存储): {ws_err}")
async def _handle_byod_query(db, conversation, employee_id, content):
"""处理 BYOD 自备电脑补贴查询。
在 H5 聊天消息流中拦截 BYOD 关键词后执行资格检查,并以 byod_card
卡片消息形式推送给员工端(前端 MessageBubble 据 msg_type 渲染
ByodSubsidyCard)。
流程:
1. 通过 WecomService 获取员工岗位(position
2. 与 BYOD 资格清单匹配(_match_position
3. 创建 byod_card 类型 AI 消息并落库
4. 经 WS 推送 ai_reply 给员工端(携带 extra_data.byod_result
5. 广播 new_message + conversation_updated 给坐席端(与 _persist_and_push 一致)
Args:
db: 异步 DB sessionprocess_h5_ai_reply 的 factory session
conversation: 当前会话对象(Conversation
employee_id: 员工企微 UserID
content: 用户消息原文(用于日志)
"""
# 延迟导入避免循环依赖(byod 模块注册路由时可能引用 app.main)
from app.api.byod import _match_position, BYOD_APPLICATION_URL, BYOD_NOTES, BYOD_REGISTER_NOTES
from app.services.wecom_service import WecomService
# 1. 获取员工岗位(企微通讯录 API)
position = ""
try:
wecom_service = WecomService()
try:
user_info = await wecom_service.get_user_info(employee_id)
position = user_info.get("position", "")
finally:
await wecom_service.close()
except Exception as e:
logger.error(f"BYOD: 获取员工岗位失败: {e}")
# 2. 岗位匹配(返回: 是否匹配, 匹配岗位, 匹配类别)
eligible, matched_pos, matched_category = _match_position(position)
# 3. 构建 BYOD 结果数据
# 字段与前端 ByodSubsidyCard.vue props 完全一致:
# eligible / position / matched_category / application_url / notes / reason
byod_result = {
"eligible": eligible,
"has_subsidy": eligible,
"position": position,
"matched_category": matched_category,
"application_url": BYOD_APPLICATION_URL, # 所有岗位都提供链接
"notes": BYOD_NOTES if eligible else BYOD_REGISTER_NOTES,
"reason": (
"" if eligible
else f"您的岗位「{position}」不在自备电脑补贴资格清单中,可进行自备电脑登记(无补贴)"
),
}
# 4. 展示文本(AI 气泡的 content,卡片下方不直接展示,但会话列表/坐席端可见)
if eligible:
display_text = f"您岗位为「{position}」,符合自备电脑补贴申请资格"
else:
display_text = f"您岗位为「{position}」,可进行自备电脑登记(无补贴)"
# 5. 创建 AI 消息(byod_card 类型,携带 byod_result
ai_message = Message(
conversation_id=conversation.id,
sender_type="ai",
sender_id="ai_bot",
sender_name="Duckula(达寇拉)",
content=display_text,
msg_type="byod_card",
extra_data={"byod_result": byod_result},
is_read=True,
)
db.add(ai_message)
await db.flush()
# 6. 更新会话状态(计数 + 时间,BYOD 视为一次实质性 AI 回复)
conversation.ai_substantive_reply_count += 1
conversation.updated_at = datetime.now()
db.add(conversation)
await db.flush()
await db.commit()
# 7. 推送 ai_reply 给员工端(前端据 msg_type="byod_card" 渲染卡片)
await ws_manager.broadcast_to_employees([employee_id], {
"type": "ai_reply",
"data": {
"message_id": str(ai_message.id),
"conversation_id": str(conversation.id),
"sender_type": "ai",
"sender_id": "ai_bot",
"sender_name": "Duckula(达寇拉)",
"content": display_text,
"msg_type": "byod_card",
"extra_data": {"byod_result": byod_result},
"is_guidance": False,
"ai_reply_count": conversation.ai_substantive_reply_count,
"can_call_agent": conversation.ai_substantive_reply_count >= 3,
"conversation_status": conversation.status,
},
})
# 8. 广播坐席端(new_message + conversation_updated,与 _persist_and_push 一致)
try:
await ws_manager.broadcast({
"type": "new_message",
"data": {
"conversation_id": str(conversation.id),
"message_id": str(ai_message.id),
"sender_type": "ai",
"sender_id": "ai_bot",
"sender_name": "Duckula(达寇拉)",
"content": display_text,
"msg_type": "byod_card",
"extra_data": {"byod_result": byod_result},
},
})
await ws_manager.broadcast({
"type": "conversation_updated",
"data": {
"conversation_id": str(conversation.id),
"status": conversation.status,
"assigned_agent_id": (
str(conversation.assigned_agent_id)
if conversation.assigned_agent_id else None
),
},
})
except Exception as ws_err:
logger.warning(f"BYOD: WS 广播给坐席失败: {ws_err}")
logger.info(
f"BYOD 查询完成: employee_id={employee_id}, position={position}, "
f"eligible={eligible}, matched_category={matched_category}"
)
async def _handle_routing(
db,
conversation: Conversation,
employee_id: str,
content: str,
) -> bool:
"""处理非IT业务路由推荐。
在 H5 聊天消息流中拦截路由关键词后调用 Dify 统一意图识别,
判定为 non_it_routing 且 routing_confidence ≥ 阈值时发送名片三段式消息。
流程:
1. 调用 Dify 统一意图识别(detect_routing_intent
2. 检查 intent_type == "non_it_routing" && routing_confidence ≥ 阈值
→ YES: 查询联系人 → 发送名片三段式消息 → 记录路由事件 → 返回 True
→ NO: 返回 False(继续走正常 AI 流程)
3. Dify 调用失败 → 关键词降级兜底(按 ROUTING_KEYWORD_TO_CATEGORY 映射)
Args:
db: 异步 DB sessionprocess_h5_ai_reply 的 factory session
conversation: 当前会话对象(Conversation
employee_id: 员工企微 UserID
content: 用户消息原文
Returns:
bool: True 表示已发送路由名片(应 return 中断后续流程),
False 表示未触发路由(继续走正常 AI 流程)
"""
from app.config import settings
# 1. 调用 Dify 统一意图识别
try:
result = await detect_routing_intent(content, employee_id)
intent_type = result.get("intent_type", "chitchat")
business_category = result.get("business_category")
routing_confidence = result.get("routing_confidence", 0.0)
# 如果是审批意图,不拦截(让审批流程处理)
if intent_type == "approval":
return False
logger.info(
f"路由意图检测(Dify): intent_type={intent_type}, "
f"business_category={business_category}, "
f"routing_confidence={routing_confidence}"
)
# 2. 检查是否触发路由推荐
threshold = settings.routing_confidence_threshold
if intent_type != "non_it_routing" or routing_confidence < threshold:
# 置信度不足或非路由意图,走正常 AI 流程
return False
if not business_category:
logger.warning("路由意图为 non_it_routing 但 business_category 为空,跳过")
return False
except Exception as e:
logger.warning(f"Dify 路由意图识别失败,降级为关键词匹配: {e}")
# 3. 降级为关键词匹配
business_category = _keyword_fallback_category(content)
if not business_category:
# 关键词也未命中,走正常 AI 流程
return False
routing_confidence = 0.75 # 降级兜底给一个略高于阈值的置信度
logger.info(f"路由意图检测(兜底): business_category={business_category}")
# 4. 查询联系人
contact = await get_contact_by_category(db, business_category)
if not contact:
logger.warning(f"未找到 {business_category} 类别的联系人,跳过路由推荐")
return False
# 5. 构建路由说明文本
category_display = business_category.replace("行政-物业", "物业")
reason = (
f"您的问题属于{category_display}业务范畴,不在IT服务台服务范围内 😊\n\n"
f"为您推荐{category_display}服务相关联系人,您可以直接点击名片联系TA:"
)
# 6. 发送名片三段式消息
await send_contact_card(
db=db,
conversation=conversation,
employee_id=employee_id,
contact=contact,
reason=reason,
business_category=business_category,
routing_confidence=routing_confidence,
)
# 7. 记录路由事件(P1
await record_routing_event(
db=db,
conversation_id=str(conversation.id),
employee_id=employee_id,
message_content=content,
business_category=business_category,
routing_confidence=routing_confidence,
contact=contact,
)
return True
async def _enrich_with_last_ai_context(db, conversation_id: str, content: str) -> str:
"""为简短回复拼接对话上下文,弥补 Dify 工作流缺少「对话历史」节点。
v2.3 改进(相对于 v2.2):
- 移除 15 字符硬限制 → 50 字符宽松阈值(问句/换行/长消息自动跳过)
- 查询最近 10 条消息(用户+AI)→ 构建完整对话摘要,含用户原始问题
- 不再依赖 extra_data.options 判断,对所有简短回复尝试拼接
- 跳过刚保存的当前消息避免重复(当前内容已作为 query 单独传给 Dify)
触发条件:消息不含问号/换行、长度 <= 50 字符 → 可能是选项选择/简短回答。
后续:Dify 工作流配置对话历史节点后,可将 `MAX_CONTEXT_LENGTH` 设为 0 来禁用此修复。
返回:拼接后的消息(如果不需要拼接则返回原内容)
"""
# 宽松的启发式判断:不含问号、不含换行、<= 50 字符 → 可能是简短回答
if "?" in content or "" in content or "\n" in content or len(content) > 50:
return content
try:
# 1. 查询最近 10 条消息(按时间倒序索取最新),含用户和 AI
stmt = (
select(Message.content, Message.sender_type, Message.created_at)
.where(Message.conversation_id == conversation_id)
.order_by(Message.created_at.desc())
.limit(10)
)
result = await db.execute(stmt)
rows = list(result.all())
if not rows or len(rows) < 2:
return content # 消息太少,无法构建有意义上下文
# 2. 反转顺序:最早 → 最新
rows.reverse()
# 3. 跳过最后一条员工消息 → 即刚刚保存的当前消息(避免在上下文中重复)
# content 已作为 query 单独发给 Dify,不应出现在上下文中
if rows and rows[-1][1] == "employee":
rows = rows[:-1]
if len(rows) < 2:
return content # 去掉当前消息后没剩几条,不拼接
# 4. 构建对话摘要(最多保留最近 8 条,避免 prompt 过长)
context_lines = []
for row_text, row_sender, _ in rows[-8:]:
if not row_text:
continue
role = "用户" if row_sender == "employee" else "AI助手"
context_lines.append(f"{role}: {row_text}")
if not context_lines:
return content
context = "\n".join(context_lines)
return (
f"【对话上下文】\n{context}\n\n"
f"请根据以上对话历史回答用户的以下消息:{content}"
)
except Exception as e:
logger.warning(f"上下文拼接失败,使用原消息: {e}")
return content
async def _push_asset_recommends(
db,
employee_id: str,
message: str,
dify_result: dict,
):
"""v3.0 资产推荐推送 - 独立于对话的运维触达通道
功能:
1. L1: 从关键词匹配资产(与当前问题相关)
2. L2: 从画像触发运维提醒(与问题无关)
3. L3: 角色通用资源推荐
Args:
db: 数据库会话
employee_id: 员工 ID
message: 用户消息(用于关键词匹配)
dify_result: Dify 返回结果(包含 intent 等信息)
"""
try:
asset_service = get_asset_recommend_service()
profile_service = get_employee_profile_service()
# 1. L1: 关键词匹配(从用户消息中提取关键词)
l1_recs = asset_service.match_keywords(message)
# 2. L2+L3: 画像匹配(需要获取员工画像)
# 为避免每次都调用第三方 API,先尝试获取画像
# 画像获取失败时只推送 L1
profile = None
try:
profile = await profile_service.get_profile(employee_id)
profile_dict = {
'huorong_version': profile.huorong_version,
'huorong_virusdb_date': profile.huorong_virusdb_date,
'huorong_offline_days': profile.huorong_offline_days,
'unionsoft_patches_missing': profile.unionsoft_patches_missing,
'unionsoft_violations': profile.unionsoft_violations,
}
l2_recs = asset_service.match_profile_triggers(profile_dict)
# L3: 角色通用推荐
role = profile.position or ''
l3_recs = asset_service.get_by_role(role)
for rec in l3_recs:
rec.layer = 'L3'
rec.layer_label = '常用资源'
rec.relevance = 'low'
except Exception as e:
logger.warning(f"[AssetRecommend] 获取画像失败: {e}")
l2_recs = []
l3_recs = []
# 3. 合并所有推荐(去重)
all_recs = l1_recs + l2_recs + l3_recs
if not all_recs:
logger.debug(f"[AssetRecommend] 无推荐: employee={employee_id}")
return
# 4. 构建 WS 消息并推送(添加异常处理避免影响主流程)
try:
ws_msg = asset_service.build_ws_message(all_recs)
await ws_manager.broadcast_to_employees([employee_id], ws_msg)
logger.info(
f"[AssetRecommend] 已推送: employee={employee_id}, "
f"L1={len(l1_recs)}, L2={len(l2_recs)}, L3={len(l3_recs)}"
)
except Exception as ws_err:
logger.warning(f"[AssetRecommend] WS推送失败(不影响主流程): {ws_err}")
except Exception as e:
logger.error(f"[AssetRecommend] 推送失败: {e}", exc_info=True)
# 资产推荐失败不影响主对话流程
# =============================================================================
# 管线步骤函数(v4.0 批次 3process_h5_ai_reply 管线化重构)
# =============================================================================
# 设计:主函数从 200+ 行/11 对 try/except 收敛为 ~40 行编排代码,
# 每步一个函数,步骤内部自管异常,主函数零嵌套。
# 行为承诺:与原 v3.2 实现外部行为一致(仅结构调整 + v3.0 降级结果补 type 字段)。
# =============================================================================
async def _step_load_conversation(db, conversation_id: str):
"""步骤1:加载会话(重试 3 次,处理事务未提交竞态)。"""
for attempt in range(3):
conversation = await db.get(Conversation, conversation_id)
if conversation:
return conversation
if attempt < 2:
await asyncio.sleep(0.5)
# 刷新 session 以看到已提交的数据
await db.rollback()
logger.warning(f"后台 AI 任务:会话不存在(重试3次后) {conversation_id}")
return None
async def _step_byod_intercept(db, conversation, employee_id, content, msg_type) -> bool:
"""步骤2:BYOD 关键词拦截。命中返回 True(终止管线)。"""
if msg_type == "text" and _byod_keyword_prefilter(content):
await _handle_byod_query(db, conversation, employee_id, content)
return True
return False
async def _step_routing_intercept(db, conversation, employee_id, content, msg_type) -> bool:
"""步骤3:非IT业务路由拦截。命中并发送名片返回 True(终止管线)。"""
if msg_type != "text" or not routing_keyword_prefilter(content):
return False
return await _handle_routing(db, conversation, employee_id, content)
async def _step_local_quick_reply(db, conversation, employee_id, content, msg_type, dify_conversation_id) -> bool:
"""步骤4:本地快判断(打招呼/呼叫人工)。命中返回 True(终止管线)。"""
if msg_type != "text":
return False
ai_handler = get_shared_ai_handler()
if not (ai_handler.is_greeting(content) or ai_handler.is_call_human(content)):
return False
result = await ai_handler.handle_message(
content=content,
dify_conversation_id=dify_conversation_id,
user_id=employee_id,
)
await _persist_and_push(
db, conversation, employee_id, result.content,
result.is_guidance, result.should_count,
result.should_transfer, result.dify_conversation_id,
)
return True
async def _step_enrich_image(db, content, msg_type, media_url, conversation_id, employee_id) -> str:
"""步骤5:图片消息增强(VisionService 分析 -> 视觉描述融合)。失败降级为原文。"""
if msg_type != "image" or not media_url:
return content
logger.info(f"图片消息检测: conversation={conversation_id}, media_url={media_url}")
try:
enriched = await _enrich_image_content(
db=db, media_url=media_url, original_content=content,
conversation_id=conversation_id, employee_id=employee_id,
)
logger.info(
f"图片内容增强完成: original_len={len(content)}, enriched_len={len(enriched)}"
)
return enriched
except Exception as vision_err:
logger.error(f"VisionService 处理失败,降级为纯文本: {vision_err}")
# 降级:使用原始 contentAI 会收到 "[图片] 截图" 这样的占位符
return content
async def _step_graph_shortcut(db, conversation, employee_id, enriched_content, msg_type) -> bool:
"""步骤6:Neo4j 图谱短路。命中返回 True(终止管线)。异常降级继续。"""
if msg_type != "text" or not enriched_content:
return False
try:
from app.services.graph_query_service import get_graph_query_service
from app.services.neo4j_client import get_neo4j_client
neo4j_client = await get_neo4j_client()
if not neo4j_client:
return False
graph_service = await get_graph_query_service(neo4j_client)
solution = (
await graph_service.find_solution_by_question(enriched_content)
if graph_service else None
)
if not solution:
return False
logger.info(
f"图谱命中: question={enriched_content[:30]}, solution={solution.action_name}"
)
# 直接返回图谱解决方案,跳过 Dify 调用
await _persist_and_push_solution(db, conversation, employee_id, solution)
return True
except Exception as graph_err:
# 图谱查询失败不阻断,继续原有 Dify 流程
import traceback
logger.warning(f"图谱查询异常(降级继续): {graph_err}\n{traceback.format_exc()}")
return False
def _build_keyword_fallback_result(matched_card, conversation, result=None, source="keyword_fallback"):
"""构建关键词降级结果(v3.0 超时降级 / v3.1 无action降级共用)。
v4.0 批次 3 合并:统一补 "type": "approval_card"v3.0 原缺此字段,
前端 approve-direct-card 渲染依赖它,属于缺陷修复)。
"""
options = matched_card.get("options") or []
return {
"text": f"我来帮您提交{matched_card.get('title', '审批')},请点击下方卡片。",
"action": {
"type": "approval_card",
"card_data": matched_card,
"url": options[0].get("url", "") if options else "",
},
"options": None,
"hit": True,
"conversation_id": (result or {}).get("conversation_id") or conversation.dify_conversation_id,
"is_structured": True,
"diagnosis_stage": "recommending" if result else None,
"response_time_ms": (result or {}).get("response_time_ms", 0),
"source": source,
}
async def _step_notify_failure(conversation_id: str, employee_id: str, message: str):
"""步骤:统一失败通知(ai_reply_failed WS 推送)。"""
try:
await ws_manager.broadcast_to_employees([employee_id], {
"type": "ai_reply_failed",
"data": {"conversation_id": conversation_id, "message": message},
})
except Exception:
# 推送失败也无所谓,员工端 3 秒轮询兜底
pass
async def _step_call_dify(db, conversation, employee_id, conversation_id, content, enriched_content, dify_conversation_id):
"""步骤7Dify 主推理(thinking 指示器 + 30s 超时 + 关键词降级)。
返回结构化结果 dict
超时且降级失败时内部推送 ai_reply_failed 并返回 None(终止管线);
超时但降级成功时已推送降级卡片,返回 None(终止管线)。
"""
# 1. 推送 "正在思考..." 指示器(员工气泡动画 + 坐席状态指示)
try:
await ws_manager.broadcast_to_employees([employee_id], {
"type": "ai_thinking",
"data": {"conversation_id": conversation_id},
})
except Exception as thinking_err:
logger.warning(f"推送 AI 思考指示器失败: {thinking_err}")
try:
await ws_manager.broadcast({
"type": "ai_thinking",
"data": {"conversation_id": conversation_id, "employee_id": employee_id},
})
except Exception:
pass # 坐席端通知失败不影响主流程
# 2. 启动延迟 "仍在思考" 后台任务(15 秒后触发,Dify 返回则取消)
async def _push_still_thinking():
await asyncio.sleep(15)
await ws_manager.broadcast_to_employees([employee_id], {
"type": "ai_thinking",
"data": {"conversation_id": conversation_id, "status": "still_thinking"},
})
thinking_task = asyncio.create_task(_push_still_thinking())
# 3. 调用 Difyblocking + 30s 硬超时)
ai_handler = get_shared_ai_handler()
try:
result = await asyncio.wait_for(
ai_handler.ai_service.get_structured_reply(
message=enriched_content,
conversation_id=dify_conversation_id or conversation.dify_conversation_id,
user_id=employee_id,
),
timeout=30,
)
except asyncio.TimeoutError:
# v3.0: Dify 超时 -> 关键词降级匹配
thinking_task.cancel()
logger.warning(f"Dify 30 秒超时: conversation={conversation_id},尝试关键词降级")
from app.services.approval_matcher import get_approval_matcher
matched_card = get_approval_matcher().match_by_keywords(content)
if matched_card:
logger.info(f"[Fallback] 关键词降级成功: {matched_card.get('title')}")
fallback_result = _build_keyword_fallback_result(matched_card, conversation)
try:
await _persist_and_push_structured(db, conversation, employee_id, fallback_result)
return None # 已推送降级卡片,终止管线
except Exception as fallback_err:
logger.error(f"[Fallback] 降级推送失败: {fallback_err}")
# 降级也失败 -> 建议转人工
await _step_notify_failure(conversation_id, employee_id, "AI 响应时间较长,建议转人工坐席处理。")
return None
# 4. 取消 "仍在思考" 任务(Dify 已返回)
thinking_task.cancel()
try:
await thinking_task # 等待 task 真正取消,避免 warning
except asyncio.CancelledError:
pass
# 5. v3.1: Dify 返回但无 action(如诊断 escalate-> 关键词降级
if not result.get("action"):
from app.services.approval_matcher import get_approval_matcher
matched_card = get_approval_matcher().match_by_keywords(content)
if matched_card:
logger.info(f"[Fallback-v3.1] Dify 无action但关键词命中: {matched_card.get('title')}")
result = _build_keyword_fallback_result(matched_card, conversation, result, "keyword_fallback_v3")
return result
async def _step_persist(db, conversation, employee_id, result):
"""步骤8:持久化 + 双 WS 推送(ai_reply + dynamic_recommend)。异常不中断管线。"""
try:
await _persist_and_push_structured(db, conversation, employee_id, result)
except Exception as persist_err:
logger.error(f"[Persist] AI回复持久化失败: {persist_err}", exc_info=True)
async def _step_assets(db, employee_id, content, result):
"""步骤9:资产推荐推送(L1/L2/L3 分层,独立运维触达通道)。异常不中断管线。"""
try:
await _push_asset_recommends(db, employee_id, content, result)
except Exception as asset_err:
logger.error(f"[Asset] 资产推荐推送失败: {asset_err}", exc_info=True)
# =============================================================================
# 管线编排主函数(v4.0 批次 3 重构版)
# =============================================================================
async def process_h5_ai_reply(
conversation_id: str,
employee_id: str,
content: str,
dify_conversation_id=None,
msg_type: str = "text",
media_url: str = None,
):
"""H5 发送消息后的 AI 回复处理(asyncio.create_task 入口)。
v4.0 批次 3 管线化重构:原 12 步/11 对 try/except 编排为 9 个步骤函数,
主函数仅最外层 1 个 try/except,行为与 v3.2 外部表现一致。
流程:
1. _step_load_conversation 加载会话(重试3次)
2. _step_byod_intercept BYOD 关键词拦截
3. _step_routing_intercept 非IT业务路由拦截(名片推荐)
4. _step_local_quick_reply 本地快判断(打招呼/呼叫人工)
5. _step_enrich_image 图片消息 VisionService 增强
5b. _enrich_with_last_ai_context 简短回复上下文拼接(v2.3)
6. _step_graph_shortcut Neo4j 图谱短路
7. _step_call_dify Dify 主推理(含超时/无action关键词降级)
8. _step_persist 持久化 + 双 WS 推送
9. _step_assets 资产推荐推送
Args:
conversation_id: 会话 ID
employee_id: 员工企微 UserID
content: 消息文本内容
dify_conversation_id: Dify 会话 ID(用于多轮上下文)
msg_type: 消息类型(text/image/file),默认 text
media_url: 媒体文件 URL(图片消息时使用)
"""
factory = _get_session_factory()
async with factory() as db:
try:
conversation = await _step_load_conversation(db, conversation_id)
if not conversation:
return
# 前置拦截管线(任一命中即终止)
if await _step_byod_intercept(db, conversation, employee_id, content, msg_type):
return
if await _step_routing_intercept(db, conversation, employee_id, content, msg_type):
return
if await _step_local_quick_reply(db, conversation, employee_id, content, msg_type, dify_conversation_id):
return
# 内容增强管线
enriched_content = await _step_enrich_image(
db, content, msg_type, media_url, conversation_id, employee_id,
)
enriched_content = await _enrich_with_last_ai_context(
db, conversation_id, enriched_content,
)
# 图谱短路
if await _step_graph_shortcut(db, conversation, employee_id, enriched_content, msg_type):
return
# 主推理(含 thinking 指示器 + 超时/无action 降级)
result = await _step_call_dify(
db, conversation, employee_id, conversation_id,
content, enriched_content, dify_conversation_id,
)
if result is None:
return # 降级路径已全部处理(超时转人工 或 已推送降级卡片)
# 后置处理管线
await _step_persist(db, conversation, employee_id, result)
await _step_assets(db, employee_id, content, result)
except Exception as e:
import traceback
logger.error(f"后台 AI 任务异常: {e}\n堆栈跟踪:\n{traceback.format_exc()}", exc_info=True)
await _step_notify_failure(conversation_id, employee_id, "AI 服务异常,请转人工坐席或稍后重试。")