Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -85,6 +85,7 @@ class WsStreamingDevice : public Device
static PropertyObjectPtr createDefaultConfig();

void removed() override;
void removedNoLock() override;

DeviceInfoPtr onGetInfo() override;

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -16,6 +16,10 @@

#pragma once

#include <memory>

#include <boost/asio/any_io_executor.hpp>

#include <opendaq/opendaq.h>

#include <ws-streaming/local_signal.hpp>
Expand All @@ -32,7 +36,8 @@ class WsStreamingListener
WsStreamingListener(
IContext *context,
ISignal *signal,
wss::local_signal *localSignal);
std::shared_ptr<wss::local_signal> localSignal,
boost::asio::any_io_executor executor);

~WsStreamingListener() override;

Expand All @@ -50,11 +55,12 @@ class WsStreamingListener

private:

SignalPtr _signal;
InputPortConfigPtr _port;
DataDescriptorPtr _lastDescriptor;
wss::local_signal& _localSignal;
bool _ruleType;
SignalPtr _signal;
InputPortConfigPtr _port;
DataDescriptorPtr _lastDescriptor;
std::shared_ptr<wss::local_signal> _localSignal;
boost::asio::any_io_executor _executor;
bool _ruleType;
};

END_NAMESPACE_OPENDAQ_WEBSOCKET_STREAMING
Original file line number Diff line number Diff line change
Expand Up @@ -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<wss::local_signal>(name, metadata))
, openDaqSignal(openDaqSignal)
{
}

wss::local_signal localSignal;
std::shared_ptr<wss::local_signal> localSignal;
SignalPtr openDaqSignal;
InputPortNotificationsPtr listener;
SignalPtr domainSignal;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -83,9 +83,12 @@ PropertyObjectPtr WsStreamingDevice::createDefaultConfig()
void WsStreamingDevice::removed()
{
streamingEvents.clear();
}

void WsStreamingDevice::removedNoLock()
{
streaming.release();

auto lock = getRecursiveConfigLock2();
Device::removed();
}

Expand All @@ -101,6 +104,9 @@ void WsStreamingDevice::onSignalAvailable(
wss::remote_signal_ptr domainSignal,
const DataDescriptorPtr& descriptor)
{
auto lock = getRecursiveConfigLock2();
if (this->objPtr.template asPtr<IRemovable>(true).isRemoved() || !streaming.assigned())
return;
daq::MirroredSignalConfigPtr openDaqDomainSignal;

if (domainSignal)
Expand Down Expand Up @@ -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;
Expand Down
56 changes: 40 additions & 16 deletions shared/libraries/websocket_streaming/src/ws_streaming_listener.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -15,6 +15,9 @@
*/

#include <cstdint>
#include <utility>

#include <boost/asio/post.hpp>

#include <opendaq/opendaq.h>

Expand All @@ -34,17 +37,22 @@ BEGIN_NAMESPACE_OPENDAQ_WEBSOCKET_STREAMING
WsStreamingListener::WsStreamingListener(
IContext *context,
ISignal *signal,
wss::local_signal *localSignal)
std::shared_ptr<wss::local_signal> localSignal,
boost::asio::any_io_executor executor)
: _signal(signal)
, _port(
InputPort(
context,
nullptr,
String("ws-streaming")))
, _lastDescriptor(_signal.getDescriptor())
, _localSignal(*localSignal)
, _localSignal(std::move(localSignal))
, _executor(std::move(executor))
{
_localSignal.set_metadata(
// 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,
_lastDescriptor));
Expand Down Expand Up @@ -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)
Expand All @@ -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;
}
Expand Down
13 changes: 7 additions & 6 deletions shared/libraries/websocket_streaming/src/ws_streaming_server.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -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);
}

Expand All @@ -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()
Expand All @@ -239,11 +239,12 @@ void WsStreamingServer::createListener(const SignalPtr& signal)
streamableSignal.listener = createWithImplementation<IInputPortNotifications, WsStreamingListener>(
this->template thisPtr<ComponentPtr>().getContext(),
signal,
&streamableSignal.localSignal);
streamableSignal.localSignal,
_server.executor());
reinterpret_cast<WsStreamingListener *>(streamableSignal.listener.getObject())->start();
});

streamableSignal.localSignal.on_unsubscribed.connect([
streamableSignal.localSignal->on_unsubscribed.connect([
=,
&streamableSignal,
signal_id = signal.getGlobalId().toStdString()
Expand All @@ -252,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(
Expand Down Expand Up @@ -344,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);
}

Expand Down
Loading