"""
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.
Invalid key.
RyzGateway
Dashboard
Gateway status and routing health
—
Providers
—
Pools
—
Keys
Provider health
No providers configured
Providers
API endpoints with rate limit configuration
Provider ID is auto-generated from the name and used in model routing. Use provider-id/exact-model-name to route directly. Click Fetch Models on a provider to auto-import available models.
Provider
Base URL
Rate limit
Models
Status
Loading...
Model Pools
Generic model names that route across multiple providers
Pool names are generic model names exposed in /v1/models. The gateway routes requests using your chosen strategy and automatically skips rate-limited endpoints.
Loading...
API Keys
Keys issued to clients for /v1/ access
Name
Key
Created
Loading...
Usage docs
How to use RyzGateway in your apps
Base URL
Point any OpenAI-compatible client at this URL.
Model name formats
provider-id/exact-model-name — Direct passthrough.