@@ -760,11 +760,8 @@ msgid "Prompt is too long. Maximum length is %(max_length)s characters." msgstr "Промпт слишком длинный. Максимальная длина — %(max_length)s символов." #: ml_model/exceptions.py:161 -msgid "" -"Service is currently unavailable due to high demand. Please try again later" -msgstr "" -"Сервис временно недоступен из-за высокой нагрузки. Пожалуйста, попробуйте " -"позже" +msgid "Service is temporarily unavailable. Please try again later" +msgstr "Сервис временно недоступен. Пожалуйста, попробуйте позже" #: ml_model/exceptions.py:169 #, python-format @@ -4,6 +4,7 @@ from typing import Any import httpx from django.conf import settings +from ml_model.exceptions import OpenAIResponseError, ServiceTemporaryUnavailableError from poller.models import Proxy type OpenAIEvent = dict[str, Any] @@ -26,6 +27,8 @@ class OpenAIStreamMixin: input_tokens, output_tokens = yield from self._stream_request( client, 'POST', json=payload, state=state ) + except ServiceTemporaryUnavailableError: + raise except Exception: if not state['response_id']: raise @@ -95,7 +98,9 @@ class OpenAIStreamMixin: return self._get_response_id(event) case 'response.output_text.delta': return self._get_delta(event) - case 'response.completed' | 'response.incomplete' | 'response.failed': + case 'response.failed': + raise ServiceTemporaryUnavailableError from self._get_response_error(event) + case 'response.completed' | 'response.incomplete': return self._get_usage(event) case _: return '' @@ -109,3 +114,11 @@ class OpenAIStreamMixin: def _get_response_id(self, event: OpenAIEvent) -> str: return event.get('response', {}).get('id', '') + + def _get_response_error(self, event: OpenAIEvent) -> OpenAIResponseError: + error = event.get('error') or (event.get('response') or {}).get('error') or {} + return OpenAIResponseError( + event_type=event.get('type', 'unknown'), + code=error.get('code', 'unknown'), + message=error.get('message', 'Unknown OpenAI response error'), + ) @@ -162,6 +162,23 @@ class ServiceHighDemandError(Exception): return _('Service is currently unavailable due to high demand. Please try again later') +class ServiceTemporaryUnavailableError(Exception): + def __str__(self) -> str: + return _('Service is temporarily unavailable. Please try again later') + + +class OpenAIResponseError(Exception): + def __init__(self, event_type: str, code: str, message: str) -> None: + self.event_type = event_type + self.code = code + self.message = message + + def __str__(self) -> str: + return ( + f'OpenAI streaming error: event_type={self.event_type}, code={self.code}, message={self.message}' + ) + + class PaidPlanRequiredError(Exception): def __init__(self, feature: str) -> None: self.feature = feature @@ -1,3 +1,4 @@ +import logging from decimal import Decimal from celery import shared_task @@ -14,6 +15,8 @@ from tools.chats.services.sse_chunk_service import SSEChunkService from tools.chats.services.sse_store import PublicSSEStoreService, SSEStoreService from tools.public_api.models import APIKey, APIStore +logger = logging.getLogger(__name__) + def _run_stream( store: SSEStoreService, @@ -51,6 +54,7 @@ def _run_stream( event_id += 1 store.push(SSEChunkService.done(event_id, exc.value or ''), ttl=settings.SSE_DONE_STREAM_TTL) except Exception as exc: + logger.exception(f'Model streaming failed: {(exc.__cause__ or exc)!r}') if event_id < 2: message.is_sent = False message.save(update_fields=['is_sent'])