Msg: High-Performance C++20 Pub/Sub Messaging
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:5orzstd: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 withnone,lz4, andzstdon 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::spanfor 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_ptrpointer 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 (
cargodetection for Zenoh) or configurable system packages (-DMSG_USE_SYSTEM_DEPS=ON).