From 703408c80858f9f608f54e15d6cb6bd83f49d071 Mon Sep 17 00:00:00 2001 From: "mark.tian" Date: Wed, 3 Dec 2025 21:56:06 +0800 Subject: [PATCH] =?UTF-8?q?=E6=8E=A5=E5=8F=A3=E8=AF=B7=E6=B1=82=E4=BD=93?= =?UTF-8?q?=E6=94=B9=E7=94=A8dict=E5=A4=84=E7=90=86=20=E4=BF=9D=E5=AD=98?= =?UTF-8?q?=E6=AF=8F=E6=AC=A1=E9=80=9A=E8=AF=9D=E4=BF=A1=E6=81=AF=E5=88=B0?= =?UTF-8?q?=E5=BA=93?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- app/database.py | 1 + app/models.py | 54 -------- app/routes.py | 327 ++++++++++-------------------------------------- 3 files changed, 65 insertions(+), 317 deletions(-) diff --git a/app/database.py b/app/database.py index f7307b6..773f1ae 100644 --- a/app/database.py +++ b/app/database.py @@ -33,6 +33,7 @@ class CallbackData(Base): task_id = Column(String(100), nullable=False, comment="任务ID") status = Column(Integer, nullable=False, comment="状态") status_description = Column(String(200), nullable=False, comment="状态描述") + raw_data = Column(Text, nullable=False, comment="原始回调数据JSON字符串") created_at = Column(DateTime, default=datetime.now) # 为手机号字段添加索引,提高查询性能 diff --git a/app/models.py b/app/models.py index 4e5f478..f9c6443 100644 --- a/app/models.py +++ b/app/models.py @@ -1,60 +1,6 @@ from pydantic import BaseModel, Field from typing import List, Optional - -class NumberData(BaseModel): - number: str - province: str - city: str - operator: str - - -class Group(BaseModel): - id: int - name: str - - -class Task(BaseModel): - id: str - name: str - - -class User(BaseModel): - id: str - name: str - - -class CustomerData(BaseModel): - name: str - email: str - company: Optional[str] = None - extra: Optional[str] = None - - -class CallbackItem(BaseModel): - bill: int - duration: int - callid: str - calldate: str - number: str - numberid: str - customer_id: str - status: int - status_str: str - user_id: str - type: int - number_data: NumberData - group: Group - task: Task - user: User - customer_data: CustomerData - - -class CallbackRequest(BaseModel): - count: int = Field(..., ge=0, description="通话失败个数") - data: List[CallbackItem] - - class CallbackResponse(BaseModel): success: bool message: str diff --git a/app/routes.py b/app/routes.py index 27fe442..4177dd8 100644 --- a/app/routes.py +++ b/app/routes.py @@ -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