File size: 2,058 Bytes
cd0c7a9 | 1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 | import json
import logging
from typing import AsyncGenerator
from litellm import acompletion
from app.config import settings
from app.ai.llm_client import llm_client
from app.ai.prompts import get_prompt
logger = logging.getLogger(__name__)
async def interpret_stream(pipeline_type: str, context: dict) -> AsyncGenerator[str, None]:
providers = llm_client.get_providers()
if not providers:
yield _error_event("No LLM API keys configured. AI interpretation unavailable.")
return
prompt = llm_client.build_prompt(pipeline_type, context)
last_error = None
for provider in providers:
try:
response = await acompletion(
model=provider["model"],
messages=[{"role": "user", "content": prompt}],
temperature=0.3,
max_tokens=2000,
stream=True,
timeout=25,
api_key=provider["api_key"],
)
async for chunk in response:
if chunk.choices and chunk.choices[0].delta.content:
yield _chunk_event(chunk.choices[0].delta.content)
yield _done_event({"model": provider["model"], "pipeline_type": pipeline_type})
return
except Exception as e:
last_error = e
logger.warning("LLM provider %s failed: %s", provider["name"], e)
continue
msg = str(last_error) if last_error else "All providers failed"
if "organization_restricted" in msg or "Organization has been restricted" in msg:
yield _error_event("AI interpretation is temporarily unavailable due to a provider restriction. Please try again later.")
else:
yield _error_event(f"AI interpretation failed: {msg}")
def _chunk_event(text: str) -> str:
return f"data: {json.dumps({'chunk': text})}\n\n"
def _done_event(meta: dict) -> str:
return f"data: {json.dumps({'done': True, 'meta': meta})}\n\n"
def _error_event(msg: str) -> str:
return f"data: {json.dumps({'error': msg})}\n\n"
|