134 lines
6.0 KiB
Python
134 lines
6.0 KiB
Python
import os
|
|
|
|
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 LANES, choose_for_lane
|
|
from services.routing_policy import load_policy
|
|
|
|
router = APIRouter(prefix="/v1")
|
|
|
|
# Antwortsprache für IDE-/Lane-Traffic: die Coding-Modelle antworten sonst englisch
|
|
# (User-Anforderung 03.07.2026). Leerer String (MC_GATEWAY_LANG_DIRECTIVE="") schaltet ab.
|
|
_LANG_DIRECTIVE = os.environ.get(
|
|
"MC_GATEWAY_LANG_DIRECTIVE",
|
|
"Antworte dem Nutzer grundsätzlich auf Deutsch (Erklärungen, Pläne, Rückfragen, "
|
|
"Zusammenfassungen) — auch wenn die Frage oder Tool-Anweisungen englisch sind. "
|
|
"Quellcode, Bezeichner und Shell-Befehle bleiben unverändert.")
|
|
|
|
|
|
def _inject_language(body: dict, alias: str) -> None:
|
|
"""Deutsch-Direktive anhängen. An die ERSTE System-Message (viele Chat-Templates
|
|
erwarten nur eine), sonst als neue System-Message. `hermes` ausgenommen — Lucys
|
|
Persona (SOUL.md) regelt die Sprache selbst."""
|
|
if not _LANG_DIRECTIVE or alias == "hermes":
|
|
return
|
|
msgs = body.get("messages")
|
|
if not isinstance(msgs, list):
|
|
return
|
|
first_sys = next((m for m in msgs if isinstance(m, dict) and m.get("role") == "system"), None)
|
|
if first_sys is None:
|
|
msgs.insert(0, {"role": "system", "content": _LANG_DIRECTIVE})
|
|
elif isinstance(first_sys.get("content"), str):
|
|
first_sys["content"] = first_sys["content"].rstrip() + "\n\n" + _LANG_DIRECTIVE
|
|
elif isinstance(first_sys.get("content"), list):
|
|
first_sys["content"].append({"type": "text", "text": _LANG_DIRECTIVE})
|
|
|
|
# Virtuelle Lanes, die der Gateway zusätzlich zu den echten Modellen als „Modell" anbietet.
|
|
_LANE_LABELS = {"coding": "Coding (Router → coder/heavy/fast)", "chat": "Chat (Router → fast/heavy)"}
|
|
|
|
|
|
@router.get("/models")
|
|
async def models(request: Request):
|
|
client = request.app.state.gw_client # geteilter Keep-Alive-Client (siehe app.py lifespan)
|
|
r = await client.get(f"{LLAMA_SWAP_URL}/v1/models", timeout=10.0)
|
|
data = r.json()
|
|
# Lanes ganz oben einblenden, damit IDEs einfach „coding"/„chat" wählen können.
|
|
lanes = [{"id": lane, "object": "model", "owned_by": "mc2-router",
|
|
"description": _LANE_LABELS.get(lane, lane)} for lane in LANES]
|
|
# Kontextlänge je Modell mitliefern (aus der llama-swap-Config geparst). Ohne sie
|
|
# budgetieren Clients blind — Hermes-Subagents nahmen 256k an, schickten passende
|
|
# max_tokens und rissen damit den echten Server-Kontext (Radar-Lauf 02.07.).
|
|
# Rollen-Aliase (heavy/coder/hermes …) tauchen bei llama-swap NICHT als Einträge auf,
|
|
# Clients fragen aber genau damit an → als eigene Einträge einblenden.
|
|
ctx_map: dict[str, int] = {}
|
|
alias_entries: list[dict] = []
|
|
try:
|
|
from services import llamaswap
|
|
for m in llamaswap.list_models():
|
|
ctx = m.get("ctx")
|
|
if ctx:
|
|
for api_id in m.get("api_ids", []):
|
|
ctx_map[api_id] = ctx
|
|
for alias in m.get("aliases", []):
|
|
entry = {"id": alias, "object": "model", "owned_by": "mc2-alias",
|
|
"description": f"Alias für {m['name']}"}
|
|
if ctx:
|
|
entry["context_length"] = ctx
|
|
alias_entries.append(entry)
|
|
except Exception:
|
|
pass
|
|
if isinstance(data, dict) and isinstance(data.get("data"), list):
|
|
for entry in data["data"]:
|
|
if (ctx := ctx_map.get(entry.get("id"))):
|
|
entry.setdefault("context_length", ctx)
|
|
data["data"] = lanes + alias_entries + data["data"]
|
|
return JSONResponse(data, 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.lower() in ("auto", "chat", "coding"):
|
|
lane = requested.lower()
|
|
alias, reason = choose_for_lane(lane, body)
|
|
body["model"] = alias
|
|
routed = {"x-mc-routed-to": alias, "x-mc-route-reason": reason, "x-mc-lane": lane}
|
|
else:
|
|
alias = requested
|
|
routed = {"x-mc-routed-to": requested}
|
|
# fast-Spur: Thinking aus für flotte Antworten (sofern Client es nicht selbst setzt).
|
|
pol = load_policy()
|
|
if pol["fast_no_think"] and alias == pol["fast"] and "chat_template_kwargs" not in body:
|
|
body["chat_template_kwargs"] = {"enable_thinking": False}
|
|
_inject_language(body, alias)
|
|
url = f"{LLAMA_SWAP_URL}{path}"
|
|
|
|
client = request.app.state.gw_client # geteilter Keep-Alive-Client (siehe app.py lifespan)
|
|
if body.get("stream"):
|
|
req = client.build_request("POST", url, json=body, timeout=None)
|
|
r = await client.send(req, stream=True)
|
|
if r.status_code != 200:
|
|
await r.aread()
|
|
try:
|
|
resp_json = r.json()
|
|
except Exception:
|
|
resp_json = {"error": {"message": r.text, "type": "upstream_error"}}
|
|
return JSONResponse(resp_json, status_code=r.status_code, headers=routed)
|
|
|
|
async def gen():
|
|
try:
|
|
async for chunk in r.aiter_raw():
|
|
record_stream_chunk(chunk, alias)
|
|
yield chunk
|
|
finally:
|
|
await r.aclose()
|
|
return StreamingResponse(gen(), media_type="text/event-stream", headers=routed)
|
|
|
|
r = await client.post(url, json=body, timeout=600.0)
|
|
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)
|