DUNE-DAQ
DUNE Trigger and Data Acquisition software
Loading...
Searching...
No Matches
dunedaq::oks::OksPipeline Class Reference

#include <pipeline.hpp>

Classes

struct  Worker

Public Member Functions

 OksPipeline (size_t size)
 ~OksPipeline ()
void waitForCompletion ()
void addJob (OksJob *job)

Private Types

typedef std::shared_ptr< Worker > WorkerPtr

Private Member Functions

bool getJob (OksJob *&job)

Private Attributes

boost::mutex m_mutex
boost::condition m_condition
boost::barrier m_barrier
boost::thread_group m_pool
std::vector< WorkerPtr > m_workers
std::queue< OksJob * > m_jobs

Friends

struct Worker

Detailed Description

Definition at line 31 of file pipeline.hpp.

Member Typedef Documentation

◆ WorkerPtr

typedef std::shared_ptr<Worker> dunedaq::oks::OksPipeline::WorkerPtr
private

Definition at line 74 of file pipeline.hpp.

Constructor & Destructor Documentation

◆ OksPipeline()

dunedaq::oks::OksPipeline::OksPipeline ( size_t size)

Definition at line 79 of file pipeline.cpp.

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}
std::vector< WorkerPtr > m_workers
Definition pipeline.hpp:86
std::shared_ptr< Worker > WorkerPtr
Definition pipeline.hpp:74
boost::thread_group m_pool
Definition pipeline.hpp:85
boost::barrier m_barrier
Definition pipeline.hpp:84
FELIX Initialization std::string initerror FELIX queue timed std::string queuename Unexpected chunk size
Worker(OksPipeline &pipeline)
Definition pipeline.hpp:45

◆ ~OksPipeline()

dunedaq::oks::OksPipeline::~OksPipeline ( )

Definition at line 89 of file pipeline.cpp.

90{
92 for (size_t i = 0; i < m_workers.size(); ++i) {
93 m_workers[i]->shutdown();
94 }
95 m_pool.join_all();
96}

Member Function Documentation

◆ addJob()

void dunedaq::oks::OksPipeline::addJob ( OksJob * job)

Definition at line 99 of file pipeline.cpp.

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}
std::queue< OksJob * > m_jobs
Definition pipeline.hpp:87

◆ getJob()

bool dunedaq::oks::OksPipeline::getJob ( OksJob *& job)
private

Definition at line 126 of file pipeline.cpp.

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}
boost::condition m_condition
Definition pipeline.hpp:83

◆ waitForCompletion()

void dunedaq::oks::OksPipeline::waitForCompletion ( )

Definition at line 111 of file pipeline.cpp.

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}

◆ Worker

friend struct Worker
friend

Definition at line 77 of file pipeline.hpp.

Member Data Documentation

◆ m_barrier

boost::barrier dunedaq::oks::OksPipeline::m_barrier
private

Definition at line 84 of file pipeline.hpp.

◆ m_condition

boost::condition dunedaq::oks::OksPipeline::m_condition
private

Definition at line 83 of file pipeline.hpp.

◆ m_jobs

std::queue<OksJob*> dunedaq::oks::OksPipeline::m_jobs
private

Definition at line 87 of file pipeline.hpp.

◆ m_mutex

boost::mutex dunedaq::oks::OksPipeline::m_mutex
private

Definition at line 82 of file pipeline.hpp.

◆ m_pool

boost::thread_group dunedaq::oks::OksPipeline::m_pool
private

Definition at line 85 of file pipeline.hpp.

◆ m_workers

std::vector<WorkerPtr> dunedaq::oks::OksPipeline::m_workers
private

Definition at line 86 of file pipeline.hpp.


The documentation for this class was generated from the following files: