ADTF
Loading...
Searching...
No Matches
signal_listening_intf.h
Go to the documentation of this file.
1
7#pragma once
8
10
12#include <unordered_set>
13#include <vector>
14#include <chrono>
15#include <optional>
17
18namespace adtf
19{
20namespace services
21{
22namespace ant
23{
24
30class ADTF3_DEPRECATED("The Signal Registry Service and its interfaces are deprecated, as it has been superseeded by the use of the Qt5 Stream Display.")
31 ISignalListening: public adtf::ucom::IObject
32{
33 public:
35 ADTF_IID(ISignalListening, "signal_listening.ant.services.adtf.iid");
36
37 protected:
39 ~ISignalListening() = default;
40
41 public:
45 class ISignalsListener
46 {
47 public:
52 virtual void SignalAdded(const ISignalRegistry::tSignalAttributes& sSignalAttributes) = 0;
53
58 virtual void SignalRemoved(const ISignalRegistry::tSignalAttributes& sSignalAttributes) = 0;
59 };
60
64 class ISignalListener
65 {
66 public:
72 virtual void SignalUpdated(ISignalRegistry::tSignalID nSignalID, const ISignalRegistry::tSignalValue& sValue) = 0;
73 };
74
75 public:
83 virtual tResult RegisterSignalsListener(ISignalsListener& oListener) = 0;
84
90 virtual tResult UnregisterSignalsListener(ISignalsListener& oListener) = 0;
91
98 virtual tResult RequestSignalUpdates(ISignalRegistry::tSignalID nSignalID, ISignalListener& oListener) = 0;
99
106 virtual tResult CancelSignalUpdates(ISignalRegistry::tSignalID nSignalID, ISignalListener& oListener) = 0;
107};
108
109}
110
111namespace flash
112{
113
117class ADTF3_DEPRECATED("The Signal Registry Service and its interfaces are deprecated, as it has been superseeded by the use of the Qt5 Stream Display.")
118 ISignalListening: public ant::ISignalListening
119{
120 public:
122 ADTF_IID(ISignalListening, "signal_listening.flash.services.adtf.iid");
123
124 protected:
126 ~ISignalListening() = default;
127
128 public:
129 class ISignalListener
130 {
131 public:
137 virtual void SignalUpdated(ISignalRegistry::tSignalID nSignalID, const adtf::services::flash::ISignalRegistry::tSignalValueNs& sValue)
138 {
139 return SignalUpdated(nSignalID, adtf::services::ant::ISignalRegistry::tSignalValue{base::flash::duration_cast<tTimeStamp>(sValue.nTimeStamp),
140 sValue.f64Value,
141 sValue.strTextValue});
142 }
143 virtual void SignalUpdated(ISignalRegistry::tSignalID /*nSignalID*/, const adtf::services::ant::ISignalRegistry::tSignalValue& /*sValue*/)
144 {}
145 };
146
147 public:
148 using ant::ISignalListening::RequestSignalUpdates;
149 using ant::ISignalListening::CancelSignalUpdates;
150
157 virtual tResult RequestSignalUpdates(ISignalRegistry::tSignalID nSignalID, ISignalListener& oListener) = 0;
158
165 virtual tResult CancelSignalUpdates(ISignalRegistry::tSignalID nSignalID, ISignalListener& oListener) = 0;
166};
167
168}
169
170using flash::ISignalListening;
171
172namespace testing
173{
174
175namespace ant
176{
177
178struct cTestSignalValueNs : public ISignalRegistry::tSignalValueNs
179{
180 cTestSignalValueNs() = default;
181 cTestSignalValueNs(const cTestSignalValueNs&) = default;
182 cTestSignalValueNs& operator=(const cTestSignalValueNs&) = default;
183 cTestSignalValueNs(cTestSignalValueNs&&) = default;
184 cTestSignalValueNs& operator=(cTestSignalValueNs&&) = default;
185 cTestSignalValueNs(const ISignalRegistry::tSignalValueNs& oSignalValue):
186 ISignalRegistry::tSignalValueNs(oSignalValue)
187 {
188 if (oSignalValue.strTextValue)
189 {
190 m_strTextValueCopy = oSignalValue.strTextValue;
191 }
192 }
193 std::string m_strTextValueCopy;
194};
195
196class cTestListener: public ISignalListening::ISignalsListener,
197 public ISignalListening::ISignalListener
198{
199 public:
200 using tSignal = std::pair<ISignalRegistry::tSignalAttributes, std::vector<cTestSignalValueNs>>;
201 using tSignals = std::map<adtf::util::cString, tSignal>;
202
203 public:
204 cTestListener()
205 {
206 m_oUpdateCount = 0;
207 THROW_IF_FAILED(_runtime->GetObject(m_pSignalListening));
208 THROW_IF_FAILED(m_pSignalListening->RegisterSignalsListener(*this));
209 }
210
211 virtual ~cTestListener()
212 {
213 m_pSignalListening->UnregisterSignalsListener(*this);
214
215 for (const auto& strSignal: m_oRequestedSignals)
216 {
217 const auto itSignal = m_oSignals.find(strSignal.c_str());
218 if (itSignal != m_oSignals.end())
219 {
220 m_pSignalListening->CancelSignalUpdates(itSignal->second.first.nSignalID, *this);
221 }
222 }
223 }
224
225 virtual void SignalAdded(const ISignalRegistry::tSignalAttributes& sSignalAttributes) override
226 {
227 {
228 std::scoped_lock oSignalsGuard(m_oSignalsMutex);
229
230 m_oIdMap[sSignalAttributes.nSignalID] = m_oSignals.emplace(adtf::util::cString(sSignalAttributes.strName), tSignal{sSignalAttributes, {}}).first;
231
232 if (m_oRequestedSignals.find(sSignalAttributes.strName) != m_oRequestedSignals.end())
233 {
234 THROW_IF_FAILED(m_pSignalListening->RequestSignalUpdates(sSignalAttributes.nSignalID, *this));
235 }
236 }
237
238 m_oSignalsChanged.notify_all();
239 }
240
241 virtual void SignalRemoved(const ISignalRegistry::tSignalAttributes& sSignalAttributes) override
242 {
243 {
244 std::scoped_lock oGuard(m_oSignalsMutex);
245 m_oSignals.erase(sSignalAttributes.strName);
246 }
247
248 m_oSignalsChanged.notify_all();
249 }
250
251 virtual void SignalUpdated(ISignalRegistry::tSignalID nSignalID, const ISignalRegistry::tSignalValueNs& sValue) override
252 {
253 std::scoped_lock oGuard(m_oSignalsMutex);
254 ++m_oUpdateCount;
255 m_oIdMap.at(nSignalID)->second.second.push_back(cTestSignalValueNs(sValue));
256 }
257
258 void ForEachSignal(std::function<void(typename tSignals::reference&)> fnCallback)
259 {
260 std::scoped_lock oSignalsGuard(m_oSignalsMutex);
261 for (auto& oSignal : m_oSignals)
262 {
263 fnCallback(oSignal);
264 }
265 }
266
267 tSignals GetCurrentSignals()
268 {
269 std::scoped_lock oSignalsGuard(m_oSignalsMutex);
270 return m_oSignals;
271 }
272
273 size_t GetUpdateCount(const char* strSignal)
274 {
275 std::scoped_lock oSignalsGuard(m_oSignalsMutex);
276 return m_oSignals.at(strSignal).second.size();
277 }
278
279 void RequestUpdates(const char* strSignal)
280 {;
281 std::optional<ISignalRegistry::tSignalID> nId;
282
283 {
284 std::scoped_lock oSignalsGuard(m_oSignalsMutex);
285 m_oRequestedSignals.insert(strSignal);
286
287 const auto itSignal = m_oSignals.find(strSignal);
288 if (itSignal != m_oSignals.end())
289 {
290 nId = itSignal->second.first.nSignalID;
291 }
292 }
293
294 if (nId)
295 {
296 THROW_IF_FAILED(m_pSignalListening->RequestSignalUpdates(*nId, *this));
297 }
298 }
299
300 void CancelUpdates(const char* strSignal)
301 {
302 std::optional<ISignalRegistry::tSignalID> nId;
303
304 {
305 std::scoped_lock oSignalsGuard(m_oSignalsMutex);
306 m_oRequestedSignals.erase(strSignal);
307 const auto itSignal = m_oSignals.find(strSignal);
308 if (itSignal != m_oSignals.end())
309 {
310 nId = itSignal->second.first.nSignalID;
311 }
312 }
313
314 if (nId)
315 {
316 m_pSignalListening->CancelSignalUpdates(*nId, *this);
317 }
318 }
319
320 bool WaitForUpdates(std::chrono::seconds nMaxSeconds, size_t nUpdateCountEachSignal)
321 {
322 auto oBeginTimePoint = std::chrono::steady_clock::now();
323 while (std::chrono::duration_cast<std::chrono::seconds>(std::chrono::steady_clock::now() - oBeginTimePoint) < nMaxSeconds)
324 {
325 const size_t oCurrentCount = m_oUpdateCount;
326 if (oCurrentCount < nUpdateCountEachSignal * m_oSignals.size())
327 {
328 std::this_thread::sleep_for(std::chrono::milliseconds(100));
329 }
330 else
331 {
332 return true;
333 }
334 }
335 const size_t oLastCount = m_oUpdateCount;
336 return (oLastCount >= nUpdateCountEachSignal * m_oSignals.size());
337 }
338
339 bool WaitForSignal(const char* strName, std::chrono::nanoseconds tmTimeout)
340 {
341 std::unique_lock oLock(m_oSignalsMutex);
342 return m_oSignalsChanged.wait_for(oLock, tmTimeout, [&] { return m_oSignals.count(strName) > 0 ;});
343 }
344
345 private:
346 adtf::ucom::object_ptr<ISignalListening> m_pSignalListening;
347 mutable std::recursive_mutex m_oSignalsMutex;
348 std::condition_variable_any m_oSignalsChanged;
349 tSignals m_oSignals;
350 std::unordered_map<ISignalRegistry::tSignalID, tSignals::iterator> m_oIdMap;
351 std::unordered_set<std::string> m_oRequestedSignals;
352 std::atomic<size_t> m_oUpdateCount;
353};
354
355}
356
357using ant::cTestListener;
358using ant::cTestSignalValueNs;
359
360}
361
362}
363
364}
#define ADTF3_DEPRECATED(_depr_message_)
Definition adtf_base_deprecated.h:27
#define ADTF_IID(_interface, _striid)
Definition adtf_iid.h:19
A_UTILS_NS::cResult tResult
For backwards compatibility and to bring latest version into scope.
Definition result.h:736
#define THROW_IF_FAILED(s)
throws if the expression returns a failed tResult
Definition result.h:88
DestinationTimeStamp duration_cast(const SourceTimeStamp &) noexcept
Duration cast base template to converted between different time resolution.
Definition chrono.h:155
Namespace for all service interfaces provided since v3.0.
Definition rpc_object_server_registry_intf.h:82
Namespace for all service interfaces provided since v3.5.
Definition signal_listening_intf.h:112
Namespace for all service testing functionality provided since v3.0.
Definition signal_listening_intf.h:176
Namespace for all service testing functionality.
Definition signal_listening_intf.h:173
Namespace for a summary of all service interfaces provided by ADTF.
Definition rpc_object_server_registry_intf.h:80
ant::IObject IObject
Alias always bringing the latest version of ant::IObject into scope.
Definition object_intf.h:109
Namespace for all functionality provided by ADTF and its SDKs.
Definition adtf_client_connector.h:14
Definition signal_listening_intf.h:179