# ============================================================================= # 企微IT智能服务台 — FastAPI 应用入口 # ============================================================================= # 说明:FastAPI 应用的主入口文件,负责: # 1. 创建 FastAPI 应用实例 # 2. 配置 CORS 跨域支持 # 3. 挂载 API 路由 # 4. 注册全局异常处理器 # 5. 添加启动事件(初始化默认数据) # 6. 提供健康检查端点 # ============================================================================= import json import logging import os from contextlib import asynccontextmanager from fastapi import FastAPI, Request from fastapi.responses import JSONResponse from fastapi.middleware.cors import CORSMiddleware from sqlalchemy import select, text # 导入配置(读取环境变量) from app.config import settings # 导入路由汇总 from app.api.router import api_router # 导入共享服务生命周期管理 from app.dependencies import init_shared_services, cleanup_shared_services # 导入异常处理器和异常类 from app.utils.response import AppException, app_exception_handler, success_response # 导入日志配置(运行期日志文件接线:standard SOP 第三阶段 · 运行期日志 D 页依赖) from app.utils.logging_config import setup_logging # 导入定时任务 from app.tasks.reminder_task import check_unreplied_sessions # 配置日志格式 logging.basicConfig( level=logging.INFO, format="[%(asctime)s] [%(levelname)s] [%(name)s] %(message)s", datefmt="%Y-%m-%d %H:%M:%S", ) logger = logging.getLogger(__name__) # -------------------------------------------------------------------------- # 开发模式判定(模块级 helper,避免在 create_app 内每次重复 import) # -------------------------------------------------------------------------- def _is_dev_mode() -> bool: """检查是否启用了开发模式(DEV_MODE=true)。 三个检查源(任一为 true 即启用): 1. 环境变量 DEV_MODE=true(最高优先级,Docker 注入) 2. settings.dev_mode(从 .env.dev 读) 3. DEBUG 模式 + 本地主机(最严格) 注意:此函数与 backend/app/api/dev_auth.py 内的 _dev_mode_enabled() 逻辑一致, 这里用于"是否挂载 dev_auth 路由",那里用于"端点内是否放行"。 """ import os env_val = os.getenv("DEV_MODE", "").lower() == "true" if env_val: return True if getattr(settings, "dev_mode", False): return True return False # -------------------------------------------------------------------------- # 应用生命周期管理(启动和关闭事件) # -------------------------------------------------------------------------- @asynccontextmanager async def lifespan(app: FastAPI): """应用生命周期管理。 在应用启动时执行初始化操作(如插入默认数据), 在应用关闭时执行清理操作。 """ # ===== 启动事件 ===== # 运行期日志接线(standard SOP 第三阶段 · 运行期日志 D 页依赖): # 在 stdout handler 之外,将结构化 JSON 日志写入 settings.RUNTIME_LOG_DIR/wecom-it-desk.log # (RotatingFileHandler:单文件上限 20MB,保留 5 个备份,复用 JSONFormatter)。 # 目录不可创建/写入时自动降级为仅 stdout,不阻塞启动。 # 注:setup_logging 会清空 root logger 已有 handler(含模块级 basicConfig), # 并以 console(stdout) + 文件(JSON) 重新配置,避免重复输出。 setup_logging( level=os.getenv("LOG_LEVEL", "INFO"), json_format=True, log_dir=settings.RUNTIME_LOG_DIR, ) logger.info("🚀 智能IT服务平台启动中...") # 校验关键配置项(防止生产环境忘记配置导致静默失败) _validate_config() # 初始化共享服务实例(Redis/AIService/WecomService/AIHandler) # 这些实例在应用运行期间复用,避免每次请求重新创建导致资源泄漏 await init_shared_services() # 自动建表(开发阶段,生产环境应用 Alembic 迁移) await _auto_create_tables() # 初始化默认数据 await _init_default_data() # 启动超时提醒定时任务 _start_scheduler() logger.info("✅ 智能IT服务平台启动完成") yield # 应用运行中 # ===== 关闭事件 ===== logger.info("👋 智能IT服务平台关闭中...") # 停止超时提醒定时任务 _stop_scheduler() # 清理共享服务实例(关闭 Redis 连接、httpx 连接池等) await cleanup_shared_services() logger.info("✅ 智能IT服务平台已关闭") # -------------------------------------------------------------------------- # 定时任务调度器 # -------------------------------------------------------------------------- # 全局调度器实例 _scheduler = None async def expire_pending_suggestions(): """将超时未处理的 pending 建议标记为 expired。 扫描所有 status='pending' 且创建时间超过 72 小时的建议, 将其状态更新为 expired(终态,不可逆)。 执行频率:每小时 1 次 超时阈值:72 小时 扫描范围:仅 status='pending'(不含 queued) """ from app.database import _get_session_factory try: async_session_factory = _get_session_factory() async with async_session_factory() as db: # 将创建超过 72 小时的 pending 建议标记为 expired result = await db.execute( text( "UPDATE knowledge_suggestions " "SET status = 'expired', updated_at = NOW() " "WHERE status = 'pending' " "AND created_at < NOW() - INTERVAL '72 hours'" ) ) await db.commit() expired_count = result.rowcount if expired_count > 0: logger.info( f"⏰ 过期处理完成:{expired_count} 条 pending 建议已标记为 expired" f"(超时 72 小时)" ) else: logger.debug("过期处理完成:无待过期的 pending 建议") except Exception as e: logger.error(f"过期处理定时任务执行失败: {e}") def _start_scheduler(): """启动定时任务调度器。 启动 APScheduler 调度器,注册超时提醒定时任务。 每 30 秒检查一次超时未回复的会话。 """ global _scheduler if _scheduler is not None: logger.warning("调度器已启动,跳过") return try: from apscheduler.schedulers.asyncio import AsyncIOScheduler _scheduler = AsyncIOScheduler() # 注册超时提醒任务(每 30 秒执行一次) # _scheduler.add_job( check_unreplied_sessions, 'interval', seconds=30, id='check_unreplied_sessions', name='检查超时未回复会话', replace_existing=True, ) # 注册过期处理任务(每 1 小时执行一次,将超时 72 小时的 pending 建议标记为 expired) # _scheduler.add_job( expire_pending_suggestions, 'interval', hours=1, id='expire_pending_suggestions', name='过期处理:超时72小时的pending建议标记为expired', replace_existing=True, ) # 注册每日测验题目生成任务(每天 03:00 执行) from apscheduler.triggers.cron import CronTrigger from app.tasks.quiz_generation_task import run_daily_quiz_generation # _scheduler.add_job( run_daily_quiz_generation, CronTrigger(hour=3, minute=0), id='daily_quiz_generation', name='每日测验题目自动生成', replace_existing=True, misfire_grace_time=3600, # 错过执行窗口1小时内仍可补执行 ) # # # 注册Token异常检测任务(每 5 分钟执行一次) # from app.tasks.token_anomaly_detection import detect_token_ip_anomaly # _scheduler.add_job( # detect_token_ip_anomaly, # "interval", # minutes=5, # id="detect_token_ip_anomaly", # name="Token异常检测:多IP使用", # replace_existing=True, # ) # _scheduler.start() # logger.info("✅ 超时提醒定时任务已启动(每 30 秒执行一次)") # logger.info("✅ 过期处理定时任务已启动(每 1 小时执行一次,超时阈值 72 小时)") # logger.info("✅ 每日题目生成任务已启动(每天 03:00 执行)") # logger.info("✅ Token异常检测任务已启动(每 5 分钟执行一次)") except Exception as e: logger.error(f"启动定时任务调度器失败: {e}") # 定时任务启动失败不阻塞应用启动 def _stop_scheduler(): """停止定时任务调度器。 在应用关闭时调用,确保定时任务正确关闭。 """ global _scheduler if _scheduler is None: return try: _scheduler.shutdown(wait=False) _scheduler = None logger.info("✅ 超时提醒定时任务已停止") except Exception as e: logger.error(f"停止定时任务调度器失败: {e}") # -------------------------------------------------------------------------- # 配置校验(启动时检查关键配置项是否为占位符) # -------------------------------------------------------------------------- # 占位符列表:这些默认值在 config.py 中设置,生产环境必须替换 _PLACEHOLDER_VALUES = { "wecom_corp_id": "ww1234567890abcdef", "wecom_secret": "your-agent-secret", "wecom_token": "your-callback-token", "wecom_encoding_aes_key": "your-aes-key-43-characters-long-encoding-key", } def _validate_config(): """校验关键配置项是否为占位符。 生产环境部署时,如果忘记修改 config.py 中的占位符值, 会导致 AES 解密静默失败、企微 API 调用 400 等问题。 此函数在启动时检查这些关键配置,输出醒目警告。 """ warnings = [] for key, placeholder in _PLACEHOLDER_VALUES.items(): actual_value = getattr(settings, key, "") if actual_value == placeholder: warnings.append(f" ⚠️ {key} = '{placeholder}' (未配置!)") if warnings: logger.warning( "━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━\n" "⚠️ 检测到以下关键配置仍为占位符,请修改 .env 或环境变量:\n" + "\n".join(warnings) + "\n" " 企微回调消息将无法正常解密!\n" " 参考 .env.example 或项目部署手册进行配置。\n" "━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━" ) else: logger.info("✅ 关键配置校验通过") # -------------------------------------------------------------------------- # 自动建表(开发阶段使用) # -------------------------------------------------------------------------- async def _auto_create_tables(): """自动创建所有数据库表。 开发阶段使用,根据模型定义自动创建表。 生产环境应使用 Alembic 迁移来管理表结构变更。 工作原理: 1. 获取 engine(懒加载) 2. 通过 Base.metadata 收集所有模型定义 3. 执行 CREATE TABLE IF NOT EXISTS """ from app.database import _get_engine, Base # 导入所有模型,确保 Base.metadata 知道所有表的定义 # 如果不导入,Base.metadata 里只有基类,不会建任何表 import app.models # noqa: F401 engine = _get_engine() async with engine.begin() as conn: # checkfirst=True: 只创建不存在的表,不会覆盖已有表和数据 await conn.run_sync(Base.metadata.create_all, checkfirst=True) logger.info("数据库表检查/创建完成") # -------------------------------------------------------------------------- # 初始化默认数据 # -------------------------------------------------------------------------- async def _init_default_data(): """初始化默认数据。 当数据库表为空时,插入预置配置数据,包括: 1. system_configs — 系统配置(关键词、阈值、话术等) 2. funny_phrases — 趣味话术 3. quick_reply_templates — 快速回复模板 4. approval_links — 审批流程链接 5. software_downloads — 软件下载入口 6. (dev 模式)demo_conversations — 演示用会话,让前端有数据可发 只在表为空时插入,避免重复插入。 """ from app.database import _get_session_factory from app.models.system_config import SystemConfig from app.models.funny_phrase import FunnyPhrase from app.models.quick_reply_template import QuickReplyTemplate from app.models.approval_link import ApprovalLink from app.models.software_download import SoftwareDownload from app.config import settings async_session_factory = _get_session_factory() async with async_session_factory() as db: try: # 1. 初始化系统配置 await _init_system_configs(db, SystemConfig) # 2. 初始化趣味话术 await _init_funny_phrases(db, FunnyPhrase) # 3. 初始化快速回复模板 await _init_quick_reply_templates(db, QuickReplyTemplate) # 4. 初始化审批流程链接 await _init_approval_links(db, ApprovalLink) # 5. 初始化软件下载入口 await _init_software_downloads(db, SoftwareDownload) # 6. v0.7.1 task #86 — RBAC 5 角色种子(细粒度权限) # 行为: 已有角色更新 permissions,缺则新建 from app.data.seed_rbac import seed_rbac_roles await seed_rbac_roles(db) # 6.1 阶段5 — 自动化场景配置种子(P0 四场景默认启用) await _init_scenario_configs(db) # 7. 初始化超级管理员(环境变量配置) from app.services.admin_user_service import init_super_admin await init_super_admin(db) # 7.1 测验题库种子数据(Dify 生成 70 知识题 + 手写诊断模板) from app.data.seed_quiz import seed_quiz_data await seed_quiz_data(db) # 8. (dev 模式)初始化 demo 会话,让前端有数据可发 # 真因:之前没建,前端硬编码的 conv-001 调 POST /messages 返 "会话不存在" 3003 if getattr(settings, 'dev_mode', False) or os.getenv('DEV_MODE', '').lower() == 'true': await _init_demo_conversations(db) await db.commit() logger.info("默认数据初始化完成") except Exception as e: await db.rollback() logger.error(f"默认数据初始化失败: {e}") async def _init_demo_conversations(db): """(dev 模式专用)建 5 条 demo 会话,让前端有数据可测。 涵盖各种状态: - ai_handling: AI 正在处理(2 条,不同员工) - queued: 等坐席接手 - serving: 坐席服务中 - resolved: 已结单 只在 conversations 表为空时建,避免重复。 """ from app.models.conversation import Conversation existing = (await db.execute(select(Conversation).limit(1))).scalar_one_or_none() if existing: logger.info("demo 会话已存在,跳过") return import uuid as _uuid from datetime import datetime, timezone, timedelta now = datetime.now(timezone.utc) demo_convs = [ { "id": "conv-001", "corp_id": "wwa8c87970b2011f41", "employee_id": "dev-user-001", "employee_name": "张三(普通员工)", "department": "财务部", "position": "会计", "level": "P5", "status": "ai_handling", "is_vip": False, "is_pinned": False, "is_todo": False, "urgency_score": 30, "tags": ["财务", "IT"], "assigned_agent_id": None, "collaborating_agent_ids": [], "participants": [], "ai_substantive_reply_count": 0, "impact_scope": 1, "is_blocking": False, "emotion_state": "normal", "dify_conversation_id": None, "last_message_at": now - timedelta(minutes=2), "last_message_summary": "想问下 VPN 怎么连", }, { "id": "conv-002", "corp_id": "wwa8c87970b2011f41", "employee_id": "dev-user-001", "employee_name": "张三(普通员工)", "department": "财务部", "position": "会计", "level": "P5", "status": "queued", "is_vip": False, "is_pinned": True, "is_todo": True, "urgency_score": 70, "tags": ["紧急", "VPN"], "assigned_agent_id": None, "collaborating_agent_ids": [], "participants": [], "ai_substantive_reply_count": 2, "impact_scope": 3, "is_blocking": True, "emotion_state": "worried", "dify_conversation_id": "dify-conv-002", "last_message_at": now - timedelta(minutes=5), "last_message_summary": "VPN 连不上,影响工作", }, { "id": "conv-003", "corp_id": "wwa8c87970b2011f41", "employee_id": "dev-multi-001", "employee_name": "周八(多角色测试)", "department": "测试部", "position": "测试工程师", "level": "P6", "status": "serving", "is_vip": True, "is_pinned": False, "is_todo": False, "urgency_score": 50, "tags": ["软件安装"], "assigned_agent_id": "dev-agent-001", "collaborating_agent_ids": [], "participants": [], "ai_substantive_reply_count": 1, "impact_scope": 1, "is_blocking": False, "emotion_state": "normal", "dify_conversation_id": "dify-conv-003", "last_message_at": now - timedelta(minutes=10), "last_message_summary": "需要装 WPS 专业版", }, { "id": "conv-004", "corp_id": "wwa8c87970b2011f41", "employee_id": "dev-supervisor-001", "employee_name": "王五(部门主管)", "department": "信息技术部", "position": "主管", "level": "M3", "status": "serving", "is_vip": True, "is_pinned": True, "is_todo": False, "urgency_score": 80, "tags": ["系统升级"], "assigned_agent_id": "dev-agent-001", "collaborating_agent_ids": ["dev-admin-001"], "participants": [], "ai_substantive_reply_count": 3, "impact_scope": 50, "is_blocking": True, "emotion_state": "urgent", "dify_conversation_id": "dify-conv-004", "last_message_at": now - timedelta(minutes=15), "last_message_summary": "ERP 系统升级咨询", }, { "id": "conv-005", "corp_id": "wwa8c87970b2011f41", "employee_id": "dev-security-001", "employee_name": "赵六(安全团队)", "department": "信息安全部", "position": "安全工程师", "level": "P7", "status": "resolved", "is_vip": False, "is_pinned": False, "is_todo": False, "urgency_score": 20, "tags": ["安全"], "assigned_agent_id": "dev-agent-001", "collaborating_agent_ids": [], "participants": [], "ai_substantive_reply_count": 5, "impact_scope": 1, "is_blocking": False, "emotion_state": "normal", "dify_conversation_id": "dify-conv-005", "last_message_at": now - timedelta(hours=2), "last_message_summary": "已处理:密码策略咨询", }, ] for data in demo_convs: db.add(Conversation(**data)) logger.info(f"已初始化 {len(demo_convs)} 条 demo 会话(仅 dev 模式)") async def _init_system_configs(db, SystemConfig): """初始化系统配置项。""" from sqlalchemy import select, func count_stmt = select(func.count(SystemConfig.id)) result = await db.execute(count_stmt) count = result.scalar() or 0 if count > 0: logger.debug(f"system_configs 已有 {count} 条数据,跳过初始化") return configs = [ SystemConfig(config_key="hand_raise_keywords", config_value=json.dumps(["转人工", "人工", "人工服务", "真人", "客服", "帮我转人工", "找人工"], ensure_ascii=False), description="举手触发关键词"), SystemConfig(config_key="emotion_keywords_angry", config_value=json.dumps(["崩溃", "愤怒", "投诉", "差劲", "垃圾", "太差了", "受不了"], ensure_ascii=False), description="愤怒情绪关键词"), SystemConfig(config_key="emotion_keywords_urgent", config_value=json.dumps(["急", "紧急", "马上", "立刻", "赶紧", "十万火急", "快点"], ensure_ascii=False), description="紧急情绪关键词"), SystemConfig(config_key="emotion_keywords_worried", config_value=json.dumps(["担心", "害怕", "出错", "丢失", "完蛋", "糟糕"], ensure_ascii=False), description="担忧情绪关键词"), SystemConfig(config_key="intervene_round_threshold", config_value="3", description="需介入追问轮次阈值"), SystemConfig(config_key="urgency_base_keyword_score", config_value="1", description="关键词匹配基础加分"), SystemConfig(config_key="urgency_emotion_bonus", config_value="1", description="情绪标记加成分"), SystemConfig(config_key="urgency_vip_bonus", config_value="1", description="VIP加成分"), SystemConfig(config_key="urgency_repeat_bonus", config_value="1", description="重复追问加成分"), SystemConfig(config_key="polling_interval_seconds", config_value="3", description="坐席轮询间隔(秒)"), SystemConfig(config_key="access_token_buffer_seconds", config_value="300", description="access_token提前刷新时间(秒)"), SystemConfig(config_key="emergency_mode", config_value="false", description="应急模式开关(true=启用员工服务通道,智能服务台降级)"), ] db.add_all(configs) await db.flush() logger.info(f"初始化 system_configs: {len(configs)} 条") async def _init_funny_phrases(db, FunnyPhrase): """初始化趣味话术。""" from sqlalchemy import select, func count_stmt = select(func.count(FunnyPhrase.id)) result = await db.execute(count_stmt) count = result.scalar() or 0 if count > 0: logger.debug(f"funny_phrases 已有 {count} 条数据,跳过初始化") return phrases = [ FunnyPhrase(scene="shake", content="少主,这就为您去摇人,稍等...", tone="亲切", sort_order=1), FunnyPhrase(scene="keyword", content="收到!这就帮您摇位大神来", tone="稍正式", sort_order=1), FunnyPhrase(scene="waiting", content="人还在路上,别急别急~", tone="安抚", sort_order=1), FunnyPhrase(scene="connected", content="人摇来了!IT坐席为您服务", tone="明确交接", sort_order=1), FunnyPhrase(scene="timeout", content="坐席都在忙,不过AI还在呢,要不先聊聊?我再继续摇", tone="降级安抚", sort_order=1), FunnyPhrase(scene="vip", content="这就帮您安排专家,请稍候", tone="正式", sort_order=1), ] db.add_all(phrases) await db.flush() logger.info(f"初始化 funny_phrases: {len(phrases)} 条") async def _init_quick_reply_templates(db, QuickReplyTemplate): """初始化快速回复模板。""" from sqlalchemy import select, func count_stmt = select(func.count(QuickReplyTemplate.id)) result = await db.execute(count_stmt) count = result.scalar() or 0 if count > 0: logger.debug(f"quick_reply_templates 已有 {count} 条数据,跳过初始化") return templates = [ QuickReplyTemplate(category="账号", title="密码重置", content="您好{employee_name},您的密码重置链接已发送至您的企业邮箱,请在30分钟内完成操作。", variables=["employee_name"], sort_order=1), QuickReplyTemplate(category="账号", title="账号解锁", content="您好,您的账号已解锁,请5分钟后重新尝试登录。如仍有问题请联系智能IT服务台。", variables=[], sort_order=2), QuickReplyTemplate(category="网络", title="VPN连接指引", content="请按以下步骤操作:1.打开VPN客户端 2.选择\u201c公司内网\u201d 3.输入域账号密码 4.点击连接。详细图文教程请查看右侧\u201c操作步骤\u201d。", variables=[], sort_order=3), QuickReplyTemplate(category="网络", title="WiFi连接", content="公司WiFi名称:Office-5G,密码请咨询前台或查看工位标签。", variables=[], sort_order=4), QuickReplyTemplate(category="软件", title="软件安装申请", content="您好,软件安装需要提交审批申请。请在右侧\u201c审批流程\u201d中点击\u201c软件安装申请\u201d链接提交。", variables=[], sort_order=5), QuickReplyTemplate(category="硬件", title="设备报修", content="您好,设备报修请提交工单。请在右侧\u201c审批流程\u201d中点击\u201c设备报修\u201d链接提交,IT会在24小时内联系您。", variables=[], sort_order=6), QuickReplyTemplate(category="通用", title="会话结束", content="您好,请问还有其他问题吗?如无其他问题,我将结束本次服务。祝您工作顺利!", variables=[], sort_order=7), QuickReplyTemplate(category="通用", title="稍等回复", content="收到,我正在为您查询,请稍等片刻。", variables=[], sort_order=8), ] db.add_all(templates) await db.flush() logger.info(f"初始化 quick_reply_templates: {len(templates)} 条") async def _init_approval_links(db, ApprovalLink): """初始化审批流程链接。""" from sqlalchemy import select, func count_stmt = select(func.count(ApprovalLink.id)) result = await db.execute(count_stmt) count = result.scalar() or 0 if count > 0: logger.debug(f"approval_links 已有 {count} 条数据,跳过初始化") return links = [ # v0.5.2:一站式运维平台真实工单链接(域名 devops.dc.servyou-it.com,已实现企微免登录) # v0.5.3 更新:去掉 "IT设备升级与硬件维修" (申请单冲突,后续移除) ApprovalLink(category="IT", title="零信任(原VPN)账号申请", url="https://devops.dc.servyou-it.com/ITSM/workflow/service/createTicket?name=%E5%91%98%E5%B7%A5%E9%9B%B6%E4%BF%A1%E4%BB%BB%EF%BC%88%E5%8E%9FVPN%EF%BC%89%E8%B4%A6%E5%8F%B7%E7%94%B3%E8%AF%B7IT", sort_order=1), ApprovalLink(category="IT", title="活动与会议技术支持", url="https://devops.dc.servyou-it.com/ITSM/workflow/service/createTicket?name=%E6%B4%BB%E5%8A%A8%E4%B8%8E%E4%BC%9A%E8%AE%AE%E6%8A%80%E6%9C%AF%E6%94%AF%E6%8C%81", sort_order=2), # sort_order=3 故意空缺:旧版本是"IT设备升级与硬件维修",已与一站式运维平台冲突,不再提供 ApprovalLink(category="IT", title="员工IT支持与故障报修", url="https://devops.dc.servyou-it.com/ITSM/workflow/service/createTicket?name=%E5%91%98%E5%B7%A5IT%E6%94%AF%E6%8C%81%E4%B8%8E%E6%95%85%E9%9A%9C%E6%8A%A5%E4%BF%AE", sort_order=4), ApprovalLink(category="IT", title="终端设备网络准入申请", url="https://devops.dc.servyou-it.com/ITSM/workflow/service/createTicket?name=%E7%BB%88%E7%AB%AF%E8%AE%BE%E5%A4%87%E7%BD%91%E7%BB%9C%E5%87%86%E5%85%A5%E7%94%B3%E8%AF%B7", sort_order=5), ApprovalLink(category="IT", title="公共邮箱账号申请", url="https://devops.dc.servyou-it.com/ITSM/workflow/service/createTicket?name=%E5%85%AC%E5%85%B1%E9%82%AE%E7%AE%B1%E8%B4%A6%E5%8F%B7%E7%94%B3%E8%AF%B7", sort_order=6), # HR / 行政 / 财务 占位(待后续接入真实流程) ApprovalLink(category="HR", title="入职手续", url="https://审批系统地址/onboarding", sort_order=7), ApprovalLink(category="HR", title="离职手续", url="https://审批系统地址/offboarding", sort_order=8), ApprovalLink(category="行政", title="办公用品申领", url="https://审批系统地址/office-supplies", sort_order=9), ApprovalLink(category="财务", title="报销申请", url="https://审批系统地址/reimbursement", sort_order=10), ] db.add_all(links) await db.flush() logger.info(f"初始化 approval_links: {len(links)} 条") async def _init_software_downloads(db, SoftwareDownload): """初始化软件下载入口。""" from sqlalchemy import select, func count_stmt = select(func.count(SoftwareDownload.id)) result = await db.execute(count_stmt) count = result.scalar() or 0 if count > 0: logger.debug(f"software_downloads 已有 {count} 条数据,跳过初始化") return downloads = [ SoftwareDownload(category="办公", name="企业微信", version="最新版", platform="全平台", download_url="https://work.weixin.qq.com/#download", sort_order=1), SoftwareDownload(category="办公", name="WPS Office", version="12.1", platform="Windows/Mac", download_url="https://www.wps.cn/download", sort_order=2), SoftwareDownload(category="办公", name="Microsoft Teams", version="最新版", platform="全平台", download_url="https://www.microsoft.com/teams/download", sort_order=3), SoftwareDownload(category="开发", name="VS Code", version="1.90", platform="Windows/Mac/Linux", download_url="https://code.visualstudio.com/download", sort_order=4), SoftwareDownload(category="开发", name="Git", version="2.45", platform="Windows/Mac", download_url="https://git-scm.com/download", sort_order=5), SoftwareDownload(category="安全", name="公司VPN客户端", version="3.2", platform="Windows/Mac", download_url="https://内部下载地址/vpn-client", sort_order=6), SoftwareDownload(category="工具", name="7-Zip", version="24.06", platform="Windows", download_url="https://www.7-zip.org/download", sort_order=7), SoftwareDownload(category="工具", name="PDF阅读器", version="最新版", platform="Windows/Mac", download_url="https://get.adobe.com/reader/", sort_order=8), ] db.add_all(downloads) await db.flush() logger.info(f"初始化 software_downloads: {len(downloads)} 条") async def _init_scenario_configs(db): """初始化自动化场景配置(P0 四场景默认启用)。""" from sqlalchemy import func, select from app.models.automation import ScenarioConfig from app.services.automation import DEFAULT_SCENARIO_CONFIGS count = (await db.execute(select(func.count(ScenarioConfig.id)))).scalar() or 0 if count > 0: logger.debug(f"auto_scenario_configs 已有 {count} 条,跳过初始化") return for key, cfg in DEFAULT_SCENARIO_CONFIGS.items(): db.add( ScenarioConfig( scenario_key=key, name=cfg.get("name", key), description=cfg.get("description", ""), enabled=cfg.get("enabled", True), trigger_conditions=cfg.get("trigger_conditions"), actions=cfg.get("actions"), approval_strategy=cfg.get("approval_strategy"), ) ) await db.flush() logger.info(f"初始化 auto_scenario_configs: {len(DEFAULT_SCENARIO_CONFIGS)} 条") # -------------------------------------------------------------------------- # 创建 FastAPI 应用 # -------------------------------------------------------------------------- def create_app() -> FastAPI: """创建并配置 FastAPI 应用实例。 使用工厂函数模式,方便测试时创建不同的应用实例。 Returns: FastAPI: 配置好的应用实例 """ # 创建 FastAPI 实例 # lifespan: 应用生命周期管理(启动/关闭事件) app = FastAPI( title="智能IT服务平台", description="基于企微自建应用消息API的IT服务坐席系统", version="1.0.0", lifespan=lifespan, ) # ---------------------------------------------------------------------- # 管理端 IP 白名单中间件(仅 production 生效) # ---------------------------------------------------------------------- # 三端认证重构 AUTH-02:管理后台仅限内网/VPN + IP 白名单 # 非生产环境自动跳过,方便本地开发 # ---------------------------------------------------------------------- from app.middleware.admin_ip_whitelist import AdminIPWhitelistMiddleware app.add_middleware(AdminIPWhitelistMiddleware) # ---------------------------------------------------------------------- # 配置 CORS(跨域资源共享) # ---------------------------------------------------------------------- # 为什么需要 CORS:前端和后端运行在不同端口,浏览器会阻止跨域请求 # allow_origins: 允许的前端地址列表 # allow_credentials: 允许携带 Cookie # allow_methods: 允许的 HTTP 方法(仅允许必要的方法) # allow_headers: 允许的请求头(仅允许必要的头) # ---------------------------------------------------------------------- app.add_middleware( CORSMiddleware, allow_origins=settings.cors_origins_list, allow_credentials=True, allow_methods=["GET", "POST", "PUT", "DELETE", "OPTIONS"], allow_headers=["Authorization", "Content-Type", "X-Employee-Id"], ) # ---------------------------------------------------------------------- # 速率限制(防止暴力破解和 DDoS) # ---------------------------------------------------------------------- # slowapi 为每个 IP 维护请求计数器(默认内存后端) # 登录接口严格限制(防暴力破解),普通接口宽松限制(防滥用) # ---------------------------------------------------------------------- from slowapi import Limiter, _rate_limit_exceeded_handler from slowapi.util import get_remote_address from slowapi.errors import RateLimitExceeded from starlette.responses import JSONResponse as RateLimitJSONResponse from app.utils.response import error_response as _rl_error_response # 速率限制器:按客户端 IP 维度限制 # 移除 env_file=None 参数:slowapi 0.1.9 不支持该参数 # python-dotenv 已在应用启动时处理 .env 文件 limiter = Limiter(key_func=get_remote_address) # 注册速率限制超限处理器 app.state.limiter = limiter @app.exception_handler(RateLimitExceeded) async def rate_limit_handler(request, exc: RateLimitExceeded): """速率限制超限响应:返回 429 状态码和友好提示。""" return RateLimitJSONResponse( status_code=429, content=_rl_error_response(429, f"请求过于频繁,请 {exc.detail} 后重试"), ) # ---------------------------------------------------------------------- # 注册全局异常处理器 # ---------------------------------------------------------------------- # 当业务逻辑抛出 AppException 时,自动转换为统一响应格式 # ---------------------------------------------------------------------- app.add_exception_handler(AppException, app_exception_handler) # ---------------------------------------------------------------------- # 注册兜底异常处理器(捕获所有未预期的异常,避免裸 500) # ---------------------------------------------------------------------- # 数据库连接失败、Redis 异常、第三方库错误等非 AppException 异常 # 都会被捕获并返回统一格式的错误响应,同时记录详细日志 # ---------------------------------------------------------------------- import traceback from fastapi.responses import JSONResponse from app.utils.response import error_response @app.exception_handler(Exception) async def catch_all_exception_handler(request, exc): """兜底异常处理器:捕获所有未预期异常。 安全处理: - 详细异常信息记录到日志(供排查) - 响应只返回通用错误信息(避免泄露内部细节) """ # 记录完整错误堆栈(用于排查问题) logger.error(f"未预期异常: {exc}\n{traceback.format_exc()}") # 返回统一格式的错误响应(HTTP 200 + 业务错误码) # 安全:响应不包含具体异常信息,仅返回通用消息 return JSONResponse( status_code=200, content=error_response(1005, "服务器内部错误,请稍后重试或联系管理员") ) # ---------------------------------------------------------------------- # 请求日志 + 兜底异常中间件 # ---------------------------------------------------------------------- # 使用中间件而非 exception_handler 来捕获所有异常 # 原因:FastAPI 的 @app.exception_handler(Exception) 在某些情况下 # 无法捕获异常(如依赖注入 yield 阶段的异常),而中间件更可靠 # ---------------------------------------------------------------------- import traceback as tb_module from starlette.requests import Request from starlette.responses import Response as StarletteResponse, JSONResponse as StarJSONResponse from app.utils.response import error_response as _error_response @app.middleware("http") async def catch_errors_and_log(request: Request, call_next): """请求日志 + 兜底异常中间件。 1. 记录每个请求的方法、路径、状态码 2. 捕获所有未处理异常,返回统一格式的 JSON 错误响应 """ # 使用 print 而非 logger,确保输出立即可见(调试阶段) print(f">>> [MW] 收到请求: {request.method} {request.url.path}", flush=True) try: response: StarletteResponse = await call_next(request) print(f"<<< [MW] 响应完成: {request.method} {request.url.path} → {response.status_code}", flush=True) return response except Exception as e: # 捕获所有未处理异常(包括依赖注入阶段的异常) # 安全:详细日志仅记录,响应不泄露异常信息 error_tb = tb_module.format_exc() print(f"!!! [MW] 未捕获异常: {request.method} {request.url.path}\n{error_tb}", flush=True) logger.error(f"!!! 未捕获异常: {request.method} {request.url.path}\n{error_tb}") # 返回统一格式的 JSON 错误响应(HTTP 200 + 业务错误码 1005) # 安全:响应不包含具体异常信息 return StarJSONResponse( status_code=200, content=_error_response(1005, "服务器内部错误,请稍后重试或联系管理员"), ) # ---------------------------------------------------------------------- # 挂载 API 路由 # ---------------------------------------------------------------------- # 注意:nginx 已经通过 location /api/ 处理了前缀路由, # 请求到达后端时 /api/ 已被 strip,因此此处不需要再加 /api 前缀 app.include_router(api_router) # ---------------------------------------------------------------------- # 开发模式 Mock OAuth(仅 DEV_MODE=true 时挂载) # ---------------------------------------------------------------------- # ⚠️ 生产环境严禁启用(DEV_MODE=false 或不设置) # 挂载的端点: # GET /api/dev/login — Mock 登录,跳过企微 OAuth 直接返回 token # GET /api/dev/users — 列出预设 dev 用户 # GET /api/dev/health — dev 模式状态自检 # 即使挂载了,每个端点内部也会再 _dev_mode_enabled() 二次校验 # ---------------------------------------------------------------------- if _is_dev_mode(): from app.api.dev_auth import router as dev_auth_router app.include_router(dev_auth_router) logger.warning( "🧪 ━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━\n" "🧪 DEV_MODE 已启用 - Mock OAuth 端点已挂载\n" "🧪 仅供本地开发测试使用,生产环境必须关闭!\n" "🧪 端点列表:\n" "🧪 GET /api/dev/login - Mock 登录\n" "🧪 GET /api/dev/users - 列出预设用户\n" "🧪 GET /api/dev/health - dev 模式状态\n" "🧪 ━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━" ) # ---------------------------------------------------------------------- # 挂载 WebSocket 路由 # ---------------------------------------------------------------------- # WebSocket 端点不挂 /api 前缀,直接注册在根路径 # 原因:WebSocket 不是 REST API,前端通过 /ws/{agent_id} 连接 # Vite 开发服务器单独配置了 /ws 的 WebSocket 代理 # ---------------------------------------------------------------------- from app.api.ws import router as ws_router app.include_router(ws_router) # 阶段5 自动化闭环:专用 WebSocket 通道 /ws/automation/{session_id} from app.api.automation import ws_router as automation_ws_router app.include_router(automation_ws_router) # ---------------------------------------------------------------------- # 诊断端点(调试用,生产环境删除) # ---------------------------------------------------------------------- @app.get("/test-ping", tags=["诊断"]) async def test_ping(): """简单测试 — 不依赖数据库和 Redis""" return success_response(data={"message": "pong"}) @app.get("/test-error", tags=["诊断"]) async def test_error(): """测试异常处理 — 故意抛出异常""" raise Exception("这是故意抛出的测试异常") # ---------------------------------------------------------------------- # 健康检查端点 # ---------------------------------------------------------------------- # 用于 Docker 健康检查和负载均衡探针 # 返回简单的 JSON 表示服务正在运行 @app.get("/health", tags=["系统"]) async def health_check(): """健康检查端点。 返回服务运行状态,用于: - Docker 健康检查 - 负载均衡探针 - 监控系统检测服务是否存活 """ return {"status": "ok", "service": "wecom-it-smart-desk"} @app.get("/ready", tags=["系统"]) async def readiness_check(): """就绪检查端点。 检查服务依赖(DB + Redis),不调用企微 API(避免阻塞)。 用于 K8s readinessProbe。 """ try: # 检查数据库 from app.database import _get_engine engine = _get_engine() async with engine.connect() as conn: await conn.execute(text("SELECT 1")) db_status = "ok" except Exception as e: db_status = f"error: {str(e)}" try: # 检查 Redis from app.config import settings redis_client = settings.create_redis_client() await redis_client.ping() redis_status = "ok" except Exception as e: redis_status = f"error: {str(e)}" if db_status == "ok" and redis_status == "ok": return {"status": "ready", "db": db_status, "redis": redis_status} else: return JSONResponse( status_code=503, content={"status": "not_ready", "db": db_status, "redis": redis_status} ) @app.get("/metrics", tags=["系统"]) async def metrics(): """指标端点。 返回服务运行指标,用于 Prometheus 采集。 """ import psutil return { "status": "ok", "metrics": { "cpu_percent": psutil.cpu_percent(interval=0.1), "memory_percent": psutil.virtual_memory().percent, "disk_percent": psutil.disk_usage("/").percent, } } @app.get("/version", tags=["系统"]) async def version(): """版本信息端点。 返回服务版本信息。 """ import subprocess try: git_hash = subprocess.check_output( ["git", "rev-parse", "HEAD"], cwd=app_root, text=True ).strip()[:8] except Exception: git_hash = "unknown" return { "service": "wecom-it-smart-desk", "version": "1.1.0", "build": git_hash, } # ---------------------------------------------------------------------- # 打印所有已注册的路由(调试用) # ---------------------------------------------------------------------- routes_info = [] for route in app.routes: if hasattr(route, 'methods') and hasattr(route, 'path'): routes_info.append(f" {', '.join(route.methods)} {route.path}") if routes_info: logger.info(f"已注册路由 ({len(routes_info)} 个):\n" + "\n".join(routes_info)) else: logger.warning("⚠️ 没有注册任何路由!") return app # 创建应用实例(uvicorn 通过 app.main:app 引用此对象) app = create_app()