use super::discriminants::*;
use super::execute::{DescriptorVersionEntry, 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 ShufflePushRequest {
pub shuffle_id: u64,
pub part: u32,
pub side: u8,
pub num_parts: u32,
pub producer_count: u32,
}
#[derive(Debug, Clone, rkyv::Archive, rkyv::Serialize, rkyv::Deserialize)]
pub struct ShufflePushChunk {
pub payload: Vec<u8>,
}
#[derive(Debug, Clone, rkyv::Archive, rkyv::Serialize, rkyv::Deserialize)]
pub struct ShufflePushEnd {
pub error: Option<TypedClusterError>,
}
#[derive(Debug, Clone, rkyv::Archive, rkyv::Serialize, rkyv::Deserialize)]
pub struct PartNodeEntry {
pub part: u32,
pub node_id: u64,
}
#[derive(Debug, Clone, rkyv::Archive, rkyv::Serialize, rkyv::Deserialize)]
pub struct ShuffleProduceRequest {
pub shuffle_id: u64,
pub side: u8,
pub num_parts: u32,
pub producer_count: u32,
pub keys: Vec<String>,
pub part_node_map: Vec<PartNodeEntry>,
pub plan_bytes: Vec<u8>,
pub tenant_id: u64,
pub database_id: u64,
pub deadline_remaining_ms: u64,
pub trace_id: [u8; 16],
pub descriptor_versions: Vec<DescriptorVersionEntry>,
}
#[derive(Debug, Clone, rkyv::Archive, rkyv::Serialize, rkyv::Deserialize)]
pub struct ShuffleProduceResponse {
pub error: Option<TypedClusterError>,
pub read_version_lsn: u64,
}
#[derive(Debug, Clone, rkyv::Archive, rkyv::Serialize, rkyv::Deserialize)]
pub struct JoinKeyPair {
pub left: String,
pub right: String,
}
#[derive(Debug, Clone, rkyv::Archive, rkyv::Serialize, rkyv::Deserialize)]
pub struct ShuffleConsumeRequest {
pub shuffle_id: u64,
pub part: u32,
pub on: Vec<JoinKeyPair>,
pub join_type: String,
pub limit: u64,
pub probe_qualifier: String,
pub index_qualifier: String,
pub tenant_id: u64,
pub database_id: u64,
pub deadline_remaining_ms: u64,
pub trace_id: [u8; 16],
}
#[derive(Debug, Clone, rkyv::Archive, rkyv::Serialize, rkyv::Deserialize)]
pub struct ShuffleConsumeResponse {
pub rows: Vec<u8>,
pub error: Option<TypedClusterError>,
}
#[derive(Debug, Clone, rkyv::Archive, rkyv::Serialize, rkyv::Deserialize)]
pub struct SortKey {
pub column: String,
pub ascending: bool,
}
#[derive(Debug, Clone, rkyv::Archive, rkyv::Serialize, rkyv::Deserialize)]
pub struct ShuffleAggregateConsumeRequest {
pub shuffle_id: u64,
pub part: u32,
pub group_by: Vec<String>,
pub aggregates_bytes: Vec<u8>,
pub having: Vec<u8>,
pub limit: u64,
pub sort_keys: Vec<SortKey>,
pub tenant_id: u64,
pub database_id: u64,
pub deadline_remaining_ms: u64,
pub trace_id: [u8; 16],
}
#[derive(Debug, Clone, rkyv::Archive, rkyv::Serialize, rkyv::Deserialize)]
pub struct ShuffleAggregateConsumeResponse {
pub rows: Vec<u8>,
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_shuffle_push_req(msg: &ShufflePushRequest, out: &mut Vec<u8>) -> Result<()> {
write_frame(RPC_SHUFFLE_PUSH_REQ, &to_bytes!(msg)?, out)
}
pub(super) fn encode_shuffle_push_chunk(msg: &ShufflePushChunk, out: &mut Vec<u8>) -> Result<()> {
write_frame(RPC_SHUFFLE_PUSH_CHUNK, &to_bytes!(msg)?, out)
}
pub(super) fn encode_shuffle_push_end(msg: &ShufflePushEnd, out: &mut Vec<u8>) -> Result<()> {
write_frame(RPC_SHUFFLE_PUSH_END, &to_bytes!(msg)?, out)
}
pub(super) fn decode_shuffle_push_req(payload: &[u8]) -> Result<RaftRpc> {
Ok(RaftRpc::ShufflePushRequest(from_bytes!(
payload,
ShufflePushRequest,
"ShufflePushRequest"
)?))
}
pub(super) fn decode_shuffle_push_chunk(payload: &[u8]) -> Result<RaftRpc> {
Ok(RaftRpc::ShufflePushChunk(from_bytes!(
payload,
ShufflePushChunk,
"ShufflePushChunk"
)?))
}
pub(super) fn decode_shuffle_push_end(payload: &[u8]) -> Result<RaftRpc> {
Ok(RaftRpc::ShufflePushEnd(from_bytes!(
payload,
ShufflePushEnd,
"ShufflePushEnd"
)?))
}
pub(super) fn encode_shuffle_produce_req(
msg: &ShuffleProduceRequest,
out: &mut Vec<u8>,
) -> Result<()> {
write_frame(RPC_SHUFFLE_PRODUCE_REQ, &to_bytes!(msg)?, out)
}
pub(super) fn encode_shuffle_produce_resp(
msg: &ShuffleProduceResponse,
out: &mut Vec<u8>,
) -> Result<()> {
write_frame(RPC_SHUFFLE_PRODUCE_RESP, &to_bytes!(msg)?, out)
}
pub(super) fn decode_shuffle_produce_req(payload: &[u8]) -> Result<RaftRpc> {
Ok(RaftRpc::ShuffleProduceRequest(from_bytes!(
payload,
ShuffleProduceRequest,
"ShuffleProduceRequest"
)?))
}
pub(super) fn decode_shuffle_produce_resp(payload: &[u8]) -> Result<RaftRpc> {
Ok(RaftRpc::ShuffleProduceResponse(from_bytes!(
payload,
ShuffleProduceResponse,
"ShuffleProduceResponse"
)?))
}
pub(super) fn encode_shuffle_consume_req(
msg: &ShuffleConsumeRequest,
out: &mut Vec<u8>,
) -> Result<()> {
write_frame(RPC_SHUFFLE_CONSUME_REQ, &to_bytes!(msg)?, out)
}
pub(super) fn encode_shuffle_consume_resp(
msg: &ShuffleConsumeResponse,
out: &mut Vec<u8>,
) -> Result<()> {
write_frame(RPC_SHUFFLE_CONSUME_RESP, &to_bytes!(msg)?, out)
}
pub(super) fn decode_shuffle_consume_req(payload: &[u8]) -> Result<RaftRpc> {
Ok(RaftRpc::ShuffleConsumeRequest(from_bytes!(
payload,
ShuffleConsumeRequest,
"ShuffleConsumeRequest"
)?))
}
pub(super) fn decode_shuffle_consume_resp(payload: &[u8]) -> Result<RaftRpc> {
Ok(RaftRpc::ShuffleConsumeResponse(from_bytes!(
payload,
ShuffleConsumeResponse,
"ShuffleConsumeResponse"
)?))
}
pub(super) fn encode_shuffle_agg_consume_req(
msg: &ShuffleAggregateConsumeRequest,
out: &mut Vec<u8>,
) -> Result<()> {
write_frame(RPC_SHUFFLE_AGG_CONSUME_REQ, &to_bytes!(msg)?, out)
}
pub(super) fn encode_shuffle_agg_consume_resp(
msg: &ShuffleAggregateConsumeResponse,
out: &mut Vec<u8>,
) -> Result<()> {
write_frame(RPC_SHUFFLE_AGG_CONSUME_RESP, &to_bytes!(msg)?, out)
}
pub(super) fn decode_shuffle_agg_consume_req(payload: &[u8]) -> Result<RaftRpc> {
Ok(RaftRpc::ShuffleAggregateConsumeRequest(from_bytes!(
payload,
ShuffleAggregateConsumeRequest,
"ShuffleAggregateConsumeRequest"
)?))
}
pub(super) fn decode_shuffle_agg_consume_resp(payload: &[u8]) -> Result<RaftRpc> {
Ok(RaftRpc::ShuffleAggregateConsumeResponse(from_bytes!(
payload,
ShuffleAggregateConsumeResponse,
"ShuffleAggregateConsumeResponse"
)?))
}
#[cfg(test)]
mod tests {
use super::*;
fn roundtrip_req(req: ShufflePushRequest) -> ShufflePushRequest {
let rpc = RaftRpc::ShufflePushRequest(req);
let encoded = super::super::encode(&rpc).unwrap();
match super::super::decode(&encoded).unwrap() {
RaftRpc::ShufflePushRequest(r) => r,
other => panic!("expected ShufflePushRequest, got {other:?}"),
}
}
fn roundtrip_chunk(chunk: ShufflePushChunk) -> ShufflePushChunk {
let rpc = RaftRpc::ShufflePushChunk(chunk);
let encoded = super::super::encode(&rpc).unwrap();
match super::super::decode(&encoded).unwrap() {
RaftRpc::ShufflePushChunk(c) => c,
other => panic!("expected ShufflePushChunk, got {other:?}"),
}
}
fn roundtrip_end(end: ShufflePushEnd) -> ShufflePushEnd {
let rpc = RaftRpc::ShufflePushEnd(end);
let encoded = super::super::encode(&rpc).unwrap();
match super::super::decode(&encoded).unwrap() {
RaftRpc::ShufflePushEnd(e) => e,
other => panic!("expected ShufflePushEnd, got {other:?}"),
}
}
#[test]
fn roundtrip_shuffle_push_request() {
let req = ShufflePushRequest {
shuffle_id: 0xDEAD_BEEF_1234_5678,
part: 7,
side: 1,
num_parts: 16,
producer_count: 3,
};
let decoded = roundtrip_req(req.clone());
assert_eq!(decoded.shuffle_id, req.shuffle_id);
assert_eq!(decoded.part, 7);
assert_eq!(decoded.side, 1);
assert_eq!(decoded.num_parts, 16);
assert_eq!(decoded.producer_count, 3);
}
#[test]
fn roundtrip_shuffle_push_request_build_side() {
let req = ShufflePushRequest {
shuffle_id: 1,
part: 0,
side: 0,
num_parts: 1,
producer_count: 1,
};
let decoded = roundtrip_req(req);
assert_eq!(decoded.side, 0);
assert_eq!(decoded.num_parts, 1);
assert_eq!(decoded.producer_count, 1);
}
#[test]
fn roundtrip_shuffle_push_chunk_payload() {
let chunk = ShufflePushChunk {
payload: vec![0x93, 0x01, 0x02, 0x03],
};
let decoded = roundtrip_chunk(chunk.clone());
assert_eq!(decoded.payload, chunk.payload);
}
#[test]
fn roundtrip_shuffle_push_chunk_empty_payload() {
let decoded = roundtrip_chunk(ShufflePushChunk { payload: vec![] });
assert!(decoded.payload.is_empty());
}
#[test]
fn roundtrip_shuffle_push_end_clean_eof() {
let decoded = roundtrip_end(ShufflePushEnd { error: None });
assert!(decoded.error.is_none());
}
#[test]
fn roundtrip_shuffle_push_end_terminal_error() {
let decoded = roundtrip_end(ShufflePushEnd {
error: Some(TypedClusterError::Internal {
code: 0xABCD,
message: "shuffle producer failed mid-flight".into(),
}),
});
match decoded.error {
Some(TypedClusterError::Internal { code, message }) => {
assert_eq!(code, 0xABCD);
assert!(message.contains("shuffle producer"));
}
other => panic!("expected Internal, got {other:?}"),
}
}
fn roundtrip_produce_req(req: ShuffleProduceRequest) -> ShuffleProduceRequest {
let rpc = RaftRpc::ShuffleProduceRequest(req);
let encoded = super::super::encode(&rpc).unwrap();
match super::super::decode(&encoded).unwrap() {
RaftRpc::ShuffleProduceRequest(r) => r,
other => panic!("expected ShuffleProduceRequest, got {other:?}"),
}
}
fn roundtrip_produce_resp(resp: ShuffleProduceResponse) -> ShuffleProduceResponse {
let rpc = RaftRpc::ShuffleProduceResponse(resp);
let encoded = super::super::encode(&rpc).unwrap();
match super::super::decode(&encoded).unwrap() {
RaftRpc::ShuffleProduceResponse(r) => r,
other => panic!("expected ShuffleProduceResponse, got {other:?}"),
}
}
#[test]
fn roundtrip_shuffle_produce_request() {
let req = ShuffleProduceRequest {
shuffle_id: 0x1234_5678_9ABC_DEF0,
side: 1,
num_parts: 4,
producer_count: 2,
keys: vec!["k".into(), "tenant.id".into()],
part_node_map: vec![
PartNodeEntry {
part: 0,
node_id: 7,
},
PartNodeEntry {
part: 1,
node_id: 9,
},
PartNodeEntry {
part: 2,
node_id: 7,
},
PartNodeEntry {
part: 3,
node_id: 9,
},
],
plan_bytes: vec![0xDE, 0xAD, 0xBE, 0xEF],
tenant_id: 42,
database_id: 3,
deadline_remaining_ms: 7000,
trace_id: [5u8; 16],
descriptor_versions: vec![DescriptorVersionEntry {
collection: "orders".into(),
version: 11,
}],
};
let decoded = roundtrip_produce_req(req.clone());
assert_eq!(decoded.shuffle_id, req.shuffle_id);
assert_eq!(decoded.side, 1);
assert_eq!(decoded.num_parts, 4);
assert_eq!(decoded.producer_count, 2);
assert_eq!(decoded.keys, vec!["k".to_string(), "tenant.id".to_string()]);
assert_eq!(decoded.part_node_map.len(), 4);
assert_eq!(decoded.part_node_map[2].part, 2);
assert_eq!(decoded.part_node_map[2].node_id, 7);
assert_eq!(decoded.plan_bytes, vec![0xDE, 0xAD, 0xBE, 0xEF]);
assert_eq!(decoded.tenant_id, 42);
assert_eq!(decoded.database_id, 3);
assert_eq!(decoded.deadline_remaining_ms, 7000);
assert_eq!(decoded.trace_id, [5u8; 16]);
assert_eq!(decoded.descriptor_versions.len(), 1);
assert_eq!(decoded.descriptor_versions[0].collection, "orders");
assert_eq!(decoded.descriptor_versions[0].version, 11);
}
#[test]
fn roundtrip_shuffle_produce_request_empty_keys_and_map() {
let req = ShuffleProduceRequest {
shuffle_id: 1,
side: 0,
num_parts: 1,
producer_count: 1,
keys: vec![],
part_node_map: vec![],
plan_bytes: vec![],
tenant_id: 0,
database_id: 0,
deadline_remaining_ms: 1000,
trace_id: [0u8; 16],
descriptor_versions: vec![],
};
let decoded = roundtrip_produce_req(req);
assert!(decoded.keys.is_empty());
assert!(decoded.part_node_map.is_empty());
assert!(decoded.descriptor_versions.is_empty());
}
#[test]
fn roundtrip_shuffle_produce_response_clean() {
let decoded = roundtrip_produce_resp(ShuffleProduceResponse {
error: None,
read_version_lsn: 0xABCD_1234,
});
assert!(decoded.error.is_none());
assert_eq!(
decoded.read_version_lsn, 0xABCD_1234,
"producer read-version LSN roundtrips on the produce reply"
);
}
#[test]
fn roundtrip_shuffle_produce_response_error() {
let decoded = roundtrip_produce_resp(ShuffleProduceResponse {
error: Some(TypedClusterError::Internal {
code: 0x55,
message: "produce scan failed".into(),
}),
read_version_lsn: 0,
});
assert_eq!(
decoded.read_version_lsn, 0,
"a failed produce carries no read-version LSN"
);
match decoded.error {
Some(TypedClusterError::Internal { code, message }) => {
assert_eq!(code, 0x55);
assert!(message.contains("produce scan"));
}
other => panic!("expected Internal, got {other:?}"),
}
}
fn roundtrip_consume_req(req: ShuffleConsumeRequest) -> ShuffleConsumeRequest {
let rpc = RaftRpc::ShuffleConsumeRequest(req);
let encoded = super::super::encode(&rpc).unwrap();
match super::super::decode(&encoded).unwrap() {
RaftRpc::ShuffleConsumeRequest(r) => r,
other => panic!("expected ShuffleConsumeRequest, got {other:?}"),
}
}
fn roundtrip_consume_resp(resp: ShuffleConsumeResponse) -> ShuffleConsumeResponse {
let rpc = RaftRpc::ShuffleConsumeResponse(resp);
let encoded = super::super::encode(&rpc).unwrap();
match super::super::decode(&encoded).unwrap() {
RaftRpc::ShuffleConsumeResponse(r) => r,
other => panic!("expected ShuffleConsumeResponse, got {other:?}"),
}
}
#[test]
fn roundtrip_shuffle_consume_request() {
let req = ShuffleConsumeRequest {
shuffle_id: 0x0FED_CBA9_8765_4321,
part: 3,
on: vec![
JoinKeyPair {
left: "lk".into(),
right: "rk".into(),
},
JoinKeyPair {
left: "tenant".into(),
right: "tenant_id".into(),
},
],
join_type: "inner".into(),
limit: 1234,
probe_qualifier: "l".into(),
index_qualifier: "r".into(),
tenant_id: 9,
database_id: 4,
deadline_remaining_ms: 8000,
trace_id: [3u8; 16],
};
let decoded = roundtrip_consume_req(req.clone());
assert_eq!(decoded.shuffle_id, req.shuffle_id);
assert_eq!(decoded.part, 3);
assert_eq!(decoded.on.len(), 2);
assert_eq!(decoded.on[0].left, "lk");
assert_eq!(decoded.on[0].right, "rk");
assert_eq!(decoded.on[1].left, "tenant");
assert_eq!(decoded.on[1].right, "tenant_id");
assert_eq!(decoded.join_type, "inner");
assert_eq!(decoded.limit, 1234);
assert_eq!(decoded.probe_qualifier, "l");
assert_eq!(decoded.index_qualifier, "r");
assert_eq!(decoded.tenant_id, 9);
assert_eq!(decoded.database_id, 4);
assert_eq!(decoded.deadline_remaining_ms, 8000);
assert_eq!(decoded.trace_id, [3u8; 16]);
}
#[test]
fn roundtrip_shuffle_consume_request_empty_keys() {
let req = ShuffleConsumeRequest {
shuffle_id: 1,
part: 0,
on: vec![],
join_type: "left".into(),
limit: u64::MAX,
probe_qualifier: String::new(),
index_qualifier: String::new(),
tenant_id: 0,
database_id: 0,
deadline_remaining_ms: 1000,
trace_id: [0u8; 16],
};
let decoded = roundtrip_consume_req(req);
assert!(decoded.on.is_empty());
assert_eq!(decoded.limit, u64::MAX);
assert_eq!(decoded.join_type, "left");
}
#[test]
fn roundtrip_shuffle_consume_response_rows() {
let decoded = roundtrip_consume_resp(ShuffleConsumeResponse {
rows: vec![0x92, 0x01, 0x02],
error: None,
});
assert_eq!(decoded.rows, vec![0x92, 0x01, 0x02]);
assert!(decoded.error.is_none());
}
#[test]
fn roundtrip_shuffle_consume_response_error() {
let decoded = roundtrip_consume_resp(ShuffleConsumeResponse {
rows: vec![],
error: Some(TypedClusterError::DeadlineExceeded { elapsed_ms: 8000 }),
});
assert!(decoded.rows.is_empty());
match decoded.error {
Some(TypedClusterError::DeadlineExceeded { elapsed_ms }) => {
assert_eq!(elapsed_ms, 8000);
}
other => panic!("expected DeadlineExceeded, got {other:?}"),
}
}
fn roundtrip_agg_consume_req(
req: ShuffleAggregateConsumeRequest,
) -> ShuffleAggregateConsumeRequest {
let rpc = RaftRpc::ShuffleAggregateConsumeRequest(req);
let encoded = super::super::encode(&rpc).unwrap();
match super::super::decode(&encoded).unwrap() {
RaftRpc::ShuffleAggregateConsumeRequest(r) => r,
other => panic!("expected ShuffleAggregateConsumeRequest, got {other:?}"),
}
}
fn roundtrip_agg_consume_resp(
resp: ShuffleAggregateConsumeResponse,
) -> ShuffleAggregateConsumeResponse {
let rpc = RaftRpc::ShuffleAggregateConsumeResponse(resp);
let encoded = super::super::encode(&rpc).unwrap();
match super::super::decode(&encoded).unwrap() {
RaftRpc::ShuffleAggregateConsumeResponse(r) => r,
other => panic!("expected ShuffleAggregateConsumeResponse, got {other:?}"),
}
}
#[test]
fn roundtrip_shuffle_agg_consume_request() {
let req = ShuffleAggregateConsumeRequest {
shuffle_id: 0x0FED_CBA9_8765_4321,
part: 2,
group_by: vec!["k".into(), "region".into()],
aggregates_bytes: vec![0x91, 0x01, 0x02],
having: vec![0xC0],
limit: 4321,
sort_keys: vec![
SortKey {
column: "k".into(),
ascending: true,
},
SortKey {
column: "total".into(),
ascending: false,
},
],
tenant_id: 9,
database_id: 4,
deadline_remaining_ms: 8000,
trace_id: [3u8; 16],
};
let decoded = roundtrip_agg_consume_req(req.clone());
assert_eq!(decoded.shuffle_id, req.shuffle_id);
assert_eq!(decoded.part, 2);
assert_eq!(
decoded.group_by,
vec!["k".to_string(), "region".to_string()]
);
assert_eq!(decoded.aggregates_bytes, vec![0x91, 0x01, 0x02]);
assert_eq!(decoded.having, vec![0xC0]);
assert_eq!(decoded.limit, 4321);
assert_eq!(decoded.sort_keys.len(), 2);
assert_eq!(decoded.sort_keys[0].column, "k");
assert!(decoded.sort_keys[0].ascending);
assert_eq!(decoded.sort_keys[1].column, "total");
assert!(!decoded.sort_keys[1].ascending);
assert_eq!(decoded.tenant_id, 9);
assert_eq!(decoded.database_id, 4);
assert_eq!(decoded.deadline_remaining_ms, 8000);
assert_eq!(decoded.trace_id, [3u8; 16]);
}
#[test]
fn roundtrip_shuffle_agg_consume_request_empty() {
let req = ShuffleAggregateConsumeRequest {
shuffle_id: 1,
part: 0,
group_by: vec![],
aggregates_bytes: vec![],
having: vec![],
limit: u64::MAX,
sort_keys: vec![],
tenant_id: 0,
database_id: 0,
deadline_remaining_ms: 1000,
trace_id: [0u8; 16],
};
let decoded = roundtrip_agg_consume_req(req);
assert!(decoded.group_by.is_empty());
assert!(decoded.aggregates_bytes.is_empty());
assert!(decoded.sort_keys.is_empty());
assert_eq!(decoded.limit, u64::MAX);
}
#[test]
fn roundtrip_shuffle_agg_consume_response_rows() {
let decoded = roundtrip_agg_consume_resp(ShuffleAggregateConsumeResponse {
rows: vec![0x92, 0x01, 0x02],
error: None,
});
assert_eq!(decoded.rows, vec![0x92, 0x01, 0x02]);
assert!(decoded.error.is_none());
}
#[test]
fn roundtrip_shuffle_agg_consume_response_error() {
let decoded = roundtrip_agg_consume_resp(ShuffleAggregateConsumeResponse {
rows: vec![],
error: Some(TypedClusterError::DeadlineExceeded { elapsed_ms: 9000 }),
});
assert!(decoded.rows.is_empty());
match decoded.error {
Some(TypedClusterError::DeadlineExceeded { elapsed_ms }) => {
assert_eq!(elapsed_ms, 9000);
}
other => panic!("expected DeadlineExceeded, got {other:?}"),
}
}
}