DUNE-DAQ
DUNE Trigger and Data Acquisition software
Toggle main menu visibility
Loading...
Searching...
No Matches
dunedaq
sourcecode
oks
src
pipeline.cpp
Go to the documentation of this file.
1
// DUNE DAQ modification notice:
2
// This file has been modified from the original ATLAS oks source for the DUNE DAQ project.
3
// Fork baseline commit: oks-08-03-04 (2022-04-14).
4
// Renamed since fork: no.
5
7
//
8
// copy from IPC
9
//
11
12
#include "
oks/pipeline.hpp
"
13
#include <boost/bind/bind.hpp>
14
15
namespace
dunedaq
{
16
namespace
oks
{
17
18
bool
19
OksPipeline::Worker::setJob
(
OksJob
* job)
20
{
21
boost::mutex::scoped_lock lock(
m_mutex
);
22
if
(!
m_running
&&
m_idle
) {
23
m_job
= job;
24
m_idle
=
false
;
25
m_condition
.notify_one();
26
return
true
;
27
}
28
return
false
;
29
}
30
31
void
32
OksPipeline::Worker::run
()
33
{
34
{
35
boost::mutex::scoped_lock lock(
m_mutex
);
36
m_pipeline
.m_barrier.wait();
37
}
38
while
(
true
) {
39
{
40
boost::mutex::scoped_lock lock(
m_mutex
);
41
m_running
=
true
;
42
}
43
while
(!
m_idle
) {
44
m_job
->run();
45
delete
m_job
;
46
m_idle
= !
m_pipeline
.getJob(
m_job
);
47
}
48
{
49
boost::mutex::scoped_lock lock(
m_mutex
);
50
m_running
=
false
;
51
if
(
m_shutdown
)
52
break
;
53
if
(
m_stop
) {
54
m_pipeline
.m_barrier.wait();
55
m_stop
=
false
;
56
}
57
if
(
m_idle
)
58
m_condition
.wait(lock);
59
}
60
}
61
}
62
63
void
64
OksPipeline::Worker::shutdown
()
65
{
66
boost::mutex::scoped_lock lock(
m_mutex
);
67
m_shutdown
=
true
;
68
m_condition
.notify_one();
69
}
70
71
void
72
OksPipeline::Worker::stop
()
73
{
74
boost::mutex::scoped_lock lock(
m_mutex
);
75
m_stop
=
true
;
76
m_condition
.notify_one();
77
}
78
79
OksPipeline::OksPipeline
(
size_t
size
)
80
:
m_barrier
(
size
+ 1)
81
{
82
for
(
size_t
i = 0; i <
size
; ++i) {
83
m_workers
.push_back(
WorkerPtr
(
new
Worker
(*
this
)));
84
m_pool
.create_thread(boost::bind(&
Worker::run
,
m_workers
.back().get()));
85
}
86
m_barrier
.wait();
87
}
88
89
OksPipeline::~OksPipeline
()
90
{
91
waitForCompletion
();
92
for
(
size_t
i = 0; i <
m_workers
.size(); ++i) {
93
m_workers
[i]->shutdown();
94
}
95
m_pool
.join_all();
96
}
97
98
void
99
OksPipeline::addJob
(
OksJob
* job)
100
{
101
for
(
size_t
i = 0; i <
m_workers
.size(); ++i) {
102
if
(
m_workers
[i]->setJob(job)) {
103
return
;
104
}
105
}
106
boost::mutex::scoped_lock lock(
m_mutex
);
107
m_jobs
.push(job);
108
}
109
110
void
111
OksPipeline::waitForCompletion
()
112
{
113
{
114
boost::mutex::scoped_lock lock(
m_mutex
);
115
while
(!
m_jobs
.empty()) {
116
m_condition
.wait(lock);
117
}
118
}
119
for
(
size_t
i = 0; i <
m_workers
.size(); ++i) {
120
m_workers
[i]->stop();
121
}
122
m_barrier
.wait();
123
}
124
125
bool
126
OksPipeline::getJob
(
OksJob
*& job)
127
{
128
boost::mutex::scoped_lock lock(
m_mutex
);
129
if
(!
m_jobs
.empty()) {
130
job =
m_jobs
.front();
131
m_jobs
.pop();
132
m_condition
.notify_one();
133
return
true
;
134
}
135
return
false
;
136
}
137
138
}
// namespace oks
139
}
// namespace dunedaq
dunedaq::oks::OksJob
Definition
pipeline.hpp:25
dunedaq::oks::OksPipeline::m_condition
boost::condition m_condition
Definition
pipeline.hpp:83
dunedaq::oks::OksPipeline::waitForCompletion
void waitForCompletion()
Definition
pipeline.cpp:111
dunedaq::oks::OksPipeline::m_jobs
std::queue< OksJob * > m_jobs
Definition
pipeline.hpp:87
dunedaq::oks::OksPipeline::m_workers
std::vector< WorkerPtr > m_workers
Definition
pipeline.hpp:86
dunedaq::oks::OksPipeline::m_mutex
boost::mutex m_mutex
Definition
pipeline.hpp:82
dunedaq::oks::OksPipeline::getJob
bool getJob(OksJob *&job)
Definition
pipeline.cpp:126
dunedaq::oks::OksPipeline::WorkerPtr
std::shared_ptr< Worker > WorkerPtr
Definition
pipeline.hpp:74
dunedaq::oks::OksPipeline::m_pool
boost::thread_group m_pool
Definition
pipeline.hpp:85
dunedaq::oks::OksPipeline::~OksPipeline
~OksPipeline()
Definition
pipeline.cpp:89
dunedaq::oks::OksPipeline::addJob
void addJob(OksJob *job)
Definition
pipeline.cpp:99
dunedaq::oks::OksPipeline::OksPipeline
OksPipeline(size_t size)
Definition
pipeline.cpp:79
dunedaq::oks::OksPipeline::m_barrier
boost::barrier m_barrier
Definition
pipeline.hpp:84
dunedaq::oks
Definition
SchemaCommand.hpp:17
dunedaq
The DUNE-DAQ namespace.
Definition
cib_utilities.cpp:5
dunedaq::size
FELIX Initialization std::string initerror FELIX queue timed std::string queuename Unexpected chunk size
Definition
FelixIssues.hpp:30
pipeline.hpp
dunedaq::oks::OksPipeline::Worker::m_stop
bool m_stop
Definition
pipeline.hpp:69
dunedaq::oks::OksPipeline::Worker::shutdown
void shutdown()
Definition
pipeline.cpp:64
dunedaq::oks::OksPipeline::Worker::stop
void stop()
Definition
pipeline.cpp:72
dunedaq::oks::OksPipeline::Worker::m_running
bool m_running
Definition
pipeline.hpp:70
dunedaq::oks::OksPipeline::Worker::m_mutex
boost::mutex m_mutex
Definition
pipeline.hpp:65
dunedaq::oks::OksPipeline::Worker::setJob
bool setJob(OksJob *job)
Definition
pipeline.cpp:19
dunedaq::oks::OksPipeline::Worker::Worker
Worker(OksPipeline &pipeline)
Definition
pipeline.hpp:45
dunedaq::oks::OksPipeline::Worker::m_condition
boost::condition m_condition
Definition
pipeline.hpp:66
dunedaq::oks::OksPipeline::Worker::m_idle
bool m_idle
Definition
pipeline.hpp:67
dunedaq::oks::OksPipeline::Worker::m_pipeline
OksPipeline & m_pipeline
Definition
pipeline.hpp:64
dunedaq::oks::OksPipeline::Worker::m_job
OksJob * m_job
Definition
pipeline.hpp:71
dunedaq::oks::OksPipeline::Worker::m_shutdown
bool m_shutdown
Definition
pipeline.hpp:68
dunedaq::oks::OksPipeline::Worker::run
void run()
Definition
pipeline.cpp:32
Generated on
for DUNE-DAQ by
1.18.0