增加接口日志输出
This commit is contained in:
116
app/routes.py
116
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)}"
|
||||
|
||||
Reference in New Issue
Block a user