Files
logos-libp2p-module/src/plugin.cpp
T

351 lines
12 KiB
C++
Raw Normal View History

2026-02-23 11:37:43 -03:00
#include "plugin.h"
2026-05-07 14:34:51 -04:00
#include <chrono>
#include <cstdio>
2026-02-23 11:37:43 -03:00
#include <cstring>
2026-05-07 14:34:51 -04:00
#include <thread>
2026-02-23 11:37:43 -03:00
2026-05-07 14:34:51 -04:00
using json = nlohmann::json;
namespace {
std::string defaultListenAddr(int transport) {
if (transport == LIBP2P_TRANSPORT_QUIC) {
return "/ip4/127.0.0.1/udp/0/quic-v1";
}
return "/ip4/127.0.0.1/tcp/0";
}
}
2026-05-07 14:34:51 -04:00
void Libp2pModuleImpl::promiseCallback(int ret, const char* msg, size_t len, void* userData) {
auto* p = static_cast<SyncPromise*>(userData);
SyncResult r;
r.ok = (ret == RET_OK);
r.message = (msg && len > 0) ? std::string(msg, len) : std::string();
p->set_value(std::move(r));
delete p;
}
void Libp2pModuleImpl::promiseBufferCallback(int ret, const uint8_t* data, size_t dataLen,
const char* msg, size_t len, void* userData) {
auto* p = static_cast<SyncPromise*>(userData);
SyncResult r;
r.ok = (ret == RET_OK);
r.message = (msg && len > 0) ? std::string(msg, len) : std::string();
if (data && dataLen > 0) {
r.buffer.assign(data, data + dataLen);
}
p->set_value(std::move(r));
delete p;
}
void Libp2pModuleImpl::emitEventSafe(const std::string& name, const std::string& data) const {
if (emitEvent) {
emitEvent(name, data);
}
}
Libp2pModuleImpl::Libp2pModuleImpl(const Libp2pModuleOptions& options)
: ctx(nullptr)
2026-02-23 11:37:43 -03:00
{
std::memset(&m_libp2pConfig, 0, sizeof(m_libp2pConfig));
2026-02-23 11:37:43 -03:00
2026-04-27 09:56:03 -03:00
m_libp2pConfig.mount_gossipsub = options.mountGossipsub ? 1 : 0;
m_libp2pConfig.gossipsub_trigger_self = options.gossipsubTriggerSelf ? 1 : 0;
2026-02-23 11:37:43 -03:00
m_libp2pConfig.max_connections = options.maxConnections;
m_libp2pConfig.max_in = options.maxInConnections;
m_libp2pConfig.max_out = options.maxOutConnections;
m_libp2pConfig.max_conns_per_peer = options.maxConnsPerPeer;
2026-03-02 18:55:16 -03:00
m_libp2pConfig.autonat = options.autonat ? 1 : 0;
m_libp2pConfig.autonat_v2 = options.autonatV2 ? 1 : 0;
m_libp2pConfig.autonat_v2_server = options.autonatV2Server ? 1 : 0;
m_libp2pConfig.circuit_relay = options.circuitRelay ? 1 : 0;
2026-03-14 10:12:36 +00:00
m_libp2pConfig.circuit_relay_client = options.circuitRelayClient ? 1 : 0;
m_libp2pConfig.transport = options.transport;
2026-03-02 18:55:16 -03:00
2026-05-07 14:34:51 -04:00
m_addrs = options.addrs;
if (m_addrs.empty()) {
m_addrs.push_back(defaultListenAddr(options.transport));
2026-03-02 18:55:16 -03:00
}
2026-06-22 13:22:57 -03:00
m_libp2pConfig.addrs = m_addrsPtr.data();
m_libp2pConfig.addrsLen = static_cast<int>(m_addrsPtr.size());
2026-03-02 18:55:16 -03:00
m_addrsPtr.reserve(m_addrs.size());
for (const auto& addr : m_addrs) {
m_addrsPtr.push_back(addr.c_str());
}
m_libp2pConfig.addrs = m_addrsPtr.data();
m_libp2pConfig.addrsLen = static_cast<int>(m_addrsPtr.size());
2026-05-07 14:34:51 -04:00
if (!options.bootstrapNodes.empty()) {
m_peerIdStorage.reserve(options.bootstrapNodes.size());
m_addrStorage.reserve(options.bootstrapNodes.size());
m_addrPtrStorage.reserve(options.bootstrapNodes.size());
m_bootstrapCNodes.reserve(options.bootstrapNodes.size());
2026-03-02 18:55:16 -03:00
2026-05-07 14:34:51 -04:00
for (const auto& [peerId, addrs] : options.bootstrapNodes) {
m_peerIdStorage.push_back(peerId);
m_addrStorage.push_back(addrs);
2026-02-23 11:37:43 -03:00
2026-05-07 14:34:51 -04:00
std::vector<const char*> ptrs;
ptrs.reserve(addrs.size());
for (const auto& a : m_addrStorage.back()) {
ptrs.push_back(a.c_str());
2026-02-23 11:37:43 -03:00
}
2026-05-07 14:34:51 -04:00
m_addrPtrStorage.push_back(std::move(ptrs));
2026-02-23 11:37:43 -03:00
libp2p_bootstrap_node_t node{};
2026-05-07 14:34:51 -04:00
node.peerId = m_peerIdStorage.back().c_str();
node.multiaddrs = m_addrPtrStorage.back().data();
node.multiaddrsLen = static_cast<int>(m_addrPtrStorage.back().size());
2026-02-23 11:37:43 -03:00
m_bootstrapCNodes.push_back(node);
}
m_libp2pConfig.kad_bootstrap_nodes = m_bootstrapCNodes.data();
2026-05-07 14:34:51 -04:00
m_libp2pConfig.kad_bootstrap_nodes_len = static_cast<int>(m_bootstrapCNodes.size());
2026-02-23 11:37:43 -03:00
}
2026-04-27 09:56:03 -03:00
m_libp2pConfig.mount_kad = options.mountKad ? 1 : 0;
m_libp2pConfig.mount_service_discovery = options.mountServiceDiscovery ? 1 : 0;
2026-02-23 11:37:43 -03:00
2026-05-07 14:34:51 -04:00
auto keyResult = newPrivateKey();
if (!keyResult.success) {
fprintf(stderr, "libp2p_new_private_key failed: %s\n", keyResult.error.c_str());
return;
2026-02-23 11:37:43 -03:00
}
2026-05-07 14:34:51 -04:00
std::string keyStr = keyResult.value.get<std::string>();
m_privKey.assign(keyStr.begin(), keyStr.end());
2026-02-23 11:37:43 -03:00
2026-06-12 12:10:24 -03:00
m_libp2pConfig.priv_key.data = m_privKey.data();
2026-05-07 14:34:51 -04:00
m_libp2pConfig.priv_key.dataLen = static_cast<int>(m_privKey.size());
2026-02-23 11:37:43 -03:00
2026-05-07 14:34:51 -04:00
auto* p = new SyncPromise();
auto f = p->get_future();
2026-02-23 11:37:43 -03:00
2026-05-07 14:34:51 -04:00
ctx = libp2p_new(&m_libp2pConfig, &Libp2pModuleImpl::promiseCallback, p);
2026-02-23 11:37:43 -03:00
2026-06-12 12:10:24 -03:00
auto r = awaitResult(f, 5000);
if (!r.ok) {
fprintf(stderr, "libp2p_new failed: %s\n", r.message.c_str());
}
2026-05-07 14:34:51 -04:00
if (!ctx) {
fprintf(stderr, "libp2p_new returned null context\n");
2026-02-23 11:37:43 -03:00
}
}
2026-05-07 14:34:51 -04:00
Libp2pModuleImpl::~Libp2pModuleImpl() {
2026-06-12 12:10:24 -03:00
try {
std::vector<uint64_t> streamIds;
{
std::shared_lock<std::shared_mutex> lock(m_streamsLock);
streamIds.reserve(m_streams.size());
for (const auto& [id, _] : m_streams) {
streamIds.push_back(id);
}
2026-05-07 14:34:51 -04:00
}
2026-02-23 11:37:43 -03:00
2026-06-12 12:10:24 -03:00
// Bounded worker pool — each streamRelease awaits up to 10s.
if (!streamIds.empty()) {
const size_t workers = std::min<size_t>(
streamIds.size(),
std::max<unsigned>(std::thread::hardware_concurrency(), 4u));
2026-02-23 11:37:43 -03:00
2026-06-12 12:10:24 -03:00
std::atomic<size_t> next{0};
std::vector<std::thread> pool;
pool.reserve(workers);
for (size_t i = 0; i < workers; ++i) {
pool.emplace_back([this, &streamIds, &next] {
try {
for (;;) {
size_t idx = next.fetch_add(1, std::memory_order_relaxed);
if (idx >= streamIds.size()) return;
streamRelease(streamIds[idx]);
}
} catch (...) {}
});
}
for (auto& t : pool) {
if (t.joinable()) t.join();
}
}
2026-02-23 11:37:43 -03:00
2026-06-12 12:10:24 -03:00
if (ctx) {
auto* p = new SyncPromise();
auto f = p->get_future();
libp2p_destroy(ctx, &Libp2pModuleImpl::promiseCallback, p);
awaitResult(f, 5000);
ctx = nullptr;
}
} catch (...) {}
2026-02-23 11:37:43 -03:00
}
2026-05-07 14:34:51 -04:00
StdLogosResult Libp2pModuleImpl::start() {
2026-06-12 12:10:24 -03:00
return callSync("Failed to start libp2p", [&](SyncPromise* p) {
return libp2p_start(ctx, &Libp2pModuleImpl::promiseCallback, p);
});
2026-02-23 11:37:43 -03:00
}
2026-05-07 14:34:51 -04:00
StdLogosResult Libp2pModuleImpl::stop() {
2026-06-12 12:10:24 -03:00
return callSync("Failed to stop libp2p", [&](SyncPromise* p) {
return libp2p_stop(ctx, &Libp2pModuleImpl::promiseCallback, p);
});
2026-05-07 14:34:51 -04:00
}
2026-02-23 11:37:43 -03:00
2026-05-07 14:34:51 -04:00
StdLogosResult Libp2pModuleImpl::publicKey() {
2026-06-12 12:10:24 -03:00
return callSyncWith("Failed to get public key",
[&](SyncPromise* p) {
return libp2p_public_key(ctx, &Libp2pModuleImpl::promiseBufferCallback, p);
},
[](const SyncResult& r) -> StdLogosResult {
return {true, std::string(r.buffer.begin(), r.buffer.end()), ""};
});
2026-05-07 14:34:51 -04:00
}
StdLogosResult Libp2pModuleImpl::newPrivateKey() {
2026-06-12 12:10:24 -03:00
// Doesn't need a ctx, so it bypasses callSync's ctx check.
2026-05-07 14:34:51 -04:00
auto* p = new SyncPromise();
auto f = p->get_future();
int ret = libp2p_new_private_key(LIBP2P_PK_SECP256K1,
&Libp2pModuleImpl::promiseBufferCallback, p);
2026-06-12 12:10:24 -03:00
if (ret != RET_OK) {
delete p;
return {false, {}, "Failed to generate private key (ret=" + std::to_string(ret) + ")"};
}
2026-05-07 14:34:51 -04:00
auto r = awaitResult(f);
if (!r.ok) return {false, {}, r.message};
return {true, std::string(r.buffer.begin(), r.buffer.end()), ""};
}
StdLogosResult Libp2pModuleImpl::toCid(const std::string& key) {
if (key.empty()) return {false, {}, "Key is empty"};
2026-06-12 12:10:24 -03:00
return callSyncWith("Failed to create CID",
[&](SyncPromise* p) {
return libp2p_create_cid(
1, "dag-pb", "sha2-256",
reinterpret_cast<const uint8_t*>(key.data()), key.size(),
&Libp2pModuleImpl::promiseCallback, p);
},
[](const SyncResult& r) -> StdLogosResult {
return {true, r.message, ""};
});
2026-02-23 11:37:43 -03:00
}
2026-05-07 14:34:51 -04:00
void Libp2pModuleImpl::eventCallback(int ret, const char* msg, size_t len, void* userData) {
auto* self = static_cast<Libp2pModuleImpl*>(userData);
if (!self) return;
2026-02-23 11:37:43 -03:00
2026-05-07 14:34:51 -04:00
std::string message = (msg && len > 0) ? std::string(msg, len) : std::string();
json j;
j["result"] = ret;
j["message"] = message;
self->emitEventSafe("libp2pEvent", j.dump());
2026-02-23 11:37:43 -03:00
}
2026-05-07 14:34:51 -04:00
bool Libp2pModuleImpl::setEventCallback() {
if (!ctx) return false;
libp2p_set_event_callback(ctx, &Libp2pModuleImpl::eventCallback, this);
2026-02-23 11:37:43 -03:00
return true;
}
2026-05-07 14:34:51 -04:00
StdLogosResult Libp2pModuleImpl::connectPeer(
const std::string& peerId,
const std::vector<std::string>& multiaddrs,
int64_t timeoutMs)
{
std::vector<const char*> addrPtrs;
addrPtrs.reserve(multiaddrs.size());
for (const auto& addr : multiaddrs) {
addrPtrs.push_back(addr.c_str());
}
2026-06-12 12:10:24 -03:00
return callSync("Failed to connect", [&](SyncPromise* p) {
return libp2p_connect(ctx, peerId.c_str(), addrPtrs.data(),
static_cast<int>(addrPtrs.size()), timeoutMs,
&Libp2pModuleImpl::promiseCallback, p);
});
2026-05-07 14:34:51 -04:00
}
StdLogosResult Libp2pModuleImpl::disconnectPeer(const std::string& peerId) {
2026-06-12 12:10:24 -03:00
return callSync("Failed to disconnect", [&](SyncPromise* p) {
return libp2p_disconnect(ctx, peerId.c_str(),
&Libp2pModuleImpl::promiseCallback, p);
});
2026-05-07 14:34:51 -04:00
}
StdLogosResult Libp2pModuleImpl::peerInfo() {
2026-06-12 12:10:24 -03:00
return callSyncWith("Failed to get peer info",
[&](SyncPromise* p) {
return libp2p_peerinfo(ctx, &Libp2pModuleImpl::promisePeerInfoCallback, p);
},
[](const SyncResult& r) { return parseJsonResponse(r.message, "peerInfo"); });
2026-05-07 14:34:51 -04:00
}
StdLogosResult Libp2pModuleImpl::connectedPeers(int64_t direction) {
2026-06-12 12:10:24 -03:00
return callSyncWith("Failed to get connected peers",
[&](SyncPromise* p) {
return libp2p_connected_peers(ctx, direction,
&Libp2pModuleImpl::promisePeersCallback, p);
},
[](const SyncResult& r) -> StdLogosResult {
if (r.message.empty()) return {true, json::array(), ""};
return parseJsonResponse(r.message, "connectedPeers");
});
2026-05-07 14:34:51 -04:00
}
StdLogosResult Libp2pModuleImpl::dial(const std::string& peerId, const std::string& proto) {
2026-06-12 12:10:24 -03:00
return callSyncWith("Failed to dial",
[&](SyncPromise* p) {
return libp2p_dial(ctx, peerId.c_str(), proto.c_str(),
&Libp2pModuleImpl::promiseConnectionCallback, p);
},
[this](const SyncResult& r) -> StdLogosResult {
auto* stream = static_cast<libp2p_stream_t*>(r.extra);
if (stream) return {true, addStream(stream), ""};
return {true, 0, ""};
});
2026-05-07 14:34:51 -04:00
}
StdLogosResult Libp2pModuleImpl::circuitRelayReserve(
const std::string& relayPeerId,
const std::vector<std::string>& relayAddrs)
{
std::vector<const char*> addrPtrs;
addrPtrs.reserve(relayAddrs.size());
for (const auto& addr : relayAddrs) {
addrPtrs.push_back(addr.c_str());
}
2026-06-12 12:10:24 -03:00
return callSyncWith("Failed to reserve relay",
[&](SyncPromise* p) {
return libp2p_circuit_relay_reserve(ctx, relayPeerId.c_str(),
addrPtrs.data(),
static_cast<int>(addrPtrs.size()),
&Libp2pModuleImpl::promiseReservationCallback, p);
},
[](const SyncResult& r) -> StdLogosResult {
if (r.message.empty()) return {true, json::array(), ""};
return parseJsonResponse(r.message, "circuitRelayReserve");
});
2026-05-07 14:34:51 -04:00
}
StdLogosResult Libp2pModuleImpl::dialCircuitRelay(
const std::string& dstPeerId,
const std::string& multiaddr,
const std::string& proto)
{
2026-06-12 12:10:24 -03:00
return callSyncWith("Failed to dial circuit relay",
[&](SyncPromise* p) {
return libp2p_dial_circuit_relay(ctx, dstPeerId.c_str(),
multiaddr.c_str(), proto.c_str(),
&Libp2pModuleImpl::promiseConnectionCallback, p);
},
[this](const SyncResult& r) -> StdLogosResult {
auto* stream = static_cast<libp2p_stream_t*>(r.extra);
if (stream) return {true, addStream(stream), ""};
return {true, 0, ""};
});
2026-05-07 14:34:51 -04:00
}