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 5cae8f01d..919437aca 100644 --- a/cpp/src/cwrapper/tsfile_cwrapper.cc +++ b/cpp/src/cwrapper/tsfile_cwrapper.cc @@ -1058,35 +1058,15 @@ 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) { - 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; -} - static ERRNO copy_table_schema(const std::shared_ptr& src, TableSchema* out_schema) { - if (!src || out_schema == nullptr) { - return common::E_TABLE_NOT_EXIST; + if (out_schema == nullptr) { + return common::E_INVALID_ARG; } *out_schema = TableSchema{}; + if (!src) { + return common::E_TABLE_NOT_EXIST; + } out_schema->table_name = strdup(src->get_table_name().c_str()); if (out_schema->table_name == nullptr) { return common::E_OOM; @@ -1104,7 +1084,18 @@ static ERRNO copy_table_schema(const std::shared_ptr& src, } const auto& measurements = src->get_measurement_schemas(); const auto& categories = src->get_column_categories(); + if (measurements.size() < static_cast(out_schema->column_num) || + categories.size() < static_cast(out_schema->column_num)) { + free_table_schema(*out_schema); + *out_schema = TableSchema{}; + return common::E_INVALID_SCHEMA; + } for (int i = 0; i < out_schema->column_num; ++i) { + if (!measurements[i]) { + free_table_schema(*out_schema); + *out_schema = TableSchema{}; + return common::E_INVALID_SCHEMA; + } out_schema->column_schemas[i].column_name = strdup(measurements[i]->measurement_name_.c_str()); if (out_schema->column_schemas[i].column_name == nullptr) { @@ -1120,17 +1111,44 @@ static ERRNO copy_table_schema(const std::shared_ptr& src, return common::E_OK; } +static void free_table_schema_array(TableSchema* schemas, size_t size) { + if (schemas == nullptr) { + return; + } + for (size_t i = 0; i < size; ++i) { + free_table_schema(schemas[i]); + } + free(schemas); +} + +static void free_device_schema_array(DeviceSchema* schemas, size_t size) { + if (schemas == nullptr) { + return; + } + for (size_t i = 0; i < size; ++i) { + free_device_schema(schemas[i]); + } + free(schemas); +} + 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) { + if (out_schema == nullptr) { return common::E_INVALID_ARG; } *out_schema = TableSchema{}; + if (reader == nullptr || table_name == nullptr) { + return common::E_INVALID_ARG; + } 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; @@ -1139,136 +1157,165 @@ ERRNO tsfile_reader_get_table_schema_checked(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); - 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)); - for (size_t i = 0; i < table_schemas.size(); i++) { - ret[i].table_name = strdup(table_schemas[i]->get_table_name().c_str()); - int column_num = table_schemas[i]->get_columns_num(); - ret[i].column_num = column_num; - ret[i].column_schemas = static_cast( - malloc(column_num * sizeof(ColumnSchema))); - auto column_schemas = table_schemas[i]->get_measurement_schemas(); - for (int j = 0; j < column_num; j++) { - ret[i].column_schemas[j].column_name = - strdup(column_schemas[j]->measurement_name_.c_str()); - ret[i].column_schemas[j].data_type = - static_cast(column_schemas[j]->data_type_); - ret[i].column_schemas[j].column_category = - static_cast( - table_schemas[i]->get_column_categories()[j]); - } +TableSchema tsfile_reader_get_table_schema(TsFileReader reader, + const char* table_name) { + TableSchema schema{}; + if (tsfile_reader_get_table_schema_checked(reader, table_name, &schema) != + common::E_OK) { + return TableSchema{}; } - *size = table_num; - return ret; + return schema; } ERRNO tsfile_reader_get_all_table_schemas_checked(TsFileReader reader, TableSchema** out_schemas, uint32_t* out_size) { + if (out_schemas != nullptr) { + *out_schemas = nullptr; + } + if (out_size != nullptr) { + *out_size = 0; + } if (reader == nullptr || out_schemas == nullptr || out_size == nullptr) { return common::E_INVALID_ARG; } - *out_schemas = nullptr; - *out_size = 0; + TableSchema* result = nullptr; + size_t initialized = 0; try { - auto schemas = static_cast(reader) - ->get_all_table_schemas(); - if (schemas.empty()) { + auto* r = static_cast(reader); + std::vector> table_schemas; + const ERRNO read_ret = r->get_all_table_schemas(table_schemas); + if (read_ret != common::E_OK) { + return read_ret; + } + if (table_schemas.empty()) { return common::E_OK; } - TableSchema* copied = static_cast( - calloc(schemas.size(), sizeof(TableSchema))); - if (copied == nullptr) { + const size_t table_num = table_schemas.size(); + if (table_num > std::numeric_limits::max()) { + return common::E_OVERFLOW; + } + result = static_cast(calloc(table_num, sizeof(*result))); + if (result == 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; + for (size_t i = 0; i < table_num; ++i) { + initialized = i + 1; + const ERRNO copy_ret = + copy_table_schema(table_schemas[i], &result[i]); + if (copy_ret != common::E_OK) { + free_table_schema_array(result, initialized); + return copy_ret; } } - *out_schemas = copied; - *out_size = static_cast(schemas.size()); + *out_schemas = result; + *out_size = static_cast(table_num); return common::E_OK; } catch (const std::bad_alloc&) { + free_table_schema_array(result, initialized); return common::E_OOM; } catch (...) { + free_table_schema_array(result, initialized); return common::E_FILE_READ_ERR; } } -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) { - return nullptr; +TableSchema* tsfile_reader_get_all_table_schemas_with_error(TsFileReader reader, + uint32_t* size, + ERRNO* error_code) { + if (size != nullptr) { + *size = 0; } - *size = static_cast(device_ids.size()); - if (device_ids.empty()) { + if (error_code == nullptr) { return nullptr; } + TableSchema* schemas = nullptr; + *error_code = + tsfile_reader_get_all_table_schemas_checked(reader, &schemas, size); + return *error_code == common::E_OK ? schemas : nullptr; +} - DeviceSchema* device_schema = static_cast( - malloc(sizeof(DeviceSchema) * device_ids.size())); - if (device_schema == nullptr) { - *size = 0; - return nullptr; +TableSchema* tsfile_reader_get_all_table_schemas(TsFileReader reader, + uint32_t* size) { + TableSchema* schemas = nullptr; + const ERRNO ret = + tsfile_reader_get_all_table_schemas_checked(reader, &schemas, size); + return ret == common::E_OK ? schemas : nullptr; +} + +ERRNO tsfile_reader_get_all_timeseries_schemas_checked( + TsFileReader reader, DeviceSchema** out_schemas, uint32_t* out_size) { + if (out_schemas != nullptr) { + *out_schemas = nullptr; + } + if (out_size != nullptr) { + *out_size = 0; + } + if (reader == nullptr || out_schemas == nullptr || out_size == nullptr) { + return common::E_INVALID_ARG; + } + DeviceSchema* result = nullptr; + size_t initialized = 0; + auto* r = static_cast(reader); + std::vector> device_ids; + const ERRNO read_ret = r->get_all_devices(device_ids); + if (read_ret != common::E_OK) { + return read_ret; + } + if (device_ids.empty()) { + return common::E_OK; + } + const size_t device_count = device_ids.size(); + if (device_count > std::numeric_limits::max()) { + return common::E_OVERFLOW; + } + result = static_cast(calloc(device_count, sizeof(*result))); + if (result == nullptr) { + return common::E_OOM; } - 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_count; ++device_index) { + initialized = device_index + 1; + const auto& device_id = device_ids[device_index]; + DeviceSchema& cur_schema = result[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_device_schema_array(result, initialized); + return common::E_OOM; + } 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_device_schema_array(result, initialized); + return ret; + } + if (schemas.empty()) { continue; } + if (schemas.size() > + static_cast(std::numeric_limits::max())) { + free_device_schema_array(result, initialized); + return common::E_OVERFLOW; + } 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_device_schema_array(result, initialized); + return common::E_OOM; + } 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_device_schema_array(result, initialized); + return common::E_OOM; + } cur_schema.timeseries_schema[i].data_type = static_cast(measurement_schema.data_type_); cur_schema.timeseries_schema[i].encoding = @@ -1278,7 +1325,17 @@ DeviceSchema* tsfile_reader_get_all_timeseries_schemas(TsFileReader reader, measurement_schema.compression_type_); } } - return device_schema; + *out_schemas = result; + *out_size = static_cast(device_count); + return common::E_OK; +} + +DeviceSchema* tsfile_reader_get_all_timeseries_schemas(TsFileReader reader, + uint32_t* size) { + DeviceSchema* schemas = nullptr; + const ERRNO ret = tsfile_reader_get_all_timeseries_schemas_checked( + reader, &schemas, size); + return ret == common::E_OK ? schemas : nullptr; } void tsfile_device_id_free_contents(DeviceID* d) { @@ -2508,55 +2565,23 @@ 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, - const char* value, TagFilterOp op, - ERRNO* err_code) { - if (err_code == nullptr) { - return nullptr; +ERRNO tsfile_tag_filter_create_checked(TsFileReader reader, + const char* table_name, + const char* column_name, + const char* value, TagFilterOp op, + TagFilterHandle* out_filter) { + if (out_filter != nullptr) { + *out_filter = nullptr; } if (reader == nullptr || table_name == nullptr || column_name == nullptr || - value == nullptr) { - *err_code = common::E_INVALID_ARG; - return nullptr; + value == nullptr || out_filter == nullptr) { + return common::E_INVALID_ARG; } auto* r = static_cast(reader); - auto schema = r->get_table_schema(table_name); - if (!schema) { - *err_code = common::E_INVALID_ARG; - return nullptr; + std::shared_ptr schema; + const int ret = r->get_table_schema(table_name, schema); + if (ret != common::E_OK) { + return ret; } storage::TagFilterBuilder builder(schema.get()); storage::Filter* filter = nullptr; @@ -2592,48 +2617,97 @@ TagFilterHandle tsfile_tag_filter_create(TsFileReader reader, filter = builder.is_not_null(column_name); break; default: - *err_code = common::E_INVALID_ARG; - return nullptr; + return common::E_INVALID_ARG; } if (filter == nullptr) { - *err_code = common::E_COLUMN_NOT_EXIST; - return nullptr; + return common::E_COLUMN_NOT_EXIST; } - *err_code = common::E_OK; - return static_cast(filter); + *out_filter = static_cast(filter); + return common::E_OK; } -TagFilterHandle tsfile_tag_filter_between(TsFileReader reader, - const char* table_name, - const char* column_name, - const char* lower, const char* upper, - bool is_not, ERRNO* err_code) { +TagFilterHandle tsfile_tag_filter_create(TsFileReader reader, + const char* table_name, + const char* column_name, + const char* value, TagFilterOp op, + ERRNO* err_code) { if (err_code == nullptr) { return nullptr; } - if (reader == nullptr || table_name == nullptr || column_name == nullptr || - lower == nullptr || upper == nullptr) { + TagFilterHandle filter = nullptr; + *err_code = tsfile_tag_filter_create_checked( + reader, table_name, column_name, value, op, &filter); + if (*err_code == common::E_TABLE_NOT_EXIST) { *err_code = common::E_INVALID_ARG; - return nullptr; + } + return *err_code == common::E_OK ? filter : nullptr; +} + +ERRNO tsfile_tag_filter_between_checked(TsFileReader reader, + const char* table_name, + const char* column_name, + const char* lower, const char* upper, + bool is_not, + TagFilterHandle* out_filter) { + if (out_filter != nullptr) { + *out_filter = nullptr; + } + if (reader == nullptr || table_name == nullptr || column_name == nullptr || + lower == nullptr || upper == nullptr || out_filter == nullptr) { + return common::E_INVALID_ARG; } auto* r = static_cast(reader); - auto schema = r->get_table_schema(table_name); - if (!schema) { - *err_code = common::E_INVALID_ARG; - return nullptr; + std::shared_ptr schema; + const int ret = r->get_table_schema(table_name, schema); + if (ret != common::E_OK) { + return ret; } storage::TagFilterBuilder builder(schema.get()); storage::Filter* filter = is_not ? builder.not_between_and(column_name, lower, upper) : builder.between_and(column_name, lower, upper); if (filter == nullptr) { - *err_code = common::E_COLUMN_NOT_EXIST; + return common::E_COLUMN_NOT_EXIST; + } + *out_filter = static_cast(filter); + return common::E_OK; +} + +TagFilterHandle tsfile_tag_filter_between(TsFileReader reader, + const char* table_name, + const char* column_name, + const char* lower, const char* upper, + bool is_not, ERRNO* err_code) { + if (err_code == nullptr) { return nullptr; } - *err_code = common::E_OK; - return static_cast(filter); + TagFilterHandle filter = nullptr; + *err_code = tsfile_tag_filter_between_checked( + reader, table_name, column_name, lower, upper, is_not, &filter); + if (*err_code == common::E_TABLE_NOT_EXIST) { + *err_code = common::E_INVALID_ARG; + } + return *err_code == common::E_OK ? filter : nullptr; } +#define DEFINE_LEGACY_TAG_FILTER_FACTORY(name, op) \ + TagFilterHandle tsfile_tag_filter_##name( \ + TsFileReader reader, const char* table_name, const char* column_name, \ + const char* value) { \ + ERRNO error_code = common::E_OK; \ + return tsfile_tag_filter_create(reader, table_name, column_name, \ + value, TAG_FILTER_##op, &error_code); \ + } + +DEFINE_LEGACY_TAG_FILTER_FACTORY(eq, EQ) +DEFINE_LEGACY_TAG_FILTER_FACTORY(neq, NEQ) +DEFINE_LEGACY_TAG_FILTER_FACTORY(lt, LT) +DEFINE_LEGACY_TAG_FILTER_FACTORY(lteq, LTEQ) +DEFINE_LEGACY_TAG_FILTER_FACTORY(gt, GT) +DEFINE_LEGACY_TAG_FILTER_FACTORY(gteq, GTEQ) + +#undef DEFINE_LEGACY_TAG_FILTER_FACTORY + TagFilterHandle tsfile_tag_filter_and(TagFilterHandle left, TagFilterHandle right) { if (!left || !right) return nullptr; diff --git a/cpp/src/cwrapper/tsfile_cwrapper.h b/cpp/src/cwrapper/tsfile_cwrapper.h index 003e59198..83e4c2dbf 100644 --- a/cpp/src/cwrapper/tsfile_cwrapper.h +++ b/cpp/src/cwrapper/tsfile_cwrapper.h @@ -1017,54 +1017,81 @@ 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 a table schema using the legacy value-returning API. + * @return A populated schema, or a zero-initialized schema when the lookup + * fails. Use tsfile_reader_get_table_schema_checked() when the error code is + * required. + * @note Release the returned schema's contents with free_table_schema(). */ TableSchema tsfile_reader_get_table_schema(TsFileReader reader, const char* table_name); -/** Retrieves one table schema and reports missing tables through ERRNO. */ +/** + * @brief Gets one table schema and returns the error code directly. + * @param out_schema Required output storage with no owned allocations. + * It is zero-initialized on failure; partial results are released internally. + * @return RET_OK on success, or the lookup, read, or allocation error. + * @note On success, release the contents with free_table_schema(*out_schema). + * The caller owns the output storage itself. + */ 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. + * @brief Gets all table schemas using the legacy pointer-returning API. + * @return Schema array, or NULL when there are no tables or an error occurs. + * Use tsfile_reader_get_all_table_schemas_checked() when the error code is + * required. + * @note size is zero on failure. Release each schema with free_table_schema(), + * then release the array with free(). */ TableSchema* tsfile_reader_get_all_table_schemas(TsFileReader reader, uint32_t* size); -/** Retrieves every table schema into a caller-freed array. */ +/** + * @brief Gets all table schemas and returns the error code directly. + * Both output parameters are required. Non-null outputs are set to NULL/0 + * before work begins and remain empty on failure. An empty file result is + * successful with RET_OK and NULL/0 outputs. + * @note On success, release each schema with free_table_schema(), then free() + * the array. Partial results are released internally on failure. + */ ERRNO tsfile_reader_get_all_table_schemas_checked(TsFileReader reader, TableSchema** out_schemas, uint32_t* out_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. + * @brief Compatibility API reporting errors through a required error_code. + * Returns NULL for an empty result or an error; error_code distinguishes them. + * size is zero on failure. Ownership is the same as the checked array API. */ 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. - * - * @return DeviceSchema list, contains timeseries info. - * @note Caller should call free_device_schema and free to free the ptr. + * @brief Gets all timeseries schemas using the legacy pointer-returning API. + * Use tsfile_reader_get_all_timeseries_schemas_checked() when the error code + * is required. + * @return Array on success, or NULL for an empty result or an error. + * @note size is zero on failure. Release each schema with free_device_schema(), + * then release the array with free(). */ DeviceSchema* tsfile_reader_get_all_timeseries_schemas(TsFileReader reader, uint32_t* size); +/** + * @brief Gets all timeseries schemas and returns the error code directly. + * Both output parameters are required. Non-null outputs are set to NULL/0 + * before work begins and remain empty on failure. No devices is a successful + * result with RET_OK and NULL/0 outputs. + * @note On success, release each schema with free_device_schema(), then free() + * the array. Partial results are released internally on failure. + */ +ERRNO tsfile_reader_get_all_timeseries_schemas_checked( + TsFileReader reader, DeviceSchema** out_schemas, uint32_t* out_size); + // ---------- Tag Filter API ---------- /** @@ -1112,39 +1139,52 @@ TagFilterHandle tsfile_tag_filter_between(TsFileReader reader, 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. + * Create a tag filter and return the error code directly. out_filter is + * required and is set to NULL on failure. Unlike the legacy create API, + * a missing table is reported as RET_TABLE_NOT_EXIST, not RET_INVALID_ARG. + * Release a successful result with tsfile_tag_filter_free(). + */ +ERRNO tsfile_tag_filter_create_checked(TsFileReader reader, + const char* table_name, + const char* column_name, + const char* value, TagFilterOp op, + TagFilterHandle* out_filter); + +/** + * Create a BETWEEN tag filter and return the error code directly. + * Output ownership and missing-table handling match the checked create API. + */ +ERRNO tsfile_tag_filter_between_checked(TsFileReader reader, + const char* table_name, + const char* column_name, + const char* lower, const char* upper, + bool is_not, + TagFilterHandle* out_filter); + +/** + * Legacy tag-filter factories returning NULL on failure. Use the checked + * create API for an error code. Free results with tsfile_tag_filter_free(). */ 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, 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..00d5a8b86 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,16 @@ class DeviceMetaIterator { int load_results_direct(); TsFileIOReader* io_reader_; - std::queue meta_index_nodes_; + // The bool tracks ownership: roots borrowed from file metadata are false; + // descendants allocated in pa_ are true and need explicit destruction. + 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..90fb71ef1 100644 --- a/cpp/src/reader/tsfile_reader.cc +++ b/cpp/src/reader/tsfile_reader.cc @@ -730,14 +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); - return table_schema; + std::shared_ptr schema; + if (get_table_schema(table_name, schema) != E_OK) { + return nullptr; + } + return 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> diff --git a/cpp/src/reader/tsfile_reader.h b/cpp/src/reader/tsfile_reader.h index 6493a4af7..368790211 100644 --- a/cpp/src/reader/tsfile_reader.h +++ b/cpp/src/reader/tsfile_reader.h @@ -256,13 +256,23 @@ class TsFileReader { TsFileProperties get_tsfile_properties(); /** - * @brief get the table schema by the table name - * - * @param table_name the table name - * @return std::shared_ptr the table schema + * @brief Legacy lookup returning null on failure. + * Use the error-reporting overload to distinguish a missing table from + * a metadata read failure. */ std::shared_ptr get_table_schema( const std::string& table_name); + + /** + * @brief Get the table schema by table name. + * + * @param table_name the table name + * @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. + */ + 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/cwrapper/c_release_test.cc b/cpp/test/cwrapper/c_release_test.cc index c8ac72346..27f02d56c 100644 --- a/cpp/test/cwrapper/c_release_test.cc +++ b/cpp/test/cwrapper/c_release_test.cc @@ -407,16 +407,45 @@ 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"); + ERRNO schema_error = RET_OK; + TableSchema schema{}; + schema_error = + tsfile_reader_get_table_schema_checked(reader, "plain_table", &schema); + ASSERT_EQ(schema_error, RET_OK); 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 = nullptr; + timeseries_schema_error = tsfile_reader_get_all_timeseries_schemas_checked( + reader, &device_schema, &size); + 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); + + TableSchema legacy_schema = + tsfile_reader_get_table_schema(reader, "plain_table"); + ASSERT_EQ(legacy_schema.column_num, 2); + free_table_schema(legacy_schema); + + uint32_t legacy_table_count = 0; + TableSchema* legacy_tables = + tsfile_reader_get_all_table_schemas(reader, &legacy_table_count); + ASSERT_EQ(legacy_table_count, 1); + ASSERT_NE(legacy_tables, nullptr); + free_table_schema(legacy_tables[0]); + free(legacy_tables); + + uint32_t legacy_device_count = 0; + DeviceSchema* legacy_devices = + tsfile_reader_get_all_timeseries_schemas(reader, &legacy_device_count); + ASSERT_EQ(legacy_device_count, 1); + ASSERT_NE(legacy_devices, nullptr); + free_device_schema(legacy_devices[0]); + free(legacy_devices); + tsfile_reader_close(reader); free_table_schema(schema); free_device_schema(*device_schema); diff --git a/cpp/test/cwrapper/cwrapper_test.cc b/cpp/test/cwrapper/cwrapper_test.cc index 1d5309fd8..c3e66f08d 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; - TableSchema* all_schema = - tsfile_reader_get_all_table_schemas(reader, &size); + uint32_t size = 0; + ERRNO all_schema_error = RET_OK; + TableSchema* all_schema = tsfile_reader_get_all_table_schemas_with_error( + 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), @@ -512,4 +498,46 @@ TEST(TagFilterCApiTest, RejectsNullInputs) { nullptr); } +TEST(CApiCompatibilityTest, KeepsLegacySchemaAndTagFilterEntryPoints) { + TableSchema legacy_schema = + tsfile_reader_get_table_schema(nullptr, "missing"); + EXPECT_EQ(legacy_schema.table_name, nullptr); + EXPECT_EQ(legacy_schema.column_num, 0); + EXPECT_EQ(legacy_schema.column_schemas, nullptr); + + uint32_t table_count = 7; + EXPECT_EQ(tsfile_reader_get_all_table_schemas(nullptr, &table_count), + nullptr); + EXPECT_EQ(table_count, 0); + + uint32_t device_count = 7; + EXPECT_EQ(tsfile_reader_get_all_timeseries_schemas(nullptr, &device_count), + nullptr); + EXPECT_EQ(device_count, 0); + + TableSchema checked_schema{}; + EXPECT_EQ(tsfile_reader_get_table_schema_checked(nullptr, "missing", + &checked_schema), + common::E_INVALID_ARG); + + TableSchema* checked_tables = nullptr; + table_count = 7; + EXPECT_EQ(tsfile_reader_get_all_table_schemas_checked( + nullptr, &checked_tables, &table_count), + common::E_INVALID_ARG); + EXPECT_EQ(checked_tables, nullptr); + EXPECT_EQ(table_count, 0); + + DeviceSchema* checked_devices = nullptr; + device_count = 7; + EXPECT_EQ(tsfile_reader_get_all_timeseries_schemas_checked( + nullptr, &checked_devices, &device_count), + common::E_INVALID_ARG); + EXPECT_EQ(checked_devices, nullptr); + EXPECT_EQ(device_count, 0); + + EXPECT_EQ(tsfile_tag_filter_eq(nullptr, "missing", "tag", "value"), + nullptr); +} + } // namespace cwrapper 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..05e62ee2d --- /dev/null +++ b/cpp/test/reader/table_read_failure_test.cc @@ -0,0 +1,469 @@ +/* + * 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 + +#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, MetadataApisPreserveReadErrors) { + for (int operation = 0; operation < 5; ++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{}; + ERRNO schema_error = common::E_OK; + if (operation == 0) { + schema_error = tsfile_reader_get_table_schema_checked( + &reader, "test", &schema); + EXPECT_EQ(schema_error, common::E_FILE_READ_ERR); + EXPECT_EQ(schema.table_name, nullptr); + EXPECT_EQ(schema.column_num, 0); + EXPECT_EQ(schema.column_schemas, nullptr); + } else if (operation == 3) { + uint32_t count = 0; + DeviceSchema* device_schemas = nullptr; + schema_error = tsfile_reader_get_all_timeseries_schemas_checked( + &reader, &device_schemas, &count); + EXPECT_EQ(schema_error, common::E_FILE_READ_ERR); + EXPECT_EQ(count, 0u); + EXPECT_EQ(device_schemas, nullptr); + } else if (operation == 4) { + uint32_t count = 0; + TableSchema* schemas = nullptr; + schema_error = tsfile_reader_get_all_table_schemas_checked( + &reader, &schemas, &count); + EXPECT_EQ(schema_error, common::E_FILE_READ_ERR); + EXPECT_EQ(count, 0u); + EXPECT_EQ(schemas, 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; + if (operation == 3) { + uint32_t count = 0; + DeviceSchema* device_schemas = nullptr; + schema_error = tsfile_reader_get_all_timeseries_schemas_checked( + &reader, &device_schemas, &count); + 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 if (operation == 4) { + uint32_t count = 0; + TableSchema* schemas = nullptr; + schema_error = tsfile_reader_get_all_table_schemas_checked( + &reader, &schemas, &count); + ASSERT_EQ(schema_error, common::E_OK); + ASSERT_NE(schemas, nullptr); + ASSERT_GT(count, 0u); + for (uint32_t i = 0; i < count; ++i) { + free_table_schema(schemas[i]); + } + free(schemas); + } else { + schema_error = tsfile_reader_get_table_schema_checked( + &reader, "test", &schema); + ASSERT_EQ(schema_error, common::E_OK); + EXPECT_STREQ(schema.table_name, "test"); + free_table_schema(schema); + } + } +} + +TEST_P(TableReadFailureTest, LegacySchemaApisReturnEmptyOnReadFailure) { + for (int operation = 0; operation < 4; ++operation) { + SCOPED_TRACE(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; + if (operation == 0) { + TableSchema schema = + tsfile_reader_get_table_schema(&reader, "test"); + EXPECT_EQ(schema.table_name, nullptr); + EXPECT_EQ(schema.column_num, 0); + EXPECT_EQ(schema.column_schemas, nullptr); + free_table_schema(schema); + } else if (operation == 1) { + uint32_t count = 7; + TableSchema* schemas = + tsfile_reader_get_all_table_schemas(&reader, &count); + EXPECT_EQ(schemas, nullptr); + EXPECT_EQ(count, 0u); + if (schemas != nullptr) { + for (uint32_t i = 0; i < count; ++i) + free_table_schema(schemas[i]); + free(schemas); + } + } else if (operation == 2) { + uint32_t count = 7; + DeviceSchema* schemas = + tsfile_reader_get_all_timeseries_schemas(&reader, &count); + EXPECT_EQ(schemas, nullptr); + EXPECT_EQ(count, 0u); + if (schemas != nullptr) { + for (uint32_t i = 0; i < count; ++i) + free_device_schema(schemas[i]); + free(schemas); + } + } else { + EXPECT_EQ(reader.get_table_schema("test"), nullptr); + } + EXPECT_TRUE(source->failed); + source->fail_at = 0; + auto schema = reader.get_table_schema("test"); + ASSERT_NE(schema, nullptr); + EXPECT_EQ(schema->get_table_name(), "test"); + } +} + +TEST_P(TableReadFailureTest, LegacyTagFactoriesReturnNullOnReadFailure) { + using Factory = TagFilterHandle (*)(TsFileReader, const char*, const char*, + const char*); + const Factory factories[] = {tsfile_tag_filter_eq, tsfile_tag_filter_neq, + tsfile_tag_filter_lt, tsfile_tag_filter_lteq, + tsfile_tag_filter_gt, tsfile_tag_filter_gteq}; + for (auto factory : factories) { + 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; + TagFilterHandle filter = factory(&reader, "test", "id0", "d0"); + EXPECT_EQ(filter, nullptr); + tsfile_tag_filter_free(filter); + EXPECT_TRUE(source->failed); + source->fail_at = 0; + filter = factory(&reader, "test", "id0", "d0"); + ASSERT_NE(filter, nullptr); + tsfile_tag_filter_free(filter); + } +} + +TEST_P(TableReadFailureTest, CheckedTagFactoriesPreserveMissingTableError) { + storage::TsFileReader reader; + ASSERT_EQ(reader.open(std::unique_ptr( + new FailingReadFile(bytes_))), + common::E_OK); + TagFilterHandle filter = &reader; + EXPECT_EQ(tsfile_tag_filter_create_checked(&reader, "missing", "id0", "d0", + TAG_FILTER_EQ, &filter), + common::E_TABLE_NOT_EXIST); + ASSERT_EQ(filter, nullptr); + filter = &reader; + EXPECT_EQ(tsfile_tag_filter_between_checked(&reader, "missing", "id0", "d0", + "d1", false, &filter), + common::E_TABLE_NOT_EXIST); + ASSERT_EQ(filter, nullptr); + + ERRNO code = common::E_OK; + EXPECT_EQ(tsfile_tag_filter_create(&reader, "missing", "id0", "d0", + TAG_FILTER_EQ, &code), + nullptr); + EXPECT_EQ(code, common::E_INVALID_ARG); + code = common::E_OK; + EXPECT_EQ(tsfile_tag_filter_between(&reader, "missing", "id0", "d0", "d1", + false, &code), + nullptr); + EXPECT_EQ(code, common::E_INVALID_ARG); +} + +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/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 9d8c75df6..e6daccfbc 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 01bd5ba57..fbc62de77 100644 --- a/cpp/tools/commands/cmd_stats.cc +++ b/cpp/tools/commands/cmd_stats.cc @@ -415,14 +415,38 @@ 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 = sorted_table_schemas(reader); - } - if (schemas.empty() || !schemas[0]) { - err << "Error: table '" << args.table << "' does not exist\n"; - return kExitUsage; + const int schemas_ret = reader.get_all_table_schemas(schemas); + if (schemas_ret != common::E_OK) { + err << "Error: failed to read table schemas: " + << error_code_message(schemas_ret) << "\n"; + return kExitFile; + } + std::sort( + schemas.begin(), schemas.end(), + [](const std::shared_ptr& lhs, + const std::shared_ptr& rhs) { + if (!lhs) return false; + if (!rhs) return true; + return lhs->get_table_name() < rhs->get_table_name(); + }); + 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/commands.h b/cpp/tools/commands/commands.h index eb786cdd4..f32bb4750 100644 --- a/cpp/tools/commands/commands.h +++ b/cpp/tools/commands/commands.h @@ -49,9 +49,10 @@ std::vector> sorted_table_schemas( std::vector collect_tree_query_paths( const ParsedArgs& args, storage::TsFileReader& reader); -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); int run_row_query(const ParsedArgs& args, storage::TsFileReader& reader, OutputFormat fmt, std::ostream& out, std::ostream& err, diff --git a/cpp/tools/commands/row_query.cc b/cpp/tools/commands/row_query.cc index 828bb6c55..420f31206 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( @@ -284,7 +292,13 @@ int run_row_query(const ParsedArgs& args, storage::TsFileReader& reader, if (is_table_model(args, reader)) { std::string table_name = args.table; if (table_name.empty()) { - auto schemas = reader.get_all_table_schemas(); + std::vector> schemas; + const int schemas_ret = reader.get_all_table_schemas(schemas); + if (schemas_ret != common::E_OK) { + err << "Error: failed to read table schemas: " + << error_code_message(schemas_ret) << "\n"; + return kExitFile; + } if (schemas.empty() || !schemas[0]) { err << "Error: no table found in file\n"; return kExitRuntime; @@ -296,15 +310,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..f63daa554 100644 --- a/go/tsfile/cgo_bridge.go +++ b/go/tsfile/cgo_bridge.go @@ -757,19 +757,24 @@ func copyTableSchemaFromC(schema *C.TableSchema) TableSchema { func (h *readerHandle) tableSchema(table string) (TableSchema, error) { tableName := cStringPtr(table) defer freeCString(tableName) + var code C.ERRNO var native C.TableSchema - code := C.tsfile_reader_get_table_schema_checked(C.TsFileReader(h.ptr), tableName, &native) + code = C.tsfile_reader_get_table_schema_checked( + C.TsFileReader(h.ptr), tableName, &native) if code != C.RET_OK { return TableSchema{}, newError("get table schema", cerrno(code)) } - defer C.free_table_schema(native) - return copyTableSchemaFromC(&native), nil + result := copyTableSchemaFromC(&native) + C.free_table_schema(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 + code = C.tsfile_reader_get_all_table_schemas_checked( + C.TsFileReader(h.ptr), &native, &count) 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 24961f85e..ce3c27612 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,15 @@ 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", + "get_all_timeseries_schemas", + ], +) def test_metadata_read_failure_is_not_an_empty_result(method): class FailingBytesIO(TrackingBytesIO): fail_reads = False @@ -420,8 +429,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 +473,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..a0a424a26 100644 --- a/python/tsfile/tsfile_cpp.pxd +++ b/python/tsfile/tsfile_cpp.pxd @@ -355,15 +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); - - 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); + ErrorCode tsfile_reader_get_table_schema_checked( + TsFileReader reader, const char * table_name, TableSchema * out_schema); + ErrorCode tsfile_reader_get_all_table_schemas_checked( + TsFileReader reader, TableSchema ** out_schemas, uint32_t * out_size); + ErrorCode tsfile_reader_get_all_timeseries_schemas_checked( + TsFileReader reader, DeviceSchema ** out_schemas, uint32_t * out_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 c8f155ae8..4f592d113 100644 --- a/python/tsfile/tsfile_py_cpp.pyx +++ b/python/tsfile/tsfile_py_cpp.pyx @@ -1236,7 +1236,11 @@ 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 = 0 + 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): @@ -1246,7 +1250,8 @@ 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) + error_code = tsfile_reader_get_all_table_schemas_checked( + reader, &schemas, &table_num) check_error(error_code) for i in range(table_num): schema_py = from_c_table_schema(schemas[i]) @@ -1257,10 +1262,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) + error_code = tsfile_reader_get_all_timeseries_schemas_checked( + reader, &schemas, &device_num) + 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)]) 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