19#include <vrs/DiskFile.h>
21#if VRS_ASYNC_DISKFILE_SUPPORTED()
29#define POSIX_AIO_SUPPORTED() (IS_MAC_PLATFORM() || IS_LINUX_PLATFORM())
31#if IS_WINDOWS_PLATFORM()
34#elif POSIX_AIO_SUPPORTED()
44#include <condition_variable>
51#include <vrs/ErrorCode.h>
52#include <vrs/VrsExport.h>
53#include <vrs/os/Platform.h>
55#define VRS_DISKFILECHUNK "AsyncDiskFileChunk"
59#if IS_WINDOWS_PLATFORM()
61using ssize_t = int64_t;
62#define O_DIRECT 0x80000000U
64struct VRS_API AsyncWindowsHandle {
65 AsyncWindowsHandle() : h_(INVALID_HANDLE_VALUE) {}
66 AsyncWindowsHandle(HANDLE h) : h_(h) {}
67 AsyncWindowsHandle(AsyncWindowsHandle&& rhs) : h_(rhs.h_) {
68 rhs.h_ = INVALID_HANDLE_VALUE;
70 AsyncWindowsHandle(AsyncWindowsHandle& rhs) : h_(rhs.h_) {}
71 AsyncWindowsHandle& operator=(AsyncWindowsHandle&& rhs) {
73 rhs.h_ = INVALID_HANDLE_VALUE;
77 bool isOpened()
const;
78 int open(
const std::string& path,
const char* modes,
int flags);
80 int pwrite(
const void* buf,
size_t count, int64_t offset,
size_t& outWriteSize);
81 int read(
void* buf,
size_t count, int64_t offset,
size_t& outReadSize);
82 int truncate(int64_t newSize);
83 int seek(int64_t pos,
int origin, int64_t& outFilepos);
86 int _readwrite(
bool readNotWrite,
void* buf,
size_t count, int64_t offset,
size_t& outSize);
89 HANDLE h_ = INVALID_HANDLE_VALUE;
92using AsyncHandle = AsyncWindowsHandle;
94struct VRS_API AsyncFileDescriptor {
95 static constexpr int INVALID_FILE_DESCRIPTOR = -1;
97 AsyncFileDescriptor() =
default;
98 explicit AsyncFileDescriptor(
int fd) : fd_(fd) {}
99 AsyncFileDescriptor(AsyncFileDescriptor&& rhs) noexcept : fd_(rhs.fd_) {
100 rhs.fd_ = INVALID_FILE_DESCRIPTOR;
102 AsyncFileDescriptor(
const AsyncFileDescriptor& rhs)
noexcept =
delete;
103 AsyncFileDescriptor& operator=(AsyncFileDescriptor&& rhs)
noexcept {
105 rhs.fd_ = INVALID_FILE_DESCRIPTOR;
108 AsyncFileDescriptor& operator=(
const AsyncFileDescriptor& rhs) =
delete;
110 bool operator==(
int fd)
const {
114 int open(
const std::string& path,
const char* modes,
int flags);
115 [[nodiscard]]
bool isOpened()
const;
116 int read(
void* ptr,
size_t bufferSize,
size_t offset,
size_t& outReadSize);
117 int truncate(int64_t newSize);
118 int seek(int64_t pos,
int origin, int64_t& outFilepos);
119 int pwrite(
const void* buf,
size_t count, off_t offset,
size_t& written);
122 int fd_ = INVALID_FILE_DESCRIPTOR;
124using AsyncHandle = AsyncFileDescriptor;
127class VRS_API AlignedBuffer {
129 void* aligned_buffer_ =
nullptr;
130 size_t capacity_ = 0;
135 AlignedBuffer(
size_t size,
size_t memalign,
size_t lenalign);
137 [[nodiscard]]
inline bool isValid()
const {
138 return aligned_buffer_ !=
nullptr;
144 static std::unique_ptr<AlignedBuffer> make(
size_t size,
size_t memalign,
size_t lenalign);
146 virtual ~AlignedBuffer();
148 [[nodiscard]]
inline size_t size()
const {
151 [[nodiscard]]
inline size_t capacity()
const {
154 [[nodiscard]]
inline bool empty()
const {
157 [[nodiscard]]
inline bool full()
const {
158 return size() == capacity();
163 [[nodiscard]]
inline void* data()
const {
164 return aligned_buffer_;
166 [[nodiscard]]
inline char* bdata()
const {
167 return reinterpret_cast<char*
>(aligned_buffer_);
173 [[nodiscard]]
bool add(
const void* buffer,
size_t size,
size_t& outCopiedSize);
177#if IS_WINDOWS_PLATFORM()
178struct VRS_API AsyncOVERLAPPED {
185class VRS_API AsyncBuffer :
public AlignedBuffer {
187 using complete_write_callback = std::function<void(ssize_t io_return,
int io_errno)>;
191 static std::unique_ptr<AsyncBuffer> make(
size_t size,
size_t memalign,
size_t lenalign);
193 ~AsyncBuffer()
override =
default;
195 void complete_write(ssize_t io_return,
int io_errno);
197 start_write(
const AsyncHandle& file, int64_t offset, complete_write_callback on_complete);
200 AsyncBuffer(
size_t size,
size_t memalign,
size_t lenalign)
201 : AlignedBuffer(size, memalign, lenalign) {}
204#if IS_WINDOWS_PLATFORM()
206 static void CompletedWriteRoutine(DWORD dwErr, DWORD cbBytesWritten, LPOVERLAPPED lpOverlapped);
207#elif POSIX_AIO_SUPPORTED()
208 struct aiocb aiocb_{};
209 static void SigEvNotifyFunction(
union sigval val);
211 complete_write_callback on_complete_ =
nullptr;
214class VRS_API AsyncDiskFileChunk {
216 AsyncDiskFileChunk() =
default;
217 AsyncDiskFileChunk(std::string path, int64_t offset, int64_t size)
218 : path_{std::move(path)}, offset_{offset}, size_{size} {}
219 AsyncDiskFileChunk(AsyncDiskFileChunk&& other)
noexcept;
222 AsyncDiskFileChunk(
const AsyncDiskFileChunk& other)
noexcept =
delete;
223 AsyncDiskFileChunk& operator=(
const AsyncDiskFileChunk& other)
noexcept =
delete;
224 AsyncDiskFileChunk& operator=(AsyncDiskFileChunk&& rhs)
noexcept =
delete;
226 ~AsyncDiskFileChunk();
228 int create(
const std::string& newpath,
const FileSpec::Extras& options);
229 int open(
bool readOnly,
const FileSpec::Extras& options);
232 [[nodiscard]]
bool eof()
const;
234 int write(
const void* buffer,
size_t count,
size_t& outWrittenSize);
235 void setSize(int64_t newSize);
237 int truncate(int64_t newSize);
238 int read(
void* buffer,
size_t count,
size_t& outReadSize);
239 [[nodiscard]] int64_t getSize()
const;
240 [[nodiscard]]
bool contains(int64_t fileOffset)
const;
241 int tell(int64_t& outFilepos)
const;
242 int seek(int64_t pos,
int origin);
243 [[nodiscard]]
const std::string& getPath()
const;
244 void setOffset(int64_t newOffset);
245 [[nodiscard]] int64_t getOffset()
const;
247 enum class IoEngine {
255 AsyncBuffer* buffer_;
258 const AsyncHandle& file_;
260 AsyncBuffer::complete_write_callback callback_;
265 AsyncBuffer::complete_write_callback callback)
266 : buffer_(buffer), file_(file), offset_(offset), callback_(std::move(callback)) {}
269 int flushWriteBuffer();
270 int ensureOpenNonDirect();
271 int ensureOpenDirect();
272 int ensureOpen_(
int requested_flags);
273 void complete_write(AsyncBuffer* buffer, ssize_t io_return,
int io_errno);
274 AsyncBuffer* get_free_buffer_locked(std::unique_lock<std::mutex>& lock);
275 AsyncBuffer* get_free_buffer();
276 void free_buffer(AsyncBuffer*& buffer);
277 void free_buffer_locked(std::unique_lock<std::mutex>& lock, AsyncBuffer*& buffer);
279 void pump_buffers_locked();
280 int alloc_write_buffers();
281 int free_write_buffers();
282 int init_parameters(
const FileSpec::Extras& options);
290 int64_t file_position_ = 0;
292 const char* file_mode_ =
nullptr;
295 int current_flags_ = 0;
297 int supported_flags_ = 0;
302 std::mutex buffers_mutex_;
304 std::condition_variable buffer_freed_cv_;
306 std::vector<AsyncBuffer*> buffers_free_;
308 std::deque<QueuedWrite> buffers_queued_;
313 size_t buffers_writing_ = 0;
316 std::vector<std::unique_ptr<AsyncBuffer>> buffers_;
319 AsyncBuffer* current_buffer_ =
nullptr;
323 std::atomic<int> async_error_ = SUCCESS;
327 IoEngine ioengine_ = IoEngine::AIO;
328 bool use_directio_ =
true;
330 size_t num_buffers_ = 0;
332 size_t buffer_size_ = 0;
336 size_t offset_align_ = 0;
338 size_t mem_align_ = 0;
int write(const string &cacheFile, const set< StreamId > &streamIds, const map< string, string > &fileTags, const map< StreamId, StreamTags > &streamTags, const vector< IndexRecord::RecordInfo > &recordIndex, bool fileHasIndex)
Definition FileDetailsCache.cpp:225
int read(const string &cacheFile, set< StreamId > &outStreamIds, map< string, string > &outFileTags, map< StreamId, StreamTags > &outStreamTags, vector< IndexRecord::RecordInfo > &outRecordIndex, bool &outFileHasIndex)
Definition FileDetailsCache.cpp:266
Definition Compressor.cpp:113