diff --git a/Makefile b/Makefile index 63fe6aa..bf6e6b4 100644 --- a/Makefile +++ b/Makefile @@ -6,10 +6,11 @@ PY := $(VENV)/bin/python PYTEST := $(PY) -m pytest COMPOSE := docker compose -f deploy/docker-compose.yaml +COMPOSE_MIXED := docker compose -f deploy/docker-compose.mixed.yaml .PHONY: help test test-backend test-router test-schema \ dev-backend dev-frontend build-frontend install-frontend \ - up down logs ps build + up down logs ps build up-mixed down-mixed logs-mixed help: @echo "Targets:" @@ -26,6 +27,10 @@ help: @echo " logs Tail logs from all services" @echo " ps Show service status" @echo " build Build images without starting" + @echo " --- mixed-engine HA (deploy/docker-compose.mixed.yaml) ---" + @echo " up-mixed vLLM + SGLang backends sharing one Postgres/router/dashboard" + @echo " down-mixed Stop + remove the mixed stack" + @echo " logs-mixed Tail logs from the mixed stack" up: $(COMPOSE) up -d --build @@ -42,6 +47,16 @@ ps: build: $(COMPOSE) build +# --- mixed-engine HA deployment (vLLM + SGLang backends, shared Postgres) --- +up-mixed: + $(COMPOSE_MIXED) up -d --build + +down-mixed: + $(COMPOSE_MIXED) down + +logs-mixed: + $(COMPOSE_MIXED) logs -f + test: test-backend test-router test-schema test-backend: diff --git a/README.md b/README.md index 9d3355c..bb62cb9 100644 --- a/README.md +++ b/README.md @@ -29,6 +29,7 @@ becomes a routable model; the router load-balances across instances; and a bundl ## Highlights - **One router controls the whole fleet** — a single OpenAI- & Anthropic-compatible origin fronts every model. Route by the `model` field across `/v1/chat/completions`, `/v1/messages`, `/v1/embeddings`, `/v1/rerank`, `/v1/score`, `/tokenize` and more; the router resolves the group and load-balances its instances, so clients never address an instance directly. +- **Two inference engines, one control plane — vLLM + SGLang** — choose the engine per model (the *Add Model* dialog has an engine selector); an engine-aware scheduler places each model on a backend that can run it, and the same router / dashboard / monitoring front both. Run a **vLLM-only** stack (`make up`) or a **mixed vLLM + SGLang** fleet (`make up-mixed`). See [docs/mixed-engine-deployment.md](docs/mixed-engine-deployment.md). - **Add a model by pasting `vllm serve …`** — parsed into a form and layered on as a dynamic overlay; the router hot-reloads, no `config.yaml` edits. - **Lifecycle + self-healing** — per-instance state machine (`stopped → starting → ready → sleeping → failed`), VRAM pre-flight guard, GPU auto-placement, crash auto-restart with backoff. - **Autoscaling with a warm-standby tier** — per group, keep `min_ready` replicas warm and scale up on queue depth (wake first, else cold-start) to `max_ready`; fold idle replicas back down `ready → sleep → stop`. vLLM **sleep mode** (level-1) frees a replica's VRAM but wakes in seconds, so scaling down needn't mean a minute-long cold start. Set it from config.yaml or the dashboard; a live Grafana dashboard + alerts are bundled. @@ -59,6 +60,11 @@ make up # build + start the whole stack `make down` stops it · `make logs` tails all services · `make ps` shows status. +**Two deployment modes:** + +- **`make up`** — the default **vLLM-only** stack. +- **`make up-mixed`** — a **vLLM + SGLang** fleet: a vLLM backend and a SGLang backend sharing one Postgres, router, dashboard and Grafana. Add a SGLang model from *Add Model → engine: `sglang`* and it is auto-placed on the SGLang backend. `make down-mixed` / `make logs-mixed` manage it. See [docs/mixed-engine-deployment.md](docs/mixed-engine-deployment.md). + ```bash curl http://localhost:8887/v1/models # router: configured model groups curl http://localhost:5000/api/models # backend: lifecycle state of each instance @@ -100,6 +106,8 @@ Request/response shapes and auth details are in [docs/API.md](docs/API.md). ## Architecture +### vLLM-only (`make up`) + ```mermaid flowchart LR Client([Clients]) @@ -131,11 +139,56 @@ The **router only routes** — the **backend owns model lifecycle**. The fronten backend, and Grafana sit behind nginx on a single origin; backend, router, and Prometheus share one network namespace so the spawned vLLM instances are reachable on `localhost`. +### Mixed vLLM + SGLang (`make up-mixed`) + +Each engine runs as its own backend container (they can't share a netns), sharing one +Postgres (scheduling / desired intent), one router, one dashboard and one monitoring stack. +Each backend publishes its ready instances as **routable addresses** to a shared file_sd that +Prometheus scrapes. See [docs/mixed-engine-deployment.md](docs/mixed-engine-deployment.md). + +```mermaid +flowchart LR + Client([Clients]) + FE["frontend
nginx · :8884"] + GF["grafana
/grafana"] + PG[("postgres
shared store · scheduling/desired")] + PR["prometheus · :9090
scrapes routable file_sd targets"] + RT["router · :8887
OpenAI-compatible LB"] + + subgraph vbe["vLLM backend (engine.Dockerfile)"] + BV["backend · :5071
NODE_ENGINES=vllm"] + VINS["vLLM instances"] + end + subgraph sbe["SGLang backend (engine-sglang.Dockerfile)"] + BS["backend · :5072
NODE_ENGINES=sglang"] + SINS["SGLang instances"] + end + + Client --> FE + FE -->|/api| BV + FE -->|/v1| RT + FE -->|/grafana| GF + BV -->|launch| VINS + BS -->|launch| SINS + RT -->|route| VINS + RT -->|route| SINS + BV <-->|leader/schedule| PG + BS <-->|converge desired| PG + PR -->|scrape| VINS + PR -->|scrape| SINS + GF -->|query| PR +``` + +The leader's **engine-aware scheduler** places each model on a backend that can run its engine; +a control action landing on the wrong node is deferred to the owning one. SGLang serves +OpenMetrics, so Prometheus stores its metrics as `sglang_*` (underscore) while vLLM keeps colons. + ## Documentation | Topic | | |---|---| | Deployment & topology | [docs/deployment.md](docs/deployment.md) | +| Mixed-engine (vLLM + SGLang) | [docs/mixed-engine-deployment.md](docs/mixed-engine-deployment.md) | | Configuration (`config.yaml`) | [docs/configuration.md](docs/configuration.md) | | Features in depth | [docs/features.md](docs/features.md) | | Monitoring (Prometheus + Grafana) | [docs/monitoring.md](docs/monitoring.md) | diff --git a/README_zh-CN.md b/README_zh-CN.md index 6eb6485..d405b7c 100644 --- a/README_zh-CN.md +++ b/README_zh-CN.md @@ -30,6 +30,7 @@ ## 功能亮點 - **一個 router 掌控整個集群** — 單一 OpenAI 與 Anthropic 相容入口統管所有模型。以 `model` 欄位路由 `/v1/chat/completions`、`/v1/messages`、`/v1/embeddings`、`/v1/rerank`、`/v1/score`、`/tokenize` 等端點;router 自動解析群組並在實例間負載平衡,客戶端永遠不直接連到單一實例。 +- **兩種推理引擎、同一個控制平面 — vLLM + SGLang** — 每顆模型可各自選引擎(*新增模型*對話框有引擎選擇器);engine-aware 排程器把每顆模型擺到「跑得動它」的 backend 上,並由同一個 router/控制台/監控統一前置。可只跑 **vLLM**(`make up`),或跑 **vLLM + SGLang 混合**集群(`make up-mixed`)。見 [docs/mixed-engine-deployment_zh-CN.md](docs/mixed-engine-deployment_zh-CN.md)。 - **貼上 `vllm serve …` 即可新增模型** — 解析成表單、以動態 overlay 疊加;router 熱重載。 - **生命週期** — 每實例狀態機(`stopped → starting → ready → sleeping → failed`)、VRAM 預檢防呆、GPU 自動擺放、崩潰指數退避自動重啟。 - **自動擴縮(含暖待命層)** — 每群組保留 `min_ready` 暖機副本,依佇列深度擴容(優先喚醒、其次冷啟)到 `max_ready`;閒置時逐階縮回 `ready → sleep → stop`。vLLM **sleep mode**(level-1)釋放副本 VRAM 但秒級喚醒,所以縮容不必付出數分鐘冷啟代價。config.yaml 或控制台皆可設定,內建即時 Grafana 面板與告警。 @@ -59,6 +60,11 @@ make up # 建置並啟動整套服務 `make down` 停止 · `make logs` 追蹤所有服務日誌 · `make ps` 看狀態。 +**兩種啟動方式:** + +- **`make up`** — 預設的**純 vLLM** 集群。 +- **`make up-mixed`** — **vLLM + SGLang** 混合集群:一個 vLLM backend 與一個 SGLang backend 共用同一顆 Postgres、router、控制台與 Grafana。從 *新增模型 → 引擎:`sglang`* 新增的模型會自動擺到 SGLang backend。對應 `make down-mixed`/`make logs-mixed`。見 [docs/mixed-engine-deployment_zh-CN.md](docs/mixed-engine-deployment_zh-CN.md)。 + ```bash curl http://localhost:8887/v1/models # router:列出設定的模型群組 curl http://localhost:5000/api/models # 後端:每個實例的生命週期狀態 @@ -98,6 +104,8 @@ curl http://localhost:8887/v1/chat/completions \ ## 架構 +### 純 vLLM(`make up`) + ```mermaid flowchart LR Client([Clients 用戶端]) @@ -129,11 +137,54 @@ flowchart LR Grafana 都在 nginx 之後以單一來源對外;backend、router、Prometheus 共用一個 network namespace,所以被拉起的 vLLM 實例可在 `localhost` 互相連到。 +### vLLM + SGLang 混合(`make up-mixed`) + +兩個引擎各跑一個 backend 容器(無法共用 netns),共用一顆 Postgres(排程/desired 意圖)、 +一個 router、一個 dashboard 與一套監控;各 backend 把自己 ready 的實例以**可路由位址**寫進 +共享 file_sd,Prometheus 一起抓。見 [docs/mixed-engine-deployment_zh-CN.md](docs/mixed-engine-deployment_zh-CN.md)。 + +```mermaid +flowchart LR + Client([Clients 用戶端]) + FE["frontend
nginx · :8884"] + GF["grafana
/grafana"] + PG[("postgres
共用 store · 排程/desired")] + PR["prometheus · :9090
抓可路由 file_sd targets"] + RT["router · :8887
OpenAI 相容負載平衡"] + + subgraph vbe["vLLM backend (engine.Dockerfile)"] + BV["backend · :5071
NODE_ENGINES=vllm"] + VINS["vLLM 實例"] + end + subgraph sbe["SGLang backend (engine-sglang.Dockerfile)"] + BS["backend · :5072
NODE_ENGINES=sglang"] + SINS["SGLang 實例"] + end + + Client --> FE + FE -->|/api| BV + FE -->|/v1| RT + FE -->|/grafana| GF + BV -->|拉起| VINS + BS -->|拉起| SINS + RT -->|路由| VINS + RT -->|路由| SINS + BV <-->|leader/排程| PG + BS <-->|收斂 desired| PG + PR -->|scrape| VINS + PR -->|scrape| SINS + GF -->|查詢| PR +``` + +leader 的 **engine-aware 排程器**把每顆模型擺到「跑得動它引擎」的 backend;落錯 node 的控制 +動作會延後給擁有者執行。SGLang 走 OpenMetrics,指標入庫為底線的 `sglang_*`,vLLM 則保留冒號。 + ## 文件 | 主題 | | |---|---| | 部署與架構 | [docs/deployment_zh-CN.md](docs/deployment_zh-CN.md) | +| 混合引擎(vLLM + SGLang) | [docs/mixed-engine-deployment_zh-CN.md](docs/mixed-engine-deployment_zh-CN.md) | | 配置(`config.yaml`) | [docs/configuration_zh-CN.md](docs/configuration_zh-CN.md) | | 功能特色(詳細) | [docs/features_zh-CN.md](docs/features_zh-CN.md) | | 監控(Prometheus + Grafana) | [docs/monitoring_zh-CN.md](docs/monitoring_zh-CN.md) | diff --git a/apps/backend/app/api/models.py b/apps/backend/app/api/models.py index 6649afc..aa7b1d6 100644 --- a/apps/backend/app/api/models.py +++ b/apps/backend/app/api/models.py @@ -23,13 +23,16 @@ SleepError, VRAMInsufficient, ) -from app.services.vllm_command import parse_vllm_command +from app.services.vllm_command import parse_command as parse_engine_command router = APIRouter(prefix="/models", tags=["models"]) class ParseRequest(BaseModel): command: str + # Which engine's CLI to parse ('vllm' | 'sglang'). Omitted = sniff from the + # command (sglang.launch_server -> sglang, else vLLM). + engine: Optional[str] = None class InstanceSpec(BaseModel): @@ -52,20 +55,18 @@ class CreateModelRequest(BaseModel): @router.get("", response_model=list[ModelView]) async def list_models(request: Request, manager: ModelManager = Depends(get_manager)): - # HA Phase 3d: a non-leader replica reports the fleet from the shared store's - # observed state (the leader/owning agents backfill it), since its own registry - # is idle. The leader (and a single-host collapsed deploy) uses its live - # registry — identical to before. - elector = getattr(request.app.state, "leader", None) - prefer_store = elector is not None and not elector.is_leader - return [ModelView(**v) for v in await manager.fleet_views(prefer_store=prefer_store)] + # HA: in Postgres/multi-node mode the fleet view comes from the shared store + # (each node backfills its *owned* observed state) — on leader and follower + # alike, since with per-node actuation (Phase 7) no single registry is complete. + # SQLite collapsed: the local registry is the truth — identical to before. + return [ModelView(**v) for v in await manager.fleet_views(prefer_store=manager.prefer_store_view())] @router.post("/parse", dependencies=[Depends(require_operator)]) async def parse_command(body: ParseRequest, manager: ModelManager = Depends(get_manager)): - """Parse a pasted vLLM command into editable fields + conflict hints.""" + """Parse a pasted vLLM / SGLang command into editable fields + conflict hints.""" try: - parsed = parse_vllm_command(body.command) + parsed = parse_engine_command(body.command, body.engine) except ValueError as e: raise HTTPException(status.HTTP_400_BAD_REQUEST, str(e)) inst = parsed["instance"] @@ -167,6 +168,13 @@ async def unload_lora(key: str, name: str, manager: ModelManager = Depends(get_m @router.get("/{key}", response_model=ModelView) async def get_model(key: str, manager: ModelManager = Depends(get_manager)): + # HA: like list_models, prefer the shared store's observed state in multi-node + # mode so a model owned by another node shows its real state (not this node's + # idle registry). Falls back to the local registry (collapsed / not in store). + if manager.prefer_store_view(): + for v in await manager.fleet_views(prefer_store=True): + if v.get("key") == key: + return ModelView(**v) try: return ModelView.from_instance(await manager.get(key)) except ModelNotFound: diff --git a/apps/backend/app/api/observability.py b/apps/backend/app/api/observability.py index c7fb8fe..5d63ddc 100644 --- a/apps/backend/app/api/observability.py +++ b/apps/backend/app/api/observability.py @@ -9,8 +9,10 @@ import json from typing import Optional +from urllib.parse import quote + from fastapi import APIRouter, Depends, HTTPException, Request, status -from fastapi.responses import StreamingResponse +from fastapi.responses import JSONResponse, StreamingResponse from app.api.deps import get_manager from app.api.schemas import ModelView @@ -62,10 +64,35 @@ async def requests_log(request: Request, model_key: Optional[str] = None, limit: return await _store(request).recent_requests(model_key=model_key, limit=limit) +async def _proxy_to_owner(request: Request, manager, key: str, suffix: str): + """If `key` runs on another node (HA), GET that node's backend API for this + node-local data (logs/metrics live as files on the owning node) and relay the + JSON. Returns None when the model is local (caller reads locally).""" + api_url = await manager.owning_node_api_url(key) + if not api_url: + return None + url = f"{api_url}/api/models/{quote(key, safe='')}/{suffix}" + headers = {} + auth = request.headers.get("authorization") + if auth: + headers["authorization"] = auth + try: + resp = await request.app.state.http_client.get( + url, params=dict(request.query_params), headers=headers, timeout=10.0 + ) + except Exception: + raise HTTPException(status.HTTP_502_BAD_GATEWAY, + f"failed to reach owning node for {key}") + return JSONResponse(status_code=resp.status_code, content=resp.json()) + + @router.get("/models/{key}/logs") async def model_logs( - key: str, tail: int = 200, manager: ModelManager = Depends(get_manager) + request: Request, key: str, tail: int = 200, manager: ModelManager = Depends(get_manager) ): + proxied = await _proxy_to_owner(request, manager, key, "logs") + if proxied is not None: + return proxied try: inst = await manager.get(key) except ModelNotFound: @@ -77,11 +104,14 @@ async def model_logs( @router.get("/models/{key}/metrics") -async def model_metrics(key: str, manager: ModelManager = Depends(get_manager)): - """vLLM startup capacity/memory/compile metrics parsed from the engine log. +async def model_metrics(request: Request, key: str, manager: ModelManager = Depends(get_manager)): + """vLLM/SGLang startup capacity/memory/compile metrics parsed from the engine log. Only meaningful once the instance is READY (the metrics are printed at the end of model loading); returns {ready: false} otherwise so the UI hides the panel.""" + proxied = await _proxy_to_owner(request, manager, key, "metrics") + if proxied is not None: + return proxied try: inst = await manager.get(key) except ModelNotFound: @@ -94,17 +124,22 @@ async def model_metrics(key: str, manager: ModelManager = Depends(get_manager)): return {"ready": True, **parse_startup_metrics(head)} -async def model_snapshot_stream(registry, interval: float = 1.0, heartbeat_every: float = 15.0): +async def model_snapshot_stream(manager, interval: float = 1.0, heartbeat_every: float = 15.0): """SSE generator: emit the full model snapshot whenever it changes. - Cheap snapshot-diff (no pub/sub bus): re-snapshot the registry each interval, - emit on change, and send a heartbeat comment if nothing changed for a while. + Cheap snapshot-diff (no pub/sub bus): re-read the fleet each interval, emit on + change, send a heartbeat if nothing changed for a while. In HA/multi-node mode + the view comes from the shared store (each node only actuates its own models, so + the local registry alone is incomplete — Phase 7); collapsed uses the registry. """ last_sig: Optional[str] = None since_emit = 0.0 while True: - snap = await registry.snapshot() - payload = [ModelView.from_instance(i).model_dump(mode="json") for i in snap] + payload = sorted( + (ModelView(**v).model_dump(mode="json") + for v in await manager.fleet_views(prefer_store=manager.prefer_store_view())), + key=lambda m: m["key"], + ) sig = json.dumps(payload, sort_keys=True, default=str) if sig != last_sig: last_sig = sig @@ -122,6 +157,6 @@ async def model_snapshot_stream(registry, interval: float = 1.0, heartbeat_every async def stream_models(request: Request): """Live model state via Server-Sent Events (replaces frontend polling).""" return StreamingResponse( - model_snapshot_stream(request.app.state.registry), + model_snapshot_stream(request.app.state.manager), media_type="text/event-stream", ) diff --git a/apps/backend/app/api/schemas.py b/apps/backend/app/api/schemas.py index d5e78db..d8d06c4 100644 --- a/apps/backend/app/api/schemas.py +++ b/apps/backend/app/api/schemas.py @@ -14,6 +14,8 @@ class ModelView(BaseModel): key: str kind: ModelKind + # Which inference engine backs this instance ("vllm" / "sglang" / …). + engine: str = "vllm" model_tag: Optional[str] = None host: str port: int @@ -35,6 +37,7 @@ def from_instance(cls, inst: ModelInstance) -> "ModelView": return cls( key=inst.key, kind=inst.kind, + engine=inst.engine, model_tag=inst.model_tag, host=inst.host, port=inst.port, diff --git a/apps/backend/app/core/settings.py b/apps/backend/app/core/settings.py index de8d311..4ce43ca 100644 --- a/apps/backend/app/core/settings.py +++ b/apps/backend/app/core/settings.py @@ -8,7 +8,7 @@ import os import socket -from dataclasses import dataclass +from dataclasses import dataclass, field from typing import Optional @@ -132,6 +132,16 @@ class BackendSettings: # The live address is heartbeated each reconcile pass; live_ttl is its lease. node_host: str = "" live_ttl: float = 30.0 + # HA Phase 7: which engines this node can actually run (determined by its image — + # the vLLM image runs vllm, the SGLang image runs sglang). Advertised in the + # node registry so the scheduler places each model on an engine-matching node. + # Empty/unset = unspecified => runs any engine (collapsed single host: unchanged). + node_engines: list[str] = field(default_factory=list) + # HA Phase 7: this node's own backend API base URL (e.g. http://node-host:5000), + # advertised in the node registry so a dashboard served by another node can proxy + # node-local requests (model logs / startup metrics) to the node that owns the + # model. Empty = not advertised (cross-node log/metrics fall back to a local read). + node_api_url: str = "" # HA Phase 3b: node-agent heartbeat. node_id reuses instance_id (the same id # used for the leader lease and live-address node_id). The agent re-registers @@ -198,6 +208,10 @@ def from_env(cls) -> "BackendSettings": or f"{socket.gethostname()}:{os.getpid()}", leader_lease_ttl=_env_float("LLMOPS_LEADER_LEASE_TTL", 15.0), node_host=os.environ.get("LLMOPS_NODE_HOST", "").strip(), + node_engines=[ + e.strip() for e in os.environ.get("LLMOPS_NODE_ENGINES", "").split(",") if e.strip() + ], + node_api_url=os.environ.get("LLMOPS_NODE_API_URL", "").strip().rstrip("/"), live_ttl=_env_float("LLMOPS_LIVE_TTL", 30.0), node_heartbeat_interval=_env_float("LLMOPS_NODE_HEARTBEAT_INTERVAL", 10.0), node_ttl=_env_float("LLMOPS_NODE_TTL", 30.0), diff --git a/apps/backend/app/llmops/instance.py b/apps/backend/app/llmops/instance.py index 75935e8..d20086b 100644 --- a/apps/backend/app/llmops/instance.py +++ b/apps/backend/app/llmops/instance.py @@ -29,6 +29,11 @@ class LaunchSpec: host: str port: int probe_url: str + # Which inference engine this spec launches ("vllm" / "sglang" / …) and the + # optional features it supports. Callers gate on capabilities, not the engine + # name. See launchers.CAP_* and docs/multi-backend-engine-design_zh-CN.md. + engine: str = "vllm" + capabilities: frozenset = field(default_factory=frozenset) model_tag: Optional[str] = None # True when launched with --enable-sleep-mode + VLLM_SERVER_DEV_MODE=1, so the # /sleep, /wake_up and /is_sleeping dev endpoints are available. @@ -51,6 +56,9 @@ class ModelInstance: port: int spec: LaunchSpec model_tag: Optional[str] = None + # Inference engine backing this instance (mirrors spec.engine); surfaced in the + # fleet view + persisted to the shared store (HA) so any replica can render it. + engine: str = "vllm" desired: Desired = Desired.STOPPED state: ModelState = ModelState.STOPPED @@ -85,6 +93,7 @@ def observed_dict(self) -> dict: return { "key": self.key, "kind": self.kind.value, + "engine": self.engine, "model_tag": self.model_tag, "host": self.host, "port": self.port, diff --git a/apps/backend/app/llmops/launchers.py b/apps/backend/app/llmops/launchers.py index a20c687..a98cc3e 100644 --- a/apps/backend/app/llmops/launchers.py +++ b/apps/backend/app/llmops/launchers.py @@ -55,8 +55,9 @@ def _write_effective_config(config) -> str: _LORA_RUNTIME_KEY = "allow_runtime_lora" # Router-only knobs that ride the shared model_config (EngineModelConfig is # extra="allow") but belong to the router, not vLLM — never pass them to -# `vllm serve` or it errors on an unknown argument. -_ROUTER_ONLY_KEYS = frozenset({"routing_strategy", "kind"}) +# `vllm serve` or it errors on an unknown argument. `engine` is launcher-meta +# (which backend to run); it too must never reach the vLLM CLI. +_ROUTER_ONLY_KEYS = frozenset({"routing_strategy", "kind", "engine"}) # Everything build_vllm_cli_args must skip (model_tag is the positional arg). _SKIP_CLI_KEYS = frozenset({"model_tag", _LORA_RUNTIME_KEY}) | _ROUTER_ONLY_KEYS @@ -126,11 +127,32 @@ def build_vllm_cli_args(model_cfg: dict) -> list[str]: return cli_args +# Engine capability flags. Callers gate optional features on *capabilities*, never +# on the engine name (`if "sleep" in caps`, never `if engine == "vllm"`), so adding +# an engine is just declaring its capability set. See docs/multi-backend-engine-design_zh-CN.md §4. +CAP_SLEEP = "sleep" # /sleep + /wake_up (warm standby, frees VRAM, process stays up) +CAP_RUNTIME_LORA = "runtime_lora" # runtime LoRA load/unload endpoints +CAP_LORA_MODULES = "lora_modules" # static --lora-modules at launch +CAP_KV_TRANSFER = "kv_transfer" # cross-instance KV cache sharing +CAP_METRICS_VLLM = "metrics_vllm" # exposes vLLM-format Prometheus metrics (waiting queue, …) +CAP_METRICS_SGLANG = "metrics_sglang" # exposes sglang:* Prometheus metrics (the router parses these) + +# Sentinel engine name for non-LLM launchers (embedding server): they aren't +# selected by an engine choice, so they register under one fixed value. +ENGINE_DEFAULT = "default" + + class Launcher(Protocol): kind: ModelKind + # Which engine this launcher serves; dispatch is keyed on (kind, engine). + engine: str + # Optional features this engine supports (see CAP_* above). Threaded onto the + # LaunchSpec so callers (autoscaler, sleep/LoRA APIs, metrics) can gate without + # re-checking the engine name. + capabilities: frozenset[str] def keys(self, config) -> list[str]: - """All instance keys this launcher's kind defines in the config.""" + """All instance keys this launcher defines in the config (its engine only).""" ... def build_spec(self, config, config_path: str, key: str) -> LaunchSpec: @@ -140,10 +162,19 @@ def build_spec(self, config, config_path: str, key: str) -> LaunchSpec: class VllmLauncher: kind = ModelKind.LLM + engine = "vllm" + capabilities = frozenset({ + CAP_SLEEP, CAP_RUNTIME_LORA, CAP_LORA_MODULES, CAP_KV_TRANSFER, CAP_METRICS_VLLM, + }) def keys(self, config) -> list[str]: out: list[str] = [] for model_tag, engine in config.LLM_engines.items(): + # Only claim groups configured for this engine. `engine` defaults to + # "vllm" (EngineModelConfig), so a config with no engine field is all + # vLLM = today's behaviour. + if getattr(engine.settings, "engine", "vllm") != self.engine: + continue for inst in engine.instances: out.append(f"{model_tag}::{inst.id}") return out @@ -204,6 +235,8 @@ def build_spec(self, config, config_path: str, key: str) -> LaunchSpec: return LaunchSpec( key=key, kind=self.kind, + engine=self.engine, + capabilities=self.capabilities, command=command, env=env, log_path=log_path, @@ -217,6 +250,8 @@ def build_spec(self, config, config_path: str, key: str) -> LaunchSpec: class EmbeddingLauncher: kind = ModelKind.EMBEDDING + engine = ENGINE_DEFAULT # not engine-selectable; one launcher for the embedding server + capabilities = frozenset() def keys(self, config) -> list[str]: emb = config.embedding_server @@ -246,6 +281,8 @@ def build_spec(self, config, config_path: str, key: str) -> LaunchSpec: return LaunchSpec( key=key, kind=self.kind, + engine=self.engine, + capabilities=self.capabilities, command=command, env=env, log_path=log_path, @@ -253,3 +290,148 @@ def build_spec(self, config, config_path: str, key: str) -> LaunchSpec: port=emb.port, probe_url=f"http://{emb.host}:{emb.port}/health", ) + + +# ---- SGLang --------------------------------------------------------------- + +# Engine-neutral typed params (EngineModelConfig) -> SGLang flag names. These three +# have different names from vLLM; everything else falls through as a kebab-cased +# -- (engine-native passthrough via extra="allow"). See §2.3 of the design doc. +_SGLANG_PARAM_MAP = { + "max_model_len": "context-length", + "gpu_memory_utilization": "mem-fraction-static", + "tensor_parallel_size": "tp-size", +} +# Keys consumed specially (model_tag/served_model_name handled up front; id is the +# key; cuda_device becomes CUDA_VISIBLE_DEVICES; enable_lora/lora_modules drive the +# LoRA flags below; allow_runtime_lora is an enable_lora trigger) + router-only knobs. +_SGLANG_SKIP_CLI_KEYS = frozenset( + {"model_tag", "served_model_name", "id", "cuda_device", + "enable_lora", "lora_modules", "enable_metrics", _LORA_RUNTIME_KEY} +) | _ROUTER_ONLY_KEYS + + +def build_sglang_cli_args(model_cfg: dict) -> list[str]: + """dict -> ``python -m sglang.launch_server`` CLI args. + + Unlike vLLM: the model is ``--model-path`` (not a positional), and we always + emit ``--served-model-name`` so ``/v1/models`` (and the router's forward_name) + is the stable ``model_tag`` rather than SGLang's default of the raw path. + + Bools are SGLang's ``store_true``: a True value emits ``--flag``; a False value + is *omitted* — there is NO ``--no-flag`` dual (vLLM's BooleanOptionalAction), + and many SGLang flags are themselves negative (``--disable-radix-cache``), so we + must not synthesise ``--no-`` forms. The three common params above are + translated; the rest pass through kebab-cased. + """ + model_tag = model_cfg.get("model_tag") + if not model_tag: + raise ValueError("model_config must provide 'model_tag'") + served = model_cfg.get("served_model_name") or model_tag + + args = ["--model-path", str(model_tag), "--served-model-name", str(served)] + + # Always expose Prometheus /metrics (vLLM does by default; SGLang needs the + # flag). The router scrapes sglang:* from it for the autoscaler signal. + args.append("--enable-metrics") + + # LoRA: SGLang needs --enable-lora at launch to accept either static adapters + # (--lora-paths NAME=PATH …) or runtime ones (POST /load_lora_adapter). Turn it + # on if any LoRA usage is configured (static modules, enable_lora, or our runtime + # toggle). Static modules use SGLang's NAME=PATH form (not vLLM's --lora-modules + # JSON). Other LoRA knobs (max_lora_rank, lora_target_modules, …) pass through + # below as plain kebab-cased flags. + lora_modules = model_cfg.get("lora_modules") or [] + if model_cfg.get("enable_lora") or model_cfg.get(_LORA_RUNTIME_KEY) or lora_modules: + args.append("--enable-lora") + paths = [f"{m['name']}={m['path']}" for m in lora_modules + if m.get("name") and m.get("path")] + if paths: + args.append("--lora-paths") + args.extend(paths) + + for key, value in model_cfg.items(): + if key in _SGLANG_SKIP_CLI_KEYS or value is None: + continue + flag = "--" + _SGLANG_PARAM_MAP.get(key, key).replace("_", "-") + if isinstance(value, bool): + if value: # store_true: True -> present; False -> omit + args.append(flag) + elif isinstance(value, list): + # SGLang multi-value flags take space-separated values (e.g. --lora-paths). + args.append(flag) + args.extend(str(v) for v in value) + elif isinstance(value, dict): + args.append(flag) + args.append(json.dumps(value, ensure_ascii=False)) + else: + args.append(flag) + args.append(str(value)) + return args + + +class SglangLauncher: + kind = ModelKind.LLM + engine = "sglang" + # Runtime + static LoRA are wired (--enable-lora / --lora-paths at launch; + # POST /load_lora_adapter — no /v1 prefix, unlike vLLM — for hot load/unload). + # No sleep (SGLang has no /sleep+/wake_up; `--sleep-on-idle` only lowers CPU, + # doesn't free VRAM) -> autoscaler degrades to ready<->stopped. metrics_sglang: + # launches with --enable-metrics and the router parses sglang:* into the same + # normalized load shape, so the autoscaler scales SGLang groups too. + # See docs/multi-backend-engine-design_zh-CN.md §5.2. + capabilities = frozenset({CAP_RUNTIME_LORA, CAP_LORA_MODULES, CAP_METRICS_SGLANG}) + + def keys(self, config) -> list[str]: + out: list[str] = [] + for model_tag, engine in config.LLM_engines.items(): + if getattr(engine.settings, "engine", "vllm") != self.engine: + continue + for inst in engine.instances: + out.append(f"{model_tag}::{inst.id}") + return out + + def build_spec(self, config, config_path: str, key: str) -> LaunchSpec: + model_tag, _, instance_id = key.partition("::") + engine = config.LLM_engines.get(model_tag) + if engine is None: + raise KeyError(f"model group '{model_tag}' not in config") + inst = next((i for i in engine.instances if i.id == instance_id), None) + if inst is None: + raise KeyError(f"instance '{instance_id}' not in group '{model_tag}'") + + # Merge shared model_config with instance overrides, mirroring VllmLauncher: + # single-GPU cuda_device -> CUDA_VISIBLE_DEVICES; the `id` field is dropped. + merged: dict = engine.settings.model_dump(by_alias=False) + merged.update(inst.model_dump()) + + env: dict[str, str] = {} + if merged.get("tensor_parallel_size", 1) == 1: + cuda_device = merged.pop("cuda_device", None) + if cuda_device is not None: + env["CUDA_VISIBLE_DEVICES"] = str(cuda_device) + merged.pop("id", None) + + # HA split deploys: optionally bind SGLang to a routable interface (0.0.0.0) + # so a router in another container/host reaches it via the advertised + # LLMOPS_NODE_HOST (instances_live). Only the bind --host changes; the local + # probe + recorded host stay localhost (which a 0.0.0.0 bind also serves). + # Empty (default) = bind the configured host = today's localhost-only. Shares + # the env var with vLLM so one setting governs both engines. + bind_host = os.environ.get("LLMOPS_VLLM_BIND_HOST", "").strip() + cli_cfg = {**merged, "host": bind_host} if bind_host else merged + command = [sys.executable, "-m", "sglang.launch_server"] + build_sglang_cli_args(cli_cfg) + log_path = os.path.join(LOG_DIR, f"{model_tag}__{instance_id}.log") + return LaunchSpec( + key=key, + kind=self.kind, + engine=self.engine, + capabilities=self.capabilities, + command=command, + env=env, + log_path=log_path, + host=inst.host, + port=inst.port, + probe_url=f"http://{inst.host}:{inst.port}/health", + model_tag=engine.settings.model_tag, + ) diff --git a/apps/backend/app/llmops/manager.py b/apps/backend/app/llmops/manager.py index 73d31c7..1776786 100644 --- a/apps/backend/app/llmops/manager.py +++ b/apps/backend/app/llmops/manager.py @@ -16,7 +16,7 @@ from app.core.settings import BackendSettings from app.llmops.events import emit_transition from app.llmops.instance import ModelInstance -from app.llmops.launchers import Launcher +from app.llmops.launchers import CAP_RUNTIME_LORA, CAP_SLEEP, Launcher from app.llmops.process import spawn_process, terminate_process_group from app.llmops.registry import ModelRegistry from app.llmops.state import Desired, ModelKind, ModelState @@ -68,6 +68,7 @@ def build_registry(config, config_path: str, launchers: list[Launcher]) -> Model ModelInstance( key=key, kind=launcher.kind, + engine=spec.engine, host=spec.host, port=spec.port, spec=spec, @@ -93,7 +94,11 @@ def __init__( notifier=None, ) -> None: self.registry = registry - self._launchers: dict[ModelKind, Launcher] = {l.kind: l for l in launchers} + # Dispatch is keyed on (kind, engine): vLLM and SGLang are both ModelKind.LLM + # but distinct launchers. Embedding registers under ENGINE_DEFAULT. + self._launchers: dict[tuple[ModelKind, str], Launcher] = { + (l.kind, l.engine): l for l in launchers + } self.http_client = http_client self.config = config self.config_path = config_path @@ -111,6 +116,54 @@ def _require(self, key: str) -> ModelInstance: raise ModelNotFound(key) return inst + def _launcher_for(self, inst: ModelInstance) -> Launcher: + """The launcher that owns an instance, by its (kind, engine).""" + return self._launchers[(inst.kind, inst.engine)] + + def _node_can_run(self, engine: str) -> bool: + """Whether THIS node can run an engine (its image has it). Empty + node_engines = unspecified = runs any (collapsed single host / single-engine + deploys: always True, so the sync actuation path below is unchanged).""" + ne = self.settings.node_engines + return not ne or engine in ne + + async def _defer_to_owner(self, inst: ModelInstance, desired: Desired) -> ModelInstance: + """HA Phase 7C: this node can't run `inst`'s engine, so don't actuate locally — + just record the intent and let the scheduler place it on an engine-matching + node, whose reconcile loop converges it. Returns the instance (state unchanged + here; the dashboard tracks progress via observed state from the owning node).""" + async with self.registry.lock: + inst.desired = desired + inst.touch() + if self.store is not None: + try: + await self.store.set_instance_desired(inst.key, desired.value) + except Exception: + logger.warning("defer: failed to persist desired for %s", inst.key, exc_info=True) + # Clear any assignment so the engine-aware scheduler places it fresh on a + # node that can actually run it (it would reassign anyway, but this avoids + # a transient wrong-node attempt). + if desired != Desired.STOPPED and hasattr(self.store, "delete_assignment"): + try: + await self.store.delete_assignment(inst.key) + except Exception: + logger.debug("defer: clear assignment failed for %s", inst.key, exc_info=True) + logger.info("Deferred %s (engine=%s) to an engine-matching node (desired=%s)", + inst.key, inst.engine, desired.value) + return inst + + def _llm_engine_capabilities(self, group: str) -> frozenset: + """Capabilities of the engine an LLM group is configured for. Callers gate + optional features (sleep, runtime LoRA, …) on these rather than the engine + name, so a new engine only needs to declare its capability set. Empty if the + group / its engine's launcher is unknown.""" + engine = self.config.LLM_engines.get(group) + if engine is None: + return frozenset() + engine_name = getattr(engine.settings, "engine", "vllm") + launcher = self._launchers.get((ModelKind.LLM, engine_name)) + return launcher.capabilities if launcher else frozenset() + async def _record(self, inst, from_state, to_state, detail=None) -> None: """Persist a state transition + dispatch any alert, via the shared funnel. Best-effort: telemetry/alerts never break ops.""" @@ -203,6 +256,25 @@ async def foreign_assignments(self) -> set[str]: alive = set() return {k for k, n in amap.items() if n != me and n in alive} + async def owning_node_api_url(self, key: str) -> Optional[str]: + """The backend API base URL of the node that owns `key`, when that node is a + *different*, alive node advertising an api_url — so node-local requests (logs, + startup metrics) can be proxied to it. None when the model is local / owner + unknown / no api_url advertised (caller then reads locally). Best-effort.""" + if self.store is None or not hasattr(self.store, "list_assignments"): + return None + try: + owner = (await self.store.list_assignments()).get(key) + if owner is None or owner == self.settings.instance_id: + return None + for n in await self.store.list_nodes(): + if n["node_id"] == owner: + url = n.get("api_url") + return url.rstrip("/") if url else None + except Exception: + logger.debug("owning_node_api_url failed for %s", key, exc_info=True) + return None + async def replay_desired(self) -> None: """On boot (after adopt_running): start instances whose persisted desired is RUNNING but which are currently STOPPED/FAILED — so a backend restart (or a @@ -253,7 +325,7 @@ async def write_prometheus_targets(self) -> bool: from app.services.prometheus_targets import build_targets, write_targets_file instances = await self.registry.snapshot() - targets = build_targets(instances) + targets = build_targets(instances, node_host=self.settings.node_host) loop = asyncio.get_event_loop() try: return await loop.run_in_executor(None, write_targets_file, path, targets) @@ -264,6 +336,17 @@ async def write_prometheus_targets(self) -> bool: async def list(self) -> list[ModelInstance]: return await self.registry.snapshot() + def prefer_store_view(self) -> bool: + """Whether model-state views should come from the shared store rather than + this node's local registry. In HA/Postgres mode every node only actuates the + models assigned to it (Phase 7), so no single registry is the full truth — + each node backfills its *owned* observed state to the store, which is the + only complete view. Use it on leader and follower alike. SQLite collapsed: + one node owns everything, so the local registry is the truth (unchanged).""" + return (self.store is not None + and getattr(self.store, "db_url", None) is not None + and hasattr(self.store, "list_instance_observed")) + async def fleet_views(self, prefer_store: bool = False) -> list[dict]: """The fleet's model-view dicts. Leader (prefer_store=False): from the live local registry — identical to before. Follower (prefer_store=True): the @@ -283,6 +366,18 @@ async def fleet_views(self, prefer_store: bool = False) -> list[dict]: base[v["key"]] = v except Exception: logger.warning("fleet_views: failed to read observed from store", exc_info=True) + # Overlay the *authoritative* desired intent (instance_desired, written + # immediately by whichever replica took the API call) over the observed + # desired, which lags until the owning node's reconcile converges. Without + # this the dashboard briefly sees desired=running + state=ready after a + # stop and flickers ready->stopped. Best-effort. + try: + desired = await self.store.list_instance_desired() + for key, want in desired.items(): + if key in base: + base[key]["desired"] = want + except Exception: + logger.debug("fleet_views: failed to overlay desired", exc_info=True) return list(base.values()) async def get(self, key: str) -> ModelInstance: @@ -357,7 +452,11 @@ async def _gpu_exists_preflight(self, key: str, spec) -> None: async def start(self, key: str, force: bool = False, reset_restart: bool = True) -> ModelInstance: inst = self._require(key) - launcher = self._launchers[inst.kind] + # HA Phase 7C: if this node can't run the engine, write intent and let an + # engine-matching node actuate it (collapsed/single-engine: always can-run). + if not self._node_can_run(inst.engine): + return await self._defer_to_owner(inst, Desired.RUNNING) + launcher = self._launcher_for(inst) # Re-resolve the spec (config may have changed) outside the lock so the # GPU pre-flight's nvidia-smi call never extends the critical section. @@ -411,6 +510,9 @@ async def start(self, key: str, force: bool = False, reset_restart: bool = True) async def stop(self, key: str) -> ModelInstance: inst = self._require(key) + # HA Phase 7C: not our engine -> record desired=stopped; the owning node stops it. + if not self._node_can_run(inst.engine): + return await self._defer_to_owner(inst, Desired.STOPPED) async with self.registry.lock: prev = inst.state @@ -462,6 +564,8 @@ async def sleep(self, key: str, level: int = 1) -> ModelInstance: The HTTP call runs outside the registry lock; on failure the desired intent is reverted so the reconciler doesn't fight a half-applied state.""" inst = self._require(key) + if not self._node_can_run(inst.engine): # HA Phase 7C: owning node sleeps it + return await self._defer_to_owner(inst, Desired.ASLEEP) if not getattr(inst.spec, "sleep_enabled", False): raise ModelConflict( f"{key} was not launched with sleep mode " @@ -501,6 +605,8 @@ async def wake(self, key: str) -> ModelInstance: """Wake a SLEEPING instance back to READY (vLLM /wake_up reloads weights to GPU). Seconds, not a cold start.""" inst = self._require(key) + if not self._node_can_run(inst.engine): # HA Phase 7C: owning node wakes it + return await self._defer_to_owner(inst, Desired.RUNNING) async with self.registry.lock: if inst.state != ModelState.SLEEPING: raise ModelConflict( @@ -617,13 +723,17 @@ async def create_overlay_model(self, group: str, instance: dict, model_config: d save_overlay(overlay, self.overlay_path) self.config = new_config - launcher = self._launchers[ModelKind.LLM] + group_engine = getattr(new_config.LLM_engines[group].settings, "engine", "vllm") + launcher = self._launchers.get((ModelKind.LLM, group_engine)) + if launcher is None: + raise ModelConflict(f"unsupported engine '{group_engine}' (no launcher registered)") spec = launcher.build_spec(self.config, self.config_path, key) async with self.registry.lock: self.registry.add( ModelInstance( key=key, kind=ModelKind.LLM, + engine=spec.engine, host=spec.host, port=spec.port, spec=spec, @@ -679,10 +789,15 @@ async def update_overlay_model(self, key: str, instance: dict, model_config: dic save_overlay(overlay, self.overlay_path) self.config = new_config - # Re-resolve the spec so the next start uses the edited values. - launcher = self._launchers[ModelKind.LLM] + # Re-resolve the spec so the next start uses the edited values. The edit may + # have changed the group's engine, so pick the launcher by the new config. + group_engine = getattr(new_config.LLM_engines[group].settings, "engine", "vllm") + launcher = self._launchers.get((ModelKind.LLM, group_engine)) + if launcher is None: + raise ModelConflict(f"unsupported engine '{group_engine}' (no launcher registered)") spec = launcher.build_spec(self.config, self.config_path, key) async with self.registry.lock: + inst.engine = spec.engine inst.host = spec.host inst.port = spec.port inst.spec = spec @@ -706,8 +821,11 @@ async def _ready_lora_targets(self, group: str) -> list[ModelInstance]: ] async def _post_lora(self, inst: ModelInstance, action: str, payload: dict) -> None: - """POST /v1/{load,unload}_lora_adapter to one instance; raise on failure.""" - url = f"http://{inst.host}:{inst.port}/v1/{action}_lora_adapter" + """POST {load,unload}_lora_adapter to one instance; raise on failure. The + path differs by engine: vLLM serves /v1/_lora_adapter, SGLang serves + /_lora_adapter (no /v1). The JSON body is the same for both.""" + prefix = "" if getattr(inst, "engine", "vllm") == "sglang" else "/v1" + url = f"http://{inst.host}:{inst.port}{prefix}/{action}_lora_adapter" resp = await self.http_client.post(url, json=payload, timeout=120.0) if resp.status_code >= 400: try: @@ -752,6 +870,11 @@ async def load_lora(self, group: str, name: str, path: str, engine = self.config.LLM_engines.get(group) if engine is None: raise ModelNotFound(group) + if CAP_RUNTIME_LORA not in self._llm_engine_capabilities(group): + engine_name = getattr(engine.settings, "engine", "vllm") + raise ModelConflict( + f"{group}'s engine ({engine_name}) does not support runtime LoRA" + ) if not getattr(engine.settings, "allow_runtime_lora", False): raise ModelConflict( f"{group} was not started with runtime LoRA updating — enable " @@ -897,13 +1020,14 @@ def resync_registry(self, new_config) -> dict[str, list[str]]: for key in added: kind, spec = desired[key] self.registry.add(ModelInstance( - key=key, kind=kind, host=spec.host, port=spec.port, spec=spec, - model_tag=spec.model_tag, log_path=spec.log_path, + key=key, kind=kind, engine=spec.engine, host=spec.host, port=spec.port, + spec=spec, model_tag=spec.model_tag, log_path=spec.log_path, )) for key in changed: kind, spec = desired[key] inst = self.registry.get(key) inst.spec, inst.host, inst.port = spec, spec.host, spec.port + inst.engine = spec.engine inst.model_tag, inst.log_path = spec.model_tag, spec.log_path return {"added": added, "removed": removed, "changed": changed} @@ -972,6 +1096,42 @@ async def import_overlay( logger.info("Imported overlay: %s", summary_keys) return summary_keys + async def sync_overlay_from_store(self) -> bool: + """HA Phase 7: pull the fleet's current overlay from the shared store and + align this node's registry to it, so a model added/edited on *any* replica + appears here too (and an engine-matching node can then actuate it). No-op + outside Postgres mode (SQLite's file is already the single shared truth) and + when nothing changed. Best-effort: never raises. Returns True if it resynced. + + Only the registry alignment for *non-running-here* keys matters in practice: + a follower adds the new STOPPED instance, then converge_desired starts it if + it's assigned here. resync_registry refuses to change a key running locally, + which is the safe behaviour (the owning node edits its own).""" + if self.store is None or getattr(self.store, "db_url", None) is None: + return False + from app.services.overlay import (build_merged_config, + hydrate_overlay_from_store, load_overlay) + try: + changed = await hydrate_overlay_from_store(self.store, self.overlay_path) + if not changed: + return False + new_config = build_merged_config(self.config_path, load_overlay(self.overlay_path)) + except Exception: + logger.debug("sync_overlay_from_store: hydrate/merge failed", exc_info=True) + return False + async with self.registry.lock: + try: + summary = self.resync_registry(new_config) + except Exception: + logger.debug("sync_overlay_from_store: resync failed", exc_info=True) + return False + self.config = new_config + if any(summary.values()): + logger.info("Synced overlay from store: %s", summary) + await self.trigger_router_reload() + return True + return False + async def stop_all(self) -> None: """Best-effort shutdown of every managed process (used at app shutdown).""" for inst in self.registry.values(): diff --git a/apps/backend/app/llmops/node_agent.py b/apps/backend/app/llmops/node_agent.py index 6c5a90d..dd3acf7 100644 --- a/apps/backend/app/llmops/node_agent.py +++ b/apps/backend/app/llmops/node_agent.py @@ -51,12 +51,21 @@ def _capacity(self) -> Optional[str]: ] return json.dumps(slim) + def _engines(self) -> Optional[str]: + """Engines this node can run, as JSON, or None when unspecified (runs any). + Set per engine image via LLMOPS_NODE_ENGINES so the scheduler places each + model on an engine-matching node.""" + engines = self.settings.node_engines + return json.dumps(engines) if engines else None + async def heartbeat_once(self) -> None: """Register/refresh this node. No-op if the store can't track nodes.""" if self.store is None or not hasattr(self.store, "upsert_node"): return await self.store.upsert_node( - self.node_id, self.hostname, self._capacity(), ttl=self.settings.node_ttl + self.node_id, self.hostname, self._capacity(), + ttl=self.settings.node_ttl, engines=self._engines(), + api_url=self.settings.node_api_url or None, ) # Housekeeping so a vanished peer's row doesn't linger past its lease. if hasattr(self.store, "prune_nodes"): diff --git a/apps/backend/app/llmops/reconciler.py b/apps/backend/app/llmops/reconciler.py index 5dc01dc..73049da 100644 --- a/apps/backend/app/llmops/reconciler.py +++ b/apps/backend/app/llmops/reconciler.py @@ -262,6 +262,58 @@ async def _backfill_node_state( logger.debug("prune_instance_observed failed", exc_info=True) +async def converge_desired( + registry: ModelRegistry, settings: BackendSettings, store, manager, foreign: set[str], +) -> None: + """Drive every *owned* instance toward its persisted desired state — the per-node + actuation step (HA Phase 7). Generalises boot-time `replay_desired` into a + continuous convergence so a node actuates whatever is assigned to it, no matter + which replica received the API/autoscaler write: + + desired=running & STOPPED -> start (a STOPPED-but-wanted instance, e.g. + just assigned to this node) + desired=stopped & live -> stop (ready/starting/sleeping) + desired=asleep & READY -> sleep + desired=running & SLEEPING -> wake + + FAILED + desired=running recovery is intentionally left to `_process_restarts` + (restart budget + backoff), so this never fights the crash-loop guard. `foreign` + keys (owned by another live node) are skipped. Collapsed single host: the one + node owns everything, so this converges the whole fleet — same end state as the + sync API path, just driven by the loop. Best-effort per key.""" + if store is None or manager is None or not hasattr(store, "list_instance_desired"): + return + try: + desired = await store.list_instance_desired() + except Exception: + logger.debug("converge_desired: failed to read desired", exc_info=True) + return + + # Snapshot (key, state) under the lock; act outside it (start/stop take the lock). + async with registry.lock: + states = { + inst.key: inst.state for inst in registry.values() + } + + for key, want in desired.items(): + if key in foreign or key not in states: + continue + state = states[key] + try: + if want == Desired.RUNNING.value and state == ModelState.STOPPED: + await manager.start(key) + elif want == Desired.STOPPED.value and state in ( + ModelState.READY, ModelState.STARTING, ModelState.SLEEPING + ): + await manager.stop(key) + elif want == Desired.ASLEEP.value and state == ModelState.READY: + await manager.sleep(key) + elif want == Desired.RUNNING.value and state == ModelState.SLEEPING: + await manager.wake(key) + except Exception: + logger.debug("converge_desired: action for %s (want=%s) failed", key, want, exc_info=True) + + async def reconcile_once( registry: ModelRegistry, http_client, settings: BackendSettings, store=None, manager=None, notifier=None, @@ -300,6 +352,10 @@ async def reconcile_once( await manager.write_prometheus_targets() if manager is not None and settings.auto_restart: await _process_restarts(registry, settings, store, manager, foreign) + # Per-node actuation (HA Phase 7): converge owned instances to their persisted + # desired state. Reuses the `foreign` set computed above. Collapsed = no-op + # beyond what the sync API already did (idempotent). + await converge_desired(registry, settings, store, manager, foreign) async def adopt_running( diff --git a/apps/backend/app/llmops/scheduler.py b/apps/backend/app/llmops/scheduler.py index 8ed8589..a8ba0e3 100644 --- a/apps/backend/app/llmops/scheduler.py +++ b/apps/backend/app/llmops/scheduler.py @@ -30,32 +30,66 @@ def node_free_vram(node: dict) -> int: return sum(max(0, g.get("memory_total", 0) - g.get("memory_used", 0)) for g in gpus) +def node_supports(node: dict, engine: str) -> bool: + """Whether a node can run a given engine. A node that doesn't advertise engines + (engines NULL/empty) is unspecified and accepts any engine — so collapsed single + host and pre-Phase-7 deploys behave exactly as before.""" + raw = node.get("engines") + if not raw: + return True + try: + return engine in json.loads(raw) + except (json.JSONDecodeError, TypeError): + return True + + def place( desired: set[str], nodes: list[dict], assignments: dict[str, str], + key_engines: Optional[dict[str, str]] = None, ) -> dict[str, str]: - """Decide assignments to (re)write. For every desired-running instance whose - current assignment is missing or points to a node not in `nodes` (dead), pick - the live node with the most free VRAM. Returns only the *changes* {key: node}; - instances already on a live node are left where they are (no churn). + """Decide assignments to (re)write. For every desired-running instance, pick the + emptiest live node that can run the instance's engine — unless it's already on a + live, engine-matching node (no churn). Returns only the *changes* {key: node}. - `nodes` is the set of currently-alive nodes (each a row with node_id+capacity). + `key_engines` maps key -> engine (default "vllm"). A key whose engine no live + node supports is left unassigned until a matching node appears — and one sitting + on a non-matching node (e.g. started on the wrong backend) is moved to a matching + one. `nodes` is the set of alive nodes (node_id + capacity + engines). """ - alive = {n["node_id"] for n in nodes} - if not alive: + if not nodes: return {} # nowhere to place; leave as-is until a node appears - # Greedy: prefer the emptiest node. Stable tiebreak on node_id for determinism. - ranked = sorted(nodes, key=lambda n: (-node_free_vram(n), n["node_id"])) - target = ranked[0]["node_id"] + key_engines = key_engines or {} + by_id = {n["node_id"]: n for n in nodes} changes: dict[str, str] = {} for key in desired: - cur = assignments.get(key) - if cur is None or cur not in alive: - changes[key] = target + engine = key_engines.get(key, "vllm") + candidates = [n for n in nodes if node_supports(n, engine)] + if not candidates: + continue # no live node can run this engine; leave unassigned + cur_node = by_id.get(assignments.get(key)) + # Keep the current placement only if its node is alive AND engine-matching. + if cur_node is not None and node_supports(cur_node, engine): + continue + # Greedy: emptiest matching node; stable tiebreak on node_id. + target = sorted(candidates, key=lambda n: (-node_free_vram(n), n["node_id"]))[0] + changes[key] = target["node_id"] return changes class Scheduler: - """Leader-only loop: keep desired instances placed on live nodes.""" + """Leader-only loop: keep desired instances placed on live, engine-matching nodes. + + `registry` (optional) resolves each instance's engine so placement can match it + to a node that can run it. Without it, engines default to "vllm" — fine for a + single-engine fleet.""" + + def __init__(self, registry=None) -> None: + self.registry = registry + + def _key_engines(self) -> dict[str, str]: + if self.registry is None: + return {} + return {inst.key: getattr(inst, "engine", "vllm") for inst in self.registry.values()} async def reschedule_once(self, store, settings) -> dict[str, str]: """One placement pass. Reads desired-running instances, live nodes and @@ -73,7 +107,7 @@ async def reschedule_once(self, store, settings) -> dict[str, str]: from app.llmops.state import Desired desired = {k for k, v in desired_map.items() if v == Desired.RUNNING.value} - changes = place(desired, nodes, assignments) + changes = place(desired, nodes, assignments, self._key_engines()) for key, node_id in changes.items(): try: await store.set_assignment(key, node_id) diff --git a/apps/backend/app/main.py b/apps/backend/app/main.py index 22ec514..c9a37d7 100644 --- a/apps/backend/app/main.py +++ b/apps/backend/app/main.py @@ -36,7 +36,7 @@ from app.core.logging import setup_logging from app.core.settings import BackendSettings from app.core.store import LLMOpsStore -from app.llmops.launchers import EmbeddingLauncher, VllmLauncher +from app.llmops.launchers import EmbeddingLauncher, SglangLauncher, VllmLauncher from app.llmops.manager import ModelManager, build_registry from app.llmops.autoscaler import autoscaler_loop from app.llmops.scheduler import Scheduler @@ -94,6 +94,18 @@ async def _gpu_poll_loop(app: FastAPI, interval: float) -> None: await asyncio.sleep(interval) +async def _overlay_sync_loop(manager, interval: float = 5.0) -> None: + """HA Phase 7: per-node loop pulling the fleet's shared overlay from the store so + a model added/edited on any replica appears in this node's registry too (then an + engine-matching node can actuate it). No-op outside Postgres/HA mode.""" + while True: + try: + await manager.sync_overlay_from_store() + except Exception: + logger.exception("overlay sync pass failed") + await asyncio.sleep(interval) + + @asynccontextmanager async def lifespan(app: FastAPI): config_path = get_config_path() @@ -113,7 +125,7 @@ async def lifespan(app: FastAPI): # Base config.yaml + dynamically-added models (overlay), merged into one view. config = build_merged_config(config_path) - launchers = [VllmLauncher(), EmbeddingLauncher()] + launchers = [VllmLauncher(), SglangLauncher(), EmbeddingLauncher()] registry = build_registry(config, config_path, launchers) # `or` (not get's default) so an env var set-but-empty (as the compose env # passes it) still falls back instead of yielding "" — an empty router_url @@ -186,23 +198,27 @@ async def lifespan(app: FastAPI): async def _on_acquire() -> None: # Restore desired state on the leader only (a follower must not also start - # models). Skips anything adopt already found alive. + # models). Skips anything adopt already found alive. The per-node reconcile + # loop's converge_desired also restores it, but this gives the leader an + # immediate boot replay rather than waiting a poll interval. if settings.replay_desired: await manager.replay_desired() + # SINGLETON loops — global decisions, one place only: scheduler (placement), + # autoscaler (desired counts), load-monitor (global router view), pruning. + # The per-node loops (reconcile/actuation + gpu-poll) run on every replica; + # see below. leader_loops[:] = [ - asyncio.create_task(reconcile_loop(registry, http_client, settings, store, manager, notifier)), - asyncio.create_task(_gpu_poll_loop(app, settings.gpu_poll_interval)), asyncio.create_task( load_monitor_loop(app, registry, http_client, router_url, settings.load_poll_interval) ), asyncio.create_task(autoscaler_loop(app, manager, settings.autoscale_interval)), asyncio.create_task( - Scheduler().run(store, settings, settings.schedule_interval) + Scheduler(registry).run(store, settings, settings.schedule_interval) ), asyncio.create_task(_audit_prune_loop(store, settings.audit_max_rows)), asyncio.create_task(_config_versions_prune_loop(store, settings.config_versions_max)), ] - logger.info("Control loops started (leader)") + logger.info("Singleton control loops started (leader)") async def _on_release() -> None: for t in leader_loops: @@ -226,9 +242,25 @@ async def _on_release() -> None: app.state.node_agent = node_agent node_agent_task = asyncio.create_task(node_agent.run()) + # PER-NODE loops (HA Phase 7): reconcile/actuation + GPU polling run on EVERY + # replica, not just the leader — actuation must happen on the host that holds + # the GPU, and each node reports its own GPU. They converge only the instances + # assigned to this node (foreign_assignments gates ownership), so collapsed + # single-host is unchanged (one node owns + converges everything). + node_loops = [ + asyncio.create_task( + reconcile_loop(registry, http_client, settings, store, manager, notifier) + ), + asyncio.create_task(_gpu_poll_loop(app, settings.gpu_poll_interval)), + asyncio.create_task(_overlay_sync_loop(manager)), + ] + try: yield finally: + for t in node_loops: + t.cancel() + await asyncio.gather(*node_loops, return_exceptions=True) node_agent_task.cancel() try: await node_agent_task diff --git a/apps/backend/app/services/prometheus_targets.py b/apps/backend/app/services/prometheus_targets.py index 48d09b3..fdb88dc 100644 --- a/apps/backend/app/services/prometheus_targets.py +++ b/apps/backend/app/services/prometheus_targets.py @@ -24,14 +24,20 @@ from app.llmops.state import ModelKind, ModelState -def build_targets(instances: Iterable[ModelInstance]) -> list[dict]: +def build_targets(instances: Iterable[ModelInstance], node_host: str = "") -> list[dict]: """Build the Prometheus file_sd target list from registry instances. - One entry per ready vLLM instance. `targets` is the scrape address + One entry per ready LLM instance. `targets` is the scrape address (`host:port`); Prometheus appends the configured metrics_path (`/metrics`). Labels carry the group/instance identity and model tag so dashboards can join on something meaningful instead of the volatile `host:port`. + `node_host` (the node's advertised routable host, LLMOPS_NODE_HOST) is used as + the scrape host when set, so a Prometheus in a *different* container/host can + reach the instance — required for the multi-backend mixed deployment. Empty + (the default / collapsed single host) keeps the instance's own host (localhost), + which a netns-sharing Prometheus scrapes exactly as before. + Sorted by address so the serialized output is stable — the writer can then skip an identical rewrite and avoid churning the file (which would otherwise nudge Prometheus to re-read it every reconcile pass). @@ -41,13 +47,17 @@ def build_targets(instances: Iterable[ModelInstance]) -> list[dict]: if inst.kind != ModelKind.LLM or inst.state != ModelState.READY: continue group, _, instance_id = inst.key.partition("::") + host = node_host or inst.host targets.append( { - "targets": [f"{inst.host}:{inst.port}"], + "targets": [f"{host}:{inst.port}"], "labels": { "group": group, "instance_id": instance_id, "model_tag": inst.model_tag or "", + # Engine lets dashboards filter vLLM (vllm:*) vs SGLang (sglang:*) + # panels, since their metric names differ. + "engine": getattr(inst, "engine", "vllm"), }, } ) diff --git a/apps/backend/app/services/vllm_command.py b/apps/backend/app/services/vllm_command.py index dc5348b..8627535 100644 --- a/apps/backend/app/services/vllm_command.py +++ b/apps/backend/app/services/vllm_command.py @@ -159,3 +159,121 @@ def parse_vllm_command(command: str) -> dict[str, Any]: "model_config": model_config, "warnings": warnings, } + + +# SGLang flag (snake_case) -> engine-neutral EngineModelConfig key. The inverse of +# launchers._SGLANG_PARAM_MAP, so a pasted SGLang command stores the same typed +# fields a vLLM one would and the launcher re-derives the SGLang flag on launch. +_SGLANG_REVERSE_MAP = { + "context_length": "max_model_len", + "mem_fraction_static": "gpu_memory_utilization", + "tp_size": "tensor_parallel_size", +} + + +def parse_sglang_command(command: str) -> dict[str, Any]: + """Return {group, instance, model_config, warnings} from a + ``python -m sglang.launch_server …`` command. + + Differs from vLLM: the model is ``--model-path`` (not positional), LoRA is + ``--lora-paths NAME=PATH …`` (not ``--lora-modules``), and the three flags + SGLang renames are mapped back to engine-neutral keys + (``--context-length`` → ``max_model_len`` …). The result carries + ``model_config.engine = "sglang"``. + """ + warnings: list[str] = [] + if not command or not command.strip(): + raise ValueError("empty command") + + tokens = shlex.split(command, comments=False) + + env: dict[str, str] = {} + while tokens and _ENV_ASSIGN.match(tokens[0]): + k, _, v = tokens[0].partition("=") + env[k] = v + tokens = tokens[1:] + + # Drop leading non-flag tokens (python, -m, sglang.launch_server, …). + rest = tokens + while rest and not rest[0].startswith("-"): + rest = rest[1:] + + flags: dict[str, Any] = {} + i = 0 + while i < len(rest): + tok = rest[i] + if tok.startswith("--"): + name = tok[2:] + if name in ("lora-paths", "lora_paths"): + # nargs='+' of NAME=PATH; consume until the next --flag. + values: list[str] = [] + j = i + 1 + while j < len(rest) and not rest[j].startswith("--"): + values.append(rest[j]) + j += 1 + flags["lora_modules"] = [_parse_lora_module(v) for v in values] + flags["enable_lora"] = True + i = j + continue + if "=" in name: + k, _, v = name.partition("=") + flags[k] = _coerce(v) + elif i + 1 < len(rest) and not rest[i + 1].startswith("-"): + flags[name] = _coerce(rest[i + 1]) + i += 1 + else: + flags[name] = True + i += 1 + + norm: dict[str, Any] = {k.replace("-", "_"): v for k, v in flags.items()} + + # model tag: --model-path (SGLang) or --model (alias). + model_tag = norm.pop("model_path", None) or norm.pop("model", None) + port = norm.pop("port", None) + host = norm.pop("host", "localhost") + served = norm.pop("served_model_name", None) + + # Map the renamed flags back to engine-neutral keys (tp-size's long alias + # --tensor-parallel-size already normalises to tensor_parallel_size). + for sgl_key, neutral in _SGLANG_REVERSE_MAP.items(): + if sgl_key in norm: + norm[neutral] = norm.pop(sgl_key) + + cuda_device: int | None = None + if env.get("CUDA_VISIBLE_DEVICES"): + first = env["CUDA_VISIBLE_DEVICES"].split(",")[0] + cuda_device = int(first) if first.isdigit() else None + + if model_tag is None: + warnings.append("could not detect a model tag (expected `--model-path `)") + if port is None: + warnings.append("no --port found; defaulting to 30000") + port = 30000 # SGLang's default server port + if served: + warnings.append("--served-model-name ignored for naming; the router uses the group name") + + group = model_tag.split("/")[-1] if model_tag else (served or "model") + instance_id = _slug(served) if served else f"{_slug(group)}-1" + + model_config: dict[str, Any] = {"model_tag": model_tag, "engine": "sglang", **norm} + + return { + "group": group, + "instance": {"id": instance_id, "host": host, "port": port, "cuda_device": cuda_device}, + "model_config": model_config, + "warnings": warnings, + } + + +def parse_command(command: str, engine: str | None = None) -> dict[str, Any]: + """Engine-aware entry point for the dashboard's paste-a-command flow. + + ``engine`` ('vllm' | 'sglang') forces the parser; when omitted it is sniffed + from the command (``sglang.launch_server`` → sglang, else vLLM). + """ + eng = (engine or "").strip().lower() + if not eng: + eng = "sglang" if "sglang" in command.lower() else "vllm" + if eng == "sglang": + return parse_sglang_command(command) + return parse_vllm_command(command) diff --git a/apps/backend/tests/api/test_observability_routes.py b/apps/backend/tests/api/test_observability_routes.py index 3df25d4..bfc4587 100644 --- a/apps/backend/tests/api/test_observability_routes.py +++ b/apps/backend/tests/api/test_observability_routes.py @@ -100,11 +100,14 @@ def test_healthz(client): async def test_sse_generator_emits_snapshot_on_change(): """Test the SSE generator directly (avoids streaming a forever-loop endpoint).""" from tests.conftest import FAKE_CONFIG + from app.core.settings import BackendSettings from app.llmops.launchers import VllmLauncher - from app.llmops.manager import build_registry + from app.llmops.manager import ModelManager, build_registry registry = build_registry(FAKE_CONFIG, "config.yaml", [VllmLauncher()]) - gen = observability.model_snapshot_stream(registry, interval=0.01) + manager = ModelManager(registry, [VllmLauncher()], None, FAKE_CONFIG, "config.yaml", + BackendSettings(), store=None) # store=None -> uses local registry + gen = observability.model_snapshot_stream(manager, interval=0.01) first = await gen.__anext__() # initial snapshot always emitted assert first.startswith("data: ") diff --git a/apps/backend/tests/unit/test_converge_desired.py b/apps/backend/tests/unit/test_converge_desired.py new file mode 100644 index 0000000..77073ea --- /dev/null +++ b/apps/backend/tests/unit/test_converge_desired.py @@ -0,0 +1,95 @@ +"""Per-node actuation (HA Phase 7): converge_desired drives owned instances toward +their persisted desired state, skipping foreign (other-node-owned) keys.""" +import pytest + +from app.core.settings import BackendSettings +from app.llmops.launchers import EmbeddingLauncher, VllmLauncher +from app.llmops.manager import build_registry +from app.llmops.reconciler import converge_desired +from app.llmops.state import Desired, ModelState +from tests.conftest import FAKE_CONFIG + +pytestmark = pytest.mark.unit + +A = "Qwen3-0.6B::qwen3" +B = "Qwen3-0.6B::qwen3-2" + + +def _registry(): + return build_registry(FAKE_CONFIG, "config.yaml", [VllmLauncher(), EmbeddingLauncher()]) + + +class _DesiredStore: + def __init__(self, desired: dict[str, str]): + self._desired = desired + + async def list_instance_desired(self) -> dict[str, str]: + return dict(self._desired) + + +class _RecordingManager: + """Records the actuation calls converge_desired makes.""" + def __init__(self): + self.calls: list[tuple[str, str]] = [] + + async def start(self, key, **kw): self.calls.append(("start", key)) + async def stop(self, key, **kw): self.calls.append(("stop", key)) + async def sleep(self, key, **kw): self.calls.append(("sleep", key)) + async def wake(self, key, **kw): self.calls.append(("wake", key)) + + +async def _run(states: dict[str, ModelState], desired: dict[str, str], foreign=frozenset()): + reg = _registry() + for key, st in states.items(): + reg.get(key).state = st + mgr = _RecordingManager() + await converge_desired(reg, BackendSettings(), _DesiredStore(desired), mgr, set(foreign)) + return mgr.calls + + +async def test_starts_stopped_but_wanted(): + calls = await _run({A: ModelState.STOPPED}, {A: Desired.RUNNING.value}) + assert calls == [("start", A)] + + +async def test_stops_live_but_unwanted(): + calls = await _run({A: ModelState.READY}, {A: Desired.STOPPED.value}) + assert calls == [("stop", A)] + + +async def test_sleeps_ready_when_asleep_wanted(): + calls = await _run({A: ModelState.READY}, {A: Desired.ASLEEP.value}) + assert calls == [("sleep", A)] + + +async def test_wakes_sleeping_when_running_wanted(): + calls = await _run({A: ModelState.SLEEPING}, {A: Desired.RUNNING.value}) + assert calls == [("wake", A)] + + +async def test_failed_is_left_to_restart_logic(): + # desired=running + FAILED must NOT be started here (that is _process_restarts' + # job, with budget/backoff) — only STOPPED-but-wanted is converged. + calls = await _run({A: ModelState.FAILED}, {A: Desired.RUNNING.value}) + assert calls == [] + + +async def test_already_converged_is_noop(): + calls = await _run({A: ModelState.READY}, {A: Desired.RUNNING.value}) + assert calls == [] + + +async def test_foreign_keys_skipped(): + # A owned by another live node -> not actuated here; B (ours) is. + calls = await _run( + {A: ModelState.STOPPED, B: ModelState.STOPPED}, + {A: Desired.RUNNING.value, B: Desired.RUNNING.value}, + foreign={A}, + ) + assert calls == [("start", B)] + + +async def test_unknown_key_skipped(): + # A desired entry with no registry instance is ignored (no crash). + calls = await _run({}, {"Ghost::x": Desired.RUNNING.value}) + assert calls == [] diff --git a/apps/backend/tests/unit/test_launchers.py b/apps/backend/tests/unit/test_launchers.py index 4441ef2..fc59640 100644 --- a/apps/backend/tests/unit/test_launchers.py +++ b/apps/backend/tests/unit/test_launchers.py @@ -2,9 +2,11 @@ import pytest -from app.llmops.launchers import (EMBEDDING_KEY, EmbeddingLauncher, +from app.llmops.launchers import (CAP_LORA_MODULES, CAP_METRICS_SGLANG, + CAP_RUNTIME_LORA, CAP_SLEEP, EMBEDDING_KEY, + ENGINE_DEFAULT, EmbeddingLauncher, SglangLauncher, VllmLauncher, _write_effective_config, - build_vllm_cli_args) + build_sglang_cli_args, build_vllm_cli_args) from app.llmops.state import ModelKind from schema import RootConfig from tests.conftest import FAKE_CONFIG @@ -12,6 +14,20 @@ pytestmark = pytest.mark.unit +def _engine_config(engine: str) -> RootConfig: + """A one-group LLM config whose engine is `engine` (no field = default).""" + mc = {"model_tag": "org/m"} + if engine is not None: + mc["engine"] = engine + return RootConfig.model_validate({ + "server": {"host": "0.0.0.0", "port": 8887}, + "LLM_engines": {"G": { + "instances": [{"id": "a", "host": "localhost", "port": 8000}], + "model_config": mc, + }}, + }) + + def _lora_config() -> RootConfig: return RootConfig.model_validate( { @@ -245,3 +261,185 @@ def test_embedding_launcher_keys_and_spec(): assert spec.env["CUDA_VISIBLE_DEVICES"] == "1" assert "PYTHONPATH" in spec.env assert spec.probe_url == "http://localhost:8005/health" + + +# ---- multi-engine abstraction (docs/multi-backend-engine-design_zh-CN.md) ---- + +def test_engine_defaults_to_vllm_when_unset(): + """A config with no `engine` field is all vLLM = historical behaviour.""" + cfg = _engine_config(engine=None) + assert cfg.LLM_engines["G"].settings.engine == "vllm" + + +def test_vllm_launcher_claims_only_its_engine(): + """keys() returns groups for this launcher's engine; a non-vllm group is + skipped so a future SGLang launcher can claim it instead.""" + v = VllmLauncher() + assert v.engine == "vllm" + assert v.keys(_engine_config(None)) == ["G::a"] # default + assert v.keys(_engine_config("vllm")) == ["G::a"] # explicit + assert v.keys(_engine_config("sglang")) == [] # not mine + + +def test_vllm_launcher_declares_capabilities_on_spec(): + v = VllmLauncher() + spec = v.build_spec(_engine_config("vllm"), "config.yaml", "G::a") + assert spec.engine == "vllm" + assert CAP_SLEEP in spec.capabilities + assert spec.capabilities == v.capabilities + + +def test_embedding_launcher_registers_under_default_engine(): + e = EmbeddingLauncher() + assert e.engine == ENGINE_DEFAULT + assert e.capabilities == frozenset() + spec = e.build_spec(FAKE_CONFIG, "config.yaml", EMBEDDING_KEY) + assert spec.engine == ENGINE_DEFAULT + + +def test_engine_is_never_passed_to_vllm_cli(): + """`engine` is launcher-meta; it must not reach `vllm serve` (unknown arg).""" + args = build_vllm_cli_args({"model_tag": "org/m", "engine": "vllm", "kind": "chat"}) + assert "--engine" not in args + assert "--kind" not in args + + +# ---- SGLang launcher (docs/multi-backend-engine-design_zh-CN.md §5.2) -------- + +def _sglang_config(extra: dict | None = None) -> RootConfig: + mc = {"model_tag": "Qwen/Qwen3-0.6B", "engine": "sglang"} + if extra: + mc.update(extra) + return RootConfig.model_validate({ + "server": {"host": "0.0.0.0", "port": 8887}, + "LLM_engines": {"S": { + "instances": [{"id": "a", "host": "localhost", "port": 8100, "cuda_device": 2}], + "model_config": mc, + }}, + }) + + +def test_sglang_args_model_path_and_served_name(): + args = build_sglang_cli_args({"model_tag": "Qwen/Qwen3-0.6B"}) + assert args[:2] == ["--model-path", "Qwen/Qwen3-0.6B"] + # served-model-name defaults to model_tag so /v1/models + forward_name are stable. + assert "--served-model-name" in args + assert args[args.index("--served-model-name") + 1] == "Qwen/Qwen3-0.6B" + + +def test_sglang_args_translate_typed_params(): + # The three common params have different SGLang flag names. + args = build_sglang_cli_args({ + "model_tag": "org/m", "max_model_len": 4096, + "gpu_memory_utilization": 0.45, "tensor_parallel_size": 2, + }) + assert args[args.index("--context-length") + 1] == "4096" + assert args[args.index("--mem-fraction-static") + 1] == "0.45" + assert args[args.index("--tp-size") + 1] == "2" + # vLLM names must NOT appear. + assert "--max-model-len" not in args + assert "--gpu-memory-utilization" not in args + + +def test_sglang_args_bool_is_store_true_no_dual(): + args = build_sglang_cli_args({ + "model_tag": "org/m", "disable_radix_cache": True, "skip_server_warmup": False, + }) + assert "--disable-radix-cache" in args # True -> present + assert "--skip-server-warmup" not in args # False -> omitted + assert "--no-skip-server-warmup" not in args # never synthesise a --no- dual + + +def test_sglang_args_skip_router_only_keys(): + args = build_sglang_cli_args( + {"model_tag": "org/m", "engine": "sglang", "kind": "chat", + "routing_strategy": "session_affinity", "dtype": "bfloat16"}) + assert "--engine" not in args and "--kind" not in args + assert "--routing-strategy" not in args + assert args[args.index("--dtype") + 1] == "bfloat16" # real flags still pass through + + +def test_sglang_args_explicit_served_name_wins(): + args = build_sglang_cli_args({"model_tag": "org/m", "served_model_name": "my-name"}) + assert args[args.index("--served-model-name") + 1] == "my-name" + + +def test_sglang_launcher_claims_only_sglang_engine(): + s = SglangLauncher() + assert s.kind == ModelKind.LLM and s.engine == "sglang" + assert s.keys(_sglang_config()) == ["S::a"] + assert s.keys(FAKE_CONFIG) == [] # FAKE_CONFIG is all vLLM + assert VllmLauncher().keys(_sglang_config()) == [] # vLLM doesn't claim sglang group + + +def test_sglang_launcher_capabilities(): + # Runtime + static LoRA and sglang:* metrics are wired; sleep is absent in + # SGLang (autoscaler degrades to ready<->stopped). + caps = SglangLauncher().capabilities + assert CAP_RUNTIME_LORA in caps and CAP_LORA_MODULES in caps + assert CAP_METRICS_SGLANG in caps + assert CAP_SLEEP not in caps + + +def test_sglang_args_always_enable_metrics(): + # /metrics must be on (vLLM exposes it by default; SGLang needs the flag) so the + # router can scrape sglang:* for the autoscaler. + args = build_sglang_cli_args({"model_tag": "org/m"}) + assert "--enable-metrics" in args + + +def test_sglang_args_lora_enable_and_paths(): + # enable_lora -> --enable-lora; static lora_modules -> SGLang NAME=PATH form. + args = build_sglang_cli_args({ + "model_tag": "org/m", "enable_lora": True, + "lora_modules": [{"name": "sql", "path": "/lora/sql"}, + {"name": "math", "path": "/lora/math"}], + "max_lora_rank": 16, + }) + assert "--enable-lora" in args + i = args.index("--lora-paths") + assert args[i + 1] == "sql=/lora/sql" and args[i + 2] == "math=/lora/math" + assert args[args.index("--max-lora-rank") + 1] == "16" # other lora knobs pass through + assert "--lora-modules" not in args # not vLLM's JSON form + + +def test_sglang_args_runtime_lora_toggle_enables_lora(): + # Our runtime toggle (allow_runtime_lora) must turn on --enable-lora so SGLang + # accepts POST /load_lora_adapter, even with no static modules. + args = build_sglang_cli_args({"model_tag": "org/m", "allow_runtime_lora": True}) + assert "--enable-lora" in args + assert "--allow-runtime-lora" not in args # the toggle is not a flag + + +def test_sglang_args_no_lora_by_default(): + args = build_sglang_cli_args({"model_tag": "org/m"}) + assert "--enable-lora" not in args and "--lora-paths" not in args + + +def test_sglang_build_spec(): + spec = SglangLauncher().build_spec(_sglang_config({"max_model_len": 4096}), "config.yaml", "S::a") + assert spec.engine == "sglang" + assert spec.command[1:3] == ["-m", "sglang.launch_server"] + assert "--model-path" in spec.command + assert spec.command[spec.command.index("--context-length") + 1] == "4096" + # single-GPU cuda_device -> env, not a CLI flag; id dropped. + assert spec.env["CUDA_VISIBLE_DEVICES"] == "2" + assert "--cuda-device" not in spec.command and "--id" not in spec.command + assert spec.probe_url == "http://localhost:8100/health" + assert spec.host == "localhost" and spec.port == 8100 + assert "--host" in spec.command and spec.command[spec.command.index("--host") + 1] == "localhost" + + +def test_sglang_bind_host_env_overrides_only_the_bind_address(monkeypatch): + # Cross-container HA: LLMOPS_VLLM_BIND_HOST binds sglang to 0.0.0.0 (--host), but + # the probe + recorded host stay localhost; routers reach it via NODE_HOST. + monkeypatch.setenv("LLMOPS_VLLM_BIND_HOST", "0.0.0.0") + spec = SglangLauncher().build_spec(_sglang_config(), "config.yaml", "S::a") + assert spec.command[spec.command.index("--host") + 1] == "0.0.0.0" # binds all + assert spec.host == "localhost" # record unchanged + assert spec.probe_url == "http://localhost:8100/health" # local probe unchanged + + +def test_sglang_binds_configured_host_by_default(): + spec = SglangLauncher().build_spec(_sglang_config(), "config.yaml", "S::a") + assert spec.command[spec.command.index("--host") + 1] == "localhost" diff --git a/apps/backend/tests/unit/test_lora_hotload.py b/apps/backend/tests/unit/test_lora_hotload.py index fe052b0..145493f 100644 --- a/apps/backend/tests/unit/test_lora_hotload.py +++ b/apps/backend/tests/unit/test_lora_hotload.py @@ -6,6 +6,7 @@ import pytest import yaml +from app.llmops.launchers import VllmLauncher from app.llmops.manager import LoraRuntimeError, ModelConflict, ModelManager, ModelNotFound from app.llmops.registry import ModelRegistry from app.llmops.state import ModelKind, ModelState @@ -64,7 +65,9 @@ def _ready_instances(ports): def _mgr(config, client, registry): - return ModelManager(registry, [], client, config, "config.yaml", + # Register the vLLM launcher so runtime-LoRA capability resolves (the gate + # checks the group's engine supports it). + return ModelManager(registry, [VllmLauncher()], client, config, "config.yaml", BackendSettings(), overlay_path="overlay.json") diff --git a/apps/backend/tests/unit/test_manager_engine.py b/apps/backend/tests/unit/test_manager_engine.py new file mode 100644 index 0000000..5183d4e --- /dev/null +++ b/apps/backend/tests/unit/test_manager_engine.py @@ -0,0 +1,309 @@ +"""Multi-engine dispatch: the manager picks a launcher by (kind, engine), and the +engine threads from config -> registry -> instance. See +docs/multi-backend-engine-design_zh-CN.md.""" +import pytest + +from app.core.settings import BackendSettings +from app.llmops.instance import LaunchSpec +from app.llmops.launchers import (CAP_RUNTIME_LORA, CAP_SLEEP, EmbeddingLauncher, + VllmLauncher) +from app.llmops.manager import ModelConflict, ModelManager, build_registry +from app.llmops.state import ModelKind +from schema import load_config + +pytestmark = pytest.mark.unit + + +class _FakeLauncher: + """A minimal second LLM engine with no optional capabilities — stands in for a + future engine to exercise capability gating + dispatch. Uses the "sglang" + engine slot (a valid schema Literal) but declares no capabilities of its own.""" + kind = ModelKind.LLM + engine = "sglang" + capabilities = frozenset() + + def keys(self, config): + out = [] + for tag, eng in config.LLM_engines.items(): + if getattr(eng.settings, "engine", "vllm") != self.engine: + continue + out.extend(f"{tag}::{i.id}" for i in eng.instances) + return out + + def build_spec(self, config, config_path, key): + tag, _, iid = key.partition("::") + inst = next(i for i in config.LLM_engines[tag].instances if i.id == iid) + return LaunchSpec( + key=key, kind=self.kind, engine=self.engine, capabilities=self.capabilities, + command=["fake-serve"], env={}, log_path="x.log", + host=inst.host, port=inst.port, probe_url=f"http://{inst.host}:{inst.port}/health", + model_tag=config.LLM_engines[tag].settings.model_tag, + ) + +# Two groups: one default (vllm), one explicitly configured for a not-yet-built +# engine (sglang). Only VllmLauncher is registered, so the sglang group is simply +# not claimed — proving keys() filters by engine rather than erroring. +CONFIG_YAML = """ +server: + port: 8887 +LLM_engines: + Qwen3-0.6B: + instances: + - id: a + host: localhost + port: 8002 + model_config: + model_tag: Qwen/Qwen3-0.6B + FutureModel: + instances: + - id: a + host: localhost + port: 8010 + model_config: + model_tag: org/future + engine: sglang +""" + + +def _manager(tmp_path): + cfg_path = tmp_path / "config.yaml" + cfg_path.write_text(CONFIG_YAML, encoding="utf-8") + config = load_config(str(cfg_path)) + launchers = [VllmLauncher(), EmbeddingLauncher()] + registry = build_registry(config, str(cfg_path), launchers) + mgr = ModelManager( + registry, launchers, None, config, str(cfg_path), + BackendSettings(), store=None, overlay_path=str(tmp_path / "overlay.json"), + ) + return mgr + + +def test_build_registry_threads_engine_and_filters_unclaimed(tmp_path): + mgr = _manager(tmp_path) + # vLLM group is registered with engine="vllm" on the instance. + vllm_inst = mgr.registry.get("Qwen3-0.6B::a") + assert vllm_inst is not None + assert vllm_inst.engine == "vllm" + assert vllm_inst.spec.engine == "vllm" + # The sglang group has no launcher registered for it -> not in the registry. + assert mgr.registry.get("FutureModel::a") is None + + +def test_launcher_for_dispatches_by_kind_and_engine(tmp_path): + mgr = _manager(tmp_path) + inst = mgr.registry.get("Qwen3-0.6B::a") + launcher = mgr._launcher_for(inst) + assert isinstance(launcher, VllmLauncher) + assert (launcher.kind, launcher.engine) == (ModelKind.LLM, "vllm") + + +def test_observed_dict_includes_engine(tmp_path): + mgr = _manager(tmp_path) + inst = mgr.registry.get("Qwen3-0.6B::a") + assert inst.observed_dict()["engine"] == "vllm" + + +async def test_create_overlay_model_records_engine(tmp_path): + """A dashboard-added model (default engine) registers with engine='vllm'.""" + mgr = _manager(tmp_path) + inst = await mgr.create_overlay_model( + "NewGroup", + {"id": "x", "host": "localhost", "port": 8020}, + {"model_tag": "org/new"}, + ) + assert inst.engine == "vllm" + assert mgr.registry.get("NewGroup::x").engine == "vllm" + + +# ---- capability gating ------------------------------------------------------ + +FAKE_ENGINE_YAML = """ +server: + port: 8887 +LLM_engines: + Plain: + instances: + - id: a + host: localhost + port: 8030 + model_config: + model_tag: org/plain + engine: sglang + allow_runtime_lora: true +""" + + +def _manager_with_fake(tmp_path): + cfg_path = tmp_path / "config.yaml" + cfg_path.write_text(FAKE_ENGINE_YAML, encoding="utf-8") + config = load_config(str(cfg_path)) + launchers = [VllmLauncher(), _FakeLauncher(), EmbeddingLauncher()] + registry = build_registry(config, str(cfg_path), launchers) + return ModelManager( + registry, launchers, None, config, str(cfg_path), + BackendSettings(), store=None, overlay_path=str(tmp_path / "overlay.json"), + ) + + +def test_capabilities_resolved_per_engine(tmp_path): + mgr = _manager_with_fake(tmp_path) + # fake engine declares no capabilities; vllm would declare sleep + lora. + assert mgr._llm_engine_capabilities("Plain") == frozenset() + assert CAP_SLEEP not in mgr._llm_engine_capabilities("Plain") + + +def test_fake_engine_group_is_dispatched_to_its_launcher(tmp_path): + mgr = _manager_with_fake(tmp_path) + inst = mgr.registry.get("Plain::a") + assert inst is not None and inst.engine == "sglang" + assert isinstance(mgr._launcher_for(inst), _FakeLauncher) + + +async def test_load_lora_rejected_when_engine_lacks_capability(tmp_path): + """Even with allow_runtime_lora set in config, an engine without the + runtime_lora capability must be refused (gate on capability, not the flag).""" + mgr = _manager_with_fake(tmp_path) + with pytest.raises(ModelConflict, match="runtime LoRA"): + await mgr.load_lora("Plain", "adapter", "repo/adapter") + + +async def test_create_overlay_model_rejects_unregistered_engine(tmp_path): + # trtllm is a valid engine name in the schema, but no launcher is registered + # for it here — the manager must refuse cleanly, not KeyError into a 500. + mgr = _manager_with_fake(tmp_path) + with pytest.raises(ModelConflict, match="unsupported engine"): + await mgr.create_overlay_model( + "Brand", {"id": "z", "host": "localhost", "port": 8040}, + {"model_tag": "org/brand", "engine": "trtllm"}, + ) + + +# ---- per-engine runtime LoRA endpoint path ---------------------------------- + +def _inst(engine: str, port: int): + from app.llmops.instance import LaunchSpec, ModelInstance + spec = LaunchSpec(key=f"G::a", kind=ModelKind.LLM, engine=engine, capabilities=frozenset(), + command=[], env={}, log_path="x", host="localhost", port=port, + probe_url=f"http://localhost:{port}/health") + return ModelInstance(key="G::a", kind=ModelKind.LLM, engine=engine, host="localhost", + port=port, spec=spec) + + +async def test_post_lora_endpoint_path_per_engine(tmp_path): + from tests.conftest import FakeHTTPClient + client = FakeHTTPClient() + mgr = ModelManager(build_registry(load_config(_write_min_cfg(tmp_path)), str(tmp_path/'c.yaml'), + [VllmLauncher()]), + [VllmLauncher()], client, load_config(_write_min_cfg(tmp_path)), + str(tmp_path/'c.yaml'), BackendSettings(), store=None, + overlay_path=str(tmp_path/'o.json')) + await mgr._post_lora(_inst("vllm", 8002), "load", {"lora_name": "x", "lora_path": "/p"}) + await mgr._post_lora(_inst("sglang", 8100), "load", {"lora_name": "x", "lora_path": "/p"}) + urls = [u for u, _ in client.posts] + assert urls[0] == "http://localhost:8002/v1/load_lora_adapter" # vLLM: /v1 prefix + assert urls[1] == "http://localhost:8100/load_lora_adapter" # SGLang: no /v1 + + +def _write_min_cfg(tmp_path): + p = tmp_path / "c.yaml" + p.write_text("server:\n port: 8887\nLLM_engines: {}\n", encoding="utf-8") + return str(p) + + +# ---- HA Phase 7C: API write-intent when this node can't run the engine ------- + +from app.llmops.launchers import SglangLauncher # noqa: E402 +from app.llmops.state import Desired, ModelState # noqa: E402 + + +class _DesiredRecordingStore: + def __init__(self): + self.desired: dict[str, str] = {} + self.deleted_assignments: list[str] = [] + + async def set_instance_desired(self, key, val): self.desired[key] = val + async def delete_assignment(self, key): self.deleted_assignments.append(key) + + +SGLANG_ONLY_YAML = """ +server: + port: 8887 +LLM_engines: + S: + instances: [{ id: a, host: localhost, port: 8100 }] + model_config: { model_tag: org/s, engine: sglang } +""" + + +def _mgr_node_engines(tmp_path, node_engines): + cfg = tmp_path / "config.yaml" + cfg.write_text(SGLANG_ONLY_YAML, encoding="utf-8") + config = load_config(str(cfg)) + launchers = [VllmLauncher(), SglangLauncher(), EmbeddingLauncher()] + registry = build_registry(config, str(cfg), launchers) + store = _DesiredRecordingStore() + mgr = ModelManager(registry, launchers, None, config, str(cfg), + BackendSettings(node_engines=node_engines), store=store, + overlay_path=str(tmp_path / "o.json")) + return mgr, store + + +async def test_start_defers_when_node_cannot_run_engine(tmp_path): + # A vLLM-only node asked to start a SGLang model: write intent, do NOT spawn. + mgr, store = _mgr_node_engines(tmp_path, ["vllm"]) + inst = await mgr.start("S::a") + assert inst.state == ModelState.STOPPED # never spawned locally + assert inst.desired == Desired.RUNNING # intent recorded + assert store.desired["S::a"] == "running" + assert "S::a" in store.deleted_assignments # cleared so scheduler re-places + + +async def test_node_can_run_all_when_unspecified(tmp_path): + # No node_engines = runs any engine -> _node_can_run True (collapsed unchanged). + mgr, _ = _mgr_node_engines(tmp_path, []) + assert mgr._node_can_run("sglang") is True + assert mgr._node_can_run("vllm") is True + + +async def test_node_can_run_respects_advertised(tmp_path): + mgr, _ = _mgr_node_engines(tmp_path, ["vllm"]) + assert mgr._node_can_run("vllm") is True + assert mgr._node_can_run("sglang") is False + + +# ---- HA Phase 7: owning-node API URL (for log/metric proxy) ------------------ + +class _AssignNodeStore: + db_url = "postgresql://x" # marks HA mode + def __init__(self, assignments, nodes): + self._a = assignments + self._n = nodes + async def list_assignments(self): return dict(self._a) + async def list_nodes(self): return list(self._n) + async def list_instance_observed(self): return [] + + +async def test_owning_node_api_url_for_foreign_model(tmp_path): + store = _AssignNodeStore( + {"S::a": "sglang-node"}, + [{"node_id": "sglang-node", "api_url": "http://sglang-backend:5000"}], + ) + cfg = tmp_path / "c.yaml"; cfg.write_text("server:\n port: 8887\nLLM_engines: {}\n") + mgr = ModelManager(build_registry(load_config(str(cfg)), str(cfg), [VllmLauncher()]), + [VllmLauncher()], None, load_config(str(cfg)), str(cfg), + BackendSettings(instance_id="vllm-node"), store=store, + overlay_path=str(tmp_path/"o.json")) + assert await mgr.owning_node_api_url("S::a") == "http://sglang-backend:5000" + + +async def test_owning_node_api_url_none_when_local(tmp_path): + store = _AssignNodeStore( + {"S::a": "vllm-node"}, # owned by us + [{"node_id": "vllm-node", "api_url": "http://vllm-backend:5000"}], + ) + cfg = tmp_path / "c.yaml"; cfg.write_text("server:\n port: 8887\nLLM_engines: {}\n") + mgr = ModelManager(build_registry(load_config(str(cfg)), str(cfg), [VllmLauncher()]), + [VllmLauncher()], None, load_config(str(cfg)), str(cfg), + BackendSettings(instance_id="vllm-node"), store=store, + overlay_path=str(tmp_path/"o.json")) + assert await mgr.owning_node_api_url("S::a") is None # local -> read locally diff --git a/apps/backend/tests/unit/test_node_agent.py b/apps/backend/tests/unit/test_node_agent.py index ac6ca07..731ccea 100644 --- a/apps/backend/tests/unit/test_node_agent.py +++ b/apps/backend/tests/unit/test_node_agent.py @@ -16,8 +16,10 @@ def __init__(self): self.nodes = {} self.pruned = 0 - async def upsert_node(self, node_id, hostname, capacity, ttl, ts=None): - self.nodes[node_id] = {"hostname": hostname, "capacity": capacity, "ttl": ttl} + async def upsert_node(self, node_id, hostname, capacity, ttl, ts=None, engines=None, + api_url=None): + self.nodes[node_id] = {"hostname": hostname, "capacity": capacity, "ttl": ttl, + "engines": engines, "api_url": api_url} async def prune_nodes(self, ts=None): self.pruned += 1 @@ -72,3 +74,21 @@ def boom(): agent = NodeAgent(store, _settings()) await agent.heartbeat_once() assert store.nodes["node-A"]["capacity"] is None + + +async def test_heartbeat_advertises_engines(monkeypatch): + import app.llmops.node_agent as na + monkeypatch.setattr(na, "get_gpu_info", lambda: []) + store = FakeNodeStore() + agent = NodeAgent(store, BackendSettings(instance_id="node-A", node_engines=["sglang"])) + await agent.heartbeat_once() + assert store.nodes["node-A"]["engines"] == '["sglang"]' + + +async def test_heartbeat_engines_null_when_unspecified(monkeypatch): + import app.llmops.node_agent as na + monkeypatch.setattr(na, "get_gpu_info", lambda: []) + store = FakeNodeStore() + agent = NodeAgent(store, BackendSettings(instance_id="node-A")) # no node_engines + await agent.heartbeat_once() + assert store.nodes["node-A"]["engines"] is None diff --git a/apps/backend/tests/unit/test_prometheus_targets.py b/apps/backend/tests/unit/test_prometheus_targets.py index 1e8732b..33be7dd 100644 --- a/apps/backend/tests/unit/test_prometheus_targets.py +++ b/apps/backend/tests/unit/test_prometheus_targets.py @@ -33,6 +33,7 @@ def test_build_targets_only_includes_ready_llm(): assert entry["labels"]["group"] == "Qwen3-0.6B" assert entry["labels"]["instance_id"] == "qwen3" assert entry["labels"]["model_tag"] == "Qwen/Qwen3-0.6B" + assert entry["labels"]["engine"] == "vllm" # default; lets dashboards filter by engine def test_build_targets_excludes_embedding_server(): @@ -97,3 +98,19 @@ async def test_manager_writes_ready_targets_when_path_set(tmp_path): assert await mgr.write_prometheus_targets() is True written = json.loads(open(path).read()) assert [t["targets"][0] for t in written] == ["localhost:8002"] + + +def test_build_targets_uses_node_host_when_set(): + # Multi-backend: a routable node_host replaces the instance's localhost so a + # Prometheus in another container can scrape it. + reg = _registry() + reg.get(HEALTHY).state = ModelState.READY + targets = build_targets(reg.values(), node_host="mixed-vllm-backend") + assert targets[0]["targets"] == ["mixed-vllm-backend:8002"] + + +def test_build_targets_defaults_to_instance_host(): + # Empty node_host (collapsed / main compose) keeps localhost — unchanged. + reg = _registry() + reg.get(HEALTHY).state = ModelState.READY + assert build_targets(reg.values())[0]["targets"] == ["localhost:8002"] diff --git a/apps/backend/tests/unit/test_scheduler.py b/apps/backend/tests/unit/test_scheduler.py index edc0755..ab3fa14 100644 --- a/apps/backend/tests/unit/test_scheduler.py +++ b/apps/backend/tests/unit/test_scheduler.py @@ -76,3 +76,54 @@ async def set_assignment(self, key, node_id, ts=None): async def test_reschedule_noop_without_store(): assert await Scheduler().reschedule_once(None, BackendSettings()) == {} + + +# ---- engine-aware placement (HA Phase 7) ------------------------------------ + +from app.llmops.scheduler import node_supports # noqa: E402 + + +def _enode(node_id, engines, *gpus): + n = _node(node_id, *gpus) + n["engines"] = json.dumps(engines) if engines is not None else None + return n + + +def test_node_supports_unspecified_runs_any(): + assert node_supports({"engines": None}, "sglang") is True + assert node_supports({}, "vllm") is True + + +def test_node_supports_matches_advertised(): + n = _enode("n", ["sglang"]) + assert node_supports(n, "sglang") is True + assert node_supports(n, "vllm") is False + + +def test_place_matches_engine_to_capable_node(): + nodes = [_enode("vllm-node", ["vllm"], (8192, 0)), + _enode("sglang-node", ["sglang"], (8192, 0))] + changes = place({"G::v", "G::s"}, nodes, {}, + key_engines={"G::v": "vllm", "G::s": "sglang"}) + assert changes["G::v"] == "vllm-node" + assert changes["G::s"] == "sglang-node" + + +def test_place_leaves_unassigned_when_no_capable_node(): + nodes = [_enode("vllm-node", ["vllm"], (8192, 0))] + changes = place({"G::s"}, nodes, {}, key_engines={"G::s": "sglang"}) + assert "G::s" not in changes # no sglang node -> stays unassigned + + +def test_place_moves_model_off_wrong_engine_node(): + # Started on the vllm node by mistake; scheduler moves it to the sglang node. + nodes = [_enode("vllm-node", ["vllm"], (8192, 0)), + _enode("sglang-node", ["sglang"], (8192, 0))] + changes = place({"G::s"}, nodes, {"G::s": "vllm-node"}, key_engines={"G::s": "sglang"}) + assert changes["G::s"] == "sglang-node" + + +def test_place_keeps_well_placed_engine_match(): + nodes = [_enode("sglang-node", ["sglang"], (8192, 0))] + changes = place({"G::s"}, nodes, {"G::s": "sglang-node"}, key_engines={"G::s": "sglang"}) + assert changes == {} # already on a matching live node -> no churn diff --git a/apps/backend/tests/unit/test_vllm_command.py b/apps/backend/tests/unit/test_vllm_command.py index aca5d21..ebe632e 100644 --- a/apps/backend/tests/unit/test_vllm_command.py +++ b/apps/backend/tests/unit/test_vllm_command.py @@ -1,6 +1,6 @@ import pytest -from app.services.vllm_command import parse_vllm_command +from app.services.vllm_command import parse_command, parse_sglang_command, parse_vllm_command pytestmark = pytest.mark.unit @@ -33,3 +33,44 @@ def test_parse_lora_modules_json_form(): ) mods = parse_vllm_command(cmd)["model_config"]["lora_modules"] assert mods == [{"name": "sql", "path": "repo/sql", "base_model_name": "org/m"}] + + +def test_parse_sglang_command_maps_flags_and_sets_engine(): + p = parse_sglang_command( + "CUDA_VISIBLE_DEVICES=1 python -m sglang.launch_server --model-path org/m " + "--port 8030 --context-length 4096 --mem-fraction-static 0.85 --tp-size 2" + ) + mc = p["model_config"] + assert mc["model_tag"] == "org/m" + assert mc["engine"] == "sglang" + # SGLang's renamed flags map back to engine-neutral keys. + assert mc["max_model_len"] == 4096 + assert mc["gpu_memory_utilization"] == 0.85 + assert mc["tensor_parallel_size"] == 2 + assert p["instance"]["port"] == 8030 + assert p["instance"]["cuda_device"] == 1 + # No SGLang-renamed keys leak through un-mapped. + assert "context_length" not in mc and "mem_fraction_static" not in mc + + +def test_parse_sglang_lora_paths(): + p = parse_sglang_command( + "python -m sglang.launch_server --model-path org/m " + "--lora-paths sql=repo/sql fin=/models/fin --port 8031" + ) + mc = p["model_config"] + assert mc["enable_lora"] is True + assert mc["lora_modules"] == [ + {"name": "sql", "path": "repo/sql"}, + {"name": "fin", "path": "/models/fin"}, + ] + assert p["instance"]["port"] == 8031 + + +def test_parse_command_dispatch_by_sniffing_and_hint(): + # Sniffed from the command text. + assert parse_command("python -m sglang.launch_server --model-path org/m")["model_config"]["engine"] == "sglang" + # vLLM path leaves engine unset (frontend defaults to vllm). + assert "engine" not in parse_command("vllm serve org/m --port 8001")["model_config"] + # Explicit hint wins even for an ambiguous command. + assert parse_command("--model-path org/m --port 8030", engine="sglang")["model_config"]["engine"] == "sglang" diff --git a/apps/frontend_llmops/src/components/AddModelDialog.vue b/apps/frontend_llmops/src/components/AddModelDialog.vue index 161ecbb..1e01d33 100644 --- a/apps/frontend_llmops/src/components/AddModelDialog.vue +++ b/apps/frontend_llmops/src/components/AddModelDialog.vue @@ -40,6 +40,24 @@ const host = ref('localhost') const port = ref(8000) const cudaDevice = ref(null) const modelTag = ref('') +// Inference engine for the whole group. Different engines map the same concepts to +// different launch flags; the launcher translates. vLLM is the default. +const engine = ref('vllm') +const ENGINE_OPTIONS = ['vllm', 'sglang', 'llamacpp', 'trtllm'] +// Which engine's command the paste box expects — drives the example placeholder and +// is sent to the parser so an ambiguous command still parses for the right engine. +const pasteEngine = ref('vllm') +const PASTE_PLACEHOLDER: Record = { + vllm: 'CUDA_VISIBLE_DEVICES=0 vllm serve Qwen/Qwen2.5-3B-Instruct --port 8020 --dtype float16 --max-model-len 4096 --gpu-memory-utilization 0.85', + sglang: 'CUDA_VISIBLE_DEVICES=0 python -m sglang.launch_server --model-path Qwen/Qwen3-0.6B --port 8030 --context-length 4096 --mem-fraction-static 0.85', +} +const pastePlaceholder = computed(() => PASTE_PLACEHOLDER[pasteEngine.value] ?? PASTE_PLACEHOLDER.vllm) +// Capability gating — the two engines support different things, so the form only +// shows what the selected engine actually has (see launchers.py CAP_* flags). +const engineIsVllm = computed(() => engine.value === 'vllm') +const engineIsSglang = computed(() => engine.value === 'sglang') +const engineHasSleep = computed(() => engine.value === 'vllm') // /sleep + /wake_up (vLLM only) +const engineHasKvShare = computed(() => engine.value === 'vllm') // OffloadingConnector (vLLM only) const params = ref<{ key: string; value: string }[]>([]) // Router-only load-balancing policy for the group. Lives in model_config but is // NOT a vLLM flag, so it's edited as its own field and kept out of the raw param @@ -164,6 +182,8 @@ function reset() { port.value = 8000 cudaDevice.value = null modelTag.value = '' + engine.value = 'vllm' + pasteEngine.value = 'vllm' params.value = [] routingStrategy.value = '' kvShared.value = false @@ -209,6 +229,7 @@ function prefillForEdit() { port.value = cfg.port cudaDevice.value = cfg.cuda_device ?? null modelTag.value = String(cfg.settings.model_tag ?? '') + engine.value = String(cfg.settings.engine ?? 'vllm') routingStrategy.value = String(cfg.settings.routing_strategy ?? '') kvShared.value = isKvShared(cfg.settings) sleepMode.value = !!cfg.settings.enable_sleep_mode @@ -216,6 +237,7 @@ function prefillForEdit() { Object.entries(cfg.settings).filter( ([k2]) => k2 !== 'model_tag' && + k2 !== 'engine' && k2 !== 'routing_strategy' && k2 !== 'kv_transfer_config' && k2 !== 'enable_sleep_mode', @@ -245,13 +267,14 @@ async function parse() { if (!command.value.trim() || parsing.value) return parsing.value = true try { - const p = await api.parseCommand(command.value) + const p = await api.parseCommand(command.value, pasteEngine.value) group.value = p.group instanceId.value = p.instance.id host.value = p.instance.host port.value = p.instance.port cudaDevice.value = p.instance.cuda_device modelTag.value = String(p.model_config.model_tag ?? '') + engine.value = String((p.model_config as Record).engine ?? 'vllm') routingStrategy.value = String( (p.model_config as Record).routing_strategy ?? '', ) @@ -261,6 +284,7 @@ async function parse() { Object.entries(p.model_config).filter( ([k]) => k !== 'model_tag' && + k !== 'engine' && k !== 'routing_strategy' && k !== 'kv_transfer_config' && k !== 'enable_sleep_mode', @@ -316,10 +340,12 @@ function onPickLora(l: LoraModule) { if (maxRank > 0) setParam('max_lora_rank', String(maxRank)) } -// ---- Tool-calling presets (see docs/vllm_auto_tool_整理.md) ---- +// ---- Tool-calling presets ---- +// Recommended (tool_call_parser, reasoning_parser) by model family — the parser +// names differ per engine (vLLM: docs/vllm_auto_tool_整理.md; SGLang: the +// `--tool-call-parser` / `--reasoning-parser` choices from sglang.launch_server). const showToolHint = ref(false) -// Recommended (tool_call_parser, reasoning_parser) by model family. -const TOOL_PRESETS = [ +const VLLM_TOOL_PRESETS = [ { label: 'Qwen2.5 / QwQ', parser: 'hermes', reasoning: '' }, { label: 'addModel.qwen3Thinking', parser: 'hermes', reasoning: 'qwen3' }, { label: 'Qwen3-Coder', parser: 'qwen3_xml', reasoning: '' }, @@ -329,14 +355,33 @@ const TOOL_PRESETS = [ { label: 'DeepSeek-V3/R1', parser: 'deepseek_v3', reasoning: '' }, { label: 'GLM-4.5/4.6', parser: 'glm45', reasoning: '' }, ] +const SGLANG_TOOL_PRESETS = [ + { label: 'addModel.toolAuto', parser: 'auto', reasoning: '' }, + { label: 'Qwen2.5 / Qwen3', parser: 'qwen25', reasoning: '' }, + { label: 'addModel.qwen3Thinking', parser: 'qwen25', reasoning: 'qwen3' }, + { label: 'Qwen3-Coder', parser: 'qwen3_coder', reasoning: '' }, + { label: 'Llama 3.x', parser: 'llama3', reasoning: '' }, + { label: 'Mistral', parser: 'mistral', reasoning: '' }, + { label: 'DeepSeek-V3', parser: 'deepseekv3', reasoning: 'deepseek-r1' }, + { label: 'GLM-4.5/4.6', parser: 'glm45', reasoning: 'glm45' }, + { label: 'Kimi K2', parser: 'kimi_k2', reasoning: '' }, + { label: 'GPT-OSS', parser: 'gpt-oss', reasoning: 'gpt-oss' }, +] +const toolPresets = computed(() => (engineIsSglang.value ? SGLANG_TOOL_PRESETS : VLLM_TOOL_PRESETS)) function setParam(key: string, value: string) { const existing = params.value.find((p) => p.key === key) if (existing) existing.value = value else params.value.push({ key, value }) } function applyToolPreset(parser: string, reasoning: string) { - setParam('enable_auto_tool_choice', 'true') - setParam('tool_call_parser', parser) + if (engineIsSglang.value) { + // SGLang enables tool calling just by setting --tool-call-parser; there's no + // --enable-auto-tool-choice (and 'auto' lets it detect from the chat template). + setParam('tool_call_parser', parser) + } else { + setParam('enable_auto_tool_choice', 'true') + setParam('tool_call_parser', parser) + } if (reasoning) setParam('reasoning_parser', reasoning) } @@ -414,6 +459,30 @@ function applyAccelPreset(name: 'latency' | 'throughput') { showAccel.value = true } +// ---- SGLang acceleration (native flags; the launcher passes these through +// kebab-cased, except gpu_memory_utilization/max_model_len which it maps to +// --mem-fraction-static / --context-length). Choices verified against +// `sglang.launch_server --help` (v0.5.x). ---- +// Cleared by "clear"; gpu_memory_utilization/max_model_len are kept (deliberate +// memory/context settings, like vLLM's gpu_memory_utilization). +const SGLANG_ACCEL_KEYS = [ + 'chunked_prefill_size', 'max_running_requests', 'schedule_policy', + 'disable_radix_cache', 'kv_cache_dtype', 'attention_backend', + 'enable_torch_compile', 'cuda_graph_max_bs', 'stream_interval', 'quantization', +] +function clearSglangAccel() { + for (const k of SGLANG_ACCEL_KEYS) clearParam(k) +} +const SGLANG_ACCEL_PRESETS: Record> = { + latency: { schedule_policy: 'lpm', chunked_prefill_size: '8192', stream_interval: '1' }, + throughput: { schedule_policy: 'lpm', chunked_prefill_size: '16384', max_running_requests: '256' }, +} +function applySglangAccelPreset(name: 'latency' | 'throughput') { + clearSglangAccel() + for (const [k, v] of Object.entries(SGLANG_ACCEL_PRESETS[name]!)) setParam(k, v) + showAccel.value = true +} + async function submit() { if (!canSubmit.value || creating.value) return creating.value = true @@ -421,12 +490,15 @@ async function submit() { lora_modules?: LoraModule[] kv_transfer_config?: KvTransferConfig enable_sleep_mode?: boolean + engine?: string } = { model_tag: modelTag.value, + // Engine for the whole group; default vllm keeps existing configs unchanged. + engine: engine.value, } for (const { key: k, value } of params.value) { - // `lora_modules`, `routing_strategy` and `kv_transfer_config` have dedicated - // editors below — never let a raw param (e.g. a stray "" from a null, or a + // `lora_modules`, `routing_strategy`, `kv_transfer_config`, `engine` have + // dedicated editors — never let a raw param (e.g. a stray "" from a null, or a // leaked key) stomp them via the generic param list. const kk = k.trim() if ( @@ -434,7 +506,8 @@ async function submit() { kk !== 'lora_modules' && kk !== 'routing_strategy' && kk !== 'kv_transfer_config' && - kk !== 'enable_sleep_mode' + kk !== 'enable_sleep_mode' && + kk !== 'engine' ) settings[kk] = coerce(value) } @@ -443,10 +516,10 @@ async function submit() { if (routingStrategy.value) settings.routing_strategy = routingStrategy.value // Cross-instance KV-cache sharing: write the OffloadingConnector preset when // the toggle is on; otherwise leave it unset so each instance keeps its own KV. - if (kvShared.value) settings.kv_transfer_config = KV_SHARE_PRESET + if (kvShared.value && engineHasKvShare.value) settings.kv_transfer_config = KV_SHARE_PRESET // Sleep-mode warm-standby tier: launch with --enable-sleep-mode + dev mode so // the instance can be slept/woken (see docs/autoscaling-design_zh-CN.md). - if (sleepMode.value) settings.enable_sleep_mode = true + if (sleepMode.value && engineHasSleep.value) settings.enable_sleep_mode = true // Mounted adapters: keep only filled rows; drop the empty base_model_name field. const cleanLoras = loras.value .filter((l) => l.name.trim() && l.path.trim()) @@ -505,16 +578,38 @@ async function submit() {
- +
+ + +
+ +
+