# Understanding ds4 Distributed Pipeline Parallelism: Architecture and Implementation

> Explore ds4's distributed pipeline parallelism architecture. Discover how it splits LLM inference across layers using a coordinator-worker model with efficient communication on Apple Silicon.

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

---

**ds4 implements distributed pipeline parallelism through a coordinator-worker model that splits LLM inference across two contiguous layer slices, communicating via a lightweight TCP-based protocol with optional RDMA support on Apple Silicon.**

The antirez/ds4 project enables large language model inference across multiple machines using a novel distributed architecture. This system implements **ds4 distributed pipeline parallelism** by dividing the computational graph into contiguous segments that execute on separate identical hardware, coordinated through a purpose-built binary protocol. The design minimizes communication overhead while maintaining a unified session API.

## Core Architecture Components

The architecture consists of three tightly-coupled layers defined primarily in [`ds4_distributed.c`](https://github.com/antirez/ds4/blob/main/ds4_distributed.c).

### The Coordinator Process

The coordinator owns the public `ds4_session` API that users interact with. It manages the engine instance, maintains a registry of connected workers, and handles per-session KV store coordination. The `ds4_dist_coordinator_state` struct encapsulates this state:

```c
ds4_dist_options dist_opt = {0};
dist_opt.role = DS4_DISTRIBUTED_COORDINATOR;
dist_opt.listen_host = "0.0.0.0";
dist_opt.listen_port = 12345;
char err[256];
if (!dist_validate_options(&dist_opt, err, sizeof(err))) {
    fprintf(stderr, "invalid options: %s\n", err);
    exit(1);
}
ds4_dist_coordinator_state *coord = calloc(1, sizeof(*coord));
coord->engine = engine;
coord->listen_fd = dist_open_listener(dist_opt.listen_host,
                                      dist_opt.listen_port,
                                      err, sizeof(err));

```

The coordinator creates a listening socket via `dist_open_listener` and accepts worker connections, then forwards inference requests and aggregates results.

### The Worker Process

Workers execute local slices of the model—specifically, a contiguous range of layers defined by `layer_start` and `layer_end`. The `ds4_dist_worker_state` and `ds4_dist_worker_session` structs represent runtime state and per-session data respectively. The worker's main event loop is `dist_worker_handle_work`, which processes incoming work frames and returns hidden states or logits.

### Transport Layer Protocol

A lightweight TCP-based framing protocol moves activation data, routing metadata, and KV snapshots between nodes. Framing helpers like `dist_write_frame_header` and `dist_read_frame_header` in [`ds4_distributed.c`](https://github.com/antirez/ds4/blob/main/ds4_distributed.c) handle message boundaries. Message types include `DS4_DIST_MSG_WORK` for task distribution and `DS4_DIST_MSG_RESULT` for returning computations.

## Model Slicing and Routing Strategy

When ds4 runs in distributed mode, the model divides into two contiguous slices. The coordinator executes the first half while the worker handles the second. During the initial handshake, both sides exchange a `ds4_dist_hello_fixed` frame (magic = `DS4_DIST_MAGIC`) containing:

- Model ID, quantization bits, and layer ranges
- Context size and printable model name
- A `ds4_dist_route_plan` specifying which layers belong to which worker

This routing plan ensures activations flow sequentially from coordinator to worker during inference.

## Runtime Execution Flow

### Session Management and KV Store Sharding

Each inference request receives a unique session ID and request ID. The coordinator creates a `ds4_dist_worker_session` and transmits a `DS4_DIST_MSG_WORK` frame. Workers store per-session KV tensors in local shards via `ds4_dist_kv_shard_file`, returning results via `DS4_DIST_MSG_RESULT` frames containing:

- `DS4_DIST_RESULT_HIDDEN_STATE` for intermediate activations
- `DS4_DIST_RESULT_LOGITS` for final outputs
- `DS4_DIST_RESULT_ACK` for acknowledgments

The coordinator merges these remote outputs back into the local session, presenting a seamless inference pipeline to the rest of the engine.

### Prefill and Activation Exchange

During prompt prefill, the coordinator sends input hidden-state activations (`input_hc`) to the worker. The `dist_write_activation_payload` function supports multiple precision formats: 32-bit floats, 16-bit half-precision, or 8-bit quantized. Workers decode these via `dist_decode_activation_payload` before executing their layer slice, maintaining API consistency regardless of distribution mode.

```c
static int dist_worker_handle_work(ds4_dist_worker_state *state,
                                   ds4_dist_worker_upstream *upstream,
                                   uint32_t bytes) {
    ds4_dist_work_fixed work;
    /* read work header */
    if (dist_read_full(upstream->fd, &work, sizeof(work)) != 1) return -1;
    dist_work_from_wire(&work);
    /* decode activation payload */
    float *input_hc;
    uint32_t hc_bytes;
    if (dist_decode_activation_payload(payload, work.input_hc_bits,
                                       work.input_hc_bytes, &input_hc,
                                       &hc_bytes, NULL, err, sizeof(err)))
        return -1;
    /* run the slice */
    ds4_session *sess = ds4_session_new(state->engine, work.layer_start,
                                        work.layer_end);
    ds4_session_set_hidden(sess, input_hc, hc_bytes);
    ds4_session_run(sess);
    /* send result back */
    ds4_dist_result_fixed res = {...};
    dist_result_to_wire(&res);
    dist_send_frame(upstream->fd, DS4_DIST_MSG_RESULT,
                   &res, sizeof(res));
}

```

## Checkpointing and Snapshot Protocol

Distributed checkpointing uses a snapshot protocol initiated by `DS4_DIST_MSG_SNAPSHOT_SAVE_REQ`. The coordinator streams KV shards in `DS4_DIST_MSG_SNAPSHOT_CHUNK` frames, while loading reverses this flow via `DS4_DIST_MSG_SNAPSHOT_LOAD_BEGIN`. These operations reuse `dist_payload_write_bytes` and `dist_payload_read_bytes` utilities, ensuring consistency with local session handling.

## Tensor-Parallel Transport with RDMA

For Apple Silicon deployments, ds4 offers an optional tensor-parallel transport layer implemented in [`ds4_tp.c`](https://github.com/antirez/ds4/blob/main/ds4_tp.c). This replaces the generic TCP control channel with low-latency RDMA gate exchange:

- `tp_hello_exchange` verifies model metadata and negotiates RDMA use
- `tp_rdma_gate_exchange` transmits decode gates via SEND/RECV operations
- `tp_rdma_big_gate_exchange` handles large prefill rows

If RDMA is unavailable, the system gracefully falls back to standard TCP transport.

```c
ds4_tp_options tp_opt = {0};
tp_opt.requested = true;
tp_opt.role = DS4_TP_LEADER;
tp_opt.transport = DS4_TP_TRANSPORT_AUTO;
tp_opt.rdma_device = "rdma_en0";
ds4_tp *tp;
if (!ds4_tp_create(&tp, &tp_opt, &identity, err, sizeof(err))) {
    fprintf(stderr, "TP init failed: %s\n", err);
    exit(1);
}

```

## Performance Monitoring and Telemetry

Both coordinator and worker embed telemetry payloads (`ds4_dist_telemetry_fixed`) in every result frame. These records track `eval_usec`, `downstream_wait_usec`, and `forward_send_usec`, enabling the coordinator to aggregate profiling data across the distributed system.

## Summary

- **ds4 distributed pipeline parallelism** uses a coordinator-worker architecture to split LLM inference across two machines executing contiguous layer slices
- The coordinator manages the public session API in [`ds4_distributed.c`](https://github.com/antirez/ds4/blob/main/ds4_distributed.c) while workers run local graph segments via `dist_worker_handle_work`
- Communication occurs via a compact binary protocol over TCP, with optional RDMA support for Apple Silicon via [`ds4_tp.c`](https://github.com/antirez/ds4/blob/main/ds4_tp.c)
- Per-session KV stores shard across nodes using `ds4_dist_kv_shard_file` and merge transparently after computation
- Snapshot protocols enable distributed checkpointing using `DS4_DIST_MSG_SNAPSHOT_CHUNK` frames and the same byte utilities as local sessions

## Frequently Asked Questions

### How does ds4 split the model across distributed nodes?

ds4 divides the LLM into two contiguous layer ranges specified in the `ds4_dist_hello_fixed` handshake. The coordinator executes layers from `layer_start` to a midpoint, while the worker processes from that midpoint to `layer_end`. This pipeline approach minimizes communication overhead by transmitting only hidden states between slices rather than full parameter sets.

### What transport protocols does ds4 support for distributed inference?

By default, ds4 uses a lightweight TCP-based framing protocol defined in [`ds4_distributed.c`](https://github.com/antirez/ds4/blob/main/ds4_distributed.c) with functions like `dist_write_frame_header` and `dist_read_frame_header`. On Apple Silicon hardware, it can leverage RDMA through the tensor-parallel transport layer ([`ds4_tp.c`](https://github.com/antirez/ds4/blob/main/ds4_tp.c)) using `tp_rdma_gate_exchange` for lower latency, falling back to TCP if RDMA is unavailable.

### How does ds4 handle KV cache management in distributed mode?

Each worker maintains local KV cache shards in `ds4_dist_kv_shard_file` indexed by session ID. During inference, the coordinator distributes work via `DS4_DIST_MSG_WORK` frames and merges results returned in `DS4_DIST_MSG_RESULT` frames. For checkpointing, the snapshot protocol streams these shards using `DS4_DIST_MSG_SNAPSHOT_CHUNK` frames between coordinator and workers.

### Can ds4 distribute inference across more than two machines?

The current architecture in [`ds4_distributed.c`](https://github.com/antirez/ds4/blob/main/ds4_distributed.c) implements a two-node coordinator-worker model. The `ds4_dist_coordinator_state` manages a single worker connection, and the protocol design assumes two contiguous slices. While the transport abstractions could theoretically extend to more nodes, the current implementation limits distribution to one coordinator and one worker.