Merge branch 'segment-overflow-2-cap' into segment-overflow-3-provisioning

This commit is contained in:
generall
2026-07-16 15:41:04 +02:00
2 changed files with 93 additions and 88 deletions

View File

@@ -687,25 +687,48 @@ impl SegmentHolder {
}
}
/// Filter CoW-move destination candidates down to appendable segments below the configured
/// size cap.
/// Appendable segments eligible as write destinations under the configured size cap, with
/// their measured sizes.
///
/// Errors with [`OperationError::OutOfAppendableCapacity`] when candidates exist but all of
/// them reached the cap; re-applying the operation after provisioning a fresh appendable
/// segment resumes where it left off (already-moved points are skipped by point version).
fn candidates_with_capacity(
/// Segments that cannot be measured right now (read-lock contention, size errors) stay
/// eligible and report `None`: the cap is a soft limit and liveness wins over precision.
/// With no cap configured, every appendable segment is eligible.
///
/// This is the single origin of [`OperationError::OutOfAppendableCapacity`], returned when
/// appendable segments exist but every measurable one reached the cap. Callers recover by
/// provisioning a fresh appendable segment and re-applying the operation; already-applied
/// points are then skipped by their point version.
fn appendable_segments_with_capacity(
&self,
candidates: &[SegmentId],
) -> OperationResult<Vec<SegmentId>> {
let eligible: Vec<_> = candidates
.iter()
.copied()
.filter(|segment_id| {
self.segment_has_capacity(*segment_id, self.max_segment_size_bytes)
})
.collect();
) -> OperationResult<Vec<(SegmentId, Option<usize>)>> {
let appendable_segments = self.appendable_segments_ids();
let mut eligible = Vec::with_capacity(appendable_segments.len());
for segment_id in &appendable_segments {
let Some(locked_segment) = self.get(*segment_id) else {
continue;
};
let size = match locked_segment.get().try_read() {
None => None,
Some(segment) => match segment.max_available_vectors_size_in_bytes() {
Ok(size) => Some(size),
Err(err) => {
log::error!("Failed to get segment size, ignoring: {err}");
None
}
},
};
let under_cap = match (size, self.max_segment_size_bytes) {
(Some(size), Some(max_segment_size_bytes)) => size < max_segment_size_bytes.get(),
_ => true,
};
if under_cap {
eligible.push((*segment_id, size));
}
}
if eligible.is_empty()
&& !candidates.is_empty()
&& !appendable_segments.is_empty()
&& let Some(max_segment_size_bytes) = self.max_segment_size_bytes
{
return Err(OperationError::OutOfAppendableCapacity {
@@ -736,48 +759,31 @@ impl SegmentHolder {
///
/// # Errors
///
/// - [`OperationError::OutOfAppendableCapacity`] when a size cap is configured and even the
/// smallest segment reached it; re-applying the operation after provisioning a fresh
/// appendable segment resumes where it left off.
/// - [`OperationError::OutOfAppendableCapacity`] when a size cap is configured and every
/// measurable segment reached it (see
/// [`appendable_segments_with_capacity`](Self::appendable_segments_with_capacity)).
/// - Service error when there are no appendable segments at all.
///
/// If capacity is not important use `random_appendable_segment` instead because it is cheaper.
pub fn smallest_appendable_segment(&self) -> OperationResult<LockedSegment> {
let segment_ids: Vec<_> = self.appendable_segments_ids();
let candidates = self.appendable_segments_with_capacity()?;
// Try a non-blocking read lock on all segments and return the smallest one
let smallest_segment = segment_ids
// Prefer the smallest measurable candidate; fall back to a random one when none could
// be measured.
let chosen = candidates
.iter()
.filter_map(|segment_id| self.get(*segment_id))
.filter_map(|locked_segment| {
match locked_segment
.get()
.try_read()
.map(|segment| segment.max_available_vectors_size_in_bytes())?
{
Ok(size) => Some((locked_segment, size)),
Err(err) => {
log::error!("Failed to get segment size, ignoring: {err}");
None
}
}
})
.min_by_key(|(_, segment_size)| *segment_size);
if let Some((smallest_segment, size)) = smallest_segment {
if let Some(max_segment_size_bytes) = self.max_segment_size_bytes
&& size >= max_segment_size_bytes.get()
{
return Err(OperationError::OutOfAppendableCapacity {
max_segment_size_bytes: max_segment_size_bytes.get(),
});
}
return Ok(LockedSegment::clone(smallest_segment));
}
.filter_map(|(segment_id, size)| size.map(|size| (*segment_id, size)))
.min_by_key(|(_, size)| *size)
.map(|(segment_id, _)| segment_id)
.or_else(|| {
candidates
.choose(&mut rand::rng())
.map(|(segment_id, _)| *segment_id)
});
// Fall back to picking a random segment
segment_ids
.choose(&mut rand::rng())
.and_then(|idx| self.appendable_segments.get(idx).cloned())
chosen
.and_then(|segment_id| self.appendable_segments.get(&segment_id))
.cloned()
.ok_or_else(|| {
OperationError::service_error("No appendable segments exist, expected at least one")
})
@@ -1116,6 +1122,13 @@ impl SegmentHolder {
// Choose random appendable segment from this
let appendable_segments = self.appendable_segments_ids();
// CoW-move destination candidates below the size cap. Computed lazily on the first point
// that actually needs a move — a batch that applies fully in place must not fail just
// because all appendable segments are full — and reused for the rest of this call.
// Callers batch points (e.g. 32 per call), so capacity is checked once per batch, not
// per point; a destination may overshoot the cap by at most one batch worth of points.
let mut eligible_segments: Option<Vec<SegmentId>> = None;
let mut applied_points: AHashSet<PointIdType> = Default::default();
let stopped = AtomicBool::new(false);
@@ -1132,14 +1145,17 @@ impl SegmentHolder {
let is_applied = if can_apply_operation {
point_operation(point_id, write_segment)?
} else {
// Re-check destination capacity per moved point: a large operation can fill the
// destination mid-way, and the next point must then go elsewhere (or the
// operation must fail with `OutOfAppendableCapacity` so the caller provisions a
// fresh appendable segment and re-applies).
let eligible_segments;
let destination_candidates = if self.max_segment_size_bytes.is_some() {
eligible_segments = self.candidates_with_capacity(&appendable_segments)?;
&eligible_segments
let destination_candidates: &[SegmentId] = if self.max_segment_size_bytes.is_some()
{
if eligible_segments.is_none() {
let eligible = self
.appendable_segments_with_capacity()?
.into_iter()
.map(|(segment_id, _size)| segment_id)
.collect();
eligible_segments = Some(eligible);
}
eligible_segments.as_ref().unwrap()
} else {
&appendable_segments
};

View File

@@ -200,26 +200,6 @@ pub fn process_field_index_operation(
/// parallel read operations are not starved.
const UPDATE_OP_CHUNK_SIZE: usize = 32;
/// Guard for the insert loops: with a size cap configured on the holder, refuse to insert into a
/// destination segment that already reached it. The resulting
/// [`OperationError::OutOfAppendableCapacity`] lets the caller provision a fresh appendable
/// segment and re-apply; points inserted so far are then skipped by their point version.
fn check_insert_capacity(
segments: &SegmentHolder,
write_segment: &RwLockWriteGuard<dyn SegmentEntry>,
) -> OperationResult<()> {
let Some(max_segment_size_bytes) = segments.max_segment_size_bytes() else {
return Ok(());
};
let size = write_segment.max_available_vectors_size_in_bytes()?;
if size >= max_segment_size_bytes.get() {
return Err(OperationError::OutOfAppendableCapacity {
max_segment_size_bytes: max_segment_size_bytes.get(),
});
}
Ok(())
}
/// Checks point id in each segment, update point if found.
/// All not found points are inserted into random segment.
/// Returns: number of updated points.
@@ -279,7 +259,6 @@ where
let segment_arc = default_write_segment.get();
let mut write_segment = segment_arc.write();
for point_id in new_point_ids {
check_insert_capacity(segments, &write_segment)?;
let point = points_map[&point_id];
res += usize::from(upsert_with_payload(
&mut write_segment,
@@ -697,7 +676,6 @@ where
let segment_arc = default_write_segment.get();
let mut write_segment = segment_arc.write();
for point_id in new_point_ids {
check_insert_capacity(segments, &write_segment)?;
res += usize::from(upsert_raw_with_payload(
&mut write_segment,
op_num,
@@ -2251,16 +2229,28 @@ mod test {
let mut non_appendable = build_segment_1(dir.path()); // points 1-5
non_appendable.appendable_flag = false;
let appendable = empty_segment(dir.path());
// Fill the only appendable segment up to the cap of 2 points.
let mut appendable = empty_segment(dir.path());
for point_id in [100u64, 101] {
appendable
.upsert_point(
10,
point_id.into(),
only_default_vector(&[1.0, 0.0, 1.0, 0.0]),
&hw_counter,
)
.unwrap();
}
let mut holder = SegmentHolder::default();
holder.add_new(non_appendable);
holder.add_new(appendable);
holder.set_max_segment_size_bytes(NonZeroUsize::new(2 * TEST_POINT_SIZE_BYTES));
// Setting payload on all 5 points CoW-moves them out of the non-appendable segment, but
// only 2 fit under the cap: the operation must fail with the recoverable capacity error
// instead of overgrowing the destination.
// Setting payload on the points of the non-appendable segment needs to CoW-move them,
// but no appendable segment is below the cap: the operation must fail with the
// recoverable capacity error instead of growing the full destination further.
let points: Vec<PointIdType> = (1..=5).map(Into::into).collect();
let err = set_payload(
&holder,
@@ -2270,14 +2260,13 @@ mod test {
&None,
&hw_counter,
)
.expect_err("the appendable segment cannot hold all 5 moved points");
.expect_err("no appendable segment below the cap can accept the moved points");
assert!(
err.is_out_of_appendable_capacity(),
"expected OutOfAppendableCapacity, got: {err}",
);
// The destination filled up to the cap exactly; the insert path now reports the same
// error instead of writing past it.
// The insert path reports the same error instead of writing past the cap.
let err = holder
.smallest_appendable_segment()
.expect_err("even the smallest appendable segment is at the cap");