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)