DUNE-DAQ
DUNE Trigger and Data Acquisition software
Loading...
Searching...
No Matches
StreamManager.cpp
Go to the documentation of this file.
1/*
2 * DUNE DAQ modification notice:
3 * This file has been modified from the original ATLAS ers source for the DUNE DAQ project.
4 * Fork baseline commit: 8267df82a4f6fe6bf02c4014923eba19eddc4614 (2020-04-14).
5 * Renamed since fork: yes (from src/StreamManager.cxx to src/StreamManager.cpp).
6 *
7 * Original copyright:
8 * Copyright (C) 2001-2020 CERN for the benefit of the ATLAS collaboration.
9 * Licensed under the Apache License, Version 2.0.
10 */
11
12/*
13 * StreamManager.cxx
14 * ERS
15 *
16 * Created by Matthias Wiesmann on 21.01.05.
17 * Modified by Serguei Kolos on 12.09.06.
18 * Copyright 2005 CERN. All rights reserved.
19 *
20 */
21
22#include <assert.h>
23#include <iostream>
24
25#include <ers/Configuration.hpp>
26#include <ers/InputStream.hpp>
27#include <ers/Issue.hpp>
28#include <ers/OutputStream.hpp>
29#include <ers/Severity.hpp>
30#include <ers/StreamFactory.hpp>
31#include <ers/StreamManager.hpp>
32#include <ers/ers.hpp>
36#include <ers/internal/Util.hpp>
38
40 BadConfiguration,
41 "The stream configuration string \"" << config << "\" has syntax errors.",
42 ((std::string)config))
43
44namespace {
48const char SEPARATOR = ',';
49
50const char* const DefaultOutputStreams[] = {
51 "lstdout", // Debug
52 "lstdout", // Log
53 "throttle,lstdout", // Information
54 "throttle,lstderr", // Warning
55 "throttle,lstderr", // Error
56 "lstderr" // Fatal
57};
58
59const char*
60get_stream_description(ers::severity severity)
61{
62 assert(ers::Debug <= severity && severity <= ers::Fatal);
63
64 std::string env_name("DUNEDAQ_ERS_");
65 env_name += ers::to_string(severity);
66 const char* env = ::getenv(env_name.c_str());
67 return env ? env : DefaultOutputStreams[severity];
68}
69
70void
71parse_stream_definition(const std::string& text, std::vector<std::string>& result)
72{
73 std::string::size_type start_p = 0, end_p = 0;
74 short brackets_open = 0;
75 while (end_p < text.length()) {
76 switch (text[end_p]) {
77 case '(':
78 ++brackets_open;
79 break;
80 case ')':
81 --brackets_open;
82 break;
83 case SEPARATOR:
84 if (!brackets_open) {
85 result.push_back(text.substr(start_p, end_p - start_p));
86 start_p = end_p + 1;
87 }
88 break;
89 default:
90 break;
91 }
92 end_p++;
93 }
94 if (brackets_open) {
95 throw ers::BadConfiguration(ERS_HERE, text);
96 }
97 if (start_p != end_p) {
98 result.push_back(text.substr(start_p, end_p - start_p));
99 }
100}
101}
102
103namespace ers {
104// Performs lazy srteam initialization. Stream instances are created
105// at the first attempt of writing to the stream
107{
108public:
110 : m_manager(manager)
111 , m_in_progress(false)
112 {
113 ;
114 }
115
116 void write(const Issue& issue)
117 {
118 ers::severity s = issue.severity();
119 std::scoped_lock lock(m_mutex);
120
121 if (!m_in_progress) {
122 m_in_progress = true;
123 } else {
124 // The issue is coming from the stream constructor
125 // We can't use ERS streams, so print it to std
126 if (s < ers::Warning)
127 std::cout << issue << std::endl;
128 else
129 std::cerr << issue << std::endl;
130 return;
131 }
132
133 if (m_manager.m_out_streams[s].get() == this) {
134 m_manager.m_out_streams[s] = std::shared_ptr<OutputStream>(m_manager.setup_stream(s));
135 }
136 m_manager.report_issue(s, issue);
137 m_in_progress = false;
138 }
139
140private:
141 std::recursive_mutex m_mutex;
144};
145
146}
147
161
166{
167 for (short ss = ers::Debug; ss <= ers::Fatal; ++ss) {
168 m_init_streams[ss] = std::make_shared<StreamInitializer>(*this);
170 }
171}
172
179
180void
182{
183 std::shared_ptr<OutputStream> head = m_out_streams[severity];
184 if (head && !head->isNull()) {
185 OutputStream* parent = head.get();
186 for (OutputStream* stream = parent; !stream->isNull(); parent = stream, stream = &parent->chained())
187 ;
188
189 parent->chained(new_stream);
190 } else {
191 m_out_streams[severity] = std::shared_ptr<OutputStream>(new_stream);
192 }
193}
194
195void
196ers::StreamManager::add_receiver(const std::string& stream, const std::string& filter, ers::IssueReceiver* receiver)
197{
198 InputStream* in = ers::StreamFactory::instance().create_in_stream(stream, filter);
199 in->set_receiver(receiver);
200
201 std::scoped_lock lock(m_mutex);
202 m_in_streams.push_back(std::shared_ptr<InputStream>(in));
203}
204
205void
206ers::StreamManager::add_receiver(const std::string& stream,
207 const std::initializer_list<std::string>& params,
208 ers::IssueReceiver* receiver)
209{
210 InputStream* in = ers::StreamFactory::instance().create_in_stream(stream, params);
211 in->set_receiver(receiver);
212
213 std::scoped_lock lock(m_mutex);
214 m_in_streams.push_back(std::shared_ptr<InputStream>(in));
215}
216
217void
219{
220 std::scoped_lock lock(m_mutex);
221 for (std::list<std::shared_ptr<InputStream>>::iterator it = m_in_streams.begin(); it != m_in_streams.end();) {
222 if ((*it)->m_receiver == receiver)
223 m_in_streams.erase(it++);
224 else
225 ++it;
226 }
227}
228
231{
232 std::string config = get_stream_description(severity);
233 std::vector<std::string> streams;
234 try {
235 parse_stream_definition(config, streams);
236 } catch (ers::BadConfiguration& ex) {
237 ERS_INTERNAL_ERROR("Configuration for the \"" << severity
238 << "\" stream is invalid. "
239 "Default configuration will be used.");
240 }
241
243
244 if (!main) {
245 std::vector<std::string> default_streams;
246 try {
247 parse_stream_definition(DefaultOutputStreams[severity], default_streams);
248 main = setup_stream(default_streams);
249 } catch (ers::BadConfiguration& ex) {
250 ERS_INTERNAL_ERROR("Can not configure the \"" << severity << "\" stream because of the following issue {" << ex
251 << "}");
252 }
253 }
254 return (main ? main : new ers::NullStream());
255}
256
258ers::StreamManager::setup_stream(const std::vector<std::string>& streams)
259{
260 size_t cnt = 0;
262 for (; cnt < streams.size(); ++cnt) {
263 main = ers::StreamFactory::instance().create_out_stream(streams[cnt]);
264 if (main)
265 break;
266 }
267
268 if (!main) {
269 return 0;
270 }
271
272 ers::OutputStream* head = main;
273 for (++cnt; cnt < streams.size(); ++cnt) {
274 ers::OutputStream* chained = ers::StreamFactory::instance().create_out_stream(streams[cnt]);
275
276 if (chained) {
277 head->chained(chained);
278 head = chained;
279 }
280 }
281
282 return main;
283}
284
289void
291{
292 ers::severity old_severity = issue.set_severity(type);
293 m_out_streams[type]->write(issue);
294 issue.set_severity(old_severity);
295} // error
296
300void
302{
303 report_issue(ers::Error, issue);
304} // error
305
310void
311ers::StreamManager::debug(const Issue& issue, int level)
312{
313 if (Configuration::instance().debug_level() >= level) {
314 ers::severity old_severity = issue.set_severity(ers::Severity(ers::Debug, level));
315 m_out_streams[ers::Debug]->write(issue);
316 issue.set_severity(old_severity);
317 }
318}
319
323void
325{
326 report_issue(ers::Fatal, issue);
327}
328
332void
334{
336}
337
341void
346
350void
352{
353 report_issue(ers::Log, issue);
354}
355
356std::ostream&
357ers::operator<<(std::ostream& out, const ers::StreamManager&)
358{
359 for (short ss = ers::Debug; ss <= ers::Fatal; ++ss) {
360 out << (ers::severity)ss << "\t\"" << get_stream_description((ers::severity)ss) << "\"" << std::endl;
361 }
362 return out;
363}
#define ERS_HERE
static Configuration & instance()
return the singleton
ERS Issue input stream interface.
void set_receiver(IssueReceiver *receiver)
ERS Issue receiver interface.
Base class for any user define issue.
Definition Issue.hpp:76
ers::Severity severity() const
severity of the issue
Definition Issue.hpp:110
ers::Severity set_severity(ers::Severity severity) const
Definition Issue.cpp:185
ERS abstract output stream interface.
virtual bool isNull() const
friend class StreamManager
OutputStream & chained()
void write(const Issue &issue)
std::recursive_mutex m_mutex
StreamManager & m_manager
StreamInitializer(StreamManager &manager)
This class manages and provides access to ERS streams.
void add_output_stream(ers::severity severity, ers::OutputStream *new_stream)
void error(const Issue &issue)
sends an issue to the error stream
std::shared_ptr< OutputStream > m_init_streams[ers::Fatal+1]
array of pointers to streams per severity
OutputStream * setup_stream(ers::severity severity)
void warning(const Issue &issue)
sends an issue to the warning stream
static StreamManager & instance()
return the singleton
std::list< std::shared_ptr< InputStream > > m_in_streams
void report_issue(ers::severity type, const Issue &issue)
void log(const Issue &issue)
sends an issue to the log stream
void fatal(const Issue &issue)
sends an issue to the fatal stream
void add_receiver(const std::string &stream, const std::string &filter, ers::IssueReceiver *receiver)
void debug(const Issue &issue, int level)
sends an Issue to the debug stream
void information(const Issue &issue)
sends an issue to the information stream
std::shared_ptr< OutputStream > m_out_streams[ers::Fatal+1]
array of pointers to streams per severity
void remove_receiver(ers::IssueReceiver *receiver)
int main(int argc, char *argv[])
#define ERS_INTERNAL_ERROR(message)
Definition macro.hpp:71
#define ERS_DECLARE_ISSUE(namespace_name, class_name, message_, attributes)
Definition macro.hpp:71
int debug_level()
Definition ers.hpp:80
std::string to_string(severity s)
std::ostream & operator<<(std::ostream &, const ers::Configuration &)
severity
Definition Severity.hpp:37
@ Debug
Definition Severity.hpp:38
@ Error
Definition Severity.hpp:42
@ Fatal
Definition Severity.hpp:43
@ Log
Definition Severity.hpp:39
@ Warning
Definition Severity.hpp:41
@ Information
Definition Severity.hpp:40
Null stream.