Spaces:
Paused
Paused
File size: 18,652 Bytes
c5c5d80 03f171d e24791f 03f171d e24791f 03f171d e24791f 03f171d e24791f 03f171d e24791f 1256f93 111e96a b6dd083 e24791f b6dd083 ea909c9 b6dd083 e24791f 03f171d e24791f b6dd083 e24791f 0656d63 e24791f 0656d63 e24791f b6dd083 e24791f 03f171d e24791f 261a815 e24791f ddd1818 03f171d 261a815 e24791f 03f171d 261a815 e24791f 261a815 e24791f 03f171d e24791f 03f171d b6dd083 e24791f 03f171d e24791f 03f171d e24791f b6dd083 e24791f c5c5d80 e24791f c5c5d80 e24791f c5c5d80 e24791f c5c5d80 e24791f 03f171d e24791f 03f171d e24791f 03f171d e24791f 03f171d b6dd083 e24791f b6dd083 e24791f 03f171d b6dd083 03f171d c5c5d80 e24791f 03f171d e24791f b6dd083 e24791f 03f171d c5c5d80 e24791f c5c5d80 e24791f c5c5d80 e24791f c5c5d80 03f171d | 1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 82 83 84 85 86 87 88 89 90 91 92 93 94 95 96 97 98 99 100 101 102 103 104 105 106 107 108 109 110 111 112 113 114 115 116 117 118 119 120 121 122 123 124 125 126 127 128 129 130 131 132 133 134 135 136 137 138 139 140 141 142 143 144 145 146 147 148 149 150 151 152 153 154 155 156 157 158 159 160 161 162 163 164 165 166 167 168 169 170 171 172 173 174 175 176 177 178 179 180 181 182 183 184 185 186 187 188 189 190 191 192 193 194 195 196 197 198 199 200 201 202 203 204 205 206 207 208 209 210 211 212 213 214 215 216 217 218 219 220 221 222 223 224 225 226 227 228 229 230 231 232 233 234 235 236 237 238 239 240 241 242 243 244 245 246 247 248 249 250 251 252 253 254 255 256 257 258 259 260 261 262 263 264 265 266 267 268 269 270 271 272 273 274 275 276 277 278 279 280 281 282 283 284 285 286 287 288 289 290 291 292 293 294 295 296 297 298 299 300 301 302 303 304 305 306 307 308 309 310 311 312 313 314 315 316 317 318 319 320 321 322 323 324 325 326 327 328 329 330 331 332 333 334 335 336 337 338 339 340 341 342 343 344 345 346 347 348 349 350 351 352 353 354 355 356 357 358 359 360 361 362 363 364 365 366 367 368 369 370 371 372 373 374 375 376 377 378 379 380 381 382 383 384 385 386 387 388 389 390 391 392 393 394 395 396 397 398 399 400 401 402 403 404 405 406 407 408 409 410 411 412 413 414 415 416 417 418 419 420 421 422 423 424 425 426 427 428 429 430 431 432 433 434 435 436 437 438 439 440 441 442 443 444 445 446 447 448 449 450 451 452 453 454 455 456 457 458 459 460 461 | import os
import subprocess
import asyncio
import time
import threading
import httpx
from fastapi import FastAPI, Request, HTTPException
from fastapi.openapi.docs import get_swagger_ui_html
from fastapi.responses import JSONResponse, StreamingResponse
app = FastAPI(title="Ornith 1.0 Inference API - CPU Agentic (single-instance)", version="2.0.0")
# ---------------------------------------------------------------------------
# Configuration
# ---------------------------------------------------------------------------
MODEL_PATH = os.getenv("MODEL_PATH", "/models/ornith-1.0-9b-Q4_K_M.gguf")
MODEL_REPO = os.getenv("MODEL_REPO", "deepreinforce-ai/Ornith-1.0-9B-GGUF")
MODEL_ALIAS = os.getenv("MODEL_ALIAS", "ornith-1.0")
MAIN_PORT = int(os.getenv("MAIN_PORT", "7860"))
# Single instance on CPU: one process gets ALL cores, concurrency via slots.
NUM_INSTANCES = int(os.getenv("NUM_INSTANCES", "1"))
BASE_PORT = int(os.getenv("BASE_PORT", "8081"))
INSTANCE_PORTS = [BASE_PORT + i for i in range(NUM_INSTANCES)]
INSTANCE_URLS = [f"http://127.0.0.1:{p}" for p in INSTANCE_PORTS]
MODEL_LOADED = [False] * NUM_INSTANCES
LLAMA_PROCESSES = [None] * NUM_INSTANCES
# CPU / inference tuning
def _effective_cpus():
"""Container CPU limit (cgroup CFS quota) — os.cpu_count() reports the HOST's
cores and ignores the cpu-basic throttle (~2 vCPU), which would oversubscribe
threads. Fall back through cgroup v2 -> v1 -> affinity -> count."""
try:
with open("/sys/fs/cgroup/cpu.max") as f: # cgroup v2
quota, period = f.read().split()
if quota != "max":
n = int(float(quota) / float(period))
if n >= 1:
return n
except Exception:
pass
try:
with open("/sys/fs/cgroup/cpu/cpu.cfs_quota_us") as f: # cgroup v1
quota = int(f.read())
with open("/sys/fs/cgroup/cpu/cpu.cfs_period_us") as f:
period = int(f.read())
if quota > 0 and period > 0 and quota // period >= 1:
return quota // period
except Exception:
pass
try:
return len(os.sched_getaffinity(0))
except Exception:
return os.cpu_count() or 4
_CPU = _effective_cpus()
# Hard ceiling on threads: even if detection over-reports (e.g. host cores leak
# through), never oversubscribe the ~2-vCPU tier. Raise CPU_THREADS_MAX if you
# move to a bigger CPU tier.
CPU_THREADS_MAX = int(os.getenv("CPU_THREADS_MAX", "2"))
CPU_THREADS = min(int(os.getenv("CPU_THREADS", str(_CPU))), CPU_THREADS_MAX)
CPU_THREADS_BATCH = min(int(os.getenv("CPU_THREADS_BATCH", str(_CPU))), CPU_THREADS_MAX)
CONTEXT_SIZE = int(os.getenv("CONTEXT_SIZE", "32768")) # total; per-slot = ctx/parallel
PARALLEL = int(os.getenv("PARALLEL", "2")) # continuous-batching slots
BATCH_SIZE = int(os.getenv("BATCH_SIZE", "512"))
UBATCH_SIZE = int(os.getenv("UBATCH_SIZE", "512"))
# Split K/V cache quantization: most of the quality loss from cache quantization
# comes from K (TurboQuant benchmarks: K-only = 6.6% of the 7.6% total perplexity
# hit), so keep K at q8_0 and take the memory saving on V. Legacy KV_CACHE_QUANT
# still overrides both. Aliases normalize stale values ("4bit", ...) so a bad env
# var can never crash llama-server on startup.
_KV_ALIASES = {"2bit": "q4_0", "4bit": "q4_0", "8bit": "q8_0",
"16bit": "f16", "fp16": "f16", "q4": "q4_0", "q8": "q8_0"}
_VALID_KV = {"f16", "q8_0", "q4_0", "q4_1", "q5_0", "q5_1", "iq4_nl"}
def _kv_type(env_name, default):
v = os.getenv(env_name, os.getenv("KV_CACHE_QUANT", default))
v = _KV_ALIASES.get(v.lower(), v)
return v if v in _VALID_KV else default
KV_CACHE_QUANT_K = _kv_type("KV_CACHE_QUANT_K", "q8_0")
KV_CACHE_QUANT_V = _kv_type("KV_CACHE_QUANT_V", "q4_0")
CACHE_REUSE = int(os.getenv("CACHE_REUSE", "256")) # min tokens for prompt-prefix reuse
FLASH_ATTN = os.getenv("FLASH_ATTN", "true").lower() == "true"
MMAP_ENABLED = os.getenv("MMAP_ENABLED", "true").lower() == "true"
MLOCK_ENABLED = os.getenv("MLOCK_ENABLED", "false").lower() == "true"
REASONING_FORMAT = os.getenv("REASONING_FORMAT", "auto") # for <think> reasoning models
# Shared HTTP client (connection pooling) — created on startup.
HTTP: httpx.AsyncClient | None = None
_rr_counter = 0
_rr_lock = asyncio.Lock()
print("=" * 60)
print("Ornith 1.0 Agentic Inference API (v2.0.0)")
print(f" detected vCPUs : {_CPU}")
print(f" instances : {NUM_INSTANCES} ports={INSTANCE_PORTS}")
print(f" threads/instance : {CPU_THREADS} (batch {CPU_THREADS_BATCH})")
print(f" context (total) : {CONTEXT_SIZE} parallel slots: {PARALLEL}")
print(f" kv-cache quant : K={KV_CACHE_QUANT_K} V={KV_CACHE_QUANT_V} flash-attn: {FLASH_ATTN}")
print(f" cache-reuse : {CACHE_REUSE}")
print(f" mmap/mlock : {MMAP_ENABLED}/{MLOCK_ENABLED}")
print("=" * 60)
# ---------------------------------------------------------------------------
# llama-server launch (flag-compat aware so we survive llama.cpp CLI changes)
# ---------------------------------------------------------------------------
def _find_binary():
for p in ("/llama.cpp/build/bin/llama-server",
"/llama.cpp/build/bin/server",
"/usr/local/bin/llama-server"):
if os.path.exists(p):
return p
try:
found = subprocess.run(["find", "/llama.cpp", "-name", "llama-server"],
capture_output=True, text=True).stdout.strip().split("\n")
if found and found[0]:
return found[0]
except Exception:
pass
return None
def _help_text(binary):
try:
r = subprocess.run([binary, "--help"], capture_output=True, text=True, timeout=30)
return (r.stdout or "") + (r.stderr or "")
except Exception:
return ""
def _download_model():
if os.path.exists(MODEL_PATH):
return True
print(f"Model missing at {MODEL_PATH}; downloading {os.path.basename(MODEL_PATH)} from {MODEL_REPO}...")
try:
from huggingface_hub import hf_hub_download
import shutil
dl = hf_hub_download(repo_id=MODEL_REPO, filename=os.path.basename(MODEL_PATH),
repo_type="model", token=os.getenv("HF_TOKEN"))
if dl != MODEL_PATH:
os.makedirs(os.path.dirname(MODEL_PATH), exist_ok=True)
shutil.copy2(dl, MODEL_PATH)
print(f"✅ Model ready at {MODEL_PATH}")
return True
except Exception as e:
print(f"❌ Model download failed: {e}")
return False
def build_cmd(binary, port, help_text):
"""Compose llama-server args, only including flags the binary actually supports."""
def has(flag):
return flag in help_text
cmd = [binary,
"--model", MODEL_PATH,
"--host", "127.0.0.1",
"--port", str(port),
"--ctx-size", str(CONTEXT_SIZE),
"--threads", str(CPU_THREADS),
"--batch-size", str(BATCH_SIZE)]
if has("--threads-batch"):
cmd += ["--threads-batch", str(CPU_THREADS_BATCH)]
if has("--ubatch-size"):
cmd += ["--ubatch-size", str(UBATCH_SIZE)]
if has("--parallel"):
cmd += ["--parallel", str(PARALLEL)]
if has("--cont-batching"):
cmd += ["--cont-batching"]
if has("--alias"):
cmd += ["--alias", MODEL_ALIAS]
# Flash attention — value form (on/off/auto) vs legacy boolean.
if FLASH_ATTN and has("--flash-attn"):
fa_line = next((l for l in help_text.splitlines() if "--flash-attn" in l), "")
if any(t in fa_line for t in ("{on", "[on", "on|off", "on,off")):
cmd += ["--flash-attn", "on"]
else:
cmd += ["--flash-attn"]
# KV cache quantization (split K/V). Quantized V requires flash-attn.
if KV_CACHE_QUANT_K != "f16" and has("--cache-type-k"):
cmd += ["--cache-type-k", KV_CACHE_QUANT_K]
if KV_CACHE_QUANT_V != "f16" and has("--cache-type-v"):
if FLASH_ATTN:
cmd += ["--cache-type-v", KV_CACHE_QUANT_V]
else:
print("⚠️ FLASH_ATTN=false: V cache falls back to f16 (quantized V "
"needs flash attention) — KV memory roughly doubles")
# Reuse cached prompt prefixes across requests: agents resend the same system
# prompt every turn, and prefill is the slow path on CPU.
if CACHE_REUSE > 0 and has("--cache-reuse"):
cmd += ["--cache-reuse", str(CACHE_REUSE)]
# Chat template + tool-calling (native OpenAI function calling).
if has("--jinja"):
cmd += ["--jinja"]
if has("--reasoning-format"):
cmd += ["--reasoning-format", REASONING_FORMAT]
if MLOCK_ENABLED and has("--mlock"):
cmd += ["--mlock"]
if has("--mmap") or has("--no-mmap"):
cmd += ["--mmap"] if MMAP_ENABLED else ["--no-mmap"]
return cmd
def start_llama_server(idx):
port = INSTANCE_PORTS[idx]
if not _download_model():
MODEL_LOADED[idx] = False
return False
binary = _find_binary()
if not binary:
print(f"Instance {idx}: ❌ llama-server binary not found")
return False
cmd = build_cmd(binary, port, _help_text(binary))
print(f"Instance {idx}: launching -> {' '.join(cmd)}")
try:
proc = subprocess.Popen(cmd, stdout=subprocess.PIPE, stderr=subprocess.STDOUT, text=True)
LLAMA_PROCESSES[idx] = proc
def log_output(p, i):
for line in p.stdout:
print(f"Instance {i} LOG: {line.rstrip()}")
threading.Thread(target=log_output, args=(proc, idx), daemon=True).start()
# Wait for readiness. Newer llama-server binds the port and returns HTTP
# 503 *while still loading*, so we must sleep on every non-200 (not only on
# connection exceptions) or we'd spin through all iterations instantly.
for _ in range(900):
if proc.poll() is not None:
print(f"Instance {idx}: ❌ exited with code {proc.returncode}")
return False
try:
r = httpx.get(f"{INSTANCE_URLS[idx]}/health", timeout=3.0)
if r.status_code == 200:
MODEL_LOADED[idx] = True
print(f"Instance {idx}: ✅ ready on port {port}")
return True
# 503 => model still loading; keep waiting.
except Exception:
pass
time.sleep(1)
print(f"Instance {idx}: ❌ startup timeout")
return False
except Exception as e:
print(f"Instance {idx}: error: {e}")
return False
# ---------------------------------------------------------------------------
# Silent-degradation guard: aggressive KV/weight quantization can break
# tool-call JSON without crashing. Verify one canned tool call parses.
# ---------------------------------------------------------------------------
SMOKE = {"ran": False, "ok": None, "detail": "pending"}
def _tool_call_smoke_test():
import json
try:
payload = {
"model": MODEL_ALIAS,
"messages": [{"role": "user",
"content": "What is the weather in Berlin? Use the tool."}],
"tools": [{"type": "function", "function": {
"name": "get_weather",
"description": "Get current weather for a city",
"parameters": {"type": "object",
"properties": {"city": {"type": "string"}},
"required": ["city"]}}}],
"tool_choice": "auto", "max_tokens": 128, "temperature": 0,
}
r = httpx.post(f"{INSTANCE_URLS[0]}/v1/chat/completions",
json=payload, timeout=300.0)
msg = r.json()["choices"][0]["message"]
calls = msg.get("tool_calls") or []
if calls:
json.loads(calls[0]["function"]["arguments"])
SMOKE.update(ran=True, ok=True,
detail=f"tool call parsed: {calls[0]['function']['name']}")
else:
SMOKE.update(ran=True, ok=False,
detail="no tool_calls in response — check quantization/chat template")
except Exception as e:
SMOKE.update(ran=True, ok=False, detail=f"error: {e}")
print(f"Tool-call smoke test: {'✅' if SMOKE['ok'] else '⚠️'} {SMOKE['detail']}")
# ---------------------------------------------------------------------------
# Instance selection (round-robin, no per-request network pre-flight)
# ---------------------------------------------------------------------------
async def pick_instance():
global _rr_counter
async with _rr_lock:
for k in range(NUM_INSTANCES):
idx = (_rr_counter + k) % NUM_INSTANCES
if MODEL_LOADED[idx]:
_rr_counter = (idx + 1) % NUM_INSTANCES
return idx
return None
# ---------------------------------------------------------------------------
# Lifecycle
# ---------------------------------------------------------------------------
@app.on_event("startup")
async def startup_event():
global HTTP
# Long read timeout: CPU generation of a full agent turn can take a while.
HTTP = httpx.AsyncClient(timeout=httpx.Timeout(connect=10.0, read=600.0, write=30.0, pool=600.0),
limits=httpx.Limits(max_connections=64, max_keepalive_connections=32))
print(f"Starting {NUM_INSTANCES} llama.cpp instance(s)...")
results = await asyncio.gather(
*[asyncio.to_thread(start_llama_server, i) for i in range(NUM_INSTANCES)],
return_exceptions=True)
loaded = sum(1 for r in results if r is True)
print(f"✅ {loaded}/{NUM_INSTANCES} instance(s) loaded")
if loaded:
threading.Thread(target=_tool_call_smoke_test, daemon=True).start()
@app.on_event("shutdown")
async def shutdown_event():
if HTTP:
await HTTP.aclose()
for p in LLAMA_PROCESSES:
if p and p.poll() is None:
p.terminate()
# ---------------------------------------------------------------------------
# Reverse-proxy core (streaming + non-streaming) to llama.cpp OpenAI endpoints
# ---------------------------------------------------------------------------
async def proxy(request: Request, path: str):
idx = await pick_instance()
if idx is None:
raise HTTPException(status_code=503, detail="No healthy instances available")
url = f"{INSTANCE_URLS[idx]}{path}"
body = await request.body()
try:
payload = await request.json()
except Exception:
payload = {}
stream = bool(payload.get("stream", False))
headers = {"Content-Type": "application/json"}
if stream:
req = HTTP.build_request("POST", url, content=body, headers=headers)
r = await HTTP.send(req, stream=True)
async def gen():
try:
async for chunk in r.aiter_raw():
yield chunk
finally:
await r.aclose()
return StreamingResponse(gen(), status_code=r.status_code,
media_type=r.headers.get("content-type", "text/event-stream"),
headers={"X-Backend-Instance": str(idx)})
else:
r = await HTTP.post(url, content=body, headers=headers)
return JSONResponse(status_code=r.status_code,
content=r.json() if r.content else {},
headers={"X-Backend-Instance": str(idx)})
# ---------------------------------------------------------------------------
# Endpoints
# ---------------------------------------------------------------------------
@app.get("/", include_in_schema=False)
async def swagger():
return get_swagger_ui_html(openapi_url=app.openapi_url, title=app.title + " - Swagger UI")
@app.get("/health")
async def health():
return {
"status": "ok" if any(MODEL_LOADED) else "degraded",
"instances": [
{"id": i, "port": INSTANCE_PORTS[i], "loaded": MODEL_LOADED[i], "url": INSTANCE_URLS[i]}
for i in range(NUM_INSTANCES)
],
"active_instances": sum(MODEL_LOADED),
"total_instances": NUM_INSTANCES,
"cpu_threads": CPU_THREADS,
"context_size": CONTEXT_SIZE,
"context_per_slot": CONTEXT_SIZE // max(PARALLEL, 1),
"parallel_slots": PARALLEL,
"kv_cache_quant_k": KV_CACHE_QUANT_K,
"kv_cache_quant_v": KV_CACHE_QUANT_V,
"cache_reuse": CACHE_REUSE,
"flash_attn": FLASH_ATTN,
"mmap_enabled": MMAP_ENABLED,
"tool_call_smoke": SMOKE,
}
@app.get("/v1/config")
async def config():
return {
"model_path": MODEL_PATH, "model_alias": MODEL_ALIAS,
"detected_vcpus": _CPU,
"instances": NUM_INSTANCES, "ports": INSTANCE_PORTS,
"cpu_threads": CPU_THREADS, "threads_batch": CPU_THREADS_BATCH,
"context_size": CONTEXT_SIZE, "parallel_slots": PARALLEL,
"batch_size": BATCH_SIZE, "ubatch_size": UBATCH_SIZE,
"kv_cache_quant_k": KV_CACHE_QUANT_K, "kv_cache_quant_v": KV_CACHE_QUANT_V,
"cache_reuse": CACHE_REUSE, "flash_attn": FLASH_ATTN,
"mmap_enabled": MMAP_ENABLED, "mlock_enabled": MLOCK_ENABLED,
"reasoning_format": REASONING_FORMAT,
"instances_status": [{"id": i, "loaded": MODEL_LOADED[i]} for i in range(NUM_INSTANCES)],
}
@app.get("/v1/models")
async def models():
"""Proxy to the backend so real capabilities (tools, etc.) are reported."""
idx = await pick_instance()
if idx is not None:
try:
r = await HTTP.get(f"{INSTANCE_URLS[idx]}/v1/models", timeout=10.0)
if r.status_code == 200:
return JSONResponse(content=r.json())
except Exception:
pass
return {"object": "list", "data": [{"id": MODEL_ALIAS, "object": "model", "owned_by": "Leon4gr45"}]}
# Native OpenAI-compatible endpoints: full passthrough (tools, tool_choice,
# response_format, streaming, logprobs, etc. all handled by llama.cpp).
@app.post("/v1/chat/completions")
async def chat_completions(request: Request):
return await proxy(request, "/v1/chat/completions")
@app.post("/v1/completions")
async def completions(request: Request):
return await proxy(request, "/v1/completions")
@app.post("/v1/embeddings")
async def embeddings(request: Request):
return await proxy(request, "/v1/embeddings")
if __name__ == "__main__":
import uvicorn
uvicorn.run(app, host="0.0.0.0", port=MAIN_PORT)
|