from fastapi import APIRouter, Request, Body, HTTPException, Depends, Path from sqlalchemy.ext.asyncio import AsyncSession from typing import Optional from app.database import get_db, CallbackFailureLog, CallbackFailureData 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") router = APIRouter() async def log_callback_request( db: AsyncSession, request: Request, site_id: str, callback_data: dict ) -> tuple[bool, Optional[int]]: """记录回调请求到数据库,返回操作是否成功和记录ID""" 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_info = request.scope['server'] server_ip = server_host server_port = server_port_info # 准备请求头信息(直接记录原始请求头) request_headers = dict(request.headers) # 准备请求体信息(直接记录原始请求体) request_body = callback_data # 记录请求URL(JSON格式) logger.info(f"🌐 请求URL: {request.url}") # 记录site_id(JSON格式) logger.info(f"📝 site_id: {site_id}") # 记录server_ip(JSON格式) server_info = { "ip": server_ip, "port": server_port } logger.info(f"🏠 server_ip: {server_info}") # 记录client_ip(JSON格式) client_info = { "ip": client_ip, "port": client_port } logger.info(f"🖥️ client_ip: {client_info}") # 记录请求头(JSON格式) logger.info(f"📋 请求头: {json.dumps(request_headers, ensure_ascii=False, indent=2)}") # 记录请求体(JSON格式) logger.info(f"📄 请求体: {json.dumps(callback_data, ensure_ascii=False, indent=2)}") try: # 保存到数据库 callback_log = CallbackFailureLog( 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=request_headers, # 保存原始请求头 request_body=request_body # 保存从request.body获取的原始请求体 ) db.add(callback_log) await db.commit() logger.info(f"✅ 回调请求记录成功保存到数据库,ID: {callback_log.id}") # 返回成功标识和记录ID return True, callback_log.id except Exception as e: logger.error(f"❌ 保存回调请求到数据库失败: {e}", exc_info=True) # 返回失败标识 return False, None async def save_callback_data_items( db: AsyncSession, callback_data_items: list, callback_failure_log_id: int ): """保存callback_data.data中的数据到数据库""" import json from datetime import datetime try: callback_data_records = [] for item in callback_data_items: # 获取手机号 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', '') # 获取用户ID user_id = item.get('user_id', '') # 获取状态信息 status = item.get('status', 0) status_description = item.get('status_str', '') # 获取通话日期,默认使用当前时间 calldate = datetime.now() if 'calldate' in item: try: # 如果calldate是字符串,尝试解析为datetime if isinstance(item['calldate'], str): calldate = datetime.fromisoformat(item['calldate'].replace('Z', '+00:00')) elif isinstance(item['calldate'], (int, float)): # 如果是时间戳,转换为datetime calldate = datetime.fromtimestamp(item['calldate']) except (ValueError, TypeError) as e: logger.warning(f"⚠️ 解析calldate失败: {item.get('calldate')}, 使用当前时间, 错误: {e}") calldate = datetime.now() # 将整个item转换为JSON字符串保存 raw_data_json = json.dumps(item, ensure_ascii=False) callback_data_record = CallbackFailureData( callback_failure_log_id=callback_failure_log_id, phone_number=phone_number, task_id=task_id, user_id=user_id, # 保存用户ID status=status, status_description=status_description, raw_data=raw_data_json, # 保存原始JSON字符串 calldate=calldate # 保存通话日期 ) callback_data_records.append(callback_data_record) if len(callback_data_records) > 0: db.add_all(callback_data_records) await db.commit() logger.info(f"✅ 成功保存 {len(callback_data_records)} 条callback_data记录到数据库") except Exception as e: logger.error(f"❌ 保存callback_data到数据库失败: {e}", exc_info=True) # 不重新抛出异常,避免影响主业务流程 @router.post("/ai-talk/callback/{siteId}/failure", response_model=CallbackResponse) async def ai_talk_callback( request: Request, callback_data: dict = Body(...), siteId: str = Path(..., description="站点ID"), db: AsyncSession = Depends(get_db) ): """ AI Talk回调接口处理 """ try: # 记录回调请求(包含siteId) success, callback_log_id = await log_callback_request(db, request, siteId, callback_data) if not success: logger.error("❌ 回调请求日志保存失败") raise HTTPException(status_code=500, detail="回调请求日志保存失败") if 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_log_id) # 检查是否启用外部API调用 if settings.external_api_enabled: # 异步调用外部接口,不等待结果(只有在成功获取callback_log_id时才调用) if success and callback_log_id: import asyncio asyncio.create_task(process_external_api_call(db, callback_log_id, siteId)) else: logger.info(f"🔌 外部API调用已禁用,直接返回成功") # 立即返回成功响应 return CallbackResponse( success=True, message="成功", processed=False, retry_count=0, site_id=siteId ) except HTTPException: raise except Exception as e: logger.error(f"💥 服务器内部错误: {e}", exc_info=True) raise HTTPException( status_code=500, detail=f"服务器内部错误: {str(e)}" )