@@ -383,24 +383,31 @@ LOGGING = { 'version': 1, 'disable_existing_loggers': False, 'formatters': { + 'pretty': { + 'format': '[{asctime}] [{levelname}] [{name}] {message}', + 'style': '{', + 'datefmt': '%d.%m.%Y %H:%M:%S', + }, 'verbose': { - 'format': '[{asctime}] {levelname} [pid={process}] [thread={thread}] [{pathname}:{lineno}] {message}', + 'format': '[{asctime}] {levelname} [{name}] [{process}:{threadName}] [{pathname}:{lineno}] {message}', 'style': '{', - 'datefmt': '%Y-%m-%d %H:%M:%S', + 'datefmt': '%Y-%m-%dT%H:%M:%SZ', }, }, 'handlers': { 'console': { 'class': 'logging.StreamHandler', - 'formatter': 'verbose', + 'formatter': 'pretty' if DEBUG else 'verbose', }, }, 'loggers': { 'UnleashClient': { + 'handlers': ['console'], 'level': 'CRITICAL', 'propagate': False, }, 'apscheduler': { + 'handlers': ['console'], 'level': 'CRITICAL', 'propagate': False, }, @@ -471,13 +478,13 @@ if CACHEOPS_REDIS: 'token_blacklist.outstandingtoken': {'ops': 'get', 'timeout': 60 * 60 * 24}, } -# UNLEASH settings -FEATURE_FLAG_API_URL = env.str('FEATURE_FLAG_API_URL') -FEATURE_FLAG_APP_NAME = env.str('FEATURE_FLAG_APP_NAME', 'staging') -FEATURE_FLAG_INSTANCE_ID = env.str('FEATURE_FLAG_INSTANCE_ID') -FEATURE_FLAG_WEBHOOK_SECRET_KEY = env.str( - 'FEATURE_FLAG_WEBHOOK_SECRET_KEY', 'FEATURE_FLAG_WEBHOOK_SECRET_KEY' -) +# Unleash settings +UNLEASH_API_URL = env.str('UNLEASH_API_URL', 'https://example.com') +UNLEASH_APP_NAME = env.str('UNLEASH_APP_NAME', 'Development') +UNLEASH_REQUEST_TIMEOUT = env.int('UNLEASH_REQUEST_TIMEOUT', 3) +UNLEASH_REQUEST_RETRIES = env.int('UNLEASH_REQUEST_RETRIES', 1) +UNLEASH_INSTANCE_ID = env.str('UNLEASH_INSTANCE_ID', '') +UNLEASH_WEBHOOK_SECRET_KEY = env.str('UNLEASH_WEBHOOK_SECRET_KEY', 'defaultsecretkey') # RECURRING SETTINGS MAX_RECURRING_ATTEMPTS = env.int('MAX_RECURRING_ATTEMPTS', 1) @@ -23,6 +23,9 @@ logger = logging.getLogger(__name__) api = NinjaAPI(title='AIR API', version='1.0.0', parser=MultiContentTypeParser(), docs_url=None) compatibility_api = NinjaAPI(title='AIR API DEBUG', version='0.0.1', docs_url=None) compatibility_api_v2 = NinjaAPI(title='AIR API DEBUG v2', version='2.0.0', docs_url=None) +public_api = NinjaAPI( + title='PUBLIC AIR API', urls_namespace='public-api-1.0.0', parser=MultiContentTypeParser(), docs_url=None +) api.add_router('users/', 'users.routes.v1.router') api.add_router('chats/', 'tools.chats.routes.v1.router') @@ -35,6 +38,8 @@ compatibility_api.add_router('ml_model/', 'ml_model.routes.v1.router') compatibility_api_v2.add_router('auth/', 'authentication.routes.v2.router') +public_api.add_router('', 'tools.public_api.routes.v1.router') + class Status(Schema): status: Literal['ok', 'dead'] @@ -97,6 +102,7 @@ urlpatterns = ( path('api/v1/api/', api.urls), path('api/v1/', compatibility_api.urls), path('api/v1/v2/', compatibility_api_v2.urls), + path('public/', public_api.urls), ] + static(settings.STATIC_URL, document_root=settings.STATIC_ROOT) + public_urlpatterns @@ -122,3 +128,4 @@ if settings.DEBUG: api.docs_url = '/docs' compatibility_api.docs_url = '/docs' compatibility_api_v2.docs_url = '/docs' + public_api.docs_url = '/docs' @@ -13,11 +13,15 @@ from lib.unleash.cache import UnleashRedisCache class UnleashFeatureFlagService(FeatureFlagService): def __init__(self) -> None: self.client = UnleashClient( - url=settings.FEATURE_FLAG_API_URL, - app_name=settings.FEATURE_FLAG_APP_NAME, - instance_id=settings.FEATURE_FLAG_INSTANCE_ID, + url=settings.UNLEASH_API_URL, + app_name=settings.UNLEASH_APP_NAME, + environment=settings.UNLEASH_APP_NAME, + instance_id=settings.UNLEASH_INSTANCE_ID, + request_timeout=settings.UNLEASH_REQUEST_TIMEOUT, + request_retries=settings.UNLEASH_REQUEST_RETRIES, cache=UnleashRedisCache(), - environment=settings.FEATURE_FLAG_APP_NAME, + disable_registration=True, + disable_metrics=True, ) def get_flag_state_by_emails(self, name: str, emails: List[Email]) -> Mapping[Email, State]: @@ -39,7 +39,7 @@ def stream_message_reconnect(request, chat_uid: UUID, offset: int = 0): return StreamingHttpResponse( sse_chat_stream.event_stream(request=request, offset=offset), content_type='text/event-stream', - headers={'Cache-Control': 'no-cache', 'Connection': 'keep-alive', 'X-Accel-Buffering': 'no'}, + headers={'Cache-Control': 'no-cache', 'X-Accel-Buffering': 'no'}, ) @@ -69,5 +69,5 @@ def stream_message(request, chat_uid: UUID, body: MessageInSchema): return StreamingHttpResponse( sse_chat_stream.event_stream(request=request), content_type='text/event-stream', - headers={'Cache-Control': 'no-cache', 'Connection': 'keep-alive', 'X-Accel-Buffering': 'no'}, + headers={'Cache-Control': 'no-cache', 'X-Accel-Buffering': 'no'}, ) @@ -66,3 +66,14 @@ class SSEStoreService: def cleanup(self) -> None: self.delete_stream() + + +class PublicSSEStoreService(SSEStoreService): + def __init__(self, message_uuid: UUID, user_uuid: UUID): + self.user_uuid = user_uuid + self.message_uuid = message_uuid + self.redis_client = self._get_redis_client() + + def _get_cache_key(self) -> str: + return f'sse:tokens:{self.user_uuid}:{self.message_uuid}:public' + \ No newline at end of file @@ -1,25 +1,28 @@ +from decimal import Decimal + from celery import shared_task from django.conf import settings +from django.db.models import F, Value +from django.db.models.functions import Greatest + from messages.models import Message +from ml_model.models import NeuronModel +from ml_model.services.base import StreamSimpleService from tools.chats.models import Chat from tools.chats.services.sse_chunk_service import SSEChunkService -from tools.chats.services.sse_store import SSEStoreService +from tools.chats.services.sse_store import PublicSSEStoreService, SSEStoreService +from tools.public_api.models import APIKey, APIStore -@shared_task(soft_time_limit=570, time_limit=600) -def event_stream_task(chat_uuid: str, message_uuid: str, user_uuid: str) -> None: - store = SSEStoreService(user_uuid=user_uuid, chat_uuid=chat_uuid) +def _run_stream(store: SSEStoreService, message_uuid: str, service: StreamSimpleService, message): stream = None event_id = 0 try: - chat = Chat.objects.select_related('model').get(pk=chat_uuid) - message = Message.objects.get(pk=message_uuid) - event_id += 1 store.push(SSEChunkService.start(event_id, message_uuid)) - stream = chat.model.service(chat).make_stream(message) + stream = service.make_stream(message) while True: token = next(stream) if not token: @@ -29,8 +32,6 @@ def event_stream_task(chat_uuid: str, message_uuid: str, user_uuid: str) -> None except StopIteration as exc: event_id += 1 store.push(SSEChunkService.done(event_id, exc.value or ''), ttl=settings.SSE_DONE_STREAM_TTL) - except (Chat.DoesNotExist, Message.DoesNotExist) as exc: - store.push(SSEChunkService.error(event_id + 1, str(exc)), ttl=settings.SSE_DONE_STREAM_TTL) except Exception as exc: if event_id < 2: message.is_sent = False @@ -40,3 +41,55 @@ def event_stream_task(chat_uuid: str, message_uuid: str, user_uuid: str) -> None finally: if stream is not None: stream.close() + + +@shared_task(soft_time_limit=570, time_limit=600) +def event_stream_task(chat_uuid: str, message_uuid: str, user_uuid: str) -> None: + store = SSEStoreService(user_uuid=user_uuid, chat_uuid=chat_uuid) + + chat = Chat.objects.select_related('model').get(pk=chat_uuid) + message = Message.objects.get(pk=message_uuid) + + service = chat.model.service(chat) + + _run_stream(store, message_uuid, service, message) + + +@shared_task(soft_time_limit=570, time_limit=600) +def public_event_stream_task( + start_user_balance: Decimal, + message_uuid: str, + user_uuid: str, + model_slug: str, + api_key_uuid: str, + debit_api_key_limit: bool, +): + store = PublicSSEStoreService(user_uuid=user_uuid, message_uuid=message_uuid) + + message = Message.objects.get(pk=message_uuid) + api_store = APIStore.objects.select_related( + 'user', + 'user__payment_plan', + 'user__payment_plan__plan', + 'user__business_account', + 'user__business_account__group', + 'user__business_account__parent_company', + 'user__business_account__parent_company__user__payment_plan', + 'user__business_account__parent_company__user__payment_plan__plan', + ).get(pk=message.object_id) + model = NeuronModel.objects.get(slug=model_slug) + + message.content_object = api_store + message.content_object.model = model + + service = model.service(api_store) + + try: + _run_stream(store, message_uuid, service, message) + finally: + if debit_api_key_limit: + spent = start_user_balance - api_store.user.balance + if spent > 0: + APIKey.objects.filter(pk=api_key_uuid).update( + token_limit=Greatest(F('token_limit') - spent, Value(Decimal('0'))), + ) @@ -0,0 +1,172 @@ +import base64 +from io import BytesIO + +import filetype +from django.db.models import Q +from django.core.files.uploadedfile import InMemoryUploadedFile +from django.utils.translation import gettext as _ +from ninja.errors import HttpError + +from ml_model.models import NeuronModel +from tools.chats.schemas import MessageInSchema +from tools.public_api.services.openai_errors import OpenAIErrorService +from tools.public_api.services.openai_stream import OpenAIStreamService + + +def _parse_body(body: dict) -> MessageInSchema: + raw = body.get('input') + if raw is None: + raise HttpError(400, _('Missing required parameter: input')) + file, lines = None, [] + if isinstance(raw, str): + content = raw.strip() + elif not isinstance(raw, list): + raise HttpError(400, _('Invalid input payload')) + else: + for item in raw: + if not isinstance(item, dict): + continue + role = item.get('role', 'user').capitalize() + ic = item.get('content') + if isinstance(ic, str): + if t := ic.strip(): + lines.append(f'[{role}] {t}') + continue + for part in ic or []: + if not isinstance(part, dict): + continue + pt = part.get('type') + if pt in ('input_text', 'text') and (t := part.get('text')) and (t := str(t).strip()): + lines.append(f'[{role}] {t}') + elif pt in ('input_image', 'image_url') and file is None: + url = part.get('image_url') + url = url.get('url', '') if isinstance(url, dict) else str(url or '') + if 'base64' in url: + buf = BytesIO(base64.b64decode(url.split('base64', 1)[-1].lstrip(','))) + kind = filetype.guess(buf.read(20)) + buf.seek(0) + file = InMemoryUploadedFile( + buf, + 'file', + f'api-file.{kind.extension if kind else "bin"}', + kind.mime if kind else 'application/octet-stream', + buf.getbuffer().nbytes, + None, + ) + content = '\n'.join(lines).strip() + if not content and not file: + raise HttpError(400, _('The request must not be empty')) + kwargs = {'content': content or '[User]'} + if file: + kwargs['file'] = file + if isinstance(md := body.get('metadata'), dict) and md: + info = dict(md) + for k in ('ttft', 'tbt'): + if isinstance(info.get(k), str): + try: + info[k] = float(info[k]) + except ValueError: + pass + kwargs['info'] = info + return MessageInSchema(**kwargs) + + +def _resolve_model(model_ref: str) -> NeuronModel: + try: + return ( + NeuronModel.objects.filter( + Q(model_modelversions__slug=model_ref) | Q(slug=model_ref), + category__slug='chat-bots', + ) + .distinct() + .get() + ) + except NeuronModel.DoesNotExist: + raise HttpError(404, _('Model not found')) + + +def _public_user_uuid(request): + from tools.public_api.routes.v1 import _get_api_key + + return _get_api_key(request, check_usage_limit=False).user.pk + + +def _to_openai( + response, + model_ref: str, + *, + request, + message_uuid: str | None = None, + starting_after: int = 0, +): + # Django 6: .streaming_content yields bytes; wrap the raw str iterator instead. + response.streaming_content = OpenAIStreamService.translate( + response._iterator, + model_ref, + message_uuid=message_uuid, + starting_after=starting_after, + user_uuid=_public_user_uuid(request), + ) + return response + + +@OpenAIErrorService.view +def openai_responses_stream(request, body: dict): + from tools.public_api.routes.v1 import public_stream_message + + if not body.get('stream'): + raise HttpError(400, _('Only streaming is supported')) + if not (model_ref := body.get('model')): + raise HttpError(400, _('You must provide a model parameter')) + model = _resolve_model(model_ref) + return _to_openai(public_stream_message(request, model.slug, _parse_body(body)), model_ref, request=request) + + +@OpenAIErrorService.view +def openai_responses_stream_reconnect( + request, + response_id: str | None = None, + *, + stream: bool = True, + starting_after: int = 0, +): + from tools.public_api.routes.v1 import public_stream_message_reconnect + + if not stream: + raise HttpError(400, _('Only streaming reconnect is supported')) + if not (message_uuid := (response_id or request.GET.get('response_id', '')).strip()): + raise HttpError(400, _('message_uuid is not provided')) + response = public_stream_message_reconnect( + request, + message_uuid, + offset=OpenAIStreamService.public_offset(starting_after), + ) + return _to_openai( + response, '', message_uuid=message_uuid, starting_after=starting_after, request=request + ) + + +from tools.public_api.routes import v1 as public_v1_routes + + +def _register_reconnect_get(path: str, *, with_response_id: bool): + if with_response_id: + + @public_v1_routes.router.get(path, tags=['openai/responses']) + def handler(request, response_id: str, stream: bool = True, starting_after: int = 0): + return openai_responses_stream_reconnect( + request, response_id, stream=stream, starting_after=starting_after + ) + else: + + @public_v1_routes.router.get(path, tags=['openai/responses']) + def handler(request, stream: bool = True, starting_after: int = 0): + return openai_responses_stream_reconnect(request, stream=stream, starting_after=starting_after) + + +for with_rid, paths in ( + (False, ('openai/v1/responses', 'openai/responses')), + (True, ('openai/v1/responses/{response_id}', 'openai/responses/{response_id}')), +): + for path in paths: + _register_reconnect_get(path, with_response_id=with_rid) @@ -0,0 +1,141 @@ +from datetime import date + +from django.http import StreamingHttpResponse +from ninja import Body, Router +from ninja.errors import HttpError + +from django.utils.translation import gettext as _ + +from messages.models import Message +from ml_model.selectors.ml_models_selector import NeuronModelSelector + +from tools.chats.schemas import MessageInSchema +from tools.chats.services.sse_chat_stream import SSEChatStreamService +from tools.chats.services.sse_store import PublicSSEStoreService +from tools.chats.tasks import public_event_stream_task + +from tools.public_api.models import APIKey, APIStore + +router = Router(auth=None, tags=['public']) + + +def _get_api_key( + request, + select_related: list[str] = None, + prefetch_related: list[str] = None, + *, + check_usage_limit: bool = True, +): + raw_api_key = request.headers.get('Authorization', '') + if not raw_api_key: + raise HttpError(401, _('No API Key in Authorization header')) + + if (split_api_key := raw_api_key.split())[0] == 'Bearer': + raw_api_key = split_api_key[-1] + + api_key = ( + APIKey.objects.select_related('user', *(select_related or [])) + .prefetch_related(*(prefetch_related or [])) + .filter(key=raw_api_key, is_deleted=False) + ).first() + + if not api_key: + raise HttpError(404, _('API key not found')) + if check_usage_limit: + if api_key.expires_at and api_key.expires_at < date.today(): + raise HttpError(401, _('API key expired')) + if api_key.token_limit is not None and api_key.token_limit < 1: + raise HttpError(403, _('API key limit exceeded')) + + return api_key + + +@router.get('text/{message_uuid}/stream/reconnect', tags=['public/text']) +def public_stream_message_reconnect(request, message_uuid: str, offset: int = 0): + api_key = _get_api_key( + request, + select_related=['user__host_account', 'user__business_account'], + check_usage_limit=False, + ) + user = api_key.user + if user.account_type not in ('business_host', 'regular', 'business_admin'): + raise HttpError(403, _('API key is not available for this account type')) + + store = PublicSSEStoreService(user_uuid=user.pk, message_uuid=message_uuid) + if not store.exists(): + raise HttpError(404, _('Stream not found')) + + sse_chat_stream = SSEChatStreamService(store) + return StreamingHttpResponse( + sse_chat_stream.event_stream(request=request, offset=offset), + content_type='text/event-stream', + headers={'Cache-Control': 'no-cache', 'X-Accel-Buffering': 'no'}, + ) + + +@router.post('text/{model_slug}/stream', tags=['public/text']) +def public_stream_message(request, model_slug: str, body: MessageInSchema): + api_key = _get_api_key( + request, + select_related=[ + 'user__host_account', + 'user__business_account', + 'user__payment_plan', + 'user__payment_plan__plan', + 'user__business_account__parent_company', + 'user__business_account__parent_company__user__payment_plan', + 'user__business_account__parent_company__user__payment_plan__plan', + ], + prefetch_related=[ + 'user__payment_plan__plan__features', + 'user__business_account__parent_company__user__payment_plan__plan__features', + ], + ) + user = api_key.user + if user.account_type not in ('business_host', 'regular', 'business_admin'): + raise HttpError(403, _('API key is not available for this account type')) + balance = user.balance + + api_store, created = APIStore.objects.get_or_create(user=user) + + selector = NeuronModelSelector(user) + model = selector.get_model_by_slug(slug=model_slug) + if model.blocked: + raise HttpError(403, _('Model is blocked by outdating or temporary block, please retry later')) + if not model.streaming: + raise HttpError(501, _('Stream not supported for this model')) + + if not body.content: + raise HttpError(400, _('The request must not be empty')) + + data = body.model_dump(include={'content', 'file', 'info'}, exclude_unset=True) + input_message = Message.objects.create( + content_object=api_store, from_model=False, from_public_api=True, **data + ) + store = PublicSSEStoreService(user_uuid=user.pk, message_uuid=input_message.pk) + store.start() + + public_event_stream_task.delay( + start_user_balance=balance, + message_uuid=str(input_message.pk), + user_uuid=str(user.pk), + model_slug=model_slug, + api_key_uuid=str(api_key.pk), + debit_api_key_limit=api_key.token_limit is not None, + ) + + sse_chat_stream = SSEChatStreamService(store) + return StreamingHttpResponse( + sse_chat_stream.event_stream(request=request), + content_type='text/event-stream', + headers={'Cache-Control': 'no-cache', 'X-Accel-Buffering': 'no'}, + ) + + +from tools.public_api.routes.providers.openai import openai_responses_stream + + +@router.post('openai/v1/responses', tags=['openai/responses']) +@router.post('openai/responses', tags=['openai/responses']) +def openai_responses(request, body: dict = Body(...)): + return openai_responses_stream(request, body) @@ -1 +1,3 @@ from .api_key import APIKeyService +from .openai_errors import OpenAIErrorService +from .openai_stream import OpenAIStreamService @@ -0,0 +1,29 @@ +from functools import wraps + +from django.http import JsonResponse +from ninja.errors import HttpError + + +class OpenAIErrorService: + TYPES = { + 400: 'invalid_request_error', 401: 'authentication_error', 403: 'permission_error', + 404: 'invalid_request_error', 409: 'invalid_request_error', 501: 'api_error', + } + + @classmethod + def response(cls, status: int, message: str) -> JsonResponse: + return JsonResponse( + {'error': {'message': str(message), 'type': cls.TYPES.get(status, 'api_error'), 'param': None, 'code': None}}, + status=status, + ) + + @classmethod + def view(cls, fn): + @wraps(fn) + def wrapper(*args, **kwargs): + try: + return fn(*args, **kwargs) + except HttpError as exc: + return cls.response(exc.status_code, exc.message) + + return wrapper @@ -0,0 +1,164 @@ +import time +from typing import Iterator +from uuid import UUID + +import orjson + +from tools.chats.services.sse_chat_stream import SSEChatStreamService +from tools.chats.services.sse_store import PublicSSEStoreService + + +class OpenAIStreamService: + SSE_HEADERS = {'Cache-Control': 'no-cache', 'X-Accel-Buffering': 'no'} + SETUP_N = 3 + _IDX = {'output_index': 0, 'content_index': 0} + + @classmethod + def public_offset(cls, starting_after: int) -> int: + return 0 if starting_after < cls.SETUP_N - 1 else starting_after - cls.SETUP_N + 2 + + @classmethod + def translate( + cls, + public_stream, + model_ref: str, + *, + message_uuid: str | None = None, + starting_after: int = 0, + user_uuid: UUID | None = None, + ) -> Iterator[str]: + min_seq = starting_after or -1 + emit_setup = starting_after < cls.SETUP_N - 1 + meta = cls._meta(message_uuid, model_ref) if message_uuid else None + setup_sent = emit_setup and bool(meta) + tokens: list[str] = [] + last_seq = min_seq + + if setup_sent: + yield from cls._setup(meta, min_seq) + + for chunk in public_stream: + if chunk == SSEChatStreamService.HEARTBEAT: + continue + event, data, eid = cls._parse(chunk) + + if event == 'start': + meta = meta or cls._meta(data['message_uuid'], model_ref) + if emit_setup and not setup_sent: + yield from cls._setup(meta, min_seq) + setup_sent = True + continue + if event == 'pending' or not meta: + continue + + ctx = {'item_id': meta['mid'], **cls._IDX} + if event == 'token' and eid and (token := data.get('content', '')): + tokens.append(token) + if (seq := cls.SETUP_N + eid - 2) > min_seq: + yield cls._sse('response.output_text.delta', seq, delta=token, **ctx) + last_seq = seq + elif event == 'done': + text = data.get('content') or ''.join(tokens) + first = cls.SETUP_N + eid - 2 if eid else cls.SETUP_N + yield from cls._close(meta, text, first, min_seq, ctx) + cls._cleanup_public_stream(user_uuid, meta['rid']) + return + elif event == 'error': + if (seq := last_seq + 1) > min_seq: + yield cls._sse( + 'response.failed', + seq, + response={ + 'id': meta['rid'], + 'object': 'response', + 'status': 'failed', + 'error': {'code': 'server_error', 'message': str(data.get('detail', ''))}, + }, + ) + cls._cleanup_public_stream(user_uuid, meta['rid']) + return + + @classmethod + def _cleanup_public_stream(cls, user_uuid: UUID | None, message_uuid: str | None) -> None: + if user_uuid and message_uuid: + PublicSSEStoreService(user_uuid=user_uuid, message_uuid=message_uuid).cleanup() + + @classmethod + def _meta(cls, message_uuid: str, model: str) -> dict: + uid = str(message_uuid) + return {'rid': uid, 'mid': f'msg_{uid.replace("-", "")}', 'ts': int(time.time()), 'model': model} + + @classmethod + def _sse(cls, event_type: str, seq: int, **fields) -> str: + payload = {'type': event_type, 'sequence_number': seq, **fields} + return f'event: {event_type}\ndata: {orjson.dumps(payload).decode()}\n\n' + + @staticmethod + def _parse(chunk: str) -> tuple[str, dict, int | None]: + event, data, eid = '', {}, None + for line in chunk.split('\n'): + if line.startswith('id:'): + eid = int(line[3:].strip()) + elif line.startswith('event:'): + event = line[6:].strip() + elif line.startswith('data:'): + data = orjson.loads(line[5:].strip()) + return event, data, eid + + @classmethod + def _resp(cls, meta: dict, status: str, output: list) -> dict: + return { + 'id': meta['rid'], + 'object': 'response', + 'created_at': meta['ts'], + 'model': meta['model'], + 'status': status, + 'output': output, + 'store': True, + 'text': {'format': {'type': 'text'}}, + } + + @classmethod + def _setup(cls, meta: dict, min_seq: int) -> Iterator[str]: + item = { + 'id': meta['mid'], + 'type': 'message', + 'status': 'in_progress', + 'role': 'assistant', + 'content': [], + } + for seq, (etype, fields) in enumerate( + ( + ('response.created', {'response': cls._resp(meta, 'in_progress', [])}), + ('response.output_item.added', {'item': item, **cls._IDX}), + ( + 'response.content_part.added', + { + 'item_id': meta['mid'], + 'part': {'type': 'output_text', 'text': '', 'annotations': []}, + **cls._IDX, + }, + ), + ) + ): + if seq > min_seq: + yield cls._sse(etype, seq, **fields) + + @classmethod + def _close(cls, meta: dict, text: str, first_seq: int, min_seq: int, ctx: dict) -> Iterator[str]: + part = {'type': 'output_text', 'text': text, 'annotations': []} + item = { + 'id': meta['mid'], + 'type': 'message', + 'status': 'completed', + 'role': 'assistant', + 'content': [part], + } + for i, (etype, fields) in enumerate( + ( + ('response.output_text.done', {'text': text, **ctx}), + ('response.completed', {'response': cls._resp(meta, 'completed', [item])}), + ) + ): + if (seq := first_seq + i) > min_seq: + yield cls._sse(etype, seq, **fields) @@ -95,9 +95,10 @@ DOMAIN=localhost PROVIDER=docker # UNLEASH -FEATURE_FLAG_API_URL=https://gitlab.kisulkens.ru/api/v4/feature_flags/unleash/243 -FEATURE_FLAG_INSTANCE_ID=glffct-ic8xsVF5eR9BaUySR-_w -FEATURE_FLAG_APP_NAME=Production +UNLEASH_API_URL=https://example.com +UNLEASH_APP_NAME=Development +UNLEASH_REQUEST_TIMEOUT=3 +UNLEASH_REQUEST_RETRIES=1 # ZROK ZROK2_API_ENDPOINT=https://zrok2.example.com @@ -108,5 +109,6 @@ COMPOSE_FILE=docker-compose.yml:docker-compose.local.yml COMPOSE_PROFILES="" # use-tunnel available BUILDKIT_PROGRESS=plain -# HTTP FRAMEWORK -GUNICORN_CMD_ARGS="-b 0.0.0.0:8000 -k gthread -w 3 -t 600 --reload" \ No newline at end of file +# SERVER +DJANGO_RUNSERVER_HIDE_WARNING=true +PYTHONWARNINGS=ignore::UserWarning:polymorphic # temporarily \ No newline at end of file @@ -7,6 +7,12 @@ x-dev-app-config: &dev-app-config services: app: <<: *dev-app-config + command: + - /bin/sh + - -c + - | + python manage.py compilemessages --locale ru_RU + python manage.py runserver 0.0.0.0:8000 ports: - "8000:8000" @@ -96,7 +96,7 @@ services: - /bin/sh - -c - | - python manage.py create_indexes --skip-system-checks + python manage.py create_indexes --skip-checks python manage.py initialize_buckets --skip-checks python manage.py collectstatic --no-input --skip-checks python manage.py migrate