Compare commits
53 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 1c5169531a | |||
| 502551ad9c | |||
| a8da4da0db | |||
| f9c25147c9 | |||
| 1f6a2936f5 | |||
| 01e070c8ef | |||
| 21830d507d | |||
| 15be4963a1 | |||
| 189458912c | |||
| 0e54312b5c | |||
| aaebe31bd8 | |||
| 53b7bbcff7 | |||
| 7864ab404b | |||
| 44e77dcb0e | |||
| 3a44141eac | |||
| 969f524962 | |||
| 5adcb4852a | |||
| 22909fb3bd | |||
| 566bb46cd4 | |||
| a20606f328 | |||
| b321e3dabd | |||
| 7886f2f523 | |||
| 28bc4bda3b | |||
| 05304aec9c | |||
| 928f4bb337 | |||
| b997a683ce | |||
| 6a5f01ff90 | |||
| 6f545fe0b6 | |||
| 1bfc00559a | |||
| bfac494b01 | |||
| 2681e7b0bf | |||
| 6cbee00dc0 | |||
| f9be3547d9 | |||
| 59609b1d8d | |||
| 9cbf2650e1 | |||
| 6e74631ecb | |||
| 555a16c310 | |||
| 5dd8436cde | |||
| 2a4db59f02 | |||
| 873f3176e8 | |||
| aeb4e0cf39 | |||
| af392f2bf0 | |||
| a6518bff03 | |||
| 54481525f4 | |||
| 18a639919c | |||
| 6932e5a449 | |||
| 939de70e48 | |||
| 3ed86d5fb3 | |||
| 5a77a89ab1 | |||
| 449c6d4875 | |||
| bea288e414 | |||
| 4052e19ff8 | |||
| 6db1c0eef0 |
@@ -142,6 +142,12 @@ wecom-it-desk-server-deploy.zip
|
||||
.workbuddy/*.log.err
|
||||
# workbuddy 记忆目录(个人上下文,不 入仓)
|
||||
.workbuddy/memory/
|
||||
# workbuddy 工作区产物(2026-08-03 补充: 之前 git add . 误入 M)
|
||||
.workbuddy/outputs/
|
||||
.workbuddy/artifacts/
|
||||
.workbuddy/automations/
|
||||
.workbuddy/tmp/
|
||||
.workbuddy/deploy-temp/
|
||||
|
||||
# =============================================================================
|
||||
# 工作树清理 (2026-07-09): 产物 / 临时 / 上传 / 调试 dump 不入仓
|
||||
@@ -233,3 +239,50 @@ backend/scripts/create_test_agent.py
|
||||
# 补充忽略 (2026-07-09 WIP 提交): 新增构建产物
|
||||
dist-deploy/
|
||||
dist-v2/
|
||||
|
||||
# 补充忽略 (2026-07-13): 临时目录 / 备份 / 截图
|
||||
.workbuddy/tmp/
|
||||
.workbuddy/automations/
|
||||
deploy-staging-ki/
|
||||
dist-old-*/
|
||||
dist_deploy/
|
||||
screenshots/
|
||||
test-screenshots/
|
||||
tools/
|
||||
chat_export/
|
||||
deliverables/
|
||||
02meiti/
|
||||
data/
|
||||
|
||||
# 补充忽略 (2026-08-03 仓库重组): 历史 dist 备份 + node_modules_old
|
||||
# 这些是早期部署流程误将 dist 目录 commit 的残留,每个 50MB+,必须不入仓
|
||||
dist.bak/
|
||||
dist.bak.*/
|
||||
dist.old/
|
||||
dist_bak*/
|
||||
dist-clean/
|
||||
dist_bak_*/
|
||||
dist_old/
|
||||
node_modules_old/
|
||||
|
||||
# 补充忽略 (2026-08-03 仓库重组): 用户运行时上传的二进制文件(src/ 前缀)
|
||||
# 原 .gitignore 只有 backend/media/* 规则,重组后需补充 src/backend/* 对应规则
|
||||
src/backend/uploads/
|
||||
src/backend/media/
|
||||
|
||||
# 补充忽略 (2026-08-03 仓库重组): 前端部署脚本生成的 bin chunk
|
||||
# agent.p[0-9].bin (坐席端部署脚本产物) + h5-v4-part[0-9].bin (H5端部署脚本产物)
|
||||
*-part*.bin
|
||||
p[0-9].bin
|
||||
*.p[0-9].bin
|
||||
|
||||
# 补充忽略 (2026-08-03 仓库重组): 前端部署脚本生成的 part 拆分文件
|
||||
# 部署脚本将大 tar/zip 拆分为 part0/part1/part2 上传 (变体多, 通配匹配)
|
||||
*.part*
|
||||
|
||||
# 补充忽略 (2026-08-03 收尾): archives/ 临时备份 + 02meiti hilo 应用数据
|
||||
# archives/ 是早期未跟踪的临时备份目录(85 个文件)
|
||||
# 02meiti/.hilo/ 是 hilo 多媒体应用数据,原 .gitignore 254 行 02meiti/ 已覆盖,
|
||||
# 但之前有 3 个文件被误加入 index (index.sqlite-shm/wal/storage.json),已 git rm --cached
|
||||
archives/
|
||||
02meiti/.hilo/
|
||||
|
||||
@@ -1,290 +0,0 @@
|
||||
{
|
||||
"timestamp": "2026-07-03T08:44:12.328961",
|
||||
"tasks": {
|
||||
"A-T1": {
|
||||
"status": "✅已完成",
|
||||
"owner": "工程师",
|
||||
"actual_end": "07-02"
|
||||
},
|
||||
"A-T2": {
|
||||
"status": "✅已完成",
|
||||
"owner": "工程师",
|
||||
"actual_end": "07-02"
|
||||
},
|
||||
"A-T3": {
|
||||
"status": "✅已完成",
|
||||
"owner": "工程师",
|
||||
"actual_end": "07-02"
|
||||
},
|
||||
"A-T4": {
|
||||
"status": "✅已完成",
|
||||
"owner": "工程师",
|
||||
"actual_end": "07-02"
|
||||
},
|
||||
"A-T5": {
|
||||
"status": "✅已完成",
|
||||
"owner": "工程师",
|
||||
"actual_end": "07-02"
|
||||
},
|
||||
"A-T6": {
|
||||
"status": "✅已完成",
|
||||
"owner": "工程师",
|
||||
"actual_end": "07-02"
|
||||
},
|
||||
"A-T7": {
|
||||
"status": "✅已完成",
|
||||
"owner": "工程师",
|
||||
"actual_end": "07-02"
|
||||
},
|
||||
"A-T8": {
|
||||
"status": "✅已完成",
|
||||
"owner": "QA",
|
||||
"actual_end": "07-02"
|
||||
},
|
||||
"A-T9": {
|
||||
"status": "✅已完成",
|
||||
"owner": "工程师",
|
||||
"actual_end": "07-02"
|
||||
},
|
||||
"A-T10": {
|
||||
"status": "✅已完成",
|
||||
"owner": "工程师",
|
||||
"actual_end": "07-02"
|
||||
},
|
||||
"A-T11": {
|
||||
"status": "✅已完成",
|
||||
"owner": "项目经理",
|
||||
"actual_end": "07-02"
|
||||
},
|
||||
"A-T12": {
|
||||
"status": "✅已完成",
|
||||
"owner": "QA",
|
||||
"actual_end": "07-02"
|
||||
},
|
||||
"A-T13": {
|
||||
"status": "✅已完成",
|
||||
"owner": "工程师",
|
||||
"actual_end": "07-02"
|
||||
},
|
||||
"A-T14": {
|
||||
"status": "🔵进行中",
|
||||
"owner": "工程师",
|
||||
"actual_end": "-"
|
||||
},
|
||||
"A-T15": {
|
||||
"status": "🔵进行中",
|
||||
"owner": "工程师",
|
||||
"actual_end": "-"
|
||||
},
|
||||
"A-T16": {
|
||||
"status": "🔵进行中",
|
||||
"owner": "QA+工程师",
|
||||
"actual_end": "-"
|
||||
},
|
||||
"A-T17": {
|
||||
"status": "🔴阻塞",
|
||||
"owner": "运维",
|
||||
"actual_end": "-"
|
||||
},
|
||||
"A-T18": {
|
||||
"status": "🔴阻塞",
|
||||
"owner": "工程师",
|
||||
"actual_end": "-"
|
||||
},
|
||||
"A-T19": {
|
||||
"status": "🔴阻塞",
|
||||
"owner": "工程师",
|
||||
"actual_end": "-"
|
||||
},
|
||||
"B-T1": {
|
||||
"status": "✅已完成",
|
||||
"owner": "工程师",
|
||||
"actual_end": "07-02"
|
||||
},
|
||||
"B-T2": {
|
||||
"status": "✅已完成",
|
||||
"owner": "工程师",
|
||||
"actual_end": "07-02"
|
||||
},
|
||||
"B-T3": {
|
||||
"status": "✅已完成",
|
||||
"owner": "工程师",
|
||||
"actual_end": "07-02"
|
||||
},
|
||||
"B-T4": {
|
||||
"status": "✅已完成",
|
||||
"owner": "工程师",
|
||||
"actual_end": "07-02"
|
||||
},
|
||||
"B-T5": {
|
||||
"status": "✅已完成",
|
||||
"owner": "工程师",
|
||||
"actual_end": "07-02"
|
||||
},
|
||||
"B-T6": {
|
||||
"status": "✅已完成",
|
||||
"owner": "工程师",
|
||||
"actual_end": "07-02"
|
||||
},
|
||||
"B-T7": {
|
||||
"status": "✅已完成",
|
||||
"owner": "工程师",
|
||||
"actual_end": "07-02"
|
||||
},
|
||||
"B-T8": {
|
||||
"status": "🟢已完成",
|
||||
"owner": "QA",
|
||||
"actual_end": "✅ 33/33 PASS"
|
||||
},
|
||||
"B-T9": {
|
||||
"status": "🔵进行中",
|
||||
"owner": "QA",
|
||||
"actual_end": "-"
|
||||
},
|
||||
"B-T10": {
|
||||
"status": "⏳等待中",
|
||||
"owner": "工程师",
|
||||
"actual_end": "-"
|
||||
},
|
||||
"B-T11": {
|
||||
"status": "🔵进行中",
|
||||
"owner": "项目经理",
|
||||
"actual_end": "-"
|
||||
},
|
||||
"B-T12": {
|
||||
"status": "⏳等待中",
|
||||
"owner": "QA",
|
||||
"actual_end": "-"
|
||||
},
|
||||
"B-T13": {
|
||||
"status": "🔵进行中",
|
||||
"owner": "工程师",
|
||||
"actual_end": "-"
|
||||
},
|
||||
"B-T14": {
|
||||
"status": "🔵进行中",
|
||||
"owner": "工程师",
|
||||
"actual_end": "-"
|
||||
},
|
||||
"B-T15": {
|
||||
"status": "🔵进行中",
|
||||
"owner": "工程师",
|
||||
"actual_end": "-"
|
||||
},
|
||||
"B-T16": {
|
||||
"status": "🔵进行中",
|
||||
"owner": "QA+工程师",
|
||||
"actual_end": "-"
|
||||
},
|
||||
"B-T17": {
|
||||
"status": "🟢已修复",
|
||||
"owner": "工程师",
|
||||
"actual_end": "-"
|
||||
},
|
||||
"C-T1": {
|
||||
"status": "✅已完成",
|
||||
"owner": "工程师",
|
||||
"actual_end": "07-02"
|
||||
},
|
||||
"C-T2": {
|
||||
"status": "✅已完成",
|
||||
"owner": "工程师",
|
||||
"actual_end": "07-02"
|
||||
},
|
||||
"C-T3": {
|
||||
"status": "✅已完成",
|
||||
"owner": "工程师",
|
||||
"actual_end": "07-02"
|
||||
},
|
||||
"C-T4": {
|
||||
"status": "✅已完成",
|
||||
"owner": "工程师",
|
||||
"actual_end": "07-02"
|
||||
},
|
||||
"C-T5": {
|
||||
"status": "✅已完成",
|
||||
"owner": "工程师",
|
||||
"actual_end": "07-02"
|
||||
},
|
||||
"C-T6": {
|
||||
"status": "✅已完成",
|
||||
"owner": "工程师",
|
||||
"actual_end": "07-02"
|
||||
},
|
||||
"C-T7": {
|
||||
"status": "✅已完成",
|
||||
"owner": "工程师",
|
||||
"actual_end": "07-02"
|
||||
},
|
||||
"C-T8": {
|
||||
"status": "✅已完成",
|
||||
"owner": "工程师",
|
||||
"actual_end": "07-02"
|
||||
},
|
||||
"C-T9": {
|
||||
"status": "✅已完成",
|
||||
"owner": "工程师",
|
||||
"actual_end": "07-02"
|
||||
},
|
||||
"C-T10": {
|
||||
"status": "✅已完成",
|
||||
"owner": "工程师",
|
||||
"actual_end": "07-02"
|
||||
},
|
||||
"C-T11": {
|
||||
"status": "✅已完成",
|
||||
"owner": "工程师",
|
||||
"actual_end": "07-02"
|
||||
},
|
||||
"C-T12": {
|
||||
"status": "✅已完成",
|
||||
"owner": "工程师",
|
||||
"actual_end": "07-02"
|
||||
},
|
||||
"C-T13": {
|
||||
"status": "✅已完成",
|
||||
"owner": "工程师",
|
||||
"actual_end": "07-02"
|
||||
},
|
||||
"C-T14": {
|
||||
"status": "✅已完成",
|
||||
"owner": "项目经理",
|
||||
"actual_end": "07-02"
|
||||
},
|
||||
"C-T15": {
|
||||
"status": "✅已完成",
|
||||
"owner": "QA",
|
||||
"actual_end": "07-02"
|
||||
},
|
||||
"C-T16": {
|
||||
"status": "✅已完成",
|
||||
"owner": "工程师",
|
||||
"actual_end": "07-02"
|
||||
},
|
||||
"C-T17": {
|
||||
"status": "✅已完成",
|
||||
"owner": "工程师",
|
||||
"actual_end": "07-02"
|
||||
},
|
||||
"C-T18": {
|
||||
"status": "✅已完成",
|
||||
"owner": "工程师",
|
||||
"actual_end": "07-02"
|
||||
},
|
||||
"C-T19": {
|
||||
"status": "✅已完成",
|
||||
"owner": "QA+工程师",
|
||||
"actual_end": "07-02"
|
||||
}
|
||||
},
|
||||
"stats": {
|
||||
"total": 55,
|
||||
"not_started": 0,
|
||||
"ready": 2,
|
||||
"in_progress": 9,
|
||||
"completed": 39,
|
||||
"blocked": 3,
|
||||
"delayed": 0,
|
||||
"waiting": 2
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,72 @@
|
||||
# 审批类型扩展、卡片直跳与同窗口导航改造 — 部署完成
|
||||
|
||||
## TL;DR
|
||||
审批流程系统已从 5 种扩展到 **12 种审批类型 / 18 个审批流程**,H5 卡片点击可直接跳转对应审批 URL;为进一步贴合企微 H5 体验,已将跳转方式从 `window.open`(新标签页)改为 `window.location.href`(同窗口导航),由企微原生提供返回按钮。后端与前端代码修改均已部署到生产环境并通过验证。
|
||||
|
||||
## 交付概览
|
||||
|
||||
| 项目 | 状态 |
|
||||
|------|------|
|
||||
| 后端 `approval.py` 18 个模板 + 12 类关键词 | 已部署,API 验证通过 |
|
||||
| 前端 H5 17 个审批卡片 + URL 直跳 | 已构建部署,页面可正常访问 |
|
||||
| 审批卡片同窗口导航改造 | 已部署:两处 `window.open` 改为 `window.location.href` |
|
||||
| 生产容器状态 | backend / nginx 均健康运行 |
|
||||
| 浏览器端验证 | H5 登录页渲染正常,无 JS 报错 |
|
||||
| 已知问题 / 遗留 | 无 |
|
||||
| Dify System Prompt v2 | 已由用户在 Dify 后台发布 |
|
||||
| 功能文档归档 | PRD v2.2 + 架构文档 v2.2 已更新 |
|
||||
|
||||
## 新增审批类型(7 种)
|
||||
1. 会议室故障报修
|
||||
2. 企业应用管理
|
||||
3. 资产变更确认
|
||||
4. 终端设备网络准入
|
||||
5. 活动与会议技术支持
|
||||
6. 员工IT支持与故障报修
|
||||
7. 公共邮箱账号申请
|
||||
|
||||
## 本次导航改造说明
|
||||
|
||||
用户提出审批页面应"内嵌打开带有返回和关闭"。经技术验证:
|
||||
- 企微审批 URL 未设 `X-Frame-Options`,理论上可被 iframe 嵌入;
|
||||
- 但 H5 生产环境配置了 `Cross-Origin-Embedder-Policy: require-corp` + CSP `default-src 'self'`,跨域 iframe 会被安全头拦截;
|
||||
- 在不修改 nginx 安全头的前提下,**方案 A(同窗口导航)**为可行方案。
|
||||
|
||||
改造点:
|
||||
- `frontend-h5/src/components/chat/ApprovalCardModal.vue` 的 `handleSelect` 中,
|
||||
两处的 `window.open(url, '_blank')` 全部改为 `window.location.href = url`;
|
||||
- 移除跳转后的 `showToast` 提示(页面立即导航离开,toast 不可见)。
|
||||
|
||||
效果:在企微 H5 webview 中点击审批卡片项,会在当前 webview 内打开审批页面,企微原生顶部返回按钮负责返回 IT 服务台。
|
||||
|
||||
## 关键文件清单
|
||||
|
||||
### 新建文档
|
||||
- `docs/02-产品需求/approval_templates.json` — 结构化审批模板数据
|
||||
- `docs/02-产品需求/dify_approval_system_prompt_v2.md` — Dify System Prompt 更新文本
|
||||
|
||||
### 代码修改(已部署)
|
||||
- `backend/app/api/approval.py` — 18 个模板、12 类关键词
|
||||
- `frontend-h5/src/components/chat/ApprovalCardModal.vue` — 12 类卡片 + URL 直跳 + 同窗口导航
|
||||
|
||||
## 验证结果
|
||||
- 容器内 `/approval/templates` 返回 **18 个模板**
|
||||
- 容器内 `/approval/keywords` 返回关键词映射正确
|
||||
- `agent-browser` 打开 `https://itsupport.servyou.com.cn/h5/` 正常渲染登录页
|
||||
- nginx 容器内 `/h5/` 返回 301,前端文件已正确部署到 `/opt/wecom-it-desk/frontend-h5/dist/`
|
||||
- 生产容器:`wecom_it_backend` healthy, `wecom_it_nginx` running
|
||||
|
||||
## 用户下一步建议
|
||||
1. ~~登录 Dify 后台粘贴 System Prompt~~ → 已完成
|
||||
2. 在企微 H5 中输入类似"我要申请会议室维修"/"公共邮箱怎么开"/"资产变更"等触发审批卡片,确认点击卡片项后**在当前 webview 内跳转**到审批页面,并可用企微顶部返回按钮回到 IT 服务台。
|
||||
3. 如需真正的自定义返回/关闭覆盖层(方案 B),需评估是否放宽 H5 的 COEP/CSP 安全头;这会影响安全级别,需单独决策。
|
||||
4. 保留 `docs/02-产品需求/approval_templates.json` 作为后续审批流程变更的数据源。
|
||||
|
||||
## 文档归档
|
||||
|
||||
| 文档 | 更新内容 |
|
||||
|------|---------|
|
||||
| `docs/02-产品需求/IT智能服务台-产品需求文档PRD-v2.md` | 新增 v2.2 增量需求(P2-07~P2-11),含 12种/18流程完整表格、企微免登录结论、关联文档索引 |
|
||||
| `docs/03-技术架构/IT智能服务台-系统架构设计文档v2.md` | 15.4 节从 6 模板扩展到 12种/18流程;新增 15.4.4~15.4.8 共 5 个子节(意图识别架构、前端卡片架构、导航方案选型、跨应用免登录、API端点);9.1 外部集成表新增运维平台条目;版本 +v2.2 |
|
||||
| `docs/02-产品需求/approval_templates.json` | 18 个模板结构化数据(已有,无需修改) |
|
||||
| `docs/02-产品需求/dify_approval_system_prompt_v2.md` | Dify System Prompt v2(已有,已发布) |
|
||||
@@ -1,12 +1,247 @@
|
||||
# 早班巡检自动化 - 执行记录
|
||||
|
||||
## 2026-07-18 09:30 执行结果
|
||||
|
||||
**数据来源**:主文档第四章 v2.8 (2026-07-14) + 独立看板 `docs/10-项目管理/项目状态看板.md` v1.0 (07-17) + 上次巡检记忆 (07-17)
|
||||
**说明**:指定路径 `docs/10-项目管理/05-项目状态看板/01-项目状态看板.md` 连续第10次不存在。独立看板v1.0已更新(#76/#82 07-17入完成区),但主文档第四章未同步
|
||||
|
||||
### 关键发现
|
||||
1. **P0待办2项**:#81敏感词检测+语气优化(阻塞14天,约07-21到期需启动)、#104结构化日志查看页(无阻塞已4天未启动);#117 Neo4j已完成但仍列P0区未清理(数据不一致持续2次巡检)
|
||||
2. **P1待办3项**:#80坐席图片预览(数据不一致持续2次——已完成区07-16 vs P1清单仍列"待排查")、#73后端文件覆盖、#86流程图review
|
||||
3. **等用户决策2项,均超3天阈值**:企微会议室Secret(自07-11,7天)、ITSM API授权(自07-11,7天)— 需PM立即关注;联软网络不通标记"暂不处理"不视为卡点
|
||||
4. **进行中0项**:主文档和独立看板均为空。上次#76已于07-17完成
|
||||
5. **数据质量问题持续**:#81编号冲突(P0敏感词 vs 已完成粘贴图片)、#80/#117双重列出、主文档"已完成"区滞后(07-17 #76/#82未入区)
|
||||
6. **07-17完成2项**:#76 ITSM工单卡片跳转(桥接页+扫码登录)、#82 H5右侧栏布局调整 — 已入独立看板v1.0
|
||||
7. **看板路径第10次缺失**:指定路径连续10次巡检不存在,建议统一看板源
|
||||
|
||||
### 全局状态
|
||||
- P0待办:2项(#81约07-21到期、#104未启动4天)
|
||||
- P1待办:3项(#80可能已完成待确认)
|
||||
- 等决策:2项(均超3天阈值,7天)
|
||||
- 进行中:0项
|
||||
|
||||
---
|
||||
|
||||
## 2026-07-17 09:30 执行结果
|
||||
|
||||
**数据来源**:主文档第四章 v2.7+ (含07-16更新) + 独立看板 `docs/10-项目管理/项目状态看板.md` v1.0 (07-17) + 记忆文件 (07-16)
|
||||
**说明**:指定路径 `docs/10-项目管理/05-项目状态看板/01-项目状态看板.md` 连续第9次不存在。发现独立看板文件在 `docs/10-项目管理/项目状态看板.md` (v1.0, 07-17新建),与主文档第四章存在数据不一致
|
||||
|
||||
### 关键发现
|
||||
1. **P0待办2项**:#81敏感词检测+语气优化(延后至约07-21到期)、#104结构化日志查看页(无阻塞可启动);#117 Neo4j已完成但仍列P0区未清理
|
||||
2. **P1待办3项**:#80坐席图片预览(已在完成区07-16但P1仍列出,数据不一致)、#73后端文件覆盖、#86流程图review
|
||||
3. **等用户决策3项,2项超3天阈值**:企微会议室Secret(≥6天自07-11)、ITSM API授权(≥6天自07-11)— 需PM立即关注
|
||||
4. **进行中1项**:#76零信任VPN卡片免登录修复(P1,今日新建计划今日完成)——仅独立看板有记录,主文档"正在做"为空
|
||||
5. **#81编号冲突**:P0"敏感词检测+语气优化"与已完成"粘贴图片边框问题"共用#81
|
||||
6. **两份看板数据不一致**:独立看板v1.0(74完成/1进行中) vs 主文档(P0/P1/等决策分区仍含已完成项)
|
||||
7. **07-16完成4项未入独立看板已完成区**:#117 Neo4j、#82 坐席500错误、#81粘贴图片边框、#80企微图片预览
|
||||
|
||||
### 全局状态
|
||||
- P0待办:2项(#81延后中、#104可启动;#117已完成未清理)
|
||||
- P1待办:3项(#80可能已完成待确认)
|
||||
- 等决策:3项(2项超3天阈值)
|
||||
- 进行中:1项(#76,仅独立看板有记录)
|
||||
|
||||
---
|
||||
|
||||
## 2026-07-16 09:30 执行结果
|
||||
|
||||
**数据来源**:主文档 v2.8 (2026-07-14) + `.workbuddy/memory/2026-07-15.md` + `.workbuddy/memory/2026-07-14.md`
|
||||
**说明**:指定看板路径 `docs/10-项目管理/05-项目状态看板/01-项目状态看板.md` 仍不存在(连续8次),状态看板在主文档第四章
|
||||
|
||||
### 关键发现
|
||||
1. **P0待办2项**:#81敏感词检测+语气优化(07-14 PM决定延后1周,原阻塞自07-04共12天,约07-21到期需启动)、#104结构化日志查看页(无阻塞可直接启动)
|
||||
2. **P1待办2项**:#73后端文件覆盖、#86流程图零依赖review
|
||||
3. **等用户决策3项,2项超3天阈值**:企微会议室Secret(阻塞≥5天自07-11)、ITSM API授权(阻塞≥5天自07-11)— 需PM立即关注
|
||||
4. **进行中0项**:看板"正在做"区为空
|
||||
5. **07-14/07-15新产出未入看板**:/h5/ 404错误修复(nginx配置)、Token多IP异常检测功能部署(T001-T003测试通过)、审批模板ID不正确两轮修复(RecommendCard+ApprovalCardModal+DB)
|
||||
6. **看板路径持续缺失**:连续8次巡检不存在
|
||||
|
||||
### 全局状态
|
||||
- P0待办:2项
|
||||
- P1待办:2项
|
||||
- 等决策:3项(2项超3天阈值)
|
||||
- 进行中:0项
|
||||
|
||||
---
|
||||
|
||||
## 2026-07-15 09:30 执行结果
|
||||
|
||||
**数据来源**:主文档 v2.8 (2026-07-14) + `.workbuddy/memory/2026-07-14.md`
|
||||
**说明**:指定看板路径 `docs/10-项目管理/05-项目状态看板/01-项目状态看板.md` 仍不存在,状态看板在主文档第四章
|
||||
|
||||
### 关键发现
|
||||
1. **P0待办2项**:#81敏感词检测+语气优化(07-14 PM决定延后1周,原阻塞自07-04共11天)、#104结构化日志查看页(07-14新排入,无阻塞可直接启动)
|
||||
2. **P1待办2项**:#73后端文件覆盖、#86流程图零依赖review
|
||||
3. **等用户决策3项,2项超3天阈值**:企微会议室Secret(阻塞≥4天自07-11)、ITSM API授权(阻塞≥4天自07-11)— 需PM立即关注
|
||||
4. **进行中0项**:看板"正在做"区为空
|
||||
5. **07-14新产出未入看板**:/h5/ 404错误修复(nginx配置)、Token多IP异常检测功能部署(T001-T003测试通过)
|
||||
6. **看板路径持续缺失**:指定路径连续7次巡检不存在
|
||||
|
||||
### 全局状态
|
||||
- P0待办:2项
|
||||
- P1待办:2项
|
||||
- 等决策:3项(2项超3天阈值)
|
||||
- 进行中:0项
|
||||
|
||||
---
|
||||
|
||||
## 2026-07-14 09:30 执行结果(已更新看板 v2.6)
|
||||
|
||||
**数据来源**:主文档 v2.5 + 实际代码检测 + PM确认
|
||||
|
||||
### 关键发现(经实际检测确认)
|
||||
1. **#48 IP白名单**:✅ 已完成。检测nginx.conf确认已配置精确内网网段(10.0.0.0/8等)
|
||||
2. **#105 摇人Bug**:✅ 已完成。代码确认不再推送企微通知栏
|
||||
3. **#107 卷挂载**:✅ 已完成。docker-compose.yml确认./app:/app/app已配置
|
||||
4. **#88 RBAC**:✅ 粗粒度已满足需求,无需细粒度。PM确认
|
||||
5. **#81 敏感词**:⏸️ 延后1周
|
||||
6. **#104 日志页**:🆕 排入本期
|
||||
7. **#75 头像**:🔄 需重新测试
|
||||
8. **火绒AccessKey**:⚠️ 07-13测试时被假值覆盖
|
||||
|
||||
### 看板更新(v2.6)
|
||||
- #48/#107/#88 移至已完成区
|
||||
- #105 从P1移除
|
||||
- 新增"等用户决策"区块
|
||||
|
||||
### 全局状态(更新后)
|
||||
- P0待办:2项(#81/#104)
|
||||
- P1待办:3项
|
||||
- 等决策:4项
|
||||
- 进行中:1项
|
||||
|
||||
---
|
||||
|
||||
## 2026-07-13 09:30 执行结果
|
||||
|
||||
**数据来源**:`docs/10-项目管理/任务说明书/IT智能服务台-项目管理主文档.md` (v2.5, 2026-07-13) + `.workbuddy/memory/2026-07-13.md` + `.workbuddy/memory/2026-07-12.md`
|
||||
**说明**:指定看板路径 `docs/10-项目管理/05-项目状态看板/01-项目状态看板.md` 仍不存在,状态看板在主文档第四章;v2.5已更新至07-13,07-12/07-13产出已纳入
|
||||
|
||||
### 关键发现
|
||||
1. **P0阻塞2项持续未推进**:#48 IP白名单收窄(阻塞≈31天,自06-13)、#81 敏感词检测+语气优化(阻塞≈10天,自07-04,但隐私正则修复已完成)— 均>3天,需PM立即关注
|
||||
2. **#105数据不一致持续4次巡检**:已完成区(07-10)+P1清单双重列出,07-10/07-11/07-12/07-13四次巡检指出至今未修正 ⚠️数据质量
|
||||
3. **#107可能已完成未更新**:07-12/07-13日志显示bind mount方案已实际使用(`./app:/app/app`),但看板仍标为in_progress
|
||||
4. **#75/#88可能已完成**:#75头像同步07-08已交付12/12测试通过;#88 RBAC v0.7.1已完成6处装饰器修复(但细粒度权限可能未完成)
|
||||
5. **07-12布局优化v2.0的8个待明确事项已清除**(07-12日志确认),从等决策清单移除
|
||||
6. **Portal /itportal/ 500错误**:07-13测试发现,nginx静态文件问题,非后端错误,未入看板
|
||||
7. **等决策项从5项减至3项**:企微会议室Secret、ITSM API授权、ITSM代办API抓包仍在阻塞
|
||||
|
||||
### 全局状态
|
||||
- P0待办:3项(2项长期阻塞,#81部分完成)
|
||||
- P1待办:5项(#105应移除,#75/#88可能已完成)
|
||||
- 等决策:3项(布局优化事项已清除)
|
||||
- 进行中:2项(#107可能已完成)
|
||||
|
||||
### PM行动项
|
||||
1. 联系网络组确认代理IP段(#48阻塞31天)⚠️紧急
|
||||
2. 确认#81敏感词"语气优化"部分是否仍需开发(隐私正则已完成)⚠️紧急
|
||||
3. 从P1清单移除#105(连续4次巡检指出)⚠️数据质量
|
||||
4. 确认#107卷挂载改造是否已完成(bind mount已实际使用)
|
||||
5. 确认#75头像同步是否已完成(07-08已交付12/12测试)
|
||||
6. 确认#88 RBAC细粒度权限是否仍需开发(6处装饰器已修复)
|
||||
7. 企微管理后台申请会议室Secret
|
||||
8. 向ITSM平台方申请app_id/app_secret
|
||||
9. 排查Portal /itportal/ 500错误并入看板
|
||||
10. #104结构化日志查看页待启动(无阻塞,可排入sprint)
|
||||
|
||||
---
|
||||
|
||||
## 2026-07-12 09:30 执行结果
|
||||
|
||||
**数据来源**:`docs/10-项目管理/任务说明书/IT智能服务台-项目管理主文档.md` (v2.4, 2026-07-10) + `.workbuddy/memory/2026-07-11.md` + `.workbuddy/memory/2026-07-12.md`
|
||||
**说明**:指定看板路径仍不存在,状态看板在主文档第四章;看板版本滞后2天,07-11/07-12大量产出未入看板
|
||||
|
||||
### 关键发现
|
||||
1. **P0阻塞2项持续未推进**:#48 IP白名单收窄(阻塞≈29天,自06-13)、#81 敏感词检测(阻塞≈8天,自07-04)— 均>3天,需PM立即关注
|
||||
2. **#105数据不一致持续3次巡检**:已完成区+P1清单双重列出,07-04/07-10/07-11三次巡检指出至今未修正
|
||||
3. **07-11/07-12大量产出未入看板**:代办事项集成(8bug修复链)、IT资产审批推送、语音识别、截图拍照、复杂场景重构、会议室预定系统(全栈部署)、知识库迭代3功能(44文件43测试通过)、知识迭代3Bug修复、坐席v9/v10部署修复、AI辅助消息框+布局优化技术设计文档
|
||||
4. **进行中2项**:#91 忘记密码 + #107 后端卷挂载改造(后者07-11已恢复卷挂载,可能已完成需确认)
|
||||
5. **等用户决策5项**:企微会议室Secret未申请、ITSM API授权待申请、ITSM代办API待抓包、布局优化v2.0的8个待明确事项、知识库迭代待确认
|
||||
6. **看板版本严重滞后**:v2.4截止07-10,07-11全日+07-12产出均未入看板
|
||||
|
||||
### 全局状态
|
||||
- P0待办:3项(2项长期阻塞)
|
||||
- P1待办:5项(1项#105已完成未清理,1项#75可能已完成)
|
||||
- 等决策:5项
|
||||
- 进行中:2项(1项可能已完成)
|
||||
|
||||
### PM行动项
|
||||
1. 联系网络组确认代理IP段(#48阻塞29天)⚠️紧急
|
||||
2. 确认敏感词库来源/语气优化范围(#81阻塞8天)⚠️紧急
|
||||
3. 从P1清单移除#105(连续3次巡检指出)⚠️数据质量
|
||||
4. 确认#75头像同步是否已完成(07-08已交付12/12测试)
|
||||
5. 确认#107卷挂载改造是否已完成(07-11已恢复卷挂载)
|
||||
6. 企微管理后台申请会议室Secret
|
||||
7. 向ITSM平台方申请app_id/app_secret
|
||||
8. 确认布局优化v2.0的8个待明确事项
|
||||
9. 更新看板至v2.5+,将07-11/07-12产出纳入已完成区
|
||||
|
||||
---
|
||||
|
||||
## 2026-07-11 09:30 执行结果
|
||||
|
||||
**数据来源**:`docs/10-项目管理/任务说明书/IT智能服务台-项目管理主文档.md` (v2.4, 2026-07-10) + `.workbuddy/memory/2026-07-11.md`
|
||||
**说明**:指定看板路径仍不存在,状态看板在主文档第四章;07-11有大量新产出未反映在看板中
|
||||
|
||||
### 关键发现
|
||||
1. **P0阻塞2项持续未推进**:#48 IP白名单收窄(阻塞≈28天,自06-13)、#81 敏感词检测(阻塞≈7天,自07-04)— 均>3天,需PM立即关注
|
||||
2. **#105数据不一致持续**:已完成区+P1清单双重列出,上次巡检已指出至今未修正
|
||||
3. **07-11大量产出未入看板**:百度ASR部署、BYOD功能、业务路由推荐(81/81)、复杂场景重构(81/81)、邀请按钮修复、通讯录同步Secret、Mac企微语音最终修复、Dify API Key更新
|
||||
4. **进行中2项**:#91 忘记密码 + #107 后端卷挂载改造(后者可能已完成,需确认)
|
||||
5. **新增阻塞项**:企微可信IP白名单(errcode 48009)、Dify Prompt更新(3个功能等待)
|
||||
6. **看板版本滞后**:主文档v2.4截止07-10,07-11全日产出来入看板
|
||||
|
||||
### 全局状态
|
||||
- P0待办:3项(2项长期阻塞)
|
||||
- P1待办:5项(1项#105已完成未清理)
|
||||
- 等决策:4项
|
||||
- 进行中:2项
|
||||
|
||||
### PM行动项
|
||||
1. 联系网络组确认代理IP段(#48阻塞28天)⚠️紧急
|
||||
2. 确认敏感词库来源(#81阻塞7天)⚠️紧急
|
||||
3. 从P1清单移除#105(连续2次巡检指出)
|
||||
4. 企微管理后台添加可信IP 218.75.34.87
|
||||
5. Dify后台更新3个Prompt(BYOD/业务路由/复杂场景)
|
||||
6. 更新看板至v2.5,将07-11产出纳入已完成区
|
||||
7. 确认#107卷挂载改造是否已完成
|
||||
|
||||
---
|
||||
|
||||
## 2026-07-10 09:30 执行结果
|
||||
|
||||
**数据来源**:`docs/10-项目管理/任务说明书/IT智能服务台-项目管理主文档.md` (v2.0, 2026-07-10)
|
||||
**说明**:指定路径 `docs/10-项目管理/05-项目状态看板/01-项目状态看板.md` 不存在,状态看板已整合至主文档第四章
|
||||
|
||||
### 关键发现
|
||||
1. **P0阻塞2项长期未推进**:#48 IP白名单收窄(阻塞≈27天,自06-13)、#81 敏感词检测(阻塞≈6天,自07-04)
|
||||
2. **数据不一致**:#105 同时出现在"已完成"和"P1重要"分区,应从P1移除
|
||||
3. **07-10大量产出**:P0认证Bug修复5个、三端部署上线、生产热修复、访问控制部署、摇人Bug修复
|
||||
4. **进行中仅1项**:#91 忘记密码-企微扫码重置
|
||||
5. **风险待处理6项**:H-9/H-11/M-6/M-7/M-8/L-8/L-9
|
||||
|
||||
### 全局状态
|
||||
- P0待办:3项(2项延后阻塞,1项待启动)
|
||||
- P1待办:4项
|
||||
- 等决策:3项
|
||||
- 进行中:1项
|
||||
|
||||
### PM行动项
|
||||
1. 联系网络组确认代理IP段(#48阻塞27天)
|
||||
2. 确认敏感词库来源(#81阻塞6天)
|
||||
3. 从P1清单移除#105(已完成)
|
||||
4. 更新任务说明书中看板路径引用
|
||||
|
||||
---
|
||||
|
||||
## 2026-07-04 09:30 执行结果
|
||||
|
||||
**数据来源**:`.taskboard-cache/任务执行状态看板_cache.json`(缓存时间 2026-07-03T08:44:12)
|
||||
**⚠️ 原始看板文件缺失**:`docs/小组任务书/任务执行状态看板.md` 不存在,本次巡检基于缓存数据 + 07-03巡检记录 + REVIEW_B_T10.md 综合分析
|
||||
**⚠️ 原始看板文件缺失**:`docs/小组任务/任务执行状态看板.md` 不存在,本次巡检基于缓存数据 + 07-03巡检记录 + REVIEW_B_T10.md 综合分析
|
||||
|
||||
### 关键发现
|
||||
1. **看板源文件丢失**:`docs/小组任务书/任务执行状态看板.md` 路径不存在,该目录也未创建,PRD中有引用但实际文件缺失
|
||||
1. **看板源文件丢失**:`docs/小组任务/任务执行状态看板.md` 路径不存在,该目录也未创建,PRD中有引用但实际文件缺失
|
||||
2. **B-T10双重可激活信号**:①依赖B-T8已完成(33/33 PASS) ②REVIEW_B_T10.md显示代码评审已于07-03通过(IS_PASS: YES),但缓存中仍为⏳等待中
|
||||
3. **3个阻塞已逾期2天**:BLOCK-17(企微SSO)、BLOCK-19(扫码登录超时)、BLOCK-20(管理后台Network Error) 均 due 07-02,现已逾期2天
|
||||
4. **BLOCK-18状态矛盾持续**:B-T17标记🟢已修复,但BLOCK-18阻塞表仍为🔵排查中(07-03已发现,至今未修正)
|
||||
@@ -20,7 +255,7 @@
|
||||
- 整体:39/55 (71%) — 若计入🟢则41/55 (75%)
|
||||
|
||||
### PM行动项
|
||||
1. **恢复看板源文件** — `docs/小组任务书/任务执行状态看板.md` 缺失,需重建
|
||||
1. **恢复看板源文件** — `docs/小组任务/任务执行状态看板.md` 缺失,需重建
|
||||
2. 通知B组激活B-T10(依赖已完成 + 评审已通过)
|
||||
3. 优先解决3个逾期阻塞(BLOCK-17/19/20),已逾期2天
|
||||
4. 确认BLOCK-18/B-T17真实状态并校正看板
|
||||
|
||||
@@ -0,0 +1,34 @@
|
||||
# 项目任务检索 - 自动化执行记录
|
||||
|
||||
## 2026-07-06 配置变更
|
||||
|
||||
### 16:24:19
|
||||
- **用户请求**:改为单次执行 + 内部循环
|
||||
- **修改内容**:
|
||||
- scheduleType: recurring → once
|
||||
- scheduledAt: 2026-07-06T17:00:00
|
||||
- prompt: 增加内部循环逻辑(12次,约24小时)
|
||||
- **循环逻辑**:
|
||||
- 首次执行立即检查任务状态
|
||||
- 无新任务则静默等待2小时
|
||||
- 有新任务立即汇报
|
||||
- 12次循环后自动退出
|
||||
|
||||
---
|
||||
|
||||
## 2026-07-06 执行摘要
|
||||
|
||||
### 执行时间
|
||||
- 10:29:02 (首次执行)
|
||||
|
||||
### 执行结果
|
||||
1. 读取项目状态看板:发现1个进行中任务(#90 坐席/管理端直接登录)
|
||||
2. 读取任务说明书:确认任务详情
|
||||
3. 任务 #90 已完成部署测试,仅剩 Code Review
|
||||
|
||||
### 用户操作
|
||||
- 用户要求将检索周期从1小时改为2小时
|
||||
- 已更新 rrule: FREQ=HOURLY;INTERVAL=2
|
||||
|
||||
### 下次执行
|
||||
- 12:30:00 左右 (每2小时执行一次)
|
||||
@@ -1,123 +0,0 @@
|
||||
# IT智能服务台 - 项目记忆
|
||||
|
||||
## 锁定的设计决策
|
||||
- **AI交互原则**:小段多回合交互,禁止一次性大段回复
|
||||
- **文档管理**:统一保存 `docs/` 目录,按类型分子目录
|
||||
- **资源申请流程**:所有资源申请→`docs/资源申请清单.md`
|
||||
- **原型已锁定**:坐席v5.3 + H5 v1.1
|
||||
- **UI偏好**:企微浅色扁平风格,accent=#07C160
|
||||
- **术语统一**:"人工"=用户呼叫坐席;"摇人"=坐席呼叫坐席
|
||||
- **双企微应用**:正式(itsupport.servyou.com.cn) + 测试(已下线)
|
||||
- **统一入口架构**:`/itportal/` 角色选择 → user/agent/admin
|
||||
- **OTP双因素认证**:admin角色访问时验证
|
||||
|
||||
## 技术架构
|
||||
- **前端**:坐席(Vue3+Element Plus) / H5(Vue3+Vant4) / 管理后台(Vue3+Element+Tailwind)
|
||||
- **后端**:FastAPI + SQLAlchemy + PostgreSQL + Redis
|
||||
- **本地开发**:Python 3.12 venv + SQLite
|
||||
- **字段映射**:后端`id`/`sender_type` → H5前端`message_id`/`message_type`,映射层在 `frontend-h5/src/api/conversation.ts` 的 `mapMessage()`
|
||||
- **WS广播**:H5发消息后通过 `ws_manager.broadcast()` 实时推送给坐席
|
||||
- **API超时**:默认20s,消息发送30s,文件上传60s
|
||||
|
||||
## 部署
|
||||
- **NAS测试**:~~itdesk.amanzac.com~~ (已下线)
|
||||
- **正式服务器**:itsupport.servyou.com.cn (10.90.5.110)
|
||||
- **堡垒机**:sxn@10.212.189.210:2222 (OTP)
|
||||
- **文件上传**:只能通过堡垒机手动上传到 `/tmp/`
|
||||
|
||||
## 外部系统集成
|
||||
- **火绒企业版**:HMAC-SHA1认证,核心接口 `_leak`(高危漏洞) / `_virus_events`(病毒事件)
|
||||
- **联软LV7000**:三层认证,核心价值 `strusername` 字段=员工→终端映射
|
||||
- **Dify**:生产 `http://yw-dify.dc.servyou-it.com/dify2openai/`
|
||||
- **RAGFlow**:生产 `http://10.80.0.85:8080/` / API `:9380`
|
||||
- **aTrust**:HMAC-SHA256,待获取API密钥
|
||||
- **映射策略**:联软(主) > aTrust(VPN辅) > eHR(静态)
|
||||
|
||||
## 管理后台
|
||||
- 路由前缀 `/api/admin/`;权限 require_admin
|
||||
- 已实现:仪表盘/功能开关/坐席管理/分配模式/快速回复审核/集成配置/会话监控/会话审计/坐席绩效/系统日志/角色管理
|
||||
- 集成三种配置模式:url_key / access_key / account_password
|
||||
|
||||
## H5端消息推送
|
||||
- 双通道:企微消息(必达) + WebSocket(即时)
|
||||
- WS端点:`/ws/h5/{employee_id}?token=xxx`
|
||||
- 降级策略:WS断连→3秒轮询
|
||||
- **本地消息缓存 (v0.7.4+)**:
|
||||
- 登录后优先加载本地缓存消息,立即显示历史记录
|
||||
- 同时异步从后端获取最新消息,合并去重后更新缓存
|
||||
- 缓存key:`h5_messages_cache`,有效期7天,最多100条/会话
|
||||
- 发送消息和轮询时自动更新缓存
|
||||
- 登出时清除缓存
|
||||
|
||||
## 近期问题修复 (2026-07)
|
||||
- **OAuth重定向计数残留**:页面刷新后`oauth_redirect_count`未重置,导致误报"登录状态异常" → 在`employee.ts` store初始化时检测有效token后自动清除计数
|
||||
- **API响应解析错误**:Axios拦截器返回`{code:0, data:{}, message}`包装格式,但部分API直接访问`response.xxx`而非`response.data.xxx` → 修正`conversation.ts`中`sendMessage`函数的响应映射
|
||||
- **数据库缺失列**:`messages`表缺少`is_recalled`列 → `ALTER TABLE messages ADD COLUMN IF NOT EXISTS is_recalled BOOLEAN DEFAULT FALSE;`
|
||||
- **数据库列类型错误**:`messages.id`列为uuid类型但代码传入varchar → `ALTER TABLE messages ALTER COLUMN id TYPE character varying(36);`
|
||||
- **Nginx部署目录**:构建产物上传到`/opt/wecom-it-desk/frontend-h5/`但nginx挂载在`/opt/wecom-it-desk/html/itdesk/` → 部署时需复制文件到正确目录
|
||||
|
||||
## 五阶段演进
|
||||
1. MVP:转人工+H5+坐席+邀请+管理后台
|
||||
2. 完整流程:WS+排队+满意度+OAuth2
|
||||
3. AI Wingman+排查流程图
|
||||
4. 知识库+数据看板
|
||||
5. 自动化闭环
|
||||
|
||||
## 堡垒机运维 (jumpserver-ops)
|
||||
|
||||
**脚本位置**:`C:\Users\simon\.workbuddy\skills\jumpserver-ops\scripts\jms_ops.py`
|
||||
|
||||
### ⚠️ 服务器操作规则(重要)
|
||||
|
||||
**在对服务器进行任何操作时,优先使用 jumpserver-ops 自动完成,而非让用户手动操作。**
|
||||
|
||||
| 操作类型 | 自动执行方式 |
|
||||
|----------|-------------|
|
||||
| 远程命令 | `python jms_ops.py exec -c "命令"` |
|
||||
| 文件上传 | `python jms_ops.py upload 本地文件 /tmp/远程路径` |
|
||||
| 文件下载 | `python jms_ops.py download /tmp/远程文件 ./本地路径` |
|
||||
|
||||
### 使用方式
|
||||
|
||||
```bash
|
||||
# 第一次执行(自动登录并缓存会话)
|
||||
python jms_ops.py exec -c "hostname"
|
||||
|
||||
# 连续测试:使用 --reuse 复用会话(30分钟内有效,2-3秒执行)
|
||||
python jms_ops.py exec -c "uptime" --reuse
|
||||
python jms_ops.py exec -c "docker ps" -c "curl -s http://localhost/api/health" --reuse
|
||||
|
||||
# 文件上传(自动根据大小选择方式)
|
||||
# - ≤10MB: base64 编码传输(快速)
|
||||
# - >10MB: elFinder Web UI(浏览器自动化)
|
||||
python jms_ops.py upload local_file.txt /tmp/remote_file.txt
|
||||
|
||||
# 文件下载
|
||||
python jms_ops.py download /tmp/remote_file.txt local_file.txt
|
||||
|
||||
# 批量命令
|
||||
python jms_ops.py batch -f commands.txt
|
||||
|
||||
# 文件传输
|
||||
python jms_ops.py upload local.conf /tmp/remote.conf
|
||||
python jms_ops.py download /remote/path ./local.conf
|
||||
```
|
||||
|
||||
### 性能
|
||||
|
||||
| 场景 | 首次执行 | --reuse 复用 |
|
||||
|------|----------|--------------|
|
||||
| 单命令 | ~13s | ~2s |
|
||||
| 3 条命令 | ~13s | ~3s |
|
||||
|
||||
### 关键参数
|
||||
- `--reuse`:复用上次会话(减少登录次数,30分钟有效)
|
||||
- `--parallel`:并行模式(每命令独立 token+会话)
|
||||
- `--cmd-timeout`:每命令超时秒数(默认 15s)
|
||||
|
||||
## 文档关联修复 (2026-07-05)
|
||||
- **起因**:2026-07-04 docs/ 重组为数字编号子目录(01-项目总览~11-历史归档),但 mkdocs.yml nav / 文档间交叉引用 / 巡检自动化路径未同步,全面断链
|
||||
- **修复**:mkdocs.yml nav 9处断链重写(移除2个归档项,纳入5份新文档)+ 11处交叉引用修复 + 索引版本号修正(v1.0→v1.3) + 巡检自动化适配
|
||||
- **关键发现**:巡检 automation-1782986180887 原依赖的"小组任务书/任务执行状态看板.md"及A/B/C三组体系(认证加固16/消息系统16/AI数据19)从未创建,每日巡检必然失败;已适配为基于 01-项目状态看板.md 的状态巡检(P0/P1/等决策/进行中)
|
||||
- **修复报告**:docs/01-项目总览/文档关联修复报告-20260705.md
|
||||
- **保留未改**:目录树展示(历史快照)、归档文档内旧路径、历史任务标题
|
||||
@@ -0,0 +1,296 @@
|
||||
---
|
||||
name: task-intake
|
||||
description: 任务接收与路由技能 - 收到任何请求时首先使用,将请求结构化为四要素(是什么/要什么/怎么做/谁来做)并路由到正确的工作流。适用于项目所有 incoming 请求的统一入口。
|
||||
agent_created: true
|
||||
version: 1.4
|
||||
date: 2026-07-10
|
||||
---
|
||||
|
||||
# Task Intake — 任务接收与路由
|
||||
|
||||
## 定位
|
||||
|
||||
项目所有 incoming 请求的**统一入口**。不是执行者,是路由器。
|
||||
|
||||
收到请求后,本技能负责:
|
||||
1. **分类** — 判断请求属于哪类任务
|
||||
2. **结构化** — 输出四要素(是什么/要什么/怎么做/谁来做)
|
||||
3. **路由** — 对照 SOP 路由表,确定执行路径
|
||||
4. **移交** — 将路由卡交给对应执行方
|
||||
|
||||
**核心原则**:task-intake 只做"想清楚"和"分对路",不做"动手干"。
|
||||
|
||||
---
|
||||
|
||||
## 触发条件
|
||||
|
||||
- ✅ 收到任何新需求/问题/任务时
|
||||
- ✅ 不确定该走什么工作流时
|
||||
- ✅ 请求类型模糊,需要先分类时
|
||||
- ❌ 已经明确知道走哪条流程时(直接执行即可,不必再过一遍 intake)
|
||||
|
||||
---
|
||||
|
||||
## 执行流程
|
||||
|
||||
### Step 1: 请求分类
|
||||
|
||||
分析请求内容,判断属于以下哪一类:
|
||||
|
||||
| 分类 | 识别特征 | 示例 |
|
||||
|------|---------|------|
|
||||
| 🏗️ 新功能开发(中大型) | 多页面/多模块、涉及后端+前端、>10个源文件 | "开发员工自助查询平台" |
|
||||
| ⚡ 新功能开发(小型) | 单页面/工具脚本、≤10个源文件 | "加一个满意度评价导出功能" |
|
||||
| 🔧 Bug 修复 | 报告明确 Bug,非新功能 | "管理后台登录报网络连接失败" |
|
||||
| 🚀 部署运维 | 部署/配置/Nginx/容器相关 | "部署管理后台前端到生产" |
|
||||
| 🩺 故障排查 | 页面打不开/502/500/接口无响应 | "H5扫码登录后页面不关闭" |
|
||||
| 🔴 应急事件 | P0/P1 级别,需立即响应 | "鉴权漏洞被利用" |
|
||||
| 🔍 代码调试 | 代码逻辑不对、行为异常 | "摇人消息没有推送到通知栏" |
|
||||
| 📊 技术评估/决策 | 需要判断值不值得做、怎么选 | "联软API对接值不值得做?" |
|
||||
| 📋 方案调研 | 需要调研后输出方案 | "火绒API方案怎么设计?" |
|
||||
| 📝 文档更新 | 更新文档/SOP/手册 | "更新故障排查手册" |
|
||||
| 🛠️ 工具沉淀 | 排查后归档脚本/工具 | "把排查脚本归到工具箱" |
|
||||
|
||||
### Step 2: 四要素结构化
|
||||
|
||||
对每个请求输出以下四要素:
|
||||
|
||||
```
|
||||
是什么:[任务分类] + [一句话描述]
|
||||
要什么:[期望产出物] + [验收标准]
|
||||
怎么做:[执行路径] + [需要的技能/工具]
|
||||
谁来做:[执行角色] + [协作方]
|
||||
```
|
||||
|
||||
**注意事项**:
|
||||
- "是什么"要精确到分类表中的具体类别
|
||||
- "要什么"必须包含可验证的产出物和验收标准,验收标准需指明验证手段(见下方验证手段分层表)
|
||||
- "怎么做"指出执行路径和工具,但不展开执行细节
|
||||
- "谁来做"明确执行方和协作方
|
||||
|
||||
**验证手段分层表**(用于"要什么"字段的验收标准):
|
||||
|
||||
| 验证类型 | 工具 | 适用场景 | 何时必须用 |
|
||||
|---------|------|---------|-----------|
|
||||
| API/后端 | curl / HTTP 请求 | 接口返回值、状态码 | 后端接口验证 |
|
||||
| 前端渲染/登录/交互 | **agent-browser** 技能 | 页面渲染、表单填写、按钮点击、键盘输入 | 涉及前端页面的修复 **必须**用 |
|
||||
| 前端诊断(F12 等效) | **agent-browser** Debug 命令 | 白屏、JS 不执行、API 异常、CSP 违规 | 前端异常排查 **必须**采集 console/errors/network |
|
||||
| 服务器状态 | jumpserver-ops | 容器状态、进程 | 部署后健康检查 |
|
||||
|
||||
**硬规则**:禁止只因 `docker logs` 无报错就断言修复。前端类修复必须 agent-browser 截图取证。前端异常排查必须采集 `errors` + `console` + `network requests`。
|
||||
|
||||
### Step 3: 路由决策
|
||||
|
||||
对照项目 SOP 路由表,确定执行路径:
|
||||
|
||||
| 输入特征 | 路由到 | 产出物 | 执行方 | 参考文档 |
|
||||
|---------|--------|--------|--------|---------|
|
||||
| 🏗️ 新功能(中大型) | 软件团队标准 SOP | PRD+架构+代码+测试 | PM→Architect→Engineer→QA | 软件团队 SOP |
|
||||
| ⚡ 新功能(小型) | 软件团队快速模式 | 代码+测试 | Engineer→QA | 软件团队 SOP |
|
||||
| 🔧 Bug 修复 | SOP §6 BugFix | 修复+验证 | Engineer→QA | SOP §6 |
|
||||
| 🚀 部署运维 | 直接执行 ⚠️ 前置检查 | 部署完成+验证 | AI+jumpserver-ops | SOP §7 工具箱 + deploy-troubleshoot Step -1 |
|
||||
| 🩺 故障排查 | deploy-troubleshoot | 定位+修复+案例 | 三步隔离法 | 故障排查手册 |
|
||||
| 🔴 应急事件 | SOP §4 应急响应 | 止血+根因 | 应急流程 | SOP §4 |
|
||||
| 🔍 代码调试 | diagnose 技能 | 根因+回归 | 六阶段调试 | diagnose SKILL.md |
|
||||
| 📊 技术评估 | Plan 模式 | 评估报告 | AI+人 | — |
|
||||
| 📋 方案调研 | Plan 模式 | 方案文档 | AI+人 | — |
|
||||
| 📝 文档更新 | 直接执行 | 文档 | AI | SOP §5 文档规范 |
|
||||
| 🛠️ 工具沉淀 | SOP §7 流程 | 工具归档+README更新 | AI | SOP §7 |
|
||||
|
||||
**路由优先级**(当请求可能匹配多个分类时):
|
||||
1. 🔴 应急事件 > 一切(先止血再说)
|
||||
2. 🩺 故障排查 > 🔧 Bug 修复(先隔离定位再修 Bug)
|
||||
3. 🏗️/⚡ 新功能 > 📊 技术评估(明确要做的不需要评估)
|
||||
4. 📝 文档更新 / 🛠️ 工具沉淀 通常作为其他任务的收尾步骤
|
||||
|
||||
### Step 3.1: 部署运维前置检查(⚠️ 涉及后端代码变更时必须执行)
|
||||
|
||||
当路由到「🚀 部署运维」且涉及后端代码变更时,**在执行部署前必须检查**:
|
||||
|
||||
#### ⛔ 硬规则:后端代码部署方式(方案 C 卷挂载,2026-07-10 上线)
|
||||
|
||||
| 变更类型 | 部署命令 | 禁止操作 | 耗时 |
|
||||
|---------|---------|---------|------|
|
||||
| `.py` 文件变更(新增/修改) | `docker compose restart backend` | ❌ `docker compose build` | ~15-30 秒 |
|
||||
| `requirements.txt` 变更 | `docker compose build backend && docker compose up -d backend` | — | ~60-90 秒 |
|
||||
| 配置文件变更(`.env`/`docker-compose.yml`) | `docker compose up -d backend` | — | ~10 秒 |
|
||||
|
||||
> **原理**:代码通过 `./app:/app/app` volume 挂载到容器,不烘焙进镜像。改代码只需 restart 让 uvicorn 重新加载,无需重建镜像。`docker compose build` 只在 Python 依赖(requirements.txt)变化时才需要。
|
||||
|
||||
| 检查项 | 命令 | 不通过时的动作 |
|
||||
|--------|------|---------------|
|
||||
| 代码目录完整性 | `for f in app/__init__.py app/main.py app/api/auth.py; do [ -f "/opt/wecom-it-desk/$f" ] && echo "PASS: $f" || echo "FAIL: $f"; done` | 上传缺失文件到 `/opt/wecom-it-desk/app/` |
|
||||
| Volume 挂载验证 | `docker exec wecom_it_backend ls /app/app/main.py` | 检查 docker-compose.yml 是否含 `./app:/app/app` 卷挂载 |
|
||||
| 代码一致性 | `HOST=$(md5sum /opt/wecom-it-desk/app/main.py \| awk '{print $1}') && CONTAINER=$(docker exec wecom_it_backend md5sum /app/app/main.py \| awk '{print $1}') && [ "$HOST" = "$CONTAINER" ] && echo PASS \| echo FAIL` | `docker compose restart backend` 重新加载代码 |
|
||||
|
||||
> **方案 C(卷挂载)已于 2026-07-10 上线**:代码不再烘焙进 Docker 镜像,通过 `./app:/app/app` volume 挂载。`backend/app/` 旧代码目录已删除。代码更新只需 `docker compose restart`,仅 `requirements.txt` 变化时才需 `docker compose build`。
|
||||
>
|
||||
> **完整检查清单**见 `deploy-troubleshoot` 技能 Step -1 和故障排查手册 §1.4。
|
||||
|
||||
### Step 4: 输出路由卡
|
||||
|
||||
```markdown
|
||||
## 任务路由卡
|
||||
|
||||
**是什么**: [任务分类] [一句话描述]
|
||||
**要什么**: [产出物] [验收标准]
|
||||
**怎么做**: [执行路径] [技能/工具]
|
||||
**谁来做**: [执行角色] [协作方]
|
||||
|
||||
**路由到**: [工作流名称]
|
||||
**预计阶段**: [阶段列表]
|
||||
**参考文档**: [SOP章节/技能/手册]
|
||||
```
|
||||
|
||||
路由卡输出后,**立即移交**给对应执行方,不在此步骤中展开执行。
|
||||
|
||||
---
|
||||
|
||||
## 与软件团队 SOP 的集成
|
||||
|
||||
当齐活林(交付总监)收到请求时:
|
||||
|
||||
```
|
||||
请求到达
|
||||
↓
|
||||
齐活林调用 task-intake
|
||||
↓
|
||||
输出路由卡
|
||||
↓
|
||||
├─ 路由到"标准SOP" → TeamCreate → PM → Architect → Engineer → QA
|
||||
├─ 路由到"快速模式" → TeamCreate → Engineer → QA
|
||||
├─ 路由到"BugFix" → TeamCreate → Engineer → QA
|
||||
├─ 路由到"故障排查" → deploy-troubleshoot → jumpserver-ops(传输)
|
||||
├─ 路由到"应急响应" → SOP §4 应急流程
|
||||
├─ 路由到"Plan模式" → 先想后做,输出评估/方案文档
|
||||
└─ 路由到"直接执行" → 文档更新/工具沉淀
|
||||
```
|
||||
|
||||
**关键**:task-intake 是齐活林判断工作流类型的**结构化工具**,替代原来的"凭经验判断"。判断结果可追溯、可复盘。
|
||||
|
||||
---
|
||||
|
||||
## 与其他技能的关系
|
||||
|
||||
```
|
||||
task-intake (路由器)
|
||||
/ | | \
|
||||
/ | | \
|
||||
deploy-troubleshoot diagnose 软件团队SOP Plan模式
|
||||
(故障排查方法论) (代码调试) (开发流程) (评估决策)
|
||||
| | |
|
||||
jumpserver-ops Bash/Read Engineer/QA
|
||||
(传输代理) (执行工具) (执行角色)
|
||||
|
|
||||
toolbox/
|
||||
(弹药库)
|
||||
```
|
||||
|
||||
- **task-intake** = 路由器,决定走哪条路
|
||||
- **deploy-troubleshoot / diagnose** = 方法论,指导怎么排查
|
||||
- **jumpserver-ops** = 传输代理,解决"怎么到服务器"
|
||||
- **toolbox/** = 弹药库,提供辅助工具
|
||||
- **软件团队 SOP** = 开发流程,指导代码实现
|
||||
- **Plan 模式** = 思考模式,用于评估/决策类任务
|
||||
|
||||
---
|
||||
|
||||
## 使用示例
|
||||
|
||||
### 示例 1: "帮我加一个满意度评价导出功能"
|
||||
|
||||
```markdown
|
||||
## 任务路由卡
|
||||
|
||||
**是什么**: ⚡ 新功能开发(小型)— 满意度评价数据导出为 Excel
|
||||
**要什么**: 导出功能代码 + QA 验证通过
|
||||
**怎么做**: 软件团队快速模式 → Engineer 实现 → QA 验证
|
||||
**谁来做**: 寇豆码(工程师) → 严过关(QA)
|
||||
|
||||
**路由到**: 软件团队快速模式
|
||||
**预计阶段**: TeamCreate → Engineer → QA
|
||||
**参考文档**: 软件团队 SOP
|
||||
```
|
||||
|
||||
### 示例 2: "管理后台登录报网络连接失败"
|
||||
|
||||
```markdown
|
||||
## 任务路由卡
|
||||
|
||||
**是什么**: 🩺 故障排查 — 管理后台登录接口无响应
|
||||
**要什么**: 故障定位 + 修复 + 验证证据(curl 接口返回 + agent-browser 登录截图)
|
||||
**怎么做**: deploy-troubleshoot 三步隔离法 → jumpserver-ops 传输
|
||||
**谁来做**: AI(排查) + jumpserver-ops(传输)
|
||||
|
||||
**路由到**: deploy-troubleshoot
|
||||
**预计阶段**: Step 0(响应头) → 三步隔离 → 修复 → 验证
|
||||
**参考文档**: 00-标准故障排查手册.md
|
||||
```
|
||||
|
||||
### 示例 3: "联软 API 对接值不值得做?"
|
||||
|
||||
```markdown
|
||||
## 任务路由卡
|
||||
|
||||
**是什么**: 📊 技术评估 — 联软 API 对接的成本收益分析
|
||||
**要什么**: 评估报告(技术可行性 + 成本 + 收益 + 风险 + 建议)
|
||||
**怎么做**: Plan 模式 → 调研 → 分析 → 输出报告
|
||||
**谁来做**: AI(调研分析) + 宋献(决策)
|
||||
|
||||
**路由到**: Plan 模式
|
||||
**预计阶段**: 调研 → 分析 → 输出评估报告 → 人工决策
|
||||
**参考文档**: 无(Plan 模式自由发挥)
|
||||
```
|
||||
|
||||
### 示例 4: "H5 扫码登录后页面不自动关闭"
|
||||
|
||||
```markdown
|
||||
## 任务路由卡
|
||||
|
||||
**是什么**: 🔧 Bug 修复 — 扫码登录成功页 JS 未执行
|
||||
**要什么**: Bug 定位 + 修复 + 回归验证(agent-browser 打开扫码页 → 截图确认 JS 执行 + 页面自动关闭)
|
||||
**怎么做**: 先 deploy-troubleshoot 排查(确认是否部署层问题)→ 如是代码层则 diagnose 调试
|
||||
**谁来做**: AI(排查) → Engineer(修复) → QA(验证)
|
||||
|
||||
**路由到**: 先故障排查,确认层级后转 BugFix
|
||||
**预计阶段**: 隔离定位 → 根因分析 → 修复 → 验证 → 案例沉淀
|
||||
**参考文档**: 00-标准故障排查手册.md + SOP §6 BugFix
|
||||
```
|
||||
|
||||
### 示例 5: "把排查脚本归到工具箱"
|
||||
|
||||
```markdown
|
||||
## 任务路由卡
|
||||
|
||||
**是什么**: 🛠️ 工具沉淀 — 排查过程产生的脚本归档
|
||||
**要什么**: 脚本归位 + README 更新 + 根目录清理
|
||||
**怎么做**: SOP §7 工具沉淀流程(评估→归档→登记→清理)
|
||||
**谁来做**: AI
|
||||
|
||||
**路由到**: 直接执行(SOP §7)
|
||||
**预计阶段**: 评估复用价值 → 归档 → 登记README → 清理
|
||||
**参考文档**: SOP §7 部署运维工具箱管理
|
||||
```
|
||||
|
||||
---
|
||||
|
||||
## 上下文隔离原则
|
||||
|
||||
task-intake 的路由卡**只传递结论,不传递思考过程**:
|
||||
|
||||
- ✅ 传递:"故障定位在 Nginx 层,证据是 curl 返回 403"
|
||||
- ❌ 不传递:"我一开始以为是后端的问题,试了 A/B/C 都不对,后来才发现..."
|
||||
|
||||
这确保下一阶段(如 diagnose 或 Engineer)拿到的是**干净的输入**,不会被前一阶段的假设和试错过程带偏。
|
||||
|
||||
---
|
||||
|
||||
## 版本历史
|
||||
|
||||
| 版本 | 日期 | 变更 |
|
||||
|------|------|------|
|
||||
| v1.0 | 2026-07-10 | 初始版本,含 11 类任务分类 + 路由表 + 5 个示例 |
|
||||
| v1.1 | 2026-07-10 | 新增验证手段分层表,示例补充 agent-browser 验证要求 |
|
||||
| v1.2 | 2026-07-10 | 验证手段分层表新增"前端诊断(F12 等效)"类型,硬规则增加 console/errors/network 采集要求 |
|
||||
| v1.3 | 2026-07-10 | 新增 Step 3.1 部署运维前置检查(代码同步 + 依赖同步),防止镜像缺文件 |
|
||||
| v1.4 | 2026-07-10 | 方案 C 上线:Step 3.1 更新为 volume 挂载验证(代码完整性+挂载状态+一致性检查) |
|
||||
|
After Width: | Height: | Size: 186 KiB |
@@ -1,365 +0,0 @@
|
||||
# =============================================================================
|
||||
# IT智能服务台 — 审批流程 API
|
||||
# =============================================================================
|
||||
# 说明:提供审批模板管理和API提交功能
|
||||
# - 模板详情获取
|
||||
# - API提交审批申请
|
||||
# - 审批状态回调处理
|
||||
# =============================================================================
|
||||
|
||||
import logging
|
||||
import os
|
||||
from typing import Optional
|
||||
|
||||
import httpx
|
||||
from fastapi import APIRouter, Depends, Query
|
||||
from pydantic import BaseModel
|
||||
import redis.asyncio as aioredis
|
||||
|
||||
from app.config import settings
|
||||
from app.utils.token_manager import ApprovalTokenManager
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
router = APIRouter()
|
||||
|
||||
# Redis客户端(依赖注入)
|
||||
async def get_redis() -> aioredis.Redis:
|
||||
"""获取Redis客户端依赖"""
|
||||
from app.main import redis_client
|
||||
return redis_client
|
||||
|
||||
|
||||
# =============================================================================
|
||||
# 审批模板配置(从环境变量读取)
|
||||
# =============================================================================
|
||||
# 环境变量:
|
||||
# APPROVAL_TEMPLATE_RESOURCE - 资源申请模板ID
|
||||
# APPROVAL_TEMPLATE_DEVICE - 设备申请模板ID
|
||||
|
||||
APPROVAL_TEMPLATE_RESOURCE = os.getenv("APPROVAL_TEMPLATE_RESOURCE", "")
|
||||
APPROVAL_TEMPLATE_DEVICE = os.getenv("APPROVAL_TEMPLATE_DEVICE", "")
|
||||
|
||||
# 动态构建审批模板配置
|
||||
APPROVAL_TEMPLATES = {}
|
||||
|
||||
if APPROVAL_TEMPLATE_RESOURCE:
|
||||
APPROVAL_TEMPLATES[APPROVAL_TEMPLATE_RESOURCE] = {
|
||||
"id": APPROVAL_TEMPLATE_RESOURCE,
|
||||
"name": "资源申请",
|
||||
"type": "jump",
|
||||
"keywords": ["申请资源", "要资源", "申请"],
|
||||
}
|
||||
|
||||
if APPROVAL_TEMPLATE_DEVICE:
|
||||
APPROVAL_TEMPLATES[APPROVAL_TEMPLATE_DEVICE] = {
|
||||
"id": APPROVAL_TEMPLATE_DEVICE,
|
||||
"name": "设备申请",
|
||||
"type": "api",
|
||||
"keywords": ["申请设备", "要设备", "电脑", "笔记本"],
|
||||
}
|
||||
|
||||
|
||||
# =============================================================================
|
||||
# Schema 定义
|
||||
# =============================================================================
|
||||
|
||||
class ApprovalTemplateResponse(BaseModel):
|
||||
"""审批模板响应"""
|
||||
id: str
|
||||
name: str
|
||||
type: str
|
||||
keywords: list[str]
|
||||
|
||||
|
||||
class ApprovalJumpRequest(BaseModel):
|
||||
"""跳转审批请求"""
|
||||
template_id: str
|
||||
employee_id: Optional[str] = None
|
||||
|
||||
|
||||
class ApprovalJumpResponse(BaseModel):
|
||||
"""跳转审批响应"""
|
||||
url: str
|
||||
template_name: str
|
||||
|
||||
|
||||
class ApprovalContentItem(BaseModel):
|
||||
"""审批表单控件内容"""
|
||||
control: str # 控件类型: Text, Textarea, Number, Money, Date, Selector, Contact, etc.
|
||||
id: str # 控件ID
|
||||
value: dict # 控件值
|
||||
|
||||
|
||||
class ApprovalSubmitRequest(BaseModel):
|
||||
"""API提交审批请求"""
|
||||
template_id: str
|
||||
employee_id: str # 申请人userid
|
||||
contents: list[ApprovalContentItem] # 表单内容
|
||||
use_template_approver: int = 1 # 1-使用模板预设流程
|
||||
|
||||
|
||||
class ApprovalSubmitResponse(BaseModel):
|
||||
"""API提交审批响应"""
|
||||
sp_no: str # 审批单号
|
||||
template_name: str
|
||||
|
||||
|
||||
class ApprovalCallbackRequest(BaseModel):
|
||||
"""审批回调请求(XML解析后的模型)"""
|
||||
sp_no: str
|
||||
sp_name: str
|
||||
template_id: str
|
||||
apply_time: int
|
||||
applyer_userid: str
|
||||
sp_status: int # 1-审批中 2-已通过 3-已驳回 4-已撤销 6-通过后撤销 7-已删除 10-已支付
|
||||
status_change_event: int # 1-提单 2-同意 3-驳回 4-转审 5-催办 6-撤销 8-通过后撤销 10-添加备注
|
||||
|
||||
|
||||
# =============================================================================
|
||||
# 企微API调用辅助函数
|
||||
# =============================================================================
|
||||
|
||||
async def get_approval_token(redis: aioredis.Redis) -> str:
|
||||
"""获取审批应用access_token"""
|
||||
manager = ApprovalTokenManager(redis)
|
||||
return await manager.get_token()
|
||||
|
||||
|
||||
async def get_template_detail(access_token: str, template_id: str) -> dict:
|
||||
"""获取审批模板详情
|
||||
|
||||
对应企微API:
|
||||
POST https://qyapi.weixin.qq.com/cgi-bin/oa/gettemplatedetail
|
||||
|
||||
返回模板内的控件构成及控件ID
|
||||
"""
|
||||
url = "https://qyapi.weixin.qq.com/cgi-bin/oa/gettemplatedetail"
|
||||
params = {"access_token": access_token}
|
||||
|
||||
async with httpx.AsyncClient(timeout=httpx.Timeout(connect=10.0, read=30.0)) as client:
|
||||
response = await client.post(url, params=params, json={"template_id": template_id})
|
||||
result = response.json()
|
||||
|
||||
if result.get("errcode") != 0:
|
||||
logger.error(f"获取模板详情失败: {result.get('errmsg')}")
|
||||
raise Exception(f"获取模板详情失败: {result.get('errmsg')}")
|
||||
|
||||
return result
|
||||
|
||||
|
||||
async def submit_approval_api(
|
||||
access_token: str,
|
||||
template_id: str,
|
||||
creator_userid: str,
|
||||
contents: list[dict],
|
||||
use_template_approver: int = 1
|
||||
) -> dict:
|
||||
"""提交审批申请
|
||||
|
||||
对应企微API:
|
||||
POST https://qyapi.weixin.qq.com/cgi-bin/oa/applyevent
|
||||
|
||||
Args:
|
||||
access_token: 审批应用access_token
|
||||
template_id: 模板ID
|
||||
creator_userid: 申请人userid
|
||||
contents: 表单控件内容列表
|
||||
use_template_approver: 1-使用模板预设流程
|
||||
|
||||
Returns:
|
||||
{"sp_no": "审批单号"}
|
||||
"""
|
||||
url = "https://qyapi.weixin.qq.com/cgi-bin/oa/applyevent"
|
||||
params = {"access_token": access_token}
|
||||
|
||||
payload = {
|
||||
"creator_userid": creator_userid,
|
||||
"template_id": template_id,
|
||||
"use_template_approver": use_template_approver,
|
||||
"apply_data": {
|
||||
"contents": contents
|
||||
}
|
||||
}
|
||||
|
||||
async with httpx.AsyncClient(timeout=httpx.Timeout(connect=10.0, read=30.0)) as client:
|
||||
response = await client.post(url, params=params, json=payload)
|
||||
result = response.json()
|
||||
|
||||
if result.get("errcode") != 0:
|
||||
logger.error(f"提交审批失败: {result.get('errmsg')}")
|
||||
raise Exception(f"提交审批失败: {result.get('errmsg')}")
|
||||
|
||||
return {"sp_no": result.get("sp_no")}
|
||||
|
||||
|
||||
# =============================================================================
|
||||
# API 端点
|
||||
# =============================================================================
|
||||
|
||||
@router.get("/approval/templates", response_model=list[ApprovalTemplateResponse])
|
||||
async def get_approval_templates():
|
||||
"""获取所有审批模板列表"""
|
||||
return list(APPROVAL_TEMPLATES.values())
|
||||
|
||||
|
||||
@router.get("/approval/templates/{template_id}", response_model=ApprovalTemplateResponse)
|
||||
async def get_approval_template(template_id: str):
|
||||
"""获取指定审批模板详情"""
|
||||
if template_id not in APPROVAL_TEMPLATES:
|
||||
from fastapi import HTTPException
|
||||
raise HTTPException(status_code=404, detail="模板不存在")
|
||||
return APPROVAL_TEMPLATES[template_id]
|
||||
|
||||
|
||||
@router.get("/approval/templates/{template_id}/detail")
|
||||
async def get_template_full_detail(
|
||||
template_id: str,
|
||||
redis: aioredis.Redis = Depends(get_redis)
|
||||
):
|
||||
"""获取审批模板完整详情(控件结构)"""
|
||||
if template_id not in APPROVAL_TEMPLATES:
|
||||
from fastapi import HTTPException
|
||||
raise HTTPException(status_code=404, detail="模板不存在")
|
||||
|
||||
try:
|
||||
token = await get_approval_token(redis)
|
||||
detail = await get_template_detail(token, template_id)
|
||||
return detail
|
||||
except Exception as e:
|
||||
from fastapi import HTTPException
|
||||
raise HTTPException(status_code=500, detail=str(e))
|
||||
|
||||
|
||||
@router.post("/approval/jump", response_model=ApprovalJumpResponse)
|
||||
async def create_approval_jump(request: ApprovalJumpRequest):
|
||||
"""生成跳转审批链接"""
|
||||
template = APPROVAL_TEMPLATES.get(request.template_id)
|
||||
if not template:
|
||||
from fastapi import HTTPException
|
||||
raise HTTPException(status_code=404, detail="模板不存在")
|
||||
|
||||
if template["type"] != "jump":
|
||||
from fastapi import HTTPException
|
||||
raise HTTPException(status_code=400, detail="该模板不支持跳转方式")
|
||||
|
||||
# 生成跳转URL(企微审批链接格式)
|
||||
jump_url = f"https://qyapi.weixin.qq.com/cgi-bin/oa/applyevent?access_token=TOKEN&template_id={request.template_id}"
|
||||
|
||||
return ApprovalJumpResponse(
|
||||
url=jump_url,
|
||||
template_name=template["name"],
|
||||
)
|
||||
|
||||
|
||||
@router.post("/approval/submit", response_model=ApprovalSubmitResponse)
|
||||
async def submit_approval(
|
||||
request: ApprovalSubmitRequest,
|
||||
redis: aioredis.Redis = Depends(get_redis)
|
||||
):
|
||||
"""API提交审批申请"""
|
||||
template = APPROVAL_TEMPLATES.get(request.template_id)
|
||||
if not template:
|
||||
from fastapi import HTTPException
|
||||
raise HTTPException(status_code=404, detail="模板不存在")
|
||||
|
||||
if template["type"] != "api":
|
||||
from fastapi import HTTPException
|
||||
raise HTTPException(status_code=400, detail="该模板不支持API提交")
|
||||
|
||||
try:
|
||||
# 1. 获取审批token
|
||||
token = await get_approval_token(redis)
|
||||
|
||||
# 2. 转换contents格式
|
||||
contents = [item.model_dump() for item in request.contents]
|
||||
|
||||
# 3. 提交审批
|
||||
result = await submit_approval_api(
|
||||
access_token=token,
|
||||
template_id=request.template_id,
|
||||
creator_userid=request.employee_id,
|
||||
contents=contents,
|
||||
use_template_approver=request.use_template_approver
|
||||
)
|
||||
|
||||
return ApprovalSubmitResponse(
|
||||
sp_no=result["sp_no"],
|
||||
template_name=template["name"],
|
||||
)
|
||||
|
||||
except Exception as e:
|
||||
from fastapi import HTTPException
|
||||
raise HTTPException(status_code=500, detail=str(e))
|
||||
|
||||
|
||||
@router.post("/approval/callback")
|
||||
async def approval_callback(
|
||||
sp_no: str = Query(...),
|
||||
sp_name: str = Query(...),
|
||||
template_id: str = Query(...),
|
||||
apply_time: int = Query(...),
|
||||
applyer_userid: str = Query(...),
|
||||
sp_status: int = Query(...),
|
||||
status_change_event: int = Query(...)
|
||||
):
|
||||
"""审批状态变化回调处理
|
||||
|
||||
对应企微审批回调事件: sys_approval_change
|
||||
|
||||
状态变化类型 (status_change_event):
|
||||
1 - 提单
|
||||
2 - 同意
|
||||
3 - 驳回
|
||||
4 - 转审
|
||||
5 - 催办
|
||||
6 - 撤销
|
||||
8 - 通过后撤销
|
||||
10 - 添加备注
|
||||
|
||||
审批单状态 (sp_status):
|
||||
1 - 审批中
|
||||
2 - 已通过
|
||||
3 - 已驳回
|
||||
4 - 已撤销
|
||||
6 - 通过后撤销
|
||||
7 - 已删除
|
||||
10 - 已支付
|
||||
"""
|
||||
logger.info(f"审批回调: sp_no={sp_no}, status={sp_status}, event={status_change_event}")
|
||||
|
||||
# TODO: 根据业务需求处理审批状态变化
|
||||
# 例如:
|
||||
# - 审批通过后,更新IT服务台待办状态
|
||||
# - 审批驳回后,通知申请人
|
||||
# - 审批撤销后,关闭相关工单
|
||||
|
||||
event_map = {
|
||||
1: "submitted",
|
||||
2: "approved",
|
||||
3: "rejected",
|
||||
4: "transferred",
|
||||
5: "reminded",
|
||||
6: "revoked",
|
||||
8: "revoked_after_approved",
|
||||
10: "commented"
|
||||
}
|
||||
|
||||
event_type = event_map.get(status_change_event, f"unknown_{status_change_event}")
|
||||
logger.info(f"审批事件类型: {event_type}")
|
||||
|
||||
return {"errcode": 0, "errmsg": "ok"}
|
||||
|
||||
|
||||
@router.get("/approval/keywords")
|
||||
async def get_approval_keywords():
|
||||
"""获取所有审批关键词(用于前端关键词检测)"""
|
||||
keywords = []
|
||||
for template in APPROVAL_TEMPLATES.values():
|
||||
for kw in template["keywords"]:
|
||||
keywords.append({
|
||||
"keyword": kw,
|
||||
"template_id": template["id"],
|
||||
"template_name": template["name"],
|
||||
"type": template["type"],
|
||||
})
|
||||
return keywords
|
||||
@@ -1,439 +0,0 @@
|
||||
# =============================================================================
|
||||
# 企微IT智能服务台 — 待办事项 API
|
||||
# =============================================================================
|
||||
# 说明:提供待办事项的 CRUD 接口
|
||||
# 接口列表:
|
||||
# GET /api/todo-items — 获取当前坐席待办列表
|
||||
# GET /api/todo-items/{id} — 获取待办详情
|
||||
# PUT /api/todo-items/{id}/status — 更新待办状态
|
||||
# Mock: 预置示例待办数据,不连接真实外部系统
|
||||
# =============================================================================
|
||||
|
||||
from datetime import datetime
|
||||
from typing import List, Optional
|
||||
|
||||
from fastapi import APIRouter, HTTPException
|
||||
from pydantic import BaseModel, Field
|
||||
|
||||
from app.utils.response import success_response, AppException
|
||||
|
||||
# 创建路由器
|
||||
router = APIRouter(prefix="/todo-items", tags=["待办事项"])
|
||||
|
||||
|
||||
# --------------------------------------------------------------------------
|
||||
# 请求/响应 Schema
|
||||
# --------------------------------------------------------------------------
|
||||
|
||||
class TodoStatusUpdateRequest(BaseModel):
|
||||
"""更新待办状态请求 Schema。"""
|
||||
status: str = Field(..., description="新状态: pending/processing/resolved")
|
||||
|
||||
|
||||
class TodoItemResponse(BaseModel):
|
||||
"""待办事项响应 Schema。"""
|
||||
id: str
|
||||
type: str
|
||||
title: str
|
||||
priority: str
|
||||
description: dict
|
||||
status: str
|
||||
assigned_agent_id: Optional[str] = None
|
||||
corp_id: str = ""
|
||||
created_at: str
|
||||
updated_at: str
|
||||
|
||||
|
||||
class TodoItemListResponse(BaseModel):
|
||||
"""待办事项列表响应 Schema。"""
|
||||
items: List[TodoItemResponse]
|
||||
total: int
|
||||
|
||||
|
||||
# --------------------------------------------------------------------------
|
||||
# Mock 数据 — 预置示例待办(共 20 条,覆盖全部类型 × 状态)
|
||||
# --------------------------------------------------------------------------
|
||||
MOCK_TODO_ITEMS: List[dict] = [
|
||||
# ========== 工单(ticket)==========
|
||||
# 待处理
|
||||
{
|
||||
"id": "todo-001",
|
||||
"type": "ticket",
|
||||
"title": "VPN连接失败 — 财务部张伟",
|
||||
"priority": "urgent",
|
||||
"description": {
|
||||
"employee_name": "张伟",
|
||||
"department": "财务部",
|
||||
"error": "VPN Error 691",
|
||||
"steps": ["检查账号状态", "重置密码", "检查VPN配置"],
|
||||
},
|
||||
"status": "pending",
|
||||
"assigned_agent_id": "agent-001",
|
||||
"corp_id": "ww1234567890",
|
||||
"created_at": "2026-06-05T09:15:00Z",
|
||||
"updated_at": "2026-06-05T09:15:00Z",
|
||||
},
|
||||
{
|
||||
"id": "todo-007",
|
||||
"type": "ticket",
|
||||
"title": "OA系统登录异常 — 人事部刘芳",
|
||||
"priority": "urgent",
|
||||
"description": {
|
||||
"employee_name": "刘芳",
|
||||
"department": "人事部",
|
||||
"error": "页面白屏,控制台报500错误",
|
||||
"affected_count": 15,
|
||||
},
|
||||
"status": "pending",
|
||||
"assigned_agent_id": "agent-001",
|
||||
"corp_id": "ww1234567890",
|
||||
"created_at": "2026-06-05T11:30:00Z",
|
||||
"updated_at": "2026-06-05T11:30:00Z",
|
||||
},
|
||||
{
|
||||
"id": "todo-009",
|
||||
"type": "ticket",
|
||||
"title": "WiFi 无法连接 — 研发部开放区",
|
||||
"priority": "urgent",
|
||||
"description": {
|
||||
"employee_name": "陈明",
|
||||
"department": "研发部",
|
||||
"error": "获取IP失败,提示无法连接到此网络",
|
||||
"location": "3楼开放区",
|
||||
},
|
||||
"status": "pending",
|
||||
"assigned_agent_id": "agent-001",
|
||||
"corp_id": "ww1234567890",
|
||||
"created_at": "2026-06-06T08:00:00Z",
|
||||
"updated_at": "2026-06-06T08:00:00Z",
|
||||
},
|
||||
{
|
||||
"id": "todo-017",
|
||||
"type": "ticket",
|
||||
"title": "鼠标失灵 — 行政部周婷",
|
||||
"priority": "normal",
|
||||
"description": {
|
||||
"employee_name": "周婷",
|
||||
"department": "行政部",
|
||||
"error": "USB鼠标间歇性失灵,更换接口无效",
|
||||
"os": "Windows 11",
|
||||
},
|
||||
"status": "pending",
|
||||
"assigned_agent_id": "agent-001",
|
||||
"corp_id": "ww1234567890",
|
||||
"created_at": "2026-06-06T09:00:00Z",
|
||||
"updated_at": "2026-06-06T09:00:00Z",
|
||||
},
|
||||
# 进行中
|
||||
{
|
||||
"id": "todo-004",
|
||||
"type": "ticket",
|
||||
"title": "邮箱容量告警 — 市场部王强",
|
||||
"priority": "high",
|
||||
"description": {
|
||||
"employee_name": "王强",
|
||||
"department": "市场部",
|
||||
"current_usage": "4.8GB / 5GB",
|
||||
"action": "协助清理或申请扩容",
|
||||
},
|
||||
"status": "processing",
|
||||
"assigned_agent_id": "agent-001",
|
||||
"corp_id": "ww1234567890",
|
||||
"created_at": "2026-06-04T14:30:00Z",
|
||||
"updated_at": "2026-06-05T08:00:00Z",
|
||||
},
|
||||
{
|
||||
"id": "todo-010",
|
||||
"type": "ticket",
|
||||
"title": "ERP系统响应慢 — 全公司反馈",
|
||||
"priority": "high",
|
||||
"description": {
|
||||
"employee_name": "多个员工",
|
||||
"department": "全公司",
|
||||
"error": "ERP首页加载超过15秒",
|
||||
"affected_count": 50,
|
||||
},
|
||||
"status": "processing",
|
||||
"assigned_agent_id": "agent-001",
|
||||
"corp_id": "ww1234567890",
|
||||
"created_at": "2026-06-05T10:00:00Z",
|
||||
"updated_at": "2026-06-05T15:00:00Z",
|
||||
},
|
||||
# 已完成
|
||||
{
|
||||
"id": "todo-011",
|
||||
"type": "ticket",
|
||||
"title": "打印机驱动安装 — 市场部赵敏",
|
||||
"priority": "normal",
|
||||
"description": {
|
||||
"employee_name": "赵敏",
|
||||
"department": "市场部",
|
||||
"device_model": "Canon LBP2900",
|
||||
"solution": "从官网下载驱动并安装,测试打印正常",
|
||||
},
|
||||
"status": "resolved",
|
||||
"assigned_agent_id": "agent-001",
|
||||
"corp_id": "ww1234567890",
|
||||
"created_at": "2026-06-01T09:00:00Z",
|
||||
"updated_at": "2026-06-02T16:00:00Z",
|
||||
},
|
||||
|
||||
# ========== 审批(approval)==========
|
||||
# 待处理
|
||||
{
|
||||
"id": "todo-002",
|
||||
"type": "approval",
|
||||
"title": "软件安装审批 — 设计部PS申请",
|
||||
"priority": "high",
|
||||
"description": {
|
||||
"employee_name": "李娜",
|
||||
"department": "设计部",
|
||||
"software": "Adobe Photoshop 2026",
|
||||
"license_type": "企业许可",
|
||||
},
|
||||
"status": "pending",
|
||||
"assigned_agent_id": "agent-001",
|
||||
"corp_id": "ww1234567890",
|
||||
"created_at": "2026-06-05T10:20:00Z",
|
||||
"updated_at": "2026-06-05T10:20:00Z",
|
||||
},
|
||||
{
|
||||
"id": "todo-005",
|
||||
"type": "approval",
|
||||
"title": "权限升级审批 — 研发部数据库访问",
|
||||
"priority": "high",
|
||||
"description": {
|
||||
"employee_name": "陈明",
|
||||
"department": "研发部",
|
||||
"target_system": "生产数据库",
|
||||
"access_level": "只读",
|
||||
"approver": "研发总监",
|
||||
},
|
||||
"status": "pending",
|
||||
"assigned_agent_id": "agent-001",
|
||||
"corp_id": "ww1234567890",
|
||||
"created_at": "2026-06-05T08:45:00Z",
|
||||
"updated_at": "2026-06-05T08:45:00Z",
|
||||
},
|
||||
{
|
||||
"id": "todo-008",
|
||||
"type": "approval",
|
||||
"title": "新员工设备采购审批 — Q3批次",
|
||||
"priority": "normal",
|
||||
"description": {
|
||||
"batch": "Q3新员工",
|
||||
"count": 5,
|
||||
"items": ["笔记本x5", "显示器x5", "键鼠套装x5"],
|
||||
"budget": "65,000元",
|
||||
},
|
||||
"status": "pending",
|
||||
"assigned_agent_id": "agent-001",
|
||||
"corp_id": "ww1234567890",
|
||||
"created_at": "2026-06-05T07:00:00Z",
|
||||
"updated_at": "2026-06-05T07:00:00Z",
|
||||
},
|
||||
{
|
||||
"id": "todo-018",
|
||||
"type": "approval",
|
||||
"title": "弹性福利审批 — 全体员工Q3",
|
||||
"priority": "normal",
|
||||
"description": {
|
||||
"applicant": "人事部",
|
||||
"type": "弹性福利",
|
||||
"budget_per_person": "3000元",
|
||||
"total_count": 120,
|
||||
},
|
||||
"status": "pending",
|
||||
"assigned_agent_id": "agent-001",
|
||||
"corp_id": "ww1234567890",
|
||||
"created_at": "2026-06-06T07:00:00Z",
|
||||
"updated_at": "2026-06-06T07:00:00Z",
|
||||
},
|
||||
# 进行中
|
||||
{
|
||||
"id": "todo-012",
|
||||
"type": "approval",
|
||||
"title": "预算审批 — IT部Q3采购",
|
||||
"priority": "high",
|
||||
"description": {
|
||||
"department": "IT部",
|
||||
"amount": "280,000元",
|
||||
"items": ["服务器x2", "防火墙x2", "交换机x4"],
|
||||
"approver": "CFO",
|
||||
},
|
||||
"status": "processing",
|
||||
"assigned_agent_id": "agent-001",
|
||||
"corp_id": "ww1234567890",
|
||||
"created_at": "2026-06-04T09:00:00Z",
|
||||
"updated_at": "2026-06-05T14:00:00Z",
|
||||
},
|
||||
# 已完成
|
||||
{
|
||||
"id": "todo-013",
|
||||
"type": "approval",
|
||||
"title": "会议室预订审批 — 销售部Q3客户拜访",
|
||||
"priority": "normal",
|
||||
"description": {
|
||||
"employee_name": "刘军",
|
||||
"department": "销售部",
|
||||
"room": "5楼大会议室",
|
||||
"time": "2026-06-10 14:00-17:00",
|
||||
"result": "已批准",
|
||||
},
|
||||
"status": "resolved",
|
||||
"assigned_agent_id": "agent-001",
|
||||
"corp_id": "ww1234567890",
|
||||
"created_at": "2026-05-28T08:00:00Z",
|
||||
"updated_at": "2026-05-29T10:00:00Z",
|
||||
},
|
||||
|
||||
# ========== 设备(device)==========
|
||||
# 待处理
|
||||
{
|
||||
"id": "todo-003",
|
||||
"type": "device",
|
||||
"title": "工位打印机故障 — 3楼A区",
|
||||
"priority": "normal",
|
||||
"description": {
|
||||
"location": "3楼A区打印间",
|
||||
"device_model": "HP LaserJet Pro M404",
|
||||
"issue": "卡纸,无法打印",
|
||||
},
|
||||
"status": "pending",
|
||||
"assigned_agent_id": "agent-001",
|
||||
"corp_id": "ww1234567890",
|
||||
"created_at": "2026-06-05T11:05:00Z",
|
||||
"updated_at": "2026-06-05T11:05:00Z",
|
||||
},
|
||||
{
|
||||
"id": "todo-014",
|
||||
"type": "device",
|
||||
"title": "核心交换机故障 — 机房",
|
||||
"priority": "urgent",
|
||||
"description": {
|
||||
"location": "机房A区",
|
||||
"device_model": "Cisco Catalyst 9300",
|
||||
"issue": "端口3-12全部down,影响2楼所有工位",
|
||||
"affected_count": 45,
|
||||
},
|
||||
"status": "pending",
|
||||
"assigned_agent_id": "agent-001",
|
||||
"corp_id": "ww1234567890",
|
||||
"created_at": "2026-06-06T00:30:00Z",
|
||||
"updated_at": "2026-06-06T00:30:00Z",
|
||||
},
|
||||
# 进行中
|
||||
{
|
||||
"id": "todo-006",
|
||||
"type": "device",
|
||||
"title": "会议室投影仪维修 — 5楼大会议室",
|
||||
"priority": "normal",
|
||||
"description": {
|
||||
"location": "5楼大会议室",
|
||||
"device_model": "Epson EB-X51",
|
||||
"issue": "投影模糊,可能灯泡老化",
|
||||
},
|
||||
"status": "processing",
|
||||
"assigned_agent_id": "agent-001",
|
||||
"corp_id": "ww1234567890",
|
||||
"created_at": "2026-06-03T16:00:00Z",
|
||||
"updated_at": "2026-06-04T10:00:00Z",
|
||||
},
|
||||
{
|
||||
"id": "todo-015",
|
||||
"type": "device",
|
||||
"title": "服务器硬盘更换 — 虚拟化集群",
|
||||
"priority": "high",
|
||||
"description": {
|
||||
"location": "机房B区",
|
||||
"device_model": "Dell R740",
|
||||
"issue": "硬盘预警,需更换并做好数据迁移",
|
||||
"affected_vms": 12,
|
||||
},
|
||||
"status": "processing",
|
||||
"assigned_agent_id": "agent-001",
|
||||
"corp_id": "ww1234567890",
|
||||
"created_at": "2026-06-05T09:00:00Z",
|
||||
"updated_at": "2026-06-05T16:00:00Z",
|
||||
},
|
||||
# 已完成
|
||||
{
|
||||
"id": "todo-016",
|
||||
"type": "device",
|
||||
"title": "员工笔记本磁盘扩容 — 人事部吴婷",
|
||||
"priority": "normal",
|
||||
"description": {
|
||||
"employee_name": "吴婷",
|
||||
"department": "人事部",
|
||||
"device_model": "ThinkPad X1 Carbon",
|
||||
"solution": "更换1TB SSD,克隆系统,测试正常",
|
||||
},
|
||||
"status": "resolved",
|
||||
"assigned_agent_id": "agent-001",
|
||||
"corp_id": "ww1234567890",
|
||||
"created_at": "2026-05-20T13:00:00Z",
|
||||
"updated_at": "2026-05-22T17:00:00Z",
|
||||
},
|
||||
]
|
||||
|
||||
|
||||
# --------------------------------------------------------------------------
|
||||
# API 接口
|
||||
# --------------------------------------------------------------------------
|
||||
|
||||
@router.get("")
|
||||
async def list_todo_items(
|
||||
status: Optional[str] = None,
|
||||
priority: Optional[str] = None,
|
||||
):
|
||||
"""获取当前坐席待办列表。
|
||||
|
||||
支持按状态和优先级过滤。
|
||||
"""
|
||||
items = MOCK_TODO_ITEMS
|
||||
|
||||
# 按状态过滤
|
||||
if status:
|
||||
items = [item for item in items if item["status"] == status]
|
||||
|
||||
# 按优先级过滤
|
||||
if priority:
|
||||
items = [item for item in items if item["priority"] == priority]
|
||||
|
||||
# 按优先级排序:urgent → high → normal
|
||||
priority_order = {"urgent": 0, "high": 1, "normal": 2}
|
||||
items = sorted(items, key=lambda x: priority_order.get(x["priority"], 3))
|
||||
|
||||
return success_response(data={
|
||||
"items": [TodoItemResponse(**item).model_dump() for item in items],
|
||||
"total": len(items),
|
||||
})
|
||||
|
||||
|
||||
@router.get("/{item_id}")
|
||||
async def get_todo_item(item_id: str):
|
||||
"""获取待办事项详情。"""
|
||||
for item in MOCK_TODO_ITEMS:
|
||||
if item["id"] == item_id:
|
||||
return success_response(data=TodoItemResponse(**item).model_dump())
|
||||
raise AppException(code=1003, message=f"待办事项 {item_id} 不存在")
|
||||
|
||||
|
||||
@router.put("/{item_id}/status")
|
||||
async def update_todo_item_status(item_id: str, request: TodoStatusUpdateRequest):
|
||||
"""更新待办事项状态。"""
|
||||
# 校验状态值
|
||||
valid_statuses = {"pending", "processing", "resolved"}
|
||||
if request.status not in valid_statuses:
|
||||
raise HTTPException(
|
||||
status_code=400,
|
||||
detail=f"无效的状态值: {request.status},合法值为: {valid_statuses}",
|
||||
)
|
||||
|
||||
for item in MOCK_TODO_ITEMS:
|
||||
if item["id"] == item_id:
|
||||
item["status"] = request.status
|
||||
item["updated_at"] = datetime.now().isoformat()
|
||||
return success_response(data=TodoItemResponse(**item).model_dump())
|
||||
|
||||
raise AppException(code=1003, message=f"待办事项 {item_id} 不存在")
|
||||
@@ -1,321 +0,0 @@
|
||||
# =============================================================================
|
||||
# 企微IT智能服务台 — WebSocket 端点
|
||||
# =============================================================================
|
||||
# 说明:提供 WebSocket 端点,供坐席前端和H5用户端建立长连接,实现实时推送。
|
||||
# 核心功能:
|
||||
# 1. 接受坐席的 WebSocket 连接请求(含 token 认证)— /ws/{agent_id}
|
||||
# 2. 接受H5员工的 WebSocket 连接请求(含 token 认证)— /ws/h5/{employee_id}
|
||||
# 3. 维持连接,监听客户端消息(主要是心跳 ping)
|
||||
# 4. 连接断开时自动清理注册信息
|
||||
# 安全(WS-01):
|
||||
# 握手时从 query param 取 token → 查 Redis 验证 → 不通过则 close(code=4001)
|
||||
# 防止未授权用户冒充坐席/员工建立 WS 连接
|
||||
#
|
||||
# 端点路径:
|
||||
# - 坐席端:/ws/{agent_id}?token=xxx
|
||||
# - H5员工端:/ws/h5/{employee_id}?token=xxx
|
||||
# 为什么不挂 /api 前缀:WebSocket 不是 REST API,不走 Vite 的 /api 代理配置
|
||||
# =============================================================================
|
||||
|
||||
import logging
|
||||
|
||||
from fastapi import APIRouter, WebSocket, WebSocketDisconnect
|
||||
|
||||
from app.services.ws_manager import manager as ws_manager
|
||||
from app.services.cache_service import cache_service
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
# WebSocket 路由器(不挂 /api 前缀,直接注册在应用根路径)
|
||||
router = APIRouter()
|
||||
|
||||
# 认证失败时的 WebSocket 关闭码
|
||||
# 4001 = 自定义码,表示"未授权"(4000+ 为应用自定义范围)
|
||||
WS_CLOSE_UNAUTHORIZED = 4001
|
||||
|
||||
|
||||
@router.websocket("/ws/{agent_id}")
|
||||
async def websocket_endpoint(
|
||||
websocket: WebSocket,
|
||||
agent_id: str,
|
||||
) -> None:
|
||||
"""坐席 WebSocket 端点主循环(含 WS-01 token 认证)。
|
||||
|
||||
做什么:
|
||||
1. 从 Authorization header 获取 token(优先)或 query param(兼容)
|
||||
2. 验证 token 有效性(查 Redis)
|
||||
3. 验证 token 与 agent_id 一致性(防冒充)
|
||||
4. 认证通过后接受连接,注册到 ConnectionManager
|
||||
5. 进入消息接收循环,处理客户端发送的消息
|
||||
6. 连接断开时清理注册信息
|
||||
|
||||
为什么需要 token 认证(WS-01):
|
||||
- 之前 /ws/{agent_id} 无任何认证,任何人知道 URL 即可冒充任意坐席
|
||||
- 攻击者可监听所有消息、发送伪造消息,是 P0 级安全漏洞
|
||||
- 修复后,必须提供与 agent_id 匹配的有效 token 才能建立连接
|
||||
|
||||
安全改进(P0-#4):
|
||||
- 优先从 Authorization: Bearer {token} header 获取 token
|
||||
- 兼容从 ?token= URL 参数获取(向后兼容)
|
||||
- 不再将 token 暴露在 URL 中,避免 access_log 泄露
|
||||
|
||||
v0.5.1 修复:移除 `request: Request` 参数(部分 Starlette 版本注入 Request 失败,
|
||||
改用 `websocket.headers` 和 `websocket.query_params` 读取 header/query)
|
||||
|
||||
Args:
|
||||
websocket: FastAPI WebSocket 对象(框架自动注入)
|
||||
agent_id: 坐席ID(从 URL 路径参数获取)
|
||||
"""
|
||||
# ======================================================================
|
||||
# WS-01: Token 认证(从 subprotocol / header / query 获取)
|
||||
# ======================================================================
|
||||
|
||||
# 步骤1: 优先从 Sec-WebSocket-Protocol (subprotocol) 获取 token,其次从 Authorization header,最后从 query(向后兼容)
|
||||
# 格式: Sec-WebSocket-Protocol: bearer.{token}
|
||||
# 说明: 浏览器原生 WebSocket API 不支持 headers 参数,但支持 subprotocols (第2参数数组)
|
||||
# 前端用 new WebSocket(url, ["bearer.{token}"]) 传递,服务端从 sec-websocket-protocol 头读取
|
||||
subprotocol = websocket.headers.get("sec-websocket-protocol", "")
|
||||
if subprotocol.startswith("bearer."):
|
||||
token = subprotocol[7:] # 去掉 "bearer." 前缀
|
||||
else:
|
||||
# 其次从 Authorization header 获取
|
||||
auth_header = websocket.headers.get("Authorization", "")
|
||||
if auth_header.startswith("Bearer "):
|
||||
token = auth_header[7:] # 去掉 "Bearer " 前缀
|
||||
else:
|
||||
# 向后兼容:从 query param 获取(即将废弃)
|
||||
token = websocket.query_params.get("token", "")
|
||||
|
||||
# 步骤2: 检查 token 是否为空
|
||||
if not token:
|
||||
# 先 accept 再 close,否则客户端收不到关闭帧
|
||||
await websocket.accept()
|
||||
await websocket.close(code=WS_CLOSE_UNAUTHORIZED, reason="Missing token")
|
||||
logger.warning(f"WebSocket 拒绝连接: agent_id={agent_id}, 原因=缺少token")
|
||||
return
|
||||
|
||||
# 步骤3: 从 Redis 查询 token 对应的坐席信息
|
||||
# Redis 中存储格式: agent:token:{token} -> agent_user_id
|
||||
# (与坐席登录 API /api/agents/login 存储格式一致)
|
||||
try:
|
||||
stored_agent_id = await cache_service.get(f"agent:token:{token}")
|
||||
except Exception as e:
|
||||
# Redis 不可用时必须拒绝连接:token 验证依赖 Redis,无法验证身份
|
||||
# 如果降级放行,攻击者可在 Redis 故障时用任意 agent_id 冒充坐席
|
||||
logger.error(f"Redis 查询失败,拒绝 WS 连接: agent_id={agent_id}, error={e}")
|
||||
await websocket.accept()
|
||||
await websocket.close(
|
||||
code=WS_CLOSE_UNAUTHORIZED,
|
||||
reason="Authentication service unavailable"
|
||||
)
|
||||
return
|
||||
|
||||
# 步骤4: 验证 token 与 agent_id 一致性
|
||||
if not stored_agent_id:
|
||||
# token 不存在(已过期或伪造)
|
||||
await websocket.accept()
|
||||
await websocket.close(code=WS_CLOSE_UNAUTHORIZED, reason="Invalid or expired token")
|
||||
logger.warning(f"WebSocket 拒绝连接: agent_id={agent_id}, 原因=token无效或已过期")
|
||||
return
|
||||
|
||||
if stored_agent_id != agent_id:
|
||||
# token 对应的坐席与请求的 agent_id 不匹配(冒充)
|
||||
await websocket.accept()
|
||||
await websocket.close(code=WS_CLOSE_UNAUTHORIZED, reason="Token-agent mismatch")
|
||||
logger.warning(
|
||||
f"WebSocket 拒绝连接: agent_id={agent_id}, "
|
||||
f"原因=token对应坐席{stored_agent_id}与请求不匹配"
|
||||
)
|
||||
return
|
||||
|
||||
# ======================================================================
|
||||
# 认证通过,建立连接
|
||||
# ======================================================================
|
||||
|
||||
# 注册连接(内部会调用 websocket.accept(),并回显协商的 subprotocol)
|
||||
await ws_manager.connect(agent_id, websocket, subprotocol=subprotocol)
|
||||
logger.info(f"坐席 WebSocket 连接已认证: agent_id={agent_id}")
|
||||
|
||||
try:
|
||||
# 消息接收循环
|
||||
# 保持连接打开,监听客户端发来的消息
|
||||
# 即使客户端不发消息,这个循环也必须保持,否则连接会关闭
|
||||
while True:
|
||||
# 等待接收客户端消息(阻塞等待)
|
||||
data = await websocket.receive_json()
|
||||
|
||||
# 处理心跳 ping
|
||||
# 前端每 30 秒发送一次 ping,后端回复 pong
|
||||
# 作用:检测连接是否存活,防止中间代理(如 Nginx)因超时断开连接
|
||||
if data.get("type") == "ping":
|
||||
await websocket.send_json({"type": "pong"})
|
||||
logger.debug(f"WebSocket 心跳: agent_id={agent_id}")
|
||||
|
||||
# 处理输入指示器 typing 事件
|
||||
# 前端在用户输入时发送 typing 事件,后端广播给同一会话的其他参与者
|
||||
elif data.get("type") == "typing":
|
||||
conversation_id = data.get("conversation_id")
|
||||
sender_name = data.get("sender_name", agent_id)
|
||||
if conversation_id:
|
||||
# 广播给所有坐席(包含 sender_type 和 sender_id,
|
||||
# 前端可据此过滤掉自己的 typing 事件)
|
||||
await ws_manager.broadcast({
|
||||
"type": "typing",
|
||||
"data": {
|
||||
"conversation_id": conversation_id,
|
||||
"sender_id": agent_id,
|
||||
"sender_name": sender_name,
|
||||
"sender_type": "agent",
|
||||
}
|
||||
})
|
||||
|
||||
else:
|
||||
# 未来可扩展处理其他类型的客户端消息
|
||||
logger.debug(
|
||||
f"WebSocket 收到未知消息: agent_id={agent_id}, "
|
||||
f"type={data.get('type', 'unknown')}"
|
||||
)
|
||||
|
||||
except WebSocketDisconnect:
|
||||
# 客户端主动断开连接(正常行为)
|
||||
# 清理 ConnectionManager 中的注册信息
|
||||
ws_manager.disconnect(agent_id)
|
||||
logger.info(f"坐席断开 WebSocket 连接: agent_id={agent_id}")
|
||||
|
||||
except Exception as e:
|
||||
# 其他异常(如网络错误、JSON 解析错误等)
|
||||
# 确保注册信息被清理
|
||||
ws_manager.disconnect(agent_id)
|
||||
logger.warning(f"WebSocket 异常断开: agent_id={agent_id}, error={e}")
|
||||
|
||||
|
||||
# ==========================================================================
|
||||
# H5员工 WebSocket 端点
|
||||
# ==========================================================================
|
||||
|
||||
@router.websocket("/ws/h5/{employee_id}")
|
||||
async def h5_websocket_endpoint(
|
||||
websocket: WebSocket,
|
||||
employee_id: str,
|
||||
) -> None:
|
||||
"""H5员工 WebSocket 端点主循环(含 token 认证)。
|
||||
|
||||
做什么:
|
||||
1. 从 Authorization header 获取 token(优先从)或 query param(兼容)
|
||||
2. 验证 employee token 有效性(查 Redis)
|
||||
3. 验证 token 与 employee_id 一致性(防冒充)
|
||||
4. 认证通过后接受连接,注册到 ConnectionManager 的员工连接表
|
||||
5. 进入消息接收循环,处理心跳 ping
|
||||
6. 连接断开时清理注册信息
|
||||
|
||||
为什么需要 H5 WS 连接:
|
||||
- H5员工需要实时接收参与者变更事件(新参与者加入、有人退出等)
|
||||
- 当前仅通过 3 秒轮询获取更新,实时性不足
|
||||
- WS 推送 + 轮询降级,双通道保证消息可达
|
||||
|
||||
安全改进(P0-#4):
|
||||
- 优先从 Authorization: Bearer {token} header 获取 token
|
||||
- 兼容从 ?token= URL 参数获取(向后兼容)
|
||||
|
||||
认证机制(与坐席端一致):
|
||||
- Redis 中存储格式: employee:token:{token} -> employee_id
|
||||
- (与H5登录 API /api/h5/mock-login 存储格式一致)
|
||||
- token 缺失、无效、过期、与 employee_id 不匹配均拒绝连接
|
||||
|
||||
v0.5.1 修复:移除 `request: Request` 参数(部分 Starlette 版本注入 Request 失败,
|
||||
改用 `websocket.headers` 和 `websocket.query_params` 读取 header/query)
|
||||
|
||||
Args:
|
||||
websocket: FastAPI WebSocket 对象(框架自动注入)
|
||||
employee_id: 员工企微 UserID(从 URL 路径参数获取)
|
||||
"""
|
||||
# ======================================================================
|
||||
# Token 认证(从 subprotocol / header / query 获取)
|
||||
# ======================================================================
|
||||
|
||||
# 步骤1: 优先从 Sec-WebSocket-Protocol (subprotocol) 获取 token,其次从 Authorization header,最后从 query(向后兼容)
|
||||
# 格式: Sec-WebSocket-Protocol: bearer.{token}
|
||||
subprotocol = websocket.headers.get("sec-websocket-protocol", "")
|
||||
if subprotocol.startswith("bearer."):
|
||||
token = subprotocol[7:] # 去掉 "bearer." 前缀
|
||||
else:
|
||||
# 其次从 Authorization header 获取
|
||||
auth_header = websocket.headers.get("Authorization", "")
|
||||
if auth_header.startswith("Bearer "):
|
||||
token = auth_header[7:] # 去掉 "Bearer " 前缀
|
||||
else:
|
||||
# 向后兼容:从 query param 获取(即将废弃)
|
||||
token = websocket.query_params.get("token", "")
|
||||
|
||||
# 步骤2: 检查 token 是否为空
|
||||
if not token:
|
||||
await websocket.accept()
|
||||
await websocket.close(code=WS_CLOSE_UNAUTHORIZED, reason="Missing token")
|
||||
logger.warning(f"H5 WebSocket 拒绝连接: employee_id={employee_id}, 原因=缺少token")
|
||||
return
|
||||
|
||||
# 步骤3: 从 Redis 查询 token 对应的员工信息
|
||||
# Redis 中存储格式: employee:token:{token} -> employee_id
|
||||
# (与H5登录 API /api/h5/mock-login 存储格式一致)
|
||||
try:
|
||||
stored_employee_id = await cache_service.get(f"employee:token:{token}")
|
||||
except Exception as e:
|
||||
# Redis 不可用时必须拒绝连接(与坐席端一致的安全策略)
|
||||
logger.error(f"Redis 查询失败,拒绝 H5 WS 连接: employee_id={employee_id}, error={e}")
|
||||
await websocket.accept()
|
||||
await websocket.close(
|
||||
code=WS_CLOSE_UNAUTHORIZED,
|
||||
reason="Authentication service unavailable"
|
||||
)
|
||||
return
|
||||
|
||||
# 步骤4: 验证 token 与 employee_id 一致性
|
||||
if not stored_employee_id:
|
||||
await websocket.accept()
|
||||
await websocket.close(code=WS_CLOSE_UNAUTHORIZED, reason="Invalid or expired token")
|
||||
logger.warning(f"H5 WebSocket 拒绝连接: employee_id={employee_id}, 原因=token无效或已过期")
|
||||
return
|
||||
|
||||
if stored_employee_id != employee_id:
|
||||
await websocket.accept()
|
||||
await websocket.close(code=WS_CLOSE_UNAUTHORIZED, reason="Token-employee mismatch")
|
||||
logger.warning(
|
||||
f"H5 WebSocket 拒绝连接: employee_id={employee_id}, "
|
||||
f"原因=token对应员工{stored_employee_id}与请求不匹配"
|
||||
)
|
||||
return
|
||||
|
||||
# ======================================================================
|
||||
# 认证通过,建立连接
|
||||
# ======================================================================
|
||||
|
||||
# 注册员工连接(内部会调用 websocket.accept(),并回显协商的 subprotocol)
|
||||
await ws_manager.connect_employee(employee_id, websocket, subprotocol=subprotocol)
|
||||
logger.info(f"H5员工 WebSocket 连接已认证: employee_id={employee_id}")
|
||||
|
||||
try:
|
||||
# 消息接收循环
|
||||
# H5员工端目前只发送心跳 ping,不需要发送 typing 等事件
|
||||
while True:
|
||||
data = await websocket.receive_json()
|
||||
|
||||
# 处理心跳 ping
|
||||
if data.get("type") == "ping":
|
||||
await websocket.send_json({"type": "pong"})
|
||||
logger.debug(f"H5 WebSocket 心跳: employee_id={employee_id}")
|
||||
|
||||
else:
|
||||
logger.debug(
|
||||
f"H5 WebSocket 收到未知消息: employee_id={employee_id}, "
|
||||
f"type={data.get('type', 'unknown')}"
|
||||
)
|
||||
|
||||
except WebSocketDisconnect:
|
||||
# 客户端主动断开连接
|
||||
ws_manager.disconnect_employee(employee_id)
|
||||
logger.info(f"H5员工断开 WebSocket 连接: employee_id={employee_id}")
|
||||
|
||||
except Exception as e:
|
||||
# 其他异常
|
||||
ws_manager.disconnect_employee(employee_id)
|
||||
logger.warning(f"H5 WebSocket 异常断开: employee_id={employee_id}, error={e}")
|
||||
@@ -1,71 +0,0 @@
|
||||
# =============================================================================
|
||||
# 企微IT智能服务台 — RBAC 角色种子数据 (v0.7.1 task #86)
|
||||
# =============================================================================
|
||||
# 启动时调用,把 5 角色 + 权限矩阵写入 roles 表
|
||||
# 兼容"角色已存在"的场景: 不重复插入,但更新 permissions
|
||||
# =============================================================================
|
||||
|
||||
import logging
|
||||
import uuid
|
||||
from datetime import datetime
|
||||
|
||||
from sqlalchemy import select
|
||||
from sqlalchemy.ext.asyncio import AsyncSession
|
||||
|
||||
from app.models.role import Role
|
||||
from app.services.rbac_service import (
|
||||
ROLE_METADATA,
|
||||
get_role_default_permissions,
|
||||
)
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
|
||||
async def seed_rbac_roles(db: AsyncSession) -> int:
|
||||
"""种子 RBAC 5 角色。
|
||||
|
||||
行为:
|
||||
1. 遍历 ROLE_METADATA
|
||||
2. 角色不存在 → 创建(UUID + 默认 permissions)
|
||||
3. 角色存在 → 更新 display_name / description / permissions
|
||||
(不动 is_default,避免影响手动设置)
|
||||
|
||||
Returns:
|
||||
int: 新建角色数
|
||||
"""
|
||||
created_count = 0
|
||||
|
||||
for role_name, meta in ROLE_METADATA.items():
|
||||
# 查询是否已存在
|
||||
stmt = select(Role).where(Role.name == role_name)
|
||||
result = await db.execute(stmt)
|
||||
role = result.scalars().first()
|
||||
|
||||
permissions = get_role_default_permissions(role_name)
|
||||
|
||||
if role:
|
||||
# 更新现有角色(不动 is_default,防止覆盖手动设置)
|
||||
role.display_name = meta["display_name"]
|
||||
role.description = meta["description"]
|
||||
role.permissions = permissions
|
||||
role.updated_at = datetime.now()
|
||||
logger.debug(f"更新角色: {role_name} ({len(permissions)} 项权限)")
|
||||
else:
|
||||
# 创建新角色
|
||||
role = Role(
|
||||
id=str(uuid.uuid4()),
|
||||
name=role_name,
|
||||
display_name=meta["display_name"],
|
||||
description=meta["description"],
|
||||
permissions=permissions,
|
||||
is_default=(meta["is_default"] == "true"),
|
||||
created_at=datetime.now(),
|
||||
updated_at=datetime.now(),
|
||||
)
|
||||
db.add(role)
|
||||
created_count += 1
|
||||
logger.info(f"创建角色: {role_name} ({len(permissions)} 项权限)")
|
||||
|
||||
await db.commit()
|
||||
logger.info(f"RBAC 角色种子完成: 新建 {created_count} 个")
|
||||
return created_count
|
||||
@@ -1,98 +0,0 @@
|
||||
# 联软LV7000配置管理
|
||||
"""
|
||||
从system_configs表读取联软API配置,构建LianruanClient实例。
|
||||
|
||||
联软配置键(前缀 integration_lianruan_):
|
||||
- integration_lianruan_base_url: 联软API地址(如 http://192.168.x.x:30098)
|
||||
- integration_lianruan_api_account: API账号
|
||||
- integration_lianruan_api_password: API密码
|
||||
- integration_lianruan_validate_key: 验证密钥(可选)
|
||||
|
||||
配置方式:管理后台 → 系统集成 → 联软LV7000 → 填入账号密码
|
||||
"""
|
||||
|
||||
import logging
|
||||
|
||||
from sqlalchemy.ext.asyncio import AsyncSession
|
||||
|
||||
from app.integrations.lianruan.client import LianruanClient
|
||||
from app.integrations.lianruan.exceptions import LianruanConfigError
|
||||
from app.models.system_config import SystemConfig
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
# 联软配置键前缀(与 admin_service INTEGRATION_DEFINITIONS 中的 key_prefix 一致)
|
||||
_PREFIX = "integration_lianruan_"
|
||||
|
||||
|
||||
async def _get_lianruan_config_value(db: AsyncSession, key_suffix: str) -> str:
|
||||
"""读取单个联软配置值。
|
||||
|
||||
Args:
|
||||
db: 数据库会话
|
||||
key_suffix: 配置键后缀(如 base_url / api_account)
|
||||
|
||||
Returns:
|
||||
str: 配置值,不存在返回空字符串
|
||||
"""
|
||||
full_key = f"{_PREFIX}{key_suffix}"
|
||||
from sqlalchemy import select
|
||||
result = await db.execute(select(SystemConfig).where(SystemConfig.key == full_key))
|
||||
config_row = result.scalar_one_or_none()
|
||||
return config_row.value if config_row else ""
|
||||
|
||||
|
||||
async def get_lianruan_config(db: AsyncSession) -> dict:
|
||||
"""从system_configs表读取联软配置。
|
||||
|
||||
Args:
|
||||
db: 数据库会话
|
||||
|
||||
Returns:
|
||||
dict: 包含 base_url / api_account / api_password / validate_key
|
||||
|
||||
Raises:
|
||||
LianruanConfigError: 配置缺失
|
||||
"""
|
||||
base_url = await _get_lianruan_config_value(db, "base_url")
|
||||
api_account = await _get_lianruan_config_value(db, "api_account")
|
||||
api_password = await _get_lianruan_config_value(db, "api_password")
|
||||
validate_key = await _get_lianruan_config_value(db, "validate_key")
|
||||
|
||||
if not base_url:
|
||||
raise LianruanConfigError("联软API未配置:缺少Base URL")
|
||||
if not api_account:
|
||||
raise LianruanConfigError("联软API未配置:缺少API账号")
|
||||
if not api_password:
|
||||
raise LianruanConfigError("联软API未配置:缺少API密码")
|
||||
|
||||
return {
|
||||
"base_url": base_url,
|
||||
"api_account": api_account,
|
||||
"api_password": api_password,
|
||||
"validate_key": validate_key,
|
||||
}
|
||||
|
||||
|
||||
async def get_lianruan_client(db: AsyncSession) -> LianruanClient:
|
||||
"""构建联软API客户端实例。
|
||||
|
||||
从system_configs表读取配置,创建LianruanClient。
|
||||
|
||||
Args:
|
||||
db: 数据库会话
|
||||
|
||||
Returns:
|
||||
LianruanClient: 已配置的联软客户端
|
||||
|
||||
Raises:
|
||||
LianruanConfigError: 配置缺失
|
||||
"""
|
||||
cfg = await get_lianruan_config(db)
|
||||
|
||||
return LianruanClient(
|
||||
base_url=cfg["base_url"],
|
||||
api_account=cfg["api_account"],
|
||||
api_password=cfg["api_password"],
|
||||
validate_key=cfg.get("validate_key", ""),
|
||||
)
|
||||
@@ -1,35 +0,0 @@
|
||||
# =============================================================================
|
||||
# RAGFlow 集成模块
|
||||
# =============================================================================
|
||||
|
||||
from .client import RagflowClient
|
||||
from .config import get_ragflow_client
|
||||
from .exceptions import (
|
||||
RagflowApiError,
|
||||
RagflowAuthError,
|
||||
RagflowConfigError,
|
||||
RagflowConnectionError,
|
||||
RagflowError,
|
||||
)
|
||||
from .models import (
|
||||
DatasetInfo,
|
||||
DocAggregate,
|
||||
DocumentInfo,
|
||||
RetrievalChunk,
|
||||
RetrievalResult,
|
||||
)
|
||||
|
||||
__all__ = [
|
||||
"RagflowClient",
|
||||
"get_ragflow_client",
|
||||
"RagflowError",
|
||||
"RagflowConfigError",
|
||||
"RagflowAuthError",
|
||||
"RagflowApiError",
|
||||
"RagflowConnectionError",
|
||||
"RetrievalChunk",
|
||||
"DocAggregate",
|
||||
"RetrievalResult",
|
||||
"DatasetInfo",
|
||||
"DocumentInfo",
|
||||
]
|
||||
@@ -1,449 +0,0 @@
|
||||
# =============================================================================
|
||||
# RAGFlow API 客户端
|
||||
# =============================================================================
|
||||
# 说明:封装 RAGFlow 知识检索引擎的 API 调用
|
||||
# 核心功能:
|
||||
# 1. 知识检索 — POST /api/v1/retrieval(核心接口)
|
||||
# 2. 数据集管理 — 列出/创建/删除知识库
|
||||
# 3. 文档管理 — 上传/列出/删除文档
|
||||
# 4. 测试连接 — 验证 API Key 是否有效
|
||||
# 认证方式:Authorization: Bearer <API_KEY>
|
||||
# 参考文档:https://ragflow.io/docs/http_api_reference
|
||||
# =============================================================================
|
||||
|
||||
import logging
|
||||
from typing import Any, Dict, List, Optional
|
||||
|
||||
import httpx
|
||||
|
||||
from .exceptions import (
|
||||
RagflowApiError,
|
||||
RagflowAuthError,
|
||||
RagflowConfigError,
|
||||
RagflowConnectionError,
|
||||
RagflowError,
|
||||
)
|
||||
from .models import (
|
||||
DatasetInfo,
|
||||
DocAggregate,
|
||||
DocumentInfo,
|
||||
RetrievalChunk,
|
||||
RetrievalResult,
|
||||
)
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
# 默认请求超时(秒)
|
||||
DEFAULT_TIMEOUT = 30.0
|
||||
|
||||
# 默认分页大小
|
||||
DEFAULT_PAGE_SIZE = 20
|
||||
|
||||
|
||||
class RagflowClient:
|
||||
"""RAGFlow API 客户端。
|
||||
|
||||
封装 RAGFlow 知识检索引擎的 API 调用,支持:
|
||||
- 知识检索(核心功能)
|
||||
- 数据集(知识库)管理
|
||||
- 文档管理
|
||||
- 连接测试
|
||||
|
||||
使用方式:
|
||||
client = RagflowClient(
|
||||
api_key="sk-xxx",
|
||||
base_url="http://10.80.0.85:9380"
|
||||
)
|
||||
result = await client.retrieval("VPN怎么连?", dataset_ids=["xxx"])
|
||||
"""
|
||||
|
||||
def __init__(
|
||||
self,
|
||||
api_key: str,
|
||||
base_url: str = "http://10.80.0.85:9380",
|
||||
timeout: float = DEFAULT_TIMEOUT,
|
||||
):
|
||||
"""初始化 RAGFlow 客户端。
|
||||
|
||||
Args:
|
||||
api_key: RAGFlow API Key(Bearer Token)
|
||||
base_url: RAGFlow API 基础地址(不含尾部斜杠)
|
||||
timeout: 默认请求超时(秒)
|
||||
|
||||
Raises:
|
||||
RagflowConfigError: API Key 为空
|
||||
"""
|
||||
if not api_key:
|
||||
raise RagflowConfigError("RAGFlow API Key 不能为空")
|
||||
|
||||
self.api_key = api_key
|
||||
self.base_url = base_url.rstrip("/")
|
||||
self.timeout = timeout
|
||||
|
||||
def _headers(self) -> Dict[str, str]:
|
||||
"""构建请求头。
|
||||
|
||||
Returns:
|
||||
Dict: 包含 Authorization 和 Content-Type 的请求头
|
||||
"""
|
||||
return {
|
||||
"Authorization": f"Bearer {self.api_key}",
|
||||
"Content-Type": "application/json",
|
||||
}
|
||||
|
||||
async def _request(
|
||||
self,
|
||||
method: str,
|
||||
path: str,
|
||||
json_data: Optional[Dict] = None,
|
||||
params: Optional[Dict] = None,
|
||||
timeout: Optional[float] = None,
|
||||
) -> Dict[str, Any]:
|
||||
"""统一请求封装。
|
||||
|
||||
Args:
|
||||
method: HTTP 方法(GET/POST/PUT/DELETE)
|
||||
path: API 路径(如 /api/v1/retrieval)
|
||||
json_data: JSON 请求体
|
||||
params: 查询参数
|
||||
timeout: 覆盖默认超时
|
||||
|
||||
Returns:
|
||||
Dict: API 响应的 JSON 数据
|
||||
|
||||
Raises:
|
||||
RagflowAuthError: 认证失败(401)
|
||||
RagflowApiError: API 返回错误
|
||||
RagflowConnectionError: 网络连接失败
|
||||
"""
|
||||
url = f"{self.base_url}{path}"
|
||||
req_timeout = timeout or self.timeout
|
||||
|
||||
try:
|
||||
async with httpx.AsyncClient() as client:
|
||||
response = await client.request(
|
||||
method=method,
|
||||
url=url,
|
||||
headers=self._headers(),
|
||||
json=json_data,
|
||||
params=params,
|
||||
timeout=req_timeout,
|
||||
)
|
||||
|
||||
# 处理 HTTP 错误
|
||||
if response.status_code == 401:
|
||||
raise RagflowAuthError("RAGFlow API Key 无效或已过期")
|
||||
|
||||
if response.status_code >= 400:
|
||||
try:
|
||||
err_body = response.json()
|
||||
err_msg = err_body.get("message", response.text)
|
||||
except Exception:
|
||||
err_msg = response.text
|
||||
raise RagflowApiError(
|
||||
code=response.status_code,
|
||||
message=f"RAGFlow API 错误 ({response.status_code}): {err_msg}",
|
||||
)
|
||||
|
||||
# 解析响应
|
||||
result = response.json()
|
||||
|
||||
# RAGFlow 统一响应格式:{code: 0, data: ..., message: ...}
|
||||
if result.get("code") != 0:
|
||||
raise RagflowApiError(
|
||||
code=result.get("code", -1),
|
||||
message=result.get("message", "未知错误"),
|
||||
)
|
||||
|
||||
return result
|
||||
|
||||
except httpx.TimeoutException:
|
||||
raise RagflowConnectionError(f"RAGFlow 请求超时 ({req_timeout}s): {path}")
|
||||
except httpx.ConnectError:
|
||||
raise RagflowConnectionError(f"RAGFlow 连接失败: {self.base_url}")
|
||||
except (RagflowAuthError, RagflowApiError, RagflowConnectionError):
|
||||
raise
|
||||
except Exception as e:
|
||||
raise RagflowError(f"RAGFlow 请求异常: {str(e)}")
|
||||
|
||||
# ==========================================================================
|
||||
# 测试连接
|
||||
# ==========================================================================
|
||||
|
||||
async def test_connection(self) -> Dict[str, Any]:
|
||||
"""测试 RAGFlow API 连接。
|
||||
|
||||
通过列出数据集(limit=1)验证 API Key 是否有效。
|
||||
|
||||
Returns:
|
||||
Dict: {success: bool, message: str}
|
||||
"""
|
||||
try:
|
||||
result = await self.list_datasets(page=1, page_size=1)
|
||||
return {
|
||||
"success": True,
|
||||
"message": f"连接成功,共 {result.get('total', 0)} 个知识库",
|
||||
}
|
||||
except RagflowAuthError:
|
||||
return {"success": False, "message": "API Key 无效或已过期"}
|
||||
except RagflowConnectionError as e:
|
||||
return {"success": False, "message": f"连接失败: {e.message}"}
|
||||
except RagflowError as e:
|
||||
return {"success": False, "message": e.message}
|
||||
|
||||
# ==========================================================================
|
||||
# 知识检索(核心接口)
|
||||
# ==========================================================================
|
||||
|
||||
async def retrieval(
|
||||
self,
|
||||
question: str,
|
||||
dataset_ids: Optional[List[str]] = None,
|
||||
document_ids: Optional[List[str]] = None,
|
||||
similarity_threshold: float = 0.2,
|
||||
vector_similarity_weight: float = 0.3,
|
||||
top_k: int = 1024,
|
||||
keyword: bool = False,
|
||||
highlight: bool = False,
|
||||
) -> RetrievalResult:
|
||||
"""知识检索 — 从知识库中搜索相关文档片段。
|
||||
|
||||
这是 RAGFlow 的核心接口,用于根据用户问题检索最相关的文本块。
|
||||
|
||||
Args:
|
||||
question: 用户查询问题
|
||||
dataset_ids: 要搜索的数据集ID列表(与 document_ids 二选一)
|
||||
document_ids: 要搜索的文档ID列表
|
||||
similarity_threshold: 最小相似度阈值(0-1),默认 0.2
|
||||
vector_similarity_weight: 向量相似度权重(0-1),默认 0.3
|
||||
top_k: 参与计算的块数量,默认 1024
|
||||
keyword: 是否启用关键词匹配,默认 False
|
||||
highlight: 是否高亮匹配术语,默认 False
|
||||
|
||||
Returns:
|
||||
RetrievalResult: 检索结果(含文本块、文档聚合、总数)
|
||||
|
||||
Raises:
|
||||
RagflowError: 检索失败
|
||||
"""
|
||||
body: Dict[str, Any] = {
|
||||
"question": question,
|
||||
"similarity_threshold": similarity_threshold,
|
||||
"vector_similarity_weight": vector_similarity_weight,
|
||||
"top_k": top_k,
|
||||
"keyword": keyword,
|
||||
"highlight": highlight,
|
||||
}
|
||||
|
||||
if dataset_ids:
|
||||
body["dataset_ids"] = dataset_ids
|
||||
if document_ids:
|
||||
body["document_ids"] = document_ids
|
||||
|
||||
result = await self._request("POST", "/api/v1/retrieval", json_data=body)
|
||||
|
||||
data = result.get("data", {})
|
||||
|
||||
# 解析文本块
|
||||
chunks = [
|
||||
RetrievalChunk.model_validate(chunk)
|
||||
for chunk in data.get("chunks", [])
|
||||
]
|
||||
|
||||
# 解析文档聚合
|
||||
doc_aggs = [
|
||||
DocAggregate.model_validate(agg)
|
||||
for agg in data.get("doc_aggs", [])
|
||||
]
|
||||
|
||||
return RetrievalResult(
|
||||
chunks=chunks,
|
||||
doc_aggs=doc_aggs,
|
||||
total=data.get("total", 0),
|
||||
)
|
||||
|
||||
# ==========================================================================
|
||||
# 数据集(知识库)管理
|
||||
# ==========================================================================
|
||||
|
||||
async def list_datasets(
|
||||
self,
|
||||
page: int = 1,
|
||||
page_size: int = DEFAULT_PAGE_SIZE,
|
||||
) -> Dict[str, Any]:
|
||||
"""列出所有数据集(知识库)。
|
||||
|
||||
Args:
|
||||
page: 页码
|
||||
page_size: 每页条数
|
||||
|
||||
Returns:
|
||||
Dict: {items: List[DatasetInfo], total: int}
|
||||
"""
|
||||
result = await self._request(
|
||||
"GET",
|
||||
"/api/v1/datasets",
|
||||
params={"page": page, "page_size": page_size},
|
||||
)
|
||||
|
||||
data = result.get("data", {})
|
||||
items = [
|
||||
DatasetInfo.model_validate(ds)
|
||||
for ds in data.get("datasets", [])
|
||||
]
|
||||
|
||||
return {"items": items, "total": data.get("total", 0)}
|
||||
|
||||
async def create_dataset(
|
||||
self,
|
||||
name: str,
|
||||
embedding_model: str = "BAAI/bge-m3@BAAI",
|
||||
chunk_method: str = "naive",
|
||||
permission: str = "me",
|
||||
) -> DatasetInfo:
|
||||
"""创建数据集(知识库)。
|
||||
|
||||
Args:
|
||||
name: 数据集名称
|
||||
embedding_model: 向量模型
|
||||
chunk_method: 分块方法(naive/qa/book/laws 等)
|
||||
permission: 权限(me/team)
|
||||
|
||||
Returns:
|
||||
DatasetInfo: 创建的数据集信息
|
||||
"""
|
||||
body = {
|
||||
"name": name,
|
||||
"embedding_model": embedding_model,
|
||||
"chunk_method": chunk_method,
|
||||
"permission": permission,
|
||||
}
|
||||
|
||||
result = await self._request("POST", "/api/v1/datasets", json_data=body)
|
||||
return DatasetInfo.model_validate(result.get("data", {}))
|
||||
|
||||
async def delete_dataset(self, dataset_ids: List[str]) -> bool:
|
||||
"""删除数据集。
|
||||
|
||||
Args:
|
||||
dataset_ids: 要删除的数据集ID列表
|
||||
|
||||
Returns:
|
||||
bool: 是否成功
|
||||
"""
|
||||
await self._request(
|
||||
"DELETE",
|
||||
"/api/v1/datasets",
|
||||
json_data={"ids": dataset_ids},
|
||||
)
|
||||
return True
|
||||
|
||||
# ==========================================================================
|
||||
# 文档管理
|
||||
# ==========================================================================
|
||||
|
||||
async def list_documents(
|
||||
self,
|
||||
dataset_id: str,
|
||||
page: int = 1,
|
||||
page_size: int = DEFAULT_PAGE_SIZE,
|
||||
) -> Dict[str, Any]:
|
||||
"""列出数据集中的文档。
|
||||
|
||||
Args:
|
||||
dataset_id: 数据集ID
|
||||
page: 页码
|
||||
page_size: 每页条数
|
||||
|
||||
Returns:
|
||||
Dict: {items: List[DocumentInfo], total: int}
|
||||
"""
|
||||
result = await self._request(
|
||||
"GET",
|
||||
f"/api/v1/datasets/{dataset_id}/documents",
|
||||
params={"page": page, "page_size": page_size},
|
||||
)
|
||||
|
||||
data = result.get("data", {})
|
||||
items = [
|
||||
DocumentInfo.model_validate(doc)
|
||||
for doc in data.get("documents", [])
|
||||
]
|
||||
|
||||
return {"items": items, "total": data.get("total", 0)}
|
||||
|
||||
async def upload_document(
|
||||
self,
|
||||
dataset_id: str,
|
||||
file_path: str,
|
||||
file_name: Optional[str] = None,
|
||||
) -> DocumentInfo:
|
||||
"""上传文档到数据集。
|
||||
|
||||
Args:
|
||||
dataset_id: 数据集ID
|
||||
file_path: 本地文件路径
|
||||
file_name: 文件名(可选,默认取 file_path 的文件名)
|
||||
|
||||
Returns:
|
||||
DocumentInfo: 上传的文档信息
|
||||
"""
|
||||
import os
|
||||
|
||||
if not os.path.exists(file_path):
|
||||
raise RagflowError(f"文件不存在: {file_path}")
|
||||
|
||||
fname = file_name or os.path.basename(file_path)
|
||||
|
||||
url = f"{self.base_url}/api/v1/datasets/{dataset_id}/documents"
|
||||
|
||||
try:
|
||||
async with httpx.AsyncClient() as client:
|
||||
with open(file_path, "rb") as f:
|
||||
response = await client.post(
|
||||
url=url,
|
||||
headers={"Authorization": f"Bearer {self.api_key}"},
|
||||
files={"file": (fname, f)},
|
||||
timeout=60.0,
|
||||
)
|
||||
|
||||
if response.status_code == 401:
|
||||
raise RagflowAuthError()
|
||||
|
||||
result = response.json()
|
||||
if result.get("code") != 0:
|
||||
raise RagflowApiError(
|
||||
code=result.get("code", -1),
|
||||
message=result.get("message", "上传失败"),
|
||||
)
|
||||
|
||||
docs = result.get("data", {}).get("documents", [])
|
||||
if docs:
|
||||
return DocumentInfo.model_validate(docs[0])
|
||||
return DocumentInfo(name=fname)
|
||||
|
||||
except (RagflowAuthError, RagflowApiError):
|
||||
raise
|
||||
except Exception as e:
|
||||
raise RagflowError(f"文档上传失败: {str(e)}")
|
||||
|
||||
async def delete_documents(
|
||||
self,
|
||||
dataset_id: str,
|
||||
document_ids: List[str],
|
||||
) -> bool:
|
||||
"""删除文档。
|
||||
|
||||
Args:
|
||||
dataset_id: 数据集ID
|
||||
document_ids: 要删除的文档ID列表
|
||||
|
||||
Returns:
|
||||
bool: 是否成功
|
||||
"""
|
||||
await self._request(
|
||||
"DELETE",
|
||||
f"/api/v1/datasets/{dataset_id}/documents",
|
||||
json_data={"ids": document_ids},
|
||||
)
|
||||
return True
|
||||
@@ -1,61 +0,0 @@
|
||||
# =============================================================================
|
||||
# RAGFlow 配置加载器
|
||||
# =============================================================================
|
||||
# 说明:从数据库 system_configs 表加载 RAGFlow 配置,创建客户端实例
|
||||
# 配置项:integration_ragflow_api_url + integration_ragflow_api_key
|
||||
|
||||
import logging
|
||||
from typing import Optional
|
||||
|
||||
from sqlalchemy import select
|
||||
from sqlalchemy.ext.asyncio import AsyncSession
|
||||
|
||||
from app.models.system_config import SystemConfig
|
||||
|
||||
from .client import RagflowClient
|
||||
from .exceptions import RagflowConfigError
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
# 默认 RAGFlow API 地址(生产环境)
|
||||
DEFAULT_RAGFLOW_BASE_URL = "http://10.80.0.85:9380"
|
||||
|
||||
|
||||
async def _get_config(db: AsyncSession, key: str) -> str:
|
||||
"""从数据库读取单个配置值。"""
|
||||
result = await db.execute(
|
||||
select(SystemConfig.config_value).where(SystemConfig.config_key == key)
|
||||
)
|
||||
row = result.scalar()
|
||||
return row if row else ""
|
||||
|
||||
|
||||
async def get_ragflow_client(db: AsyncSession) -> RagflowClient:
|
||||
"""从数据库配置创建 RAGFlow 客户端实例。
|
||||
|
||||
读取 system_configs 表中的:
|
||||
- integration_ragflow_api_url: RAGFlow API 地址
|
||||
- integration_ragflow_api_key: RAGFlow API Key
|
||||
|
||||
Args:
|
||||
db: 数据库会话
|
||||
|
||||
Returns:
|
||||
RagflowClient: 客户端实例
|
||||
|
||||
Raises:
|
||||
RagflowConfigError: 配置缺失
|
||||
"""
|
||||
api_url = await _get_config(db, "integration_ragflow_api_url")
|
||||
api_key = await _get_config(db, "integration_ragflow_api_key")
|
||||
|
||||
# 如果数据库没有配置,使用默认地址
|
||||
if not api_url:
|
||||
api_url = DEFAULT_RAGFLOW_BASE_URL
|
||||
|
||||
if not api_key:
|
||||
raise RagflowConfigError(
|
||||
"RAGFlow API Key 未配置,请在管理后台 → 集成管理 → RAGFlow 中设置"
|
||||
)
|
||||
|
||||
return RagflowClient(api_key=api_key, base_url=api_url)
|
||||
@@ -1,35 +0,0 @@
|
||||
# =============================================================================
|
||||
# RAGFlow API 异常定义
|
||||
# =============================================================================
|
||||
|
||||
|
||||
class RagflowError(Exception):
|
||||
"""RAGFlow 基础异常。"""
|
||||
def __init__(self, message: str = "RAGFlow 错误"):
|
||||
self.message = message
|
||||
super().__init__(self.message)
|
||||
|
||||
|
||||
class RagflowConfigError(RagflowError):
|
||||
"""配置错误(缺少 API Key 或 Base URL)。"""
|
||||
def __init__(self, message: str = "RAGFlow 配置缺失"):
|
||||
super().__init__(message)
|
||||
|
||||
|
||||
class RagflowAuthError(RagflowError):
|
||||
"""认证失败(API Key 无效)。"""
|
||||
def __init__(self, message: str = "RAGFlow 认证失败"):
|
||||
super().__init__(message)
|
||||
|
||||
|
||||
class RagflowApiError(RagflowError):
|
||||
"""API 调用失败(非 200 响应)。"""
|
||||
def __init__(self, code: int = 0, message: str = "RAGFlow API 错误"):
|
||||
self.code = code
|
||||
super().__init__(message)
|
||||
|
||||
|
||||
class RagflowConnectionError(RagflowError):
|
||||
"""网络连接失败。"""
|
||||
def __init__(self, message: str = "RAGFlow 连接失败"):
|
||||
super().__init__(message)
|
||||
@@ -1,110 +0,0 @@
|
||||
# =============================================================================
|
||||
# RAGFlow API 数据模型
|
||||
# =============================================================================
|
||||
# 说明:定义 RAGFlow API 请求/响应的 Pydantic 数据模型
|
||||
# 参考:https://ragflow.io/docs/http_api_reference
|
||||
|
||||
from typing import Any, Dict, List, Optional
|
||||
|
||||
from pydantic import BaseModel, Field
|
||||
|
||||
|
||||
class RetrievalChunk(BaseModel):
|
||||
"""检索返回的单个文本块。
|
||||
|
||||
Attributes:
|
||||
id: 块唯一ID
|
||||
content: 块内容文本
|
||||
document_id: 所属文档ID
|
||||
document_keyword: 所属文档名称
|
||||
similarity: 综合相似度分数
|
||||
term_similarity: 关键词相似度
|
||||
vector_similarity: 向量相似度
|
||||
highlight: 高亮标记的内容(可选)
|
||||
"""
|
||||
id: str = Field(default="", description="块唯一ID")
|
||||
content: str = Field(default="", description="块内容文本")
|
||||
document_id: str = Field(default="", description="所属文档ID")
|
||||
document_keyword: str = Field(default="", description="所属文档名称")
|
||||
similarity: float = Field(default=0.0, description="综合相似度分数")
|
||||
term_similarity: float = Field(default=0.0, description="关键词相似度")
|
||||
vector_similarity: float = Field(default=0.0, description="向量相似度")
|
||||
highlight: Optional[str] = Field(default=None, description="高亮标记的内容")
|
||||
|
||||
model_config = {"from_attributes": True}
|
||||
|
||||
|
||||
class DocAggregate(BaseModel):
|
||||
"""文档聚合统计。
|
||||
|
||||
Attributes:
|
||||
doc_id: 文档ID
|
||||
doc_name: 文档名称
|
||||
count: 命中的块数量
|
||||
"""
|
||||
doc_id: str = Field(default="", description="文档ID")
|
||||
doc_name: str = Field(default="", description="文档名称")
|
||||
count: int = Field(default=0, description="命中块数量")
|
||||
|
||||
model_config = {"from_attributes": True}
|
||||
|
||||
|
||||
class RetrievalResult(BaseModel):
|
||||
"""检索结果。
|
||||
|
||||
Attributes:
|
||||
chunks: 命中的文本块列表
|
||||
doc_aggs: 按文档聚合统计
|
||||
total: 命中总数
|
||||
"""
|
||||
chunks: List[RetrievalChunk] = Field(default_factory=list, description="命中文本块列表")
|
||||
doc_aggs: List[DocAggregate] = Field(default_factory=list, description="文档聚合统计")
|
||||
total: int = Field(default=0, description="命中总数")
|
||||
|
||||
model_config = {"from_attributes": True}
|
||||
|
||||
|
||||
class DatasetInfo(BaseModel):
|
||||
"""数据集(知识库)信息。
|
||||
|
||||
Attributes:
|
||||
id: 数据集ID
|
||||
name: 数据集名称
|
||||
chunk_method: 分块方法
|
||||
permission: 权限
|
||||
document_count: 文档数量
|
||||
embedding_model: 向量模型
|
||||
create_time: 创建时间
|
||||
update_time: 更新时间
|
||||
"""
|
||||
id: str = Field(default="", description="数据集ID")
|
||||
name: str = Field(default="", description="数据集名称")
|
||||
chunk_method: str = Field(default="naive", description="分块方法")
|
||||
permission: str = Field(default="me", description="权限")
|
||||
document_count: int = Field(default=0, description="文档数量")
|
||||
embedding_model: str = Field(default="", description="向量模型")
|
||||
create_time: Optional[str] = Field(default=None, description="创建时间")
|
||||
update_time: Optional[str] = Field(default=None, description="更新时间")
|
||||
|
||||
model_config = {"from_attributes": True}
|
||||
|
||||
|
||||
class DocumentInfo(BaseModel):
|
||||
"""文档信息。
|
||||
|
||||
Attributes:
|
||||
id: 文档ID
|
||||
name: 文档名称
|
||||
chunk_method: 分块方法
|
||||
chunk_count: 块数量
|
||||
create_time: 创建时间
|
||||
update_time: 更新时间
|
||||
"""
|
||||
id: str = Field(default="", description="文档ID")
|
||||
name: str = Field(default="", description="文档名称")
|
||||
chunk_method: str = Field(default="naive", description="分块方法")
|
||||
chunk_count: int = Field(default=0, description="块数量")
|
||||
create_time: Optional[str] = Field(default=None, description="创建时间")
|
||||
update_time: Optional[str] = Field(default=None, description="更新时间")
|
||||
|
||||
model_config = {"from_attributes": True}
|
||||
@@ -1,70 +0,0 @@
|
||||
# =============================================================================
|
||||
# 企微IT智能服务台 — 模型包初始化
|
||||
# =============================================================================
|
||||
# 说明:导出所有模型类,方便 Alembic 和其他模块统一导入
|
||||
# 注意:即使某些模型在当前文件未直接使用,也必须导入
|
||||
# 否则 Alembic 无法检测到这些模型,不会生成对应的迁移脚本
|
||||
# =============================================================================
|
||||
|
||||
from app.models.conversation import Conversation
|
||||
from app.models.message import Message
|
||||
from app.models.agent import Agent
|
||||
from app.models.quick_reply_template import QuickReplyTemplate
|
||||
from app.models.system_config import SystemConfig
|
||||
from app.models.funny_phrase import FunnyPhrase
|
||||
from app.models.approval_link import ApprovalLink
|
||||
from app.models.software_download import SoftwareDownload
|
||||
from app.models.agent_note import AgentNote
|
||||
from app.models.employee import Employee
|
||||
from app.models.todo_item import TodoItem
|
||||
from app.models.troubleshooting_template import TroubleshootingTemplate
|
||||
from app.models.config_change_log import ConfigChangeLog
|
||||
from app.models.role import Role
|
||||
from app.models.user_role import UserRole
|
||||
from app.models.role_mapping_rule import RoleMappingRule
|
||||
from app.models.conversation_evaluation import ConversationEvaluation
|
||||
from app.models.audit_log import AuditLog
|
||||
from app.models.conversation_annotation import ConversationAnnotation # P2-10 会话标注
|
||||
from app.models.knowledge_suggestion import KnowledgeSuggestion # P2-13 知识库自动迭代
|
||||
from app.models.knowledge_base import KnowledgeBase
|
||||
# 阶段5 自动化闭环模型
|
||||
from app.models.automation import (
|
||||
AutoSession,
|
||||
AutoAction,
|
||||
ApprovalTicket,
|
||||
ScenarioConfig,
|
||||
RuleVersion,
|
||||
ActionLog,
|
||||
MappingCache,
|
||||
)
|
||||
# 所有模型类的列表,方便遍历
|
||||
__all__ = [
|
||||
"Conversation",
|
||||
"Message",
|
||||
"Agent",
|
||||
"QuickReplyTemplate",
|
||||
"SystemConfig",
|
||||
"FunnyPhrase",
|
||||
"ApprovalLink",
|
||||
"SoftwareDownload",
|
||||
"AgentNote",
|
||||
"Employee",
|
||||
"TodoItem",
|
||||
"TroubleshootingTemplate",
|
||||
"ConfigChangeLog",
|
||||
"Role",
|
||||
"UserRole",
|
||||
"RoleMappingRule",
|
||||
"ConversationEvaluation",
|
||||
"AuditLog",
|
||||
"ConversationAnnotation",
|
||||
"KnowledgeSuggestion",
|
||||
"KnowledgeBase",
|
||||
"AutoSession",
|
||||
"AutoAction",
|
||||
"ApprovalTicket",
|
||||
"ScenarioConfig",
|
||||
"RuleVersion",
|
||||
"ActionLog",
|
||||
"MappingCache",
|
||||
]
|
||||
@@ -1,38 +0,0 @@
|
||||
# =============================================================================
|
||||
# 企微IT智能服务台 — 服务包初始化
|
||||
# =============================================================================
|
||||
# 说明:将 services/ 目录标记为 Python 包
|
||||
# 导出所有服务类,方便统一导入
|
||||
# =============================================================================
|
||||
|
||||
from app.services.wecom_service import WecomService
|
||||
from app.services.message_router import MessageRouter
|
||||
from app.services.scoring_service import ScoringService
|
||||
from app.services.session_service import SessionService
|
||||
from app.services.funny_phrase_service import FunnyPhraseService
|
||||
from app.services.ai_handler import AIHandler
|
||||
|
||||
# Tier0 新增服务导出
|
||||
from app.services.neo4j_client import Neo4jClient, dep_neo4j_client, get_neo4j_client
|
||||
from app.services.knowledge_iteration_service import KnowledgeIterationService, dep_knowledge_iteration_service
|
||||
from app.services.wingman_service import WingmanService
|
||||
from app.services.vision_service import VisionService
|
||||
from app.services.ragflow_ingestion_service import RagflowIngestionService
|
||||
|
||||
__all__ = [
|
||||
"WecomService",
|
||||
"MessageRouter",
|
||||
"ScoringService",
|
||||
"SessionService",
|
||||
"FunnyPhraseService",
|
||||
"AIHandler",
|
||||
# Tier0
|
||||
"Neo4jClient",
|
||||
"dep_neo4j_client",
|
||||
"get_neo4j_client",
|
||||
"KnowledgeIterationService",
|
||||
"dep_knowledge_iteration_service",
|
||||
"WingmanService",
|
||||
"VisionService",
|
||||
"RagflowIngestionService",
|
||||
]
|
||||
@@ -1,331 +0,0 @@
|
||||
# =============================================================================
|
||||
# 企微IT智能服务台 — AI 服务(Dify 接入)
|
||||
# =============================================================================
|
||||
# 做什么:封装 Dify API 调用,实现 AI 自动回复
|
||||
# 为什么:
|
||||
# - ARCHITECTURE.md 设计了 ai_handling 状态,但当前未实现
|
||||
# - 现有系统交接文档提供了 Dify API 地址和 Key
|
||||
# - 这是实现「AI 自助解决」的核心模块
|
||||
# 依赖:需要 Dify API 可达(生产环境 http://yw-dify.dc.servyou-it.com)
|
||||
# =============================================================================
|
||||
|
||||
import json
|
||||
import logging
|
||||
import asyncio
|
||||
from typing import Any, Dict, List, Optional, AsyncGenerator
|
||||
|
||||
import httpx
|
||||
|
||||
from app.config import settings
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
|
||||
class AIService:
|
||||
"""AI 服务:封装 Dify API,提供 AI 回复能力。
|
||||
|
||||
支持两种调用模式:
|
||||
1. 非流式(简单场景):一次性获取完整回复
|
||||
2. 流式(推荐):SSE 流式返回,前端可逐字显示
|
||||
|
||||
参考:现有系统交接文档
|
||||
- API URL: http://yw-dify.dc.servyou-it.com/dify2openai/v1/chat/completions
|
||||
- Key: http://yw-dify.dc.servyou-it.com/v1|app-UaTWYdBSwN6VktKQlbh5YN5H|Chat
|
||||
"""
|
||||
|
||||
def __init__(self):
|
||||
"""初始化 AI 服务。
|
||||
|
||||
做什么:从配置读取 Dify API 地址和认证信息
|
||||
为什么:集中管理 API 配置,便于切换测试/生产环境
|
||||
"""
|
||||
# Dify 兼容 OpenAI 格式的 API 端点
|
||||
self.api_url = settings.dify_api_url
|
||||
# Dify API Key(格式:base_url|app_id|app_name)
|
||||
self.api_key = settings.dify_api_key
|
||||
# 请求超时(秒)
|
||||
self.timeout = settings.dify_timeout
|
||||
|
||||
# httpx 异步客户端(复用连接池)
|
||||
self._client: Optional[httpx.AsyncClient] = None
|
||||
|
||||
async def _get_client(self) -> httpx.AsyncClient:
|
||||
"""获取或创建 httpx 异步客户端。
|
||||
|
||||
做什么:懒加载 httpx.AsyncClient,复用连接池
|
||||
为什么:避免每次请求都创建新连接,提升性能
|
||||
"""
|
||||
if self._client is None or self._client.is_closed:
|
||||
self._client = httpx.AsyncClient(
|
||||
timeout=httpx.Timeout(self.timeout),
|
||||
headers={
|
||||
"Authorization": f"Bearer {self.api_key}",
|
||||
"Content-Type": "application/json",
|
||||
}
|
||||
)
|
||||
return self._client
|
||||
|
||||
async def close(self):
|
||||
"""关闭 httpx 客户端。
|
||||
|
||||
做什么:释放连接池资源
|
||||
为什么:避免连接泄漏,尤其在长期运行的 FastAPI 应用中
|
||||
"""
|
||||
if self._client and not self._client.is_closed:
|
||||
await self._client.aclose()
|
||||
self._client = None
|
||||
logger.debug("AIService httpx client closed")
|
||||
|
||||
# --------------------------------------------------------------------------
|
||||
# 非流式调用:一次性获取 AI 完整回复
|
||||
# --------------------------------------------------------------------------
|
||||
async def get_reply(
|
||||
self,
|
||||
message: str,
|
||||
conversation_id: Optional[str] = None,
|
||||
user_id: Optional[str] = None,
|
||||
) -> Dict[str, Any]:
|
||||
"""调用 Dify API 获取 AI 回复(非流式)。
|
||||
|
||||
Args:
|
||||
message: 员工发送的消息内容
|
||||
conversation_id: 会话ID(用于 Dify 多轮对话上下文)
|
||||
user_id: 员工企微 UserID(用于 Dify 用户标识)
|
||||
|
||||
Returns:
|
||||
Dict: {
|
||||
"content": str, # AI 回复内容
|
||||
"hit": bool, # 是否命中知识库(可回复)
|
||||
"conversation_id": str, # Dify 会话ID(用于后续多轮对话)
|
||||
"usage": dict, # Token 用量(可选)
|
||||
}
|
||||
|
||||
做什么:发送消息到 Dify,解析返回内容,判断是否能回复
|
||||
为什么:
|
||||
- 非流式适合简单场景,代码简单
|
||||
- 返回结构兼容 OpenAI Chat Completions 格式
|
||||
- 通过回复内容判断是否命中知识库(有实质内容 = 命中)
|
||||
"""
|
||||
payload = {
|
||||
"model": "Chat", # Dify 应用名称(来自 API Key 格式)
|
||||
"messages": [
|
||||
{"role": "user", "content": message}
|
||||
],
|
||||
"stream": False, # 非流式
|
||||
"temperature": 0.1, # 低温度,保证回答稳定性
|
||||
}
|
||||
|
||||
# 传入 Dify 会话ID,保持多轮对话上下文
|
||||
if conversation_id:
|
||||
payload["conversation_id"] = conversation_id
|
||||
|
||||
# 传入用户标识(Dify 侧用于日志和追溯)
|
||||
if user_id:
|
||||
payload["user"] = user_id
|
||||
|
||||
try:
|
||||
client = await self._get_client()
|
||||
logger.info(f"调用 Dify API: message={message[:50]}...")
|
||||
response = await client.post(self.api_url, json=payload)
|
||||
response.raise_for_status()
|
||||
data = response.json()
|
||||
|
||||
# 解析 OpenAI 兼容格式的返回
|
||||
# 格式:{"choices": [{"message": {"content": "..."}}]}
|
||||
choices = data.get("choices", [])
|
||||
if not choices:
|
||||
logger.warning("Dify API 返回空 choices")
|
||||
return {
|
||||
"content": "",
|
||||
"hit": False,
|
||||
"conversation_id": conversation_id or "",
|
||||
"usage": {},
|
||||
}
|
||||
|
||||
reply_content = choices[0]["message"]["content"]
|
||||
|
||||
# 判断是否命中知识库:
|
||||
# 策略1:检查内容是否为空或过长(Dify 可能返回提示语)
|
||||
# 策略2:检查是否包含「抱歉」「不知道」等无法回答的特征词
|
||||
hit = self._check_knowledge_hit(reply_content)
|
||||
|
||||
# 提取 Dify 返回的 conversation_id(用于多轮对话)
|
||||
dify_conv_id = data.get("conversation_id", conversation_id or "")
|
||||
|
||||
logger.info(
|
||||
f"Dify API 返回: hit={hit}, "
|
||||
f"content_length={len(reply_content)}, "
|
||||
f"conv_id={dify_conv_id[:20] if dify_conv_id else '(new)'}"
|
||||
)
|
||||
|
||||
return {
|
||||
"content": reply_content,
|
||||
"hit": hit,
|
||||
"conversation_id": dify_conv_id,
|
||||
"usage": data.get("usage", {}),
|
||||
}
|
||||
|
||||
except httpx.TimeoutException:
|
||||
logger.error("Dify API 超时")
|
||||
return {
|
||||
"content": "⏰ AI 服务响应超时,请稍后再试或输入「IT」转人工。",
|
||||
"hit": False,
|
||||
"conversation_id": conversation_id or "",
|
||||
"usage": {},
|
||||
}
|
||||
except httpx.HTTPStatusError as e:
|
||||
logger.error(f"Dify API HTTP 错误: status={e.response.status_code}")
|
||||
return {
|
||||
"content": "⚠️ AI 服务暂时不可用,请输入「IT」转人工。",
|
||||
"hit": False,
|
||||
"conversation_id": conversation_id or "",
|
||||
"usage": {},
|
||||
}
|
||||
except Exception as e:
|
||||
logger.error(f"Dify API 调用失败: {e}")
|
||||
return {
|
||||
"content": "⚠️ AI 服务异常,请输入「IT」转人工。",
|
||||
"hit": False,
|
||||
"conversation_id": conversation_id or "",
|
||||
"usage": {},
|
||||
}
|
||||
|
||||
# --------------------------------------------------------------------------
|
||||
# 流式调用:SSE 流式返回(供 WebSocket 推送给前端)
|
||||
# --------------------------------------------------------------------------
|
||||
async def get_reply_stream(
|
||||
self,
|
||||
message: str,
|
||||
conversation_id: Optional[str] = None,
|
||||
user_id: Optional[str] = None,
|
||||
) -> AsyncGenerator[Dict[str, Any], None]:
|
||||
"""调用 Dify API 获取流式 AI 回复(SSE),逐块 yield 给调用方。
|
||||
|
||||
Yields:
|
||||
Dict: {"delta": str, "finished": bool, "conversation_id": str, "hit": bool|None}
|
||||
- 流式中间块:{"delta": 增量, "finished": False, "hit": None}
|
||||
- 终态块:{"delta": "", "finished": True, "hit": 命中判断}
|
||||
|
||||
实现:
|
||||
- stream=True 走 SSE,解析 data: {...} 行,逐块 yield delta
|
||||
- 流结束后用完整内容整体判断 hit(_check_knowledge_hit)
|
||||
容错:若 Dify 不支持流式 / 超时 / 非 SSE 格式,catch 后 fallback 到
|
||||
get_reply 非流式,yield 一次完整内容(前端退化为"整段到达",
|
||||
功能不破,仅无逐字动画)。
|
||||
"""
|
||||
payload = {
|
||||
"model": "Chat",
|
||||
"messages": [{"role": "user", "content": message}],
|
||||
"stream": True,
|
||||
"temperature": 0.1,
|
||||
}
|
||||
if conversation_id:
|
||||
payload["conversation_id"] = conversation_id
|
||||
if user_id:
|
||||
payload["user"] = user_id
|
||||
|
||||
try:
|
||||
client = await self._get_client()
|
||||
full_parts: list = []
|
||||
dify_conv_id = conversation_id or ""
|
||||
async with client.stream("POST", self.api_url, json=payload) as response:
|
||||
response.raise_for_status()
|
||||
async for line in response.aiter_lines():
|
||||
if not line:
|
||||
continue
|
||||
line = line.strip()
|
||||
if not line.startswith("data:"):
|
||||
continue
|
||||
data = line[5:].strip()
|
||||
if data == "[DONE]":
|
||||
break
|
||||
try:
|
||||
chunk = json.loads(data)
|
||||
except json.JSONDecodeError:
|
||||
continue
|
||||
# OpenAI / Dify SSE 格式:choices[0].delta.content
|
||||
try:
|
||||
delta = chunk["choices"][0]["delta"].get("content", "")
|
||||
except (KeyError, IndexError, TypeError):
|
||||
delta = ""
|
||||
if delta:
|
||||
full_parts.append(delta)
|
||||
yield {
|
||||
"delta": delta,
|
||||
"finished": False,
|
||||
"conversation_id": dify_conv_id,
|
||||
"hit": None,
|
||||
}
|
||||
# Dify 可能在流式块里给出 conversation_id
|
||||
cid = chunk.get("conversation_id")
|
||||
if cid:
|
||||
dify_conv_id = cid
|
||||
|
||||
# 流结束:用完整内容判断命中
|
||||
full_content = "".join(full_parts)
|
||||
hit = self._check_knowledge_hit(full_content) if full_content else False
|
||||
yield {
|
||||
"delta": "",
|
||||
"finished": True,
|
||||
"conversation_id": dify_conv_id,
|
||||
"hit": hit,
|
||||
}
|
||||
except Exception as e:
|
||||
# 流式不可用(dify2openai 不支持 / 超时 / 非 SSE),回退非流式
|
||||
logger.warning(f"Dify 流式失败,回退非流式: {e}")
|
||||
try:
|
||||
result = await self.get_reply(message, conversation_id, user_id)
|
||||
yield {
|
||||
"delta": result["content"],
|
||||
"finished": True,
|
||||
"conversation_id": result["conversation_id"],
|
||||
"hit": result["hit"],
|
||||
}
|
||||
except Exception as e2:
|
||||
logger.error(f"Dify 流式与非流式均失败: {e2}")
|
||||
yield {
|
||||
"delta": "⚠️ AI 服务异常,请输入「IT」转人工或稍后重试。",
|
||||
"finished": True,
|
||||
"conversation_id": conversation_id or "",
|
||||
"hit": False,
|
||||
}
|
||||
|
||||
# --------------------------------------------------------------------------
|
||||
# 判断是否命中知识库
|
||||
# --------------------------------------------------------------------------
|
||||
def _check_knowledge_hit(self, content: str) -> bool:
|
||||
"""判断 AI 回复是否命中知识库(可以回答用户问题)。
|
||||
|
||||
Args:
|
||||
content: AI 回复内容
|
||||
|
||||
Returns:
|
||||
bool: True=命中(可以回复),False=未命中(需转人工)
|
||||
|
||||
做什么:分析 AI 回复内容,判断是否能有效回答问题
|
||||
为什么:
|
||||
- Dify 在无法回答时通常会返回固定提示语
|
||||
- 参考现有系统:「抱歉,您的问题可能不在服务业务范围内」
|
||||
- 命中 = 有实质内容且不像是「无法回答」的提示
|
||||
"""
|
||||
if not content or len(content.strip()) < 5:
|
||||
return False
|
||||
|
||||
# 未命中特征词(Dify 无法回答时的典型回复)
|
||||
miss_keywords = [
|
||||
"抱歉", "对不起", "不知道", "无法回答",
|
||||
"不在服务范围内", "超出我的能力", "暂不支持",
|
||||
"请转人工", "联系管理员",
|
||||
]
|
||||
content_lower = content.lower()
|
||||
|
||||
# 如果回复中包含多个未命中特征词 → 判断为未命中
|
||||
miss_count = sum(1 for kw in miss_keywords if kw in content_lower)
|
||||
if miss_count >= 2:
|
||||
return False
|
||||
|
||||
# 如果回复长度过短(< 10 字符)且包含特征词 → 未命中
|
||||
if len(content) < 10 and any(kw in content_lower for kw in miss_keywords):
|
||||
return False
|
||||
|
||||
return True
|
||||
@@ -1,48 +0,0 @@
|
||||
# =============================================================================
|
||||
# 企微IT智能服务台 — 阶段5 自动化 意图识别
|
||||
# =============================================================================
|
||||
# 说明:调用 Dify 识别员工诉求命中哪个自动化场景;Dify 未配置时走关键词兜底,
|
||||
# 保证 P0 四个场景在无真实 Dify 环境下也能跑通闭环。
|
||||
# =============================================================================
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import logging
|
||||
from typing import Any, Dict, Optional
|
||||
|
||||
from app.integrations.dify import DifyClient, get_dify_client
|
||||
from app.integrations.factory import build_dify_client
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
|
||||
class IntentRouter:
|
||||
"""意图识别路由器。"""
|
||||
|
||||
def __init__(self, db: Any = None, audit: Any = None):
|
||||
self.db = db
|
||||
self.audit = audit
|
||||
|
||||
async def detect(self, description: str, employee_id: str = "") -> Dict[str, Any]:
|
||||
"""识别意图,返回 {scenario_key, confidence, raw, error}。
|
||||
|
||||
优先走 Dify;若 Dify 未配置或调用失败,使用关键词兜底。
|
||||
"""
|
||||
client: Optional[DifyClient] = None
|
||||
try:
|
||||
client = await build_dify_client(audit=self.audit)
|
||||
except Exception as e: # noqa: BLE001
|
||||
logger.warning(f"构建 Dify 客户端失败,转关键词兜底: {e}")
|
||||
|
||||
if client is None:
|
||||
fb = DifyClient._fallback_intent(description)
|
||||
fb["error"] = "dify_not_configured"
|
||||
return fb
|
||||
|
||||
try:
|
||||
return await client.detect_intent(description, employee_id)
|
||||
except Exception as e: # noqa: BLE001
|
||||
logger.warning(f"Dify 意图识别异常,转关键词兜底: {e}")
|
||||
fb = DifyClient._fallback_intent(description)
|
||||
fb["error"] = str(e)
|
||||
return fb
|
||||
@@ -1,486 +0,0 @@
|
||||
# =============================================================================
|
||||
# 企微IT智能服务台 — 阶段5 自动化 会话编排服务
|
||||
# =============================================================================
|
||||
# 说明:自动化会话的编排中枢,串联「意图识别 → 场景校验 → 终端映射 →
|
||||
# 动作计划生成 → 执行引擎」。同时提供会话 CRUD、转人工、结果反馈、
|
||||
# 静默关单等接口。
|
||||
#
|
||||
# 编排在后台任务中运行(run_session_in_background),API 创建会话后立即返回,
|
||||
# 进度通过 WS 实时推送。
|
||||
# =============================================================================
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import logging
|
||||
from datetime import datetime, timezone
|
||||
from typing import Any, Dict, List, Optional
|
||||
|
||||
from sqlalchemy import select
|
||||
|
||||
from app.config import settings
|
||||
from app.constants import AutomationErrorCode
|
||||
from app.database import _get_session_factory
|
||||
from app.models.automation import (
|
||||
ApprovalTicket,
|
||||
AutoAction,
|
||||
AutoSession,
|
||||
ScenarioConfig,
|
||||
)
|
||||
from app.services.automation.approval import ApprovalService
|
||||
from app.services.automation.exception_handler import AutomationException
|
||||
from app.services.automation.executor import ActionExecutor
|
||||
from app.services.automation.intent_router import IntentRouter
|
||||
from app.services.automation.mapping_resolver import MappingResolver
|
||||
from app.services.automation.progress_publisher import (
|
||||
cancel_silent_close,
|
||||
publish_progress,
|
||||
publish_takeover,
|
||||
)
|
||||
from app.services.automation import DEFAULT_SCENARIO_CONFIGS
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
# 需要终端映射的场景
|
||||
_SCENARIOS_NEED_TERMINAL = {"virus_dispose", "terminal_locate"}
|
||||
|
||||
|
||||
class AutoSessionService:
|
||||
"""自动化会话编排服务。"""
|
||||
|
||||
def __init__(self, db: Any, redis: Any = None, audit: Any = None):
|
||||
self.db = db
|
||||
self.redis = redis
|
||||
self.audit = audit
|
||||
|
||||
# --------------------------------------------------------------------------
|
||||
# 会话 CRUD
|
||||
# --------------------------------------------------------------------------
|
||||
async def create_session(
|
||||
self,
|
||||
conversation_id: Optional[str],
|
||||
employee_id: str,
|
||||
description: str,
|
||||
mode: str = "real_exec",
|
||||
) -> AutoSession:
|
||||
"""创建自动化会话(状态=created)。"""
|
||||
session = AutoSession(
|
||||
conversation_id=conversation_id,
|
||||
employee_id=employee_id,
|
||||
mode=mode,
|
||||
title=(description or "")[:200],
|
||||
status="created",
|
||||
meta={"description": description},
|
||||
)
|
||||
self.db.add(session)
|
||||
await self.db.flush()
|
||||
return session
|
||||
|
||||
async def get_session(self, session_id: str) -> Optional[AutoSession]:
|
||||
"""按 ID 取会话。"""
|
||||
stmt = select(AutoSession).where(AutoSession.id == session_id)
|
||||
return (await self.db.execute(stmt)).scalar_one_or_none()
|
||||
|
||||
async def list_sessions(
|
||||
self,
|
||||
employee_id: Optional[str] = None,
|
||||
agent_id: Optional[str] = None,
|
||||
status: Optional[str] = None,
|
||||
page: int = 1,
|
||||
page_size: int = 50,
|
||||
) -> List[AutoSession]:
|
||||
"""列出会话(支持过滤分页)。"""
|
||||
stmt = select(AutoSession)
|
||||
if employee_id:
|
||||
stmt = stmt.where(AutoSession.employee_id == employee_id)
|
||||
if agent_id:
|
||||
stmt = stmt.where(AutoSession.agent_id == agent_id)
|
||||
if status:
|
||||
stmt = stmt.where(AutoSession.status == status)
|
||||
stmt = stmt.order_by(AutoSession.created_at.desc())
|
||||
stmt = stmt.limit(page_size).offset((page - 1) * page_size)
|
||||
return list((await self.db.execute(stmt)).scalars().all())
|
||||
|
||||
async def get_session_detail(
|
||||
self, session_id: str
|
||||
) -> Optional[Dict[str, Any]]:
|
||||
"""取会话详情(含动作列表与当前待决审批单)。"""
|
||||
session = await self.get_session(session_id)
|
||||
if session is None:
|
||||
return None
|
||||
act_stmt = select(AutoAction).where(AutoAction.session_id == session_id)
|
||||
actions = list((await self.db.execute(act_stmt)).scalars().all())
|
||||
actions.sort(key=lambda a: a.action_index)
|
||||
|
||||
ticket = None
|
||||
if session.current_action_id:
|
||||
tstmt = select(ApprovalTicket).where(
|
||||
ApprovalTicket.action_id == session.current_action_id,
|
||||
ApprovalTicket.status == "pending",
|
||||
)
|
||||
ticket = (await self.db.execute(tstmt)).scalar_one_or_none()
|
||||
return {"session": session, "actions": actions, "ticket": ticket}
|
||||
|
||||
# --------------------------------------------------------------------------
|
||||
# 编排主流程
|
||||
# --------------------------------------------------------------------------
|
||||
async def start(self, session_id: str) -> None:
|
||||
"""编排:意图识别 → 场景校验 → 映射 → 计划 → 执行。"""
|
||||
session = await self.get_session(session_id)
|
||||
if session is None:
|
||||
logger.warning(f"start 会话不存在: {session_id}")
|
||||
return
|
||||
if session.status != "created":
|
||||
logger.info(f"会话已启动过,跳过: {session_id} status={session.status}")
|
||||
return
|
||||
|
||||
session.status = "running"
|
||||
await self.db.flush()
|
||||
await publish_progress(session.id, "start", "开始自动化处置")
|
||||
|
||||
description = (session.meta or {}).get("description", "")
|
||||
|
||||
# 1. 意图识别
|
||||
router = IntentRouter(self.db, audit=self.audit)
|
||||
try:
|
||||
intent = await router.detect(description, session.employee_id)
|
||||
except Exception as e: # noqa: BLE001
|
||||
intent = {"scenario_key": None, "confidence": 0.0, "error": str(e)}
|
||||
session.scenario_key = intent.get("scenario_key")
|
||||
session.confidence = float(intent.get("confidence") or 0.0)
|
||||
session.intent = intent
|
||||
await self.db.flush()
|
||||
await publish_progress(
|
||||
session.id,
|
||||
"intent",
|
||||
f"识别场景: {session.scenario_key or '未知'}(置信度 {session.confidence:.2f})",
|
||||
)
|
||||
|
||||
# 2. 置信度门槛 → 低置信度转人工
|
||||
thresholds = settings.get_automation_thresholds()
|
||||
confidence_min = float(thresholds.get("confidence_min", 0.6))
|
||||
if not session.scenario_key or session.confidence < confidence_min:
|
||||
await self._handoff(
|
||||
session, f"意图识别置信度不足({session.confidence:.2f}),转人工"
|
||||
)
|
||||
return
|
||||
|
||||
# 3. 场景配置校验
|
||||
scenario = await self._get_scenario_config(session.scenario_key)
|
||||
if scenario is None or not scenario.enabled:
|
||||
await self._handoff(
|
||||
session, f"场景未启用或未配置: {session.scenario_key}"
|
||||
)
|
||||
return
|
||||
|
||||
# 4. 终端映射(按需)
|
||||
mapping: Dict[str, Any] = {}
|
||||
if session.scenario_key in _SCENARIOS_NEED_TERMINAL:
|
||||
resolver = MappingResolver(self.db, audit=self.audit)
|
||||
mapping = await resolver.resolve(session.employee_id, session.scenario_key)
|
||||
if (
|
||||
session.scenario_key == "virus_dispose"
|
||||
and not mapping.get("client_ids")
|
||||
):
|
||||
await self._handoff(session, "未解析到目标终端,转人工")
|
||||
return
|
||||
session.meta = {**(session.meta or {}), "mapping": mapping}
|
||||
await self.db.flush()
|
||||
|
||||
# 5. 生成动作计划
|
||||
plan = self._build_actions(scenario, mapping)
|
||||
for idx, item in enumerate(plan):
|
||||
action = AutoAction(session_id=session.id, action_index=idx, **item)
|
||||
self.db.add(action)
|
||||
await self.db.flush()
|
||||
await publish_progress(
|
||||
session.id, "plan_ready", f"已生成处置方案(共 {len(plan)} 步)"
|
||||
)
|
||||
|
||||
# 6. 执行引擎
|
||||
executor = ActionExecutor(self.db, self.redis, audit=self.audit)
|
||||
await executor.run(session_id)
|
||||
|
||||
def _build_actions(
|
||||
self, scenario: ScenarioConfig, mapping: Dict[str, Any]
|
||||
) -> List[Dict[str, Any]]:
|
||||
"""从场景配置(或默认模板)生成动作计划。"""
|
||||
plan_items = scenario.actions
|
||||
if not plan_items:
|
||||
default = DEFAULT_SCENARIO_CONFIGS.get(scenario.scenario_key, {})
|
||||
plan_items = default.get("actions", [])
|
||||
|
||||
result: List[Dict[str, Any]] = []
|
||||
for it in plan_items:
|
||||
params: Dict[str, Any] = dict(it.get("params") or {})
|
||||
confirm_channel = it.get("confirm_channel")
|
||||
if mapping.get("client_ids"):
|
||||
params.setdefault("client_ids", mapping["client_ids"])
|
||||
if confirm_channel:
|
||||
params["confirm_channel"] = confirm_channel
|
||||
result.append(
|
||||
{
|
||||
"action_type": it.get("action_type", ""),
|
||||
"adapter": it.get("adapter", "internal"),
|
||||
"risk_level": it.get("risk_level", "read"),
|
||||
"title": it.get("title", it.get("action_type", "")),
|
||||
"description": it.get("description", it.get("title", "")),
|
||||
"payload": params,
|
||||
}
|
||||
)
|
||||
return result
|
||||
|
||||
async def _get_scenario_config(
|
||||
self, scenario_key: str
|
||||
) -> Optional[ScenarioConfig]:
|
||||
"""加载场景配置(DB 优先,缺失则用默认模板的开关)。"""
|
||||
stmt = select(ScenarioConfig).where(
|
||||
ScenarioConfig.scenario_key == scenario_key
|
||||
)
|
||||
config = (await self.db.execute(stmt)).scalar_one_or_none()
|
||||
if config is not None:
|
||||
return config
|
||||
# DB 无记录 → 用默认模板(默认启用)
|
||||
default = DEFAULT_SCENARIO_CONFIGS.get(scenario_key)
|
||||
if default is None:
|
||||
return None
|
||||
return ScenarioConfig(
|
||||
scenario_key=scenario_key,
|
||||
name=default.get("name", scenario_key),
|
||||
description=default.get("description", ""),
|
||||
enabled=default.get("enabled", True),
|
||||
trigger_conditions=default.get("trigger_conditions"),
|
||||
actions=default.get("actions"),
|
||||
approval_strategy=default.get("approval_strategy"),
|
||||
)
|
||||
|
||||
async def _handoff(self, session: AutoSession, reason: str) -> None:
|
||||
"""置为转人工并推送事件。"""
|
||||
session.status = "handoff"
|
||||
session.closed_by = "system(auto)"
|
||||
await self.db.flush()
|
||||
await publish_takeover(session.id, reason)
|
||||
|
||||
# --------------------------------------------------------------------------
|
||||
# 转人工 / 反馈 / 关单
|
||||
# --------------------------------------------------------------------------
|
||||
async def takeover(
|
||||
self, session_id: str, agent_id: str, note: Optional[str] = None
|
||||
) -> AutoSession:
|
||||
"""转人工接管。"""
|
||||
session = await self.get_session(session_id)
|
||||
if session is None:
|
||||
raise AutomationException(AutomationErrorCode.SESSION_NOT_FOUND)
|
||||
session.status = "handoff"
|
||||
session.agent_id = agent_id
|
||||
session.closed_by = agent_id
|
||||
await self.db.flush()
|
||||
cancel_silent_close(session_id)
|
||||
await publish_takeover(session.id, f"坐席 {agent_id} 接管:{note or ''}")
|
||||
return session
|
||||
|
||||
async def resolve_feedback(
|
||||
self, session_id: str, satisfied: bool, note: Optional[str] = None
|
||||
) -> AutoSession:
|
||||
"""员工处置结果反馈。
|
||||
|
||||
满意 → 关单;不满意 → 转人工。
|
||||
"""
|
||||
session = await self.get_session(session_id)
|
||||
if session is None:
|
||||
raise AutomationException(AutomationErrorCode.SESSION_NOT_FOUND)
|
||||
cancel_silent_close(session_id)
|
||||
if satisfied:
|
||||
session.status = "closed"
|
||||
session.closed_by = session.employee_id
|
||||
else:
|
||||
session.status = "handoff"
|
||||
session.closed_by = "employee(reject)"
|
||||
await publish_takeover(session.id, f"员工不满意,转人工:{note or ''}")
|
||||
await self.db.flush()
|
||||
return session
|
||||
|
||||
async def auto_close(self, session_id: str) -> None:
|
||||
"""静默关单(仅 resolved 态可关)。"""
|
||||
session = await self.get_session(session_id)
|
||||
if session is None:
|
||||
return
|
||||
if session.status != "resolved":
|
||||
return
|
||||
session.status = "closed"
|
||||
session.closed_by = "system(auto)"
|
||||
await self.db.flush()
|
||||
|
||||
# --------------------------------------------------------------------------
|
||||
# 场景配置管理(管理端)
|
||||
# --------------------------------------------------------------------------
|
||||
async def list_scenario_configs(self) -> List[ScenarioConfig]:
|
||||
"""列出全部场景配置。"""
|
||||
stmt = select(ScenarioConfig).order_by(ScenarioConfig.scenario_key)
|
||||
return list((await self.db.execute(stmt)).scalars().all())
|
||||
|
||||
async def upsert_scenario_config(
|
||||
self, scenario_key: str, data: Dict[str, Any], operator: str = ""
|
||||
) -> ScenarioConfig:
|
||||
"""更新或创建场景配置,并快照为规则版本。"""
|
||||
from app.models.automation import RuleVersion
|
||||
|
||||
stmt = select(ScenarioConfig).where(
|
||||
ScenarioConfig.scenario_key == scenario_key
|
||||
)
|
||||
config = (await self.db.execute(stmt)).scalar_one_or_none()
|
||||
if config is None:
|
||||
config = ScenarioConfig(scenario_key=scenario_key)
|
||||
self.db.add(config)
|
||||
|
||||
if data.get("name") is not None:
|
||||
config.name = data["name"]
|
||||
if data.get("description") is not None:
|
||||
config.description = data["description"]
|
||||
if data.get("enabled") is not None:
|
||||
config.enabled = data["enabled"]
|
||||
if data.get("trigger_conditions") is not None:
|
||||
config.trigger_conditions = data["trigger_conditions"]
|
||||
if data.get("actions") is not None:
|
||||
config.actions = data["actions"]
|
||||
if data.get("approval_strategy") is not None:
|
||||
config.approval_strategy = data["approval_strategy"]
|
||||
|
||||
await self.db.flush()
|
||||
|
||||
# 快照为规则版本
|
||||
last = (
|
||||
await self.db.execute(
|
||||
select(RuleVersion.version)
|
||||
.where(RuleVersion.scenario_key == scenario_key)
|
||||
.order_by(RuleVersion.version.desc())
|
||||
.limit(1)
|
||||
)
|
||||
).scalar_one_or_none()
|
||||
new_version = (last or 0) + 1
|
||||
snapshot = RuleVersion(
|
||||
scenario_key=scenario_key,
|
||||
version=new_version,
|
||||
content={
|
||||
"actions": config.actions,
|
||||
"approval_strategy": config.approval_strategy,
|
||||
"trigger_conditions": config.trigger_conditions,
|
||||
},
|
||||
status="published",
|
||||
canary_percent=100,
|
||||
created_by=operator or "admin",
|
||||
remark=f"更新场景配置至 v{new_version}",
|
||||
)
|
||||
self.db.add(snapshot)
|
||||
await self.db.flush()
|
||||
config.current_version_id = snapshot.id
|
||||
await self.db.flush()
|
||||
return config
|
||||
|
||||
async def list_rule_versions(
|
||||
self, scenario_key: Optional[str] = None
|
||||
) -> List[RuleVersion]:
|
||||
"""列出规则版本。"""
|
||||
stmt = select(RuleVersion)
|
||||
if scenario_key:
|
||||
stmt = stmt.where(RuleVersion.scenario_key == scenario_key)
|
||||
stmt = stmt.order_by(RuleVersion.created_at.desc())
|
||||
return list((await self.db.execute(stmt)).scalars().all())
|
||||
|
||||
async def metrics(self) -> Dict[str, Any]:
|
||||
"""汇总看板指标。"""
|
||||
from sqlalchemy import func
|
||||
|
||||
total = (
|
||||
await self.db.execute(select(func.count(AutoSession.id)))
|
||||
).scalar() or 0
|
||||
resolved = (
|
||||
await self.db.execute(
|
||||
select(func.count(AutoSession.id)).where(
|
||||
AutoSession.status == "resolved"
|
||||
)
|
||||
)
|
||||
).scalar() or 0
|
||||
handoff = (
|
||||
await self.db.execute(
|
||||
select(func.count(AutoSession.id)).where(
|
||||
AutoSession.status == "handoff"
|
||||
)
|
||||
)
|
||||
).scalar() or 0
|
||||
error = (
|
||||
await self.db.execute(
|
||||
select(func.count(AutoSession.id)).where(
|
||||
AutoSession.status == "error"
|
||||
)
|
||||
)
|
||||
).scalar() or 0
|
||||
auto_actions = (
|
||||
await self.db.execute(
|
||||
select(func.count(AutoAction.id)).where(
|
||||
AutoAction.status == "success",
|
||||
AutoAction.risk_level.in_(["read", "low"]),
|
||||
)
|
||||
)
|
||||
).scalar() or 0
|
||||
approval_actions = (
|
||||
await self.db.execute(
|
||||
select(func.count(AutoAction.id)).where(
|
||||
AutoAction.risk_level == "high"
|
||||
)
|
||||
)
|
||||
).scalar() or 0
|
||||
|
||||
by_scenario_rows = (
|
||||
await self.db.execute(
|
||||
select(AutoSession.scenario_key, func.count(AutoSession.id)).group_by(
|
||||
AutoSession.scenario_key
|
||||
)
|
||||
)
|
||||
).all()
|
||||
by_scenario = {k: v for k, v in by_scenario_rows if k}
|
||||
|
||||
return {
|
||||
"total_sessions": total,
|
||||
"resolved_sessions": resolved,
|
||||
"handoff_sessions": handoff,
|
||||
"error_sessions": error,
|
||||
"auto_executed_actions": auto_actions,
|
||||
"approval_required_actions": approval_actions,
|
||||
"by_scenario": by_scenario,
|
||||
}
|
||||
|
||||
|
||||
# --------------------------------------------------------------------------
|
||||
# 后台编排入口
|
||||
# --------------------------------------------------------------------------
|
||||
async def run_session_in_background(session_id: str) -> None:
|
||||
"""后台运行会话编排(独立 DB 会话,结束时提交/回滚)。"""
|
||||
factory = _get_session_factory()
|
||||
async with factory() as db:
|
||||
svc = AutoSessionService(db)
|
||||
try:
|
||||
await svc.start(session_id)
|
||||
await db.commit()
|
||||
except AutomationException as e:
|
||||
await db.rollback()
|
||||
logger.warning(f"编排业务异常 session={session_id}: {e.message}")
|
||||
# 标记转人工
|
||||
try:
|
||||
session = await svc.get_session(session_id)
|
||||
if session and session.status not in ("closed", "handoff"):
|
||||
session.status = "handoff"
|
||||
session.closed_by = "system(auto)"
|
||||
await db.commit()
|
||||
except Exception: # noqa: BLE001
|
||||
await db.rollback()
|
||||
except Exception as e: # noqa: BLE001
|
||||
await db.rollback()
|
||||
logger.error(f"编排未预期异常 session={session_id}: {e}")
|
||||
try:
|
||||
session = await svc.get_session(session_id)
|
||||
if session and session.status not in ("closed", "handoff"):
|
||||
session.status = "handoff"
|
||||
session.closed_by = "system(auto)"
|
||||
await db.commit()
|
||||
except Exception: # noqa: BLE001
|
||||
await db.rollback()
|
||||
@@ -1,331 +0,0 @@
|
||||
# =============================================================================
|
||||
# 企微IT智能服务台 — 内容审核服务
|
||||
# =============================================================================
|
||||
# 说明:#81 v0.6.0 内容审核 — 检测敏感词 + 提示坐席优化语气
|
||||
# 用途:坐席发送消息前自动审核,避免发送违规内容
|
||||
# 设计:基于 wordfilter 开源库 + 自定义敏感词库
|
||||
# =============================================================================
|
||||
|
||||
from dataclasses import dataclass
|
||||
from enum import Enum
|
||||
from typing import List, Optional, Tuple
|
||||
|
||||
from wordfilter import Wordfilter
|
||||
|
||||
import logging
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
|
||||
class ModerationAction(str, Enum):
|
||||
"""内容审核动作"""
|
||||
PASS = "pass" # 通过
|
||||
WARN = "warn" # 警告(允许发送,但标记)
|
||||
BLOCK = "block" # 阻断(必须修改)
|
||||
|
||||
|
||||
class ModerationCategory(str, Enum):
|
||||
"""审核分类"""
|
||||
PROFANITY = "profanity" # 脏话
|
||||
POLITICS = "politics" # 政治敏感
|
||||
PORN = "porn" # 色情
|
||||
AD = "ad" # 广告
|
||||
PRIVACY = "privacy" # 隐私泄露(身份证/电话)
|
||||
OTHER = "other" # 其他
|
||||
|
||||
|
||||
@dataclass
|
||||
class ModerationResult:
|
||||
"""审核结果"""
|
||||
action: ModerationAction
|
||||
category: Optional[ModerationCategory]
|
||||
matched_words: List[str]
|
||||
suggestion: str = ""
|
||||
|
||||
@property
|
||||
def is_blocked(self) -> bool:
|
||||
return self.action == ModerationAction.BLOCK
|
||||
|
||||
@property
|
||||
def is_warned(self) -> bool:
|
||||
return self.action == ModerationAction.WARN
|
||||
|
||||
|
||||
class ContentModerationService:
|
||||
"""内容审核服务 — 检测 + 提示。
|
||||
|
||||
设计要点:
|
||||
1. 加载 wordfilter + 自定义敏感词库
|
||||
2. 提供 3 个级别动作:pass / warn / block
|
||||
3. 返回命中的敏感词,给前端提示
|
||||
4. 异步不阻塞消息发送主流程
|
||||
"""
|
||||
|
||||
def __init__(self):
|
||||
# 初始化 wordfilter(新 API: Wordfilter() 实例,而非 init() 全局)
|
||||
self.wf = Wordfilter()
|
||||
# 加载自定义敏感词库(预留,生产环境从配置文件加载)
|
||||
self.custom_sensitive_words: List[str] = [
|
||||
# 坐席严禁发送的
|
||||
"投诉我", # 暗示员工投诉自己
|
||||
"你爱找谁找谁", # 不当推诿
|
||||
"自己不会百度吗", # 不当反问
|
||||
"这点小事", # 轻视员工问题
|
||||
# 隐私保护(后端检测,前端不知道)
|
||||
# 实际部署时从 system_config 加载
|
||||
]
|
||||
if self.custom_sensitive_words:
|
||||
self.wf.addWords(self.custom_sensitive_words)
|
||||
|
||||
# ==================================================================
|
||||
# 主入口
|
||||
# ==================================================================
|
||||
|
||||
def moderate(self, text: str) -> ModerationResult:
|
||||
"""审核文本。
|
||||
|
||||
Args:
|
||||
text: 待审核文本(坐席准备发的消息)
|
||||
|
||||
Returns:
|
||||
ModerationResult: 审核结果
|
||||
"""
|
||||
if not text or not text.strip():
|
||||
return ModerationResult(
|
||||
action=ModerationAction.PASS,
|
||||
category=None,
|
||||
matched_words=[],
|
||||
)
|
||||
|
||||
text = text.strip()
|
||||
|
||||
# 1. wordfilter 检测
|
||||
matched: List[str] = []
|
||||
if self.wf.blacklisted(text):
|
||||
# 找出具体哪些词命中
|
||||
matched = self._extract_matched(text)
|
||||
|
||||
if not matched:
|
||||
return ModerationResult(
|
||||
action=ModerationAction.PASS,
|
||||
category=None,
|
||||
matched_words=[],
|
||||
)
|
||||
|
||||
# 2. 分类(简单规则:有命中就给 warn,后续可分级)
|
||||
category = self._classify(matched)
|
||||
|
||||
# 3. 决定动作(目前策略:命中即 warn,后续可升级 block)
|
||||
# 后续决策点:是否给某些类(政治/色情)直接 block
|
||||
action = ModerationAction.WARN
|
||||
suggestion = self._generate_suggestion(category, matched)
|
||||
|
||||
logger.info(
|
||||
f"[ContentModeration] 检测到敏感词 text={text[:30]}... "
|
||||
f"matched={matched} category={category}"
|
||||
)
|
||||
|
||||
return ModerationResult(
|
||||
action=action,
|
||||
category=category,
|
||||
matched_words=matched,
|
||||
suggestion=suggestion,
|
||||
)
|
||||
|
||||
# ==================================================================
|
||||
# 隐私信息检测(基于正则,跟敏感词无关)
|
||||
# ==================================================================
|
||||
|
||||
def check_privacy_leak(self, text: str) -> List[str]:
|
||||
"""检测文本是否包含隐私信息(身份证 / 电话 / 银行卡)。
|
||||
|
||||
Returns:
|
||||
命中的隐私字段列表(描述性,如 ["phone", "id_card"])
|
||||
"""
|
||||
import re
|
||||
leaked = []
|
||||
|
||||
# 手机号(11位1开头)
|
||||
# BUGFIX: \b 和 (?<!\w) 对中文均失效(Python3 \w 含中文),
|
||||
# 改用 (?<!\d) / (?!\d) 检查数字边界——"电话13800138000" 可正确匹配
|
||||
if re.search(r"(?<!\d)1[3-9]\d{9}(?!\d)", text):
|
||||
leaked.append("phone")
|
||||
|
||||
# 身份证号(18位)
|
||||
if re.search(r"(?<!\d)\d{17}[\dXx](?!\d)", text):
|
||||
leaked.append("id_card")
|
||||
|
||||
# 银行卡(16-19位连续数字,简单判断)
|
||||
if re.search(r"(?<!\d)\d{16,19}(?!\d)", text):
|
||||
leaked.append("bank_card")
|
||||
|
||||
# 邮箱(个人邮箱,非公司邮箱)
|
||||
# 邮箱以 ASCII 字母开头,(?<!\w) 这里可用(前面不会是中文邮箱前缀)
|
||||
personal_email_pattern = (
|
||||
r"(?<!\w)[a-zA-Z0-9._%+-]+@(?!servyou-it\.com|"
|
||||
r"servyou\.com\.cn)[a-zA-Z0-9.-]+\.[a-zA-Z]{2,}(?!\w)"
|
||||
)
|
||||
if re.search(personal_email_pattern, text):
|
||||
leaked.append("personal_email")
|
||||
|
||||
return leaked
|
||||
|
||||
# ==================================================================
|
||||
# 工具方法
|
||||
# ==================================================================
|
||||
|
||||
def _extract_matched(self, text: str) -> List[str]:
|
||||
"""提取命中的敏感词。"""
|
||||
# wordfilter 没有直接的 "提取所有命中词" API,只能 replace 看
|
||||
matched = []
|
||||
# 遍历自建词库看哪些命中
|
||||
for word in self.custom_sensitive_words:
|
||||
if word in text:
|
||||
matched.append(word)
|
||||
return matched
|
||||
|
||||
def _classify(self, matched: List[str]) -> ModerationCategory:
|
||||
"""根据命中的词分类。"""
|
||||
# 简单分类:命中"投诉""爱找谁"等 → profanity
|
||||
# 后续可扩展
|
||||
return ModerationCategory.PROFANITY
|
||||
|
||||
def _generate_suggestion(
|
||||
self, category: ModerationCategory, matched: List[str]
|
||||
) -> str:
|
||||
"""生成修改建议。"""
|
||||
suggestions_map = {
|
||||
ModerationCategory.PROFANITY: (
|
||||
"建议改为更专业的表达,例如:"
|
||||
"「我理解您的问题,我们一起想办法解决」"
|
||||
),
|
||||
ModerationCategory.POLITICS: (
|
||||
"请避免讨论政治话题,保持服务专业性"
|
||||
),
|
||||
ModerationCategory.PORN: "请使用正式语言",
|
||||
ModerationCategory.AD: "请勿发送广告内容",
|
||||
ModerationCategory.PRIVACY: (
|
||||
"请勿发送员工隐私信息(电话/身份证),如需联系请走企微"
|
||||
),
|
||||
ModerationCategory.OTHER: "请检查并修改表达",
|
||||
}
|
||||
return suggestions_map.get(category, "请检查并修改表达")
|
||||
|
||||
@staticmethod
|
||||
def _get_fallback_question(keywords: List[str]) -> dict:
|
||||
"""Dify 失败时的兜底题(从预置题池随机抽一道)。
|
||||
|
||||
注意:这里写死 10 道 IT 基础题,生产环境可改成查 quiz_questions.source='manual'
|
||||
"""
|
||||
import random
|
||||
|
||||
fallback_pool = [
|
||||
{
|
||||
"question": "电脑突然黑屏,最安全的做法是?",
|
||||
"options": ["强制关机重启", "拔电源重启", "等几分钟看是否恢复", "砸电脑"],
|
||||
"correct_index": 0,
|
||||
"hint": "想想最稳妥的第一步",
|
||||
"explanation": "黑屏可能是系统卡死,强制重启通常能恢复,拔电源可能损坏硬件",
|
||||
"source": "manual",
|
||||
},
|
||||
{
|
||||
"question": "打印机不响应,首先应该检查?",
|
||||
"options": ["打印机电源", "重装系统", "换台电脑", "直接呼叫维修"],
|
||||
"correct_index": 0,
|
||||
"hint": "最基础的物理连接",
|
||||
"explanation": "80% 故障是电源/线缆问题,先排除最简单的再考虑复杂方案",
|
||||
"source": "manual",
|
||||
},
|
||||
{
|
||||
"question": "密码忘了应该怎么办?",
|
||||
"options": ["自己猜", "暴力破解", "找 IT 重置", "不用了"],
|
||||
"correct_index": 2,
|
||||
"hint": "走正规流程最安全",
|
||||
"explanation": "找 IT 重置是最快最安全的做法,自己猜可能锁账号,暴力破解违法",
|
||||
"source": "manual",
|
||||
},
|
||||
{
|
||||
"question": "无法连接公司 VPN,首选排查?",
|
||||
"options": ["检查网络是否通", "重装系统", "换电脑", "联系运营商"],
|
||||
"correct_index": 0,
|
||||
"hint": "从外到内排查",
|
||||
"explanation": "先确认能上网,再排查 VPN 客户端,最后才是公司 VPN 服务器",
|
||||
"source": "manual",
|
||||
},
|
||||
{
|
||||
"question": "Outlook 收不到邮件,先看哪里?",
|
||||
"options": ["垃圾邮件箱", "重装 Office", "换邮箱", "打电话给 IT"],
|
||||
"correct_index": 0,
|
||||
"hint": "最容易被忽略的",
|
||||
"explanation": "新邮件被误判到垃圾箱是常见原因,先看再排查服务器",
|
||||
"source": "manual",
|
||||
},
|
||||
{
|
||||
"question": "Office 软件打开慢,先做什么?",
|
||||
"options": ["清理开机启动项", "换电脑", "买新硬盘", "卸载重装"],
|
||||
"correct_index": 0,
|
||||
"hint": "性能问题先减负",
|
||||
"explanation": "开机启动项太多会拖慢所有应用,清理后再观察",
|
||||
"source": "manual",
|
||||
},
|
||||
{
|
||||
"question": "电脑提示磁盘空间不足,应该?",
|
||||
"options": ["清理回收站和临时文件", "关机", "重装系统", "不处理"],
|
||||
"correct_index": 0,
|
||||
"hint": "先释放空间再判断",
|
||||
"explanation": "90% 的情况清理回收站 + temp 目录就能解决,严重才需要重装",
|
||||
"source": "manual",
|
||||
},
|
||||
{
|
||||
"question": "网页打不开,首先排查?",
|
||||
"options": ["检查网络连接", "换浏览器", "重装系统", "砸键盘"],
|
||||
"correct_index": 0,
|
||||
"hint": "从最基础的开始",
|
||||
"explanation": "先看能不能打开其他网页,排除是网站问题还是网络问题",
|
||||
"source": "manual",
|
||||
},
|
||||
{
|
||||
"question": "U 盘插入电脑没反应,先检查?",
|
||||
"options": ["换个 USB 接口", "格式化 U 盘", "扔了", "拆电脑"],
|
||||
"correct_index": 0,
|
||||
"hint": "先排除最简单的问题",
|
||||
"explanation": "USB 接口可能松动或供电不足,先换接口试,不要先动数据",
|
||||
"source": "manual",
|
||||
},
|
||||
{
|
||||
"question": "电脑突然变卡,第一步应该?",
|
||||
"options": ["看任务管理器占用", "砸电脑", "重装系统", "关机睡觉"],
|
||||
"correct_index": 0,
|
||||
"hint": "数据先行",
|
||||
"explanation": "任务管理器能看到 CPU/内存/磁盘占用,定位是哪个进程在吃资源",
|
||||
"source": "manual",
|
||||
},
|
||||
]
|
||||
|
||||
chosen = random.choice(fallback_pool)
|
||||
return chosen
|
||||
|
||||
def add_custom_word(self, word: str) -> None:
|
||||
"""动态添加敏感词(运营后台调用)。"""
|
||||
self.wf.addWords([word])
|
||||
if word not in self.custom_sensitive_words:
|
||||
self.custom_sensitive_words.append(word)
|
||||
|
||||
def remove_custom_word(self, word: str) -> None:
|
||||
"""动态删除敏感词。"""
|
||||
# wordfilter 没有 remove API,降级用 replace 占位
|
||||
# wordfilter.remove(word) # 实际库不一定支持
|
||||
if word in self.custom_sensitive_words:
|
||||
self.custom_sensitive_words.remove(word)
|
||||
|
||||
|
||||
# 单例
|
||||
_moderation_service: Optional[ContentModerationService] = None
|
||||
|
||||
|
||||
def get_moderation_service() -> ContentModerationService:
|
||||
"""获取内容审核服务单例。"""
|
||||
global _moderation_service
|
||||
if _moderation_service is None:
|
||||
_moderation_service = ContentModerationService()
|
||||
return _moderation_service
|
||||
@@ -1,196 +0,0 @@
|
||||
# =============================================================================
|
||||
# 员工目录解析服务
|
||||
# =============================================================================
|
||||
# 说明:角色分配时,将管理员输入的「员工账号 或 姓名」解析为企微 UserID,
|
||||
# 并校验该员工确实属于企微组织架构(需求:分配角色时按姓名/账号自动转换 + 校验)。
|
||||
#
|
||||
# 数据源优先级(自动适配,无需改代码即可在权限开通后升级):
|
||||
# 1. 企微通讯录(实时):
|
||||
# - get_user_info(userid) 校验账号是否为组织内真实员工
|
||||
# - get_department_members(1, 1) 拉取全组织架构,用于「姓名 -> 账号」匹配
|
||||
# - 需要企微应用具备「通讯录读取」权限;权限不足(errcode 60011)时自动降级
|
||||
# 2. 本地 employees 表(仅登录过的员工):作为降级目录,保证功能在缺权限时仍可用
|
||||
#
|
||||
# 设计目标:无论企微权限是否齐全,分配功能都可用;权限齐全时自动获得全公司
|
||||
# 姓名搜索能力(full_directory=True),缺权限时仅覆盖已登录员工。
|
||||
# =============================================================================
|
||||
|
||||
import json
|
||||
import logging
|
||||
from typing import Any, Dict, List, Optional, Tuple
|
||||
|
||||
import redis.asyncio as aioredis
|
||||
from sqlalchemy import select
|
||||
from sqlalchemy.ext.asyncio import AsyncSession
|
||||
|
||||
from app.models.employee import Employee
|
||||
from app.services.wecom_service import WecomService
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
# 组织目录 Redis 缓存 key 与 TTL(10 分钟,避免频繁调用企微通讯录 API)
|
||||
ORG_DIRECTORY_CACHE_KEY = "wecom:org_directory"
|
||||
ORG_DIRECTORY_CACHE_TTL = 600
|
||||
|
||||
|
||||
async def get_org_directory(
|
||||
db: AsyncSession,
|
||||
redis: Optional[aioredis.Redis],
|
||||
) -> Tuple[List[Dict[str, Any]], bool]:
|
||||
"""获取「组织目录」(员工账号 + 姓名 列表),用于姓名 -> 账号匹配。
|
||||
|
||||
优先返回缓存;缓存未命中时尝试从企微通讯录拉全组织(需通讯录读取权限)。
|
||||
若企微权限不足或调用失败,降级到本地 employees 表。
|
||||
|
||||
Returns:
|
||||
(directory, full_directory)
|
||||
- directory: [{"employee_id": str, "name": str, "department": str}, ...]
|
||||
- full_directory: True=来自企微全组织(覆盖全公司);False=仅本地已登录员工
|
||||
"""
|
||||
# 1. 尝试命中缓存(缓存一定来自企微全组织,full=True)
|
||||
if redis:
|
||||
try:
|
||||
raw = await redis.get(ORG_DIRECTORY_CACHE_KEY)
|
||||
if raw:
|
||||
logger.debug("命中组织目录缓存")
|
||||
return json.loads(raw.decode("utf-8")), True
|
||||
except Exception as e:
|
||||
logger.warning(f"读取组织目录缓存失败(降级): {e}")
|
||||
|
||||
# 2. 尝试从企微通讯录拉全组织
|
||||
wecom = WecomService(redis_client=redis)
|
||||
try:
|
||||
members = await wecom.get_department_members(1, 1)
|
||||
directory = [
|
||||
{
|
||||
"employee_id": m.get("userid", ""),
|
||||
"name": m.get("name", "") or "",
|
||||
"department": ",".join(str(d) for d in (m.get("department") or [])),
|
||||
}
|
||||
for m in members
|
||||
if m.get("userid")
|
||||
]
|
||||
# 写入缓存(仅全组织结果缓存,降级结果不缓存以免长期误用)
|
||||
if redis:
|
||||
try:
|
||||
await redis.setex(
|
||||
ORG_DIRECTORY_CACHE_KEY,
|
||||
ORG_DIRECTORY_CACHE_TTL,
|
||||
json.dumps(directory, ensure_ascii=False),
|
||||
)
|
||||
except Exception as e:
|
||||
logger.warning(f"写入组织目录缓存失败: {e}")
|
||||
logger.info(f"组织目录来自企微全组织,共 {len(directory)} 人")
|
||||
return directory, True
|
||||
except Exception as e:
|
||||
err_text = str(e)
|
||||
if "60011" in err_text or "privilege" in err_text.lower():
|
||||
logger.warning("企微通讯录部门读取权限不足,降级到本地 employees 表")
|
||||
else:
|
||||
logger.warning(f"企微通讯录获取失败,降级本地: {err_text}")
|
||||
|
||||
# 3. 降级:本地 employees 表(仅登录过的员工)
|
||||
try:
|
||||
result = await db.execute(
|
||||
select(Employee.employee_id, Employee.name).where(Employee.employee_id != "")
|
||||
)
|
||||
rows = result.all()
|
||||
directory = [
|
||||
{"employee_id": r[0], "name": r[1] or "", "department": ""} for r in rows
|
||||
]
|
||||
logger.info(f"组织目录降级到本地 employees 表,共 {len(directory)} 人")
|
||||
return directory, False
|
||||
except Exception as e:
|
||||
logger.error(f"本地 employees 表查询失败: {e}")
|
||||
return [], False
|
||||
|
||||
|
||||
async def resolve_target(
|
||||
target: str,
|
||||
db: AsyncSession,
|
||||
redis: Optional[aioredis.Redis] = None,
|
||||
) -> Dict[str, Any]:
|
||||
"""将输入的「员工账号 或 姓名」解析为企微 UserID,并校验组织内存在性。
|
||||
|
||||
Returns(结构化结果,由调用方翻译为响应/异常):
|
||||
{"found": True, "employee_id": str, "name": str, "source": str}
|
||||
{"found": False, "reason": str, "suggestion": str}
|
||||
{"ambiguous": True, "candidates": [{"employee_id","name","department"}, ...]}
|
||||
"""
|
||||
target = (target or "").strip()
|
||||
if not target:
|
||||
return {
|
||||
"found": False,
|
||||
"reason": "请输入员工账号或姓名",
|
||||
"suggestion": "请填写企微员工账号或姓名后重试",
|
||||
}
|
||||
|
||||
wecom = WecomService(redis_client=redis)
|
||||
|
||||
# 1) 先尝试按 userid 实时校验(企微组织内真实员工)
|
||||
try:
|
||||
info = await wecom.get_user_info(target)
|
||||
# 成功 -> target 本身就是有效 userid
|
||||
return {
|
||||
"found": True,
|
||||
"employee_id": info.get("userid") or target,
|
||||
"name": info.get("name") or "",
|
||||
"source": "wecom_userid",
|
||||
}
|
||||
except Exception as e:
|
||||
logger.debug(f"get_user_info('{target}') 未命中(将尝试按姓名解析): {e}")
|
||||
|
||||
# 2) 按姓名(包含)解析
|
||||
directory, full = await get_org_directory(db, redis)
|
||||
ql = target.lower()
|
||||
# 精确 userid 匹配优先(目录里可能存在)
|
||||
exact = [m for m in directory if m["employee_id"] and m["employee_id"].lower() == ql]
|
||||
# 姓名包含匹配
|
||||
name_hits = [m for m in directory if m["name"] and ql in m["name"].lower()]
|
||||
matches = exact if exact else name_hits
|
||||
|
||||
if len(matches) == 1:
|
||||
m = matches[0]
|
||||
# 二次实时校验该 userid 确实在组织内(网络可用时)
|
||||
try:
|
||||
info = await wecom.get_user_info(m["employee_id"])
|
||||
return {
|
||||
"found": True,
|
||||
"employee_id": info.get("userid") or m["employee_id"],
|
||||
"name": info.get("name") or m.get("name", ""),
|
||||
"source": "wecom_name" if full else "local_name",
|
||||
}
|
||||
except Exception:
|
||||
# 实时校验失败(网络/权限),但目录里有 -> 仍可用
|
||||
return {
|
||||
"found": True,
|
||||
"employee_id": m["employee_id"],
|
||||
"name": m.get("name", ""),
|
||||
"source": "local_name",
|
||||
}
|
||||
|
||||
if len(matches) > 1:
|
||||
return {
|
||||
"ambiguous": True,
|
||||
"candidates": [
|
||||
{
|
||||
"employee_id": m["employee_id"],
|
||||
"name": m["name"],
|
||||
"department": m.get("department", ""),
|
||||
}
|
||||
for m in matches[:10]
|
||||
],
|
||||
}
|
||||
|
||||
# 未找到
|
||||
if full:
|
||||
return {
|
||||
"found": False,
|
||||
"reason": f"企微组织架构中未找到匹配「{target}」的员工",
|
||||
"suggestion": "请确认姓名/账号拼写,或改为输入员工账号",
|
||||
}
|
||||
return {
|
||||
"found": False,
|
||||
"reason": f"未找到匹配「{target}」的员工",
|
||||
"suggestion": "当前仅能按姓名搜索已登录过本系统的员工;请直接输入员工账号,或为企微应用开通「通讯录读取」权限以搜索全公司",
|
||||
}
|
||||
@@ -1,59 +0,0 @@
|
||||
# =============================================================================
|
||||
# 企微IT智能服务台 — 超时提醒服务
|
||||
# =============================================================================
|
||||
# 说明:提供发送超时提醒企微消息的功能
|
||||
# 在坐席回复后员工长时间未回复时,发送企微提醒消息
|
||||
# =============================================================================
|
||||
|
||||
import logging
|
||||
from app.config import settings
|
||||
from app.services.wecom_service import WecomService
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
# 超时提醒消息内容
|
||||
REMINDER_MESSAGE = (
|
||||
"IT服务提醒:您有新的消息未查看,"
|
||||
"咨询将在10分钟后标记为待关闭,请尽快点击处理 👉 "
|
||||
"https://itsupport.servyou.com.cn/itdesk/"
|
||||
)
|
||||
|
||||
|
||||
async def send_reminder_message(employee_id: str) -> bool:
|
||||
"""发送超时提醒企微消息。
|
||||
|
||||
当坐席回复后员工超过3分钟未回复时,发送企微消息提醒员工查看。
|
||||
|
||||
Args:
|
||||
employee_id: 员工的企微 UserID
|
||||
|
||||
Returns:
|
||||
bool: 发送是否成功
|
||||
|
||||
Raises:
|
||||
Exception: 发送失败时抛出异常
|
||||
"""
|
||||
try:
|
||||
redis_client = settings.create_redis_client()
|
||||
wecom_service = WecomService(redis_client)
|
||||
|
||||
try:
|
||||
result = await wecom_service.send_text_message(
|
||||
employee_id, REMINDER_MESSAGE
|
||||
)
|
||||
if result.get("errcode") == 0:
|
||||
logger.info(f"超时提醒发送成功: employee_id={employee_id}")
|
||||
return True
|
||||
else:
|
||||
logger.warning(
|
||||
f"超时提醒发送失败: employee_id={employee_id}, "
|
||||
f"errcode={result.get('errcode')}, errmsg={result.get('errmsg')}"
|
||||
)
|
||||
return False
|
||||
finally:
|
||||
await wecom_service.close()
|
||||
await redis_client.close()
|
||||
|
||||
except Exception as e:
|
||||
logger.error(f"发送超时提醒异常: employee_id={employee_id}, error={e}")
|
||||
raise
|
||||
@@ -1,668 +0,0 @@
|
||||
# =============================================================================
|
||||
# 企微IT智能服务台 — 企微 API 封装服务
|
||||
# =============================================================================
|
||||
# 说明:封装所有与企微服务器的交互逻辑,包括:
|
||||
# 1. access_token 管理(Redis 缓存 + 自动刷新)
|
||||
# 2. 发送消息(文本/图片/文件)
|
||||
# 3. 获取员工信息(通讯录 API)
|
||||
# 4. 上传临时素材
|
||||
# 5. OAuth2 授权换算用户身份
|
||||
# =============================================================================
|
||||
|
||||
import json
|
||||
import logging
|
||||
from typing import Any, Dict, List, Optional
|
||||
|
||||
import httpx
|
||||
import redis.asyncio as aioredis
|
||||
|
||||
from app.config import settings
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
|
||||
class WecomService:
|
||||
"""企微 API 调用服务。
|
||||
|
||||
封装所有与企微服务器的 HTTP 交互,提供异步方法。
|
||||
access_token 通过 Redis 缓存管理,避免频繁调用获取接口。
|
||||
|
||||
Attributes:
|
||||
redis: Redis 异步客户端(用于缓存 access_token)
|
||||
client: httpx 异步 HTTP 客户端
|
||||
"""
|
||||
|
||||
def __init__(self, redis_client: Optional[aioredis.Redis] = None):
|
||||
"""初始化企微服务。
|
||||
|
||||
Args:
|
||||
redis_client: Redis 异步客户端实例(可为 None,本地开发时 Redis 不可用)
|
||||
"""
|
||||
self.redis = redis_client
|
||||
# 创建 httpx 异步客户端
|
||||
# timeout: 连接超时5秒,读取超时10秒
|
||||
self.client = httpx.AsyncClient(
|
||||
timeout=httpx.Timeout(connect=5.0, read=10.0, write=10.0, pool=5.0)
|
||||
)
|
||||
# 内存缓存(Redis 不可用时的降级方案)
|
||||
self._token_cache: Optional[str] = None
|
||||
|
||||
# --------------------------------------------------------------------------
|
||||
# access_token 管理
|
||||
# --------------------------------------------------------------------------
|
||||
async def get_access_token(self) -> str:
|
||||
"""获取企微 access_token。
|
||||
|
||||
优先从 Redis 缓存获取,如果缓存不存在或即将过期则重新获取。
|
||||
access_token 有效期 7200 秒,缓存 TTL 设为 6900 秒(提前 300 秒刷新)。
|
||||
|
||||
对应企微API:
|
||||
GET https://qyapi.weixin.qq.com/cgi-bin/gettoken?corpid=ID&corpsecret=SECRET
|
||||
|
||||
Returns:
|
||||
str: access_token 字符串
|
||||
|
||||
Raises:
|
||||
Exception: 获取 access_token 失败
|
||||
"""
|
||||
# Redis 缓存 key
|
||||
cache_key = "wecom:access_token"
|
||||
|
||||
# 1. 尝试从 Redis 缓存获取
|
||||
if self.redis:
|
||||
try:
|
||||
cached_token = await self.redis.get(cache_key)
|
||||
if cached_token:
|
||||
logger.debug("从缓存获取 access_token")
|
||||
return cached_token.decode("utf-8")
|
||||
except Exception as e:
|
||||
logger.warning(f"Redis 读取失败(降级): {e}")
|
||||
|
||||
# 1b. 尝试从内存缓存获取
|
||||
if self._token_cache:
|
||||
logger.debug("从内存缓存获取 access_token")
|
||||
return self._token_cache
|
||||
|
||||
# 2. 缓存未命中,调用企微 API 获取
|
||||
logger.info("缓存未命中,调用企微API获取 access_token")
|
||||
url = "https://qyapi.weixin.qq.com/cgi-bin/gettoken"
|
||||
params = {
|
||||
"corpid": settings.wecom_corp_id,
|
||||
"corpsecret": settings.wecom_secret,
|
||||
}
|
||||
|
||||
try:
|
||||
response = await self.client.get(url, params=params)
|
||||
result = response.json()
|
||||
|
||||
# 检查企微API返回码
|
||||
if result.get("errcode") != 0:
|
||||
error_msg = result.get("errmsg", "未知错误")
|
||||
logger.error(f"获取 access_token 失败: errcode={result.get('errcode')}, errmsg={error_msg}")
|
||||
raise Exception(f"企微API错误: {error_msg}")
|
||||
|
||||
access_token = result["access_token"]
|
||||
expires_in = result.get("expires_in", 7200)
|
||||
|
||||
# 3. 缓存到 Redis,TTL = 有效期 - 300秒(提前刷新)
|
||||
buffer_seconds = 300
|
||||
cache_ttl = max(expires_in - buffer_seconds, 60) # 至少缓存 60 秒
|
||||
if self.redis:
|
||||
try:
|
||||
await self.redis.setex(cache_key, cache_ttl, access_token)
|
||||
except Exception as e:
|
||||
logger.warning(f"Redis 写入失败(降级): {e}")
|
||||
|
||||
# 3b. 同时缓存到内存
|
||||
self._token_cache = access_token
|
||||
|
||||
logger.info(f"access_token 获取成功,缓存 TTL={cache_ttl}秒")
|
||||
return access_token
|
||||
|
||||
except httpx.HTTPError as e:
|
||||
logger.error(f"获取 access_token 网络错误: {e}")
|
||||
raise Exception(f"企微API网络错误: {e}") from e
|
||||
|
||||
# --------------------------------------------------------------------------
|
||||
# 发送文本消息
|
||||
# --------------------------------------------------------------------------
|
||||
async def send_text_message(
|
||||
self, user_id: str, content: str
|
||||
) -> Dict[str, Any]:
|
||||
"""向员工发送文本消息。
|
||||
|
||||
对应企微API:
|
||||
POST https://qyapi.weixin.qq.com/cgi-bin/message/send?access_token=TOKEN
|
||||
|
||||
请求体:
|
||||
{
|
||||
"touser": "UserID",
|
||||
"msgtype": "text",
|
||||
"agentid": 1000002,
|
||||
"text": {"content": "消息内容"}
|
||||
}
|
||||
|
||||
Args:
|
||||
user_id: 员工的企微 UserID
|
||||
content: 消息内容(纯文本)
|
||||
|
||||
Returns:
|
||||
Dict[str, Any]: 企微API返回结果
|
||||
"""
|
||||
access_token = await self.get_access_token()
|
||||
url = f"https://qyapi.weixin.qq.com/cgi-bin/message/send?access_token={access_token}"
|
||||
|
||||
payload = {
|
||||
"touser": user_id,
|
||||
"msgtype": "text",
|
||||
"agentid": int(settings.wecom_agent_id),
|
||||
"text": {"content": content},
|
||||
}
|
||||
|
||||
try:
|
||||
response = await self.client.post(url, json=payload)
|
||||
result = response.json()
|
||||
|
||||
if result.get("errcode") != 0:
|
||||
logger.error(
|
||||
f"发送文本消息失败: user_id={user_id}, "
|
||||
f"errcode={result.get('errcode')}, errmsg={result.get('errmsg')}"
|
||||
)
|
||||
else:
|
||||
logger.info(f"发送文本消息成功: user_id={user_id}")
|
||||
|
||||
return result
|
||||
|
||||
except httpx.HTTPError as e:
|
||||
logger.error(f"发送文本消息网络错误: user_id={user_id}, error={e}")
|
||||
raise Exception(f"发送消息网络错误: {e}") from e
|
||||
|
||||
# --------------------------------------------------------------------------
|
||||
# 发送卡片消息
|
||||
# --------------------------------------------------------------------------
|
||||
async def send_card_message(
|
||||
self,
|
||||
user_id: str,
|
||||
title: str,
|
||||
description: str,
|
||||
url: str = "",
|
||||
btntxt: str = "详情",
|
||||
) -> Dict[str, Any]:
|
||||
"""向员工发送文本卡片消息。
|
||||
|
||||
对应企微API:
|
||||
POST https://qyapi.weixin.qq.com/cgi-bin/message/send?access_token=TOKEN
|
||||
|
||||
请求体:
|
||||
{
|
||||
"touser": "UserID",
|
||||
"msgtype": "textcard",
|
||||
"agentid": 1000002,
|
||||
"textcard": {
|
||||
"title": "标题",
|
||||
"description": "描述",
|
||||
"url": "链接",
|
||||
"btntxt": "按钮文字"
|
||||
}
|
||||
}
|
||||
|
||||
Args:
|
||||
user_id: 员工的企微 UserID
|
||||
title: 卡片标题
|
||||
description: 卡片描述
|
||||
url: 卡片点击跳转链接
|
||||
btntxt: 按钮文字(默认"详情")
|
||||
|
||||
Returns:
|
||||
Dict[str, Any]: 企微API返回结果
|
||||
"""
|
||||
access_token = await self.get_access_token()
|
||||
url_api = f"https://qyapi.weixin.qq.com/cgi-bin/message/send?access_token={access_token}"
|
||||
|
||||
payload = {
|
||||
"touser": user_id,
|
||||
"msgtype": "textcard",
|
||||
"agentid": int(settings.wecom_agent_id),
|
||||
"textcard": {
|
||||
"title": title,
|
||||
"description": description,
|
||||
"url": url,
|
||||
"btntxt": btntxt,
|
||||
},
|
||||
}
|
||||
|
||||
try:
|
||||
response = await self.client.post(url_api, json=payload)
|
||||
result = response.json()
|
||||
|
||||
if result.get("errcode") != 0:
|
||||
logger.error(
|
||||
f"发送卡片消息失败: user_id={user_id}, "
|
||||
f"errcode={result.get('errcode')}, errmsg={result.get('errmsg')}"
|
||||
)
|
||||
else:
|
||||
logger.info(f"发送卡片消息成功: user_id={user_id}")
|
||||
|
||||
return result
|
||||
|
||||
except httpx.HTTPError as e:
|
||||
logger.error(f"发送卡片消息网络错误: user_id={user_id}, error={e}")
|
||||
raise Exception(f"发送消息网络错误: {e}") from e
|
||||
|
||||
# --------------------------------------------------------------------------
|
||||
# 发送图片消息
|
||||
# --------------------------------------------------------------------------
|
||||
async def send_image_message(
|
||||
self, user_id: str, media_id: str
|
||||
) -> Dict[str, Any]:
|
||||
"""向员工发送图片消息。
|
||||
|
||||
对应企微API:
|
||||
POST https://qyapi.weixin.qq.com/cgi-bin/message/send?access_token=TOKEN
|
||||
|
||||
请求体:
|
||||
{
|
||||
"touser": "UserID",
|
||||
"msgtype": "image",
|
||||
"agentid": 1000002,
|
||||
"image": {"media_id": "MEDIA_ID"}
|
||||
}
|
||||
|
||||
注意:发送图片前需要先通过 upload_temp_media 上传图片获取 media_id。
|
||||
|
||||
Args:
|
||||
user_id: 员工的企微 UserID
|
||||
media_id: 图片媒体ID(通过上传临时素材获取)
|
||||
|
||||
Returns:
|
||||
Dict[str, Any]: 企微API返回结果
|
||||
"""
|
||||
access_token = await self.get_access_token()
|
||||
url = f"https://qyapi.weixin.qq.com/cgi-bin/message/send?access_token={access_token}"
|
||||
|
||||
payload = {
|
||||
"touser": user_id,
|
||||
"msgtype": "image",
|
||||
"agentid": int(settings.wecom_agent_id),
|
||||
"image": {"media_id": media_id},
|
||||
}
|
||||
|
||||
try:
|
||||
response = await self.client.post(url, json=payload)
|
||||
result = response.json()
|
||||
|
||||
if result.get("errcode") != 0:
|
||||
logger.error(
|
||||
f"发送图片消息失败: user_id={user_id}, "
|
||||
f"errcode={result.get('errcode')}, errmsg={result.get('errmsg')}"
|
||||
)
|
||||
else:
|
||||
logger.info(f"发送图片消息成功: user_id={user_id}")
|
||||
|
||||
return result
|
||||
|
||||
except httpx.HTTPError as e:
|
||||
logger.error(f"发送图片消息网络错误: user_id={user_id}, error={e}")
|
||||
raise Exception(f"发送消息网络错误: {e}") from e
|
||||
|
||||
# --------------------------------------------------------------------------
|
||||
# 发送文件消息
|
||||
# --------------------------------------------------------------------------
|
||||
async def send_file_message(
|
||||
self, user_id: str, media_id: str
|
||||
) -> Dict[str, Any]:
|
||||
"""向员工发送文件消息。
|
||||
|
||||
对应企微API:
|
||||
POST https://qyapi.weixin.qq.com/cgi-bin/message/send?access_token=TOKEN
|
||||
|
||||
请求体:
|
||||
{
|
||||
"touser": "UserID",
|
||||
"msgtype": "file",
|
||||
"agentid": 1000002,
|
||||
"file": {"media_id": "MEDIA_ID"}
|
||||
}
|
||||
|
||||
注意:发送文件前需要先通过 upload_temp_media 上传文件获取 media_id。
|
||||
|
||||
Args:
|
||||
user_id: 员工的企微 UserID
|
||||
media_id: 文件媒体ID(通过上传临时素材获取)
|
||||
|
||||
Returns:
|
||||
Dict[str, Any]: 企微API返回结果
|
||||
"""
|
||||
access_token = await self.get_access_token()
|
||||
url = f"https://qyapi.weixin.qq.com/cgi-bin/message/send?access_token={access_token}"
|
||||
|
||||
payload = {
|
||||
"touser": user_id,
|
||||
"msgtype": "file",
|
||||
"agentid": int(settings.wecom_agent_id),
|
||||
"file": {"media_id": media_id},
|
||||
}
|
||||
|
||||
try:
|
||||
response = await self.client.post(url, json=payload)
|
||||
result = response.json()
|
||||
|
||||
if result.get("errcode") != 0:
|
||||
logger.error(
|
||||
f"发送文件消息失败: user_id={user_id}, "
|
||||
f"errcode={result.get('errcode')}, errmsg={result.get('errmsg')}"
|
||||
)
|
||||
else:
|
||||
logger.info(f"发送文件消息成功: user_id={user_id}")
|
||||
|
||||
return result
|
||||
|
||||
except httpx.HTTPError as e:
|
||||
logger.error(f"发送文件消息网络错误: user_id={user_id}, error={e}")
|
||||
raise Exception(f"发送消息网络错误: {e}") from e
|
||||
|
||||
# --------------------------------------------------------------------------
|
||||
# 获取员工通讯录信息
|
||||
# --------------------------------------------------------------------------
|
||||
async def get_user_info(self, user_id: str) -> Dict[str, Any]:
|
||||
"""获取员工通讯录详细信息(用于 VIP 判断)。
|
||||
|
||||
对应企微API:
|
||||
GET https://qyapi.weixin.qq.com/cgi-bin/user/get?access_token=TOKEN&userid=USERID
|
||||
|
||||
返回数据包含:
|
||||
- userid: 员工UserID
|
||||
- name: 员工姓名
|
||||
- department: 部门ID列表
|
||||
- position: 岗位
|
||||
- mobile: 手机号
|
||||
- email: 邮箱
|
||||
- status: 激活状态
|
||||
|
||||
需要企微通讯录只读权限。
|
||||
|
||||
Args:
|
||||
user_id: 员工的企微 UserID
|
||||
|
||||
Returns:
|
||||
Dict[str, Any]: 员工信息字典
|
||||
|
||||
Raises:
|
||||
Exception: 获取失败
|
||||
"""
|
||||
access_token = await self.get_access_token()
|
||||
url = "https://qyapi.weixin.qq.com/cgi-bin/user/get"
|
||||
params = {
|
||||
"access_token": access_token,
|
||||
"userid": user_id,
|
||||
}
|
||||
|
||||
try:
|
||||
response = await self.client.get(url, params=params)
|
||||
result = response.json()
|
||||
|
||||
if result.get("errcode", 0) != 0:
|
||||
logger.error(
|
||||
f"获取员工信息失败: user_id={user_id}, "
|
||||
f"errcode={result.get('errcode')}, errmsg={result.get('errmsg')}"
|
||||
)
|
||||
raise Exception(f"获取员工信息失败: {result.get('errmsg')}")
|
||||
|
||||
logger.info(f"获取员工信息成功: user_id={user_id}, name={result.get('name', '')}")
|
||||
return result
|
||||
|
||||
except httpx.HTTPError as e:
|
||||
logger.error(f"获取员工信息网络错误: user_id={user_id}, error={e}")
|
||||
raise Exception(f"获取员工信息网络错误: {e}") from e
|
||||
|
||||
# --------------------------------------------------------------------------
|
||||
# 获取部门成员列表
|
||||
# --------------------------------------------------------------------------
|
||||
async def get_department_members(
|
||||
self, department_id: int = 1, fetch_child: int = 1
|
||||
) -> List[Dict[str, Any]]:
|
||||
"""获取部门成员列表。
|
||||
|
||||
对应企微API:
|
||||
GET https://qyapi.weixin.qq.com/cgi-bin/user/list?access_token=TOKEN&department_id=ID&fetch_child=1
|
||||
|
||||
Args:
|
||||
department_id: 部门ID(默认1为根部门)
|
||||
fetch_child: 是否递归获取子部门(1=是, 0=否)
|
||||
|
||||
Returns:
|
||||
List[Dict[str, Any]]: 部门成员列表
|
||||
|
||||
Raises:
|
||||
Exception: 获取失败
|
||||
"""
|
||||
access_token = await self.get_access_token()
|
||||
url = "https://qyapi.weixin.qq.com/cgi-bin/user/list"
|
||||
params = {
|
||||
"access_token": access_token,
|
||||
"department_id": department_id,
|
||||
"fetch_child": fetch_child,
|
||||
}
|
||||
|
||||
try:
|
||||
response = await self.client.get(url, params=params)
|
||||
result = response.json()
|
||||
|
||||
if result.get("errcode", 0) != 0:
|
||||
logger.error(
|
||||
f"获取部门成员失败: dept_id={department_id}, "
|
||||
f"errcode={result.get('errcode')}, errmsg={result.get('errmsg')}"
|
||||
)
|
||||
raise Exception(f"获取部门成员失败: {result.get('errmsg')}")
|
||||
|
||||
userlist = result.get("userlist", [])
|
||||
logger.info(f"获取部门成员成功: dept_id={department_id}, count={len(userlist)}")
|
||||
return userlist
|
||||
|
||||
except httpx.HTTPError as e:
|
||||
logger.error(f"获取部门成员网络错误: dept_id={department_id}, error={e}")
|
||||
raise Exception(f"获取部门成员网络错误: {e}") from e
|
||||
|
||||
# --------------------------------------------------------------------------
|
||||
# JS-SDK 票据 (v0.5.4:应急页身份检测用)
|
||||
# --------------------------------------------------------------------------
|
||||
async def get_jsapi_ticket(self) -> str:
|
||||
"""获取企微 JS-SDK 票据 jsapi_ticket。
|
||||
|
||||
对应企微API:
|
||||
GET https://qyapi.weixin.qq.com/cgi-bin/get_jsapi_ticket?access_token=TOKEN
|
||||
|
||||
jsapi_ticket 用于计算 JS-SDK 签名(sha1),让前端 wx.config/wx.agentConfig 鉴权通过。
|
||||
有效期 7200 秒,缓存到 Redis(提前 300 秒刷新)。
|
||||
|
||||
Returns:
|
||||
str: jsapi_ticket 字符串
|
||||
|
||||
Raises:
|
||||
Exception: 获取失败
|
||||
"""
|
||||
cache_key = "wecom:jsapi_ticket"
|
||||
|
||||
# 1. Redis 缓存
|
||||
if self.redis:
|
||||
try:
|
||||
cached = await self.redis.get(cache_key)
|
||||
if cached:
|
||||
logger.debug("从缓存获取 jsapi_ticket")
|
||||
return cached.decode("utf-8")
|
||||
except Exception as e:
|
||||
logger.warning(f"Redis 读取 jsapi_ticket 失败(降级): {e}")
|
||||
|
||||
# 2. 调用企微 API
|
||||
access_token = await self.get_access_token()
|
||||
url = f"https://qyapi.weixin.qq.com/cgi-bin/get_jsapi_ticket?access_token={access_token}"
|
||||
|
||||
try:
|
||||
response = await self.client.get(url)
|
||||
result = response.json()
|
||||
|
||||
if result.get("errcode", 0) != 0:
|
||||
logger.error(
|
||||
f"获取 jsapi_ticket 失败: "
|
||||
f"errcode={result.get('errcode')}, errmsg={result.get('errmsg')}"
|
||||
)
|
||||
raise Exception(f"获取 jsapi_ticket 失败: {result.get('errmsg')}")
|
||||
|
||||
ticket = result.get("ticket", "")
|
||||
expires_in = result.get("expires_in", 7200)
|
||||
|
||||
# 3. 缓存到 Redis(TTL = expires_in - 300s)
|
||||
cache_ttl = max(expires_in - 300, 60)
|
||||
if self.redis:
|
||||
try:
|
||||
await self.redis.setex(cache_key, cache_ttl, ticket)
|
||||
except Exception as e:
|
||||
logger.warning(f"Redis 写入 jsapi_ticket 失败(降级): {e}")
|
||||
|
||||
logger.info(f"jsapi_ticket 获取成功,缓存 TTL={cache_ttl}秒")
|
||||
return ticket
|
||||
|
||||
except httpx.HTTPError as e:
|
||||
logger.error(f"获取 jsapi_ticket 网络错误: {e}")
|
||||
raise Exception(f"企微API网络错误: {e}") from e
|
||||
|
||||
@staticmethod
|
||||
def generate_jsapi_signature(
|
||||
ticket: str, nonce_str: str, timestamp: int, url: str
|
||||
) -> str:
|
||||
"""生成 JS-SDK 签名(sha1)。
|
||||
|
||||
对应企微JS-SDK签名算法:
|
||||
1. 拼接:jsapi_ticket={ticket}&noncestr={nonce_str}×tamp={timestamp}&url={url}
|
||||
2. sha1(拼接字符串)
|
||||
|
||||
注意:
|
||||
- url 不含 # 及其后面部分
|
||||
- url 不含 ?
|
||||
- url 是前端调用 wx.config 的页面 URL
|
||||
|
||||
Args:
|
||||
ticket: jsapi_ticket
|
||||
nonce_str: 随机字符串(前端生成,16位)
|
||||
timestamp: 当前时间戳(秒)
|
||||
url: 当前页面 URL(不含 # 后面)
|
||||
|
||||
Returns:
|
||||
str: sha1 签名字符串(40 字符)
|
||||
"""
|
||||
import hashlib
|
||||
|
||||
# 拼接签名字符串
|
||||
raw = f"jsapi_ticket={ticket}&noncestr={nonce_str}×tamp={timestamp}&url={url}"
|
||||
# sha1 哈希
|
||||
signature = hashlib.sha1(raw.encode("utf-8")).hexdigest()
|
||||
return signature
|
||||
|
||||
# --------------------------------------------------------------------------
|
||||
# 上传临时素材
|
||||
# --------------------------------------------------------------------------
|
||||
async def upload_temp_media(
|
||||
self, media_type: str, file_data: bytes, filename: str = "upload"
|
||||
) -> str:
|
||||
"""上传临时素材(图片/文件/语音),获取 media_id。
|
||||
|
||||
对应企微API:
|
||||
POST https://qyapi.weixin.qq.com/cgi-bin/media/upload?access_token=TOKEN&type=TYPE
|
||||
|
||||
临时素材有效期 3 天,适用于发送图片/文件消息。
|
||||
|
||||
Args:
|
||||
media_type: 媒体类型(image/file/voice)
|
||||
file_data: 文件二进制数据
|
||||
filename: 文件名
|
||||
|
||||
Returns:
|
||||
str: media_id(用于发送图片/文件消息时引用)
|
||||
|
||||
Raises:
|
||||
Exception: 上传失败
|
||||
"""
|
||||
access_token = await self.get_access_token()
|
||||
url = f"https://qyapi.weixin.qq.com/cgi-bin/media/upload?access_token={access_token}&type={media_type}"
|
||||
|
||||
try:
|
||||
# 使用 multipart 上传文件
|
||||
files = {"media": (filename, file_data)}
|
||||
response = await self.client.post(url, files=files)
|
||||
result = response.json()
|
||||
|
||||
if result.get("errcode", 0) != 0:
|
||||
logger.error(
|
||||
f"上传临时素材失败: type={media_type}, "
|
||||
f"errcode={result.get('errcode')}, errmsg={result.get('errmsg')}"
|
||||
)
|
||||
raise Exception(f"上传临时素材失败: {result.get('errmsg')}")
|
||||
|
||||
media_id = result.get("media_id", "")
|
||||
logger.info(f"上传临时素材成功: type={media_type}, media_id={media_id}")
|
||||
return media_id
|
||||
|
||||
except httpx.HTTPError as e:
|
||||
logger.error(f"上传临时素材网络错误: type={media_type}, error={e}")
|
||||
raise Exception(f"上传临时素材网络错误: {e}") from e
|
||||
|
||||
# --------------------------------------------------------------------------
|
||||
# OAuth2 授权换算用户身份
|
||||
# --------------------------------------------------------------------------
|
||||
async def get_oauth_user_info(self, code: str) -> Dict[str, str]:
|
||||
"""通过 OAuth2 授权码换取员工身份信息。
|
||||
|
||||
对应企微API:
|
||||
GET https://qyapi.weixin.qq.com/cgi-bin/auth/getuserinfo?access_token=TOKEN&code=CODE
|
||||
|
||||
H5 页面通过企微 OAuth2 静默授权获取 code,后端用 code 换取员工 UserID。
|
||||
适用于 H5 用户端身份识别。
|
||||
|
||||
Args:
|
||||
code: 企微 OAuth2 授权码
|
||||
|
||||
Returns:
|
||||
Dict[str, str]: 包含 userid 和 user_ticket 的字典
|
||||
|
||||
Raises:
|
||||
Exception: 换取失败
|
||||
"""
|
||||
access_token = await self.get_access_token()
|
||||
url = "https://qyapi.weixin.qq.com/cgi-bin/auth/getuserinfo"
|
||||
params = {
|
||||
"access_token": access_token,
|
||||
"code": code,
|
||||
}
|
||||
|
||||
try:
|
||||
response = await self.client.get(url, params=params)
|
||||
result = response.json()
|
||||
|
||||
if result.get("errcode", 0) != 0:
|
||||
logger.error(
|
||||
f"OAuth2换取用户身份失败: code={code}, "
|
||||
f"errcode={result.get('errcode')}, errmsg={result.get('errmsg')}"
|
||||
)
|
||||
raise Exception(f"OAuth2换取用户身份失败: {result.get('errmsg')}")
|
||||
|
||||
user_id = result.get("userid", "")
|
||||
logger.info(f"OAuth2换取用户身份成功: userid={user_id}")
|
||||
return {
|
||||
"userid": user_id,
|
||||
"user_ticket": result.get("user_ticket", ""),
|
||||
}
|
||||
|
||||
except httpx.HTTPError as e:
|
||||
logger.error(f"OAuth2换取用户身份网络错误: code={code}, error={e}")
|
||||
raise Exception(f"OAuth2换取用户身份网络错误: {e}") from e
|
||||
|
||||
# --------------------------------------------------------------------------
|
||||
# 关闭客户端
|
||||
# --------------------------------------------------------------------------
|
||||
async def close(self) -> None:
|
||||
"""关闭 HTTP 客户端连接池。
|
||||
|
||||
应用关闭时调用,释放资源。
|
||||
"""
|
||||
await self.client.aclose()
|
||||
logger.info("WecomService HTTP 客户端已关闭")
|
||||
@@ -1,206 +0,0 @@
|
||||
# =============================================================================
|
||||
# 企微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 logging
|
||||
from datetime import datetime
|
||||
|
||||
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.ws_manager import manager as ws_manager
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
|
||||
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 process_h5_ai_reply(
|
||||
conversation_id: str,
|
||||
employee_id: str,
|
||||
content: str,
|
||||
dify_conversation_id=None,
|
||||
):
|
||||
"""H5 发送消息后的 AI 回复处理(asyncio.create_task 入口)。
|
||||
|
||||
流程:
|
||||
- 本地快判断(打招呼 / 呼叫人工)→ 同步结果,整段推送(不调 Dify)
|
||||
- 否则流式调 Dify,逐 chunk 推 ai_reply_chunk,流结束推 ai_reply 终态
|
||||
- 任意异常 → 推 ai_reply_failed,不阻塞用户
|
||||
"""
|
||||
ai_handler = get_shared_ai_handler()
|
||||
factory = _get_session_factory()
|
||||
async with factory() as db:
|
||||
try:
|
||||
conversation = await db.get(Conversation, conversation_id)
|
||||
if not conversation:
|
||||
logger.warning(f"后台 AI 任务:会话不存在 {conversation_id}")
|
||||
return
|
||||
|
||||
is_guidance = False
|
||||
should_count = False
|
||||
should_transfer = False
|
||||
new_dify_conv_id = dify_conversation_id
|
||||
full_parts: list = []
|
||||
|
||||
# 本地快判断(不打 Dify):打招呼 / 呼叫人工 → 同步路径
|
||||
if ai_handler.is_greeting(content) or ai_handler.is_call_human(content):
|
||||
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
|
||||
|
||||
# 流式调 Dify(get_reply_stream 内部已处理真 SSE / 非流式 fallback)
|
||||
# 注意:首参是 message(用户文本),不是 content
|
||||
async for chunk in ai_handler.ai_service.get_reply_stream(
|
||||
message=content,
|
||||
conversation_id=dify_conversation_id,
|
||||
user_id=employee_id,
|
||||
):
|
||||
delta = chunk.get("delta", "")
|
||||
if delta:
|
||||
full_parts.append(delta)
|
||||
await ws_manager.broadcast_to_employees([employee_id], {
|
||||
"type": "ai_reply_chunk",
|
||||
"data": {
|
||||
"conversation_id": conversation_id,
|
||||
"chunk": delta,
|
||||
},
|
||||
})
|
||||
if chunk.get("finished"):
|
||||
new_dify_conv_id = chunk.get("conversation_id") or dify_conversation_id
|
||||
hit = chunk.get("hit")
|
||||
# 命中 → 计数;未命中 → 转人工
|
||||
should_count = bool(hit)
|
||||
should_transfer = not bool(hit)
|
||||
|
||||
content_ai = "".join(full_parts)
|
||||
if not content_ai:
|
||||
# 流式无内容(极端情况),给降级提示,不转人工
|
||||
content_ai = "⚠️ AI 暂时没有返回内容,请输入「IT」转人工。"
|
||||
should_count = False
|
||||
should_transfer = False
|
||||
await _persist_and_push(
|
||||
db, conversation, employee_id, content_ai,
|
||||
is_guidance, should_count, should_transfer, new_dify_conv_id,
|
||||
)
|
||||
except Exception as e:
|
||||
logger.error(f"后台 AI 任务异常: {e}", exc_info=True)
|
||||
try:
|
||||
await ws_manager.broadcast_to_employees([employee_id], {
|
||||
"type": "ai_reply_failed",
|
||||
"data": {
|
||||
"conversation_id": conversation_id,
|
||||
"message": "⚠️ AI 服务异常,请输入「IT」转人工或稍后重试。",
|
||||
},
|
||||
})
|
||||
except Exception:
|
||||
# 推送失败也无所谓,员工端 3 秒轮询兜底
|
||||
pass
|
||||
@@ -1,100 +0,0 @@
|
||||
# =============================================================================
|
||||
# 企微IT智能服务台 — 超时提醒定时任务
|
||||
# =============================================================================
|
||||
# 说明:定时检查超时未回复的会话,发送企微提醒消息
|
||||
# 运行频率:每 30 秒执行一次
|
||||
# 超时逻辑:
|
||||
# 1. 坐席回复后 3 分钟员工未回复 -> 发送企微提醒(只发 1 次)
|
||||
# 2. 坐席回复后 10 分钟员工未回复 -> 标记会话为 pending_close
|
||||
# =============================================================================
|
||||
|
||||
import logging
|
||||
from datetime import datetime, timedelta
|
||||
|
||||
from sqlalchemy import select
|
||||
|
||||
from app.models.conversation import Conversation
|
||||
from app.services.reminder_service import send_reminder_message
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
# 超时配置(分钟)
|
||||
REMINDER_TIMEOUT_MINUTES = 3 # 未回复超时时间
|
||||
CLOSE_TIMEOUT_MINUTES = 10 # 自动待关闭时间
|
||||
|
||||
|
||||
async def check_unreplied_sessions():
|
||||
"""检查超时未回复的会话,发送提醒并标记待关闭。
|
||||
|
||||
此函数由 APScheduler 定时调用(每 30 秒)。
|
||||
执行流程:
|
||||
1. 查找需要发送提醒的会话(坐席回复超过3分钟,员工未回复且未发送过提醒)
|
||||
2. 发送企微提醒消息
|
||||
3. 标记已发送提醒
|
||||
4. 查找需要标记待关闭的会话(坐席回复超过10分钟)
|
||||
5. 更新会话状态为 pending_close
|
||||
"""
|
||||
# 导入数据库 session 工厂
|
||||
from app.database import _get_session_factory
|
||||
|
||||
async_session_factory = _get_session_factory()
|
||||
|
||||
async with async_session_factory() as db:
|
||||
try:
|
||||
# 1. 查找需要发送提醒的会话
|
||||
# 条件:active 状态 + 有坐席回复 + 超过3分钟未回复 + 未发送过提醒
|
||||
reminder_threshold = datetime.now() - timedelta(minutes=REMINDER_TIMEOUT_MINUTES)
|
||||
|
||||
reminder_stmt = select(Conversation).where(
|
||||
Conversation.status == "serving",
|
||||
Conversation.last_agent_reply_at.isnot(None),
|
||||
Conversation.last_agent_reply_at < reminder_threshold,
|
||||
Conversation.reminder_sent == False,
|
||||
)
|
||||
result = await db.execute(reminder_stmt)
|
||||
sessions_to_remind = result.scalars().all()
|
||||
|
||||
logger.info(f"发现 {len(sessions_to_remind)} 个需要发送提醒的会话")
|
||||
|
||||
# 2. 发送企微提醒消息
|
||||
for session in sessions_to_remind:
|
||||
try:
|
||||
success = await send_reminder_message(session.employee_id)
|
||||
if success:
|
||||
# 3. 标记已发送提醒
|
||||
session.reminder_sent = True
|
||||
session.reminder_sent_at = datetime.now()
|
||||
logger.info(f"会话 {session.id} 已发送提醒: employee_id={session.employee_id}")
|
||||
else:
|
||||
logger.warning(f"会话 {session.id} 提醒发送失败,跳过")
|
||||
except Exception as e:
|
||||
logger.error(f"会话 {session.id} 发送提醒异常: {e}")
|
||||
continue
|
||||
|
||||
# 4. 查找需要标记待关闭的会话
|
||||
# 条件:active 状态 + 有坐席回复 + 超过10分钟未回复
|
||||
close_threshold = datetime.now() - timedelta(minutes=CLOSE_TIMEOUT_MINUTES)
|
||||
|
||||
close_stmt = select(Conversation).where(
|
||||
Conversation.status == "serving",
|
||||
Conversation.last_agent_reply_at.isnot(None),
|
||||
Conversation.last_agent_reply_at < close_threshold,
|
||||
)
|
||||
result = await db.execute(close_stmt)
|
||||
sessions_to_close = result.scalars().all()
|
||||
|
||||
logger.info(f"发现 {len(sessions_to_close)} 个需要标记待关闭的会话")
|
||||
|
||||
# 5. 更新会话状态为 pending_close
|
||||
for session in sessions_to_close:
|
||||
session.status = "pending_close"
|
||||
logger.info(f"会话 {session.id} 已标记为待关闭: employee_id={session.employee_id}")
|
||||
|
||||
# 提交数据库变更
|
||||
await db.commit()
|
||||
logger.info("超时检查任务执行完成")
|
||||
|
||||
except Exception as e:
|
||||
await db.rollback()
|
||||
logger.error(f"超时检查任务执行异常: {e}")
|
||||
raise
|
||||
@@ -1,29 +0,0 @@
|
||||
import sqlite3
|
||||
conn = sqlite3.connect('it_smart_desk.db')
|
||||
cursor = conn.cursor()
|
||||
|
||||
# Check employee table
|
||||
cursor.execute("SELECT name FROM sqlite_master WHERE type='table' AND name='employees'")
|
||||
if cursor.fetchone():
|
||||
cursor.execute('PRAGMA table_info(employees)')
|
||||
cols = [row[1] for row in cursor.fetchall()]
|
||||
print('Employee columns:')
|
||||
for c in cols:
|
||||
print(f' {c}')
|
||||
|
||||
missing = ['it_level', 'it_level_source', 'notes']
|
||||
for m in missing:
|
||||
status = "EXISTS" if m in cols else "MISSING!"
|
||||
print(f'{m}: {status}')
|
||||
else:
|
||||
print("No employees table found")
|
||||
|
||||
# Check todo_items and troubleshooting_templates tables
|
||||
for table in ['todo_items', 'troubleshooting_templates']:
|
||||
cursor.execute(f"SELECT name FROM sqlite_master WHERE type='table' AND name='{table}'")
|
||||
if cursor.fetchone():
|
||||
print(f"\n{table} table: EXISTS")
|
||||
else:
|
||||
print(f"\n{table} table: NOT FOUND (will be auto-created by SQLAlchemy on first access)")
|
||||
|
||||
conn.close()
|
||||
@@ -1,14 +0,0 @@
|
||||
import sqlite3
|
||||
conn = sqlite3.connect('it_smart_desk.db')
|
||||
cursor = conn.cursor()
|
||||
cursor.execute('PRAGMA table_info(conversations)')
|
||||
cols = [row[1] for row in cursor.fetchall()]
|
||||
print('Columns in conversations table:')
|
||||
for c in cols:
|
||||
print(f' {c}')
|
||||
print()
|
||||
missing = ['impact_scope', 'is_blocking', 'emotion_state']
|
||||
for m in missing:
|
||||
status = "EXISTS" if m in cols else "MISSING!"
|
||||
print(f'{m}: {status}')
|
||||
conn.close()
|
||||
@@ -1,84 +0,0 @@
|
||||
# -*- coding: utf-8 -*-
|
||||
"""内容审核服务 真实验证(#81 敏感词检测 / 隐私泄露识别)
|
||||
|
||||
真实验证点(来自功能规格说明书 + 状态看板验收标准):
|
||||
- moderate("你爱找谁找谁") → WARN + matched 含该词
|
||||
- check_privacy_leak("电话13800138000") → 含 "phone"
|
||||
- 命中敏感词动作是 WARN(仅警告,不阻断发送)
|
||||
- 自定义词库为写死的若干条(生产应从配置加载,当前未接)
|
||||
"""
|
||||
import pytest
|
||||
|
||||
from app.services.content_moderation_service import (
|
||||
ContentModerationService,
|
||||
ModerationAction,
|
||||
)
|
||||
|
||||
|
||||
@pytest.fixture
|
||||
def moderation_service():
|
||||
# 直接使用构造函数(单例亦可,这里用新实例避免跨测试状态)
|
||||
return ContentModerationService()
|
||||
|
||||
|
||||
def test_moderate_returns_warn_with_matched_word(moderation_service):
|
||||
"""验收点1: 命中自定义敏感词 → WARN 且 matched 含该词"""
|
||||
result = moderation_service.moderate("你爱找谁找谁")
|
||||
assert result.action == ModerationAction.WARN
|
||||
assert "你爱找谁找谁" in result.matched_words
|
||||
|
||||
|
||||
def test_moderate_all_known_custom_words_warn(moderation_service):
|
||||
"""所有已知自定义敏感词均能命中并返回 WARN"""
|
||||
words = ["投诉我", "你爱找谁找谁", "自己不会百度吗", "这点小事"]
|
||||
for w in words:
|
||||
r = moderation_service.moderate(w)
|
||||
assert r.action == ModerationAction.WARN, f"{w} 应被 warn"
|
||||
assert w in r.matched_words, f"{w} 应在 matched 中"
|
||||
|
||||
|
||||
def test_moderate_clean_text_passes(moderation_service):
|
||||
"""正常文本 → PASS,无命中词"""
|
||||
r = moderation_service.moderate("您好,我的电脑无法开机了")
|
||||
assert r.action == ModerationAction.PASS
|
||||
assert r.matched_words == []
|
||||
|
||||
|
||||
def test_moderate_empty_string_passes(moderation_service):
|
||||
"""空字符串 → PASS"""
|
||||
r = moderation_service.moderate("")
|
||||
assert r.action == ModerationAction.PASS
|
||||
assert r.matched_words == []
|
||||
|
||||
|
||||
def test_default_action_is_warn_not_block(moderation_service):
|
||||
"""关键事实: 当前命中动作是 WARN 而非 BLOCK(仅警告、不阻断发送)"""
|
||||
r = moderation_service.moderate("自己不会百度吗")
|
||||
assert r.action != ModerationAction.BLOCK
|
||||
assert r.action == ModerationAction.WARN
|
||||
|
||||
|
||||
def test_check_privacy_leak_phone(moderation_service):
|
||||
"""验收点2: 手机号被识别为 phone"""
|
||||
leaked = moderation_service.check_privacy_leak("我的电话13800138000")
|
||||
assert "phone" in leaked
|
||||
|
||||
|
||||
def test_check_privacy_leak_id_card(moderation_service):
|
||||
"""身份证号被识别为 id_card"""
|
||||
leaked = moderation_service.check_privacy_leak("身份证11010119900307123X")
|
||||
assert "id_card" in leaked
|
||||
|
||||
|
||||
def test_check_privacy_leak_clean_text_empty(moderation_service):
|
||||
"""正常沟通内容不触发隐私识别"""
|
||||
leaked = moderation_service.check_privacy_leak("这是正常的工作沟通内容")
|
||||
assert leaked == []
|
||||
|
||||
|
||||
def test_custom_word_list_is_hardcoded(moderation_service):
|
||||
"""确认自定义词库是写死的(生产应从配置加载,当前未接)"""
|
||||
words = moderation_service.custom_sensitive_words
|
||||
assert len(words) >= 4
|
||||
for w in ["投诉我", "你爱找谁找谁", "自己不会百度吗", "这点小事"]:
|
||||
assert w in words
|
||||
|
Before Width: | Height: | Size: 16 KiB |
|
Before Width: | Height: | Size: 25 KiB |
|
Before Width: | Height: | Size: 20 KiB |
|
Before Width: | Height: | Size: 7.4 KiB |
|
Before Width: | Height: | Size: 132 KiB |
|
Before Width: | Height: | Size: 29 KiB |
|
Before Width: | Height: | Size: 19 KiB |
|
Before Width: | Height: | Size: 15 KiB |
|
Before Width: | Height: | Size: 25 KiB |
|
Before Width: | Height: | Size: 966 KiB |
|
Before Width: | Height: | Size: 9.6 KiB |
|
Before Width: | Height: | Size: 12 KiB |
|
Before Width: | Height: | Size: 4.5 KiB |
|
Before Width: | Height: | Size: 11 KiB |
|
Before Width: | Height: | Size: 21 KiB |
|
Before Width: | Height: | Size: 18 KiB |
|
Before Width: | Height: | Size: 4.9 KiB |
|
Before Width: | Height: | Size: 16 KiB |
|
Before Width: | Height: | Size: 20 KiB |
|
Before Width: | Height: | Size: 3.1 KiB |
|
Before Width: | Height: | Size: 35 KiB |
|
Before Width: | Height: | Size: 23 KiB |
|
Before Width: | Height: | Size: 12 KiB |
|
Before Width: | Height: | Size: 2.1 KiB |
|
Before Width: | Height: | Size: 7.2 KiB |
|
Before Width: | Height: | Size: 20 KiB |
|
Before Width: | Height: | Size: 14 KiB |
|
Before Width: | Height: | Size: 11 KiB |
|
Before Width: | Height: | Size: 9.3 KiB |
|
Before Width: | Height: | Size: 3.6 KiB |
|
Before Width: | Height: | Size: 18 KiB |
|
Before Width: | Height: | Size: 30 KiB |
|
Before Width: | Height: | Size: 2.8 KiB |
|
Before Width: | Height: | Size: 3.0 KiB |
|
Before Width: | Height: | Size: 3.1 KiB |
|
Before Width: | Height: | Size: 4.1 KiB |
|
Before Width: | Height: | Size: 5.2 KiB |
|
Before Width: | Height: | Size: 2.5 KiB |
|
After Width: | Height: | Size: 13 KiB |