Understanding ds4 Distributed Pipeline Parallelism: Architecture and Implementation

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.

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:

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 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.

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. 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.

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 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
  • 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 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) 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 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.

Have a question about this repo?

These articles cover the highlights, but your codebase questions are specific. Give your agent direct access to the source. Share this with your agent to get started:

Share the following with your agent to get started:
curl -s "https://instagit.com/install.md"

Works with
Claude Codex Cursor VS Code OpenClaw Any MCP Client

Maintain an open-source project? Get it listed too →