use super::discriminants::*;
use super::execute::TypedClusterError;
use super::header::write_frame;
use super::raft_rpc::RaftRpc;
use crate::error::{ClusterError, Result};
#[derive(Debug, Clone, rkyv::Archive, rkyv::Serialize, rkyv::Deserialize)]
pub struct AssignSurrogateRequest {
pub vshard_id: u32,
pub database_id: u64,
pub tenant_id: u64,
pub collection: String,
pub pk: Vec<u8>,
pub deadline_remaining_ms: u64,
pub trace_id: [u8; 16],
}
#[derive(Debug, Clone, rkyv::Archive, rkyv::Serialize, rkyv::Deserialize)]
pub struct AssignSurrogateResponse {
pub surrogate: u32,
pub error: Option<TypedClusterError>,
}
macro_rules! to_bytes {
($msg:expr) => {
rkyv::to_bytes::<rkyv::rancor::Error>($msg)
.map(|b| b.to_vec())
.map_err(|e| ClusterError::Codec {
detail: format!("rkyv serialize: {e}"),
})
};
}
macro_rules! from_bytes {
($payload:expr, $T:ty, $name:expr) => {{
let mut aligned = rkyv::util::AlignedVec::<16>::with_capacity($payload.len());
aligned.extend_from_slice($payload);
rkyv::from_bytes::<$T, rkyv::rancor::Error>(&aligned).map_err(|e| ClusterError::Codec {
detail: format!("rkyv deserialize {}: {e}", $name),
})
}};
}
pub(super) fn encode_assign_surrogate_req(
msg: &AssignSurrogateRequest,
out: &mut Vec<u8>,
) -> Result<()> {
write_frame(RPC_ASSIGN_SURROGATE_REQ, &to_bytes!(msg)?, out)
}
pub(super) fn encode_assign_surrogate_resp(
msg: &AssignSurrogateResponse,
out: &mut Vec<u8>,
) -> Result<()> {
write_frame(RPC_ASSIGN_SURROGATE_RESP, &to_bytes!(msg)?, out)
}
pub(super) fn decode_assign_surrogate_req(payload: &[u8]) -> Result<RaftRpc> {
Ok(RaftRpc::AssignSurrogateRequest(from_bytes!(
payload,
AssignSurrogateRequest,
"AssignSurrogateRequest"
)?))
}
pub(super) fn decode_assign_surrogate_resp(payload: &[u8]) -> Result<RaftRpc> {
Ok(RaftRpc::AssignSurrogateResponse(from_bytes!(
payload,
AssignSurrogateResponse,
"AssignSurrogateResponse"
)?))
}
#[cfg(test)]
mod tests {
use super::*;
fn roundtrip_req(req: AssignSurrogateRequest) -> AssignSurrogateRequest {
let rpc = RaftRpc::AssignSurrogateRequest(req);
let encoded = super::super::encode(&rpc).unwrap();
match super::super::decode(&encoded).unwrap() {
RaftRpc::AssignSurrogateRequest(r) => r,
other => panic!("expected AssignSurrogateRequest, got {other:?}"),
}
}
fn roundtrip_resp(resp: AssignSurrogateResponse) -> AssignSurrogateResponse {
let rpc = RaftRpc::AssignSurrogateResponse(resp);
let encoded = super::super::encode(&rpc).unwrap();
match super::super::decode(&encoded).unwrap() {
RaftRpc::AssignSurrogateResponse(r) => r,
other => panic!("expected AssignSurrogateResponse, got {other:?}"),
}
}
#[test]
fn roundtrip_assign_surrogate_request() {
let req = AssignSurrogateRequest {
vshard_id: 512,
database_id: 7,
tenant_id: 42,
collection: "people".into(),
pk: vec![0x61, 0x6C, 0x69, 0x63, 0x65],
deadline_remaining_ms: 5000,
trace_id: [9u8; 16],
};
let decoded = roundtrip_req(req.clone());
assert_eq!(decoded.vshard_id, 512);
assert_eq!(decoded.database_id, 7);
assert_eq!(decoded.tenant_id, 42);
assert_eq!(decoded.collection, "people");
assert_eq!(decoded.pk, vec![0x61, 0x6C, 0x69, 0x63, 0x65]);
assert_eq!(decoded.deadline_remaining_ms, 5000);
assert_eq!(decoded.trace_id, [9u8; 16]);
}
#[test]
fn roundtrip_assign_surrogate_request_empty_pk() {
let req = AssignSurrogateRequest {
vshard_id: 0,
database_id: 0,
tenant_id: 0,
collection: String::new(),
pk: vec![],
deadline_remaining_ms: 1000,
trace_id: [0u8; 16],
};
let decoded = roundtrip_req(req);
assert!(decoded.pk.is_empty());
assert!(decoded.collection.is_empty());
assert_eq!(decoded.deadline_remaining_ms, 1000);
}
#[test]
fn roundtrip_assign_surrogate_response_ok() {
let decoded = roundtrip_resp(AssignSurrogateResponse {
surrogate: 12345,
error: None,
});
assert_eq!(decoded.surrogate, 12345);
assert!(decoded.error.is_none());
}
#[test]
fn roundtrip_assign_surrogate_response_error() {
let decoded = roundtrip_resp(AssignSurrogateResponse {
surrogate: 0,
error: Some(TypedClusterError::Internal {
code: 0x7,
message: "assign-remote-surrogate not configured".into(),
}),
});
assert_eq!(decoded.surrogate, 0);
match decoded.error {
Some(TypedClusterError::Internal { code, message }) => {
assert_eq!(code, 0x7);
assert!(message.contains("assign-remote-surrogate"));
}
other => panic!("expected Internal, got {other:?}"),
}
}
}