""" RyzGateway - BYOK OpenAI-Compatible AI Gateway Single-file FastAPI app for HuggingFace Spaces deployment. Model routing: - provider/model → direct passthrough to that provider - generic-name → routed pool (round-robin, skip rate-limited) Rate limit handling: - Auto-detects 429s and marks endpoints as cooled down - Respects RPM / RPD / RPS limits configured per provider - Smart selection skips exhausted endpoints Features: - Glassmorphism UI (highly optimized for mobile touch) - Auto-fetches and imports models directly from provider's /v1/models """ import os import re import time import uuid import json import asyncio import logging from pathlib import Path from typing import Any, Dict, List, Optional, Literal from contextlib import asynccontextmanager from collections import defaultdict, deque import httpx from fastapi import FastAPI, HTTPException, Header, Request, Depends from fastapi.responses import HTMLResponse, JSONResponse, StreamingResponse from pydantic import BaseModel, Field # ── Logging ──────────────────────────────────────────────────────────────────── logging.basicConfig(level=logging.INFO, format="%(asctime)s [%(levelname)s] %(message)s") log = logging.getLogger("ryzgateway") # ── Persistence (JSON file, reloaded on each read) ───────────────────────────── DATA_FILE = Path(os.getenv("DATA_FILE", "/data/gateway.json")) def _default_data() -> dict: return { "api_keys": {}, # key_id → {key, name, created_at} "providers": {}, # provider_id → {name, base_url, api_key, rate_type, rate_limit, notes} "pool_models": {}, # pool_id → {name, members: [{provider_id, model_name, weight}], strategy} } def load_data() -> dict: try: if DATA_FILE.exists(): return json.loads(DATA_FILE.read_text()) except Exception as e: log.error(f"Failed to load data: {e}") return _default_data() def save_data(data: dict): try: DATA_FILE.parent.mkdir(parents=True, exist_ok=True) DATA_FILE.write_text(json.dumps(data, indent=2)) except Exception as e: log.error(f"Failed to save data: {e}") # ── In-memory rate limit state ───────────────────────────────────────────────── # keyed by (provider_id, model_name) class RateLimitState: def __init__(self): self._lock = asyncio.Lock() self._cooldown: Dict[str, float] = {} self._timestamps: Dict[str, deque] = defaultdict(lambda: deque(maxlen=10000)) self._rr_index: Dict[str, int] = defaultdict(int) def key(self, provider_id: str, model_name: str) -> str: return f"{provider_id}:{model_name}" async def mark_rate_limited(self, provider_id: str, model_name: str, cooldown_secs: float = 60.0): async with self._lock: k = self.key(provider_id, model_name) self._cooldown[k] = time.time() + cooldown_secs log.warning(f"Rate limited: {k} — cooling for {cooldown_secs}s") async def is_available(self, provider_id: str, model_name: str, provider_cfg: dict) -> bool: async with self._lock: k = self.key(provider_id, model_name) now = time.time() if self._cooldown.get(k, 0) > now: return False rate_type = provider_cfg.get("rate_type", "none") rate_limit = provider_cfg.get("rate_limit", 0) if rate_type == "none" or not rate_limit: return True window = {"rpm": 60, "rps": 1, "rpd": 86400}.get(rate_type, 60) cutoff = now - window dq = self._timestamps[k] while dq and dq[0] < cutoff: dq.popleft() return len(dq) < rate_limit async def record_request(self, provider_id: str, model_name: str): async with self._lock: k = self.key(provider_id, model_name) self._timestamps[k].append(time.time()) def get_cooldown_remaining(self, provider_id: str, model_name: str) -> float: k = self.key(provider_id, model_name) remaining = self._cooldown.get(k, 0) - time.time() return max(0.0, remaining) async def next_rr(self, pool_id: str, count: int) -> int: async with self._lock: idx = self._rr_index[pool_id] % count self._rr_index[pool_id] = (idx + 1) % count return idx rl_state = RateLimitState() # ── Auth ─────────────────────────────────────────────────────────────────────── ADMIN_KEY = os.getenv("ADMIN_KEY", "admin-changeme") async def require_admin(authorization: Optional[str] = Header(None)): if not authorization or not authorization.startswith("Bearer "): raise HTTPException(status_code=401, detail="Missing Authorization header") token = authorization[7:] if token != ADMIN_KEY: raise HTTPException(status_code=403, detail="Invalid admin key") return token async def require_api_key(authorization: Optional[str] = Header(None)): if not authorization or not authorization.startswith("Bearer "): raise HTTPException(status_code=401, detail="Missing Authorization header") token = authorization[7:] if token == ADMIN_KEY: return token data = load_data() for k in data["api_keys"].values(): if k["key"] == token: return token raise HTTPException(status_code=403, detail="Invalid API key") # ── Routing logic ────────────────────────────────────────────────────────────── async def resolve_endpoint(model_str: str, data: dict): if "/" in model_str: parts = model_str.split("/", 1) provider_id, model_name = parts[0], parts[1] provider = data["providers"].get(provider_id) if not provider: raise HTTPException(status_code=404, detail=f"Provider '{provider_id}' not found") return provider, model_name, provider_id pool = None pool_id = None for pid, p in data["pool_models"].items(): if p["name"] == model_str or pid == model_str: pool = p pool_id = pid break if not pool: raise HTTPException(status_code=404, detail=f"Model or pool '{model_str}' not found.") members = pool.get("members", []) if not members: raise HTTPException(status_code=503, detail=f"Pool '{model_str}' has no members") strategy = pool.get("strategy", "round_robin") if strategy == "round_robin": start = await rl_state.next_rr(pool_id, len(members)) for i in range(len(members)): m = members[(start + i) % len(members)] provider = data["providers"].get(m["provider_id"]) if not provider: continue if await rl_state.is_available(m["provider_id"], m["model_name"], provider): return provider, m["model_name"], m["provider_id"] elif strategy == "priority": for m in members: provider = data["providers"].get(m["provider_id"]) if not provider: continue if await rl_state.is_available(m["provider_id"], m["model_name"], provider): return provider, m["model_name"], m["provider_id"] elif strategy == "weighted": import random available = [] for m in members: provider = data["providers"].get(m["provider_id"]) if not provider: continue if await rl_state.is_available(m["provider_id"], m["model_name"], provider): available.append((m, provider)) if available: weights = [m[0].get("weight", 1) for m in available] chosen_m, chosen_p = random.choices(available, weights=weights, k=1)[0] return chosen_p, chosen_m["model_name"], chosen_m["provider_id"] raise HTTPException(status_code=429, detail=f"All endpoints in pool '{model_str}' are currently rate-limited.") # ── Proxy logic ──────────────────────────────────────────────────────────────── TIMEOUT = httpx.Timeout(320.0, connect=30.0, read=300.0, write=30.0) def _extract_retry_after(response: httpx.Response) -> float: ra = response.headers.get("retry-after", "") try: return float(ra) except Exception: return 60.0 async def proxy_request(provider: dict, model_name: str, provider_id: str, body: dict, stream: bool): base_url = provider["base_url"].rstrip("/") api_key = provider["api_key"] headers = { "Authorization": f"Bearer {api_key}", "Content-Type": "application/json", } body = {**body, "model": model_name} url = f"{base_url}/chat/completions" await rl_state.record_request(provider_id, model_name) if stream: async def stream_gen(): async with httpx.AsyncClient(timeout=TIMEOUT) as client: async with client.stream("POST", url, json=body, headers=headers) as resp: if resp.status_code == 429: cooldown = _extract_retry_after(resp) await rl_state.mark_rate_limited(provider_id, model_name, cooldown) yield f"data: {json.dumps({'error': 'rate_limited', 'provider': provider_id})}\n\n" return async for chunk in resp.aiter_bytes(): yield chunk return StreamingResponse(stream_gen(), media_type="text/event-stream") async with httpx.AsyncClient(timeout=TIMEOUT) as client: resp = await client.post(url, json=body, headers=headers) if resp.status_code == 429: cooldown = _extract_retry_after(resp) await rl_state.mark_rate_limited(provider_id, model_name, cooldown) raise HTTPException(status_code=429, detail=f"Upstream rate limit hit on {provider_id}/{model_name}. Retry after {cooldown}s.") if resp.status_code >= 500: raise HTTPException(status_code=502, detail=f"Upstream error from {provider_id}: {resp.status_code}") if resp.status_code >= 400: try: detail = resp.json() except Exception: detail = resp.text raise HTTPException(status_code=resp.status_code, detail=detail) return JSONResponse(content=resp.json()) # ── App ──────────────────────────────────────────────────────────────────────── @asynccontextmanager async def lifespan(app: FastAPI): log.info("RyzGateway starting up") if not DATA_FILE.exists(): save_data(_default_data()) log.info(f"Initialized empty data at {DATA_FILE}") yield log.info("RyzGateway shutting down") app = FastAPI(title="RyzGateway", lifespan=lifespan) # ── OpenAI-compatible endpoints ──────────────────────────────────────────────── class ChatMessage(BaseModel): role: str content: Any class ChatRequest(BaseModel): model: str messages: List[ChatMessage] stream: Optional[bool] = False temperature: Optional[float] = None max_tokens: Optional[int] = None top_p: Optional[float] = None model_config = {"extra": "allow"} @app.post("/v1/chat/completions") async def chat_completions(request: ChatRequest, _key: str = Depends(require_api_key)): data = load_data() provider, model_name, provider_id = await resolve_endpoint(request.model, data) body = request.model_dump(exclude_none=True, exclude={"model"}) body.update(request.model_extra or {}) return await proxy_request(provider, model_name, provider_id, body, request.stream or False) @app.get("/v1/models") async def list_models(_key: str = Depends(require_api_key)): data = load_data() now = int(time.time()) models = [] for pid, p in data["providers"].items(): for m in p.get("known_models", []): models.append({ "id": f"{pid}/{m}", "object": "model", "created": now, "owned_by": pid, }) for pool_id, pool in data["pool_models"].items(): models.append({ "id": pool["name"], "object": "model", "created": now, "owned_by": "pool", }) return {"object": "list", "data": models} # ── Admin API ────────────────────────────────────────────────────────────────── @app.get("/admin/api-keys") async def get_api_keys(_=Depends(require_admin)): data = load_data() return list(data["api_keys"].values()) @app.post("/admin/api-keys") async def create_api_key(body: dict, _=Depends(require_admin)): data = load_data() kid = str(uuid.uuid4())[:8] key = "ryz-" + str(uuid.uuid4()).replace("-", "") data["api_keys"][kid] = { "id": kid, "key": key, "name": body.get("name", "Unnamed Key"), "created_at": int(time.time()), } save_data(data) return data["api_keys"][kid] @app.delete("/admin/api-keys/{kid}") async def delete_api_key(kid: str, _=Depends(require_admin)): data = load_data() if kid not in data["api_keys"]: raise HTTPException(status_code=404, detail="Key not found") del data["api_keys"][kid] save_data(data) return {"deleted": True} # --- Providers --- @app.get("/admin/providers") async def get_providers(_=Depends(require_admin)): data = load_data() result = [] for pid, p in data["providers"].items(): entry = {**p, "id": pid} raw = entry.get("api_key", "") entry["api_key_masked"] = raw[:6] + "..." + raw[-4:] if len(raw) > 10 else "****" entry.pop("api_key", None) cooldowns = {} for m in p.get("known_models", []): cd = rl_state.get_cooldown_remaining(pid, m) if cd > 0: cooldowns[m] = round(cd, 1) entry["active_cooldowns"] = cooldowns result.append(entry) return result @app.post("/admin/providers") async def create_provider(body: dict, _=Depends(require_admin)): data = load_data() pid = body.get("id") or re.sub(r"[^a-z0-9_-]", "-", body.get("name", "provider").lower())[:32] if pid in data["providers"]: raise HTTPException(status_code=409, detail=f"Provider ID '{pid}' already exists") data["providers"][pid] = { "name": body.get("name", pid), "base_url": body.get("base_url", ""), "api_key": body.get("api_key", ""), "rate_type": body.get("rate_type", "none"), "rate_limit": int(body.get("rate_limit", 0)), "known_models": body.get("known_models", []), "notes": body.get("notes", ""), } save_data(data) return {"id": pid, **data["providers"][pid], "api_key": "[hidden]"} @app.put("/admin/providers/{pid}") async def update_provider(pid: str, body: dict, _=Depends(require_admin)): data = load_data() if pid not in data["providers"]: raise HTTPException(status_code=404, detail="Provider not found") p = data["providers"][pid] for field in ["name", "base_url", "rate_type", "notes", "known_models"]: if field in body: p[field] = body[field] if "rate_limit" in body: p["rate_limit"] = int(body["rate_limit"]) if "api_key" in body and body["api_key"] and not body["api_key"].startswith("****"): p["api_key"] = body["api_key"] save_data(data) return {"id": pid, **p, "api_key": "[hidden]"} @app.delete("/admin/providers/{pid}") async def delete_provider(pid: str, _=Depends(require_admin)): data = load_data() if pid not in data["providers"]: raise HTTPException(status_code=404, detail="Provider not found") del data["providers"][pid] save_data(data) return {"deleted": True} @app.post("/admin/providers/{pid}/fetch-models") async def fetch_provider_models(pid: str, _=Depends(require_admin)): data = load_data() if pid not in data["providers"]: raise HTTPException(status_code=404, detail="Provider not found") p = data["providers"][pid] base_url = p["base_url"].rstrip("/") api_key = p["api_key"] if not base_url or not api_key: raise HTTPException(status_code=400, detail="Provider base_url or api_key is missing.") headers = { "Authorization": f"Bearer {api_key}", "Content-Type": "application/json", } try: async with httpx.AsyncClient(timeout=15.0) as client: resp = await client.get(f"{base_url}/models", headers=headers) if resp.status_code != 200: raise HTTPException(status_code=502, detail=f"Upstream returned {resp.status_code}: {resp.text[:200]}") resp_data = resp.json() models = [] if isinstance(resp_data, dict) and "data" in resp_data: for item in resp_data["data"]: if "id" in item: models.append(item["id"]) elif isinstance(resp_data, list): for item in resp_data: if isinstance(item, dict) and "id" in item: models.append(item["id"]) elif isinstance(item, str): models.append(item) if not models: raise HTTPException(status_code=404, detail="No models found in upstream /v1/models response.") # Remove duplicates and sort unique_models = sorted(list(set(models))) p["known_models"] = unique_models save_data(data) return {"fetched_count": len(unique_models), "models": unique_models} except httpx.RequestError as e: raise HTTPException(status_code=502, detail=f"Network error fetching models: {str(e)}") # --- Pool Models --- @app.get("/admin/pools") async def get_pools(_=Depends(require_admin)): data = load_data() result = [] for pool_id, pool in data["pool_models"].items(): entry = {**pool, "id": pool_id} annotated = [] for m in pool.get("members", []): provider = data["providers"].get(m["provider_id"], {}) available = await rl_state.is_available(m["provider_id"], m["model_name"], provider) cooldown = rl_state.get_cooldown_remaining(m["provider_id"], m["model_name"]) annotated.append({**m, "available": available, "cooldown_remaining": round(cooldown, 1)}) entry["members"] = annotated result.append(entry) return result @app.post("/admin/pools") async def create_pool(body: dict, _=Depends(require_admin)): data = load_data() pool_id = body.get("id") or re.sub(r"[^a-z0-9_-]", "-", body.get("name", "pool").lower())[:32] if pool_id in data["pool_models"]: raise HTTPException(status_code=409, detail=f"Pool ID '{pool_id}' already exists") data["pool_models"][pool_id] = { "name": body.get("name", pool_id), "strategy": body.get("strategy", "round_robin"), "members": body.get("members", []), } save_data(data) return {"id": pool_id, **data["pool_models"][pool_id]} @app.put("/admin/pools/{pool_id}") async def update_pool(pool_id: str, body: dict, _=Depends(require_admin)): data = load_data() if pool_id not in data["pool_models"]: raise HTTPException(status_code=404, detail="Pool not found") p = data["pool_models"][pool_id] for field in ["name", "strategy", "members"]: if field in body: p[field] = body[field] save_data(data) return {"id": pool_id, **p} @app.delete("/admin/pools/{pool_id}") async def delete_pool(pool_id: str, _=Depends(require_admin)): data = load_data() if pool_id not in data["pool_models"]: raise HTTPException(status_code=404, detail="Pool not found") del data["pool_models"][pool_id] save_data(data) return {"deleted": True} # --- Status --- @app.get("/admin/status") async def get_status(_=Depends(require_admin)): data = load_data() return { "providers": len(data["providers"]), "pools": len(data["pool_models"]), "api_keys": len(data["api_keys"]), } # ── Dashboard (Mobile-Optimized Glassmorphism UI) ────────────────────────────── DASHBOARD_HTML = r""" RyzGateway

RyzGateway

Enter your admin key to access the dashboard.

""" @app.get("/", response_class=HTMLResponse) async def dashboard(): return DASHBOARD_HTML if __name__ == "__main__": import uvicorn uvicorn.run("app:app", host="0.0.0.0", port=7860, reload=False)