Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
169 changes: 98 additions & 71 deletions components/spider-execution-manager/src/client/grpc/liveness.rs
Original file line number Diff line number Diff line change
Expand Up @@ -6,13 +6,14 @@ use async_trait::async_trait;
use spider_core::types::id::{ExecutionManagerId, SessionId};
use spider_proto_rust::storage::{
self,
execution_manager_liveness_error,
execution_manager_liveness_service_client::ExecutionManagerLivenessServiceClient,
register_execution_manager_response,
update_execution_manager_heartbeat_response,
};
use spider_utils::grpc::client::ConnectionPool;
use tonic::transport::{Channel, Endpoint};
use tonic::{
Code,
Status,
transport::{Channel, Endpoint},
};

use crate::client::liveness::{LivenessClient, LivenessResponseError, RegistrationResponse};

Expand Down Expand Up @@ -59,7 +60,7 @@ impl LivenessClient for GrpcLivenessClient {
.get_client()
.register_execution_manager(request)
.await
.map_err(to_transport_error)?
.map_err(|status| status_to_error(&status))?
.into_inner();

register_response_to_result(response)
Expand All @@ -77,28 +78,27 @@ impl LivenessClient for GrpcLivenessClient {
.get_client()
.update_execution_manager_heartbeat(request)
.await
.map_err(to_transport_error)?
.map_err(|status| status_to_error(&status))?
.into_inner();

heartbeat_response_to_result(response)
Ok(heartbeat_response_to_result(response))
}
}

impl From<storage::ExecutionManagerLivenessError> for LivenessResponseError {
fn from(error: storage::ExecutionManagerLivenessError) -> Self {
match execution_manager_liveness_error::ErrCode::try_from(error.err_code) {
Ok(execution_manager_liveness_error::ErrCode::MarkedDead) => Self::MarkedDead,
Ok(execution_manager_liveness_error::ErrCode::InvalidInput) => {
Self::IllegalId(error.message)
}
Ok(
execution_manager_liveness_error::ErrCode::Server
| execution_manager_liveness_error::ErrCode::Unspecified,
) => Self::Transport(error.message),
Err(error) => Self::Transport(format!(
"unknown execution manager liveness error kind: {error}"
)),
}
/// Maps an execution-manager-liveness gRPC [`Status`] to a [`LivenessResponseError`].
///
/// # Returns
///
/// The [`LivenessResponseError`] for `status`'s code:
///
/// * [`LivenessResponseError::MarkedDead`] for `FAILED_PRECONDITION`.
/// * [`LivenessResponseError::IllegalId`] for `INVALID_ARGUMENT`.
/// * [`LivenessResponseError::Transport`] for any other code.
fn status_to_error(status: &Status) -> LivenessResponseError {
match status.code() {
Code::FailedPrecondition => LivenessResponseError::MarkedDead,
Code::InvalidArgument => LivenessResponseError::IllegalId(status.message().to_owned()),
_ => LivenessResponseError::Transport(status.message().to_owned()),
}
}

Expand All @@ -109,38 +109,24 @@ impl From<storage::ExecutionManagerLivenessError> for LivenessResponseError {
fn register_response_to_result(
response: storage::RegisterExecutionManagerResponse,
) -> Result<RegistrationResponse, LivenessResponseError> {
match response.result {
Some(register_execution_manager_response::Result::Registration(registration)) => {
Ok(RegistrationResponse {
em_id: ExecutionManagerId::from(registration.execution_manager_id),
session_id: registration.session_id,
})
}
Some(register_execution_manager_response::Result::Error(error)) => Err(error.into()),
None => Err(LivenessResponseError::Transport(
"register execution manager response missing result".to_owned(),
)),
}
let registration = response.registration.ok_or_else(|| {
LivenessResponseError::Transport(
"register execution manager response missing registration".to_owned(),
)
})?;
Ok(RegistrationResponse {
em_id: ExecutionManagerId::from(registration.execution_manager_id),
session_id: registration.session_id,
})
}

/// # Returns
///
/// [`storage::UpdateExecutionManagerHeartbeatResponse`] converted into
/// [`Result<SessionId, LivenessResponseError>`].
fn heartbeat_response_to_result(
/// The [`SessionId`] carried by `response`.
const fn heartbeat_response_to_result(
response: storage::UpdateExecutionManagerHeartbeatResponse,
) -> Result<SessionId, LivenessResponseError> {
match response.result {
Some(update_execution_manager_heartbeat_response::Result::SessionId(session_id)) => {
Ok(session_id)
}
Some(update_execution_manager_heartbeat_response::Result::Error(error)) => {
Err(error.into())
}
None => Err(LivenessResponseError::Transport(
"update execution manager heartbeat response missing result".to_owned(),
)),
}
) -> SessionId {
response.session_id
}

/// Converts a displayable transport-layer error into [`LivenessResponseError::Transport`].
Expand All @@ -165,12 +151,10 @@ mod tests {
const EM_ID: ExecutionManagerId = ExecutionManagerId::from(5);

let response = storage::RegisterExecutionManagerResponse {
result: Some(register_execution_manager_response::Result::Registration(
storage::ExecutionManagerRegistration {
execution_manager_id: EM_ID.get(),
session_id: SESSION_ID,
},
)),
registration: Some(storage::ExecutionManagerRegistration {
execution_manager_id: EM_ID.get(),
session_id: SESSION_ID,
}),
};

let registration = register_response_to_result(response)
Expand All @@ -185,36 +169,79 @@ mod tests {
);
}

#[test]
fn register_response_to_result_rejects_missing_registration() {
let response = storage::RegisterExecutionManagerResponse { registration: None };

assert!(matches!(
register_response_to_result(response),
Err(LivenessResponseError::Transport(_))
));
}

#[test]
fn register_response_to_result_accepts_zero_session_id() {
let response = storage::RegisterExecutionManagerResponse {
registration: Some(storage::ExecutionManagerRegistration {
execution_manager_id: 5,
session_id: 0,
}),
};

let registration = register_response_to_result(response)
.expect("registration response conversion should succeed");

assert_eq!(registration.session_id, 0);
}

#[test]
fn heartbeat_response_to_result_returns_session_id() {
const SESSION_ID: SessionId = 9;

let response = storage::UpdateExecutionManagerHeartbeatResponse {
result: Some(
update_execution_manager_heartbeat_response::Result::SessionId(SESSION_ID),
),
session_id: SESSION_ID,
};

let session_id = heartbeat_response_to_result(response)
.expect("heartbeat response conversion should succeed");
let session_id = heartbeat_response_to_result(response);

assert_eq!(session_id, SESSION_ID);
}

#[test]
fn liveness_storage_error_maps_invalid_input_to_illegal_id() {
const ERROR_MSG: &str = "bad em id";
fn heartbeat_response_to_result_accepts_zero_session_id() {
let response = storage::UpdateExecutionManagerHeartbeatResponse { session_id: 0 };

let error = storage::ExecutionManagerLivenessError {
err_code: execution_manager_liveness_error::ErrCode::InvalidInput.into(),
message: ERROR_MSG.to_owned(),
};
assert_eq!(heartbeat_response_to_result(response), 0);
}

match LivenessResponseError::from(error) {
LivenessResponseError::IllegalId(message) => {
assert_eq!(message, ERROR_MSG);
}
error => panic!("unexpected liveness response error: {error:?}"),
#[test]
fn status_maps_failed_precondition_to_marked_dead() {
let status = tonic::Status::failed_precondition("already dead");

assert!(matches!(
status_to_error(&status),
LivenessResponseError::MarkedDead
));
}

#[test]
fn status_maps_invalid_argument_to_illegal_id() {
const ERROR_MSG: &str = "bad em id";
let status = tonic::Status::invalid_argument(ERROR_MSG);

match status_to_error(&status) {
LivenessResponseError::IllegalId(message) => assert_eq!(message, ERROR_MSG),
error => panic!("unexpected liveness status mapping: {error:?}"),
}
}

#[test]
fn status_maps_other_codes_to_transport() {
let status = tonic::Status::internal("boom");

assert!(matches!(
status_to_error(&status),
LivenessResponseError::Transport(_)
));
}
}
Loading
Loading