qdrant-edge 0.7.0

A lightweight, in-process vector search engine designed for embedded devices, autonomous systems, and mobile agents.
Documentation
use std::path::Path;
use std::sync::Arc;

use crate::common::counter::hardware_counter::HardwareCounterCell;
use crate::common::save_on_disk::SaveOnDisk;
use crate::common::storage_version::StorageVersion;
use crate::common::types::PointOffsetType;
use parking_lot::{RwLockUpgradableReadGuard, RwLockWriteGuard};
use crate::segment::common::operation_error::OperationResult;
use crate::segment::entry::ReadSegmentEntry as _;
use crate::segment::segment::SegmentVersion;
use crate::segment::types::SegmentConfig;

use crate::shard::locked_segment::LockedSegment;
use crate::shard::payload_index_schema::PayloadIndexSchema;
use crate::shard::proxy_segment::UnsyncedProxySegment;
use crate::shard::segment_holder::locked::UpdatesGuard;
use crate::shard::segment_holder::{SegmentHolder, SegmentId};
use crate::shard::snapshots::snapshot_manifest::SnapshotManifest;

impl SegmentHolder {
    pub fn snapshot_manifest(&self) -> OperationResult<SnapshotManifest> {
        let mut manifest = SnapshotManifest::default();

        for (_, segment) in self.iter() {
            let segment_manifest = segment.get().read().get_segment_manifest()?;
            manifest.add(segment_manifest);
        }

        Ok(manifest)
    }

    /// Proxy all shard segments for [`proxy_all_segments_and_apply`].
    #[allow(clippy::type_complexity)]
    pub fn proxy_all_segments<'a>(
        segments_lock: RwLockUpgradableReadGuard<'a, SegmentHolder>,
        segments_path: &Path,
        segment_config: Option<SegmentConfig>,
        payload_index_schema: Arc<SaveOnDisk<PayloadIndexSchema>>,
        deferred_internal_id: Option<PointOffsetType>,
    ) -> OperationResult<(
        Vec<(SegmentId, LockedSegment)>,
        SegmentId,
        RwLockUpgradableReadGuard<'a, SegmentHolder>,
    )> {
        // This counter will be used to measure operations on temp segment,
        // which is part of internal process and can be ignored
        let hw_counter = HardwareCounterCell::disposable();

        // Create temporary appendable segment to direct all proxy writes into
        let tmp_segment = segments_lock.build_tmp_segment(
            segments_path,
            segment_config,
            payload_index_schema,
            deferred_internal_id,
            false,
        )?;

        // List all segments we want to snapshot
        let segment_ids = segments_lock.segment_ids();

        // Create proxy for all segments
        let mut new_proxies = Vec::with_capacity(segment_ids.len());
        for segment_id in segment_ids {
            let segment = segments_lock.get(segment_id).unwrap();
            let proxy = UnsyncedProxySegment::new(segment.clone());

            // Write segment is fresh, so it has no operations
            // Operation with number 0 will be applied
            proxy.replicate_field_indexes(0, &hw_counter, &tmp_segment)?;
            new_proxies.push((segment_id, proxy));
        }

        // Save segment version once all payload indices have been converted
        // If this ends up not being saved due to a crash, the segment will not be used
        match &tmp_segment {
            LockedSegment::Original(segment) => {
                let segment_path = &segment.read().segment_path;
                SegmentVersion::save(segment_path)?;
            }
            LockedSegment::Proxy(_) => unreachable!(),
        }

        // Replace all segments with proxies
        // We cannot fail past this point to prevent only having some segments proxified
        let mut proxies = Vec::with_capacity(new_proxies.len());
        let mut write_segments = RwLockUpgradableReadGuard::upgrade(segments_lock);
        for (segment_id, proxy) in new_proxies {
            // Some points might have been changed in the underlying segment before we upgraded the
            // `write_segments` lock, so the wrapped segment is only frozen now. Finalizing here
            // syncs `deleted_mask` from the now-immutable segment — the type-state guarantees this
            // happens exactly once and cannot be skipped.
            let proxy = proxy.finalize();

            // Replicate field indexes the second time, because optimized segments could have
            // been changed. The probability is small, though, so we can afford this operation
            // under the full collection write lock
            let op_num = proxy.version();
            if let Err(err) = proxy.replicate_field_indexes(op_num, &hw_counter, &tmp_segment) {
                log::error!("Failed to replicate proxy segment field indexes, ignoring: {err}");
            }

            // We must keep existing segment IDs because ongoing optimizations might depend on the mapping being the same
            write_segments.replace(segment_id, proxy)?;
            let locked_proxy_segment = write_segments
                .get(segment_id)
                .cloned()
                .expect("failed to get segment from segment holder we just swapped in");
            proxies.push((segment_id, locked_proxy_segment));
        }

        // Make sure at least one appendable segment exists
        let temp_segment_id = write_segments.add_new_locked(tmp_segment);

        let segments_lock = RwLockWriteGuard::downgrade_to_upgradable(write_segments);

        Ok((proxies, temp_segment_id, segments_lock))
    }

    /// Try to unproxy a single shard segment for [`proxy_all_segments_and_apply`].
    ///
    /// # Warning
    ///
    /// If unproxying fails an error is returned with the lock and the proxy is left behind in the
    /// shard holder.
    pub fn try_unproxy_segment<'a>(
        segments_lock: RwLockUpgradableReadGuard<'a, SegmentHolder>,
        segment_id: SegmentId,
        proxy_segment: LockedSegment,
        updates_guard: UpdatesGuard<'a>,
    ) -> Result<
        RwLockUpgradableReadGuard<'a, SegmentHolder>,
        RwLockUpgradableReadGuard<'a, SegmentHolder>,
    > {
        // We must propagate all changes in the proxy into their wrapped segments, as we'll put the
        // wrapped segment back into the segment holder. This can be an expensive step,
        // so it is important, that we don't block reads while doing this.

        let proxy_segment = match proxy_segment {
            LockedSegment::Proxy(proxy_segment) => proxy_segment,
            LockedSegment::Original(_) => {
                log::warn!(
                    "Unproxying segment {segment_id} that is not proxified, that is unexpected, skipping",
                );
                return Err(segments_lock);
            }
        };

        // propagate changes to wrapped segment with segment holder read lock
        {
            if let Err(err) = proxy_segment.write().propagate_to_wrapped() {
                log::error!(
                    "Propagating proxy segment {segment_id} changes to wrapped segment failed, ignoring: {err}",
                );
            }
        }

        let mut write_segments = RwLockUpgradableReadGuard::upgrade(segments_lock);

        let wrapped_segment = proxy_segment.read().wrapped_segment.clone();
        write_segments.replace(segment_id, wrapped_segment).unwrap();

        drop(updates_guard); // Release updates lock as soon as possible

        // Downgrade write lock to read and give it back
        Ok(RwLockWriteGuard::downgrade_to_upgradable(write_segments))
    }

    /// Unproxy all shard segments for [`proxy_all_segments_and_apply`].
    pub fn unproxy_all_segments(
        segments_lock: RwLockUpgradableReadGuard<SegmentHolder>,
        proxies: Vec<(SegmentId, LockedSegment)>,
        tmp_segment_id: SegmentId,
        updates_guard: UpdatesGuard<'_>,
    ) -> OperationResult<()> {
        // We must propagate all changes in the proxy into their wrapped segments, as we'll put the
        // wrapped segment back into the segment holder. This can be an expensive step,
        // so it is important, that we don't block reads while doing this.

        // propagate changes to wrapped segment with segment holder read lock
        proxies
            .iter()
            .filter_map(|(segment_id, proxy_segment)| match proxy_segment {
                LockedSegment::Proxy(proxy_segment) => Some((segment_id, proxy_segment)),
                LockedSegment::Original(_) => None,
            }).for_each(|(proxy_id, proxy_segment)| {
            if let Err(err) = proxy_segment.write().propagate_to_wrapped() {
                log::error!("Propagating proxy segment {proxy_id} changes to wrapped segment failed, ignoring: {err}");
            }
        });

        // Swap out each proxy with wrapped segment once changes are propagated
        let mut write_segments = RwLockUpgradableReadGuard::upgrade(segments_lock);
        for (segment_id, proxy_segment) in proxies {
            match proxy_segment {
                LockedSegment::Proxy(proxy_segment) => {
                    let wrapped_segment = proxy_segment.read().wrapped_segment.clone();
                    write_segments.replace(segment_id, wrapped_segment)?;
                }
                // If already unproxied, do nothing
                LockedSegment::Original(_) => {}
            }
        }

        debug_assert!(
            write_segments.get(tmp_segment_id).is_some(),
            "temp segment must exist",
        );

        // Remove temporary appendable segment, if we don't need it anymore
        write_segments.remove_segment_if_not_needed(tmp_segment_id)?;

        drop(updates_guard); // Release updates lock as soon as possible

        Ok(())
    }
}