Fix ForwardProxyShard::update error handling (#1638)

Refactor `ForwardProxyShard::update` to separate failure of the local shard from the failure of the remote shard.

TODO:
- Figure a way to return the result of the local shard (regardless of status of remote shard)...
- ...and also handle the failure of remote shard
This commit is contained in:
Roman Titov
2023-04-11 14:25:19 +02:00
committed by Andrey Vasnetsov
parent 54ecc5d7b7
commit ce39ce7cf3
4 changed files with 36 additions and 4 deletions
+16
View File
@@ -388,6 +388,8 @@ pub enum CollectionError {
shards_failed: u32,
first_err: Box<CollectionError>,
},
#[error("Remote shard on {peer_id} failed during forward proxy operation: {error}")]
ForwardProxyError { peer_id: PeerId, error: Box<Self> },
}
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 {
Self::ForwardProxyError {
peer_id,
error: Box::new(error.into()),
}
}
pub fn remote_peer_id(&self) -> Option<PeerId> {
match self {
Self::ForwardProxyError { peer_id, .. } => Some(*peer_id),
_ => None,
}
}
}
impl From<SystemTimeError> for CollectionError {
@@ -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`
+7 -1
View File
@@ -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);
@@ -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<CollectionError> for StorageError {
CollectionError::BadShardSelection { description } => {
StorageError::BadRequest { description }
}
CollectionError::ForwardProxyError { error, .. } => {
let full_description = format!("{error}");
StorageError::from_inconsistent_shard_failure(*error, full_description)
}
}
}
}