diff --git a/.gitignore b/.gitignore index f9cce9bd6..cc9e88d10 100644 --- a/.gitignore +++ b/.gitignore @@ -33,6 +33,7 @@ python/tsfile/*so* python/tsfile/*dll* python/tsfile/*dylib* python/tsfile/*.h +!python/tsfile/python_random_access_read_file.h python/tsfile/*.cpp python/tsfile/**/*.cpp python/tsfile/**/*so* diff --git a/cpp/README-zh.md b/cpp/README-zh.md index 82b2acf34..e4c70bfc0 100644 --- a/cpp/README-zh.md +++ b/cpp/README-zh.md @@ -175,6 +175,9 @@ storage::set_write_thread_count(4); ### 本地文件读取后端 +`LocalRandomAccessReadFile` 是 `RandomAccessReadFile` 的本地文件实现, +对应的公开头文件为 `file/local_random_access_read_file.h`。 + Reader 可以为本地文件选择内存映射 I/O 或传统的定位读取路径。配置会在 reader 打开文件时确定,因此修改配置不会影响已经打开的 reader。 diff --git a/cpp/README.md b/cpp/README.md index 96d5795d5..de7ba9cd6 100644 --- a/cpp/README.md +++ b/cpp/README.md @@ -350,6 +350,9 @@ By default, parallel write is enabled when the machine has more than one CPU cor ### Local File Read Backend +`LocalRandomAccessReadFile` implements `RandomAccessReadFile` for local files. Its +public header is `file/local_random_access_read_file.h`. + Readers can use memory-mapped I/O or the traditional positioned-read path for local files. The setting is captured when a reader opens a file, so changing it does not affect readers that are already open. diff --git a/cpp/bench_mark/bench_mark_src/read_backend_benchmark.cc b/cpp/bench_mark/bench_mark_src/read_backend_benchmark.cc index b2bb35312..9810972af 100644 --- a/cpp/bench_mark/bench_mark_src/read_backend_benchmark.cc +++ b/cpp/bench_mark/bench_mark_src/read_backend_benchmark.cc @@ -30,7 +30,7 @@ #include #include "common/global.h" -#include "file/read_file.h" +#include "file/local_random_access_read_file.h" #include "file/tsfile_io_reader.h" #include "reader/result_set.h" #include "reader/tsfile_reader.h" @@ -71,7 +71,8 @@ struct QueryPlan { int row_count; }; -typedef std::vector> OpenFiles; +typedef std::vector> + OpenFiles; uint64_t update_checksum(uint64_t checksum, const std::vector& buffer, int32_t read_len) { @@ -201,7 +202,8 @@ int build_query_plan(storage::TsFileReader& reader, QueryPlan& plan, bool open_files(const std::vector& paths, OpenFiles& files) { files.clear(); for (size_t i = 0; i < paths.size(); ++i) { - std::unique_ptr file(new storage::ReadFile()); + std::unique_ptr file( + new storage::LocalRandomAccessReadFile()); const int ret = file->open(paths[i]); if (ret != common::E_OK) { std::cerr << "failed to open " << paths[i] << ": error " << ret @@ -219,7 +221,7 @@ Result run_sequential(const OpenFiles& files) { const std::chrono::steady_clock::time_point start = std::chrono::steady_clock::now(); for (size_t file_index = 0; file_index < files.size(); ++file_index) { - storage::ReadFile& file = *files[file_index]; + storage::LocalRandomAccessReadFile& file = *files[file_index]; for (int64_t offset = 0; offset < file.file_size(); offset += kSequentialBlockSize) { int32_t read_len = 0; @@ -251,7 +253,7 @@ Result run_random(const OpenFiles& files, int32_t requested_block_size, const std::chrono::steady_clock::time_point start = std::chrono::steady_clock::now(); for (size_t file_index = 0; file_index < files.size(); ++file_index) { - storage::ReadFile& file = *files[file_index]; + storage::LocalRandomAccessReadFile& file = *files[file_index]; const uint64_t file_size = static_cast(file.file_size()); const int32_t block_size = static_cast(std::min( file_size, static_cast(requested_block_size))); @@ -411,7 +413,7 @@ Result run_concurrent(const std::vector& paths) { workers.reserve(paths.size()); for (size_t file_index = 0; file_index < paths.size(); ++file_index) { workers.push_back(std::thread([&, file_index]() { - storage::ReadFile file; + storage::LocalRandomAccessReadFile file; if (file.open(paths[file_index]) != common::E_OK) { per_file[file_index].success = false; return; diff --git a/cpp/src/common/config/config.h b/cpp/src/common/config/config.h index 03f2f02b8..1293ebecb 100644 --- a/cpp/src/common/config/config.h +++ b/cpp/src/common/config/config.h @@ -25,7 +25,7 @@ namespace common { -/** Backend used by local ReadFile instances. */ +/** Backend used by local LocalRandomAccessReadFile instances. */ enum class FileReadBackend : uint8_t { AUTO = 0, MMAP = 1, diff --git a/cpp/src/common/global.h b/cpp/src/common/global.h index 5ef7d6d40..f1f6df19c 100644 --- a/cpp/src/common/global.h +++ b/cpp/src/common/global.h @@ -200,9 +200,9 @@ FORCE_INLINE bool get_parallel_write_enabled() { } // Select the backend used by subsequently opened local files. Existing -// ReadFile instances retain the backend selected when they were opened. This -// setting deliberately lives outside exported ConfigValue so adding it does -// not change that public data structure's ABI. +// LocalRandomAccessReadFile instances retain the backend selected when opened. +// This setting deliberately lives outside exported ConfigValue so adding it +// does not change that public data structure's ABI. extern int set_file_read_backend(FileReadBackend backend); extern FileReadBackend get_file_read_backend(); diff --git a/cpp/src/cwrapper/tsfile_cwrapper.cc b/cpp/src/cwrapper/tsfile_cwrapper.cc index c7ac821ef..de17ba49b 100644 --- a/cpp/src/cwrapper/tsfile_cwrapper.cc +++ b/cpp/src/cwrapper/tsfile_cwrapper.cc @@ -874,8 +874,30 @@ TableSchema tsfile_reader_get_table_schema(TsFileReader reader, TableSchema* tsfile_reader_get_all_table_schemas(TsFileReader reader, uint32_t* size) { + ERRNO error_code = common::E_OK; + return tsfile_reader_get_all_table_schemas_with_error(reader, size, + &error_code); +} + +TableSchema* tsfile_reader_get_all_table_schemas_with_error(TsFileReader reader, + uint32_t* size, + ERRNO* error_code) { + if (size != nullptr) { + *size = 0; + } + if (error_code == nullptr) { + return nullptr; + } + *error_code = common::E_INVALID_ARG; + if (reader == nullptr || size == nullptr) { + return nullptr; + } auto* r = static_cast(reader); - auto table_schemas = r->get_all_table_schemas(); + std::vector> table_schemas; + *error_code = r->get_all_table_schemas(table_schemas); + if (*error_code != common::E_OK || table_schemas.empty()) { + return nullptr; + } size_t table_num = table_schemas.size(); TableSchema* ret = static_cast(malloc(sizeof(TableSchema) * table_num)); @@ -1594,7 +1616,11 @@ ERRNO tsfile_reader_get_all_devices(TsFileReader reader, DeviceID** out_devices, *out_devices = nullptr; *out_length = 0; auto* r = static_cast(reader); - const auto ids = r->get_all_devices(); + std::vector> ids; + const int ret = r->get_all_devices(ids); + if (ret != common::E_OK) { + return ret; + } if (ids.empty()) { return common::E_OK; } diff --git a/cpp/src/cwrapper/tsfile_cwrapper.h b/cpp/src/cwrapper/tsfile_cwrapper.h index 282f918b8..8c0d4db84 100644 --- a/cpp/src/cwrapper/tsfile_cwrapper.h +++ b/cpp/src/cwrapper/tsfile_cwrapper.h @@ -980,10 +980,23 @@ TableSchema tsfile_reader_get_table_schema(TsFileReader reader, * * @return TableSchema, contains table and column info. * @note Caller should call free_table_schema and free to free the ptr. + * @note Use tsfile_reader_get_all_table_schemas_with_error to distinguish + * metadata read failures from an empty schema list. */ TableSchema* tsfile_reader_get_all_table_schemas(TsFileReader reader, uint32_t* size); +/** + * @brief Gets all table schemas and reports metadata read failures. + * @return Schema array, or NULL when there are no tables or on error. Check + * error_code to distinguish these cases; size is zero on error. + * @note Caller must free each schema with free_table_schema, then free the + * array. + */ +TableSchema* tsfile_reader_get_all_table_schemas_with_error(TsFileReader reader, + uint32_t* size, + ERRNO* error_code); + /** * @brief Gets all timeseries schema in the tsfile. * diff --git a/cpp/src/file/read_file.cc b/cpp/src/file/local_random_access_read_file.cc similarity index 84% rename from cpp/src/file/read_file.cc rename to cpp/src/file/local_random_access_read_file.cc index 28a33850f..d617202a0 100644 --- a/cpp/src/file/read_file.cc +++ b/cpp/src/file/local_random_access_read_file.cc @@ -17,7 +17,7 @@ * under the License. */ -#include "file/read_file.h" +#include "file/local_random_access_read_file.h" #include #include @@ -59,7 +59,7 @@ uint64_t generation_hash(uint64_t size, int64_t mtime_ns) { } } // namespace -ReadFile::ReadFile() +LocalRandomAccessReadFile::LocalRandomAccessReadFile() : file_path_(), fd_(-1), file_size_(-1), @@ -79,7 +79,8 @@ ReadFile::ReadFile() { } -int ReadFile::generation(uint64_t& size, uint64_t& fingerprint) const { +int LocalRandomAccessReadFile::generation(uint64_t& size, + uint64_t& fingerprint) const { if (!is_opened()) { return E_FILE_READ_ERR; } @@ -93,8 +94,8 @@ int ReadFile::generation(uint64_t& size, uint64_t& fingerprint) const { return get_file_generation_from_descriptor(size, fingerprint); } -int ReadFile::get_file_generation_from_descriptor(uint64_t& size, - uint64_t& fingerprint) const { +int LocalRandomAccessReadFile::get_file_generation_from_descriptor( + uint64_t& size, uint64_t& fingerprint) const { int64_t mtime_ns = 0; #ifdef _WIN32 if (fd_ < 0) { @@ -139,7 +140,7 @@ int ReadFile::get_file_generation_from_descriptor(uint64_t& size, return E_OK; } -void ReadFile::close() { +void LocalRandomAccessReadFile::close() { unmap_file(); if (fd_ >= 0) { ::close(fd_); @@ -155,7 +156,7 @@ void ReadFile::close() { #endif } -int ReadFile::open(const std::string& file_path) { +int LocalRandomAccessReadFile::open(const std::string& file_path) { int ret = E_OK; close(); file_path_ = file_path; @@ -238,7 +239,7 @@ int ReadFile::open(const std::string& file_path) { return ret; } -int ReadFile::get_file_size(int64_t& file_size) { +int LocalRandomAccessReadFile::get_file_size(int64_t& file_size) { #ifdef _WIN32 struct __stat64 s; if (_fstat64(fd_, &s) < 0) { @@ -258,7 +259,7 @@ int ReadFile::get_file_size(int64_t& file_size) { return E_OK; } -int ReadFile::map_file() { +int LocalRandomAccessReadFile::map_file() { DBUG_EXECUTE_IF("read_file_mmap_fail", return E_FILE_MAP_ERR;); DBUG_EXECUTE_IF("read_file_mmap_unsupported", return E_NOT_SUPPORT;); @@ -354,7 +355,7 @@ int ReadFile::map_file() { return E_OK; } -void ReadFile::unmap_file() { +void LocalRandomAccessReadFile::unmap_file() { if (mapped_data_ != nullptr) { #ifdef _WIN32 UnmapViewOfFile(mapped_data_); @@ -372,63 +373,12 @@ void ReadFile::unmap_file() { mapped_size_ = 0; } -int ReadFile::check_file_magic() { - int ret = E_OK; - if (file_size_ < MIN_FILE_SIZE) { - ret = E_TSFILE_CORRUPTED; - LOGE("tsfile" << file_path_.c_str() - << "is corrupted, file_size=" << file_size_); - } else { - char buf[MAGIC_STRING_TSFILE_LEN]; - int32_t read_len = 0; - // file header magic - memset(buf, 0, MAGIC_STRING_TSFILE_LEN); - if (RET_FAIL(read(0, buf, MAGIC_STRING_TSFILE_LEN, read_len))) { - } else if (read_len != MAGIC_STRING_TSFILE_LEN) { - ret = E_TSFILE_CORRUPTED; - } else if (memcmp(buf, MAGIC_STRING_TSFILE, MAGIC_STRING_TSFILE_LEN) != - 0) { - ret = E_TSFILE_CORRUPTED; - } - if (IS_FAIL(ret)) { - return ret; - } - - char version = 0; - if (RET_FAIL(read(MAGIC_STRING_TSFILE_LEN, &version, 1, read_len))) { - } else if (read_len != 1) { - ret = E_TSFILE_CORRUPTED; - } else { - file_version_ = static_cast(version); - // Version 3 remains readable for backward compatibility; version - // 4 is the current writer format. Other values are not safely - // interpretable and must be reported as an input failure before - // metadata parsing begins. - if (file_version_ != 3 && - file_version_ != static_cast(VERSION_NUM_BYTE)) { - ret = E_UNSUPPORTED_VERSION; - } - } - if (IS_FAIL(ret)) { - return ret; - } - - // file footer magic - memset(buf, 0, MAGIC_STRING_TSFILE_LEN); - if (RET_FAIL(read(file_size_ - MAGIC_STRING_TSFILE_LEN, buf, - MAGIC_STRING_TSFILE_LEN, read_len))) { - } else if (read_len != MAGIC_STRING_TSFILE_LEN) { - ret = E_TSFILE_CORRUPTED; - } else if (memcmp(buf, MAGIC_STRING_TSFILE, MAGIC_STRING_TSFILE_LEN) != - 0) { - ret = E_TSFILE_CORRUPTED; - } - } - return ret; +int LocalRandomAccessReadFile::check_file_magic() { + return validate_tsfile(*this, &file_version_); } -int ReadFile::read(int64_t offset, char* buf, int32_t buf_size, - int32_t& read_len) { +int LocalRandomAccessReadFile::read(int64_t offset, char* buf, int32_t buf_size, + int32_t& read_len) { read_len = 0; if (offset < 0 || buf_size < 0 || (buf == nullptr && buf_size > 0)) { return E_INVALID_ARG; diff --git a/cpp/src/file/read_file.h b/cpp/src/file/local_random_access_read_file.h similarity index 74% rename from cpp/src/file/read_file.h rename to cpp/src/file/local_random_access_read_file.h index 0c7500d48..4d74254fd 100644 --- a/cpp/src/file/read_file.h +++ b/cpp/src/file/local_random_access_read_file.h @@ -17,8 +17,8 @@ * under the License. */ -#ifndef FILE_READ_FILE_H -#define FILE_READ_FILE_H +#ifndef FILE_LOCAL_RANDOM_ACCESS_READ_FILE_H +#define FILE_LOCAL_RANDOM_ACCESS_READ_FILE_H #include @@ -26,26 +26,30 @@ #include #include "common/config/config.h" +#include "file/random_access_read_file.h" #include "utils/errno_define.h" #include "utils/util_define.h" namespace storage { -class ReadFile { +class LocalRandomAccessReadFile : public RandomAccessReadFile { public: - ReadFile(); - ~ReadFile() { destroy(); } - ReadFile(const ReadFile&) = delete; - ReadFile& operator=(const ReadFile&) = delete; + LocalRandomAccessReadFile(); + ~LocalRandomAccessReadFile() override { destroy(); } + LocalRandomAccessReadFile(const LocalRandomAccessReadFile&) = delete; + LocalRandomAccessReadFile& operator=(const LocalRandomAccessReadFile&) = + delete; void destroy() { close(); } int open(const std::string& file_path); - FORCE_INLINE bool is_opened() const { + FORCE_INLINE bool is_opened() const override { return fd_ >= 0 || mapped_data_ != nullptr; } - FORCE_INLINE int64_t file_size() const { return file_size_; } - FORCE_INLINE const std::string& file_path() const { return file_path_; } + FORCE_INLINE int64_t file_size() const override { return file_size_; } + FORCE_INLINE const std::string& file_path() const override { + return file_path_; + } FORCE_INLINE unsigned char file_version() const { return file_version_; } FORCE_INLINE common::FileReadBackend active_backend() const { return active_backend_; @@ -55,17 +59,17 @@ class ReadFile { * MMAP readers return the generation captured before closing the file * descriptor. */ - int generation(uint64_t& size, uint64_t& fingerprint) const; + int generation(uint64_t& size, uint64_t& fingerprint) const override; /* * try to reader @buf_size bytes from @offset of this file * into @buf. @read_len return the actual len reader. */ int read(int64_t offset, char* buf, int32_t buf_size, - int32_t& ret_read_len); + int32_t& ret_read_len) override; // open()/close() must not race with read() or generation(). Concurrent // reads are supported after the file has been opened. - void close(); + void close() override; private: int get_file_size(int64_t& file_size); @@ -97,4 +101,4 @@ class ReadFile { }; } // end namespace storage -#endif // FILE_READ_FILE_H +#endif // FILE_LOCAL_RANDOM_ACCESS_READ_FILE_H diff --git a/cpp/src/file/random_access_read_file.cc b/cpp/src/file/random_access_read_file.cc new file mode 100644 index 000000000..240538709 --- /dev/null +++ b/cpp/src/file/random_access_read_file.cc @@ -0,0 +1,80 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ + +#include "file/random_access_read_file.h" + +#include + +#include "common/tsfile_common.h" +#include "utils/errno_define.h" + +namespace storage { + +int validate_tsfile(RandomAccessReadFile& file, unsigned char* file_version) { + static const int64_t MIN_FILE_SIZE = 2 * MAGIC_STRING_TSFILE_LEN + 1; + if (!file.is_opened()) { + return common::E_FILE_READ_ERR; + } + if (file.file_size() < MIN_FILE_SIZE) { + return common::E_TSFILE_CORRUPTED; + } + + char buffer[MAGIC_STRING_TSFILE_LEN]; + int32_t read_size = 0; + int ret = file.read(0, buffer, MAGIC_STRING_TSFILE_LEN, read_size); + if (ret != common::E_OK) { + return ret; + } + if (read_size != MAGIC_STRING_TSFILE_LEN || + std::memcmp(buffer, MAGIC_STRING_TSFILE, MAGIC_STRING_TSFILE_LEN) != + 0) { + return common::E_TSFILE_CORRUPTED; + } + + char version = 0; + ret = file.read(MAGIC_STRING_TSFILE_LEN, &version, 1, read_size); + if (ret != common::E_OK) { + return ret; + } + if (read_size != 1) { + return common::E_TSFILE_CORRUPTED; + } + const unsigned char parsed_version = static_cast(version); + if (parsed_version != 3 && + parsed_version != static_cast(VERSION_NUM_BYTE)) { + return common::E_UNSUPPORTED_VERSION; + } + + ret = file.read(file.file_size() - MAGIC_STRING_TSFILE_LEN, buffer, + MAGIC_STRING_TSFILE_LEN, read_size); + if (ret != common::E_OK) { + return ret; + } + if (read_size != MAGIC_STRING_TSFILE_LEN || + std::memcmp(buffer, MAGIC_STRING_TSFILE, MAGIC_STRING_TSFILE_LEN) != + 0) { + return common::E_TSFILE_CORRUPTED; + } + if (file_version != nullptr) { + *file_version = parsed_version; + } + return common::E_OK; +} + +} // namespace storage diff --git a/cpp/src/file/random_access_read_file.h b/cpp/src/file/random_access_read_file.h new file mode 100644 index 000000000..04c2d25d7 --- /dev/null +++ b/cpp/src/file/random_access_read_file.h @@ -0,0 +1,56 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ + +#ifndef FILE_RANDOM_ACCESS_READ_FILE_H +#define FILE_RANDOM_ACCESS_READ_FILE_H + +#include + +#include + +namespace storage { + +/** + * A readable, seek-independent byte source used by the TsFile reader stack. + * + * Implementations must continue short underlying reads until the requested + * range is filled or EOF is reached. Concurrent read() calls must be safe + * after initialization; close() must not race with read() or generation(). + */ +class RandomAccessReadFile { + public: + virtual ~RandomAccessReadFile() = default; + + virtual bool is_opened() const = 0; + virtual int64_t file_size() const = 0; + virtual const std::string& file_path() const = 0; + virtual int generation(uint64_t& size, uint64_t& fingerprint) const = 0; + + /** Read up to @p size bytes at @p offset without exposing a cursor. */ + virtual int read(int64_t offset, char* buffer, int32_t size, + int32_t& read_size) = 0; + virtual void close() = 0; +}; + +int validate_tsfile(RandomAccessReadFile& file, + unsigned char* file_version = nullptr); + +} // namespace storage + +#endif // FILE_RANDOM_ACCESS_READ_FILE_H diff --git a/cpp/src/file/tsfile_io_reader.cc b/cpp/src/file/tsfile_io_reader.cc index 41bcd6b93..9b4d8aef4 100644 --- a/cpp/src/file/tsfile_io_reader.cc +++ b/cpp/src/file/tsfile_io_reader.cc @@ -22,6 +22,7 @@ #include #include "common/allocator/alloc_base.h" +#include "file/local_random_access_read_file.h" #include "reader/prepared_series.h" using namespace common; @@ -29,14 +30,15 @@ using namespace common; namespace storage { int TsFileIOReader::init(const std::string& file_path) { int ret = E_OK; - read_file_ = new ReadFile; + LocalRandomAccessReadFile* local_file = new LocalRandomAccessReadFile; + read_file_ = local_file; read_file_created_ = true; - if (RET_FAIL(read_file_->open(file_path))) { + if (RET_FAIL(local_file->open(file_path))) { } return ret; } -int TsFileIOReader::init(ReadFile* read_file) { +int TsFileIOReader::init(RandomAccessReadFile* read_file) { if (IS_NULL(read_file)) { ASSERT(false); return E_INVALID_ARG; @@ -49,7 +51,7 @@ int TsFileIOReader::init(ReadFile* read_file) { void TsFileIOReader::reset() { if (read_file_ != nullptr) { if (read_file_created_) { - read_file_->destroy(); + read_file_->close(); delete read_file_; } read_file_ = nullptr; @@ -92,9 +94,9 @@ int TsFileIOReader::alloc_ssi(std::shared_ptr device_id, } namespace { -int load_exact_timeseries_index(ReadFile* read_file, uint64_t offset, - uint32_t length, PageArena& arena, - TimeseriesIndex*& index) { +int load_exact_timeseries_index(RandomAccessReadFile* read_file, + uint64_t offset, uint32_t length, + PageArena& arena, TimeseriesIndex*& index) { if (read_file == nullptr || length == 0 || length > static_cast(std::numeric_limits::max()) || offset > static_cast(std::numeric_limits::max()) || @@ -336,17 +338,18 @@ int TsFileIOReader::alloc_multi_ssi( // Batch-load all measurement TimeseriesIndex entries in a single pass over // the index tree — cheaper than N separate binary searches + file reads. - std::vector, int64_t>> - all_leaves; - if (RET_FAIL(get_all_leaf(top_node, all_leaves, ssi_pa))) { - ssi->destroy(); - mem_free(ssi); - ssi = nullptr; - return ret; - } std::vector all_ts_idxs; - if (RET_FAIL( - do_load_all_timeseries_index(all_leaves, ssi_pa, all_ts_idxs))) { + { + // Leaf entry deleters access arena-owned nodes. Release the entries + // before any error path destroys the iterator and its arena. + std::vector, int64_t>> + all_leaves; + if (RET_FAIL(get_all_leaf(top_node, all_leaves, ssi_pa))) { + } else { + ret = do_load_all_timeseries_index(all_leaves, ssi_pa, all_ts_idxs); + } + } + if (RET_FAIL(ret)) { ssi->destroy(); mem_free(ssi); ssi = nullptr; @@ -391,7 +394,9 @@ int TsFileIOReader::get_device_timeseries_meta_without_chunk_meta( std::shared_ptr device_id, std::vector& timeseries_indexs, PageArena& pa) { int ret = E_OK; - load_tsfile_meta_if_necessary(); + if (RET_FAIL(load_tsfile_meta_if_necessary())) { + return ret; + } std::shared_ptr meta_index_entry; int64_t end_offset; std::vector, int64_t>> @@ -412,7 +417,9 @@ int TsFileIOReader::get_device_timeseries_meta_by_offset( int64_t start_offset, int64_t end_offset, std::vector& timeseries_indexs, PageArena& pa) { int ret = E_OK; - load_tsfile_meta_if_necessary(); + if (RET_FAIL(load_tsfile_meta_if_necessary())) { + return ret; + } std::vector, int64_t>> meta_index_entry_list; @@ -433,6 +440,8 @@ int TsFileIOReader::get_device_timeseries_meta_by_offset( if (RET_FAIL(read_file_->read(start_offset, data_buf, read_size, ret_read_len))) { return ret; + } else if (ret_read_len != read_size) { + return E_FILE_READ_ERR; } if (RET_FAIL(top_node->deserialize_from(data_buf, read_size))) { return ret; @@ -446,7 +455,9 @@ int TsFileIOReader::get_device_timeseries_meta_by_offset( } } - get_all_leaf(top_node, meta_index_entry_list, pa); + if (RET_FAIL(get_all_leaf(top_node, meta_index_entry_list, pa))) { + return ret; + } if (RET_FAIL(do_load_all_timeseries_index(meta_index_entry_list, pa, timeseries_indexs))) { @@ -669,6 +680,8 @@ int TsFileIOReader::get_cached_device_node(std::shared_ptr device_id, if (RET_FAIL(read_file_->read(start_offset, data_buf.get(), read_size, ret_read_len))) { return ret; + } else if (ret_read_len != read_size) { + return E_FILE_READ_ERR; } CachedDeviceNode cached; @@ -841,6 +854,8 @@ int TsFileIOReader::load_all_measurement_index_entry( MetaIndexNode::self_deleter); if (RET_FAIL(read_file_->read(start_offset, data_buf, read_size, ret_read_len))) { + } else if (ret_read_len != read_size) { + ret = E_FILE_READ_ERR; } else if (RET_FAIL(top_node->deserialize_from(data_buf, read_size))) { } #if DEBUG_SE @@ -851,7 +866,7 @@ int TsFileIOReader::load_all_measurement_index_entry( #endif // 2. search from top_node in top-down way if (IS_SUCC(ret)) { - get_all_leaf(top_node, ret_measurement_index_entry, pa); + ret = get_all_leaf(top_node, ret_measurement_index_entry, pa); } if (ret == E_NOT_EXIST) { ret = E_MEASUREMENT_NOT_EXIST; @@ -876,6 +891,9 @@ int TsFileIOReader::read_device_meta_index(int64_t start_offset, device_meta_index = new (m_idx_node_buf) MetaIndexNode(&pa); if (RET_FAIL(read_file_->read(start_offset, data_buf, read_size, ret_read_len))) { + return ret; + } else if (ret_read_len != read_size) { + return E_FILE_READ_ERR; } if (!leaf) { ret = device_meta_index->device_deserialize_from(data_buf, read_size); @@ -917,6 +935,8 @@ int TsFileIOReader::get_timeseries_indexes( if (RET_FAIL(read_file_->read(start_offset, data_buf, read_size, ret_read_len))) { return ret; + } else if (ret_read_len != read_size) { + return E_FILE_READ_ERR; } else if (RET_FAIL(top_node->deserialize_from(data_buf, read_size))) { return ret; } @@ -924,7 +944,10 @@ int TsFileIOReader::get_timeseries_indexes( bool is_aligned = is_aligned_device(top_node); TimeseriesIndex* timeseries_index = nullptr; if (is_aligned) { - get_time_column_metadata(top_node, timeseries_index, pa); + if (RET_FAIL( + get_time_column_metadata(top_node, timeseries_index, pa))) { + return ret; + } } int64_t idx = 0; @@ -1017,8 +1040,9 @@ int TsFileIOReader::search_from_internal_node( int32_t ret_read_len = 0; if (RET_FAIL(read_file_->read(index_entry->get_offset(), data_buf, read_size, ret_read_len))) { + return ret; } else if (read_size != ret_read_len) { - return E_TSFILE_CORRUPTED; + return E_FILE_READ_ERR; } if (!is_device) { ret = cur_level_index_node->deserialize_from(data_buf, read_size); @@ -1087,6 +1111,9 @@ int TsFileIOReader::get_time_column_metadata( return ret; } } + if (ret_read_len != end_idx - start_idx) { + return E_FILE_READ_ERR; + } buffer.wrap_from(ti_buf, end_idx - start_idx); void* buf = pa.alloc(sizeof(TimeseriesIndex)); if (IS_NULL(buf)) { @@ -1108,10 +1135,15 @@ int TsFileIOReader::get_time_column_metadata( if (RET_FAIL(read_file_->read(start_idx, ti_buf, end_idx - start_idx, ret_read_len))) { return ret; + } else if (ret_read_len != end_idx - start_idx) { + return E_FILE_READ_ERR; } std::shared_ptr meta_index_node = std::make_shared(&pa); - meta_index_node->deserialize_from(ti_buf, end_idx - start_idx); + if (RET_FAIL(meta_index_node->deserialize_from(ti_buf, + end_idx - start_idx))) { + return ret; + } return get_time_column_metadata(meta_index_node, ret_timeseries_index, pa); } @@ -1132,6 +1164,8 @@ int TsFileIOReader::do_load_timeseries_index( } if (RET_FAIL( read_file_->read(start_offset, ti_buf, read_size, ret_read_len))) { + } else if (ret_read_len != read_size) { + ret = E_FILE_READ_ERR; } else { ByteStream bs; bs.wrap_from(ti_buf, read_size); @@ -1218,6 +1252,8 @@ int TsFileIOReader::do_load_all_timeseries_index( if (RET_FAIL(read_file_->read(start_offset, ti_buf, read_size, ret_read_len))) { return ret; + } else if (ret_read_len != read_size) { + return E_FILE_READ_ERR; } ByteStream bs; bs.wrap_from(ti_buf, read_size); @@ -1296,13 +1332,16 @@ int TsFileIOReader::get_all_leaf( read_file_->read(index_node->children_[i]->get_offset(), data_buf, read_size, ret_read_len))) { } else if (read_size != ret_read_len) { - ret = E_TSFILE_CORRUPTED; + ret = E_FILE_READ_ERR; } else if (RET_FAIL(cur_level_index_node->deserialize_from( data_buf, read_size))) { } else { ret = get_all_leaf(cur_level_index_node, index_node_entry_list, pa); } + if (RET_FAIL(ret)) { + return ret; + } } } return ret; @@ -1344,7 +1383,7 @@ int TsFileIOReader::get_next_page(TsBlock *ret_tsblock) } int TsFileIOReader::init_first_chunk_reader(ChunkMeta *cm, - ReadFile *read_file, + RandomAccessReadFile *read_file, const ColumnDesc &col_desc) { ASSERT(!chunk_reader_.has_more_data()); diff --git a/cpp/src/file/tsfile_io_reader.h b/cpp/src/file/tsfile_io_reader.h index c4883549a..2f98a753c 100644 --- a/cpp/src/file/tsfile_io_reader.h +++ b/cpp/src/file/tsfile_io_reader.h @@ -26,7 +26,7 @@ #include #include "common/tsblock/tsblock.h" -#include "file/read_file.h" +#include "file/random_access_read_file.h" #include "reader/chunk_reader.h" #include "reader/filter/filter.h" #include "reader/tsfile_series_scan_iterator.h" @@ -55,7 +55,7 @@ class TsFileIOReader { device_node_cache_pa_.init(512, common::MOD_TSFILE_READER); } - // Free only the ReadFile we own (created by init(const std::string&)). + // Free only the local source we own (created by init(const std::string&)). // Without an explicit destructor that raw pointer leaks whenever a // TsFileIOReader value goes out of scope without an explicit reset() (e.g. // a stack instance in a test). We deliberately do NOT call reset() here: @@ -68,7 +68,7 @@ class TsFileIOReader { // reset() leaves read_file_ == nullptr, so this never double-frees. ~TsFileIOReader() { if (read_file_created_ && read_file_ != nullptr) { - read_file_->destroy(); + read_file_->close(); delete read_file_; read_file_ = nullptr; } @@ -76,7 +76,7 @@ class TsFileIOReader { int init(const std::string& file_path); - int init(ReadFile* read_file); + int init(RandomAccessReadFile* read_file); void reset(); @@ -112,9 +112,15 @@ class TsFileIOReader { std::string get_file_path() const { return read_file_->file_path(); } + int get_tsfile_meta(TsFileMeta*& tsfile_meta) { + const int ret = load_tsfile_meta_if_necessary(); + tsfile_meta = ret == common::E_OK ? &tsfile_meta_ : nullptr; + return ret; + } + // Raw read access for callers that need to parse structures the metadata // index does not carry, e.g. the chunk header at a ChunkMeta offset. - ReadFile* get_read_file() const { return read_file_; } + RandomAccessReadFile* get_read_file() const { return read_file_; } TsFileMeta* get_tsfile_meta() { load_tsfile_meta_if_necessary(); @@ -237,7 +243,7 @@ class TsFileIOReader { static std::string device_node_cache_key( const std::shared_ptr& device_id); - ReadFile* read_file_; + RandomAccessReadFile* read_file_; common::PageArena tsfile_meta_page_arena_; TsFileMeta tsfile_meta_; bool tsfile_meta_ready_; diff --git a/cpp/src/reader/aligned_chunk_reader.cc b/cpp/src/reader/aligned_chunk_reader.cc index a9f759b67..1c5b838d7 100644 --- a/cpp/src/reader/aligned_chunk_reader.cc +++ b/cpp/src/reader/aligned_chunk_reader.cc @@ -33,7 +33,7 @@ using namespace common; namespace storage { -int AlignedChunkReader::init(ReadFile* read_file, String m_name, +int AlignedChunkReader::init(RandomAccessReadFile* read_file, String m_name, TSDataType data_type, Filter* time_filter) { read_file_ = read_file; measurement_name_.shallow_copy_from(m_name); @@ -1384,7 +1384,8 @@ int AlignedChunkReader::decode_time_page_with(const ChunkPageInfo& page_info, if (heap) common::mem_free(compressed_buf); return ret; } - // ReadFile::read() returns E_OK + short read_len on EOF; uncompressing + // RandomAccessReadFile::read() returns E_OK + short read_len on EOF; + // uncompressing // page_info.time_compressed_size from a buffer with uninitialised tail // bytes would feed garbage to the decompressor. if (read_len != static_cast(page_info.time_compressed_size)) { diff --git a/cpp/src/reader/aligned_chunk_reader.h b/cpp/src/reader/aligned_chunk_reader.h index 45dd2ef99..5f6b0b94e 100644 --- a/cpp/src/reader/aligned_chunk_reader.h +++ b/cpp/src/reader/aligned_chunk_reader.h @@ -24,7 +24,7 @@ #include "common/tsfile_common.h" #include "compress/compressor.h" #include "encoding/decoder.h" -#include "file/read_file.h" +#include "file/random_access_read_file.h" #include "reader/filter/filter.h" #include "reader/ichunk_reader.h" @@ -120,7 +120,7 @@ class AlignedChunkReader : public IChunkReader { time_uncompressed_buf_(nullptr), value_uncompressed_buf_(nullptr), cur_value_index(-1) {} - int init(ReadFile* read_file, common::String m_name, + int init(RandomAccessReadFile* read_file, common::String m_name, common::TSDataType data_type, Filter* time_filter) override; void reset() override; void destroy() override; @@ -276,7 +276,7 @@ class AlignedChunkReader : public IChunkReader { int32_t row_limit) const; private: - ReadFile* read_file_; + RandomAccessReadFile* read_file_; // ── Single-value mode fields (kept for backward compat) ────────────── ChunkMeta* time_chunk_meta_; ChunkMeta* value_chunk_meta_; diff --git a/cpp/src/reader/chunk_reader.cc b/cpp/src/reader/chunk_reader.cc index b88b070ef..271b5c206 100644 --- a/cpp/src/reader/chunk_reader.cc +++ b/cpp/src/reader/chunk_reader.cc @@ -27,8 +27,8 @@ using namespace common; namespace storage { -int ChunkReader::init(ReadFile* read_file, String m_name, TSDataType data_type, - Filter* time_filter) { +int ChunkReader::init(RandomAccessReadFile* read_file, String m_name, + TSDataType data_type, Filter* time_filter) { read_file_ = read_file; measurement_name_.shallow_copy_from(m_name); time_decoder_ = DecoderFactory::alloc_time_decoder(); diff --git a/cpp/src/reader/chunk_reader.h b/cpp/src/reader/chunk_reader.h index b92be8169..b54301199 100644 --- a/cpp/src/reader/chunk_reader.h +++ b/cpp/src/reader/chunk_reader.h @@ -24,7 +24,7 @@ #include "common/tsfile_common.h" #include "compress/compressor.h" #include "encoding/decoder.h" -#include "file/read_file.h" +#include "file/random_access_read_file.h" #include "reader/filter/filter.h" #include "reader/ichunk_reader.h" @@ -48,7 +48,7 @@ class ChunkReader : public IChunkReader { time_in_(), value_in_(), uncompressed_buf_(nullptr) {} - int init(ReadFile* read_file, common::String m_name, + int init(RandomAccessReadFile* read_file, common::String m_name, common::TSDataType data_type, Filter* time_filter) override; void reset() override; void destroy() override; @@ -127,7 +127,7 @@ class ChunkReader : public IChunkReader { Filter* filter); private: - ReadFile* read_file_; + RandomAccessReadFile* read_file_; ChunkMeta* chunk_meta_; common::String measurement_name_; ChunkHeader chunk_header_; diff --git a/cpp/src/reader/ichunk_reader.h b/cpp/src/reader/ichunk_reader.h index 32985cfd2..1e6fe3997 100644 --- a/cpp/src/reader/ichunk_reader.h +++ b/cpp/src/reader/ichunk_reader.h @@ -23,7 +23,7 @@ #include "common/allocator/my_string.h" #include "compress/compressor.h" #include "encoding/decoder.h" -#include "file/read_file.h" +#include "file/random_access_read_file.h" #include "reader/filter/filter.h" namespace storage { @@ -31,7 +31,7 @@ namespace storage { class IChunkReader { public: IChunkReader() {} - virtual int init(ReadFile* read_file, common::String m_name, + virtual int init(RandomAccessReadFile* read_file, common::String m_name, common::TSDataType data_type, Filter* time_filter) { return common::E_OK; } diff --git a/cpp/src/reader/qds_without_timegenerator.cc b/cpp/src/reader/qds_without_timegenerator.cc index fcd38762d..f9ce7903c 100644 --- a/cpp/src/reader/qds_without_timegenerator.cc +++ b/cpp/src/reader/qds_without_timegenerator.cc @@ -70,7 +70,9 @@ int QDSWithoutTimeGenerator::init_prepared( value_iters_.resize(1); row_record_ = new RowRecord(2); index_lookup_.insert({column_name, 1}); - load_next_tsblock(0, true); + if (RET_FAIL(load_next_tsblock(0, true))) { + return ret; + } remaining_offset_ = ssi->get_row_offset(); if (!table_aligned) { remaining_limit_ = ssi->get_row_limit(); @@ -161,7 +163,9 @@ int QDSWithoutTimeGenerator::init_internal(TsFileIOReader* io_reader, value_iters_.resize(path_count); for (size_t i = 0; i < path_count; i++) { - load_next_tsblock(i, true); + if (RET_FAIL(load_next_tsblock(i, true))) { + return ret; + } // Prefer the type carried by the value iterator, but fall back to the // timeseries-index type captured before load_next_tsblock() when no // TsBlock was produced (e.g. limit==0 skips every row, or the series is @@ -222,6 +226,10 @@ void QDSWithoutTimeGenerator::close() { } int QDSWithoutTimeGenerator::next(bool& has_next) { + has_next = false; + if (read_error_ != E_OK) { + return read_error_; + } // For single path, apply offset/limit at row level. if (is_single_path_) { while (true) { @@ -264,7 +272,10 @@ int QDSWithoutTimeGenerator::next(bool& has_next) { heap_time_.insert(std::make_pair(timev, idx)); time_iters_[idx]->next(); } else { - load_next_tsblock(idx, false); + read_error_ = load_next_tsblock(idx, false); + if (read_error_ != E_OK) { + return read_error_; + } } if (skip_row) { @@ -319,7 +330,11 @@ int QDSWithoutTimeGenerator::next(bool& has_next) { // Pass merge_cursor (current time) as min_time_hint // to help SSI skip chunks/pages that are entirely before // the current merge position. - load_next_tsblock_with_hint(iter->second, false, time); + read_error_ = + load_next_tsblock_with_hint(iter->second, false, time); + if (read_error_ != E_OK) { + return read_error_; + } } std::multimap::iterator cur = iter; iter++; // cppcheck-suppress postfixOperator diff --git a/cpp/src/reader/qds_without_timegenerator.h b/cpp/src/reader/qds_without_timegenerator.h index eae2a6495..1c0fc9f03 100644 --- a/cpp/src/reader/qds_without_timegenerator.h +++ b/cpp/src/reader/qds_without_timegenerator.h @@ -44,7 +44,8 @@ class QDSWithoutTimeGenerator : public ResultSet { remaining_offset_(0), remaining_limit_(-1), owned_time_filter_(nullptr), - is_single_path_(false) {} + is_single_path_(false), + read_error_(common::E_OK) {} ~QDSWithoutTimeGenerator() { close(); } int init(TsFileIOReader* io_reader, QueryExpression* qe); int init(TsFileIOReader* io_reader, QueryExpression* qe, int offset, @@ -80,6 +81,9 @@ class QDSWithoutTimeGenerator : public ResultSet { int remaining_limit_; Filter* owned_time_filter_; bool is_single_path_; + // A failed block load leaves the merge iterators unusable. Keep reporting + // the failure until the result is closed, including subsequent next calls. + int read_error_; }; } // namespace storage diff --git a/cpp/src/reader/query_executor.h b/cpp/src/reader/query_executor.h index 553b4f392..33d44251a 100644 --- a/cpp/src/reader/query_executor.h +++ b/cpp/src/reader/query_executor.h @@ -23,7 +23,7 @@ #include "common/row_record.h" #include "expression.h" -#include "file/read_file.h" +#include "file/random_access_read_file.h" #include "reader/tsfile_series_scan_iterator.h" namespace storage { @@ -40,7 +40,8 @@ class QueryExecutor { } } - // virtual int init(QueryExpression *query_expr, ReadFile *read_file) { + // virtual int init(QueryExpression *query_expr, + // RandomAccessReadFile *read_file) { // ASSERT(false); return 0; }; virtual RowRecord* execute() { diff --git a/cpp/src/reader/table_query_executor.h b/cpp/src/reader/table_query_executor.h index 4cbd8ea3d..11ccbbd23 100644 --- a/cpp/src/reader/table_query_executor.h +++ b/cpp/src/reader/table_query_executor.h @@ -21,6 +21,7 @@ #include "common/schema.h" #include "expression.h" +#include "file/random_access_read_file.h" #include "imeta_data_querier.h" #include "reader/block/device_ordered_tsblock_reader.h" #include "reader/block/tsblock_reader.h" @@ -46,7 +47,8 @@ class TableQueryExecutor { table_query_ordering_(table_query_ordering), block_size_(block_size), return_mode_(RETURN_ROW) {} - TableQueryExecutor(ReadFile* read_file, const int batch_size = -1) { + TableQueryExecutor(RandomAccessReadFile* read_file, + const int batch_size = -1) { tsfile_io_reader_ = new TsFileIOReader(); tsfile_io_reader_->init(read_file); meta_data_querier_ = new MetadataQuerier(tsfile_io_reader_); @@ -92,4 +94,4 @@ class TableQueryExecutor { } // namespace storage -#endif // READER_TABLE_QUERY_EXECUTOR_H \ No newline at end of file +#endif // READER_TABLE_QUERY_EXECUTOR_H diff --git a/cpp/src/reader/tsfile_executor.cc b/cpp/src/reader/tsfile_executor.cc index bdae6c863..5ba971e70 100644 --- a/cpp/src/reader/tsfile_executor.cc +++ b/cpp/src/reader/tsfile_executor.cc @@ -42,7 +42,7 @@ TsFileExecutor::TsFileExecutor() TsFileExecutor::~TsFileExecutor() {} -int TsFileExecutor::init(ReadFile* read_file) { +int TsFileExecutor::init(RandomAccessReadFile* read_file) { int ret = E_OK; io_reader_.reset(); if (RET_FAIL(io_reader_.init(read_file))) { diff --git a/cpp/src/reader/tsfile_executor.h b/cpp/src/reader/tsfile_executor.h index 86fe581ab..cd862b836 100644 --- a/cpp/src/reader/tsfile_executor.h +++ b/cpp/src/reader/tsfile_executor.h @@ -22,7 +22,7 @@ #include #include -#include "file/read_file.h" +#include "file/random_access_read_file.h" #include "query_executor.h" #include "result_set.h" @@ -33,7 +33,7 @@ class TsFileExecutor // : public QueryExecutor public: TsFileExecutor(); ~TsFileExecutor(); - int init(ReadFile* read_file); + int init(RandomAccessReadFile* read_file); int init(const std::string& file_path); int execute(QueryExpression* query_expr, ResultSet*& ret_qds); int execute(QueryExpression* query_expr, ResultSet*& ret_qds, int offset, @@ -54,6 +54,9 @@ class TsFileExecutor // : public QueryExecutor int64_t start_time, int64_t end_time, int offset, int limit, ResultSet*& ret_qds); void destroy_query_data_set(ResultSet* qds); + int get_tsfile_meta(TsFileMeta*& tsfile_meta) { + return io_reader_.get_tsfile_meta(tsfile_meta); + } TsFileMeta* get_tsfile_meta() { return io_reader_.get_tsfile_meta(); } TsFileIOReader* get_tsfile_io_reader() { return &io_reader_; } diff --git a/cpp/src/reader/tsfile_reader.cc b/cpp/src/reader/tsfile_reader.cc index 4d482d186..b3e6b59a1 100644 --- a/cpp/src/reader/tsfile_reader.cc +++ b/cpp/src/reader/tsfile_reader.cc @@ -20,12 +20,13 @@ #include #include +#include #include #include "common/allocator/byte_stream.h" #include "common/schema.h" #include "common/tsfile_common.h" -#include "file/read_file.h" +#include "file/local_random_access_read_file.h" #include "filter/time_operator.h" #include "tsfile_executor.h" @@ -55,7 +56,7 @@ int parse_paths(const std::vector& path_list, int get_all_device_entries(std::vector& entries, std::shared_ptr index_node, - ReadFile* read_file, PageArena& pa) { + RandomAccessReadFile* read_file, PageArena& pa) { int ret = E_OK; if (index_node == nullptr) { return ret; @@ -95,6 +96,8 @@ int get_all_device_entries(std::vector& entries, }); if (RET_FAIL(read_file->read(start_offset, data_buf, read_size, ret_read_len))) { + } else if (ret_read_len != read_size) { + ret = E_FILE_READ_ERR; } else if (RET_FAIL(top_node->device_deserialize_from(data_buf, read_size))) { } else { @@ -121,22 +124,46 @@ TsFileReader::TsFileReader() TsFileReader::~TsFileReader() { close(); } int TsFileReader::open(const std::string& file_path) { + std::unique_ptr read_file( + new LocalRandomAccessReadFile()); int ret = E_OK; - read_file_ = new storage::ReadFile; - tsfile_executor_ = new storage::TsFileExecutor(); // Keep reader diagnostics in the caller's error channel. Printing here // would leak an unstructured line to process stdout/stderr before the CLI // can attach its stable error text and exit code. - if (RET_FAIL(read_file_->open(file_path))) { - } else if (RET_FAIL(tsfile_executor_->init(read_file_))) { + if (RET_FAIL(read_file->open(file_path))) { + return ret; } - return ret; + const unsigned char file_version = read_file->file_version(); + return open_source(std::move(read_file), file_version); +} + +int TsFileReader::open(std::unique_ptr read_file) { + if (read_file == nullptr || !read_file->is_opened()) { + return E_INVALID_ARG; + } + unsigned char file_version = 0; + const int ret = validate_tsfile(*read_file, &file_version); + if (ret != E_OK) { + return ret; + } + return open_source(std::move(read_file), file_version); } -unsigned char TsFileReader::get_file_version() const { - return read_file_ == nullptr ? 0 : read_file_->file_version(); +int TsFileReader::open_source(std::unique_ptr read_file, + unsigned char file_version) { + close(); + read_file_ = std::move(read_file); + file_version_ = file_version; + tsfile_executor_ = new storage::TsFileExecutor(); + const int ret = tsfile_executor_->init(read_file_.get()); + if (ret != E_OK) { + close(); + } + return ret; } +unsigned char TsFileReader::get_file_version() const { return file_version_; } + int TsFileReader::ensure_table_query_executor(int batch_size) { if (table_query_executor_ != nullptr && table_query_executor_batch_size_ == batch_size) { @@ -148,7 +175,8 @@ int TsFileReader::ensure_table_query_executor(int batch_size) { table_query_executor_ = nullptr; } - table_query_executor_ = new TableQueryExecutor(read_file_, batch_size); + table_query_executor_ = + new TableQueryExecutor(read_file_.get(), batch_size); table_query_executor_batch_size_ = batch_size; return E_OK; } @@ -165,9 +193,9 @@ int TsFileReader::close() { } if (read_file_ != nullptr) { read_file_->close(); - delete read_file_; - read_file_ = nullptr; + read_file_.reset(); } + file_version_ = 0; return ret; } @@ -206,9 +234,9 @@ int TsFileReader::query(const std::string& table_name, ResultSet*& result_set, Filter* tag_filter, int batch_size) { int ret = E_OK; - TsFileMeta* tsfile_meta = tsfile_executor_->get_tsfile_meta(); - if (tsfile_meta == nullptr) { - return E_FILE_READ_ERR; + TsFileMeta* tsfile_meta = nullptr; + if (RET_FAIL(tsfile_executor_->get_tsfile_meta(tsfile_meta))) { + return ret; } std::shared_ptr table_schema = tsfile_meta->table_schemas_.at(to_lower(table_name)); @@ -276,9 +304,9 @@ int TsFileReader::queryByRow(const std::string& table_name, int offset, int limit, ResultSet*& result_set, Filter* tag_filter, int batch_size) { int ret = E_OK; - TsFileMeta* tsfile_meta = tsfile_executor_->get_tsfile_meta(); - if (tsfile_meta == nullptr) { - return E_FILE_READ_ERR; + TsFileMeta* tsfile_meta = nullptr; + if (RET_FAIL(tsfile_executor_->get_tsfile_meta(tsfile_meta))) { + return ret; } auto it = tsfile_meta->table_schemas_.find(to_lower(table_name)); if (it == tsfile_meta->table_schemas_.end() || it->second == nullptr) { @@ -297,11 +325,14 @@ int TsFileReader::query_table_on_tree( const std::vector& measurement_names, int64_t star_time, int64_t end_time, ResultSet*& result_set) { int ret = E_OK; - TsFileMeta* tsfile_meta = tsfile_executor_->get_tsfile_meta(); - if (tsfile_meta == nullptr) { - return E_FILE_READ_ERR; + TsFileMeta* tsfile_meta = nullptr; + if (RET_FAIL(tsfile_executor_->get_tsfile_meta(tsfile_meta))) { + return ret; + } + std::vector> device_ids; + if (RET_FAIL(get_all_devices(device_ids))) { + return ret; } - auto device_ids = this->get_all_device_ids(); std::vector> satisfied_device_ids; std::unordered_set measurement_names_set_to_query; size_t device_max_len = 0; @@ -309,7 +340,9 @@ int TsFileReader::query_table_on_tree( if (measurement_names.empty()) { for (auto& device_name : device_ids) { std::vector schemas; - this->get_timeseries_schema(device_name, schemas); + if (RET_FAIL(get_timeseries_schema(device_name, schemas))) { + return ret; + } satisfied_device_ids.push_back(device_name); for (auto& schema : schemas) { measurement_names_set_to_query.insert(schema.measurement_name_); @@ -325,7 +358,9 @@ int TsFileReader::query_table_on_tree( measurement_names.begin(), measurement_names.end()); for (auto& device_name : device_ids) { std::vector schemas; - this->get_timeseries_schema(device_name, schemas); + if (RET_FAIL(get_timeseries_schema(device_name, schemas))) { + return ret; + } bool device_has_required_measurement_names = false; for (auto& schema : schemas) { @@ -393,16 +428,8 @@ std::vector> TsFileReader::get_all_devices( } std::vector> TsFileReader::get_all_device_ids() { - TsFileMeta* tsfile_meta = tsfile_executor_->get_tsfile_meta(); std::vector> device_ids; - if (tsfile_meta != nullptr) { - PageArena pa; - pa.init(512, MOD_TSFILE_READER); - for (auto entry : tsfile_meta->table_metadata_index_node_map_) { - auto index_node = entry.second; - get_all_devices(device_ids, index_node, pa); - } - } + get_all_devices(device_ids); return device_ids; } @@ -410,6 +437,28 @@ std::vector> TsFileReader::get_all_devices() { return get_all_device_ids(); } +int TsFileReader::get_all_devices( + std::vector>& device_ids) { + device_ids.clear(); + if (tsfile_executor_ == nullptr) { + return E_INVALID_ARG; + } + TsFileMeta* tsfile_meta = nullptr; + int ret = tsfile_executor_->get_tsfile_meta(tsfile_meta); + if (ret != E_OK) { + return ret; + } + PageArena pa; + pa.init(512, MOD_TSFILE_READER); + for (const auto& entry : tsfile_meta->table_metadata_index_node_map_) { + if (RET_FAIL(get_all_devices(device_ids, entry.second, pa))) { + device_ids.clear(); + return ret; + } + } + return E_OK; +} + int TsFileReader::get_all_devices( std::vector>& device_ids, std::shared_ptr index_node, PageArena& pa) { @@ -445,10 +494,15 @@ int TsFileReader::get_all_devices( if (RET_FAIL(read_file_->read(start_offset, data_buf, read_size, ret_read_len))) { + return ret; + } else if (ret_read_len != read_size) { + return E_FILE_READ_ERR; } else if (RET_FAIL(top_node->device_deserialize_from( data_buf, read_size))) { - } else { - ret = get_all_devices(device_ids, top_node, pa); + return ret; + } else if (RET_FAIL( + get_all_devices(device_ids, top_node, pa))) { + return ret; } } } @@ -462,7 +516,8 @@ namespace { // chunk header on disk: ChunkMeta::deserialize_from() reads nothing but // offset_of_chunk_header_, so a ChunkMeta obtained from the metadata index // never carries them. Read the header back from the file instead. -int read_chunk_header_codec(ReadFile* read_file, int64_t chunk_header_offset, +int read_chunk_header_codec(RandomAccessReadFile* read_file, + int64_t chunk_header_offset, size_t measurement_name_len, common::TSEncoding& encoding, common::CompressionType& compression) { @@ -618,8 +673,8 @@ DeviceTimeseriesMetadataMap TsFileReader::get_timeseries_metadata() { pa.init(512, MOD_TSFILE_READER); std::vector entries; for (auto& table_entry : tsfile_meta->table_metadata_index_node_map_) { - if (get_all_device_entries(entries, table_entry.second, read_file_, - pa) != E_OK) { + if (get_all_device_entries(entries, table_entry.second, + read_file_.get(), pa) != E_OK) { return result; } } @@ -672,12 +727,26 @@ std::shared_ptr TsFileReader::get_table_schema( std::vector> TsFileReader::get_all_table_schemas() { - TsFileMeta* file_metadata = tsfile_executor_->get_tsfile_meta(); std::vector> table_schemas; + get_all_table_schemas(table_schemas); + return table_schemas; +} + +int TsFileReader::get_all_table_schemas( + std::vector>& table_schemas) { + table_schemas.clear(); + if (tsfile_executor_ == nullptr) { + return E_INVALID_ARG; + } + TsFileMeta* file_metadata = nullptr; + const int ret = tsfile_executor_->get_tsfile_meta(file_metadata); + if (ret != E_OK) { + return ret; + } for (const auto& table_schema : file_metadata->table_schemas_) { table_schemas.push_back(table_schema.second); } - return table_schemas; + return E_OK; } } // namespace storage diff --git a/cpp/src/reader/tsfile_reader.h b/cpp/src/reader/tsfile_reader.h index f07ef3679..1c2458ed0 100644 --- a/cpp/src/reader/tsfile_reader.h +++ b/cpp/src/reader/tsfile_reader.h @@ -20,15 +20,16 @@ #ifndef READER_TSFILE_READER_H #define READER_TSFILE_READER_H +#include + #include "common/row_record.h" #include "common/tsfile_common.h" #include "expression.h" -#include "file/read_file.h" +#include "file/random_access_read_file.h" #include "reader/prepared_series.h" #include "reader/table_query_executor.h" namespace storage { class TsFileExecutor; -class ReadFile; class ResultSet; struct MeasurementSchema; } // namespace storage @@ -49,6 +50,10 @@ class TsFileReader { public: TsFileReader(); ~TsFileReader(); + TsFileReader(const TsFileReader&) = delete; + TsFileReader& operator=(const TsFileReader&) = delete; + TsFileReader(TsFileReader&& other) = delete; + TsFileReader& operator=(TsFileReader&& other) = delete; /** * @brief open the tsfile * @@ -56,6 +61,12 @@ class TsFileReader { * @return Returns 0 on success, or a non-zero error code on failure. */ int open(const std::string& file_path); + /** + * @brief open an initialized random-access source + * + * The reader takes ownership of @p read_file. + */ + int open(std::unique_ptr read_file); /** * @brief close the tsfile, this method should be called after the * query is finished @@ -203,6 +214,9 @@ class TsFileReader { */ std::vector> get_all_devices(); + /** Error-reporting overload. The output is empty on failure. */ + int get_all_devices(std::vector>& device_ids); + /** * @brief get the timeseries schema by the device id and measurement name * @@ -251,7 +265,13 @@ class TsFileReader { */ std::vector> get_all_table_schemas(); + /** Error-reporting overload. The output is empty on failure. */ + int get_all_table_schemas( + std::vector>& table_schemas); + private: + int open_source(std::unique_ptr read_file, + unsigned char file_version); int ensure_table_query_executor(int batch_size); int get_timeseries_metadata_impl( std::shared_ptr device_id, @@ -259,10 +279,11 @@ class TsFileReader { int get_all_devices(std::vector>& device_ids, std::shared_ptr index_node, common::PageArena& pa); - storage::ReadFile* read_file_; + std::unique_ptr read_file_; storage::TsFileExecutor* tsfile_executor_; storage::TableQueryExecutor* table_query_executor_; int table_query_executor_batch_size_ = -1; + unsigned char file_version_ = 0; common::PageArena tsfile_reader_meta_pa_; // Test-only hook for the unbounded-arena-growth regression check. friend class TsFileReaderMetaArenaTest; diff --git a/cpp/src/reader/tsfile_series_scan_iterator.cc b/cpp/src/reader/tsfile_series_scan_iterator.cc index 8c677493e..8d1cb1179 100644 --- a/cpp/src/reader/tsfile_series_scan_iterator.cc +++ b/cpp/src/reader/tsfile_series_scan_iterator.cc @@ -32,8 +32,9 @@ using namespace common; namespace storage { int TsFileSeriesScanIterator::init_prepared( - const std::shared_ptr& prepared, ReadFile* read_file, - Filter* time_filter, common::PageArena& data_pa) { + const std::shared_ptr& prepared, + RandomAccessReadFile* read_file, Filter* time_filter, + common::PageArena& data_pa) { if (prepared == nullptr || prepared->index() == nullptr || read_file == nullptr) { return E_INVALID_ARG; @@ -68,7 +69,8 @@ int TsFileSeriesScanIterator::init_prepared( int TsFileSeriesScanIterator::init_prepared_multi( const std::vector>& prepared, - ReadFile* read_file, Filter* time_filter, common::PageArena& data_pa) { + RandomAccessReadFile* read_file, Filter* time_filter, + common::PageArena& data_pa) { if (prepared.empty() || prepared.front() == nullptr || read_file == nullptr) { return E_INVALID_ARG; diff --git a/cpp/src/reader/tsfile_series_scan_iterator.h b/cpp/src/reader/tsfile_series_scan_iterator.h index d6ea0e29d..449f240b4 100644 --- a/cpp/src/reader/tsfile_series_scan_iterator.h +++ b/cpp/src/reader/tsfile_series_scan_iterator.h @@ -26,7 +26,7 @@ #include "aligned_chunk_reader.h" #include "common/tsblock/tsblock.h" -#include "file/read_file.h" +#include "file/random_access_read_file.h" #include "file/tsfile_io_reader.h" #include "reader/chunk_reader.h" #include "reader/filter/filter.h" @@ -59,8 +59,9 @@ class TsFileSeriesScanIterator { row_limit_(-1) {} ~TsFileSeriesScanIterator() { destroy(); } int init(std::shared_ptr device_id, - const std::string& measurement_name, ReadFile* read_file, - Filter* time_filter, common::PageArena& data_pa) { + const std::string& measurement_name, + RandomAccessReadFile* read_file, Filter* time_filter, + common::PageArena& data_pa) { ASSERT(read_file != nullptr); device_id_ = device_id; measurement_name_ = measurement_name; @@ -70,11 +71,12 @@ class TsFileSeriesScanIterator { return common::E_OK; } int init_prepared(const std::shared_ptr& prepared, - ReadFile* read_file, Filter* time_filter, + RandomAccessReadFile* read_file, Filter* time_filter, common::PageArena& data_pa); int init_prepared_multi( const std::vector>& prepared, - ReadFile* read_file, Filter* time_filter, common::PageArena& data_pa); + RandomAccessReadFile* read_file, Filter* time_filter, + common::PageArena& data_pa); void destroy(); /** @@ -205,7 +207,7 @@ class TsFileSeriesScanIterator { common::TsBlock* alloc_tsblock_multi(); private: - ReadFile* read_file_; + RandomAccessReadFile* read_file_; std::shared_ptr device_id_; std::string measurement_name_; diff --git a/cpp/test/file/read_file_test.cc b/cpp/test/file/local_random_access_read_file_test.cc similarity index 86% rename from cpp/test/file/read_file_test.cc rename to cpp/test/file/local_random_access_read_file_test.cc index bc031fc98..b29bc446f 100644 --- a/cpp/test/file/read_file_test.cc +++ b/cpp/test/file/local_random_access_read_file_test.cc @@ -17,7 +17,7 @@ * under the License. */ -#include "file/read_file.h" +#include "file/local_random_access_read_file.h" #include @@ -70,7 +70,7 @@ class InjectionGuard { const char* point_; }; -class ReadFileBackendTest : public ::testing::Test { +class LocalRandomAccessReadFileTest : public ::testing::Test { protected: void SetUp() override { content_ = "TsFile"; @@ -102,14 +102,14 @@ class ReadFileBackendTest : public ::testing::Test { std::string content_; }; -TEST_F(ReadFileBackendTest, PreadIsTheDefaultBackend) { +TEST_F(LocalRandomAccessReadFileTest, PreadIsTheDefaultBackend) { EXPECT_EQ(common::get_file_read_backend(), common::FileReadBackend::PREAD); } -TEST_F(ReadFileBackendTest, PreadPreservesPositionedReadBehavior) { +TEST_F(LocalRandomAccessReadFileTest, PreadPreservesPositionedReadBehavior) { ASSERT_EQ(common::set_file_read_backend(common::FileReadBackend::PREAD), common::E_OK); - storage::ReadFile file; + storage::LocalRandomAccessReadFile file; ASSERT_EQ(file.open(file_name_), common::E_OK); EXPECT_TRUE(file.is_opened()); EXPECT_EQ(file.active_backend(), common::FileReadBackend::PREAD); @@ -132,10 +132,11 @@ TEST_F(ReadFileBackendTest, PreadPreservesPositionedReadBehavior) { EXPECT_EQ(read_len, 0); } -TEST_F(ReadFileBackendTest, MmapReadsBoundedRangesAndReleasesResources) { +TEST_F(LocalRandomAccessReadFileTest, + MmapReadsBoundedRangesAndReleasesResources) { ASSERT_EQ(common::set_file_read_backend(common::FileReadBackend::MMAP), common::E_OK); - storage::ReadFile file; + storage::LocalRandomAccessReadFile file; ASSERT_EQ(file.open(file_name_), common::E_OK); EXPECT_TRUE(file.is_opened()); EXPECT_EQ(file.active_backend(), common::FileReadBackend::MMAP); @@ -161,10 +162,10 @@ TEST_F(ReadFileBackendTest, MmapReadsBoundedRangesAndReleasesResources) { EXPECT_EQ(std::remove(file_name_.c_str()), 0); } -TEST_F(ReadFileBackendTest, AutoPrefersMmapAndCloseIsIdempotent) { +TEST_F(LocalRandomAccessReadFileTest, AutoPrefersMmapAndCloseIsIdempotent) { ASSERT_EQ(common::set_file_read_backend(common::FileReadBackend::AUTO), common::E_OK); - storage::ReadFile file; + storage::LocalRandomAccessReadFile file; ASSERT_EQ(file.open(file_name_), common::E_OK); EXPECT_EQ(file.active_backend(), common::FileReadBackend::MMAP); @@ -178,50 +179,53 @@ TEST_F(ReadFileBackendTest, AutoPrefersMmapAndCloseIsIdempotent) { EXPECT_EQ(file.active_backend(), common::FileReadBackend::PREAD); } -TEST_F(ReadFileBackendTest, AutoFallsBackButRequiredMmapReportsFailure) { +TEST_F(LocalRandomAccessReadFileTest, + AutoFallsBackButRequiredMmapReportsFailure) { InjectionGuard mmap_failure("read_file_mmap_fail"); ASSERT_EQ(common::set_file_read_backend(common::FileReadBackend::AUTO), common::E_OK); - storage::ReadFile automatic; + storage::LocalRandomAccessReadFile automatic; ASSERT_EQ(automatic.open(file_name_), common::E_OK); EXPECT_EQ(automatic.active_backend(), common::FileReadBackend::PREAD); automatic.close(); ASSERT_EQ(common::set_file_read_backend(common::FileReadBackend::MMAP), common::E_OK); - storage::ReadFile required; + storage::LocalRandomAccessReadFile required; EXPECT_EQ(required.open(file_name_), common::E_FILE_MAP_ERR); EXPECT_FALSE(required.is_opened()); } -TEST_F(ReadFileBackendTest, AutoFallsBackButRequiredMmapReportsUnsupported) { +TEST_F(LocalRandomAccessReadFileTest, + AutoFallsBackButRequiredMmapReportsUnsupported) { InjectionGuard mmap_unsupported("read_file_mmap_unsupported"); ASSERT_EQ(common::set_file_read_backend(common::FileReadBackend::AUTO), common::E_OK); - storage::ReadFile automatic; + storage::LocalRandomAccessReadFile automatic; ASSERT_EQ(automatic.open(file_name_), common::E_OK); EXPECT_EQ(automatic.active_backend(), common::FileReadBackend::PREAD); automatic.close(); ASSERT_EQ(common::set_file_read_backend(common::FileReadBackend::MMAP), common::E_OK); - storage::ReadFile required; + storage::LocalRandomAccessReadFile required; EXPECT_EQ(required.open(file_name_), common::E_NOT_SUPPORT); EXPECT_FALSE(required.is_opened()); } -TEST_F(ReadFileBackendTest, EmptyFileIsRejectedBeforeMapping) { +TEST_F(LocalRandomAccessReadFileTest, EmptyFileIsRejectedBeforeMapping) { write_file(empty_file_name_, ""); ASSERT_EQ(common::set_file_read_backend(common::FileReadBackend::MMAP), common::E_OK); - storage::ReadFile file; + storage::LocalRandomAccessReadFile file; EXPECT_EQ(file.open(empty_file_name_), common::E_TSFILE_CORRUPTED); EXPECT_FALSE(file.is_opened()); } -TEST_F(ReadFileBackendTest, InvalidConfigurationDoesNotChangeBackend) { +TEST_F(LocalRandomAccessReadFileTest, + InvalidConfigurationDoesNotChangeBackend) { ASSERT_EQ(common::set_file_read_backend(common::FileReadBackend::PREAD), common::E_OK); EXPECT_EQ( diff --git a/cpp/test/file/utf8_path_test.cc b/cpp/test/file/utf8_path_test.cc index a1a93d495..da801875e 100644 --- a/cpp/test/file/utf8_path_test.cc +++ b/cpp/test/file/utf8_path_test.cc @@ -24,7 +24,7 @@ #include #include "common/tsfile_common.h" -#include "file/read_file.h" +#include "file/local_random_access_read_file.h" #include "file/restorable_tsfile_io_writer.h" #include "file/utf8_file_open.h" #include "file/write_file.h" @@ -43,9 +43,8 @@ CWriteFile write_file_new(const char* pathname, int32_t* err_code); void free_write_file(CWriteFile* write_file); } -// Note: storage::WriteFile / storage::ReadFile are named after Win32 API -// functions that declares at global scope, so the class names are -// always written qualified here. +// storage::WriteFile shares its name with a global Win32 API function, so +// storage classes are explicitly qualified here. using namespace common; namespace { @@ -84,9 +83,10 @@ std::wstring WideName() { #endif // Smallest byte sequence that counts as a complete TsFile: head magic, the -// current version byte, then the tail magic. This is what ReadFile::open() -// requires (>= MIN_FILE_SIZE bytes, magic at both ends) and also what -// RestorableTsFileIOWriter's self check treats as complete. +// current version byte, then the tail magic. This is what +// LocalRandomAccessReadFile::open() requires (>= MIN_FILE_SIZE bytes, magic at +// both ends) and also what RestorableTsFileIOWriter's self check treats as +// complete. // // Built on first use rather than at file scope: VERSION_NUM_BYTE is defined in // another translation unit, so a file-scope initializer would depend on static @@ -147,7 +147,7 @@ bool CreateFixtureByWidePath() { } bool MinimalTsFileIsUnchanged() { - storage::ReadFile read_file; + storage::LocalRandomAccessReadFile read_file; std::string buf(MinimalTsFile().size(), '\0'); int32_t read_len = 0; const bool opened = read_file.open(Utf8Name()) == E_OK; @@ -212,22 +212,24 @@ TEST_F(Utf8PathTest, ExistingUtf8PathIsRejectedByCWrapper) { EXPECT_TRUE(MinimalTsFileIsUnchanged()); } -// ReadFile must find a file that exists on disk under a non-ASCII name. -TEST_F(Utf8PathTest, ReadFileOpensUtf8Path) { +// LocalRandomAccessReadFile must find a file that exists on disk under a +// non-ASCII name. +TEST_F(Utf8PathTest, LocalRandomAccessReadFileOpensUtf8Path) { ASSERT_TRUE(CreateFixtureByWidePath()); ASSERT_TRUE(ExistsByWidePath()); - storage::ReadFile read_file; + storage::LocalRandomAccessReadFile read_file; EXPECT_EQ(read_file.open(Utf8Name()), E_OK) << "an existing file with a non-ASCII name could not be opened"; EXPECT_TRUE(read_file.is_opened()); read_file.close(); } -// Round trip through both classes: what WriteFile wrote, ReadFile must read. -// The wide-path check matters even though the round trip alone would succeed -// while both sides are equally broken -- a consistently mangled name still -// round trips, so only the on-disk name proves the bytes were honoured. +// Round trip through both classes: what WriteFile wrote, +// LocalRandomAccessReadFile must read. The wide-path check matters even though +// the round trip alone would succeed while both sides are equally broken -- a +// consistently mangled name still round trips, so only the on-disk name proves +// the bytes were honoured. TEST_F(Utf8PathTest, Utf8PathRoundTripsBetweenWriteAndRead) { storage::WriteFile write_file; ASSERT_EQ(write_file.create(Utf8Name(), O_WRONLY | O_CREAT | O_TRUNC, 0666), @@ -241,7 +243,7 @@ TEST_F(Utf8PathTest, Utf8PathRoundTripsBetweenWriteAndRead) { ASSERT_TRUE(ExistsByWidePath()) << "the round trip used a name that is not the requested UTF-8 path"; - storage::ReadFile read_file; + storage::LocalRandomAccessReadFile read_file; ASSERT_EQ(read_file.open(Utf8Name()), E_OK); EXPECT_EQ(read_file.file_size(), static_cast(content.size())); diff --git a/cpp/test/reader/chunk_reader_resource_test.cc b/cpp/test/reader/chunk_reader_resource_test.cc index 11d456a4d..e9d4193c0 100644 --- a/cpp/test/reader/chunk_reader_resource_test.cc +++ b/cpp/test/reader/chunk_reader_resource_test.cc @@ -26,6 +26,7 @@ #include #include "common/allocator/alloc_base.h" +#include "file/local_random_access_read_file.h" #include "reader/aligned_chunk_reader.h" namespace storage { @@ -49,7 +50,7 @@ class TempTsFile { std::string path_; }; -AlignedChunkReader* allocate_reader(ReadFile* read_file) { +AlignedChunkReader* allocate_reader(LocalRandomAccessReadFile* read_file) { void* memory = common::mem_alloc(sizeof(AlignedChunkReader), common::MOD_CHUNK_READER); if (memory == nullptr) { @@ -72,7 +73,7 @@ void free_reader(AlignedChunkReader* reader) { TEST(ChunkReaderResourceTest, AlignedInitialShortReadReleasesBuffer) { TempTsFile temp_file("aligned_initial_short_read.tsfile"); - ReadFile read_file; + LocalRandomAccessReadFile read_file; ASSERT_EQ(read_file.open(temp_file.path()), common::E_OK); AlignedChunkReader* reader = allocate_reader(&read_file); ASSERT_NE(reader, nullptr); @@ -93,7 +94,7 @@ TEST(ChunkReaderResourceTest, AlignedInitialShortReadReleasesBuffer) { TEST(ChunkReaderResourceTest, MultiAlignedInitialShortReadReleasesBuffer) { TempTsFile temp_file("multi_aligned_initial_short_read.tsfile"); - ReadFile read_file; + LocalRandomAccessReadFile read_file; ASSERT_EQ(read_file.open(temp_file.path()), common::E_OK); AlignedChunkReader* reader = allocate_reader(&read_file); ASSERT_NE(reader, nullptr); diff --git a/cpp/test/reader/table_view/table_model_encoding_compression_compatibility_test.cc b/cpp/test/reader/table_view/table_model_encoding_compression_compatibility_test.cc index c52cd0515..09ed3e49b 100644 --- a/cpp/test/reader/table_view/table_model_encoding_compression_compatibility_test.cc +++ b/cpp/test/reader/table_view/table_model_encoding_compression_compatibility_test.cc @@ -46,6 +46,7 @@ #include "common/db_common.h" #include "common/schema.h" #include "common/tablet.h" +#include "file/local_random_access_read_file.h" #include "file/tsfile_io_reader.h" #include "file/write_file.h" #include "reader/table_result_set.h" @@ -655,7 +656,7 @@ void AssertOnWireCodec(const std::string& directory, ASSERT_NE(nullptr, value_chunks); ASSERT_GT(value_chunks->size(), 0U); - ReadFile read_file; + LocalRandomAccessReadFile read_file; ASSERT_EQ(E_OK, read_file.open(JoinPath(directory, fixture_case.file_name))); for (auto cursor = value_chunks->begin(); cursor != value_chunks->end(); diff --git a/cpp/test/reader/table_view/tsfile_reader_table_test.cc b/cpp/test/reader/table_view/tsfile_reader_table_test.cc index be0a6f64c..d8ce24409 100644 --- a/cpp/test/reader/table_view/tsfile_reader_table_test.cc +++ b/cpp/test/reader/table_view/tsfile_reader_table_test.cc @@ -465,7 +465,7 @@ TEST_F(TsFileTableReaderTest, ReadNonExistColumn) { tsfile_table_writer->flush(); tsfile_table_writer->close(); - TsFileReader reader = TsFileReader(); + TsFileReader reader; reader.open(write_file_.get_file_path()); ResultSet* ret = nullptr; std::vector column_names = {"non-exist-column"}; @@ -501,7 +501,7 @@ TEST_F(TsFileTableReaderTest, TestDecoder) { ASSERT_EQ(ret_, common::E_OK); ret_ = tsfile_table_writer_->close(); ASSERT_EQ(ret_, common::E_OK); - TsFileReader reader = TsFileReader(); + TsFileReader reader; reader.open(write_file_.get_file_path()); ResultSet* ret = nullptr; int ret_value = diff --git a/cpp/test/reader/tsfile_reader_test.cc b/cpp/test/reader/tsfile_reader_test.cc index 80aa161da..314d69c6d 100644 --- a/cpp/test/reader/tsfile_reader_test.cc +++ b/cpp/test/reader/tsfile_reader_test.cc @@ -22,9 +22,13 @@ #include #include +#include +#include #include +#include #include #include +#include #include #include @@ -32,6 +36,7 @@ #include "common/schema.h" #include "common/tablet.h" #include "common/tsblock/tsblock.h" +#include "file/random_access_read_file.h" #include "file/tsfile_io_reader.h" #include "file/tsfile_io_writer.h" #include "file/write_file.h" @@ -45,6 +50,15 @@ using namespace storage; using namespace common; +static_assert(!std::is_copy_constructible::value, + "TsFileReader must not be copy constructible"); +static_assert(!std::is_copy_assignable::value, + "TsFileReader must not be copy assignable"); +static_assert(!std::is_move_constructible::value, + "TsFileReader must not be move constructible"); +static_assert(!std::is_move_assignable::value, + "TsFileReader must not be move assignable"); + TEST(TsFileSeriesScanIteratorTest, ConsumeRowOffsetSaturates) { storage::TsFileSeriesScanIterator ssi; ssi.set_row_range(/*offset=*/10, /*limit=*/-1); @@ -140,6 +154,345 @@ class TsFileReaderTest : public ::testing::Test { } }; +namespace { + +class InMemoryRandomAccessReadFile : public RandomAccessReadFile { + public: + explicit InMemoryRandomAccessReadFile(std::vector bytes) + : bytes_(std::move(bytes)), opened_(true), name_("memory://test") {} + + bool is_opened() const override { return opened_; } + + int64_t file_size() const override { + return static_cast(bytes_.size()); + } + + const std::string& file_path() const override { return name_; } + + int generation(uint64_t& size, uint64_t& fingerprint) const override { + if (!opened_) { + return E_FILE_READ_ERR; + } + size = bytes_.size(); + fingerprint = 0; + return E_OK; + } + + int read(int64_t offset, char* buffer, int32_t size, + int32_t& read_size) override { + read_size = 0; + if (!opened_ || offset < 0 || size < 0 || + (buffer == nullptr && size > 0)) { + return E_INVALID_ARG; + } + if (size == 0 || static_cast(offset) >= bytes_.size()) { + return E_OK; + } + const size_t available = bytes_.size() - static_cast(offset); + read_size = static_cast( + std::min(available, static_cast(size))); + std::memcpy(buffer, bytes_.data() + offset, + static_cast(read_size)); + return E_OK; + } + + void close() override { opened_ = false; } + + private: + std::vector bytes_; + bool opened_; + std::string name_; +}; + +class ShortMetadataReadFile : public InMemoryRandomAccessReadFile { + public: + explicit ShortMetadataReadFile(const std::vector& bytes) + : InMemoryRandomAccessReadFile(bytes) {} + + int read(int64_t offset, char* buffer, int32_t size, + int32_t& read_size) override { + int ret = + InMemoryRandomAccessReadFile::read(offset, buffer, size, read_size); + if (ret == E_OK && offset >= metadata_offset && + ++read_count == short_read_at && read_size > 0) { + // Keep the entire buffer initialized with valid bytes so a missing + // length check fails deterministically instead of crashing in the + // parser. Only the reported prefix is valid under the read + // contract. + read_size = report_zero ? 0 : read_size - 1; + } + return ret; + } + + int read_count = 0; + int short_read_at = -1; + bool report_zero = false; + int64_t metadata_offset = 0; +}; + +} // namespace + +class MetadataReadLengthTest : public TsFileReaderTest, + public ::testing::WithParamInterface { + protected: + void SetUp() override { + TsFileReaderTest::SetUp(); + saved_index_degree_ = g_config_value_.max_degree_of_index_node_; + ASSERT_EQ(set_max_degree_of_index_node(2), E_OK); + } + + void TearDown() override { + set_max_degree_of_index_node(saved_index_degree_); + TsFileReaderTest::TearDown(); + } + + uint32_t saved_index_degree_ = 0; +}; + +TEST_P(MetadataReadLengthTest, RejectsIncompleteRangesBeforeParsing) { + for (const bool aligned : {false, true}) { + const std::string device = aligned ? "root.aligned" : "root.unaligned"; + TsRecord record(100, device); + for (int column = 0; column < GetParam(); ++column) { + const std::string name = "value" + std::to_string(column); + MeasurementSchema schema(name, INT32, PLAIN, UNCOMPRESSED); + ASSERT_EQ(aligned + ? tsfile_writer_->register_aligned_timeseries(device, + schema) + : tsfile_writer_->register_timeseries(device, schema), + E_OK); + record.add_point(name, static_cast(42)); + } + ASSERT_EQ(aligned ? tsfile_writer_->write_record_aligned(record) + : tsfile_writer_->write_record(record), + E_OK); + } + ASSERT_EQ(tsfile_writer_->flush(), E_OK); + ASSERT_EQ(tsfile_writer_->close(), E_OK); + std::ifstream input(file_name_, std::ios::binary); + ASSERT_TRUE(input.is_open()); + const std::vector bytes((std::istreambuf_iterator(input)), + std::istreambuf_iterator()); + + for (const bool aligned : {false, true}) { + auto device = std::make_shared( + aligned ? "root.aligned" : "root.unaligned"); + // Exercise the public metadata entry points, including both single + // series (cached node) and aligned multi-series allocation. + for (const std::string operation : + {"device_node", "by_offset", "all_series", "selected_series", + "single_scan", "multi_scan"}) { + if (!aligned && operation == "multi_scan") continue; + int metadata_reads = 0; + for (int fail_at = 0; fail_at <= metadata_reads; ++fail_at) { + for (const bool report_zero : {false, true}) { + SCOPED_TRACE(::testing::Message() + << "aligned=" << aligned << " operation=" + << operation << " read=" << fail_at + << " zero=" << report_zero); + ShortMetadataReadFile source(bytes); + TsFileIOReader reader; + ASSERT_EQ(reader.init(&source), E_OK); + TsFileMeta* meta = nullptr; + ASSERT_EQ(reader.get_tsfile_meta(meta), E_OK); + std::shared_ptr entry; + int64_t end_offset = 0; + ASSERT_EQ(reader.load_device_index_entry( + std::make_shared(device), + entry, end_offset), + E_OK); + source.read_count = 0; + source.metadata_offset = meta->meta_offset_; + source.short_read_at = fail_at; + source.report_zero = report_zero; + PageArena pa; + pa.init(512, MOD_TSFILE_READER); + std::vector indexes; + int ret = E_OK; + if (operation == "device_node") { + MetaIndexNode* node = nullptr; + ret = reader.read_device_meta_index( + entry->get_offset(), end_offset, pa, node, true); + if (node != nullptr) node->~MetaIndexNode(); + } else if (operation == "by_offset") { + ret = reader.get_device_timeseries_meta_by_offset( + entry->get_offset(), end_offset, indexes, pa); + } else if (operation == "all_series") { + ret = + reader + .get_device_timeseries_meta_without_chunk_meta( + device, indexes, pa); + } else if (operation == "selected_series") { + indexes.resize(1, nullptr); + ret = reader.get_timeseries_indexes(device, {"value0"}, + indexes, pa); + } else { + TsFileSeriesScanIterator* ssi = nullptr; + ret = operation == "single_scan" + ? reader.alloc_ssi(device, "value0", ssi, pa) + : reader.alloc_multi_ssi(device, {"value0"}, + ssi, pa); + reader.revert_ssi(ssi); + } + if (fail_at == 0) { + ASSERT_EQ(ret, E_OK); + metadata_reads = source.read_count; + ASSERT_GT(metadata_reads, 0); + } else { + EXPECT_EQ(ret, E_FILE_READ_ERR); + EXPECT_EQ(source.read_count, fail_at); + } + } + } + } + } +} + +// One column keeps the measurement root a leaf; five columns force internal +// nodes and exercise recursive traversal and aligned time-index descent. +INSTANTIATE_TEST_SUITE_P(LeafAndInternalIndexes, MetadataReadLengthTest, + ::testing::Values(1, 5)); + +class DeviceIndexReadTest : public TsFileReaderTest { + protected: + void SetUp() override { + TsFileReaderTest::SetUp(); + saved_index_degree_ = g_config_value_.max_degree_of_index_node_; + ASSERT_EQ(set_max_degree_of_index_node(2), E_OK); + // Five devices force an internal device index, not just measurement + // indexes within one device. + for (int i = 0; i < 5; ++i) { + const std::string device = "root.sg.d" + std::to_string(i); + ASSERT_EQ(tsfile_writer_->register_timeseries( + device, MeasurementSchema("value", INT32, PLAIN, + UNCOMPRESSED)), + E_OK); + TsRecord record(100, device); + record.add_point("value", static_cast(42)); + ASSERT_EQ(tsfile_writer_->write_record(record), E_OK); + } + ASSERT_EQ(tsfile_writer_->flush(), E_OK); + ASSERT_EQ(tsfile_writer_->close(), E_OK); + std::ifstream input(file_name_, std::ios::binary); + ASSERT_TRUE(input.is_open()); + bytes_.assign(std::istreambuf_iterator(input), + std::istreambuf_iterator()); + } + + void TearDown() override { + set_max_degree_of_index_node(saved_index_degree_); + TsFileReaderTest::TearDown(); + std::remove(file_name_.c_str()); + } + + void open_reader(TsFileReader& reader, ShortMetadataReadFile*& source) { + source = new ShortMetadataReadFile(bytes_); + ASSERT_EQ(reader.open(std::unique_ptr(source)), + E_OK); + // Warm the file footer so subsequent reads start at the device index. + std::vector> devices; + ASSERT_EQ(reader.get_all_devices(devices), E_OK); + ASSERT_EQ(devices.size(), 5u); + source->read_count = 0; + } + + uint32_t saved_index_degree_ = 0; + std::vector bytes_; +}; + +TEST_F(DeviceIndexReadTest, MetadataEnumerationRejectsShortDeviceIndex) { + for (int fail_at : {0, 1, 2}) { + for (bool report_zero : {false, true}) { + SCOPED_TRACE(::testing::Message() << fail_at << ":" << report_zero); + TsFileReader reader; + ShortMetadataReadFile* source = nullptr; + ASSERT_NO_FATAL_FAILURE(open_reader(reader, source)); + source->short_read_at = fail_at; + source->report_zero = report_zero; + auto metadata = reader.get_timeseries_metadata(); + if (fail_at == 0) { + ASSERT_EQ(metadata.size(), 5u); + } else { + // This legacy map-returning API has no error-code output, but + // must stop before parsing a short node or returning metadata. + EXPECT_TRUE(metadata.empty()); + EXPECT_EQ(source->read_count, fail_at); + } + } + } +} + +TEST_F(DeviceIndexReadTest, TreeTableQueryPropagatesDeviceAndSchemaReadErrors) { + for (const auto& measurements : + {std::vector{}, std::vector{"value"}}) { + for (bool fail_schema : {false, true}) { + for (bool report_zero : {false, true}) { + SCOPED_TRACE(::testing::Message() + << measurements.size() << ":" << fail_schema << ":" + << report_zero); + TsFileReader reader; + ShortMetadataReadFile* source = nullptr; + ASSERT_NO_FATAL_FAILURE(open_reader(reader, source)); + std::vector> devices; + ASSERT_EQ(reader.get_all_devices(devices), E_OK); + const int device_reads = source->read_count; + ASSERT_GT(device_reads, 0); + source->read_count = 0; + source->short_read_at = fail_schema ? device_reads + 1 : 1; + source->report_zero = report_zero; + ResultSet* result = nullptr; + EXPECT_EQ( + reader.query_table_on_tree(measurements, 0, 200, result), + E_FILE_READ_ERR); + EXPECT_EQ(source->read_count, source->short_read_at); + if (result != nullptr) reader.destroy_query_data_set(result); + } + } + } +} + +TEST_F(TsFileReaderTest, ReadsThroughRandomAccessReadFile) { + const std::string device = "root.sg.device"; + const std::string measurement = "temperature"; + ASSERT_EQ(tsfile_writer_->register_timeseries( + device, MeasurementSchema(measurement, TSDataType::INT32, + TSEncoding::PLAIN, + CompressionType::UNCOMPRESSED)), + E_OK); + TsRecord record(100, device); + record.add_point(measurement, static_cast(42)); + ASSERT_EQ(tsfile_writer_->write_record(record), E_OK); + ASSERT_EQ(tsfile_writer_->flush(), E_OK); + ASSERT_EQ(tsfile_writer_->close(), E_OK); + + std::ifstream input(file_name_, std::ios::binary); + ASSERT_TRUE(input.is_open()); + std::vector bytes((std::istreambuf_iterator(input)), + std::istreambuf_iterator()); + ASSERT_FALSE(bytes.empty()); + + std::unique_ptr source( + new InMemoryRandomAccessReadFile(std::move(bytes))); + TsFileReader reader; + ASSERT_EQ(reader.open(std::move(source)), E_OK); + EXPECT_EQ(reader.get_file_version(), + static_cast(VERSION_NUM_BYTE)); + + std::vector paths = {device + "." + measurement}; + ResultSet* result = nullptr; + ASSERT_EQ(reader.query(paths, 0, 200, result), E_OK); + ASSERT_NE(result, nullptr); + bool has_next = false; + ASSERT_EQ(result->next(has_next), E_OK); + ASSERT_TRUE(has_next); + EXPECT_EQ(result->get_value(2), 42); + ASSERT_EQ(result->next(has_next), E_OK); + EXPECT_FALSE(has_next); + + reader.destroy_query_data_set(result); + reader.close(); +} + TEST_F(TsFileReaderTest, ResultSetMetadata) { std::string device_path = "device1"; std::string measurement_name = "temperature"; @@ -355,6 +708,26 @@ TEST_F(TsFileReaderTest, GetTimeseriesSchemaUsesLastChunkCodec) { EXPECT_EQ(schemas[0].encoding_, TS_2DIFF); EXPECT_EQ(schemas[0].compression_type_, UNCOMPRESSED); reader.close(); + + // The stored codec must also be read through non-local backends. + std::ifstream input(path, std::ios::binary); + ASSERT_TRUE(input.is_open()); + std::vector bytes((std::istreambuf_iterator(input)), + std::istreambuf_iterator()); + input.close(); + std::unique_ptr source( + new InMemoryRandomAccessReadFile(std::move(bytes))); + ASSERT_EQ(reader.open(std::move(source)), E_OK); + schemas.clear(); + ASSERT_EQ(reader.get_timeseries_schema( + std::make_shared(device_name), schemas), + E_OK); + ASSERT_EQ(schemas.size(), 1u); + EXPECT_EQ(schemas[0].measurement_name_, measurement_name); + EXPECT_EQ(schemas[0].data_type_, INT32); + EXPECT_EQ(schemas[0].encoding_, TS_2DIFF); + EXPECT_EQ(schemas[0].compression_type_, UNCOMPRESSED); + reader.close(); remove(path.c_str()); } diff --git a/cpp/test/tools/model_format_e2e_test.cc b/cpp/test/tools/model_format_e2e_test.cc index 95a239b4f..f782ce569 100644 --- a/cpp/test/tools/model_format_e2e_test.cc +++ b/cpp/test/tools/model_format_e2e_test.cc @@ -479,8 +479,9 @@ TEST(IndependentFixtures, EmptyTreeAndInputFailuresHaveExactDiagnostics) { std::remove(unsupported.c_str()); // Directories are special input paths, not TsFiles. The check is made in - // ReadFile before magic parsing, so this diagnostic is stable across all - // read commands and does not leak parser output to stdout. + // LocalRandomAccessReadFile before magic parsing, so this diagnostic is + // stable across all read commands and does not leak parser output to + // stdout. for (const std::string& command : commands) { std::remove("tsfile_cli_input_error.csv"); expect_cli_exact(input_args(command, "."), 2, "", diff --git a/cpp/test/writer/table_view/tsfile_writer_table_test.cc b/cpp/test/writer/table_view/tsfile_writer_table_test.cc index 2dd9b5643..eeb82dbf4 100644 --- a/cpp/test/writer/table_view/tsfile_writer_table_test.cc +++ b/cpp/test/writer/table_view/tsfile_writer_table_test.cc @@ -194,7 +194,7 @@ TEST_F(TsFileWriterTableTest, WithoutTagAndMultiPage) { tsfile_table_writer->flush(); tsfile_table_writer->close(); - TsFileReader reader = TsFileReader(); + TsFileReader reader; reader.open(write_file_.get_file_path()); ResultSet* ret = nullptr; int ret_value = reader.query("test_table", {"value"}, 0, 50, ret); @@ -358,7 +358,7 @@ TEST_F(TsFileWriterTableTest, EmptyTagWrite) { tsfile_table_writer->flush(); tsfile_table_writer->close(); - TsFileReader reader = TsFileReader(); + TsFileReader reader; reader.open(write_file_.get_file_path()); ResultSet* ret = nullptr; int ret_value = @@ -462,7 +462,7 @@ TEST_F(TsFileWriterTableTest, WriteAndReadSimple) { tsfile_table_writer->flush(); tsfile_table_writer->close(); - TsFileReader reader = TsFileReader(); + TsFileReader reader; reader.open(write_file_.get_file_path()); ResultSet* ret = nullptr; std::vector column_names = {"device", "VALUE"}; @@ -588,7 +588,7 @@ TEST_F(TsFileWriterTableTest, WriteWithNullAndEmptyTag) { delete table_schema; - auto reader = TsFileReader(); + TsFileReader reader; reader.open(write_file_.get_file_path()); ResultSet* ret = nullptr; int ret_value = @@ -762,7 +762,7 @@ TEST_F(TsFileWriterTableTest, WriteDataWithEmptyField) { delete table_schema; - auto reader = TsFileReader(); + TsFileReader reader; reader.open(write_file_.get_file_path()); ResultSet* ret = nullptr; int ret_value = reader.query( @@ -870,7 +870,7 @@ TEST_F(TsFileWriterTableTest, MultiDatatypes) { delete table_schema; - auto reader = TsFileReader(); + TsFileReader reader; reader.open(write_file_.get_file_path()); ResultSet* ret = nullptr; int ret_value = reader.query("testTable", measurement_names, 0, 100, ret); @@ -971,7 +971,7 @@ TEST_F(TsFileWriterTableTest, DiffCodecTypes) { delete table_schema; - auto reader = TsFileReader(); + TsFileReader reader; reader.open(write_file_.get_file_path()); ResultSet* ret = nullptr; int ret_value = reader.query("testTable", measurement_names, 0, 100, ret); @@ -1088,7 +1088,7 @@ TEST_F(TsFileWriterTableTest, EncodingConfigIntegration) { ASSERT_EQ(tsfile_table_writer->close(), E_OK); // 5. Verify read data matches what was written - auto reader = TsFileReader(); + TsFileReader reader; reader.open(write_file_.get_file_path()); ResultSet* ret = nullptr; int ret_value = diff --git a/python/README-zh.md b/python/README-zh.md index 9224a8bcd..9d322370d 100644 --- a/python/README-zh.md +++ b/python/README-zh.md @@ -104,3 +104,21 @@ set_tsfile_config({"file_read_backend_": FileReadBackend.AUTO}) 该配置只影响之后打开的 reader。通过内存映射后端打开文件期间,请勿修改或 截断该文件。 + +## 可定位的二进制文件对象 + +`TsFileReader` 也可以直接接收可定位的二进制文件对象。因此,通过 `fsspec` +等库打开远程文件后,无需先把整个文件复制到本地即可读取。 + +```python +import fsspec +from tsfile import TsFileReader + +with fsspec.open("s3://bucket/example.tsfile", "rb") as source: + with TsFileReader(source) as reader: + result = reader.query_table("table_name", ["column_name"]) +``` + +文件对象必须提供 `seek()`、`tell()` 和带明确长度的二进制 `read(size)`。 +该对象仍由调用方管理:`TsFileReader` 会在使用期间保持其存活、恢复其游标位置, +但不会关闭它。本地 `FileReadBackend` 配置不适用于文件对象。 diff --git a/python/README.md b/python/README.md index d02d9c54c..22e063166 100644 --- a/python/README.md +++ b/python/README.md @@ -99,3 +99,23 @@ set_tsfile_config({"file_read_backend_": FileReadBackend.AUTO}) The setting only affects readers opened afterward. Do not modify or truncate a file while it is open through the memory-mapped backend. + +## Seekable binary file objects + +`TsFileReader` also accepts a seekable binary file object. This allows remote +files opened by libraries such as `fsspec` to be read without first copying the +whole file to local storage. + +```python +import fsspec +from tsfile import TsFileReader + +with fsspec.open("s3://bucket/example.tsfile", "rb") as source: + with TsFileReader(source) as reader: + result = reader.query_table("table_name", ["column_name"]) +``` + +The object must provide `seek()`, `tell()`, and explicitly sized binary +`read(size)` operations. The caller owns the object: `TsFileReader` keeps it +alive while in use, preserves its cursor position, and does not close it. +The local `FileReadBackend` setting does not apply to file objects. diff --git a/python/pom.xml b/python/pom.xml index 047315215..f6332b3ea 100644 --- a/python/pom.xml +++ b/python/pom.xml @@ -196,6 +196,9 @@ *.h *.cpp + + python_random_access_read_file.h + tsfile.egg-info diff --git a/python/pyproject.toml b/python/pyproject.toml index 8cd7d00c0..cd28b3d54 100644 --- a/python/pyproject.toml +++ b/python/pyproject.toml @@ -65,6 +65,7 @@ include = ["tsfile*"] [tool.setuptools.package-data] tsfile = [ + "python_random_access_read_file.h", "*.pxd", "*.pxi", "*.so", diff --git a/python/setup.py b/python/setup.py index e74c3296f..c7b4cceaa 100644 --- a/python/setup.py +++ b/python/setup.py @@ -274,7 +274,12 @@ def finalize_options(self): exts = [ Extension("tsfile.dataset._merge", ["tsfile/dataset/_merge.pyx"], **merge_common), Extension("tsfile.tsfile_py_cpp", ["tsfile/tsfile_py_cpp.pyx"], **common), - Extension("tsfile.tsfile_reader", ["tsfile/tsfile_reader.pyx"], **common), + Extension( + "tsfile.tsfile_reader", + ["tsfile/tsfile_reader.pyx", "tsfile/python_random_access_read_file.cc"], + depends=["tsfile/python_random_access_read_file.h"], + **common, + ), Extension("tsfile.tsfile_writer", ["tsfile/tsfile_writer.pyx"], **common), ] diff --git a/python/tests/test_reader_sources.py b/python/tests/test_reader_sources.py new file mode 100644 index 000000000..24961f85e --- /dev/null +++ b/python/tests/test_reader_sources.py @@ -0,0 +1,463 @@ +# Licensed to the Apache Software Foundation (ASF) under one +# or more contributor license agreements. See the NOTICE file +# distributed with this work for additional information +# regarding copyright ownership. The ASF licenses this file +# to you under the Apache License, Version 2.0 (the +# "License"); you may not use this file except in compliance +# with the License. You may obtain a copy of the License at +# +# http://www.apache.org/licenses/LICENSE-2.0 +# +# Unless required by applicable law or agreed to in writing, +# software distributed under the License is distributed on an +# "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +# KIND, either express or implied. See the License for the +# specific language governing permissions and limitations +# under the License. +# + +"""Reader input sources: paths, file objects, ownership, and read failures.""" + +import gc +import io +import os +import subprocess +import sys +import weakref +from pathlib import Path +from typing import Optional + +import pytest +import numpy as np + +from tsfile import ( + Field, + RowRecord, + TSDataType, + TimeseriesSchema, + TsFileReader, + TsFileWriter, +) +from tsfile.exceptions import FileOpenError, FileReadError + +RESOURCES = Path(__file__).parent / "resources" + + +class TrackingBytesIO(io.BytesIO): + def __init__(self, data: bytes, max_chunk_size: Optional[int] = None): + super().__init__(data) + self.max_chunk_size = max_chunk_size + self.read_sizes = [] + self.close_calls = 0 + + def read(self, size=-1): + if size < 0: + raise AssertionError("TsFileReader must use explicitly sized reads") + self.read_sizes.append(size) + if self.max_chunk_size is not None: + size = min(size, self.max_chunk_size) + return super().read(size) + + def close(self): + self.close_calls += 1 + super().close() + + +def collect_table_rows(source): + with TsFileReader(source) as reader: + result = reader.query_table("test", ["s0", "s2"]) + try: + rows = [] + while result.next(): + rows.append(tuple(result.get_value_by_index(i) for i in range(1, 4))) + return rows + finally: + result.close() + + +def collect_tree_rows(source): + with TsFileReader(source) as reader: + result = reader.query_timeseries( + "root.ln.wf01.wt01", + ["temperature", "status"], + 0, + (1 << 63) - 1, + ) + try: + rows = [] + while result.next(): + rows.append(tuple(result.get_value_by_index(i) for i in range(1, 4))) + return rows + finally: + result.close() + + +class CustomPath: + def __init__(self, path): + self.path = path + + def __fspath__(self): + return self.path + + +@pytest.mark.parametrize( + "make_source", + [str, Path, os.fsencode, CustomPath, lambda p: CustomPath(os.fsencode(p)), np.str_], + ids=["str", "pathlib", "bytes", "pathlike-str", "pathlike-bytes", "numpy-str"], +) +def test_path_sources_match_string_path(make_source): + path = RESOURCES / "simple_table_t1.tsfile" + expected = collect_table_rows(str(path)) + assert expected + assert collect_table_rows(make_source(str(path))) == expected + + +@pytest.mark.skipif(os.name != "posix", reason="POSIX filesystem byte paths") +@pytest.mark.parametrize( + "make_source", + [lambda p: p, CustomPath, lambda p: Path(os.fsdecode(p))], + ids=["bytes", "pathlike-bytes", "pathlib"], +) +def test_non_utf8_missing_path_reaches_native_open(tmp_path, make_source): + path = os.fsencode(tmp_path) + b"/missing-\xff.tsfile" + with pytest.raises(FileOpenError): + TsFileReader(make_source(path)) + + +@pytest.mark.skipif( + not sys.platform.startswith("linux"), + reason="Linux filenames permit arbitrary non-NUL bytes; macOS rejects them", +) +def test_non_utf8_filesystem_path_is_preserved(tmp_path): + fixture = RESOURCES / "simple_table_t1.tsfile" + path = os.fsencode(tmp_path) + b"/sample-\xff.tsfile" + with open(path, "wb") as output: + output.write(fixture.read_bytes()) + assert collect_table_rows(path) == collect_table_rows(str(fixture)) + + +def test_file_object_table_query_handles_short_reads_and_preserves_cursor(): + path = RESOURCES / "simple_table_t1.tsfile" + expected = collect_table_rows(str(path)) + source = TrackingBytesIO(path.read_bytes(), max_chunk_size=7) + source.seek(11) + + actual = collect_table_rows(source) + + assert actual == expected + assert source.tell() == 11 + assert source.read_sizes + assert source.close_calls == 0 + + +def test_file_object_initializes_native_runtime_in_fresh_process(): + path = RESOURCES / "simple_table_t1.tsfile" + script = """ +import io +import sys +from pathlib import Path + +from tsfile import TsFileReader + +source = io.BytesIO(Path(sys.argv[1]).read_bytes()) +source.seek(11) +with TsFileReader(source) as reader: + result = reader.query_table("test", ["s0"]) + try: + assert result.next() + assert result.get_value_by_index(1) == 1760106020000 + assert result.get_value_by_index(2) == "a" + finally: + result.close() +assert source.tell() == 11 +""" + + completed = subprocess.run( + [sys.executable, "-c", script, str(path)], + cwd=Path(__file__).parents[1], + capture_output=True, + text=True, + ) + + assert completed.returncode == 0, completed.stderr + + +def test_file_object_multi_field_query_returns_expected_row(): + path = RESOURCES / "simple_table_t1.tsfile" + source = io.BytesIO(path.read_bytes()) + source.seek(19) + + with TsFileReader(source) as reader: + result = reader.query_table("test", ["s2", "s3"]) + try: + assert result.next() + assert result.get_value_by_index(1) == 1760106020000 + assert result.get_value_by_index(2) == 1010 + assert result.get_value_by_index(3) == 2.0 + finally: + result.close() + + assert source.tell() == 19 + + +def test_reader_retains_source_until_close_then_releases_it(): + path = RESOURCES / "simple_table_t1.tsfile" + source = TrackingBytesIO(path.read_bytes()) + source_ref = weakref.ref(source) + reader = TsFileReader(source) + + del source + gc.collect() + assert source_ref() is not None + + reader.close() + gc.collect() + assert source_ref() is None + + +def test_file_object_tree_query_handles_short_reads(): + path = RESOURCES / "simple_tree.tsfile" + expected = collect_tree_rows(str(path)) + source = TrackingBytesIO(path.read_bytes(), max_chunk_size=5) + + assert collect_tree_rows(source) == expected + + +def test_source_without_file_methods_is_rejected(): + with pytest.raises(TypeError, match="seekable binary file object"): + TsFileReader(object()) + + +@pytest.mark.parametrize( + "failure,fail_restore", + [ + ("initial_tell", False), + ("seek_end", False), + ("size_tell", False), + ("restore", True), + ("seek_end", True), + ("size_tell", True), + ], +) +def test_size_probe_errors_map_to_file_open_error(failure, fail_restore): + class FailingProbeSource(TrackingBytesIO): + tell_calls = 0 + + def __init__(self, data): + super().__init__(data) + self.seek_calls = [] + + def tell(self): + self.tell_calls += 1 + if (failure == "initial_tell" and self.tell_calls == 1) or ( + failure == "size_tell" and self.tell_calls == 2 + ): + raise OSError("size probe tell failed") + return super().tell() + + def seek(self, offset, whence=os.SEEK_SET): + self.seek_calls.append((offset, whence)) + if whence == os.SEEK_SET and fail_restore: + raise RuntimeError("cursor restoration failed") + position = super().seek(offset, whence) + if whence == os.SEEK_END and failure == "seek_end": + # A source may move its cursor before reporting failure. + raise OSError("size probe seek failed") + return position + + data = (RESOURCES / "simple_table_t1.tsfile").read_bytes() + source = FailingProbeSource(data) + io.BytesIO.seek(source, 13) + + with pytest.raises(FileOpenError): + TsFileReader(source) + + assert not source.closed + assert source.close_calls == 0 + assert source.read_sizes == [] + if failure == "initial_tell": + assert source.seek_calls == [] + else: + assert source.seek_calls == [(0, os.SEEK_END), (13, os.SEEK_SET)] + assert io.BytesIO.tell(source) == (len(data) if fail_restore else 13) + + +def test_source_read_error_during_open_preserves_cursor(): + path = RESOURCES / "simple_table_t1.tsfile" + + class FailingBytesIO(TrackingBytesIO): + fail_reads = False + + def read(self, size=-1): + if self.fail_reads: + raise OSError("remote read failed") + return super().read(size) + + source = FailingBytesIO(path.read_bytes()) + source.seek(13) + source.fail_reads = True + + with pytest.raises(FileReadError): + TsFileReader(source) + + assert source.tell() == 13 + assert source.close_calls == 0 + + +def test_source_read_error_during_query_preserves_cursor_and_source(): + path = RESOURCES / "simple_table_t1.tsfile" + + class FailingBytesIO(TrackingBytesIO): + fail_reads = False + + def read(self, size=-1): + if self.fail_reads: + raise OSError("remote read failed") + return super().read(size) + + source = FailingBytesIO(path.read_bytes()) + source.seek(17) + + with TsFileReader(source) as reader: + source.fail_reads = True + with pytest.raises(FileReadError): + result = reader.query_table("test", ["s0"]) + try: + result.next() + finally: + result.close() + + assert source.tell() == 17 + assert source.close_calls == 0 + + +@pytest.mark.parametrize("measurements", [["temperature"], ["temperature", "status"]]) +@pytest.mark.parametrize("stage", ["initial_block", "next_block"]) +def test_tree_query_propagates_data_block_read_errors(tmp_path, measurements, stage): + class FailingBytesIO(TrackingBytesIO): + read_calls = 0 + fail_at = None + failed = False + + def read(self, size=-1): + self.read_calls += 1 + if self.fail_at is not None and self.read_calls >= self.fail_at: + self.failed = True + raise OSError("data block read failed") + return super().read(size) + + path = tmp_path / "multiple_chunks.tsfile" + device = "root.sg.device" + with TsFileWriter(str(path)) as writer: + for name in measurements: + writer.register_timeseries(device, TimeseriesSchema(name, TSDataType.INT64)) + for chunk in range(2): + for timestamp in range(chunk * 100, (chunk + 1) * 100): + writer.write_row_record( + RowRecord( + device, + timestamp, + [ + Field(name, timestamp, TSDataType.INT64) + for name in measurements + ], + ) + ) + writer.flush() + data = path.read_bytes() + + def query(reader): + return reader.query_timeseries(device, measurements, 0, (1 << 63) - 1) + + # Locate the initial and subsequent block reads without fixing their call + # numbers: opening and metadata reads may change as the reader evolves. + baseline = FailingBytesIO(data) + with TsFileReader(baseline) as reader: + with query(reader) as result: + initial_block_read = baseline.read_calls + rows = 0 + while result.next(): + rows += 1 + assert rows == 200 + assert baseline.read_calls > initial_block_read + + source = FailingBytesIO(data) + source.seek(17) + with TsFileReader(source) as reader: + if stage == "initial_block": + source.fail_at = initial_block_read + with pytest.raises(FileReadError): + with query(reader): + pass + else: + with query(reader) as result: + source.fail_at = source.read_calls + 1 + with pytest.raises(FileReadError): + while result.next(): + pass + # A failed result must not become a successful EOF on retry. + source.fail_at = None + with pytest.raises(FileReadError): + result.next() + assert source.failed + assert source.tell() == 17 + assert source.close_calls == 0 + + +@pytest.mark.parametrize("method", ["get_all_devices", "get_all_table_schemas"]) +def test_metadata_read_failure_is_not_an_empty_result(method): + class FailingBytesIO(TrackingBytesIO): + fail_reads = False + failed = False + + def read(self, size=-1): + if self.fail_reads: + self.failed = True + raise OSError("metadata read failed") + return super().read(size) + + source = FailingBytesIO((RESOURCES / "simple_table_t1.tsfile").read_bytes()) + source.seek(17) + with TsFileReader(source) as reader: + source.fail_reads = True + with pytest.raises(FileReadError): + getattr(reader, method)() + assert source.failed + assert source.tell() == 17 + assert source.close_calls == 0 + + +def test_device_index_read_failure_does_not_return_partial_devices(tmp_path): + class FailingBytesIO(TrackingBytesIO): + fail_index_reads = False + index_reads = 0 + failed = False + + def read(self, size=-1): + if self.fail_index_reads: + self.index_reads += 1 + if self.index_reads == 2: + self.failed = True + raise OSError("device index read failed") + return super().read(size) + + # More than two default index nodes (256 children each), so a failure in + # the second node must not be hidden by a successful read of the third. + path = tmp_path / "device_index.tsfile" + with TsFileWriter(str(path)) as writer: + for index in range(513): + device = f"root.sg.d{index:04d}" + writer.register_timeseries(device, TimeseriesSchema("s", TSDataType.INT64)) + writer.write_row_record( + RowRecord(device, 0, [Field("s", index, TSDataType.INT64)]) + ) + source = FailingBytesIO(path.read_bytes()) + source.seek(17) + with TsFileReader(source) as reader: + assert len(reader.get_all_devices()) == 513 + source.fail_index_reads = True + with pytest.raises(FileReadError): + reader.get_all_devices() + assert source.failed + assert source.tell() == 17 + assert source.close_calls == 0 diff --git a/python/tsfile/python_random_access_read_file.cc b/python/tsfile/python_random_access_read_file.cc new file mode 100644 index 000000000..9ed4ecd46 --- /dev/null +++ b/python/tsfile/python_random_access_read_file.cc @@ -0,0 +1,310 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ + +#include "python_random_access_read_file.h" + +#include +#include +#include +#include +#include +#include +#include +#include + +#include "file/random_access_read_file.h" +#include "reader/tsfile_reader.h" +#include "utils/errno_define.h" + +namespace { + +class PythonSourceLock { + public: + explicit PythonSourceLock(std::mutex& mutex) + : gil_state_(PyGILState_Ensure()), lock_(mutex, std::defer_lock) { + PyThreadState* thread_state = PyEval_SaveThread(); + lock_.lock(); + PyEval_RestoreThread(thread_state); + } + + ~PythonSourceLock() { + lock_.unlock(); + PyGILState_Release(gil_state_); + } + + private: + PyGILState_STATE gil_state_; + std::unique_lock lock_; +}; + +bool has_callable_attribute(PyObject* source, const char* name) { + PyObject* attribute = PyObject_GetAttrString(source, name); + if (attribute == nullptr) { + PyErr_Clear(); + return false; + } + const bool callable = PyCallable_Check(attribute) != 0; + Py_DECREF(attribute); + return callable; +} + +bool seek_to(PyObject* source, int64_t offset, int whence) { + PyObject* result = PyObject_CallMethod( + source, "seek", "Li", static_cast(offset), whence); + if (result == nullptr) { + return false; + } + Py_DECREF(result); + return true; +} + +bool tell_position(PyObject* source, int64_t& position) { + PyObject* result = + PyObject_CallMethod(source, const_cast("tell"), nullptr); + if (result == nullptr) { + return false; + } + const long long value = PyLong_AsLongLong(result); + Py_DECREF(result); + if (value == -1 && PyErr_Occurred()) { + return false; + } + position = static_cast(value); + return true; +} + +void restore_position_preserving_error(PyObject* source, int64_t position) { + PyObject* error_type = nullptr; + PyObject* error_value = nullptr; + PyObject* traceback = nullptr; + PyErr_Fetch(&error_type, &error_value, &traceback); + if (!seek_to(source, position, SEEK_SET)) { + PyErr_Clear(); + } + PyErr_Restore(error_type, error_value, traceback); +} + +std::string source_name(PyObject* source) { + std::string name(""); + PyObject* value = PyObject_GetAttrString(source, "name"); + if (value == nullptr) { + PyErr_Clear(); + return name; + } + if (PyUnicode_Check(value)) { + const char* text = PyUnicode_AsUTF8(value); + if (text != nullptr) { + name.assign(text); + } else { + PyErr_Clear(); + } + } else if (PyBytes_Check(value)) { + char* text = nullptr; + Py_ssize_t length = 0; + if (PyBytes_AsStringAndSize(value, &text, &length) == 0) { + name.assign(text, static_cast(length)); + } else { + PyErr_Clear(); + } + } + Py_DECREF(value); + return name; +} + +class PythonRandomAccessReadFile : public storage::RandomAccessReadFile { + public: + PythonRandomAccessReadFile(PyObject* source, int64_t size, std::string name) + : source_(source), size_(size), name_(std::move(name)) { + Py_INCREF(source_); + } + + ~PythonRandomAccessReadFile() override { close(); } + + bool is_opened() const override { + PythonSourceLock lock(mutex_); + return source_ != nullptr; + } + + int64_t file_size() const override { return size_; } + + const std::string& file_path() const override { return name_; } + + int generation(uint64_t& size, uint64_t& fingerprint) const override { + PythonSourceLock lock(mutex_); + if (source_ == nullptr) { + return common::E_FILE_READ_ERR; + } + size = static_cast(size_); + fingerprint = 0; + return common::E_OK; + } + + int read(int64_t offset, char* buffer, int32_t size, + int32_t& read_size) override { + read_size = 0; + if (offset < 0 || size < 0 || (buffer == nullptr && size > 0)) { + return common::E_INVALID_ARG; + } + + PythonSourceLock lock(mutex_); + if (source_ == nullptr) { + return common::E_FILE_READ_ERR; + } + if (size == 0 || offset >= size_) { + return common::E_OK; + } + + const int64_t available = size_ - offset; + const int32_t requested = static_cast( + std::min(available, static_cast(size))); + + int ret = common::E_OK; + int64_t saved_position = 0; + const bool has_saved_position = tell_position(source_, saved_position); + if (!has_saved_position || !seek_to(source_, offset, SEEK_SET)) { + PyErr_Clear(); + ret = common::E_FILE_READ_ERR; + } + + while (ret == common::E_OK && read_size < requested) { + const int32_t remaining = requested - read_size; + PyObject* chunk = + PyObject_CallMethod(source_, "read", "i", remaining); + if (chunk == nullptr) { + PyErr_Clear(); + ret = common::E_FILE_READ_ERR; + break; + } + + Py_buffer view; + if (PyObject_GetBuffer(chunk, &view, PyBUF_CONTIG_RO) != 0) { + PyErr_Clear(); + Py_DECREF(chunk); + ret = common::E_FILE_READ_ERR; + break; + } + if (view.len < 0 || view.len > remaining) { + ret = common::E_FILE_READ_ERR; + } else if (view.len == 0) { + PyBuffer_Release(&view); + Py_DECREF(chunk); + break; + } else { + std::memcpy(buffer + read_size, view.buf, + static_cast(view.len)); + read_size += static_cast(view.len); + } + PyBuffer_Release(&view); + Py_DECREF(chunk); + } + + if (has_saved_position && !seek_to(source_, saved_position, SEEK_SET)) { + PyErr_Clear(); + ret = common::E_FILE_READ_ERR; + } + return ret; + } + + void close() override { + PythonSourceLock lock(mutex_); + if (source_ == nullptr) { + return; + } + Py_DECREF(source_); + source_ = nullptr; + } + + private: + PyObject* source_; + int64_t size_; + std::string name_; + mutable std::mutex mutex_; +}; + +} // namespace + +void* create_tsfile_reader_from_python_file(PyObject* source, + int32_t* error_code) { + if (error_code == nullptr) { + PyErr_SetString(PyExc_ValueError, "error_code must not be null"); + return nullptr; + } + *error_code = common::E_OK; + if (source == nullptr || !has_callable_attribute(source, "seek") || + !has_callable_attribute(source, "tell") || + !has_callable_attribute(source, "read")) { + *error_code = common::E_INVALID_ARG; + PyErr_SetString(PyExc_TypeError, + "source must be a seekable binary file object"); + return nullptr; + } + + int64_t original_position = 0; + if (!tell_position(source, original_position)) { + *error_code = common::E_FILE_OPEN_ERR; + PyErr_Clear(); + return nullptr; + } + if (!seek_to(source, 0, SEEK_END)) { + *error_code = common::E_FILE_OPEN_ERR; + restore_position_preserving_error(source, original_position); + PyErr_Clear(); + return nullptr; + } + + int64_t size = 0; + if (!tell_position(source, size)) { + *error_code = common::E_FILE_OPEN_ERR; + restore_position_preserving_error(source, original_position); + PyErr_Clear(); + return nullptr; + } + if (!seek_to(source, original_position, SEEK_SET)) { + *error_code = common::E_FILE_OPEN_ERR; + PyErr_Clear(); + return nullptr; + } + if (size < 0) { + *error_code = common::E_INVALID_ARG; + PyErr_SetString(PyExc_ValueError, + "source returned a negative file size"); + return nullptr; + } + + try { + *error_code = storage::libtsfile_init(); + if (*error_code != common::E_OK) { + return nullptr; + } + std::unique_ptr random_access_read_file( + new PythonRandomAccessReadFile(source, size, source_name(source))); + std::unique_ptr reader( + new storage::TsFileReader()); + *error_code = reader->open(std::move(random_access_read_file)); + if (*error_code != common::E_OK) { + return nullptr; + } + return reader.release(); + } catch (const std::bad_alloc&) { + *error_code = common::E_OOM; + } catch (...) { + *error_code = common::E_FILE_OPEN_ERR; + } + return nullptr; +} diff --git a/python/tsfile/python_random_access_read_file.h b/python/tsfile/python_random_access_read_file.h new file mode 100644 index 000000000..78b68cbfd --- /dev/null +++ b/python/tsfile/python_random_access_read_file.h @@ -0,0 +1,29 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ + +#ifndef PYTHON_TSFILE_PYTHON_RANDOM_ACCESS_READ_FILE_H +#define PYTHON_TSFILE_PYTHON_RANDOM_ACCESS_READ_FILE_H + +#include +#include + +void* create_tsfile_reader_from_python_file(PyObject* source, + int32_t* error_code); + +#endif // PYTHON_TSFILE_PYTHON_RANDOM_ACCESS_READ_FILE_H diff --git a/python/tsfile/tsfile_cpp.pxd b/python/tsfile/tsfile_cpp.pxd index 5fbb0311d..131300361 100644 --- a/python/tsfile/tsfile_cpp.pxd +++ b/python/tsfile/tsfile_cpp.pxd @@ -360,6 +360,8 @@ cdef extern from "cwrapper/tsfile_cwrapper.h": TableSchema * tsfile_reader_get_all_table_schemas(TsFileReader reader, uint32_t * size); + TableSchema * tsfile_reader_get_all_table_schemas_with_error( + TsFileReader reader, uint32_t * size, ErrorCode * error_code); DeviceSchema * tsfile_reader_get_all_timeseries_schemas(TsFileReader reader, uint32_t * size); @@ -435,7 +437,7 @@ cdef extern from "cwrapper/tsfile_cwrapper.h": ErrorCode* err_code) # resultSet : get data from resultSet - bint tsfile_result_set_next(ResultSet result_set, ErrorCode * err_code); + bint tsfile_result_set_next(ResultSet result_set, ErrorCode * err_code) nogil; bint tsfile_result_set_is_null_by_index(ResultSet result_set, uint32_t column_index); bint tsfile_result_set_is_null_by_name(ResultSet result_set, const char * column_name); void free_tsfile_result_set(ResultSet * result_set); diff --git a/python/tsfile/tsfile_py_cpp.pyx b/python/tsfile/tsfile_py_cpp.pyx index e10c402c6..c8f155ae8 100644 --- a/python/tsfile/tsfile_py_cpp.pyx +++ b/python/tsfile/tsfile_py_cpp.pyx @@ -20,6 +20,7 @@ from datetime import date as date_type from .date_utils import parse_date_to_int from .tsfile_cpp cimport * +import os import pandas as pd import numpy as np @@ -784,7 +785,7 @@ cdef TsFileWriter tsfile_writer_new_c(object pathname, uint64_t memory_threshold cdef TsFileReader tsfile_reader_new_c(object pathname) except NULL: cdef ErrorCode errno = 0 cdef TsFileReader reader = NULL - cdef bytes encoded_path = PyUnicode_AsUTF8String(pathname) + cdef bytes encoded_path = os.fsencode(pathname) cdef const char * c_path = encoded_path reader = tsfile_reader_new(c_path, &errno) check_error(errno) @@ -1240,11 +1241,13 @@ cdef object get_table_schema(TsFileReader reader, object table_name): cdef object get_all_table_schema(TsFileReader reader): cdef uint32_t table_num = 0 + cdef ErrorCode error_code = 0 cdef TableSchema * schemas cdef int i table_schemas = {} - schemas = tsfile_reader_get_all_table_schemas(reader, &table_num) + schemas = tsfile_reader_get_all_table_schemas_with_error(reader, &table_num, &error_code) + check_error(error_code) for i in range(table_num): schema_py = from_c_table_schema(schemas[i]) table_schemas.update([(schema_py.get_table_name(), schema_py)]) diff --git a/python/tsfile/tsfile_reader.pyx b/python/tsfile/tsfile_reader.pyx index c4354799c..f6509a4b7 100644 --- a/python/tsfile/tsfile_reader.pyx +++ b/python/tsfile/tsfile_reader.pyx @@ -18,6 +18,7 @@ #cython: language_level=3 +import os import weakref from typing import List, Optional, Dict @@ -25,6 +26,7 @@ import pandas as pd from libc.string cimport strlen from libc.stdlib cimport free, malloc from cpython.bytes cimport PyBytes_FromStringAndSize +from cpython.ref cimport PyObject from libc.string cimport memset import pyarrow as pa from libc.stdint cimport INT64_MIN, INT64_MAX, uint32_t, uintptr_t @@ -37,6 +39,10 @@ from .date_utils import parse_int_to_date from .tsfile_cpp cimport * from .tsfile_py_cpp cimport * +cdef extern from "python_random_access_read_file.h": + TsFileReader create_tsfile_reader_from_python_file( + PyObject* source, int32_t* error_code) except? NULL + cdef class ResultSetPy: """ Get data from a query result. When reader run a query, a query handler will return. @@ -84,8 +90,10 @@ cdef class ResultSetPy: :return: boolean, true means get next rows. """ cdef ErrorCode code = 0 + cdef bint has_next self.check_result_set_invalid() - has_next = tsfile_result_set_next(self.result, &code) + with nogil: + has_next = tsfile_result_set_next(self.result, &code) check_error(code) return has_next @@ -344,7 +352,7 @@ cdef class PreparedSeriesPy: cdef class TsFileReaderPy: """ - Cython wrapper class for interacting with TsFileReader C implementation. + Cython wrapper class for interacting with the native TsFileReader. Provides a Pythonic interface to read and query time series data from TsFiles. """ @@ -355,16 +363,22 @@ cdef class TsFileReaderPy: cdef TsFileReader reader cdef object activate_result_set_list - def __init__(self, pathname): + def __init__(self, source): """ - Initialize a TsFile reader for the specified file path. + Initialize from a str/bytes/PathLike path or seekable binary file object. """ self.reader = NULL self.activate_result_set_list = weakref.WeakSet() - self.init_reader(pathname) + self.init_reader(source) - cdef init_reader(self, pathname): - self.reader = tsfile_reader_new_c(pathname) + cdef init_reader(self, source): + cdef ErrorCode error_code = 0 + if isinstance(source, (str, bytes, os.PathLike)): + self.reader = tsfile_reader_new_c(os.fspath(source)) + return + self.reader = create_tsfile_reader_from_python_file( + source, &error_code) + check_error(error_code) def query_table(self, table_name : str, column_names : List[str], start_time : int = INT64_MIN, end_time : int = INT64_MAX,