From 85cd08d3838d5ee3a635ce269ba50439c99e1f84 Mon Sep 17 00:00:00 2001 From: Jalin Wang Date: Wed, 17 Jun 2026 14:18:12 +0800 Subject: [PATCH] feat(storage): expose copy-on-write mmap flag and fix MMAP_POPULATE flag placement This PR exposes a copy-on-write mmap option through the public StorageOptions API and fixes a bug where the MAP_POPULATE flag was applied to the wrong mmap() argument. --- src/ailego/io/file.cc | 19 +- src/core/interface/index.cc | 11 + src/include/zvec/core/interface/index_param.h | 8 +- tests/core/interface/index_interface_test.cc | 263 +++++++++++++++++- 4 files changed, 290 insertions(+), 11 deletions(-) diff --git a/src/ailego/io/file.cc b/src/ailego/io/file.cc index fa7b36f..9b23d68 100644 --- a/src/ailego/io/file.cc +++ b/src/ailego/io/file.cc @@ -312,18 +312,20 @@ ssize_t File::offset(void) const { void *File::MemoryMap(NativeHandle handle, ssize_t off, size_t len, int opts) { int prot = ((opts & File::MMAP_READONLY) ? PROT_READ : PROT_READ | PROT_WRITE); + int flags = (opts & File::MMAP_SHARED) ? MAP_SHARED : MAP_PRIVATE; #if defined(MAP_POPULATE) if (opts & File::MMAP_POPULATE) { - prot |= MAP_POPULATE; + flags |= MAP_POPULATE; } #endif - int flags = (opts & File::MMAP_SHARED) ? MAP_SHARED : MAP_PRIVATE; + #if defined(MAP_HUGETLB) if (opts & File::MMAP_HUGE_PAGE) { flags |= MAP_HUGETLB; } #endif + void *addr = mmap(nullptr, len, prot, flags, handle, off); ailego_null_if_false(addr != MAP_FAILED); @@ -344,14 +346,13 @@ void *File::MemoryMap(size_t len, int opts) { #if defined(MAP_ANONYMOUS) int prot = ((opts & File::MMAP_READONLY) ? PROT_READ : PROT_READ | PROT_WRITE); - -#if defined(MAP_POPULATE) - if (opts & File::MMAP_POPULATE) { - prot |= MAP_POPULATE; - } -#endif int flags = (opts & File::MMAP_SHARED) ? MAP_SHARED | MAP_ANONYMOUS : MAP_PRIVATE | MAP_ANONYMOUS; +#if defined(MAP_POPULATE) + if (opts & File::MMAP_POPULATE) { + flags |= MAP_POPULATE; + } +#endif #if defined(MAP_HUGETLB) if (opts & File::MMAP_HUGE_PAGE) { flags |= MAP_HUGETLB; @@ -776,4 +777,4 @@ void File::MemoryWarmup(void *addr, size_t len) { } } // namespace ailego -} // namespace zvec \ No newline at end of file +} // namespace zvec diff --git a/src/core/interface/index.cc b/src/core/interface/index.cc index 00e4b4b..332d452 100644 --- a/src/core/interface/index.cc +++ b/src/core/interface/index.cc @@ -17,6 +17,7 @@ #include #include #include "mixed_reducer/mixed_reducer_params.h" +#include "utility/utility_params.h" namespace zvec::core_interface { @@ -255,6 +256,16 @@ int Index::Open(const std::string &file_path, StorageOptions storage_options) { // storage_params.set("proxima.mmap_file.storage.memory_warmup", true); // storage_params.set("proxima.mmap_file.storage.segment_meta_capacity", // 1024); + + storage_params.set(core::MMAPFILE_STORAGE_COPY_ON_WRITE, + storage_options.copy_on_write); + // force_flush must be enabled with copy_on_write so that + // IndexMapping::flush() actually persists data written via file_.write() to + // disk; without it, data would be lost. See IndexMapping::flush() in + // index_mapping.cc. + storage_params.set(core::MMAPFILE_STORAGE_FORCE_FLUSH, + storage_options.copy_on_write); + switch (storage_options.type) { case StorageOptions::StorageType::kMMAP: { storage_ = core::IndexFactory::CreateStorage("MMapFileStorage"); diff --git a/src/include/zvec/core/interface/index_param.h b/src/include/zvec/core/interface/index_param.h index 2a91cd1..186c160 100644 --- a/src/include/zvec/core/interface/index_param.h +++ b/src/include/zvec/core/interface/index_param.h @@ -41,6 +41,12 @@ struct StorageOptions { StorageType type = StorageType::kNone; bool create_new = false; bool read_only = false; + + // Only meaningful when type == kMMAP. + // false: MAP_SHARED. Writes through mmap auto-persist to the file. + // true : MAP_PRIVATE on a writable file. Flush/close forces dirty pages + // back to disk via explicit pwrite. + bool copy_on_write = false; }; struct MergeOptions { @@ -444,4 +450,4 @@ struct DiskAnnIndexParam : public BaseIndexParam { bool omit_empty_value = false) const override; }; -} // namespace zvec::core_interface \ No newline at end of file +} // namespace zvec::core_interface diff --git a/tests/core/interface/index_interface_test.cc b/tests/core/interface/index_interface_test.cc index 85671d5..5fcfe37 100644 --- a/tests/core/interface/index_interface_test.cc +++ b/tests/core/interface/index_interface_test.cc @@ -186,6 +186,267 @@ TEST(IndexInterface, General) { .build()); } +TEST(IndexInterface, CopyOnWrite) { + constexpr uint32_t kDimension = 64; + constexpr uint32_t kNumVectors = 50; + const std::string index_name{"test_cow.index"}; + + auto make_vec = [&](uint32_t seed) { + std::vector v(kDimension, 0.0f); + v[seed % kDimension] = 1.0f; + return v; + }; + + auto func = [&](const BaseIndexParam::Pointer ¶m, + const BaseIndexQueryParam::Pointer &query_param) { + zvec::test_util::RemoveTestFiles(index_name); + + // Phase 1: build the index with shared mmap (writeable shared mapping) + // since the COW mode isn't used as the initial ingest path here. + { + auto index = IndexFactory::CreateAndInitIndex(*param); + ASSERT_NE(nullptr, index); + ASSERT_EQ( + 0, index->Open(index_name, {StorageOptions::StorageType::kMMAP, + /*create_new=*/true, /*read_only=*/false, + /*copy_on_write=*/false})); + + std::vector> vecs; + vecs.reserve(kNumVectors); + for (uint32_t i = 0; i < kNumVectors; ++i) { + vecs.emplace_back(make_vec(i)); + VectorData vd; + vd.vector = DenseVector{vecs.back().data()}; + ASSERT_EQ(0, index->Add(vd, /*key=*/100 + i)); + } + ASSERT_EQ(0, index->Train()); + ASSERT_EQ(0, index->Close()); + } + + // Phase 2: reopen with COW mmap. Search and Fetch must succeed against + // the persisted file. + { + auto index = IndexFactory::CreateAndInitIndex(*param); + ASSERT_NE(nullptr, index); + ASSERT_EQ( + 0, index->Open(index_name, {StorageOptions::StorageType::kMMAP, + /*create_new=*/false, /*read_only=*/true, + /*copy_on_write=*/true})); + + for (uint32_t i = 0; i < kNumVectors; ++i) { + auto target = make_vec(i); + VectorData query; + query.vector = DenseVector{target.data()}; + SearchResult result; + ASSERT_EQ(0, index->Search(query, query_param, &result)); + ASSERT_FALSE(result.doc_list_.empty()); + ASSERT_EQ(100u + i, result.doc_list_[0].key()); + + VectorDataBuffer fetched; + ASSERT_EQ(0, index->Fetch(100 + i, &fetched)); + auto *fetched_ptr = reinterpret_cast( + std::get(fetched.vector_buffer).data.data()); + ASSERT_FLOAT_EQ(1.0f, fetched_ptr[i % kDimension]); + } + ASSERT_EQ(0, index->Close()); + } + + // Phase 3: reopen with shared mmap to confirm the file is intact after + // the COW session. + { + auto index = IndexFactory::CreateAndInitIndex(*param); + ASSERT_NE(nullptr, index); + ASSERT_EQ( + 0, index->Open(index_name, {StorageOptions::StorageType::kMMAP, + /*create_new=*/false, /*read_only=*/true, + /*copy_on_write=*/false})); + + auto target = make_vec(13); + VectorData query; + query.vector = DenseVector{target.data()}; + SearchResult result; + ASSERT_EQ(0, index->Search(query, query_param, &result)); + ASSERT_FALSE(result.doc_list_.empty()); + ASSERT_EQ(113u, result.doc_list_[0].key()); + ASSERT_EQ(0, index->Close()); + } + + // Phase 4: repeated open/close under COW mmap must not lose entries. + for (int cycle = 0; cycle < 3; ++cycle) { + auto index = IndexFactory::CreateAndInitIndex(*param); + ASSERT_NE(nullptr, index); + ASSERT_EQ( + 0, index->Open(index_name, {StorageOptions::StorageType::kMMAP, + /*create_new=*/false, /*read_only=*/true, + /*copy_on_write=*/true})); + uint32_t i = static_cast(cycle * 5 + 2); + auto target = make_vec(i); + VectorData query; + query.vector = DenseVector{target.data()}; + SearchResult result; + ASSERT_EQ(0, index->Search(query, query_param, &result)); + ASSERT_FALSE(result.doc_list_.empty()); + ASSERT_EQ(100u + i, result.doc_list_[0].key()); + ASSERT_EQ(0, index->Close()); + } + + // Phase 5: open in COW mmap (writable MAP_PRIVATE with forced flush). + // Without performing writes the close path still exercises the pwrite + // branch with no dirty pages, which must not corrupt the file. + { + auto index = IndexFactory::CreateAndInitIndex(*param); + ASSERT_NE(nullptr, index); + ASSERT_EQ( + 0, index->Open(index_name, {StorageOptions::StorageType::kMMAP, + /*create_new=*/false, /*read_only=*/true, + /*copy_on_write=*/true})); + + auto target = make_vec(21); + VectorData query; + query.vector = DenseVector{target.data()}; + SearchResult result; + ASSERT_EQ(0, index->Search(query, query_param, &result)); + ASSERT_FALSE(result.doc_list_.empty()); + ASSERT_EQ(121u, result.doc_list_[0].key()); + ASSERT_EQ(0, index->Close()); + } + + // Phase 6: reopen with shared mmap to confirm Phase 5's open/close left + // the file intact. + { + auto index = IndexFactory::CreateAndInitIndex(*param); + ASSERT_NE(nullptr, index); + ASSERT_EQ( + 0, index->Open(index_name, {StorageOptions::StorageType::kMMAP, + /*create_new=*/false, /*read_only=*/true, + /*copy_on_write=*/false})); + for (uint32_t i = 0; i < kNumVectors; ++i) { + auto target = make_vec(i); + VectorData query; + query.vector = DenseVector{target.data()}; + SearchResult result; + ASSERT_EQ(0, index->Search(query, query_param, &result)); + ASSERT_FALSE(result.doc_list_.empty()); + ASSERT_EQ(100u + i, result.doc_list_[0].key()); + } + ASSERT_EQ(0, index->Close()); + } + + zvec::test_util::RemoveTestFiles(index_name); + }; + + func(FlatIndexParamBuilder() + .WithMetricType(MetricType::kInnerProduct) + .WithDataType(DataType::DT_FP32) + .WithDimension(kDimension) + .WithIsSparse(false) + .Build(), + FlatQueryParamBuilder().with_topk(5).with_fetch_vector(false).build()); + + func(HNSWIndexParamBuilder() + .WithMetricType(MetricType::kInnerProduct) + .WithDataType(DataType::DT_FP32) + .WithDimension(kDimension) + .WithIsSparse(false) + .WithEFConstruction(100) + .Build(), + HNSWQueryParamBuilder() + .with_topk(5) + .with_fetch_vector(false) + .with_ef_search(20) + .build()); + + func(VamanaIndexParamBuilder() + .WithMetricType(MetricType::kInnerProduct) + .WithDataType(DataType::DT_FP32) + .WithDimension(kDimension) + .WithIsSparse(false) + .WithMaxDegree(32) + .WithSearchListSize(64) + .WithAlpha(1.2f) + .Build(), + VamanaQueryParamBuilder() + .with_topk(5) + .with_fetch_vector(false) + .with_ef_search(32) + .build()); + + // Flat-only durability check for COW mmap: writes performed under + // MAP_PRIVATE must be pwrite-flushed back and visible after a shared-mmap + // reopen. Flat is used because Add/Flush against a previously-built file is + // straightforward to reason about for this storage layer. + { + const std::string persist_index{"test_cow_persist.index"}; + zvec::test_util::RemoveTestFiles(persist_index); + auto persist_param = FlatIndexParamBuilder() + .WithMetricType(MetricType::kInnerProduct) + .WithDataType(DataType::DT_FP32) + .WithDimension(kDimension) + .WithIsSparse(false) + .Build(); + auto persist_query = + FlatQueryParamBuilder().with_topk(5).with_fetch_vector(false).build(); + + { + auto index = IndexFactory::CreateAndInitIndex(*persist_param); + ASSERT_NE(nullptr, index); + ASSERT_EQ(0, index->Open(persist_index, + {StorageOptions::StorageType::kMMAP, + /*create_new=*/true, /*read_only=*/false, + /*copy_on_write=*/false})); + auto v0 = make_vec(0); + VectorData vd; + vd.vector = DenseVector{v0.data()}; + ASSERT_EQ(0, index->Add(vd, /*key=*/500)); + ASSERT_EQ(0, index->Train()); + ASSERT_EQ(0, index->Close()); + } + + // Add a new vector through COW mmap and explicitly Flush so + // dirty private pages are written back to the file. + { + auto index = IndexFactory::CreateAndInitIndex(*persist_param); + ASSERT_NE(nullptr, index); + ASSERT_EQ(0, index->Open(persist_index, + {StorageOptions::StorageType::kMMAP, + /*create_new=*/false, /*read_only=*/false, + /*copy_on_write=*/true})); + auto v1 = make_vec(1); + VectorData vd; + vd.vector = DenseVector{v1.data()}; + ASSERT_EQ(0, index->Add(vd, /*key=*/501)); + ASSERT_EQ(0, index->Flush()); + ASSERT_EQ(0, index->Close()); + } + + // Reopen with shared mmap: the entry written in COW mode must be durable + // on disk. + { + auto index = IndexFactory::CreateAndInitIndex(*persist_param); + ASSERT_NE(nullptr, index); + ASSERT_EQ(0, index->Open(persist_index, + {StorageOptions::StorageType::kMMAP, + /*create_new=*/false, /*read_only=*/true, + /*copy_on_write=*/false})); + auto target = make_vec(1); + VectorData query; + query.vector = DenseVector{target.data()}; + SearchResult result; + ASSERT_EQ(0, index->Search(query, persist_query, &result)); + ASSERT_FALSE(result.doc_list_.empty()); + ASSERT_EQ(501u, result.doc_list_[0].key()); + + VectorDataBuffer fetched; + ASSERT_EQ(0, index->Fetch(501, &fetched)); + auto *fetched_ptr = reinterpret_cast( + std::get(fetched.vector_buffer).data.data()); + ASSERT_FLOAT_EQ(1.0f, fetched_ptr[1 % kDimension]); + ASSERT_EQ(0, index->Close()); + } + zvec::test_util::RemoveTestFiles(persist_index); + } +} + TEST(IndexInterface, BufferGeneral) { zvec::ailego::MemoryLimitPool::get_instance().init(100 * 1024 * 1024); constexpr uint32_t kDimension = 64; @@ -2121,4 +2382,4 @@ TEST(IndexInterface, IsDirtyBufferPool) { #if defined(__GNUC__) || defined(__GNUG__) #pragma GCC diagnostic pop -#endif \ No newline at end of file +#endif