DUNE-DAQ
DUNE Trigger and Data Acquisition software
Loading...
Searching...
No Matches
InfoGatherer.hpp
Go to the documentation of this file.
1
11
12#ifndef TIMINGLIBS_SRC_INFOGATHERER_HPP_
13#define TIMINGLIBS_SRC_INFOGATHERER_HPP_
14
17
18#include "ers/Issue.hpp"
19#include "logging/Logging.hpp" // NOTE: if ISSUES ARE DECLARED BEFORE include logging/Logging.hpp, TLOG_DEBUG<<issue wont work.
20
22#include "iomanager/Sender.hpp"
23
24#include "nlohmann/json.hpp"
25
26#include <functional>
27#include <future>
28#include <list>
29#include <memory>
30#include <shared_mutex>
31#include <string>
32
33namespace dunedaq {
34
39 GatherThreadingIssue, // Issue Class Name
40 "Gather Threading Issue detected: " << err, // Message
41 ((std::string)err)) // Message parameters
42
45 " Failed to send send " << device << " device info to " << destination << ".",
46 ((std::string)device)((std::string)destination))
47namespace timinglibs {
48
53class InfoGatherer
54{
55public:
61 explicit InfoGatherer(std::function<void(InfoGatherer&)> gather_data,
62 uint gather_interval,
63 const std::string& device_name,
64 int op_mon_level)
65 : m_run_gathering(false)
66 , m_gathering_thread(nullptr)
67 , m_gather_interval(gather_interval)
68 , m_device_name(device_name)
69 , m_last_gathered_time(0)
70 , m_op_mon_level(op_mon_level)
71 , m_gather_data(gather_data)
72 , m_device_info_connection_id(device_name + "_info")
73 , m_hw_info_sender(nullptr)
74 , m_sent_counter(0)
75 , m_failed_to_send_counter(0)
76 , m_queue_timeout(1)
77 {
78 // m_info_collector = std::make_unique<opmonlib::InfoCollector>();
79 m_hw_info_sender = iomanager::IOManager::get()->get_sender<nlohmann::json>(m_device_info_connection_id);
80 }
81
82 virtual ~InfoGatherer()
83 {
84 if (run_gathering())
85 stop_gathering_thread();
86 }
87
88 InfoGatherer(const InfoGatherer&) = delete;
89 InfoGatherer& operator=(const InfoGatherer&) = delete;
90 InfoGatherer(InfoGatherer&&) = delete;
91 InfoGatherer& operator=(InfoGatherer&&) = delete;
92
97 void start_gathering_thread(const std::string& name = "noname")
98 {
99 if (run_gathering()) {
100 ers::warning(GatherThreadingIssue(ERS_HERE,
101 "Attempted to start gathering thread "
102 "when it is already supposed to be running!"));
103 return;
104 }
105 m_run_gathering = true;
106 m_gathering_thread.reset(new std::thread([&] { m_gather_data(*this); }));
107 auto handle = m_gathering_thread->native_handle();
108 auto rc = pthread_setname_np(handle, name.c_str());
109 if (rc != 0) {
110 std::ostringstream s;
111 s << "The name " << name << " provided for the thread is too long.";
112 ers::warning(GatherThreadingIssue(ERS_HERE, s.str()));
113 }
114 }
115
122 void stop_gathering_thread()
123 {
124 if (!run_gathering()) {
125 ers::warning(GatherThreadingIssue(ERS_HERE,
126 "Attempted to stop gathering thread "
127 "when it is not supposed to be running!"));
128 return;
129 }
130 m_run_gathering = false;
131 if (m_gathering_thread->joinable()) {
132 try {
133 m_gathering_thread->join();
134 } catch (std::system_error const& e) {
135 throw GatherThreadingIssue(ERS_HERE, std::string("Error while joining gathering thread, ") + e.what());
136 }
137 } else {
138 throw GatherThreadingIssue(ERS_HERE, "Thread not in joinable state during working thread stop!");
139 }
140 }
141
146 bool run_gathering() const { return m_run_gathering.load(); }
147
148 void update_gather_interval(uint new_gather_interval) { m_gather_interval.store(new_gather_interval); }
149 uint get_gather_interval() const { return m_gather_interval.load(); }
150
151 void update_last_gathered_time(int64_t last_time) { m_last_gathered_time.store(last_time); }
152 time_t get_last_gathered_time() const { return m_last_gathered_time.load(); }
153
154 std::string get_device_name() const { return m_device_name; }
155
156 int get_op_mon_level() const { return m_op_mon_level; }
157
158 template<class DSGN>
159 void collect_info_from_device(const DSGN& device)
160 {
161 std::unique_lock info_collector_lock(m_info_collector_mutex);
162 m_device_info.reset(new timing::timingfirmwareinfo::TimingDeviceInfo());
163 device.get_info(*m_device_info);
164 update_last_gathered_time(std::time(nullptr));
165 send_device_info();
166 }
167
168 // void add_info_to_collector(std::string label, opmonlib::InfoCollector& ic)
169 // {
170 // std::unique_lock info_collector_lock(m_info_collector_mutex);
171 // if (m_info_collector->is_empty()) {
172 // TLOG_DEBUG(3) << "skipping add info for gatherer: " << get_device_name()
173 // << " with gathered time: " << get_last_gathered_time() << " and level " << get_op_mon_level();
174 // } else {
175 // ic.add(label, *m_info_collector);
176 // }
177 // m_info_collector = std::make_unique<opmonlib::InfoCollector>();
178 // update_last_gathered_time(0);
179 // }
180
181private:
182 void send_device_info()
183 {
184 // if (m_info_collector->is_empty())
185 //{
186 // TLOG_DEBUG(3) << "skipping sending info for gatherer: " << get_device_name() << ", collector empty.";
187 // return;
188 // }
189
190 if (!m_hw_info_sender) {
191 TLOG_DEBUG(3) << "skipping sending info for gatherer: " << get_device_name();
192 return;
193 }
194
195 nlohmann::json info;
196 to_json(info, *m_device_info);
197 bool was_successfully_sent = false;
198 while (!was_successfully_sent) {
199 try {
200 m_hw_info_sender->send(std::move(info), m_queue_timeout);
201 TLOG_DEBUG(4) << "sent " << get_device_name() << " info";
202 ++m_sent_counter;
203 was_successfully_sent = true;
204 } catch (const dunedaq::iomanager::TimeoutExpired& excpt) {
205 ers::error(DeviceInfoSendFailed(ERS_HERE, m_device_name, m_device_info_connection_id));
206 ++m_failed_to_send_counter;
207 }
208 }
209 }
210
211protected:
212 std::atomic<bool> m_run_gathering;
213 std::unique_ptr<std::thread> m_gathering_thread;
214 std::atomic<uint> m_gather_interval;
215 mutable std::shared_mutex m_mon_data_mutex;
216 std::string m_device_name;
217 std::atomic<time_t> m_last_gathered_time;
218 int m_op_mon_level;
219 // std::unique_ptr<opmonlib::InfoCollector> m_info_collector;
220 std::unique_ptr<timing::timingfirmwareinfo::TimingDeviceInfo> m_device_info;
221 mutable std::mutex m_info_collector_mutex;
222 std::function<void(InfoGatherer&)> m_gather_data;
223 std::string m_device_info_connection_id;
225 std::shared_ptr<sink_t> m_hw_info_sender;
226 std::atomic<uint> m_sent_counter;
227 std::atomic<uint> m_failed_to_send_counter;
228 std::chrono::milliseconds m_queue_timeout;
229};
230
231} // namespace timinglibs
232} // namespace dunedaq
233
234#endif // TIMINGLIBS_SRC_INFOGATHERER_HPP_
235
236// Local Variables:
237// c-basic-offset: 2
238// End:
#define ERS_HERE
static std::shared_ptr< IOManager > get()
Definition IOManager.hpp:40
#define TLOG_DEBUG(lvl,...)
Definition Logging.hpp:116
The DUNE-DAQ namespace.
ERS_DECLARE_ISSUE(cibmodules, CIBCommunicationError, " CIB Hardware Communication Error: "<< descriptor,((std::string) descriptor)) ERS_DECLARE_ISSUE(cibmodules
void warning(const Issue &issue)
Definition ers.hpp:150
void error(const Issue &issue)
Definition ers.hpp:101