Skip to main content

nodedb_cluster/rpc_codec/
reservation.rs

1// SPDX-License-Identifier: BUSL-1.1
2
3//! Routed reserve-read / release-reservation one-shot RPCs (Calvin OLLP).
4//!
5//! Assign-only reserve and ack-only release for the Calvin OLLP dependent-read
6//! path: a coordinator that wants to reserve or release a read lock on the
7//! SEQUENCER-GROUP LEADER sends one of these requests; the leader mutates its
8//! local scheduler state and replies with exactly one response frame. Both are
9//! one-shot request/response — no streaming.
10//!
11//! The domain payloads (`LockKeyWire`, `TxnIdWire`, `ReleaseReason` — see
12//! [`crate::calvin::types::lock_wire`]) are msgpack-only, not rkyv. They ride
13//! these rkyv envelope structs as opaque pre-encoded bytes (mirroring
14//! `tx_class_bytes` in [`super::calvin_submit`]) and are decoded at the
15//! hook boundary in the host crate.
16//!
17//! Discriminants 41/42 (`ReserveRead`) and 43/44 (`ReleaseReservation`) are
18//! permanently assigned to these variants.
19
20use super::discriminants::*;
21use super::execute::TypedClusterError;
22use super::header::write_frame;
23use super::raft_rpc::RaftRpc;
24use crate::error::{ClusterError, Result};
25
26// ── Wire types ──────────────────────────────────────────────────────────────
27
28/// Coordinator → sequencer-leader routed reserve-read request.
29///
30/// Carries the `LockKey` as opaque msgpack bytes (`lock_key_bytes`); the
31/// leader decodes it and assign-only reserves the read lock, minting an owner
32/// if `owner_bytes` is `None`. The `deadline_remaining_ms` / `trace_id` fields
33/// mirror [`SubmitCalvinInboxRequest`](super::calvin_submit::SubmitCalvinInboxRequest)
34/// so the leader-side handler shares the same deadline / tracing prologue as
35/// the other one-shot RPCs; the leader bounds the reserve by
36/// `deadline_remaining_ms`.
37///
38/// Cross-version safety: new optional fields should be added as `Option<T>`.
39#[derive(Debug, Clone, rkyv::Archive, rkyv::Serialize, rkyv::Deserialize)]
40pub struct ReserveReadRequest {
41    /// The `LockKey` to reserve, encoded with `zerompk::to_msgpack_vec`.
42    pub lock_key_bytes: Vec<u8>,
43    pub vshard: u32,
44    /// Pre-assigned owner (`TxnIdWire`, msgpack-encoded), when the reservation
45    /// is being made on behalf of an already-known owner. `None` means the
46    /// leader mints a new owner.
47    pub owner_bytes: Option<Vec<u8>>,
48    /// Deadline budget remaining for the reserve on the leader (ms).
49    pub deadline_remaining_ms: u64,
50    pub trace_id: [u8; 16],
51}
52
53/// Terminal reply to a [`ReserveReadRequest`].
54///
55/// `error: None` means the leader reserved the read lock; `owner_bytes`
56/// carries the minted (or confirmed) owner (`TxnIdWire`, msgpack-encoded).
57/// `error: Some(e)` means the reserve failed (lock conflict, the `LockKey`
58/// failed to decode, or no reserve-read hook is configured); `owner_bytes` is
59/// `None` in that case.
60#[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/// Coordinator → sequencer-leader routed release-reservation request.
67///
68/// Carries the owner (`TxnIdWire`) and release reason (`ReleaseReason`) as
69/// opaque msgpack bytes; the leader decodes both and releases the
70/// reservation. Ack-only — there is no data payload to return beyond success
71/// or a typed error.
72///
73/// Cross-version safety: new optional fields should be added as `Option<T>`.
74#[derive(Debug, Clone, rkyv::Archive, rkyv::Serialize, rkyv::Deserialize)]
75pub struct ReleaseReservationRequest {
76    /// The owner (`TxnIdWire`) releasing the reservation, encoded with
77    /// `zerompk::to_msgpack_vec`.
78    pub owner_bytes: Vec<u8>,
79    pub vshard: u32,
80    /// The release reason (`ReleaseReason`), encoded with
81    /// `zerompk::to_msgpack_vec`.
82    pub reason_bytes: Vec<u8>,
83    /// Deadline budget remaining for the release on the leader (ms).
84    pub deadline_remaining_ms: u64,
85    pub trace_id: [u8; 16],
86}
87
88/// Terminal reply to a [`ReleaseReservationRequest`].
89///
90/// `error: None` means the leader released the reservation (ack). `error:
91/// Some(e)` means the release failed (unknown owner, either payload failed to
92/// decode, or no release-reservation hook is configured).
93#[derive(Debug, Clone, rkyv::Archive, rkyv::Serialize, rkyv::Deserialize)]
94pub struct ReleaseReservationResponse {
95    pub error: Option<TypedClusterError>,
96}
97
98// ── Codec ────────────────────────────────────────────────────────────────────
99
100macro_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}