From 453c7afec5fdc350d4598595b305d0b590d047f1 Mon Sep 17 00:00:00 2001 From: Santiago Quiroz upegui Date: Thu, 20 Aug 2026 00:24:56 -0500 Subject: [PATCH] =?UTF-8?q?Historia=20t=C3=A9cnica:=20Cierre=20de=20backlo?= =?UTF-8?q?g=20v3=20=E2=80=94=20failover,=20compresi=C3=B3n=20sem=C3=A1nti?= =?UTF-8?q?ca,=20bot=20Telegram,=20RPC=20multi-host,=20logs=20litellm=20y?= =?UTF-8?q?=20deploy=20Docker?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Dominio: - providers_service: pick_provider — provider efectivo por request con failover: si el destino con api_base local/LAN no responde (chequeo TCP 0.4s), cae al primer fallback alcanzable de fallback_provider_ids; retorna flag es-provider-activo para decidir litellm vs directo. - Nuevo compression_service: compresión semántica opt-in — al acercarse al límite de contexto resume la mitad vieja del historial con el provider activo (sin partir pares tool_use/tool_result); cualquier fallo cae al truncado clásico. - Nuevo telegram_bot: relay opt-in por long-polling al gateway (TELEGRAM_BOT_TOKEN + allowlist de chat ids; allowlist vacía = inerte). - llamacpp_service: rpc_servers en local_launch (--rpc, granja multi-PC); tail_file movido a core/utils. - models/provider.py: fallback_provider_ids en el registry. Aplicación: - messages.py y openai_compat.py: resolución vía pick_provider; anthropic vía litellm solo si es el provider activo; hook de compresión semántica en el punto de truncado. - api/providers.py: routing GET/PUT extendido con fallback_provider_ids (None = no tocar). - api/proxy.py: GET /proxy/logs (tail de litellm-out/err). - api/settings.py: auth-info expone semantic_compression. Infraestructura: - settings_service.write_env_key: aplica en caliente (os.environ + cache_clear). - Frontend: campo failover en RoutingPanel, toggle de compresión semántica y viewer de logs litellm en Settings, proxyApi.getLogs, tipos actualizados. - main.py: task del bot Telegram en lifespan. Configuración: - Dockerfile multi-stage (node build + python slim + litellm[proxy]), docker-compose con volumen persistente, .dockerignore; sección de deploy en README; imagen construida y smoke-testeada localmente (/api/health 200 en contenedor). - config.py: flag semantic_compression; bump version a 2.12.0. Pruebas: - test_failover.py (6), test_compression_service.py (13), test_telegram_bot.py (13); test_routing/test_messages_native/test_openai_compat actualizados a pick_provider. Cobertura global del proyecto: suite backend 147 passed; tsc frontend sin errores. --- .dockerignore | 13 ++ CLAUDE.md | 4 +- Dockerfile | 23 ++ README.md | 12 + backend/app/api/messages.py | 27 ++- backend/app/api/openai_compat.py | 15 +- backend/app/api/providers.py | 11 +- backend/app/api/proxy.py | 13 ++ backend/app/api/settings.py | 1 + backend/app/core/config.py | 4 + backend/app/core/utils.py | 14 ++ backend/app/main.py | 2 + backend/app/models/provider.py | 2 + backend/app/services/compression_service.py | 97 +++++++++ backend/app/services/llamacpp_service.py | 21 +- backend/app/services/providers_service.py | 64 +++++- backend/app/services/settings_service.py | 6 + backend/app/services/telegram_bot.py | 230 ++++++++++++++++++++ backend/tests/test_compression_service.py | 186 ++++++++++++++++ backend/tests/test_failover.py | 200 +++++++++++++++++ backend/tests/test_messages_native.py | 11 +- backend/tests/test_openai_compat.py | 8 +- backend/tests/test_routing.py | 3 +- backend/tests/test_telegram_bot.py | 79 +++++++ docker-compose.yml | 13 ++ frontend/package-lock.json | 4 +- frontend/package.json | 2 +- frontend/src/components/RoutingPanel.tsx | 20 +- frontend/src/pages/Settings.tsx | 44 +++- frontend/src/services/api.ts | 1 + frontend/src/types/provider.ts | 1 + 31 files changed, 1079 insertions(+), 52 deletions(-) create mode 100644 .dockerignore create mode 100644 Dockerfile create mode 100644 backend/app/services/compression_service.py create mode 100644 backend/app/services/telegram_bot.py create mode 100644 backend/tests/test_compression_service.py create mode 100644 backend/tests/test_failover.py create mode 100644 backend/tests/test_telegram_bot.py create mode 100644 docker-compose.yml diff --git a/.dockerignore b/.dockerignore new file mode 100644 index 0000000..e612f74 --- /dev/null +++ b/.dockerignore @@ -0,0 +1,13 @@ +.git +build +dist +frontend/node_modules +frontend/dist +backend/__pycache__ +**/__pycache__ +**/*.pyc +.playwright-mcp +docs +*.png +*.jar +*.html diff --git a/CLAUDE.md b/CLAUDE.md index 02cf5b6..c0a0897 100644 --- a/CLAUDE.md +++ b/CLAUDE.md @@ -61,7 +61,9 @@ pyinstaller bipolar-code.spec # run from repo root **Provider registry** persists in `{config_dir}/providers.json`. Built-in providers include `copilot`, `anthropic`, `lmstudio`, `nvidia_nim`, `openrouter`, `deepseek`, `ollama` and `llamacpp`. litellm always exposes the aliases `claude-sonnet-4-6`, `claude-opus-4-6`, `gpt-4o` regardless of the active backend. Providers with `anthropic_native: true` (llama-server, LM Studio ≥0.4.1, Ollama 2026+) receive `/v1/messages` verbatim — no litellm, no OAI translation. The `llamacpp` provider spawns a managed local `llama-server` (Vulkan multi-GPU, port 4002) via `/api/llamacpp/*`. `/v1/chat/completions` on :8000 exposes the active provider as an OpenAI-compatible BYOK endpoint (VS Code Copilot Chat, Cursor, Cline). All `/v1/*` routes require auth (`ui_api_key`, or legacy `proxy_api_key`). -**Scenario routing**: `ProviderRegistry.routing_rules` (UI: Providers → "Routing por escenario") route each request by requested model name — first match wins; `pattern` = case-insensitive substring, `min_tokens` = longContext threshold. E.g. `haiku` → small local model, `opus` → real Anthropic (routed anthropic goes DIRECT to api.anthropic.com, not through litellm), 60k+ tokens → long-context provider. `local_launch.router_mode` runs llama-server without `--model` serving every GGUF in the models dir with dynamic load/unload (`--models-dir`). +**Scenario routing**: `ProviderRegistry.routing_rules` (UI: Providers → "Routing por escenario") route each request by requested model name — first match wins; `pattern` = case-insensitive substring, `min_tokens` = longContext threshold. E.g. `haiku` → small local model, `opus` → real Anthropic (routed anthropic goes DIRECT to api.anthropic.com, not through litellm), 60k+ tokens → long-context provider. `local_launch.router_mode` runs llama-server without `--model` serving every GGUF in the models dir with dynamic load/unload (`--models-dir`). `local_launch.rpc_servers` adds remote `ggml-rpc-server` workers (`--rpc`, multi-PC farm). + +**Failover**: `ProviderRegistry.fallback_provider_ids` — if the effective provider has a local/LAN api_base and its port doesn't answer (0.4s TCP check in `pick_provider`), the request falls to the first reachable fallback. **Semantic compression** (`SEMANTIC_COMPRESSION=true`, toggle in Settings): near the context limit, `compression_service` summarizes the old half of the conversation via the active provider instead of truncating; any failure falls back to plain truncation. **Telegram bot** (opt-in): `TELEGRAM_BOT_TOKEN` + `TELEGRAM_ALLOWED_CHAT_IDS` env vars — long-polling relay to the active provider; empty allowlist = inert. **Docker**: `docker compose up -d --build` deploys the gateway; llama-server stays on the host. **Platform guards**: `providers_service._start_litellm` uses PowerShell on Windows and `subprocess.Popen(start_new_session=True)` on Linux/macOS. `proxy_service._set_user_env` writes the Windows registry only on `sys.platform == "win32"`; it always writes `~/.claude/settings.json`. diff --git a/Dockerfile b/Dockerfile new file mode 100644 index 0000000..7bf579f --- /dev/null +++ b/Dockerfile @@ -0,0 +1,23 @@ +# bipolar-code — imagen de gateway (backend + UI compilada). +# llama-server NO va dentro: corre en el host con las GPUs; apunta el provider +# llamacpp a http://host.docker.internal:4002. +FROM node:20-alpine AS frontend +WORKDIR /app/frontend +COPY frontend/package.json frontend/package-lock.json ./ +RUN npm ci +COPY frontend/ ./ +RUN npm run build + +FROM python:3.11-slim +WORKDIR /app/backend +COPY backend/requirements.txt . +RUN pip install --no-cache-dir -r requirements.txt "litellm[proxy]>=1.60" +COPY backend/ . +COPY --from=frontend /app/frontend/dist /app/frontend/dist + +ENV LITELLM_CONFIG_DIR=/data \ + PYTHONUNBUFFERED=1 +VOLUME ["/data"] +EXPOSE 8000 + +CMD ["uvicorn", "app.main:app", "--host", "0.0.0.0", "--port", "8000"] diff --git a/README.md b/README.md index fb29870..291c9ba 100644 --- a/README.md +++ b/README.md @@ -29,6 +29,18 @@ Web UI para gestionar un proxy [LiteLLM](https://github.com/BerriAI/litellm) — --- +## Deploy con Docker (servidor remoto / headless) + +```bash +docker compose up -d --build +``` + +- Gateway + UI en `:8000`; config persistida en el volumen `bipolar-data`. +- `llama-server` NO va dentro del contenedor (necesita las GPUs del host): córrelo en el host y apunta el provider `llamacpp` a `http://host.docker.internal:4002`. +- Fuera de tu LAN usa una VPN (Tailscale/WireGuard) — no expongas el puerto a internet. + +--- + ## Dev Setup (código fuente) ### Requisitos diff --git a/backend/app/api/messages.py b/backend/app/api/messages.py index 9690e76..3fef862 100644 --- a/backend/app/api/messages.py +++ b/backend/app/api/messages.py @@ -11,7 +11,7 @@ from app.core.config import get_settings from app.core.logging import get_logger from app.core.utils import sanitize_error as _sanitize_error -from app.services import providers_service, token_service, usage_tracker +from app.services import compression_service, providers_service, token_service, usage_tracker from app.services.pricing_service import estimate_cost log = get_logger(__name__) @@ -326,20 +326,27 @@ async def messages_passthrough(request: Request): ctx_window = token_service.get_context_window(model) used = token_service.count_tokens(messages) - active = providers_service.get_active_provider() - routed_model: str | None = None - route = providers_service.resolve_route(model, used) - if route: - active, routed_model = route + active, routed_model, is_active_provider = await providers_service.pick_provider(model, used) active_provider_id = active.id if active else "unknown" - # Provider anthropic RUTEADO va directo a api.anthropic.com: litellm corre - # con el config del provider activo, no del destino de la regla - is_native = bool(active and (active.anthropic_native or (route and active.litellm_prefix == "anthropic"))) + # anthropic vía litellm SOLO si es el provider activo configurado (litellm + # corre con SU config); ruteado o failover → directo a api.anthropic.com + is_native = bool(active and (active.anthropic_native or (active.litellm_prefix == "anthropic" and not is_active_provider))) is_anthropic = bool(active and active.litellm_prefix == "anthropic" and not is_native) truncated = False if ctx_window > 0 and used >= int(ctx_window * 0.9): - messages = token_service.truncate_messages(messages, ctx_window) + compressed = None + if settings.semantic_compression and active: + compressed = await compression_service.compress_messages(messages, active, model) + if compressed: + messages = compressed + log.info( + "semantic_compression_applied", + before_tokens=used, + after_tokens=token_service.count_tokens(messages), + ) + else: + messages = token_service.truncate_messages(messages, ctx_window) body["messages"] = messages truncated = True diff --git a/backend/app/api/openai_compat.py b/backend/app/api/openai_compat.py index aaa3232..f9cfd78 100644 --- a/backend/app/api/openai_compat.py +++ b/backend/app/api/openai_compat.py @@ -62,13 +62,14 @@ def _record_usage(provider_id: str, model: str, usage: dict) -> None: async def chat_completions(request: Request): settings = get_settings() body = await request.json() - active = providers_service.get_active_provider() - routed_model = None - route = providers_service.resolve_route(str(body.get("model", ""))) - # Routing en la superficie OAI: solo destinos OpenAI-compat (un destino - # anthropic requeriría traducir el formato, cosa que esta ruta no hace) - if route and route[0].litellm_prefix != "anthropic": - active, routed_model = route + active, routed_model, is_active_provider = await providers_service.pick_provider( + str(body.get("model", "")) + ) + # Un destino anthropic NO activo requeriría traducir el formato OAI→Anthropic + # (litellm corre con el config del activo): en ese caso se ignora la ruta + if active and active.litellm_prefix == "anthropic" and not is_active_provider: + active = providers_service.get_active_provider() + routed_model = None provider_id = active.id if active else "unknown" url, headers, model = resolve_target(active, settings) diff --git a/backend/app/api/providers.py b/backend/app/api/providers.py index c33bddf..88f4815 100644 --- a/backend/app/api/providers.py +++ b/backend/app/api/providers.py @@ -102,20 +102,27 @@ def list_providers(): @router.get("/routing") def get_routing(): registry = providers_service.load_registry() - return {"enabled": registry.routing_enabled, "rules": registry.routing_rules} + return { + "enabled": registry.routing_enabled, + "rules": registry.routing_rules, + "fallback_provider_ids": registry.fallback_provider_ids, + } class RoutingUpdate(BaseModel): enabled: bool rules: list[RoutingRule] = [] + fallback_provider_ids: Optional[list[str]] = None # None = no tocar @router.put("/routing") def set_routing(body: RoutingUpdate): unknown = [r.provider_id for r in body.rules if not providers_service.get_provider(r.provider_id)] + if body.fallback_provider_ids: + unknown += [pid for pid in body.fallback_provider_ids if not providers_service.get_provider(pid)] if unknown: raise HTTPException(status_code=400, detail=f"Providers no registrados: {unknown}") - result = providers_service.set_routing(body.enabled, body.rules) + result = providers_service.set_routing(body.enabled, body.rules, body.fallback_provider_ids) log.info("routing_updated", enabled=body.enabled, rules=len(body.rules)) return result diff --git a/backend/app/api/proxy.py b/backend/app/api/proxy.py index ea44472..96e7bc2 100644 --- a/backend/app/api/proxy.py +++ b/backend/app/api/proxy.py @@ -65,3 +65,16 @@ async def set_proxy_route(body: RouteRequest): except Exception as e: log.error('proxy_set_route_error', error=str(e), requested=body.mode) raise HTTPException(status_code=500, detail=str(e)) + + +@router.get('/logs') +async def litellm_logs(lines: int = 80): + from pathlib import Path + from app.core.config import get_settings + from app.core.utils import tail_file + + config_dir = Path(get_settings().litellm_config_dir) + merged = [] + for name, tag in (('litellm-out.log', 'out'), ('litellm-err.log', 'err')): + merged += [f'[{tag}] {line}' for line in tail_file(config_dir / name, lines)] + return {'logs': merged[-lines:]} diff --git a/backend/app/api/settings.py b/backend/app/api/settings.py index 2a56f86..f726dbd 100644 --- a/backend/app/api/settings.py +++ b/backend/app/api/settings.py @@ -38,6 +38,7 @@ async def get_auth_info(): "rate_limit_rpm": s.rate_limit_rpm, "allowed_origins": s.allowed_origins, "proxy_base_url": "http://:8000", + "semantic_compression": s.semantic_compression, } diff --git a/backend/app/core/config.py b/backend/app/core/config.py index a5a720f..b915fe5 100644 --- a/backend/app/core/config.py +++ b/backend/app/core/config.py @@ -49,6 +49,10 @@ class Settings(BaseSettings): # Rate limiting: máx requests por minuto por IP (0 = desactivado) rate_limit_rpm: int = 120 + # Compresión semántica: resumir historial viejo con el provider activo + # al acercarse al límite de contexto, en vez de solo truncar + semantic_compression: bool = False + # Anthropic anthropic_api_key: str = "" anthropic_real_api_key: str = "" diff --git a/backend/app/core/utils.py b/backend/app/core/utils.py index 3c3724f..7427b2e 100644 --- a/backend/app/core/utils.py +++ b/backend/app/core/utils.py @@ -1,4 +1,18 @@ import re +from pathlib import Path + + +def tail_file(path: Path, lines: int) -> list[str]: + """Últimas N líneas leyendo solo el bloque final (64KB) del archivo.""" + try: + with open(path, "rb") as f: + f.seek(0, 2) + size = f.tell() + f.seek(max(0, size - 65536)) + data = f.read().decode("utf-8", errors="replace") + return data.splitlines()[-lines:] + except OSError: + return [] def sanitize_error(msg: str) -> str: diff --git a/backend/app/main.py b/backend/app/main.py index 581a2b9..b210514 100644 --- a/backend/app/main.py +++ b/backend/app/main.py @@ -65,8 +65,10 @@ async def _copilot_token_refresh_loop(): async def lifespan(app: FastAPI): from app.services import usage_tracker from app.services.llamacpp_service import autostart_if_configured + from app.services.telegram_bot import run_telegram_bot await usage_tracker.init_db() asyncio.create_task(autostart_if_configured()) + asyncio.create_task(run_telegram_bot()) task = asyncio.create_task(_copilot_token_refresh_loop()) try: yield diff --git a/backend/app/models/provider.py b/backend/app/models/provider.py index ecfcb46..d767171 100644 --- a/backend/app/models/provider.py +++ b/backend/app/models/provider.py @@ -57,3 +57,5 @@ class ProviderRegistry(BaseModel): providers: list[Provider] = [] routing_enabled: bool = False routing_rules: list[RoutingRule] = [] + # Failover: si el provider efectivo (local) no responde, probar estos en orden + fallback_provider_ids: list[str] = [] diff --git a/backend/app/services/compression_service.py b/backend/app/services/compression_service.py new file mode 100644 index 0000000..95c74f0 --- /dev/null +++ b/backend/app/services/compression_service.py @@ -0,0 +1,97 @@ +""" +Compresión semántica de conversaciones: resume la parte vieja del historial +con el provider activo en vez de solo truncar. Opt-in (SEMANTIC_COMPRESSION=true). +Cualquier fallo devuelve None y el caller cae al truncado clásico. +""" +import os + +import httpx + +from app.core.logging import get_logger +from app.models.provider import Provider + +log = get_logger(__name__) + +KEEP_RECENT_MESSAGES = 8 +_MAX_SOURCE_CHARS = 60000 + +_SUMMARY_PROMPT = ( + "Resume la siguiente conversación entre un usuario y un asistente de código. " + "Preserva: decisiones tomadas, archivos tocados y sus cambios, errores encontrados, " + "y el estado actual de la tarea. Sé compacto (máximo ~800 palabras). " + "Responde SOLO con el resumen.\n\n" +) + + +def message_to_text(message: dict) -> str: + content = message.get("content", "") + if isinstance(content, str): + return content + parts = [] + for block in content if isinstance(content, list) else []: + btype = block.get("type", "") + if btype == "text": + parts.append(block.get("text", "")) + elif btype == "tool_use": + parts.append(f"[tool_use: {block.get('name', '?')}]") + elif btype == "tool_result": + parts.append("[tool_result]") + return "\n".join(p for p in parts if p) + + +def _starts_with_tool_result(message: dict) -> bool: + content = message.get("content") + if not isinstance(content, list) or not content: + return False + return content[0].get("type") == "tool_result" + + +def split_for_compression(messages: list[dict], keep_recent: int = KEEP_RECENT_MESSAGES) -> tuple[list[dict], list[dict]]: + """(viejos_a_resumir, recientes_intactos) sin partir un par tool_use/tool_result.""" + if len(messages) <= keep_recent: + return [], messages + cut = len(messages) - keep_recent + while cut > 0 and _starts_with_tool_result(messages[cut]): + cut -= 1 + return messages[:cut], messages[cut:] + + +def build_summary_request(old_messages: list[dict], model: str) -> dict: + transcript = "\n".join(f"{m.get('role', '?')}: {message_to_text(m)}" for m in old_messages) + return { + "model": model, + "messages": [{"role": "user", "content": _SUMMARY_PROMPT + transcript[:_MAX_SOURCE_CHARS]}], + "max_tokens": 1500, + "stream": False, + } + + +async def compress_messages(messages: list[dict], provider: Provider, model: str) -> list[dict] | None: + old, recent = split_for_compression(messages) + if not old: + return None + + url = f"{provider.api_base.rstrip('/')}/chat/completions" + api_key = os.environ.get(provider.auth_env_var, "") if provider.auth_env_var else "" + headers = {"Authorization": f"Bearer {api_key or 'no-key'}", "Content-Type": "application/json"} + if provider.extra_headers: + headers.update(provider.extra_headers) + body = build_summary_request(old, provider.active_model or model) + + try: + async with httpx.AsyncClient(timeout=httpx.Timeout(connect=5.0, read=90.0, write=5.0, pool=5.0)) as client: + resp = await client.post(url, json=body, headers=headers) + resp.raise_for_status() + summary = resp.json()["choices"][0]["message"]["content"] + except Exception as e: + log.warning("semantic_compression_failed", provider=provider.id, error=str(e)[:200]) + return None + + if not summary or not str(summary).strip(): + return None + + summary_message = { + "role": "user", + "content": [{"type": "text", "text": f"[Resumen de la conversación previa]\n{summary}"}], + } + return [summary_message] + recent diff --git a/backend/app/services/llamacpp_service.py b/backend/app/services/llamacpp_service.py index e3a0feb..7bb4192 100644 --- a/backend/app/services/llamacpp_service.py +++ b/backend/app/services/llamacpp_service.py @@ -131,6 +131,12 @@ def build_cmdline(provider: Provider, devices: list[dict]) -> list[str]: if ratios: cmd += ["--tensor-split", ",".join(str(r) for r in ratios)] + rpc_servers = launch.get("rpc_servers") or [] + if isinstance(rpc_servers, list) and rpc_servers: + # llama.cpp RPC: workers remotos (ggml-rpc-server) se suman como + # devices — la granja multi-PC del backlog v3 + cmd += ["--rpc", ",".join(str(s) for s in rpc_servers)] + extra = launch.get("extra_args", []) if isinstance(extra, list): cmd += [str(a) for a in extra] @@ -258,22 +264,11 @@ def _read_pid() -> int | None: return None -def _tail_file(path: Path, lines: int) -> list[str]: - try: - with open(path, "rb") as f: - f.seek(0, 2) - size = f.tell() - f.seek(max(0, size - 65536)) - data = f.read().decode("utf-8", errors="replace") - return data.splitlines()[-lines:] - except OSError: - return [] - - def tail_logs(lines: int = 80) -> list[str]: + from app.core.utils import tail_file merged = [] for name, tag in (("llamacpp-out.log", "out"), ("llamacpp-err.log", "err")): - merged += [f"[{tag}] {line}" for line in _tail_file(_config_dir() / name, lines)] + merged += [f"[{tag}] {line}" for line in tail_file(_config_dir() / name, lines)] return merged[-lines:] diff --git a/backend/app/services/providers_service.py b/backend/app/services/providers_service.py index d867a7c..f05ec1e 100644 --- a/backend/app/services/providers_service.py +++ b/backend/app/services/providers_service.py @@ -205,13 +205,19 @@ def get_provider(provider_id: str) -> Optional[Provider]: return next((p for p in registry.providers if p.id == provider_id), None) -def set_routing(enabled: bool, rules: list) -> dict: +def set_routing(enabled: bool, rules: list, fallback_provider_ids: Optional[list[str]] = None) -> dict: with _registry_lock: registry = load_registry() registry.routing_enabled = enabled registry.routing_rules = rules + if fallback_provider_ids is not None: + registry.fallback_provider_ids = fallback_provider_ids save_registry(registry) - return {"enabled": enabled, "rules": rules} + return { + "enabled": enabled, + "rules": rules, + "fallback_provider_ids": registry.fallback_provider_ids, + } def resolve_route(model_name: str, prompt_tokens: int = 0) -> Optional[tuple[Provider, str]]: @@ -231,6 +237,60 @@ def resolve_route(model_name: str, prompt_tokens: int = 0) -> Optional[tuple[Pro return None +def _base_host_port(api_base: str) -> tuple[str, int]: + host_port = api_base.split("//")[-1].split("/")[0] + host, _, port = host_port.partition(":") + return host, int(port) if port.isdigit() else (443 if api_base.startswith("https") else 80) + + +def _is_local_base(api_base: str) -> bool: + host, _ = _base_host_port(api_base) + return host in ("127.0.0.1", "localhost") or host.startswith("192.168.") or host.startswith("10.") + + +async def _is_reachable(api_base: str, timeout: float = 0.4) -> bool: + host, port = _base_host_port(api_base) + try: + _, writer = await asyncio.wait_for(asyncio.open_connection(host, port), timeout) + writer.close() + await writer.wait_closed() + return True + except (OSError, asyncio.TimeoutError, ValueError): + return False + + +async def pick_provider(model_name: str, prompt_tokens: int = 0) -> tuple[Optional[Provider], Optional[str], bool]: + """Provider efectivo para un request: routing → failover si el destino local + no responde. Retorna (provider, model_override, es_el_provider_activo) — + el flag decide si un provider anthropic puede ir vía litellm (solo el activo: + litellm corre con SU config).""" + registry = load_registry() + route = resolve_route(model_name, prompt_tokens) + primary = route[0] if route else get_provider(registry.active_provider_id) + routed_model = route[1] if route else None + + if not registry.fallback_provider_ids or not primary: + return primary, routed_model, bool(primary and primary.id == registry.active_provider_id) + + candidates: list[Provider] = [primary] + for pid in registry.fallback_provider_ids: + if pid != primary.id: + fallback = next((p for p in registry.providers if p.id == pid), None) + if fallback: + candidates.append(fallback) + + for candidate in candidates: + # Solo se chequea reachability de bases locales/LAN; las cloud se asumen arriba + if not _is_local_base(candidate.api_base) or await _is_reachable(candidate.api_base): + is_active = candidate.id == registry.active_provider_id + if candidate.id != primary.id: + log.warning("provider_failover", primary=primary.id, fallback=candidate.id) + return candidate, None, is_active + return candidate, routed_model, is_active + + return primary, routed_model, primary.id == registry.active_provider_id + + def get_active_provider() -> Optional[Provider]: registry = load_registry() return get_provider(registry.active_provider_id) diff --git a/backend/app/services/settings_service.py b/backend/app/services/settings_service.py index e07d681..fae404d 100644 --- a/backend/app/services/settings_service.py +++ b/backend/app/services/settings_service.py @@ -48,3 +48,9 @@ def write_env_key(key: str, value: str) -> None: content = content.rstrip("\n") + f"\n{key}={value}\n" log.info("env_key_added", key=key) env_path.write_text(content, encoding="utf-8") + + # Aplicar en caliente: settings está lru_cache'd y pydantic no relee el .env + import os + from app.core.config import get_settings + os.environ[key] = value + get_settings.cache_clear() diff --git a/backend/app/services/telegram_bot.py b/backend/app/services/telegram_bot.py new file mode 100644 index 0000000..60610a0 --- /dev/null +++ b/backend/app/services/telegram_bot.py @@ -0,0 +1,230 @@ +import asyncio +import os + +import httpx + +from app.core.config import get_settings +from app.core.logging import get_logger + + +log = get_logger(__name__) + +_GATEWAY_URL = "http://127.0.0.1:8000/v1/chat/completions" +_POLL_TIMEOUT_SECONDS = 25 +_RETRY_DELAY_SECONDS = 5 + +_chat_histories: dict[int, list[dict]] = {} +_logged_disallowed_chat_ids: set[int] = set() +_empty_allowlist_warning_logged = False + + +def parse_allowed_chat_ids(raw: str) -> set[int]: + """Convierte una lista CSV en identificadores de chat válidos.""" + chat_ids = set() + for entry in raw.split(","): + try: + chat_ids.add(int(entry.strip())) + except ValueError: + continue + return chat_ids + + +def chunk_text(text: str, limit: int = 4096) -> list[str]: + """Divide texto conservando líneas completas siempre que sea posible.""" + if not text: + return [] + if limit <= 0: + raise ValueError("limit debe ser mayor que cero") + + chunks = [] + current = "" + for line in text.splitlines(keepends=True): + if len(line) > limit: + if current: + chunks.append(current) + current = "" + while len(line) > limit: + chunks.append(line[:limit]) + line = line[limit:] + + if len(current) + len(line) <= limit: + current += line + else: + if current: + chunks.append(current) + current = line + + if current: + chunks.append(current) + return chunks + + +def build_chat_body( + history: list[dict], + user_text: str, + model: str = "claude-sonnet-4-6", +) -> dict: + """Construye el cuerpo OpenAI para una conversación de Telegram.""" + return { + "model": model, + "messages": history + [{"role": "user", "content": user_text}], + "stream": False, + } + + +def trim_history(history: list[dict], max_turns: int = 10) -> list[dict]: + """Conserva los mensajes correspondientes a los turnos más recientes.""" + max_messages = max(0, max_turns * 2) + if len(history) <= max_messages: + return history + if max_messages == 0: + return [] + return history[-max_messages:] + + +def _log_empty_allowlist_once() -> None: + """Registra una sola vez que el bot no tiene chats autorizados.""" + global _empty_allowlist_warning_logged + if _empty_allowlist_warning_logged: + return + log.warning("telegram_bot_empty_allowlist") + _empty_allowlist_warning_logged = True + + +def _log_disallowed_chat_once(chat_id: int) -> None: + """Registra una sola vez cada chat ignorado por la lista de acceso.""" + if chat_id in _logged_disallowed_chat_ids: + return + log.debug("telegram_bot_chat_not_allowed", chat_id=chat_id) + _logged_disallowed_chat_ids.add(chat_id) + + +async def _get_updates( + client: httpx.AsyncClient, + telegram_url: str, + offset: int, +) -> list[dict]: + """Obtiene el siguiente lote de actualizaciones de Telegram.""" + response = await client.get( + f"{telegram_url}/getUpdates", + params={"timeout": _POLL_TIMEOUT_SECONDS, "offset": offset}, + ) + response.raise_for_status() + payload = response.json() + updates = payload.get("result", []) + if not payload.get("ok") or not isinstance(updates, list): + raise ValueError("Respuesta inválida de Telegram") + return updates + + +async def _request_completion( + client: httpx.AsyncClient, + history: list[dict], + user_text: str, +) -> str: + """Solicita una respuesta al gateway local de modelos.""" + response = await client.post( + _GATEWAY_URL, + headers={"x-api-key": get_settings().ui_api_key}, + json=build_chat_body(history, user_text), + ) + response.raise_for_status() + content = response.json()["choices"][0]["message"]["content"] + if not isinstance(content, str): + raise ValueError("Respuesta inválida del gateway local") + return content + + +async def _send_reply( + client: httpx.AsyncClient, + telegram_url: str, + chat_id: int, + text: str, +) -> None: + """Envía una respuesta de texto plano a Telegram.""" + for chunk in chunk_text(text): + response = await client.post( + f"{telegram_url}/sendMessage", + json={"chat_id": chat_id, "text": chunk}, + ) + response.raise_for_status() + + +async def _handle_update( + client: httpx.AsyncClient, + telegram_url: str, + update: dict, + allowed_chat_ids: set[int], +) -> None: + """Procesa un mensaje de texto autorizado y conserva su historial.""" + message = update.get("message") + if not isinstance(message, dict) or not isinstance(message.get("text"), str): + return + + chat = message.get("chat") + if not isinstance(chat, dict) or not isinstance(chat.get("id"), int): + return + + chat_id = chat["id"] + if chat_id not in allowed_chat_ids: + _log_disallowed_chat_once(chat_id) + return + + user_text = message["text"] + history = _chat_histories.get(chat_id, []) + assistant_text = await _request_completion(client, history, user_text) + _chat_histories[chat_id] = trim_history( + history + + [ + {"role": "user", "content": user_text}, + {"role": "assistant", "content": assistant_text}, + ] + ) + await _send_reply(client, telegram_url, chat_id, assistant_text) + + +async def run_telegram_bot() -> None: + """Ejecuta el bot opt-in de Telegram mediante long polling.""" + token = os.environ.get("TELEGRAM_BOT_TOKEN", "").strip() + if not token: + log.info("telegram_bot_disabled_no_token") + return + + allowed_chat_ids = parse_allowed_chat_ids( + os.environ.get("TELEGRAM_ALLOWED_CHAT_IDS", "") + ) + if not allowed_chat_ids: + _log_empty_allowlist_once() + return + + telegram_url = f"https://api.telegram.org/bot{token}" + timeout = httpx.Timeout(connect=10.0, read=300.0, write=30.0, pool=10.0) + offset = 0 + log.info("telegram_bot_started", allowed_chats=len(allowed_chat_ids)) + + while True: + try: + async with httpx.AsyncClient(timeout=timeout) as client: + while True: + updates = await _get_updates(client, telegram_url, offset) + for update in updates: + update_id = update.get("update_id") + if isinstance(update_id, int): + offset = max(offset, update_id + 1) + await _handle_update( + client, + telegram_url, + update, + allowed_chat_ids, + ) + except asyncio.CancelledError: + raise + except httpx.HTTPError as exc: + log.warning( + "telegram_bot_network_failed", + error=type(exc).__name__, + ) + await asyncio.sleep(_RETRY_DELAY_SECONDS) + except Exception as exc: + log.warning("telegram_bot_loop_failed", error=str(exc)) + await asyncio.sleep(_RETRY_DELAY_SECONDS) diff --git a/backend/tests/test_compression_service.py b/backend/tests/test_compression_service.py new file mode 100644 index 0000000..b90e770 --- /dev/null +++ b/backend/tests/test_compression_service.py @@ -0,0 +1,186 @@ +from unittest.mock import AsyncMock, Mock + +import httpx +import pytest + +from app.models.provider import Provider +from app.services import compression_service + + +def _provider() -> Provider: + return Provider( + id="test-provider", + name="Test Provider", + api_base="http://localhost:4000/v1", + active_model="provider-model", + ) + + +def test_message_to_text_passes_through_string_content(): + assert compression_service.message_to_text({"content": "plain text"}) == "plain text" + + +def test_message_to_text_joins_text_blocks(): + message = { + "content": [ + {"type": "text", "text": "first"}, + {"type": "text", "text": "second"}, + ] + } + + assert compression_service.message_to_text(message) == "first\nsecond" + + +def test_message_to_text_formats_tool_use_block(): + message = {"content": [{"type": "tool_use", "name": "read_file"}]} + + assert compression_service.message_to_text(message) == "[tool_use: read_file]" + + +def test_message_to_text_formats_tool_result_block(): + message = {"content": [{"type": "tool_result", "content": "ignored"}]} + + assert compression_service.message_to_text(message) == "[tool_result]" + + +def test_message_to_text_returns_empty_string_for_empty_list(): + assert compression_service.message_to_text({"content": []}) == "" + + +def test_split_for_compression_keeps_all_messages_when_at_or_below_limit(): + messages = [{"role": "user", "content": str(index)} for index in range(3)] + + old, recent = compression_service.split_for_compression(messages, keep_recent=3) + + assert old == [] + assert recent is messages + + +def test_split_for_compression_moves_cut_back_before_tool_pair(): + messages = [ + {"role": "user", "content": "old"}, + {"role": "assistant", "content": "older"}, + { + "role": "assistant", + "content": [{"type": "tool_use", "name": "read_file"}], + }, + { + "role": "user", + "content": [{"type": "tool_result", "content": "result"}], + }, + {"role": "assistant", "content": "latest"}, + ] + + old, recent = compression_service.split_for_compression(messages, keep_recent=2) + + assert old == messages[:2] + assert recent == messages[2:] + + +def test_split_for_compression_uses_normal_split_sizes(): + messages = [{"role": "user", "content": str(index)} for index in range(6)] + + old, recent = compression_service.split_for_compression(messages, keep_recent=2) + + assert old == messages[:4] + assert recent == messages[4:] + + +def test_build_summary_request_sets_options_and_includes_roles(): + request = compression_service.build_summary_request( + [ + {"role": "user", "content": "question"}, + {"role": "assistant", "content": "answer"}, + ], + "summary-model", + ) + + assert request["model"] == "summary-model" + assert request["stream"] is False + assert request["max_tokens"] == 1500 + assert "user: question\nassistant: answer" in request["messages"][0]["content"] + + +def test_build_summary_request_caps_source_at_60000_characters(): + request = compression_service.build_summary_request( + [{"role": "user", "content": "x" * 70_000}], + "summary-model", + ) + + content = request["messages"][0]["content"] + assert len(content) <= len(compression_service._SUMMARY_PROMPT) + 60_000 + + +@pytest.mark.asyncio +async def test_compress_messages_returns_summary_followed_by_recent_messages(monkeypatch): + messages = [ + {"role": "user", "content": f"message {index}"} + for index in range(10) + ] + response = Mock() + response.json.return_value = { + "choices": [{"message": {"content": "resumen"}}] + } + client = AsyncMock() + client.post.return_value = response + async_client = Mock() + async_client.return_value.__aenter__ = AsyncMock(return_value=client) + async_client.return_value.__aexit__ = AsyncMock(return_value=False) + monkeypatch.setattr(compression_service.httpx, "AsyncClient", async_client) + + result = await compression_service.compress_messages( + messages, + _provider(), + "requested-model", + ) + + assert result is not None + assert result[1:] == messages[-compression_service.KEEP_RECENT_MESSAGES :] + assert result[0]["role"] == "user" + assert result[0]["content"][0]["text"].startswith( + "[Resumen de la conversación previa]" + ) + assert result[0]["content"][0]["text"].endswith("resumen") + + +@pytest.mark.asyncio +async def test_compress_messages_returns_none_without_calling_http_when_short(monkeypatch): + async_client = Mock() + monkeypatch.setattr(compression_service.httpx, "AsyncClient", async_client) + messages = [ + {"role": "user", "content": f"message {index}"} + for index in range(compression_service.KEEP_RECENT_MESSAGES) + ] + + result = await compression_service.compress_messages( + messages, + _provider(), + "requested-model", + ) + + assert result is None + async_client.assert_not_called() + + +@pytest.mark.asyncio +async def test_compress_messages_returns_none_on_http_error(monkeypatch): + messages = [ + {"role": "user", "content": f"message {index}"} + for index in range(compression_service.KEEP_RECENT_MESSAGES + 1) + ] + response = Mock() + response.raise_for_status.side_effect = httpx.HTTPError("request failed") + client = AsyncMock() + client.post.return_value = response + async_client = Mock() + async_client.return_value.__aenter__ = AsyncMock(return_value=client) + async_client.return_value.__aexit__ = AsyncMock(return_value=False) + monkeypatch.setattr(compression_service.httpx, "AsyncClient", async_client) + + result = await compression_service.compress_messages( + messages, + _provider(), + "requested-model", + ) + + assert result is None diff --git a/backend/tests/test_failover.py b/backend/tests/test_failover.py new file mode 100644 index 0000000..4e1fa7e --- /dev/null +++ b/backend/tests/test_failover.py @@ -0,0 +1,200 @@ +from unittest.mock import AsyncMock + +import pytest + +from app.models.provider import Provider, ProviderRegistry, RoutingRule +from app.services import providers_service + + +def _provider(provider_id: str, api_base: str, active_model: str = "active-model") -> Provider: + return Provider( + id=provider_id, + name=provider_id, + api_base=api_base, + active_model=active_model, + ) + + +def _patch_registry(monkeypatch, registry: ProviderRegistry) -> None: + monkeypatch.setattr(providers_service, "load_registry", lambda: registry) + + +@pytest.mark.asyncio +async def test_pick_provider_without_fallback_returns_active_without_reachability_check( + monkeypatch, +): + active = _provider("active", "http://127.0.0.1:4000/v1") + registry = ProviderRegistry( + active_provider_id=active.id, + providers=[active], + fallback_provider_ids=[], + ) + _patch_registry(monkeypatch, registry) + is_reachable = AsyncMock() + monkeypatch.setattr(providers_service, "_is_reachable", is_reachable) + + provider, model_override, is_active = await providers_service.pick_provider( + "requested-model" + ) + + assert provider is active + assert model_override is None + assert is_active is True + is_reachable.assert_not_awaited() + + +@pytest.mark.asyncio +async def test_pick_provider_fails_over_without_leaking_routed_model(monkeypatch): + primary = _provider("primary", "http://127.0.0.1:4000/v1") + fallback = _provider("fallback", "http://localhost:4001/v1") + registry = ProviderRegistry( + active_provider_id=primary.id, + providers=[primary, fallback], + routing_enabled=True, + routing_rules=[ + RoutingRule( + pattern="sonnet", + provider_id=primary.id, + model="routed-primary-model", + ) + ], + fallback_provider_ids=[fallback.id], + ) + _patch_registry(monkeypatch, registry) + is_reachable = AsyncMock( + side_effect=lambda api_base: api_base == fallback.api_base + ) + monkeypatch.setattr(providers_service, "_is_reachable", is_reachable) + + provider, model_override, is_active = await providers_service.pick_provider( + "claude-sonnet" + ) + + assert provider is fallback + assert model_override is None + assert is_active is False + assert [call.args[0] for call in is_reachable.await_args_list] == [ + primary.api_base, + fallback.api_base, + ] + + +@pytest.mark.asyncio +async def test_pick_provider_returns_reachable_routed_primary_with_routed_model( + monkeypatch, +): + primary = _provider("primary", "http://127.0.0.1:4000/v1") + fallback = _provider("fallback", "http://localhost:4001/v1") + registry = ProviderRegistry( + active_provider_id=primary.id, + providers=[primary, fallback], + routing_enabled=True, + routing_rules=[ + RoutingRule( + pattern="haiku", + provider_id=primary.id, + model="routed-primary-model", + ) + ], + fallback_provider_ids=[fallback.id], + ) + _patch_registry(monkeypatch, registry) + is_reachable = AsyncMock(return_value=True) + monkeypatch.setattr(providers_service, "_is_reachable", is_reachable) + + provider, model_override, is_active = await providers_service.pick_provider( + "claude-haiku" + ) + + assert provider is primary + assert model_override == "routed-primary-model" + assert is_active is True + is_reachable.assert_awaited_once_with(primary.api_base) + + +@pytest.mark.asyncio +async def test_pick_provider_assumes_non_local_primary_is_reachable(monkeypatch): + primary = _provider("cloud", "https://api.example.com/v1") + fallback = _provider("fallback", "http://localhost:4001/v1") + registry = ProviderRegistry( + active_provider_id=primary.id, + providers=[primary, fallback], + fallback_provider_ids=[fallback.id], + ) + _patch_registry(monkeypatch, registry) + is_reachable = AsyncMock() + monkeypatch.setattr(providers_service, "_is_reachable", is_reachable) + + provider, model_override, is_active = await providers_service.pick_provider( + "requested-model" + ) + + assert provider is primary + assert model_override is None + assert is_active is True + is_reachable.assert_not_awaited() + + +@pytest.mark.asyncio +async def test_pick_provider_returns_primary_when_all_candidates_are_unreachable( + monkeypatch, +): + primary = _provider("primary", "http://127.0.0.1:4000/v1") + fallback = _provider("fallback", "http://localhost:4001/v1") + registry = ProviderRegistry( + active_provider_id=primary.id, + providers=[primary, fallback], + fallback_provider_ids=[fallback.id], + ) + _patch_registry(monkeypatch, registry) + is_reachable = AsyncMock(return_value=False) + monkeypatch.setattr(providers_service, "_is_reachable", is_reachable) + + provider, model_override, is_active = await providers_service.pick_provider( + "requested-model" + ) + + assert provider is primary + assert provider is not None + assert model_override is None + assert is_active is True + assert is_reachable.await_count == 2 + + +@pytest.mark.parametrize( + ("active_provider_id", "expected_is_active"), + [ + ("fallback", True), + ("other", False), + ], +) +@pytest.mark.asyncio +async def test_pick_provider_sets_is_active_from_returned_candidate_id( + monkeypatch, + active_provider_id, + expected_is_active, +): + primary = _provider("routed-primary", "http://127.0.0.1:4000/v1") + fallback = _provider("fallback", "https://fallback.example.com/v1") + other = _provider("other", "https://other.example.com/v1") + registry = ProviderRegistry( + active_provider_id=active_provider_id, + providers=[primary, fallback, other], + routing_enabled=True, + routing_rules=[ + RoutingRule(pattern="sonnet", provider_id=primary.id) + ], + fallback_provider_ids=[fallback.id], + ) + _patch_registry(monkeypatch, registry) + is_reachable = AsyncMock(return_value=False) + monkeypatch.setattr(providers_service, "_is_reachable", is_reachable) + + provider, model_override, is_active = await providers_service.pick_provider( + "claude-sonnet" + ) + + assert provider is fallback + assert model_override is None + assert is_active is expected_is_active + is_reachable.assert_awaited_once_with(primary.api_base) diff --git a/backend/tests/test_messages_native.py b/backend/tests/test_messages_native.py index 2657560..0091ce5 100644 --- a/backend/tests/test_messages_native.py +++ b/backend/tests/test_messages_native.py @@ -53,8 +53,8 @@ def _post_messages(client, provider): 'data: {"type": "message_stop"}', ] with patch( - "app.api.messages.providers_service.get_active_provider", - return_value=provider, + "app.api.messages.providers_service.pick_provider", + new=AsyncMock(return_value=(provider, None, False)), ), patch("app.api.messages.httpx.AsyncClient") as mock_client: instance = mock_client.return_value instance.__aenter__ = AsyncMock(return_value=instance) @@ -113,11 +113,8 @@ def test_routing_overrides_active_provider(client): } lines = ['data: {"type": "message_stop"}'] with patch( - "app.api.messages.providers_service.get_active_provider", - return_value=_native_provider(), - ), patch( - "app.api.messages.providers_service.resolve_route", - return_value=(routed, "qwen3-4b"), + "app.api.messages.providers_service.pick_provider", + new=AsyncMock(return_value=(routed, "qwen3-4b", False)), ), patch("app.api.messages.httpx.AsyncClient") as mock_client: instance = mock_client.return_value instance.__aenter__ = AsyncMock(return_value=instance) diff --git a/backend/tests/test_openai_compat.py b/backend/tests/test_openai_compat.py index d539120..509fc26 100644 --- a/backend/tests/test_openai_compat.py +++ b/backend/tests/test_openai_compat.py @@ -62,8 +62,8 @@ def test_chat_completions_rewrites_model_and_forwards(client): upstream.status_code = 200 upstream.json = lambda: {"choices": [], "usage": {"prompt_tokens": 3, "completion_tokens": 5}} with patch( - "app.api.openai_compat.providers_service.get_active_provider", - return_value=_provider(), + "app.api.openai_compat.providers_service.pick_provider", + new=AsyncMock(return_value=(_provider(), None, True)), ), patch("app.api.openai_compat.httpx.AsyncClient") as mock_client: instance = mock_client.return_value instance.__aenter__ = AsyncMock(return_value=instance) @@ -83,8 +83,8 @@ def test_chat_completions_upstream_error_relayed(client): upstream.status_code = 400 upstream.json = lambda: {"error": {"message": "bad"}} with patch( - "app.api.openai_compat.providers_service.get_active_provider", - return_value=_provider(), + "app.api.openai_compat.providers_service.pick_provider", + new=AsyncMock(return_value=(_provider(), None, True)), ), patch("app.api.openai_compat.httpx.AsyncClient") as mock_client: instance = mock_client.return_value instance.__aenter__ = AsyncMock(return_value=instance) diff --git a/backend/tests/test_routing.py b/backend/tests/test_routing.py index ec59b70..a085edd 100644 --- a/backend/tests/test_routing.py +++ b/backend/tests/test_routing.py @@ -182,6 +182,7 @@ def test_get_routing_returns_enabled_and_rules(client, monkeypatch): assert response.json() == { "enabled": True, "rules": [rule.model_dump()], + "fallback_provider_ids": [], } mock_load_registry.assert_called_once_with() mock_set_routing.assert_not_called() @@ -245,6 +246,6 @@ def test_put_routing_calls_set_routing_and_returns_result(client, monkeypatch): assert response.json() == service_result mock_get_provider.assert_called_once_with(target.id) mock_set_routing.assert_called_once_with( - True, [RoutingRule(**rule_payload)] + True, [RoutingRule(**rule_payload)], None ) mock_load_registry.assert_not_called() diff --git a/backend/tests/test_telegram_bot.py b/backend/tests/test_telegram_bot.py new file mode 100644 index 0000000..a92ad16 --- /dev/null +++ b/backend/tests/test_telegram_bot.py @@ -0,0 +1,79 @@ +from app.services import telegram_bot + + +def test_parse_allowed_chat_ids_parses_valid_csv(): + assert telegram_bot.parse_allowed_chat_ids("123,456,-789") == {123, 456, -789} + + +def test_parse_allowed_chat_ids_tolerates_spaces_and_empty_entries(): + assert telegram_bot.parse_allowed_chat_ids(" 123, , 456 ,") == {123, 456} + + +def test_parse_allowed_chat_ids_skips_garbage_entries(): + assert telegram_bot.parse_allowed_chat_ids("123,nope,4.5,456") == {123, 456} + + +def test_parse_allowed_chat_ids_empty_string_returns_empty_set(): + assert telegram_bot.parse_allowed_chat_ids("") == set() + + +def test_chunk_text_short_text_returns_single_chunk(): + assert telegram_bot.chunk_text("respuesta breve") == ["respuesta breve"] + + +def test_chunk_text_preserves_whole_lines_when_possible(): + text = "uno\ndos\ntres" + + assert telegram_bot.chunk_text(text, limit=8) == ["uno\ndos\n", "tres"] + + +def test_chunk_text_hard_splits_a_line_longer_than_limit(): + assert telegram_bot.chunk_text("abcdefghij", limit=4) == ["abcd", "efgh", "ij"] + + +def test_chunk_text_empty_text_returns_empty_list(): + assert telegram_bot.chunk_text("") == [] + + +def test_build_chat_body_uses_default_model(): + body = telegram_bot.build_chat_body([], "hola") + + assert body["model"] == "claude-sonnet-4-6" + + +def test_build_chat_body_prepends_history_before_new_user_message(): + history = [ + {"role": "user", "content": "pregunta anterior"}, + {"role": "assistant", "content": "respuesta anterior"}, + ] + + body = telegram_bot.build_chat_body(history, "pregunta nueva") + + assert body["messages"] == history + [ + {"role": "user", "content": "pregunta nueva"} + ] + + +def test_build_chat_body_disables_streaming(): + body = telegram_bot.build_chat_body([], "hola") + + assert body["stream"] is False + + +def test_trim_history_under_limit_returns_history_untouched(): + history = [ + {"role": "user", "content": "hola"}, + {"role": "assistant", "content": "hola"}, + ] + + result = telegram_bot.trim_history(history, max_turns=2) + + assert result is history + + +def test_trim_history_over_limit_keeps_most_recent_turns(): + history = [{"role": "user", "content": str(index)} for index in range(8)] + + result = telegram_bot.trim_history(history, max_turns=2) + + assert result == history[-4:] diff --git a/docker-compose.yml b/docker-compose.yml new file mode 100644 index 0000000..3aaf2bb --- /dev/null +++ b/docker-compose.yml @@ -0,0 +1,13 @@ +services: + bipolar-code: + build: . + ports: + - "8000:8000" + volumes: + - bipolar-data:/data + extra_hosts: + - "host.docker.internal:host-gateway" + restart: unless-stopped + +volumes: + bipolar-data: diff --git a/frontend/package-lock.json b/frontend/package-lock.json index 6a15779..4c6a8c3 100644 --- a/frontend/package-lock.json +++ b/frontend/package-lock.json @@ -1,12 +1,12 @@ { "name": "bipolar-code-frontend", - "version": "2.11.0", + "version": "2.12.0", "lockfileVersion": 3, "requires": true, "packages": { "": { "name": "bipolar-code-frontend", - "version": "2.11.0", + "version": "2.12.0", "dependencies": { "@tanstack/react-query": "^5.40.0", "axios": "^1.7.2", diff --git a/frontend/package.json b/frontend/package.json index 0fd1f0e..ba7901e 100644 --- a/frontend/package.json +++ b/frontend/package.json @@ -1 +1 @@ -{"name":"bipolar-code-frontend","private":true,"version":"2.11.0","type":"module","scripts":{"dev":"vite","build":"tsc -b && vite build","preview":"vite preview","test":"vitest run"},"dependencies":{"@tanstack/react-query":"^5.40.0","axios":"^1.7.2","react":"^18.3.1","react-dom":"^18.3.1","react-router-dom":"^6.23.1","recharts":"^3.8.1"},"devDependencies":{"@testing-library/jest-dom":"^6.4.6","@testing-library/react":"^16.0.0","@types/react":"^18.3.3","@types/react-dom":"^18.3.0","@vitejs/plugin-react":"^4.3.1","autoprefixer":"^10.4.19","jsdom":"^29.0.2","postcss":"^8.4.39","tailwindcss":"^3.4.4","typescript":"^5.4.5","vite":"^5.3.1","vitest":"^1.6.0"}} \ No newline at end of file +{"name":"bipolar-code-frontend","private":true,"version":"2.12.0","type":"module","scripts":{"dev":"vite","build":"tsc -b && vite build","preview":"vite preview","test":"vitest run"},"dependencies":{"@tanstack/react-query":"^5.40.0","axios":"^1.7.2","react":"^18.3.1","react-dom":"^18.3.1","react-router-dom":"^6.23.1","recharts":"^3.8.1"},"devDependencies":{"@testing-library/jest-dom":"^6.4.6","@testing-library/react":"^16.0.0","@types/react":"^18.3.3","@types/react-dom":"^18.3.0","@vitejs/plugin-react":"^4.3.1","autoprefixer":"^10.4.19","jsdom":"^29.0.2","postcss":"^8.4.39","tailwindcss":"^3.4.4","typescript":"^5.4.5","vite":"^5.3.1","vitest":"^1.6.0"}} \ No newline at end of file diff --git a/frontend/src/components/RoutingPanel.tsx b/frontend/src/components/RoutingPanel.tsx index fa6c75e..250fbbc 100644 --- a/frontend/src/components/RoutingPanel.tsx +++ b/frontend/src/components/RoutingPanel.tsx @@ -16,17 +16,23 @@ export function RoutingPanel({ providers }: RoutingPanelProps) { const { data } = useQuery({ queryKey: ['routing'], queryFn: routingApi.get }) const [enabled, setEnabled] = useState(false) const [rules, setRules] = useState([]) + const [fallbackCsv, setFallbackCsv] = useState('') const [dirty, setDirty] = useState(false) useEffect(() => { if (data && !dirty) { setEnabled(data.enabled) setRules(data.rules) + setFallbackCsv((data.fallback_provider_ids ?? []).join(', ')) } }, [data, dirty]) const save = useMutation({ - mutationFn: () => routingApi.set({ enabled, rules: rules.filter(r => r.provider_id) }), + mutationFn: () => routingApi.set({ + enabled, + rules: rules.filter(r => r.provider_id), + fallback_provider_ids: fallbackCsv.split(',').map(s => s.trim()).filter(Boolean), + }), onSuccess: () => { setDirty(false) qc.invalidateQueries({ queryKey: ['routing'] }) @@ -99,6 +105,18 @@ export function RoutingPanel({ providers }: RoutingPanelProps) { ))} +
+ +
+
+ + {/* Compresión semántica */} +
+ +
+ + {/* Logs de litellm */} +
+ + {showLitellmLogs && ( +
+              {(litellmLogs?.logs ?? []).join('\n') || 'Sin logs todavía'}
+            
+ )} +
) diff --git a/frontend/src/services/api.ts b/frontend/src/services/api.ts index ddc4538..245fa61 100644 --- a/frontend/src/services/api.ts +++ b/frontend/src/services/api.ts @@ -53,6 +53,7 @@ export const proxyApi = { start: () => api.post<{ started: boolean; provider: string }>('/proxy/start').then(r => r.data), getRoute: () => api.get<{ mode: string; litellm_running: boolean; proxy_status: any }>('/proxy/route').then(r => r.data), setRoute: (mode: 'direct' | 'proxy') => api.post('/proxy/route', { mode }).then(r => r.data), + getLogs: (lines = 80) => api.get<{ logs: string[] }>('/proxy/logs', { params: { lines } }).then(r => r.data), } export const providersApi = { diff --git a/frontend/src/types/provider.ts b/frontend/src/types/provider.ts index 725bcf7..bdc5013 100644 --- a/frontend/src/types/provider.ts +++ b/frontend/src/types/provider.ts @@ -95,6 +95,7 @@ export interface RoutingRule { export interface RoutingConfig { enabled: boolean rules: RoutingRule[] + fallback_provider_ids?: string[] } export interface ProviderRegistry {