DUNE-DAQ
DUNE Trigger and Data Acquisition software
Loading...
Searching...
No Matches
MonitorableObject.cpp
Go to the documentation of this file.
1
8
9#include "logging/Logging.hpp"
10#include <NullOpMonFacility.hpp>
12#include <opmonlib/Utils.hpp>
13
14#include <google/protobuf/util/time_util.h>
15
16#include <chrono>
17
21#define TRACE_NAME "MonitorableObject" // NOLINT
22enum
23{
26};
27
28using namespace dunedaq::opmonlib;
29
30std::shared_ptr<OpMonFacility> MonitorableObject::s_default_facility = std::make_shared<NullOpMonFacility>();
31
32void
34{
35
36 std::lock_guard<std::mutex> lock(m_node_mutex);
37
38 // check if the name is already present to ensure uniqueness
39 auto it = m_nodes.find(name);
40 if (it != m_nodes.end()) {
41 // This not desired because names are suppposed to be unique
42 // But if the pointer is expired, there is no harm in override it
43 if (it->second.expired()) {
44 ers::warning(NonUniqueNodeName(ERS_HERE, name, to_string(get_opmon_id())));
45 } else {
46 throw NonUniqueNodeName(ERS_HERE, name, to_string(get_opmon_id()));
47 }
48 }
49
50 m_nodes[name] = p;
51
52 p->m_opmon_name = name;
53 p->inherit_parent_properties(*this);
54
55 TLOG() << "Node " << name << " registered to " << to_string(get_opmon_id());
56}
57
58void
59MonitorableObject::publish(google::protobuf::Message&& m, CustomOrigin&& co, OpMonLevel l) const noexcept
60{
61
62 auto timestamp = google::protobuf::util::TimeUtil::GetCurrentTime();
63
64 auto start_time = std::chrono::high_resolution_clock::now();
65
67 TLOG_DEBUG(TLVL_LEVEL_SUPPRESSION) << "Metric " << m.GetTypeName() << " ignored because of the level";
69 return;
70 }
71
72 auto e = to_entry(m, co);
73
74 if (e.data().empty()) {
75 ers::warning(EntryWithNoData(ERS_HERE, e.measurement()));
76 return;
77 }
78
79 *e.mutable_origin() = get_opmon_id();
80
81 *e.mutable_time() = timestamp;
82
83 // this pointer is always garanteed to be filled, even if with a null Facility.
84 // But the facility can fail
85 try {
86 m_facility.load()->publish(std::move(e));
88 } catch (const OpMonPublishFailure& e) {
89 ers::error(e);
91 }
92
93 auto stop_time = std::chrono::high_resolution_clock::now();
94
95 auto duration = std::chrono::duration_cast<std::chrono::microseconds>(stop_time - start_time);
96 m_cpu_us_counter += duration.count();
97}
98
101{
102
103 auto start_time = std::chrono::high_resolution_clock::now();
104
105 TLOG_DEBUG(TLVL_MONITORING_STEPS) << "Collecting data from " << to_string(get_opmon_id());
108
109 info.set_n_invalid_links(0);
110
111 try {
113 } catch (const ers::Issue& i) {
115 auto cause_ptr = i.cause();
116 while (cause_ptr) {
118 cause_ptr = cause_ptr->cause();
119 }
120 ers::error(ErrorWhileCollecting(ERS_HERE, to_string(get_opmon_id()), i));
121 } catch (const std::exception& e) {
123 ers::error(ErrorWhileCollecting(ERS_HERE, to_string(get_opmon_id()), e));
124 } catch (...) {
126 ers::error(ErrorWhileCollecting(ERS_HERE, to_string(get_opmon_id())));
127 }
128
129 info.set_n_published_measurements(m_published_counter.exchange(0));
130 info.set_n_ignored_measurements(m_ignored_counter.exchange(0));
131 info.set_n_errors(m_error_counter.exchange(0));
132 if (info.n_published_measurements() > 0) {
133 info.set_n_publishing_nodes(1);
134 }
135 info.set_cpu_elapsed_time_us(m_cpu_us_counter.exchange(0));
136
137 std::lock_guard<std::mutex> lock(m_node_mutex);
138
139 info.set_n_registered_nodes(m_nodes.size());
140
141 unsigned int n_invalid_links = 0;
142
143 for (auto it = m_nodes.begin(); it != m_nodes.end();) {
144
145 auto ptr = it->second.lock();
146
147 if (ptr) {
148 auto child_info = ptr->collect(); // MR: can we make this an async? There is no point to wait all done here
149 info.set_n_registered_nodes(info.n_registered_nodes() + child_info.n_registered_nodes());
150 info.set_n_publishing_nodes(info.n_publishing_nodes() + child_info.n_publishing_nodes());
151 info.set_n_invalid_links(info.n_invalid_links() + child_info.n_invalid_links());
152 info.set_n_published_measurements(info.n_published_measurements() + child_info.n_published_measurements());
153 info.set_n_ignored_measurements(info.n_ignored_measurements() + child_info.n_ignored_measurements());
154 info.set_n_errors(info.n_errors() + child_info.n_errors());
155 info.set_cpu_elapsed_time_us(info.cpu_elapsed_time_us() + child_info.cpu_elapsed_time_us());
156 }
157
158 // prune the dead links
159 if (it->second.expired()) {
160 it = m_nodes.erase(it);
161 ++n_invalid_links;
162 } else {
163 ++it;
164 }
165 }
166
167 info.set_n_invalid_links(info.n_invalid_links() + n_invalid_links);
168
169 auto stop_time = std::chrono::high_resolution_clock::now();
170
171 auto duration = std::chrono::duration_cast<std::chrono::microseconds>(stop_time - start_time);
172 info.set_clockwall_elapsed_time_us(duration.count());
173
174 return info;
175}
176
177void
179{
180
181 m_opmon_level = l;
182
183 std::lock_guard<std::mutex> lock(m_node_mutex);
184 for (const auto& [key, wp] : m_nodes) {
185 auto p = wp.lock();
186 if (p) {
187 p->set_opmon_level(l);
188 }
189 }
190}
191
192void
194{
195
196 m_facility.store(parent.m_facility);
197 m_parent_id = parent.get_opmon_id();
199
200 std::lock_guard<std::mutex> lock(m_node_mutex);
201
202 for (const auto& [key, wp] : m_nodes) {
203
204 auto p = wp.lock();
205 if (p) {
206 p->inherit_parent_properties(*this);
207 }
208 }
209}
#define ERS_HERE
@ TLVL_MONITORING_STEPS
@ TLVL_LEVEL_SUPPRESSION
std::shared_ptr< MonitorableObject > NewNodePtr
void inherit_parent_properties(const MonitorableObject &parent)
std::atomic< metric_counter_t > m_published_counter
std::atomic< metric_counter_t > m_ignored_counter
opmon::MonitoringTreeInfo collect() noexcept
std::atomic< metric_counter_t > m_error_counter
void set_opmon_level(OpMonLevel) noexcept
std::atomic< facility_ptr_t > m_facility
void register_node(ElementId name, NewNodePtr)
void publish(google::protobuf::Message &&, CustomOrigin &&co={}, OpMonLevel l=to_level(EntryOpMonLevel::kDefault)) const noexcept
static bool publishable_metric(OpMonLevel entry, OpMonLevel system) noexcept
std::atomic< time_counter_t > m_cpu_us_counter
std::map< ElementId, NodePtr > m_nodes
MonitorableObject(const MonitorableObject &)=delete
Base class for any user define issue.
Definition Issue.hpp:76
const Issue * cause() const
return the cause Issue of this Issue
Definition Issue.hpp:100
#define TLOG_DEBUG(lvl,...)
Definition Logging.hpp:116
#define TLOG(...)
Definition macro.hpp:21
std::invoke_result< decltype(&dunedaq::confmodel::OpMonConf::get_level), dunedaq::confmodel::OpMonConf >::type OpMonLevel
dunedaq::opmon::OpMonEntry to_entry(const google::protobuf::Message &m, const CustomOrigin &co)
Definition Utils.cpp:21
std::string to_string(const dunedaq::opmon::OpMonId &)
Definition Utils.cpp:166
std::map< std::string, std::string > CustomOrigin
Definition Utils.hpp:46
Cannot add TPSet with start_time
void warning(const Issue &issue)
Definition ers.hpp:150
void error(const Issue &issue)
Definition ers.hpp:101