Files
wecom_it_smart_desk/backend/app/services/triage_service.py
Simon 449c6d4875 feat: 2026-07-12~13 全量更新 - AI对话链路改造+H5 v4/v5+坐席端v5+上下文感知诊断+知识库迭代3
## H5 员工端 v4 (2026-07-13 00:48 已部署)
- 人工按钮三态文案统一为"人工坐席"
- 按钮位置移至发送键和语音按钮上方(垂直堆叠)
- 点按钮直接调 store.shakeAgent(),删除 CallAgentModal 弹窗动画
- 截图快捷键提示改为"截图->粘贴:Alt+Shift+A-Ctrl+V ---> Ctrl+V"
- 移动端隐藏截图提示(CSS 媒体查询)
- AI转人工提示改为"已为您呼叫人工坐席,请稍等!"
- 坐席接入提示改为"坐席正在查看您的信息,请等待处理回复!"
- 删除"摇铃呼叫坐席"入口和文案
- 删除孤儿组件 MessageList.vue + shake 动画 CSS

## H5 员工端 v5 (2026-07-13 02:08 已部署)
- RightPanel v2.1:删除"软件安装"和"资源权限"标签页
- 移除标签栏,智能推荐(DynamicRecommend)直接展示
- 删除 SoftwareDownloads/ApprovalLinks 引用和相关 CSS

## AI 对话链路全栈改造 Phase 1-6 (已部署)
- Phase 1: Dify JSON输出 + 后端blocking解析 + 双WS推送 + 错误降级
- Phase 2: 关键词收窄(~25强意图词) + 两级分类Prompt + 删除前端checkApprovalIntent
- Phase 3: WS扩展(ai_thinking+dynamic_recommend) + ai_structured气泡 + RightPanel v2 + 选项回传
- Phase 4: VisionService接入 + 图片消息融合(5秒窗口) + 降级策略
- Phase 5: 坐席端ai_thinking指示器 + ai_structured/byod_card渲染 + handleNewMessage修复
- Phase 6: diagnosis_stage(6值) + response_time_ms计时 + 慢响应告警(>10s)

## 坐席端 v5 (2026-07-13 01:38 已部署)
- ai_structured/byod_card 只读渲染
- AI思考指示器 UI
- handleNewMessage 透传 msg_type/extra_data 修复
- 布局优化v2.0: QuickReplyBar L1+L2悬浮 + ReplyBox左右分区 + 右栏260/560px切换
- 键盘快捷键v2.3: 纯数字路由 + ESC分层撤销 + Shift+Space用event.code

## 上下文感知智能诊断闭环 (2026-07-12 已部署)
- 三层诊断(API→Script→AI) + 三段排队(VIP→info_locked→not locked)
- 答题插队 + 五场景关闭
- 迁移052(6表+6列) + queue_service + quiz_service + closing_service
- H5前端: QueueWaiting + RightPanel双Tab + InputBar三态 + ResolveConfirmCard
- 坐席前端: pending_close结单流程 + 信息锁定(Dify步骤完成+有效回答率≥70%)

## 知识库迭代3 (2026-07-12 已部署)
- 分诊交互(H5+坐席+Dify独立应用)
- 拓扑预览(ECharts只读)
- 代答排除(4种匹配器: keyword/regex/intent/category)
- 迁移051 + 44文件43测试通过

## 后端变更
- 6个Python文件改造(h5_ai_task.py/h5.py/ai_service.py/closing_service.py等)
- funny_phrase_service.py: shake/connected/keyword 默认文案更新
- session_service.py: 企微消息文案同步
- 新增: queue.py/quiz.py/triage.py/exclusion_rules.py 等API端点
- 新增: diagnostic.py/quiz.py/triage_session.py 等模型
- 新增: closing_service/queue_service/quiz_service/triage_service 等服务

## 文档更新
- CHANGELOG.md: 新增 [未发布] 区全部变更记录
- 项目管理主文档 v2.5: 新增v0.7.3版本 + 已完成看板 + 最近搞定
- 版本记录: 新增v0.7.3条目
- AI对话链路实施计划: Phase 1-6 全部标记已实施
- 新增架构图/时序图/类图(mermaid)

## 部署路径修正
- 服务器项目根路径: /opt/wecom-it-desk/
- 所有前端dist均为ro bind mount,只能在宿主机源路径操作
- 服务器nginx /h5/ 是静态文件服务(非proxy_pass)
- elFinder上传二进制不可靠(MD5不匹配),改用base64分块上传
2026-07-13 02:17:03 +08:00

991 lines
34 KiB
Python
Raw Permalink 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智能服务台 — 分诊业务逻辑服务
# =============================================================================
# 说明:分诊交互的核心业务逻辑,包括:
# 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