use crate::error::{ClusterError, Result};
use crate::rpc_codec::{
self, AssignSurrogateResponse, RaftRpc, ReleaseReservationResponse, ReserveReadResponse,
ShuffleAggregateConsumeResponse, ShuffleConsumeResponse, SubmitCalvinInboxResponse,
SubmitCalvinTxnResponse, auth_envelope,
};
use crate::transport::auth_context::AuthContext;
use crate::transport::rpc_handler::RaftRpcHandler;
pub(super) async fn try_handle_oneshot_rpc<H: RaftRpcHandler>(
handler: &H,
request: RaftRpc,
send: &mut quinn::SendStream,
auth: &AuthContext,
) -> Result<Option<RaftRpc>> {
if let RaftRpc::ShuffleProduceRequest(req) = request {
let resp = handler.on_shuffle_produce(req).await;
let resp_rpc = RaftRpc::ShuffleProduceResponse(resp);
let resp_inner = rpc_codec::encode(&resp_rpc)?;
let resp_seq = auth.peer_seq_out.next();
let mut resp_envelope =
Vec::with_capacity(auth_envelope::ENVELOPE_OVERHEAD + resp_inner.len());
auth_envelope::write_envelope(
auth.local_node_id,
resp_seq,
&resp_inner,
&auth.mac_key,
&mut resp_envelope,
)?;
send.write_all(&resp_envelope)
.await
.map_err(|e| ClusterError::Transport {
detail: format!("write shuffle produce response: {e}"),
})?;
send.finish().map_err(|e| ClusterError::Transport {
detail: format!("finish shuffle produce response: {e}"),
})?;
return Ok(None);
}
if let RaftRpc::ShuffleConsumeRequest(req) = request {
let resp: ShuffleConsumeResponse = handler.on_shuffle_consume(req).await;
let resp_rpc = RaftRpc::ShuffleConsumeResponse(resp);
let resp_inner = rpc_codec::encode(&resp_rpc)?;
let resp_seq = auth.peer_seq_out.next();
let mut resp_envelope =
Vec::with_capacity(auth_envelope::ENVELOPE_OVERHEAD + resp_inner.len());
auth_envelope::write_envelope(
auth.local_node_id,
resp_seq,
&resp_inner,
&auth.mac_key,
&mut resp_envelope,
)?;
send.write_all(&resp_envelope)
.await
.map_err(|e| ClusterError::Transport {
detail: format!("write shuffle consume response: {e}"),
})?;
send.finish().map_err(|e| ClusterError::Transport {
detail: format!("finish shuffle consume response: {e}"),
})?;
return Ok(None);
}
if let RaftRpc::ShuffleAggregateConsumeRequest(req) = request {
let resp: ShuffleAggregateConsumeResponse = handler.on_shuffle_aggregate(req).await;
let resp_rpc = RaftRpc::ShuffleAggregateConsumeResponse(resp);
let resp_inner = rpc_codec::encode(&resp_rpc)?;
let resp_seq = auth.peer_seq_out.next();
let mut resp_envelope =
Vec::with_capacity(auth_envelope::ENVELOPE_OVERHEAD + resp_inner.len());
auth_envelope::write_envelope(
auth.local_node_id,
resp_seq,
&resp_inner,
&auth.mac_key,
&mut resp_envelope,
)?;
send.write_all(&resp_envelope)
.await
.map_err(|e| ClusterError::Transport {
detail: format!("write shuffle aggregate consume response: {e}"),
})?;
send.finish().map_err(|e| ClusterError::Transport {
detail: format!("finish shuffle aggregate consume response: {e}"),
})?;
return Ok(None);
}
if let RaftRpc::AssignSurrogateRequest(req) = request {
let resp: AssignSurrogateResponse = handler.on_assign_surrogate(req).await;
let resp_rpc = RaftRpc::AssignSurrogateResponse(resp);
let resp_inner = rpc_codec::encode(&resp_rpc)?;
let resp_seq = auth.peer_seq_out.next();
let mut resp_envelope =
Vec::with_capacity(auth_envelope::ENVELOPE_OVERHEAD + resp_inner.len());
auth_envelope::write_envelope(
auth.local_node_id,
resp_seq,
&resp_inner,
&auth.mac_key,
&mut resp_envelope,
)?;
send.write_all(&resp_envelope)
.await
.map_err(|e| ClusterError::Transport {
detail: format!("write assign surrogate response: {e}"),
})?;
send.finish().map_err(|e| ClusterError::Transport {
detail: format!("finish assign surrogate response: {e}"),
})?;
return Ok(None);
}
if let RaftRpc::SubmitCalvinTxnRequest(req) = request {
let resp: SubmitCalvinTxnResponse = handler.on_submit_calvin_txn(req).await;
let resp_rpc = RaftRpc::SubmitCalvinTxnResponse(resp);
let resp_inner = rpc_codec::encode(&resp_rpc)?;
let resp_seq = auth.peer_seq_out.next();
let mut resp_envelope =
Vec::with_capacity(auth_envelope::ENVELOPE_OVERHEAD + resp_inner.len());
auth_envelope::write_envelope(
auth.local_node_id,
resp_seq,
&resp_inner,
&auth.mac_key,
&mut resp_envelope,
)?;
send.write_all(&resp_envelope)
.await
.map_err(|e| ClusterError::Transport {
detail: format!("write submit calvin txn response: {e}"),
})?;
send.finish().map_err(|e| ClusterError::Transport {
detail: format!("finish submit calvin txn response: {e}"),
})?;
return Ok(None);
}
if let RaftRpc::SubmitCalvinInboxRequest(req) = request {
let resp: SubmitCalvinInboxResponse = handler.on_submit_calvin_inbox(req).await;
let resp_rpc = RaftRpc::SubmitCalvinInboxResponse(resp);
let resp_inner = rpc_codec::encode(&resp_rpc)?;
let resp_seq = auth.peer_seq_out.next();
let mut resp_envelope =
Vec::with_capacity(auth_envelope::ENVELOPE_OVERHEAD + resp_inner.len());
auth_envelope::write_envelope(
auth.local_node_id,
resp_seq,
&resp_inner,
&auth.mac_key,
&mut resp_envelope,
)?;
send.write_all(&resp_envelope)
.await
.map_err(|e| ClusterError::Transport {
detail: format!("write submit calvin inbox response: {e}"),
})?;
send.finish().map_err(|e| ClusterError::Transport {
detail: format!("finish submit calvin inbox response: {e}"),
})?;
return Ok(None);
}
if let RaftRpc::ReserveReadRequest(req) = request {
let resp: ReserveReadResponse = handler.on_reserve_read(req).await;
let resp_rpc = RaftRpc::ReserveReadResponse(resp);
let resp_inner = rpc_codec::encode(&resp_rpc)?;
let resp_seq = auth.peer_seq_out.next();
let mut resp_envelope =
Vec::with_capacity(auth_envelope::ENVELOPE_OVERHEAD + resp_inner.len());
auth_envelope::write_envelope(
auth.local_node_id,
resp_seq,
&resp_inner,
&auth.mac_key,
&mut resp_envelope,
)?;
send.write_all(&resp_envelope)
.await
.map_err(|e| ClusterError::Transport {
detail: format!("write reserve read response: {e}"),
})?;
send.finish().map_err(|e| ClusterError::Transport {
detail: format!("finish reserve read response: {e}"),
})?;
return Ok(None);
}
if let RaftRpc::ReleaseReservationRequest(req) = request {
let resp: ReleaseReservationResponse = handler.on_release_reservation(req).await;
let resp_rpc = RaftRpc::ReleaseReservationResponse(resp);
let resp_inner = rpc_codec::encode(&resp_rpc)?;
let resp_seq = auth.peer_seq_out.next();
let mut resp_envelope =
Vec::with_capacity(auth_envelope::ENVELOPE_OVERHEAD + resp_inner.len());
auth_envelope::write_envelope(
auth.local_node_id,
resp_seq,
&resp_inner,
&auth.mac_key,
&mut resp_envelope,
)?;
send.write_all(&resp_envelope)
.await
.map_err(|e| ClusterError::Transport {
detail: format!("write release reservation response: {e}"),
})?;
send.finish().map_err(|e| ClusterError::Transport {
detail: format!("finish release reservation response: {e}"),
})?;
return Ok(None);
}
if let RaftRpc::TimeoutNowRequest(req) = request {
handler.on_timeout_now(req).await;
send.finish().map_err(|e| ClusterError::Transport {
detail: format!("finish timeout_now (no-response): {e}"),
})?;
return Ok(None);
}
Ok(Some(request))
}