RTC Toolkit 6.0.0
Loading...
Searching...
No Matches
ddsPubThread.hpp
Go to the documentation of this file.
1
11#ifndef RTCTK_REUSABLECOMPONENT_TELREPUB_DDS_PUB_THREAD_HPP
12#define RTCTK_REUSABLECOMPONENT_TELREPUB_DDS_PUB_THREAD_HPP
13
14#include "agnostictopicif.hpp"
15#include "mudpiProcessor.hpp"
16#include "queue.hpp"
17#include <cstdint>
18#include <fastdds/dds/publisher/DataWriter.hpp>
19#include <fmt/core.h>
20#include <fmt/printf.h>
21#include <memory>
28#include <tbb/concurrent_queue.h>
29
30#include <numapp/thread.hpp>
31
32#include <chrono>
33#include <iostream>
34#include <utility>
35
36namespace rtctk::telRepub {
37
38using namespace rtctk::componentFramework;
39
44 numapp::NumaPolicies thread_policies;
45};
46
51 std::uint32_t queue_size = 200u;
52};
53
59class PubThread {
60public:
69 PubThread(DataWriter& data_writer,
71 const numapp::NumaPolicies& thread_policies,
72 std::uint32_t sample_id_increment);
73
77 virtual ~PubThread();
78
82 virtual void RunAsync();
83
87 void IdleAsync();
88
89 std::string GetTopicName();
90
91 virtual bool Monitor();
92
93protected:
94 enum class State { Idle, Run, Exit };
95
96 std::string m_topic_name;
97
98 log4cplus::Logger m_logger;
99
100 std::atomic<State> m_requested_state;
101 std::thread m_thread;
102 std::string m_thread_name;
103
104 DataWriter& m_data_writer;
105
107
112
114
119 // @note: OLDB does not support uint64_t which we really want, so we have to use signed 64.
120 perfc::CounterI64 m_pc_published_samples;
121 perfc::ScopedRegistration m_pc_published_samples_reg;
122
123 perfc::CounterI64 m_pc_publication_errors;
124 perfc::ScopedRegistration m_pc_publication_errors_reg;
125
127 perfc::ScopedRegistration m_pc_last_sample_id_published_reg;
128
129 std::unique_ptr<DurationMonitor<>> m_dur_mon;
130
132
133 static std::atomic<std::uint16_t> s_count; // instance index
134}; // PubThread
135
141class PubThreadMudpi : public PubThread {
142public:
151 PubThreadMudpi(DataWriter& data_writer,
154 WranglerFunction&& wrangler);
155
159 virtual ~PubThreadMudpi();
160
164 void RunAsync() override;
165
170
171 std::shared_ptr<Queue> GetQueue();
172
173 bool Monitor() override;
174
175private:
181 void ProcessingLoop();
182
183 std::shared_ptr<Queue> m_frame_queue;
184 llnetio::mudpi::TopicId m_mudpi_topic_id;
185 MudpiProcessor m_mudpi_processor;
186 WranglerFunction m_wrangler;
187};
188
194class PubThreadSim : public PubThread {
195public:
206 PubThreadSim(DataWriter& data_writer,
208 const numapp::NumaPolicies& thread_policies,
209 std::uint32_t sample_id_increment,
210 std::uint16_t sim_freq,
211 std::uint32_t sim_payload_size);
212
216 virtual ~PubThreadSim();
217
218private:
224 void SimProcessingLoop();
225
226 AgnosticTopic m_sim_msg;
227
228 std::uint16_t m_sim_freq;
229 std::uint32_t m_sim_payload_size; // bytes
230};
231
232} // namespace rtctk::telRepub
233
234#endif // RTCTK_REUSABLECOMPONENT_TELREPUB_DDS_PUB_THREAD_HPP
UDP Buffer management.
Models a single alert source that can be set or cleared.
Definition alertServiceIf.hpp:47
Component metrics interface.
Definition componentMetricsIf.hpp:163
Container class that holds services of any type.
Definition serviceContainer.hpp:38
Processing MUDPI data received by UDP receiver: rtctk::telRepub::UdpReceiver.
Definition mudpiProcessor.hpp:55
bool Monitor() override
Definition ddsPubThread.cpp:165
MudpiProcessor & GetMudpiProcessor()
Get MudpiProcessor.
Definition ddsPubThread.cpp:157
std::shared_ptr< Queue > GetQueue()
Definition ddsPubThread.cpp:161
virtual ~PubThreadMudpi()
Joins publisher thread and prints some statistics information.
Definition ddsPubThread.cpp:145
void RunAsync() override
Enable publishing data.
Definition ddsPubThread.cpp:152
PubThreadMudpi(DataWriter &data_writer, componentFramework::ServiceContainer &service, CfgPubThreadMudpi &cfg, WranglerFunction &&wrangler)
Spawns publisher thread and sets up performance counters for monitoring.
Definition ddsPubThread.cpp:119
virtual ~PubThreadSim()
Joins publisher thread and prints some statistics information.
Definition ddsPubThread.cpp:275
PubThreadSim(DataWriter &data_writer, componentFramework::ServiceContainer &service, const numapp::NumaPolicies &thread_policies, std::uint32_t sample_id_increment, std::uint16_t sim_freq, std::uint32_t sim_payload_size)
Spawns publisher thread and sets up performance counters for monitoring.
Definition ddsPubThread.cpp:241
DataWriter & m_data_writer
Definition ddsPubThread.hpp:104
std::string m_thread_name
Definition ddsPubThread.hpp:102
perfc::ScopedRegistration m_pc_last_sample_id_published_reg
Definition ddsPubThread.hpp:127
componentFramework::AlertSource m_wrangler_alert
Definition ddsPubThread.hpp:113
virtual bool Monitor()
Definition ddsPubThread.cpp:114
std::string m_topic_name
Definition ddsPubThread.hpp:96
void IdleAsync()
Disable publishing data.
Definition ddsPubThread.cpp:103
std::string GetTopicName()
Definition ddsPubThread.cpp:110
virtual void RunAsync()
Enable publishing data.
Definition ddsPubThread.cpp:95
PubThread(DataWriter &data_writer, componentFramework::ServiceContainer &service, const numapp::NumaPolicies &thread_policies, std::uint32_t sample_id_increment)
Spawns publisher thread and sets up performance counters for monitoring.
Definition ddsPubThread.cpp:20
componentFramework::AlertSource m_topic_publish_alert
Alerts.
Definition ddsPubThread.hpp:111
State
Definition ddsPubThread.hpp:94
@ Run
Definition ddsPubThread.hpp:94
@ Idle
Definition ddsPubThread.hpp:94
@ Exit
Definition ddsPubThread.hpp:94
perfc::CounterI64 m_pc_published_samples
Diverse Performance counters.
Definition ddsPubThread.hpp:120
perfc::CounterI64 m_pc_last_sample_id_published
Definition ddsPubThread.hpp:126
perfc::ScopedRegistration m_pc_publication_errors_reg
Definition ddsPubThread.hpp:124
static std::atomic< std::uint16_t > s_count
Definition ddsPubThread.hpp:133
ComponentMetricsIf & m_metrics
Definition ddsPubThread.hpp:106
std::thread m_thread
Definition ddsPubThread.hpp:101
virtual ~PubThread()
Joins publisher thread and prints some statistics information.
Definition ddsPubThread.cpp:75
log4cplus::Logger m_logger
Definition ddsPubThread.hpp:98
std::atomic< State > m_requested_state
Definition ddsPubThread.hpp:100
perfc::CounterI64 m_pc_publication_errors
Definition ddsPubThread.hpp:123
perfc::ScopedRegistration m_pc_published_samples_reg
Definition ddsPubThread.hpp:121
std::unique_ptr< DurationMonitor<> > m_dur_mon
Definition ddsPubThread.hpp:129
std::uint32_t m_expected_sample_id_increment
Definition ddsPubThread.hpp:131
Header file for ComponentMetricsIf.
Declares some common DDS functionality.
Header file for Duration Monitor.
Logging Support Library based on log4cplus.
MUDPI processor: check and aggregate MUDPI payload to a single topic and put to the queue for publish...
Definition ddsPubThread.cpp:16
std::function< std::error_code( const gsl::span< const gsl::span< const uint8_t > >, std::vector< uint8_t > &)> WranglerFunction
The wrangler function that is called with a span of spans containing the payload and an vector where ...
Definition wrangler.hpp:27
UDP Buffer management.
Wrangler: User extension point.
Structure to hold MudpiProcessor's configuration.
Definition mudpiProcessor.hpp:39
Structure to hold PubThreadMudpi's configuration.
Definition ddsPubThread.hpp:50
std::uint32_t queue_size
Definition ddsPubThread.hpp:51
Structure to hold PubThread's configuration.
Definition ddsPubThread.hpp:43
numapp::NumaPolicies thread_policies
Definition ddsPubThread.hpp:44