ADTF
Loading...
Searching...
No Matches
samplestream.h
Go to the documentation of this file.
1
7#pragma once
8#include "samplestream_intf.h"
11#include "streamtype.h"
14#include "triggerpipe_intf.h"
15#include "triggerpipe.h"
19
20#include <adtfbase/adtf_base.h>
22
23#include <mutex>
24#include <optional>
25#include <string>
26#include <memory>
27#include <unordered_map>
28#include <vector>
29
30namespace adtf
31{
32namespace streaming
33{
34
35namespace private_implementations
36{
37class cStreamingRequestsProxy;
38}
39
40namespace detail
41{
42
53 public ucom::ant::enable_object_ptr_from_this<ISampleStream>,
55{
56public:
61 cSampleStreamBase(cTriggerPipeItemBase& oPipeItem, ISampleStream& oSelf) noexcept;
63
64 cSampleStreamBase() = delete;
65 cSampleStreamBase(const cSampleStreamBase&) = delete;
66 cSampleStreamBase& operator=(const cSampleStreamBase&) = delete;
67
68public: // implements the sample stream
72 tResult SetStreamError(const tResult& oError) noexcept;
90 ISampleStream::IPushReadEventSink*& pPushEventSink,
92 size_t szQueueSize) noexcept;
105 size_t szQueueSize) noexcept;
106
107public: // implements the trigger pipe source
110 ITriggerPipeItem::tPriority ui32Prio) noexcept;
113
114public: // implements hollow::IInternalBindingProxy and vision::IBlockableSampleStream
118 void BlockSampleStreaming(bool bBlock) noexcept;
120 bool IsSampleStreamingBlocked() const noexcept;
121
122public: // implements wolverine::ISampleStreamStatistics, @see cStatistics
124 wolverine::ISampleStreamStatistics::tSampleStreamStatistic GetStatistics() noexcept;
125
127 void EnableSampleSizeStatistics(bool bEnable) noexcept;
128
130 void ResetStatistics() noexcept;
131
132public: // policies
138 void SetForwardAllTriggers(bool bForwardAll) noexcept;
139
145 void SetSubStreamFilter(const char* strSubStream) noexcept;
146
147private:
148 class cOutStream;
149 class cPushStream;
150
151 struct tItemFilter;
152
153 struct tSink;
154 using tSinkList = std::vector<tSink>;
155
156 struct tDirectChild;
157 using tDirectChildren = std::vector<tDirectChild>;
158
159 class cStatistics;
160
161private:
162 void UpdateStreamTypeHandler() noexcept;
163 void HandleNewStreamType(const ucom::ant::iobject_ptr<const IStreamType>& pStreamType) noexcept override;
164
165 void Push(const ucom::ant::iobject_ptr<const ISample>& pSample) noexcept;
166 void Push(const ucom::ant::iobject_ptr<const IStreamType>& pStreamType) noexcept;
171 bool ResolveSubstreamFilter(const ucom::ant::iobject_ptr<const IStreamType>& pStreamType,
172 ucom::ant::object_ptr<const IStreamType>& pSubstreamType,
173 uint32_t& nSubstreamId) noexcept;
176 void ReportUnresolvedSubstreamFilter() noexcept;
177
178 tResult AttachBindingProxy(const ucom::ant::iobject_ptr<ISampleStream>& pSampleStreamTo) noexcept;
179 void DetachBindingProxy(const ucom::ant::iobject_ptr<ISampleStream>& pSampleStreamTo) noexcept;
180
181 void HandleSubItemTriggerError(tResult oError) noexcept;
182 bool IsInFoldChain(const cSampleStreamBase* pStream) const noexcept;
183 const cSampleStreamBase* FindFilterAtOrAbove() const noexcept;
184 const cSampleStreamBase* FindFilterAtOrBelow() const noexcept;
185 bool HasWriter() const noexcept;
186 bool IsAttached() const noexcept;
187 bool IsAttachedToOther(const cSampleStreamBase* pStream) const noexcept;
188 void SetFoldedInto(cSampleStreamBase* pStream) noexcept;
189
190 tDirectChildren::iterator FindDirectChild(const cSampleStreamBase* pChild) noexcept;
191 tDirectChildren::iterator FindTriggerChild(const ITriggerPipeItem* pItem) noexcept;
192 tDirectChildren::iterator AcquireDirectChild(cSampleStreamBase* pChild) noexcept;
193
194 void ApplyPolicy(tSink& sSink) noexcept;
195 static const void* GetSinkKey(const tSink& sSink) noexcept;
196 bool HasDirectReader() const noexcept;
197 void CollectStatistics(cStatistics& oRoot, const tSink& sFromParent, bool bTriggerReachable) noexcept;
198
199 void ApplySubstreamFilter(tItemFilter& sFilter) noexcept;
201 const tSinkList& RebuildSinks(bool bRoot) noexcept;
202 cSampleStreamBase& GetFoldRoot() noexcept;
203 void NotifyTopologyChanged() noexcept;
204 void FlushInitialTypes() noexcept;
205 void UnregisterPushSink(cPushStream* pSink) noexcept;
206 void DropChild(cSampleStreamBase* pChild) noexcept;
207
209 std::string GetFullName() const noexcept;
210
211private:
212 // This only stores the hierachy of the trigger pipe but is no longer used to run it.
213 cTriggerPipeItemBase& m_oPipeItem;
214 ISampleStream& m_oSelf;
215
216 bool m_bForwardAllTrigger = true;
217 bool m_bSamplesBlocked = false;
218 std::string m_strSubstreamFilter;
221 ucom::ant::object_ptr<const IStreamType> m_pFilteredType;
222 ucom::ant::object_ptr<const IStreamType> m_pSubstreamType;
223 std::optional<uint32_t> m_nSubstreamIdFilter;
226 tResult m_nUnresolvedError = ERR_NOERROR;
227
230 ucom::ant::object_ptr<IBindingProxy> m_pInternalBindingProxy;
232 private_implementations::cStreamingRequestsProxy* m_pRequestsProxy = nullptr;
233 ucom::ant::object_ptr<IStreamingRequest> m_pRequest;
234
235 cSampleStreamBase* m_pFoldedInto = nullptr;
238 bool m_bInitialTypesPending = false;
241 ucom::ant::object_ptr<const IStreamType> m_pLastType;
242 mutable std::mutex m_oTypeMutex;
243 ucom::ant::weak_object_ptr<ISampleOutStream> m_pCurrentWriter;
244
246 tSinkList m_oSinks;
248 std::mutex m_oSinkMutex;
249 tDirectChildren m_oDirectChildren;
250
251 std::unique_ptr<cStatistics> m_pStatistics;
252};
253
254} // namespace detail
255
269template<typename... INTERFACES>
271 public ucom::object<INamedGraphObject,
276 INTERFACES...>>,
277 public detail::cSampleStreamBase
278{
279 using base_type = ucom::object<INamedGraphObject,
284 INTERFACES...>>;
285
286public:
288 sample_stream(): detail::cSampleStreamBase(static_cast<detail::cTriggerPipeItemBase&>(*this), *this)
289 {
290 }
291
296 sample_stream(const char* strName): sample_stream()
297 {
298 base_type::SetName(strName);
299 }
300
301public: // implements ant::ISampleStream
302 tResult GetType(ucom::ant::iobject_ptr<const IStreamType>& pStreamType) const noexcept override
303 {
304 return detail::cSampleStreamBase::GetType(pStreamType);
305 }
306
312 tTimeStamp GetTime() const noexcept override
313 {
314 return 0;
315 }
316
317 tResult SetStreamError(const tResult& oError) noexcept override
318 {
320 }
321
322 tResult AttachRouting(const ucom::ant::iobject_ptr<ISampleStream>& pSampleStreamTo) noexcept override
323 {
324 return detail::cSampleStreamBase::AttachRouting(pSampleStreamTo);
325 }
326
327 tResult DetachRouting(const ucom::ant::iobject_ptr<ISampleStream>& pSampleStreamTo) noexcept override
328 {
329 return detail::cSampleStreamBase::DetachRouting(pSampleStreamTo);
330 }
331
337 tResult Open([[maybe_unused]] const char* strName,
339 [[maybe_unused]] const ucom::ant::iobject_ptr<const IStreamType>& pInitialAcceptedStreamType,
340 ISampleStream::IPushReadEventSink*& pPushEventSink,
342 size_t szQueueSize) noexcept override
343 {
344 return detail::cSampleStreamBase::Open(pInStream, pPushEventSink, ui32Mode, szQueueSize);
345 }
346
354 tResult Open([[maybe_unused]] const char* strName,
357 size_t szQueueSize) noexcept override
358 {
359 return detail::cSampleStreamBase::Open(pOutStream, ui32Mode, szQueueSize);
360 }
361
362public: // implements ant::ITriggerPipeSource
364 ITriggerPipeItem::tPriority ui32Prio) noexcept override
365 {
366 return detail::cSampleStreamBase::RegisterSubItem(pSubRun, ui32Prio);
367 }
368
373
380 ITriggerPipeItem::tPriority /* ui32Prio */) noexcept override
381 {
382 RETURN_ERROR(ERR_NOT_SUPPORTED);
383 }
384
385public: // implements hollow::IInternalBindingProxy
390
391public: // implements vision::IBlockableSampleStream
392 void BlockSampleStreaming(bool bBlock) noexcept override
393 {
395 }
396
397 bool IsSampleStreamingBlocked() const noexcept override
398 {
400 }
401
402public: // implements wolverine::ISampleStreamStatistics
412
414 void EnableSampleSizeStatistics(bool bEnable) noexcept override
415 {
417 }
418
420 void ResetStatistics() noexcept override
421 {
423 }
424};
425
430{
431public:
433 ADTF_CLASS_ID_NAME(cSampleStream, "sample_stream.streaming.adtf.cid", "Sample Stream");
434
435public:
437};
438
439} // namespace streaming
440} // namespace adtf
A_UTILS_NS::cResult tResult
For backwards compatibility and to bring latest version into scope.
Definition result.h:775
#define RETURN_ERROR(code)
Return specific error code, which requires the calling function's return type to be tResult.
Definition result.h:42
tMode
Definition samplestreamaccess_intf.h:26
Definition interface_binding_proxy_intf.h:22
Definition named_graph_object_intf.h:25
Definition sampleoutstream_intf.h:29
Definition samplestream_intf.h:30
Definition sample_intf.h:34
Defines access methods for the interface of a Stream Type.
Definition streamtype_intf.h:96
Definition triggerpipe_intf.h:103
tPriority
Definition triggerpipe_intf.h:39
Definition samplestream.h:430
ADTF_CLASS_ID_NAME(cSampleStream, "sample_stream.streaming.adtf.cid", "Sample Stream")
Implements adtf::ucom::IClassInfo.
tResult Open(ucom::ant::iobject_ptr< ISampleOutStream > &pOutStream, ISampleStreamAccess::tMode ui32Mode, size_t szQueueSize) noexcept
void EnableSampleSizeStatistics(bool bEnable) noexcept
tResult RegisterSubItem(const ucom::ant::iobject_ptr< ITriggerPipeItem > &pSubRun, ITriggerPipeItem::tPriority ui32Prio) noexcept
tResult UnregisterSubItem(const ucom::ant::iobject_ptr< ITriggerPipeItem > &pSubRun) noexcept
tResult AttachRouting(const ucom::ant::iobject_ptr< ISampleStream > &pSampleStreamTo) noexcept
tResult GetInternalBindingProxy(ucom::ant::iobject_ptr< IBindingProxy > &pBindingProxy) noexcept
cSampleStreamBase(cTriggerPipeItemBase &oPipeItem, ISampleStream &oSelf) noexcept
tResult SetStreamError(const tResult &oError) noexcept
tResult GetType(ucom::ant::iobject_ptr< const IStreamType > &pStreamType) const noexcept
void SetForwardAllTriggers(bool bForwardAll) noexcept
tResult Open(ucom::ant::iobject_ptr< ISampleInStream > &pInStream, ISampleStream::IPushReadEventSink *&pPushEventSink, ISampleStreamAccess::tMode ui32Mode, size_t szQueueSize) noexcept
tResult DetachRouting(const ucom::ant::iobject_ptr< ISampleStream > &pSampleStreamTo) noexcept
wolverine::ISampleStreamStatistics::tSampleStreamStatistic GetStatistics() noexcept
void SetSubStreamFilter(const char *strSubStream) noexcept
bool IsSampleStreamingBlocked() const noexcept
void BlockSampleStreaming(bool bBlock) noexcept
Definition streamingrequests_intf.h:98
Definition streamingrequests_intf.h:26
Definition streamingrequests_intf.h:45
sample_stream(const char *strName)
Definition samplestream.h:296
tResult ChangePriority(const ucom::ant::iobject_ptr< ITriggerPipeItem > &, ITriggerPipeItem::tPriority) noexcept override
Definition samplestream.h:379
tResult AttachRouting(const ucom::ant::iobject_ptr< ISampleStream > &pSampleStreamTo) noexcept override
Definition samplestream.h:322
tResult GetType(ucom::ant::iobject_ptr< const IStreamType > &pStreamType) const noexcept override
Definition samplestream.h:302
void BlockSampleStreaming(bool bBlock) noexcept override
Definition samplestream.h:392
tResult Open(const char *strName, ucom::ant::iobject_ptr< ISampleInStream > &pInStream, const ucom::ant::iobject_ptr< const IStreamType > &pInitialAcceptedStreamType, ISampleStream::IPushReadEventSink *&pPushEventSink, ISampleStreamAccess::tMode ui32Mode, size_t szQueueSize) noexcept override
Definition samplestream.h:337
void ResetStatistics() noexcept override
Definition samplestream.h:420
tResult DetachRouting(const ucom::ant::iobject_ptr< ISampleStream > &pSampleStreamTo) noexcept override
Definition samplestream.h:327
bool IsSampleStreamingBlocked() const noexcept override
Definition samplestream.h:397
tResult Open(const char *strName, ucom::ant::iobject_ptr< ISampleOutStream > &pOutStream, ISampleStreamAccess::tMode ui32Mode, size_t szQueueSize) noexcept override
Definition samplestream.h:354
tResult RegisterSubItem(const ucom::ant::iobject_ptr< ITriggerPipeItem > &pSubRun, ITriggerPipeItem::tPriority ui32Prio) noexcept override
Definition samplestream.h:363
tResult GetInternalBindingProxy(ucom::ant::iobject_ptr< IBindingProxy > &pBindingProxy) noexcept override
Definition samplestream.h:386
tResult SetStreamError(const tResult &oError) noexcept override
Definition samplestream.h:317
tTimeStamp GetTime() const noexcept override
Definition samplestream.h:312
void EnableSampleSizeStatistics(bool bEnable) noexcept override
Definition samplestream.h:414
tResult UnregisterSubItem(const ucom::ant::iobject_ptr< ITriggerPipeItem > &pSubRun) noexcept override
Definition samplestream.h:369
sample_stream()
Creates a sample stream without a name.
Definition samplestream.h:288
wolverine::ISampleStreamStatistics::tSampleStreamStatistic GetStatistics() noexcept override
Definition samplestream.h:408
Definition triggerpipe.h:265
Definition samplestream_intf.h:155
Samplestream statistics Interface. This interface describes statistical information of an ISampleStre...
Definition samplestream_statistics_intf.h:24
adtf::services::wolverine::ISampleStreamTracer::tTraceData tSampleStreamStatistic
structure containing the current samplestream statistics for samplestream input and output
Definition samplestream_statistics_intf.h:34
Safely retrieve a valid object_ptr<> instance to *this when all we have is *this.
Definition object_ptr_utilities.h:40
Base object pointer to realize binary compatible reference counting in interface methods.
Definition object_ptr_intf.h:112
Definition object.h:397
Namespace for all functionality of the ADTF Streaming SDK provided since v3.0.
Definition bindingproxyoutport.h:16
Namespace for all functionality of the ADTF Streaming SDK provided since v3.7.
Definition filtergraphport.h:580
Namespace for all functionality of the ADTF Streaming SDK provided since v3.21.
Definition request_forwarding.h:17
Namespace for all functionality of the ADTF Streaming SDK provided since v3.22.
Definition samplestream_statistics_intf.h:17
Namespace for the ADTF Streaming SDK.
Definition bindingproxyinport.h:14
Namespace for the ADTF uCOM SDK.
Definition adtf_system.h:324
Namespace for all functionality provided by ADTF and its SDKs.
Definition adtf_client_connector.h:14