接口请求体改用dict处理

保存每次通话信息到库
This commit is contained in:
mark.tian
2025-12-03 21:56:06 +08:00
parent 70444d5665
commit 703408c808
3 changed files with 65 additions and 317 deletions

View File

@@ -1,15 +1,12 @@
from fastapi import APIRouter, Request, HTTPException, Depends, Path
from fastapi import APIRouter, Request, Body, HTTPException, Depends, Path
from sqlalchemy.ext.asyncio import AsyncSession
from sqlalchemy import func, select
import httpx
import asyncio
from typing import Dict, Any, Optional
from typing import Optional
from app.database import get_db, CallbackLog, ExternalApiLog, CallbackData
from app.models import CallbackRequest, CallbackResponse
from app.redis_lock import redis_manager
from app.database import get_db, CallbackLog, CallbackData
from app.models import CallbackResponse
from app.config import settings
from app.logger import get_logger
from app.external_api_processor import process_external_api_call
logger = get_logger("routes")
@@ -19,7 +16,8 @@ router = APIRouter()
async def log_callback_request(
db: AsyncSession,
request: Request,
site_id: str
site_id: str,
callback_data: dict
) -> Optional[int]:
"""记录回调请求到数据库,返回记录ID"""
import json
@@ -38,21 +36,9 @@ async def log_callback_request(
# 准备请求头信息(直接记录原始请求头)
request_headers = dict(request.headers)
# 获取原始请求体
import json as json_module
try:
# 获取原始请求体字节并解码为字符串
body_bytes = await request.body()
body_str = body_bytes.decode('utf-8')
# 尝试解析为JSON对象,如果失败则使用原始字符串
try:
request_body = json_module.loads(body_str)
except json_module.JSONDecodeError:
request_body = body_str
except Exception as e:
logger.warning(f"⚠️ 读取请求体失败: {e}")
request_body = "Unable to read request body"
# 准备请求体信息(直接记录原始请求体)
request_body = callback_data
# 记录请求URL(JSON格式)
logger.info(f"🌐 请求URL: {request.url}")
@@ -78,7 +64,7 @@ async def log_callback_request(
logger.info(f"📋 请求头: {json.dumps(request_headers, ensure_ascii=False, indent=2)}")
# 记录请求体(JSON格式)
logger.info(f"📄 请求体: {json.dumps(request_body, ensure_ascii=False, indent=2)}")
logger.info(f"📄 请求体: {json.dumps(callback_data, ensure_ascii=False, indent=2)}")
try:
# 保存到数据库
@@ -105,63 +91,39 @@ async def log_callback_request(
return None
async def check_phone_number_threshold(
db: AsyncSession,
phone_number: str
) -> tuple[bool, bool]:
"""
检查手机号在数据库中的出现次数是否超过阈值
Args:
db: 数据库会话
phone_number: 要检查的手机号
Returns:
tuple[是否超过阈值, 是否查询失败]
"""
try:
# 查询手机号在数据库中出现的次数
count_query = select(func.count(CallbackData.id)).where(
CallbackData.phone_number == phone_number
)
result = await db.execute(count_query)
phone_count = result.scalar() or 0
logger.info(f"📊 手机号 {phone_number} 在数据库中出现次数: {phone_count}")
# 判断是否超过阈值
exceeds_threshold = phone_count >= settings.count_threshold
if exceeds_threshold:
logger.info(f"✅ 手机号 {phone_number} 出现次数 {phone_count} >= {settings.count_threshold},超过阈值")
return exceeds_threshold, False
except Exception as e:
logger.error(f"❌ 查询手机号 {phone_number} 失败: {e}", exc_info=True)
# 查询失败时返回False,不影响主流程
return False, True
async def save_callback_data_items(
db: AsyncSession,
callback_data_items: list
):
"""保存callback_data.data中的数据到数据库"""
import json
try:
callback_data_records = []
for item in callback_data_items:
if not hasattr(item, 'number_data') or not item.number_data:
continue
if not hasattr(item.number_data, 'number') or not item.number_data.number:
# 获取手机号
number_data = item.get('number_data', {})
phone_number = number_data.get('number')
if not phone_number:
continue
# 获取任务ID
task = item.get('task', {})
task_id = task.get('id', '')
# 获取状态信息
status = item.get('status', 0)
status_description = item.get('status_str', '')
# 将整个item转换为JSON字符串保存
raw_data_json = json.dumps(item, ensure_ascii=False)
callback_data_record = CallbackData(
phone_number=item.number_data.number,
task_id=item.task.id,
status=item.status,
status_description=item.status_str
phone_number=phone_number,
task_id=task_id,
status=status,
status_description=status_description,
raw_data=raw_data_json # 保存原始JSON字符串
)
callback_data_records.append(callback_data_record)
@@ -176,133 +138,10 @@ async def save_callback_data_items(
# 不重新抛出异常,避免影响主业务流程
async def log_external_api_request(
db: AsyncSession,
callback_logs_id: int,
request_url: str,
request_headers: Dict[str, Any],
request_body: Dict[str, Any],
response_status: Optional[int],
response_headers: Optional[Dict[str, Any]],
response_body: Optional[str],
retry_count: int = 0
):
"""记录外部API请求到数据库"""
external_api_log = ExternalApiLog(
callback_logs_id=callback_logs_id,
request_url=request_url,
request_headers=request_headers,
request_body=request_body,
response_status=response_status,
response_headers=response_headers,
response_body=response_body,
retry_count=retry_count
)
db.add(external_api_log)
await db.commit()
async def call_external_api_with_retry(
db: AsyncSession,
request_body: Dict[str, Any],
max_retries: int = None,
callback_logs_id: int = None
) -> tuple[bool, int]:
"""
调用外部API并支持重试机制
Args:
db: 数据库会话
request_body: 请求体
max_retries: 最大重试次数
Returns:
tuple[是否成功, 实际重试次数]
"""
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"
}
for attempt in range(1, max_retries+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,
headers=headers,
json=request_body
)
logger.debug(f"📥 外部API响应: status={response.status_code}")
# 记录每次尝试的结果
await log_external_api_request(
db=db,
callback_logs_id=callback_logs_id,
request_url=settings.external_api_url,
request_headers=headers,
request_body=request_body,
response_status=response.status_code,
response_headers=dict(response.headers),
response_body=response.text,
retry_count=attempt
)
# 检查响应状态
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
# 否则等待一段时间后重试
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,
callback_logs_id=callback_logs_id,
request_url=settings.external_api_url,
request_headers=headers,
request_body=request_body,
response_status=None,
response_headers=None,
response_body=f"RequestError: {str(e)}",
retry_count=attempt
)
# 如果是最后一次尝试,直接返回失败
if attempt == max_retries:
logger.error(f"❌ 外部API调用最终失败: {e}")
return False, attempt
# 否则等待一段时间后重试
wait_time = 1 * (attempt + 1) # 递增延迟
logger.info(f"⏳ 等待 {wait_time}秒 后重试...")
await asyncio.sleep(wait_time)
return False, max_retries
@router.post("/ai-talk/callback/{siteId}/failure", response_model=CallbackResponse)
async def ai_talk_callback(
callback_data: CallbackRequest,
request: Request,
callback_data: dict = Body(...),
siteId: str = Path(..., description="站点ID"),
db: AsyncSession = Depends(get_db)
):
@@ -315,80 +154,42 @@ async def ai_talk_callback(
logger.error("❌ callback_data 参数为空")
raise HTTPException(status_code=400, detail="callback_data 参数不能为空")
if not isinstance(callback_data.data, list):
logger.error(f"❌ data 参数类型无效: {type(callback_data.data)}")
if not isinstance(callback_data.get('data'), list):
logger.error(f"❌ data 参数类型无效: {type(callback_data.get('data'))}")
raise HTTPException(status_code=400, detail="data 参数必须是数组")
# 记录回调请求(包含siteId)
callback_log_id = await log_callback_request(db, request, siteId)
callback_log_id = await log_callback_request(db, request, siteId, callback_data)
# 使用分布式锁防止并发调用,使用callback_log_id作为唯一值
async with redis_manager.create_lock(f"ai_talk_callback_{callback_log_id}"):
logger.debug(f"🔒 获取Redis锁成功: ai_talk_callback_{callback_log_id}")
# 保存callback_data.data中的数据
data_list = callback_data.get('data', [])
if data_list and len(data_list) > 0:
await save_callback_data_items(db, data_list)
# 保存callback_data.data中的数据
if callback_data.data and len(callback_data.data) > 0:
await save_callback_data_items(db, callback_data.data)
# 检查是否启用外部API调用
if not settings.external_api_enabled:
logger.info(f"🔌 外部API调用已禁用,直接返回成功")
return CallbackResponse(
success=True,
message="外部API调用已禁用,直接返回",
processed=False,
retry_count=0,
site_id=siteId
)
# 提取所有手机号并去重
phone_numbers_set = set()
for item in callback_data.data:
phone_numbers_set.add(item.number_data.number)
phone_numbers = list(phone_numbers_set)
logger.info(f"📱 提取到的手机号列表(去重后): {phone_numbers}")
# 检查去重后的手机号是否超过 settings.count_threshold,超过直接返回
for phone_number in phone_numbers:
# 检查手机号是否超过阈值
exceeds_threshold, query_failed = await check_phone_number_threshold(db, phone_number)
# 如果查询出错,跳过此手机号
if query_failed:
logger.warning(f"⚠️ 查询手机号 {phone_number} 失败,跳过处理")
continue
if exceeds_threshold:
logger.info(f"✅ 手机号 {phone_number} 出现次数超过阈值,跳过处理")
continue
# 检查是否启用外部API调用
if not settings.external_api_enabled:
logger.info(f"🔌 外部API调用已禁用,直接返回成功")
return CallbackResponse(
success=True,
message="成功",
processed=False,
retry_count=0,
site_id=siteId
)
# 异步调用外部接口,不等待结果
import asyncio
asyncio.create_task(process_external_api_call(db, callback_data, siteId, callback_log_id))
# 立即返回成功响应
return CallbackResponse(
success=True,
message="成功",
processed=False,
retry_count=0,
site_id=siteId
)
# 手机号出现的次数少于settings.count_threshold,调用外部API
# 调用外部API并支持重试
request_body = callback_data.model_dump()
success, retry_count = await call_external_api_with_retry(
db=db,
request_body=request_body,
max_retries=settings.external_api_retry_max,
callback_logs_id=callback_log_id
)
if success:
logger.info(f"✅ 外部API调用成功,重试次数: {retry_count}")
return CallbackResponse(
success=True,
message="外部API调用成功",
processed=True,
retry_count=retry_count,
site_id=siteId
)
else:
logger.error(f"❌ 外部API调用失败,已重试{retry_count}次")
raise HTTPException(
status_code=500,
detail=f"外部API调用失败,已重试{retry_count}次"
)
except HTTPException:
raise