diff --git a/src/db/index/column/vector_column/combined_vector_column_indexer.cc b/src/db/index/column/vector_column/combined_vector_column_indexer.cc index f64001e..fd79be3 100644 --- a/src/db/index/column/vector_column/combined_vector_column_indexer.cc +++ b/src/db/index/column/vector_column/combined_vector_column_indexer.cc @@ -116,6 +116,18 @@ Result CombinedVectorColumnIndexer::Search( } } + std::unique_ptr group_by; + if (query_params.group_by) { + auto group_by_func = query_params.group_by->group_by; + auto block_offset = block_offsets_[i]; + group_by = std::make_unique( + query_params.group_by->group_topk, query_params.group_by->group_count, + [group_by_func = std::move(group_by_func), + block_offset](uint64_t block_doc_id) { + return group_by_func(block_doc_id + block_offset); + }); + } + vector_column_params::QueryParams modified_query_params{ query_params.data_type, query_params.dimension, @@ -123,12 +135,7 @@ Result CombinedVectorColumnIndexer::Search( filter, query_params.fetch_vector, query_params.query_params, - query_params.group_by - ? std::make_unique( - query_params.group_by->group_topk, - query_params.group_by->group_count, - query_params.group_by->group_by) - : nullptr, + std::move(group_by), {}, need_refine ? std::shared_ptr( new vector_column_params::RefinerParam{ @@ -270,4 +277,4 @@ CombinedVectorColumnIndexer::Fetch(uint32_t segment_doc_id) const { return indexer->Fetch(target_block_doc_id); } -} // namespace zvec \ No newline at end of file +} // namespace zvec diff --git a/src/db/index/column/vector_column/combined_vector_column_indexer.h b/src/db/index/column/vector_column/combined_vector_column_indexer.h index 9235791..c407b9b 100644 --- a/src/db/index/column/vector_column/combined_vector_column_indexer.h +++ b/src/db/index/column/vector_column/combined_vector_column_indexer.h @@ -38,7 +38,6 @@ class CombinedVectorColumnIndexer { const vector_column_params::VectorData &vector_data, const vector_column_params::QueryParams &query_params); - // doc_id is segment local id virtual Result Fetch( uint32_t segment_doc_id) const; @@ -81,4 +80,4 @@ class CombinedVectorColumnIndexer { uint64_t min_doc_id_{0}; }; -} // namespace zvec \ No newline at end of file +} // namespace zvec diff --git a/src/db/index/segment/segment.cc b/src/db/index/segment/segment.cc index 9bc87b8..42752f5 100644 --- a/src/db/index/segment/segment.cc +++ b/src/db/index/segment/segment.cc @@ -49,7 +49,6 @@ #include "db/index/column/fts_column/fts_rocksdb_merge.h" #include "db/index/column/fts_column/fts_types.h" #include "db/index/column/inverted_column/inverted_indexer.h" -#include "db/index/column/vector_column/engine_helper.hpp" #include "db/index/column/vector_column/vector_column_indexer.h" #include "db/index/column/vector_column/vector_column_params.h" #include "db/index/common/index_filter.h" @@ -61,11 +60,7 @@ #include "db/index/storage/mmap_forward_store.h" #include "db/index/storage/store_helper.h" #include "db/index/storage/wal/wal_file.h" -#include "zvec/ailego/container/params.h" -#include "zvec/core/framework/index_factory.h" -#include "zvec/core/framework/index_meta.h" #include "zvec/core/framework/index_provider.h" -#include "zvec/core/framework/index_reformer.h" #include "column_merging_reader.h" #include "sql_expr_parser.h" @@ -225,10 +220,10 @@ class SegmentImpl : public Segment, Status destroy() override; TablePtr fetch(const std::vector &columns, - const std::vector &indices) const override; + const std::vector &segment_doc_ids) const override; ExecBatchPtr fetch(const std::vector &columns, - int index) const override; + int segment_doc_id) const override; RecordBatchReaderPtr scan( const std::vector &columns) const override; @@ -302,7 +297,7 @@ class SegmentImpl : public Segment, Status append_wal(const Doc &doc); Status update_version(uint32_t delete_snapshot_path_suffix); - Result get_global_doc_id(uint32_t local_id) const; + Result get_global_doc_id(uint32_t segment_doc_id) const; BlockID allocate_block_id(); @@ -323,12 +318,12 @@ class SegmentImpl : public Segment, TablePtr fetch_normal(const std::vector &columns, const std::shared_ptr &result_schema, - const std::vector &indices) const; + const std::vector &segment_doc_ids) const; // For performance tuning TablePtr fetch_perf(const std::vector &columns, const std::shared_ptr &result_schema, - const std::vector &indices) const; + const std::vector &segment_doc_ids) const; void fresh_persist_chunked_array(); @@ -416,7 +411,6 @@ class SegmentImpl : public Segment, class SegmentImpl::CombinedRecordBatchReader : public arrow::RecordBatchReader { public: CombinedRecordBatchReader( - std::shared_ptr segment, std::vector> readers, const std::vector &columns); @@ -427,14 +421,12 @@ class SegmentImpl::CombinedRecordBatchReader : public arrow::RecordBatchReader { arrow::Status ReadNext(std::shared_ptr *batch) override; private: - std::shared_ptr segment_; std::vector> readers_; - std::vector offsets_; std::shared_ptr projected_schema_; - bool need_local_doc_id_ = false; + bool emit_segment_row_id_ = false; size_t current_reader_index_; - size_t local_doc_id_; - int local_doc_id_col_index_ = -1; + uint64_t next_segment_row_id_to_emit_; + int segment_row_id_output_col_index_ = -1; }; //////////////////////////////////////////////////////////////////////////////////// @@ -1433,7 +1425,7 @@ Doc::Ptr SegmentImpl::Fetch( const auto &block_offsets = get_persist_block_offsets(BlockType::VECTOR_INDEX, field->name()); auto block_offset = block_offsets[block_idx]; - auto local_row = segment_doc_id - block_offset; + auto block_doc_id = segment_doc_id - block_offset; auto column_name = field->name(); auto iter = vector_indexers_.find(column_name); @@ -1445,12 +1437,12 @@ Doc::Ptr SegmentImpl::Fetch( continue; } auto vector_indexer = vector_indexers[block_idx]; - auto fetch_result = vector_indexer->Fetch(local_row); + auto fetch_result = vector_indexer->Fetch(block_doc_id); if (!fetch_result) { LOG_ERROR( - "vector indexer fetch failed, local_row: %d, block_idx: %d, " + "vector indexer fetch failed, block_doc_id: %d, block_idx: %d, " "segment_doc_id: %d", - local_row, block_idx, segment_doc_id); + block_doc_id, block_idx, segment_doc_id); return nullptr; } const auto &vector_buffer = fetch_result.value(); @@ -1472,18 +1464,18 @@ Doc::Ptr SegmentImpl::Fetch( p_block_offsets.empty() ? 0 : p_block_offsets.back() + p_block_metas.back().doc_count_; - int local_row = segment_doc_id - mem_block_offset; + int block_doc_id = segment_doc_id - mem_block_offset; auto column_name = field->name(); auto iter = memory_vector_indexers_.find(column_name); if (iter != memory_vector_indexers_.end()) { auto vector_indexer = iter->second; - auto fetch_result = vector_indexer->Fetch(local_row); + auto fetch_result = vector_indexer->Fetch(block_doc_id); if (!fetch_result.has_value()) { LOG_ERROR( "vector indexer fetch failed, column: %s, doc_count: %lu, " - "mem_block_offset: %d, local_row: %d", + "mem_block_offset: %d, block_doc_id: %d", field->name().c_str(), vector_indexer->doc_count(), - mem_block_offset, local_row); + mem_block_offset, block_doc_id); continue; } const auto &vector_buffer = fetch_result.value(); @@ -2339,17 +2331,17 @@ bool SegmentImpl::validate(const std::vector &columns) const { TablePtr SegmentImpl::fetch_perf( const std::vector &columns, const std::shared_ptr &result_schema, - const std::vector &indices) const { + const std::vector &segment_doc_ids) const { std::vector> chunk_arrays; chunk_arrays.resize(columns.size()); - bool need_local_doc_id = false; - size_t local_doc_id_col_index = 0; + bool has_segment_row_id_column = false; + size_t segment_row_id_col_index = 0; for (size_t i = 0; i < columns.size(); ++i) { if (columns[i] == LOCAL_ROW_ID) { - need_local_doc_id = true; - local_doc_id_col_index = i; + has_segment_row_id_column = true; + segment_row_id_col_index = i; chunk_arrays[i] = nullptr; continue; } @@ -2358,18 +2350,19 @@ TablePtr SegmentImpl::fetch_perf( std::vector> result_arrays(columns.size()); - std::vector> indices_in_table; - for (const auto &target_index : indices) { + // Parallel to segment_doc_ids: each pair is (chunk_index, row_index_in_chunk) + std::vector> chunk_row_indices_for_ids; + for (const auto segment_doc_id : segment_doc_ids) { auto it = std::upper_bound(chunk_offsets_.begin(), chunk_offsets_.end(), - target_index); + segment_doc_id); if (it == chunk_offsets_.begin()) { - LOG_ERROR("Target index %d is out of bounds", target_index); + LOG_ERROR("Segment doc ID %d is out of bounds", segment_doc_id); return nullptr; } int chunk_index = static_cast(std::distance(chunk_offsets_.begin(), it) - 1); - int64_t index_in_chunk = target_index - chunk_offsets_[chunk_index]; - indices_in_table.emplace_back(chunk_index, index_in_chunk); + int64_t row_index_in_chunk = segment_doc_id - chunk_offsets_[chunk_index]; + chunk_row_indices_for_ids.emplace_back(chunk_index, row_index_in_chunk); } for (size_t i = 0; i < columns.size(); ++i) { @@ -2378,8 +2371,8 @@ TablePtr SegmentImpl::fetch_perf( } const auto &source_column = chunk_arrays[i]; std::shared_ptr array; - auto status = - BuildArrayFromIndicesWithType(source_column, indices_in_table, &array); + auto status = BuildArrayFromIndicesWithType( + source_column, chunk_row_indices_for_ids, &array); if (!status.ok()) { LOG_ERROR("BuildArrayFromIndices failed: %s", status.ToString().c_str()); return nullptr; @@ -2387,11 +2380,11 @@ TablePtr SegmentImpl::fetch_perf( result_arrays[i] = array; } - if (need_local_doc_id) { + if (has_segment_row_id_column) { std::vector values; - values.reserve(indices.size()); - for (const auto idx : indices) { - values.push_back(idx); + values.reserve(segment_doc_ids.size()); + for (const auto segment_doc_id : segment_doc_ids) { + values.push_back(segment_doc_id); } arrow::UInt64Builder builder; @@ -2406,26 +2399,26 @@ TablePtr SegmentImpl::fetch_perf( LOG_ERROR("Failed to finish builder: %s", s.message().c_str()); return nullptr; } - result_arrays[local_doc_id_col_index] = array; + result_arrays[segment_row_id_col_index] = array; } return arrow::Table::Make(result_schema, result_arrays, - static_cast(indices.size())); + static_cast(segment_doc_ids.size())); } TablePtr SegmentImpl::fetch_normal( const std::vector &columns, const std::shared_ptr &result_schema, - const std::vector &indices) const { + const std::vector &segment_doc_ids) const { // Store scalars per column: column_index -> (output_row, scalar) std::vector>>> column_results(columns.size()); - // Collect local_doc_id values if needed - std::vector> local_doc_id_values; + // Collect segment-local row IDs when LOCAL_ROW_ID is requested. + std::vector> segment_row_id_values; - // Group fetch requests by block: block_index -> {column -> [(output_row, - // local_row)]} + // Group fetch requests by block: + // block_index -> {column -> [(output_row, block_row)]} // block_index >= 0: persisted store // block_index == -1: memory store std::map>>> @@ -2436,26 +2429,27 @@ TablePtr SegmentImpl::fetch_normal( const auto &block_offsets = get_persist_block_offsets(BlockType::SCALAR); const auto &block_metas = get_persist_block_metas(BlockType::SCALAR); - // Phase 1: Map each (doc_id, column) to its block and local row - for (int output_row = 0; output_row < static_cast(indices.size()); - ++output_row) { - int doc_id = indices[output_row]; + // Phase 1: Map each (segment_doc_id, column) to its block and block-local + // row. + for (int output_row = 0; + output_row < static_cast(segment_doc_ids.size()); ++output_row) { + int segment_doc_id = segment_doc_ids[output_row]; for (size_t col_index = 0; col_index < columns.size(); ++col_index) { const std::string &col = columns[col_index]; if (col == LOCAL_ROW_ID) { - local_doc_id_values.emplace_back(output_row, doc_id); + segment_row_id_values.emplace_back(output_row, segment_doc_id); continue; } int offset_idx = -1; - int block_index = - find_persist_block_id(BlockType::SCALAR, doc_id, col, &offset_idx); + int block_index = find_persist_block_id(BlockType::SCALAR, segment_doc_id, + col, &offset_idx); - int local_row = -1; + int block_row = -1; if (block_index != -1 && offset_idx > -1 && offset_idx < static_cast(block_offsets.size())) { - local_row = doc_id - block_offsets[offset_idx]; - block_request_map[block_index][col].emplace_back(output_row, local_row); + block_row = segment_doc_id - block_offsets[offset_idx]; + block_request_map[block_index][col].emplace_back(output_row, block_row); continue; } @@ -2467,15 +2461,17 @@ TablePtr SegmentImpl::fetch_normal( : block_offsets.back() + block_metas.back().doc_count_; const auto &mem_block = segment_meta_->writing_forward_block().value(); - if (mem_offset <= doc_id && - doc_id < mem_offset + static_cast(mem_block.doc_count_)) { - local_row = doc_id - mem_offset; - block_request_map[-1][col].emplace_back(output_row, local_row); + if (mem_offset <= segment_doc_id && + segment_doc_id < + mem_offset + static_cast(mem_block.doc_count_)) { + block_row = segment_doc_id - mem_offset; + block_request_map[-1][col].emplace_back(output_row, block_row); continue; } } - LOG_ERROR("Document ID %d not found in segment %d", doc_id, meta()->id()); + LOG_ERROR("Segment doc ID %d not found in segment %d", segment_doc_id, + meta()->id()); return nullptr; } } @@ -2483,7 +2479,7 @@ TablePtr SegmentImpl::fetch_normal( // Phase 2: Execute batched fetch per block for (const auto &[block_index, col_to_rows] : block_request_map) { std::vector fetch_columns; - std::vector fetch_local_rows; + std::vector fetch_block_rows; std::vector> output_to_result_index; // (output_row, result_pos) @@ -2493,20 +2489,20 @@ TablePtr SegmentImpl::fetch_normal( } // all column has same output size, here just take first column - for (const auto &[output_row, local_row] : + for (const auto &[output_row, block_row] : col_to_rows.at(fetch_columns[0])) { - fetch_local_rows.push_back(local_row); + fetch_block_rows.push_back(block_row); output_to_result_index.emplace_back( - output_row, static_cast(fetch_local_rows.size() - 1)); + output_row, static_cast(fetch_block_rows.size() - 1)); } std::shared_ptr block_table; if (block_index >= 0 && block_index < static_cast(persist_stores_.size())) { block_table = - persist_stores_[block_index]->fetch(fetch_columns, fetch_local_rows); + persist_stores_[block_index]->fetch(fetch_columns, fetch_block_rows); } else if (block_index == -1 && memory_store_) { - block_table = memory_store_->fetch(fetch_columns, fetch_local_rows); + block_table = memory_store_->fetch(fetch_columns, fetch_block_rows); } if (!block_table || block_table->num_rows() == 0) { @@ -2530,7 +2526,7 @@ TablePtr SegmentImpl::fetch_normal( } auto flat_array = flat_array_res.ValueOrDie(); - for (size_t j = 0; j < fetch_local_rows.size(); ++j) { + for (size_t j = 0; j < fetch_block_rows.size(); ++j) { auto scalar_result = flat_array->GetScalar(j); if (!scalar_result.ok()) continue; int output_row = output_to_result_index[j].first; @@ -2543,14 +2539,14 @@ TablePtr SegmentImpl::fetch_normal( // Phase 3: Construct result arrays std::vector> result_arrays(columns.size()); - bool need_local_doc_id = false; - size_t local_doc_id_col_index = -1; + bool has_segment_row_id_column = false; + size_t segment_row_id_col_index = -1; for (size_t col_index = 0; col_index < columns.size(); ++col_index) { const std::string &col = columns[col_index]; if (col == LOCAL_ROW_ID) { - need_local_doc_id = true; - local_doc_id_col_index = col_index; + has_segment_row_id_column = true; + segment_row_id_col_index = col_index; continue; } @@ -2558,7 +2554,7 @@ TablePtr SegmentImpl::fetch_normal( std::sort(result_vec.begin(), result_vec.end()); std::vector> ordered_scalars; - for (int i = 0; i < static_cast(indices.size()); ++i) { + for (int i = 0; i < static_cast(segment_doc_ids.size()); ++i) { auto it = std::find_if( result_vec.begin(), result_vec.end(), [i](const std::pair> &p) { @@ -2582,13 +2578,13 @@ TablePtr SegmentImpl::fetch_normal( } } - // Add LOCAL_ROW_ID array if requested - if (need_local_doc_id) { - std::sort(local_doc_id_values.begin(), local_doc_id_values.end()); + // Add segment-local values for the LOCAL_ROW_ID column. + if (has_segment_row_id_column) { + std::sort(segment_row_id_values.begin(), segment_row_id_values.end()); std::vector values; - values.reserve(local_doc_id_values.size()); - for (const auto &[row, id] : local_doc_id_values) { - values.push_back(id); + values.reserve(segment_row_id_values.size()); + for (const auto &[row, segment_row_id] : segment_row_id_values) { + values.push_back(segment_row_id); } arrow::UInt64Builder builder; @@ -2603,7 +2599,7 @@ TablePtr SegmentImpl::fetch_normal( LOG_ERROR("Failed to finish builder: %s", s.message().c_str()); return nullptr; } - result_arrays[local_doc_id_col_index] = std::move(array); + result_arrays[segment_row_id_col_index] = std::move(array); } // Wrap arrays into ChunkedArray and build final table @@ -2614,11 +2610,11 @@ TablePtr SegmentImpl::fetch_normal( } return arrow::Table::Make(result_schema, result_columns, - static_cast(indices.size())); + static_cast(segment_doc_ids.size())); } TablePtr SegmentImpl::fetch(const std::vector &columns, - const std::vector &indices) const { + const std::vector &segment_doc_ids) const { if (!validate(columns)) { return nullptr; } @@ -2649,8 +2645,8 @@ TablePtr SegmentImpl::fetch(const std::vector &columns, auto result_schema = std::make_shared(fields); - // Early return for empty indices - if (indices.empty()) { + // Early return for empty segment doc IDs. + if (segment_doc_ids.empty()) { arrow::ArrayVector empty_arrays; for (const auto &field : fields) { empty_arrays.push_back(arrow::MakeEmptyArray(field->type()).ValueOrDie()); @@ -2664,13 +2660,13 @@ TablePtr SegmentImpl::fetch(const std::vector &columns, } if (use_fetch_perf_) { - return fetch_perf(columns, result_schema, indices); + return fetch_perf(columns, result_schema, segment_doc_ids); } - return fetch_normal(columns, result_schema, indices); + return fetch_normal(columns, result_schema, segment_doc_ids); } ExecBatchPtr SegmentImpl::fetch(const std::vector &columns, - int doc_id) const { + int segment_doc_id) const { if (columns.empty()) { LOG_ERROR("Empty columns"); return nullptr; @@ -2709,12 +2705,12 @@ ExecBatchPtr SegmentImpl::fetch(const std::vector &columns, if (is_in_single_persist_store) { int offset_idx = -1; - int block_index = find_persist_block_id(BlockType::SCALAR, doc_id, + int block_index = find_persist_block_id(BlockType::SCALAR, segment_doc_id, columns[0], &offset_idx); if (block_index != -1 && offset_idx > -1 && offset_idx < static_cast(block_offsets.size())) { - int local_row = doc_id - block_offsets[offset_idx]; - return persist_stores_[block_index]->fetch(columns, local_row); + int block_row = segment_doc_id - block_offsets[offset_idx]; + return persist_stores_[block_index]->fetch(columns, block_row); } // Check memory store @@ -2725,14 +2721,15 @@ ExecBatchPtr SegmentImpl::fetch(const std::vector &columns, : block_offsets.back() + block_metas.back().doc_count_; const auto &mem_block = segment_meta_->writing_forward_block().value(); - if (mem_offset <= doc_id && - doc_id < mem_offset + static_cast(mem_block.doc_count_)) { - int local_row = doc_id - mem_offset; - return memory_store_->fetch(columns, local_row); + if (mem_offset <= segment_doc_id && + segment_doc_id < + mem_offset + static_cast(mem_block.doc_count_)) { + int block_row = segment_doc_id - mem_offset; + return memory_store_->fetch(columns, block_row); } } } else { - auto table = fetch(columns, std::vector{doc_id}); + auto table = fetch(columns, std::vector{segment_doc_id}); if (table) { std::vector datums; for (const auto &col : table->columns()) { @@ -2749,7 +2746,7 @@ ExecBatchPtr SegmentImpl::fetch(const std::vector &columns, } } - LOG_ERROR("Document ID %d not found in persist segment", doc_id); + LOG_ERROR("Segment doc ID %d not found in persist segment", segment_doc_id); return nullptr; } @@ -2767,6 +2764,8 @@ RecordBatchReaderPtr SegmentImpl::scan( std::map, std::vector>> block_groups; + bool emit_segment_row_id = + std::find(columns.begin(), columns.end(), LOCAL_ROW_ID) != columns.end(); for (size_t i = 0; i < scalar_blocks.size() && i < persist_stores_.size(); ++i) { @@ -2779,6 +2778,9 @@ RecordBatchReaderPtr SegmentImpl::scan( interested_cols.push_back(col); } } + if (interested_cols.empty() && emit_segment_row_id) { + interested_cols.push_back(GLOBAL_DOC_ID); + } if (interested_cols.empty()) { continue; @@ -2794,7 +2796,16 @@ RecordBatchReaderPtr SegmentImpl::scan( } if (memory_store_ && memory_store_->num_rows() > 0) { - auto reader = memory_store_->scan(columns); + std::vector memory_scan_columns; + for (const auto &col : columns) { + if (col != LOCAL_ROW_ID) { + memory_scan_columns.push_back(col); + } + } + if (memory_scan_columns.empty()) { + memory_scan_columns.push_back(GLOBAL_DOC_ID); + } + auto reader = memory_store_->scan(memory_scan_columns); if (reader) { auto &mem_block = segment_meta_->writing_forward_block().value(); auto key = std::make_pair(mem_block.min_doc_id(), mem_block.max_doc_id()); @@ -2834,8 +2845,8 @@ RecordBatchReaderPtr SegmentImpl::scan( } } - return std::make_shared( - shared_from_this(), std::move(merged_readers), columns); + return std::make_shared(std::move(merged_readers), + columns); } @@ -2844,13 +2855,11 @@ RecordBatchReaderPtr SegmentImpl::scan( //////////////////////////////////////////////////////////////////////////////////// SegmentImpl::CombinedRecordBatchReader::CombinedRecordBatchReader( - std::shared_ptr segment, std::vector> readers, const std::vector &columns) - : segment_(segment), - readers_(std::move(readers)), + : readers_(std::move(readers)), current_reader_index_(0), - local_doc_id_(0) { + next_segment_row_id_to_emit_(0) { if (!readers_.empty()) { auto schema = readers_[0]->schema(); std::vector> selected_fields; @@ -2859,8 +2868,8 @@ SegmentImpl::CombinedRecordBatchReader::CombinedRecordBatchReader( if (col_name == LOCAL_ROW_ID) { selected_fields.push_back( arrow::field(LOCAL_ROW_ID, arrow::uint64(), false)); - need_local_doc_id_ = true; - local_doc_id_col_index_ = static_cast(i); + emit_segment_row_id_ = true; + segment_row_id_output_col_index_ = static_cast(i); } else { if (auto field = schema->GetFieldByName(col_name); field) { selected_fields.push_back(field); @@ -2869,17 +2878,6 @@ SegmentImpl::CombinedRecordBatchReader::CombinedRecordBatchReader( } projected_schema_ = arrow::schema(selected_fields); - - auto segment_meta = segment_->meta(); - const auto &blocks = segment_meta->persisted_blocks(); - for (const auto &block : blocks) { - if (block.type() != BlockType::SCALAR) continue; - offsets_.push_back(block.min_doc_id_); - } - if (segment_meta->has_writing_forward_block()) { - const auto &mem_block = segment_meta->writing_forward_block().value(); - offsets_.push_back(mem_block.min_doc_id_); - } } } @@ -2899,24 +2897,25 @@ arrow::Status SegmentImpl::CombinedRecordBatchReader::ReadNext( return status; } - if (need_local_doc_id_ && *batch) { + if (emit_segment_row_id_ && *batch) { auto num_rows = (*batch)->num_rows(); arrow::UInt64Builder builder; ARROW_RETURN_NOT_OK(builder.Reserve(num_rows)); for (int64_t i = 0; i < num_rows; ++i) { - builder.UnsafeAppend(local_doc_id_++); + builder.UnsafeAppend(next_segment_row_id_to_emit_++); } - std::shared_ptr local_id_array; - ARROW_RETURN_NOT_OK(builder.Finish(&local_id_array)); + std::shared_ptr segment_row_id_array; + ARROW_RETURN_NOT_OK(builder.Finish(&segment_row_id_array)); auto result = - (*batch)->AddColumn(local_doc_id_col_index_, + (*batch)->AddColumn(segment_row_id_output_col_index_, projected_schema_->GetFieldByName(LOCAL_ROW_ID), - std::move(local_id_array)); - if (result.ok()) { - *batch = std::move(result.ValueOrDie()); + std::move(segment_row_id_array)); + if (!result.ok()) { + return result.status(); } + *batch = std::move(result.ValueOrDie()); } if (*batch) { @@ -2924,9 +2923,6 @@ arrow::Status SegmentImpl::CombinedRecordBatchReader::ReadNext( } current_reader_index_++; - if (current_reader_index_ < readers_.size()) { - local_doc_id_ = offsets_[current_reader_index_]; - } } *batch = nullptr; @@ -4404,14 +4400,13 @@ BlockID SegmentImpl::allocate_block_id() { return block_id_allocator_.fetch_add(1); } -Result SegmentImpl::get_global_doc_id(uint32_t local_id) const { +Result SegmentImpl::get_global_doc_id(uint32_t segment_doc_id) const { std::lock_guard lock(seg_mtx_); - if (local_id >= doc_ids_.size()) { + if (segment_doc_id >= doc_ids_.size()) { return tl::make_unexpected( - Status::InvalidArgument("local_id out of range")); + Status::InvalidArgument("segment_doc_id out of range")); } - // global doc_id - return doc_ids_[local_id]; + return doc_ids_[segment_doc_id]; } @@ -4679,4 +4674,4 @@ Result> SegmentImpl::fts_search( return std::move(ret.value()); } -} // namespace zvec \ No newline at end of file +} // namespace zvec diff --git a/src/db/index/segment/segment.h b/src/db/index/segment/segment.h index 3b21c64..e31f072 100644 --- a/src/db/index/segment/segment.h +++ b/src/db/index/segment/segment.h @@ -64,9 +64,10 @@ class Segment { virtual SegmentMeta::Ptr meta() const = 0; + // Count documents visible to an optional global-doc-ID filter. virtual uint64_t doc_count(const IndexFilter::Ptr filter = nullptr) = 0; - // for collection + // ---- Schema and index mutation ----------------------------------------- virtual Status add_column(FieldSchema::Ptr column_schema, const std::string &expression, const AddColumnOptions &options) = 0; @@ -84,7 +85,6 @@ class Segment { std::unordered_map *quant_vector_indexers) = 0; - // defined in segment.h cause it needs to access block_id generator virtual Status create_vector_index( const std::string &column, const IndexParams::Ptr &index_params, int concurrency, SegmentMeta::Ptr *new_segment_meta, @@ -111,13 +111,11 @@ class Segment { virtual bool all_vector_index_ready() const = 0; - // defined in segment.h cause it needs to access block_id generator virtual Status create_scalar_index( const std::vector &columns, const IndexParams::Ptr &index_params, SegmentMeta::Ptr *new_segment_meta, InvertedIndexer::Ptr *new_scalar_indexer) = 0; - // defined in segment.h cause it needs to access block_id generator virtual Status drop_scalar_index( const std::vector &columns, SegmentMeta::Ptr *new_segment_meta, @@ -127,6 +125,7 @@ class Segment { const CollectionSchema &schema, const SegmentMeta::Ptr &segment_meta, const InvertedIndexer::Ptr &scalar_indexer) = 0; + // ---- Data operations ---------------------------------------------------- virtual Status Insert(Doc &doc) = 0; virtual Status Upsert(Doc &doc) = 0; @@ -142,53 +141,51 @@ class Segment { &output_fields = std::nullopt, bool include_vector = true) = 0; - // for sqlengine virtual TablePtr fetch(const std::vector &columns, - const std::vector &indices) const = 0; + const std::vector &segment_doc_ids) const = 0; virtual ExecBatchPtr fetch(const std::vector &columns, - int index) const = 0; + int segment_doc_id) const = 0; - // caller should hold segment shared_ptr for segment handle the indexer's - // lifetime + // Keep Segment alive while consuming the returned reader. virtual RecordBatchReaderPtr scan( const std::vector &columns) const = 0; - // caller hold segment shared_ptr for segment handle the indexer's lifetime + // ---- Index accessors ---------------------------------------------------- + // Keep Segment alive while using returned indexers. virtual CombinedVectorColumnIndexer::Ptr get_combined_vector_indexer( const std::string &field_name) const = 0; - // caller hold segment shared_ptr for segment handle the indexer's lifetime virtual CombinedVectorColumnIndexer::Ptr get_quant_combined_vector_indexer( const std::string &field_name) const = 0; - // caller hold segment shared_ptr for segment handle the indexer's lifetime virtual std::vector get_vector_indexer( const std::string &field_name) const = 0; virtual std::vector get_quant_vector_indexer( const std::string &field_name) const = 0; - // caller hold segment shared_ptr for segment handle the indexer's lifetime virtual InvertedColumnIndexer::Ptr get_scalar_indexer( const std::string &field_name) const = 0; - // caller hold segment shared_ptr for segment handle the indexer's lifetime virtual fts::FtsColumnIndexerPtr get_fts_indexer( const std::string &field_name) const = 0; + // ---- Index queries and filters ----------------------------------------- virtual Result> fts_search( const std::string &field_name, const fts::FtsAstNode &ast, const fts::FtsQueryParams ¶ms) = 0; + // Returned filter is evaluated with segment-local row IDs. It translates the + // local row ID to a global doc ID before consulting the delete store. virtual const IndexFilter::Ptr get_filter() = 0; - // for others + // ---- Persistence and lifecycle ----------------------------------------- virtual Status flush() = 0; + virtual Status dump() = 0; - // only mark need_destroyed virtual Status destroy() = 0; }; -} // namespace zvec \ No newline at end of file +} // namespace zvec diff --git a/tests/db/index/segment/segment_row_id_test.cc b/tests/db/index/segment/segment_row_id_test.cc new file mode 100644 index 0000000..c64ccd5 --- /dev/null +++ b/tests/db/index/segment/segment_row_id_test.cc @@ -0,0 +1,234 @@ +// Copyright 2025-present the zvec project +// +// Licensed 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 "db/common/constants.h" +#include "segment_test_fixture.h" + +using namespace zvec; + + +namespace { + + +struct LocalRowIdProjection { + std::vector columns; + int row_id_column_index; +}; + +const std::vector &LocalRowIdProjections() { + static const std::vector projections = { + {{LOCAL_ROW_ID, "id", "name"}, 0}, + {{"id", LOCAL_ROW_ID, "name"}, 1}, + {{"id", "name", LOCAL_ROW_ID}, 2}, + }; + return projections; +} + +void ExpectSingleRowLocalRowID(const ExecBatchPtr &batch, int column_index, + uint64_t expected) { + ASSERT_TRUE(batch != nullptr); + EXPECT_EQ(batch->length, 1); + EXPECT_EQ(batch->values.size(), 3); + + auto local_row_id_scalar = batch->values[column_index].scalar(); + ASSERT_TRUE(local_row_id_scalar != nullptr); + auto local_row_id_value = + std::dynamic_pointer_cast(local_row_id_scalar); + ASSERT_TRUE(local_row_id_value != nullptr); + EXPECT_EQ(local_row_id_value->value, expected); +} + +void ExpectScanLocalRowIDColumn(const RecordBatchReaderPtr &reader, + int column_index, uint32_t expected_rows) { + ASSERT_TRUE(reader != nullptr); + ASSERT_TRUE(reader->schema() != nullptr); + + std::shared_ptr batch; + uint32_t total_doc = 0; + while (true) { + auto status = reader->ReadNext(&batch); + ASSERT_TRUE(status.ok()) << status.ToString(); + if (batch == nullptr) break; + + ASSERT_GT(batch->num_columns(), column_index); + EXPECT_EQ(batch->column(column_index)->type()->id(), arrow::Type::UINT64); + EXPECT_EQ(batch->column_name(column_index), LOCAL_ROW_ID); + + total_doc += batch->num_rows(); + } + EXPECT_EQ(total_doc, expected_rows); +} + +void ExpectFetchedLocalRowIDs( + const TablePtr &table, int column_index, + const std::vector &expected_segment_doc_ids) { + ASSERT_TRUE(table != nullptr); + EXPECT_EQ(table->num_columns(), 3); + EXPECT_EQ(table->num_rows(), + static_cast(expected_segment_doc_ids.size())); + + auto field = table->schema()->field(column_index); + EXPECT_EQ(field->name(), LOCAL_ROW_ID); + + auto id_column = table->column(column_index); + auto id_array = + std::dynamic_pointer_cast(id_column->chunk(0)); + ASSERT_TRUE(id_array != nullptr); + + std::vector actual_ids; + actual_ids.reserve(id_array->length()); + for (int64_t i = 0; i < id_array->length(); ++i) { + actual_ids.push_back(id_array->Value(i)); + } + + std::vector expected_u64_ids(expected_segment_doc_ids.begin(), + expected_segment_doc_ids.end()); + EXPECT_EQ(actual_ids, expected_u64_ids) + << "LOCAL_ROW_ID values don't match expected order"; +} + + +} // namespace + + +TEST_P(SegmentTest, FetchSingleRowWithLocalRowIDInRequestedPosition) { + auto segment = test::TestHelper::CreateSegmentWithDoc( + col_path_, *schema_, 0, 0, id_map_, delete_store_, version_manager_, + options_, 0, 10); + ASSERT_TRUE(segment != nullptr); + + for (const auto &projection : LocalRowIdProjections()) { + ExecBatchPtr batch = segment->fetch(projection.columns, 4); + ExpectSingleRowLocalRowID(batch, projection.row_id_column_index, 4); + } +} + + +TEST_P(SegmentTest, ScanAndFetchPreserveLocalRowIDColumnPosition) { + auto segment = test::TestHelper::CreateSegmentWithDoc( + col_path_, *schema_, 0, 0, id_map_, delete_store_, version_manager_, + options_, 0, 10); + ASSERT_TRUE(segment != nullptr); + + std::vector segment_doc_ids = {0, 3, 6, 1, 0}; + for (const auto &projection : LocalRowIdProjections()) { + auto reader = segment->scan(projection.columns); + ExpectScanLocalRowIDColumn(reader, projection.row_id_column_index, 10); + + auto table = segment->fetch(projection.columns, segment_doc_ids); + ExpectFetchedLocalRowIDs(table, projection.row_id_column_index, + segment_doc_ids); + } +} + + +TEST_P(SegmentTest, ScanLocalRowIDIsSegmentLocal) { + options_.max_buffer_size_ = 1 * 1024; + + auto segment = test::TestHelper::CreateSegmentWithDoc( + col_path_, *schema_, 0, 100, id_map_, delete_store_, version_manager_, + options_, 100, 25); + ASSERT_TRUE(segment != nullptr); + + auto reader = segment->scan({LOCAL_ROW_ID, "id"}); + ASSERT_TRUE(reader != nullptr); + + std::shared_ptr batch; + std::vector actual_ids; + while (true) { + auto status = reader->ReadNext(&batch); + ASSERT_TRUE(status.ok()) << status.ToString(); + if (batch == nullptr) break; + + ASSERT_EQ(batch->num_columns(), 2); + auto id_array = + std::dynamic_pointer_cast(batch->column(0)); + ASSERT_TRUE(id_array != nullptr); + for (int64_t i = 0; i < id_array->length(); ++i) { + actual_ids.push_back(id_array->Value(i)); + } + } + + std::vector expected_ids; + for (uint64_t i = 0; i < 25; ++i) { + expected_ids.push_back(i); + } + EXPECT_EQ(actual_ids, expected_ids); +} + + +TEST_P(SegmentTest, ScanOnlyLocalRowIDDoesNotExposeGlobalDocID) { + auto segment = test::TestHelper::CreateSegmentWithDoc( + col_path_, *schema_, 0, 100, id_map_, delete_store_, version_manager_, + options_, 100, 10); + ASSERT_TRUE(segment != nullptr); + + auto reader = segment->scan({LOCAL_ROW_ID}); + ASSERT_TRUE(reader != nullptr); + ASSERT_TRUE(reader->schema() != nullptr); + ASSERT_EQ(reader->schema()->num_fields(), 1); + EXPECT_EQ(reader->schema()->field(0)->name(), LOCAL_ROW_ID); + + std::shared_ptr batch; + std::vector actual_ids; + while (true) { + auto status = reader->ReadNext(&batch); + ASSERT_TRUE(status.ok()) << status.ToString(); + if (batch == nullptr) break; + + ASSERT_EQ(batch->num_columns(), 1); + EXPECT_EQ(batch->column_name(0), LOCAL_ROW_ID); + auto id_array = + std::dynamic_pointer_cast(batch->column(0)); + ASSERT_TRUE(id_array != nullptr); + for (int64_t i = 0; i < id_array->length(); ++i) { + actual_ids.push_back(id_array->Value(i)); + } + } + + std::vector expected_ids; + for (uint64_t i = 0; i < 10; ++i) { + expected_ids.push_back(i); + } + EXPECT_EQ(actual_ids, expected_ids); +} + + +TEST_P(SegmentTest, DocCountDeleteFilterWithNonZeroGlobalDocID) { + auto segment = test::TestHelper::CreateSegmentWithDoc( + col_path_, *schema_, 0, 100, id_map_, delete_store_, version_manager_, + options_, 0, 10); + ASSERT_TRUE(segment != nullptr); + + auto status = segment->Delete("pk_5"); + EXPECT_TRUE(status.ok()) << "Delete by pk failed: " << status.message(); + + status = segment->Delete(103); + EXPECT_TRUE(status.ok()) << "Delete by global doc id failed: " + << status.message(); + + EXPECT_EQ(segment->doc_count(), 10); + EXPECT_EQ(segment->doc_count(delete_store_->make_filter()), 8); +} + + +INSTANTIATE_TEST_SUITE_P(MMapTest, SegmentTest, testing::Values(true, false)); diff --git a/tests/db/index/segment/segment_test.cc b/tests/db/index/segment/segment_test.cc index b3cd2f9..9582d2b 100644 --- a/tests/db/index/segment/segment_test.cc +++ b/tests/db/index/segment/segment_test.cc @@ -12,8 +12,8 @@ // See the License for the specific language governing permissions and // limitations under the License. -#include #include +#include // NOLINT(misc-include-cleaner): protects from the private/public test macro below. #define private public #define protected public #include "db/index/segment/segment.h" @@ -34,94 +34,16 @@ #include "db/index/common/delete_store.h" #include "db/index/common/id_map.h" #include "db/index/common/version_manager.h" -#include "db/index/storage/store_helper.h" #include "db/index/storage/wal/wal_file.h" +#include "segment_test_fixture.h" #include "utils/utils.h" #include "zvec/db/options.h" using namespace zvec; -class SegmentTest : public testing::TestWithParam { - protected: - void SetUp() override { - ailego::LoggerBroker::SetLevel(ailego::Logger::LEVEL_INFO); - - FileHelper::RemoveDirectory(col_path); - FileHelper::CreateDirectory(col_path); - - zvec::ailego::MemoryLimitPool::get_instance().init(MIN_MEMORY_LIMIT_BYTES); - - std::string idmap_path = - FileHelper::MakeFilePath(col_path, FileID::ID_FILE, 0); - id_map = IDMap::CreateAndOpen(col_name, idmap_path, true, false); - if (id_map == nullptr) { - throw std::runtime_error("Failed to create id map"); - } - - std::string delete_store_path = - FileHelper::MakeFilePath(col_path, FileID::DELETE_FILE, 0); - delete_store = std::make_shared(col_name); - - schema = - test::TestHelper::CreateSchemaWithScalarIndex(false, false, col_name); - - schema->add_field( - std::make_shared("id", DataType::INT32, false)); - schema->add_field( - std::make_shared("name", DataType::STRING, false)); - schema->add_field( - std::make_shared("age", DataType::UINT32, false)); - - schema->add_field( - std::make_shared("binary", DataType::BINARY, false)); - - schema->add_field(std::make_shared( - "array_binary", DataType::ARRAY_BINARY, false)); - - bool enable_mmap = GetParam(); - - Version version; - version.set_schema(*schema); - version.set_enable_mmap(enable_mmap); - auto version_manager_tmp = VersionManager::Create(col_path, version); - if (!version_manager_tmp.has_value()) { - throw std::runtime_error("Failed to create version manager"); - } - - version_manager = version_manager_tmp.value(); - - // default options - options.read_only_ = false; - options.enable_mmap_ = enable_mmap; - options.max_buffer_size_ = 64 * 1024 * 1024; - } - - void TearDown() override { - id_map.reset(); - delete_store.reset(); - version_manager.reset(); - - // FileHelper::RemoveDirectory(col_path); - } - - public: - std::string GetColPath() { - return col_path; - } - - protected: - std::string col_name = "test_segment"; - std::string col_path = "./test_collection"; - IDMap::Ptr id_map; - DeleteStore::Ptr delete_store; - VersionManager::Ptr version_manager; - CollectionSchema::Ptr schema; - SegmentOptions options; -}; - TEST_P(SegmentTest, EmptySchema) { auto segment = test::TestHelper::CreateSegmentWithDoc( - col_path, *schema, 0, 0, id_map, delete_store, version_manager, options, + col_path_, *schema_, 0, 0, id_map_, delete_store_, version_manager_, options_, 0, 0); ASSERT_TRUE(segment != nullptr); EXPECT_EQ(segment->id(), 0); @@ -131,10 +53,10 @@ TEST_P(SegmentTest, EmptySchema) { TEST_P(SegmentTest, General) { - options.max_buffer_size_ = 1 * 1024; + options_.max_buffer_size_ = 1 * 1024; auto segment = test::TestHelper::CreateSegmentWithDoc( - col_path, *schema, 0, 0, id_map, delete_store, version_manager, options, + col_path_, *schema_, 0, 0, id_map_, delete_store_, version_manager_, options_, 0, 25); ASSERT_TRUE(segment != nullptr); @@ -155,9 +77,10 @@ TEST_P(SegmentTest, General) { } EXPECT_EQ(total_doc, 25); - std::vector indices = {0, 3, 6, 1, 0, 14, 12, 21}; + std::vector segment_doc_ids = {0, 3, 6, 1, 0, 14, 12, 21}; auto combined_table = segment->fetch( - {LOCAL_ROW_ID, "id", "name", "age", "binary", "array_binary"}, indices); + {LOCAL_ROW_ID, "id", "name", "age", "binary", "array_binary"}, + segment_doc_ids); ASSERT_TRUE(combined_table != nullptr); EXPECT_EQ(combined_table->num_columns(), 6); EXPECT_EQ(combined_table->num_rows(), 8); @@ -170,7 +93,7 @@ TEST_P(SegmentTest, General) { auto id_array = std::dynamic_pointer_cast(id_column->chunk(0)); - std::vector &expected_ids = indices; + std::vector &expected_ids = segment_doc_ids; std::vector actual_ids; for (int i = 0; i < id_array->length(); ++i) { @@ -183,13 +106,13 @@ TEST_P(SegmentTest, General) { TEST_P(SegmentTest, InsertMoreData) { auto segment = test::TestHelper::CreateSegmentWithDoc( - col_path, *schema, 0, 0, id_map, delete_store, version_manager, options, + col_path_, *schema_, 0, 0, id_map_, delete_store_, version_manager_, options_, 0, 0); ASSERT_TRUE(segment != nullptr); uint64_t MAX_DOC = 1000; auto start = std::chrono::system_clock::now(); - test::TestHelper::SegmentInsertDoc(segment, *schema, 0, MAX_DOC); + test::TestHelper::SegmentInsertDoc(segment, *schema_, 0, MAX_DOC); auto end = std::chrono::system_clock::now(); auto cost = std::chrono::duration_cast(end - start) .count(); @@ -210,24 +133,24 @@ TEST_P(SegmentTest, InsertMoreData) { TEST_P(SegmentTest, InsertScalarTypes) { auto tmp_schema = - test::TestHelper::CreateSchemaWithScalarIndex(true, true, col_name); + test::TestHelper::CreateSchemaWithScalarIndex(true, true, col_name_); auto invert_params = std::make_shared(false); - schema->add_field(std::make_shared("binary", DataType::BINARY, + schema_->add_field(std::make_shared("binary", DataType::BINARY, false, invert_params)); - schema->add_field(std::make_shared( + schema_->add_field(std::make_shared( "array_binary", DataType::ARRAY_BINARY, false, invert_params)); auto segment = test::TestHelper::CreateSegmentWithDoc( - col_path, *schema, 0, 0, id_map, delete_store, version_manager, options, + col_path_, *schema_, 0, 0, id_map_, delete_store_, version_manager_, options_, 0, 10); ASSERT_TRUE(segment != nullptr); } TEST_P(SegmentTest, InsertVectorTypes) { auto tmp_schema = test::TestHelper::CreateSchemaWithVectorIndex( - false, col_name, + false, col_name_, std::make_shared(MetricType::IP, 16, 20, QuantizeType::FP16)); @@ -235,17 +158,17 @@ TEST_P(SegmentTest, InsertVectorTypes) { int doc_count = 100; { auto segment = test::TestHelper::CreateSegmentWithDoc( - col_path, *tmp_schema, 0, 0, id_map, delete_store, version_manager, - options, 0, doc_count); + col_path_, *tmp_schema, 0, 0, id_map_, delete_store_, version_manager_, + options_, 0, doc_count); ASSERT_TRUE(segment != nullptr); } // Open { - Version v = version_manager->get_current_version(); + Version v = version_manager_->get_current_version(); auto result = - Segment::Open(col_path, *tmp_schema, *v.writing_segment_meta(), id_map, - delete_store, version_manager, options); + Segment::Open(col_path_, *tmp_schema, *v.writing_segment_meta(), id_map_, + delete_store_, version_manager_, options_); ASSERT_TRUE(result.has_value()); auto segment = result.value(); @@ -256,7 +179,7 @@ TEST_P(SegmentTest, InsertVectorTypes) { TEST_P(SegmentTest, FetchByGlobalDocID) { auto segment = test::TestHelper::CreateSegmentWithDoc( - col_path, *schema, 0, 0, id_map, delete_store, version_manager, options, + col_path_, *schema_, 0, 0, id_map_, delete_store_, version_manager_, options_, 0, 1); ASSERT_TRUE(segment != nullptr); @@ -269,7 +192,7 @@ TEST_P(SegmentTest, FetchByGlobalDocID) { TEST_P(SegmentTest, FetchSingleRow) { int doc_count = 10; auto segment = test::TestHelper::CreateSegmentWithDoc( - col_path, *schema, 0, 0, id_map, delete_store, version_manager, options, + col_path_, *schema_, 0, 0, id_map_, delete_store_, version_manager_, options_, 0, doc_count); ASSERT_TRUE(segment != nullptr); @@ -296,22 +219,23 @@ TEST_P(SegmentTest, FetchSingleRowWithPersistStore) { int doc_count = 1000; { auto segment = test::TestHelper::CreateSegmentWithDoc( - col_path, *schema, 0, 0, id_map, delete_store, version_manager, options, + col_path_, *schema_, 0, 0, id_map_, delete_store_, version_manager_, options_, 0, doc_count); ASSERT_TRUE(segment != nullptr); } // Open { - Version v = version_manager->get_current_version(); - SegmentOptions options; - options.read_only_ = false; - auto result = Segment::Open(col_path, *schema, *v.writing_segment_meta(), - id_map, delete_store, version_manager, options); + Version v = version_manager_->get_current_version(); + SegmentOptions open_options; + open_options.read_only_ = false; + auto result = Segment::Open(col_path_, *schema_, *v.writing_segment_meta(), + id_map_, delete_store_, version_manager_, + open_options); ASSERT_TRUE(result.has_value()); auto segment = result.value(); - test::TestHelper::SegmentInsertDoc(segment, *schema, doc_count, + test::TestHelper::SegmentInsertDoc(segment, *schema_, doc_count, doc_count * 2); auto func = [&](int index) -> void { @@ -335,7 +259,7 @@ TEST_P(SegmentTest, FetchSingleRowWithPersistStore) { TEST_P(SegmentTest, FetchSingleRowWithUserID) { auto segment = test::TestHelper::CreateSegmentWithDoc( - col_path, *schema, 0, 0, id_map, delete_store, version_manager, options, + col_path_, *schema_, 0, 0, id_map_, delete_store_, version_manager_, options_, 0, 10); ASSERT_TRUE(segment != nullptr); @@ -352,7 +276,7 @@ TEST_P(SegmentTest, FetchSingleRowWithUserID) { TEST_P(SegmentTest, FetchSingleRowWithGlobalDocID) { auto segment = test::TestHelper::CreateSegmentWithDoc( - col_path, *schema, 0, 0, id_map, delete_store, version_manager, options, + col_path_, *schema_, 0, 0, id_map_, delete_store_, version_manager_, options_, 0, 10); ASSERT_TRUE(segment != nullptr); @@ -367,218 +291,9 @@ TEST_P(SegmentTest, FetchSingleRowWithGlobalDocID) { global_doc_id_scalar) != nullptr); } -TEST_P(SegmentTest, FetchSingleRowWithLocalRowID) { - auto segment = test::TestHelper::CreateSegmentWithDoc( - col_path, *schema, 0, 0, id_map, delete_store, version_manager, options, - 0, 10); - ASSERT_TRUE(segment != nullptr); - - ExecBatchPtr batch = segment->fetch({LOCAL_ROW_ID, "id", "name"}, 4); - ASSERT_TRUE(batch != nullptr); - EXPECT_EQ(batch->length, 1); - EXPECT_EQ(batch->values.size(), 3); - - auto local_doc_id_scalar = batch->values[0].scalar(); - ASSERT_TRUE(local_doc_id_scalar != nullptr); - EXPECT_TRUE(std::dynamic_pointer_cast( - local_doc_id_scalar) != nullptr); - auto local_doc_id_value = - std::dynamic_pointer_cast(local_doc_id_scalar); - EXPECT_EQ(local_doc_id_value->value, 4); -} - -TEST_P(SegmentTest, FetchSingleRowWithLocalRowIDMiddle) { - auto segment = test::TestHelper::CreateSegmentWithDoc( - col_path, *schema, 0, 0, id_map, delete_store, version_manager, options, - 0, 10); - ASSERT_TRUE(segment != nullptr); - - ExecBatchPtr batch = segment->fetch({"id", LOCAL_ROW_ID, "name"}, 4); - ASSERT_TRUE(batch != nullptr); - EXPECT_EQ(batch->length, 1); - EXPECT_EQ(batch->values.size(), 3); - - auto local_doc_id_scalar = batch->values[1].scalar(); - ASSERT_TRUE(local_doc_id_scalar != nullptr); - EXPECT_TRUE(std::dynamic_pointer_cast( - local_doc_id_scalar) != nullptr); - auto local_doc_id_value = - std::dynamic_pointer_cast(local_doc_id_scalar); - EXPECT_EQ(local_doc_id_value->value, 4); -} - -TEST_P(SegmentTest, FetchSingleRowWithLocalRowIDEnd) { - auto segment = test::TestHelper::CreateSegmentWithDoc( - col_path, *schema, 0, 0, id_map, delete_store, version_manager, options, - 0, 10); - ASSERT_TRUE(segment != nullptr); - - ExecBatchPtr batch = segment->fetch({"id", "name", LOCAL_ROW_ID}, 4); - ASSERT_TRUE(batch != nullptr); - EXPECT_EQ(batch->length, 1); - EXPECT_EQ(batch->values.size(), 3); - - auto local_doc_id_scalar = batch->values[2].scalar(); - ASSERT_TRUE(local_doc_id_scalar != nullptr); - EXPECT_TRUE(std::dynamic_pointer_cast( - local_doc_id_scalar) != nullptr); - auto local_doc_id_value = - std::dynamic_pointer_cast(local_doc_id_scalar); - EXPECT_EQ(local_doc_id_value->value, 4); -} - -TEST_P(SegmentTest, CheckOrderWithLocalRowID) { - auto segment = test::TestHelper::CreateSegmentWithDoc( - col_path, *schema, 0, 0, id_map, delete_store, version_manager, options, - 0, 10); - ASSERT_TRUE(segment != nullptr); - - auto combined_reader = segment->scan({LOCAL_ROW_ID, "id", "name"}); - ASSERT_TRUE(combined_reader != nullptr); - EXPECT_TRUE(combined_reader->schema() != nullptr); - - std::shared_ptr batch; - uint32_t total_doc = 0; - while (true) { - auto status = combined_reader->ReadNext(&batch); - if (status.ok() == false) break; - if (batch == nullptr) break; - EXPECT_EQ(batch->num_columns(), 3); - EXPECT_EQ(batch->column(0)->type()->id(), arrow::Type::UINT64); - EXPECT_EQ(batch->column_name(0), LOCAL_ROW_ID); - total_doc += batch->num_rows(); - } - EXPECT_EQ(total_doc, 10); - - - std::vector indices = {0, 3, 6, 1, 0}; - auto combined_table = segment->fetch({LOCAL_ROW_ID, "id", "name"}, indices); - ASSERT_TRUE(combined_table != nullptr); - EXPECT_EQ(combined_table->num_columns(), 3); - EXPECT_EQ(combined_table->num_rows(), 5); - - auto field = combined_table->schema()->field(0); - EXPECT_EQ(field->name(), LOCAL_ROW_ID); - - // Get data from the LOCAL_ROW_ID column for each row - auto id_column = combined_table->column(0); - auto id_array = - std::dynamic_pointer_cast(id_column->chunk(0)); - - std::vector &expected_ids = indices; - std::vector actual_ids; - - for (int i = 0; i < id_array->length(); ++i) { - actual_ids.push_back(id_array->Value(i)); - } - - EXPECT_EQ(actual_ids, expected_ids) - << "ID column values don't match expected order"; -} - -TEST_P(SegmentTest, CheckOrderWithLocalRowIDMiddle) { - auto segment = test::TestHelper::CreateSegmentWithDoc( - col_path, *schema, 0, 0, id_map, delete_store, version_manager, options, - 0, 10); - ASSERT_TRUE(segment != nullptr); - - auto combined_reader = segment->scan({"id", LOCAL_ROW_ID, "name"}); - ASSERT_TRUE(combined_reader != nullptr); - EXPECT_TRUE(combined_reader->schema() != nullptr); - - std::shared_ptr batch; - uint32_t total_doc = 0; - while (true) { - auto status = combined_reader->ReadNext(&batch); - if (status.ok() == false) break; - if (batch == nullptr) break; - - EXPECT_EQ(batch->num_columns(), 3); - EXPECT_EQ(batch->column(1)->type()->id(), arrow::Type::UINT64); - EXPECT_EQ(batch->column_name(1), LOCAL_ROW_ID); - - total_doc += batch->num_rows(); - } - EXPECT_EQ(total_doc, 10); - - std::vector indices = {0, 3, 6, 1, 0}; - auto combined_table = segment->fetch({"id", LOCAL_ROW_ID, "name"}, indices); - ASSERT_TRUE(combined_table != nullptr); - EXPECT_EQ(combined_table->num_columns(), 3); - EXPECT_EQ(combined_table->num_rows(), 5); - - auto field = combined_table->schema()->field(1); - EXPECT_EQ(field->name(), LOCAL_ROW_ID); - - // Get data from the LOCAL_ROW_ID column for each row - auto id_column = combined_table->column(1); - auto id_array = - std::dynamic_pointer_cast(id_column->chunk(0)); - - std::vector &expected_ids = indices; - std::vector actual_ids; - - for (int i = 0; i < id_array->length(); ++i) { - actual_ids.push_back(id_array->Value(i)); - } - - EXPECT_EQ(actual_ids, expected_ids) - << "ID column values don't match expected order"; -} - -TEST_P(SegmentTest, CheckOrderWithLocalRowIDEnd) { - auto segment = test::TestHelper::CreateSegmentWithDoc( - col_path, *schema, 0, 0, id_map, delete_store, version_manager, options, - 0, 10); - ASSERT_TRUE(segment != nullptr); - - auto combined_reader = segment->scan({"id", "name", LOCAL_ROW_ID}); - ASSERT_TRUE(combined_reader != nullptr); - EXPECT_TRUE(combined_reader->schema() != nullptr); - - std::shared_ptr batch; - uint32_t total_doc = 0; - while (true) { - auto status = combined_reader->ReadNext(&batch); - if (status.ok() == false) break; - if (batch == nullptr) break; - - EXPECT_EQ(batch->num_columns(), 3); - EXPECT_EQ(batch->column(2)->type()->id(), arrow::Type::UINT64); - EXPECT_EQ(batch->column_name(2), LOCAL_ROW_ID); - - total_doc += batch->num_rows(); - } - EXPECT_EQ(total_doc, 10); - - std::vector indices = {0, 3, 6, 1, 0}; - auto combined_table = segment->fetch({"id", "name", LOCAL_ROW_ID}, indices); - ASSERT_TRUE(combined_table != nullptr); - EXPECT_EQ(combined_table->num_columns(), 3); - EXPECT_EQ(combined_table->num_rows(), 5); - - auto field = combined_table->schema()->field(2); - EXPECT_EQ(field->name(), LOCAL_ROW_ID); - - // Get data from the LOCAL_ROW_ID column for each row - auto id_column = combined_table->column(2); - auto id_array = - std::dynamic_pointer_cast(id_column->chunk(0)); - - std::vector &expected_ids = indices; - std::vector actual_ids; - - for (int i = 0; i < id_array->length(); ++i) { - actual_ids.push_back(id_array->Value(i)); - } - - EXPECT_EQ(actual_ids, expected_ids) - << "ID column values don't match expected order"; -} - TEST_P(SegmentTest, FetchSingleRowWithNegativeIndex) { auto segment = test::TestHelper::CreateSegmentWithDoc( - col_path, *schema, 0, 0, id_map, delete_store, version_manager, options, + col_path_, *schema_, 0, 0, id_map_, delete_store_, version_manager_, options_, 0, 10); ASSERT_TRUE(segment != nullptr); @@ -588,7 +303,7 @@ TEST_P(SegmentTest, FetchSingleRowWithNegativeIndex) { TEST_P(SegmentTest, FetchSingleRowWithOutOfRangeIndex) { auto segment = test::TestHelper::CreateSegmentWithDoc( - col_path, *schema, 0, 0, id_map, delete_store, version_manager, options, + col_path_, *schema_, 0, 0, id_map_, delete_store_, version_manager_, options_, 0, 10); ASSERT_TRUE(segment != nullptr); @@ -598,7 +313,7 @@ TEST_P(SegmentTest, FetchSingleRowWithOutOfRangeIndex) { TEST_P(SegmentTest, FetchSingleRowWithInvalidColumn) { auto segment = test::TestHelper::CreateSegmentWithDoc( - col_path, *schema, 0, 0, id_map, delete_store, version_manager, options, + col_path_, *schema_, 0, 0, id_map_, delete_store_, version_manager_, options_, 0, 10); ASSERT_TRUE(segment != nullptr); @@ -608,7 +323,7 @@ TEST_P(SegmentTest, FetchSingleRowWithInvalidColumn) { TEST_P(SegmentTest, FetchSingleRowWithEmptyColumns) { auto segment = test::TestHelper::CreateSegmentWithDoc( - col_path, *schema, 0, 0, id_map, delete_store, version_manager, options, + col_path_, *schema_, 0, 0, id_map_, delete_store_, version_manager_, options_, 0, 10); ASSERT_TRUE(segment != nullptr); @@ -618,7 +333,7 @@ TEST_P(SegmentTest, FetchSingleRowWithEmptyColumns) { TEST_P(SegmentTest, FetchSingleRowFromEmptySegment) { auto segment = test::TestHelper::CreateSegmentWithDoc( - col_path, *schema, 0, 0, id_map, delete_store, version_manager, options, + col_path_, *schema_, 0, 0, id_map_, delete_store_, version_manager_, options_, 0, 0); ASSERT_TRUE(segment != nullptr); @@ -628,7 +343,7 @@ TEST_P(SegmentTest, FetchSingleRowFromEmptySegment) { TEST_P(SegmentTest, FetchSingleRowWithBinaryFields) { auto segment = test::TestHelper::CreateSegmentWithDoc( - col_path, *schema, 0, 0, id_map, delete_store, version_manager, options, + col_path_, *schema_, 0, 0, id_map_, delete_store_, version_manager_, options_, 0, 10); ASSERT_TRUE(segment != nullptr); @@ -653,24 +368,24 @@ TEST_P(SegmentTest, Recover) { int doc_count = 100; { auto segment = test::TestHelper::CreateSegmentWithDoc( - col_path, *schema, 0, 0, id_map, delete_store, version_manager, options, + col_path_, *schema_, 0, 0, id_map_, delete_store_, version_manager_, options_, 0, doc_count); ASSERT_TRUE(segment != nullptr); } // simulate wal file { - Version v = version_manager->get_current_version(); + Version v = version_manager_->get_current_version(); auto writing_block_id = v.writing_segment_meta()->writing_forward_block_->id(); - auto wal_file = FileHelper::MakeWalPath(col_path, 0, writing_block_id); + auto wal_file = FileHelper::MakeWalPath(col_path_, 0, writing_block_id); WalOptions wal_option{0, true}; WalFilePtr wal_file_; WalFile::CreateAndOpen(wal_file, wal_option, &wal_file_); ASSERT_TRUE(wal_file_ != nullptr); for (int i = doc_count; i < doc_count + 100; i++) { - Doc doc = test::TestHelper::CreateDoc(i, *schema); + Doc doc = test::TestHelper::CreateDoc(i, *schema_); doc.set_operator(Operator::INSERT); std::vector buf = doc.serialize(); auto ret = wal_file_->append(std::string(buf.begin(), buf.end())); @@ -678,7 +393,7 @@ TEST_P(SegmentTest, Recover) { } for (int i = 0; i < doc_count; i++) { - Doc doc = test::TestHelper::CreateDoc(i, *schema); + Doc doc = test::TestHelper::CreateDoc(i, *schema_); doc.set_doc_id(i); // global doc id doc.set_operator(Operator::UPDATE); std::vector buf = doc.serialize(); @@ -687,7 +402,7 @@ TEST_P(SegmentTest, Recover) { } for (int i = 0; i < doc_count; i++) { - Doc doc = test::TestHelper::CreateDoc(i, *schema); + Doc doc = test::TestHelper::CreateDoc(i, *schema_); doc.set_operator(Operator::UPSERT); std::vector buf = doc.serialize(); auto ret = wal_file_->append(std::string(buf.begin(), buf.end())); @@ -695,7 +410,7 @@ TEST_P(SegmentTest, Recover) { } for (int i = 0; i < doc_count; i++) { - Doc doc = test::TestHelper::CreateDoc(i, *schema); + Doc doc = test::TestHelper::CreateDoc(i, *schema_); doc.set_doc_id(i + 300); // global doc id doc.set_operator(Operator::DELETE); std::vector buf = doc.serialize(); @@ -706,11 +421,12 @@ TEST_P(SegmentTest, Recover) { // recover { - Version v = version_manager->get_current_version(); - SegmentOptions options; - options.read_only_ = false; - auto result = Segment::Open(col_path, *schema, *v.writing_segment_meta(), - id_map, delete_store, version_manager, options); + Version v = version_manager_->get_current_version(); + SegmentOptions open_options; + open_options.read_only_ = false; + auto result = Segment::Open(col_path_, *schema_, *v.writing_segment_meta(), + id_map_, delete_store_, version_manager_, + open_options); ASSERT_TRUE(result.has_value()); auto segment = result.value(); @@ -728,8 +444,7 @@ TEST_P(SegmentTest, Recover) { // Why 400 ? because in segment we just mark deleted doc EXPECT_EQ(total_doc, 400); - // auto filter = segment->get_filter(); - auto filter = delete_store->make_filter(); + auto filter = delete_store_->make_filter(); auto actual_doc_count = segment->doc_count(filter); EXPECT_EQ(actual_doc_count, 100); } @@ -737,16 +452,16 @@ TEST_P(SegmentTest, Recover) { TEST_P(SegmentTest, UpdateDoc) { auto segment = test::TestHelper::CreateSegmentWithDoc( - col_path, *schema, 0, 0, id_map, delete_store, version_manager, options, + col_path_, *schema_, 0, 0, id_map_, delete_store_, version_manager_, options_, 0, 10); ASSERT_TRUE(segment != nullptr); // before update - uint64_t count = segment->doc_count(segment->get_filter()); + uint64_t count = segment->doc_count(delete_store_->make_filter()); EXPECT_EQ(count, 10); // Create a new document to update - Doc update_doc = test::TestHelper::CreateDoc(5, *schema); + Doc update_doc = test::TestHelper::CreateDoc(5, *schema_); update_doc.set("name", "updated_name"); update_doc.set("age", 99); @@ -755,7 +470,7 @@ TEST_P(SegmentTest, UpdateDoc) { EXPECT_TRUE(status.ok()) << "Update failed: " << status.message(); // after update - count = segment->doc_count(segment->get_filter()); + count = segment->doc_count(delete_store_->make_filter()); EXPECT_EQ(count, 10); // Fetch the updated document and verify changes @@ -769,23 +484,23 @@ TEST_P(SegmentTest, UpdateDoc) { TEST_P(SegmentTest, UpdateDocBatch) { int doc_count = 10; auto segment = test::TestHelper::CreateSegmentWithDoc( - col_path, *schema, 0, 0, id_map, delete_store, version_manager, options, + col_path_, *schema_, 0, 0, id_map_, delete_store_, version_manager_, options_, 0, doc_count); ASSERT_TRUE(segment != nullptr); // before update - uint64_t count = segment->doc_count(segment->get_filter()); + uint64_t count = segment->doc_count(delete_store_->make_filter()); EXPECT_EQ(count, doc_count); // Create a new document to update for (int i = 0; i < doc_count; i++) { - Doc update_doc = test::TestHelper::CreateDoc(i, *schema); + Doc update_doc = test::TestHelper::CreateDoc(i, *schema_); // Update the document auto status = segment->Update(update_doc); EXPECT_TRUE(status.ok()) << "Update failed: " << status.message(); } // after update - count = segment->doc_count(segment->get_filter()); + count = segment->doc_count(delete_store_->make_filter()); EXPECT_EQ(count, doc_count); // Fetch the updated document and verify changes @@ -798,12 +513,12 @@ TEST_P(SegmentTest, UpdateDocBatch) { TEST_P(SegmentTest, DeleteDoc) { auto segment = test::TestHelper::CreateSegmentWithDoc( - col_path, *schema, 0, 0, id_map, delete_store, version_manager, options, + col_path_, *schema_, 0, 0, id_map_, delete_store_, version_manager_, options_, 0, 10); ASSERT_TRUE(segment != nullptr); // before update - uint64_t count = segment->doc_count(segment->get_filter()); + uint64_t count = segment->doc_count(delete_store_->make_filter()); EXPECT_EQ(count, 10); // Delete a document by primary key @@ -811,7 +526,7 @@ TEST_P(SegmentTest, DeleteDoc) { EXPECT_TRUE(status.ok()) << "Delete by pk failed: " << status.message(); // after delete - count = segment->doc_count(segment->get_filter()); + count = segment->doc_count(delete_store_->make_filter()); EXPECT_EQ(count, 9); // Delete a document by global doc id @@ -819,19 +534,19 @@ TEST_P(SegmentTest, DeleteDoc) { EXPECT_TRUE(status.ok()) << "Delete by global doc id failed: " << status.message(); - count = segment->doc_count(segment->get_filter()); + count = segment->doc_count(delete_store_->make_filter()); EXPECT_EQ(count, 8); } TEST_P(SegmentTest, DeleteBatch) { int doc_count = 10; auto segment = test::TestHelper::CreateSegmentWithDoc( - col_path, *schema, 0, 0, id_map, delete_store, version_manager, options, + col_path_, *schema_, 0, 0, id_map_, delete_store_, version_manager_, options_, 0, doc_count); ASSERT_TRUE(segment != nullptr); // before update - uint64_t count = segment->doc_count(segment->get_filter()); + uint64_t count = segment->doc_count(delete_store_->make_filter()); EXPECT_EQ(count, doc_count); for (int i = 0; i < doc_count; i++) { @@ -840,29 +555,29 @@ TEST_P(SegmentTest, DeleteBatch) { } // after delete - count = segment->doc_count(segment->get_filter()); + count = segment->doc_count(delete_store_->make_filter()); EXPECT_EQ(count, 0); } TEST_P(SegmentTest, UpsertDoc) { auto segment = test::TestHelper::CreateSegmentWithDoc( - col_path, *schema, 0, 0, id_map, delete_store, version_manager, options, + col_path_, *schema_, 0, 0, id_map_, delete_store_, version_manager_, options_, 0, 5); ASSERT_TRUE(segment != nullptr); // before update - uint64_t count = segment->doc_count(segment->get_filter()); + uint64_t count = segment->doc_count(delete_store_->make_filter()); EXPECT_EQ(count, 5); // Upsert an existing document - Doc upsert_doc1 = test::TestHelper::CreateDoc(3, *schema); + Doc upsert_doc1 = test::TestHelper::CreateDoc(3, *schema_); upsert_doc1.set("name", "upserted_name"); auto status = segment->Upsert(upsert_doc1); EXPECT_TRUE(status.ok()) << "Upsert existing doc failed: " << status.message(); - count = segment->doc_count(segment->get_filter()); + count = segment->doc_count(delete_store_->make_filter()); EXPECT_EQ(count, 5); // Verify the update @@ -871,12 +586,12 @@ TEST_P(SegmentTest, UpsertDoc) { EXPECT_EQ(ret_doc->get("name"), "upserted_name"); // Upsert a new document - Doc upsert_doc2 = test::TestHelper::CreateDoc(6, *schema); + Doc upsert_doc2 = test::TestHelper::CreateDoc(6, *schema_); upsert_doc2.set("name", "new_upserted_doc"); status = segment->Upsert(upsert_doc2); EXPECT_TRUE(status.ok()) << "Upsert new doc failed: " << status.message(); - count = segment->doc_count(segment->get_filter()); + count = segment->doc_count(delete_store_->make_filter()); EXPECT_EQ(count, 6); // Verify the new document was inserted @@ -888,31 +603,31 @@ TEST_P(SegmentTest, UpsertDoc) { TEST_P(SegmentTest, UpsertDocBatch) { int doc_count = 10; auto segment = test::TestHelper::CreateSegmentWithDoc( - col_path, *schema, 0, 0, id_map, delete_store, version_manager, options, + col_path_, *schema_, 0, 0, id_map_, delete_store_, version_manager_, options_, 0, doc_count); ASSERT_TRUE(segment != nullptr); // before update - uint64_t count = segment->doc_count(segment->get_filter()); + uint64_t count = segment->doc_count(delete_store_->make_filter()); EXPECT_EQ(count, doc_count); for (int i = 0; i < doc_count; i++) { // Upsert existing document - Doc upsert_doc1 = test::TestHelper::CreateDoc(i, *schema); + Doc upsert_doc1 = test::TestHelper::CreateDoc(i, *schema_); upsert_doc1.set("name", "upserted_name" + std::to_string(i)); auto status = segment->Upsert(upsert_doc1); EXPECT_TRUE(status.ok()) << "Upsert existing doc failed: " << status.message(); // Upsert new document - Doc upsert_doc2 = test::TestHelper::CreateDoc(doc_count + i, *schema); + Doc upsert_doc2 = test::TestHelper::CreateDoc(doc_count + i, *schema_); upsert_doc2.set("name", "new_upserted_doc" + std::to_string(i)); status = segment->Upsert(upsert_doc2); EXPECT_TRUE(status.ok()) << "Upsert new doc failed: " << status.message(); } - count = segment->doc_count(segment->get_filter()); + count = segment->doc_count(delete_store_->make_filter()); EXPECT_EQ(count, doc_count * 2); int incr_idx = 0; @@ -933,7 +648,7 @@ TEST_P(SegmentTest, UpsertDocBatch) { TEST_P(SegmentTest, Flush) { auto segment = test::TestHelper::CreateSegmentWithDoc( - col_path, *schema, 0, 0, id_map, delete_store, version_manager, options, + col_path_, *schema_, 0, 0, id_map_, delete_store_, version_manager_, options_, 0, 100); ASSERT_TRUE(segment != nullptr); @@ -944,7 +659,7 @@ TEST_P(SegmentTest, Flush) { TEST_P(SegmentTest, FlushAfterInsert) { auto segment = test::TestHelper::CreateSegmentWithDoc( - col_path, *schema, 0, 0, id_map, delete_store, version_manager, options, + col_path_, *schema_, 0, 0, id_map_, delete_store_, version_manager_, options_, 0, 100); ASSERT_TRUE(segment != nullptr); @@ -952,7 +667,7 @@ TEST_P(SegmentTest, FlushAfterInsert) { auto status = segment->flush(); EXPECT_TRUE(status.ok()) << "Flush failed: " << status.message(); - test::TestHelper::SegmentInsertDoc(segment, *schema, 100, 150); + test::TestHelper::SegmentInsertDoc(segment, *schema_, 100, 150); ASSERT_EQ(segment->doc_count(), 150); @@ -960,7 +675,7 @@ TEST_P(SegmentTest, FlushAfterInsert) { auto ret_doc = segment->Fetch(i); EXPECT_TRUE(ret_doc != nullptr); - Doc verify_doc = test::TestHelper::CreateDoc(i, *schema); + Doc verify_doc = test::TestHelper::CreateDoc(i, *schema_); auto vv = verify_doc.get>("dense_fp32").value(); auto v = ret_doc->get>("dense_fp32").value(); for (uint32_t j = 0; j < vv.size(); j++) { @@ -971,7 +686,7 @@ TEST_P(SegmentTest, FlushAfterInsert) { TEST_P(SegmentTest, Dump) { auto segment = test::TestHelper::CreateSegmentWithDoc( - col_path, *schema, 0, 0, id_map, delete_store, version_manager, options, + col_path_, *schema_, 0, 0, id_map_, delete_store_, version_manager_, options_, 0, 100); ASSERT_TRUE(segment != nullptr); @@ -986,7 +701,7 @@ TEST_P(SegmentTest, Dump) { TEST_P(SegmentTest, DocCount) { auto segment = test::TestHelper::CreateSegmentWithDoc( - col_path, *schema, 0, 0, id_map, delete_store, version_manager, options, + col_path_, *schema_, 0, 0, id_map_, delete_store_, version_manager_, options_, 0, 50); ASSERT_TRUE(segment != nullptr); @@ -1000,21 +715,21 @@ TEST_P(SegmentTest, DocCount) { segment->Delete("pk_30"); // Get document count again - count = segment->doc_count(segment->get_filter()); + count = segment->doc_count(delete_store_->make_filter()); EXPECT_EQ(count, 47); } // TEST_P(SegmentTest, Insert100WData) { -// options.max_buffer_size_ = 8 * 1024 * 1024; +// options_.max_buffer_size_ = 8 * 1024 * 1024; // auto segment = test::TestHelper::CreateSegmentWithDoc( -// col_path, *schema, 0, 0, id_map, delete_store, version_manager, -// options, 0, 0); +// col_path_, *schema_, 0, 0, id_map_, delete_store_, version_manager_, +// options_, 0, 0); // ASSERT_TRUE(segment != nullptr); // uint64_t MAX_DOC = 1000000; // auto start = std::chrono::system_clock::now(); -// test::TestHelper::SegmentInsertDoc(segment, *schema, 0, MAX_DOC); +// test::TestHelper::SegmentInsertDoc(segment, *schema_, 0, MAX_DOC); // auto end = std::chrono::system_clock::now(); // auto cost = std::chrono::duration_cast(end - // start) @@ -1042,18 +757,18 @@ TEST_P(SegmentTest, DocCount) { // } TEST_P(SegmentTest, CombinedVectorColumnIndexer) { - options.max_buffer_size_ = 10 * 1024; + options_.max_buffer_size_ = 10 * 1024; auto segment = test::TestHelper::CreateSegmentWithDoc( - col_path, *schema, 0, 0, id_map, delete_store, version_manager, options, + col_path_, *schema_, 0, 0, id_map_, delete_store_, version_manager_, options_, 0, 0); ASSERT_TRUE(segment != nullptr); uint64_t MAX_DOC = 1000; - test::TestHelper::SegmentInsertDoc(segment, *schema, 0, MAX_DOC); + test::TestHelper::SegmentInsertDoc(segment, *schema_, 0, MAX_DOC); - Doc new_doc = test::TestHelper::CreateDoc(1000, *schema); + Doc new_doc = test::TestHelper::CreateDoc(1000, *schema_); auto status = segment->Insert(new_doc); ASSERT_TRUE(status.ok()); @@ -1075,7 +790,7 @@ TEST_P(SegmentTest, CombinedVectorColumnIndexer) { } // query - auto dense_fp32_field = schema->get_field("dense_fp32"); + auto dense_fp32_field = schema_->get_field("dense_fp32"); auto query_vector = new_doc.get>("dense_fp32").value(); auto query = vector_column_params::VectorData{ vector_column_params::DenseVector{query_vector.data()}}; @@ -1102,7 +817,7 @@ TEST_P(SegmentTest, CombinedVectorColumnIndexer) { } TEST_P(SegmentTest, CombinedVectorColumnIndexerWithQuantVectorIndex) { - options.max_buffer_size_ = 10 * 1024; + options_.max_buffer_size_ = 10 * 1024; auto tmp_schema = test::TestHelper::CreateSchemaWithVectorIndex( false, "demo", @@ -1110,15 +825,15 @@ TEST_P(SegmentTest, CombinedVectorColumnIndexerWithQuantVectorIndex) { QuantizeType::FP16)); auto segment = test::TestHelper::CreateSegmentWithDoc( - col_path, *tmp_schema, 0, 0, id_map, delete_store, version_manager, - options, 0, 0); + col_path_, *tmp_schema, 0, 0, id_map_, delete_store_, version_manager_, + options_, 0, 0); ASSERT_TRUE(segment != nullptr); uint64_t MAX_DOC = 1000; - test::TestHelper::SegmentInsertDoc(segment, *schema, 0, MAX_DOC); + test::TestHelper::SegmentInsertDoc(segment, *schema_, 0, MAX_DOC); - Doc new_doc = test::TestHelper::CreateDoc(1000, *schema); + Doc new_doc = test::TestHelper::CreateDoc(1000, *schema_); auto status = segment->Insert(new_doc); ASSERT_TRUE(status.ok()); @@ -1141,7 +856,7 @@ TEST_P(SegmentTest, CombinedVectorColumnIndexerWithQuantVectorIndex) { } // query - auto dense_fp32_field = schema->get_field("dense_fp32"); + auto dense_fp32_field = schema_->get_field("dense_fp32"); auto query_vector = new_doc.get>("dense_fp32").value(); auto query = vector_column_params::VectorData{ vector_column_params::DenseVector{query_vector.data()}}; @@ -1172,28 +887,28 @@ TEST_P(SegmentTest, CombinedVectorColumnIndexerWithQuantVectorIndex) { } TEST_P(SegmentTest, CombinedVectorColumnIndexerQueryWithPks) { - options.max_buffer_size_ = 10 * 1024; + options_.max_buffer_size_ = 10 * 1024; auto tmp_schema = test::TestHelper::CreateSchemaWithVectorIndex( false, "demo", std::make_shared(MetricType::IP)); auto segment = test::TestHelper::CreateSegmentWithDoc( - col_path, *tmp_schema, 0, 0, id_map, delete_store, version_manager, - options, 0, 0); + col_path_, *tmp_schema, 0, 0, id_map_, delete_store_, version_manager_, + options_, 0, 0); ASSERT_TRUE(segment != nullptr); uint64_t MAX_DOC = 1000; - test::TestHelper::SegmentInsertDoc(segment, *schema, 0, MAX_DOC); + test::TestHelper::SegmentInsertDoc(segment, *schema_, 0, MAX_DOC); auto combined_indexer = segment->get_combined_vector_indexer("dense_fp32"); ASSERT_TRUE(combined_indexer != nullptr); - Doc verify_doc = test::TestHelper::CreateDoc(999, *schema); + Doc verify_doc = test::TestHelper::CreateDoc(999, *schema_); std::vector> bf_pks = { {10, 20, 30, 40, 50, 60, 70, 80, 90, 999}}; // query - auto dense_fp32_field = schema->get_field("dense_fp32"); + auto dense_fp32_field = schema_->get_field("dense_fp32"); auto query_vector = verify_doc.get>("dense_fp32").value(); auto query = vector_column_params::VectorData{ vector_column_params::DenseVector{query_vector.data()}}; @@ -1232,7 +947,7 @@ TEST_P(SegmentTest, CombinedVectorColumnIndexerQueryWithPks) { TEST_P(SegmentTest, ConcurrentInsertOperations) { auto segment = test::TestHelper::CreateSegmentWithDoc( - col_path, *schema, 0, 0, id_map, delete_store, version_manager, options, + col_path_, *schema_, 0, 0, id_map_, delete_store_, version_manager_, options_, 0, 0); ASSERT_TRUE(segment != nullptr); @@ -1245,7 +960,7 @@ TEST_P(SegmentTest, ConcurrentInsertOperations) { threads.emplace_back([&, t]() { for (int i = 0; i < docs_per_thread; ++i) { int doc_id = t * docs_per_thread + i; - Doc doc = test::TestHelper::CreateDoc(doc_id, *schema); + Doc doc = test::TestHelper::CreateDoc(doc_id, *schema_); auto status = segment->Insert(doc); EXPECT_TRUE(status.ok()) << "Thread " << t << " insert failed for doc " << doc_id; @@ -1265,7 +980,7 @@ TEST_P(SegmentTest, ConcurrentInsertOperations) { TEST_P(SegmentTest, ConcurrentMixedOperations) { auto segment = test::TestHelper::CreateSegmentWithDoc( - col_path, *schema, 0, 0, id_map, delete_store, version_manager, options, + col_path_, *schema_, 0, 0, id_map_, delete_store_, version_manager_, options_, 0, 100); ASSERT_TRUE(segment != nullptr); @@ -1274,7 +989,7 @@ TEST_P(SegmentTest, ConcurrentMixedOperations) { // Thread 1: Insert new documents threads.emplace_back([&]() { for (int i = 100; i < 120; ++i) { - Doc doc = test::TestHelper::CreateDoc(i, *schema); + Doc doc = test::TestHelper::CreateDoc(i, *schema_); auto status = segment->Insert(doc); EXPECT_TRUE(status.ok() || status.code() == StatusCode::ALREADY_EXISTS); } @@ -1283,7 +998,7 @@ TEST_P(SegmentTest, ConcurrentMixedOperations) { // Thread 2: Update existing documents threads.emplace_back([&]() { for (int i = 0; i < 50; i += 5) { - Doc doc = test::TestHelper::CreateDoc(i, *schema); + Doc doc = test::TestHelper::CreateDoc(i, *schema_); doc.set("name", "updated_concurrent_" + std::to_string(i)); auto status = segment->Update(doc); EXPECT_TRUE(status.ok() || status.code() == StatusCode::NOT_FOUND); @@ -1307,11 +1022,11 @@ TEST_P(SegmentTest, ConcurrentMixedOperations) { // corner cases TEST_P(SegmentTest, DuplicateInsert) { auto segment = test::TestHelper::CreateSegmentWithDoc( - col_path, *schema, 0, 0, id_map, delete_store, version_manager, options, + col_path_, *schema_, 0, 0, id_map_, delete_store_, version_manager_, options_, 0, 0); ASSERT_TRUE(segment != nullptr); - Doc doc1 = test::TestHelper::CreateDoc(0, *schema); + Doc doc1 = test::TestHelper::CreateDoc(0, *schema_); auto status1 = segment->Insert(doc1); EXPECT_TRUE(status1.ok()) << "First insert failed: " << status1.message(); @@ -1336,7 +1051,7 @@ TEST_P(SegmentTest, DuplicateInsert) { TEST_P(SegmentTest, DuplicateDelete) { auto segment = test::TestHelper::CreateSegmentWithDoc( - col_path, *schema, 0, 0, id_map, delete_store, version_manager, options, + col_path_, *schema_, 0, 0, id_map_, delete_store_, version_manager_, options_, 0, 5); ASSERT_TRUE(segment != nullptr); @@ -1353,7 +1068,7 @@ TEST_P(SegmentTest, DuplicateDelete) { TEST_P(SegmentTest, DeleteNonExistentDoc) { auto segment = test::TestHelper::CreateSegmentWithDoc( - col_path, *schema, 0, 0, id_map, delete_store, version_manager, options, + col_path_, *schema_, 0, 0, id_map_, delete_store_, version_manager_, options_, 0, 5); ASSERT_TRUE(segment != nullptr); @@ -1363,11 +1078,11 @@ TEST_P(SegmentTest, DeleteNonExistentDoc) { TEST_P(SegmentTest, UpdateNonExistentDoc) { auto segment = test::TestHelper::CreateSegmentWithDoc( - col_path, *schema, 0, 0, id_map, delete_store, version_manager, options, + col_path_, *schema_, 0, 0, id_map_, delete_store_, version_manager_, options_, 0, 5); ASSERT_TRUE(segment != nullptr); - Doc doc = test::TestHelper::CreateDoc(999, *schema); + Doc doc = test::TestHelper::CreateDoc(999, *schema_); doc.set("name", "non_existent_doc"); auto status = segment->Update(doc); @@ -1376,25 +1091,25 @@ TEST_P(SegmentTest, UpdateNonExistentDoc) { TEST_P(SegmentTest, UpsertNonExistentDoc) { auto segment = test::TestHelper::CreateSegmentWithDoc( - col_path, *schema, 0, 0, id_map, delete_store, version_manager, options, + col_path_, *schema_, 0, 0, id_map_, delete_store_, version_manager_, options_, 0, 5); ASSERT_TRUE(segment != nullptr); - Doc doc = test::TestHelper::CreateDoc(999, *schema); + Doc doc = test::TestHelper::CreateDoc(999, *schema_); doc.set("name", "new_upserted_doc"); auto status = segment->Upsert(doc); EXPECT_TRUE(status.ok()) << "Upsert non-existent doc should succeed: " << status.message(); - auto filter = segment->get_filter(); + auto filter = delete_store_->make_filter(); uint64_t count = segment->doc_count(filter); EXPECT_EQ(count, 6); } TEST_P(SegmentTest, ScanWithEmptyColumns) { auto segment = test::TestHelper::CreateSegmentWithDoc( - col_path, *schema, 0, 0, id_map, delete_store, version_manager, options, + col_path_, *schema_, 0, 0, id_map_, delete_store_, version_manager_, options_, 0, 5); ASSERT_TRUE(segment != nullptr); @@ -1404,7 +1119,7 @@ TEST_P(SegmentTest, ScanWithEmptyColumns) { TEST_P(SegmentTest, ScanWithInvalidColumns) { auto segment = test::TestHelper::CreateSegmentWithDoc( - col_path, *schema, 0, 0, id_map, delete_store, version_manager, options, + col_path_, *schema_, 0, 0, id_map_, delete_store_, version_manager_, options_, 0, 10); ASSERT_TRUE(segment != nullptr); @@ -1415,7 +1130,7 @@ TEST_P(SegmentTest, ScanWithInvalidColumns) { TEST_P(SegmentTest, FetchNonExistentDoc) { auto segment = test::TestHelper::CreateSegmentWithDoc( - col_path, *schema, 0, 0, id_map, delete_store, version_manager, options, + col_path_, *schema_, 0, 0, id_map_, delete_store_, version_manager_, options_, 0, 5); ASSERT_TRUE(segment != nullptr); @@ -1423,36 +1138,36 @@ TEST_P(SegmentTest, FetchNonExistentDoc) { EXPECT_TRUE(doc == nullptr) << "Fetch non-existent doc should return nullptr"; } -TEST_P(SegmentTest, FetchWithInvalidIndices) { +TEST_P(SegmentTest, FetchWithInvalidSegmentDocIDs) { auto segment = test::TestHelper::CreateSegmentWithDoc( - col_path, *schema, 0, 0, id_map, delete_store, version_manager, options, + col_path_, *schema_, 0, 0, id_map_, delete_store_, version_manager_, options_, 0, 5); ASSERT_TRUE(segment != nullptr); - std::vector invalid_indices = {999, 1000}; - auto table = segment->fetch({"id", "name"}, invalid_indices); + std::vector invalid_segment_doc_ids = {999, 1000}; + auto table = segment->fetch({"id", "name"}, invalid_segment_doc_ids); ASSERT_TRUE(table == nullptr); } TEST_P(SegmentTest, FetchWithInvalidColumns) { auto segment = test::TestHelper::CreateSegmentWithDoc( - col_path, *schema, 0, 0, id_map, delete_store, version_manager, options, + col_path_, *schema_, 0, 0, id_map_, delete_store_, version_manager_, options_, 0, 10); ASSERT_TRUE(segment != nullptr); // Try to fetch with invalid column name - std::vector indices = {0, 1, 2}; - auto table = segment->fetch({"invalid_column"}, indices); + std::vector segment_doc_ids = {0, 1, 2}; + auto table = segment->fetch({"invalid_column"}, segment_doc_ids); EXPECT_TRUE(table == nullptr); } TEST_P(SegmentTest, InsertEmptyDocWithNullableSchema) { - auto nullable_schema = test::TestHelper::CreateNormalSchema(true, col_name); + auto nullable_schema = test::TestHelper::CreateNormalSchema(true, col_name_); auto segment = test::TestHelper::CreateSegmentWithDoc( - col_path, *nullable_schema, 0, 0, id_map, delete_store, version_manager, - options, 0, 0); + col_path_, *nullable_schema, 0, 0, id_map_, delete_store_, version_manager_, + options_, 0, 0); ASSERT_TRUE(segment != nullptr); Doc empty_doc; @@ -1463,7 +1178,7 @@ TEST_P(SegmentTest, InsertEmptyDocWithNullableSchema) { TEST_P(SegmentTest, MultipleDuplicateDeletes) { auto segment = test::TestHelper::CreateSegmentWithDoc( - col_path, *schema, 0, 0, id_map, delete_store, version_manager, options, + col_path_, *schema_, 0, 0, id_map_, delete_store_, version_manager_, options_, 0, 5); ASSERT_TRUE(segment != nullptr); @@ -1475,64 +1190,64 @@ TEST_P(SegmentTest, MultipleDuplicateDeletes) { EXPECT_FALSE(status.ok()) << "Delete iteration " << i << " should fail"; } - auto filter = segment->get_filter(); + auto filter = delete_store_->make_filter(); uint64_t count = segment->doc_count(filter); EXPECT_EQ(count, 4); } TEST_P(SegmentTest, FetchWithTwoVectorFields) { - schema->add_field(std::make_shared( + schema_->add_field(std::make_shared( "dense2_fp32", DataType::VECTOR_FP32, 128, false, std::make_shared(MetricType::IP))); int doc_count = 1000; auto segment = test::TestHelper::CreateSegmentWithDoc( - col_path, *schema, 0, 0, id_map, delete_store, version_manager, options, + col_path_, *schema_, 0, 0, id_map_, delete_store_, version_manager_, options_, 0, doc_count); ASSERT_TRUE(segment != nullptr); segment.reset(); - version_manager.reset(); - id_map->flush(); - id_map.reset(); + version_manager_.reset(); + id_map_->flush(); + id_map_.reset(); std::string delete_store_path = - FileHelper::MakeFilePath(col_path, FileID::DELETE_FILE, 0); - delete_store->flush(delete_store_path); - delete_store.reset(); + FileHelper::MakeFilePath(col_path_, FileID::DELETE_FILE, 0); + delete_store_->flush(delete_store_path); + delete_store_.reset(); - auto recover_version_manager = VersionManager::Recovery(col_path); + auto recover_version_manager = VersionManager::Recovery(col_path_); auto recover_version_mgr = recover_version_manager.value(); ASSERT_TRUE(recover_version_mgr != nullptr); auto v = recover_version_mgr->get_current_version(); // idmap - std::string idmap_path = FileHelper::MakeFilePath(col_path, FileID::ID_FILE, + std::string idmap_path = FileHelper::MakeFilePath(col_path_, FileID::ID_FILE, v.id_map_path_suffix()); - IDMap::Ptr recover_id_map = std::make_shared(col_name); + IDMap::Ptr recover_id_map = std::make_shared(col_name_); auto status = recover_id_map->open(idmap_path, false, false); ASSERT_TRUE(status.ok()); - delete_store_path = FileHelper::MakeFilePath(col_path, FileID::DELETE_FILE, + delete_store_path = FileHelper::MakeFilePath(col_path_, FileID::DELETE_FILE, v.delete_snapshot_path_suffix()); auto recover_delete_store = - DeleteStore::CreateAndLoad(col_name, delete_store_path); + DeleteStore::CreateAndLoad(col_name_, delete_store_path); ASSERT_TRUE(recover_delete_store != nullptr); int incr_doc_count = 1000; - auto result = Segment::Open(col_path, *schema, *v.writing_segment_meta(), + auto result = Segment::Open(col_path_, *schema_, *v.writing_segment_meta(), recover_id_map, recover_delete_store, - recover_version_mgr, options); + recover_version_mgr, options_); ASSERT_TRUE(result.has_value()); segment = std::move(result).value(); ASSERT_TRUE(segment != nullptr); auto s = test::TestHelper::SegmentInsertDoc( - segment, *schema, doc_count, doc_count + incr_doc_count, false); + segment, *schema_, doc_count, doc_count + incr_doc_count, false); ASSERT_TRUE(s.ok()); for (int i = 0; i < doc_count + incr_doc_count; i++) { - auto expect_doc = test::TestHelper::CreateDoc(i, *schema); + auto expect_doc = test::TestHelper::CreateDoc(i, *schema_); auto ret_doc = segment->Fetch(i); if (*ret_doc != expect_doc) { std::cout << " ret_doc: " << ret_doc->to_string() << std::endl; @@ -1545,9 +1260,9 @@ TEST_P(SegmentTest, FetchWithTwoVectorFields) { TEST_P(SegmentTest, FetchPerf) { // create segment int doc_count = 1000; - options.max_buffer_size_ = 100 * 1024; + options_.max_buffer_size_ = 100 * 1024; auto segment = test::TestHelper::CreateSegmentWithDoc( - col_path, *schema, 0, 0, id_map, delete_store, version_manager, options, + col_path_, *schema_, 0, 0, id_map_, delete_store_, version_manager_, options_, 0, doc_count); ASSERT_TRUE(segment != nullptr); @@ -1555,50 +1270,50 @@ TEST_P(SegmentTest, FetchPerf) { auto writing_segment_meta = segment->meta(); // convert writing segment meta to persisted segment meta - Version version = version_manager->get_current_version(); + Version version = version_manager_->get_current_version(); writing_segment_meta->remove_writing_forward_block(); auto s = version.add_persisted_segment_meta(writing_segment_meta); ASSERT_TRUE(s.ok()); - s = version_manager->apply(version); + s = version_manager_->apply(version); ASSERT_TRUE(s.ok()); - s = version_manager->flush(); + s = version_manager_->flush(); ASSERT_TRUE(s.ok()); segment.reset(); - version_manager.reset(); - id_map->flush(); - id_map.reset(); + version_manager_.reset(); + id_map_->flush(); + id_map_.reset(); std::string delete_store_path = - FileHelper::MakeFilePath(col_path, FileID::DELETE_FILE, 0); - delete_store->flush(delete_store_path); - delete_store.reset(); + FileHelper::MakeFilePath(col_path_, FileID::DELETE_FILE, 0); + delete_store_->flush(delete_store_path); + delete_store_.reset(); - auto recover_version_manager = VersionManager::Recovery(col_path); + auto recover_version_manager = VersionManager::Recovery(col_path_); auto recover_version_mgr = recover_version_manager.value(); ASSERT_TRUE(recover_version_mgr != nullptr); Version v = recover_version_mgr->get_current_version(); const auto &persist_metas = v.persisted_segment_metas(); // idmap - std::string idmap_path = FileHelper::MakeFilePath(col_path, FileID::ID_FILE, + std::string idmap_path = FileHelper::MakeFilePath(col_path_, FileID::ID_FILE, v.id_map_path_suffix()); - IDMap::Ptr recover_id_map = std::make_shared(col_name); + IDMap::Ptr recover_id_map = std::make_shared(col_name_); auto status = recover_id_map->open(idmap_path, false, false); ASSERT_TRUE(status.ok()); - delete_store_path = FileHelper::MakeFilePath(col_path, FileID::DELETE_FILE, + delete_store_path = FileHelper::MakeFilePath(col_path_, FileID::DELETE_FILE, v.delete_snapshot_path_suffix()); auto recover_delete_store = - DeleteStore::CreateAndLoad(col_name, delete_store_path); + DeleteStore::CreateAndLoad(col_name_, delete_store_path); ASSERT_TRUE(recover_delete_store != nullptr); // open persist segment - options.read_only_ = true; + options_.read_only_ = true; auto result = - Segment::Open(col_path, *schema, *persist_metas[0], recover_id_map, - recover_delete_store, recover_version_mgr, options); + Segment::Open(col_path_, *schema_, *persist_metas[0], recover_id_map, + recover_delete_store, recover_version_mgr, options_); ASSERT_TRUE(result.has_value()); segment = std::move(result).value(); ASSERT_TRUE(segment != nullptr); @@ -1608,13 +1323,13 @@ TEST_P(SegmentTest, FetchPerf) { "int32 + 1", AddColumnOptions()); EXPECT_TRUE(s.ok()); - std::vector indices = {0, 3, 6, 1, 0, 501, 999}; + std::vector segment_doc_ids = {0, 3, 6, 1, 0, 501, 999}; auto func = [&](const std::vector columns, int local_row_id_idx) -> void { - auto combined_table = segment->fetch(columns, indices); + auto combined_table = segment->fetch(columns, segment_doc_ids); ASSERT_TRUE(combined_table != nullptr); EXPECT_EQ(combined_table->num_columns(), columns.size()); - EXPECT_EQ(combined_table->num_rows(), indices.size()); + EXPECT_EQ(combined_table->num_rows(), segment_doc_ids.size()); auto field = combined_table->schema()->field(local_row_id_idx); EXPECT_EQ(field->name(), LOCAL_ROW_ID); @@ -1624,7 +1339,7 @@ TEST_P(SegmentTest, FetchPerf) { auto id_array = std::dynamic_pointer_cast(id_column->chunk(0)); - std::vector &expected_ids = indices; + std::vector &expected_ids = segment_doc_ids; std::vector actual_ids; for (int i = 0; i < id_array->length(); ++i) { @@ -1649,10 +1364,10 @@ TEST_P(SegmentTest, FetchPerf) { TEST_P(SegmentTest, AddColumn) { // create segment - options.max_buffer_size_ = 10 * 1024 * 1024; + options_.max_buffer_size_ = 10 * 1024 * 1024; int doc_count = 1000; auto segment = test::TestHelper::CreateSegmentWithDoc( - col_path, *schema, 0, 0, id_map, delete_store, version_manager, options, + col_path_, *schema_, 0, 0, id_map_, delete_store_, version_manager_, options_, 0, doc_count); ASSERT_TRUE(segment != nullptr); @@ -1665,50 +1380,50 @@ TEST_P(SegmentTest, AddColumn) { auto writing_segment_meta = segment->meta(); // convert writing segment meta to persisted segment meta - Version version = version_manager->get_current_version(); + Version version = version_manager_->get_current_version(); writing_segment_meta->remove_writing_forward_block(); s = version.add_persisted_segment_meta(writing_segment_meta); ASSERT_TRUE(s.ok()); - s = version_manager->apply(version); + s = version_manager_->apply(version); ASSERT_TRUE(s.ok()); - s = version_manager->flush(); + s = version_manager_->flush(); ASSERT_TRUE(s.ok()); segment.reset(); - version_manager.reset(); - id_map->flush(); - id_map.reset(); + version_manager_.reset(); + id_map_->flush(); + id_map_.reset(); std::string delete_store_path = - FileHelper::MakeFilePath(col_path, FileID::DELETE_FILE, 0); - delete_store->flush(delete_store_path); - delete_store.reset(); + FileHelper::MakeFilePath(col_path_, FileID::DELETE_FILE, 0); + delete_store_->flush(delete_store_path); + delete_store_.reset(); - auto recover_version_manager = VersionManager::Recovery(col_path); + auto recover_version_manager = VersionManager::Recovery(col_path_); auto recover_version_mgr = recover_version_manager.value(); ASSERT_TRUE(recover_version_mgr != nullptr); Version v = recover_version_mgr->get_current_version(); const auto &persist_metas = v.persisted_segment_metas(); // idmap - std::string idmap_path = FileHelper::MakeFilePath(col_path, FileID::ID_FILE, + std::string idmap_path = FileHelper::MakeFilePath(col_path_, FileID::ID_FILE, v.id_map_path_suffix()); - IDMap::Ptr recover_id_map = std::make_shared(col_name); + IDMap::Ptr recover_id_map = std::make_shared(col_name_); auto status = recover_id_map->open(idmap_path, false, false); ASSERT_TRUE(status.ok()); - delete_store_path = FileHelper::MakeFilePath(col_path, FileID::DELETE_FILE, + delete_store_path = FileHelper::MakeFilePath(col_path_, FileID::DELETE_FILE, v.delete_snapshot_path_suffix()); auto recover_delete_store = - DeleteStore::CreateAndLoad(col_name, delete_store_path); + DeleteStore::CreateAndLoad(col_name_, delete_store_path); ASSERT_TRUE(recover_delete_store != nullptr); // open persist segment - options.read_only_ = true; + options_.read_only_ = true; auto result = - Segment::Open(col_path, *schema, *persist_metas[0], recover_id_map, - recover_delete_store, recover_version_mgr, options); + Segment::Open(col_path_, *schema_, *persist_metas[0], recover_id_map, + recover_delete_store, recover_version_mgr, options_); ASSERT_TRUE(result.has_value()); segment = std::move(result).value(); ASSERT_TRUE(segment != nullptr); @@ -1766,7 +1481,7 @@ TEST_P(SegmentTest, AddColumn) { } EXPECT_EQ(total_doc, doc_count); - auto new_schema = *schema; + auto new_schema = *schema_; new_schema.add_field(field_schema); auto check_doc = [&](int doc_count) { @@ -1867,13 +1582,14 @@ TEST_P(SegmentTest, AddColumn) { for (auto &[column_name, field_schema] : test_column_schemas) { auto expressions = test_expressions[column_name]; for (auto &expression : expressions) { - std::string col_name = column_name + "_" + - std::to_string(ailego::Crc32c::Hash( - expression.data(), expression.size())); + std::string generated_col_name = + column_name + "_" + + std::to_string( + ailego::Crc32c::Hash(expression.data(), expression.size())); auto new_field_schema = std::make_shared( field_schema->name(), field_schema->data_type(), field_schema->nullable(), field_schema->index_params()); - new_field_schema->set_name(col_name); + new_field_schema->set_name(generated_col_name); func(new_field_schema, expression); } } @@ -1881,53 +1597,53 @@ TEST_P(SegmentTest, AddColumn) { TEST_P(SegmentTest, AddNullableColumnWithoutExpressionMultiBlock) { // Use small buffer to force multiple scalar blocks within a single segment - options.max_buffer_size_ = 1 * 1024; + options_.max_buffer_size_ = 1 * 1024; int doc_count = 100; auto segment = test::TestHelper::CreateSegmentWithDoc( - col_path, *schema, 0, 0, id_map, delete_store, version_manager, options, + col_path_, *schema_, 0, 0, id_map_, delete_store_, version_manager_, options_, 0, doc_count); ASSERT_TRUE(segment != nullptr); segment->dump(); auto writing_segment_meta = segment->meta(); - Version version = version_manager->get_current_version(); + Version version = version_manager_->get_current_version(); writing_segment_meta->remove_writing_forward_block(); auto s = version.add_persisted_segment_meta(writing_segment_meta); ASSERT_TRUE(s.ok()); - s = version_manager->apply(version); + s = version_manager_->apply(version); ASSERT_TRUE(s.ok()); - s = version_manager->flush(); + s = version_manager_->flush(); ASSERT_TRUE(s.ok()); segment.reset(); - version_manager.reset(); - id_map->flush(); - id_map.reset(); + version_manager_.reset(); + id_map_->flush(); + id_map_.reset(); std::string delete_store_path = - FileHelper::MakeFilePath(col_path, FileID::DELETE_FILE, 0); - delete_store->flush(delete_store_path); - delete_store.reset(); + FileHelper::MakeFilePath(col_path_, FileID::DELETE_FILE, 0); + delete_store_->flush(delete_store_path); + delete_store_.reset(); - auto recover_version_manager = VersionManager::Recovery(col_path); + auto recover_version_manager = VersionManager::Recovery(col_path_); auto recover_version_mgr = recover_version_manager.value(); ASSERT_TRUE(recover_version_mgr != nullptr); Version v = recover_version_mgr->get_current_version(); const auto &persist_metas = v.persisted_segment_metas(); - std::string idmap_path = FileHelper::MakeFilePath(col_path, FileID::ID_FILE, + std::string idmap_path = FileHelper::MakeFilePath(col_path_, FileID::ID_FILE, v.id_map_path_suffix()); - IDMap::Ptr recover_id_map = std::make_shared(col_name); + IDMap::Ptr recover_id_map = std::make_shared(col_name_); auto status = recover_id_map->open(idmap_path, false, false); ASSERT_TRUE(status.ok()); - delete_store_path = FileHelper::MakeFilePath(col_path, FileID::DELETE_FILE, + delete_store_path = FileHelper::MakeFilePath(col_path_, FileID::DELETE_FILE, v.delete_snapshot_path_suffix()); auto recover_delete_store = - DeleteStore::CreateAndLoad(col_name, delete_store_path); + DeleteStore::CreateAndLoad(col_name_, delete_store_path); ASSERT_TRUE(recover_delete_store != nullptr); // Verify we have multiple scalar blocks @@ -1939,10 +1655,10 @@ TEST_P(SegmentTest, AddNullableColumnWithoutExpressionMultiBlock) { } ASSERT_GT(scalar_block_count, 1); - options.read_only_ = true; + options_.read_only_ = true; auto result = - Segment::Open(col_path, *schema, *persist_metas[0], recover_id_map, - recover_delete_store, recover_version_mgr, options); + Segment::Open(col_path_, *schema_, *persist_metas[0], recover_id_map, + recover_delete_store, recover_version_mgr, options_); ASSERT_TRUE(result.has_value()); segment = std::move(result).value(); ASSERT_TRUE(segment != nullptr); @@ -1954,14 +1670,16 @@ TEST_P(SegmentTest, AddNullableColumnWithoutExpressionMultiBlock) { {"add_uint32_null", DataType::UINT32}, {"add_int64_null", DataType::INT64}, }; - for (auto &[col_name, data_type] : nullable_types) { + for (auto &[nullable_col_name, data_type] : nullable_types) { auto field_schema = - std::make_shared(col_name, data_type, true); + std::make_shared(nullable_col_name, data_type, true); s = segment->add_column(field_schema, "", AddColumnOptions()); - ASSERT_TRUE(s.ok()) << "Failed to add nullable column " << col_name << ": " - << s.message(); + ASSERT_TRUE(s.ok()) + << "Failed to add nullable column " << nullable_col_name << ": " + << s.message(); - auto combined_reader = segment->scan({"id", "name", "age", col_name}); + auto combined_reader = + segment->scan({"id", "name", "age", nullable_col_name}); ASSERT_TRUE(combined_reader != nullptr); std::shared_ptr batch; uint32_t total_doc = 0; @@ -1978,53 +1696,53 @@ TEST_P(SegmentTest, AddNullableColumnWithoutExpressionMultiBlock) { TEST_P(SegmentTest, AddColumnWithExpressionMultiBlock) { // Use small buffer to force multiple scalar blocks within a single segment - options.max_buffer_size_ = 1 * 1024; + options_.max_buffer_size_ = 1 * 1024; int doc_count = 100; auto segment = test::TestHelper::CreateSegmentWithDoc( - col_path, *schema, 0, 0, id_map, delete_store, version_manager, options, + col_path_, *schema_, 0, 0, id_map_, delete_store_, version_manager_, options_, 0, doc_count); ASSERT_TRUE(segment != nullptr); segment->dump(); auto writing_segment_meta = segment->meta(); - Version version = version_manager->get_current_version(); + Version version = version_manager_->get_current_version(); writing_segment_meta->remove_writing_forward_block(); auto s = version.add_persisted_segment_meta(writing_segment_meta); ASSERT_TRUE(s.ok()); - s = version_manager->apply(version); + s = version_manager_->apply(version); ASSERT_TRUE(s.ok()); - s = version_manager->flush(); + s = version_manager_->flush(); ASSERT_TRUE(s.ok()); segment.reset(); - version_manager.reset(); - id_map->flush(); - id_map.reset(); + version_manager_.reset(); + id_map_->flush(); + id_map_.reset(); std::string delete_store_path = - FileHelper::MakeFilePath(col_path, FileID::DELETE_FILE, 0); - delete_store->flush(delete_store_path); - delete_store.reset(); + FileHelper::MakeFilePath(col_path_, FileID::DELETE_FILE, 0); + delete_store_->flush(delete_store_path); + delete_store_.reset(); - auto recover_version_manager = VersionManager::Recovery(col_path); + auto recover_version_manager = VersionManager::Recovery(col_path_); auto recover_version_mgr = recover_version_manager.value(); ASSERT_TRUE(recover_version_mgr != nullptr); Version v = recover_version_mgr->get_current_version(); const auto &persist_metas = v.persisted_segment_metas(); - std::string idmap_path = FileHelper::MakeFilePath(col_path, FileID::ID_FILE, + std::string idmap_path = FileHelper::MakeFilePath(col_path_, FileID::ID_FILE, v.id_map_path_suffix()); - IDMap::Ptr recover_id_map = std::make_shared(col_name); + IDMap::Ptr recover_id_map = std::make_shared(col_name_); auto status = recover_id_map->open(idmap_path, false, false); ASSERT_TRUE(status.ok()); - delete_store_path = FileHelper::MakeFilePath(col_path, FileID::DELETE_FILE, + delete_store_path = FileHelper::MakeFilePath(col_path_, FileID::DELETE_FILE, v.delete_snapshot_path_suffix()); auto recover_delete_store = - DeleteStore::CreateAndLoad(col_name, delete_store_path); + DeleteStore::CreateAndLoad(col_name_, delete_store_path); ASSERT_TRUE(recover_delete_store != nullptr); // Verify we have multiple scalar blocks @@ -2036,10 +1754,10 @@ TEST_P(SegmentTest, AddColumnWithExpressionMultiBlock) { } ASSERT_GT(scalar_block_count, 1); - options.read_only_ = true; + options_.read_only_ = true; auto result = - Segment::Open(col_path, *schema, *persist_metas[0], recover_id_map, - recover_delete_store, recover_version_mgr, options); + Segment::Open(col_path_, *schema_, *persist_metas[0], recover_id_map, + recover_delete_store, recover_version_mgr, options_); ASSERT_TRUE(result.has_value()); segment = std::move(result).value(); ASSERT_TRUE(segment != nullptr); @@ -2086,52 +1804,52 @@ TEST_P(SegmentTest, AddColumnWithExpressionMultiBlock) { } TEST_P(SegmentTest, AlterColumnMultiBlock) { - options.max_buffer_size_ = 1 * 1024; + options_.max_buffer_size_ = 1 * 1024; int doc_count = 100; auto segment = test::TestHelper::CreateSegmentWithDoc( - col_path, *schema, 0, 0, id_map, delete_store, version_manager, options, + col_path_, *schema_, 0, 0, id_map_, delete_store_, version_manager_, options_, 0, doc_count); ASSERT_TRUE(segment != nullptr); segment->dump(); auto writing_segment_meta = segment->meta(); - Version version = version_manager->get_current_version(); + Version version = version_manager_->get_current_version(); writing_segment_meta->remove_writing_forward_block(); auto s = version.add_persisted_segment_meta(writing_segment_meta); ASSERT_TRUE(s.ok()); - s = version_manager->apply(version); + s = version_manager_->apply(version); ASSERT_TRUE(s.ok()); - s = version_manager->flush(); + s = version_manager_->flush(); ASSERT_TRUE(s.ok()); segment.reset(); - version_manager.reset(); - id_map->flush(); - id_map.reset(); + version_manager_.reset(); + id_map_->flush(); + id_map_.reset(); std::string delete_store_path = - FileHelper::MakeFilePath(col_path, FileID::DELETE_FILE, 0); - delete_store->flush(delete_store_path); - delete_store.reset(); + FileHelper::MakeFilePath(col_path_, FileID::DELETE_FILE, 0); + delete_store_->flush(delete_store_path); + delete_store_.reset(); - auto recover_version_manager = VersionManager::Recovery(col_path); + auto recover_version_manager = VersionManager::Recovery(col_path_); auto recover_version_mgr = recover_version_manager.value(); ASSERT_TRUE(recover_version_mgr != nullptr); Version v = recover_version_mgr->get_current_version(); const auto &persist_metas = v.persisted_segment_metas(); - std::string idmap_path = FileHelper::MakeFilePath(col_path, FileID::ID_FILE, + std::string idmap_path = FileHelper::MakeFilePath(col_path_, FileID::ID_FILE, v.id_map_path_suffix()); - IDMap::Ptr recover_id_map = std::make_shared(col_name); + IDMap::Ptr recover_id_map = std::make_shared(col_name_); auto status = recover_id_map->open(idmap_path, false, false); ASSERT_TRUE(status.ok()); - delete_store_path = FileHelper::MakeFilePath(col_path, FileID::DELETE_FILE, + delete_store_path = FileHelper::MakeFilePath(col_path_, FileID::DELETE_FILE, v.delete_snapshot_path_suffix()); auto recover_delete_store = - DeleteStore::CreateAndLoad(col_name, delete_store_path); + DeleteStore::CreateAndLoad(col_name_, delete_store_path); ASSERT_TRUE(recover_delete_store != nullptr); int scalar_block_count = 0; @@ -2142,10 +1860,10 @@ TEST_P(SegmentTest, AlterColumnMultiBlock) { } ASSERT_GT(scalar_block_count, 1); - options.read_only_ = true; + options_.read_only_ = true; auto result = - Segment::Open(col_path, *schema, *persist_metas[0], recover_id_map, - recover_delete_store, recover_version_mgr, options); + Segment::Open(col_path_, *schema_, *persist_metas[0], recover_id_map, + recover_delete_store, recover_version_mgr, options_); ASSERT_TRUE(result.has_value()); segment = std::move(result).value(); ASSERT_TRUE(segment != nullptr); @@ -2170,52 +1888,52 @@ TEST_P(SegmentTest, AlterColumnMultiBlock) { } TEST_P(SegmentTest, DropColumnMultiBlock) { - options.max_buffer_size_ = 1 * 1024; + options_.max_buffer_size_ = 1 * 1024; int doc_count = 100; auto segment = test::TestHelper::CreateSegmentWithDoc( - col_path, *schema, 0, 0, id_map, delete_store, version_manager, options, + col_path_, *schema_, 0, 0, id_map_, delete_store_, version_manager_, options_, 0, doc_count); ASSERT_TRUE(segment != nullptr); segment->dump(); auto writing_segment_meta = segment->meta(); - Version version = version_manager->get_current_version(); + Version version = version_manager_->get_current_version(); writing_segment_meta->remove_writing_forward_block(); auto s = version.add_persisted_segment_meta(writing_segment_meta); ASSERT_TRUE(s.ok()); - s = version_manager->apply(version); + s = version_manager_->apply(version); ASSERT_TRUE(s.ok()); - s = version_manager->flush(); + s = version_manager_->flush(); ASSERT_TRUE(s.ok()); segment.reset(); - version_manager.reset(); - id_map->flush(); - id_map.reset(); + version_manager_.reset(); + id_map_->flush(); + id_map_.reset(); std::string delete_store_path = - FileHelper::MakeFilePath(col_path, FileID::DELETE_FILE, 0); - delete_store->flush(delete_store_path); - delete_store.reset(); + FileHelper::MakeFilePath(col_path_, FileID::DELETE_FILE, 0); + delete_store_->flush(delete_store_path); + delete_store_.reset(); - auto recover_version_manager = VersionManager::Recovery(col_path); + auto recover_version_manager = VersionManager::Recovery(col_path_); auto recover_version_mgr = recover_version_manager.value(); ASSERT_TRUE(recover_version_mgr != nullptr); Version v = recover_version_mgr->get_current_version(); const auto &persist_metas = v.persisted_segment_metas(); - std::string idmap_path = FileHelper::MakeFilePath(col_path, FileID::ID_FILE, + std::string idmap_path = FileHelper::MakeFilePath(col_path_, FileID::ID_FILE, v.id_map_path_suffix()); - IDMap::Ptr recover_id_map = std::make_shared(col_name); + IDMap::Ptr recover_id_map = std::make_shared(col_name_); auto status = recover_id_map->open(idmap_path, false, false); ASSERT_TRUE(status.ok()); - delete_store_path = FileHelper::MakeFilePath(col_path, FileID::DELETE_FILE, + delete_store_path = FileHelper::MakeFilePath(col_path_, FileID::DELETE_FILE, v.delete_snapshot_path_suffix()); auto recover_delete_store = - DeleteStore::CreateAndLoad(col_name, delete_store_path); + DeleteStore::CreateAndLoad(col_name_, delete_store_path); ASSERT_TRUE(recover_delete_store != nullptr); int scalar_block_count = 0; @@ -2226,10 +1944,10 @@ TEST_P(SegmentTest, DropColumnMultiBlock) { } ASSERT_GT(scalar_block_count, 1); - options.read_only_ = true; + options_.read_only_ = true; auto result = - Segment::Open(col_path, *schema, *persist_metas[0], recover_id_map, - recover_delete_store, recover_version_mgr, options); + Segment::Open(col_path_, *schema_, *persist_metas[0], recover_id_map, + recover_delete_store, recover_version_mgr, options_); ASSERT_TRUE(result.has_value()); segment = std::move(result).value(); ASSERT_TRUE(segment != nullptr); @@ -2257,58 +1975,58 @@ TEST_P(SegmentTest, DropColumnMultiBlock) { } TEST_P(SegmentTest, AddNullableThenAlterDropMultiBlock) { - options.max_buffer_size_ = 1 * 1024; + options_.max_buffer_size_ = 1 * 1024; int doc_count = 100; auto segment = test::TestHelper::CreateSegmentWithDoc( - col_path, *schema, 0, 0, id_map, delete_store, version_manager, options, + col_path_, *schema_, 0, 0, id_map_, delete_store_, version_manager_, options_, 0, doc_count); ASSERT_TRUE(segment != nullptr); segment->dump(); auto writing_segment_meta = segment->meta(); - Version version = version_manager->get_current_version(); + Version version = version_manager_->get_current_version(); writing_segment_meta->remove_writing_forward_block(); auto s = version.add_persisted_segment_meta(writing_segment_meta); ASSERT_TRUE(s.ok()); - s = version_manager->apply(version); + s = version_manager_->apply(version); ASSERT_TRUE(s.ok()); - s = version_manager->flush(); + s = version_manager_->flush(); ASSERT_TRUE(s.ok()); segment.reset(); - version_manager.reset(); - id_map->flush(); - id_map.reset(); + version_manager_.reset(); + id_map_->flush(); + id_map_.reset(); std::string delete_store_path = - FileHelper::MakeFilePath(col_path, FileID::DELETE_FILE, 0); - delete_store->flush(delete_store_path); - delete_store.reset(); + FileHelper::MakeFilePath(col_path_, FileID::DELETE_FILE, 0); + delete_store_->flush(delete_store_path); + delete_store_.reset(); - auto recover_version_manager = VersionManager::Recovery(col_path); + auto recover_version_manager = VersionManager::Recovery(col_path_); auto recover_version_mgr = recover_version_manager.value(); ASSERT_TRUE(recover_version_mgr != nullptr); Version v = recover_version_mgr->get_current_version(); const auto &persist_metas = v.persisted_segment_metas(); - std::string idmap_path = FileHelper::MakeFilePath(col_path, FileID::ID_FILE, + std::string idmap_path = FileHelper::MakeFilePath(col_path_, FileID::ID_FILE, v.id_map_path_suffix()); - IDMap::Ptr recover_id_map = std::make_shared(col_name); + IDMap::Ptr recover_id_map = std::make_shared(col_name_); auto status = recover_id_map->open(idmap_path, false, false); ASSERT_TRUE(status.ok()); - delete_store_path = FileHelper::MakeFilePath(col_path, FileID::DELETE_FILE, + delete_store_path = FileHelper::MakeFilePath(col_path_, FileID::DELETE_FILE, v.delete_snapshot_path_suffix()); auto recover_delete_store = - DeleteStore::CreateAndLoad(col_name, delete_store_path); + DeleteStore::CreateAndLoad(col_name_, delete_store_path); ASSERT_TRUE(recover_delete_store != nullptr); - options.read_only_ = true; + options_.read_only_ = true; auto result = - Segment::Open(col_path, *schema, *persist_metas[0], recover_id_map, - recover_delete_store, recover_version_mgr, options); + Segment::Open(col_path_, *schema_, *persist_metas[0], recover_id_map, + recover_delete_store, recover_version_mgr, options_); ASSERT_TRUE(result.has_value()); segment = std::move(result).value(); ASSERT_TRUE(segment != nullptr); @@ -2369,7 +2087,7 @@ TEST_P(SegmentTest, AlterColumn) { // create segment int doc_count = 1000; auto segment = test::TestHelper::CreateSegmentWithDoc( - col_path, *schema, 0, 0, id_map, delete_store, version_manager, options, + col_path_, *schema_, 0, 0, id_map_, delete_store_, version_manager_, options_, 0, doc_count); ASSERT_TRUE(segment != nullptr); @@ -2383,27 +2101,27 @@ TEST_P(SegmentTest, AlterColumn) { auto writing_segment_meta = segment->meta(); // convert writing segment meta to persisted segment meta - Version version = version_manager->get_current_version(); + Version version = version_manager_->get_current_version(); writing_segment_meta->remove_writing_forward_block(); s = version.add_persisted_segment_meta(writing_segment_meta); ASSERT_TRUE(s.ok()); - s = version_manager->apply(version); + s = version_manager_->apply(version); ASSERT_TRUE(s.ok()); - s = version_manager->flush(); + s = version_manager_->flush(); ASSERT_TRUE(s.ok()); segment.reset(); - version_manager.reset(); - id_map->flush(); - id_map.reset(); + version_manager_.reset(); + id_map_->flush(); + id_map_.reset(); std::string delete_store_path = - FileHelper::MakeFilePath(col_path, FileID::DELETE_FILE, 0); - delete_store->flush(delete_store_path); - delete_store.reset(); + FileHelper::MakeFilePath(col_path_, FileID::DELETE_FILE, 0); + delete_store_->flush(delete_store_path); + delete_store_.reset(); - auto recover_version_manager = VersionManager::Recovery(col_path); + auto recover_version_manager = VersionManager::Recovery(col_path_); auto recover_version_mgr = recover_version_manager.value(); ASSERT_TRUE(recover_version_mgr != nullptr); @@ -2411,23 +2129,23 @@ TEST_P(SegmentTest, AlterColumn) { const auto &persist_metas = v.persisted_segment_metas(); // idmap - std::string idmap_path = FileHelper::MakeFilePath(col_path, FileID::ID_FILE, + std::string idmap_path = FileHelper::MakeFilePath(col_path_, FileID::ID_FILE, v.id_map_path_suffix()); - IDMap::Ptr recover_id_map = std::make_shared(col_name); + IDMap::Ptr recover_id_map = std::make_shared(col_name_); auto status = recover_id_map->open(idmap_path, false, false); ASSERT_TRUE(status.ok()); - delete_store_path = FileHelper::MakeFilePath(col_path, FileID::DELETE_FILE, + delete_store_path = FileHelper::MakeFilePath(col_path_, FileID::DELETE_FILE, v.delete_snapshot_path_suffix()); auto recover_delete_store = - DeleteStore::CreateAndLoad(col_name, delete_store_path); + DeleteStore::CreateAndLoad(col_name_, delete_store_path); ASSERT_TRUE(recover_delete_store != nullptr); // open persist segment - options.read_only_ = true; + options_.read_only_ = true; auto result = - Segment::Open(col_path, *schema, *persist_metas[0], recover_id_map, - recover_delete_store, recover_version_mgr, options); + Segment::Open(col_path_, *schema_, *persist_metas[0], recover_id_map, + recover_delete_store, recover_version_mgr, options_); ASSERT_TRUE(result.has_value()); segment = std::move(result).value(); ASSERT_TRUE(segment != nullptr); @@ -2473,7 +2191,7 @@ TEST_P(SegmentTest, AlterColumn) { // std::string column_name = "int32"; for (auto &dest_column : test_alter_columns) { if (column_name == dest_column) continue; - auto field_schema = schema->get_field(dest_column); + auto field_schema = schema_->get_field(dest_column); auto new_field_schema = std::make_shared(*field_schema); new_field_schema->set_name(column_name); func(column_name, new_field_schema); @@ -2485,7 +2203,7 @@ TEST_P(SegmentTest, DropColumn) { // create segment int doc_count = 1000; auto segment = test::TestHelper::CreateSegmentWithDoc( - col_path, *schema, 0, 0, id_map, delete_store, version_manager, options, + col_path_, *schema_, 0, 0, id_map_, delete_store_, version_manager_, options_, 0, doc_count); ASSERT_TRUE(segment != nullptr); @@ -2496,50 +2214,50 @@ TEST_P(SegmentTest, DropColumn) { auto writing_segment_meta = segment->meta(); // convert writing segment meta to persisted segment meta - Version version = version_manager->get_current_version(); + Version version = version_manager_->get_current_version(); writing_segment_meta->remove_writing_forward_block(); s = version.add_persisted_segment_meta(writing_segment_meta); ASSERT_TRUE(s.ok()); - s = version_manager->apply(version); + s = version_manager_->apply(version); ASSERT_TRUE(s.ok()); - s = version_manager->flush(); + s = version_manager_->flush(); ASSERT_TRUE(s.ok()); segment.reset(); - version_manager.reset(); - id_map->flush(); - id_map.reset(); + version_manager_.reset(); + id_map_->flush(); + id_map_.reset(); std::string delete_store_path = - FileHelper::MakeFilePath(col_path, FileID::DELETE_FILE, 0); - delete_store->flush(delete_store_path); - delete_store.reset(); + FileHelper::MakeFilePath(col_path_, FileID::DELETE_FILE, 0); + delete_store_->flush(delete_store_path); + delete_store_.reset(); - auto recover_version_manager = VersionManager::Recovery(col_path); + auto recover_version_manager = VersionManager::Recovery(col_path_); auto recover_version_mgr = recover_version_manager.value(); ASSERT_TRUE(recover_version_mgr != nullptr); Version v = recover_version_mgr->get_current_version(); const auto &persist_metas = v.persisted_segment_metas(); // idmap - std::string idmap_path = FileHelper::MakeFilePath(col_path, FileID::ID_FILE, + std::string idmap_path = FileHelper::MakeFilePath(col_path_, FileID::ID_FILE, v.id_map_path_suffix()); - IDMap::Ptr recover_id_map = std::make_shared(col_name); + IDMap::Ptr recover_id_map = std::make_shared(col_name_); auto status = recover_id_map->open(idmap_path, false, false); ASSERT_TRUE(status.ok()); - delete_store_path = FileHelper::MakeFilePath(col_path, FileID::DELETE_FILE, + delete_store_path = FileHelper::MakeFilePath(col_path_, FileID::DELETE_FILE, v.delete_snapshot_path_suffix()); auto recover_delete_store = - DeleteStore::CreateAndLoad(col_name, delete_store_path); + DeleteStore::CreateAndLoad(col_name_, delete_store_path); ASSERT_TRUE(recover_delete_store != nullptr); // open persist segment - options.read_only_ = true; + options_.read_only_ = true; auto result = - Segment::Open(col_path, *schema, *persist_metas[0], recover_id_map, - recover_delete_store, recover_version_mgr, options); + Segment::Open(col_path_, *schema_, *persist_metas[0], recover_id_map, + recover_delete_store, recover_version_mgr, options_); ASSERT_TRUE(result.has_value()); segment = std::move(result).value(); ASSERT_TRUE(segment != nullptr); diff --git a/tests/db/index/segment/segment_test_fixture.h b/tests/db/index/segment/segment_test_fixture.h new file mode 100644 index 0000000..140f930 --- /dev/null +++ b/tests/db/index/segment/segment_test_fixture.h @@ -0,0 +1,105 @@ +// Copyright 2025-present the zvec project +// +// Licensed 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. + + +#pragma once + +#include +#include +#include +#include +#include +#include "db/common/file_helper.h" +#include "db/index/common/delete_store.h" +#include "db/index/common/id_map.h" +#include "db/index/common/version_manager.h" +#include "utils/utils.h" +#include "zvec/db/options.h" + + +namespace zvec { + + +class SegmentTest : public testing::TestWithParam { + protected: + void SetUp() override { + ailego::LoggerBroker::SetLevel(ailego::Logger::LEVEL_WARN); + zvec::ailego::MemoryLimitPool::get_instance().init(MIN_MEMORY_LIMIT_BYTES); + + FileHelper::RemoveDirectory(col_path_); + FileHelper::CreateDirectory(col_path_); + + auto id_map_path = FileHelper::MakeFilePath(col_path_, FileID::ID_FILE, 0); + id_map_ = IDMap::CreateAndOpen(col_name_, id_map_path, true, false); + if (id_map_ == nullptr) { + throw std::runtime_error("Failed to create id map"); + } + + delete_store_ = std::make_shared(col_name_); + + schema_ = + test::TestHelper::CreateSchemaWithScalarIndex(false, false, col_name_); + schema_->add_field( + std::make_shared("id", DataType::INT32, false)); + schema_->add_field( + std::make_shared("name", DataType::STRING, false)); + schema_->add_field( + std::make_shared("age", DataType::UINT32, false)); + schema_->add_field( + std::make_shared("binary", DataType::BINARY, false)); + schema_->add_field(std::make_shared( + "array_binary", DataType::ARRAY_BINARY, false)); + + bool enable_mmap = GetParam(); + + Version version; + version.set_schema(*schema_); + version.set_enable_mmap(enable_mmap); + auto version_manager_tmp = VersionManager::Create(col_path_, version); + if (!version_manager_tmp.has_value()) { + throw std::runtime_error("Failed to create version manager"); + } + + version_manager_ = version_manager_tmp.value(); + + options_.read_only_ = false; + options_.enable_mmap_ = enable_mmap; + options_.max_buffer_size_ = 64 * 1024 * 1024; + } + + void TearDown() override { + id_map_.reset(); + delete_store_.reset(); + version_manager_.reset(); + + FileHelper::RemoveDirectory(col_path_); + } + + public: + std::string GetColPath() { + return col_path_; + } + + protected: + std::string col_name_ = "test_segment"; + std::string col_path_ = "./test_collection"; + IDMap::Ptr id_map_; + DeleteStore::Ptr delete_store_; + VersionManager::Ptr version_manager_; + CollectionSchema::Ptr schema_; + SegmentOptions options_; +}; + + +} // namespace zvec diff --git a/tests/db/sqlengine/mock_segment.h b/tests/db/sqlengine/mock_segment.h index 4892b2c..ccdb619 100644 --- a/tests/db/sqlengine/mock_segment.h +++ b/tests/db/sqlengine/mock_segment.h @@ -333,7 +333,7 @@ class MockSegment : public Segment { } ExecBatchPtr fetch(const std::vector &columns, - int index) const override { + int segment_doc_id) const override { LOG_ERROR("Not implemented"); return nullptr; }