"""Implementacje dostawców LLM (LOG-31) — na samym httpx, bez SDK. Świadomie bez bibliotek `openai` / `anthropic`: lokalny serwer modelu (Ollama, vLLM, llama.cpp) i OpenAI mówią **tym samym** protokołem `/chat/completions`, więc jedna implementacja obsługuje oba — różni je tylko adres i klucz. Anthropic ma własny kształt `/v1/messages`, stąd druga klasa. Mniej zależności, mniej powierzchni ataku, pełna kontrola nad tym, co wychodzi z sieci. **Gwarancja niepustej odpowiedzi.** Horoskop ma powstać niezależnie od objętości promptu, więc `generate()` nie jest pojedynczym strzałem, tylko pętlą: 1. wyślij turę z policzonym limitem wyjścia, 2. jeśli model urwał na limicie — dopisz turę „kontynuuj" i sklej tekst, 3. jeśli tura nie dała ani znaku tekstu — ponów z podpowiedzią, 4. dopiero brak tekstu po wszystkich próbach jest błędem (z diagnostyką). Kontynuacja jest pewniejsza niż jedno wielkie żądanie: każda tura mieści się w timeoucie HTTP, a długość odpowiedzi przestaje być ograniczona jedną turą. **Anthropic i myślenie.** Modele Claude potrafią mieć włączone myślenie, którego tokeny liczą się do `max_tokens`. Przy ciasnym limicie cała tura potrafi wyjść jako same bloki `thinking` z pustym tekstem — dokładnie ten objaw, który zgłoszono. Traktujemy taką turę jak ucięcie i kontynuujemy, zamiast zwracać pustkę. """ from __future__ import annotations import os import time import httpx from app.llm.base import Completion, LLMError, LLMProvider RETRY_STATUSES = {429, 500, 502, 503, 504} MAX_ATTEMPTS = 3 # ile razy wolno poprosić model o dokończenie urwanej odpowiedzi MAX_CONTINUATIONS = 12 # ile tokenów zamawiać na jedną turę — mieści się w timeoucie, a pętla i tak # dociągnie resztę; zbyt duża wartość ryzykuje zerwanie połączenia w trakcie TURN_TOKENS_CAP = 16_000 _CONTINUE = ( "Kontynuuj dokładnie od miejsca, w którym przerwałeś — nie powtarzaj tego, " "co już napisałeś, i nie zaczynaj od nowa. Jeśli skończyłeś całą odpowiedź, " "napisz wyłącznie: KONIEC" ) _NUDGE = ( "Nie otrzymałem żadnej treści. Napisz odpowiedź zgodnie z powyższym poleceniem, " "zaczynając od razu od treści horoskopu." ) _DONE_MARKER = "KONIEC" def _post_with_retry(url: str, headers: dict, payload: dict, timeout: float) -> dict: """POST z ponawianiem i backoffem — chroni przed chwilowym 429/5xx.""" last: Exception | None = None for attempt in range(MAX_ATTEMPTS): try: with httpx.Client(timeout=timeout) as client: r = client.post(url, json=payload, headers=headers) if r.status_code in RETRY_STATUSES and attempt < MAX_ATTEMPTS - 1: time.sleep(2 ** attempt) continue if r.status_code >= 400: raise LLMError(f"Model odpowiedział błędem {r.status_code}: {r.text[:300]}") return r.json() except httpx.TimeoutException as e: last = e if attempt < MAX_ATTEMPTS - 1: time.sleep(2 ** attempt) continue raise LLMError( f"Model nie odpowiedział w czasie {timeout:.0f}s. Zwiększ LLM_TIMEOUT " f"albo zmniejsz budżet promptu." ) from e except httpx.HTTPError as e: last = e if attempt < MAX_ATTEMPTS - 1: time.sleep(2 ** attempt) continue raise LLMError(f"Nie udało się połączyć z modelem: {e}") from e raise LLMError(f"Nie udało się wywołać modelu: {last}") def _merge_usage(total: dict, turn: dict) -> dict: """Sumuje zużycie tokenów przez wszystkie tury jednej odpowiedzi.""" for key, value in (turn or {}).items(): if isinstance(value, int): total[key] = total.get(key, 0) + value return total def _join(parts: list[str]) -> str: return "".join(parts).strip() def _explain_empty(turns: int, usage: dict, stop: str | None) -> str: detail = [] if stop: detail.append(f"powód zakończenia: {stop}") for key in ("completion_tokens", "output_tokens"): if usage.get(key) is not None: detail.append(f"tokeny odpowiedzi: {usage[key]}") break suffix = f" ({', '.join(detail)})" if detail else "" return ( f"Model nie zwrócił żadnej treści po {turns} próbach{suffix}. " f"Najczęstsza przyczyna: prompt wypełnił okno kontekstu i nie zostało miejsca " f"na odpowiedź. Zmniejsz budżet promptu albo wybierz model z większym oknem." ) class _Driver: """Wspólna pętla: tura → ewentualna kontynuacja → sklejony tekst. Podklasy dostarczają tylko `_turn()` — reszta (kontynuacje, ponawianie pustej tury, sumowanie zużycia) jest identyczna dla obu protokołów. """ name: str model: str leaves_lan: bool def _turn(self, messages: list[dict], max_tokens: int): """(tekst, czy_ucięta, zużycie, nazwa_modelu, powód_zakończenia).""" raise NotImplementedError def generate(self, prompt: str, max_tokens: int) -> Completion: messages: list[dict] = [{"role": "user", "content": prompt}] parts: list[str] = [] usage: dict = {} model_name = self.model remaining = max(max_tokens, 256) stop: str | None = None turns = 0 nudged = False while turns <= MAX_CONTINUATIONS: turns += 1 budget = max(256, min(remaining, TURN_TOKENS_CAP)) text, truncated, turn_usage, model_name, stop = self._turn(messages, budget) _merge_usage(usage, turn_usage) remaining -= budget chunk = text.strip() if chunk: if chunk.endswith(_DONE_MARKER): # model zgłasza koniec parts.append(("\n" if parts else "") + chunk[: -len(_DONE_MARKER)].rstrip()) break parts.append(("\n" if parts else "") + chunk) if not truncated: break elif not truncated: # pusta i NIE ucięta: jedna próba z podpowiedzią, potem koniec if nudged or parts: break nudged = True messages = messages + [ {"role": "assistant", "content": "…"}, {"role": "user", "content": _NUDGE}, ] continue # ucięta (także tura złożona z samego myślenia) — poproś o dokończenie if remaining < 256: break messages = [ {"role": "user", "content": prompt}, {"role": "assistant", "content": _join(parts) or "…"}, {"role": "user", "content": _CONTINUE}, ] final = _join(parts) if not final: raise LLMError(_explain_empty(turns, usage, stop)) usage["turns"] = turns return Completion(text=final, model=model_name, provider=self.name, leaves_lan=self.leaves_lan, usage=usage) def count_tokens(self, prompt: str) -> int: """Szacunek tokenów promptu. Dostawcy z własnym licznikiem nadpisują.""" return int(len(prompt) / 3.6) class ChatCompletionsProvider(_Driver, LLMProvider): """Protokół OpenAI `/chat/completions` — lokalny serwer modelu ORAZ OpenAI.""" def __init__(self, name: str, base_url: str, model: str, api_key: str = "", timeout: float = 120.0, leaves_lan: bool = True) -> None: self.name = name self.base_url = base_url.rstrip("/") self.model = model self.api_key = api_key self.timeout = timeout self.leaves_lan = leaves_lan def _headers(self) -> dict: h = {"Content-Type": "application/json"} if self.api_key: h["Authorization"] = f"Bearer {self.api_key}" return h def _turn(self, messages: list[dict], max_tokens: int): data = _post_with_retry( f"{self.base_url}/chat/completions", self._headers(), {"model": self.model, "max_tokens": max_tokens, "messages": messages}, self.timeout, ) try: choice = data["choices"][0] text = choice["message"].get("content") or "" except (KeyError, IndexError, TypeError) as e: raise LLMError(f"Nieoczekiwany kształt odpowiedzi modelu: {str(data)[:300]}") from e stop = choice.get("finish_reason") return (text, stop == "length", data.get("usage") or {}, data.get("model", self.model), stop) def health(self) -> dict: info = {"provider": self.name, "model": self.model, "leaves_lan": self.leaves_lan} try: with httpx.Client(timeout=min(self.timeout, 10.0)) as client: r = client.get(f"{self.base_url}/models", headers=self._headers()) info["status"] = "ok" if r.status_code < 400 else f"http {r.status_code}" except httpx.HTTPError as e: info["status"] = f"down: {e}" return info class AnthropicProvider(_Driver, LLMProvider): """Protokół Anthropic `/v1/messages`.""" leaves_lan = True def __init__(self, base_url: str, model: str, api_key: str = "", timeout: float = 120.0) -> None: self.name = "anthropic" self.base_url = base_url.rstrip("/") self.model = model self.api_key = api_key self.timeout = timeout def _headers(self) -> dict: return { "Content-Type": "application/json", "x-api-key": self.api_key, "anthropic-version": "2023-06-01", } def _thinking(self) -> dict: """Konfiguracja myślenia. Domyślnie adaptacyjne — podnosi jakość tekstu. UWAGA: tokeny myślenia liczą się do `max_tokens`, więc przy ciasnym limicie cała tura potrafi wyjść jako samo myślenie z pustym tekstem. Pętla kontynuacji to obsługuje, ale ANTHROPIC_THINKING=off wyłącza myślenie, gdy zależy nam na przewidywalnym zużyciu tokenów. """ mode = os.getenv("ANTHROPIC_THINKING", "adaptive").lower() if mode in ("off", "disabled", "0", "false"): return {"thinking": {"type": "disabled"}} return { "thinking": {"type": "adaptive"}, "output_config": {"effort": os.getenv("ANTHROPIC_EFFORT", "high")}, } def _turn(self, messages: list[dict], max_tokens: int): if not self.api_key: raise LLMError("Brak ANTHROPIC_API_KEY — dostawca anthropic wymaga klucza.") payload = {"model": self.model, "max_tokens": max_tokens, "messages": messages} payload.update(self._thinking()) data = _post_with_retry(f"{self.base_url}/v1/messages", self._headers(), payload, self.timeout) try: blocks = data["content"] text = "".join(b.get("text", "") for b in blocks if b.get("type") == "text") except (KeyError, TypeError) as e: raise LLMError(f"Nieoczekiwany kształt odpowiedzi modelu: {str(data)[:300]}") from e stop = data.get("stop_reason") # tura złożona z samego myślenia = budżet poszedł na rozumowanie; traktujemy # jak ucięcie, żeby pętla poprosiła o treść zamiast zwrócić pustkę thinking_only = not text.strip() and any( b.get("type") in ("thinking", "redacted_thinking") for b in blocks ) return (text, stop == "max_tokens" or thinking_only, data.get("usage") or {}, data.get("model", self.model), stop) def count_tokens(self, prompt: str) -> int: """Dokładny licznik Anthropic — nie szacunek. Od tego zależy, czy po zmieszczeniu promptu zostanie miejsce na odpowiedź.""" if not self.api_key: return super().count_tokens(prompt) try: data = _post_with_retry( f"{self.base_url}/v1/messages/count_tokens", self._headers(), {"model": self.model, "messages": [{"role": "user", "content": prompt}]}, min(self.timeout, 30.0), ) return int(data.get("input_tokens") or super().count_tokens(prompt)) except LLMError: return super().count_tokens(prompt) def health(self) -> dict: return { "provider": self.name, "model": self.model, "leaves_lan": True, "status": "ok (klucz ustawiony)" if self.api_key else "brak ANTHROPIC_API_KEY", }