12#ifndef TIMINGLIBS_SRC_INFOGATHERER_HPP_
13#define TIMINGLIBS_SRC_INFOGATHERER_HPP_
24#include "nlohmann/json.hpp"
30#include <shared_mutex>
40 "Gather Threading Issue detected: " << err,
45 " Failed to send send " << device <<
" device info to " << destination <<
".",
46 ((std::string)device)((std::string)destination))
61 explicit InfoGatherer(std::function<
void(InfoGatherer&)> gather_data,
63 const std::string& device_name,
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)
75 , m_failed_to_send_counter(0)
82 virtual ~InfoGatherer()
85 stop_gathering_thread();
88 InfoGatherer(
const InfoGatherer&) =
delete;
89 InfoGatherer& operator=(
const InfoGatherer&) =
delete;
90 InfoGatherer(InfoGatherer&&) =
delete;
91 InfoGatherer& operator=(InfoGatherer&&) =
delete;
97 void start_gathering_thread(
const std::string& name =
"noname")
99 if (run_gathering()) {
101 "Attempted to start gathering thread "
102 "when it is already supposed to be running!"));
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());
110 std::ostringstream s;
111 s <<
"The name " << name <<
" provided for the thread is too long.";
122 void stop_gathering_thread()
124 if (!run_gathering()) {
126 "Attempted to stop gathering thread "
127 "when it is not supposed to be running!"));
130 m_run_gathering =
false;
131 if (m_gathering_thread->joinable()) {
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());
138 throw GatherThreadingIssue(
ERS_HERE,
"Thread not in joinable state during working thread stop!");
146 bool run_gathering()
const {
return m_run_gathering.load(); }
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(); }
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(); }
154 std::string get_device_name()
const {
return m_device_name; }
156 int get_op_mon_level()
const {
return m_op_mon_level; }
159 void collect_info_from_device(
const DSGN& device)
161 std::unique_lock info_collector_lock(m_info_collector_mutex);
163 device.get_info(*m_device_info);
164 update_last_gathered_time(std::time(
nullptr));
182 void send_device_info()
190 if (!m_hw_info_sender) {
191 TLOG_DEBUG(3) <<
"skipping sending info for gatherer: " << get_device_name();
196 to_json(info, *m_device_info);
197 bool was_successfully_sent =
false;
198 while (!was_successfully_sent) {
200 m_hw_info_sender->send(std::move(info), m_queue_timeout);
201 TLOG_DEBUG(4) <<
"sent " << get_device_name() <<
" info";
203 was_successfully_sent =
true;
204 }
catch (
const dunedaq::iomanager::TimeoutExpired& excpt) {
206 ++m_failed_to_send_counter;
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;
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;
static std::shared_ptr< IOManager > get()
#define TLOG_DEBUG(lvl,...)
ERS_DECLARE_ISSUE(cibmodules, CIBCommunicationError, " CIB Hardware Communication Error: "<< descriptor,((std::string) descriptor)) ERS_DECLARE_ISSUE(cibmodules
void warning(const Issue &issue)
void error(const Issue &issue)