DUNE-DAQ
DUNE Trigger and Data Acquisition software
Loading...
Searching...
No Matches
ZeroCopyRecordingRequestHandlerModel.hxx
Go to the documentation of this file.
1// Declarations for ZeroCopyRecordingRequestHandlerModel
2
3namespace dunedaq {
4namespace datahandlinglibs {
5
6// Special configuration that checks LB alignment and O_DIRECT flag on output file
7template<class ReadoutType, class LatencyBufferType>
8void
10{
12
13 if (data_rec_conf != nullptr) {
14 if (!data_rec_conf->get_output_file().empty()) {
16 inherited::m_sourceid.subsystem = ReadoutType::subsystem;
17
18 // Check for alignment restrictions for filesystem block size. (XFS default: 4096)
19 if (inherited::m_latency_buffer->get_alignment_size() == 0 ||
20 sizeof(ReadoutType) * inherited::m_latency_buffer->size() % 4096) {
21 ers::error(ConfigurationError(ERS_HERE, inherited::m_sourceid, "Latency buffer is not 4kB aligned"));
22 }
23
24 // Check for sensible stream chunk size
25 inherited::m_stream_buffer_size = data_rec_conf->get_streaming_buffer_size();
26 if (inherited::m_stream_buffer_size % 4096 != 0) {
28 ConfigurationError(ERS_HERE, inherited::m_sourceid, "Streaming chunk size is not divisible by 4kB!"));
29 }
30
31 // Prepare filename with full path
32 std::string file_full_path =
33 data_rec_conf->get_output_file() + inherited::m_sourceid.to_string() + std::string(".bin");
34 inherited::m_output_file = file_full_path;
35
36 // RS: This will need to go away with the SNB store handler!
37 if (std::remove(file_full_path.c_str()) == 0) {
38 TLOG(TLVL_WORK_STEPS) << "Removed existing output file from previous run: " << file_full_path;
39 }
40
41 m_oflag = O_CREAT | O_WRONLY;
42 if (data_rec_conf->get_use_o_direct()) {
43 m_oflag |= O_DIRECT;
44 }
45 m_fd = ::open(file_full_path.c_str(), m_oflag, 0644);
46 if (m_fd == -1) {
47 TLOG() << "Failed to open file!";
48 throw ConfigurationError(ERS_HERE, inherited::m_sourceid, "Failed to open file!");
49 }
51
52 } else { // no output dir specified
53 TLOG(TLVL_WORK_STEPS) << "No output path is specified in data recorder config. Recording feature is inactive.";
54 }
55 } else {
56 TLOG(TLVL_WORK_STEPS) << "No recording config object specified. Recording feature is inactive.";
57 }
58
60}
61
62// Special record command that writes to files from memory aligned LBs
63template<class ReadoutType, class LatencyBufferType>
64void
66 const appfwk::DAQModule::CommandData_t& cmdargs)
67{
68 if (inherited::m_recording.load()) {
70 CommandError(ERS_HERE, inherited::m_sourceid, "A recording is still running, no new recording was started!"));
71 return;
72 }
73
74 // FIXME: Recording parameters to be clarified!
75 int recording_time_sec = 0;
76 if (cmdargs.contains("duration")) {
77 recording_time_sec = cmdargs["duration"];
78 } else {
80 CommandError(ERS_HERE, inherited::m_sourceid, "A recording command with missing duration field received!"));
81 }
82 if (recording_time_sec == 0) {
84 ERS_HERE, inherited::m_sourceid, "Recording for 0 seconds requested. Recording command is ignored!"));
85 return;
86 }
87
89 [&](int duration) {
90 size_t chunk_size = inherited::m_stream_buffer_size;
91 size_t alignment_size = inherited::m_latency_buffer->get_alignment_size();
92 TLOG() << "Start recording for " << duration << " second(s)" << std::endl;
93 inherited::m_recording.exchange(true);
94 auto start_of_recording = std::chrono::high_resolution_clock::now();
95 auto current_time = start_of_recording;
97
98 const char* current_write_pointer = nullptr;
99 const char* start_of_buffer_pointer =
100 reinterpret_cast<const char*>(inherited::m_latency_buffer->start_of_buffer()); // NOLINT
101 const char* current_end_pointer;
102 const char* end_of_buffer_pointer =
103 reinterpret_cast<const char*>(inherited::m_latency_buffer->end_of_buffer()); // NOLINT
104
105 size_t bytes_written = 0;
106 size_t failed_writes = 0;
107
108 while (std::chrono::duration_cast<std::chrono::seconds>(current_time - start_of_recording).count() < duration) {
110 size_t considered_chunks_in_loop = 0;
111
112 // Wait for potential running cleanup to finish first
113 {
114 std::unique_lock<std::mutex> lock(inherited::m_cv_mutex);
115 inherited::m_cv.wait(lock, [&] { return !inherited::m_cleanup_requested; });
116 }
117 inherited::m_cv.notify_all();
118
119 // Some frames have to be skipped to start copying from an aligned piece of memory
120 // These frames cannot be written without O_DIRECT as this would mess up the alignment of the write pointer
121 // into the target file
123 auto begin = inherited::m_latency_buffer->begin();
124 if (begin == inherited::m_latency_buffer->end()) {
125 // There are no elements in the buffer, update time and try again
126 current_time = std::chrono::high_resolution_clock::now();
127 continue;
128 }
129 inherited::m_next_timestamp_to_record = begin->get_timestamp();
130 size_t skipped_frames = 0;
131 while (reinterpret_cast<std::uintptr_t>(&(*begin)) % alignment_size) { // NOLINT
132 ++begin;
133 skipped_frames++;
134 if (!begin.good()) {
135 // We reached the end of the buffer without finding an aligned element
136 // Reset the next timestamp to record and try again
137 current_time = std::chrono::high_resolution_clock::now();
139 continue;
140 }
141 }
142 TLOG() << "Skipped " << skipped_frames << " frames";
143 current_write_pointer = reinterpret_cast<const char*>(&(*begin)); // NOLINT
144 }
145
146 current_end_pointer = reinterpret_cast<const char*>(inherited::m_latency_buffer->back()); // NOLINT
147
148 // Break the loop from time to time to update the timestamp and check if we should stop recording
149 while (considered_chunks_in_loop < 100) {
150 auto iptr = reinterpret_cast<std::uintptr_t>(current_write_pointer); // NOLINT
151 if (iptr % alignment_size) {
152 // This should never happen
153 TLOG() << "Error: Write pointer is not aligned";
154 }
155 bool failed_write = false;
156 if (current_write_pointer + chunk_size < current_end_pointer) {
157 // We can write a whole chunk to file
158 failed_write |= !::write(m_fd, current_write_pointer, chunk_size);
159 if (!failed_write) {
160 bytes_written += chunk_size;
161 }
162 current_write_pointer += chunk_size;
163 } else if (current_end_pointer < current_write_pointer) {
164 if (current_write_pointer + chunk_size < end_of_buffer_pointer) {
165 // Write whole chunk to file
166 failed_write |= !::write(m_fd, current_write_pointer, chunk_size);
167 if (!failed_write) {
168 bytes_written += chunk_size;
169 }
170 current_write_pointer += chunk_size;
171 } else {
172 // Write the last bit of the buffer without using O_DIRECT as it possibly doesn't fulfill the
173 // alignment requirement
174 fcntl(m_fd, F_SETFL, O_CREAT | O_WRONLY);
175 failed_write |= !::write(m_fd, current_write_pointer, end_of_buffer_pointer - current_write_pointer);
176 fcntl(m_fd, F_SETFL, m_oflag);
177 if (!failed_write) {
178 bytes_written += end_of_buffer_pointer - current_write_pointer;
179 }
180 current_write_pointer = start_of_buffer_pointer;
181 }
182 }
183
184 if (current_write_pointer == end_of_buffer_pointer) {
185 current_write_pointer = start_of_buffer_pointer;
186 }
187
188 if (failed_write) {
189 ++failed_writes;
191 }
192 considered_chunks_in_loop++;
193 // This expression is "a bit" complicated as it finds the last frame that was written to file completely
195 reinterpret_cast<const ReadoutType*>( // NOLINT
196 start_of_buffer_pointer +
197 (((current_write_pointer - start_of_buffer_pointer) / ReadoutType::fixed_payload_size) *
198 ReadoutType::fixed_payload_size))
199 ->get_timestamp();
200 }
201 }
202 current_time = std::chrono::high_resolution_clock::now();
203 }
204
205 // Complete writing the last frame to file
206 if (current_write_pointer != nullptr) {
207 const char* last_started_frame =
208 start_of_buffer_pointer +
209 (((current_write_pointer - start_of_buffer_pointer) / ReadoutType::fixed_payload_size) *
210 ReadoutType::fixed_payload_size);
211 if (last_started_frame != current_write_pointer) {
212 fcntl(m_fd, F_SETFL, O_CREAT | O_WRONLY);
213 if (!::write(m_fd,
214 current_write_pointer,
215 (last_started_frame + ReadoutType::fixed_payload_size) - current_write_pointer)) {
217 } else {
218 bytes_written += (last_started_frame + ReadoutType::fixed_payload_size) - current_write_pointer;
219 }
220 }
221 }
222 ::close(m_fd);
223
224 inherited::m_next_timestamp_to_record = std::numeric_limits<uint64_t>::max(); // NOLINT (build/unsigned)
225
226 TLOG() << "Stopped recording, wrote " << bytes_written << " bytes. Failed write count: " << failed_writes;
227 inherited::m_recording.exchange(false);
228 },
229 recording_time_sec);
230}
231
232} // namespace datahandlinglibs
233} // namespace dunedaq
#define ERS_HERE
const dunedaq::appmodel::RequestHandler * get_request_handler() const
Get "request_handler" relationship value.
const dunedaq::appmodel::DataHandlerConf * get_module_configuration() const
Get "module_configuration" relationship value.
uint32_t get_source_id() const
Get "source_id" attribute value.
const dunedaq::appmodel::DataRecorderConf * get_data_recorder() const
Get "data_recorder" relationship value.
void conf(const dunedaq::appmodel::DataHandlerModule *)
void record(const appfwk::DAQModule::CommandData_t &args) override
#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 error(const Issue &issue)
Definition ers.hpp:101