diff --git a/.api.env b/.api.env index 5f00152..3ef3f5a 100644 --- a/.api.env +++ b/.api.env @@ -27,4 +27,4 @@ PIXVERSE_API_KEY=sk-2c785f1f77ace5b4f39cb3d4dc5ef554 INTERNAL_TELEGRAM_SECRET=do8asyd0h21uodh2od2hkdmbzxc2349ASAX80scasokcu23ked2zxc LIARA_API_URL =https://ai.liara.ir/api/68eb653bb55873971e0d46d6/v1 LIARA_API_KEY=eyJhbGciOiJIUzI1NiIsInR5cCI6IkpXVCJ9.eyJrZXkiOiI2OThhZGIxMWQzODZhNWVmODNlMTg2YzAiLCJ0eXBlIjoiYWlfa2V5IiwiaWF0IjoxNzcwNzA3NzMwfQ.ckoR00Uxt8W4DmLCGW24P46Zq-et0Yhtu2xej5P0ClQ -KIA_API_KEY =7ed05fadbea55c56b43c69c2061472ac \ No newline at end of file +KIA_API_KEY =9cc6da3bef9560efcb593382beda8774 \ No newline at end of file diff --git a/src/services/chatbot.py b/src/services/chatbot.py index c9865c5..2485547 100644 --- a/src/services/chatbot.py +++ b/src/services/chatbot.py @@ -448,11 +448,16 @@ async def send_message( input_messages = [human_msg] async for chunk, metadata in app.astream({"messages": input_messages}, config, stream_mode="messages"): - if isinstance(chunk, AIMessage): - if chunk.usage_metadata: - input_tokens += chunk.usage_metadata["input_tokens"] - output_tokens += chunk.usage_metadata["output_tokens"] - yield json.dumps({"content": chunk.content}, ensure_ascii=False) + try: + if isinstance(chunk, AIMessage): + if chunk.usage_metadata: + input_tokens += chunk.usage_metadata["input_tokens"] + output_tokens += chunk.usage_metadata["output_tokens"] + yield json.dumps({"content": chunk.content}, ensure_ascii=False) + except Exception as e: + print(f"[send_message] Error processing chunk: {e}") + yield json.dumps({"error": True, "detail": str(e)}, ensure_ascii=False) + return await user.save() diff --git a/src/services/kiaai/chat.py b/src/services/kiaai/chat.py index 399ac03..d3ed16a 100644 --- a/src/services/kiaai/chat.py +++ b/src/services/kiaai/chat.py @@ -1,4 +1,5 @@ import json +import asyncio import httpx from langchain_core.messages import ( AIMessage, @@ -89,40 +90,50 @@ class KiaAIService: model: str, messages: list[BaseMessage], reasoning_effort: str | None = None, + max_retries: int = 3, + retry_delay: float = 2.0, ) -> AIMessage: """ Call KIA API and return an assembled AIMessage. - - Args: - model: Model name with optional 'kia/' prefix (e.g. 'kia/gpt-5-2') - messages: LangChain message list - reasoning_effort: 'low' | 'medium' | 'high' — only for supported models + Retries up to max_retries times on 500 server errors. """ model_name = model.removeprefix("kia/") url = f"{KIA_BASE_URL}/{model_name}/v1/chat/completions" - payload: dict = { - "messages": _to_openai_messages(messages), - } - - # Add reasoning_effort only for models that support it - if reasoning_effort and model_name in REASONING_MODELS: - payload["reasoning_effort"] = reasoning_effort + openai_messages = _to_openai_messages(messages) print(f"[KiaAI] POST {url}") - print(f"[KiaAI] messages_count={len(messages)} | reasoning_effort={reasoning_effort}") - - openai_messages = _to_openai_messages(messages) + print(f"[KiaAI] messages_count={len(openai_messages)} | reasoning_effort={reasoning_effort}") for i, m in enumerate(openai_messages): - content_preview = m["content"][:120] if isinstance(m["content"], str) else str(m["content"])[:120] - print(f"[KiaAI] msg[{i}] role={m['role']} | content={content_preview!r}") + preview = m["content"][:120] if isinstance(m["content"], str) else str(m["content"])[:120] + print(f"[KiaAI] msg[{i}] role={m['role']} | content={preview!r}") - payload: dict = {"messages": openai_messages} - - # Add reasoning_effort only for models that support it + payload: dict = { + "model": model_name, + "messages": openai_messages, + "stream": True, + } if reasoning_effort and model_name in REASONING_MODELS: payload["reasoning_effort"] = reasoning_effort + last_error: Exception | None = None + + for attempt in range(1, max_retries + 1): + if attempt > 1: + print(f"[KiaAI] Retry attempt {attempt}/{max_retries} after {retry_delay}s...") + await asyncio.sleep(retry_delay) + + try: + result = await self._do_request(url, payload) + return result + except RuntimeError as e: + last_error = e + print(f"[KiaAI] Attempt {attempt} failed: {e}") + + raise last_error + + async def _do_request(self, url: str, payload: dict) -> AIMessage: + """Execute a single SSE streaming request to KIA and assemble AIMessage.""" full_content = "" usage: dict = {} @@ -140,8 +151,6 @@ class KiaAIService: resp.raise_for_status() async for line in resp.aiter_lines(): - print(f"[KiaAI] RAW LINE: {line!r}") - if not line.strip(): continue @@ -157,14 +166,12 @@ class KiaAIService: print(f"[KiaAI] Skipping non-JSON line: {raw_data!r}") continue - # KIA server error inside stream - if chunk.get("code") and chunk["code"] != 200: - error_msg = chunk.get("msg", "Unknown KIA error") - print(f"[KiaAI] Server error in stream: code={chunk['code']} msg={error_msg!r}") + # KIA server error inside stream body + if chunk.get("code") and int(chunk["code"]) >= 500: + error_msg = chunk.get("msg", "Unknown KIA server error") + print(f"[KiaAI] Server error chunk: code={chunk['code']} msg={error_msg!r}") raise RuntimeError(f"KIA API error {chunk['code']}: {error_msg}") - print(f"[KiaAI] CHUNK keys={list(chunk.keys())} choices={chunk.get('choices')}") - choices = chunk.get("choices") or [] for choice in choices: delta = choice.get("delta") or {}