From fbcdc0cf135eb1494f37b2b78bc9789be46bd069 Mon Sep 17 00:00:00 2001 From: Andrey Vasnetsov Date: Fri, 14 Nov 2025 16:12:29 +0100 Subject: [PATCH] Implement ReplicatePoints in grpc (#7536) * implement ReplicatePoints in grpc * move filter into a separate file --- lib/api/src/grpc/proto/collections.proto | 12 + lib/api/src/grpc/proto/common.proto | 142 ++++ lib/api/src/grpc/proto/points.proto | 145 +--- .../grpc/proto/points_internal_service.proto | 1 + lib/api/src/grpc/qdrant.rs | 662 +++++++++--------- lib/api/src/grpc/validate.rs | 7 + lib/collection/src/operations/conversions.rs | 27 +- 7 files changed, 527 insertions(+), 469 deletions(-) create mode 100644 lib/api/src/grpc/proto/common.proto diff --git a/lib/api/src/grpc/proto/collections.proto b/lib/api/src/grpc/proto/collections.proto index 6aeedd1e2a..df9c27bd2c 100644 --- a/lib/api/src/grpc/proto/collections.proto +++ b/lib/api/src/grpc/proto/collections.proto @@ -4,6 +4,7 @@ package qdrant; option csharp_namespace = "Qdrant.Client.Grpc"; import "json_with_int.proto"; +import "common.proto"; enum Datatype { Default = 0; @@ -12,6 +13,10 @@ enum Datatype { Float16 = 3; } +// --------------------------------------------- +// ------------- Collection Config ------------- +// --------------------------------------------- + message VectorParams { uint64 size = 1; // Size of the vectors Distance distance = 2; // Distance function used for comparing vectors @@ -724,6 +729,12 @@ message RestartTransfer { ShardTransferMethod method = 4; } +message ReplicatePoints { + ShardKey from_shard_key = 1; // Source shard key + ShardKey to_shard_key = 2; // Target shard key + optional Filter filter = 3; // If set - only points matching the filter will be replicated +} + enum ShardTransferMethod { StreamRecords = 0; // Stream shard records in batches Snapshot = 1; // Snapshot the shard and recover it on the target peer @@ -758,6 +769,7 @@ message UpdateCollectionClusterSetupRequest { CreateShardKey create_shard_key = 7; DeleteShardKey delete_shard_key = 8; RestartTransfer restart_transfer = 9; + ReplicatePoints replicate_points = 10; } optional uint64 timeout = 6; // Wait timeout for operation commit in seconds, if not specified - default value will be supplied } diff --git a/lib/api/src/grpc/proto/common.proto b/lib/api/src/grpc/proto/common.proto new file mode 100644 index 0000000000..70da91726b --- /dev/null +++ b/lib/api/src/grpc/proto/common.proto @@ -0,0 +1,142 @@ +syntax = "proto3"; +package qdrant; + +option csharp_namespace = "Qdrant.Client.Grpc"; +option java_outer_classname = "Points"; + +import "google/protobuf/timestamp.proto"; + +message PointId { + oneof point_id_options { + uint64 num = 1; // Numerical ID of the point + string uuid = 2; // UUID + } +} + +message GeoPoint { + double lon = 1; + double lat = 2; +} + +message Filter { + repeated Condition should = 1; // At least one of those conditions should match + repeated Condition must = 2; // All conditions must match + repeated Condition must_not = 3; // All conditions must NOT match + optional MinShould min_should = 4; // At least minimum amount of given conditions should match +} + +message MinShould { + repeated Condition conditions = 1; + uint64 min_count = 2; +} + +message Condition { + oneof condition_one_of { + FieldCondition field = 1; + IsEmptyCondition is_empty = 2; + HasIdCondition has_id = 3; + Filter filter = 4; + IsNullCondition is_null = 5; + NestedCondition nested = 6; + HasVectorCondition has_vector = 7; + } +} + +message IsEmptyCondition { + string key = 1; +} + +message IsNullCondition { + string key = 1; +} + +message HasIdCondition { + repeated PointId has_id = 1; +} + +message HasVectorCondition { + string has_vector = 1; +} + +message NestedCondition { + string key = 1; // Path to nested object + Filter filter = 2; // Filter condition +} + +message FieldCondition { + string key = 1; + Match match = 2; // Check if point has field with a given value + Range range = 3; // Check if points value lies in a given range + GeoBoundingBox geo_bounding_box = 4; // Check if points geolocation lies in a given area + GeoRadius geo_radius = 5; // Check if geo point is within a given radius + ValuesCount values_count = 6; // Check number of values for a specific field + GeoPolygon geo_polygon = 7; // Check if geo point is within a given polygon + DatetimeRange datetime_range = 8; // Check if datetime is within a given range + optional bool is_empty = 9; // Check if field is empty + optional bool is_null = 10; // Check if field is null +} + +message Match { + oneof match_value { + string keyword = 1; // Match string keyword + int64 integer = 2; // Match integer + bool boolean = 3; // Match boolean + string text = 4; // Match text + RepeatedStrings keywords = 5; // Match multiple keywords + RepeatedIntegers integers = 6; // Match multiple integers + RepeatedIntegers except_integers = 7; // Match any other value except those integers + RepeatedStrings except_keywords = 8; // Match any other value except those keywords + string phrase = 9; // Match phrase text + string text_any = 10; // Match any word in the text + } +} + +message RepeatedStrings { + repeated string strings = 1; +} + +message RepeatedIntegers { + repeated int64 integers = 1; +} + +message Range { + optional double lt = 1; + optional double gt = 2; + optional double gte = 3; + optional double lte = 4; +} + +message DatetimeRange { + optional google.protobuf.Timestamp lt = 1; + optional google.protobuf.Timestamp gt = 2; + optional google.protobuf.Timestamp gte = 3; + optional google.protobuf.Timestamp lte = 4; +} + +message GeoBoundingBox { + GeoPoint top_left = 1; // north-west corner + GeoPoint bottom_right = 2; // south-east corner +} + +message GeoRadius { + GeoPoint center = 1; // Center of the circle + float radius = 2; // In meters +} + +message GeoLineString { + repeated GeoPoint points = 1; // Ordered sequence of GeoPoints representing the line +} + +// For a valid GeoPolygon, both the exterior and interior GeoLineStrings must consist of a minimum of 4 points. +// Additionally, the first and last points of each GeoLineString must be the same. +message GeoPolygon { + GeoLineString exterior = 1; // The exterior line bounds the surface + repeated GeoLineString interiors = 2; // Interior lines (if present) bound holes within the surface +} + +message ValuesCount { + optional uint64 lt = 1; + optional uint64 gt = 2; + optional uint64 gte = 3; + optional uint64 lte = 4; +} \ No newline at end of file diff --git a/lib/api/src/grpc/proto/points.proto b/lib/api/src/grpc/proto/points.proto index 2b6c5e7430..fd50fa93e5 100644 --- a/lib/api/src/grpc/proto/points.proto +++ b/lib/api/src/grpc/proto/points.proto @@ -4,6 +4,7 @@ package qdrant; option csharp_namespace = "Qdrant.Client.Grpc"; import "collections.proto"; +import "common.proto"; import "google/protobuf/timestamp.proto"; import "json_with_int.proto"; @@ -31,17 +32,6 @@ message ReadConsistency { } } -// --------------------------------------------- -// ------------- Point Id Requests ------------- -// --------------------------------------------- - -message PointId { - oneof point_id_options { - uint64 num = 1; // Numerical ID of the point - string uuid = 2; // UUID - } -} - message SparseIndices { repeated uint32 data = 1; } @@ -1075,133 +1065,6 @@ message SearchMatrixOffsetsResponse { optional Usage usage = 3; } -// --------------------------------------------- -// ------------- Filter Conditions ------------- -// --------------------------------------------- - -message Filter { - repeated Condition should = 1; // At least one of those conditions should match - repeated Condition must = 2; // All conditions must match - repeated Condition must_not = 3; // All conditions must NOT match - optional MinShould min_should = 4; // At least minimum amount of given conditions should match -} - -message MinShould { - repeated Condition conditions = 1; - uint64 min_count = 2; -} - -message Condition { - oneof condition_one_of { - FieldCondition field = 1; - IsEmptyCondition is_empty = 2; - HasIdCondition has_id = 3; - Filter filter = 4; - IsNullCondition is_null = 5; - NestedCondition nested = 6; - HasVectorCondition has_vector = 7; - } -} - -message IsEmptyCondition { - string key = 1; -} - -message IsNullCondition { - string key = 1; -} - -message HasIdCondition { - repeated PointId has_id = 1; -} - -message HasVectorCondition { - string has_vector = 1; -} - -message NestedCondition { - string key = 1; // Path to nested object - Filter filter = 2; // Filter condition -} - -message FieldCondition { - string key = 1; - Match match = 2; // Check if point has field with a given value - Range range = 3; // Check if points value lies in a given range - GeoBoundingBox geo_bounding_box = 4; // Check if points geolocation lies in a given area - GeoRadius geo_radius = 5; // Check if geo point is within a given radius - ValuesCount values_count = 6; // Check number of values for a specific field - GeoPolygon geo_polygon = 7; // Check if geo point is within a given polygon - DatetimeRange datetime_range = 8; // Check if datetime is within a given range - optional bool is_empty = 9; // Check if field is empty - optional bool is_null = 10; // Check if field is null -} - -message Match { - oneof match_value { - string keyword = 1; // Match string keyword - int64 integer = 2; // Match integer - bool boolean = 3; // Match boolean - string text = 4; // Match text - RepeatedStrings keywords = 5; // Match multiple keywords - RepeatedIntegers integers = 6; // Match multiple integers - RepeatedIntegers except_integers = 7; // Match any other value except those integers - RepeatedStrings except_keywords = 8; // Match any other value except those keywords - string phrase = 9; // Match phrase text - string text_any = 10; // Match any word in the text - } -} - -message RepeatedStrings { - repeated string strings = 1; -} - -message RepeatedIntegers { - repeated int64 integers = 1; -} - -message Range { - optional double lt = 1; - optional double gt = 2; - optional double gte = 3; - optional double lte = 4; -} - -message DatetimeRange { - optional google.protobuf.Timestamp lt = 1; - optional google.protobuf.Timestamp gt = 2; - optional google.protobuf.Timestamp gte = 3; - optional google.protobuf.Timestamp lte = 4; -} - -message GeoBoundingBox { - GeoPoint top_left = 1; // north-west corner - GeoPoint bottom_right = 2; // south-east corner -} - -message GeoRadius { - GeoPoint center = 1; // Center of the circle - float radius = 2; // In meters -} - -message GeoLineString { - repeated GeoPoint points = 1; // Ordered sequence of GeoPoints representing the line -} - -// For a valid GeoPolygon, both the exterior and interior GeoLineStrings must consist of a minimum of 4 points. -// Additionally, the first and last points of each GeoLineString must be the same. -message GeoPolygon { - GeoLineString exterior = 1; // The exterior line bounds the surface - repeated GeoLineString interiors = 2; // Interior lines (if present) bound holes within the surface -} - -message ValuesCount { - optional uint64 lt = 1; - optional uint64 gt = 2; - optional uint64 gte = 3; - optional uint64 lte = 4; -} - // --------------------------------------------- // -------------- Points Selector -------------- // --------------------------------------------- @@ -1229,12 +1092,6 @@ message PointStruct { optional Vectors vectors = 4; } - -message GeoPoint { - double lon = 1; - double lat = 2; -} - // --------------------------------------------- // ----------- Measurements collector ---------- // --------------------------------------------- diff --git a/lib/api/src/grpc/proto/points_internal_service.proto b/lib/api/src/grpc/proto/points_internal_service.proto index 154b1740a9..9e7feefce0 100644 --- a/lib/api/src/grpc/proto/points_internal_service.proto +++ b/lib/api/src/grpc/proto/points_internal_service.proto @@ -1,6 +1,7 @@ syntax = "proto3"; import "points.proto"; +import "common.proto"; package qdrant; option csharp_namespace = "Qdrant.Client.Grpc"; diff --git a/lib/api/src/grpc/qdrant.rs b/lib/api/src/grpc/qdrant.rs index 141027cb77..42902dff8d 100644 --- a/lib/api/src/grpc/qdrant.rs +++ b/lib/api/src/grpc/qdrant.rs @@ -99,6 +99,328 @@ impl NullValue { } } } +#[derive(serde::Serialize)] +#[allow(clippy::derive_partial_eq_without_eq)] +#[derive(Clone, PartialEq, ::prost::Message)] +pub struct PointId { + #[prost(oneof = "point_id::PointIdOptions", tags = "1, 2")] + pub point_id_options: ::core::option::Option, +} +/// Nested message and enum types in `PointId`. +pub mod point_id { + #[derive(serde::Serialize)] + #[allow(clippy::derive_partial_eq_without_eq)] + #[derive(Clone, PartialEq, ::prost::Oneof)] + pub enum PointIdOptions { + /// Numerical ID of the point + #[prost(uint64, tag = "1")] + Num(u64), + /// UUID + #[prost(string, tag = "2")] + Uuid(::prost::alloc::string::String), + } +} +#[derive(serde::Serialize)] +#[allow(clippy::derive_partial_eq_without_eq)] +#[derive(Clone, PartialEq, ::prost::Message)] +pub struct GeoPoint { + #[prost(double, tag = "1")] + pub lon: f64, + #[prost(double, tag = "2")] + pub lat: f64, +} +#[derive(validator::Validate)] +#[derive(serde::Serialize)] +#[allow(clippy::derive_partial_eq_without_eq)] +#[derive(Clone, PartialEq, ::prost::Message)] +pub struct Filter { + /// At least one of those conditions should match + #[prost(message, repeated, tag = "1")] + #[validate(nested)] + pub should: ::prost::alloc::vec::Vec, + /// All conditions must match + #[prost(message, repeated, tag = "2")] + #[validate(nested)] + pub must: ::prost::alloc::vec::Vec, + /// All conditions must NOT match + #[prost(message, repeated, tag = "3")] + #[validate(nested)] + pub must_not: ::prost::alloc::vec::Vec, + /// At least minimum amount of given conditions should match + #[prost(message, optional, tag = "4")] + #[validate(nested)] + pub min_should: ::core::option::Option, +} +#[derive(validator::Validate)] +#[derive(serde::Serialize)] +#[allow(clippy::derive_partial_eq_without_eq)] +#[derive(Clone, PartialEq, ::prost::Message)] +pub struct MinShould { + #[prost(message, repeated, tag = "1")] + #[validate(nested)] + pub conditions: ::prost::alloc::vec::Vec, + #[prost(uint64, tag = "2")] + pub min_count: u64, +} +#[derive(validator::Validate)] +#[derive(serde::Serialize)] +#[allow(clippy::derive_partial_eq_without_eq)] +#[derive(Clone, PartialEq, ::prost::Message)] +pub struct Condition { + #[prost(oneof = "condition::ConditionOneOf", tags = "1, 2, 3, 4, 5, 6, 7")] + #[validate(nested)] + pub condition_one_of: ::core::option::Option, +} +/// Nested message and enum types in `Condition`. +pub mod condition { + #[derive(serde::Serialize)] + #[allow(clippy::derive_partial_eq_without_eq)] + #[derive(Clone, PartialEq, ::prost::Oneof)] + pub enum ConditionOneOf { + #[prost(message, tag = "1")] + Field(super::FieldCondition), + #[prost(message, tag = "2")] + IsEmpty(super::IsEmptyCondition), + #[prost(message, tag = "3")] + HasId(super::HasIdCondition), + #[prost(message, tag = "4")] + Filter(super::Filter), + #[prost(message, tag = "5")] + IsNull(super::IsNullCondition), + #[prost(message, tag = "6")] + Nested(super::NestedCondition), + #[prost(message, tag = "7")] + HasVector(super::HasVectorCondition), + } +} +#[derive(serde::Serialize)] +#[allow(clippy::derive_partial_eq_without_eq)] +#[derive(Clone, PartialEq, ::prost::Message)] +pub struct IsEmptyCondition { + #[prost(string, tag = "1")] + pub key: ::prost::alloc::string::String, +} +#[derive(serde::Serialize)] +#[allow(clippy::derive_partial_eq_without_eq)] +#[derive(Clone, PartialEq, ::prost::Message)] +pub struct IsNullCondition { + #[prost(string, tag = "1")] + pub key: ::prost::alloc::string::String, +} +#[derive(serde::Serialize)] +#[allow(clippy::derive_partial_eq_without_eq)] +#[derive(Clone, PartialEq, ::prost::Message)] +pub struct HasIdCondition { + #[prost(message, repeated, tag = "1")] + pub has_id: ::prost::alloc::vec::Vec, +} +#[derive(serde::Serialize)] +#[allow(clippy::derive_partial_eq_without_eq)] +#[derive(Clone, PartialEq, ::prost::Message)] +pub struct HasVectorCondition { + #[prost(string, tag = "1")] + pub has_vector: ::prost::alloc::string::String, +} +#[derive(validator::Validate)] +#[derive(serde::Serialize)] +#[allow(clippy::derive_partial_eq_without_eq)] +#[derive(Clone, PartialEq, ::prost::Message)] +pub struct NestedCondition { + /// Path to nested object + #[prost(string, tag = "1")] + pub key: ::prost::alloc::string::String, + /// Filter condition + #[prost(message, optional, tag = "2")] + #[validate(nested)] + pub filter: ::core::option::Option, +} +#[derive(serde::Serialize)] +#[allow(clippy::derive_partial_eq_without_eq)] +#[derive(Clone, PartialEq, ::prost::Message)] +pub struct FieldCondition { + #[prost(string, tag = "1")] + pub key: ::prost::alloc::string::String, + /// Check if point has field with a given value + #[prost(message, optional, tag = "2")] + pub r#match: ::core::option::Option, + /// Check if points value lies in a given range + #[prost(message, optional, tag = "3")] + pub range: ::core::option::Option, + /// Check if points geolocation lies in a given area + #[prost(message, optional, tag = "4")] + pub geo_bounding_box: ::core::option::Option, + /// Check if geo point is within a given radius + #[prost(message, optional, tag = "5")] + pub geo_radius: ::core::option::Option, + /// Check number of values for a specific field + #[prost(message, optional, tag = "6")] + pub values_count: ::core::option::Option, + /// Check if geo point is within a given polygon + #[prost(message, optional, tag = "7")] + pub geo_polygon: ::core::option::Option, + /// Check if datetime is within a given range + #[prost(message, optional, tag = "8")] + pub datetime_range: ::core::option::Option, + /// Check if field is empty + #[prost(bool, optional, tag = "9")] + pub is_empty: ::core::option::Option, + /// Check if field is null + #[prost(bool, optional, tag = "10")] + pub is_null: ::core::option::Option, +} +#[derive(serde::Serialize)] +#[allow(clippy::derive_partial_eq_without_eq)] +#[derive(Clone, PartialEq, ::prost::Message)] +pub struct Match { + #[prost(oneof = "r#match::MatchValue", tags = "1, 2, 3, 4, 5, 6, 7, 8, 9, 10")] + pub match_value: ::core::option::Option, +} +/// Nested message and enum types in `Match`. +pub mod r#match { + #[derive(serde::Serialize)] + #[allow(clippy::derive_partial_eq_without_eq)] + #[derive(Clone, PartialEq, ::prost::Oneof)] + pub enum MatchValue { + /// Match string keyword + #[prost(string, tag = "1")] + Keyword(::prost::alloc::string::String), + /// Match integer + #[prost(int64, tag = "2")] + Integer(i64), + /// Match boolean + #[prost(bool, tag = "3")] + Boolean(bool), + /// Match text + #[prost(string, tag = "4")] + Text(::prost::alloc::string::String), + /// Match multiple keywords + #[prost(message, tag = "5")] + Keywords(super::RepeatedStrings), + /// Match multiple integers + #[prost(message, tag = "6")] + Integers(super::RepeatedIntegers), + /// Match any other value except those integers + #[prost(message, tag = "7")] + ExceptIntegers(super::RepeatedIntegers), + /// Match any other value except those keywords + #[prost(message, tag = "8")] + ExceptKeywords(super::RepeatedStrings), + /// Match phrase text + #[prost(string, tag = "9")] + Phrase(::prost::alloc::string::String), + /// Match any word in the text + #[prost(string, tag = "10")] + TextAny(::prost::alloc::string::String), + } +} +#[derive(serde::Serialize)] +#[allow(clippy::derive_partial_eq_without_eq)] +#[derive(Clone, PartialEq, ::prost::Message)] +pub struct RepeatedStrings { + #[prost(string, repeated, tag = "1")] + pub strings: ::prost::alloc::vec::Vec<::prost::alloc::string::String>, +} +#[derive(serde::Serialize)] +#[allow(clippy::derive_partial_eq_without_eq)] +#[derive(Clone, PartialEq, ::prost::Message)] +pub struct RepeatedIntegers { + #[prost(int64, repeated, tag = "1")] + pub integers: ::prost::alloc::vec::Vec, +} +#[derive(serde::Serialize)] +#[allow(clippy::derive_partial_eq_without_eq)] +#[derive(Clone, PartialEq, ::prost::Message)] +pub struct Range { + #[prost(double, optional, tag = "1")] + pub lt: ::core::option::Option, + #[prost(double, optional, tag = "2")] + pub gt: ::core::option::Option, + #[prost(double, optional, tag = "3")] + pub gte: ::core::option::Option, + #[prost(double, optional, tag = "4")] + pub lte: ::core::option::Option, +} +#[derive(validator::Validate)] +#[derive(serde::Serialize)] +#[allow(clippy::derive_partial_eq_without_eq)] +#[derive(Clone, PartialEq, ::prost::Message)] +pub struct DatetimeRange { + #[prost(message, optional, tag = "1")] + #[validate(custom(function = "crate::grpc::validate::validate_timestamp"))] + pub lt: ::core::option::Option<::prost_wkt_types::Timestamp>, + #[prost(message, optional, tag = "2")] + #[validate(custom(function = "crate::grpc::validate::validate_timestamp"))] + pub gt: ::core::option::Option<::prost_wkt_types::Timestamp>, + #[prost(message, optional, tag = "3")] + #[validate(custom(function = "crate::grpc::validate::validate_timestamp"))] + pub gte: ::core::option::Option<::prost_wkt_types::Timestamp>, + #[prost(message, optional, tag = "4")] + #[validate(custom(function = "crate::grpc::validate::validate_timestamp"))] + pub lte: ::core::option::Option<::prost_wkt_types::Timestamp>, +} +#[derive(serde::Serialize)] +#[allow(clippy::derive_partial_eq_without_eq)] +#[derive(Clone, PartialEq, ::prost::Message)] +pub struct GeoBoundingBox { + /// north-west corner + #[prost(message, optional, tag = "1")] + pub top_left: ::core::option::Option, + /// south-east corner + #[prost(message, optional, tag = "2")] + pub bottom_right: ::core::option::Option, +} +#[derive(serde::Serialize)] +#[allow(clippy::derive_partial_eq_without_eq)] +#[derive(Clone, PartialEq, ::prost::Message)] +pub struct GeoRadius { + /// Center of the circle + #[prost(message, optional, tag = "1")] + pub center: ::core::option::Option, + /// In meters + #[prost(float, tag = "2")] + pub radius: f32, +} +#[derive(serde::Serialize)] +#[allow(clippy::derive_partial_eq_without_eq)] +#[derive(Clone, PartialEq, ::prost::Message)] +pub struct GeoLineString { + /// Ordered sequence of GeoPoints representing the line + #[prost(message, repeated, tag = "1")] + pub points: ::prost::alloc::vec::Vec, +} +/// For a valid GeoPolygon, both the exterior and interior GeoLineStrings must consist of a minimum of 4 points. +/// Additionally, the first and last points of each GeoLineString must be the same. +#[derive(validator::Validate)] +#[derive(serde::Serialize)] +#[allow(clippy::derive_partial_eq_without_eq)] +#[derive(Clone, PartialEq, ::prost::Message)] +pub struct GeoPolygon { + /// The exterior line bounds the surface + #[prost(message, optional, tag = "1")] + #[validate( + custom(function = "crate::grpc::validate::validate_geo_polygon_exterior") + )] + pub exterior: ::core::option::Option, + /// Interior lines (if present) bound holes within the surface + #[prost(message, repeated, tag = "2")] + #[validate( + custom(function = "crate::grpc::validate::validate_geo_polygon_interiors") + )] + pub interiors: ::prost::alloc::vec::Vec, +} +#[derive(serde::Serialize)] +#[allow(clippy::derive_partial_eq_without_eq)] +#[derive(Clone, PartialEq, ::prost::Message)] +pub struct ValuesCount { + #[prost(uint64, optional, tag = "1")] + pub lt: ::core::option::Option, + #[prost(uint64, optional, tag = "2")] + pub gt: ::core::option::Option, + #[prost(uint64, optional, tag = "3")] + pub gte: ::core::option::Option, + #[prost(uint64, optional, tag = "4")] + pub lte: ::core::option::Option, +} #[derive(validator::Validate)] #[derive(serde::Serialize)] #[allow(clippy::derive_partial_eq_without_eq)] @@ -1526,6 +1848,20 @@ pub struct RestartTransfer { #[prost(enumeration = "ShardTransferMethod", tag = "4")] pub method: i32, } +#[derive(serde::Serialize)] +#[allow(clippy::derive_partial_eq_without_eq)] +#[derive(Clone, PartialEq, ::prost::Message)] +pub struct ReplicatePoints { + /// Source shard key + #[prost(message, optional, tag = "1")] + pub from_shard_key: ::core::option::Option, + /// Target shard key + #[prost(message, optional, tag = "2")] + pub to_shard_key: ::core::option::Option, + /// If set - only points matching the filter will be replicated + #[prost(message, optional, tag = "3")] + pub filter: ::core::option::Option, +} #[derive(validator::Validate)] #[derive(serde::Serialize)] #[allow(clippy::derive_partial_eq_without_eq)] @@ -1578,7 +1914,7 @@ pub struct UpdateCollectionClusterSetupRequest { pub timeout: ::core::option::Option, #[prost( oneof = "update_collection_cluster_setup_request::Operation", - tags = "2, 3, 4, 5, 7, 8, 9" + tags = "2, 3, 4, 5, 7, 8, 9, 10" )] #[validate(nested)] pub operation: ::core::option::Option< @@ -1605,6 +1941,8 @@ pub mod update_collection_cluster_setup_request { DeleteShardKey(super::DeleteShardKey), #[prost(message, tag = "9")] RestartTransfer(super::RestartTransfer), + #[prost(message, tag = "10")] + ReplicatePoints(super::ReplicatePoints), } } #[derive(serde::Serialize)] @@ -4202,27 +4540,6 @@ pub mod read_consistency { #[derive(serde::Serialize)] #[allow(clippy::derive_partial_eq_without_eq)] #[derive(Clone, PartialEq, ::prost::Message)] -pub struct PointId { - #[prost(oneof = "point_id::PointIdOptions", tags = "1, 2")] - pub point_id_options: ::core::option::Option, -} -/// Nested message and enum types in `PointId`. -pub mod point_id { - #[derive(serde::Serialize)] - #[allow(clippy::derive_partial_eq_without_eq)] - #[derive(Clone, PartialEq, ::prost::Oneof)] - pub enum PointIdOptions { - /// Numerical ID of the point - #[prost(uint64, tag = "1")] - Num(u64), - /// UUID - #[prost(string, tag = "2")] - Uuid(::prost::alloc::string::String), - } -} -#[derive(serde::Serialize)] -#[allow(clippy::derive_partial_eq_without_eq)] -#[derive(Clone, PartialEq, ::prost::Message)] pub struct SparseIndices { #[prost(uint32, repeated, tag = "1")] pub data: ::prost::alloc::vec::Vec, @@ -6666,298 +6983,6 @@ pub struct SearchMatrixOffsetsResponse { #[derive(serde::Serialize)] #[allow(clippy::derive_partial_eq_without_eq)] #[derive(Clone, PartialEq, ::prost::Message)] -pub struct Filter { - /// At least one of those conditions should match - #[prost(message, repeated, tag = "1")] - #[validate(nested)] - pub should: ::prost::alloc::vec::Vec, - /// All conditions must match - #[prost(message, repeated, tag = "2")] - #[validate(nested)] - pub must: ::prost::alloc::vec::Vec, - /// All conditions must NOT match - #[prost(message, repeated, tag = "3")] - #[validate(nested)] - pub must_not: ::prost::alloc::vec::Vec, - /// At least minimum amount of given conditions should match - #[prost(message, optional, tag = "4")] - #[validate(nested)] - pub min_should: ::core::option::Option, -} -#[derive(validator::Validate)] -#[derive(serde::Serialize)] -#[allow(clippy::derive_partial_eq_without_eq)] -#[derive(Clone, PartialEq, ::prost::Message)] -pub struct MinShould { - #[prost(message, repeated, tag = "1")] - #[validate(nested)] - pub conditions: ::prost::alloc::vec::Vec, - #[prost(uint64, tag = "2")] - pub min_count: u64, -} -#[derive(validator::Validate)] -#[derive(serde::Serialize)] -#[allow(clippy::derive_partial_eq_without_eq)] -#[derive(Clone, PartialEq, ::prost::Message)] -pub struct Condition { - #[prost(oneof = "condition::ConditionOneOf", tags = "1, 2, 3, 4, 5, 6, 7")] - #[validate(nested)] - pub condition_one_of: ::core::option::Option, -} -/// Nested message and enum types in `Condition`. -pub mod condition { - #[derive(serde::Serialize)] - #[allow(clippy::derive_partial_eq_without_eq)] - #[derive(Clone, PartialEq, ::prost::Oneof)] - pub enum ConditionOneOf { - #[prost(message, tag = "1")] - Field(super::FieldCondition), - #[prost(message, tag = "2")] - IsEmpty(super::IsEmptyCondition), - #[prost(message, tag = "3")] - HasId(super::HasIdCondition), - #[prost(message, tag = "4")] - Filter(super::Filter), - #[prost(message, tag = "5")] - IsNull(super::IsNullCondition), - #[prost(message, tag = "6")] - Nested(super::NestedCondition), - #[prost(message, tag = "7")] - HasVector(super::HasVectorCondition), - } -} -#[derive(serde::Serialize)] -#[allow(clippy::derive_partial_eq_without_eq)] -#[derive(Clone, PartialEq, ::prost::Message)] -pub struct IsEmptyCondition { - #[prost(string, tag = "1")] - pub key: ::prost::alloc::string::String, -} -#[derive(serde::Serialize)] -#[allow(clippy::derive_partial_eq_without_eq)] -#[derive(Clone, PartialEq, ::prost::Message)] -pub struct IsNullCondition { - #[prost(string, tag = "1")] - pub key: ::prost::alloc::string::String, -} -#[derive(serde::Serialize)] -#[allow(clippy::derive_partial_eq_without_eq)] -#[derive(Clone, PartialEq, ::prost::Message)] -pub struct HasIdCondition { - #[prost(message, repeated, tag = "1")] - pub has_id: ::prost::alloc::vec::Vec, -} -#[derive(serde::Serialize)] -#[allow(clippy::derive_partial_eq_without_eq)] -#[derive(Clone, PartialEq, ::prost::Message)] -pub struct HasVectorCondition { - #[prost(string, tag = "1")] - pub has_vector: ::prost::alloc::string::String, -} -#[derive(validator::Validate)] -#[derive(serde::Serialize)] -#[allow(clippy::derive_partial_eq_without_eq)] -#[derive(Clone, PartialEq, ::prost::Message)] -pub struct NestedCondition { - /// Path to nested object - #[prost(string, tag = "1")] - pub key: ::prost::alloc::string::String, - /// Filter condition - #[prost(message, optional, tag = "2")] - #[validate(nested)] - pub filter: ::core::option::Option, -} -#[derive(serde::Serialize)] -#[allow(clippy::derive_partial_eq_without_eq)] -#[derive(Clone, PartialEq, ::prost::Message)] -pub struct FieldCondition { - #[prost(string, tag = "1")] - pub key: ::prost::alloc::string::String, - /// Check if point has field with a given value - #[prost(message, optional, tag = "2")] - pub r#match: ::core::option::Option, - /// Check if points value lies in a given range - #[prost(message, optional, tag = "3")] - pub range: ::core::option::Option, - /// Check if points geolocation lies in a given area - #[prost(message, optional, tag = "4")] - pub geo_bounding_box: ::core::option::Option, - /// Check if geo point is within a given radius - #[prost(message, optional, tag = "5")] - pub geo_radius: ::core::option::Option, - /// Check number of values for a specific field - #[prost(message, optional, tag = "6")] - pub values_count: ::core::option::Option, - /// Check if geo point is within a given polygon - #[prost(message, optional, tag = "7")] - pub geo_polygon: ::core::option::Option, - /// Check if datetime is within a given range - #[prost(message, optional, tag = "8")] - pub datetime_range: ::core::option::Option, - /// Check if field is empty - #[prost(bool, optional, tag = "9")] - pub is_empty: ::core::option::Option, - /// Check if field is null - #[prost(bool, optional, tag = "10")] - pub is_null: ::core::option::Option, -} -#[derive(serde::Serialize)] -#[allow(clippy::derive_partial_eq_without_eq)] -#[derive(Clone, PartialEq, ::prost::Message)] -pub struct Match { - #[prost(oneof = "r#match::MatchValue", tags = "1, 2, 3, 4, 5, 6, 7, 8, 9, 10")] - pub match_value: ::core::option::Option, -} -/// Nested message and enum types in `Match`. -pub mod r#match { - #[derive(serde::Serialize)] - #[allow(clippy::derive_partial_eq_without_eq)] - #[derive(Clone, PartialEq, ::prost::Oneof)] - pub enum MatchValue { - /// Match string keyword - #[prost(string, tag = "1")] - Keyword(::prost::alloc::string::String), - /// Match integer - #[prost(int64, tag = "2")] - Integer(i64), - /// Match boolean - #[prost(bool, tag = "3")] - Boolean(bool), - /// Match text - #[prost(string, tag = "4")] - Text(::prost::alloc::string::String), - /// Match multiple keywords - #[prost(message, tag = "5")] - Keywords(super::RepeatedStrings), - /// Match multiple integers - #[prost(message, tag = "6")] - Integers(super::RepeatedIntegers), - /// Match any other value except those integers - #[prost(message, tag = "7")] - ExceptIntegers(super::RepeatedIntegers), - /// Match any other value except those keywords - #[prost(message, tag = "8")] - ExceptKeywords(super::RepeatedStrings), - /// Match phrase text - #[prost(string, tag = "9")] - Phrase(::prost::alloc::string::String), - /// Match any word in the text - #[prost(string, tag = "10")] - TextAny(::prost::alloc::string::String), - } -} -#[derive(serde::Serialize)] -#[allow(clippy::derive_partial_eq_without_eq)] -#[derive(Clone, PartialEq, ::prost::Message)] -pub struct RepeatedStrings { - #[prost(string, repeated, tag = "1")] - pub strings: ::prost::alloc::vec::Vec<::prost::alloc::string::String>, -} -#[derive(serde::Serialize)] -#[allow(clippy::derive_partial_eq_without_eq)] -#[derive(Clone, PartialEq, ::prost::Message)] -pub struct RepeatedIntegers { - #[prost(int64, repeated, tag = "1")] - pub integers: ::prost::alloc::vec::Vec, -} -#[derive(serde::Serialize)] -#[allow(clippy::derive_partial_eq_without_eq)] -#[derive(Clone, PartialEq, ::prost::Message)] -pub struct Range { - #[prost(double, optional, tag = "1")] - pub lt: ::core::option::Option, - #[prost(double, optional, tag = "2")] - pub gt: ::core::option::Option, - #[prost(double, optional, tag = "3")] - pub gte: ::core::option::Option, - #[prost(double, optional, tag = "4")] - pub lte: ::core::option::Option, -} -#[derive(validator::Validate)] -#[derive(serde::Serialize)] -#[allow(clippy::derive_partial_eq_without_eq)] -#[derive(Clone, PartialEq, ::prost::Message)] -pub struct DatetimeRange { - #[prost(message, optional, tag = "1")] - #[validate(custom(function = "crate::grpc::validate::validate_timestamp"))] - pub lt: ::core::option::Option<::prost_wkt_types::Timestamp>, - #[prost(message, optional, tag = "2")] - #[validate(custom(function = "crate::grpc::validate::validate_timestamp"))] - pub gt: ::core::option::Option<::prost_wkt_types::Timestamp>, - #[prost(message, optional, tag = "3")] - #[validate(custom(function = "crate::grpc::validate::validate_timestamp"))] - pub gte: ::core::option::Option<::prost_wkt_types::Timestamp>, - #[prost(message, optional, tag = "4")] - #[validate(custom(function = "crate::grpc::validate::validate_timestamp"))] - pub lte: ::core::option::Option<::prost_wkt_types::Timestamp>, -} -#[derive(serde::Serialize)] -#[allow(clippy::derive_partial_eq_without_eq)] -#[derive(Clone, PartialEq, ::prost::Message)] -pub struct GeoBoundingBox { - /// north-west corner - #[prost(message, optional, tag = "1")] - pub top_left: ::core::option::Option, - /// south-east corner - #[prost(message, optional, tag = "2")] - pub bottom_right: ::core::option::Option, -} -#[derive(serde::Serialize)] -#[allow(clippy::derive_partial_eq_without_eq)] -#[derive(Clone, PartialEq, ::prost::Message)] -pub struct GeoRadius { - /// Center of the circle - #[prost(message, optional, tag = "1")] - pub center: ::core::option::Option, - /// In meters - #[prost(float, tag = "2")] - pub radius: f32, -} -#[derive(serde::Serialize)] -#[allow(clippy::derive_partial_eq_without_eq)] -#[derive(Clone, PartialEq, ::prost::Message)] -pub struct GeoLineString { - /// Ordered sequence of GeoPoints representing the line - #[prost(message, repeated, tag = "1")] - pub points: ::prost::alloc::vec::Vec, -} -/// For a valid GeoPolygon, both the exterior and interior GeoLineStrings must consist of a minimum of 4 points. -/// Additionally, the first and last points of each GeoLineString must be the same. -#[derive(validator::Validate)] -#[derive(serde::Serialize)] -#[allow(clippy::derive_partial_eq_without_eq)] -#[derive(Clone, PartialEq, ::prost::Message)] -pub struct GeoPolygon { - /// The exterior line bounds the surface - #[prost(message, optional, tag = "1")] - #[validate( - custom(function = "crate::grpc::validate::validate_geo_polygon_exterior") - )] - pub exterior: ::core::option::Option, - /// Interior lines (if present) bound holes within the surface - #[prost(message, repeated, tag = "2")] - #[validate( - custom(function = "crate::grpc::validate::validate_geo_polygon_interiors") - )] - pub interiors: ::prost::alloc::vec::Vec, -} -#[derive(serde::Serialize)] -#[allow(clippy::derive_partial_eq_without_eq)] -#[derive(Clone, PartialEq, ::prost::Message)] -pub struct ValuesCount { - #[prost(uint64, optional, tag = "1")] - pub lt: ::core::option::Option, - #[prost(uint64, optional, tag = "2")] - pub gt: ::core::option::Option, - #[prost(uint64, optional, tag = "3")] - pub gte: ::core::option::Option, - #[prost(uint64, optional, tag = "4")] - pub lte: ::core::option::Option, -} -#[derive(validator::Validate)] -#[derive(serde::Serialize)] -#[allow(clippy::derive_partial_eq_without_eq)] -#[derive(Clone, PartialEq, ::prost::Message)] pub struct PointsSelector { #[prost(oneof = "points_selector::PointsSelectorOneOf", tags = "1, 2")] #[validate(nested)] @@ -6997,15 +7022,6 @@ pub struct PointStruct { #[validate(nested)] pub vectors: ::core::option::Option, } -#[derive(serde::Serialize)] -#[allow(clippy::derive_partial_eq_without_eq)] -#[derive(Clone, PartialEq, ::prost::Message)] -pub struct GeoPoint { - #[prost(double, tag = "1")] - pub lon: f64, - #[prost(double, tag = "2")] - pub lat: f64, -} /// --- /// /// ## ----------- Measurements collector ---------- diff --git a/lib/api/src/grpc/validate.rs b/lib/api/src/grpc/validate.rs index 416043a08d..f7fc6cfce0 100644 --- a/lib/api/src/grpc/validate.rs +++ b/lib/api/src/grpc/validate.rs @@ -108,6 +108,7 @@ impl Validate for grpc::update_collection_cluster_setup_request::Operation { Operation::CreateShardKey(op) => op.validate(), Operation::DeleteShardKey(op) => op.validate(), Operation::RestartTransfer(op) => op.validate(), + Operation::ReplicatePoints(op) => op.validate(), } } } @@ -181,6 +182,12 @@ impl Validate for grpc::RestartTransfer { } } +impl Validate for grpc::ReplicatePoints { + fn validate(&self) -> Result<(), ValidationErrors> { + Ok(()) + } +} + impl Validate for grpc::condition::ConditionOneOf { fn validate(&self) -> Result<(), ValidationErrors> { use grpc::condition::ConditionOneOf; diff --git a/lib/collection/src/operations/conversions.rs b/lib/collection/src/operations/conversions.rs index 090e7a806f..d20e91a4a6 100644 --- a/lib/collection/src/operations/conversions.rs +++ b/lib/collection/src/operations/conversions.rs @@ -18,13 +18,13 @@ use segment::common::operation_error::OperationError; use segment::data_types::modifier::Modifier; use segment::data_types::vectors::{VectorInternal, VectorStructInternal}; use segment::types::{ - Distance, HnswConfig, MultiVectorConfig, QuantizationConfig, StrictModeConfigOutput, + Distance, Filter, HnswConfig, MultiVectorConfig, QuantizationConfig, StrictModeConfigOutput, WithPayloadInterface, }; use shard::retrieve::record_internal::RecordInternal; use tonic::Status; -use super::cluster_ops::ReshardingDirection; +use super::cluster_ops::{ReplicatePoints, ReplicatePointsOperation, ReshardingDirection}; use super::consistency_params::ReadConsistency; use super::types::{ CollectionConfig, ContextExamplePair, CoreSearchRequest, Datatype, DiscoverRequestInternal, @@ -1679,6 +1679,29 @@ impl TryFrom for ClusterOperations { }, }) } + Operation::ReplicatePoints(op) => { + let api::grpc::qdrant::ReplicatePoints { + from_shard_key, + to_shard_key, + filter, + } = op; + + ClusterOperations::ReplicatePoints(ReplicatePointsOperation { + replicate_points: ReplicatePoints { + filter: filter.map(Filter::try_from).transpose()?, + from_shard_key: from_shard_key + .and_then(convert_shard_key_from_grpc) + .ok_or_else(|| { + Status::invalid_argument("from_shard_key is not specified") + })?, + to_shard_key: to_shard_key + .and_then(convert_shard_key_from_grpc) + .ok_or_else(|| { + Status::invalid_argument("to_shard_key is not specified") + })?, + }, + }) + } }) } }