# Coordinator-to-Worker Protocol in DS4 Distributed Inference: Custom Binary Framing Over TCP/RDMA

> Discover the custom binary protocol DS4 uses for coordinator-to-worker communication. Learn about its 12-byte header, TCP/RDMA support, and efficient distributed inference.

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

---

**The DS4 distributed inference system uses a custom binary protocol with a 12-byte framed header (magic number, message type, payload length) running over TCP sockets with optional RDMA fallback.**

This lightweight protocol enables efficient communication between a central coordinator (`ds4_server`) and multiple worker agents (`ds4_agent`) during distributed neural network inference. The implementation lives 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), with transport abstraction in [`ds4_tp.c`](https://github.com/antirez/ds4/blob/main/ds4_tp.c).

## Protocol Overview: The DS4 Distributed Frame Format

Every message exchanged between coordinator and worker follows a strict binary framing scheme. The protocol uses **network byte order (big-endian)** for all multi-byte fields.

### Header Structure

The frame header is defined by `ds4_dist_frame_header`:

```c
typedef struct {
    uint32_t magic;      /* 0x44533444 ("DS4D") */
    uint32_t type;       /* Message type identifier */
    uint32_t bytes;      /* Payload length (excludes header) */
} ds4_dist_frame_header;

```

This 12-byte header precedes every payload. The **magic number `0x44533444`** serves dual purposes: validating the stream hasn't desynchronized and detecting host byte-order mismatches between coordinator and worker.

### Core Message Types

The protocol defines enumerated message kinds in [`ds4_distributed.h`](https://github.com/antirez/ds4/blob/main/ds4_distributed.h):

| Constant | Value | Purpose |
|----------|-------|---------|
| `DS4_DIST_MSG_HELLO` | 1 | Initial handshake with worker capabilities |
| `DS4_DIST_MSG_ERROR` | 2 | Error propagation from worker to coordinator |
| `DS4_DIST_MSG_WORK` | 3 | Inference work assignment (layer range, tokens) |
| `DS4_DIST_MSG_RESULT` | 4 | Computed logits and hidden states |
| `DS4_DIST_MSG_SNAPSHOT_*` | 5+ | Checkpoint state synchronization |

## Low-Level Framing Functions

### Writing Frame Headers

The `dist_write_frame_header()` function serializes the header with proper byte-order conversion:

```c
int dist_write_frame_header(int fd, uint32_t type, uint32_t bytes) {
    ds4_dist_frame_header h;
    h.magic = htonl(DS4_DIST_MAGIC);
    h.type  = htonl(type);
    h.bytes = htonl(bytes);
    return dist_write_full(fd, &h, sizeof(h));
}

```

This pattern appears throughout [`ds4_distributed.c`](https://github.com/antirez/ds4/blob/main/ds4_distributed.c) whenever the coordinator dispatches work or workers return results.

### Reading and Validating Frames

Symmetrically, `dist_read_frame_header()` handles incoming frames:

```c
int dist_read_frame_header(int fd, uint32_t *type, uint32_t *bytes,
                           char *err, size_t errlen) {
    ds4_dist_frame_header h;
    if (dist_read_full(fd, &h, sizeof(h)) != 0) {
        snprintf(err, errlen, "Read error: %s", strerror(errno));
        return -1;
    }
    uint32_t magic = ntohl(h.magic);
    if (magic != DS4_DIST_MAGIC) {
        snprintf(err, errlen, "Bad magic: 0x%08x", magic);
        return -1;
    }
    *type  = ntohl(h.type);
    *bytes = ntohl(h.bytes);
    return 0;
}

```

The magic validation catches protocol version mismatches or connection corruption early.

## Transport Layer: TCP with Optional RDMA

### Socket Initialization

Both coordinator and worker endpoints configure **low-latency TCP** via `dist_set_socket_low_latency()`:

```c
void dist_set_socket_low_latency(int fd) {
    int yes = 1;
    setsockopt(fd, IPPROTO_TCP, TCP_NODELAY, &yes, sizeof(yes));
    /* Large buffers for bulk tensor transfer */
    int bufsize = 4 * 1024 * 1024;
    setsockopt(fd, SOL_SOCKET, SO_SNDBUF, &bufsize, sizeof(bufsize));
    setsockopt(fd, SOL_SOCKET, SO_RCVBUF, &bufsize, sizeof(bufsize));
}

```

This eliminates Nagle's algorithm delays and accommodates the large payloads typical in transformer inference.

### RDMA Negotiation

When both endpoints detect RDMA-capable hardware, the protocol upgrades transport. In [`ds4_tp.c`](https://github.com/antirez/ds4/blob/main/ds4_tp.c), the `rdma_ok` flag exchanges during HELLO:

```c
/* During handshake: exchange RDMA capability */
if (peer_caps & DS4_CAP_RDMA) {
    tp->rdma_ctx = ds4_rdma_init(fd);
    if (tp->rdma_ctx) {
        tp->use_rdma = 1;
        /* Subsequent large tensor transfers bypass TCP stack */
    }
}

```

If RDMA initialization fails, the connection transparently falls back to TCP.

## Practical Communication Patterns

### Coordinator Dispatching Work

A typical work assignment flows through these steps:

1. **Serialize work descriptor** (layer range, token IDs, attention state)
2. **Write WORK frame header**
3. **Stream payload** (possibly zero-copy via `sendfile()` for large tensors)
4. **Block awaiting RESULT frame**

```c
/* Coordinator side: dispatch inference slice */
void dist_send_work(int fd, ds4_work_request *req) {
    size_t payload_len = serialize_work(req, buffer);
    
    dist_write_frame_header(fd, DS4_DIST_MSG_WORK, payload_len);
    dist_write_full(fd, buffer, payload_len);
    
    /* Immediate read of response (synchronous per-worker) */
    uint32_t type, bytes;
    dist_read_frame_header(fd, &type, &bytes, err, sizeof(err));
    assert(type == DS4_DIST_MSG_RESULT);
    /* ... process result ... */
}

```

### Worker Processing Loop

Workers run an event loop that demultiplexes incoming frames:

```c
void worker_event_loop(int coord_fd) {
    while (!terminate) {
        uint32_t type, bytes;
        if (dist_read_frame_header(coord_fd, &type, &bytes, err, sizeof(err)) < 0)
            break;
            
        switch (type) {
        case DS4_DIST_MSG_WORK:
            handle_work(coord_fd, bytes);
            break;
        case DS4_DIST_MSG_SNAPSHOT_LOAD:
            restore_checkpoint(coord_fd, bytes);
            break;
        default:
            send_error(coord_fd, "Unknown message type");
        }
    }
}

```

The `handle_work()` function deserializes the request, executes the assigned transformer layers via `ds4_compute_slice()`, then returns a RESULT frame with computed outputs.

### Error Propagation

Workers report failures through dedicated ERROR frames, carrying human-readable strings plus optional error codes:

```c
void dist_send_error(int fd, const char *fmt, ...) {
    char msg[256];
    va_list ap;
    va_start(ap, fmt);
    vsnprintf(msg, sizeof(msg), fmt, ap);
    va_end(ap);
    
    dist_write_frame_header(fd, DS4_DIST_MSG_ERROR, strlen(msg));
    dist_write_full(fd, msg, strlen(msg));
}

```

The coordinator treats ERROR frames as fatal for that worker connection, triggering reconnection or work reassignment.

## Key Source Files

| File | Responsibility | Notable Functions |
|------|---------------|-------------------|
| [`ds4_distributed.c`](https://github.com/antirez/ds4/blob/main/ds4_distributed.c) | Core protocol implementation | `dist_write_frame_header()`, `dist_read_frame_header()`, `dist_set_socket_low_latency()` |
| [`ds4_distributed.h`](https://github.com/antirez/ds4/blob/main/ds4_distributed.h) | Protocol constants and structures | `DS4_DIST_MAGIC`, `ds4_dist_frame_header`, message type enums |
| [`ds4_tp.c`](https://github.com/antirez/ds4/blob/main/ds4_tp.c) | Transport abstraction (TCP/RDMA) | `ds4_tp_init()`, RDMA capability negotiation |
| [`ds4_server.c`](https://github.com/antirez/ds4/blob/main/ds4_server.c) | Coordinator entry point | Connection acceptance, worker pool management |
| [`ds4_agent.c`](https://github.com/antirez/ds4/blob/main/ds4_agent.c) | Worker entry point | Connection to coordinator, work execution loop |

## Summary

- **Framing**: 12-byte header with magic `0x44533444`, type identifier, payload length
- **Transport**: TCP with `TCP_NODELAY` and large buffers; transparent RDMA upgrade when available
- **Message types**: HELLO (handshake), WORK (assignment), RESULT (output), ERROR (failure propagation), SNAPSHOT_* (checkpointing)
- **Byte order**: Network (big-endian) throughout; conversion via `htonl()`/`ntohl()`
- **Error handling**: Magic validation on every frame; explicit ERROR messages for semantic failures

## Frequently Asked Questions

### Does DS4 use HTTP, gRPC, or another standard protocol for distributed inference?

No. DS4 implements a **custom binary protocol** optimized for low-latency tensor transmission. Standard protocols introduce parsing overhead and framing inefficiencies unsuitable for sub-millisecond inference coordination. The magic-number framing provides fast stream validation without JSON or Protobuf parsing costs.

### How does the protocol detect version mismatches between coordinator and worker?

The **magic number validation** catches immediate corruption, while **HELLO message exchange** carries version metadata. The `ds4_dist_hello_fixed` structure includes `model_id`, `quant_bits`, and protocol capability flags. If the coordinator detects incompatible settings, it sends an ERROR frame and closes the connection before any work assignment.

### What happens when RDMA is available but initialization fails?

The transport layer in [`ds4_tp.c`](https://github.com/antirez/ds4/blob/main/ds4_tp.c) implements **graceful degradation**. If `ds4_rdma_init()` returns NULL (hardware unavailable, permissions denied, or configuration error), the connection continues with standard TCP. The `tp->use_rdma` flag remains zero, and all subsequent `dist_write_full()` calls use the socket file descriptor directly. This fallback path requires no changes to framing or message handling.