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"