Files
mission-control-v2/backend/routers/gateway_proxy.py
T
Hitonabi 341ea870bb Refactor: Zentrales Logging + robuste Token-Erfassung (Phase 2)
- app.py: logging.basicConfig (MC_LOG_LEVEL, INFO default) als eine
  Konfiguration für alle Module.
- Neuer services/gateway_stream.py: SSE-/Non-Stream-usage-Parsing aus dem
  gateway_proxy-Router extrahiert; robuster Zeilenparser mit Debug-Logging
  statt verschluckter Exceptions. Router ist jetzt dünn.
- token_stats.py: In-Memory-Cache + gedrosseltes Flushen (5s) + atexit-Flush
  statt Write-pro-Request; atomarer Write (.tmp -> replace); thread-safe.
- agent.py/discover.py: stille `except Exception: pass` durch gezieltes
  log.debug/warning ersetzt; ungenutzten yaml-Import entfernt.

Verifiziert: Stream-Parsing (Summen + per-Modell), malformed-Chunk übersteht,
flush schreibt; app importiert sauber.

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
2026-06-26 14:32:15 +02:00

58 lines
2.1 KiB
Python

import httpx
from fastapi import APIRouter, Request
from fastapi.responses import JSONResponse, StreamingResponse
from config import LLAMA_SWAP_URL
from services.gateway_stream import record_stream_chunk, record_usage
from services.router_logic import FAST, FAST_NO_THINK, choose_model
router = APIRouter(prefix="/v1")
@router.get("/models")
async def models():
async with httpx.AsyncClient(timeout=10) as c:
r = await c.get(f"{LLAMA_SWAP_URL}/v1/models")
return JSONResponse(r.json(), status_code=r.status_code)
async def _proxy(path: str, request: Request):
body = await request.json()
requested = str(body.get("model") or "auto")
if requested == "auto":
alias, reason = choose_model(body)
body["model"] = alias
routed = {"x-mc-routed-to": alias, "x-mc-route-reason": reason}
else:
alias = requested
routed = {"x-mc-routed-to": requested}
# fast-Spur: Thinking aus für flotte Antworten (sofern Client es nicht selbst setzt).
if FAST_NO_THINK and alias == FAST and "chat_template_kwargs" not in body:
body["chat_template_kwargs"] = {"enable_thinking": False}
url = f"{LLAMA_SWAP_URL}{path}"
if body.get("stream"):
async def gen():
async with httpx.AsyncClient(timeout=None) as c:
async with c.stream("POST", url, json=body) as r:
async for chunk in r.aiter_raw():
record_stream_chunk(chunk, alias)
yield chunk
return StreamingResponse(gen(), media_type="text/event-stream", headers=routed)
async with httpx.AsyncClient(timeout=600) as c:
r = await c.post(url, json=body)
resp_json = r.json()
record_usage(resp_json.get("usage") if isinstance(resp_json, dict) else None, alias)
return JSONResponse(resp_json, status_code=r.status_code, headers=routed)
@router.post("/chat/completions")
async def chat_completions(request: Request):
return await _proxy("/v1/chat/completions", request)
@router.post("/completions")
async def completions(request: Request):
return await _proxy("/v1/completions", request)