DUNE-DAQ
DUNE Trigger and Data Acquisition software
Loading...
Searching...
No Matches
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
15namespace dunedaq {
16namespace oks {
17
18bool
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
31void
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
63void
65{
66 boost::mutex::scoped_lock lock(m_mutex);
67 m_shutdown = true;
68 m_condition.notify_one();
69}
70
71void
73{
74 boost::mutex::scoped_lock lock(m_mutex);
75 m_stop = true;
76 m_condition.notify_one();
77}
78
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
90{
92 for (size_t i = 0; i < m_workers.size(); ++i) {
93 m_workers[i]->shutdown();
94 }
95 m_pool.join_all();
96}
97
98void
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
110void
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
125bool
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
boost::condition m_condition
Definition pipeline.hpp:83
std::queue< OksJob * > m_jobs
Definition pipeline.hpp:87
std::vector< WorkerPtr > m_workers
Definition pipeline.hpp:86
bool getJob(OksJob *&job)
Definition pipeline.cpp:126
std::shared_ptr< Worker > WorkerPtr
Definition pipeline.hpp:74
boost::thread_group m_pool
Definition pipeline.hpp:85
void addJob(OksJob *job)
Definition pipeline.cpp:99
boost::barrier m_barrier
Definition pipeline.hpp:84
The DUNE-DAQ namespace.
FELIX Initialization std::string initerror FELIX queue timed std::string queuename Unexpected chunk size
Worker(OksPipeline &pipeline)
Definition pipeline.hpp:45