From 91c398dc4b2868fae00990ddd58a4080784890da Mon Sep 17 00:00:00 2001 From: ColinLee Date: Thu, 24 Sep 2026 17:26:47 +0800 Subject: [PATCH 1/2] fix(cpp): propagate table read failures without silent data loss Propagate metadata and device-index read errors through table queries and tag filters instead of returning empty or partial results or misleading errors. Preserve terminal result-set failures and reject unexpected short reads while allowing EOF-limited prefetches. Use checked schema lookups in the C/Python bindings and release tag filters when query creation fails. Add fault-injection regression coverage for row and batch queries, tag filtering, and leaf/internal device indexes. Validation: C++ 953 passed, 3 skipped; Python 356 passed, 1 skipped. Spotless, Black, and whitespace checks passed. --- cpp/src/cwrapper/tsfile_cwrapper.cc | 50 ++- cpp/src/reader/aligned_chunk_reader.cc | 31 ++ .../block/device_ordered_tsblock_reader.cc | 5 +- cpp/src/reader/chunk_reader.cc | 9 + cpp/src/reader/device_meta_iterator.cc | 64 ++-- cpp/src/reader/device_meta_iterator.h | 10 +- cpp/src/reader/meta_data_querier.cc | 4 +- cpp/src/reader/meta_data_querier.h | 1 - cpp/src/reader/table_query_executor.cc | 25 +- cpp/src/reader/table_result_set.cc | 22 +- cpp/src/reader/table_result_set.h | 4 + cpp/src/reader/task/device_task_iterator.cc | 4 +- cpp/src/reader/task/device_task_iterator.h | 2 +- cpp/src/reader/tsfile_reader.cc | 22 +- cpp/src/reader/tsfile_reader.h | 4 + cpp/test/reader/table_read_failure_test.cc | 321 ++++++++++++++++++ python/tests/test_reader_sources.py | 97 +++++- python/tsfile/tsfile_cpp.pxd | 2 + python/tsfile/tsfile_py_cpp.pyx | 5 +- python/tsfile/tsfile_reader.pyx | 8 +- 20 files changed, 608 insertions(+), 82 deletions(-) create mode 100644 cpp/test/reader/table_read_failure_test.cc diff --git a/cpp/src/cwrapper/tsfile_cwrapper.cc b/cpp/src/cwrapper/tsfile_cwrapper.cc index 5cae8f01d..9f0bbe82b 100644 --- a/cpp/src/cwrapper/tsfile_cwrapper.cc +++ b/cpp/src/cwrapper/tsfile_cwrapper.cc @@ -1060,25 +1060,9 @@ int tsfile_result_set_metadata_get_column_num(ResultSetMetaData result_set) { TableSchema tsfile_reader_get_table_schema(TsFileReader reader, const char* table_name) { - auto* r = static_cast(reader); - auto table_shcema = r->get_table_schema(table_name); - TableSchema ret_schema; - ret_schema.table_name = strdup(table_shcema->get_table_name().c_str()); - int column_num = table_shcema->get_columns_num(); - ret_schema.column_num = column_num; - ret_schema.column_schemas = - static_cast(malloc(sizeof(ColumnSchema) * column_num)); - for (int i = 0; i < column_num; i++) { - auto column_schema = table_shcema->get_measurement_schemas()[i]; - ret_schema.column_schemas[i].column_name = - strdup(column_schema->measurement_name_.c_str()); - ret_schema.column_schemas[i].data_type = - static_cast(column_schema->data_type_); - ret_schema.column_schemas[i].column_category = - static_cast( - table_shcema->get_column_categories()[i]); - } - return ret_schema; + TableSchema schema{}; + tsfile_reader_get_table_schema_checked(reader, table_name, &schema); + return schema; } static ERRNO copy_table_schema(const std::shared_ptr& src, @@ -1128,9 +1112,13 @@ ERRNO tsfile_reader_get_table_schema_checked(TsFileReader reader, } *out_schema = TableSchema{}; try { - auto schema = + std::shared_ptr schema; + const int ret = static_cast(reader)->get_table_schema( - table_name); + table_name, schema); + if (ret != common::E_OK) { + return ret; + } return copy_table_schema(schema, out_schema); } catch (const std::bad_alloc&) { return common::E_OOM; @@ -2553,9 +2541,13 @@ TagFilterHandle tsfile_tag_filter_create(TsFileReader reader, return nullptr; } auto* r = static_cast(reader); - auto schema = r->get_table_schema(table_name); - if (!schema) { - *err_code = common::E_INVALID_ARG; + std::shared_ptr schema; + const int ret = r->get_table_schema(table_name, schema); + if (ret != common::E_OK) { + // Preserve the existing missing-table contract, but not at the cost + // of disguising metadata I/O errors as invalid filter arguments. + *err_code = + ret == common::E_TABLE_NOT_EXIST ? common::E_INVALID_ARG : ret; return nullptr; } storage::TagFilterBuilder builder(schema.get()); @@ -2617,9 +2609,13 @@ TagFilterHandle tsfile_tag_filter_between(TsFileReader reader, return nullptr; } auto* r = static_cast(reader); - auto schema = r->get_table_schema(table_name); - if (!schema) { - *err_code = common::E_INVALID_ARG; + std::shared_ptr schema; + const int ret = r->get_table_schema(table_name, schema); + if (ret != common::E_OK) { + // Preserve the existing missing-table contract, but not at the cost + // of disguising metadata I/O errors as invalid filter arguments. + *err_code = + ret == common::E_TABLE_NOT_EXIST ? common::E_INVALID_ARG : ret; return nullptr; } storage::TagFilterBuilder builder(schema.get()); diff --git a/cpp/src/reader/aligned_chunk_reader.cc b/cpp/src/reader/aligned_chunk_reader.cc index 1c5b838d7..915712522 100644 --- a/cpp/src/reader/aligned_chunk_reader.cc +++ b/cpp/src/reader/aligned_chunk_reader.cc @@ -259,6 +259,13 @@ int AlignedChunkReader::load_by_aligned_meta(ChunkMeta* time_chunk_meta, ret = read_file_->read(time_chunk_meta_->offset_of_chunk_header_, time_file_data_buf, file_data_time_buf_size_, ret_read_len); + if (IS_SUCC(ret) && + ret_read_len < + UTIL_MIN(static_cast(file_data_time_buf_size_), + read_file_->file_size() - + time_chunk_meta_->offset_of_chunk_header_)) { + ret = E_FILE_READ_ERR; + } if (!IS_SUCC(ret)) { mem_free(time_file_data_buf); return ret; @@ -286,6 +293,13 @@ int AlignedChunkReader::load_by_aligned_meta(ChunkMeta* time_chunk_meta, ret = read_file_->read(value_chunk_meta_->offset_of_chunk_header_, value_file_data_buf, file_data_value_buf_size_, ret_read_len); + if (IS_SUCC(ret) && + ret_read_len < + UTIL_MIN(static_cast(file_data_value_buf_size_), + read_file_->file_size() - + value_chunk_meta_->offset_of_chunk_header_)) { + ret = E_FILE_READ_ERR; + } if (!IS_SUCC(ret)) { mem_free(value_file_data_buf); return ret; @@ -478,6 +492,9 @@ int AlignedChunkReader::read_from_file_and_rewrap( int ret_read_len = 0; if (RET_FAIL( read_file_->read(offset, file_data_buf, read_size, ret_read_len))) { + } else if (ret_read_len < UTIL_MIN(static_cast(read_size), + read_file_->file_size() - offset)) { + ret = E_FILE_READ_ERR; } else { in_stream_.wrap_from(file_data_buf, ret_read_len); #ifdef DEBUG_SE @@ -1216,6 +1233,13 @@ int AlignedChunkReader::load_by_aligned_meta_multi( ret = read_file_->read(time_chunk_meta_->offset_of_chunk_header_, time_file_data_buf, file_data_time_buf_size_, ret_read_len); + if (IS_SUCC(ret) && + ret_read_len < + UTIL_MIN(static_cast(file_data_time_buf_size_), + read_file_->file_size() - + time_chunk_meta_->offset_of_chunk_header_)) { + ret = E_FILE_READ_ERR; + } if (!IS_SUCC(ret)) { mem_free(time_file_data_buf); return ret; @@ -1264,6 +1288,13 @@ int AlignedChunkReader::load_by_aligned_meta_multi( ret = read_file_->read(col->chunk_meta->offset_of_chunk_header_, vbuf, col->file_data_buf_size, ret_read_len); + if (IS_SUCC(ret) && + ret_read_len < + UTIL_MIN(static_cast(col->file_data_buf_size), + read_file_->file_size() - + col->chunk_meta->offset_of_chunk_header_)) { + ret = E_FILE_READ_ERR; + } if (!IS_SUCC(ret)) { mem_free(vbuf); return ret; diff --git a/cpp/src/reader/block/device_ordered_tsblock_reader.cc b/cpp/src/reader/block/device_ordered_tsblock_reader.cc index 38cb63f7e..5c30d9bc6 100644 --- a/cpp/src/reader/block/device_ordered_tsblock_reader.cc +++ b/cpp/src/reader/block/device_ordered_tsblock_reader.cc @@ -51,7 +51,10 @@ int DeviceOrderedTsBlockReader::has_next(bool& has_next) { has_next = false; return common::E_OK; } - if (!device_task_iterator_->has_next()) { + if (RET_FAIL(device_task_iterator_->has_next(has_next))) { + return ret; + } + if (!has_next) { break; } DeviceQueryTask* task = nullptr; diff --git a/cpp/src/reader/chunk_reader.cc b/cpp/src/reader/chunk_reader.cc index 271b5c206..abeff8375 100644 --- a/cpp/src/reader/chunk_reader.cc +++ b/cpp/src/reader/chunk_reader.cc @@ -114,6 +114,12 @@ int ChunkReader::load_by_meta(ChunkMeta* meta) { } ret = read_file_->read(chunk_meta_->offset_of_chunk_header_, file_data_buf, file_data_buf_size_, ret_read_len); + if (IS_SUCC(ret) && + ret_read_len < UTIL_MIN(static_cast(file_data_buf_size_), + read_file_->file_size() - + chunk_meta_->offset_of_chunk_header_)) { + ret = E_FILE_READ_ERR; + } if (!IS_SUCC(ret)) { mem_free(file_data_buf); return ret; @@ -275,6 +281,9 @@ int ChunkReader::read_from_file_and_rewrap(int want_size) { int ret_read_len = 0; if (RET_FAIL( read_file_->read(offset, file_data_buf, read_size, ret_read_len))) { + } else if (ret_read_len < UTIL_MIN(static_cast(read_size), + read_file_->file_size() - offset)) { + ret = E_FILE_READ_ERR; } else { in_stream_.wrap_from(file_data_buf, ret_read_len); // DEBUG_hex_dump_buf("wrapped buf = ", file_data_buf, 256); diff --git a/cpp/src/reader/device_meta_iterator.cc b/cpp/src/reader/device_meta_iterator.cc index 6edb413eb..4bfc437e4 100644 --- a/cpp/src/reader/device_meta_iterator.cc +++ b/cpp/src/reader/device_meta_iterator.cc @@ -35,33 +35,48 @@ void DeviceMetaIterator::destroy_remaining_cached_devices() { DeviceMetaIterator::~DeviceMetaIterator() { destroy_remaining_cached_devices(); + while (!meta_index_nodes_.empty()) { + auto pending = meta_index_nodes_.front(); + meta_index_nodes_.pop(); + if (pending.second) { + pending.first->~MetaIndexNode(); + } + } pa_.destroy(); } -bool DeviceMetaIterator::has_next() { +int DeviceMetaIterator::has_next(bool& has_next) { + has_next = false; + if (read_error_ != common::E_OK) { + return read_error_; + } if (!result_cache_.empty()) { - return true; + has_next = true; + return common::E_OK; } if (direct_device_id_ != nullptr) { if (direct_lookup_done_) { - return false; - } - if (load_results_direct() != common::E_OK) { - return false; + return common::E_OK; } - return !result_cache_.empty(); + read_error_ = load_results_direct(); + } else { + read_error_ = load_results(); } - - if (load_results() != common::E_OK) { - return false; + if (read_error_ == common::E_OK) { + has_next = !result_cache_.empty(); } - return !result_cache_.empty(); + return read_error_; } int DeviceMetaIterator::next( std::pair, MetaIndexNode*>& ret_meta) { - if (!has_next()) { + bool available = false; + const int ret = has_next(available); + if (ret != common::E_OK) { + return ret; + } + if (!available) { return common::E_NO_MORE_DATA; } @@ -71,20 +86,23 @@ int DeviceMetaIterator::next( } int DeviceMetaIterator::load_results() { - int root_num = meta_index_nodes_.size(); while (!meta_index_nodes_.empty()) { - auto meta_data_index_node = meta_index_nodes_.front(); + auto pending = meta_index_nodes_.front(); meta_index_nodes_.pop(); - const auto& node_type = meta_data_index_node->node_type_; - if (node_type == MetaIndexNodeType::LEAF_DEVICE) { - load_leaf_device(meta_data_index_node); - } else if (node_type == MetaIndexNodeType::INTERNAL_DEVICE) { - load_internal_node(meta_data_index_node); + auto* node = pending.first; + int ret = common::E_OK; + if (node->node_type_ == MetaIndexNodeType::LEAF_DEVICE) { + ret = load_leaf_device(node); + } else if (node->node_type_ == MetaIndexNodeType::INTERNAL_DEVICE) { + ret = load_internal_node(node); } else { - return common::E_INVALID_NODE_TYPE; + ret = common::E_INVALID_NODE_TYPE; + } + if (pending.second) { + node->~MetaIndexNode(); } - if (root_num-- <= 0) { - meta_data_index_node->~MetaIndexNode(); + if (ret != common::E_OK) { + return ret; } } return common::E_OK; @@ -136,7 +154,7 @@ int DeviceMetaIterator::load_internal_node(MetaIndexNode* meta_index_node) { start_offset, end_offset, pa_, child_node, false))) { return ret; } else { - meta_index_nodes_.push(child_node); + meta_index_nodes_.push({child_node, true}); } } return ret; diff --git a/cpp/src/reader/device_meta_iterator.h b/cpp/src/reader/device_meta_iterator.h index 9f42819d3..c9c6e38fe 100644 --- a/cpp/src/reader/device_meta_iterator.h +++ b/cpp/src/reader/device_meta_iterator.h @@ -41,7 +41,7 @@ class DeviceMetaIterator { // A valid schema-only table has no device index. Treat a null root as // an empty iterator instead of dereferencing it during has_next(). if (meat_index_node != nullptr) { - meta_index_nodes_.push(meat_index_node); + meta_index_nodes_.push({meat_index_node, false}); } pa_.init(512, common::MOD_DEVICE_META_ITER); try_setup_direct_lookup(meat_index_node); @@ -54,7 +54,7 @@ class DeviceMetaIterator { id_filter_(id_filter), direct_lookup_done_(false) { for (auto meta_index_node : meta_index_node_list) { - meta_index_nodes_.push(meta_index_node); + meta_index_nodes_.push({meta_index_node, false}); } should_split_device_name = true; pa_.init(512, common::MOD_DEVICE_META_ITER); @@ -64,7 +64,7 @@ class DeviceMetaIterator { void destroy_remaining_cached_devices(); - bool has_next(); + int has_next(bool& has_next); int next(std::pair, MetaIndexNode*>& ret_meta); @@ -77,13 +77,15 @@ class DeviceMetaIterator { int load_results_direct(); TsFileIOReader* io_reader_; - std::queue meta_index_nodes_; + // Roots are borrowed from file metadata; descendant nodes belong to pa_. + std::queue> meta_index_nodes_; std::queue, MetaIndexNode*>> result_cache_; const Filter* id_filter_; common::PageArena pa_; bool should_split_device_name; + int read_error_ = common::E_OK; bool direct_lookup_done_; std::shared_ptr direct_device_id_; MetaIndexNode* direct_root_node_ = nullptr; diff --git a/cpp/src/reader/meta_data_querier.cc b/cpp/src/reader/meta_data_querier.cc index 0accbdde9..caf609c23 100644 --- a/cpp/src/reader/meta_data_querier.cc +++ b/cpp/src/reader/meta_data_querier.cc @@ -25,7 +25,6 @@ namespace storage { MetadataQuerier::MetadataQuerier(TsFileIOReader* tsfile_io_reader) : io_reader_(tsfile_io_reader) { - file_metadata_ = io_reader_->get_tsfile_meta(); device_chunk_meta_cache_ = std::unique_ptr< common::Cache>, std::mutex>>( @@ -69,8 +68,7 @@ MetadataQuerier::get_chunk_metadata_map(const std::vector& paths) const { } int MetadataQuerier::get_whole_file_metadata(TsFileMeta* tsfile_meta) const { - tsfile_meta = io_reader_->get_tsfile_meta(); - return common::E_OK; + return io_reader_->get_tsfile_meta(tsfile_meta); } void MetadataQuerier::load_chunk_metadatas(const std::vector& paths) { diff --git a/cpp/src/reader/meta_data_querier.h b/cpp/src/reader/meta_data_querier.h index be575323f..4867e58f5 100644 --- a/cpp/src/reader/meta_data_querier.h +++ b/cpp/src/reader/meta_data_querier.h @@ -70,7 +70,6 @@ class MetadataQuerier : public IMetadataQuerier { private: TsFileIOReader* io_reader_; - TsFileMeta* file_metadata_; std::unique_ptr< common::Cache*/ std::vector>, std::mutex>> diff --git a/cpp/src/reader/table_query_executor.cc b/cpp/src/reader/table_query_executor.cc index ec3807b31..5d708baae 100644 --- a/cpp/src/reader/table_query_executor.cc +++ b/cpp/src/reader/table_query_executor.cc @@ -28,7 +28,11 @@ int TableQueryExecutor::query(const std::string& table_name, Filter* field_filter, ResultSet*& ret_qds) { int ret = common::E_OK; TsFileMeta* file_metadata = nullptr; - file_metadata = tsfile_io_reader_->get_tsfile_meta(); + ret_qds = nullptr; + if (RET_FAIL(tsfile_io_reader_->get_tsfile_meta(file_metadata))) { + delete time_filter; + return ret; + } common::PageArena pa; pa.init(512, common::MOD_TSFILE_READER); MetaIndexNode* table_root = nullptr; @@ -98,7 +102,11 @@ int TableQueryExecutor::query(const std::string& table_name, ResultSet*& ret_qds) { int ret = common::E_OK; TsFileMeta* file_metadata = nullptr; - file_metadata = tsfile_io_reader_->get_tsfile_meta(); + ret_qds = nullptr; + if (RET_FAIL(tsfile_io_reader_->get_tsfile_meta(file_metadata))) { + delete time_filter; + return ret; + } common::PageArena pa; pa.init(512, common::MOD_TSFILE_READER); MetaIndexNode* table_root = nullptr; @@ -165,13 +173,20 @@ int TableQueryExecutor::query_on_tree( common::PageArena pa; pa.init(512, common::MOD_TSFILE_READER); int ret = common::E_OK; - TsFileMeta* file_meta = tsfile_io_reader_->get_tsfile_meta(); + ret_qds = nullptr; + TsFileMeta* file_meta = nullptr; + if (RET_FAIL(tsfile_io_reader_->get_tsfile_meta(file_meta))) { + delete time_filter; + return ret; + } std::unordered_set table_inodes; for (auto const& device : devices) { - MetaIndexNode* table_inode; + MetaIndexNode* table_inode = nullptr; if (RET_FAIL(file_meta->get_table_metaindex_node( device->get_table_name(), table_inode))) { - }; + delete time_filter; + return ret; + } table_inodes.insert(table_inode); } diff --git a/cpp/src/reader/table_result_set.cc b/cpp/src/reader/table_result_set.cc index 1a8d2a687..c77d39633 100644 --- a/cpp/src/reader/table_result_set.cc +++ b/cpp/src/reader/table_result_set.cc @@ -39,6 +39,21 @@ void TableResultSet::init() { TableResultSet::~TableResultSet() { close(); } int TableResultSet::next(bool& has_next) { + has_next = false; + if (read_error_ != common::E_OK) { + return read_error_; + } + const int ret = next_internal(has_next); + if (ret != common::E_OK) { + read_error_ = ret; + has_next = false; + row_ready_ = false; + row_materialized_ = false; + } + return ret; +} + +int TableResultSet::next_internal(bool& has_next) { if (return_mode_ != RETURN_ROW) { return tsblock_reader_->has_next(has_next); } @@ -177,6 +192,9 @@ std::shared_ptr TableResultSet::get_metadata() { int TableResultSet::get_next_tsblock(common::TsBlock*& block) { int ret = common::E_OK; block = nullptr; + if (read_error_ != common::E_OK) { + return read_error_; + } if (return_mode_ == RETURN_ROW) { return common::E_INVALID_ARG; @@ -184,7 +202,7 @@ int TableResultSet::get_next_tsblock(common::TsBlock*& block) { bool has_next = false; if (RET_FAIL(tsblock_reader_->has_next(has_next))) { - return ret; + return read_error_ = ret; } if (!has_next) { @@ -192,7 +210,7 @@ int TableResultSet::get_next_tsblock(common::TsBlock*& block) { } if (RET_FAIL(tsblock_reader_->next(tsblock_))) { - return ret; + return read_error_ = ret; } if (tsblock_ == nullptr) { diff --git a/cpp/src/reader/table_result_set.h b/cpp/src/reader/table_result_set.h index d92072934..b42b69c3f 100644 --- a/cpp/src/reader/table_result_set.h +++ b/cpp/src/reader/table_result_set.h @@ -61,6 +61,7 @@ class TableResultSet : public ResultSet { private: void init(); + int next_internal(bool& has_next); // Lazy materialization: fill row_record_ from the current row when a // caller actually requests the RowRecord (or a non-fast accessor). void materialize_current_row(); @@ -74,6 +75,9 @@ class TableResultSet : public ResultSet { std::vector data_types_; const int return_mode_; bool closed_ = false; + // A failed read may have advanced device/page state. Never resume it as + // EOF. + int read_error_ = common::E_OK; // True when row_iterator_ points at a row that hasn't been consumed yet. bool row_ready_ = false; // True when row_record_ has been populated for the current row. diff --git a/cpp/src/reader/task/device_task_iterator.cc b/cpp/src/reader/task/device_task_iterator.cc index e22fefb06..5b20f1715 100644 --- a/cpp/src/reader/task/device_task_iterator.cc +++ b/cpp/src/reader/task/device_task_iterator.cc @@ -25,8 +25,8 @@ void DeviceTaskIterator::flush_remaining_device_meta_cache() { device_meta_iterator_->destroy_remaining_cached_devices(); } -bool DeviceTaskIterator::has_next() const { - return device_meta_iterator_->has_next(); +int DeviceTaskIterator::has_next(bool& has_next) const { + return device_meta_iterator_->has_next(has_next); } int DeviceTaskIterator::next(DeviceQueryTask*& task) { diff --git a/cpp/src/reader/task/device_task_iterator.h b/cpp/src/reader/task/device_task_iterator.h index cc5a75562..ad86dbc1d 100644 --- a/cpp/src/reader/task/device_task_iterator.h +++ b/cpp/src/reader/task/device_task_iterator.h @@ -72,7 +72,7 @@ class DeviceTaskIterator { void flush_remaining_device_meta_cache(); - bool has_next() const; + int has_next(bool& has_next) const; int next(DeviceQueryTask*& task); diff --git a/cpp/src/reader/tsfile_reader.cc b/cpp/src/reader/tsfile_reader.cc index f860a06e5..6997cb389 100644 --- a/cpp/src/reader/tsfile_reader.cc +++ b/cpp/src/reader/tsfile_reader.cc @@ -730,16 +730,26 @@ ResultSet* TsFileReader::read_timeseries( std::shared_ptr TsFileReader::get_table_schema( const std::string& table_name) { - TsFileMeta* file_metadata = tsfile_executor_->get_tsfile_meta(); std::shared_ptr table_schema; - // A schema-only table has no device-level metadata index. Schema lookup - // must therefore be independent of the presence of data pages; callers - // can still construct an empty result set from the returned schema. - if (file_metadata == nullptr) return table_schema; - file_metadata->get_table_schema(to_lower(table_name), table_schema); + get_table_schema(table_name, table_schema); return table_schema; } +int TsFileReader::get_table_schema(const std::string& table_name, + std::shared_ptr& table_schema) { + table_schema.reset(); + 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; + } + // Schema-only tables have no device index, but still have a valid schema. + return file_metadata->get_table_schema(to_lower(table_name), table_schema); +} + std::vector> TsFileReader::get_all_table_schemas() { std::vector> table_schemas; diff --git a/cpp/src/reader/tsfile_reader.h b/cpp/src/reader/tsfile_reader.h index 6493a4af7..11bb72913 100644 --- a/cpp/src/reader/tsfile_reader.h +++ b/cpp/src/reader/tsfile_reader.h @@ -263,6 +263,10 @@ class TsFileReader { */ std::shared_ptr get_table_schema( const std::string& table_name); + + /** Error-reporting overload. The output is null on failure. */ + int get_table_schema(const std::string& table_name, + std::shared_ptr& table_schema); /** * @brief get all table schemas in the tsfile * diff --git a/cpp/test/reader/table_read_failure_test.cc b/cpp/test/reader/table_read_failure_test.cc new file mode 100644 index 000000000..4d2ed13ce --- /dev/null +++ b/cpp/test/reader/table_read_failure_test.cc @@ -0,0 +1,321 @@ +/* + * 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 a + * + * 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 + +#include +#include +#include +#include +#include +#include +#include +#include +#include +#include + +#include "common/config/config.h" +#include "common/tablet.h" +#include "cwrapper/tsfile_cwrapper.h" +#include "file/write_file.h" +#include "reader/filter/tag_filter.h" +#include "reader/table_result_set.h" +#include "reader/tsfile_reader.h" +#include "writer/tsfile_table_writer.h" + +namespace { + +class FailingReadFile : public storage::RandomAccessReadFile { + public: + explicit FailingReadFile(const std::vector& bytes) : bytes_(bytes) {} + bool is_opened() const override { return opened_; } + int64_t file_size() const override { return bytes_.size(); } + const std::string& file_path() const override { return path_; } + int generation(uint64_t& size, uint64_t& fingerprint) const override { + size = bytes_.size(); + fingerprint = 0; + return common::E_OK; + } + int read(int64_t offset, char* buffer, int32_t size, + int32_t& read_size) override { + ++reads; + read_size = 0; + if (fail_at > 0 && (persistent ? reads >= fail_at : reads == fail_at)) { + failed = true; + return short_read ? common::E_OK : common::E_FILE_READ_ERR; + } + if (offset < 0 || size < 0) return common::E_INVALID_ARG; + if (offset >= static_cast(bytes_.size())) return common::E_OK; + read_size = static_cast( + std::min(size, bytes_.size() - offset)); + std::memcpy(buffer, bytes_.data() + offset, read_size); + return common::E_OK; + } + void close() override { opened_ = false; } + + int reads = 0; + int fail_at = 0; + bool persistent = false; + bool short_read = false; + bool failed = false; + + private: + const std::vector& bytes_; + bool opened_ = true; + std::string path_ = "memory://table-read-failure"; +}; + +// One tag exercises TagEq's direct lookup; two tags exercise filtered +// traversal. Two devices fit in a leaf; five force multiple levels of internal +// index nodes. +class TableReadFailureTest + : public ::testing::TestWithParam> { + protected: + void SetUp() override { + storage::libtsfile_init(); + saved_index_degree_ = common::g_config_value_.max_degree_of_index_node_; + ASSERT_EQ(storage::set_max_degree_of_index_node(2), common::E_OK); + filename_ = + "table_read_failure_" + + std::to_string( + std::chrono::steady_clock::now().time_since_epoch().count()) + + ".tsfile"; + storage::WriteFile file; + int flags = O_WRONLY | O_CREAT | O_TRUNC; +#ifdef _WIN32 + flags |= O_BINARY; +#endif + ASSERT_EQ(file.create(filename_, flags, 0666), common::E_OK); + std::vector columns; + for (int i = 0; i < std::get<0>(GetParam()); ++i) { + columns.emplace_back("id" + std::to_string(i), common::STRING, + common::UNCOMPRESSED, common::PLAIN, + common::ColumnCategory::TAG); + } + columns.emplace_back("value", common::INT64, common::UNCOMPRESSED, + common::PLAIN, common::ColumnCategory::FIELD); + storage::TableSchema schema("test", columns); + storage::TsFileTableWriter writer(&file, &schema); + storage::Tablet tablet( + "test", schema.get_measurement_names(), schema.get_data_types(), + schema.get_column_categories(), std::get<1>(GetParam()) * 30); + for (int device = 0; device < std::get<1>(GetParam()); ++device) { + for (int t = 0; t < 30; ++t) { + const int row = device * 30 + t; + const std::string id = "d" + std::to_string(device); + ASSERT_EQ(tablet.add_timestamp(row, t), common::E_OK); + ASSERT_EQ(tablet.add_value(row, "id0", id.c_str()), + common::E_OK); + if (std::get<0>(GetParam()) == 2) { + ASSERT_EQ(tablet.add_value(row, "id1", "tag"), + common::E_OK); + } + ASSERT_EQ( + tablet.add_value(row, "value", static_cast(row)), + common::E_OK); + } + } + ASSERT_EQ(writer.write_table(tablet), common::E_OK); + ASSERT_EQ(writer.flush(), common::E_OK); + ASSERT_EQ(writer.close(), common::E_OK); + std::ifstream input(filename_, std::ios::binary); + ASSERT_TRUE(input.is_open()); + bytes_.assign(std::istreambuf_iterator(input), + std::istreambuf_iterator()); + } + + void TearDown() override { + storage::set_max_degree_of_index_node(saved_index_degree_); + std::remove(filename_.c_str()); + storage::libtsfile_destroy(); + } + + struct ScanResult { + int ret = common::E_OK; + int rows = 0; + int open_reads = 0; + int reads = 0; + bool failed = false; + }; + + void scan(bool filter, int batch_size, int fail_at, bool persistent, + bool short_read, ScanResult& out) { + storage::TsFileReader reader; + auto* source = new FailingReadFile(bytes_); + source->fail_at = fail_at; + source->persistent = persistent; + source->short_read = short_read; + ASSERT_EQ( + reader.open(std::unique_ptr(source)), + common::E_OK); + out.open_reads = source->reads; + storage::TagEq eq(1, "d0"); + storage::ResultSet* result = nullptr; + out.ret = reader.query("test", {"id0", "value"}, 0, INT64_MAX, result, + filter ? &eq : nullptr, batch_size); + if (out.ret == common::E_OK) { + ASSERT_NE(result, nullptr); + bool next = false; + while ((out.ret = result->next(next)) == common::E_OK && next) { + if (batch_size == 0) { + ++out.rows; + } else { + common::TsBlock* block = nullptr; + out.ret = result->get_next_tsblock(block); + if (out.ret != common::E_OK) break; + ASSERT_NE(block, nullptr); + out.rows += block->get_row_count(); + } + } + } + out.reads = source->reads; + out.failed = source->failed; + if (out.ret != common::E_OK && result != nullptr) { + source->fail_at = 0; + bool next = true; + EXPECT_EQ(result->next(next), out.ret); + EXPECT_FALSE(next); + if (batch_size != 0) { + common::TsBlock* block = nullptr; + EXPECT_EQ(result->get_next_tsblock(block), out.ret); + EXPECT_EQ(block, nullptr); + } + EXPECT_EQ(source->reads, out.reads); + } + if (result != nullptr) reader.destroy_query_data_set(result); + } + + std::vector bytes_; + std::string filename_; + uint32_t saved_index_degree_ = 0; +}; + +TEST_P(TableReadFailureTest, EveryReadFailureReachesCaller) { + for (bool filter : {false, true}) { + for (int batch_size : {0, 16}) { + ScanResult baseline; + ASSERT_NO_FATAL_FAILURE( + scan(filter, batch_size, 0, false, false, baseline)); + ASSERT_EQ(baseline.ret, common::E_OK); + ASSERT_EQ(baseline.rows, + filter ? 30 : std::get<1>(GetParam()) * 30); + for (bool persistent : {false, true}) { + for (bool short_read : {false, true}) { + for (int fail_at = baseline.open_reads + 1; + fail_at <= baseline.reads; ++fail_at) { + SCOPED_TRACE(::testing::Message() + << filter << ":" << batch_size << ":" + << persistent << ":" << short_read << ":" + << fail_at); + ScanResult failed; + ASSERT_NO_FATAL_FAILURE(scan(filter, batch_size, + fail_at, persistent, + short_read, failed)); + EXPECT_TRUE(failed.failed); + EXPECT_EQ(failed.ret, common::E_FILE_READ_ERR); + } + } + } + } + } +} + +TEST_P(TableReadFailureTest, CheckedSchemaAndTagFactoriesPreserveReadErrors) { + for (int operation = 0; operation < 3; ++operation) { + storage::TsFileReader reader; + auto* source = new FailingReadFile(bytes_); + ASSERT_EQ( + reader.open(std::unique_ptr(source)), + common::E_OK); + source->fail_at = source->reads + 1; + source->persistent = true; + ::TableSchema schema{}; + if (operation == 0) { + EXPECT_EQ(tsfile_reader_get_table_schema_checked(&reader, "test", + &schema), + common::E_FILE_READ_ERR); + EXPECT_EQ(schema.table_name, nullptr); + } else { + ERRNO error = common::E_OK; + TagFilterHandle filter = + operation == 1 + ? tsfile_tag_filter_create(&reader, "test", "id0", "d0", + TAG_FILTER_EQ, &error) + : tsfile_tag_filter_between(&reader, "test", "id0", "d0", + "d0", false, &error); + EXPECT_EQ(error, common::E_FILE_READ_ERR); + EXPECT_EQ(filter, nullptr); + if (filter != nullptr) tsfile_tag_filter_free(filter); + } + EXPECT_TRUE(source->failed); + // Metadata failures must not poison the reader or cache an empty + // schema. + source->fail_at = 0; + ASSERT_EQ( + tsfile_reader_get_table_schema_checked(&reader, "test", &schema), + common::E_OK); + EXPECT_STREQ(schema.table_name, "test"); + free_table_schema(schema); + } +} + +INSTANTIATE_TEST_SUITE_P(LeafAndInternalDeviceIndexes, TableReadFailureTest, + ::testing::Combine(::testing::Values(1, 2), + ::testing::Values(2, 5))); + +class FailingBlockReader : public storage::TsBlockReader { + public: + int has_next(bool& next) override { + next = true; + return common::E_OK; + } + int next(common::TsBlock*& block) override { + ++reads; + block = nullptr; + return common::E_FILE_READ_ERR; + } + void close() override {} + int reads = 0; +}; + +TEST(TableResultReadFailureTest, FailureWhileFetchingBlockIsTerminal) { + storage::libtsfile_init(); + for (int mode : {storage::RETURN_ROW, storage::RETURN_BATCH}) { + auto* source = new FailingBlockReader(); + storage::TableResultSet result( + std::unique_ptr(source), {"value"}, + {common::INT64}, mode); + bool next = true; + common::TsBlock* block = nullptr; + if (mode == storage::RETURN_ROW) { + EXPECT_EQ(result.next(next), common::E_FILE_READ_ERR); + EXPECT_FALSE(next); + } else { + EXPECT_EQ(result.get_next_tsblock(block), common::E_FILE_READ_ERR); + EXPECT_EQ(block, nullptr); + } + EXPECT_EQ(result.next(next), common::E_FILE_READ_ERR); + EXPECT_FALSE(next); + EXPECT_EQ(result.get_next_tsblock(block), common::E_FILE_READ_ERR); + EXPECT_EQ(block, nullptr); + EXPECT_EQ(source->reads, 1); + } + storage::libtsfile_destroy(); +} + +} // namespace diff --git a/python/tests/test_reader_sources.py b/python/tests/test_reader_sources.py index 24961f85e..654692fa7 100644 --- a/python/tests/test_reader_sources.py +++ b/python/tests/test_reader_sources.py @@ -39,6 +39,7 @@ TsFileWriter, ) from tsfile.exceptions import FileOpenError, FileReadError +from tsfile.tag_filter import BetweenTagFilter, ComparisonTagFilter RESOURCES = Path(__file__).parent / "resources" @@ -404,7 +405,9 @@ def query(reader): assert source.close_calls == 0 -@pytest.mark.parametrize("method", ["get_all_devices", "get_all_table_schemas"]) +@pytest.mark.parametrize( + "method", ["get_all_devices", "get_all_table_schemas", "get_table_schema"] +) def test_metadata_read_failure_is_not_an_empty_result(method): class FailingBytesIO(TrackingBytesIO): fail_reads = False @@ -420,8 +423,11 @@ def read(self, size=-1): source.seek(17) with TsFileReader(source) as reader: source.fail_reads = True + args = ("test",) if method == "get_table_schema" else () with pytest.raises(FileReadError): - getattr(reader, method)() + getattr(reader, method)(*args) + source.fail_reads = False + assert getattr(reader, method)(*args) is not None assert source.failed assert source.tell() == 17 assert source.close_calls == 0 @@ -461,3 +467,90 @@ def read(self, size=-1): assert source.failed assert source.tell() == 17 assert source.close_calls == 0 + + +@pytest.mark.parametrize("method", ["query_table", "query_table_by_row"]) +@pytest.mark.parametrize("tag_kind", [None, "eq", "between"]) +@pytest.mark.parametrize("batch_size", [0, 16]) +@pytest.mark.parametrize("short_read", [False, True]) +@pytest.mark.parametrize("persistent", [False, True]) +def test_table_queries_propagate_every_read_failure( + method, tag_kind, batch_size, short_read, persistent +): + 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 + if persistent + else self.read_calls == self.fail_at + ): + self.failed = True + if short_read: + return b"" + raise OSError("table read failed") + return super().read(size) + + data = (RESOURCES / "simple_table_t1.tsfile").read_bytes() + + def query(reader): + tag_filter = None + if tag_kind == "eq": + tag_filter = ComparisonTagFilter("s0", "a", ComparisonTagFilter.EQ) + elif tag_kind == "between": + tag_filter = BetweenTagFilter("s0", "a", "a") + return getattr(reader, method)( + "test", ["s0", "s2"], tag_filter=tag_filter, batch_size=batch_size + ) + + def consume(result): + rows = 0 + if batch_size: + while True: + batch = result.read_arrow_record_batch() + if batch is None: + break + rows += batch.num_rows + else: + while result.next(): + rows += 1 + return rows + + baseline = FailingBytesIO(data) + with TsFileReader(baseline) as reader: + open_reads = baseline.read_calls + with query(reader) as result: + assert consume(result) == (60 if tag_kind is None else 30) + + # Sweep actual reads rather than hard-coding call numbers: metadata may + # be prefetched or cached differently as the implementation evolves. + for fail_at in range(open_reads + 1, baseline.read_calls + 1): + source = FailingBytesIO(data) + source.seek(17) + with TsFileReader(source) as reader: + source.fail_at = fail_at + result = None + try: + with pytest.raises(FileReadError): + result = query(reader) + consume(result) + assert source.failed, f"read {fail_at} was not reached" + if result is not None: + # Recovery of the source must not turn a failed result + # into successful EOF or resume it after skipped devices. + source.fail_at = None + with pytest.raises(FileReadError): + result.next() + if batch_size: + with pytest.raises(FileReadError): + result.read_arrow_record_batch() + finally: + if result is not None: + result.close() + assert source.tell() == 17 + assert not source.closed + assert source.close_calls == 0 diff --git a/python/tsfile/tsfile_cpp.pxd b/python/tsfile/tsfile_cpp.pxd index b0b2d6b78..e97314d15 100644 --- a/python/tsfile/tsfile_cpp.pxd +++ b/python/tsfile/tsfile_cpp.pxd @@ -357,6 +357,8 @@ cdef extern from "cwrapper/tsfile_cwrapper.h": TableSchema tsfile_reader_get_table_schema(TsFileReader reader, const char * table_name); + ErrorCode tsfile_reader_get_table_schema_checked( + TsFileReader reader, const char * table_name, TableSchema * out_schema); TableSchema * tsfile_reader_get_all_table_schemas(TsFileReader reader, uint32_t * size); diff --git a/python/tsfile/tsfile_py_cpp.pyx b/python/tsfile/tsfile_py_cpp.pyx index c8f155ae8..05fa3cc98 100644 --- a/python/tsfile/tsfile_py_cpp.pyx +++ b/python/tsfile/tsfile_py_cpp.pyx @@ -1236,7 +1236,10 @@ cdef ResultSet tsfile_reader_query_table_with_tag_filter_c(TsFileReader reader, cdef object get_table_schema(TsFileReader reader, object table_name): cdef bytes table_name_bytes = PyUnicode_AsUTF8String(table_name) cdef const char * table_name_c = table_name_bytes - cdef TableSchema schema = tsfile_reader_get_table_schema(reader, table_name_c) + cdef TableSchema schema + cdef ErrorCode code = tsfile_reader_get_table_schema_checked( + reader, table_name_c, &schema) + check_error(code) return from_c_table_schema(schema) cdef object get_all_table_schema(TsFileReader reader): diff --git a/python/tsfile/tsfile_reader.pyx b/python/tsfile/tsfile_reader.pyx index f6509a4b7..bc094e8b7 100644 --- a/python/tsfile/tsfile_reader.pyx +++ b/python/tsfile/tsfile_reader.pyx @@ -391,8 +391,10 @@ cdef class TsFileReaderPy: """ cdef ResultSet result cdef TagFilterHandle c_tag_filter = NULL + pyresult = ResultSetPy(self) if tag_filter is not None: c_tag_filter = self._build_c_tag_filter(table_name.lower(), tag_filter) + pyresult._tag_filter_handle = c_tag_filter if batch_size <= 0: result = tsfile_reader_query_table_with_tag_filter_c( self.reader, table_name.lower(), @@ -403,8 +405,6 @@ cdef class TsFileReaderPy: self.reader, table_name.lower(), [column_name.lower() for column_name in column_names], start_time, end_time, c_tag_filter, batch_size) - pyresult = ResultSetPy(self) - pyresult._tag_filter_handle = c_tag_filter pyresult.init_c(result, table_name) self.activate_result_set_list.add(pyresult) return pyresult @@ -494,13 +494,13 @@ cdef class TsFileReaderPy: """ cdef ResultSet result cdef TagFilterHandle c_tag_filter = NULL + pyresult = ResultSetPy(self) if tag_filter is not None: c_tag_filter = self._build_c_tag_filter(table_name.lower(), tag_filter) + pyresult._tag_filter_handle = c_tag_filter result = tsfile_reader_query_table_by_row_c(self.reader, table_name.lower(), [column_name.lower() for column_name in column_names], offset, limit, c_tag_filter, batch_size) - pyresult = ResultSetPy(self) - pyresult._tag_filter_handle = c_tag_filter pyresult.init_c(result, table_name) self.activate_result_set_list.add(pyresult) return pyresult From 9bace1e2129801f94159ebcbd4ca23a3c1c48ed5 Mon Sep 17 00:00:00 2001 From: ColinLee Date: Thu, 24 Sep 2026 19:23:31 +0800 Subject: [PATCH 2/2] refactor(cpp): unify schema APIs around error reporting Remove the legacy no-error table schema and tag factory entry points. Make table and timeseries schema lookups return allocated objects with ERRNO output parameters, including read failures instead of silently returning empty or partial schemas. Update the C++ reader, CLI, examples, Python bindings, Go cgo bridge, and tests to the unified API. Add fault-injection coverage for all-timeseries schema reads. Validation: C++ 953 passed, 3 skipped; Python 357 passed, 1 skipped; Go tests passed. Formatting and whitespace checks passed. --- cpp/examples/cpp_examples/demo_read.cpp | 4 +- cpp/src/cwrapper/tsfile_cwrapper.cc | 195 ++++++++---------- cpp/src/cwrapper/tsfile_cwrapper.h | 90 ++------ cpp/src/reader/tsfile_reader.cc | 7 - cpp/src/reader/tsfile_reader.h | 10 +- cpp/test/cwrapper/c_release_test.cc | 17 +- cpp/test/cwrapper/cwrapper_test.cc | 24 +-- cpp/test/reader/table_read_failure_test.cc | 49 +++-- .../table_view/tsfile_reader_table_test.cc | 5 +- cpp/tools/commands/cmd_count.cc | 15 +- cpp/tools/commands/cmd_stats.cc | 23 ++- cpp/tools/commands/row_query.cc | 50 +++-- go/tsfile/cgo_bridge.go | 16 +- python/tests/test_reader_sources.py | 8 +- python/tsfile/tsfile_cpp.pxd | 15 +- python/tsfile/tsfile_py_cpp.pyx | 20 +- 16 files changed, 264 insertions(+), 284 deletions(-) diff --git a/cpp/examples/cpp_examples/demo_read.cpp b/cpp/examples/cpp_examples/demo_read.cpp index 91eb7c461..90e24c676 100644 --- a/cpp/examples/cpp_examples/demo_read.cpp +++ b/cpp/examples/cpp_examples/demo_read.cpp @@ -18,6 +18,7 @@ */ #include +#include #include #include @@ -40,7 +41,8 @@ int demo_read() { columns.emplace_back("id2"); columns.emplace_back("s1"); - auto table_schema = reader.get_table_schema(table_name); + std::shared_ptr table_schema; + HANDLE_ERROR(reader.get_table_schema(table_name, table_schema)); storage::Filter* tag_filter1 = storage::TagFilterBuilder(table_schema.get()).eq("id1", "id1_filed_1"); storage::Filter* tag_filter2 = diff --git a/cpp/src/cwrapper/tsfile_cwrapper.cc b/cpp/src/cwrapper/tsfile_cwrapper.cc index 9f0bbe82b..30c2b9318 100644 --- a/cpp/src/cwrapper/tsfile_cwrapper.cc +++ b/cpp/src/cwrapper/tsfile_cwrapper.cc @@ -1058,13 +1058,6 @@ int tsfile_result_set_metadata_get_column_num(ResultSetMetaData result_set) { return result_set.column_num; } -TableSchema tsfile_reader_get_table_schema(TsFileReader reader, - const char* table_name) { - TableSchema schema{}; - tsfile_reader_get_table_schema_checked(reader, table_name, &schema); - return schema; -} - static ERRNO copy_table_schema(const std::shared_ptr& src, TableSchema* out_schema) { if (!src || out_schema == nullptr) { @@ -1104,39 +1097,48 @@ static ERRNO copy_table_schema(const std::shared_ptr& src, return common::E_OK; } -ERRNO tsfile_reader_get_table_schema_checked(TsFileReader reader, - const char* table_name, - TableSchema* out_schema) { - if (reader == nullptr || table_name == nullptr || out_schema == nullptr) { - return common::E_INVALID_ARG; +TableSchema* tsfile_reader_get_table_schema(TsFileReader reader, + const char* table_name, + ERRNO* error_code) { + if (error_code == nullptr) { + return nullptr; + } + *error_code = common::E_INVALID_ARG; + if (reader == nullptr || table_name == nullptr) { + return nullptr; } - *out_schema = TableSchema{}; try { std::shared_ptr schema; const int ret = static_cast(reader)->get_table_schema( table_name, schema); if (ret != common::E_OK) { - return ret; + *error_code = ret; + return nullptr; + } + auto* result = static_cast(malloc(sizeof(TableSchema))); + if (result == nullptr) { + *error_code = common::E_OOM; + return nullptr; + } + *error_code = copy_table_schema(schema, result); + if (*error_code != common::E_OK) { + free(result); + return nullptr; } - return copy_table_schema(schema, out_schema); + return result; } catch (const std::bad_alloc&) { - return common::E_OOM; + *error_code = common::E_OOM; + return nullptr; } catch (...) { - return common::E_FILE_READ_ERR; + *error_code = common::E_FILE_READ_ERR; + return nullptr; } } 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) { + uint32_t* size, + ERRNO* error_code) { if (size != nullptr) { *size = 0; } @@ -1177,86 +1179,84 @@ TableSchema* tsfile_reader_get_all_table_schemas_with_error(TsFileReader reader, return ret; } -ERRNO tsfile_reader_get_all_table_schemas_checked(TsFileReader reader, - TableSchema** out_schemas, - uint32_t* out_size) { - if (reader == nullptr || out_schemas == nullptr || out_size == nullptr) { - return common::E_INVALID_ARG; +DeviceSchema* tsfile_reader_get_all_timeseries_schemas(TsFileReader reader, + uint32_t* size, + ERRNO* error_code) { + if (size != nullptr) { + *size = 0; } - *out_schemas = nullptr; - *out_size = 0; - try { - auto schemas = static_cast(reader) - ->get_all_table_schemas(); - if (schemas.empty()) { - return common::E_OK; - } - TableSchema* copied = static_cast( - calloc(schemas.size(), sizeof(TableSchema))); - if (copied == nullptr) { - return common::E_OOM; - } - for (size_t i = 0; i < schemas.size(); ++i) { - ERRNO ret = copy_table_schema(schemas[i], &copied[i]); - if (ret != common::E_OK) { - for (size_t j = 0; j < i; ++j) { - free_table_schema(copied[j]); - } - free(copied); - return ret; - } - } - *out_schemas = copied; - *out_size = static_cast(schemas.size()); - return common::E_OK; - } catch (const std::bad_alloc&) { - return common::E_OOM; - } catch (...) { - return common::E_FILE_READ_ERR; + if (error_code == nullptr) { + return nullptr; + } + *error_code = common::E_INVALID_ARG; + if (reader == nullptr || size == nullptr) { + return nullptr; } -} - -DeviceSchema* tsfile_reader_get_all_timeseries_schemas(TsFileReader reader, - uint32_t* size) { auto* r = static_cast(reader); - auto device_ids = r->get_all_device_ids(); - if (size == nullptr) { + std::vector> device_ids; + *error_code = r->get_all_devices(device_ids); + if (*error_code != common::E_OK) { return nullptr; } - *size = static_cast(device_ids.size()); if (device_ids.empty()) { + *error_code = common::E_OK; return nullptr; } - DeviceSchema* device_schema = static_cast( - malloc(sizeof(DeviceSchema) * device_ids.size())); - if (device_schema == nullptr) { - *size = 0; + auto* device_schemas = static_cast( + calloc(device_ids.size(), sizeof(DeviceSchema))); + if (device_schemas == nullptr) { + *error_code = common::E_OOM; return nullptr; } + auto free_partial = [&]() { + for (size_t i = 0; i < device_ids.size(); ++i) { + free_device_schema(device_schemas[i]); + } + free(device_schemas); + }; - size_t device_index = 0; - for (const auto& device_id : device_ids) { - DeviceSchema& cur_schema = device_schema[device_index++]; + for (size_t device_index = 0; device_index < device_ids.size(); + ++device_index) { + const auto& device_id = device_ids[device_index]; + DeviceSchema& cur_schema = device_schemas[device_index]; std::string device_name = device_id == nullptr ? "" : device_id->get_device_name(); cur_schema.device_name = strdup(device_name.c_str()); - cur_schema.timeseries_num = 0; - cur_schema.timeseries_schema = nullptr; + if (cur_schema.device_name == nullptr) { + free_partial(); + *error_code = common::E_OOM; + return nullptr; + } std::vector schemas; - int ret = r->get_timeseries_schema(device_id, schemas); - if (ret != common::E_OK || schemas.empty()) { + const int ret = r->get_timeseries_schema(device_id, schemas); + if (ret != common::E_OK) { + free_partial(); + *error_code = ret; + return nullptr; + } + if (schemas.empty()) { continue; } cur_schema.timeseries_num = static_cast(schemas.size()); cur_schema.timeseries_schema = static_cast( - malloc(sizeof(TimeseriesSchema) * schemas.size())); + calloc(schemas.size(), sizeof(TimeseriesSchema))); + if (cur_schema.timeseries_schema == nullptr) { + free_partial(); + *error_code = common::E_OOM; + return nullptr; + } for (size_t i = 0; i < schemas.size(); ++i) { const auto& measurement_schema = schemas[i]; cur_schema.timeseries_schema[i].timeseries_name = strdup(measurement_schema.measurement_name_.c_str()); + if (cur_schema.timeseries_schema[i].timeseries_name == nullptr) { + free_partial(); + *error_code = common::E_OOM; + return nullptr; + } cur_schema.timeseries_schema[i].data_type = static_cast(measurement_schema.data_type_); cur_schema.timeseries_schema[i].encoding = @@ -1266,7 +1266,9 @@ DeviceSchema* tsfile_reader_get_all_timeseries_schemas(TsFileReader reader, measurement_schema.compression_type_); } } - return device_schema; + *size = static_cast(device_ids.size()); + *error_code = common::E_OK; + return device_schemas; } void tsfile_device_id_free_contents(DeviceID* d) { @@ -2496,37 +2498,6 @@ ResultSet _tsfile_reader_query_device(TsFileReader reader, // ============== Tag Filter API Implementation ============== -// Helper macro to avoid repetition in tag filter factory functions. -// The shared_ptr must stay alive while TagFilterBuilder accesses the schema. -// Every C-API entry must validate its pointers: a null reader would deref -// during the static_cast, and null table/column/value would feed std::string -// a null pointer (UB / crash). -// The function-name suffix and the TagFilterBuilder method are always the same -// operator, so the macro takes a single argument used for both. -#define DEFINE_TAG_FILTER_FACTORY(op) \ - TagFilterHandle tsfile_tag_filter_##op( \ - TsFileReader reader, const char* table_name, const char* column_name, \ - const char* value) { \ - if (reader == nullptr || table_name == nullptr || \ - column_name == nullptr || value == nullptr) { \ - return nullptr; \ - } \ - auto* r = static_cast(reader); \ - auto schema = r->get_table_schema(table_name); \ - if (!schema) return nullptr; \ - storage::TagFilterBuilder builder(schema.get()); \ - return builder.op(column_name, value); \ - } - -DEFINE_TAG_FILTER_FACTORY(eq) -DEFINE_TAG_FILTER_FACTORY(neq) -DEFINE_TAG_FILTER_FACTORY(lt) -DEFINE_TAG_FILTER_FACTORY(lteq) -DEFINE_TAG_FILTER_FACTORY(gt) -DEFINE_TAG_FILTER_FACTORY(gteq) - -#undef DEFINE_TAG_FILTER_FACTORY - TagFilterHandle tsfile_tag_filter_create(TsFileReader reader, const char* table_name, const char* column_name, diff --git a/cpp/src/cwrapper/tsfile_cwrapper.h b/cpp/src/cwrapper/tsfile_cwrapper.h index 003e59198..8f747ece0 100644 --- a/cpp/src/cwrapper/tsfile_cwrapper.h +++ b/cpp/src/cwrapper/tsfile_cwrapper.h @@ -1017,33 +1017,14 @@ int tsfile_result_set_metadata_get_column_num(ResultSetMetaData result_set); // const char* device_id); /** - * @brief Gets specific table's schema in the tsfile. - * - * @return TableSchema, contains table and column info. - * @note Caller should call free_table_schema to free the tableschema. + * @brief Gets one table schema and reports lookup or read failures. + * @return Allocated TableSchema, or NULL on error. Check error_code to + * distinguish a missing table from a metadata read failure. + * @note Caller must call free_table_schema(*schema), then free(schema). */ -TableSchema tsfile_reader_get_table_schema(TsFileReader reader, - const char* table_name); - -/** Retrieves one table schema and reports missing tables through ERRNO. */ -ERRNO tsfile_reader_get_table_schema_checked(TsFileReader reader, - const char* table_name, - TableSchema* out_schema); -/** - * @brief Gets all table schema in the tsfile. - * - * @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); - -/** Retrieves every table schema into a caller-freed array. */ -ERRNO tsfile_reader_get_all_table_schemas_checked(TsFileReader reader, - TableSchema** out_schemas, - uint32_t* out_size); +TableSchema* tsfile_reader_get_table_schema(TsFileReader reader, + const char* table_name, + ERRNO* error_code); /** * @brief Gets all table schemas and reports metadata read failures. @@ -1052,18 +1033,20 @@ ERRNO tsfile_reader_get_all_table_schemas_checked(TsFileReader reader, * @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); +TableSchema* tsfile_reader_get_all_table_schemas(TsFileReader reader, + uint32_t* size, + ERRNO* error_code); /** - * @brief Gets all timeseries schema in the tsfile. - * - * @return DeviceSchema list, contains timeseries info. - * @note Caller should call free_device_schema and free to free the ptr. + * @brief Gets all timeseries schemas and reports metadata read failures. + * @return Allocated schema array, or NULL when there are no devices or on + * error. Check error_code to distinguish these cases. + * @note Caller must free each schema with free_device_schema, then free the + * array. */ DeviceSchema* tsfile_reader_get_all_timeseries_schemas(TsFileReader reader, - uint32_t* size); + uint32_t* size, + ERRNO* error_code); // ---------- Tag Filter API ---------- @@ -1111,45 +1094,6 @@ TagFilterHandle tsfile_tag_filter_between(TsFileReader reader, const char* lower, const char* upper, bool is_not, ERRNO* err_code); -/** - * @brief Create a tag equality filter: column == value. - * - * @param reader [in] Valid TsFileReader handle (used to resolve column index). - * @param table_name [in] Target table name. - * @param column_name [in] Tag column name. - * @param value [in] Value to compare against. - * @return TagFilterHandle on success, NULL on failure. - */ -TagFilterHandle tsfile_tag_filter_eq(TsFileReader reader, - const char* table_name, - const char* column_name, - const char* value); - -TagFilterHandle tsfile_tag_filter_neq(TsFileReader reader, - const char* table_name, - const char* column_name, - const char* value); - -TagFilterHandle tsfile_tag_filter_lt(TsFileReader reader, - const char* table_name, - const char* column_name, - const char* value); - -TagFilterHandle tsfile_tag_filter_lteq(TsFileReader reader, - const char* table_name, - const char* column_name, - const char* value); - -TagFilterHandle tsfile_tag_filter_gt(TsFileReader reader, - const char* table_name, - const char* column_name, - const char* value); - -TagFilterHandle tsfile_tag_filter_gteq(TsFileReader reader, - const char* table_name, - const char* column_name, - const char* value); - /** * @brief Logical AND of two tag filters. Takes ownership of left and right. */ diff --git a/cpp/src/reader/tsfile_reader.cc b/cpp/src/reader/tsfile_reader.cc index 6997cb389..de4f34167 100644 --- a/cpp/src/reader/tsfile_reader.cc +++ b/cpp/src/reader/tsfile_reader.cc @@ -728,13 +728,6 @@ ResultSet* TsFileReader::read_timeseries( return nullptr; } -std::shared_ptr TsFileReader::get_table_schema( - const std::string& table_name) { - std::shared_ptr table_schema; - get_table_schema(table_name, table_schema); - return table_schema; -} - int TsFileReader::get_table_schema(const std::string& table_name, std::shared_ptr& table_schema) { table_schema.reset(); diff --git a/cpp/src/reader/tsfile_reader.h b/cpp/src/reader/tsfile_reader.h index 11bb72913..f422ca7c7 100644 --- a/cpp/src/reader/tsfile_reader.h +++ b/cpp/src/reader/tsfile_reader.h @@ -256,15 +256,13 @@ class TsFileReader { TsFileProperties get_tsfile_properties(); /** - * @brief get the table schema by the table name + * @brief Get the table schema by table name. * * @param table_name the table name - * @return std::shared_ptr the table schema + * @param[out] table_schema the resolved schema, null on failure + * @return Returns 0 on success, E_TABLE_NOT_EXIST when the table is + * absent, or a non-zero read error code on metadata failure. */ - std::shared_ptr get_table_schema( - const std::string& table_name); - - /** Error-reporting overload. The output is null on failure. */ int get_table_schema(const std::string& table_name, std::shared_ptr& table_schema); /** diff --git a/cpp/test/cwrapper/c_release_test.cc b/cpp/test/cwrapper/c_release_test.cc index c8ac72346..0eedf9445 100644 --- a/cpp/test/cwrapper/c_release_test.cc +++ b/cpp/test/cwrapper/c_release_test.cc @@ -407,18 +407,25 @@ TEST_F(CReleaseTest, TsFileWriterConfTest) { free_write_file(&file); TsFileReader reader = tsfile_reader_new("plain_file.tsfile", &err_no); ASSERT_EQ(RET_OK, err_no); - TableSchema schema = tsfile_reader_get_table_schema(reader, "plain_table"); - ASSERT_EQ(schema.column_num, 2); + ERRNO schema_error = RET_OK; + TableSchema* schema = + tsfile_reader_get_table_schema(reader, "plain_table", &schema_error); + ASSERT_EQ(schema_error, RET_OK); + ASSERT_NE(schema, nullptr); + ASSERT_EQ(schema->column_num, 2); uint32_t size = 0; - DeviceSchema* device_schema = - tsfile_reader_get_all_timeseries_schemas(reader, &size); + ERRNO timeseries_schema_error = RET_OK; + DeviceSchema* device_schema = tsfile_reader_get_all_timeseries_schemas( + reader, &size, ×eries_schema_error); + ASSERT_EQ(timeseries_schema_error, RET_OK); ASSERT_EQ(1, size); ASSERT_EQ(1, device_schema->timeseries_num); ASSERT_EQ(device_schema->timeseries_schema[0].encoding, TS_ENCODING_PLAIN); ASSERT_EQ(device_schema->timeseries_schema[0].compression, TS_COMPRESSION_UNCOMPRESSED); tsfile_reader_close(reader); - free_table_schema(schema); + free_table_schema(*schema); + free(schema); free_device_schema(*device_schema); free(device_schema); free(column_list[0]); diff --git a/cpp/test/cwrapper/cwrapper_test.cc b/cpp/test/cwrapper/cwrapper_test.cc index 1d5309fd8..80db7cec1 100644 --- a/cpp/test/cwrapper/cwrapper_test.cc +++ b/cpp/test/cwrapper/cwrapper_test.cc @@ -328,10 +328,13 @@ TEST_F(CWrapperTest, WriterFlushTabletAndReadData) { row++; } ASSERT_EQ(row, num_timestamp); - uint32_t size; + uint32_t size = 0; + ERRNO all_schema_error = RET_OK; TableSchema* all_schema = - tsfile_reader_get_all_table_schemas(reader, &size); + tsfile_reader_get_all_table_schemas(reader, &size, &all_schema_error); + ASSERT_EQ(all_schema_error, RET_OK); ASSERT_EQ(1, size); + ASSERT_NE(all_schema, nullptr); ASSERT_EQ(std::string(all_schema[0].table_name), std::string(schema.table_name)); ASSERT_EQ(all_schema[0].column_num, schema.column_num); @@ -465,23 +468,6 @@ TEST(TagFilterCApiTest, RejectsNullInputs) { const char* col = "c"; const char* val = "v"; - EXPECT_EQ(tsfile_tag_filter_eq(nullptr, table, col, val), nullptr); - EXPECT_EQ(tsfile_tag_filter_eq(reinterpret_cast(1), nullptr, - col, val), - nullptr); - EXPECT_EQ(tsfile_tag_filter_eq(reinterpret_cast(1), table, - nullptr, val), - nullptr); - EXPECT_EQ(tsfile_tag_filter_eq(reinterpret_cast(1), table, - col, nullptr), - nullptr); - - EXPECT_EQ(tsfile_tag_filter_neq(nullptr, table, col, val), nullptr); - EXPECT_EQ(tsfile_tag_filter_lt(nullptr, table, col, val), nullptr); - EXPECT_EQ(tsfile_tag_filter_lteq(nullptr, table, col, val), nullptr); - EXPECT_EQ(tsfile_tag_filter_gt(nullptr, table, col, val), nullptr); - EXPECT_EQ(tsfile_tag_filter_gteq(nullptr, table, col, val), nullptr); - ERRNO err = common::E_OK; EXPECT_EQ( tsfile_tag_filter_create(nullptr, table, col, val, TAG_FILTER_EQ, &err), diff --git a/cpp/test/reader/table_read_failure_test.cc b/cpp/test/reader/table_read_failure_test.cc index 4d2ed13ce..393456cad 100644 --- a/cpp/test/reader/table_read_failure_test.cc +++ b/cpp/test/reader/table_read_failure_test.cc @@ -235,8 +235,8 @@ TEST_P(TableReadFailureTest, EveryReadFailureReachesCaller) { } } -TEST_P(TableReadFailureTest, CheckedSchemaAndTagFactoriesPreserveReadErrors) { - for (int operation = 0; operation < 3; ++operation) { +TEST_P(TableReadFailureTest, MetadataApisPreserveReadErrors) { + for (int operation = 0; operation < 4; ++operation) { storage::TsFileReader reader; auto* source = new FailingReadFile(bytes_); ASSERT_EQ( @@ -244,12 +244,21 @@ TEST_P(TableReadFailureTest, CheckedSchemaAndTagFactoriesPreserveReadErrors) { common::E_OK); source->fail_at = source->reads + 1; source->persistent = true; - ::TableSchema schema{}; + TableSchema* schema = nullptr; + ERRNO schema_error = common::E_OK; if (operation == 0) { - EXPECT_EQ(tsfile_reader_get_table_schema_checked(&reader, "test", - &schema), - common::E_FILE_READ_ERR); - EXPECT_EQ(schema.table_name, nullptr); + schema = + tsfile_reader_get_table_schema(&reader, "test", &schema_error); + EXPECT_EQ(schema_error, common::E_FILE_READ_ERR); + EXPECT_EQ(schema, nullptr); + } else if (operation == 3) { + uint32_t count = 0; + DeviceSchema* device_schemas = + tsfile_reader_get_all_timeseries_schemas(&reader, &count, + &schema_error); + EXPECT_EQ(schema_error, common::E_FILE_READ_ERR); + EXPECT_EQ(count, 0u); + EXPECT_EQ(device_schemas, nullptr); } else { ERRNO error = common::E_OK; TagFilterHandle filter = @@ -266,11 +275,27 @@ TEST_P(TableReadFailureTest, CheckedSchemaAndTagFactoriesPreserveReadErrors) { // Metadata failures must not poison the reader or cache an empty // schema. source->fail_at = 0; - ASSERT_EQ( - tsfile_reader_get_table_schema_checked(&reader, "test", &schema), - common::E_OK); - EXPECT_STREQ(schema.table_name, "test"); - free_table_schema(schema); + if (operation == 3) { + uint32_t count = 0; + DeviceSchema* device_schemas = + tsfile_reader_get_all_timeseries_schemas(&reader, &count, + &schema_error); + ASSERT_EQ(schema_error, common::E_OK); + ASSERT_NE(device_schemas, nullptr); + ASSERT_GT(count, 0u); + for (uint32_t i = 0; i < count; ++i) { + free_device_schema(device_schemas[i]); + } + free(device_schemas); + } else { + schema = + tsfile_reader_get_table_schema(&reader, "test", &schema_error); + ASSERT_EQ(schema_error, common::E_OK); + ASSERT_NE(schema, nullptr); + EXPECT_STREQ(schema->table_name, "test"); + free_table_schema(*schema); + free(schema); + } } } 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 b261f9cec..4f81ab380 100644 --- a/cpp/test/reader/table_view/tsfile_reader_table_test.cc +++ b/cpp/test/reader/table_view/tsfile_reader_table_test.cc @@ -393,7 +393,10 @@ TEST_F(TsFileTableReaderTest, TableModelGetSchema) { } } - auto table_schema = reader.get_table_schema("testtable0"); + std::shared_ptr table_schema; + ASSERT_EQ(reader.get_table_schema("testtable0", table_schema), + common::E_OK); + ASSERT_NE(table_schema, nullptr); ASSERT_EQ(table_schema->get_table_name(), "testtable0"); for (int i = 0; i < 5; i++) { ASSERT_EQ(table_schema->get_data_types()[i], TSDataType::STRING); diff --git a/cpp/tools/commands/cmd_count.cc b/cpp/tools/commands/cmd_count.cc index dca81a5f8..95bd0aa02 100644 --- a/cpp/tools/commands/cmd_count.cc +++ b/cpp/tools/commands/cmd_count.cc @@ -93,11 +93,16 @@ int collect_table_count(const ParsedArgs& args, storage::TsFileReader& reader, TableCountSummary& summary, std::ostream& err, bool require_all_measurements) { std::string table_name = storage::to_lower(args.table); - std::shared_ptr schema = - reader.get_table_schema(table_name); - if (!schema) { - err << "Error: table '" << args.table << "' does not exist\n"; - return kExitUsage; + std::shared_ptr schema; + const int schema_ret = reader.get_table_schema(table_name, schema); + if (schema_ret != common::E_OK) { + if (schema_ret == common::E_TABLE_NOT_EXIST) { + err << "Error: table '" << args.table << "' does not exist\n"; + return kExitUsage; + } + err << "Error: failed to read schema for table '" << args.table + << "': " << error_code_message(schema_ret) << "\n"; + return kExitFile; } summary.table_name = schema->get_table_name(); diff --git a/cpp/tools/commands/cmd_stats.cc b/cpp/tools/commands/cmd_stats.cc index f0ae7eee7..f5c26d9a0 100644 --- a/cpp/tools/commands/cmd_stats.cc +++ b/cpp/tools/commands/cmd_stats.cc @@ -415,14 +415,25 @@ int cmd_table_stats(const ParsedArgs& args, storage::TsFileReader& reader, OutputFormat fmt, std::ostream& out, std::ostream& err) { std::vector> schemas; if (!args.table.empty()) { - schemas.push_back( - reader.get_table_schema(storage::to_lower(args.table))); + std::shared_ptr schema; + const int schema_ret = + reader.get_table_schema(storage::to_lower(args.table), schema); + if (schema_ret != common::E_OK) { + if (schema_ret == common::E_TABLE_NOT_EXIST) { + err << "Error: table '" << args.table << "' does not exist\n"; + return kExitUsage; + } + err << "Error: failed to read schema for table '" << args.table + << "': " << error_code_message(schema_ret) << "\n"; + return kExitFile; + } + schemas.push_back(schema); } else { schemas = reader.get_all_table_schemas(); - } - if (schemas.empty() || !schemas[0]) { - err << "Error: table '" << args.table << "' does not exist\n"; - return kExitUsage; + if (schemas.empty() || !schemas[0]) { + err << "Error: table '" << args.table << "' does not exist\n"; + return kExitUsage; + } } if (args.table.empty()) { diff --git a/cpp/tools/commands/row_query.cc b/cpp/tools/commands/row_query.cc index 828bb6c55..b9eecf0f4 100644 --- a/cpp/tools/commands/row_query.cc +++ b/cpp/tools/commands/row_query.cc @@ -152,16 +152,23 @@ int resolve_tree_paths(const ParsedArgs& args, storage::TsFileReader& reader, } // namespace -std::unique_ptr build_table_tag_filter( - const ParsedArgs& args, storage::TsFileReader& reader, - const std::string& table_name, std::ostream& err) { +int build_table_tag_filter(const ParsedArgs& args, + storage::TsFileReader& reader, + const std::string& table_name, std::ostream& err, + std::unique_ptr& ret_filter) { if (!args.has_tag_filter) { - return std::unique_ptr(); + return kExitOk; } - auto schema = reader.get_table_schema(table_name); - if (!schema) { - err << "Error: no schema found for table " << table_name << "\n"; - return std::unique_ptr(); + std::shared_ptr schema; + const int schema_ret = reader.get_table_schema(table_name, schema); + if (schema_ret != common::E_OK) { + if (schema_ret == common::E_TABLE_NOT_EXIST) { + err << "Error: no schema found for table " << table_name << "\n"; + return kExitUsage; + } + err << "Error: failed to read schema for table " << table_name << ": " + << error_code_message(schema_ret) << "\n"; + return kExitFile; } storage::TagFilterBuilder builder(schema.get()); @@ -182,7 +189,7 @@ std::unique_ptr build_table_tag_filter( } catch (const std::regex_error&) { err << "Error: invalid regular expression for TAG '" << spec.column << "'\n"; - return std::unique_ptr(); + return kExitUsage; } { int tag_order = schema->find_id_column_order(spec.column); @@ -204,7 +211,7 @@ std::unique_ptr build_table_tag_filter( if (filter == nullptr) { err << "Error: invalid tag filter column '" << spec.column << "' for table " << table_name << "\n"; - return std::unique_ptr(); + return kExitUsage; } if (!combined) { combined.reset(filter); @@ -216,7 +223,8 @@ std::unique_ptr build_table_tag_filter( combined.release(), filter)); } } - return combined; + ret_filter = std::move(combined); + return kExitOk; } std::vector collect_tree_query_paths( @@ -296,15 +304,27 @@ int run_row_query(const ParsedArgs& args, storage::TsFileReader& reader, } table_name = schemas[0]->get_table_name(); } - auto table_schema = reader.get_table_schema(table_name); + std::shared_ptr table_schema; + const int schema_ret = + reader.get_table_schema(table_name, table_schema); + if (schema_ret != common::E_OK) { + if (schema_ret == common::E_TABLE_NOT_EXIST) { + err << "Error: table '" << table_name << "' does not exist\n"; + return kExitUsage; + } + err << "Error: failed to read schema for table '" << table_name + << "': " << error_code_message(schema_ret) << "\n"; + return kExitFile; + } std::vector cols; int selection_ret = resolve_table_fields(args, table_schema, cols, err); if (selection_ret != kExitOk) { return selection_ret; } - tag_filter = build_table_tag_filter(args, reader, table_name, err); - if (args.has_tag_filter && tag_filter == nullptr) { - return kExitUsage; + int filter_ret = + build_table_tag_filter(args, reader, table_name, err, tag_filter); + if (filter_ret != kExitOk) { + return filter_ret; } if (push_down) { qret = reader.queryByRow( diff --git a/go/tsfile/cgo_bridge.go b/go/tsfile/cgo_bridge.go index 5208c7219..123218dd1 100644 --- a/go/tsfile/cgo_bridge.go +++ b/go/tsfile/cgo_bridge.go @@ -757,19 +757,25 @@ func copyTableSchemaFromC(schema *C.TableSchema) TableSchema { func (h *readerHandle) tableSchema(table string) (TableSchema, error) { tableName := cStringPtr(table) defer freeCString(tableName) - var native C.TableSchema - code := C.tsfile_reader_get_table_schema_checked(C.TsFileReader(h.ptr), tableName, &native) + var code C.ERRNO + native := C.tsfile_reader_get_table_schema(C.TsFileReader(h.ptr), tableName, &code) if code != C.RET_OK { return TableSchema{}, newError("get table schema", cerrno(code)) } - defer C.free_table_schema(native) - return copyTableSchemaFromC(&native), nil + if native == nil { + return TableSchema{}, newError("get table schema", ErrFileRead.Code) + } + result := copyTableSchemaFromC(native) + C.free_table_schema(*native) + C.free(unsafe.Pointer(native)) + return result, nil } func (h *readerHandle) allTableSchemas() ([]TableSchema, error) { var native *C.TableSchema var count C.uint32_t - code := C.tsfile_reader_get_all_table_schemas_checked(C.TsFileReader(h.ptr), &native, &count) + var code C.ERRNO + native = C.tsfile_reader_get_all_table_schemas(C.TsFileReader(h.ptr), &count, &code) if code != C.RET_OK { return nil, newError("get all table schemas", cerrno(code)) } diff --git a/python/tests/test_reader_sources.py b/python/tests/test_reader_sources.py index 654692fa7..ce3c27612 100644 --- a/python/tests/test_reader_sources.py +++ b/python/tests/test_reader_sources.py @@ -406,7 +406,13 @@ def query(reader): @pytest.mark.parametrize( - "method", ["get_all_devices", "get_all_table_schemas", "get_table_schema"] + "method", + [ + "get_all_devices", + "get_all_table_schemas", + "get_table_schema", + "get_all_timeseries_schemas", + ], ) def test_metadata_read_failure_is_not_an_empty_result(method): class FailingBytesIO(TrackingBytesIO): diff --git a/python/tsfile/tsfile_cpp.pxd b/python/tsfile/tsfile_cpp.pxd index e97314d15..4ceebb1e7 100644 --- a/python/tsfile/tsfile_cpp.pxd +++ b/python/tsfile/tsfile_cpp.pxd @@ -355,17 +355,12 @@ cdef extern from "cwrapper/tsfile_cwrapper.h": char ** sensor_name, uint32_t sensor_num, int64_t start_time, int64_t end_time, ErrorCode *err_code) - TableSchema tsfile_reader_get_table_schema(TsFileReader reader, - const char * table_name); - ErrorCode tsfile_reader_get_table_schema_checked( - TsFileReader reader, const char * table_name, TableSchema * out_schema); - - TableSchema * tsfile_reader_get_all_table_schemas(TsFileReader reader, - uint32_t * size); - TableSchema * tsfile_reader_get_all_table_schemas_with_error( + TableSchema * tsfile_reader_get_table_schema( + TsFileReader reader, const char * table_name, ErrorCode * error_code); + TableSchema * tsfile_reader_get_all_table_schemas( + TsFileReader reader, uint32_t * size, ErrorCode * error_code); + DeviceSchema * tsfile_reader_get_all_timeseries_schemas( TsFileReader reader, uint32_t * size, ErrorCode * error_code); - DeviceSchema * tsfile_reader_get_all_timeseries_schemas(TsFileReader reader, - uint32_t * size); void tsfile_device_id_free_contents(DeviceID * d) diff --git a/python/tsfile/tsfile_py_cpp.pyx b/python/tsfile/tsfile_py_cpp.pyx index 05fa3cc98..6bd42cb89 100644 --- a/python/tsfile/tsfile_py_cpp.pyx +++ b/python/tsfile/tsfile_py_cpp.pyx @@ -1236,11 +1236,16 @@ cdef ResultSet tsfile_reader_query_table_with_tag_filter_c(TsFileReader reader, cdef object get_table_schema(TsFileReader reader, object table_name): cdef bytes table_name_bytes = PyUnicode_AsUTF8String(table_name) cdef const char * table_name_c = table_name_bytes - cdef TableSchema schema - cdef ErrorCode code = tsfile_reader_get_table_schema_checked( - reader, table_name_c, &schema) + cdef TableSchema * schema + cdef ErrorCode code = 0 + schema = tsfile_reader_get_table_schema( + reader, table_name_c, &code) check_error(code) - return from_c_table_schema(schema) + if schema == NULL: + raise RuntimeError("tsfile_reader_get_table_schema returned NULL") + schema_py = from_c_table_schema(schema[0]) + free(schema) + return schema_py cdef object get_all_table_schema(TsFileReader reader): cdef uint32_t table_num = 0 @@ -1249,7 +1254,7 @@ cdef object get_all_table_schema(TsFileReader reader): cdef int i table_schemas = {} - schemas = tsfile_reader_get_all_table_schemas_with_error(reader, &table_num, &error_code) + schemas = tsfile_reader_get_all_table_schemas(reader, &table_num, &error_code) check_error(error_code) for i in range(table_num): schema_py = from_c_table_schema(schemas[i]) @@ -1260,10 +1265,13 @@ cdef object get_all_table_schema(TsFileReader reader): cdef object get_all_timeseries_schema(TsFileReader reader): cdef uint32_t device_num = 0 cdef DeviceSchema * schemas + cdef ErrorCode error_code = 0 cdef int i device_schemas = {} - schemas = tsfile_reader_get_all_timeseries_schemas(reader, &device_num) + schemas = tsfile_reader_get_all_timeseries_schemas( + reader, &device_num, &error_code) + check_error(error_code) for i in range(device_num): schema_py = from_c_device_schema(schemas[i]) device_schemas.update([(schema_py.get_device_name(), schema_py)])