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 @@ -53,7 +53,8 @@ BEGIN_NAMESPACE_OPENDAQ_WEBSOCKET_STREAMING
* their metadata, and registered via addToAvailableSignals() once a descriptor can be built.
* The fetch subscription is briefly kept afterwards so an immediate application subscribe can
* take it over without wire traffic (devices may process a back-to-back unsubscribe/subscribe
* pair out of order); a sweep timer releases it if unused.
* pair out of order); a sweep timer releases it if unused. A takeover replays the cached
* descriptor to openDAQ.
*
* Once registered with openDAQ, the onAddSignal() and onRemoveSignal() functions are implemented
* to manage the subscription state of each known signal. When data is received for an active
Expand Down Expand Up @@ -174,6 +175,9 @@ class WsStreaming : public Streaming
/*! @brief Registers a signal with openDAQ, marks its initial-fetch subscription as held for takeover and publishes signals deferred on it. */
void publishSignalEntry(const std::shared_ptr<WsStreamingRemoteSignalEntry>& entry);

/*! @brief Pushes the entry's cached descriptor into openDAQ as descriptor-changed events, propagating to signals that use it as their domain. */
void emitDescriptorChangedEvents(const std::shared_ptr<WsStreamingRemoteSignalEntry>& entry);

/*! @brief Checks whether any signal is still awaiting initial metadata or deferred on an unpublished domain signal. */
bool anyInitialFetchPending() const;

Expand Down
34 changes: 20 additions & 14 deletions shared/libraries/websocket_streaming/src/ws_streaming.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -165,6 +165,9 @@ void WsStreaming::subscribeRemoteSignal(const std::string& signalId)
signalIt->second->fetchState = FetchState::None;
signalIt->second->isSubscribed = true;
triggerSubscribeAck(signalId, true);

if (signalIt->second->descriptor.assigned())
emitDescriptorChangedEvents(signalIt->second);
return;
}

Expand Down Expand Up @@ -408,20 +411,7 @@ void WsStreaming::onRemoteSignalMetadataChanged(std::weak_ptr<WsStreamingRemoteS
}

if (entry->descriptor.assigned() && entry->isPublished)
{
auto packet = DataDescriptorChangedEventPacket(entry->descriptor, nullptr);
onPacket(entry->ptr->id(), packet);

// propagate to signals that use this signal as their domain
packet = DataDescriptorChangedEventPacket(nullptr, entry->descriptor);
for (const auto& [id, dataSignalEntry] : signals)
{
if (dataSignalEntry != entry &&
dataSignalEntry->domainEntry == entry &&
dataSignalEntry->isPublished)
onPacket(dataSignalEntry->ptr->id(), packet);
}
}
emitDescriptorChangedEvents(entry);
if (entry->descriptor.assigned() && !entry->isPublished)
{
// defer until the domain publishes: publishSignalEntry() then publishes this signal too;
Expand All @@ -439,6 +429,22 @@ void WsStreaming::onRemoteSignalMetadataChanged(std::weak_ptr<WsStreamingRemoteS
}
}

void WsStreaming::emitDescriptorChangedEvents(const std::shared_ptr<WsStreamingRemoteSignalEntry>& entry)
{
auto packet = DataDescriptorChangedEventPacket(entry->descriptor, nullptr);
onPacket(entry->ptr->id(), packet);

// propagate to signals that use this signal as their domain
packet = DataDescriptorChangedEventPacket(nullptr, entry->descriptor);
for (const auto& [id, dataSignalEntry] : signals)
{
if (dataSignalEntry != entry &&
dataSignalEntry->domainEntry == entry &&
dataSignalEntry->isPublished)
onPacket(dataSignalEntry->ptr->id(), packet);
}
}

void WsStreaming::publishSignalEntry(const std::shared_ptr<WsStreamingRemoteSignalEntry>& entry)
{
LOG_I("Signal {} is now ready, publishing it", entry->ptr->id());
Expand Down
Loading