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.
This commit is contained in:
parent
e12bad5e37
commit
85cd08d383
|
|
@ -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
|
||||
} // namespace zvec
|
||||
|
|
|
|||
|
|
@ -17,6 +17,7 @@
|
|||
#include <zvec/core/framework/index_storage.h>
|
||||
#include <zvec/core/interface/index.h>
|
||||
#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");
|
||||
|
|
|
|||
|
|
@ -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
|
||||
} // namespace zvec::core_interface
|
||||
|
|
|
|||
|
|
@ -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<float> 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<std::vector<float>> 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<const float *>(
|
||||
std::get<DenseVectorBuffer>(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<uint32_t>(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<const float *>(
|
||||
std::get<DenseVectorBuffer>(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
|
||||
#endif
|
||||
|
|
|
|||
Loading…
Reference in New Issue