DUNE-DAQ
DUNE Trigger and Data Acquisition software
Loading...
Searching...
No Matches
LocalStream.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/LocalStream.cxx to src/LocalStream.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 * LocalStream.cxx
14 * ERS
15 *
16 * Created by Serguei Kolos on 21.01.05.
17 * Copyright 2005 CERN. All rights reserved.
18 *
19 */
20#include <ers/LocalStream.hpp>
21#include <ers/StreamManager.hpp>
23
28ers::LocalStream&
29ers::LocalStream::instance()
30{
33 static ers::LocalStream* instance = ers::SingletonCreator<ers::LocalStream>::create();
34
35 return *instance;
36}
37
41ers::LocalStream::LocalStream()
42 : m_terminated(false)
43{
44}
45
46ers::LocalStream::~LocalStream()
47{
48 remove_issue_catcher();
49}
50
51void
52ers::LocalStream::remove_issue_catcher()
53{
54 std::unique_ptr<std::thread> catcher;
55 {
56 std::unique_lock lock(m_mutex);
57 if (!m_issue_catcher_thread.get()) {
58 return;
59 }
60 m_terminated = true;
61 m_condition.notify_one();
62 catcher.swap(m_issue_catcher_thread);
63 }
64
65 catcher->join();
66}
67
68void
69ers::LocalStream::thread_wrapper()
70{
71 std::unique_lock lock(m_mutex);
72 m_catcher_thread_id = std::this_thread::get_id();
73 while (!m_terminated) {
74 m_condition.wait(lock, [this]() { return !m_issues.empty() || m_terminated; });
75
76 while (!m_terminated && !m_issues.empty()) {
77 ers::Issue* issue = m_issues.front();
78 m_issues.pop();
79
80 lock.unlock();
81 m_issue_catcher(*issue);
82 delete issue;
83 lock.lock();
84 }
85 }
86 m_catcher_thread_id = {};
87 m_terminated = false;
88}
89
91ers::LocalStream::set_issue_catcher(const std::function<void(const ers::Issue&)>& catcher)
92{
93 std::unique_lock lock(m_mutex);
94 if (m_issue_catcher_thread.get()) {
95 throw ers::IssueCatcherAlreadySet(ERS_HERE);
96 }
97 m_issue_catcher = catcher;
98 m_issue_catcher_thread.reset(new std::thread(std::bind(&ers::LocalStream::thread_wrapper, this)));
99
100 return new ers::IssueCatcherHandler;
101}
102
103void
104ers::LocalStream::report_issue(ers::severity type, const ers::Issue& issue)
105{
106 if (m_issue_catcher_thread.get() && m_catcher_thread_id != std::this_thread::get_id()) {
107 ers::Issue* clone = issue.clone();
108 clone->set_severity(type);
109 std::unique_lock lock(m_mutex);
110 m_issues.push(clone);
111 m_condition.notify_one();
112 } else {
113 StreamManager::instance().report_issue(type, issue);
114 }
115}
116
117void
118ers::LocalStream::error(const ers::Issue& issue)
119{
120 report_issue(ers::Error, issue);
121}
122
123void
124ers::LocalStream::fatal(const ers::Issue& issue)
125{
126 report_issue(ers::Fatal, issue);
127}
128
129void
130ers::LocalStream::warning(const ers::Issue& issue)
131{
132 report_issue(ers::Warning, issue);
133}
#define ERS_HERE
Implements issue catcher lifetime management.
Base class for any user define issue.
Definition Issue.hpp:76
virtual Issue * clone() const =0
ers::Severity set_severity(ers::Severity severity) const
Definition Issue.cpp:185
severity
Definition Severity.hpp:37
@ Error
Definition Severity.hpp:42
@ Fatal
Definition Severity.hpp:43
@ Warning
Definition Severity.hpp:41