use serde::{Deserialize, Serialize};
#[derive(Debug, Serialize, Deserialize)]
pub struct RvMetadataRequest {
pub handle: RvHandleWire,
}
#[derive(Debug, Serialize, Deserialize)]
pub struct DataMetadata {
pub total_len: u64,
pub refcount: u32,
pub pinned: bool,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
pub struct RdmaOffer {
pub backends: Vec<String>,
}
#[derive(Debug, Serialize, Deserialize)]
pub struct RvAcquireRequest {
pub handle: RvHandleWire,
#[serde(default)]
pub rdma: Option<RdmaOffer>,
}
#[derive(Debug, Serialize, Deserialize)]
pub enum AcquireResponse {
Ready {
lease_id: u64,
transfer_id: u64,
total_len: u64,
chunk_size: u32,
chunk_count: u32,
},
Rdma {
lease_id: u64,
descriptor: Vec<u8>,
#[serde(default)]
lease_timeout_ms: u64,
},
}
#[derive(Debug, Serialize, Deserialize)]
pub struct RvPullRequest {
pub transfer_id: u64,
pub chunk_index: u32,
}
#[derive(Debug, Serialize, Deserialize)]
pub struct RvRefRequest {
pub handle: RvHandleWire,
}
#[derive(Debug, Serialize, Deserialize)]
pub struct RvDetachRequest {
pub handle: RvHandleWire,
pub lease_id: u64,
}
#[derive(Debug, Serialize, Deserialize)]
pub struct RvReleaseRequest {
pub handle: RvHandleWire,
pub lease_id: u64,
}
#[derive(Debug, Serialize, Deserialize)]
pub struct RvLeaseRenewRequest {
pub handle: RvHandleWire,
pub lease_id: u64,
}
#[derive(Debug, Clone, Copy, Serialize, Deserialize)]
pub struct RvHandleWire {
pub hi: u64,
pub lo: u64,
}
impl RvHandleWire {
pub fn from_handle(handle: crate::DataHandle) -> Self {
let raw = handle.as_u128();
Self {
hi: (raw >> 64) as u64,
lo: raw as u64,
}
}
pub fn to_handle(self) -> crate::DataHandle {
crate::DataHandle::from_u128(((self.hi as u128) << 64) | (self.lo as u128))
}
}
#[derive(Debug, Serialize, Deserialize)]
pub struct RvError {
pub message: String,
}
#[cfg(test)]
mod tests {
use super::*;
mod old_wire {
use serde::{Deserialize, Serialize};
#[derive(Debug, Serialize, Deserialize)]
pub struct RvHandleWire {
pub hi: u64,
pub lo: u64,
}
#[derive(Debug, Serialize, Deserialize)]
pub struct RvAcquireRequest {
pub handle: RvHandleWire,
}
#[derive(Debug, Serialize, Deserialize)]
pub enum AcquireResponse {
Ready {
lease_id: u64,
transfer_id: u64,
total_len: u64,
chunk_size: u32,
chunk_count: u32,
},
Rdma {
lease_id: u64,
descriptor: Vec<u8>,
},
}
}
fn handle() -> RvHandleWire {
RvHandleWire { hi: 7, lo: 42 }
}
#[test]
fn an_old_consumers_acquire_carries_no_offer() {
let old = serde_json::to_vec(&old_wire::RvAcquireRequest {
handle: old_wire::RvHandleWire { hi: 7, lo: 42 },
})
.expect("serialize the old shape");
assert!(
!String::from_utf8_lossy(&old).contains("rdma"),
"the old shape must not mention the field this test is about"
);
let new: RvAcquireRequest = serde_json::from_slice(&old).expect("deserialize");
assert!(new.rdma.is_none());
assert_eq!(new.handle.hi, 7);
assert_eq!(new.handle.lo, 42);
}
#[test]
fn an_old_owner_ignores_a_new_consumers_offer() {
let new = serde_json::to_vec(&RvAcquireRequest {
handle: handle(),
rdma: Some(RdmaOffer {
backends: vec!["ucx".to_string()],
}),
})
.expect("serialize");
let old: old_wire::RvAcquireRequest =
serde_json::from_slice(&new).expect("an old owner must still parse a new acquire");
assert_eq!(old.handle.hi, 7);
assert_eq!(old.handle.lo, 42);
}
#[test]
fn an_explicit_null_offer_is_no_offer() {
let json = br#"{"handle":{"hi":7,"lo":42},"rdma":null}"#;
let req: RvAcquireRequest = serde_json::from_slice(json).expect("deserialize");
assert!(req.rdma.is_none());
}
#[test]
fn an_offer_of_unknown_backends_still_parses() {
let json = br#"{"handle":{"hi":1,"lo":2},"rdma":{"backends":["nixl","libfabric"]}}"#;
let req: RvAcquireRequest = serde_json::from_slice(json).expect("deserialize");
assert_eq!(
req.rdma.expect("offer").backends,
vec!["nixl".to_string(), "libfabric".to_string()]
);
}
#[test]
fn an_rdma_response_without_a_lease_timeout_means_no_deadline() {
let old = serde_json::to_vec(&old_wire::AcquireResponse::Rdma {
lease_id: 9,
descriptor: vec![1, 2, 3],
})
.expect("serialize the old shape");
let new: AcquireResponse = serde_json::from_slice(&old).expect("deserialize");
match new {
AcquireResponse::Rdma {
lease_id,
descriptor,
lease_timeout_ms,
} => {
assert_eq!(lease_id, 9);
assert_eq!(descriptor, vec![1, 2, 3]);
assert_eq!(lease_timeout_ms, 0, "a missing timeout must mean none");
}
other => panic!("expected Rdma, got {other:?}"),
}
}
#[test]
fn a_new_rdma_response_is_readable_by_an_old_consumer() {
let new = serde_json::to_vec(&AcquireResponse::Rdma {
lease_id: 9,
descriptor: vec![1, 2, 3],
lease_timeout_ms: 30_000,
})
.expect("serialize");
let old: old_wire::AcquireResponse = serde_json::from_slice(&new).expect("deserialize");
match old {
old_wire::AcquireResponse::Rdma {
lease_id,
descriptor,
} => {
assert_eq!(lease_id, 9);
assert_eq!(descriptor, vec![1, 2, 3]);
}
other => panic!("expected Rdma, got {other:?}"),
}
}
#[test]
fn the_ready_variant_is_unchanged_in_both_directions() {
let new = serde_json::to_vec(&AcquireResponse::Ready {
lease_id: 1,
transfer_id: 2,
total_len: 3,
chunk_size: 4,
chunk_count: 5,
})
.expect("serialize");
let old: old_wire::AcquireResponse = serde_json::from_slice(&new).expect("old reads new");
assert!(matches!(old, old_wire::AcquireResponse::Ready { .. }));
let old_bytes = serde_json::to_vec(&old_wire::AcquireResponse::Ready {
lease_id: 1,
transfer_id: 2,
total_len: 3,
chunk_size: 4,
chunk_count: 5,
})
.expect("serialize");
let new: AcquireResponse = serde_json::from_slice(&old_bytes).expect("new reads old");
assert!(matches!(
new,
AcquireResponse::Ready {
lease_id: 1,
transfer_id: 2,
total_len: 3,
chunk_size: 4,
chunk_count: 5,
}
));
}
}