diff --git a/src/custom_handlers.cpp b/src/custom_handlers.cpp index ed53662..4043b04 100644 --- a/src/custom_handlers.cpp +++ b/src/custom_handlers.cpp @@ -39,6 +39,7 @@ StdLogosResult Libp2pModuleImpl::mountProtocol(const std::string& proto) { // Without emitEvent, protocolHandler would register a stream that no caller // could ever read, close, or release — leaking the stream. if (!emitEvent) return {false, {}, "emitEvent must be set before mounting a protocol"}; + publishEmitEvent(); auto handlerCtx = std::make_unique(); handlerCtx->instance = this; diff --git a/src/gossipsub.cpp b/src/gossipsub.cpp index 7252d2b..3c20b7f 100644 --- a/src/gossipsub.cpp +++ b/src/gossipsub.cpp @@ -43,6 +43,7 @@ void Libp2pModuleImpl::gossipsubResultCallback(int ret, const char* msg, size_t StdLogosResult Libp2pModuleImpl::gossipsubSubscribe(const std::string& topic) { if (!ctx) return {false, {}, "No libp2p context"}; + publishEmitEvent(); auto subCtx = std::make_unique(); subCtx->instance = this; @@ -67,6 +68,7 @@ StdLogosResult Libp2pModuleImpl::gossipsubSubscribe(const std::string& topic) { StdLogosResult Libp2pModuleImpl::gossipsubUnsubscribe(const std::string& topic) { if (!ctx) return {false, {}, "No libp2p context"}; + publishEmitEvent(); SubscribeCtx* ctxPtr = nullptr; { diff --git a/src/plugin.cpp b/src/plugin.cpp index 9c03176..8dcfd76 100644 --- a/src/plugin.cpp +++ b/src/plugin.cpp @@ -38,10 +38,21 @@ void Libp2pModuleImpl::promiseBufferCallback(int ret, const uint8_t* data, size_ delete p; } +void Libp2pModuleImpl::publishEmitEvent() { + std::unique_lock lock(m_emitEventLock); + m_emitEventSnapshot = emitEvent; +} + void Libp2pModuleImpl::emitEventSafe(const std::string& name, const std::string& data) const { - if (emitEvent) { - emitEvent(name, data); + EmitEventFn fn; + { + std::shared_lock lock(m_emitEventLock); + fn = m_emitEventSnapshot; } + if (!fn) { + return; + } + fn(name, data); } Libp2pModuleImpl::Libp2pModuleImpl(const Libp2pModuleOptions& options) @@ -184,6 +195,7 @@ Libp2pModuleImpl::~Libp2pModuleImpl() { } StdLogosResult Libp2pModuleImpl::start() { + publishEmitEvent(); return callSync("Failed to start libp2p", [&](SyncPromise* p) { return libp2p_start(ctx, &Libp2pModuleImpl::promiseCallback, p); }); @@ -249,6 +261,7 @@ void Libp2pModuleImpl::eventCallback(int ret, const char* msg, size_t len, void* bool Libp2pModuleImpl::setEventCallback() { if (!ctx) return false; + publishEmitEvent(); libp2p_set_event_callback(ctx, &Libp2pModuleImpl::eventCallback, this); return true; } diff --git a/src/plugin.h b/src/plugin.h index ff50ac2..549cd1f 100644 --- a/src/plugin.h +++ b/src/plugin.h @@ -377,6 +377,12 @@ private: static void mountCompleteCallback(int ret, const char* msg, size_t len, void* userData); static void eventCallback(int ret, const char* msg, size_t len, void* userData); + using EmitEventFn = std::function; + + // Lock-guarded snapshot of `emitEvent`, taken on the caller thread before any worker can emit, so worker threads never read the public field unsynchronized. + mutable std::shared_mutex m_emitEventLock; + EmitEventFn m_emitEventSnapshot; + void publishEmitEvent(); void emitEventSafe(const std::string& name, const std::string& data) const; // Wraps the new-promise / invoke / await / clean-up dance shared by every