Distributed systems, high-performance computing clusters, robotics platforms, and real-time simulations rarely rely on a single networking transport. One subsystem might require ultra-low-latency point-to-point UDP multicast, another relies on an enterprise message broker like NATS, a third coordinates through Eclipse Zenoh, while a fourth pushes cache updates through Redis.

Msg is a modern C++20 library designed to decouple application code from transport mechanics. It unifies transport protocols, in-flight compression, and message serialization behind a clean, URI-driven publish/subscribe interface.


URI-Based Configuration

Rather than writing transport-specific initialization boilerplate, msg configures endpoints via a concise URI schema:

<transport>;<compression>;<address>

or simplified 2-part format (defaulting publishers to none compression):

<transport>;<address>

Supported Transports

Transport Identifier Protocol Typical Use Case
ZeroMQ Pub/Sub zmq_pub_sub TCP Reliable inter-process and node-to-node telemetry
ZeroMQ Radio/Dish zmq_radio_dish UDP High-rate multicast and loss-tolerant sensor streams
NATS nats TCP Resilient cloud messaging and microservice pub/sub
Redis redis TCP In-memory pub/sub and persistent message cache
Eclipse Zenoh zenoh TCP / UDP / Peer Brokerless peer-to-peer, routed mesh, robotics & edge

Transparent 2-Bit Compression Engine

In traditional messaging libraries, publishers and subscribers must agree in advance on compression algorithms. msg features a transparent sender-driven compression model:

  • 1-Byte Message Header: Each outgoing message prepends a lightweight header containing a 2-bit compression identifier:
    • 0b00 (none): Uncompressed data (dispatched to subscriber callbacks with zero copies).
    • 0b01 (lz4): Microsecond-speed LZ4 compression for high-bandwidth real-time streams.
    • 0b10 (zstd / zstd:<level>): High-ratio Facebook Zstandard compression (e.g. zstd:5 or zstd:20) for bandwidth-constrained links.
    • 0b11: Reserved for future expansion.
  • Dynamic Subscriber Decompression: Subscribers inspect the 2-bit header per frame and decompress dynamically via direct static functions with zero heap allocations and zero virtual dispatch.
  • Heterogeneous Streams: A subscriber (e.g., initialized with zenoh;tcp://127.0.0.1:7447) can simultaneously receive from multiple publishers sending with none, lz4, and zstd on the same topic without any configuration changes.

Code Examples

1. Zenoh Pub/Sub with Transparent LZ4 Compression

Publishing with LZ4 compression while the subscriber receives transparently:

#include "msg/msg.hpp"
#include <iostream>
#include <span>
#include <string>

int main() {
    const std::string endpoint = "tcp://127.0.0.1:7447";

    // Subscriber does not need to specify compression:
    // It dynamically inspects the 2-bit header and decompresses incoming frames!
    msg::Sub sub("zenoh;" + endpoint, "telemetry", [](std::span<const uint8_t> payload) {
        std::string_view message(reinterpret_cast<const char*>(payload.data()), payload.size());
        std::cout << "Received telemetry: " << message << "\n";
    });

    // Publisher configures LZ4 compression in-flight
    msg::Pub pub("zenoh;lz4;" + endpoint, "telemetry");

    std::string payload = "SensorPacket: lat=38.2910 lon=-76.5421 alt=1200.5";
    pub.send({reinterpret_cast<const uint8_t*>(payload.data()), payload.size()});

    return 0;
}

2. Strongly-Typed Google FlatBuffers

msg natively integrates with Google FlatBuffers for zero-copy schema serialization:

#include "msg/msg.hpp"
#include "person_generated.h"

int main() {
    // Strongly-typed FlatBuffer subscriber
    msg::FbsSub<Person>("nats;nats://127.0.0.1:4222", [](const Person& person) {
        std::cout << "Received Person ID: " << person.id() 
                  << ", Name: " << person.name()->string_view() << "\n";
    });

    // Strongly-typed FlatBuffer publisher with Zstandard compression
    msg::FbsPub<Person> pub("nats;zstd:10;nats://127.0.0.1:4222");

    PersonT person;
    person.id = 12893;
    person.name = "Alice";
    person.email = "alice@example.com";

    pub.send(person);

    return 0;
}

3. JSON Pub/Sub

#include "msg/msg.hpp"
#include "nlohmann/json.hpp"
#include <iostream>

int main() {
    msg::JsonSub sub("redis;127.0.0.1:6379", "sensors", [](const nlohmann::json& data) {
        std::cout << "Sensor ID: " << data["id"] << ", Value: " << data["value"] << "\n";
    });

    msg::JsonPub pub("redis;127.0.0.1:6379", "sensors");
    pub.send({{"id", 42}, {"value", 98.6}});

    return 0;
}

Performance & Architecture Highlights

  • Zero Unnecessary Allocations: Hot transmission paths utilize C++20 std::span for zero-copy views. Uncompressed payloads dispatch straight into subscriber callbacks with zero intermediate copies.
  • Direct Static Decompression: Decompression engines execute as pure static functions, avoiding std::shared_ptr pointer dereferences and virtual method table lookups on the message receive path.
  • High-Throughput Validated: Benchmarked on local Zenoh loopback at >1.1 Million messages/sec and >1.1 GB/sec throughput with zero packet drops.
  • Dependency Management: Automatically managed via CPM.cmake with automatic system toolchain discovery (cargo detection for Zenoh) or configurable system packages (-DMSG_USE_SYSTEM_DEPS=ON).