Skip to content

Latest commit

 

History

History
462 lines (388 loc) · 25 KB

File metadata and controls

462 lines (388 loc) · 25 KB

Several machines as one coderai

Complete orchestration, distribution and escalation of remotizable advanced inference — this is the page where the three words earn their keep: orchestration (nodes as engines, pools), distribution (one model over several machines), escalation (hosts you own, then pods rented by the second).

Three things are possible once coderai runs on more than one machine, and they answer three different needs:

Need Mechanism Where it is set
More cards, one admin: a model runs somewhere on the network, requests are routed there by capability and load Cluster nodes — other coderai installs used as engines of this one Settings → Cluster; per model: Engine / card
One GGUF too big for any single machine llama.cpp RPC — cards on other machines join this machine's layer split Settings → Cluster → RPC servers; per model: Cards on other machines
One HF model too big for any single machine, served with continuous batching vLLM over Ray (pipeline parallel across nodes) / SGLang multi-node Settings → vLLM / ktransformers; per model: the vLLM / kt block
Several machines serving the same capability or the same host model, with failover Pools of remotes and hosts Settings → Remote capabilities; per model: More machines

Everything is visible on Admin → Cluster: engines here, nodes, RPC servers, hosts, capability remotes and pods, each with whether it answers now.

The Models page never lets you save a pair that cannot work: pick a backend and the engines that cannot run it are greyed out; pin an engine and the backends it cannot run are greyed out; an impossible pair already selected shows a red note and Save refuses (codai/cluster/compat.py is the table, validate_engine_pin applies it again server-side).


1. Cluster nodes — another coderai as an engine of this one

The front already runs N engine processes it cannot see inside: it polls /internal/engine-state, hands each the models it owns, proxies by capability and load. A node is that, one network hop away: a whole coderai (its own front, cards, engines, thermal protection, admin) that this front treats as one more engine.

On the node

Nothing to install. Create an API token on its Tokens page — that is what the head holds. Optionally, in Settings → Cluster:

  • This node's name — what the head calls it (defaults to the hostname); also the key of a model's per-node path.
  • Let other heads use this install as a node — on by default; off makes /cluster/* answer 401 to everyone.

If the node serves HTTPS with a self-signed certificate, copy its PEM: the head pastes it in the node row (verify = pem).

Zero-config: a shared token and mDNS

The manual way above (a token per node, its URL on the head) still works and is what you use across subnets or VPNs. On one LAN there is a shorter way — the same on every box, head and nodes alike, Settings → Cluster:

  1. Cluster token — one secret, typed on every install. A node accepts it on its /cluster/* endpoints exactly like one of its own API tokens; a head uses it as the API key of every node it discovers, and of any listed node whose token field is blank.
  2. Find nodes on the LAN (mDNS) — the install announces itself as a _coderai._tcp.local. service (name, version, port, capabilities) and browses for the others. The announcement carries an HMAC fingerprint of the token, never the token: peers with the same fingerprint are members, everyone else is merely seen. The Cluster page lists both.
  3. Auto-join members as nodes (on by default) — every member that serves becomes a cluster node of this head, with no row to fill in. It leaves when its announcement stops; a node listed by hand under the same name is never duplicated.
"cluster": {"enabled": true, "token": "…", "discovery": true, "auto_join": true}

Every box can have all three on: each is then a head of the others and a node of the others, and a model pinned to a name runs where that name is.

mDNS is link-local multicast: it does not cross routers, and a container only sees it on the host's network — coderai-docker --host-network (the server then listens on the host port directly; RPC and Ray ports need no mapping either). The zeroconf package is in the images and in requirements.txt; without it the Cluster page says discovery is off.

On the head

Settings → Cluster → Use the nodes below as engines, one row per node: name, URL, the node's token, how to verify its certificate (system for a public CA, pem with the pasted certificate, off for plain http / a LAN you trust), and optionally a capability list that narrows what the head will send it (blank = whatever the node reports).

"cluster": {
  "enabled": true,
  "nodes": [
    {"name": "box2", "url": "https://box2:8776", "api_key": "…", "verify": "pem",
     "ca_pem": "-----BEGIN CERTIFICATE-----…", "capabilities": []}
  ]
}

The head then:

  • polls GET /cluster/state on the node (its front answers from its own registry — cheap, never waits on a busy GPU there): loaded models, VRAM summed over its engines, tasks, and the union of its engines' capabilities (a node with an NVIDIA and a Radeon engine offers both transformers and gguf; the node's own front then picks the card);
  • assigns models to it like to any engine — a per-model Engine / card pin by node name, or the default engine, or round-robin among compatible engines — and pushes the assignment with the models' entries (POST /cluster/reload-config); a model the node's catalogue lacks is registered in memory on every engine there, under a path that exists on the node, else the HuggingFace id / URL coderai can recover for it (the hub cache path encodes the repo). Set Path on the node on the model when neither applies;
  • proxies requests there with the node's token, never the caller's credentials; the node applies its own queue, rate limits and thermal rules;
  • forwards Load / Unload from the head's Models page to the node (/cluster/model-load, /cluster/model-unload — the node runs the action under a short-lived admin session of its own).

Nodes appear in the Engine / card select as name (cluster node), on the Tasks page tiles, and on the Cluster page with their engines, VRAM and last error (401: the node refused the token, 404: not a coderai front, a connection error…). Changing the node list in Settings takes effect on the next poll; no restart.

What a node is not: it never rents pods on the head's behalf and never nests nodes of its own — the head sees the node's cards, not the node's remotes.

2. One GGUF over several machines — llama.cpp RPC

llama.cpp's RPC backend makes a card on another machine one more ggml device: register host:port where an rpc-server listens and the usual layer split (tensor_split) spreads the model over local and remote cards alike. Per token only the activations between the layers on each side cross the wire, so decode over a LAN is workable; loading and prompt processing are wire-bound. 1 GbE hurts (~100 MB/s: a 20 GB slice takes minutes to load and long prompts crawl); 10 GbE is fine. The protocol has no authentication: bind it to a LAN or WireGuard address, never the internet.

The machine lending its cards

Settings → Cluster → RPC servers this machine runs: one row per process — bind host/port, the ggml device (CUDA0, Vulkan1 … blank = the server's default), an optional memory cap, threads for a CPU server. The front starts them, restarts them if they die (up to a crash-loop limit), and advertises advertise_host:port on /cluster/state, so a head lists them on its Models page. A machine that is also a cluster node does this from the same Settings section; a machine that is only lending cards runs coderai the same way and simply has no models.

The binary is built by packaging/build-rpc-server.sh from the same llama.cpp the bundled llama-cpp-python vendors (the RPC protocol is versioned; a client and a server from different commits refuse each other). build.sh runs it after building llama-cpp-python; the OCI image and the coderai-llama pod image ship it. The llama image also runs it beside the API when started with CODERAI_RPC_SERVER_PORT=50052 (CODERAI_RPC_SERVER_DEVICE, CODERAI_RPC_SERVER_MEM_MB, CODERAI_RPC_SERVER_ONLY=1 for a card-only container) — a pod with a direct TCP port, or a host you run the image on.

"cluster": {"rpc_servers": [
  {"name": "3090", "host": "10.0.0.2", "port": 50052, "device": "CUDA0", "mem_gb": 0},
  {"name": "rx580", "host": "10.0.0.2", "port": 50053, "device": "Vulkan1"}
]}

The model

On the model's page, Cards on other machines: host:port, host2:port. Known servers (this machine's and every node's) are suggested. The model is split by definition; the ratio in Weight distribution reads local cards first, then these servers in the listed order, e.g. 0.6,0.4 for one local card and one RPC server. Blank ratio = proportional to free memory, RPC devices included (the VRAM and Speed strategies both see them).

{"path": "/AI/guffcache/big-Q4_K_M.gguf", "rpc_servers": "10.0.0.2:50052, 10.0.0.3:50052",
 "tensor_split": "0.4,0.3,0.3"}

How the split is made (same page, Cards on other machines): layer (default) gives each device whole layers, and per token only the activations between consecutive layers cross to the next device — what Ethernet carries well. Row is llama.cpp's tensor parallelism: every matrix is split by rows over all the devices, local cards and RPC servers alike, so they work on every layer together — at the price of an all-reduce per layer. On one PCIe bus it is faster; over a network it wants a fabric measured in tens of Gb/s (10 GbE is the floor, RDMA the point), and on 1 GbE it is slower than layer split. "split_mode": "row" in models.json.

Under the hood (codai/backends/ggml_rpc.py): the RPC registration is process-wide and permanent, so from the first registration on every load in that engine gets an explicit device list — this machine's cards plus exactly the servers the model asked for. A model that names servers the build cannot use fails with the reason (built without GGML_RPC, rpc-server at … did not answer) rather than quietly loading on local cards only. A model on a cluster node with rpc_servers set spreads from that node; the servers are addressed from there.

The bundled llama-cpp-python must have been built with -DGGML_RPC=ON (build.sh and the images do; an older install: rebuild it).

3. One HF model over several machines — vLLM on Ray, SGLang multi-node

vLLM and SGLang spread a model over machines on their own; coderai does the choreography: starts the peers, tells them where rank 0 is, waits for them, launches, and tears everything down with the service.

On ordinary Ethernet keep tensor parallel inside a machine (an all-reduce every layer wants NVLink/PCIe) and make pipeline parallel the number of machines (one activation transfer per layer boundary). Tensor parallel may span the nodes — set it larger than this machine's cards and Ray places the ranks across the cluster, the two combining as tensor × pipeline GPUs in total — which pays off on 10 GbE and up and is the right choice on RDMA / InfiniBand (NCCL_IB_*, NCCL_SOCKET_IFNAME in the node's start command environment). Every machine must run the same build — the coderai-vllm image on all of them is the easy way.

vLLM

Settings → vLLM → Several machines over Ray sets the defaults; a model's own vLLM block (visible when its backend is vLLM) overrides them: tensor parallel, pipeline parallel, executor, an existing Ray cluster to join, memory fraction, extra args, and the nodes — one per line, name | start command | stop command | gpus. With nodes listed (or pipeline

1, or executor = ray) coderai:

  1. starts a Ray head here (ray start --head --port 6379 --node-ip-address <advertise host>), unless ray_address names an existing cluster;
  2. runs each node's start command with {ray_address}, {head} and {port} filled in — e.g. ssh box2 docker run -d --rm --name vllm-worker --gpus all --network host -e CODERAI_RAY_ADDRESS={ray_address} ghcr.io/nextime/coderai-vllm:latest (the image joins as a Ray worker when CODERAI_RAY_ADDRESS is set);
  3. waits until Ray reports tensor × pipeline GPUs (nodes_ready_timeout_s);
  4. launches vLLM with --pipeline-parallel-size N --distributed-executor-backend ray.

Stopping the service runs the nodes' stop commands and ray stop.

{"path": "Org/Huge-70B", "backend": "vllm",
 "vllm": {"tensor_parallel_size": 2, "pipeline_parallel_size": 2,
          "nodes": [{"name": "box2", "gpus": 2,
                     "start_cmd": "ssh box2 docker run -d --rm --name vllm-worker --gpus all --network host -e CODERAI_RAY_ADDRESS={ray_address} ghcr.io/nextime/coderai-vllm:latest",
                     "stop_cmd": "ssh box2 docker stop vllm-worker"}]}}

SGLang (ktransformers)

Settings → ktransformers → Several machines, or the model's kt block: nnodes, tp_size (GPUs across all nodes), the rank-0 address (dist_init_addr, default <advertise host>:20000), and the node commands for ranks 1… — name | start command | stop command, templated with {rank}, {nnodes}, {dist_init_addr}, {model}. Ranks 1… are started first; rank 0 launches here with --nnodes N --node-rank 0 --dist-init-addr.

4. Every generation capability — distributing the work

An image, video, audio, embedding or OCR request often carries more than one unit of work. Tick Distribute the work of one request across engines and nodes on the model (or "distribute": {"enabled": true} in models.json) and the front splits it over every engine that has the model — here and on cluster nodes — sends the parts concurrently and merges the answers:

Endpoint Split by Merge
/v1/images/generations, /v1/video/generations n (each part gets its own seed offset) images/videos in order
/v1/embeddings the input list vectors re-indexed, usage summed
/v1/rerank the documents list scores re-indexed, sorted, top_n applied
/v1/audio/speech sentences of input spoken parts joined with ffmpeg into the requested format
/v1/audio/transcriptions time windows cut at silences (chunk_seconds) text joined, segment/word timestamps shifted (json, verbose_json, text; srt/vtt and diarization stay on one engine)
/v1/ocr/batch files results in order

Options: which engines/nodes (nodes, blank = every one that can serve the capability, the one already holding the model first), max_parts, min_items (below it nothing is split), chunk_seconds. A part that fails is retried once on another engine before the request fails. Files a node generates are addressed through the head (/v1/files/… rewritten and fetched from the node). Streaming requests are never split. One generation is never split either — that is what the next section and the RPC / vLLM sections are for.

5. Parts of a video pipeline on other machines

A video model is several models run one after the other — text encoder (UMT5-XXL, 11 GB), one or two denoising experts (Wan 2.2: high-noise then low-noise, ~10 GB each at 4-bit), VAE — and what passes between them is small: prompt embeddings, a latent tensor of a few MB. So each part can live on a different machine and a generation becomes a relay over the LAN, which lets a 14B dual-expert model run where no single card holds it all.

Per model, Pipeline parts on other machines ("components"):

{"components": {"text_encoder": "box2", "low_noise": "box3", "vae": "box3"}}
  • text_encoder — the prompt is encoded there (the head never runs its own text encoder for this model);
  • low_noise — the head runs the high-noise steps, hands the latents over at the expert boundary, the node runs the rest;
  • vae — the latents are decoded there and the node returns the finished video (the head's post-processing — upscale, audio, subtitles — is then skipped, since the node's answer is the answer).

Each named node needs the same model configured. The node runs the ordinary /v1/video/generations with a _handoff block (encode, denoise from a step, decode), so loading, acceleration presets and LoRAs happen exactly as for a normal generation there. Wan T2V and I2V pipelines; other pipelines run on one machine with a note in the log.

6. LoRA / QLoRA training over several machines

Every LoRA trained on a model can run on this machine and on cluster nodes at once: each machine trains on its own card with its share of the samples, and after every step the gradients of the adapter — a few MB — are averaged across all machines (synchronous data parallel over torch.distributed, gloo by default, nccl on request), so every machine holds the same adapter throughout and the result on this machine is the whole cluster's work. Steps are the same; each step now sees N samples instead of one.

  • Per model: LoRA training nodes (lora_train_nodes) on the base model's page; or per job: "nodes": ["box2", "box3"] on POST /v1/loras/train (also {url, api_key} blocks for machines outside cluster.nodes), dist_backend.
  • The machine that takes the job is rank 0: it resolves the images (a character/environment profile exists only there), sends the same job with the images inlined to each node, then trains. Peers save nothing that outlives the job and are cancelled when rank 0 aborts.
  • The nodes need the same base model. Rank 0's advertise host and a free port are the rendezvous (cluster.advertise_host when the default route's address is not the one the nodes can reach).
  • Covers the in-process trainers: SD1.5, SDXL, Z-Image, the flow-DiT family (LTX-2 …) and Wan, 4-bit QLoRA included.

7. Prefix-cache-aware routing

A model resident in more than one place — two local engines, a local engine and a node, replicas after a fan-out — is served by the least busy holder, local before remote on a tie. For a conversation, though, the better place is where its earlier turns went: that engine's KV / prefix cache (the multi-slot cache of the GGUF and HF backends, vLLM's prefix caching) already holds the shared opening, and the prompt is processed from the first new token rather than from scratch.

So the front keeps a small map conversation → engine and prefers that engine while it is alive, capable and still holds the model. The key is the one the engines use for their own slot affinity: an X-Session-Id header or the OpenAI user field when the client sends one, else a hash of the conversation's stable opening (system prompt + first user turn, or the prompt head). A pin still wins; a preference never routes to an engine that cannot serve. Settings → Cluster → Prefix-cache routing (on by default, TTL 30 min): cluster.prefix_affinity, cluster.prefix_affinity_ttl_s.

8. Pools — several machines for one capability or one host model

  • Remote capabilities (Settings → Remote capabilities): a capability's URL field takes several URLs, comma-separated. They form a pool: the first healthy, least-busy one takes each request; a URL that fails is marked down for 30 s and the request moves to the next one.
  • backend: host models: More machines, one per line url | token | start command | stop command (blank fields inherit the ones above). Same pool: healthy + fewest in flight wins, a dead one is skipped, an on-demand one is started only when none is up, and one that dies mid-request hands the request to another.

9. Operating it — metrics, usage, logs, recovery

  • GET /metrics — Prometheus exposition, no client library. Per completed request: coderai_requests_total and coderai_request_seconds_{sum,count} labelled by API key name (the name on the Tokens page; session for a browser; anonymous; cluster for a head), model, kind, status and the engine or node that served it. Gauges: coderai_engine_up / _inflight / _vram_free_bytes / _vram_total_bytes / _loaded_models {engine,backend,remote}, coderai_node_up{node}, coderai_node_recoveries_total, coderai_rpc_server_up{endpoint}, coderai_discovery_peers / _members, coderai_runpod_spend_usd{model,period}, coderai_runpod_pods, coderai_requests_active, coderai_info{version}. Guarded like the other telemetry (a session, an API token, or the cluster token) unless Settings → Server → Serve /metrics without a credential (server.metrics_public). A scrape config:

    scrape_configs:
      - job_name: coderai
        static_configs: [{targets: ["box1:8776", "box2:8776"]}]
        authorization: {credentials: "<an API token or the cluster token>"}
  • Usage by key / by model — the same counters as tables on the Cluster page (GET /admin/api/usage): requests, errors, wall time, kinds, where they ran. Since the front started; Prometheus is the durable store.

  • Node logs — the Cluster page shows each node's recent output (its front's lines and every engine line it re-emits) with the log button on the node row, and this front's own with the button on the engines card: GET /cluster/logs?lines=N on the node (token-guarded), GET /admin/api/cluster/logs?node=NAME on the head.

  • Recovery — a node that stops answering keeps its assignment; when it answers again (it very likely restarted, and the entries the head pushed lived only in its memory) the head pushes its assigned models again and they load there without anyone clicking (coderai_node_recoveries_total). Local engines have had this since the supervisor: a crashed engine is respawned, a crash-looping one quarantined, and its models come back with the next assignment push. A discovered node that vanishes from mDNS is dropped and re-added when it reappears.

10. Apple Silicon, MLX and other accelerators — what fits

CoderAI's images are x86-64, CUDA + Vulkan; nothing of CoderAI runs natively on macOS or on Ascend/Hygon/MThreads cards, and a container on those has no GPU. What already works, and what would be needed:

  • A Mac as a text node today: mlx_lm.server (or LM Studio, Ollama) is OpenAI-compatible, so it is a remote capability for text (Settings → Remote capabilities, its URL) or a backend: host model with no start command. It is then one more place a text request can go — with pools, failover and metering, but without CoderAI's model page, VRAM accounting, or fan-out awareness. This is the cheapest useful step and needs no code.
  • A real MLX node would be a small Python package (coderai-node-mlx: the front's /cluster/state + /v1/chat/completions + /v1/embeddings over mlx-lm, later mlx-audio / diffusers-mlx) installed with pip on the Mac and joining by token and mDNS like any node. Two to three days; the node protocol is already the only contract a node must honour. That is the point where "a stack of Macs" stops being exo's story alone.
  • Tensor parallel over Thunderbolt/RDMA (exo's headline) is MLX-only and macOS-only; nothing to port, and the same reasoning that keeps our cross-machine split at layer/pipeline granularity applies on Ethernet.
  • Ascend / Hygon / MThreads / Iluvatar / MetaX (GPUStack's list) each come with a vLLM fork (vllm-ascend, …) or a llama.cpp backend (CANN, MUSA). The cheapest route is the one we already have for anything that speaks OpenAI: run the vendor's vLLM as a service_url / host model. Native support means building their llama.cpp backends into the wheel and a per-vendor image; do it when there is a card to test on, not before.

Reference

  • Config: cluster section in config.json — enabled, nodes, node_name, serve, poll_timeout_s, rpc_servers, rpc_bin, advertise_host, token, discovery, auto_join, prefix_affinity, prefix_affinity_ttl_s; server.metrics_public; vllm.pipeline_parallel_size / distributed_executor_backend / ray_address / ray_port / nodes / nodes_ready_timeout_s; ktransformers.nnodes / dist_init_addr / tp_size / nodes.
  • Per model (models.json): engine (a node name), node_path / node_paths, rpc_servers, split_mode, vllm block, kt block, host.hosts, distribute block, components block, lora_train_nodes.
  • Node API (a token of the node, or the cluster token): GET /cluster/state, GET /cluster/rpc-servers, GET /cluster/logs, POST /cluster/reload-config, POST /cluster/model-load, POST /cluster/model-unload.
  • Head: GET /admin/api/cluster (what the Cluster page shows), GET /admin/api/cluster/logs?node=, GET /admin/api/usage; GET /metrics.
  • Code: codai/cluster/ (nodes, discovery, rpc, multinode, compat, overview, fanout, components, ddp), codai/frontproxy/affinity.py, codai/frontproxy/metrics.py, codai/backends/ggml_rpc.py, codai/backends/overrides.py, codai/frontproxy/engine_supervisor.py (remote engines), codai/api/remote_gateway.EndpointPool, codai/api/host_worker.HostPool.