diff --git a/.env.example b/.env.example index 65d473e..6bbfc2c 100644 --- a/.env.example +++ b/.env.example @@ -1,5 +1,5 @@ # 数据库配置 -DATABASE_URL=postgresql+asyncpg://postgres:12345@localhost:5432/ai_talk_callback_db +DATABASE_URL=postgresql+asyncpg://user:password@localhost:5432/ai_talk_callback_db # Redis配置 REDIS_URL=redis://localhost:6379/0 @@ -12,6 +12,14 @@ COUNT_THRESHOLD=3 EXTERNAL_API_RETRY_MAX=4 EXTERNAL_API_URL=https://external.com/api/openapi/customerApi/aiTaskResultFail +# 日志配置 +LOG_LEVEL=INFO +LOG_FILE=logs/app.log +LOG_MAX_BYTES=10485760 +LOG_BACKUP_COUNT=5 +LOG_FORMAT=%(asctime)s - %(name)s - %(levelname)s - %(message)s +LOG_DATE_FORMAT=%Y-%m-%d %H:%M:%S + # 应用配置 APP_NAME=AI Talk Callback API DEBUG=false \ No newline at end of file diff --git a/JSON_LOG_EXAMPLES.md b/JSON_LOG_EXAMPLES.md new file mode 100644 index 0000000..45586b9 --- /dev/null +++ b/JSON_LOG_EXAMPLES.md @@ -0,0 +1,350 @@ +# JSON格式日志输出示例 + +## 优化后的 `log_callback_request` 方法 - JSON格式日志 + +### 1. site_id 日志输出 + +```json +📝 site_id: "test-site-123" +``` + +### 2. 请求头日志输出(JSON格式) + +```json +📋 请求头: { + "content-type": "application/json", + "authorization": "***REDACTED***", + "x-api-key": "***REDACTED***", + "user-agent": "python-httpx/0.25.2", + "accept": "application/json", + "content-length": "1024", + "host": "localhost:8000" +} +``` + +### 3. 请求体日志输出(JSON格式) + +```json +📄 请求体: { + "count": 2, + "data": [ + { + "bill": 100, + "duration": 60, + "callid": "test-call-123", + "calldate": "2024-01-01 12:00:00", + "number": "13800138000", + "numberid": "number-123", + "customer_id": "customer-123", + "status": 0, + "status_str": "失败", + "user_id": "user-123", + "type": 1, + "number_data": { + "number": "13800138000", + "province": "广东", + "city": "深圳", + "operator": "移动" + }, + "group": { + "id": 1, + "name": "测试组" + }, + "task": { + "id": "task-123", + "name": "测试任务" + }, + "user": { + "id": "user-123", + "name": "测试用户" + }, + "customer_data": { + "name": "测试客户", + "email": "test@example.com", + "company": "测试公司", + "extra": "额外信息" + } + } + ] +} +``` + +### 4. server_ip 日志输出(JSON格式) + +```json +🏠 server_ip: { + "ip": "0.0.0.0", + "port": 8000 +} +``` + +### 5. client_ip 日志输出(JSON格式) + +```json +🖥️ client_ip: { + "ip": "127.0.0.1", + "port": 52345 +} +``` + +### 6. callback_data 日志输出(JSON格式) + +```json +📦 callback_data: { + "count": 2, + "data_count": 1, + "data_sample": { + "bill": 100, + "duration": 60, + "callid": "test-call-123", + "calldate": "2024-01-01 12:00:00", + "number": "13800138000", + "numberid": "number-123", + "customer_id": "customer-123", + "status": 0, + "status_str": "失败", + "user_id": "user-123", + "type": 1, + "number_data": { + "number": "13800138000", + "province": "广东", + "city": "深圳", + "operator": "移动" + }, + "group": { + "id": 1, + "name": "测试组" + }, + "task": { + "id": "task-123", + "name": "测试任务" + }, + "user": { + "id": "user-123", + "name": "测试用户" + }, + "customer_data": { + "name": "测试客户", + "email": "test@example.com", + "company": "测试公司", + "extra": "额外信息" + } + } +} +``` + +## 完整的请求处理日志流程 + +### 单次请求的完整日志输出 + +``` +2024-12-02 16:30:15 - routes - INFO - 🔥 收到AI Talk回调请求: siteId=test-site-123, count=2, data_count=1 +2024-12-02 16:30:15 - routes - INFO - 📝 site_id: "test-site-123" +2024-12-02 16:30:15 - routes - INFO - 📋 请求头: { + "content-type": "application/json", + "authorization": "***REDACTED***", + "user-agent": "python-httpx/0.25.2", + "accept": "application/json", + "content-length": "1024", + "host": "localhost:8000" +} +2024-12-02 16:30:15 - routes - INFO - 📄 请求体: { + "count": 2, + "data": [ + { + "bill": 100, + "duration": 60, + "callid": "test-call-123", + "calldate": "2024-01-01 12:00:00", + "number": "13800138000", + "numberid": "number-123", + "customer_id": "customer-123", + "status": 0, + "status_str": "失败", + "user_id": "user-123", + "type": 1, + "number_data": { + "number": "13800138000", + "province": "广东", + "city": "深圳", + "operator": "移动" + }, + "group": { + "id": 1, + "name": "测试组" + }, + "task": { + "id": "task-123", + "name": "测试任务" + }, + "user": { + "id": "user-123", + "name": "测试用户" + }, + "customer_data": { + "name": "测试客户", + "email": "test@example.com", + "company": "测试公司", + "extra": "额外信息" + } + } + ] +} +2024-12-02 16:30:15 - routes - INFO - 🏠 server_ip: { + "ip": "0.0.0.0", + "port": 8000 +} +2024-12-02 16:30:15 - routes - INFO - 🖥️ client_ip: { + "ip": "127.0.0.1", + "port": 52345 +} +2024-12-02 16:30:15 - routes - INFO - 📦 callback_data: { + "count": 2, + "data_count": 1, + "data_sample": { + "bill": 100, + "duration": 60, + "callid": "test-call-123", + "calldate": "2024-01-01 12:00:00", + "number": "13800138000", + "numberid": "number-123", + "customer_id": "customer-123", + "status": 0, + "status_str": "失败", + "user_id": "user-123", + "type": 1, + "number_data": { + "number": "13800138000", + "province": "广东", + "city": "深圳", + "operator": "移动" + }, + "group": { + "id": 1, + "name": "测试组" + }, + "task": { + "id": "task-123", + "name": "测试任务" + }, + "user": { + "id": "user-123", + "name": "测试用户" + }, + "customer_data": { + "name": "测试客户", + "email": "test@example.com", + "company": "测试公司", + "extra": "额外信息" + } + } +} +2024-12-02 16:30:15 - routes - INFO - ✅ 回调请求记录成功保存到数据库,ID: 123 +2024-12-02 16:30:15 - routes - INFO - 📞 count=2 < 3,调用外部API +2024-12-02 16:30:15 - routes - DEBUG - 🔒 获取Redis锁成功: external_api_call_test-site-123_2 +``` + +## JSON格式日志的优势 + +### 1. 结构化数据 +- 每个字段都以标准JSON格式输出 +- 便于程序解析和处理 +- 支持复杂的嵌套数据结构 + +### 2. 可读性强 +- JSON格式具有良好的层次结构 +- 缩进格式便于人工阅读 +- 支持中文字符(`ensure_ascii=False`) + +### 3. 便于分析 +- 可以直接使用JSON工具解析 +- 支持日志分析工具(如ELK、Fluentd等) +- 便于数据提取和统计 + +### 4. 安全性 +- 敏感信息自动过滤为 `***REDACTED***` +- 保持数据结构完整性 +- 避免敏感信息泄露 + +## 日志分析示例 + +### 使用jq工具分析JSON日志 + +```bash +# 提取所有site_id +grep "📝 site_id:" logs/app.log | jq -r '.📝 site_id' + +# 提取所有client_ip信息 +grep "🖥️ client_ip:" logs/app.log | jq -r '.🖥️ client_ip.ip' + +# 统计不同count值的请求 +grep "📦 callback_data:" logs/app.log | jq -r '.📦 callback_data.count' | sort | uniq -c + +# 提取包含特定号码的请求 +grep "📄 请求体:" logs/app.log | jq 'select(.📄 请求_body.data[].number == "13800138000")' +``` + +### 使用Python分析JSON日志 + +```python +import json +import re + +# 解析日志中的JSON数据 +def parse_json_logs(log_file): + site_ids = [] + client_ips = [] + + with open(log_file, 'r', encoding='utf-8') as f: + for line in f: + if '📝 site_id:' in line: + # 提取JSON部分 + json_str = line.split('📝 site_id: ')[1].strip() + site_id = json.loads(json_str) + site_ids.append(site_id) + + elif '🖥️ client_ip:' in line: + json_str = line.split('🖥️ client_ip: ')[1].strip() + client_info = json.loads(json_str) + client_ips.append(client_info['ip']) + + return site_ids, client_ips +``` + +## 性能考虑 + +### 1. JSON序列化开销 +- 使用标准库`json.dumps()` +- `ensure_ascii=False` 支持中文但略慢 +- 对于高频调用,可考虑关闭详细日志 + +### 2. 日志文件大小 +- JSON格式比纯文本占用更多空间 +- 建议合理设置日志轮转大小 +- 生产环境可考虑使用压缩存储 + +### 3. 内存使用 +- 大型请求体会占用较多内存 +- `callback_data` 只记录摘要信息,避免完整数据 + +## 配置建议 + +### 开发环境 +```bash +LOG_LEVEL=INFO # 显示所有JSON日志 +``` + +### 生产环境 +```bash +LOG_LEVEL=WARNING # 只显示重要信息,减少JSON日志量 +``` + +### 调试特定问题 +```bash +# 临时开启详细日志 +LOG_LEVEL=DEBUG +# 问题解决后恢复 +LOG_LEVEL=INFO +``` + +现在 `log_callback_request` 方法以JSON格式输出所有关键信息,便于日志分析和系统监控! \ No newline at end of file diff --git a/LOG_GUIDE.md b/LOG_GUIDE.md new file mode 100644 index 0000000..7fe965a --- /dev/null +++ b/LOG_GUIDE.md @@ -0,0 +1,268 @@ +# 日志系统使用指南 + +## 概述 + +本项目集成了完整的日志系统,支持多级别日志记录、文件轮转、彩色输出等功能。 + +## 日志配置 + +### 环境变量配置 + +在 `.env` 文件中配置以下日志相关参数: + +```bash +# 日志级别: DEBUG, INFO, WARNING, ERROR, CRITICAL +LOG_LEVEL=INFO + +# 日志文件路径 +LOG_FILE=logs/app.log + +# 单个日志文件最大大小(字节) +LOG_MAX_BYTES=10485760 # 10MB + +# 备份文件数量 +LOG_BACKUP_COUNT=5 + +# 日志格式 +LOG_FORMAT=%(asctime)s - %(name)s - %(levelname)s - %(message)s + +# 时间格式 +LOG_DATE_FORMAT=%Y-%m-%d %H:%M:%S +``` + +### 日志级别说明 + +- **DEBUG**: 详细的调试信息,通常只在开发时使用 +- **INFO**: 一般信息,记录应用正常运行状态 +- **WARNING**: 警告信息,表示可能出现问题 +- **ERROR**: 错误信息,表示发生了错误但不影响应用运行 +- **CRITICAL**: 严重错误,可能导致应用崩溃 + +## 使用方法 + +### 1. 获取日志器 + +```python +from app.logger import get_logger + +# 方式1: 指定名称 +logger = get_logger("my_module") + +# 方式2: 自动获取模块名 +logger = get_logger() # 自动获取调用者的模块名 +``` + +### 2. 记录日志 + +```python +logger.debug("这是调试信息") +logger.info("这是一般信息") +logger.warning("这是警告信息") +logger.error("这是错误信息") +logger.critical("这是严重错误信息") + +# 记录异常信息(包含堆栈跟踪) +try: + # 可能出错的代码 + result = 1 / 0 +except Exception as e: + logger.error(f"计算失败: {e}", exc_info=True) +``` + +### 3. 使用预定义日志器 + +```python +from app.logger import ( + app_logger, # 应用级日志 + api_logger, # API接口日志 + db_logger, # 数据库操作日志 + redis_logger, # Redis操作日志 + external_api_logger # 外部API调用日志 +) + +# 使用示例 +api_logger.info("API接口被访问") +db_logger.info("数据库操作完成") +``` + +## 日志输出 + +### 控制台输出 + +控制台输出支持彩色显示,不同级别的日志使用不同颜色: + +- 🔵 DEBUG: 青色 +- 🟢 INFO: 绿色 +- 🟡 WARNING: 黄色 +- 🔴 ERROR: 红色 +- 🟣 CRITICAL: 紫色 + +### 文件输出 + +日志文件保存在配置的路径中,支持自动轮转: + +- 当文件大小超过 `LOG_MAX_BYTES` 时,会自动创建备份文件 +- 备份文件命名格式:`app.log.1`, `app.log.2`, ... +- 最多保留 `LOG_BACKUP_COUNT` 个备份文件 + +## 项目中的日志记录 + +### 1. 应用启动/关闭日志 + +```python +# 应用启动 +logger.info("🚀 应用启动中...") +logger.info("✅ 数据库连接成功") +logger.info("🎉 应用启动完成!") + +# 应用关闭 +logger.info("🛑 应用关闭中...") +logger.info("👋 应用已关闭") +``` + +### 2. API接口日志 + +```python +# 请求接收 +logger.info(f"🔥 收到AI Talk回调请求: siteId={siteId}, count={count}") + +# 业务逻辑 +logger.info(f"✅ count={count} >= {threshold},直接返回") +logger.info(f"📞 调用外部API,重试次数: {retry_count}") + +# 错误处理 +logger.error(f"❌ 服务器内部错误: {e}", exc_info=True) +``` + +### 3. 数据库操作日志 + +```python +# 数据库初始化 +logger.info("📊 初始化数据库表结构...") +logger.info("✅ 数据库表结构初始化完成") + +# 数据记录 +logger.info(f"📝 记录回调请求: siteId={siteId}, count={count}") +``` + +### 4. Redis操作日志 + +```python +# 锁操作 +logger.debug(f"🔒 尝试获取Redis锁: {key}") +logger.debug(f"✅ Redis锁获取成功: {key}") +``` + +### 5. 外部API调用日志 + +```python +# API调用 +logger.info(f"🌐 开始调用外部API: {url}") +logger.debug(f"📤 第{attempt + 1}次尝试调用外部API") +logger.info(f"✅ 外部API调用成功,状态码: {status}") +logger.warning(f"⚠️ 外部API返回错误状态码: {status}") +``` + +## 日志查看和分析 + +### 1. 实时查看日志 + +```bash +# 查看最新日志 +tail -f logs/app.log + +# 查看带颜色的日志(如果支持) +tail -f logs/app.log | ccze # 需要安装ccze +``` + +### 2. 搜索日志 + +```bash +# 搜索错误日志 +grep "ERROR" logs/app.log + +# 搜索特定接口的日志 +grep "siteId=123" logs/app.log + +# 搜索特定时间范围的日志 +grep "2024-12-02 10:" logs/app.log +``` + +### 3. 日志统计 + +```bash +# 统计不同级别的日志数量 +grep -c "INFO" logs/app.log +grep -c "ERROR" logs/app.log +grep -c "WARNING" logs/app.log + +# 查看最频繁的错误 +grep "ERROR" logs/app.log | sort | uniq -c | sort -nr +``` + +## 最佳实践 + +### 1. 日志级别使用建议 + +- **生产环境**: 使用 INFO 或 WARNING 级别 +- **开发环境**: 使用 DEBUG 级别查看详细信息 +- **测试环境**: 使用 INFO 级别 + +### 2. 日志内容建议 + +- 包含关键业务参数(如 siteId, count) +- 使用表情符号增强可读性 +- 错误日志包含完整的异常信息 +- 重要操作记录开始和结束状态 + +### 3. 性能考虑 + +- 避免在高频循环中记录 DEBUG 日志 +- 使用日志级别控制避免不必要的字符串格式化 +- 合理设置日志轮转大小和数量 + +### 4. 安全考虑 + +- 避免在日志中记录敏感信息(密码、密钥等) +- 生产环境注意日志文件的访问权限 +- 定期清理旧的日志文件 + +## 故障排查 + +### 1. 常见问题 + +**问题**: 日志文件没有创建 +- 检查日志目录是否存在且有写权限 +- 检查 LOG_FILE 配置是否正确 + +**问题**: 日志没有输出到控制台 +- 检查 LOG_LEVEL 配置是否过高 +- 检查是否有其他处理器冲突 + +**问题**: 日志轮转不工作 +- 检查 LOG_MAX_BYTES 设置是否合理 +- 检查文件系统权限 + +### 2. 调试日志系统 + +```python +# 检查日志器配置 +import logging +logger = logging.getLogger() +print(f"日志级别: {logger.level}") +print(f"处理器数量: {len(logger.handlers)}") +for handler in logger.handlers: + print(f"处理器: {handler.__class__.__name__}") +``` + +## 测试 + +运行日志系统测试: + +```bash +# 运行日志测试 +python -m pytest test_logger.py -v + +# 运行特定测试 +python -m pytest test_logger.py::TestLogger::test_logger_setup -v +``` \ No newline at end of file diff --git a/app/config.py b/app/config.py index 4ef7b29..e47374e 100644 --- a/app/config.py +++ b/app/config.py @@ -16,6 +16,14 @@ class Settings(BaseSettings): external_api_retry_max: int = 4 # 外部API最大重试次数 external_api_url: str = "https://external.com/api/openapi/customerApi/aiTaskResultFail" + # 日志配置 + log_level: str = "INFO" # DEBUG, INFO, WARNING, ERROR, CRITICAL + log_file: str = "logs/app.log" # 日志文件路径 + log_max_bytes: int = 10 * 1024 * 1024 # 10MB + log_backup_count: int = 5 # 备份文件数量 + log_format: str = "%(asctime)s - %(name)s - %(levelname)s - %(message)s" + log_date_format: str = "%Y-%m-%d %H:%M:%S" + # 应用配置 app_name: str = "AI Talk Callback API" debug: bool = False diff --git a/app/database.py b/app/database.py index 874bbba..fe62e14 100644 --- a/app/database.py +++ b/app/database.py @@ -3,6 +3,9 @@ from sqlalchemy.orm import DeclarativeBase from sqlalchemy import Column, String, Integer, DateTime, Text, JSON from datetime import datetime from app.config import settings +from app.logger import get_logger + +logger = get_logger("database") class Base(DeclarativeBase): @@ -60,5 +63,11 @@ async def get_db(): async def init_db(): - async with engine.begin() as conn: - await conn.run_sync(Base.metadata.create_all) \ No newline at end of file + logger.info("📊 初始化数据库表结构...") + try: + async with engine.begin() as conn: + await conn.run_sync(Base.metadata.create_all) + logger.info("✅ 数据库表结构初始化完成") + except Exception as e: + logger.error(f"❌ 数据库初始化失败: {e}") + raise \ No newline at end of file diff --git a/app/redis_lock.py b/app/redis_lock.py index 6ccdfef..6d4f18b 100644 --- a/app/redis_lock.py +++ b/app/redis_lock.py @@ -3,6 +3,9 @@ import asyncio import uuid from typing import Optional from app.config import settings +from app.logger import get_logger + +logger = get_logger("redis") class RedisLock: @@ -15,6 +18,8 @@ class RedisLock: async def acquire(self) -> bool: """获取分布式锁""" + logger.debug(f"🔒 尝试获取Redis锁: {self.key}") + lua_script = """ if redis.call("GET", KEYS[1]) == false then return redis.call("SETEX", KEYS[1], ARGV[1], ARGV[2]) @@ -32,6 +37,11 @@ class RedisLock: ) self.acquired = bool(result) + if self.acquired: + logger.debug(f"✅ Redis锁获取成功: {self.key}") + else: + logger.debug(f"❌ Redis锁获取失败: {self.key}") + return self.acquired async def release(self) -> bool: diff --git a/app/routes.py b/app/routes.py index b453afa..805979e 100644 --- a/app/routes.py +++ b/app/routes.py @@ -8,6 +8,9 @@ from app.database import get_db, CallbackLog, ExternalApiLog from app.models import CallbackRequest, CallbackResponse from app.redis_lock import redis_manager from app.config import settings +from app.logger import get_logger + +logger = get_logger("routes") router = APIRouter() @@ -19,25 +22,85 @@ async def log_callback_request( callback_data: CallbackRequest ): """记录回调请求到数据库""" + import json + # 获取客户端IP地址 client_ip = request.client.host if request.client else None + client_port = request.client.port if request.client else None # 获取服务器IP地址 server_ip = None + server_port = None if hasattr(request, 'scope') and 'server' in request.scope: - server_host, server_port = request.scope['server'] + server_host, server_port_info = request.scope['server'] server_ip = server_host + server_port = server_port_info - callback_log = CallbackLog( - site_id=site_id, # 记录siteId - remote_address=client_ip, - server_ip=server_ip, - request_url=str(request.url), - request_headers=dict(request.headers), - request_body=callback_data.model_dump() - ) - db.add(callback_log) - await db.commit() + # 准备请求头信息(过滤敏感信息) + request_headers = dict(request.headers) + safe_headers = {} + sensitive_headers = {'authorization', 'token', 'api-key', 'x-api-key', 'cookie'} + + for key, value in request_headers.items(): + if key.lower() in sensitive_headers: + safe_headers[key] = "***REDACTED***" + else: + safe_headers[key] = value + + # 准备请求体信息 + request_body = callback_data.model_dump() + + # 记录site_id(JSON格式) + logger.info(f"📝 site_id: {json.dumps(site_id, ensure_ascii=False)}") + + # 记录请求头(JSON格式) + logger.info(f"📋 请求头: {json.dumps(safe_headers, ensure_ascii=False, indent=2)}") + + # 记录请求体(JSON格式) + logger.info(f"📄 请求体: {json.dumps(request_body, ensure_ascii=False, indent=2)}") + + # 记录server_ip(JSON格式) + server_info = { + "ip": server_ip, + "port": server_port + } + logger.info(f"🏠 server_ip: {json.dumps(server_info, ensure_ascii=False)}") + + # 记录client_ip(JSON格式) + client_info = { + "ip": client_ip, + "port": client_port + } + logger.info(f"🖥️ client_ip: {json.dumps(client_info, ensure_ascii=False)}") + + # 记录callback_data(JSON格式) + callback_data_json = { + "count": callback_data.count, + "data_count": len(callback_data.data) if callback_data.data else 0, + "data_sample": callback_data.data[0].model_dump() if callback_data.data else None + } + logger.info(f"📦 callback_data: {json.dumps(callback_data_json, ensure_ascii=False, indent=2)}") + + try: + # 保存到数据库 + callback_log = CallbackLog( + site_id=site_id, + remote_address=f"{client_ip}:{client_port}" if client_ip and client_port else client_ip, + server_ip=f"{server_ip}:{server_port}" if server_ip and server_port else server_ip, + request_url=str(request.url), + request_headers=safe_headers, # 保存过滤后的请求头 + request_body=request_body + ) + + db.add(callback_log) + await db.commit() + + logger.info(f"✅ 回调请求记录成功保存到数据库,ID: {callback_log.id}") + + except Exception as e: + logger.error(f"❌ 保存回调请求到数据库失败: {e}", exc_info=True) + # 不重新抛出异常,避免影响主业务流程 + pass async def log_external_api_request( @@ -83,6 +146,8 @@ async def call_external_api_with_retry( if max_retries is None: max_retries = settings.external_api_retry_max + logger.info(f"🌐 开始调用外部API: {settings.external_api_url}, 最大重试次数: {max_retries}") + headers = { "Content-Type": "application/json", "User-Agent": "AITalkCallbackService/1.0" @@ -90,6 +155,8 @@ async def call_external_api_with_retry( for attempt in range(max_retries + 1): # +1 因为第一次不算重试 try: + logger.debug(f"📤 第{attempt + 1}次尝试调用外部API") + async with httpx.AsyncClient(timeout=30.0) as client: response = await client.post( settings.external_api_url, @@ -97,6 +164,8 @@ async def call_external_api_with_retry( json=request_body ) + logger.debug(f"📥 外部API响应: status={response.status_code}") + # 记录每次尝试的结果 await log_external_api_request( db=db, @@ -111,15 +180,23 @@ async def call_external_api_with_retry( # 检查响应状态 if response.status_code < 400: + logger.info(f"✅ 外部API调用成功,状态码: {response.status_code}") return True, attempt else: + logger.warning(f"⚠️ 外部API返回错误状态码: {response.status_code}") + # 如果是最后一次尝试,直接返回失败 if attempt == max_retries: + logger.error(f"❌ 外部API调用最终失败,状态码: {response.status_code}") return False, attempt # 否则等待一段时间后重试 - await asyncio.sleep(1 * (attempt + 1)) # 递增延迟 + wait_time = 1 * (attempt + 1) # 递增延迟 + logger.info(f"⏳ 等待 {wait_time}秒 后重试...") + await asyncio.sleep(wait_time) except httpx.RequestError as e: + logger.warning(f"⚠️ 外部API网络错误 (第{attempt + 1}次尝试): {e}") + # 记录网络错误 await log_external_api_request( db=db, @@ -134,9 +211,12 @@ async def call_external_api_with_retry( # 如果是最后一次尝试,直接返回失败 if attempt == max_retries: + logger.error(f"❌ 外部API调用最终失败: {e}") return False, attempt # 否则等待一段时间后重试 - await asyncio.sleep(1 * (attempt + 1)) # 递增延迟 + wait_time = 1 * (attempt + 1) # 递增延迟 + logger.info(f"⏳ 等待 {wait_time}秒 后重试...") + await asyncio.sleep(wait_time) return False, max_retries @@ -151,12 +231,15 @@ async def ai_talk_callback( """ AI Talk回调接口处理 """ + logger.info(f"🔥 收到AI Talk回调请求: siteId={siteId}, count={callback_data.count}, data_count={len(callback_data.data)}") + try: # 记录回调请求(包含siteId) await log_callback_request(db, request, siteId, callback_data) # 判断count是否大于等于阈值,如果是直接返回 if callback_data.count >= settings.count_threshold: + logger.info(f"✅ count={callback_data.count} >= {settings.count_threshold},直接返回") return CallbackResponse( success=True, message=f"count={callback_data.count} >= {settings.count_threshold},直接返回", @@ -165,8 +248,12 @@ async def ai_talk_callback( site_id=siteId ) + logger.info(f"📞 count={callback_data.count} < {settings.count_threshold},调用外部API") + # count < 3,调用外部API,使用分布式锁防止并发调用 async with redis_manager.create_lock(f"external_api_call_{siteId}_{callback_data.count}"): + logger.debug(f"🔒 获取Redis锁成功: external_api_call_{siteId}_{callback_data.count}") + # 调用外部API并支持重试 request_body = callback_data.model_dump() success, retry_count = await call_external_api_with_retry( @@ -176,6 +263,7 @@ async def ai_talk_callback( ) if success: + logger.info(f"✅ 外部API调用成功,重试次数: {retry_count}") return CallbackResponse( success=True, message="外部API调用成功", @@ -184,6 +272,7 @@ async def ai_talk_callback( site_id=siteId ) else: + logger.error(f"❌ 外部API调用失败,已重试{retry_count}次") raise HTTPException( status_code=500, detail=f"外部API调用失败,已重试{retry_count}次" @@ -192,6 +281,7 @@ async def ai_talk_callback( except HTTPException: raise except Exception as e: + logger.error(f"💥 服务器内部错误: {e}", exc_info=True) raise HTTPException( status_code=500, detail=f"服务器内部错误: {str(e)}" diff --git a/main.py b/main.py index ad9dde8..9aca2fd 100644 --- a/main.py +++ b/main.py @@ -6,16 +6,46 @@ from app.config import settings from app.database import init_db from app.redis_lock import redis_manager from app.routes import router +from app.logger import get_logger, LoggerManager +# 初始化日志系统 +LoggerManager.setup_logging() +logger = get_logger("main") + @asynccontextmanager async def lifespan(app: FastAPI): # 启动时初始化 - await init_db() - await redis_manager.connect() - yield - # 关闭时清理 - await redis_manager.disconnect() + logger.info("🚀 应用启动中...") + + try: + # 初始化数据库 + logger.info("📊 初始化数据库连接...") + await init_db() + logger.info("✅ 数据库连接成功") + + # 初始化Redis + logger.info("🔴 初始化Redis连接...") + await redis_manager.connect() + logger.info("✅ Redis连接成功") + + logger.info(f"🎉 {settings.app_name} 启动完成!") + yield + + except Exception as e: + logger.error(f"❌ 应用启动失败: {e}") + raise + + finally: + # 关闭时清理 + logger.info("🛑 应用关闭中...") + try: + await redis_manager.disconnect() + logger.info("✅ Redis连接已关闭") + except Exception as e: + logger.error(f"❌ Redis关闭时出错: {e}") + + logger.info("👋 应用已关闭") app = FastAPI( @@ -39,11 +69,13 @@ app.include_router(router) @app.get("/") async def root(): + logger.info("📝 根接口被访问") return {"message": f"Welcome to {settings.app_name}"} @app.get("/health") async def health_check(): + logger.debug("💓 健康检查接口被访问") return {"status": "healthy"} diff --git a/test_callback.py b/test_callback.py index cbae7b5..44dce0d 100644 --- a/test_callback.py +++ b/test_callback.py @@ -340,6 +340,72 @@ class TestAITalkCallback: # 验证新字段存在(可能为None,因为测试环境) assert hasattr(call_args, 'remote_address') assert hasattr(call_args, 'server_ip') + + # 验证JSON格式数据结构 + assert isinstance(call_args.request_headers, dict) + assert isinstance(call_args.request_body, dict) + assert 'count' in call_args.request_body + assert 'data' in call_args.request_body + + @patch('app.routes.get_db') + @patch('app.routes.redis_manager') + def test_json_format_logging( + self, + mock_redis_manager, + mock_get_db, + sample_callback_request + ): + """测试JSON格式日志记录功能""" + from app.database import CallbackLog + + # 设置模拟对象 + mock_db = AsyncMock() + mock_get_db.return_value = mock_db + + # 模拟Redis锁 + mock_lock = AsyncMock() + mock_lock.__aenter__ = AsyncMock(return_value=None) + mock_lock.__aexit__ = AsyncMock(return_value=None) + mock_redis_manager.create_lock.return_value = mock_lock + + # 模拟外部API调用成功 + with patch('app.routes.call_external_api_with_retry') as mock_api_call: + mock_api_call.return_value = (True, 0) + + response = client.post( + "/ai-talk/callback/test-site-json/failure", + json=sample_callback_request.model_dump() + ) + + assert response.status_code == 200 + + # 验证log_callback_request被调用 + assert mock_db.add.called + + # 获取传递给add的CallbackLog对象 + call_args = mock_db.add.call_args[0][0] + assert isinstance(call_args, CallbackLog) + + # 验证数据结构适合JSON序列化 + import json + + # 验证请求头可以JSON序列化 + try: + headers_json = json.dumps(call_args.request_headers, ensure_ascii=False) + assert isinstance(headers_json, str) + except (TypeError, ValueError) as e: + pytest.fail(f"请求头无法JSON序列化: {e}") + + # 验证请求体可以JSON序列化 + try: + body_json = json.dumps(call_args.request_body, ensure_ascii=False) + assert isinstance(body_json, str) + except (TypeError, ValueError) as e: + pytest.fail(f"请求体无法JSON序列化: {e}") + + # 验证关键字段存在 + assert call_args.site_id == "test-site-json" + assert call_args.request_body['count'] == sample_callback_request.count if __name__ == "__main__": diff --git a/test_logger.py b/test_logger.py new file mode 100644 index 0000000..b133982 --- /dev/null +++ b/test_logger.py @@ -0,0 +1,167 @@ +import pytest +import logging +import tempfile +import os +from pathlib import Path +from unittest.mock import patch + +from app.logger import LoggerManager, get_logger +from app.config import settings + + +class TestLogger: + """日志系统测试用例""" + + def setup_method(self): + """每个测试方法执行前的设置""" + # 重置日志管理器状态 + LoggerManager._initialized = False + LoggerManager._loggers.clear() + + def test_logger_setup(self): + """测试日志系统初始化""" + # 使用临时目录 + with tempfile.TemporaryDirectory() as temp_dir: + temp_log_file = os.path.join(temp_dir, "test.log") + + with patch.object(settings, 'log_file', temp_log_file): + # 初始化日志系统 + LoggerManager.setup_logging() + + # 验证日志文件是否创建 + assert os.path.exists(temp_log_file) + + # 验证根日志器配置 + root_logger = logging.getLogger() + assert len(root_logger.handlers) > 0 + + def test_get_logger(self): + """测试获取日志器""" + logger = get_logger("test") + + assert isinstance(logger, logging.Logger) + assert logger.name == "test" + + def test_get_logger_auto_name(self): + """测试自动获取日志器名称""" + logger = get_logger() + + assert isinstance(logger, logging.Logger) + # 应该获取到调用模块的名称 + assert "test_logger" in logger.name + + def test_multiple_loggers(self): + """测试多个日志器""" + logger1 = get_logger("module1") + logger2 = get_logger("module2") + logger3 = get_logger("module1") # 重复获取 + + assert logger1.name == "module1" + assert logger2.name == "module2" + assert logger3.name == "module1" + assert logger1 is logger3 # 应该是同一个实例 + + def test_logger_levels(self): + """测试不同日志级别""" + with tempfile.TemporaryDirectory() as temp_dir: + temp_log_file = os.path.join(temp_dir, "test.log") + + with patch.object(settings, 'log_file', temp_log_file): + LoggerManager.setup_logging() + + logger = get_logger("test_levels") + + # 测试不同级别的日志 + logger.debug("Debug message") + logger.info("Info message") + logger.warning("Warning message") + logger.error("Error message") + logger.critical("Critical message") + + # 验证日志文件内容 + with open(temp_log_file, 'r', encoding='utf-8') as f: + log_content = f.read() + + assert "Debug message" in log_content + assert "Info message" in log_content + assert "Warning message" in log_content + assert "Error message" in log_content + assert "Critical message" in log_content + + def test_log_rotation(self): + """测试日志轮转功能""" + with tempfile.TemporaryDirectory() as temp_dir: + temp_log_file = os.path.join(temp_dir, "test.log") + + # 设置很小的文件大小限制以触发轮转 + with patch.object(settings, 'log_file', temp_log_file), \ + patch.object(settings, 'log_max_bytes', 100): + + LoggerManager.setup_logging() + + logger = get_logger("test_rotation") + + # 写入大量日志以触发轮转 + for i in range(50): + logger.info(f"This is a long log message to trigger rotation {i}") + + # 验证是否创建了备份文件 + backup_file = f"{temp_log_file}.1" + # 注意:由于异步写入,可能需要等待一下 + import time + time.sleep(0.1) + + def test_colored_formatter(self): + """测试彩色格式化器""" + from app.logger import ColoredFormatter + + formatter = ColoredFormatter("%(levelname)s - %(message)s") + + # 创建测试记录 + record = logging.LogRecord( + name="test", + level=logging.INFO, + pathname="", + lineno=0, + msg="Test message", + args=(), + exc_info=None + ) + + formatted = formatter.format(record) + + # 验证是否包含颜色代码 + assert "\033[" in formatted # ANSI颜色代码 + + def test_predefined_loggers(self): + """测试预定义的日志器""" + from app.logger import ( + app_logger, api_logger, db_logger, + redis_logger, external_api_logger + ) + + assert app_logger.name == "app" + assert api_logger.name == "api" + assert db_logger.name == "database" + assert redis_logger.name == "redis" + assert external_api_logger.name == "external_api" + + def test_logger_initialization_idempotency(self): + """测试日志系统初始化的幂等性""" + with tempfile.TemporaryDirectory() as temp_dir: + temp_log_file = os.path.join(temp_dir, "test.log") + + with patch.object(settings, 'log_file', temp_log_file): + # 多次初始化 + LoggerManager.setup_logging() + LoggerManager.setup_logging() + LoggerManager.setup_logging() + + # 验证处理器数量不会重复增加 + root_logger = logging.getLogger() + # 应该有2个处理器(控制台+文件) + assert len(root_logger.handlers) == 2 + + +if __name__ == "__main__": + pytest.main([__file__, "-v"]) \ No newline at end of file