DUNE-DAQ
DUNE Trigger and Data Acquisition software
Loading...
Searching...
No Matches
DAPHNEEthFrameProcessor.cpp
Go to the documentation of this file.
1
9
13
14#include "confmodel/GeoId.hpp"
15
16#include <atomic>
17#include <functional>
18#include <memory>
19#include <string>
20
23
25DUNE_DAQ_TYPESTRING(std::vector<dunedaq::trigger::TriggerPrimitiveTypeAdapter>, "TriggerPrimitiveVector")
26
27namespace dunedaq {
28namespace fdreadoutlibs {
29
30void
32{
33 TLOG() << "Looking for TP sink...";
34
35 for (auto output : conf->get_outputs()) {
36 TLOG() << "On outputs... (" << output->UID() << "," << output->get_data_type() << ")";
37 try {
38 if (output->get_data_type() == "TriggerPrimitiveVector") {
39 TLOG() << "Found TP sink.";
40 m_tp_sink = get_iom_sender<std::vector<trigger::TriggerPrimitiveTypeAdapter>>(output->UID());
41 TLOG() << " SINK INITIALIZED for TriggerPrimitives with UID : " << output->UID();
42 }
43 } catch (const ers::Issue& excpt) {
44 ers::error(datahandlinglibs::ResourceQueueError(ERS_HERE, "tp", "DefaultRequestHandlerModel", excpt));
45 }
46 }
47
48 TLOG() << "Registering processing tasks...";
49 inherited::add_preprocess_task(std::bind(&DAPHNEEthFrameProcessor::timestamp_check, this, std::placeholders::_1));
50
52 if (dp == nullptr) {
53 TLOG() << " PDS Data processor does not exist.";
54 } else {
55 auto proc_conf = dp->cast<appmodel::PDSRawDataProcessor>();
56 if (proc_conf == nullptr) {
57 TLOG() << "PDS RawDataProcessor does not exist.";
58 } else {
59 m_def_adc_intg_thresh = proc_conf->get_default_adc_intg_thresh();
60
61 auto geo_id = conf->get_geo_id();
62 if (geo_id != nullptr) {
63 m_det_id = geo_id->get_detector_id();
64 m_crate_id = geo_id->get_crate_id();
65 m_slot_id = geo_id->get_slot_id();
66 m_stream_id = geo_id->get_stream_id();
67 }
68
69 m_channel_map = dunedaq::detchannelmaps::make_pds_map(proc_conf->get_channel_map());
70 const std::vector<unsigned int> channel_mask_vec = proc_conf->get_channel_mask();
71
72 for (int chan = 0; chan < 48; chan++) { // 40 physical PDS channel 8 not. 0->7 contain light info, 8,9, additional
73 // info. 10-17 light, 18,19 not etc...
74 trgdataformats::channel_t off_channel = m_channel_map->get_offline_channel_from_det_crate_slot_stream_chan(
76 if (std::find(channel_mask_vec.begin(), channel_mask_vec.end(), off_channel) != channel_mask_vec.end())
77 m_channel_mask_set.insert(
78 off_channel); // m_channel_mask will be a vector fille with random chanel which need to be masked.
79 }
80
82 // Extract TPs back as a pre-processing task, due to LatencyBuffer post-proc issues using SkipList.
83 inherited::add_preprocess_task(std::bind(&DAPHNEEthFrameProcessor::extract_tps, this, std::placeholders::_1));
84 }
85 }
86 }
87
88 TLOG() << "Calling parent conf.";
90}
91
92void
93DAPHNEEthFrameProcessor::start(const appfwk::DAQModule::CommandData_t& args)
94{
95 // Reset timestamp check
96 m_previous_ts = 0;
97 m_current_ts = 0;
100
101 // Reset stats
102 m_t0 = std::chrono::high_resolution_clock::now();
103 m_num_new_tps.exchange(0);
104
105 inherited::start(args);
106}
107void
108DAPHNEEthFrameProcessor::stop(const appfwk::DAQModule::CommandData_t& args)
109{
110 inherited::stop(args);
111}
112
115void
117{
118
119 /*
120 for (size_t i=0; i<types::kDAPHNENumFrames; i++){
121 auto df_ptr = reinterpret_cast<dunedaq::fddetdataformats::DAPHNEEthFrame*>(fp);
122
123 if(df_ptr[i].get_timestamp() > 0xFFFFFFFFFFFF0000 || df_ptr[i].get_timestamp() < 0xFFFF){
124 ers::warning(PDSUnphysicalFrameTimestamp(ERS_HERE, df_ptr[i].get_timestamp(), df_ptr[i].get_channel(), i));
125 // Force the TS to 0
126 df_ptr[i].daq_header.timestamp_1 = df_ptr[i].daq_header.timestamp_2 = 0;
127 }
128 }
129 */
130
131 // Acquire timestamp
133 uint64_t k_clock_frequency = 62500000; // NOLINT(build/unsigned)
134 TLOG_DEBUG(TLVL_FRAME_RECEIVED) << "Received DAPHNE frame timestamp value of " << m_current_ts << " ticks (..."
135 << std::fixed << std::setprecision(8)
136 << (static_cast<double>(m_current_ts % (k_clock_frequency * 1000)) /
137 static_cast<double>(k_clock_frequency))
138 << " sec)"; // NOLINT
139
140 if (m_ts_error_ctr > 1000) {
141 if (!m_problem_reported) {
142 std::cout << "*** Data Integrity ERROR *** Timestamp continuity is completely broken! "
143 << "Something is wrong with the FE source or with the configuration!\n";
144 m_problem_reported = true;
145 }
146 }
147
150}
151
155void
157{
158 // check error fields
159}
160
161void
163{
164
165 if (!fp || fp == nullptr) {
166 return;
167 }
168
169 /*
170 auto nonconstframeptr = const_cast<frameptr>(fp);
171 auto df_ptr = reinterpret_cast<dunedaq::fddetdataformats::DAPHNEEthFrame*>((uint8_t*)nonconstframeptr); // NOLINT
172 std::vector<trigger::TriggerPrimitiveTypeAdapter> ttpp;
173
174 for (size_t i=0; i<types::kDAPHNENumFrames; i++)
175 {
176 for(size_t j=0; j<fddetdataformats::DAPHNEEthFrame::PeakDescriptorData::max_peaks;j++)
177 {
178 if(df_ptr[i].peaks_data.is_found(j))
179 {
180 int ch = m_channel_map->get_offline_channel_from_det_crate_slot_stream_chan(df_ptr[i].daq_header.det_id,
181 df_ptr[i].daq_header.crate_id, df_ptr[i].daq_header.slot_id, df_ptr[i].daq_header.link_id, df_ptr[i].get_channel());
182 if (std::binary_search(m_channel_mask_set.begin(), m_channel_mask_set.end(), ch)) continue;
183 if (df_ptr[i].peaks_data.get_adc_integral(j) < m_def_adc_intg_thresh) continue;
184
185
186 trigger::TriggerPrimitiveTypeAdapter tpa;
187 tpa.tp = peak_to_tp(df_ptr[i],j);// this is the trigger primitive
188 //check for timestamps that are due to frame timestamps ~ ts=0, and ignore these peaks
189 if(tpa.tp.time_start > 0xFFFFFFFFFFFF0000 || tpa.tp.time_start < 0xFFFF){
190 ers::warning(PDSPeakIgnored(ERS_HERE, tpa.tp.time_start, tpa.tp.channel, i, j));
191 continue;
192 }
193
194 tpa.tp.detid = df_ptr->daq_header.det_id;
195 ttpp.push_back(tpa);
196 }
197 }
198 }
199
200 int num_new_tps = ttpp.size();
201 if (num_new_tps > 0) {
202
203 const auto s_ts_begin = ttpp.front().tp.time_start;
204 const auto channel_begin = ttpp.front().tp.channel;
205 const auto s_ts_end = ttpp.back().tp.time_start;
206 const auto channel_end = ttpp.back().tp.channel;
207
208 if (!m_tp_sink->try_send(std::move(ttpp), iomanager::Sender::s_no_block)) {
209 ers::warning(FailedToSendTPVector(ERS_HERE, s_ts_begin, channel_begin, s_ts_end, channel_end));
210 m_tps_send_failed += num_new_tps;
211 } else {
212 m_num_new_tps += num_new_tps;
213 }
214 }
215 */
216 return;
217}
218
219void
221{
222
223 // right now, just fill some basic tp info...
225 auto now = std::chrono::high_resolution_clock::now();
226 int num_new_tps = m_num_new_tps.exchange(0);
227 int num_new_tps_suppressed_too_long = 0; // not relevant for PDS TPs
228 int num_new_tps_send_failed = m_tps_send_failed.exchange(0);
229 double seconds = std::chrono::duration_cast<std::chrono::microseconds>(now - m_t0).count() / 1000000.;
230 TLOG_DEBUG(TLVL_BOOKKEEPING) << "TP rate: " << std::to_string(num_new_tps / seconds / 1000.) << " [kHz]";
231 TLOG_DEBUG(TLVL_BOOKKEEPING) << "Total new TPs: " << num_new_tps;
232
234 tp_info.set_rate_tp_hits(num_new_tps / seconds / 1000.);
235
236 tp_info.set_num_tps_sent(num_new_tps);
237 tp_info.set_num_tps_suppressed_too_long(num_new_tps_suppressed_too_long);
238 tp_info.set_num_tps_send_failed(num_new_tps_send_failed);
239
240 publish(std::move(tp_info));
241
242 m_t0 = now;
243 }
244
246}
247
248} // namespace fdreadoutlibs
249} // namespace dunedaq
@ TLVL_FRAME_RECEIVED
#define ERS_HERE
#define DUNE_DAQ_TYPESTRING(Type, typestring)
Declare the datatype_to_string method for the given type.
const dunedaq::appmodel::DataProcessor * get_data_processor() const
Get "data_processor" relationship value.
const dunedaq::appmodel::DataHandlerConf * get_module_configuration() const
Get "module_configuration" relationship value.
const dunedaq::confmodel::GeoId * get_geo_id() const
Get "geo_id" relationship value.
const std::vector< const dunedaq::confmodel::Connection * > & get_outputs() const
Get "outputs" relationship value. Output connections from this module.
void start(const appfwk::DAQModule::CommandData_t &) override
void stop(const appfwk::DAQModule::CommandData_t &) override
std::chrono::time_point< std::chrono::high_resolution_clock > m_t0
std::shared_ptr< detchannelmaps::PDSChannelMap > m_channel_map
std::shared_ptr< iomanager::SenderConcept< std::vector< trigger::TriggerPrimitiveTypeAdapter > > > m_tp_sink
void stop(const appfwk::DAQModule::CommandData_t &args) override
Stop operation.
void start(const appfwk::DAQModule::CommandData_t &args) override
Start operation.
void conf(const appmodel::DataHandlerModule *conf) override
Set the emulator mode, if active, timestamps of processed packets are overwritten with new ones.
void publish(google::protobuf::Message &&, CustomOrigin &&co={}, OpMonLevel l=to_level(EntryOpMonLevel::kDefault)) const noexcept
Base class for any user define issue.
Definition Issue.hpp:76
#define TLOG_DEBUG(lvl,...)
Definition Logging.hpp:116
#define TLOG(...)
Definition macro.hpp:21
The DUNE-DAQ namespace.
void error(const Issue &issue)
Definition ers.hpp:101