nodedb_cluster/rpc_codec/
reservation.rs1use super::discriminants::*;
21use super::execute::TypedClusterError;
22use super::header::write_frame;
23use super::raft_rpc::RaftRpc;
24use crate::error::{ClusterError, Result};
25
26#[derive(Debug, Clone, rkyv::Archive, rkyv::Serialize, rkyv::Deserialize)]
40pub struct ReserveReadRequest {
41 pub lock_key_bytes: Vec<u8>,
43 pub vshard: u32,
44 pub owner_bytes: Option<Vec<u8>>,
48 pub deadline_remaining_ms: u64,
50 pub trace_id: [u8; 16],
51}
52
53#[derive(Debug, Clone, rkyv::Archive, rkyv::Serialize, rkyv::Deserialize)]
61pub struct ReserveReadResponse {
62 pub owner_bytes: Option<Vec<u8>>,
63 pub error: Option<TypedClusterError>,
64}
65
66#[derive(Debug, Clone, rkyv::Archive, rkyv::Serialize, rkyv::Deserialize)]
75pub struct ReleaseReservationRequest {
76 pub owner_bytes: Vec<u8>,
79 pub vshard: u32,
80 pub reason_bytes: Vec<u8>,
83 pub deadline_remaining_ms: u64,
85 pub trace_id: [u8; 16],
86}
87
88#[derive(Debug, Clone, rkyv::Archive, rkyv::Serialize, rkyv::Deserialize)]
94pub struct ReleaseReservationResponse {
95 pub error: Option<TypedClusterError>,
96}
97
98macro_rules! to_bytes {
101 ($msg:expr) => {
102 rkyv::to_bytes::<rkyv::rancor::Error>($msg)
103 .map(|b| b.to_vec())
104 .map_err(|e| ClusterError::Codec {
105 detail: format!("rkyv serialize: {e}"),
106 })
107 };
108}
109
110macro_rules! from_bytes {
111 ($payload:expr, $T:ty, $name:expr) => {{
112 let mut aligned = rkyv::util::AlignedVec::<16>::with_capacity($payload.len());
113 aligned.extend_from_slice($payload);
114 rkyv::from_bytes::<$T, rkyv::rancor::Error>(&aligned).map_err(|e| ClusterError::Codec {
115 detail: format!("rkyv deserialize {}: {e}", $name),
116 })
117 }};
118}
119
120pub(super) fn encode_reserve_read_req(msg: &ReserveReadRequest, out: &mut Vec<u8>) -> Result<()> {
121 write_frame(RPC_RESERVE_READ_REQ, &to_bytes!(msg)?, out)
122}
123pub(super) fn encode_reserve_read_resp(msg: &ReserveReadResponse, out: &mut Vec<u8>) -> Result<()> {
124 write_frame(RPC_RESERVE_READ_RESP, &to_bytes!(msg)?, out)
125}
126
127pub(super) fn decode_reserve_read_req(payload: &[u8]) -> Result<RaftRpc> {
128 Ok(RaftRpc::ReserveReadRequest(from_bytes!(
129 payload,
130 ReserveReadRequest,
131 "ReserveReadRequest"
132 )?))
133}
134pub(super) fn decode_reserve_read_resp(payload: &[u8]) -> Result<RaftRpc> {
135 Ok(RaftRpc::ReserveReadResponse(from_bytes!(
136 payload,
137 ReserveReadResponse,
138 "ReserveReadResponse"
139 )?))
140}
141
142pub(super) fn encode_release_reservation_req(
143 msg: &ReleaseReservationRequest,
144 out: &mut Vec<u8>,
145) -> Result<()> {
146 write_frame(RPC_RELEASE_RESERVATION_REQ, &to_bytes!(msg)?, out)
147}
148pub(super) fn encode_release_reservation_resp(
149 msg: &ReleaseReservationResponse,
150 out: &mut Vec<u8>,
151) -> Result<()> {
152 write_frame(RPC_RELEASE_RESERVATION_RESP, &to_bytes!(msg)?, out)
153}
154
155pub(super) fn decode_release_reservation_req(payload: &[u8]) -> Result<RaftRpc> {
156 Ok(RaftRpc::ReleaseReservationRequest(from_bytes!(
157 payload,
158 ReleaseReservationRequest,
159 "ReleaseReservationRequest"
160 )?))
161}
162pub(super) fn decode_release_reservation_resp(payload: &[u8]) -> Result<RaftRpc> {
163 Ok(RaftRpc::ReleaseReservationResponse(from_bytes!(
164 payload,
165 ReleaseReservationResponse,
166 "ReleaseReservationResponse"
167 )?))
168}
169
170#[cfg(test)]
171mod tests {
172 use super::*;
173
174 fn roundtrip_reserve_req(req: ReserveReadRequest) -> ReserveReadRequest {
175 let rpc = RaftRpc::ReserveReadRequest(req);
176 let encoded = super::super::encode(&rpc).unwrap();
177 match super::super::decode(&encoded).unwrap() {
178 RaftRpc::ReserveReadRequest(r) => r,
179 other => panic!("expected ReserveReadRequest, got {other:?}"),
180 }
181 }
182
183 fn roundtrip_reserve_resp(resp: ReserveReadResponse) -> ReserveReadResponse {
184 let rpc = RaftRpc::ReserveReadResponse(resp);
185 let encoded = super::super::encode(&rpc).unwrap();
186 match super::super::decode(&encoded).unwrap() {
187 RaftRpc::ReserveReadResponse(r) => r,
188 other => panic!("expected ReserveReadResponse, got {other:?}"),
189 }
190 }
191
192 fn roundtrip_release_req(req: ReleaseReservationRequest) -> ReleaseReservationRequest {
193 let rpc = RaftRpc::ReleaseReservationRequest(req);
194 let encoded = super::super::encode(&rpc).unwrap();
195 match super::super::decode(&encoded).unwrap() {
196 RaftRpc::ReleaseReservationRequest(r) => r,
197 other => panic!("expected ReleaseReservationRequest, got {other:?}"),
198 }
199 }
200
201 fn roundtrip_release_resp(resp: ReleaseReservationResponse) -> ReleaseReservationResponse {
202 let rpc = RaftRpc::ReleaseReservationResponse(resp);
203 let encoded = super::super::encode(&rpc).unwrap();
204 match super::super::decode(&encoded).unwrap() {
205 RaftRpc::ReleaseReservationResponse(r) => r,
206 other => panic!("expected ReleaseReservationResponse, got {other:?}"),
207 }
208 }
209
210 #[test]
211 fn roundtrip_reserve_read_request() {
212 let req = ReserveReadRequest {
213 lock_key_bytes: vec![0x01, 0x02, 0x03],
214 vshard: 4,
215 owner_bytes: None,
216 deadline_remaining_ms: 10_000,
217 trace_id: [9u8; 16],
218 };
219 let decoded = roundtrip_reserve_req(req);
220 assert_eq!(decoded.lock_key_bytes, vec![0x01, 0x02, 0x03]);
221 assert_eq!(decoded.vshard, 4);
222 assert!(decoded.owner_bytes.is_none());
223 assert_eq!(decoded.deadline_remaining_ms, 10_000);
224 assert_eq!(decoded.trace_id, [9u8; 16]);
225 }
226
227 #[test]
228 fn roundtrip_reserve_read_request_with_owner() {
229 let req = ReserveReadRequest {
230 lock_key_bytes: vec![],
231 vshard: 0,
232 owner_bytes: Some(vec![0xAA, 0xBB]),
233 deadline_remaining_ms: 0,
234 trace_id: [0u8; 16],
235 };
236 let decoded = roundtrip_reserve_req(req);
237 assert_eq!(decoded.owner_bytes, Some(vec![0xAA, 0xBB]));
238 }
239
240 #[test]
241 fn roundtrip_reserve_read_response_ok() {
242 let decoded = roundtrip_reserve_resp(ReserveReadResponse {
243 owner_bytes: Some(vec![0x0a, 0x0b]),
244 error: None,
245 });
246 assert_eq!(decoded.owner_bytes, Some(vec![0x0a, 0x0b]));
247 assert!(decoded.error.is_none());
248 }
249
250 #[test]
251 fn roundtrip_reserve_read_response_error() {
252 let decoded = roundtrip_reserve_resp(ReserveReadResponse {
253 owner_bytes: None,
254 error: Some(TypedClusterError::Internal {
255 code: 0,
256 message: "reserve-read not configured".into(),
257 }),
258 });
259 assert!(decoded.owner_bytes.is_none());
260 match decoded.error {
261 Some(TypedClusterError::Internal { code, message }) => {
262 assert_eq!(code, 0);
263 assert!(message.contains("reserve-read"));
264 }
265 other => panic!("expected Internal, got {other:?}"),
266 }
267 }
268
269 #[test]
270 fn roundtrip_release_reservation_request() {
271 let req = ReleaseReservationRequest {
272 owner_bytes: vec![0x01, 0x02],
273 vshard: 2,
274 reason_bytes: vec![0x03],
275 deadline_remaining_ms: 5_000,
276 trace_id: [3u8; 16],
277 };
278 let decoded = roundtrip_release_req(req);
279 assert_eq!(decoded.owner_bytes, vec![0x01, 0x02]);
280 assert_eq!(decoded.vshard, 2);
281 assert_eq!(decoded.reason_bytes, vec![0x03]);
282 assert_eq!(decoded.deadline_remaining_ms, 5_000);
283 assert_eq!(decoded.trace_id, [3u8; 16]);
284 }
285
286 #[test]
287 fn roundtrip_release_reservation_response_ok() {
288 let decoded = roundtrip_release_resp(ReleaseReservationResponse { error: None });
289 assert!(decoded.error.is_none());
290 }
291
292 #[test]
293 fn roundtrip_release_reservation_response_error() {
294 let decoded = roundtrip_release_resp(ReleaseReservationResponse {
295 error: Some(TypedClusterError::Internal {
296 code: 0,
297 message: "release-reservation not configured".into(),
298 }),
299 });
300 match decoded.error {
301 Some(TypedClusterError::Internal { code, message }) => {
302 assert_eq!(code, 0);
303 assert!(message.contains("release-reservation"));
304 }
305 other => panic!("expected Internal, got {other:?}"),
306 }
307 }
308}