Files
wecom_it_smart_desk/backend/app/api/h5.py
T

2108 lines
81 KiB
Python
Raw Normal View History

# =============================================================================
# 企微IT智能服务台 — H5 用户端 API
# =============================================================================
# 说明:H5 用户端的接口,包括:
# 1. GET /api/h5/oauth/authorize — 获取企微OAuth2授权URL
# 2. POST /api/h5/oauth/callback — OAuth2回调,返回token+用户信息
# 3. GET /api/h5/me — 获取当前用户详细信息
# 4. GET /api/h5/user — 获取当前用户信息(兼容旧接口)
# 5. GET /api/h5/conversations/current — 获取当前会话
# 6. POST /api/h5/conversations/current/messages — 用户发送消息
# 7. GET /api/h5/conversations/current/messages/poll — 用户轮询新消息
# 8. POST /api/h5/conversations/current/shake — 举手(敲桌子呼叫坐席)
# 9. GET /api/h5/approval-links — 获取审批流程链接
# 10. GET /api/h5/software-downloads — 获取软件下载列表
#
# 重构记录(2026-06):
# - 移除 _get_redis() 手动创建 Redis 模式,改用 DI 共享实例
# - 移除本地打招呼/呼叫人工检测逻辑,改用 AIHandler 统一处理
# - 移除本地 AI 调用/计数/降级逻辑,改用 AIHandler 统一处理
# - 所有服务实例通过 FastAPI Depends 注入,不再手动创建/关闭
# =============================================================================
import json
import logging
import re
import secrets
from datetime import datetime
from typing import Any, Dict, List, Optional
from urllib.parse import quote
from uuid import UUID
import redis
import redis.asyncio as aioredis
from fastapi import APIRouter, Depends, Header, Query, Request
from slowapi import Limiter
from slowapi.util import get_remote_address
from sqlalchemy import select
from sqlalchemy.ext.asyncio import AsyncSession
# 速率限制器实例
# 移除 env_file=None 参数:slowapi 0.1.9 不支持该参数
# python-dotenv 已在应用启动时处理 .env 文件
limiter = Limiter(key_func=get_remote_address)
from app.config import settings
from app.database import get_db
from app.utils.env_gating import is_production
from app.dependencies import dep_redis, dep_wecom_service
from app.models.approval_link import ApprovalLink
from app.models.conversation import Conversation
from app.models.message import Message
from app.models.software_download import SoftwareDownload
from app.schemas.h5 import (
ApprovalLinkResponse,
OAuthCallbackRequest,
ShakeRequest,
SoftwareDownloadResponse,
)
from app.schemas.conversation import ConversationResponse, InviteParticipantRequest, JoinConversationRequest
from app.schemas.message import MessageResponse
import asyncio
from app.tasks.h5_ai_task import process_h5_ai_reply
from app.services.funny_phrase_service import FunnyPhraseService
from app.services.ws_manager import manager as ws_manager
from app.services.wecom_service import WecomService
from app.services.employee_directory import (
count_tree_employees,
get_org_directory,
get_org_tree_cached,
)
from app.utils.response import AppException, ERR_UNAUTHORIZED, success_response
from app.services.closing_service import ClosingService
from pydantic import BaseModel, Field
logger = logging.getLogger(__name__)
# 创建路由器
router = APIRouter()
# H5 员工端 token TTL:8小时(与坐席端一致)
EMPLOYEE_TOKEN_TTL_SECONDS = 8 * 60 * 60 # 8小时
# --------------------------------------------------------------------------
# 辅助:检测请求是否来自企微 WebView(后端第二道防线)
# --------------------------------------------------------------------------
# 企微桌面端 UA 示例:Mozilla/5.0 ... wxwork/4.1.22 ...
# 企微移动端 UA 示例:Mozilla/5.0 (iPhone ... MicroMessenger/7.x ... wxwork/3.x ...
_WEWORK_UA_RE = re.compile(r"wxwork", re.IGNORECASE)
def _require_wework_ua(request: Request) -> None:
"""校验请求 User-Agent 是否来自企微 WebView。
仅生产环境强制校验(env_gating.is_production()):
非企微环境的 OAuth2 请求直接拒绝。
本地开发 / dev / test 环境跳过检测,方便调试。
Args:
request: FastAPI Request 对象,用于读取 User-Agent
Raises:
AppException: 非企微环境且处于生产环境时抛出 4003 错误
"""
# 仅生产环境强制校验;非生产环境(dev/test/本地)一律放行
if not is_production():
return
ua = request.headers.get("user-agent", "")
if not _WEWORK_UA_RE.search(ua):
raise AppException(4003, "请在企业微信中访问此服务")
# --------------------------------------------------------------------------
# 辅助:从请求头获取员工ID(旧版,仅作为过渡期兼容)
# --------------------------------------------------------------------------
def _get_employee_id(
x_employee_id: Optional[str] = Header(None, alias="X-Employee-Id"),
) -> str:
"""从请求头获取员工ID(旧版兼容,已废弃)。
安全警告:此方法直接信任请求头中的明文员工ID,任何人可伪造身份。
仅在本地开发环境(mock_login_enabled=true)下允许使用。
生产环境必须使用 _get_current_employeeBearer Token 认证)。
Args:
x_employee_id: 请求头中的员工ID
Returns:
str: 员工企微 UserID
Raises:
AppException: 未提供员工ID 或 生产环境禁用
"""
# 生产环境禁止使用明文头认证(可被任意伪造)
if not settings.mock_login_enabled:
raise AppException(
4001,
"X-Employee-Id 认证方式已在生产环境禁用,请使用 OAuth2 登录"
)
if not x_employee_id:
raise ERR_UNAUTHORIZED
return x_employee_id
# --------------------------------------------------------------------------
# 辅助:从 Bearer Token 获取当前员工ID(新版,替换 _get_employee_id
# --------------------------------------------------------------------------
async def _get_current_employee(
authorization: Optional[str] = Header(None, alias="Authorization"),
x_employee_id: Optional[str] = Header(None, alias="X-Employee-Id"),
redis_client: Optional[aioredis.Redis] = Depends(dep_redis),
) -> str:
"""从请求头中提取员工身份(认证依赖)。
认证优先级:
1. Bearer Token(生产环境):从 Redis 查找对应的 employee_id
2. X-Employee-Id 头(开发降级):直接读取 employee_id(仅本地开发使用)
Token 存储格式:
Redis key: employee:token:{token}
Redis value: employee_id (企微 UserID)
重构说明:不再手动创建/关闭 Redis 客户端,改用 DI 注入共享实例。
Args:
authorization: 请求头中的 Authorization 字段(格式:Bearer token
x_employee_id: 请求头中的 X-Employee-Id 字段(开发降级用)
redis_client: 共享 Redis 客户端(DI 注入)
Returns:
str: 员工企微 UserID
Raises:
AppException: 未授权(无有效认证)
"""
# =====================================================================
# 方式1Bearer Token 认证(生产环境)
# =====================================================================
if authorization:
token = authorization.replace("Bearer ", "") if authorization.startswith("Bearer ") else authorization
if token:
if redis_client:
try:
employee_id_bytes = await redis_client.get(f"employee:token:{token}")
if employee_id_bytes:
# Redis 返回 bytes,需要解码
return employee_id_bytes.decode("utf-8") if isinstance(employee_id_bytes, bytes) else employee_id_bytes
except AppException:
raise
except redis.exceptions.TimeoutError as e:
logger.error(f"Redis 连接超时,无法验证 Bearer Token: {e}")
# Redis 不可用时:开发模式下降级使用 X-Employee-Id 明文头
if x_employee_id and settings.mock_login_enabled:
logger.warning("Redis 超时,降级使用 X-Employee-Id 明文认证(仅开发环境)")
return x_employee_id
# 生产环境:返回明确的业务错误码,而非让框架转 500
raise AppException(code=1003, message="认证服务暂不可用,请稍后重试")
except Exception as e:
logger.error(f"Redis 读取失败,无法验证 Bearer Token: {e}")
# Redis 不可用时:开发模式下降级使用 X-Employee-Id 明文头
if x_employee_id and settings.mock_login_enabled:
logger.warning("Redis 异常,降级使用 X-Employee-Id 明文认证(仅开发环境)")
return x_employee_id
# 生产环境:返回明确的业务错误码,而非让框架转 500
raise AppException(code=1003, message="认证服务暂不可用,请稍后重试")
else:
# redis_client 为 Nonedep_redis 返回了 NoneRedis 连接创建失败)
logger.error("Redis 客户端不可用,无法验证 Bearer Token")
if x_employee_id and settings.mock_login_enabled:
logger.warning("Redis 不可用,降级使用 X-Employee-Id 明文认证(仅开发环境)")
return x_employee_id
raise AppException(code=1003, message="认证服务暂不可用,请稍后重试")
# =====================================================================
# 方式2X-Employee-Id 明文头(仅开发环境,生产环境禁用)
# =====================================================================
# 安全说明:X-Employee-Id 可被任意 HTTP 客户端伪造,不能用于身份认证
# 仅在 MOCK_LOGIN_ENABLED=true 时允许,方便本地开发调试
if x_employee_id and settings.mock_login_enabled:
return x_employee_id
raise ERR_UNAUTHORIZED
# --------------------------------------------------------------------------
# GET /api/h5/oauth/authorize — 获取企微OAuth2授权URL
# --------------------------------------------------------------------------
@router.get("/h5/oauth/authorize")
async def get_oauth_authorize_url(
request: Request,
redirect_uri: Optional[str] = Query(None, description="OAuth2回调地址(可选,默认使用请求来源域名/h5/)"),
request_host: Optional[str] = Header(None, alias="Host"),
):
"""获取企微OAuth2授权URL。
前端调用此接口获取完整的企微OAuth2授权链接,
然后跳转到该链接进行静默授权。
授权流程:
1. 前端请求此接口获取授权URL
2. 前端跳转到授权URL
3. 企微自动重定向到 redirect_uri?code=CODE&state=STATE
4. 前端拿到code,调用 POST /api/h5/oauth/callback
Args:
request: FastAPI Request 对象(用于 UA 检测)
redirect_uri: 自定义回调地址(可选)
request_host: 请求的 Host 头(自动获取,用于构造默认回调地址)
Returns:
Dict: 统一响应格式,包含 authorize_url 字段
"""
# 后端第二道防线:非企微环境拒绝授权
_require_wework_ua(request)
corp_id = settings.wecom_corp_id
# 确定回调地址:优先使用参数传入的,否则根据 Host 头构造
if redirect_uri:
encoded_redirect = quote(redirect_uri, safe="")
elif request_host:
# 从 Host 头构造回调地址(支持 http 和 https)
scheme = "https" # 企微H5应用通常使用 https
encoded_redirect = quote(f"{scheme}://{request_host}/itdesk/", safe="")
else:
# 最终降级:使用配置中的 CORS 源地址
default_origin = settings.cors_origins_list[0] if settings.cors_origins_list else "https://localhost"
encoded_redirect = quote(f"{default_origin}/itdesk/", safe="")
# 构造企微OAuth2静默授权URLsnsapi_base:用户无感知)
# 企业微信 OAuth2 地址(注意是 open.work.weixin.qq.com
authorize_url = (
f"https://open.work.weixin.qq.com/connect/oauth2/authorize"
f"?appid={corp_id}"
f"&redirect_uri={encoded_redirect}"
f"&response_type=code"
f"&scope=snsapi_base"
f"&state=STATE"
f"#wechat_redirect"
)
return success_response(data={"authorize_url": authorize_url})
# --------------------------------------------------------------------------
# POST /api/h5/oauth/callback — OAuth2 回调
# --------------------------------------------------------------------------
@router.post("/h5/oauth/callback")
@limiter.limit("20/minute") # OAuth 回调限流:正常用户不会频繁触发
async def oauth_callback(
request: Request,
body: OAuthCallbackRequest,
db: AsyncSession = Depends(get_db),
redis_client: Optional[aioredis.Redis] = Depends(dep_redis),
wecom_service: WecomService = Depends(dep_wecom_service),
):
"""企微 OAuth2 授权回调。
H5 页面通过企微 OAuth2 静默授权获取 code,后端用 code 换取员工身份。
成功后生成 Bearer Token 存入 Redis,返回 token + 员工信息。
重构说明:不再手动创建/关闭 Redis 和 WecomService,改用 DI 注入共享实例。
流程:
1. 前端跳转企微授权页面
2. 企微回调到 H5 页面并携带 code
3. H5 前端将 code 发给后端
4. 后端用 code 调用企微 API 换取员工 UserID
5. 后端获取员工详细信息(姓名、部门、岗位等)
6. 生成 Bearer Token 存入 Redis
7. 返回 token + 员工信息
Args:
request: FastAPI Request 对象(用于 UA 检测)
body: OAuth2 回调请求体(包含 code
db: 数据库会话
redis_client: 共享 Redis 客户端(DI 注入)
wecom_service: 共享企微服务(DI 注入)
Returns:
Dict: 统一响应格式,包含 token 和员工信息
"""
# 后端第二道防线:非企微环境拒绝回调
_require_wework_ua(request)
try:
# 1. 用 code 换取员工身份
user_info = await wecom_service.get_oauth_user_info(body.code)
employee_id = user_info.get("userid", "")
if not employee_id:
raise AppException(2007, "OAuth2授权失败:未获取到员工ID")
# 2. 获取员工详细信息
employee_name = ""
department = ""
position = ""
avatar = ""
# 2.1 获取员工详细信息(包含头像)
try:
detail = await wecom_service.get_user_info(employee_id)
employee_name = detail.get("name", "")
# department 返回的是部门ID列表,取第一个部门名称需要额外API调用
# 简化处理:将部门ID列表转为逗号分隔的字符串
dept_ids = detail.get("department", [])
department = ",".join(str(d) for d in dept_ids) if dept_ids else ""
position = detail.get("position", "")
avatar = detail.get("avatar", "")
# 调试日志:检查企微API返回的完整数据
logger.info(f"企微用户详情返回: employee_id={employee_id}, avatar={avatar[:50] if avatar else '(空)'}...")
except Exception as e:
logger.warning(f"获取员工详细信息失败: employee_id={employee_id}, error={e}")
# 2.2 将员工信息保存到数据库(包含头像)
try:
from app.models.employee import Employee
from sqlalchemy import select
stmt = select(Employee).where(
Employee.employee_id == employee_id,
Employee.corp_id == settings.wecom_corp_id
)
result = await db.execute(stmt)
employee = result.scalars().first()
if employee:
# 更新已有记录
employee.name = employee_name
employee.department = department
employee.position = position
# 【FE-UA-005 优化】每次登录强制更新头像URL + 清缓存(统一走 avatar_service
# 头像更新失败不阻塞登录(sync_employee_avatar 内部已容错)
if avatar:
try:
from app.services.avatar_service import sync_employee_avatar
await sync_employee_avatar(db, redis_client, employee_id, avatar)
except Exception as e:
logger.warning(f"同步员工头像失败(不阻塞登录): employee_id={employee_id}, error={e}")
else:
# 创建新记录
employee = Employee(
corp_id=settings.wecom_corp_id,
employee_id=employee_id,
name=employee_name,
department=department,
position=position,
avatar=avatar,
)
db.add(employee)
logger.info(f"创建员工记录(含头像): employee_id={employee_id}")
await db.commit()
except Exception as e:
logger.warning(f"保存员工信息到数据库失败: employee_id={employee_id}, error={e}")
# 不阻塞登录流程,继续执行
# 3. 生成 Bearer Token(与坐席端一致:secrets.token_urlsafe(32)
token = secrets.token_urlsafe(32)
# 4. Token 存入 Rediskey: employee:token:{token}, value: employee_id, TTL 8小时)
if redis_client:
try:
await redis_client.setex(
f"employee:token:{token}",
EMPLOYEE_TOKEN_TTL_SECONDS,
employee_id,
)
except Exception as e:
logger.warning(f"Redis 写入失败(token 不会持久化): {e}")
# 5. 缓存员工基本信息到 Redis(用于快速读取,避免频繁调用企微API)
employee_info_cache = {
"employee_id": employee_id,
"employee_name": employee_name,
"department": department,
"position": position,
"avatar": avatar,
}
try:
await redis_client.setex(
f"employee:info:{employee_id}",
EMPLOYEE_TOKEN_TTL_SECONDS,
json.dumps(employee_info_cache, ensure_ascii=False),
)
except Exception as e:
logger.warning(f"员工信息缓存写入失败(不阻塞流程): {e}")
logger.info(f"OAuth2授权成功: employee_id={employee_id}, name={employee_name}")
# 6. 返回 token + 员工信息
return success_response(
data={
"employee_id": employee_id,
"employee_name": employee_name,
"token": token,
"department": department,
"position": position,
"avatar": avatar,
}
)
except AppException:
raise
except Exception as e:
logger.error(f"OAuth2回调处理失败: {e}")
raise AppException(2007, f"OAuth2授权失败: {e}")
# --------------------------------------------------------------------------
# GET /api/h5/oauth/sns-callback — 企微 OAuth2 静默授权回调(302 重定向版)
# --------------------------------------------------------------------------
@router.get("/h5/oauth/sns-callback")
async def oauth_sns_callback(
request: Request,
code: str = Query(..., description="企微 OAuth2 授权码"),
state: Optional[str] = Query(None, description="透传参数(保留兼容,未使用)"),
db: AsyncSession = Depends(get_db),
redis_client: Optional[aioredis.Redis] = Depends(dep_redis),
wecom_service: WecomService = Depends(dep_wecom_service),
):
"""企微 OAuth2 静默授权回调(snsapi_base → 302 带 ?token=)。
适用于 snsapi_base 静默授权:企微回调到此端点并携带 code,
后端用 code 换取员工身份 → 生成 Bearer Token → 302 重定向到
H5 前端页面,并在 URL 上附带 ?token=,供前端镜像到 localStorage。
仅生产环境强制 UA 校验(与 _require_wework_ua 一致,使用 env_gating)。
Args:
code: 企微授权码
state: 透传参数(未使用,保留兼容)
db: 数据库会话
redis_client: 共享 Redis 客户端(DI 注入)
wecom_service: 共享企微服务(DI 注入)
Returns:
RedirectResponse -> {scheme}://{host}/itdesk/?token={token}
"""
# 仅生产环境强制 UA 校验
if is_production():
ua = request.headers.get("user-agent", "")
if not _WEWORK_UA_RE.search(ua):
raise AppException(4003, "请在企业微信中访问此服务")
# 1. 用 code 换取员工身份
user_info = await wecom_service.get_oauth_user_info(code)
employee_id = user_info.get("userid", "")
if not employee_id:
raise AppException(2007, "OAuth2授权失败:未获取到员工ID")
# 2. 获取员工详细信息(姓名、部门、岗位、头像)
employee_name = ""
department = ""
position = ""
avatar = ""
try:
detail = await wecom_service.get_user_info(employee_id)
employee_name = detail.get("name", "")
dept_ids = detail.get("department", [])
department = ",".join(str(d) for d in dept_ids) if dept_ids else ""
position = detail.get("position", "")
avatar = detail.get("avatar", "")
except Exception as e:
logger.warning(f"获取员工详细信息失败: employee_id={employee_id}, error={e}")
# 3. 落库 / 更新员工信息(含头像)
try:
from app.models.employee import Employee
stmt = select(Employee).where(
Employee.employee_id == employee_id,
Employee.corp_id == settings.wecom_corp_id,
)
result = await db.execute(stmt)
employee = result.scalars().first()
if employee:
employee.name = employee_name
employee.department = department
employee.position = position
if avatar:
try:
from app.services.avatar_service import sync_employee_avatar
await sync_employee_avatar(db, redis_client, employee_id, avatar)
except Exception as e:
logger.warning(f"同步员工头像失败(不阻塞登录): employee_id={employee_id}, error={e}")
else:
employee = Employee(
corp_id=settings.wecom_corp_id,
employee_id=employee_id,
name=employee_name,
department=department,
position=position,
avatar=avatar,
)
db.add(employee)
await db.commit()
except Exception as e:
logger.warning(f"保存员工信息到数据库失败: employee_id={employee_id}, error={e}")
# 4. 生成 Bearer Token 并写入 Redis
token = secrets.token_urlsafe(32)
if redis_client:
try:
await redis_client.setex(
f"employee:token:{token}",
EMPLOYEE_TOKEN_TTL_SECONDS,
employee_id,
)
except Exception as e:
logger.warning(f"Redis 写入失败(token 不会持久化): {e}")
employee_info_cache = {
"employee_id": employee_id,
"employee_name": employee_name,
"department": department,
"position": position,
"avatar": avatar,
}
try:
await redis_client.setex(
f"employee:info:{employee_id}",
EMPLOYEE_TOKEN_TTL_SECONDS,
json.dumps(employee_info_cache, ensure_ascii=False),
)
except Exception as e:
logger.warning(f"员工信息缓存写入失败(不阻塞流程): {e}")
logger.info(f"OAuth2 sns-callback 授权成功: employee_id={employee_id}, name={employee_name}")
# 5. 302 重定向到 H5 前端页面,附带 ?token= 供前端镜像到 localStorage
from fastapi.responses import RedirectResponse
host = request.headers.get("host", "")
scheme = "https"
landing = "/itdesk/"
redirect_url = f"{scheme}://{host}{landing}?token={token}"
return RedirectResponse(url=redirect_url)
# --------------------------------------------------------------------------
# POST /api/h5/mock-login — Mock 登录(测试阶段,跳过 OAuth2)
# --------------------------------------------------------------------------
@router.post("/h5/mock-login")
@limiter.limit("5/minute") # Mock 登录严格限流:每IP每分钟最多5次
async def mock_login(
request: Request,
body: dict,
redis_client: Optional[aioredis.Redis] = Depends(dep_redis),
):
"""Mock 登录(测试阶段使用,跳过企微 OAuth2)。
仅当后端配置 MOCK_LOGIN_ENABLED=true 时可用。
直接通过员工 ID 生成 Bearer Token,并存入 Redis。
返回格式与 OAuth2 回调完全一致。
Args:
body: 请求体 { employee_id: str, employee_name: str }
redis_client: 共享 Redis 客户端(DI 注入)
Returns:
Dict: 统一响应格式,包含 token 和员工信息
"""
if not settings.mock_login_enabled:
raise AppException(2007, "Mock 登录未启用,请联系管理员")
employee_id = body.get("employee_id", "").strip()
employee_name = body.get("employee_name", "测试用户").strip()
if not employee_id:
raise AppException(2007, "请提供 employee_id")
# 生成 Bearer Token
token = secrets.token_urlsafe(32)
# Token 存入 Rediskey: employee:token:{token}, value: employee_id, TTL 8小时)
if redis_client:
try:
await redis_client.setex(
f"employee:token:{token}",
EMPLOYEE_TOKEN_TTL_SECONDS,
employee_id,
)
except Exception as e:
logger.warning(f"Redis 写入失败(token 不会持久化): {e}")
# 缓存员工基本信息到 Redis
employee_info_cache = {
"employee_id": employee_id,
"employee_name": employee_name,
"department": "IT部",
"position": "测试岗位",
"avatar": "",
}
try:
await redis_client.setex(
f"employee:info:{employee_id}",
EMPLOYEE_TOKEN_TTL_SECONDS,
json.dumps(employee_info_cache, ensure_ascii=False),
)
except Exception as e:
logger.warning(f"员工信息缓存写入失败(不阻塞流程): {e}")
logger.info(f"Mock 登录成功: employee_id={employee_id}, name={employee_name}")
return success_response(
data={
"employee_id": employee_id,
"employee_name": employee_name,
"token": token,
"department": "IT部",
"position": "测试岗位",
"avatar": "",
}
)
# --------------------------------------------------------------------------
# GET /api/h5/me — 获取当前用户详细信息
# --------------------------------------------------------------------------
@router.get("/h5/me")
async def get_current_employee_info(
employee_id: str = Depends(_get_current_employee),
db: AsyncSession = Depends(get_db),
redis_client: Optional[aioredis.Redis] = Depends(dep_redis),
wecom_service: WecomService = Depends(dep_wecom_service),
):
"""获取当前登录员工的详细信息。
需要在请求头中携带有效的 Bearer token。
优先从 Redis 缓存读取,缓存不存在则调用企微API获取。
重构说明:不再手动创建/关闭 Redis 和 WecomService,改用 DI 注入共享实例。
Args:
employee_id: 员工企微 UserID(通过认证依赖注入)
db: 数据库会话
redis_client: 共享 Redis 客户端(DI 注入)
wecom_service: 共享企微服务(DI 注入)
Returns:
Dict: 统一响应格式,包含员工详细信息
"""
# 1. 优先从 Redis 缓存读取
if redis_client:
try:
cached_info = await redis_client.get(f"employee:info:{employee_id}")
if cached_info:
info_str = cached_info.decode("utf-8") if isinstance(cached_info, bytes) else cached_info
info = json.loads(info_str)
# 补充 is_vip 字段
info["is_vip"] = False
return success_response(data=info)
except Exception as e:
logger.warning(f"从Redis读取员工信息缓存失败: {e}")
# 2. 缓存不存在,调用企微API获取
try:
detail = await wecom_service.get_user_info(employee_id)
employee_name = detail.get("name", "")
dept_ids = detail.get("department", [])
department = ",".join(str(d) for d in dept_ids) if dept_ids else ""
position = detail.get("position", "")
avatar = detail.get("avatar", "")
mobile = detail.get("mobile", "")
email = detail.get("email", "")
# 写入缓存
employee_info = {
"employee_id": employee_id,
"employee_name": employee_name,
"department": department,
"position": position,
"mobile": mobile,
"email": email,
"avatar": avatar,
"is_vip": False,
}
if redis_client:
try:
await redis_client.setex(
f"employee:info:{employee_id}",
EMPLOYEE_TOKEN_TTL_SECONDS,
json.dumps(employee_info, ensure_ascii=False),
)
except Exception:
pass
return success_response(data=employee_info)
except AppException:
raise
except Exception as e:
logger.error(f"获取员工信息失败: employee_id={employee_id}, error={e}")
raise AppException(2006, f"获取员工信息失败: {e}")
# --------------------------------------------------------------------------
# GET /api/h5/user — 获取当前用户信息(兼容旧接口)
# --------------------------------------------------------------------------
@router.get("/h5/user")
async def get_current_user(
employee_id: str = Depends(_get_current_employee),
db: AsyncSession = Depends(get_db),
):
"""获取当前用户信息。
通过 Bearer Token 认证后获取员工信息。
Args:
employee_id: 员工企微 UserID(通过认证依赖注入)
db: 数据库会话
Returns:
Dict: 统一响应格式,包含员工信息
"""
# 尝试从会话记录获取员工信息
stmt = select(Conversation).where(
Conversation.employee_id == employee_id
).order_by(Conversation.created_at.desc())
result = await db.execute(stmt)
latest_conv = result.scalars().first()
user_info = {
"employee_id": employee_id,
"employee_name": latest_conv.employee_name if latest_conv else "",
"department": latest_conv.department if latest_conv else "",
"position": latest_conv.position if latest_conv else "",
"is_vip": latest_conv.is_vip if latest_conv else False,
}
return success_response(data=user_info)
# --------------------------------------------------------------------------
# GET /api/h5/conversations/current — 获取当前会话
# --------------------------------------------------------------------------
@router.get("/h5/conversations/current")
async def get_current_conversation(
employee_id: str = Depends(_get_current_employee),
db: AsyncSession = Depends(get_db),
):
"""获取当前用户的活跃会话。
查找员工当前状态为 ai_handling、queued 或 serving 的会话。
如果没有活跃会话,返回空数据。
Args:
employee_id: 员工企微 UserID
db: 数据库会话
Returns:
Dict: 统一响应格式,包含会话信息
"""
stmt = select(Conversation).where(
Conversation.employee_id == employee_id,
Conversation.status.in_(["ai_handling", "queued", "serving"]),
).order_by(Conversation.created_at.desc())
result = await db.execute(stmt)
conversation = result.scalars().first()
if not conversation:
return success_response(data=None)
conv_data = ConversationResponse.model_validate(conversation).model_dump()
# 附加「是否可以呼叫坐席」标志(AI实质性回复 >= 3)
conv_data["can_call_agent"] = conversation.ai_substantive_reply_count >= 3
conv_data["ai_substantive_reply_count"] = conversation.ai_substantive_reply_count
return success_response(data=conv_data)
# --------------------------------------------------------------------------
# POST /api/h5/conversations/current/messages — 用户发送消息
# --------------------------------------------------------------------------
# 消息处理逻辑(2026-06 重构,使用 AIHandler 统一处理):
# 1. AIHandler 检测打招呼 → 引导描述问题,不计数
# 2. AIHandler 检测呼叫人工 → 拦截引导,不计数
# 3. AIHandler 调用 Dify API 获取 AI 回复
# - 命中 → AI 回复,ai_substantive_reply_count +1
# - 未命中 → 转 queued,返回转人工提示
# - 异常 → 降级模板回复,不计数,不转人工
# 4. 计数 >= 3 时,前端自动显示「呼叫坐席」按钮
# --------------------------------------------------------------------------
@router.post("/h5/conversations/current/messages")
async def h5_send_message(
body: dict,
employee_id: str = Depends(_get_current_employee),
db: AsyncSession = Depends(get_db),
):
"""H5 用户发送消息(含 AI 回复与计数)。
重构说明:AI 调用逻辑已统一至 AIHandler,此接口仅负责:
1. 会话管理(查找/创建)
2. 消息持久化
3. 根据 AIHandler 返回结果更新会话状态和计数
4. 返回响应
Args:
body: 消息请求体(包含 content
employee_id: 员工企微 UserID
db: 数据库会话
ai_handler: AI 处理器(DI 注入,统一 AI 调用逻辑)
Returns:
Dict: 统一响应格式,包含用户消息和 AI 回复
"""
content = body.get("content", "")
if not content:
raise AppException(1001, "消息内容不能为空")
# 支持非文本消息类型(image/file 等)
msg_type = body.get("msg_type", "text") # 消息内容类型:text/image/file
media_url = body.get("media_url") # 图片/文件 URL
file_name = body.get("file_name") # 文件名
file_size = body.get("file_size") # 文件大小(字节)
# 1. 查找或创建会话(新会话默认 ai_handling
stmt = select(Conversation).where(
Conversation.employee_id == employee_id,
Conversation.status.in_(["ai_handling", "queued", "serving"]),
).order_by(Conversation.created_at.desc())
result = await db.execute(stmt)
conversation = result.scalars().first()
if not conversation:
conversation = Conversation(
employee_id=employee_id,
status="ai_handling", # 先让 AI 尝试回答
urgency_score=1,
tags={},
ai_substantive_reply_count=0,
last_message_at=datetime.now(),
last_message_summary=content[:256],
)
db.add(conversation)
await db.flush()
# 2. 创建用户消息记录
message = Message(
conversation_id=conversation.id,
sender_type="employee",
sender_id=employee_id,
content=content,
msg_type=msg_type,
media_url=media_url,
file_name=file_name,
file_size=file_size,
is_read=False,
)
db.add(message)
# 更新会话信息
conversation.last_message_at = datetime.now()
conversation.last_message_summary = content[:256]
conversation.updated_at = datetime.now()
db.add(conversation)
await db.flush()
# 3. 广播用户消息给坐席端(员工端靠乐观更新已显示自己消息)
# 为什么:坐席端需实时看到员工新消息,仅依赖 3 秒轮询会有延迟
try:
await ws_manager.broadcast({
"type": "new_message",
"data": {
"conversation_id": str(conversation.id),
"message_id": str(message.id),
"sender_type": "employee",
"sender_id": employee_id,
"sender_name": "",
"content": content,
"msg_type": msg_type,
"urgency_score": conversation.urgency_score,
"tags": conversation.tags,
},
})
except Exception as ws_err:
# WS 广播失败不阻塞消息存储,只记录 warning
logger.warning(f"WS 广播用户消息失败(消息已存储): {ws_err}")
# 4. 提交当前事务,确保后台任务能读到刚创建的 conversation/message
# 为什么:asyncio.create_task 立即运行,若 HTTP 事务未提交,
# 后台 DB session 会报"会话不存在"race condition
await db.commit()
# 5. 启动后台 AI 任务(异步,不阻塞 HTTP 返回)
# 为什么:AI 推理(Dify)慢(3~15s),放后台经 WS 流式推回,
# 发送接口瞬时返回,前端不再卡"发送中"
# 约束:后台任务使用独立 DB session,且需单 worker(见 h5_ai_task.py
# v2.1Phase 4):传递 msg_type 和 media_url,支持图片消息 VisionService 分析
asyncio.create_task(
process_h5_ai_reply(
conversation_id=str(conversation.id),
employee_id=employee_id,
content=content,
dify_conversation_id=conversation.dify_conversation_id,
msg_type=msg_type,
media_url=media_url,
)
)
# 5. 立即返回用户消息(AI 回复经 WS 异步推送,不在此同步返回)
user_msg_data = MessageResponse.model_validate(message).model_dump()
return success_response(
data={
"user_message": user_msg_data,
"ai_reply": 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,
}
)
# --------------------------------------------------------------------------
# GET /api/h5/conversations/current/messages — 用户获取消息列表(历史消息)
# --------------------------------------------------------------------------
@router.get("/h5/conversations/current/messages")
async def h5_get_messages(
limit: int = Query(50, description="每页消息数量,默认50"),
before: Optional[str] = Query(None, description="获取此消息ID之前的消息(向上翻页)"),
employee_id: str = Depends(_get_current_employee),
db: AsyncSession = Depends(get_db),
):
"""H5 用户获取消息列表(历史消息)。
前端在进入会话或切换会话时调用,获取完整的消息历史记录。
支持分页向上翻页(通过 before 参数)。
Args:
limit: 每页消息数量(默认50)
before: 消息ID,获取此消息之前的消息(向上翻页)
employee_id: 员工企微 UserID
db: 数据库会话
Returns:
Dict: 统一响应格式,包含消息列表和 has_more 标志
"""
# 查找当前会话
stmt = select(Conversation).where(
Conversation.employee_id == employee_id,
Conversation.status.in_(["ai_handling", "queued", "serving"]),
).order_by(Conversation.created_at.desc())
result = await db.execute(stmt)
conversation = result.scalars().first()
if not conversation:
return success_response(data={"items": [], "has_more": False})
# 查询消息列表
msg_stmt = select(Message).where(
Message.conversation_id == conversation.id
).order_by(Message.created_at.desc()).limit(limit)
# 如果指定了 before,获取此消息之前的消息
if before:
try:
from uuid import UUID as UUIDType
UUIDType(before) # 仅校验格式
# 查询 before 消息的创建时间
before_stmt = select(Message.created_at).where(
Message.id == str(before)
)
before_result = await db.execute(before_stmt)
before_time = before_result.scalar_one_or_none()
if before_time:
msg_stmt = msg_stmt.where(Message.created_at < before_time)
except ValueError:
pass # 无效的UUID格式,忽略 before 参数
msg_result = await db.execute(msg_stmt)
messages = list(msg_result.scalars().all())
# 反转顺序(按时间正序返回)
messages.reverse()
items = [MessageResponse.model_validate(m).model_dump() for m in messages]
# 判断是否还有更多:查询的消息数是否等于 limit
has_more = len(messages) == limit
return success_response(data={"items": items, "has_more": has_more})
# --------------------------------------------------------------------------
# GET /api/h5/conversations/current/messages/poll — 用户轮询新消息
# --------------------------------------------------------------------------
@router.get("/h5/conversations/current/messages/poll")
async def h5_poll_messages(
after_message_id: Optional[str] = Query(None, description="返回此消息ID之后的新消息"),
employee_id: str = Depends(_get_current_employee),
db: AsyncSession = Depends(get_db),
):
"""H5 用户轮询新消息。
前端定时调用获取坐席回复的新消息。
Args:
after_message_id: 上次轮询的最后一消息ID
employee_id: 员工企微 UserID
db: 数据库会话
Returns:
Dict: 统一响应格式,包含新消息列表
"""
# 查找当前会话
stmt = select(Conversation).where(
Conversation.employee_id == employee_id,
Conversation.status.in_(["ai_handling", "queued", "serving"]),
).order_by(Conversation.created_at.desc())
result = await db.execute(stmt)
conversation = result.scalars().first()
if not conversation:
return success_response(data={"items": [], "has_more": False})
# 查询新消息
msg_stmt = select(Message).where(
Message.conversation_id == conversation.id
).order_by(Message.created_at.asc())
if after_message_id:
# 校验 UUID 格式,然后转字符串(兼容 SQLite/PG 的 String(36) 列,避免类型不匹配)
from uuid import UUID as UUIDType
try:
UUIDType(after_message_id) # 仅校验
except ValueError:
# 无效的UUID格式,返回空列表
items = []
return success_response(data={"items": items, "has_more": False})
# 必须用字符串比较,Message.id 在 DB 里是 String(36)/VARCHAR,
# 传 UUID 对象会被 SQLAlchemy 推断成 UUID 类型 → PostgreSQL 报
# "operator does not exist: character varying = uuid"
after_stmt = select(Message.created_at).where(
Message.id == str(after_message_id)
)
after_result = await db.execute(after_stmt)
after_time = after_result.scalar_one_or_none()
if after_time:
msg_stmt = msg_stmt.where(Message.created_at > after_time)
msg_result = await db.execute(msg_stmt)
messages = list(msg_result.scalars().all())
items = [MessageResponse.model_validate(m).model_dump() for m in messages]
return success_response(data={"items": items, "has_more": False})
# --------------------------------------------------------------------------
# POST /api/h5/conversations/current/shake — 举手/敲桌子呼叫坐席
# --------------------------------------------------------------------------
@router.post("/h5/conversations/current/shake")
async def shake(
body: ShakeRequest,
db: AsyncSession = Depends(get_db),
wecom_service: Optional[WecomService] = Depends(dep_wecom_service),
):
"""举手(敲桌子呼叫坐席)。
前端按钮从「摇人🔔」改为「敲桌子👊👊」,后端端点保持不变。
重构说明:不再手动创建/关闭 Redis 和 WecomService,改用 DI 注入共享实例。
流程:
1. 查找或创建会话
2. 设置举手标记
3. 获取趣味话术
4. 发送系统消息
5. 通过企微 API 发送话术给员工
6. 返回会话信息和话术
Args:
body: 举手请求体(包含 employee_id 和 employee_name
db: 数据库会话
wecom_service: 共享企微服务(DI 注入)
Returns:
Dict: 统一响应格式,包含会话信息和趣味话术
"""
employee_id = body.employee_id
employee_name = body.employee_name
# 1. 查找或创建会话
stmt = select(Conversation).where(
Conversation.employee_id == employee_id,
Conversation.status.in_(["ai_handling", "queued", "serving"]),
).order_by(Conversation.created_at.desc())
result = await db.execute(stmt)
conversation = result.scalars().first()
if not conversation:
# 无活跃会话 → 拒绝,必须先与 AI 互动(前端按钮此时不应出现,这是后端兜底)
raise AppException(
1003,
"请先描述您的问题,Duckula(达寇拉)需要先帮您分析。至少互动3轮后才能呼叫人工坐席哦~"
)
# ================================================================
# 决策 E2/E3:紧急关键词直通检测
# ================================================================
# 检查最近消息是否包含紧急关键词
# 如果命中 → 绕过3轮AI互动限制,直接进入排队/分配(紧急直通)
from app.services.triage_service import URGENCY_HIGH_KEYWORDS
last_msg_summary = (conversation.last_message_summary or "").lower()
is_emergency = any(kw in last_msg_summary for kw in URGENCY_HIGH_KEYWORDS)
# 前置校验:必须满足 AI 实质性回复 >= 3 次才能呼叫坐席
# 例外:紧急关键词命中时绕过此限制(决策 E2 紧急直通)
if not is_emergency and conversation.ai_substantive_reply_count < 3:
raise AppException(
1003,
"请先描述您的问题,Duckula(达寇拉)需要先帮您分析。至少互动3轮后才能呼叫人工坐席哦~"
)
# ================================================================
# 决策 E4:未梳理提醒
# ================================================================
# 如果信息未锁定(info_locked=False),返回 needs_info_confirm 标记
# 前端据此弹窗提示"完成信息梳理可进入快速通道"
needs_info_confirm = not conversation.info_locked
# 更新员工姓名
if employee_name and not conversation.employee_name:
conversation.employee_name = employee_name
# 设置举手标记
tags = dict(conversation.tags) if conversation.tags else {}
tags["hand_raise"] = True
if is_emergency:
tags["emergency_direct_connect"] = True # 紧急直通标记
conversation.tags = tags
conversation.urgency_score = max(conversation.urgency_score, 2)
if is_emergency:
conversation.urgency_score = max(conversation.urgency_score, 5) # 紧急直通设最高紧急度
conversation.last_message_at = datetime.now()
conversation.updated_at = datetime.now()
db.add(conversation)
await db.flush()
# 2. 获取趣味话术
funny_phrase_service = FunnyPhraseService(db)
is_vip = conversation.is_vip
phrase = await funny_phrase_service.get_phrase("shake", is_vip=is_vip)
# 3. 创建系统消息
system_msg = Message(
conversation_id=conversation.id,
sender_type="system",
sender_id="system",
sender_name="系统",
content=phrase,
msg_type="system",
is_read=True,
)
db.add(system_msg)
await db.flush()
# 4. 消息仅存储到数据库,由前端通过 WebSocket/轮询在 H5 页面内展示
# (不再通过企微应用消息推送,避免出现在通知栏)
# 5. 自动分配空闲坐席
from app.services.session_service import SessionService
from app.services.ws_manager import manager as ws_manager
from app.models.agent import Agent
assigned_agent_id: Optional[str] = None
assign_result: str = "queued"
# 查找在线且未满负荷的坐席(按当前负载升序,取第一个)
stmt = select(Agent).where(
Agent.status == "online",
Agent.current_load < Agent.max_load
).order_by(Agent.current_load).limit(1)
result = await db.execute(stmt)
available_agent = result.scalars().first()
if available_agent:
# 找到空闲坐席,分配给该会话
try:
session_service = SessionService(db, wecom_service)
await session_service.assign_agent(conversation.id, available_agent.user_id)
assigned_agent_id = available_agent.user_id
assign_result = "assigned"
logger.info(f"自动分配坐席: conv_id={conversation.id}, agent={assigned_agent_id}")
except Exception as e:
logger.warning(f"自动分配坐席失败: {e}")
assign_result = "assign_failed"
else:
# 无空闲坐席,进入排队(会话状态保持 queued,由 AI 未命中时自动处理)
assign_result = "queued"
logger.info(f"无空闲坐席,会话进入排队: conv_id={conversation.id}")
# 6. 广播 new_conversation 事件通知所有坐席
try:
await ws_manager.broadcast({
"type": "new_conversation",
"data": {
"conversation_id": str(conversation.id),
"employee_id": employee_id,
"employee_name": employee_name or "未知用户",
"urgency_score": conversation.urgency_score,
"hand_raise": True,
"assigned_agent_id": assigned_agent_id,
"assign_result": assign_result,
"is_emergency": is_emergency, # 紧急直通标记
"info_locked": conversation.info_locked, # 信息梳理状态
"queue_priority": conversation.queue_priority, # 答题插队优先级
}
})
except Exception as e:
logger.warning(f"WebSocket广播失败(不阻塞流程): {e}")
logger.info(f"举手触发: employee_id={employee_id}, conv_id={conversation.id}, assign_result={assign_result}, emergency={is_emergency}")
# 7. 返回会话信息和话术
conv_data = ConversationResponse.model_validate(conversation).model_dump()
return success_response(
data={
"conversation": conv_data,
"funny_phrase": phrase,
"assign_result": assign_result,
"assigned_agent_id": assigned_agent_id,
"is_emergency": is_emergency, # 紧急直通标记
"needs_info_confirm": needs_info_confirm, # 信息未梳理提醒
"info_locked": conversation.info_locked, # 当前信息锁定状态
}
)
# --------------------------------------------------------------------------
# POST /api/h5/conversations/current/call-agent — 摇人按钮触发转人工
# --------------------------------------------------------------------------
@router.post("/h5/conversations/current/call-agent")
async def call_agent(
body: ShakeRequest,
db: AsyncSession = Depends(get_db),
wecom_service: Optional[WecomService] = Depends(dep_wecom_service),
):
"""摇人按钮 - 呼叫坐席。
用户点击摇人按钮后,触发转人工流程:
1. 查找当前会话
2. 校验AI回复次数 >= 3(与shake一致)
3. 将会话状态改为 queued(排队中)
4. 尝试分配空闲坐席
5. 发送系统消息通知用户
6. 通过企微消息通知坐席
Args:
body: 呼叫坐席请求体(包含 employee_id 和 employee_name
db: 数据库会话
wecom_service: 共享企微服务(DI 注入)
Returns:
Dict: 包含会话信息和排队状态
"""
from app.services.session_service import SessionService
employee_id = body.employee_id
employee_name = body.employee_name
# 1. 查找当前活跃会话
stmt = select(Conversation).where(
Conversation.employee_id == employee_id,
Conversation.status.in_(["ai_handling", "queued", "serving"]),
).order_by(Conversation.created_at.desc())
result = await db.execute(stmt)
conversation = result.scalars().first()
if not conversation:
raise AppException(
code=1003,
message="请先描述您的问题,Duckula(达寇拉)需要先帮您分析。至少互动3轮后才能呼叫人工坐席哦~"
)
# 2. 前置校验:必须满足 AI 实质性回复 >= 3 次
if conversation.ai_substantive_reply_count < 3:
raise AppException(
code=1003,
message="请先描述您的问题,Duckula(达寇拉)需要先帮您分析。至少互动3轮后才能呼叫人工坐席哦~"
)
# 更新员工姓名
if employee_name and not conversation.employee_name:
conversation.employee_name = employee_name
# 3. 将会话状态改为 queued(排队中)
conversation.status = "queued"
conversation.last_message_at = datetime.now()
conversation.updated_at = datetime.now()
# 设置紧急度加分
tags = dict(conversation.tags) if conversation.tags else {}
tags["user_called_agent"] = True # 标记用户主动呼叫
conversation.tags = tags
db.add(conversation)
await db.flush()
# 4. 尝试分配空闲坐席
session_service = SessionService(db)
assigned_agent = await session_service.auto_assign_agent(conversation.id)
# 5. 获取趣味话术
funny_phrase_service = FunnyPhraseService(db)
is_vip = conversation.is_vip
phrase = await funny_phrase_service.get_phrase("transfer", is_vip=is_vip)
# 6. 创建系统消息
system_content = phrase
if assigned_agent:
system_content = f"{phrase}\n\n为您服务的是:{assigned_agent.name}"
conversation.status = "serving"
conversation.assigned_agent_id = assigned_agent.user_id
system_msg = Message(
conversation_id=conversation.id,
sender_type="system",
sender_id="system",
sender_name="系统",
content=system_content,
msg_type="system",
is_read=True,
)
db.add(system_msg)
# 7. 消息仅存储到数据库,由前端通过 WebSocket/轮询在 H5 页面内展示
# (不再通过企微应用消息推送,避免出现在通知栏)
# 8. 如果分配了坐席,通知坐席有新会话
if assigned_agent and wecom_service:
try:
notify_phrase = f"新会话:{employee_name} 呼叫人工服务,请及时接单"
# 获取坐席的userid并发送通知(需要坐席绑定企微)
# 此处简化处理,仅记录日志
logger.info(f"分配坐席: agent_id={assigned_agent.id}, employee_id={employee_id}")
except Exception as e:
logger.warning(f"坐席通知失败: {e}")
await db.commit()
logger.info(f"呼叫坐席: employee_id={employee_id}, conv_id={conversation.id}, agent_id={assigned_agent.id if assigned_agent else 'None'}")
# 9. 返回结果
conv_data = ConversationResponse.model_validate(conversation).model_dump()
return success_response(
data={
"conversation": conv_data,
"status": conversation.status,
"queue_position": 1 if not assigned_agent else None,
"estimated_wait_seconds": 30 if not assigned_agent else 0,
}
)
# --------------------------------------------------------------------------
# GET /api/h5/conversations/current/queue-status — 查询排队状态
# --------------------------------------------------------------------------
@router.get("/h5/conversations/current/queue-status")
async def get_queue_status(
employee_id: str = Query(..., description="员工ID"),
db: AsyncSession = Depends(get_db),
):
"""查询当前排队状态(三段排序版)。
排队三段排序:
1. VIP段(is_vip=true
2. 已梳理段(info_locked=true
3. 待梳理段(info_locked=false
段内排序:queue_priority DESC → urgency_score DESC → created_at ASC
Args:
employee_id: 员工ID
Returns:
Dict: 排队状态信息(含段位、位置、预估等待时间)
"""
from app.services.queue_service import get_queue_service
# 1. 查找该员工的排队会话
stmt = select(Conversation).where(
Conversation.employee_id == employee_id,
Conversation.status == "queued",
).order_by(Conversation.created_at.asc())
result = await db.execute(stmt)
conversation = result.scalars().first()
if not conversation:
# 不在排队中
return success_response(data={
"in_queue": False,
"status": None,
"queue_position": None,
"estimated_wait_seconds": 0,
})
# 2. 使用 QueueService 计算三段排序位置
queue_service = get_queue_service()
queue_info = await queue_service.calculate_queue_position(db, conversation)
return success_response(data={
"in_queue": True,
"status": conversation.status,
"conversation_id": str(conversation.id),
"queue_position": queue_info["position"],
"estimated_wait_seconds": queue_info["estimated_wait_sec"],
"segment": queue_info["segment"],
"segment_label": queue_info["segment_label"],
"ahead_count": queue_info["ahead_count"],
"queue_priority": queue_info["queue_priority"],
"info_locked": conversation.info_locked,
"is_vip": conversation.is_vip,
})
# --------------------------------------------------------------------------
# POST /api/h5/conversations/current/cancel-queue — 取消排队
# --------------------------------------------------------------------------
@router.post("/h5/conversations/current/cancel-queue")
async def cancel_queue(
body: ShakeRequest,
db: AsyncSession = Depends(get_db),
):
"""取消排队。
用户主动取消排队,释放排队位置。
Args:
body: 包含 employee_id
Returns:
Dict: 操作结果
"""
employee_id = body.employee_id
# 1. 查找排队中的会话
stmt = select(Conversation).where(
Conversation.employee_id == employee_id,
Conversation.status == "queued",
)
result = await db.execute(stmt)
conversation = result.scalars().first()
if not conversation:
raise AppException(code=1004, message="您当前不在排队中")
# 2. 将会话状态改回 ai_handling
conversation.status = "ai_handling"
conversation.updated_at = datetime.now()
# 移除用户主动呼叫标记
tags = dict(conversation.tags) if conversation.tags else {}
tags.pop("user_called_agent", None)
conversation.tags = tags
db.add(conversation)
await db.commit()
logger.info(f"取消排队: employee_id={employee_id}, conv_id={conversation.id}")
return success_response(data={
"message": "已取消排队,会话将继续由AI服务",
})
# --------------------------------------------------------------------------
# GET /api/h5/approval-links — 获取审批流程链接
# --------------------------------------------------------------------------
@router.get("/h5/approval-links")
async def get_approval_links(
category: Optional[str] = Query(None, description="按分类过滤: IT/HR/行政/财务"),
db: AsyncSession = Depends(get_db),
):
"""获取审批流程链接。
从 approval_links 表读取,支持按分类过滤。
用于 H5 用户端 AI 助手面板。
Args:
category: 按分类过滤(可选)
db: 数据库会话
Returns:
Dict: 统一响应格式,包含审批链接列表
"""
stmt = select(ApprovalLink).order_by(ApprovalLink.sort_order)
if category:
stmt = stmt.where(ApprovalLink.category == category)
result = await db.execute(stmt)
links = list(result.scalars().all())
items = [ApprovalLinkResponse.model_validate(link).model_dump() for link in links]
return success_response(data={"items": items})
# --------------------------------------------------------------------------
# GET /api/h5/software-downloads — 获取软件下载列表
# --------------------------------------------------------------------------
@router.get("/h5/software-downloads")
async def get_software_downloads(
category: Optional[str] = Query(None, description="按分类过滤: 办公/开发/安全/工具"),
db: AsyncSession = Depends(get_db),
):
"""获取软件下载列表。
从 software_downloads 表读取,支持按分类过滤。
用于 H5 用户端 AI 助手面板。
Args:
category: 按分类过滤(可选)
db: 数据库会话
Returns:
Dict: 统一响应格式,包含软件下载列表
"""
stmt = select(SoftwareDownload).order_by(SoftwareDownload.sort_order)
if category:
stmt = stmt.where(SoftwareDownload.category == category)
result = await db.execute(stmt)
downloads = list(result.scalars().all())
items = [SoftwareDownloadResponse.model_validate(d).model_dump() for d in downloads]
return success_response(data={"items": items})
# ==========================================================================
# 邀请功能 H5 专用端点(P0-09~P0-11
# ==========================================================================
# 说明:为 H5 员工端提供带认证的参与者管理接口
# 认证:使用 _get_current_employee 依赖,验证 Bearer Token
# 安全:校验 token 对应的 employee_id 与请求体一致,防止冒充
# --------------------------------------------------------------------------
from app.services.session_service import SessionService
@router.post("/h5/conversations/{conversation_id}/join")
async def h5_join_conversation(
conversation_id: str,
employee_id: str = Depends(_get_current_employee),
db: AsyncSession = Depends(get_db),
):
"""H5员工加入会话(带认证)。
做什么:被邀请人点击企微卡片链接后,通过此接口加入会话
为什么:需要验证请求者身份,防止冒充其他员工加入会话
认证:Bearer Token → employee_id,与请求体中的 employee_id 校验一致性
副作用:
- 更新参与者的 joined 状态
- 在会话中创建系统消息
- WebSocket 广播参与者变更
Args:
conversation_id: 会话ID
employee_id: 当前登录员工ID(从 Token 认证获取)
db: 数据库会话
Returns:
Dict: 统一响应格式,包含更新后的会话信息
"""
session_service = SessionService(db)
conversation = await session_service.join_conversation(
conversation_id=conversation_id,
employee_id=employee_id,
)
response_data = ConversationResponse.model_validate(conversation).model_dump()
return success_response(data=response_data)
@router.post("/h5/conversations/{conversation_id}/leave-participant")
async def h5_leave_participant(
conversation_id: str,
employee_id: str = Depends(_get_current_employee),
db: AsyncSession = Depends(get_db),
):
"""H5员工退出会话(带认证)。
做什么:参与者主动退出会话
为什么:需要验证请求者身份,防止冒充其他员工退出
认证:Bearer Token → employee_id
副作用:
- 在会话中创建系统消息
- WebSocket 广播参与者变更
Args:
conversation_id: 会话ID
employee_id: 当前登录员工ID(从 Token 认证获取)
db: 数据库会话
Returns:
Dict: 统一响应格式,包含更新后的会话信息
"""
session_service = SessionService(db)
conversation = await session_service.leave_as_participant(
conversation_id=conversation_id,
employee_id=employee_id,
)
response_data = ConversationResponse.model_validate(conversation).model_dump()
return success_response(data=response_data)
@router.get("/h5/conversations/{conversation_id}/participants")
async def h5_get_participants(
conversation_id: str,
employee_id: str = Depends(_get_current_employee),
db: AsyncSession = Depends(get_db),
):
"""获取会话参与者列表(带认证 + 参与者权限校验)。
做什么:返回指定会话的所有参与者信息
为什么:H5员工端需要查看参与者面板
认证:Bearer Token → employee_id
权限:仅会话发起人(employee_id)或已加入的参与者可查看,防止越权读取他会话信息
P0 安全修复(2026-06-14 评审):
此前仅校验"用户已登录",未校验"是否属于本会话",存在数据泄露风险——
任意已登录员工可枚举 conversation_id 读取其他会话的参与者名单。
Args:
conversation_id: 会话ID
employee_id: 当前登录员工ID(从 Token 认证获取)
db: 数据库会话
Returns:
Dict: 统一响应格式,包含参与者列表
Raises:
ERR_CONVERSATION_NOT_FOUND: 会话不存在
AppException(4003): 当前员工不是会话发起人/参与者
"""
stmt = select(Conversation).where(Conversation.id == conversation_id)
result = await db.execute(stmt)
conversation = result.scalars().first()
if not conversation:
from app.utils.response import ERR_CONVERSATION_NOT_FOUND
raise ERR_CONVERSATION_NOT_FOUND
# P0-1 修复:校验当前员工是否有权查看本会话的参与者
# 权限规则:会话发起人 OR 已加入的参与者
is_creator = conversation.employee_id == employee_id
is_participant = any(
p.get("id") == employee_id
for p in (conversation.participants or [])
)
if not (is_creator or is_participant):
raise AppException(4003, "您不是该会话的参与者,无权查看")
participants = conversation.participants or []
return success_response(data={"participants": participants})
# ==========================================================================
# H5 员工端目录与邀请功能(P0-09~P0-11 扩展)
# ==========================================================================
# 说明:为 H5 员工端提供员工搜索、组织架构树、邀请参与者接口
# 认证:使用 _get_current_employee 依赖,验证 Bearer Token
# 复用:get_org_directory() 服务(与坐席端 agent_directory 共享同一数据源)
# --------------------------------------------------------------------------
@router.get("/h5/employees/search")
async def h5_search_employees(
keyword: str = Query("", description="搜索关键词(姓名或工号,模糊匹配)"),
employee_id: str = Depends(_get_current_employee),
db: AsyncSession = Depends(get_db),
redis: Optional[aioredis.Redis] = Depends(dep_redis),
):
"""H5 员工端模糊搜索员工(按姓名或工号)。
复用 get_org_directory() 获取组织目录(优先企微全组织,降级本地 employees 表),
在内存中做大小写不敏感的模糊匹配。排除当前登录员工自己。
Args:
keyword: 搜索关键词(为空时返回空列表,避免无意义全量返回)
employee_id: 当前登录员工ID(从 Token 认证获取)
db: 数据库会话
redis: Redis 客户端(可选,用于企微目录缓存读取)
Returns:
Dict: 统一响应格式,data 为员工列表
[{ "id": "userid", "name": "姓名", "department": "部门" }, ...]
"""
kw = (keyword or "").strip()
if not kw:
return success_response(data=[])
try:
# 获取组织目录(含 30 分钟 Redis 缓存 + 本地降级)
directory, _ = await get_org_directory(db, redis)
kw_lower = kw.lower()
results: List[Dict[str, Any]] = []
for emp in directory:
# 排除当前员工自己(避免邀请自己加入会话)
if emp.get("employee_id") == employee_id:
continue
name = emp.get("name", "") or ""
emp_id = emp.get("employee_id", "") or ""
# 模糊匹配:姓名 或 工号(employee_id)包含关键词(大小写不敏感)
if kw_lower in name.lower() or kw_lower in emp_id.lower():
results.append({
"id": emp_id,
"name": name,
"department": emp.get("department", ""),
})
logger.info(f"H5员工搜索: keyword='{kw}', 命中 {len(results)}")
return success_response(data=results)
except AppException:
raise
except Exception as e:
logger.error(f"H5员工搜索异常: {e}", exc_info=True)
raise AppException(1005, f"搜索失败: {str(e)}")
@router.get("/h5/org/tree")
async def h5_get_org_tree(
employee_id: str = Depends(_get_current_employee),
db: AsyncSession = Depends(get_db),
redis: Optional[aioredis.Redis] = Depends(dep_redis),
):
"""H5 员工端获取组织架构树(部门层级 + 每个部门下的员工列表)。
利用企微 department/list 返回的 parentid 字段构建真正的层级树,
不再按部门名扁平分组。部门ID加 ``dept_`` 前缀作为唯一 key
避免同名部门合并。员工可以出现在其所属的所有部门下。
树结构示例(多层级):
[
{
"id": "dept_1",
"label": "公司",
"dept_id": 1,
"parentid": 0,
"children": [
{
"id": "dept_2",
"label": "研发一部",
"dept_id": 2,
"parentid": 1,
"children": [
{"id": "zhangsan", "label": "张三", "isLeaf": true, "department": "研发一部"}
]
}
]
}
]
性能优化:
- 树构建结果独立缓存(key: wecom:org_tree:h5TTL 30 分钟)
- 缓存中包含所有员工,读取后过滤掉当前登录员工自己
Args:
employee_id: 当前登录员工ID
db: 数据库会话
redis: Redis 客户端
Returns:
Dict: 统一响应格式,data 为树节点列表
"""
try:
# 获取组织架构树(含独立缓存 + 排除当前员工)
tree = await get_org_tree_cached(db, redis, "h5", employee_id)
total_employees = count_tree_employees(tree)
logger.info(f"H5组织架构树: {len(tree)} 个顶层节点, 共 {total_employees}")
return success_response(data=tree)
except AppException:
raise
except Exception as e:
logger.error(f"H5组织架构树获取异常: {e}", exc_info=True)
raise AppException(1005, f"获取组织架构树失败: {str(e)}")
@router.post("/h5/conversations/{conversation_id}/invite-participant")
async def h5_invite_participant(
conversation_id: str,
body: InviteParticipantRequest,
employee_id: str = Depends(_get_current_employee),
db: AsyncSession = Depends(get_db),
wecom_service: WecomService = Depends(dep_wecom_service),
):
"""H5 员工端邀请参与者加入会话(P0-09 扩展)。
做什么:会话发起人(员工)邀请其他员工参与当前会话
为什么:复杂IT问题可能需要业务方同事补充信息,员工可自行邀请
认证:Bearer Token → employee_id,作为 inviter_agent_id 传入
权限:后端 session_service.invite_participants 已改为允许主责坐席或会话发起人
副作用:
- 向被邀请人发送企微卡片通知(含「加入会话」按钮)
- 在会话中创建系统消息
- WebSocket 广播参与者变更
Args:
conversation_id: 会话ID
body: 邀请请求(含被邀请人列表 + 历史共享模式)
employee_id: 当前登录员工ID(从 Token 认证获取,作为邀请人)
db: 数据库会话
wecom_service: 共享企微服务(DI 注入,发送卡片通知用)
Returns:
Dict: 统一响应格式,包含更新后的会话信息
"""
session_service = SessionService(db, wecom_service=wecom_service)
conversation = await session_service.invite_participants(
conversation_id=conversation_id,
inviter_agent_id=employee_id,
participants=[p.model_dump() for p in body.participants],
history_mode=body.history_mode,
)
response_data = ConversationResponse.model_validate(conversation).model_dump()
return success_response(data=response_data)
# ==========================================================================
# 关闭机制 API(决策 G1-G5
# ==========================================================================
# 五种关闭场景的 H5 端点:
# POST /api/h5/conversations/current/resolve — 员工确认AI已解决
# POST /api/h5/conversations/current/close — 员工主动关闭
# POST /api/h5/conversations/current/resolve/confirm — 员工确认坐席结单
# POST /api/h5/conversations/current/resolve/reject — 员工拒绝坐席结单
# POST /api/h5/conversations/current/reopen — 24h内重开
# ==========================================================================
class SelfResolveRequest(BaseModel):
"""员工确认AI已解决请求体。"""
resolve_summary: Optional[str] = Field(None, description="解决摘要(可选)")
class EmployeeCloseRequest(BaseModel):
"""员工主动关闭请求体。"""
close_reason: Optional[str] = Field(None, description="关闭原因(可选)")
class ResolveConfirmRequest(BaseModel):
"""员工确认/拒绝坐席结单请求体。"""
action: str = Field(..., description="confirm=确认, reject=拒绝")
reason: Optional[str] = Field(None, description="拒绝原因(拒绝时可选)")
class ReopenRequest(BaseModel):
"""重开会话请求体。"""
original_conversation_id: str = Field(..., description="原会话ID")
# --------------------------------------------------------------------------
# POST /api/h5/conversations/current/resolve — 员工确认AI已解决
# --------------------------------------------------------------------------
@router.post("/h5/conversations/current/resolve")
async def h5_self_resolve(
body: SelfResolveRequest,
employee_id: str = Depends(_get_current_employee),
db: AsyncSession = Depends(get_db),
):
"""员工确认AI已解决问题(AI自助场景)。
触发场景:
- 对话流中"已解决"确认卡片按钮
- AI检测到关闭关键词后推送的确认卡片
状态转换:ai_handling → resolved
关闭方:employee / 关闭方式:ai_self
Args:
body: 请求体(可选 resolve_summary
employee_id: 当前登录员工ID
db: 数据库会话
Returns:
Dict: 统一响应格式,包含已关闭的会话信息
"""
closing_service = ClosingService(db)
conversation = await closing_service.employee_self_resolve(
employee_id=employee_id,
resolve_summary=body.resolve_summary,
)
await db.commit()
response_data = ConversationResponse.model_validate(conversation).model_dump()
return success_response(data=response_data)
# --------------------------------------------------------------------------
# POST /api/h5/conversations/current/close — 员工主动关闭
# --------------------------------------------------------------------------
@router.post("/h5/conversations/current/close")
async def h5_employee_close(
body: EmployeeCloseRequest,
employee_id: str = Depends(_get_current_employee),
db: AsyncSession = Depends(get_db),
):
"""员工主动关闭会话。
适用场景:
- 问题自行解决,不需要AI或坐席帮助
- 不想继续等待
- 问题已通过其他渠道解决
状态转换:任意活跃状态 → resolved
关闭方:employee / 关闭方式:employee_initiative
Args:
body: 请求体(可选 close_reason
employee_id: 当前登录员工ID
db: 数据库会话
Returns:
Dict: 统一响应格式,包含已关闭的会话信息
"""
closing_service = ClosingService(db)
conversation = await closing_service.employee_initiative_close(
employee_id=employee_id,
close_reason=body.close_reason,
)
await db.commit()
response_data = ConversationResponse.model_validate(conversation).model_dump()
return success_response(data=response_data)
# --------------------------------------------------------------------------
# POST /api/h5/conversations/current/resolve/confirm — 员工确认/拒绝坐席结单
# --------------------------------------------------------------------------
@router.post("/h5/conversations/current/resolve/confirm")
async def h5_resolve_confirm(
body: ResolveConfirmRequest,
employee_id: str = Depends(_get_current_employee),
db: AsyncSession = Depends(get_db),
):
"""员工确认或拒绝坐席的结单请求。
坐席发起结单后,会话进入 pending_close 状态,
员工通过此端点确认或拒绝。
- confirm: pending_close → resolved(坐席结单+员工确认)
- reject: pending_close → serving(恢复服务)
- 5分钟内不响应:系统自动关闭
Args:
body: 请求体(action=confirm/reject, reason=拒绝原因)
employee_id: 当前登录员工ID
db: 数据库会话
Returns:
Dict: 统一响应格式,包含更新后的会话信息
"""
closing_service = ClosingService(db)
if body.action == "confirm":
conversation = await closing_service.employee_confirm_resolve(employee_id)
message = "结单确认成功,会话已关闭。"
elif body.action == "reject":
conversation = await closing_service.employee_reject_resolve(
employee_id, reason=body.reason
)
message = "已为您恢复服务,坐席将继续处理。"
else:
raise AppException(1008, f"无效的action: {body.action},应为 confirm 或 reject")
await db.commit()
response_data = ConversationResponse.model_validate(conversation).model_dump()
response_data["message"] = message
return success_response(data=response_data)
# --------------------------------------------------------------------------
# POST /api/h5/conversations/current/reopen — 24h内重开已关闭会话
# --------------------------------------------------------------------------
@router.post("/h5/conversations/current/reopen")
async def h5_reopen(
body: ReopenRequest,
employee_id: str = Depends(_get_current_employee),
db: AsyncSession = Depends(get_db),
):
"""24小时内重开已关闭的会话。
创建新会话并关联原会话ID,用于上下文继承。
新会话状态为 ai_handling,复用原会话的员工信息。
限制条件:
- 原会话必须已关闭(status=resolved
- 距离关闭不超过24小时
- 重开后新会话关联原会话的 reference_conversation_id
Args:
body: 请求体(original_conversation_id
employee_id: 当前登录员工ID
db: 数据库会话
Returns:
Dict: 统一响应格式,包含新创建的会话信息
"""
closing_service = ClosingService(db)
new_conversation = await closing_service.reopen_conversation(
employee_id=employee_id,
original_conversation_id=body.original_conversation_id,
)
await db.commit()
response_data = ConversationResponse.model_validate(new_conversation).model_dump()
response_data["is_reopen"] = True
response_data["message"] = "问题已重新接入,请描述您遇到的情况。"
return success_response(data=response_data)
# ==========================================================================
# GET /api/h5/it-health — IT 健康信息
# ==========================================================================
# 说明:返回当前登录员工终端的 IT 健康信息,包括设备基本信息、
# CPU/内存/磁盘使用率、安全检查状态、合规检查状态。
# 数据来源:联软(设备信息) + 火绒(安全状态) + 资产服务(资产编号)
# 降级策略:联软/火绒未配置时返回 Mock 数据
# ==========================================================================
@router.get("/h5/it-health")
async def h5_get_it_health(
employee_id: str = Depends(_get_current_employee),
db: AsyncSession = Depends(get_db),
):
"""获取当前员工终端的 IT 健康信息。
从联软、火绒、资产服务聚合数据,返回设备信息和安全状态。
如果联软/火绒集成未配置,返回 Mock 降级数据。
Args:
employee_id: 员工企微 UserID(通过认证依赖注入)
db: 数据库会话(读取集成配置)
Returns:
Dict: 统一响应格式,包含:
- current_device: 当前设备信息(设备名/IP/MAC/OS/CPU/内存/磁盘/安全检查)
- other_devices: 其他设备列表
- data_source: "real"(真实数据)或 "mock"(降级数据)
- generated_at: 生成时间
"""
from app.services.it_health_service import ITHealthService
service = ITHealthService(db)
result = await service.get_it_health(employee_id)
return success_response(data=result)