pub mod segment_entry;
pub mod snapshot_entry;
mod vector_name_changes;
#[cfg(test)]
mod tests;
use std::borrow::Cow;
use ahash::AHashMap;
use crate::common::bitvec::BitVec;
use crate::common::counter::hardware_counter::HardwareCounterCell;
use crate::common::types::PointOffsetType;
use itertools::Itertools as _;
use crate::segment::common::operation_error::OperationResult;
use crate::segment::types::*;
pub use self::vector_name_changes::{IntendedVector, ProxyVectorNameChanges};
use crate::shard::locked_segment::LockedSegment;
pub type DeletedPoints = AHashMap<PointIdType, ProxyDeletedPoint>;
#[derive(Debug)]
pub struct ProxySegment {
pub wrapped_segment: LockedSegment,
deleted_mask: Option<BitVec>,
changed_indexes: ProxyIndexChanges,
changed_vector_names: ProxyVectorNameChanges,
deleted_points: DeletedPoints,
deleted_deferred_count: usize,
wrapped_config: SegmentConfig,
version: SeqNumberType,
}
#[must_use = "an UnsyncedProxySegment must be turned into a ProxySegment via `.finalize()`"]
#[derive(Debug)]
pub struct UnsyncedProxySegment(ProxySegment);
impl UnsyncedProxySegment {
pub fn new(segment: LockedSegment) -> Self {
if matches!(segment, LockedSegment::Proxy(_)) {
log::debug!("Double proxy segment creation");
}
let (wrapped_config, version) = {
let read_segment = segment.get().read();
(read_segment.config().clone(), read_segment.version())
};
UnsyncedProxySegment(ProxySegment {
wrapped_segment: segment,
deleted_mask: None,
changed_indexes: ProxyIndexChanges::default(),
changed_vector_names: ProxyVectorNameChanges::default(),
deleted_points: AHashMap::new(),
deleted_deferred_count: 0,
wrapped_config,
version,
})
}
pub fn finalize(mut self) -> ProxySegment {
self.0.sync_deleted_mask();
self.0
}
pub fn wrapped_segment(&self) -> &LockedSegment {
&self.0.wrapped_segment
}
pub fn replicate_field_indexes(
&self,
op_num: SeqNumberType,
hw_counter: &HardwareCounterCell,
segment_to_update: &LockedSegment,
) -> OperationResult<()> {
self.0
.replicate_field_indexes(op_num, hw_counter, segment_to_update)
}
}
impl ProxySegment {
#[cfg(feature = "testing")]
pub fn new(segment: LockedSegment) -> Self {
UnsyncedProxySegment::new(segment).finalize()
}
fn sync_deleted_mask(&mut self) {
match &self.wrapped_segment {
LockedSegment::Original(raw_segment) => {
self.deleted_mask = Some(raw_segment.read().get_deleted_points_bitvec());
}
LockedSegment::Proxy(_) => {
}
}
}
pub fn replicate_field_indexes(
&self,
op_num: SeqNumberType,
hw_counter: &HardwareCounterCell,
segment_to_update: &LockedSegment,
) -> OperationResult<()> {
let existing_indexes = segment_to_update.get().read().get_indexed_fields();
let expected_indexes = self.wrapped_segment.get().read().get_indexed_fields();
for (expected_field, expected_schema) in &expected_indexes {
let existing_schema = existing_indexes.get(expected_field);
if existing_schema != Some(expected_schema) {
if existing_schema.is_some() {
segment_to_update
.get()
.write()
.delete_field_index(op_num, expected_field)?;
}
segment_to_update.get().write().create_field_index(
op_num,
expected_field,
Some(expected_schema),
hw_counter,
)?;
}
}
for existing_field in existing_indexes.keys() {
if !expected_indexes.contains_key(existing_field) {
segment_to_update
.get()
.write()
.delete_field_index(op_num, existing_field)?;
}
}
Ok(())
}
fn set_deleted_offset(&mut self, point_offset: Option<PointOffsetType>) -> bool {
match (&mut self.deleted_mask, point_offset) {
(Some(deleted_mask), Some(point_offset)) => {
if deleted_mask.len() <= point_offset as usize {
deleted_mask.resize(point_offset as usize + 1, false);
}
deleted_mask.set(point_offset as usize, true);
true
}
_ => false,
}
}
fn add_deleted_points_condition_to_filter(
filter: Option<Cow<'_, Filter>>,
deleted_points: impl IntoIterator<Item = PointIdType>,
) -> Filter {
let wrapper_condition = Condition::HasId(HasIdCondition::from_iter(deleted_points));
match filter {
None => Filter::new_must_not(wrapper_condition),
Some(f) => {
let mut new_filter = f.into_owned();
let new_must_not = match new_filter.must_not {
None => Some(vec![wrapper_condition]),
Some(mut conditions) => {
conditions.push(wrapper_condition);
Some(conditions)
}
};
new_filter.must_not = new_must_not;
new_filter
}
}
}
pub fn propagate_to_wrapped(&mut self) -> OperationResult<()> {
let wrapped_segment = self.wrapped_segment.get();
let mut wrapped_segment = wrapped_segment.upgradable_read();
{
let op_num = wrapped_segment.version();
if !self.changed_indexes.is_empty() {
wrapped_segment.with_upgraded(|wrapped_segment| {
for (field_name, change) in self.changed_indexes.iter_ordered() {
debug_assert!(
change.version() >= op_num,
"proxied index change should have newer version than segment",
);
match change {
ProxyIndexChange::Create(schema, version) => {
wrapped_segment.create_field_index(
*version,
field_name,
Some(schema),
&HardwareCounterCell::disposable(), )?;
}
ProxyIndexChange::Delete(version) => {
wrapped_segment.delete_field_index(*version, field_name)?;
}
ProxyIndexChange::DeleteIfIncompatible(version, schema) => {
wrapped_segment.delete_field_index_if_incompatible(
*version, field_name, schema,
)?;
}
}
}
OperationResult::Ok(())
})?;
self.changed_indexes.clear();
}
}
{
if !self.changed_vector_names.is_empty() {
wrapped_segment.with_upgraded(|wrapped_segment| {
for (vector_name, intent) in self.changed_vector_names.iter_ordered() {
match intent {
IntendedVector::Absent { version } => {
wrapped_segment.delete_vector_name(*version, vector_name)?;
}
IntendedVector::Present {
config,
version,
supersedes_wrapped,
} => {
if *supersedes_wrapped {
wrapped_segment.delete_vector_name(*version, vector_name)?;
}
wrapped_segment.create_vector_name(
*version,
vector_name,
config,
)?;
}
}
}
OperationResult::Ok(())
})?;
self.changed_vector_names.clear();
}
}
{
if !self.deleted_points.is_empty() {
wrapped_segment.with_upgraded(|wrapped_segment| {
for (point_id, versions) in self.deleted_points.iter() {
wrapped_segment.delete_point(
versions.operation_version,
*point_id,
&HardwareCounterCell::disposable(), )?;
}
OperationResult::Ok(())
})?;
self.deleted_points.clear();
self.deleted_deferred_count = 0;
}
}
Ok(())
}
pub fn get_deleted_points(&self) -> &DeletedPoints {
&self.deleted_points
}
pub fn get_index_changes(&self) -> &ProxyIndexChanges {
&self.changed_indexes
}
pub fn get_vector_name_changes(&self) -> &ProxyVectorNameChanges {
&self.changed_vector_names
}
}
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub struct ProxyDeletedPoint {
pub local_version: SeqNumberType,
pub operation_version: SeqNumberType,
}
#[derive(Debug, Default)]
pub struct ProxyIndexChanges {
changes: AHashMap<PayloadKeyType, ProxyIndexChange>,
}
impl ProxyIndexChanges {
pub fn insert(&mut self, key: PayloadKeyType, change: ProxyIndexChange) {
self.changes.insert(key, change);
}
pub fn remove(&mut self, key: &PayloadKeyType) {
self.changes.remove(key);
}
pub fn len(&self) -> usize {
self.changes.len()
}
pub fn is_empty(&self) -> bool {
self.changes.is_empty()
}
pub fn clear(&mut self) {
self.changes.clear();
}
pub fn iter_ordered(&self) -> impl Iterator<Item = (&PayloadKeyType, &ProxyIndexChange)> {
self.changes
.iter()
.sorted_by_key(|(_, change)| change.version())
}
pub fn iter_unordered(&self) -> impl Iterator<Item = (&PayloadKeyType, &ProxyIndexChange)> {
self.changes.iter()
}
pub fn merge(&mut self, other: &Self) {
for (key, change) in &other.changes {
self.changes.insert(key.clone(), change.clone());
}
}
}
#[derive(Debug, Clone)]
pub enum ProxyIndexChange {
Create(PayloadFieldSchema, SeqNumberType),
Delete(SeqNumberType),
DeleteIfIncompatible(SeqNumberType, PayloadFieldSchema),
}
impl ProxyIndexChange {
pub fn version(&self) -> SeqNumberType {
match self {
ProxyIndexChange::Create(_, version) => *version,
ProxyIndexChange::Delete(version) => *version,
ProxyIndexChange::DeleteIfIncompatible(version, _) => *version,
}
}
}