""" 回调处理服务模块 包含与回调处理相关的业务逻辑函数 """ import json import time import httpx from datetime import datetime from typing import Optional, Tuple, Dict, Any from sqlalchemy import text from fastapi import Request from sqlalchemy.ext.asyncio import AsyncSession from app.config import settings from app.logger import get_callback_service_logger from app.database import CallbackFailureLog logger = get_callback_service_logger() async def log_callback_request( db: AsyncSession, request: Request, site_id: str, callback_data: dict ) -> bool: """记录回调请求到数据库,返回操作是否成功""" 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}") # 返回成功标识 return True except Exception as e: logger.error(f"❌ 保存回调请求到数据库失败: {e}", exc_info=True) # 返回失败标识 return False def save_callback_data_items( conn, callback_data_items: list, callback_failure_log_id: int ) -> bool: """保存callback_data.data中的数据到数据库 Returns: bool: 保存是否成功,True表示成功,False表示失败或跳过保存 """ try: logger.info(f"📊 共传入callback_log_id {callback_failure_log_id} 的 {len(callback_data_items)} 条数据") # 检查是否已经存在该callback_log_id的数据 existing_data_result = conn.execute( text("SELECT * FROM callback_failure_data WHERE callback_failure_log_id = :log_id"), {"log_id": callback_failure_log_id} ) existing_data_list = existing_data_result.fetchall() if existing_data_list: logger.warning(f"⚠️ callback_log_id {callback_failure_log_id} 已存在 {len(existing_data_list)} 条数据,跳过保存") return False 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) # 直接插入数据库 conn.execute( text(""" INSERT INTO callback_failure_data (callback_failure_log_id, phone_number, task_id, user_id, status, status_description, raw_data, calldate) VALUES (:log_id, :phone_number, :task_id, :user_id, :status, :status_description, :raw_data, :calldate) """), { "log_id": callback_failure_log_id, "phone_number": phone_number, "task_id": task_id, "user_id": user_id, "status": status, "status_description": status_description, "raw_data": raw_data_json, "calldate": calldate } ) conn.commit() logger.info(f"✅ 成功保存 {len(callback_data_items)} 条callback_data记录到数据库") return True except Exception as e: logger.error(f"❌ 保存callback_data到数据库失败: {e}", exc_info=True) # 不重新抛出异常,避免影响主业务流程 return False def get_uncompleted_callback_log(conn) -> Tuple[bool, Optional[Tuple[int, str, str, str]]]: """获取一条未完成的回调请求(按创建时间取最小值) Returns: Tuple[bool, Optional[Tuple[int, str, str, str]]]: - 第一个值表示查询是否成功(True表示成功,False表示失败) - 第二个值为回调日志数据元组或None """ try: # 查询一条未完成的回调日志(按创建时间升序排列,取第一条) result = conn.execute( text(""" SELECT id, site_id, request_headers, request_body FROM callback_failure_logs WHERE is_completed = false ORDER BY created_at ASC LIMIT 1 """) ) row = result.fetchone() if not row: logger.info("📋 没有找到未完成的回调日志") return True, None callback_log_id, site_id, request_headers, request_body = row request_headers_json = json.dumps(request_headers, ensure_ascii=False) request_body_json = json.dumps(request_body, ensure_ascii=False) return True, (callback_log_id, site_id, request_headers_json, request_body_json) except Exception as e: logger.error(f"❌ 查询未完成回调日志失败: {e}", exc_info=True) return False, None def get_callback_log_data(conn, callback_log_id: int) -> Tuple[bool, Optional[str], Optional[str]]: """获取回调日志数据""" try: # 查询回调日志 result = conn.execute( text("SELECT request_headers, request_body FROM callback_failure_logs WHERE id = :log_id"), {"log_id": callback_log_id} ) row = result.fetchone() if not row: logger.error(f"❌ 未找到回调日志: {callback_log_id}") return False, None, None request_headers_json = json.dumps(row[0], ensure_ascii=False) request_body_json = json.dumps(row[1], ensure_ascii=False) return True, request_headers_json, request_body_json except Exception as e: logger.error(f"❌ 获取回调日志数据失败: {e}", exc_info=True) return False, None, None def call_external_api_with_retry( conn, request_body: dict, request_headers: Dict[str, str], max_retries: int, callback_failure_log_id: int ) -> Tuple[bool, int]: """带重试的调用外部API接口""" for attempt in range(1, max_retries + 1): try: logger.info(f"🌐 尝试调用外部API接口,第{attempt}次") with httpx.Client(timeout=30.0) as client: response = client.post( settings.external_api_url, json=request_body, headers=request_headers ) # 记录推送日志 _log_dtc_push_call( conn=conn, callback_failure_log_id=callback_failure_log_id, request_url=settings.external_api_url, request_headers=request_headers, request_body=request_body, response_status=response.status_code, response_headers=dict(response.headers), response_body=response.text, retry_count=attempt - 1 ) if response.status_code == 200: logger.info(f"✅ 调用外部API接口成功,状态码: {response.status_code}") return True, attempt - 1 else: logger.warning(f"⚠️ 外部API返回非成功状态码: {response.status_code}") # 如果是客户端错误(4xx),不重试 if 400 <= response.status_code < 500: logger.error(f"❌ 客户端错误,不重试: {response.status_code}") return False, attempt - 1 except httpx.TimeoutException: logger.warning(f"⏰ 调用外部API接口超时,第{attempt}次尝试") except httpx.RequestError as e: logger.warning(f"🌐 调用外部API接口请求错误,第{attempt}次尝试: {e}") except Exception as e: logger.error(f"❌ 调用外部API接口异常,第{attempt}次尝试: {e}", exc_info=True) # 如果不是最后一次尝试,等待一段时间再重试 if attempt < max_retries: time.sleep(2 ** attempt) # 指数退避 # 所有重试都失败了 logger.error(f"❌ 调用外部API接口失败,已重试{max_retries}次") return False, max_retries def mark_callback_log_completed(conn, callback_log_id: int) -> bool: """标记CallbackFailureLog记录为已完成""" try: # 先检查推送日志中是否有响应成功的记录 success_result = conn.execute( text(""" SELECT COUNT(*) as success_count FROM external_api_logs WHERE callback_failure_log_id = :log_id AND response_status = 200 """), {"log_id": callback_log_id} ) success_count = success_result.fetchone()[0] if success_count == 0: logger.warning(f"⚠️ 回调日志 {callback_log_id} 没有成功的推送记录,不标记为完成") return False logger.info(f"✅ 回调日志 {callback_log_id} 找到 {success_count} 条成功推送记录,开始标记为完成") # 更新指定日志记录为已完成 result = conn.execute( text(""" UPDATE callback_failure_logs SET is_completed = true WHERE id = :log_id AND is_completed = false """), {"log_id": callback_log_id} ) # 检查是否真的更新了记录 if result.rowcount == 0: logger.warning(f"⚠️ 回调日志 {callback_log_id} 已标记为完成或不存在") return False conn.commit() logger.info(f"✅ 回调日志 {callback_log_id} 已标记为完成") return True except Exception as e: logger.error(f"❌ 标记回调日志 {callback_log_id} 完成失败: {e}", exc_info=True) conn.rollback() return False def get_related_records_by_unique_data_list( conn, unique_data_list: list, callback_log_id: int, limit_count: int ) -> Tuple[bool, list]: """根据unique_data_list中的数据查询相关记录,按创建时间排序 Args: conn: 数据库连接 unique_data_list: 包含手机号、task_id、user_id的数据项列表 callback_log_id: 回调日志ID,作为过滤条件 limit_count: 获取记录数量限制 Returns: Tuple[bool, list]: - 第一个值表示查询是否成功(True表示成功,False表示失败) - 第二个值为查询到的相关记录列表,失败时返回空列表 """ related_records = [] try: # 遍历unique_data_list查询相关记录 for item in unique_data_list: # 提取手机号 number_data = item.get('number_data', {}) phone_number = number_data.get('number') # 提取task_id task = item.get('task', {}) task_id = task.get('id', '') # 提取user_id user_id = item.get('user_id', '') if phone_number and task_id and user_id: result = conn.execute( text(""" SELECT phone_number, task_id, user_id, created_at FROM callback_failure_data WHERE phone_number = :phone_number AND task_id = :task_id AND user_id = :user_id AND callback_failure_log_id = :callback_log_id ORDER BY created_at DESC LIMIT :limit_count """), { "phone_number": phone_number, "task_id": task_id, "user_id": user_id, "callback_log_id": callback_log_id, "limit_count": limit_count } ) records = result.fetchall() if records: related_records.extend([{ "phone_number": record[0], "task_id": record[1], "user_id": record[2], "created_at": record[3] } for record in records]) logger.info(f"📊 查询到相关记录数量: {len(related_records)} (阈值: {limit_count})") if related_records: for i, record in enumerate(related_records): logger.info(f"📋 记录{i+1}: 手机号={record['phone_number']}, task_id={record['task_id']}, 用户ID={record['user_id']}, 创建时间={record['created_at']}") return True, related_records except Exception as e: logger.error(f"❌ 查询相关记录失败: {e}", exc_info=True) return False, [] def _log_dtc_push_call( conn, callback_failure_log_id: int, request_url: str, request_headers: Dict[str, Any], request_body: Dict[str, Any], response_status: int, response_headers: Dict[str, str], response_body: str, retry_count: int ) -> bool: """记录DTC推送调用日志""" try: # 提取手机号从请求体中 phone_number = None if 'data' in request_body and isinstance(request_body['data'], list) and len(request_body['data']) > 0: first_item = request_body['data'][0] if 'number_data' in first_item and 'number' in first_item['number_data']: phone_number = first_item['number_data']['number'] # 记录到external_api_logs表 conn.execute( text(""" INSERT INTO external_api_logs ( callback_failure_log_id, phone_number, request_url, request_headers, request_body, response_status, response_headers, response_body, retry_count, created_at ) VALUES ( :callback_failure_log_id, :phone_number, :request_url, :request_headers, :request_body, :response_status, :response_headers, :response_body, :retry_count, datetime('now') ) """), { "callback_failure_log_id": callback_failure_log_id, "phone_number": phone_number, "request_url": request_url, "request_headers": json.dumps(request_headers, ensure_ascii=False), "request_body": json.dumps(request_body, ensure_ascii=False), "response_status": response_status, "response_headers": json.dumps(response_headers, ensure_ascii=False), "response_body": response_body, "retry_count": retry_count } ) conn.commit() logger.debug(f"📝 DTC推送调用日志记录成功,状态码: {response_status}") return True except Exception as e: logger.error(f"❌ 记录DTC推送调用日志失败: {e}", exc_info=True) conn.rollback() return False