fix: 修复 OpenAI 和 Gemini API 调用重试逻辑及日志记录

This commit is contained in:
yinpeng
2024-12-21 02:27:59 +08:00
parent 80bcaf5cd4
commit 33a5f9d89f
9 changed files with 192 additions and 168 deletions
+1
View File
@@ -36,6 +36,7 @@ share/python-wheels/
.installed.cfg .installed.cfg
*.egg *.egg
MANIFEST MANIFEST
.idea/
# PyInstaller # PyInstaller
# Usually these files are written by a python script from a template # Usually these files are written by a python script from a template
+6 -4
View File
@@ -1,14 +1,13 @@
from email.header import Header
from fastapi import APIRouter, Depends from fastapi import APIRouter, Depends
from fastapi.responses import StreamingResponse from fastapi.responses import StreamingResponse
from app.core.config import settings
from app.core.logger import get_gemini_logger
from app.core.security import SecurityService from app.core.security import SecurityService
from app.schemas.gemini_models import GeminiRequest
from app.services.chat_service import ChatService from app.services.chat_service import ChatService
from app.services.key_manager import KeyManager from app.services.key_manager import KeyManager
from app.services.model_service import ModelService from app.services.model_service import ModelService
from app.schemas.gemini_models import GeminiRequest
from app.core.config import settings
from app.core.logger import get_gemini_logger
router = APIRouter(prefix="/gemini/v1beta") router = APIRouter(prefix="/gemini/v1beta")
logger = get_gemini_logger() logger = get_gemini_logger()
@@ -19,6 +18,7 @@ key_manager = KeyManager(settings.API_KEYS)
model_service = ModelService(settings.MODEL_SEARCH) model_service = ModelService(settings.MODEL_SEARCH)
chat_service = ChatService(base_url=settings.BASE_URL, key_manager=key_manager) chat_service = ChatService(base_url=settings.BASE_URL, key_manager=key_manager)
@router.get("/models") @router.get("/models")
async def list_models( async def list_models(
key: str = None, key: str = None,
@@ -31,6 +31,7 @@ async def list_models(
logger.info(f"Using API key: {api_key}") logger.info(f"Using API key: {api_key}")
return model_service.get_gemini_models(api_key) return model_service.get_gemini_models(api_key)
@router.post("/models/{model_name}:generateContent") @router.post("/models/{model_name}:generateContent")
async def generate_content( async def generate_content(
model_name: str, model_name: str,
@@ -66,6 +67,7 @@ async def generate_content(
if retries >= MAX_RETRIES: if retries >= MAX_RETRIES:
logger.error(f"Max retries ({MAX_RETRIES}) reached. Raising error") logger.error(f"Max retries ({MAX_RETRIES}) reached. Raising error")
@router.post("/models/{model_name}:streamGenerateContent") @router.post("/models/{model_name}:streamGenerateContent")
async def stream_generate_content( async def stream_generate_content(
model_name: str, model_name: str,
+5 -5
View File
@@ -48,9 +48,9 @@ async def chat_completion(
api_key = await key_manager.get_next_working_key() api_key = await key_manager.get_next_working_key()
logger.info(f"Using API key: {api_key}") logger.info(f"Using API key: {api_key}")
retries = 0 retries = 0
MAX_RETRIES = 3 max_retries = 3
while retries < MAX_RETRIES: while retries < max_retries:
try: try:
response = await chat_service.create_chat_completion( response = await chat_service.create_chat_completion(
request=request, request=request,
@@ -64,13 +64,13 @@ async def chat_completion(
except Exception as e: except Exception as e:
logger.warning( logger.warning(
f"API call failed with error: {str(e)}. Attempt {retries + 1} of {MAX_RETRIES}" f"API call failed with error: {str(e)}. Attempt {retries + 1} of {max_retries}"
) )
api_key = await key_manager.handle_api_failure(api_key) api_key = await key_manager.handle_api_failure(api_key)
logger.info(f"Switched to new API key: {api_key}") logger.info(f"Switched to new API key: {api_key}")
retries += 1 retries += 1
if retries >= MAX_RETRIES: if retries >= max_retries:
logger.error(f"Max retries ({MAX_RETRIES}) reached. Raising error") logger.error(f"Max retries ({max_retries}) reached. Raising error")
raise raise
+1
View File
@@ -9,6 +9,7 @@ class Settings(BaseSettings):
MODEL_SEARCH: List[str] = ["gemini-2.0-flash-exp"] MODEL_SEARCH: List[str] = ["gemini-2.0-flash-exp"]
TOOLS_CODE_EXECUTION_ENABLED: bool = False TOOLS_CODE_EXECUTION_ENABLED: bool = False
SHOW_SEARCH_LINK: bool = True SHOW_SEARCH_LINK: bool = True
class Config: class Config:
env_file = ".env" env_file = ".env"
+16
View File
@@ -15,6 +15,7 @@ COLORS = {
# Windows系统启用ANSI支持 # Windows系统启用ANSI支持
if platform.system() == 'Windows': if platform.system() == 'Windows':
import ctypes import ctypes
kernel32 = ctypes.windll.kernel32 kernel32 = ctypes.windll.kernel32
kernel32.SetConsoleMode(kernel32.GetStdHandle(-11), 7) kernel32.SetConsoleMode(kernel32.GetStdHandle(-11), 7)
@@ -23,6 +24,7 @@ class ColoredFormatter(logging.Formatter):
""" """
自定义的日志格式化器,添加颜色支持 自定义的日志格式化器,添加颜色支持
""" """
def format(self, record): def format(self, record):
# 获取对应级别的颜色代码 # 获取对应级别的颜色代码
color = COLORS.get(record.levelname, '') color = COLORS.get(record.levelname, '')
@@ -30,6 +32,7 @@ class ColoredFormatter(logging.Formatter):
record.levelname = f"{color}{record.levelname}\033[0m" record.levelname = f"{color}{record.levelname}\033[0m"
return super().format(record) return super().format(record)
# 日志格式 # 日志格式
FORMATTER = ColoredFormatter( FORMATTER = ColoredFormatter(
"%(asctime)s - %(name)s - %(levelname)s - [%(filename)s:%(lineno)d] - %(message)s" "%(asctime)s - %(name)s - %(levelname)s - [%(filename)s:%(lineno)d] - %(message)s"
@@ -44,7 +47,11 @@ LOG_LEVELS = {
"critical": logging.CRITICAL, "critical": logging.CRITICAL,
} }
class Logger: class Logger:
def __init__(self):
pass
_loggers: Dict[str, logging.Logger] = {} _loggers: Dict[str, logging.Logger] = {}
@staticmethod @staticmethod
@@ -82,30 +89,39 @@ class Logger:
""" """
return Logger._loggers.get(name) return Logger._loggers.get(name)
# 预定义的loggers # 预定义的loggers
def get_openai_logger(): def get_openai_logger():
return Logger.setup_logger("openai") return Logger.setup_logger("openai")
def get_gemini_logger(): def get_gemini_logger():
return Logger.setup_logger("gemini") return Logger.setup_logger("gemini")
def get_chat_logger(): def get_chat_logger():
return Logger.setup_logger("chat") return Logger.setup_logger("chat")
def get_model_logger(): def get_model_logger():
return Logger.setup_logger("model") return Logger.setup_logger("model")
def get_security_logger(): def get_security_logger():
return Logger.setup_logger("security") return Logger.setup_logger("security")
def get_key_manager_logger(): def get_key_manager_logger():
return Logger.setup_logger("key_manager") return Logger.setup_logger("key_manager")
def get_main_logger(): def get_main_logger():
return Logger.setup_logger("main") return Logger.setup_logger("main")
def get_embeddings_logger(): def get_embeddings_logger():
return Logger.setup_logger("embeddings") return Logger.setup_logger("embeddings")
def get_request_logger(): def get_request_logger():
return Logger.setup_logger("request") return Logger.setup_logger("request")
+1 -1
View File
@@ -3,9 +3,9 @@ from starlette.middleware.base import BaseHTTPMiddleware
import json import json
from app.core.logger import get_request_logger from app.core.logger import get_request_logger
logger = get_request_logger() logger = get_request_logger()
# 添加中间件类 # 添加中间件类
class RequestLoggingMiddleware(BaseHTTPMiddleware): class RequestLoggingMiddleware(BaseHTTPMiddleware):
async def dispatch(self, request: Request, call_next): async def dispatch(self, request: Request, call_next):
+5 -2
View File
@@ -1,9 +1,12 @@
from typing import List, Optional, Dict, Any, Literal from typing import List, Optional, Dict, Any, Literal
from pydantic import BaseModel from pydantic import BaseModel
class SafetySetting(BaseModel): class SafetySetting(BaseModel):
category: Optional[Literal["HARM_CATEGORY_HATE_SPEECH", "HARM_CATEGORY_DANGEROUS_CONTENT", "HARM_CATEGORY_HARASSMENT", "HARM_CATEGORY_SEXUALLY_EXPLICIT"]] = None category: Optional[Literal[
threshold: Optional[Literal["HARM_BLOCK_THRESHOLD_UNSPECIFIED", "BLOCK_LOW_AND_ABOVE", "BLOCK_MEDIUM_AND_ABOVE","BLOCK_ONLY_HIGH","BLOCK_NONE","OFF"]] = None "HARM_CATEGORY_HATE_SPEECH", "HARM_CATEGORY_DANGEROUS_CONTENT", "HARM_CATEGORY_HARASSMENT", "HARM_CATEGORY_SEXUALLY_EXPLICIT"]] = None
threshold: Optional[Literal[
"HARM_BLOCK_THRESHOLD_UNSPECIFIED", "BLOCK_LOW_AND_ABOVE", "BLOCK_MEDIUM_AND_ABOVE", "BLOCK_ONLY_HIGH", "BLOCK_NONE", "OFF"]] = None
class GenerationConfig(BaseModel): class GenerationConfig(BaseModel):
+55 -57
View File
@@ -11,12 +11,7 @@ from app.schemas.openai_models import ChatRequest
logger = get_chat_logger() logger = get_chat_logger()
class ChatService: def convert_messages_to_gemini_format(messages: list) -> list:
def __init__(self, base_url: str, key_manager=None):
self.base_url = base_url
self.key_manager = key_manager
def convert_messages_to_gemini_format(self, messages: list) -> list:
"""Convert OpenAI message format to Gemini format""" """Convert OpenAI message format to Gemini format"""
converted_messages = [] converted_messages = []
for msg in messages: for msg in messages:
@@ -60,6 +55,23 @@ class ChatService:
return converted_messages return converted_messages
def format_execution_result(result_data: dict) -> str:
"""格式化执行结果输出"""
outcome = result_data.get("outcome", "")
output = result_data.get("output", "").strip()
return f"""\n【执行结果】\n> outcome: {outcome}\n\n【输出结果】\n```plaintext\n{output}\n```\n"""
def create_search_link(web):
return f'\n- [{web["title"]}]({web["uri"]})'
class ChatService:
def __init__(self, base_url: str, key_manager=None):
self.base_url = base_url
self.key_manager = key_manager
def convert_gemini_response_to_openai( def convert_gemini_response_to_openai(
self, self,
response: Dict[str, Any], response: Dict[str, Any],
@@ -82,28 +94,17 @@ class ChatService:
elif "codeExecution" in parts[0]: elif "codeExecution" in parts[0]:
text = self.format_code_block(parts[0]["codeExecution"]) text = self.format_code_block(parts[0]["codeExecution"])
elif "executableCodeResult" in parts[0]: elif "executableCodeResult" in parts[0]:
text = self.format_execution_result( text = format_execution_result(
parts[0]["executableCodeResult"] parts[0]["executableCodeResult"]
) )
elif "codeExecutionResult" in parts[0]: elif "codeExecutionResult" in parts[0]:
text = self.format_execution_result( text = format_execution_result(
parts[0]["codeExecutionResult"] parts[0]["codeExecutionResult"]
) )
else: else:
text = "" text = ""
if ( text = self.add_search_link_text(model, candidate, text)
settings.SHOW_SEARCH_LINK
and model.endswith("-search")
and "groundingMetadata" in candidate
and "groundingChunks" in candidate["groundingMetadata"]
):
groundingChunks = candidate["groundingMetadata"]["groundingChunks"]
text += "\n\n---\n\n"
text += f"**【引用来源】**\n\n"
for _, groundingChunk in enumerate(groundingChunks, 1):
if "web" in groundingChunk:
text += self.create_search_link(groundingChunk["web"])
else: else:
text = "" text = ""
@@ -150,18 +151,7 @@ class ChatService:
if response.get("candidates"): if response.get("candidates"):
text = response["candidates"][0]["content"]["parts"][0]["text"] text = response["candidates"][0]["content"]["parts"][0]["text"]
candidate = response["candidates"][0] candidate = response["candidates"][0]
if ( text = self.add_search_link_text(model, candidate, text)
settings.SHOW_SEARCH_LINK
and model.endswith("-search")
and "groundingMetadata" in candidate
and "groundingChunks" in candidate["groundingMetadata"]
):
groundingChunks = candidate["groundingMetadata"]["groundingChunks"]
text += "\n\n---\n\n"
text += f"**【引用来源】**\n\n"
for _, groundingChunk in enumerate(groundingChunks, 1):
if "web" in groundingChunk:
text += self.create_search_link(groundingChunk["web"])
res["choices"][0]["message"]["content"] = text res["choices"][0]["message"]["content"] = text
return res return res
else: else:
@@ -173,6 +163,23 @@ class ChatService:
res["choices"][0]["message"]["content"] = f"Error converting Gemini response: {str(e)}" res["choices"][0]["message"]["content"] = f"Error converting Gemini response: {str(e)}"
return res return res
def add_search_link_text(self, model, candidate, text):
if (
settings.SHOW_SEARCH_LINK
and model.endswith("-search")
and "groundingMetadata" in candidate
and "groundingChunks" in candidate["groundingMetadata"]
):
grounding_chunks = candidate["groundingMetadata"]["groundingChunks"]
text += "\n\n---\n\n"
text += f"**【引用来源】**\n\n"
for _, grounding_chunk in enumerate(grounding_chunks, 1):
if "web" in grounding_chunk:
text += create_search_link(grounding_chunk["web"])
return text
else:
return text
async def create_chat_completion( async def create_chat_completion(
self, self,
request: ChatRequest, request: ChatRequest,
@@ -210,7 +217,7 @@ class ChatService:
gemini_model = model[:-7] # Remove -search suffix gemini_model = model[:-7] # Remove -search suffix
else: else:
gemini_model = model gemini_model = model
gemini_messages = self.convert_messages_to_gemini_format(messages) gemini_messages = convert_messages_to_gemini_format(messages)
if not stream: if not stream:
# 非流式模式下,移除代码执行工具 # 非流式模式下,移除代码执行工具
@@ -247,26 +254,26 @@ class ChatService:
if stream: if stream:
async def generate(): async def generate():
retries = 0 retries = 0
MAX_RETRIES = 3 max_retries = 3
current_api_key = api_key current_api_key = api_key
while retries < MAX_RETRIES: while retries < max_retries:
try: try:
timeout = httpx.Timeout( timeout = httpx.Timeout(
60.0, read=60.0 60.0, read=60.0
) # 连接超时60秒,读取超时60秒 ) # 连接超时60秒,读取超时60秒
async with httpx.AsyncClient(timeout=timeout) as client: async with httpx.AsyncClient(timeout=timeout) as async_client:
stream_url = f"https://generativelanguage.googleapis.com/v1beta/models/{gemini_model}:streamGenerateContent?alt=sse&key={current_api_key}" stream_url = f"https://generativelanguage.googleapis.com/v1beta/models/{gemini_model}:streamGenerateContent?alt=sse&key={current_api_key}"
async with client.stream( async with async_client.stream(
"POST", stream_url, json=payload "POST", stream_url, json=payload
) as response: ) as async_response:
if response.status_code != 200: if async_response.status_code != 200:
error_content = await response.read() error_content = await async_response.read()
error_msg = error_content.decode("utf-8") error_msg = error_content.decode("utf-8")
logger.error( logger.error(
f"API error: {response.status_code}, {error_msg}" f"API error: {async_response.status_code}, {error_msg}"
) )
if retries < MAX_RETRIES - 1: if retries < max_retries - 1:
current_api_key = ( current_api_key = (
await self.key_manager.handle_api_failure( await self.key_manager.handle_api_failure(
current_api_key current_api_key
@@ -276,12 +283,12 @@ class ChatService:
continue continue
else: else:
logger.error( logger.error(
f"Max retries reached. Final error: {response.status_code}, {error_msg}" f"Max retries reached. Final error: {async_response.status_code}, {error_msg}"
) )
yield f"data: {json.dumps({'error': f'API error: {response.status_code}, {error_msg}'})}\n\n" yield f"data: {json.dumps({'error': f'API error: {async_response.status_code}, {error_msg}'})}\n\n"
return return
async for line in response.aiter_lines(): async for line in async_response.aiter_lines():
if line.startswith("data: "): if line.startswith("data: "):
try: try:
chunk = json.loads(line[6:]) chunk = json.loads(line[6:])
@@ -297,7 +304,7 @@ class ChatService:
yield f"data: {json.dumps(openai_chunk)}\n\n" yield f"data: {json.dumps(openai_chunk)}\n\n"
except json.JSONDecodeError: except json.JSONDecodeError:
continue continue
yield f"data: {json.dumps(self.convert_gemini_response_to_openai({}, model,stream=True, finish_reason='stop'))}\n\n" yield f"data: {json.dumps(self.convert_gemini_response_to_openai({}, model, stream=True, finish_reason='stop'))}\n\n"
yield "data: [DONE]\n\n" yield "data: [DONE]\n\n"
return return
@@ -305,7 +312,7 @@ class ChatService:
logger.warning( logger.warning(
f"Read timeout occurred, attempting retry {retries + 1}" f"Read timeout occurred, attempting retry {retries + 1}"
) )
if retries < MAX_RETRIES - 1: if retries < max_retries - 1:
current_api_key = await self.key_manager.handle_api_failure( current_api_key = await self.key_manager.handle_api_failure(
current_api_key current_api_key
) )
@@ -322,7 +329,7 @@ class ChatService:
logger.exception( logger.exception(
f"Stream error: {str(e)}, attempting retry {retries + 1}" f"Stream error: {str(e)}, attempting retry {retries + 1}"
) )
if retries < MAX_RETRIES - 1: if retries < max_retries - 1:
current_api_key = await self.key_manager.handle_api_failure( current_api_key = await self.key_manager.handle_api_failure(
current_api_key current_api_key
) )
@@ -362,12 +369,6 @@ class ChatService:
return f"""\n【代码执行】\n```{language}\n{code}\n```\n""" return f"""\n【代码执行】\n```{language}\n{code}\n```\n"""
def format_execution_result(self, result_data: dict) -> str:
"""格式化执行结果输出"""
outcome = result_data.get("outcome", "")
output = result_data.get("output", "").strip()
return f"""\n【执行结果】\n> outcome: {outcome}\n\n【输出结果】\n```plaintext\n{output}\n```\n"""
async def generate_content( async def generate_content(
self, model_name: str, request: GeminiRequest, api_key: str self, model_name: str, request: GeminiRequest, api_key: str
) -> dict: ) -> dict:
@@ -451,6 +452,3 @@ class ChatService:
retries += 1 retries += 1
continue continue
raise raise
def create_search_link(self, web):
return f'\n- [{web["title"]}]({web["uri"]})'
+5 -2
View File
@@ -1,5 +1,8 @@
from typing import Union, List
import openai import openai
from typing import Union, List, Dict, Any from openai.types import CreateEmbeddingResponse
from app.core.logger import get_embeddings_logger from app.core.logger import get_embeddings_logger
logger = get_embeddings_logger() logger = get_embeddings_logger()
@@ -11,7 +14,7 @@ class EmbeddingService:
async def create_embedding( async def create_embedding(
self, input_text: Union[str, List[str]], model: str, api_key: str self, input_text: Union[str, List[str]], model: str, api_key: str
) -> Dict[str, Any]: ) -> CreateEmbeddingResponse:
"""Create embeddings using OpenAI API""" """Create embeddings using OpenAI API"""
try: try:
client = openai.OpenAI(api_key=api_key, base_url=self.base_url) client = openai.OpenAI(api_key=api_key, base_url=self.base_url)