DUNE-DAQ
DUNE Trigger and Data Acquisition software
Toggle main menu visibility
Loading...
Searching...
No Matches
dunedaq
sourcecode
datahandlinglibs
include
datahandlinglibs
models
DataSubscriberModel.hpp
Go to the documentation of this file.
1
8
#ifndef DATAHANDLINGLIBS_SRC_SOURCEMODEL_HPP_
9
#define DATAHANDLINGLIBS_SRC_SOURCEMODEL_HPP_
10
11
#include "
datahandlinglibs/concepts/SourceConcept.hpp
"
12
#include "
datahandlinglibs/opmon/datahandling_info.pb.h
"
13
14
#include "
iomanager/IOManager.hpp
"
15
#include "
iomanager/Receiver.hpp
"
16
#include "
iomanager/Sender.hpp
"
17
#include "
logging/Logging.hpp
"
18
// #include "utilities/ReusableThread.hpp"
19
20
#include "
confmodel/Connection.hpp
"
21
22
// #include <folly/ProducerConsumerQueue.h>
23
// #include <nlohmann/json.hpp>
24
25
// #include <atomic>
26
// #include <memory>
27
// #include <mutex>
28
#include <functional>
29
30
namespace
dunedaq::datahandlinglibs
{
31
32
template
<
class
PayloadType>
33
class
DataSubscriberModel
:
public
SourceConcept
34
{
35
public
:
36
using
inherited
=
SourceConcept
;
37
42
DataSubscriberModel
()
43
:
SourceConcept
()
44
{
45
}
46
~DataSubscriberModel
() {}
47
48
void
init
(
const
confmodel::DaqModule
* cfg)
override
49
{
50
if
(cfg->
get_outputs
().size() != 1) {
51
throw
datahandlinglibs::InitializationError(
ERS_HERE
,
"Only 1 output supported for subscribers"
);
52
}
53
m_data_sender
= get_iom_sender<PayloadType>(cfg->
get_outputs
()[0]->UID());
54
55
if
(cfg->
get_inputs
().size() != 1) {
56
throw
datahandlinglibs::InitializationError(
ERS_HERE
,
"Only 1 input supported for subscribers"
);
57
}
58
m_data_receiver
= get_iom_receiver<PayloadType>(cfg->
get_inputs
()[0]->UID());
59
}
60
61
void
start
()
62
{
63
m_packets
= 0;
64
m_sum_packets
= 0;
65
m_dropped_packets
= 0;
66
m_data_receiver
->add_callback(std::bind(&
DataSubscriberModel::handle_payload
,
this
, std::placeholders::_1));
67
}
68
69
void
stop
() {
m_data_receiver
->remove_callback(); }
70
71
bool
handle_payload
(PayloadType& message)
// NOLINT(build/unsigned)
72
{
73
++
m_packets
;
74
++
m_sum_packets
;
75
if
(!
m_data_sender
->try_send(std::move(message),
iomanager::Sender::s_no_block
)) {
76
++
m_dropped_packets
;
77
}
78
return
true
;
79
}
80
81
protected
:
82
virtual
void
generate_opmon_data
()
override
83
{
84
opmon::DataSourceInfo
info;
85
info.set_num_packets(
m_packets
.exchange(0));
86
info.set_sum_packets(
m_sum_packets
);
87
info.set_num_dropped_packets(
m_dropped_packets
.exchange(0));
88
89
this->
publish
(std::move(info));
90
}
91
92
private
:
93
using
source_t
=
dunedaq::iomanager::ReceiverConcept<PayloadType>
;
94
std::shared_ptr<source_t>
m_data_receiver
;
95
96
using
sink_t
=
dunedaq::iomanager::SenderConcept<PayloadType>
;
97
std::shared_ptr<sink_t>
m_data_sender
;
98
99
std::atomic<uint64_t>
m_packets
{ 0 };
100
std::atomic<uint64_t>
m_sum_packets
{ 0 };
101
std::atomic<uint64_t>
m_dropped_packets
{ 0 };
102
};
103
104
}
// namespace dunedaq::datahandlinglibs
105
106
#endif
// DATAHANDLINGLIBS_SRC_SOURCEMODEL_HPP_
IOManager.hpp
ERS_HERE
#define ERS_HERE
Definition
LocalContext.hpp:127
dunedaq::confmodel::DaqModule
Definition
DaqModule.hpp:34
dunedaq::confmodel::DaqModule::get_inputs
const std::vector< const dunedaq::confmodel::Connection * > & get_inputs() const
Get "inputs" relationship value. List of connections to/from this module.
Definition
DaqModule.hpp:111
dunedaq::confmodel::DaqModule::get_outputs
const std::vector< const dunedaq::confmodel::Connection * > & get_outputs() const
Get "outputs" relationship value. Output connections from this module.
Definition
DaqModule.hpp:138
dunedaq::datahandlinglibs::DataSubscriberModel::init
void init(const confmodel::DaqModule *cfg) override
Definition
DataSubscriberModel.hpp:48
dunedaq::datahandlinglibs::DataSubscriberModel::generate_opmon_data
virtual void generate_opmon_data() override
Definition
DataSubscriberModel.hpp:82
dunedaq::datahandlinglibs::DataSubscriberModel::~DataSubscriberModel
~DataSubscriberModel()
Definition
DataSubscriberModel.hpp:46
dunedaq::datahandlinglibs::DataSubscriberModel::source_t
dunedaq::iomanager::ReceiverConcept< PayloadType > source_t
Definition
DataSubscriberModel.hpp:93
dunedaq::datahandlinglibs::DataSubscriberModel::m_sum_packets
std::atomic< uint64_t > m_sum_packets
Definition
DataSubscriberModel.hpp:100
dunedaq::datahandlinglibs::DataSubscriberModel::m_dropped_packets
std::atomic< uint64_t > m_dropped_packets
Definition
DataSubscriberModel.hpp:101
dunedaq::datahandlinglibs::DataSubscriberModel::sink_t
dunedaq::iomanager::SenderConcept< PayloadType > sink_t
Definition
DataSubscriberModel.hpp:96
dunedaq::datahandlinglibs::DataSubscriberModel::m_packets
std::atomic< uint64_t > m_packets
Definition
DataSubscriberModel.hpp:99
dunedaq::datahandlinglibs::DataSubscriberModel::m_data_receiver
std::shared_ptr< source_t > m_data_receiver
Definition
DataSubscriberModel.hpp:94
dunedaq::datahandlinglibs::DataSubscriberModel::stop
void stop()
Definition
DataSubscriberModel.hpp:69
dunedaq::datahandlinglibs::DataSubscriberModel::m_data_sender
std::shared_ptr< sink_t > m_data_sender
Definition
DataSubscriberModel.hpp:97
dunedaq::datahandlinglibs::DataSubscriberModel::inherited
SourceConcept inherited
Definition
DataSubscriberModel.hpp:36
dunedaq::datahandlinglibs::DataSubscriberModel::handle_payload
bool handle_payload(PayloadType &message)
Definition
DataSubscriberModel.hpp:71
dunedaq::datahandlinglibs::DataSubscriberModel::DataSubscriberModel
DataSubscriberModel()
SourceModel Constructor.
Definition
DataSubscriberModel.hpp:42
dunedaq::datahandlinglibs::DataSubscriberModel::start
void start()
Definition
DataSubscriberModel.hpp:61
dunedaq::datahandlinglibs::SourceConcept::SourceConcept
SourceConcept()
Definition
SourceConcept.hpp:28
dunedaq::datahandlinglibs::opmon::DataSourceInfo
Definition
datahandling_info.pb.h:268
dunedaq::iomanager::ReceiverConcept
Definition
Receiver.hpp:48
dunedaq::iomanager::SenderConcept
Definition
Sender.hpp:44
dunedaq::iomanager::Sender::s_no_block
static constexpr timeout_t s_no_block
Definition
Sender.hpp:26
dunedaq::opmonlib::MonitorableObject::publish
void publish(google::protobuf::Message &&, CustomOrigin &&co={}, OpMonLevel l=to_level(EntryOpMonLevel::kDefault)) const noexcept
Definition
MonitorableObject.cpp:59
datahandling_info.pb.h
SourceConcept.hpp
Connection.hpp
Receiver.hpp
Sender.hpp
Logging.hpp
dunedaq::datahandlinglibs
Definition
DataHandlingConcept.hpp:16
Generated on
for DUNE-DAQ by
1.18.0