2026-07-23 14:19:54 +02:00
|
|
|
# Tutorial 4: Custom Protocol Handlers
|
|
|
|
|
|
|
|
|
|
In the [previous tutorial](tutorial_3_connecting_peers.md), we used the
|
|
|
|
|
built-in Ping protocol to exchange data between peers. But real
|
|
|
|
|
applications need their own protocols!
|
|
|
|
|
|
|
|
|
|
This tutorial shows you how to mount a custom protocol on a node so
|
|
|
|
|
that it can handle incoming streams from peers who dial that protocol.
|
|
|
|
|
|
|
|
|
|
## How Custom Protocols Work
|
|
|
|
|
|
|
|
|
|
A protocol in libp2p is identified by a **protocol ID string** — a
|
|
|
|
|
`/`-separated path like `/myapp/chat/1.0.0`. When a remote peer dials
|
|
|
|
|
this protocol ID, your node receives a new stream.
|
|
|
|
|
|
|
|
|
|
To handle incoming streams you:
|
|
|
|
|
1. Call `mountProtocol()` to register a protocol ID
|
|
|
|
|
2. Set an `emitEvent` callback that listens for `"protocolStream"` events
|
|
|
|
|
3. Read from and write to the stream in the event handler ⚠️
|
|
|
|
|
(see performance notes below)
|
|
|
|
|
|
|
|
|
|
The stream lifecycle on the server side is:
|
|
|
|
|
1. Receive `protocolStream` event with a `streamId`
|
|
|
|
|
2. Read data from the stream
|
|
|
|
|
3. Write data to the stream (optional)
|
|
|
|
|
4. Call `streamRelease()` when done with the stream
|
|
|
|
|
|
|
|
|
|
> **Note**: Unlike the dialing side, the server side does **not** call
|
|
|
|
|
> `streamClose()` — the peer that initiated the stream is responsible
|
|
|
|
|
> for closing it.
|
2026-07-23 20:37:33 +02:00
|
|
|
|
|
|
|
|
-----------
|
|
|
|
|
|
2026-07-23 14:19:54 +02:00
|
|
|
```cpp
|
|
|
|
|
#include <cstdio>
|
|
|
|
|
#include <chrono>
|
|
|
|
|
#include <condition_variable>
|
|
|
|
|
#include <cstdint>
|
|
|
|
|
#include <mutex>
|
|
|
|
|
#include <string>
|
|
|
|
|
#include <vector>
|
|
|
|
|
#include "plugin.h"
|
|
|
|
|
|
|
|
|
|
using json = nlohmann::json;
|
|
|
|
|
|
|
|
|
|
```
|
|
|
|
|
|
|
|
|
|
## Defining our custom protocol
|
|
|
|
|
|
|
|
|
|
We'll create an **echo protocol**: the server reads a length-prefixed
|
|
|
|
|
message and echoes it back to the client. This is a common pattern
|
|
|
|
|
for request-response protocols.
|
|
|
|
|
|
|
|
|
|
```cpp
|
|
|
|
|
const std::string kEchoProtocol = "/examples/echo/1.0.0";
|
|
|
|
|
|
|
|
|
|
int main()
|
|
|
|
|
{
|
|
|
|
|
printf("=== Tutorial 4: Custom Protocol Handlers ===\n\n");
|
|
|
|
|
|
2026-08-11 10:04:04 -03:00
|
|
|
setLogLevel("fatal");
|
2026-08-01 20:22:42 +02:00
|
|
|
|
2026-07-23 14:19:54 +02:00
|
|
|
```
|
|
|
|
|
|
|
|
|
|
## Step 1: Create and start two nodes
|
|
|
|
|
```cpp
|
|
|
|
|
Libp2pModuleOptions optsA, optsB;
|
|
|
|
|
optsA.addrs = {"/ip4/127.0.0.1/tcp/9290"};
|
|
|
|
|
optsB.addrs = {"/ip4/127.0.0.1/tcp/9291"};
|
|
|
|
|
|
|
|
|
|
Libp2pModuleImpl nodeA(optsA);
|
|
|
|
|
Libp2pModuleImpl nodeB(optsB);
|
|
|
|
|
|
2026-07-24 16:14:38 +02:00
|
|
|
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;
|
|
|
|
|
}
|
2026-07-23 14:19:54 +02:00
|
|
|
printf("Both nodes started\n");
|
|
|
|
|
|
|
|
|
|
```
|
|
|
|
|
|
|
|
|
|
## Step 2: Set up the protocol handler on Node A
|
|
|
|
|
|
|
|
|
|
We define an `emitEvent` callback on Node A. Whenever a remote peer
|
|
|
|
|
dials our protocol, a `"protocolStream"` event fires with a JSON
|
|
|
|
|
payload containing the `streamId`.
|
|
|
|
|
|
|
|
|
|
For simplicity, this example reads one message, echoes it, then
|
|
|
|
|
releases the stream. A real application would likely spawn a
|
|
|
|
|
background thread to handle concurrent streams.
|
|
|
|
|
```cpp
|
|
|
|
|
struct ServerState {
|
|
|
|
|
std::mutex mtx;
|
|
|
|
|
std::condition_variable cv;
|
|
|
|
|
uint64_t streamId = 0;
|
|
|
|
|
bool ready = false;
|
|
|
|
|
};
|
|
|
|
|
ServerState server;
|
|
|
|
|
|
|
|
|
|
nodeA.emitEvent = [&](const std::string& name, const std::string& data) {
|
|
|
|
|
if (name != "protocolStream") return;
|
|
|
|
|
|
|
|
|
|
auto j = json::parse(data);
|
|
|
|
|
uint64_t sid = j["streamId"].get<uint64_t>();
|
|
|
|
|
|
|
|
|
|
{
|
|
|
|
|
std::lock_guard<std::mutex> lock(server.mtx);
|
|
|
|
|
server.streamId = sid;
|
|
|
|
|
server.ready = true;
|
|
|
|
|
}
|
|
|
|
|
server.cv.notify_one();
|
|
|
|
|
};
|
|
|
|
|
|
|
|
|
|
```
|
|
|
|
|
|
|
|
|
|
Register the echo protocol on Node A:
|
|
|
|
|
```cpp
|
|
|
|
|
printf("Mounting protocol '%s' on Node A...\n", kEchoProtocol.c_str());
|
2026-07-24 16:14:38 +02:00
|
|
|
StdLogosResult mountRes = nodeA.mountProtocol(kEchoProtocol);
|
|
|
|
|
if (!mountRes.success) {
|
|
|
|
|
fprintf(stderr, "Failed to mount protocol: %s\n",
|
|
|
|
|
mountRes.error.c_str());
|
2026-07-23 14:19:54 +02:00
|
|
|
return 1;
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
```
|
|
|
|
|
|
|
|
|
|
## Step 3: Get Node A's address and connect Node B
|
|
|
|
|
```cpp
|
2026-07-24 16:14:38 +02:00
|
|
|
StdLogosResult infoARes = nodeA.peerInfo();
|
2026-07-23 19:29:14 +02:00
|
|
|
if (!infoARes.success) {
|
|
|
|
|
fprintf(stderr, "Failed to get Node A info: %s\n",
|
|
|
|
|
infoARes.error.c_str());
|
|
|
|
|
return 1;
|
|
|
|
|
}
|
|
|
|
|
auto infoA = infoARes.value;
|
2026-07-23 14:19:54 +02:00
|
|
|
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");
|
2026-07-24 16:14:38 +02:00
|
|
|
StdLogosResult connectRes = nodeB.connectPeer(peerIdA, addrsA, 5000);
|
|
|
|
|
if (!connectRes.success) {
|
|
|
|
|
fprintf(stderr, "Failed to connect: %s\n",
|
|
|
|
|
connectRes.error.c_str());
|
2026-07-23 14:19:54 +02:00
|
|
|
return 1;
|
|
|
|
|
}
|
|
|
|
|
printf("Connected\n");
|
|
|
|
|
|
|
|
|
|
```
|
|
|
|
|
|
|
|
|
|
## Step 4: Node B dials the echo protocol
|
|
|
|
|
|
|
|
|
|
When Node B dials our custom protocol, Node A's protocol handler
|
|
|
|
|
fires, and Node A receives a new stream.
|
|
|
|
|
```cpp
|
|
|
|
|
printf("Node B dialing '%s'...\n", kEchoProtocol.c_str());
|
2026-07-24 16:14:38 +02:00
|
|
|
StdLogosResult dialRes = nodeB.dial(peerIdA, kEchoProtocol);
|
2026-07-23 14:19:54 +02:00
|
|
|
if (!dialRes.success) {
|
|
|
|
|
fprintf(stderr, "Dial failed: %s\n", dialRes.error.c_str());
|
|
|
|
|
return 1;
|
|
|
|
|
}
|
|
|
|
|
uint64_t clientStreamId = dialRes.value.get<uint64_t>();
|
|
|
|
|
printf("Node B client stream id: %llu\n",
|
|
|
|
|
(unsigned long long)clientStreamId);
|
|
|
|
|
|
|
|
|
|
```
|
|
|
|
|
|
|
|
|
|
## Step 5: Wait for Node A to receive the stream
|
|
|
|
|
```cpp
|
|
|
|
|
uint64_t serverStreamId = 0;
|
|
|
|
|
{
|
|
|
|
|
std::unique_lock<std::mutex> lock(server.mtx);
|
|
|
|
|
if (!server.cv.wait_for(lock, std::chrono::seconds(5),
|
|
|
|
|
[&] { return server.ready; })) {
|
|
|
|
|
fprintf(stderr, "Timed out waiting for incoming stream\n");
|
|
|
|
|
return 1;
|
|
|
|
|
}
|
|
|
|
|
serverStreamId = server.streamId;
|
|
|
|
|
}
|
|
|
|
|
printf("Node A received stream id: %llu\n",
|
|
|
|
|
(unsigned long long)serverStreamId);
|
|
|
|
|
|
|
|
|
|
```
|
|
|
|
|
|
|
|
|
|
## Step 6: Node B sends a message
|
|
|
|
|
```cpp
|
|
|
|
|
std::string message = "Hello from Node B!";
|
|
|
|
|
printf("Node B sending: \"%s\"\n", message.c_str());
|
2026-07-24 16:14:38 +02:00
|
|
|
StdLogosResult writeRes = nodeB.streamWriteLp(clientStreamId, message);
|
|
|
|
|
if (!writeRes.success) {
|
|
|
|
|
fprintf(stderr, "Write failed: %s\n", writeRes.error.c_str());
|
2026-07-23 14:19:54 +02:00
|
|
|
return 1;
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
```
|
|
|
|
|
|
|
|
|
|
## Step 7: Node A reads the message and echoes it back
|
|
|
|
|
```cpp
|
2026-07-24 16:14:38 +02:00
|
|
|
StdLogosResult readRes = nodeA.streamReadLp(serverStreamId, 4096);
|
2026-07-23 14:19:54 +02:00
|
|
|
if (!readRes.success) {
|
|
|
|
|
fprintf(stderr, "Node A read failed: %s\n",
|
|
|
|
|
readRes.error.c_str());
|
|
|
|
|
return 1;
|
|
|
|
|
}
|
|
|
|
|
std::string received = base64Decode(readRes.value.get<std::string>());
|
|
|
|
|
printf("Node A received: \"%s\"\n", received.c_str());
|
|
|
|
|
|
|
|
|
|
```
|
|
|
|
|
|
|
|
|
|
Echo it back:
|
|
|
|
|
```cpp
|
2026-07-24 16:14:38 +02:00
|
|
|
StdLogosResult echoWriteRes = nodeA.streamWriteLp(serverStreamId, received);
|
|
|
|
|
if (!echoWriteRes.success) {
|
|
|
|
|
fprintf(stderr, "Node A echo write failed: %s\n",
|
|
|
|
|
echoWriteRes.error.c_str());
|
2026-07-23 14:19:54 +02:00
|
|
|
return 1;
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
```
|
|
|
|
|
|
|
|
|
|
## Step 8: Node B reads the echo
|
|
|
|
|
```cpp
|
2026-07-24 16:14:38 +02:00
|
|
|
StdLogosResult echoRes = nodeB.streamReadLp(clientStreamId, 4096);
|
2026-07-23 14:19:54 +02:00
|
|
|
if (!echoRes.success) {
|
|
|
|
|
fprintf(stderr, "Node B read echo failed: %s\n",
|
|
|
|
|
echoRes.error.c_str());
|
|
|
|
|
return 1;
|
|
|
|
|
}
|
|
|
|
|
std::string echo = base64Decode(echoRes.value.get<std::string>());
|
|
|
|
|
printf("Node B received echo: \"%s\"\n", echo.c_str());
|
|
|
|
|
|
|
|
|
|
if (echo != message) {
|
|
|
|
|
fprintf(stderr, "Echo mismatch! Got '%s'\n", echo.c_str());
|
|
|
|
|
return 1;
|
|
|
|
|
}
|
|
|
|
|
printf("Echo verified successfully!\n");
|
|
|
|
|
|
|
|
|
|
```
|
|
|
|
|
|
|
|
|
|
## Step 9: Clean up
|
|
|
|
|
|
|
|
|
|
The initiator (the peer that dialed) is responsible for closing the
|
|
|
|
|
stream. The responder only needs to release their handle:
|
|
|
|
|
```cpp
|
|
|
|
|
// Server side (responder) - just release the handle
|
|
|
|
|
nodeA.streamRelease(serverStreamId);
|
|
|
|
|
|
|
|
|
|
// Client side (initiator) - close with EOF, then release
|
|
|
|
|
nodeB.streamCloseWithEOF(clientStreamId);
|
|
|
|
|
nodeB.streamRelease(clientStreamId);
|
|
|
|
|
|
|
|
|
|
nodeA.stop();
|
|
|
|
|
nodeB.stop();
|
|
|
|
|
|
|
|
|
|
printf("\n=== Tutorial 4 Complete ===\n");
|
|
|
|
|
|
|
|
|
|
return 0;
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
```
|
|
|
|
|
|
|
|
|
|
## Key Takeaways
|
|
|
|
|
|
|
|
|
|
- `mountProtocol()` registers a handler on the server side
|
|
|
|
|
- `emitEvent` with the `"protocolStream"` event delivers incoming streams
|
|
|
|
|
- The dialing side uses `streamClose()`/`streamCloseWithEOF()`
|
|
|
|
|
- The server side uses `streamRelease()` (no close)
|
|
|
|
|
- Use `streamWriteLp()` / `streamReadLp()` for length-prefixed messages
|
|
|
|
|
|
|
|
|
|
## Important: Performance Considerations
|
|
|
|
|
|
|
|
|
|
For production code, avoid blocking reads/writes in event handlers.
|
|
|
|
|
Instead, use asynchronous patterns where:
|
|
|
|
|
1. Event handler quickly passes stream to a queue/worker
|
|
|
|
|
2. Separate mechanism handles the actual I/O
|
|
|
|
|
3. Event loop stays responsive for other connections
|
|
|
|
|
|
|
|
|
|
The tutorial is simplified for learning; real implementations
|
|
|
|
|
should use non-blocking patterns to maintain system responsiveness.
|
|
|
|
|
|
|
|
|
|
## Run tutorial
|
|
|
|
|
|
|
|
|
|
```bash
|
|
|
|
|
./build/tutorial/tutorial_4_custom_protocol
|
|
|
|
|
```
|
|
|
|
|
---
|
|
|
|
|
|
2026-07-23 20:04:23 +02:00
|
|
|
<p align="center"><a href="tutorial_3_connecting_peers.md">← Connecting Peers and Exchanging Data</a> | <a href="tutorial_5_kademlia_basics.md">Kademlia DHT Basics →</a></p>
|