Add vLLM as a first-class managed engine backend (LLM + OCR Surya2)

vLLM integrated exactly like ds4/ktransformers — a managed external engine in an
ISOLATED venv (it pins torch 2.13/cu13, conflicting with the main venv), proxied over
its OpenAI HTTP API, selected per-model via a `backend: vllm` pin or the vllm.model_id
alias (never auto-claimed). Continuous batching for high aggregate throughput.

- config.py: VllmConfig (+ Config field, from_dict, to_dict)
- codai/api/vllm_worker.py: managed vllm.entrypoints.openai.api_server subprocess in the
  isolated venv (venv resolver: config>/opt/coderai/vllm_venv>/cache>~/.coderai), /v1/models
  health gate, auto_build
- codai/backends/vllm.py: VllmBackend OpenAI proxy (mirrors ktransformers)
- manager: get_active_vllm_config, "vllm" in _ENGINE_BACKENDS, _vllm_name_claims=False,
  vllm_should_handle, load branch, text-accept, /v1/models surfacing
- front-proxy: required_capability (vllm pin+alias), _DEFAULT_CAPS vllm on GPU nodes,
  assignment/engine_supervisor/app threading + reload list
- admin: routes get/set, settings.html vLLM card, models.html dropdown option
- requirements-vllm.txt (vllm==0.27.1); docs/vllm.md marked implemented; __version__ 0.1.86

OCR Surya2-via-vLLM: ocr.surya_serve = local|vllm|llamacpp. In vllm mode the Surya engine
serves surya_model (datalab-to/surya-ocr-2) through the vLLM backend and attaches via
SURYA_INFERENCE_URL — the correct path for the latest "Surya2" VLM (llama-cpp-python's
server hit a recurrent/hybrid KV-slot bug on surya-2).
Co-Authored-By: 's avatarClaude Opus 4.8 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01Mw2KQiswmD69T45fTfjKwW
parent 2654d749
......@@ -16,7 +16,7 @@
# Canonical product version for CoderAI — single source of truth. Both the API
# metadata and the admin web UI read from here.
__version__ = "0.1.85"
__version__ = "0.1.86"
# Configure the CUDA caching allocator BEFORE torch is imported anywhere.
# expandable_segments lets the allocator return freed pages to the driver even
......
......@@ -3505,6 +3505,23 @@ def build_settings_dict(c, gpu_cards):
"extra_env": c.ktransformers.extra_env,
"auto_build": c.ktransformers.auto_build,
},
"vllm": {
"enabled": c.vllm.enabled,
"venv": c.vllm.venv,
"model_path": c.vllm.model_path,
"model_id": c.vllm.model_id,
"host": c.vllm.host,
"port": c.vllm.port,
"ctx": c.vllm.ctx,
"gpu_memory_utilization": c.vllm.gpu_memory_utilization,
"tensor_parallel_size": c.vllm.tensor_parallel_size,
"max_num_seqs": c.vllm.max_num_seqs,
"dtype": c.vllm.dtype,
"quantization": c.vllm.quantization,
"extra_args": c.vllm.extra_args,
"extra_env": c.vllm.extra_env,
"auto_build": c.vllm.auto_build,
},
"ocr": {
"enabled": c.ocr.enabled,
"default_engine": c.ocr.default_engine,
......@@ -3531,6 +3548,9 @@ def build_settings_dict(c, gpu_cards):
"surya_langs": c.ocr.surya_langs,
"surya_venv": c.ocr.surya_venv,
"surya_auto_build": c.ocr.surya_auto_build,
"surya_serve": c.ocr.surya_serve,
"surya_model": c.ocr.surya_model,
"surya_server_url": c.ocr.surya_server_url,
"detect_mode": c.ocr.detect_mode,
"detect_model_path": c.ocr.detect_model_path,
"detect_conf": c.ocr.detect_conf,
......@@ -3987,6 +4007,35 @@ async def api_save_settings(request: Request, username: str = Depends(require_ad
if "auto_build" in d:
c.ktransformers.auto_build = bool(d["auto_build"])
if "vllm" in data:
d = data["vllm"]
v = c.vllm
v.enabled = bool(d.get("enabled", v.enabled))
if "venv" in d: v.venv = (d.get("venv") or "").strip()
if "model_path" in d: v.model_path = (d.get("model_path") or "").strip()
if "model_id" in d: v.model_id = (d.get("model_id") or v.model_id or "vllm").strip()
if "host" in d: v.host = (d.get("host") or "127.0.0.1").strip()
if "port" in d:
try: v.port = int(d.get("port") or 0)
except (TypeError, ValueError): pass
if "ctx" in d:
try: v.ctx = max(512, int(d.get("ctx") or v.ctx))
except (TypeError, ValueError): pass
if "gpu_memory_utilization" in d:
try: v.gpu_memory_utilization = float(d.get("gpu_memory_utilization") or v.gpu_memory_utilization)
except (TypeError, ValueError): pass
if "tensor_parallel_size" in d:
try: v.tensor_parallel_size = max(1, int(d.get("tensor_parallel_size") or 1))
except (TypeError, ValueError): pass
if "max_num_seqs" in d:
try: v.max_num_seqs = max(0, int(d.get("max_num_seqs") or 0))
except (TypeError, ValueError): pass
if "dtype" in d: v.dtype = (d.get("dtype") or "").strip()
if "quantization" in d: v.quantization = (d.get("quantization") or "").strip()
if "extra_args" in d: v.extra_args = (d.get("extra_args") or "").strip()
if "extra_env" in d: v.extra_env = (d.get("extra_env") or "").strip()
if "auto_build" in d: v.auto_build = bool(d["auto_build"])
if "ocr" in data:
d = data["ocr"]
o = c.ocr
......@@ -4029,6 +4078,9 @@ async def api_save_settings(request: Request, username: str = Depends(require_ad
if "surya_langs" in d: o.surya_langs = (d.get("surya_langs") or "it").strip()
if "surya_venv" in d: o.surya_venv = (d.get("surya_venv") or "").strip()
if "surya_auto_build" in d: o.surya_auto_build = bool(d["surya_auto_build"])
if "surya_serve" in d: o.surya_serve = (d.get("surya_serve") or "local").strip()
if "surya_model" in d: o.surya_model = (d.get("surya_model") or "").strip()
if "surya_server_url" in d: o.surya_server_url = (d.get("surya_server_url") or "").strip()
# detection
if "detect_mode" in d: o.detect_mode = (d.get("detect_mode") or "off").strip()
if "detect_model_path" in d: o.detect_model_path = (d.get("detect_model_path") or "").strip()
......
......@@ -555,6 +555,7 @@ window.__DEFAULT_WHISPER_SERVER_PATH__ = {{ default_whisper_server_path|tojson }
<option value="ds4">ds4 (DeepSeek-V4 CUDA engine)</option>
<option value="k3">k3 (Kimi-K3 / kimi-k3-in-c)</option>
<option value="kt">ktransformers (SGLang, CPU+GPU)</option>
<option value="vllm">vLLM (high-concurrency, CUDA)</option>
</select>
</div>
<div class="form-row" style="margin:0">
......
......@@ -756,6 +756,58 @@
</div>
</div>
<!-- vLLM (high-concurrency) -->
<div class="card">
<div class="card-title">vLLM (high-concurrency)</div>
<p class="form-hint" style="margin-bottom:.6rem">Serve an LLM (or VLM) through <a href="https://github.com/vllm-project/vllm" target="_blank" rel="noopener">vLLM</a>'s OpenAI server with <b>continuous batching</b> — far higher aggregate throughput than serialized single-instance backends. coderai launches <code>vllm.entrypoints.openai.api_server</code> as a managed subprocess and proxies to it. vLLM pins its own torch/CUDA, so it runs in an <b>isolated venv</b> (CUDA-only). Selected per model via a <code>backend: vllm</code> pin or the model id below. Also serves Surya2 for OCR.</p>
<div class="form-row">
<label style="display:flex;align-items:center;gap:.5rem;cursor:pointer">
<input type="checkbox" id="s-vllm-enabled" onchange="toggleVllmFields()">
<span style="font-size:13px;font-weight:500">Enable vLLM</span>
</label>
</div>
<div id="vllm-fields" style="display:none">
<div class="form-group">
<label class="form-label">Model id (alias + --served-model-name)</label>
<input type="text" id="s-vllm-model-id" class="form-input" placeholder="vllm">
<span class="form-hint">This id routes to vLLM. Also matched by any model whose config <code>backend</code> is <code>vllm</code>.</span>
</div>
<div class="form-group">
<label class="form-label">Model path (HF dir or repo id)</label>
<input type="text" id="s-vllm-model-path" class="form-input" placeholder="/nvme/Qwen3-8B or Qwen/Qwen3-8B">
</div>
<div class="form-row" style="display:flex;gap:1rem;flex-wrap:wrap">
<div class="form-group" style="flex:1;min-width:130px"><label class="form-label">Max model len</label><input type="number" id="s-vllm-ctx" class="form-input" min="512" placeholder="32768"></div>
<div class="form-group" style="flex:1;min-width:130px"><label class="form-label">GPU mem util</label><input type="number" id="s-vllm-gmu" class="form-input" min="0.1" max="1" step="0.05" placeholder="0.90"></div>
<div class="form-group" style="flex:1;min-width:110px"><label class="form-label">TP size</label><input type="number" id="s-vllm-tp" class="form-input" min="1" placeholder="1"></div>
<div class="form-group" style="flex:1;min-width:120px"><label class="form-label">Max num seqs</label><input type="number" id="s-vllm-mns" class="form-input" min="0" placeholder="0 = default"></div>
</div>
<div class="form-row" style="display:flex;gap:1rem;flex-wrap:wrap">
<div class="form-group" style="flex:1;min-width:120px"><label class="form-label">dtype</label><input type="text" id="s-vllm-dtype" class="form-input" placeholder="auto"></div>
<div class="form-group" style="flex:1;min-width:120px"><label class="form-label">Quantization</label><input type="text" id="s-vllm-quant" class="form-input" placeholder="none (awq/gptq/fp8)"></div>
<div class="form-group" style="flex:1;min-width:90px"><label class="form-label">Port</label><input type="number" id="s-vllm-port" class="form-input" min="0" placeholder="0 = auto"></div>
</div>
<div class="form-group">
<label class="form-label">Isolated venv dir</label>
<input type="text" id="s-vllm-venv" class="form-input" placeholder="(auto: /opt/coderai/vllm_venv, else /cache/vllm_venv)">
</div>
<div class="form-group">
<label class="form-label">Extra args</label>
<input type="text" id="s-vllm-extra-args" class="form-input" placeholder="--enable-prefix-caching --swap-space 4">
</div>
<div class="form-group">
<label class="form-label">Extra subprocess env</label>
<input type="text" id="s-vllm-extra-env" class="form-input" placeholder="CUDA_VISIBLE_DEVICES=0">
</div>
<div class="form-group">
<label class="toggle-row" style="display:flex;align-items:center;gap:.5rem;cursor:pointer">
<input type="checkbox" id="s-vllm-auto-build">
<span>Auto-build the isolated venv (<code>pip install vllm</code>) on first use</span>
</label>
</div>
</div>
</div>
<!-- OCR (dedicated document OCR engines) -->
<div class="card">
<div class="card-title">OCR (document → text)</div>
......@@ -1073,6 +1125,10 @@ function toggleOcrFields(){
document.getElementById('ocr-fields').style.display =
document.getElementById('s-ocr-enabled').checked ? 'block' : 'none';
}
function toggleVllmFields(){
document.getElementById('vllm-fields').style.display =
document.getElementById('s-vllm-enabled').checked ? 'block' : 'none';
}
// ---- OCR isolated-venv build (paddle/surya) ----
const _ocrBuildPollers = {};
......@@ -1582,6 +1638,24 @@ async function loadSettings(){
document.getElementById('s-kt-extra-env').value = kt.extra_env ?? '';
toggleKtFields();
// vLLM
const vll = d.vllm || {};
document.getElementById('s-vllm-enabled').checked = !!vll.enabled;
document.getElementById('s-vllm-model-id').value = vll.model_id ?? 'vllm';
document.getElementById('s-vllm-model-path').value = vll.model_path ?? '';
document.getElementById('s-vllm-ctx').value = vll.ctx ?? 32768;
document.getElementById('s-vllm-gmu').value = vll.gpu_memory_utilization ?? 0.90;
document.getElementById('s-vllm-tp').value = vll.tensor_parallel_size ?? 1;
document.getElementById('s-vllm-mns').value = vll.max_num_seqs ?? 0;
document.getElementById('s-vllm-dtype').value = vll.dtype ?? '';
document.getElementById('s-vllm-quant').value = vll.quantization ?? '';
document.getElementById('s-vllm-port').value = vll.port ?? 0;
document.getElementById('s-vllm-venv').value = vll.venv ?? '';
document.getElementById('s-vllm-extra-args').value = vll.extra_args ?? '';
document.getElementById('s-vllm-extra-env').value = vll.extra_env ?? '';
document.getElementById('s-vllm-auto-build').checked = !!vll.auto_build;
toggleVllmFields();
// OCR subsystem
const ocr = d.ocr || {};
document.getElementById('s-ocr-enabled').checked = !!ocr.enabled;
......@@ -1743,6 +1817,22 @@ async function saveSettings(){
auto_build: document.getElementById('s-kt-auto-build').checked,
extra_env: document.getElementById('s-kt-extra-env').value.trim(),
},
vllm:{
enabled: document.getElementById('s-vllm-enabled').checked,
model_id: document.getElementById('s-vllm-model-id').value.trim() || 'vllm',
model_path: document.getElementById('s-vllm-model-path').value.trim(),
ctx: parseInt(document.getElementById('s-vllm-ctx').value) || 32768,
gpu_memory_utilization: parseFloat(document.getElementById('s-vllm-gmu').value) || 0.90,
tensor_parallel_size: parseInt(document.getElementById('s-vllm-tp').value) || 1,
max_num_seqs: parseInt(document.getElementById('s-vllm-mns').value) || 0,
dtype: document.getElementById('s-vllm-dtype').value.trim(),
quantization: document.getElementById('s-vllm-quant').value.trim(),
port: parseInt(document.getElementById('s-vllm-port').value) || 0,
venv: document.getElementById('s-vllm-venv').value.trim(),
extra_args: document.getElementById('s-vllm-extra-args').value.trim(),
extra_env: document.getElementById('s-vllm-extra-env').value.trim(),
auto_build: document.getElementById('s-vllm-auto-build').checked,
},
ocr:{
enabled: document.getElementById('s-ocr-enabled').checked,
default_engine: document.getElementById('s-ocr-default-engine').value || 'paddle',
......
# CoderAI - OpenAI-compatible API server
# Copyright (C) 2026 Stefy Lanza <stefy@nexlab.net>
#
# This program is free software: you can redistribute it and/or modify
# it under the terms of the GNU General Public License as published by
# the Free Software Foundation, either version 3 of the License, or
# (at your option) any later version.
"""Fully-managed vLLM worker — the vLLM OpenAI server, driven as a subprocess.
vLLM (https://github.com/vllm-project/vllm) exposes an OpenAI-compatible HTTP server
(``python -m vllm.entrypoints.openai.api_server``) with continuous batching + paged KV,
so — like :mod:`codai.api.kt_worker` — coderai launches it as a managed subprocess,
health-checks ``/v1/models``, and :mod:`codai.backends.vllm` proxies to it.
vLLM pins its own torch/CUDA (e.g. torch 2.13 / cu13), which conflicts with the main
coderai venv, so it runs in an ISOLATED venv (``vllm.venv``) — same pattern as the OCR
Paddle/Surya engines. This module owns the venv + process lifecycle only. Also reused by
the OCR subsystem to serve Surya2 (a VLM) for :mod:`codai.ocr.surya`.
"""
import collections
import os
import shlex
import socket
import subprocess
import threading
import time
from typing import Optional
_lock = threading.RLock()
_services: dict[str, dict] = {} # svc_key -> {"proc","port","url"}
# Repo-root requirements for the isolated vLLM venv (auto-build).
_REQ = os.path.join(os.path.dirname(__file__), "..", "..", "requirements-vllm.txt")
def resolve_venv_dir(cfg) -> str:
"""Isolated vLLM venv: config > baked /opt/coderai/vllm_venv > /cache mount > ~/.coderai."""
configured = (getattr(cfg, "venv", "") or "").strip()
if configured:
return configured
baked = "/opt/coderai/vllm_venv"
if os.path.isdir(baked):
return baked
cache = os.environ.get("CODERAI_CACHE_DIR") or ("/cache" if os.path.isdir("/cache") else "")
if cache and os.path.isdir(cache):
return os.path.join(cache, "vllm_venv")
return os.path.expanduser("~/.coderai/vllm_venv")
def _venv_python(cfg) -> str:
return os.path.join(os.path.expanduser(resolve_venv_dir(cfg)), "bin", "python")
def _free_port() -> int:
s = socket.socket(socket.AF_INET, socket.SOCK_STREAM)
s.bind(("127.0.0.1", 0))
port = s.getsockname()[1]
s.close()
return port
def _pump_logs(proc, tail):
for line in proc.stdout:
line = line.rstrip()
if line:
tail.append(line)
print(f"[vllm] {line}", flush=True)
def _health_ok(url: str) -> bool:
import requests
try:
r = requests.get(url + "/v1/models", timeout=3)
return r.status_code == 200
except Exception:
return False
def ensure_built(cfg) -> str:
"""Ensure the isolated vLLM venv exists; return its python. Builds it if auto_build."""
py = _venv_python(cfg)
if os.path.isfile(py):
return py
if not getattr(cfg, "auto_build", False):
raise RuntimeError(
f"vLLM isolated venv not found at {os.path.dirname(os.path.dirname(py))}. "
f"Build it (python3 -m venv <dir> && <dir>/bin/pip install -r "
f"requirements-vllm.txt) or enable vllm.auto_build.")
venv_dir = os.path.expanduser(resolve_venv_dir(cfg))
req = os.path.abspath(_REQ)
print(f"[vllm] creating isolated venv at {venv_dir} …", flush=True)
import sys as _sys
try:
os.makedirs(os.path.dirname(venv_dir) or ".", exist_ok=True)
subprocess.run([_sys.executable, "-m", "venv", venv_dir], check=True)
subprocess.run([py, "-m", "pip", "install", "-U", "pip"], check=True)
if os.path.isfile(req):
subprocess.run([py, "-m", "pip", "install", "-r", req], check=True)
else:
subprocess.run([py, "-m", "pip", "install", "vllm"], check=True)
except subprocess.CalledProcessError as exc:
raise RuntimeError(f"vLLM: venv build failed: {exc}")
if not os.path.isfile(py):
raise RuntimeError("vLLM: venv python missing after build")
return py
def resolve_service_key(cfg, model_path: Optional[str] = None):
mp = os.path.expanduser((model_path or getattr(cfg, "model_path", "") or "").strip())
key = mp or (getattr(cfg, "model_id", "vllm") or "vllm")
return mp, key
def _launch_cmd(py, cfg, host: str, port: int, model_path: str,
served_name: Optional[str] = None) -> list:
mid = served_name or (getattr(cfg, "model_id", "vllm") or "vllm")
cmd = [py, "-m", "vllm.entrypoints.openai.api_server",
"--host", host, "--port", str(port),
"--model", model_path,
"--served-model-name", mid]
ctx = int(getattr(cfg, "ctx", 0) or 0)
if ctx > 0:
cmd += ["--max-model-len", str(ctx)]
gmu = float(getattr(cfg, "gpu_memory_utilization", 0) or 0)
if gmu > 0:
cmd += ["--gpu-memory-utilization", str(gmu)]
tp = int(getattr(cfg, "tensor_parallel_size", 0) or 0)
if tp > 0:
cmd += ["--tensor-parallel-size", str(tp)]
mns = int(getattr(cfg, "max_num_seqs", 0) or 0)
if mns > 0:
cmd += ["--max-num-seqs", str(mns)]
dtype = (getattr(cfg, "dtype", "") or "").strip()
if dtype:
cmd += ["--dtype", dtype]
quant = (getattr(cfg, "quantization", "") or "").strip()
if quant:
cmd += ["--quantization", quant]
extra = (getattr(cfg, "extra_args", "") or "").strip()
if extra:
cmd += shlex.split(extra)
return cmd
def ensure_service(cfg, model_path: Optional[str] = None,
served_name: Optional[str] = None,
ready_timeout: float = 3600.0) -> str:
"""Launch (or reuse) a vLLM OpenAI server for a model; return its base URL.
``model_path``/``served_name`` override the config (used by the OCR subsystem to serve
surya-2 on its own vLLM instance alongside any LLM instance)."""
resolved, svc_key = resolve_service_key(cfg, model_path)
if served_name:
svc_key = f"{svc_key}|{served_name}"
with _lock:
svc = _services.get(svc_key)
if svc and svc["proc"].poll() is None and _health_ok(svc["url"]):
return svc["url"]
if svc:
_services.pop(svc_key, None)
py = ensure_built(cfg)
model = resolved or (getattr(cfg, "model_path", "") or "").strip()
if not model:
raise RuntimeError(
"vLLM: no model resolved. Set vllm.model_path (an HF model dir or id).")
host = (getattr(cfg, "host", "127.0.0.1") or "127.0.0.1").strip()
port = int(getattr(cfg, "port", 0) or 0) or _free_port()
url_host = "127.0.0.1" if host in ("0.0.0.0", "") else host
url = f"http://{url_host}:{port}"
cmd = _launch_cmd(py, cfg, host, port, model, served_name)
env = os.environ.copy()
# vLLM's venv bundles its own torch/CUDA; make sure its libs win.
vlib = os.path.join(os.path.expanduser(resolve_venv_dir(cfg)),
"lib", "python3.13", "site-packages", "nvidia")
extra_env = (getattr(cfg, "extra_env", "") or "").strip()
applied = {}
if extra_env:
for tok in shlex.split(extra_env):
if "=" in tok:
k, v = tok.split("=", 1)
if k.strip():
env[k.strip()] = v; applied[k.strip()] = v
print(f"[vllm] launching: {' '.join(cmd)}"
+ (f" ({' '.join(f'{k}={v}' for k, v in applied.items())})" if applied else ""),
flush=True)
tail = collections.deque(maxlen=80)
proc = subprocess.Popen(cmd, stdout=subprocess.PIPE, stderr=subprocess.STDOUT,
text=True, bufsize=1, env=env)
threading.Thread(target=_pump_logs, args=(proc, tail), daemon=True).start()
_services[svc_key] = {"proc": proc, "port": port, "url": url}
def _tail_msg():
joined = " | ".join(list(tail)[-6:]).strip()
return f". Last output: {joined}" if joined else ""
deadline = time.time() + ready_timeout
while time.time() < deadline:
if proc.poll() is not None:
stop_service(svc_key)
raise RuntimeError(
f"vLLM exited (code {proc.returncode}) before becoming ready" + _tail_msg())
if _health_ok(url):
print(f"[vllm] service ready for {svc_key} at {url}", flush=True)
return url
time.sleep(2)
stop_service(svc_key)
raise RuntimeError(f"vLLM for {svc_key} did not become ready in time" + _tail_msg())
def stop_service(svc_key: str) -> None:
with _lock:
svc = _services.pop(svc_key, None)
if not svc:
return
proc = svc["proc"]
if proc.poll() is None:
try:
proc.terminate(); proc.wait(timeout=10)
except Exception:
pass
if proc.poll() is None:
try:
proc.kill()
except Exception:
pass
print(f"[vllm] service for {svc_key} stopped", flush=True)
def stop_all() -> None:
for k in list(_services.keys()):
stop_service(k)
import atexit as _atexit
_atexit.register(stop_all)
This diff is collapsed.
......@@ -514,6 +514,39 @@ class KtransformersConfig:
auto_build: bool = False # pip-install SGLang+kt-kernel if missing (heavy; off by default)
@dataclass
class VllmConfig:
"""vLLM high-concurrency backend — the vLLM OpenAI server, driven as a subprocess.
vLLM (https://github.com/vllm-project/vllm) provides continuous batching + paged KV for
far higher aggregate throughput than serialized single-instance backends. Like ds4/kt,
coderai launches ``python -m vllm.entrypoints.openai.api_server`` as a managed
subprocess and proxies ``/v1/chat/completions`` to it (:mod:`codai.backends.vllm`).
vLLM pins its own torch/CUDA (e.g. torch 2.13 / cu13), which conflicts with the main
coderai venv, so it runs in an ISOLATED venv (``venv``; blank → baked
/opt/coderai/vllm_venv, else the /cache mount, else ~/.coderai/vllm_venv), built from
requirements-vllm.txt. Selected PER MODEL via a ``backend: "vllm"`` pin or the
``model_id`` alias — never by a broad name marker (it would collide with every engine).
Also reused by the OCR subsystem to serve Surya2 (a VLM) with continuous batching.
"""
enabled: bool = False
venv: str = "" # isolated venv dir; blank = auto (baked/cache/home)
model_path: str = "" # HF model dir or repo id (--model)
model_id: str = "vllm" # id/alias that routes to vllm (and --served-model-name)
host: str = "127.0.0.1"
port: int = 0 # 0 = auto-pick a free port
ctx: int = 32768 # --max-model-len
gpu_memory_utilization: float = 0.90 # --gpu-memory-utilization
tensor_parallel_size: int = 1 # --tensor-parallel-size
max_num_seqs: int = 0 # --max-num-seqs (0 = vLLM default)
dtype: str = "" # --dtype (blank = auto; e.g. bfloat16/float16)
quantization: str = "" # --quantization (blank = none; e.g. awq/gptq/fp8)
extra_args: str = "" # extra flags for the api_server
extra_env: str = "" # free-form KEY=VALUE env for the subprocess
auto_build: bool = False # create the isolated venv + pip install vllm if missing
@dataclass
class OcrConfig:
"""Dedicated OCR subsystem configuration.
......@@ -570,6 +603,14 @@ class OcrConfig:
surya_langs: str = "it"
surya_venv: str = "" # isolated venv dir; blank = ~/.coderai/surya_venv
surya_auto_build: bool = False # create the venv + pip install requirements-surya.txt on first use
# Surya serving mode: "local" = classic det+recognition on torch in the isolated venv
# (surya-ocr <=0.17); "vllm"/"llamacpp" = the latest "Surya2" VLM served by an external
# OpenAI server that Surya attaches to (SURYA_INFERENCE_URL). "vllm" reuses coderai's
# vLLM backend to serve `surya_model` (continuous batching); "llamacpp" attaches to a
# llama.cpp server at `surya_server_url`.
surya_serve: str = "local" # local | vllm | llamacpp
surya_model: str = "datalab-to/surya-ocr-2" # HF checkpoint for the served (vllm) backend
surya_server_url: str = "" # external OpenAI server URL (llamacpp/manual); blank = auto
# --- stamp / signature detection --- [O3]
detect_mode: str = "off" # off|layout|detector|both
......@@ -607,6 +648,7 @@ class Config:
colibri: ColibriConfig = field(default_factory=ColibriConfig)
k3: K3Config = field(default_factory=K3Config)
ktransformers: KtransformersConfig = field(default_factory=KtransformersConfig)
vllm: VllmConfig = field(default_factory=VllmConfig)
ocr: OcrConfig = field(default_factory=OcrConfig)
compaction: CompactionConfig = field(default_factory=CompactionConfig)
broker: BrokerConfig = field(default_factory=BrokerConfig)
......@@ -796,6 +838,7 @@ class ConfigManager:
colibri=_dc(ColibriConfig, config_data.get("colibri", {})),
k3=_dc(K3Config, config_data.get("k3", {})),
ktransformers=_dc(KtransformersConfig, config_data.get("ktransformers", {})),
vllm=_dc(VllmConfig, config_data.get("vllm", {})),
ocr=_dc(OcrConfig, config_data.get("ocr", {})),
compaction=_dc(CompactionConfig, config_data.get("compaction", {})),
broker=_dc(BrokerConfig, config_data.get("broker", {})),
......@@ -1026,6 +1069,23 @@ class ConfigManager:
"extra_env": self.config.ktransformers.extra_env,
"auto_build": self.config.ktransformers.auto_build,
},
"vllm": {
"enabled": self.config.vllm.enabled,
"venv": self.config.vllm.venv,
"model_path": self.config.vllm.model_path,
"model_id": self.config.vllm.model_id,
"host": self.config.vllm.host,
"port": self.config.vllm.port,
"ctx": self.config.vllm.ctx,
"gpu_memory_utilization": self.config.vllm.gpu_memory_utilization,
"tensor_parallel_size": self.config.vllm.tensor_parallel_size,
"max_num_seqs": self.config.vllm.max_num_seqs,
"dtype": self.config.vllm.dtype,
"quantization": self.config.vllm.quantization,
"extra_args": self.config.vllm.extra_args,
"extra_env": self.config.vllm.extra_env,
"auto_build": self.config.vllm.auto_build,
},
"ocr": {
"enabled": self.config.ocr.enabled,
"default_engine": self.config.ocr.default_engine,
......@@ -1052,6 +1112,9 @@ class ConfigManager:
"surya_langs": self.config.ocr.surya_langs,
"surya_venv": self.config.ocr.surya_venv,
"surya_auto_build": self.config.ocr.surya_auto_build,
"surya_serve": self.config.ocr.surya_serve,
"surya_model": self.config.ocr.surya_model,
"surya_server_url": self.config.ocr.surya_server_url,
"detect_mode": self.config.ocr.detect_mode,
"detect_model_path": self.config.ocr.detect_model_path,
"detect_conf": self.config.ocr.detect_conf,
......
......@@ -908,7 +908,7 @@ class FrontProxy:
new = cm.config
for f in ("server", "backend", "models", "offload", "vulkan", "image",
"whisper", "archive", "thermal", "jobs", "enhance", "ds4", "colibri",
"k3", "ktransformers",
"k3", "ktransformers", "vllm",
"compaction", "broker", "system_prompt", "tools_closer_prompt",
"grammar_guided", "parser", "tmp_dir"):
if hasattr(new, f):
......@@ -922,6 +922,7 @@ class FrontProxy:
colibri = getattr(self.config, "colibri", None)
k3 = getattr(self.config, "k3", None)
kt = getattr(self.config, "ktransformers", None)
vllm = getattr(self.config, "vllm", None)
info = self._model_info(model)
cap = _router.required_capability(
model, path=path,
......@@ -933,7 +934,9 @@ class FrontProxy:
k3_model_id=getattr(k3, "model_id", None) if k3 else None,
k3_enabled=bool(getattr(k3, "enabled", False)) if k3 else False,
kt_model_id=getattr(kt, "model_id", None) if kt else None,
kt_enabled=bool(getattr(kt, "enabled", False)) if kt else False)
kt_enabled=bool(getattr(kt, "enabled", False)) if kt else False,
vllm_model_id=getattr(vllm, "model_id", None) if vllm else None,
vllm_enabled=bool(getattr(vllm, "enabled", False)) if vllm else False)
# The name heuristic can't see that a bare alias (e.g. '…-q4_k_m', no
# literal 'gguf') backs a .gguf file, so it falls through to
# 'transformers' (CUDA-only) and the request never reaches a Vulkan/AMD
......
......@@ -54,7 +54,7 @@ def _route_key(entry):
return None
def _required_cap(entry, ds4_cfg, colibri_cfg=None, k3_cfg=None, kt_cfg=None):
def _required_cap(entry, ds4_cfg, colibri_cfg=None, k3_cfg=None, kt_cfg=None, vllm_cfg=None):
from codai.frontproxy.router import required_capability
path = _entry_path(entry) or ""
backend = entry.get("backend") if isinstance(entry, dict) else None
......@@ -67,11 +67,13 @@ def _required_cap(entry, ds4_cfg, colibri_cfg=None, k3_cfg=None, kt_cfg=None):
k3_model_id=getattr(k3_cfg, "model_id", None) if k3_cfg else None,
k3_enabled=bool(getattr(k3_cfg, "enabled", False)) if k3_cfg else False,
kt_model_id=getattr(kt_cfg, "model_id", None) if kt_cfg else None,
kt_enabled=bool(getattr(kt_cfg, "enabled", False)) if kt_cfg else False)
kt_enabled=bool(getattr(kt_cfg, "enabled", False)) if kt_cfg else False,
vllm_model_id=getattr(vllm_cfg, "model_id", None) if vllm_cfg else None,
vllm_enabled=bool(getattr(vllm_cfg, "enabled", False)) if vllm_cfg else False)
def compute_assignment(engines, models_path, default_engine=None, ds4_cfg=None,
colibri_cfg=None, k3_cfg=None, kt_cfg=None):
colibri_cfg=None, k3_cfg=None, kt_cfg=None, vllm_cfg=None):
"""Return {engine_name: [model_identifiers]} — each model owned by one engine."""
assignment = {e.name: [] for e in engines}
if not engines or not models_path:
......@@ -91,7 +93,7 @@ def compute_assignment(engines, models_path, default_engine=None, ds4_cfg=None,
ident = _route_key(entry)
if not ident or ident in seen:
continue
cap = _required_cap(entry, ds4_cfg, colibri_cfg, k3_cfg, kt_cfg)
cap = _required_cap(entry, ds4_cfg, colibri_cfg, k3_cfg, kt_cfg, vllm_cfg)
candidates = [e for e in engines if e.can_serve(cap)]
if not candidates:
continue # nothing can run it — leave unassigned
......
......@@ -158,8 +158,9 @@ class EngineSupervisor:
colibri = getattr(self.config, "colibri", None)
k3 = getattr(self.config, "k3", None)
kt = getattr(self.config, "ktransformers", None)
vllm = getattr(self.config, "vllm", None)
assignment = compute_assignment(engines, self.models_path,
default_engine, ds4, colibri, k3, kt)
default_engine, ds4, colibri, k3, kt, vllm)
for e in engines:
owned = assignment.get(e.name, [])
e.assigned_models = set(owned) # the front's router enforces this
......@@ -608,8 +609,9 @@ class EngineSupervisor:
colibri = getattr(self.config, "colibri", None)
k3 = getattr(self.config, "k3", None)
kt = getattr(self.config, "ktransformers", None)
vllm = getattr(self.config, "vllm", None)
assignment = compute_assignment(real, self.models_path,
default_engine, ds4, colibri, k3, kt)
default_engine, ds4, colibri, k3, kt, vllm)
except Exception as exc:
print(f"[front] live reassignment skipped: {exc}", flush=True)
assignment = {}
......
......@@ -44,11 +44,11 @@ def _short_stem(key: str) -> str:
# k3 (kimi-k3-in-c) is a CPU engine → any node can host it; kt (ktransformers/SGLang)
# is CPU+GPU → the GPU-capable nodes.
_DEFAULT_CAPS = {
"nvidia": {"transformers", "gguf", "whisper", "ds4", "colibri", "k3", "kt"},
"cuda": {"transformers", "gguf", "whisper", "ds4", "colibri", "k3", "kt"},
"nvidia": {"transformers", "gguf", "whisper", "ds4", "colibri", "k3", "kt", "vllm"},
"cuda": {"transformers", "gguf", "whisper", "ds4", "colibri", "k3", "kt", "vllm"},
"vulkan": {"gguf", "whisper", "k3"},
"opencl": {"gguf", "whisper", "k3"},
"auto": {"transformers", "gguf", "whisper", "ds4", "colibri", "k3", "kt"},
"auto": {"transformers", "gguf", "whisper", "ds4", "colibri", "k3", "kt", "vllm"},
}
......
......@@ -80,7 +80,9 @@ def required_capability(model: Optional[str], path: Optional[str] = None,
k3_model_id: Optional[str] = None,
k3_enabled: bool = False,
kt_model_id: Optional[str] = None,
kt_enabled: bool = False) -> Optional[str]:
kt_enabled: bool = False,
vllm_model_id: Optional[str] = None,
vllm_enabled: bool = False) -> Optional[str]:
"""The capability an engine must have to serve this request.
* ``whisper`` — whisper.cpp STT (transcription endpoint or a
......@@ -101,7 +103,7 @@ def required_capability(model: Optional[str], path: Optional[str] = None,
return "whisper"
# 1. Explicit engine-backend pin wins (authoritative), like the manager resolver.
b = (backend or "").lower()
if b in ("colibri", "ds4", "k3", "kt"):
if b in ("colibri", "ds4", "k3", "kt", "vllm"):
return b
m = (model or "").lower()
......@@ -110,6 +112,8 @@ def required_capability(model: Optional[str], path: Optional[str] = None,
return bool(mid) and (m == mid or m.split("/")[-1] == mid)
# 2. An enabled engine's model_id alias.
if vllm_enabled and _alias(vllm_model_id):
return "vllm"
if kt_enabled and _alias(kt_model_id):
return "kt"
if k3_enabled and _alias(k3_model_id):
......
......@@ -80,6 +80,17 @@ def get_active_ktransformers_config():
return None
def get_active_vllm_config():
"""Return the active VllmConfig from the server config, or None."""
try:
from codai.admin.routes import config_manager
if config_manager is not None and config_manager.config is not None:
return config_manager.config.vllm
except Exception:
pass
return None
_GGUF_ARCH_CACHE: Dict[tuple, str] = {}
......@@ -202,7 +213,7 @@ def _gguf_architecture(path: str):
# --------------------------------------------------------------------------- #
# Ordered by arbitration preference for ambiguous names. New engines append here.
_ENGINE_BACKENDS = ("ds4", "colibri", "k3", "kt")
_ENGINE_BACKENDS = ("ds4", "colibri", "k3", "kt", "vllm")
def _engine_config(engine: str):
......@@ -212,6 +223,7 @@ def _engine_config(engine: str):
"colibri": get_active_colibri_config,
"k3": get_active_k3_config,
"kt": get_active_ktransformers_config,
"vllm": get_active_vllm_config,
}.get(engine)
return getter() if getter else None
......@@ -318,11 +330,18 @@ def _kt_name_claims(model_name: str) -> bool:
return False
def _vllm_name_claims(model_name: str) -> bool:
"""vLLM serves any HF/GGUF model, so it must NEVER auto-claim by a broad name marker.
Routes ONLY via an explicit ``backend: "vllm"`` pin or its configured ``model_id``."""
return False
_ENGINE_NAME_CLAIMS = {
"ds4": _ds4_name_claims,
"colibri": _colibri_name_claims,
"k3": _k3_name_claims,
"kt": _kt_name_claims,
"vllm": _vllm_name_claims,
}
......@@ -382,6 +401,11 @@ def kt_should_handle(model_name: str) -> bool:
return resolve_engine_backend(model_name) == "kt"
def vllm_should_handle(model_name: str) -> bool:
"""True when vllm is the resolved engine backend for ``model_name``."""
return resolve_engine_backend(model_name) == "vllm"
def _trim_cpu_ram() -> None:
"""Return freed CPU heap memory to the OS (and let the kernel reclaim swap).
......@@ -554,6 +578,17 @@ class ModelManager:
self.tool_parser = ModelParserAdapter(model_name=model_name)
return
# vLLM: when enabled, proxy matching models to the managed vLLM OpenAI server
# (isolated venv, continuous batching). Pin/alias-selected only.
if vllm_should_handle(model_name):
from codai.backends.vllm import VllmBackend
print(f"Routing '{model_name}' to vLLM backend")
self.backend_type = "vllm"
self.backend = VllmBackend(get_active_vllm_config())
self.backend.load_model(model_name, **kwargs)
self.tool_parser = ModelParserAdapter(model_name=model_name)
return
available = detect_available_backends()
# Check if model is a GGUF file. The name alone isn't reliable: a gguf's
......@@ -2217,6 +2252,11 @@ class MultiModelManager:
if model_type in (None, "text") and kt_should_handle(requested_or_resolved):
return True
# vLLM served models: accept for text when vllm is the resolver's pick
# (explicit backend:vllm pin or the vllm model_id alias).
if model_type in (None, "text") and vllm_should_handle(requested_or_resolved):
return True
# If a model_type is specified, reject models registered under a
# different type (e.g. an image GGUF requested via /v1/chat/completions).
if model_type:
......@@ -5126,6 +5166,11 @@ class MultiModelManager:
mid = getattr(kt_cfg, "model_id", "ktransformers") or "ktransformers"
_add(mid, "text", {"backend": "kt"})
vllm_cfg = get_active_vllm_config()
if vllm_cfg is not None and getattr(vllm_cfg, "enabled", False):
mid = getattr(vllm_cfg, "model_id", "vllm") or "vllm"
_add(mid, "text", {"backend": "vllm"})
return models
......
......@@ -75,6 +75,10 @@ class SubprocessOcrEngine(OcrEngine):
def _worker_opts(self) -> dict:
return {}
def _worker_env(self) -> dict:
"""Extra environment variables for the worker subprocess (overridable)."""
return {}
def _pip_install_cmd(self, py: str, req: str):
"""Return the pip command list to build the venv. Overridable (e.g. paddle index)."""
return [py, "-m", "pip", "install", "-r", req]
......@@ -118,6 +122,7 @@ class SubprocessOcrEngine(OcrEngine):
def load(self) -> None:
py = self._ensure_venv()
env = dict(os.environ)
env.update(self._worker_env())
try:
self._proc = subprocess.Popen(
[py, _WORKER, self._worker_engine],
......
......@@ -36,10 +36,44 @@ class SuryaEngine(SubprocessOcrEngine):
def _worker_opts(self) -> dict:
return {"langs": self.cfg.surya_langs or "it"}
def _serve_mode(self) -> str:
return (getattr(self.cfg, "surya_serve", "local") or "local").strip().lower()
def _worker_env(self) -> dict:
"""For the served ("Surya2") modes, tell the worker to attach to an external
OpenAI server instead of running detection+recognition locally."""
mode = self._serve_mode()
if mode not in ("vllm", "llamacpp"):
return {}
url = self._server_url
if not url:
return {}
return {"SURYA_INFERENCE_BACKEND": mode, "SURYA_INFERENCE_URL": url}
def load(self) -> None:
if not self.cfg.surya_accept_license:
raise OcrError(
"Surya is license-gated. Set ocr.surya_accept_license = true to use it.",
status=400,
)
self._server_url = ""
mode = self._serve_mode()
if mode == "vllm":
# Serve the Surya2 VLM checkpoint via coderai's vLLM backend (continuous
# batching) and attach Surya to it. Requires vllm enabled/auto-buildable.
from codai.api import vllm_worker
from codai.models.manager import get_active_vllm_config
vcfg = get_active_vllm_config()
if vcfg is None:
raise OcrError("Surya vllm mode needs the vLLM backend configured", status=400)
model = (getattr(self.cfg, "surya_model", "") or "datalab-to/surya-ocr-2").strip()
base = vllm_worker.ensure_service(vcfg, model_path=model, served_name=model)
self._server_url = base.rstrip("/") + "/v1"
elif mode == "llamacpp":
url = (getattr(self.cfg, "surya_server_url", "") or "").strip()
if not url:
raise OcrError(
"Surya llamacpp mode needs ocr.surya_server_url (a running llama-server "
"OpenAI endpoint).", status=400)
self._server_url = url
super().load()
......@@ -233,15 +233,27 @@ class SuryaWorker:
self._langs = [l.strip() for l in str(opts.get("langs", "it")).split(",") if l.strip()] or ["it"]
def load(self):
from surya.detection import DetectionPredictor
# 0.17.x: RecognitionPredictor(FoundationPredictor()) — a shared VLM foundation.
try:
from surya.foundation import FoundationPredictor
from surya.recognition import RecognitionPredictor
self._rec = RecognitionPredictor(FoundationPredictor())
self._det = DetectionPredictor()
self._mode = "predictor"
return
except Exception:
pass
# 0.6–0.16: RecognitionPredictor() with no args.
try:
from surya.recognition import RecognitionPredictor
from surya.detection import DetectionPredictor
self._rec = RecognitionPredictor()
self._det = DetectionPredictor()
self._mode = "predictor"
return
except Exception:
pass
# very old: functional run_ocr API.
from surya.ocr import run_ocr
from surya.model.detection.model import (
load_model as ldm, load_processor as ldp)
......@@ -252,10 +264,29 @@ class SuryaWorker:
self._rec = lrm(); self._rec_proc = lrp()
self._mode = "run_ocr"
def _predict(self, img):
# Predictor API drifted across classic surya versions:
# 0.17.x: rec(images, task_names, det_predictor) (task_names defaults internally)
# 0.6.x : rec(images, langs, det_predictor)
attempts = [
lambda: self._rec([img]), # 0.20+ (VLM full-page, backend does it)
lambda: self._rec([img], full_page=True), # 0.20+ variant
lambda: self._rec([img], det_predictor=self._det), # 0.17.x (task_names default)
lambda: self._rec([img], ["ocr_with_boxes"], self._det), # 0.17.x explicit
lambda: self._rec([img], [self._langs], self._det), # 0.6.x (langs)
]
last = None
for call in attempts:
try:
return call()
except (TypeError, ValueError, AssertionError) as e:
last = e
raise last
def ocr(self, img):
img = img.convert("RGB")
if self._mode == "predictor":
preds = self._rec([img], [self._langs], self._det)
preds = self._predict(img)
else:
preds = self._run_ocr([img], [self._langs], self._det, self._det_proc, self._rec, self._rec_proc)
res = preds[0] if preds else None
......
# vLLM backend (future option)
# vLLM backend
> **Status: NOT YET IMPLEMENTED — design note / decision record.**
> For now, parallel throughput is achieved with **multiple llama.cpp instances**
> of the same model (see the manager's multi-instance loading). This document
> captures why and how we would add a vLLM backend later.
> **Status: IMPLEMENTED (v0.1.86).** vLLM is a first-class managed engine backend,
> wired exactly like ds4/ktransformers: a subprocess in an ISOLATED venv, proxied over
> its OpenAI HTTP API, selected per-model via a `backend: "vllm"` pin or the `vllm.model_id`
> alias (never auto-claimed). Config: `codai/config.py::VllmConfig`; worker:
> `codai/api/vllm_worker.py`; proxy: `codai/backends/vllm.py`; manager wiring +
> front-proxy capability (`vllm` on GPU nodes) + admin card. It also serves **Surya2**
> for the OCR subsystem (`ocr.surya_serve = "vllm"`).
>
> Engine layering: nvidia/cuda/vulkan/opencl are **in-process base engines**; vLLM (like
> ds4/colibri/k3/kt) is a **managed external engine** — it must be out-of-process because
> it pins its own torch/CUDA (torch 2.13 / cu13) that conflicts with the main venv.
## Why add vLLM
......
......@@ -13,9 +13,11 @@
# with coderai's GPLv3); some releases add a revenue-capped commercial clause — check the
# version you install for commercial use.
# ---------------------------------------------------------------------------
# PIN to a classic (pre-"Surya2") version. surya-ocr >=~0.15 rearchitected to a VLM that
# requires an EXTERNAL vLLM-in-Docker or llama-server backend (SpawnError without them) —
# unusable inside coderai's isolated venv. 0.6.4 runs det+recognition locally on torch,
# matching the RecognitionPredictor(images, langs, det_predictor) API the worker uses.
surya-ocr==0.6.4
# PIN to the NEWEST classic (pre-"Surya2") version. surya-ocr >=0.20 rearchitected to a
# VLM that needs an EXTERNAL backend — vLLM (defaults to a vllm-openai Docker image),
# llama-server (LLAMA_CPP_BINARY), or any OpenAI endpoint (openai_client). That is the
# right long-term path (serve Surya2 via coderai's planned vLLM/llama.cpp backend), but it
# does NOT run standalone in this isolated venv. 0.17.1 is the last version that runs
# det+recognition locally on torch (GPU), so it's the classic pin. See docs/ocr.md.
surya-ocr==0.17.1
pypdfium2>=4.20.0
# ---------------------------------------------------------------------------
# CoderAI — vLLM high-concurrency backend, ISOLATED venv.
#
# vLLM pins its own torch/CUDA (e.g. torch==2.13.0 / cu13) which conflicts with the main
# coderai venv, so it runs in a dedicated virtualenv + managed subprocess (same isolation
# pattern as the OCR Paddle/Surya engines and ktransformers/SGLang). coderai launches
# `python -m vllm.entrypoints.openai.api_server` from this venv and proxies to it
# (codai/backends/vllm.py). Also serves Surya2 (a VLM) for the OCR subsystem.
#
# python3 -m venv <dir> && <dir>/bin/pip install -r requirements-vllm.txt
# (or set vllm.auto_build = true to have coderai build it on first use)
#
# vLLM's wheel pulls its exact torch/torchvision/torchaudio + CUDA runtime, so do NOT pin
# torch here — let vLLM choose the matching versions. Pin only vLLM itself.
# ---------------------------------------------------------------------------
vllm==0.27.1
Markdown is supported
0% or
You are about to add 0 people to the discussion. Proceed with caution.
Finish editing this message first!
Please register or to comment