# What Protocol Is Used for Distributed Coordinator-Worker Communication in DS4?

> Discover the custom lightweight binary protocol DS4 uses for distributed coordinator-worker communication over TCP sockets. Learn more about this efficient protocol.

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

---

**DS4 implements a custom lightweight binary protocol that runs over plain TCP sockets, identified by the magic constant "DS4D" and defined in [`ds4_distributed.c`](https://github.com/antirez/ds4/blob/main/ds4_distributed.c) and [`ds4_distributed.h`](https://github.com/antirez/ds4/blob/main/ds4_distributed.h).**

The DS4 inference engine requires efficient distributed coordination between a central coordinator and multiple worker nodes. Rather than adopting an external messaging framework, the project defines its own **DS4 binary protocol**—a purpose-built, low-latency transport layer optimized for neural network inference workloads. This protocol handles everything from worker registration to activation tensor exchange across the network.

## How the DS4 Protocol Works Over TCP

The distributed coordinator-worker communication in DS4 relies on standard TCP sockets enhanced for performance. According to the antirez/ds4 source code, the protocol layer configures several low-latency socket options before any data flows.

### TCP Transport Configuration

In [`ds4_distributed.c`](https://github.com/antirez/ds4/blob/main/ds4_distributed.c), the function `dist_set_socket_low_latency` (lines 27–45) prepares each connection:

```c
/* From ds4_distributed.c:27-45 */
static void dist_set_socket_low_latency(int fd) {
    int yes = 1;
    /* Disable Nagle's algorithm for minimal latency */
    setsockopt(fd, IPPROTO_TCP, TCP_NODELAY, &yes, sizeof(yes));
    /* Enable TCP keep-alive to detect dead peers */
    setsockopt(fd, SOL_SOCKET, SO_KEEPALIVE, &yes, sizeof(yes));
    /* Additional buffer size tuning available via config */
}

```

This ensures that **coordinator-worker communication** achieves predictable, low-latency behavior even under high-frequency message exchange typical of distributed inference.

## Message Framing: The DS4D Header Format

Every message in the DS4 protocol begins with a fixed 12-byte header structure defined in [`ds4_distributed.c`](https://github.com/antirez/ds4/blob/main/ds4_distributed.c) (lines 72–76):

```c
typedef struct {
    uint32_t magic;      /* 0x44533444 "DS4D" in ASCII */
    uint32_t type;       /* Message type identifier */
    uint32_t size;       /* Payload size in bytes */
} ds4_dist_frame_header;

```

The magic number `0x44533444u` spells "DS4D" when interpreted as ASCII, providing a quick validity check for protocol compliance. The `dist_write_frame_header` function (lines 130–136) serializes this header to the wire:

```c
int dist_write_frame_header(int fd, uint32_t type, uint32_t size) {
    ds4_dist_frame_header hdr = {
        .magic = htonl(0x44533444u),
        .type  = htonl(type),
        .size  = htonl(size)
    };
    return dist_write_full(fd, &hdr, sizeof(hdr));
}

```

## Core Message Types for Coordinator-Worker Communication

The DS4 protocol defines a compact set of message types for **distributed coordinator-worker communication**, enumerated at the top of [`ds4_distributed.c`](https://github.com/antirez/ds4/blob/main/ds4_distributed.c) (lines 44–53):

| Constant | Value | Purpose |
|----------|-------|---------|
| `DS4_DIST_MSG_HELLO` | 1 | Worker initialization and capability announcement |
| `DS4_DIST_MSG_ERROR` | 2 | Error reporting from either side |
| `DS4_DIST_MSG_WORK` | 3 | Work request dispatched from coordinator to worker |
| `DS4_DIST_MSG_RESULT` | 4 | Inference result returned from worker to coordinator |
| `DS4_DIST_MSG_SNAPSHOT_*` | 5–9 | KV-cache snapshot exchange for state persistence |

### HELLO Message: Worker Registration

When a worker connects, it sends a `HELLO` message containing its model assignment and layer range. The fixed payload structure `ds4_dist_hello_fixed` includes network-converted fields:

```c
ds4_dist_hello_fixed hello = {
    .model_id    = htonl(model_id),
    .quant_bits  = htonl(quant_bits),
    .layer_start = htonl(layer_start),
    .layer_end   = htonl(layer_end),
    /* ... additional capability fields ... */
    .listen_port = htonl(listen_port),
    .model_name_len = htonl(strlen(model_name))
};

dist_write_frame_header(fd, DS4_DIST_MSG_HELLO, sizeof(hello));
dist_write_full(fd, &hello, sizeof(hello));
dist_write_full(fd, model_name, strlen(model_name));

```

### WORK Message: Dispatching Inference Tasks

The coordinator sends `WORK` messages to assign token processing across specific layer ranges. The `ds4_dist_work_fixed` struct (converted via `dist_work_to_wire`) contains routing hashes, token positions, and compression settings:

```c
ds4_dist_work_fixed work = {
    .model_id        = htonl(state->model_id),
    .session_hi      = htonl(session_hi),
    .session_lo      = htonl(session_lo),
    /* Hash identifiers for request routing */
    .prefix_hash_hi  = htonl(prefix_hash_hi),
    .prefix_hash_lo  = htonl(prefix_hash_lo),
    /* Layer assignment and token range */
    .layer_start     = htonl(layer_start),
    .layer_end       = htonl(layer_end),
    .n_tokens        = htonl(n_tokens),
    /* Activation compression parameters */
    .input_hc_bits   = htonl(input_hc_bits)
};

dist_send_work_frame(fd, &work, tokens, input_hc, route_blob);

```

### Reading and Dispatching Messages

Both coordinator and worker use the same framing logic. The `dist_read_frame_header` function parses incoming headers and drives message dispatch:

```c
uint32_t type, bytes;
char err[256];

int rc = dist_read_frame_header(fd, &type, &bytes, err, sizeof(err));
if (rc <= 0) {
    fprintf(stderr, "Failed to read header: %s\n", err);
    return -1;
}

switch (type) {
    case DS4_DIST_MSG_HELLO:   /* handle registration */ break;
    case DS4_DIST_MSG_WORK:    /* process inference task */ break;
    case DS4_DIST_MSG_RESULT:  /* collect output tensors */ break;
    case DS4_DIST_MSG_ERROR:   /* handle failure */ break;
    /* snapshot cases for state management */
}

```

## Activation Compression in the Protocol

A distinctive feature of DS4's **distributed communication protocol** is built-in support for quantized activation transfer. The `bits` field in work messages selects between:

- **32-bit** IEEE-754 floats (full precision)
- **16-bit** half-precision (FP16/BF16)
- **8-bit** E4M3 format (reduced bandwidth)

The conversion utilities `dist_write_activation_payload` and `dist_decode_activation_payload` handle the packing and unpacking transparently, allowing workers to trade numerical precision for network throughput without protocol changes.

## Connection Lifecycle and Socket Management

The DS4 protocol implementation provides symmetrical connection primitives:

- **Coordinator side**: `dist_open_listener` creates the accepting socket; [`ds4_server.c`](https://github.com/antirez/ds4/blob/main/ds4_server.c) manages the accept loop and worker pool.
- **Worker side**: `dist_connect_endpoint` establishes outbound connections to the coordinator.

Both paths converge on the same framing and message handling code, ensuring consistent **coordinator-worker protocol** behavior regardless of direction.

## Summary

- **DS4 uses a custom binary protocol** named "DS4D" (magic `0x44533444`) for all distributed communication.
- **Transport layer** runs on standard TCP with `TCP_NODELAY`, keep-alive, and configurable buffers—implemented in [`ds4_distributed.c`](https://github.com/antirez/ds4/blob/main/ds4_distributed.c) lines 27–45.
- **12-byte fixed header** carries magic, message type, and payload size for unambiguous framing.
- **Five core message types** cover the full lifecycle: `HELLO`, `ERROR`, `WORK`, `RESULT`, and snapshot variants.
- **Network byte order conversion** (`htonl`/`ntohl`) ensures cross-platform compatibility.
- **Built-in quantization** supports 32/16/8-bit activation tensors to optimize bandwidth.

## Frequently Asked Questions

### Is the DS4 protocol based on HTTP, gRPC, or another standard?

No. DS4 implements its own lightweight binary protocol from scratch. The design prioritizes minimal overhead and predictable latency over compatibility with existing frameworks. The protocol is defined entirely within [`ds4_distributed.c`](https://github.com/antirez/ds4/blob/main/ds4_distributed.c) and [`ds4_distributed.h`](https://github.com/antirez/ds4/blob/main/ds4_distributed.h) without external dependencies.

### How does DS4 ensure message boundaries are preserved over TCP?

Through fixed-size header framing. Every message starts with a 12-byte `ds4_dist_frame_header` containing the payload length. The receiver first reads exactly 12 bytes, extracts the size field, then reads the remaining payload—eliminating any ambiguity about where one message ends and the next begins.

### Can the DS4 protocol run over TLS or other secure transports?

The current implementation in [`ds4_distributed.c`](https://github.com/antirez/ds4/blob/main/ds4_distributed.c) uses plain TCP sockets. There is no built-in TLS support; security would need to be provided at the network layer (VPN, WireGuard, etc.) or by modifying the socket creation code to wrap connections with OpenSSL or similar libraries.

### What happens if a worker sends an unrecognized message type?

The coordinator's dispatch switch statement handles known types explicitly. Unrecognized types fall through to error handling paths that typically log the anomaly and close the connection. The protocol includes a dedicated `DS4_DIST_MSG_ERROR` type for propagating failure information bidirectionally.