# Setting up Pipeline Parallelism for Multi-Machine Distributed Inference in ds4

> Learn how to set up pipeline parallelism for multi-machine distributed inference in ds4. ds4 splits LLMs across GPUs, keeping them busy during data flow for efficient inference.

- Repository: [Salvatore Sanfilippo/ds4](https://github.com/antirez/ds4)
- Tags: how-to-guide
- Published: 2026-08-08

---

**The ds4 repository implements a distributed inference runtime that splits large language models across multiple machines and executes a prefill pipeline to keep every GPU busy while data flows from one slice to the next.**

This guide walks through configuring **pipeline parallelism for multi-machine distributed inference in ds4**, the lightweight inference engine created by Salvatore Sanfilippo. By distributing model layers across separate workers and coordinating them through a central scheduler, ds4 maximizes hardware utilization for long-context generation.

## Architecture of the ds4 Distributed Inference Runtime

The distributed machinery lives primarily in **[`ds4_distributed.c`](https://github.com/antirez/ds4/blob/main/ds4_distributed.c)**, which implements a coordinator-worker pattern with low-latency TCP transport and quantized activation packing.

### The Coordinator Role

The **coordinator** accepts client sessions, breaks prompts into chunks, and streams work to workers in a pipelined fashion. Key functions include:

- **`dist_coordinator_can_pipeline_prefill`** (lines 3421-3428): Validates whether a prompt can be split into chunks that fit the workers' prefetch windows.
- **`dist_coordinator_prefill_prompt`** (lines 3795-3804): Creates the pre-fill work packets and manages the dispatch queue.

The coordinator maintains a `dist_coordinator.state.prefill_chunk` value (set via `--dist-prefill-chunk`) that determines the maximum tokens per chunk, typically defaulting to 4096.

### The Worker Role

Each **worker** holds a contiguous **model slice** (defined by `layer_start` and `layer_end`) and executes that slice for incoming chunks. The entry point is **`dist_worker_handle_work`** (lines 4967-4972), which processes `WORK` frames from the coordinator and returns `RESULT` frames containing hidden states or logits.

Workers register themselves by sending a `HELLO` frame (`DS4_DIST_MSG_HELLO`) to the coordinator upon startup, using the wire format defined in `dist_hello_to_wire` and `dist_hello_from_wire`.

### Transport Layer and Activation Packing

The transport layer uses raw TCP sockets with Nagle's algorithm disabled and keep-alive enabled for low-latency communication. The function **`dist_set_socket_low_latency`** (lines 1028-1043) configures these socket options, while **`dist_connect_endpoint`** (lines 1444-1465) handles connection establishment.

To reduce bandwidth, ds4 supports **activation quantization** via **`dist_write_activation_payload`** (lines 925-964) and **`dist_decode_activation_payload`** (lines 1005-1025). Activations can be transmitted as 32-bit floats, 16-bit half-precision (`f16`), or 8-bit (`f8_e4m3`) values depending on the `bits` parameter.

## Deploying a Multi-Machine Pipeline Cluster

Follow these steps to launch a distributed inference cluster across multiple physical machines.

1. **Build ds4** on every machine with GPU support enabled:

   ```bash
   make clean && make -j
   ```

2. **Launch the Coordinator** on the machine hosting the first model slice:

   ```bash
   ./ds4_server --dist-role coordinator \
                --dist-listen 0.0.0.0:8000 \
                --model /path/to/model.gguf \
                --dist-route worker1:8001,worker2:8002 \
                --dist-prefill-chunk 2048
   ```

   The `--dist-route` flag constructs a `ds4_dist_route_plan` that maps each layer slice to a specific worker host and port.

3. **Launch Worker processes** on remaining machines, specifying their layer ranges:

   ```bash
   # Worker 1 (layers 0-12)

   ./ds4_server --dist-role worker \
                --dist-listen 0.0.0.0:8001 \
                --model /path/to/model.gguf \
                --dist-layers 0-12
   
   # Worker 2 (layers 13-24)

   ./ds4_server --dist-role worker \
                --dist-listen 0.0.0.0:8002 \
                --model /path/to/model.gguf \
                --dist-layers 13-24
   ```

4. **Connect a client** to the coordinator:

   ```bash
   ./ds4_cli --host coordinator-host:8000 \
             --prompt "Explain the architecture" \
             --max-tokens 512
   ```

### Tuning Environment Variables

Fine-tune the pipeline behavior by setting these variables before launching:

- **`DS4_DIST_PREFILL_SEND_DEPTH`**: Controls the maximum number of outstanding chunks in flight (default: `2`). Increasing this value improves pipeline utilization but consumes more memory.
- **`DS4_DIST_SOCKET_BUFFER_MB`**: Sets the TCP socket buffer size in megabytes (default: `128`).
- **`DS4_DIST_DECODE_PROFILE`**: Enables trace output for decode-time profiling when set to `1`.

These variables are read by helper functions in [`ds4_distributed.c`](https://github.com/antirez/ds4/blob/main/ds4_distributed.c) during startup.

## How the Prefill Pipeline Works Internally

Understanding the data flow helps optimize chunk sizes and cluster topology.

### Chunking and Hash Validation

When a prompt arrives, the coordinator splits it into chunks of up to `prefill_chunk` tokens. Each chunk carries a **prefix hash** computed by `dist_token_hash_prefix`, allowing workers to validate KV cache consistency and reject chunks that would produce incorrect states.

### Pipelined Dispatch Mechanics

While the first worker processes chunk 0, the coordinator immediately dispatches chunk 1 to the next available worker. This **staggered execution** continues until the pipeline depth reaches `dist_prefill_send_depth`. The coordinator aggregates partial results from `RESULT` frames and forwards the final output to the client.

## Code Reference Examples

### Checking Pipeline Eligibility (Coordinator)

The coordinator uses this logic to determine if pipelining is feasible:

```c
/* ds4_distributed.c – lines 3421-3428 */
static bool dist_coordinator_can_pipeline_prefill(
        ds4_dist_coordinator_state *state,
        const ds4_dist_route_plan *plan,
        ds4_session *session,
        uint32_t prompt_len,
        uint32_t chunk_cap)
{
    /* Pipeline only if the prompt can be split into chunks that fit the
       worker's prefetch window (prefill_window) and the coordinator's
       maximum chunk size. */
    return (prompt_len > chunk_cap && state->prefill_window > 0);
}

```

### Activation Packing (Worker)

Workers quantize activations before transmission to minimize network overhead:

```c
/* ds4_distributed.c – lines 925-964 */
int dist_write_activation_payload(int fd,
                                  const float *src,
                                  uint64_t values,
                                  uint32_t bits)
{
    bits = dist_activation_bits_or_default(bits);
    if (bits == 32u) return dist_write_full(fd, src, values * sizeof(float));
    /* pack to f16 or f8_e4m3 based on bits parameter */
    /* ... */
}

```

## Summary

- **ds4** implements pipeline parallelism through a coordinator-worker architecture defined in [`ds4_distributed.c`](https://github.com/antirez/ds4/blob/main/ds4_distributed.c).
- The **coordinator** splits prompts into chunks and pipelines them across workers using the `dist_coordinator_prefill_prompt` function.
- **Workers** execute specific layer ranges (`--dist-layers`) and communicate via quantized activations (32/16/8-bit) to minimize bandwidth.
- Deployment requires starting a coordinator with `--dist-route` and workers with matching `--dist-layers` ranges.
- Tune performance using `DS4_DIST_PREFILL_SEND_DEPTH` and `--dist-prefill-chunk` based on your network latency and GPU memory.

## Frequently Asked Questions

### What hardware requirements are needed for ds4 distributed inference?

Each machine requires a CUDA-capable or Metal-capable GPU with sufficient VRAM to hold its assigned model slice. The coordinator can run on CPU-only hardware, but workers need GPU acceleration to execute their layer ranges efficiently. All machines must have low-latency network connectivity (10 Gbps or higher recommended) to prevent pipeline stalls.

### How does ds4 handle network latency between machines?

The transport layer disables Nagle's algorithm and enables TCP keep-alive through `dist_set_socket_low_latency` to minimize latency. Additionally, the `DS4_DIST_PREFILL_SEND_DEPTH` environment variable allows you to increase the number of chunks in flight, effectively hiding network latency by ensuring workers always have data ready to process while waiting for previous results.

### Can I mix different GPU types in a ds4 cluster?

Yes, provided all GPUs support the required precision formats (float32, float16, or the specific quantization scheme used). When mixing hardware, assign `--dist-layers` ranges proportional to each GPU's memory capacity and compute speed. Slower workers become bottlenecks, so balance layer counts accordingly.

### What is the optimal chunk size for --dist-prefill-chunk?

The optimal size depends on your model's hidden dimension and network bandwidth. Smaller chunks (512-1024 tokens) reduce latency for short prompts but increase overhead. Larger chunks (2048-4096 tokens) improve throughput for long sequences but require more VRAM per worker. Start with 2048 and adjust based on the `DS4_DIST_DECODE_PROFILE` traces to minimize idle time between chunks.