6.5 KiB
Tutorial 7: GossipSub - Polling for Messages
GossipSub is a pub/sub protocol that lets peers broadcast messages to everyone subscribed to a topic. It's the foundation for many decentralized applications, from chat rooms to blockchain transaction propagation.
In this tutorial we'll:
- Subscribe two nodes to the same topic
- Wait for the GossipSub mesh to form
- Publish a message from one node
- Receive it on the other node by polling
How GossipSub Polling Works
- Peers subscribe to topics by calling
gossipsubSubscribe(topic) - libp2p builds a mesh, which is a set of peer connections per topic
- When a peer publishes to a topic, the message is forwarded through the mesh to all subscribers
- Subscribers call
gossipsubNextMessage(topic, timeout)to wait for the next message
Note
: The GossipSub mesh takes some time to form after both peers have subscribed. A short delay of 1-2 seconds is usually enough.
Note
: Messages wait in a per-topic queue bounded by
gossipsubQueueMaxMessages(1024) andgossipsubQueueMaxBytes(4 MiB). Past either bound the newest message is dropped and counted inlibp2p_module_gossipsub_queue_dropped_total. Poll fast enough to keep that counter flat.
#include <chrono>
#include <cstdio>
#include <string>
#include <thread>
#include <vector>
#include "plugin.h"
int main()
{
printf("=== Tutorial 7: GossipSub - Polling for Messages ===\n\n");
setLogLevel("fatal");
Step 1: Create two nodes with GossipSub enabled
GossipSub is enabled by default (mountGossipsub: true).
We keep it explicit here for clarity.
Libp2pModuleOptions optsA, optsB;
optsA.addrs = {"/ip4/127.0.0.1/tcp/9590"};
optsA.mountGossipsub = true;
optsB.addrs = {"/ip4/127.0.0.1/tcp/9591"};
optsB.mountGossipsub = true;
Libp2pModuleImpl nodeA(optsA);
Libp2pModuleImpl nodeB(optsB);
StdLogosResult startARes = nodeA.start();
if (!startARes.success) {
fprintf(stderr, "Node A failed: %s\n", startARes.error.c_str());
return 1;
}
StdLogosResult startBRes = nodeB.start();
if (!startBRes.success) {
fprintf(stderr, "Node B failed: %s\n", startBRes.error.c_str());
return 1;
}
printf("Both nodes started\n");
Step 2: Connect the nodes
GossipSub needs connectivity between peers to build the mesh.
StdLogosResult infoARes = nodeA.peerInfo();
if (!infoARes.success) {
fprintf(stderr, "Failed to get Node A info: %s\n",
infoARes.error.c_str());
return 1;
}
auto infoA = infoARes.value;
std::string peerIdA = infoA["peerId"].get<std::string>();
std::vector<std::string> addrsA;
for (const auto& a : infoA["addrs"])
addrsA.push_back(a.get<std::string>());
printf("Connecting Node B to Node A...\n");
StdLogosResult connectRes = nodeB.connectPeer(peerIdA, addrsA, 5000);
if (!connectRes.success) {
fprintf(stderr, "Failed to connect: %s\n",
connectRes.error.c_str());
return 1;
}
printf("Connected\n");
Step 3: Both nodes subscribe to a topic
A topic is just a string identifier. Peers who subscribe to the same topic will receive each other's messages.
std::string topic = "chat-room-1";
printf("Node B subscribing to topic: \"%s\"\n", topic.c_str());
StdLogosResult subscribeBRes = nodeB.gossipsubSubscribe(topic);
if (!subscribeBRes.success) {
fprintf(stderr, "Node B subscribe failed: %s\n",
subscribeBRes.error.c_str());
return 1;
}
printf("Node A subscribing to topic: \"%s\"\n", topic.c_str());
StdLogosResult subscribeARes = nodeA.gossipsubSubscribe(topic);
if (!subscribeARes.success) {
fprintf(stderr, "Node A subscribe failed: %s\n",
subscribeARes.error.c_str());
return 1;
}
Step 4: Wait for the GossipSub mesh to form
After subscribing, libp2p needs time to discover subscribers and build the message-forwarding mesh.
printf("Waiting for GossipSub mesh to form...\n");
std::this_thread::sleep_for(std::chrono::milliseconds(2000));
Step 5: Publish a message
Node A publishes a message to the topic. It will be forwarded to all subscribers, including Node B.
std::string payload = "Hello from Node A via GossipSub!";
printf("\nNode A publishing: \"%s\"\n", payload.c_str());
printf(" Topic: \"%s\"\n", topic.c_str());
StdLogosResult publishRes = nodeA.gossipsubPublish(topic, payload);
if (!publishRes.success) {
fprintf(stderr, "Publish failed: %s\n", publishRes.error.c_str());
return 1;
}
printf("Message published!\n");
Step 6: Receive the message on Node B
gossipsubNextMessage() blocks until a message arrives or the
timeout expires. The timeout is in milliseconds.
printf("\nNode B waiting for message...\n");
StdLogosResult res = nodeB.gossipsubNextMessage(topic, 3000);
if (!res.success) {
fprintf(stderr, "Node B did not receive any messages: %s\n",
res.error.c_str());
return 1;
}
std::string received = res.value.get<std::string>();
printf("Node B received: \"%s\"\n", received.c_str());
if (received == payload) {
printf("Message verified!\n");
} else {
fprintf(stderr,
"Message content differs (expected: \"%s\", got: \"%s\")\n",
payload.c_str(), received.c_str());
return 1;
}
Step 7: Unsubscribe and clean up
printf("\nUnsubscribing...\n");
if (!nodeB.gossipsubUnsubscribe(topic).success) {
fprintf(stderr, "Node B unsubscribe failed\n");
return 1;
}
if (!nodeA.gossipsubUnsubscribe(topic).success) {
fprintf(stderr, "Node A unsubscribe failed\n");
return 1;
}
nodeA.stop();
nodeB.stop();
printf("\n=== Tutorial 7 Complete ===\n");
return 0;
}
Key Takeaways
- GossipSub provides topic-based pub/sub messaging
- Both publisher and subscriber must subscribe to the topic
- Allow 1-2 seconds for the mesh to form
- Use
gossipsubNextMessage(topic, timeout)for polling - The per-topic queue is bounded; watch
libp2p_module_gossipsub_queue_dropped_total - Always unsubscribe and stop cleanly
Run tutorial
./build/tutorial/tutorial_7_gossipsub_polling
← Kademlia Provider Records | GossipSub - Event Callback Messages →