33#ifndef DFMODULES_INCLUDE_DFMODULES_FILEDATASTOREIMPL_HPP_
34#define DFMODULES_INCLUDE_DFMODULES_FILEDATASTOREIMPL_HPP_
48#include "boost/date_time/posix_time/posix_time.hpp"
49#include "boost/lexical_cast.hpp"
57#include <sys/statvfs.h>
68 FileDataStoreImplBadConfiguration,
69 appfwk::GeneralDAQModuleIssue,
70 "Construction of the FileDataStoreImpl base class failed due to faulty configuration",
71 ((std::string)name), )
76 "Selected operation mode \"" << selected_operation
77 <<
"\" is NOT supported. Please update the configuration file.",
79 ((
std::
string)selected_operation))
84 "A problem was encountered when opening or closing file \"" << filename <<
"\"",
86 ((
std::
string)filename))
91 "The specified output destination, \"" << output_path
92 <<
"\", is not a valid file system path on this server.",
94 ((
std::
string)output_path))
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 <<
".",
103 ((
std::
string)path)((
size_t)free_bytes)((
size_t)needed_bytes)((
std::
string)criteria))
106 TimeSliceAlreadyExists,
108 "The TimeSlice record for timeslice #" << timeslice_number <<
" already exists.",
117concept FileHandleConcept =
requires(T file_handle,
118 const T const_file_handle,
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>;
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>;
136template<FileHandleConcept FileHandleClass>
137class FileDataStoreImpl :
public DataStore
146 static constexpr size_t s_unset_record_number{ std::numeric_limits<size_t>::max() };
153 explicit FileDataStoreImpl(std::string
const& name,
154 std::shared_ptr<appfwk::ConfigurationManager> mcfg,
155 std::string
const& writer_name)
157 , m_file_handle{
nullptr }
160 , m_writer_identifier{ writer_name }
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) }
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()
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 }
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() }
186 if (!m_config_params || !m_session) {
187 throw FileDataStoreImplBadConfiguration(
ERS_HERE, get_name());
190 if (m_operation_mode !=
"one-event-per-file" && m_operation_mode !=
"all-per-file") {
197 struct statvfs vfs_results;
198 int retval = statvfs(m_path.c_str(), &vfs_results);
204 virtual void open_new_file(
const std::string& unique_filename) = 0;
208 std::string get_application_name()
const noexcept {
return m_writer_identifier; }
210 unsigned get_compression_level()
const noexcept
212 return m_compression_level;
217 auto& get_file_handle() {
return m_file_handle; }
219 size_t get_file_index()
const noexcept {
return m_file_index.load(); }
221 const std::string& get_offline_data_stream()
const noexcept {
return m_offline_data_stream; }
223 const std::string& get_operational_environment()
const noexcept {
return m_operational_environment; }
225 bool get_run_is_for_test_purposes()
const noexcept {
return m_run_is_for_test_purposes; }
241 throw_if_insufficient_space_for_object(tr_size,
"trigger record");
250 open_file_if_needed(full_filename);
251 }
catch (std::exception
const& excpt) {
252 throw FileOperationProblem(
ERS_HERE, get_name(), full_filename, excpt);
255 throw FileOperationProblem(
ERS_HERE, get_name(), full_filename);
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();
264 m_new_bytes += m_total_file_size - m_previous_file_size;
266 m_previous_file_size.store(m_total_file_size.load());
279 throw_if_insufficient_space_for_object(ts_size,
"time slice");
288 open_file_if_needed(full_filename);
289 }
catch (std::exception
const& excpt) {
290 throw FileOperationProblem(
ERS_HERE, get_name(), full_filename, excpt);
293 throw FileOperationProblem(
ERS_HERE, get_name(), full_filename);
299 if (m_file_handle->timeslice_already_exists(ts)) {
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();
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);
313 m_new_bytes += m_total_file_size - m_previous_file_size;
315 m_previous_file_size.store(m_total_file_size.load());
329 m_run_number = run_number;
330 m_run_is_for_test_purposes = run_is_for_test_purposes;
332 struct statvfs vfs_results;
333 TLOG_DEBUG(TLVL_BASIC) << get_name() <<
": Preparing to get the statvfs results for path: \"" << m_path <<
"\"";
335 int retval = statvfs(m_path.c_str(), &vfs_results);
336 TLOG_DEBUG(TLVL_BASIC) << get_name() <<
": statvfs return code is " << retval;
338 throw InvalidOutputPath(
ERS_HERE, get_name(), m_path);
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");
352 m_uncompressed_raw_data_size = 0;
353 m_current_record_number = s_unset_record_number;
366 std::string open_filename = m_file_handle->get_file_name();
368 m_file_handle.reset();
370 }
catch (std::exception
const& excpt) {
372 throw FileOperationProblem(
ERS_HERE, get_name(), open_filename, excpt);
376 throw FileOperationProblem(
ERS_HERE, get_name(), open_filename);
381 FileDataStoreImpl(
const FileDataStoreImpl&) =
delete;
382 FileDataStoreImpl& operator=(
const FileDataStoreImpl&) =
delete;
383 FileDataStoreImpl(FileDataStoreImpl&&) =
delete;
384 FileDataStoreImpl& operator=(FileDataStoreImpl&&) =
delete;
387 void generate_opmon_data()
override
389 opmon::FileDataStoreImplInfo info;
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 } });
404 void throw_if_insufficient_space_for_object(
size_t obj_size,
const std::string& obj_name)
const;
407 size_t get_free_space(
const std::string& the_path)
const;
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)
414 float compression_factor{ 1.0 };
416 if (m_compression_level != 0 && m_recorded_size != 0) {
420 static_cast<float>(m_file_handle->get_uncompressed_raw_data_size()) / m_file_handle->get_total_file_size();
423 float size_of_next_write = size_of_object_to_write / compression_factor;
425 if ((m_total_file_size + size_of_next_write) > m_max_file_size && m_recorded_size > 0) {
428 m_uncompressed_raw_data_size = 0;
429 m_previous_file_size.store(0);
434 if (m_operation_mode ==
"one-event-per-file" && current_record_number != s_unset_record_number &&
435 current_record_number != object_record_number) {
441 void open_file_if_needed(
const std::string& file_name)
444 if (!m_file_handle || m_basic_name_of_open_file.compare(file_name)) {
448 std::string open_filename = m_file_handle->get_file_name();
450 m_file_handle.reset();
451 }
catch (std::exception
const& excpt) {
452 throw FileOperationProblem(
ERS_HERE, get_name(), open_filename, excpt);
455 throw FileOperationProblem(
ERS_HERE, get_name(), open_filename);
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));
471 size_t ufn_len = unique_filename.length();
472 size_t extension_length =
473 m_file_handle->get_file_name_extension().size() + 1;
474 if (ufn_len > extension_length + 1) {
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);
483 TLOG_DEBUG(TLVL_BASIC) << get_name() <<
": going to open file " << unique_filename;
484 m_basic_name_of_open_file = file_name;
486 open_new_file(unique_filename);
489 TLOG_DEBUG(TLVL_BASIC) << get_name() <<
": Pointer file to " << m_basic_name_of_open_file
490 <<
" was already opened";
494 std::unique_ptr<FileHandleClass> m_file_handle;
499 std::atomic<size_t> m_file_index;
500 const std::string m_writer_identifier;
506 unsigned m_compression_level;
508 const std::string m_operational_environment;
509 const std::string m_offline_data_stream;
510 bool m_run_is_for_test_purposes;
512 std::string m_basic_name_of_open_file;
515 std::atomic<size_t> m_recorded_size;
518 std::atomic<size_t> m_uncompressed_raw_data_size;
521 std::atomic<size_t> m_previous_file_size = 0;
524 std::atomic<size_t> m_total_file_size;
529 size_t m_current_record_number;
532 std::atomic<uint64_t> m_new_bytes;
533 std::atomic<uint64_t> m_new_objects;
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;
#define TLOG_DEBUG(lvl,...)
An ERS Issue for DataStore creation failure.
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)