修复部分情况下,假流式无法发出空白回复的问题。简化空白回复的发送逻辑。
This commit is contained in:
sanjusss
2025-07-25 19:44:01 +08:00
parent a6558b4668
commit 3d6b5063d5
+11 -26
View File
@@ -358,39 +358,24 @@ class OpenAIChatService:
logger.info( logger.info(
f"Fake streaming enabled for model: {model}. Calling non-streaming endpoint." f"Fake streaming enabled for model: {model}. Calling non-streaming endpoint."
) )
keep_sending_empty_data = True
async def send_empty_data_locally() -> AsyncGenerator[str, None]:
"""定期发送空数据以保持连接"""
while keep_sending_empty_data:
await asyncio.sleep(settings.FAKE_STREAM_EMPTY_DATA_INTERVAL_SECONDS)
if keep_sending_empty_data:
empty_chunk = self.response_handler.handle_response({}, model, stream=True, finish_reason='stop', usage_metadata=None)
yield f"data: {json.dumps(empty_chunk)}\n\n"
logger.debug("Sent empty data chunk for fake stream heartbeat.")
empty_data_generator = send_empty_data_locally()
api_response_task = asyncio.create_task( api_response_task = asyncio.create_task(
self.api_client.generate_content(payload, model, api_key) self.api_client.generate_content(payload, model, api_key)
) )
i = 0
try: try:
while not api_response_task.done(): while not api_response_task.done():
try: i = i + 1
next_empty_chunk = await asyncio.wait_for( """定期发送空数据以保持连接"""
empty_data_generator.__anext__(), timeout=0.1 if i >= settings.FAKE_STREAM_EMPTY_DATA_INTERVAL_SECONDS :
) i = 0
yield next_empty_chunk empty_chunk = self.response_handler.handle_response({}, model, stream=True, finish_reason='stop', usage_metadata=None)
except asyncio.TimeoutError: yield f"data: {json.dumps(empty_chunk)}\n\n"
pass logger.debug("Sent empty data chunk for fake stream heartbeat.")
except ( await asyncio.sleep(1)
StopAsyncIteration
):
break
response = await api_response_task
finally: finally:
keep_sending_empty_data = False response = await api_response_task
if response and response.get("candidates"): if response and response.get("candidates"):
response = self.response_handler.handle_response(response, model, stream=True, finish_reason='stop', usage_metadata=response.get("usageMetadata", {})) response = self.response_handler.handle_response(response, model, stream=True, finish_reason='stop', usage_metadata=response.get("usageMetadata", {}))