DS4 Distributed Worker-Coordinator Communication Protocol: Binary TCP Implementation
TLDR: DS4 implements a custom lightweight binary protocol over TCP sockets that uses 12-byte headers with a "DS4D" magic constant to frame messages between coordinators and workers, supporting five core message types including handshake, work dispatch, and result return.
The antirez/ds4 repository implements a distributed inference runtime that relies on a bespoke communication layer to coordinate neural network workloads across multiple nodes. Understanding the distributed worker-coordinator communication protocol is essential for debugging network issues, extending the system, or integrating custom transport layers. This protocol is defined entirely within the distributed transport layer found in ds4_distributed.c and ds4_distributed.h.
Transport Layer and Socket Configuration
The foundation of the DS4 distributed worker-coordinator communication protocol is built on standard TCP sockets configured for low-latency inference workloads. In ds4_distributed.c, the dist_set_socket_low_latency function (lines 27-45) configures TCP_NODELAY to disable Nagle's algorithm, enables TCP keep-alive, and sets custom buffer sizes to minimize latency between nodes.
Coordinators initialize listening sockets via dist_open_listener, while workers establish outbound connections using dist_connect_endpoint. Both paths use identical framing logic, ensuring symmetric communication patterns regardless of connection direction.
Binary Message Framing and Header Structure
Every packet in the DS4 protocol starts with a fixed 12-byte header defined by the ds4_dist_frame_header struct (lines 72-76 in ds4_distributed.c). This header contains three critical fields: a magic number 0x44533444u (representing "DS4D" in ASCII), a 32-bit message type identifier, and the payload size in bytes.
The dist_write_frame_header function (lines 130-136) serializes this header to the wire, while dist_read_frame_header parses incoming frames. This framing mechanism enables the protocol to handle arbitrary binary payloads without delimiters, solving the TCP stream boundary problem inherent in raw socket communication.
Core Message Types
The protocol defines five primary message categories as unsigned constants (lines 44-53 in ds4_distributed.c):
- DS4_DIST_MSG_HELLO (1): Initial handshake where workers advertise their model ID, layer range, and quantization capabilities.
- DS4_DIST_MSG_ERROR (2): Error reporting for protocol violations or runtime failures.
- DS4_DIST_MSG_WORK (3): Work requests from coordinator to worker containing activation tensors and routing information.
- DS4_DIST_MSG_RESULT (4): Result messages from worker to coordinator carrying computed outputs.
- DS4_DIST_MSG_SNAPSHOT_ (5-9):* KV-snapshot exchange messages for distributed state saving and loading.
Data Serialization and Network Byte Order
Each message type uses a tightly-packed C struct for its payload. For example, ds4_dist_hello_fixed contains model metadata fields, while ds4_dist_work_fixed carries tensor dimensions and hash identifiers. Before transmission, all multi-byte fields are converted to network byte order using htonl, with corresponding ntohl conversions on receipt.
The dist_work_to_wire and dist_work_from_wire helper functions handle these transformations, ensuring cross-platform compatibility between coordinator and worker nodes regardless of host endianness.
Connection Lifecycle and Communication Flow
The distributed worker-coordinator communication protocol follows a strict request-response pattern. Workers initiate connections and immediately transmit a HELLO message containing their capabilities:
/* Build a hello payload */
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),
.has_output = htonl(has_output),
.has_hidden = htonl(has_hidden),
.ctx_size = htonl(ctx_size),
.n_layers = htonl(n_layers),
.listen_port = htonl(listen_port),
.model_name_len = htonl(strlen(model_name))
};
/* Write the frame header (type = DS4_DIST_MSG_HELLO) */
dist_write_frame_header(fd, DS4_DIST_MSG_HELLO, sizeof(hello));
/* Send the fixed part and then the model name string */
dist_write_full(fd, &hello, sizeof(hello));
dist_write_full(fd, model_name, strlen(model_name));
Coordinators dispatch work using the dist_send_work_frame function, which serializes a ds4_dist_work_fixed struct followed by variable-length activation data. On the receiving side, the coordinator uses dist_read_frame_header to dispatch messages based on type:
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;
}
/* Dispatch based on the message type */
switch (type) {
case DS4_DIST_MSG_HELLO:
// handle hello
break;
case DS4_DIST_MSG_WORK:
// handle work request
break;
case DS4_DIST_MSG_RESULT:
// handle result
break;
/* … other cases … */
}
Workers process the computation and return results via DS4_DIST_MSG_RESULT frames using identical framing logic.
Activation Compression and Payload Optimization
The protocol supports dynamic quantization of activation tensors to reduce bandwidth. The bits field in the work message specifies whether activations are transmitted as 32-bit floats, 16-bit half-precision values, or 8-bit "e4m3" format. The dist_write_activation_payload and dist_decode_activation_payload functions handle compression and decompression, allowing the distributed system to trade computational overhead for network efficiency.
Summary
- DS4 uses a custom binary protocol over TCP with
TCP_NODELAYenabled for low-latency communication. - The 12-byte frame header includes a "DS4D" magic constant, message type, and payload length for unambiguous message boundaries.
- Five core message types handle handshake (HELLO), error reporting, work distribution (WORK), result collection (RESULT), and state synchronization (SNAPSHOT).
- All numeric fields use network byte order conversion (
htonl/ntohl) to ensure endianness compatibility. - Activation tensors support 8-bit, 16-bit, and 32-bit precision modes to optimize bandwidth utilization.
Frequently Asked Questions
Is the DS4 protocol based on HTTP or gRPC?
No, DS4 implements a custom binary protocol specifically optimized for low-latency inference workloads. Unlike HTTP-based protocols, DS4 uses raw TCP sockets with a lightweight 12-byte header framing mechanism defined in ds4_distributed.c, eliminating the parsing overhead and verbosity associated with text-based protocols.
How does DS4 handle message boundaries over TCP?
DS4 solves the TCP stream boundary problem through explicit message framing. Every message begins with a fixed 12-byte header containing the payload length, allowing the receiver to know exactly how many bytes to read. The dist_read_frame_header function in ds4_distributed.c implements this logic, ensuring complete message delivery before dispatching to type-specific handlers.
What compression formats does DS4 support for activation tensors?
The protocol supports three numeric precisions for activation data: 32-bit IEEE-754 floats, 16-bit half-precision floats, and 8-bit "e4m3" floating-point format. The coordinator specifies the desired precision via the bits field in the WORK message, and workers use dist_decode_activation_payload to interpret the compressed data correctly.
Where is the protocol implementation located in the source tree?
The complete distributed worker-coordinator communication protocol is implemented in ds4_distributed.c with public API declarations in ds4_distributed.h. The coordinator-specific server logic resides in ds4_server.c, which utilizes these protocol primitives to manage worker connections and dispatch inference tasks.
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:
curl -s "https://instagit.com/install.md" Maintain an open-source project? Get it listed too →