DUNE-DAQ
DUNE Trigger and Data Acquisition software
Loading...
Searching...
No Matches
SourceEmulatorModel.hxx
Go to the documentation of this file.
1// Declarations for SourceEmulatorModel
2
5
6namespace dunedaq {
7namespace datahandlinglibs {
8
9void
11{
12 // TLOG() << "Generate random ADC patterns" ;
13 std::srand(source_id * 12345);
14 m_channel.reserve(size);
15 for (int i = 0; i < size; i++) {
16 int random_ch = std::rand() % 64;
17 m_channel.push_back(random_ch);
18 }
19}
20
21template<class ReadoutType>
22void
33
34template<class ReadoutType>
35void
38{
39 if (m_is_configured) {
40 TLOG_DEBUG(TLVL_WORK_STEPS) << "This emulator is already configured!";
41 } else {
42 // m_conf = args.get<module_conf_t>();
43 // m_link_conf = link_conf.get<link_conf_t>();
44
45 std::mt19937 mt(rand()); // NOLINT(runtime/threadsafe_fn)
46 std::uniform_real_distribution<double> dis(0.0, 1.0);
47
48 m_sourceid.id = link_conf->get_source_id();
49 m_sourceid.subsystem = ReadoutType::subsystem;
50
51 m_crateid = link_conf->get_geo_id()->get_crate_id();
52 m_slotid = link_conf->get_geo_id()->get_slot_id();
53 m_linkid = link_conf->get_geo_id()->get_stream_id();
54
55 m_t0_now = emu_params->get_set_t0();
56 m_file_source = std::make_unique<FileSourceBuffer>(emu_params->get_input_file_size_limit(), sizeof(ReadoutType));
57 try {
58 m_file_source->read(emu_params->get_data_file_name());
59 } catch (const ers::Issue& ex) {
60 ers::fatal(ex);
61 throw ConfigurationError(ERS_HERE, m_sourceid, "", ex);
62 }
64 if (m_dropout_rate == 0.0) {
65 m_dropouts = std::vector<bool>(1);
66 } else {
67 m_dropouts = std::vector<bool>(m_dropouts_length);
68 }
69 for (size_t i = 0; i < m_dropouts.size(); ++i) {
70 m_dropouts[i] = dis(mt) >= m_dropout_rate;
71 }
72
76 m_error_bit_generator.generate();
77
78 // Generate random ADC pattern
80 auto vec_size = emu_params->get_random_population_size();
82 TLOG() << "Generated pattern.";
83 m_pattern_generator.generate(m_sourceid.id, vec_size);
84
85 if (emu_params->get_TP_rate_per_channel() != 0) {
86 TLOG() << "TP rate per channel multiplier (base of 100 Hz/ch): " << emu_params->get_TP_rate_per_channel();
87 // Define time to wait when adding an ADC above threshold
88 // Adding a hit every 9768 gives a total Sent TP rate of approx 100 Hz/wire with WIBEth
90 }
91 }
92
93 m_is_configured = true;
94 }
95 // Configure thread:
96 m_producer_thread.set_name("fakeprod", m_sourceid.id);
97}
98
99template<class ReadoutType>
100void
101SourceEmulatorModel<ReadoutType>::start(const appfwk::DAQModule::CommandData_t& /*args*/)
102{
104 TLOG_DEBUG(TLVL_WORK_STEPS) << "Starting threads...";
105 // FIXME: don't know where to take the slowdown from... m_rate_limiter = std::make_unique<RateLimiter>(m_rate_khz /
106 // m_link_conf.slowdown);
107 m_rate_limiter = std::make_unique<RateLimiter>(m_rate_khz);
108 // m_stats_thread.set_work(&SourceEmulatorModel<ReadoutType>::run_stats, this);
110}
111
112template<class ReadoutType>
113void
114SourceEmulatorModel<ReadoutType>::stop(const appfwk::DAQModule::CommandData_t& /*args*/)
115{
116 while (!m_producer_thread.get_readiness()) {
117 std::this_thread::sleep_for(std::chrono::milliseconds(100));
118 }
119}
120
121template<class ReadoutType>
122void
124{
126 info.set_sum_packets(m_packet_count_tot.load());
127 info.set_num_packets(m_packet_count.exchange(0));
128
129 this->publish(std::move(info));
130}
131
132template<class ReadoutType>
133void
135{
136 TLOG_DEBUG(TLVL_WORK_STEPS) << "Data generation thread " << m_this_link_number << " started";
137
138 // pthread_setname_np(pthread_self(), get_name().c_str());
139
140 uint offset = 0; // NOLINT(build/unsigned)
141 auto& source = m_file_source->get();
142
143 uint num_elem = m_file_source->num_elements();
144 if (num_elem == 0) {
145 TLOG_DEBUG(TLVL_WORK_STEPS) << "No elements to read from buffer! Sleeping...";
146 std::this_thread::sleep_for(std::chrono::milliseconds(100));
147 num_elem = m_file_source->num_elements();
148 }
149
150 auto rptr = reinterpret_cast<ReadoutType*>(source.data()); // NOLINT
151
152 // set the initial timestamp to a configured value, otherwise just use the timestamp from the header
153 uint64_t ts_0 = rptr->get_timestamp(); // NOLINT(build/unsigned)
154 if (m_t0_now) {
155 auto time_now = std::chrono::system_clock::now().time_since_epoch();
156 uint64_t current_time = // NOLINT (build/unsigned)
157 std::chrono::duration_cast<std::chrono::microseconds>(time_now).count();
158 // FIXME: where do I get the clockspeed from?
159 // ts_0 = (m_conf.clock_speed_hz / 100000) * current_time;
160 ts_0 = 625 * current_time / 10;
161 }
162 TLOG_DEBUG(TLVL_BOOKKEEPING) << "Using first timestamp: " << ts_0;
163 uint64_t timestamp = ts_0; // NOLINT(build/unsigned)
164 int dropout_index = 0;
165 uint64_t number_pattern_hits_generated = 0;
166 // 64 total channels, placing on slot 0 gives 64 available slots.
167 const uint64_t max_tps_per_frame = 64;
168
169 while (m_run_marker.load()) {
170 // TLOG() << "Generating " << m_frames_per_tick << " for TS " << timestamp;
171 for (uint16_t i = 0; i < m_frames_per_tick; i++) {
172 // Which element to push to the buffer
173 if (offset == num_elem || (offset + 1) * sizeof(ReadoutType) > source.size()) {
174 offset = 0;
175 }
176
177 bool create_frame = m_dropouts[dropout_index]; // NOLINT(runtime/threadsafe_fn)
178 dropout_index = (dropout_index + 1) % m_dropouts.size();
179 if (create_frame) {
180 ReadoutType payload;
181 // Memcpy from file buffer to flat char array
182 ::memcpy(static_cast<void*>(&payload),
183 static_cast<void*>(source.data() + offset * sizeof(ReadoutType)),
184 sizeof(ReadoutType));
185
186 // Fake timestamp
187 payload.fake_timestamps(timestamp, m_time_tick_diff);
188
189 // Fake geoid
190 payload.fake_geoid(m_crateid, m_slotid, m_linkid);
191
192 // Introducing frame errors
193 std::vector<uint16_t> frame_errs; // NOLINT(build/unsigned)
194 for (size_t i = 0; i < rptr->get_num_frames(); ++i) {
195 frame_errs.push_back(m_error_bit_generator.next());
196 }
197 payload.fake_frame_errors(&frame_errs);
198
200 uint64_t number_pattern_hits_expected = (timestamp - ts_0) / m_time_to_wait;
201
202 // Calculate how many TPs to generate in this frame
203 uint64_t tps_this_frame = 0;
204 if (number_pattern_hits_expected > number_pattern_hits_generated) {
205 tps_this_frame = number_pattern_hits_expected - number_pattern_hits_generated;
206 }
207
208 if (tps_this_frame > max_tps_per_frame) {
209 tps_this_frame = max_tps_per_frame;
210 }
211
212 // Distribute TPs across channels via the pattern generator.
213 for (uint64_t tp_idx = 0; tp_idx < tps_this_frame; ++tp_idx) {
214 int channel = m_pattern_generator.get_channel_number();
215 // The pattern generator draws channel in range 0-63, current
216 // behaviour for frame type with 32 channels is silent dropping.
217 try {
218 payload.fake_adc_pattern(channel);
219 } catch (const std::out_of_range&) {
220 }
221 }
222
223 // Count the number of patterns attempted to inject. Prevents
224 // expected - generated deficit from accumulating in case of
225 // injection failure.
226 number_pattern_hits_generated += tps_this_frame;
227 }
228
229 // send it
230 try {
231 (*m_raw_data_callback)(std::move(payload));
232 } catch (ers::Issue& excpt) {
233 ers::warning(CannotWriteToQueue(ERS_HERE, m_sourceid, "raw data input queue", excpt));
234 // std::runtime_error("Queue timed out...");
235 }
236
237 // Count packet and limit rate if needed.
238 ++offset;
241 }
242 }
243 timestamp += m_time_tick_diff * rptr->get_num_frames();
244
245 m_rate_limiter->limit();
246 }
247 TLOG_DEBUG(TLVL_WORK_STEPS) << "Data generation thread " << m_sourceid.to_string() << " finished";
248}
249
250} // namespace datahandlinglibs
251} // namespace dunedaq
#define ERS_HERE
uint32_t get_random_population_size() const
Get "random_population_size" attribute value.
bool get_generate_periodic_adc_pattern() const
Get "generate_periodic_adc_pattern" attribute value.
bool get_set_t0() const
Get "set_t0" attribute value. Set first timestamp to now.
uint32_t get_input_file_size_limit() const
Get "input_file_size_limit" attribute value.
float get_frame_error_rate_hz() const
Get "frame_error_rate_hz" attribute value.
const std::string & get_data_file_name() const
Get "data_file_name" attribute value.
float get_TP_rate_per_channel() const
Get "TP_rate_per_channel" attribute value. TP rate per channel in units of 100 Hz.
const dunedaq::confmodel::GeoId * get_geo_id() const
Get "geo_id" relationship value.
uint32_t get_source_id() const
Get "source_id" attribute value.
uint32_t get_stream_id() const
Get "stream_id" attribute value.
Definition GeoId.hpp:196
uint32_t get_slot_id() const
Get "slot_id" attribute value.
Definition GeoId.hpp:165
uint32_t get_crate_id() const
Get "crate_id" attribute value.
Definition GeoId.hpp:134
static std::shared_ptr< DataMoveCallbackRegistry > get()
const appmodel::DataMoveCallbackConf * m_sink_conf
void stop(const appfwk::DAQModule::CommandData_t &)
void conf(const confmodel::DetectorStream *stream_conf, const appmodel::StreamEmulationParameters *emu_conf)
void start(const appfwk::DAQModule::CommandData_t &)
std::shared_ptr< std::function< void(ReadoutType &&)> > m_raw_data_callback
std::unique_ptr< FileSourceBuffer > m_file_source
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
double offset
#define TLOG_DEBUG(lvl,...)
Definition Logging.hpp:116
#define TLOG(...)
Definition macro.hpp:21
The DUNE-DAQ namespace.
FELIX Initialization std::string initerror FELIX queue timed std::string queuename Unexpected chunk size
void warning(const Issue &issue)
Definition ers.hpp:150
void fatal(const Issue &issue)
Definition ers.hpp:111