From 3054d552db8741a80e66e44d2bc989d2187f2069 Mon Sep 17 00:00:00 2001 From: Viacheslau Kalenikau Date: Thu, 2 Jul 2026 14:29:12 +0300 Subject: [PATCH 1/3] WsStreamingDevice mutex --- .../include/websocket_streaming/ws_streaming_device.h | 1 + .../websocket_streaming/src/ws_streaming_device.cpp | 9 ++++++++- 2 files changed, 9 insertions(+), 1 deletion(-) diff --git a/shared/libraries/websocket_streaming/include/websocket_streaming/ws_streaming_device.h b/shared/libraries/websocket_streaming/include/websocket_streaming/ws_streaming_device.h index 05db676..93fffad 100644 --- a/shared/libraries/websocket_streaming/include/websocket_streaming/ws_streaming_device.h +++ b/shared/libraries/websocket_streaming/include/websocket_streaming/ws_streaming_device.h @@ -85,6 +85,7 @@ class WsStreamingDevice : public Device static PropertyObjectPtr createDefaultConfig(); void removed() override; + void removedNoLock() override; DeviceInfoPtr onGetInfo() override; diff --git a/shared/libraries/websocket_streaming/src/ws_streaming_device.cpp b/shared/libraries/websocket_streaming/src/ws_streaming_device.cpp index ef5060c..0705742 100644 --- a/shared/libraries/websocket_streaming/src/ws_streaming_device.cpp +++ b/shared/libraries/websocket_streaming/src/ws_streaming_device.cpp @@ -83,9 +83,12 @@ PropertyObjectPtr WsStreamingDevice::createDefaultConfig() void WsStreamingDevice::removed() { streamingEvents.clear(); +} +void WsStreamingDevice::removedNoLock() +{ streaming.release(); - + auto lock = getRecursiveConfigLock2(); Device::removed(); } @@ -101,6 +104,9 @@ void WsStreamingDevice::onSignalAvailable( wss::remote_signal_ptr domainSignal, const DataDescriptorPtr& descriptor) { + auto lock = getRecursiveConfigLock2(); + if (this->objPtr.template asPtr(true).isRemoved() || !streaming.assigned()) + return; daq::MirroredSignalConfigPtr openDaqDomainSignal; if (domainSignal) @@ -134,6 +140,7 @@ void WsStreamingDevice::onSignalAvailable( void WsStreamingDevice::onSignalUnavailable(wss::remote_signal_ptr signal) { + auto lock = getRecursiveConfigLock2(); auto it = streamingSignals.find(signal->id()); if (it == streamingSignals.end()) return; From 0eb0fb09b8cb024ecbe524711021e24454116688 Mon Sep 17 00:00:00 2001 From: Viacheslau Kalenikau Date: Fri, 3 Jul 2026 12:38:24 +0300 Subject: [PATCH 2/3] WsStreamingListener: move signal data publising to the executor marshal local_signal set_metadata/publish_data onto the server strand so peer I/O never races client disconnect --- .../ws_streaming_listener.h | 16 +++--- .../src/ws_streaming_listener.cpp | 52 ++++++++++++++----- .../src/ws_streaming_server.cpp | 3 +- 3 files changed, 50 insertions(+), 21 deletions(-) diff --git a/shared/libraries/websocket_streaming/include/websocket_streaming/ws_streaming_listener.h b/shared/libraries/websocket_streaming/include/websocket_streaming/ws_streaming_listener.h index 0dd5ad2..88aca97 100644 --- a/shared/libraries/websocket_streaming/include/websocket_streaming/ws_streaming_listener.h +++ b/shared/libraries/websocket_streaming/include/websocket_streaming/ws_streaming_listener.h @@ -16,6 +16,8 @@ #pragma once +#include + #include #include @@ -32,7 +34,8 @@ class WsStreamingListener WsStreamingListener( IContext *context, ISignal *signal, - wss::local_signal *localSignal); + wss::local_signal *localSignal, + boost::asio::any_io_executor executor); ~WsStreamingListener() override; @@ -50,11 +53,12 @@ class WsStreamingListener private: - SignalPtr _signal; - InputPortConfigPtr _port; - DataDescriptorPtr _lastDescriptor; - wss::local_signal& _localSignal; - bool _ruleType; + SignalPtr _signal; + InputPortConfigPtr _port; + DataDescriptorPtr _lastDescriptor; + wss::local_signal& _localSignal; + boost::asio::any_io_executor _executor; + bool _ruleType; }; END_NAMESPACE_OPENDAQ_WEBSOCKET_STREAMING diff --git a/shared/libraries/websocket_streaming/src/ws_streaming_listener.cpp b/shared/libraries/websocket_streaming/src/ws_streaming_listener.cpp index e0f9a04..7f3fc2f 100644 --- a/shared/libraries/websocket_streaming/src/ws_streaming_listener.cpp +++ b/shared/libraries/websocket_streaming/src/ws_streaming_listener.cpp @@ -15,6 +15,9 @@ */ #include +#include + +#include #include @@ -34,7 +37,8 @@ BEGIN_NAMESPACE_OPENDAQ_WEBSOCKET_STREAMING WsStreamingListener::WsStreamingListener( IContext *context, ISignal *signal, - wss::local_signal *localSignal) + wss::local_signal *localSignal, + boost::asio::any_io_executor executor) : _signal(signal) , _port( InputPort( @@ -43,7 +47,11 @@ WsStreamingListener::WsStreamingListener( String("ws-streaming"))) , _lastDescriptor(_signal.getDescriptor()) , _localSignal(*localSignal) + , _executor(std::move(executor)) { + // The listener is constructed on the streaming endpoint's strand (inside the on_subscribed handler) + // so touching _localSignal directly here is safe. All later access happens from + // packetReceived() on the acquisition thread and is send to _executor instead. _localSignal.set_metadata( descriptorToMetadata( signal, @@ -109,22 +117,34 @@ void WsStreamingListener::onDataPacketReceived(DataPacketPtr packet) auto descriptor = packet.getDataDescriptor(); + // wss::local_signal/connection/peer are not thread-safe + // packetReceived() runs on the acquisition thread + // so marshal every local_signal call onto _executor if (descriptor != _lastDescriptor) { - _localSignal.set_metadata( - descriptorToMetadata( - _signal, - descriptor)); + auto metadata = descriptorToMetadata(_signal, descriptor); + boost::asio::post( + _executor, + [localSignal = &_localSignal, metadata = std::move(metadata)]() + { + localSignal->set_metadata(metadata); + }); _lastDescriptor = descriptor; } if (packet.getRawDataSize()) - _localSignal.publish_data( - offset, - packet.getSampleCount(), - packet.getRawData(), - packet.getRawDataSize()); + // Capture the packet by value so its raw buffer stays alive until the handler runs + boost::asio::post( + _executor, + [localSignal = &_localSignal, packet, offset]() + { + localSignal->publish_data( + offset, + packet.getSampleCount(), + packet.getRawData(), + packet.getRawDataSize()); + }); } void WsStreamingListener::onEventPacketReceived(EventPacketPtr packet) @@ -138,10 +158,14 @@ void WsStreamingListener::onEventPacketReceived(EventPacketPtr packet) if (valueDescriptorChanged && newValueDescriptor != _lastDescriptor) { - _localSignal.set_metadata( - descriptorToMetadata( - _signal, - newValueDescriptor)); + // Marshal onto the streaming endpoint's strand (see onDataPacketReceived). + auto metadata = descriptorToMetadata(_signal, newValueDescriptor); + boost::asio::post( + _executor, + [localSignal = &_localSignal, metadata = std::move(metadata)]() + { + localSignal->set_metadata(metadata); + }); _lastDescriptor = newValueDescriptor; } diff --git a/shared/libraries/websocket_streaming/src/ws_streaming_server.cpp b/shared/libraries/websocket_streaming/src/ws_streaming_server.cpp index 109fbf6..f6d4a90 100644 --- a/shared/libraries/websocket_streaming/src/ws_streaming_server.cpp +++ b/shared/libraries/websocket_streaming/src/ws_streaming_server.cpp @@ -239,7 +239,8 @@ void WsStreamingServer::createListener(const SignalPtr& signal) streamableSignal.listener = createWithImplementation( this->template thisPtr().getContext(), signal, - &streamableSignal.localSignal); + &streamableSignal.localSignal, + _server.executor()); reinterpret_cast(streamableSignal.listener.getObject())->start(); }); From 8d6a4196942b6ceefae615ea9ce76b8d0890ea1a Mon Sep 17 00:00:00 2001 From: Viacheslav Kalenikov Date: Mon, 6 Jul 2026 13:37:02 +0200 Subject: [PATCH 3/3] own local_signal via shared_ptr (async publish_data/set_metadata posts hold it alive) --- .../websocket_streaming/ws_streaming_listener.h | 6 ++++-- .../websocket_streaming/ws_streaming_server.h | 4 ++-- .../src/ws_streaming_listener.cpp | 12 ++++++------ .../websocket_streaming/src/ws_streaming_server.cpp | 12 ++++++------ 4 files changed, 18 insertions(+), 16 deletions(-) diff --git a/shared/libraries/websocket_streaming/include/websocket_streaming/ws_streaming_listener.h b/shared/libraries/websocket_streaming/include/websocket_streaming/ws_streaming_listener.h index 88aca97..d61ec37 100644 --- a/shared/libraries/websocket_streaming/include/websocket_streaming/ws_streaming_listener.h +++ b/shared/libraries/websocket_streaming/include/websocket_streaming/ws_streaming_listener.h @@ -16,6 +16,8 @@ #pragma once +#include + #include #include @@ -34,7 +36,7 @@ class WsStreamingListener WsStreamingListener( IContext *context, ISignal *signal, - wss::local_signal *localSignal, + std::shared_ptr localSignal, boost::asio::any_io_executor executor); ~WsStreamingListener() override; @@ -56,7 +58,7 @@ class WsStreamingListener SignalPtr _signal; InputPortConfigPtr _port; DataDescriptorPtr _lastDescriptor; - wss::local_signal& _localSignal; + std::shared_ptr _localSignal; boost::asio::any_io_executor _executor; bool _ruleType; }; diff --git a/shared/libraries/websocket_streaming/include/websocket_streaming/ws_streaming_server.h b/shared/libraries/websocket_streaming/include/websocket_streaming/ws_streaming_server.h index 80d48a4..11c0f72 100644 --- a/shared/libraries/websocket_streaming/include/websocket_streaming/ws_streaming_server.h +++ b/shared/libraries/websocket_streaming/include/websocket_streaming/ws_streaming_server.h @@ -115,12 +115,12 @@ class WsStreamingServer : public Server const std::string& name, const wss::metadata& metadata, SignalPtr openDaqSignal) - : localSignal(name, metadata) + : localSignal(std::make_shared(name, metadata)) , openDaqSignal(openDaqSignal) { } - wss::local_signal localSignal; + std::shared_ptr localSignal; SignalPtr openDaqSignal; InputPortNotificationsPtr listener; SignalPtr domainSignal; diff --git a/shared/libraries/websocket_streaming/src/ws_streaming_listener.cpp b/shared/libraries/websocket_streaming/src/ws_streaming_listener.cpp index 7f3fc2f..98c7803 100644 --- a/shared/libraries/websocket_streaming/src/ws_streaming_listener.cpp +++ b/shared/libraries/websocket_streaming/src/ws_streaming_listener.cpp @@ -37,7 +37,7 @@ BEGIN_NAMESPACE_OPENDAQ_WEBSOCKET_STREAMING WsStreamingListener::WsStreamingListener( IContext *context, ISignal *signal, - wss::local_signal *localSignal, + std::shared_ptr localSignal, boost::asio::any_io_executor executor) : _signal(signal) , _port( @@ -46,13 +46,13 @@ WsStreamingListener::WsStreamingListener( nullptr, String("ws-streaming"))) , _lastDescriptor(_signal.getDescriptor()) - , _localSignal(*localSignal) + , _localSignal(std::move(localSignal)) , _executor(std::move(executor)) { // The listener is constructed on the streaming endpoint's strand (inside the on_subscribed handler) // so touching _localSignal directly here is safe. All later access happens from // packetReceived() on the acquisition thread and is send to _executor instead. - _localSignal.set_metadata( + _localSignal->set_metadata( descriptorToMetadata( signal, _lastDescriptor)); @@ -125,7 +125,7 @@ void WsStreamingListener::onDataPacketReceived(DataPacketPtr packet) auto metadata = descriptorToMetadata(_signal, descriptor); boost::asio::post( _executor, - [localSignal = &_localSignal, metadata = std::move(metadata)]() + [localSignal = _localSignal, metadata = std::move(metadata)]() { localSignal->set_metadata(metadata); }); @@ -137,7 +137,7 @@ void WsStreamingListener::onDataPacketReceived(DataPacketPtr packet) // Capture the packet by value so its raw buffer stays alive until the handler runs boost::asio::post( _executor, - [localSignal = &_localSignal, packet, offset]() + [localSignal = _localSignal, packet, offset]() { localSignal->publish_data( offset, @@ -162,7 +162,7 @@ void WsStreamingListener::onEventPacketReceived(EventPacketPtr packet) auto metadata = descriptorToMetadata(_signal, newValueDescriptor); boost::asio::post( _executor, - [localSignal = &_localSignal, metadata = std::move(metadata)]() + [localSignal = _localSignal, metadata = std::move(metadata)]() { localSignal->set_metadata(metadata); }); diff --git a/shared/libraries/websocket_streaming/src/ws_streaming_server.cpp b/shared/libraries/websocket_streaming/src/ws_streaming_server.cpp index f6d4a90..f9c36ac 100644 --- a/shared/libraries/websocket_streaming/src/ws_streaming_server.cpp +++ b/shared/libraries/websocket_streaming/src/ws_streaming_server.cpp @@ -208,7 +208,7 @@ void WsStreamingServer::createListener(const SignalPtr& signal) if (domainSignal != it->second.domainSignal) { - _server.remove_local_signal(it->second.localSignal); + _server.remove_local_signal(*it->second.localSignal); _localSignals.erase(it); } @@ -230,7 +230,7 @@ void WsStreamingServer::createListener(const SignalPtr& signal) streamableSignal.domainSignal = domainSignal; - streamableSignal.localSignal.on_subscribed.connect([ + streamableSignal.localSignal->on_subscribed.connect([ =, &streamableSignal, signal_id = signal.getGlobalId().toStdString() @@ -239,12 +239,12 @@ void WsStreamingServer::createListener(const SignalPtr& signal) streamableSignal.listener = createWithImplementation( this->template thisPtr().getContext(), signal, - &streamableSignal.localSignal, + streamableSignal.localSignal, _server.executor()); reinterpret_cast(streamableSignal.listener.getObject())->start(); }); - streamableSignal.localSignal.on_unsubscribed.connect([ + streamableSignal.localSignal->on_unsubscribed.connect([ =, &streamableSignal, signal_id = signal.getGlobalId().toStdString() @@ -253,7 +253,7 @@ void WsStreamingServer::createListener(const SignalPtr& signal) streamableSignal.listener.release(); }); - _server.add_local_signal(streamableSignal.localSignal); + _server.add_local_signal(*streamableSignal.localSignal); } void WsStreamingServer::onClientConnected( @@ -345,7 +345,7 @@ void WsStreamingServer::rescan() if (it->second.openDaqSignal.isRemoved()) { auto jt = it++; - _server.remove_local_signal(jt->second.localSignal); + _server.remove_local_signal(*jt->second.localSignal); _localSignals.erase(jt); }