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

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 and ds4_distributed.h, with transport abstraction in 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:

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:

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:

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 whenever the coordinator dispatches work or workers return results.

Reading and Validating Frames

Symmetrically, dist_read_frame_header() handles incoming frames:

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():

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, the rdma_ok flag exchanges during HELLO:

/* 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
/* 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:

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:

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 Core protocol implementation dist_write_frame_header(), dist_read_frame_header(), dist_set_socket_low_latency()
ds4_distributed.h Protocol constants and structures DS4_DIST_MAGIC, ds4_dist_frame_header, message type enums
ds4_tp.c Transport abstraction (TCP/RDMA) ds4_tp_init(), RDMA capability negotiation
ds4_server.c Coordinator entry point Connection acceptance, worker pool management
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 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.

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 →