DUNE-DAQ
DUNE Trigger and Data Acquisition software
Loading...
Searching...
No Matches
BufferedFileWriter.hpp
Go to the documentation of this file.
1
11#ifndef DATAHANDLINGLIBS_INCLUDE_DATAHANDLINGLIBS_UTILS_BUFFEREDFILEWRITER_HPP_
12#define DATAHANDLINGLIBS_INCLUDE_DATAHANDLINGLIBS_UTILS_BUFFEREDFILEWRITER_HPP_
13
16
17#include "logging/Logging.hpp"
18
19#include <boost/align/aligned_allocator.hpp>
20#include <boost/iostreams/device/file_descriptor.hpp>
21#include <boost/iostreams/filter/lzma.hpp>
22#include <boost/iostreams/filter/zlib.hpp>
23#include <boost/iostreams/filter/zstd.hpp>
24#include <boost/iostreams/filtering_stream.hpp>
25#include <boost/iostreams/stream.hpp>
26#include <boost/iostreams/stream_buffer.hpp>
27
28#include <fcntl.h>
29#include <fstream>
30#include <iostream>
31#include <limits>
32#include <string>
33#include <unistd.h>
34
36
37namespace dunedaq {
38namespace datahandlinglibs {
46template<size_t Alignment = 4096>
48{
49 using io_sink_t = boost::iostreams::file_descriptor_sink;
50 using aligned_allocator_t = boost::alignment::aligned_allocator<io_sink_t::char_type, Alignment>;
52 boost::iostreams::filtering_stream<boost::iostreams::output, char, std::char_traits<char>, aligned_allocator_t>;
53
54public:
67 BufferedFileWriter(std::string filename,
68 size_t buffer_size,
69 std::string compression_algorithm = "None",
70 bool use_o_direct = true)
71 {
72 open(filename, buffer_size, compression_algorithm, use_o_direct);
73 }
74
79
84 {
85 if (m_is_open)
86 close();
87 }
88
93
106 void open(std::string filename,
107 size_t buffer_size,
108 std::string compression_algorithm = "None",
109 bool use_o_direct = true)
110 {
111 m_use_o_direct = use_o_direct;
112 if (m_is_open) {
113 close();
114 }
115
117 m_buffer_size = buffer_size;
118 m_compression_algorithm = compression_algorithm;
119 auto oflag = O_CREAT | O_WRONLY;
120 if (m_use_o_direct) {
121 oflag = oflag | O_DIRECT;
122 }
123
124 m_fd = ::open(m_filename.c_str(), oflag, 0644);
125 if (m_fd == -1) {
126 throw BufferedReaderWriterCannotOpenFile(ERS_HERE, m_filename);
127 }
128
129 m_sink = io_sink_t(m_fd, boost::iostreams::file_descriptor_flags::close_handle);
130 if (m_compression_algorithm == "zstd") {
131 TLOG_DEBUG(TLVL_WORK_STEPS) << "Using zstd compression" << std::endl;
132 m_output_stream.push(boost::iostreams::zstd_compressor(boost::iostreams::zstd::best_speed));
133 } else if (m_compression_algorithm == "lzma") {
134 TLOG_DEBUG(TLVL_WORK_STEPS) << "Using lzma compression" << std::endl;
135 m_output_stream.push(boost::iostreams::lzma_compressor(boost::iostreams::lzma::best_speed));
136 } else if (m_compression_algorithm == "zlib") {
137 TLOG_DEBUG(TLVL_WORK_STEPS) << "Using zlib compression" << std::endl;
138 m_output_stream.push(boost::iostreams::zlib_compressor(boost::iostreams::zlib::best_speed));
139 } else if (m_compression_algorithm == "None") {
140 TLOG_DEBUG(TLVL_WORK_STEPS) << "Running without compression" << std::endl;
141 } else {
143 "Non-recognized compression algorithm: " + m_compression_algorithm);
144 }
145
147 m_is_open = true;
148 }
149
154 bool is_open() const { return m_is_open; }
155
161 bool write(const char* memory, const size_t size)
162 {
163 if (!m_is_open)
164 return false;
165 m_output_stream.write(memory, size); // NOLINT
166 return !m_output_stream.bad();
167 }
168
172 void close()
173 {
174 // Set the file descriptor to not use O_DIRECT. This is necessary because the write size has to be aligned for
175 // O_DIRECT to succeed. This is not guaranteed for the data that remains in the buffer.
176 fcntl(m_fd, F_SETFL, O_CREAT | O_WRONLY);
177 m_output_stream.reset();
178 m_is_open = false;
179 }
180
185 void flush()
186 {
187 // Set the file descriptor to not use O_DIRECT. This is necessary because the write size has to be aligned for
188 // O_DIRECT to succeed. This is not guaranteed for the data that remains in the buffer.
189 fcntl(m_fd, F_SETFL, O_CREAT | O_WRONLY);
190 // This does not flush the compressor as it is not flushable
191 m_output_stream.flush();
192 // Activate O_DIRECT again
193 auto oflag = O_CREAT | O_WRONLY;
194 if (m_use_o_direct) {
195 oflag = oflag | O_DIRECT;
196 }
197 fcntl(m_fd, F_SETFL, oflag);
198 }
199
200private:
201 // Config parameters
202 std::string m_filename;
205
206 // Internals
207 int m_fd;
210 bool m_is_open = false;
211 bool m_use_o_direct = true;
212};
213
214} // namespace datahandlinglibs
215} // namespace dunedaq
216
217#endif // DATAHANDLINGLIBS_INCLUDE_DATAHANDLINGLIBS_UTILS_BUFFEREDFILEWRITER_HPP_
#define ERS_HERE
BufferedFileWriter(const BufferedFileWriter &)=delete
BufferedFileWriter is not copy-constructible.
boost::alignment::aligned_allocator< io_sink_t::char_type, Alignment > aligned_allocator_t
void open(std::string filename, size_t buffer_size, std::string compression_algorithm="None", bool use_o_direct=true)
BufferedFileWriter & operator=(BufferedFileWriter &&)=delete
BufferedFileWriter is not move-assignable.
boost::iostreams::filtering_stream< boost::iostreams::output, char, std::char_traits< char >, aligned_allocator_t > filtering_ostream_t
bool write(const char *memory, const size_t size)
boost::iostreams::file_descriptor_sink io_sink_t
BufferedFileWriter(std::string filename, size_t buffer_size, std::string compression_algorithm="None", bool use_o_direct=true)
BufferedFileWriter(BufferedFileWriter &&)=delete
BufferedFileWriter is not move-constructible.
BufferedFileWriter & operator=(const BufferedFileWriter &)=delete
BufferedFileWriter is not copy-assginable.
#define TLOG_DEBUG(lvl,...)
Definition Logging.hpp:116
The DUNE-DAQ namespace.
FELIX Initialization std::string initerror FELIX queue timed std::string queuename Unexpected chunk size
SourceID[" << sourceid << "] Command daqdataformats::SourceID Readout Initialization std::string initerror BufferedReaderWriterConfigurationError