Fixed local data sources growing on disk with every change (#1034)

This commit is contained in:
Thorsten Sommer authored and GitHub committed 2026-10-09 20:48:27 +02:00
1 parent 56215aa6dd
commit 5bd3bb08df
4 files changed
+263 -1

No files matched your search

+217
View File
@@ -32,6 +32,8 @@ const STORE_INITIALIZATION_MARKER: &str = "store_name.txt";
const STORE_INITIALIZATION_MARKER_TEMP: &str = "store_name.tmp";
const STORE_DISPLAY_NAME_MARKER: &str = "data_source_name.txt";
const STORE_DISPLAY_NAME_MARKER_TEMP: &str = "data_source_name.tmp";
const WAL_DIRECTORY: &str = "wal";
const CLOSED_WAL_SEGMENT_PREFIX: &str = "closed-";
/// Marks a response whose store exists on disk but cannot be opened. The .NET side keys its repair
/// offer off this value instead of parsing `issue`, so rewording the message stays harmless.
@@ -241,6 +243,7 @@ impl QdrantEdgeDatabase {
}
let shard = if is_initialized {
compact_wal(&path, store_name);
match EdgeShard::load(&path, None) {
Ok(shard) => shard,
Err(error) => return Err(StoreUnreadableError::new(store_name, &path, error).into()),
@@ -282,6 +285,7 @@ impl QdrantEdgeDatabase {
return Ok(None);
}
compact_wal(&path, store_name);
let shard = match EdgeShard::load(&path, None) {
Ok(shard) => shard,
Err(error) => return Err(StoreUnreadableError::new(store_name, &path, error).into()),
@@ -424,6 +428,16 @@ impl QdrantEdgeDatabase {
info!(Source = "Qdrant Edge"; "Optimized vector store '{}'.", store_name);
}
shard.flush()?;
// Temporary compatibility shim, see compact_wal: without this, the WAL of a long session
// keeps growing until the next start. The store has to be unloaded first, because Windows
// cannot delete memory-mapped files. The next request loads it again.
let path = self.store_path(store_name)?;
if has_obsolete_wal_segments(&path) {
drop(self.shards.remove(store_name));
compact_wal(&path, store_name);
}
Ok(())
}
@@ -778,6 +792,79 @@ fn remove_partial_store(path: &Path) -> String {
}
}
// Temporary compatibility shim until a qdrant-edge release acknowledges its WAL by itself:
// documentation/compatibility-shims/2026-10-qdrant-edge-wal-compaction.md
//
// Qdrant Edge writes every update to its WAL, but never acknowledges the WAL and never replays it,
// so the WAL grows with every update. We flush after every update, which means the segments already
// hold everything the WAL contains. The store must not be loaded while this runs.
fn compact_wal(store_path: &Path, store_name: &str) {
match remove_obsolete_wal_segments(&store_path.join(WAL_DIRECTORY)) {
Ok((0, _)) => {},
Ok((removed_segments, removed_bytes)) => info!(Source = "Qdrant Edge"; "Removed {removed_segments} obsolete WAL segment(s) with {:.1} MiB from vector store '{store_name}'.", removed_bytes as f64 / (1024.0 * 1024.0)),
Err(error) => warn!(Source = "Qdrant Edge"; "Could not remove the obsolete WAL segments of vector store '{store_name}': {error}"),
}
}
/// Whether the store's WAL holds segments which `compact_wal` would remove.
fn has_obsolete_wal_segments(store_path: &Path) -> bool {
wal_file_names(&store_path.join(WAL_DIRECTORY))
.is_ok_and(|file_names| !obsolete_wal_segments(&file_names).is_empty())
}
/// Removes the obsolete closed segments, oldest first, and returns their number and size. When a
/// removal fails, the segments left over are still contiguous, so Qdrant Edge can open the WAL.
fn remove_obsolete_wal_segments(wal_path: &Path) -> std::io::Result<(usize, u64)> {
let mut removed_segments = 0;
let mut removed_bytes = 0;
for file_name in obsolete_wal_segments(&wal_file_names(wal_path)?) {
let segment_path = wal_path.join(file_name);
let segment_size = fs::metadata(&segment_path)?.len();
fs::remove_file(&segment_path)?;
removed_segments += 1;
removed_bytes += segment_size;
}
Ok((removed_segments, removed_bytes))
}
fn wal_file_names(wal_path: &Path) -> std::io::Result<Vec<String>> {
if !wal_path.is_dir() {
return Ok(Vec::new());
}
// Every file counts: a closed segment missing from this list would leave a gap behind.
let mut file_names = Vec::new();
for entry in fs::read_dir(wal_path)? {
let entry = entry?;
if entry.file_type()?.is_file() && let Ok(file_name) = entry.file_name().into_string() {
file_names.push(file_name);
}
}
Ok(file_names)
}
/// The closed WAL segments which may go, oldest first. That is all of them except the newest:
/// Qdrant Edge numbers its operations from the first closed segment and would start again at zero
/// without one. The segments then ignore every later update of a point they already know.
fn obsolete_wal_segments(file_names: &[String]) -> Vec<String> {
let mut closed_segments = file_names
.iter()
.filter_map(|file_name| {
let start_index = file_name.strip_prefix(CLOSED_WAL_SEGMENT_PREFIX)?.parse::<u64>().ok()?;
Some((start_index, file_name.clone()))
})
.collect::<Vec<_>>();
closed_segments.sort_unstable_by_key(|(start_index, _)| *start_index);
closed_segments.pop();
closed_segments
.into_iter()
.map(|(_, file_name)| file_name)
.collect()
}
fn validate_vector_size(vector_size: usize) -> QdrantEdgeResult<()> {
if vector_size == 0 {
return Err("Vector size must be greater than zero.".into());
@@ -946,6 +1033,7 @@ fn validate_store_name(store_name: &str) -> QdrantEdgeResult<()> {
#[cfg(test)]
mod tests {
use super::*;
use qdrant_edge::WalOptions;
#[test]
fn validate_store_name_allows_safe_store_names() {
@@ -1081,6 +1169,91 @@ mod tests {
);
}
#[test]
fn obsolete_wal_segments_keeps_the_newest_closed_segment() {
let file_names = ["closed-120", "open-7", "closed-0", "tmp-open-8", "closed-40", ".wal", "first-index", "closed-abc"].map(String::from);
// Sorted by number, not by name: as text, "closed-40" would come after "closed-120".
assert_eq!(obsolete_wal_segments(&file_names), vec!["closed-0", "closed-40"]);
assert!(obsolete_wal_segments(&["closed-0", "open-1"].map(String::from)).is_empty());
assert!(obsolete_wal_segments(&["open-1"].map(String::from)).is_empty());
}
#[test]
fn a_reloaded_store_drops_obsolete_wal_segments_and_keeps_its_points() {
let test_directory = test_directory("wal-reload");
let (database, point_ids) = store_with_several_closed_wal_segments(&test_directory);
let store_path = database.store_path(SEARCH_STORE).unwrap();
drop(database);
let mut database = QdrantEdgeDatabase::new(test_directory.clone());
let found = database.search_embedding(SEARCH_STORE, search_vector(), point_ids.len(), None).unwrap();
assert_eq!(closed_wal_segments(&store_path).len(), 1);
assert_eq!(
found_point_ids(&found).into_iter().collect::<HashSet<_>>(),
point_ids.iter().map(String::as_str).collect::<HashSet<_>>(),
"the points live in the segments, so removing WAL segments must not lose any of them"
);
drop(database);
fs::remove_dir_all(test_directory).unwrap();
}
#[test]
fn an_update_after_wal_compaction_replaces_the_existing_point() {
let test_directory = test_directory("wal-update");
let (database, point_ids) = store_with_several_closed_wal_segments(&test_directory);
drop(database);
// The newest point carries the highest operation number so far. Should the numbering start
// again after the compaction, the segment would ignore this update as outdated.
let newest_point_id = point_ids.last().unwrap();
let mut updated_point = test_point(newest_point_id, vec![0.0, 0.0, 1.0]);
updated_point.text = "Updated text.".to_string();
let mut database = QdrantEdgeDatabase::new(test_directory.clone());
database.insert_embedding(SEARCH_STORE, vec![updated_point]).unwrap();
let found = database
.search_embedding(SEARCH_STORE, search_vector(), 1, Some(vec![newest_point_id.clone()]))
.unwrap();
assert_eq!(found.len(), 1);
assert_eq!(found[0].text, "Updated text.");
drop(database);
fs::remove_dir_all(test_directory).unwrap();
}
#[test]
fn optimize_removes_obsolete_wal_segments_of_a_loaded_store() {
let test_directory = test_directory("wal-optimize");
let (mut database, point_ids) = store_with_several_closed_wal_segments(&test_directory);
let store_path = database.store_path(SEARCH_STORE).unwrap();
database.optimize_store(SEARCH_STORE).unwrap();
assert_eq!(closed_wal_segments(&store_path).len(), 1);
// The next request loads the store again, with all its points and still taking updates.
let newest_point_id = point_ids.last().unwrap();
let mut updated_point = test_point(newest_point_id, vec![0.0, 0.0, 1.0]);
updated_point.text = "Updated text.".to_string();
database.insert_embedding(SEARCH_STORE, vec![updated_point]).unwrap();
let found = database.search_embedding(SEARCH_STORE, search_vector(), point_ids.len(), None).unwrap();
assert_eq!(
found_point_ids(&found).into_iter().collect::<HashSet<_>>(),
point_ids.iter().map(String::as_str).collect::<HashSet<_>>()
);
let updated = found.iter().find(|result| &result.point_id == newest_point_id).unwrap();
assert_eq!(updated.text, "Updated text.");
drop(database);
fs::remove_dir_all(test_directory).unwrap();
}
const SEARCH_STORE: &str = "rag_6cc665a82b1e4d42bc748015b7b391ec";
const POINT_X: &str = "0b5f1e8a-3c2d-4e6f-9a1b-7c8d9e0f1a2b";
const POINT_Y: &str = "1c6a2f9b-4d3e-4f70-8b2c-8d9e0f1a2b3c";
@@ -1113,6 +1286,50 @@ mod tests {
database
}
/// A store with tiny WAL segments, filled one point at a time until the WAL has closed several
/// of them. Returns the database with the store still loaded, and the IDs in insertion order.
fn store_with_several_closed_wal_segments(test_directory: &Path) -> (QdrantEdgeDatabase, Vec<String>) {
const SMALL_WAL_SEGMENT_CAPACITY: usize = 8 * 1024;
const MAX_POINTS: usize = 1_000;
let mut database = QdrantEdgeDatabase::new(test_directory.to_path_buf());
let store_path = database.store_path(SEARCH_STORE).unwrap();
fs::create_dir_all(&store_path).unwrap();
// The production config, except for the WAL. Qdrant Edge keeps these options in the store,
// so they stay in effect whenever the store gets loaded again.
let config = EdgeConfig {
wal_options: Some(WalOptions {
segment_capacity: SMALL_WAL_SEGMENT_CAPACITY,
..WalOptions::default()
}),
..edge_config(3)
};
drop(EdgeShard::new(&store_path, config).unwrap());
write_store_initialization_marker(&store_path, SEARCH_STORE).unwrap();
let mut point_ids = Vec::new();
while closed_wal_segments(&store_path).len() < 3 {
assert!(point_ids.len() < MAX_POINTS, "the WAL never closed enough segments");
let point_id = Uuid::from_u128(point_ids.len() as u128 + 1).to_string();
database
.insert_embedding(SEARCH_STORE, vec![test_point(&point_id, vec![1.0, point_ids.len() as f32, 1.0])])
.unwrap();
point_ids.push(point_id);
}
(database, point_ids)
}
fn closed_wal_segments(store_path: &Path) -> Vec<String> {
fs::read_dir(store_path.join(WAL_DIRECTORY))
.unwrap()
.map(|entry| entry.unwrap().file_name().into_string().unwrap())
.filter(|file_name| file_name.starts_with(CLOSED_WAL_SEGMENT_PREFIX))
.collect()
}
/// Closest to the point on the x axis, then the one on the y axis, then the one on the z axis.
fn search_vector() -> Vec<f32> {
vec![1.0, 0.5, 0.0]