DUNE-DAQ
DUNE Trigger and Data Acquisition software
Loading...
Searching...
No Matches
FileDataStoreImpl.hpp
Go to the documentation of this file.
1
32
33#ifndef DFMODULES_INCLUDE_DFMODULES_FILEDATASTOREIMPL_HPP_
34#define DFMODULES_INCLUDE_DFMODULES_FILEDATASTOREIMPL_HPP_
35
38
42#include "confmodel/Session.hpp"
43
44#include "appfwk/DAQModule.hpp"
46#include "logging/Logging.hpp" // NOTE: if ISSUES ARE DECLARED BEFORE include logging/Logging.hpp, TLOG_DEBUG<<issue wont work.
47
48#include "boost/date_time/posix_time/posix_time.hpp"
49#include "boost/lexical_cast.hpp"
50
51#include <concepts>
52#include <cstdlib>
53#include <functional>
54#include <limits>
55#include <memory>
56#include <string>
57#include <sys/statvfs.h>
58#include <utility>
59#include <vector>
60
61namespace dunedaq {
62
63// Disable coverage checking LCOV_EXCL_START
68 FileDataStoreImplBadConfiguration,
69 appfwk::GeneralDAQModuleIssue,
70 "Construction of the FileDataStoreImpl base class failed due to faulty configuration",
71 ((std::string)name), )
72
76 "Selected operation mode \"" << selected_operation
77 << "\" is NOT supported. Please update the configuration file.",
78 ((std::string)name),
79 ((std::string)selected_operation))
80
82 FileOperationProblem,
84 "A problem was encountered when opening or closing file \"" << filename << "\"",
85 ((std::string)name),
86 ((std::string)filename))
87
89 InvalidOutputPath,
91 "The specified output destination, \"" << output_path
92 << "\", is not a valid file system path on this server.",
93 ((std::string)name),
94 ((std::string)output_path))
95
97 InsufficientDiskSpace,
99 "There is insufficient free space on the disk associated with output file path \""
100 << path << "\". There are " << free_bytes << " bytes free, and the "
101 << "required minimum is " << needed_bytes << " bytes based on " << criteria << ".",
102 ((std::string)name),
103 ((std::string)path)((size_t)free_bytes)((size_t)needed_bytes)((std::string)criteria))
104
106 TimeSliceAlreadyExists,
108 "The TimeSlice record for timeslice #" << timeslice_number << " already exists.",
109 ((std::string)name),
110 ((daqdataformats::timeslice_number_t)timeslice_number))
111
112// Re-enable coverage checking LCOV_EXCL_STOP
113namespace dfmodules {
114
115// A class which satisfies the FileHandleConcept should implement format-specific writes of objects
116template<typename T>
117concept FileHandleConcept = requires(T file_handle,
118 const T const_file_handle,
120 const daqdataformats::TimeSlice& ts) {
121 { file_handle.write(tr) } -> std::same_as<void>;
122 { file_handle.write(ts) } -> std::same_as<void>;
123 { file_handle.timeslice_already_exists(ts) } -> std::convertible_to<bool>;
124 { const_file_handle.get_file_name() } -> std::convertible_to<std::string>;
125 { const_file_handle.get_file_name_extension() } -> std::convertible_to<std::string>; // Don't include the "."
126 { const_file_handle.get_recorded_size() } -> std::convertible_to<size_t>;
127 { const_file_handle.get_uncompressed_raw_data_size() } -> std::convertible_to<size_t>;
128 { const_file_handle.get_total_file_size() } -> std::convertible_to<size_t>;
129};
130
136template<FileHandleConcept FileHandleClass> // E.g., SummaryTextDataWriter
137class FileDataStoreImpl : public DataStore
138{
139
140public:
141 enum
142 {
143 TLVL_BASIC = 2
144 };
145
146 static constexpr size_t s_unset_record_number{ std::numeric_limits<size_t>::max() };
147
153 explicit FileDataStoreImpl(std::string const& name,
154 std::shared_ptr<appfwk::ConfigurationManager> mcfg,
155 std::string const& writer_name)
156 : DataStore(name)
157 , m_file_handle{ nullptr }
158 , m_run_number{ 0 }
159 , m_file_index{ 0 }
160 , m_writer_identifier{ writer_name }
161 , m_config_params{ mcfg ? mcfg->get_dal<appmodel::DataStoreConf>(name) : nullptr }
162 , m_session{ mcfg ? mcfg->get_session() : nullptr }
163 , m_compression_level{ m_config_params ? m_config_params->get_compression_level() : static_cast<unsigned>(0) }
164 // NOLINT(build/unsigned)
165 , m_operational_environment{ m_session ? m_session->get_detector_configuration()->get_op_env() : "unavailable" }
166 , m_offline_data_stream{ m_session ? m_session->get_detector_configuration()->get_offline_data_stream()
167 : "unavailable" }
168 , m_run_is_for_test_purposes{ false }
169 , m_basic_name_of_open_file{ "" }
170 , m_recorded_size{ 0 }
171 , m_uncompressed_raw_data_size{ 0 }
172 , m_previous_file_size{ 0 }
173 , m_total_file_size{ 0 }
174 , m_current_record_number{ s_unset_record_number }
175 , m_new_bytes{ 0 }
176 , m_new_objects{ 0 }
177 , m_operation_mode{ m_config_params ? m_config_params->get_mode() : "unavailable" }
178 , m_path{ m_config_params ? m_config_params->get_directory_path() : "unavailable" }
179 , m_max_file_size{ m_config_params ? m_config_params->get_max_file_size() : std::numeric_limits<size_t>::max() }
180 , m_disable_unique_suffix{ m_config_params ? m_config_params->get_disable_unique_filename_suffix() : false }
181 , m_free_space_safety_factor_for_write{ m_config_params ? m_config_params->get_free_space_safety_factor()
182 : std::numeric_limits<float>::max() }
183 {
184 TLOG_DEBUG(TLVL_BASIC) << get_name();
185
186 if (!m_config_params || !m_session) {
187 throw FileDataStoreImplBadConfiguration(ERS_HERE, get_name());
188 }
189
190 if (m_operation_mode != "one-event-per-file" && m_operation_mode != "all-per-file") {
191
192 throw InvalidOperationMode(ERS_HERE, get_name(), m_operation_mode);
193 }
194
195 // 05-Apr-2022, KAB: added warning message when the output destination
196 // is not a valid directory.
197 struct statvfs vfs_results;
198 int retval = statvfs(m_path.c_str(), &vfs_results);
199 if (retval != 0) {
200 ers::warning(InvalidOutputPath(ERS_HERE, get_name(), m_path));
201 }
202 }
203
204 virtual void open_new_file(const std::string& unique_filename) = 0;
205
206 // Getter functions which can be used by the implementation of open_new_file
207
208 std::string get_application_name() const noexcept { return m_writer_identifier; }
209
210 unsigned get_compression_level() const noexcept
211 { // NOLINT(build/unsigned)
212 return m_compression_level;
213 }
214
215 const appmodel::DataStoreConf& get_configuration() const noexcept { return *m_config_params; }
216
217 auto& get_file_handle() { return m_file_handle; }
218
219 size_t get_file_index() const noexcept { return m_file_index.load(); }
220
221 const std::string& get_offline_data_stream() const noexcept { return m_offline_data_stream; }
222
223 const std::string& get_operational_environment() const noexcept { return m_operational_environment; }
224
225 bool get_run_is_for_test_purposes() const noexcept { return m_run_is_for_test_purposes; }
226
227 daqdataformats::run_number_t get_run_number() const noexcept { return m_run_number; }
228
229 const confmodel::Session& get_session() const noexcept { return *m_session; }
230
237 void write(const daqdataformats::TriggerRecord& tr) override
238 {
239 size_t tr_size = tr.get_total_size_bytes();
240
241 throw_if_insufficient_space_for_object(tr_size, "trigger record");
242
243 increment_file_index_if_needed(tr_size, tr.get_header_ref().get_trigger_number(), m_current_record_number);
244
245 m_current_record_number = tr.get_header_ref().get_trigger_number();
246
247 std::string full_filename = get_file_name(tr.get_header_ref().get_run_number());
248
249 try {
250 open_file_if_needed(full_filename);
251 } catch (std::exception const& excpt) {
252 throw FileOperationProblem(ERS_HERE, get_name(), full_filename, excpt);
253 } catch (...) { // NOLINT(runtime/exceptions)
254 // NOLINT here because we *ARE* re-throwing the exception!
255 throw FileOperationProblem(ERS_HERE, get_name(), full_filename);
256 }
257
258 // write the record
259 m_file_handle->write(tr);
260 m_recorded_size = m_file_handle->get_recorded_size();
261 m_uncompressed_raw_data_size = m_file_handle->get_uncompressed_raw_data_size();
262 m_total_file_size = m_file_handle->get_total_file_size();
263
264 m_new_bytes += m_total_file_size - m_previous_file_size;
265 ++m_new_objects;
266 m_previous_file_size.store(m_total_file_size.load());
267 }
268
276 void write(const daqdataformats::TimeSlice& ts) override
277 {
278 size_t ts_size = ts.get_total_size_bytes();
279 throw_if_insufficient_space_for_object(ts_size, "time slice");
280
281 increment_file_index_if_needed(ts_size, ts.get_header().timeslice_number, m_current_record_number);
282
283 m_current_record_number = ts.get_header().timeslice_number;
284
285 std::string full_filename = get_file_name(ts.get_header().run_number);
286
287 try {
288 open_file_if_needed(full_filename);
289 } catch (std::exception const& excpt) {
290 throw FileOperationProblem(ERS_HERE, get_name(), full_filename, excpt);
291 } catch (...) { // NOLINT(runtime/exceptions)
292 // NOLINT here because we *ARE* re-throwing the exception!
293 throw FileOperationProblem(ERS_HERE, get_name(), full_filename);
294 }
295
296 // write the record
297 try {
298
299 if (m_file_handle->timeslice_already_exists(ts)) {
300 throw TimeSliceAlreadyExists(ERS_HERE, get_name(), ts.get_header().timeslice_number);
301 }
302
303 m_file_handle->write(ts);
304 m_recorded_size = m_file_handle->get_recorded_size();
305 m_uncompressed_raw_data_size = m_file_handle->get_uncompressed_raw_data_size();
306 m_total_file_size = m_file_handle->get_total_file_size();
307
308 } catch (TimeSliceAlreadyExists const& excpt) {
309 std::string msg = "writing a time slice to file " + m_file_handle->get_file_name();
310 throw IgnorableDataStoreProblem(ERS_HERE, get_name(), msg, excpt);
311 }
312
313 m_new_bytes += m_total_file_size - m_previous_file_size;
314 ++m_new_objects;
315 m_previous_file_size.store(m_total_file_size.load());
316 }
317
327 void prepare_for_run(daqdataformats::run_number_t run_number, bool run_is_for_test_purposes) override
328 {
329 m_run_number = run_number;
330 m_run_is_for_test_purposes = run_is_for_test_purposes;
331
332 struct statvfs vfs_results;
333 TLOG_DEBUG(TLVL_BASIC) << get_name() << ": Preparing to get the statvfs results for path: \"" << m_path << "\"";
334
335 int retval = statvfs(m_path.c_str(), &vfs_results);
336 TLOG_DEBUG(TLVL_BASIC) << get_name() << ": statvfs return code is " << retval;
337 if (retval != 0) {
338 throw InvalidOutputPath(ERS_HERE, get_name(), m_path);
339 }
340
341 size_t free_space = vfs_results.f_bsize * vfs_results.f_bavail;
342 TLOG_DEBUG(TLVL_BASIC) << get_name() << ": Free space on disk with path \"" << m_path << "\" is " << free_space
343 << " bytes. This will be compared with the maximum size of a single file ("
344 << m_max_file_size << ") as a simple test to see if there is enough free space.";
345 if (free_space < m_max_file_size) {
346 throw InsufficientDiskSpace(
347 ERS_HERE, get_name(), m_path, free_space, m_max_file_size, "the configured maximum size of a single file");
348 }
349
350 m_file_index = 0;
351 m_recorded_size = 0;
352 m_uncompressed_raw_data_size = 0;
353 m_current_record_number = s_unset_record_number;
354 }
355
363 void finish_with_run(daqdataformats::run_number_t /*run_number*/) override
364 {
365 if (m_file_handle) {
366 std::string open_filename = m_file_handle->get_file_name();
367 try {
368 m_file_handle.reset();
369 m_run_number = 0;
370 } catch (std::exception const& excpt) {
371 m_run_number = 0;
372 throw FileOperationProblem(ERS_HERE, get_name(), open_filename, excpt);
373 } catch (...) { // NOLINT(runtime/exceptions)
374 m_run_number = 0;
375 // NOLINT here because we *ARE* re-throwing the exception!
376 throw FileOperationProblem(ERS_HERE, get_name(), open_filename);
377 }
378 }
379 }
380
381 FileDataStoreImpl(const FileDataStoreImpl&) = delete;
382 FileDataStoreImpl& operator=(const FileDataStoreImpl&) = delete;
383 FileDataStoreImpl(FileDataStoreImpl&&) = delete;
384 FileDataStoreImpl& operator=(FileDataStoreImpl&&) = delete;
385
386protected:
387 void generate_opmon_data() override
388 {
389 opmon::FileDataStoreImplInfo info;
390
391 info.set_new_bytes_output(m_new_bytes.exchange(0));
392 info.set_new_written_object(m_new_objects.exchange(0));
393 info.set_bytes_in_file(m_total_file_size.load());
394 info.set_written_files(m_file_index.load());
395 publish(std::move(info), { { "path", m_path } });
396 }
397
398private:
399 // Translates various available parameters (directory path, file extension, etc.) into the appropriate filename.
400 std::string get_file_name(daqdataformats::run_number_t run_number) const;
401
402 // Throws a RetryableDataStoreProblem if (size of the object)*(free space safety factor) exceeds available space in
403 // m_path
404 void throw_if_insufficient_space_for_object(size_t obj_size, const std::string& obj_name) const;
405
406 // Available space in the path whose name is passed to get_free_space
407 size_t get_free_space(const std::string& the_path) const;
408
409 // Check if a new file should be opened for the record and increment m_file_index if so
410 void increment_file_index_if_needed(const size_t size_of_object_to_write,
411 const size_t object_record_number,
412 const size_t current_record_number)
413 {
414 float compression_factor{ 1.0 };
415
416 if (m_compression_level != 0 && m_recorded_size != 0) {
417 // Without compression, the uncompressed raw data size is approximately the total file size, so it
418 // serves as an approximation of what would have been written without compression
419 compression_factor =
420 static_cast<float>(m_file_handle->get_uncompressed_raw_data_size()) / m_file_handle->get_total_file_size();
421 }
422
423 float size_of_next_write = size_of_object_to_write / compression_factor;
424
425 if ((m_total_file_size + size_of_next_write) > m_max_file_size && m_recorded_size > 0) {
426 ++m_file_index;
427 m_recorded_size = 0;
428 m_uncompressed_raw_data_size = 0;
429 m_previous_file_size.store(0);
430 return;
431 }
432
433 // JCF, 06-27-2026: TODO: probably need to reset m_recorded_size, etc., as is done right above
434 if (m_operation_mode == "one-event-per-file" && current_record_number != s_unset_record_number &&
435 current_record_number != object_record_number) {
436 ++m_file_index;
437 return;
438 }
439 }
440
441 void open_file_if_needed(const std::string& file_name)
442 {
443
444 if (!m_file_handle || m_basic_name_of_open_file.compare(file_name)) {
445
446 // close an existing open file
447 if (m_file_handle) {
448 std::string open_filename = m_file_handle->get_file_name();
449 try {
450 m_file_handle.reset();
451 } catch (std::exception const& excpt) {
452 throw FileOperationProblem(ERS_HERE, get_name(), open_filename, excpt);
453 } catch (...) { // NOLINT(runtime/exceptions)
454 // NOLINT here because we *ARE* re-throwing the exception!
455 throw FileOperationProblem(ERS_HERE, get_name(), open_filename);
456 }
457 }
458
459 // 04-Feb-2021, KAB: adding unique substrings to the filename
460 // 05-Feb-2026, KAB: moved this block of code *after* the block that closes an
461 // existing open file. When the two blocks were executed in the opposite order,
462 // it was possible that the closing of the currently-open file could take a
463 // non-trivial amount of time, and then the timestamp in the filename (determined
464 // here) and the timestamp in the creation_timestamp HDF5 file Attribute
465 // (determined inside the HDF5RawDataFile constructor) could disagree.
466 std::string unique_filename = file_name;
467 if (!m_disable_unique_suffix) {
468 time_t now{ time(nullptr) };
469 std::string file_creation_timestamp = boost::posix_time::to_iso_string(boost::posix_time::from_time_t(now));
470 // timestamp substring
471 size_t ufn_len = unique_filename.length();
472 size_t extension_length =
473 m_file_handle->get_file_name_extension().size() + 1; // + 1 for the "." before the extension
474 if (ufn_len > extension_length + 1) { // this gives us some confidence that we have at least one character
475 // preceding the extension (e.g., x.hdf5)
476 std::string timestamp_substring = "_" + file_creation_timestamp;
477 TLOG_DEBUG(TLVL_BASIC) << get_name() << ": timestamp substring for filename: " << timestamp_substring;
478 unique_filename.insert(ufn_len - extension_length, timestamp_substring);
479 }
480 }
481
482 // opening file for the first time OR something changed in the name or the way of opening the file
483 TLOG_DEBUG(TLVL_BASIC) << get_name() << ": going to open file " << unique_filename;
484 m_basic_name_of_open_file = file_name;
485
486 open_new_file(unique_filename);
487
488 } else {
489 TLOG_DEBUG(TLVL_BASIC) << get_name() << ": Pointer file to " << m_basic_name_of_open_file
490 << " was already opened";
491 }
492 }
493
494 std::unique_ptr<FileHandleClass> m_file_handle;
495
496 daqdataformats::run_number_t m_run_number;
497
498 // Total number of generated files
499 std::atomic<size_t> m_file_index;
500 const std::string m_writer_identifier;
501
502 const appmodel::DataStoreConf* m_config_params;
503
504 const confmodel::Session* m_session;
505
506 unsigned m_compression_level; // NOLINT(build/unsigned)
507
508 const std::string m_operational_environment;
509 const std::string m_offline_data_stream;
510 bool m_run_is_for_test_purposes;
511
512 std::string m_basic_name_of_open_file;
513
514 // Size of data being written, excluding metadata
515 std::atomic<size_t> m_recorded_size;
516
517 // Theoretical, "uncompressed" size of data being written, excluding metadata
518 std::atomic<size_t> m_uncompressed_raw_data_size;
519
520 // Used for tracking the "delta" of the current write
521 std::atomic<size_t> m_previous_file_size = 0;
522
523 // Total size of the file, including raw data, metadata, and free space
524 std::atomic<size_t> m_total_file_size;
525
526 // Record number for the record that is currently being written out
527 // This is only useful for long-readout windows, in which there may
528 // be multiple calls to write()
529 size_t m_current_record_number;
530
531 // incremental written data
532 std::atomic<uint64_t> m_new_bytes; // NOLINT(build/unsigned)
533 std::atomic<uint64_t> m_new_objects; // NOLINT(build/unsigned)
534
535 const std::string m_operation_mode;
536 const std::string m_path;
537 const size_t m_max_file_size;
538 const bool m_disable_unique_suffix;
539 float m_free_space_safety_factor_for_write;
540};
541
542} // namespace dfmodules
543} // namespace dunedaq
544
546
547#endif // DFMODULES_INCLUDE_DFMODULES_FILEDATASTOREIMPL_HPP_
#define ERS_HERE
C++ Representation of a DUNE TimeSlice, consisting of a TimeSliceHeader object and a vector of pointe...
Definition TimeSlice.hpp:27
size_t get_total_size_bytes() const
Get size of timeslice from underlying TimeSliceHeader and Fragments.
Definition TimeSlice.hpp:78
TimeSliceHeader get_header() const
Get a copy of the TimeSliceHeader struct.
Definition TimeSlice.hpp:44
C++ Representation of a DUNE TriggerRecord, consisting of a TriggerRecordHeader object and a vector o...
size_t get_total_size_bytes() const
Get size of trigger record from underlying TriggerRecordHeader and Fragments.
const TriggerRecordHeader & get_header_ref() const
Get a handle to the TriggerRecordHeader.
#define TLOG_DEBUG(lvl,...)
Definition Logging.hpp:116
An ERS Issue for DataStore creation failure.
Definition DataStore.hpp:91
The DUNE-DAQ namespace.
GeneralDAQModuleIssue
Definition DAQModule.hpp:72
ERS_DECLARE_ISSUE_BASE(timinglibs, ProgressUpdate, appfwk::GeneralDAQModuleIssue, message,((std::string) name),((std::string) message)) ERS_DECLARE_ISSUE_BASE(timinglibs
void warning(const Issue &issue)
Definition ers.hpp:150
timeslice_number_t timeslice_number
Slice number of this TimeSlice within the stream.