diff --git a/lib/collection/src/operations/types.rs b/lib/collection/src/operations/types.rs index b2686a9a00..4f88d9ec56 100644 --- a/lib/collection/src/operations/types.rs +++ b/lib/collection/src/operations/types.rs @@ -388,6 +388,8 @@ pub enum CollectionError { shards_failed: u32, first_err: Box, }, + #[error("Remote shard on {peer_id} failed during forward proxy operation: {error}")] + ForwardProxyError { peer_id: PeerId, error: Box }, } impl CollectionError { @@ -409,6 +411,20 @@ impl CollectionError { pub fn bad_shard_selection(description: String) -> CollectionError { CollectionError::BadShardSelection { description } } + + pub fn forward_proxy_error(peer_id: PeerId, error: impl Into) -> Self { + Self::ForwardProxyError { + peer_id, + error: Box::new(error.into()), + } + } + + pub fn remote_peer_id(&self) -> Option { + match self { + Self::ForwardProxyError { peer_id, .. } => Some(*peer_id), + _ => None, + } + } } impl From for CollectionError { diff --git a/lib/collection/src/shards/forward_proxy_shard.rs b/lib/collection/src/shards/forward_proxy_shard.rs index ba7156960f..b71c75ed30 100644 --- a/lib/collection/src/shards/forward_proxy_shard.rs +++ b/lib/collection/src/shards/forward_proxy_shard.rs @@ -11,8 +11,8 @@ use tokio::sync::Mutex; use crate::operations::point_ops::{PointOperations, PointStruct, PointSyncOperation}; use crate::operations::types::{ - CollectionInfo, CollectionResult, CountRequest, CountResult, PointRequest, Record, - SearchRequestBatch, UpdateResult, + CollectionError, CollectionInfo, CollectionResult, CountRequest, CountResult, PointRequest, + Record, SearchRequestBatch, UpdateResult, }; use crate::operations::{CollectionUpdateOperations, CreateIndex, FieldIndexOperations}; use crate::shards::local_shard::LocalShard; @@ -152,7 +152,10 @@ impl ShardOperation for ForwardProxyShard { // during the transfer restart and finalization. local_shard.update(operation.clone(), wait).await?; - self.remote_shard.update(operation, false).await + self.remote_shard + .update(operation, false) + .await + .map_err(|err| CollectionError::forward_proxy_error(self.remote_shard.peer_id, err)) } /// Forward read-only `scroll_by` to `wrapped_shard` diff --git a/lib/collection/src/shards/replica_set.rs b/lib/collection/src/shards/replica_set.rs index bd868dfb81..77731ed667 100644 --- a/lib/collection/src/shards/replica_set.rs +++ b/lib/collection/src/shards/replica_set.rs @@ -1180,6 +1180,7 @@ impl ShardReplicaSet { } _ => {} } + log::debug!( "Deactivating peer {} because of failed update of shard {}:{}", peer_id, @@ -1323,7 +1324,12 @@ impl ShardReplicaSet { .get() .update(operation.clone(), wait) .await - .map_err(|err| (self.this_peer_id(), err)) + .map_err(|err| { + let peer_id = + err.remote_peer_id().unwrap_or_else(|| self.this_peer_id()); + + (peer_id, err) + }) }; let remote_updates = join_all(remote_futures); diff --git a/lib/storage/src/content_manager/errors.rs b/lib/storage/src/content_manager/errors.rs index 96724ef26b..c49de720a1 100644 --- a/lib/storage/src/content_manager/errors.rs +++ b/lib/storage/src/content_manager/errors.rs @@ -79,6 +79,9 @@ impl StorageError { CollectionError::BadShardSelection { .. } => StorageError::BadRequest { description: overriding_description, }, + CollectionError::ForwardProxyError { error, .. } => { + Self::from_inconsistent_shard_failure(*error, overriding_description) + } } } } @@ -109,6 +112,10 @@ impl From for StorageError { CollectionError::BadShardSelection { description } => { StorageError::BadRequest { description } } + CollectionError::ForwardProxyError { error, .. } => { + let full_description = format!("{error}"); + StorageError::from_inconsistent_shard_failure(*error, full_description) + } } } }