use std::sync::Arc;
use crate::rendezvous::protocol::{
AcquireResponse, RvAcquireRequest, RvDetachRequest, RvLeaseRenewRequest, RvMetadataRequest,
RvPullRequest, RvRefRequest, RvReleaseRequest,
};
use crate::rendezvous::store::{DEFAULT_CHUNK_SIZE, DataStore, LeaseOutcome};
pub fn create_rv_metadata_handler(store: Arc<DataStore>) -> crate::messenger::Handler {
crate::messenger::Handler::typed_unary(
"_rv_metadata",
move |ctx: crate::messenger::TypedContext<RvMetadataRequest>| {
let handle = ctx.input.handle.to_handle();
let (_, local_id) = handle.unpack();
match store.metadata(local_id) {
Some(meta) => Ok(meta),
None => anyhow::bail!("rendezvous handle not found: {handle}"),
}
},
)
.build()
}
pub fn create_rv_acquire_handler(store: Arc<DataStore>) -> crate::messenger::Handler {
crate::messenger::Handler::typed_unary(
"_rv_acquire",
move |ctx: crate::messenger::TypedContext<RvAcquireRequest>| {
let handle = ctx.input.handle.to_handle();
let (_, local_id) = handle.unpack();
let lease_id = store
.acquire_read_lock(local_id)
.ok_or_else(|| anyhow::anyhow!("rendezvous handle not found: {handle}"))?;
let total_len = store
.get_total_len(local_id)
.ok_or_else(|| anyhow::anyhow!("slot vanished after lock acquire"))?;
#[cfg(all(target_os = "linux", feature = "ucx"))]
if let Some(response) = rdma_response(
&store,
local_id,
lease_id,
total_len,
ctx.input.rdma.as_ref(),
) {
return Ok(response);
}
let (transfer_id, chunk_size, chunk_count) = store
.create_transfer(local_id, lease_id, DEFAULT_CHUNK_SIZE)
.ok_or_else(|| anyhow::anyhow!("slot vanished after lock acquire"))?;
Ok(AcquireResponse::Ready {
lease_id,
transfer_id,
total_len,
chunk_size,
chunk_count,
})
},
)
.build()
}
#[cfg(all(target_os = "linux", feature = "ucx"))]
fn rdma_response(
store: &Arc<DataStore>,
local_id: u64,
lease_id: u64,
total_len: u64,
offer: Option<&crate::rendezvous::protocol::RdmaOffer>,
) -> Option<AcquireResponse> {
use crate::observability::RdmaPathReason;
use crate::rendezvous::store::StageMode;
let decline = |reason: RdmaPathReason| -> Option<AcquireResponse> {
store.record_path(reason);
None
};
let Some(offer) = offer else {
return decline(RdmaPathReason::NoOffer);
};
let Some(rdma) = store.rdma() else {
return decline(RdmaPathReason::NotConfigured);
};
if !rdma.config.enabled {
return decline(RdmaPathReason::KillSwitch);
}
if store.stage_mode(local_id) != Some(StageMode::Pinned) {
return decline(RdmaPathReason::NotPinned);
}
if total_len < rdma.config.rdma_min_bytes {
return decline(RdmaPathReason::BelowMin);
}
let backend = rdma.backend;
if !offer.backends.iter().any(|name| name == backend.key()) {
return decline(RdmaPathReason::NoOffer);
}
let descriptor = store
.with_pinned(local_id, |slot| {
(slot.backend() == backend).then(|| slot.descriptor())?
})
.flatten();
let Some(bytes) = descriptor.and_then(|d| d.encode()) else {
return decline(RdmaPathReason::NotPinned);
};
let lease_timeout_ms = u64::try_from(rdma.config.lease_timeout.as_millis()).unwrap_or(u64::MAX);
if lease_timeout_ms != 0 {
store.set_lease_deadline(lease_id, local_id, rdma.config.lease_timeout);
}
store.record_path(RdmaPathReason::Ok);
Some(AcquireResponse::Rdma {
lease_id,
descriptor: bytes,
lease_timeout_ms,
})
}
pub fn create_rv_pull_handler(store: Arc<DataStore>) -> crate::messenger::Handler {
crate::messenger::Handler::unary_handler(
"_rv_pull",
move |ctx: crate::messenger::Context| -> crate::messenger::UnifiedResponse {
let req: RvPullRequest = serde_json::from_slice(&ctx.payload)?;
match store.get_chunk(req.transfer_id, req.chunk_index) {
Some(chunk) => Ok(Some(chunk)),
None => anyhow::bail!(
"chunk not found: transfer_id={}, chunk_index={}",
req.transfer_id,
req.chunk_index
),
}
},
)
.build()
}
pub fn create_rv_ref_handler(store: Arc<DataStore>) -> crate::messenger::Handler {
crate::messenger::Handler::unary_handler(
"_rv_ref",
move |ctx: crate::messenger::Context| -> crate::messenger::UnifiedResponse {
let req: RvRefRequest = serde_json::from_slice(&ctx.payload)?;
let handle = req.handle.to_handle();
let (_, local_id) = handle.unpack();
if !store.ref_increment(local_id) {
anyhow::bail!("_rv_ref: handle not found: {handle}");
}
Ok(None)
},
)
.build()
}
pub fn create_rv_lease_renew_handler(store: Arc<DataStore>) -> crate::messenger::Handler {
crate::messenger::Handler::am_handler(
"_rv_lease_renew",
move |ctx: crate::messenger::Context| {
let req: RvLeaseRenewRequest = serde_json::from_slice(&ctx.payload)?;
let handle = req.handle.to_handle();
let (_, local_id) = handle.unpack();
match store.lease_slot(req.lease_id) {
Some(expected_local_id) if expected_local_id == local_id => {
if !store.renew_lease(req.lease_id) {
tracing::debug!(
lease = req.lease_id,
%handle,
"_rv_lease_renew: lease carries no deadline to renew"
);
}
}
Some(expected_local_id) => {
tracing::warn!(
"_rv_lease_renew: lease {} maps to slot {}, not {}",
req.lease_id,
expected_local_id,
local_id,
);
}
None => {
tracing::debug!(
lease = req.lease_id,
%handle,
"_rv_lease_renew: lease already ended"
);
}
}
Ok(())
},
)
.build()
}
pub fn create_rv_detach_handler(store: Arc<DataStore>) -> crate::messenger::Handler {
crate::messenger::Handler::am_handler("_rv_detach", move |ctx: crate::messenger::Context| {
let req: RvDetachRequest = serde_json::from_slice(&ctx.payload)?;
let handle = req.handle.to_handle();
let (_, local_id) = handle.unpack();
match store.consume_lease(req.lease_id, local_id) {
LeaseOutcome::Consumed => {
store.release_read_lock(local_id);
store.remove_transfers_by_lease(req.lease_id);
}
LeaseOutcome::Mismatch { actual } => {
tracing::warn!(
"_rv_detach: lease {} maps to slot {}, not {}; ignored",
req.lease_id,
actual,
local_id,
);
}
LeaseOutcome::Unknown => {
tracing::warn!(
"_rv_detach: invalid or already-consumed lease {} for {handle}",
req.lease_id,
);
}
}
Ok(())
})
.build()
}
pub fn create_rv_release_handler(store: Arc<DataStore>) -> crate::messenger::Handler {
crate::messenger::Handler::am_handler("_rv_release", move |ctx: crate::messenger::Context| {
let req: RvReleaseRequest = serde_json::from_slice(&ctx.payload)?;
let handle = req.handle.to_handle();
let (_, local_id) = handle.unpack();
match store.consume_lease(req.lease_id, local_id) {
LeaseOutcome::Consumed => {
store.release_read_lock(local_id);
store.remove_transfers_by_lease(req.lease_id);
let should_free = store.ref_decrement(local_id);
if should_free {
store.try_free(local_id);
tracing::debug!("_rv_release: freed slot for {handle}");
}
}
LeaseOutcome::Mismatch { actual } => {
tracing::warn!(
"_rv_release: lease {} maps to slot {}, not {}; ignored",
req.lease_id,
actual,
local_id,
);
}
LeaseOutcome::Unknown => {
tracing::warn!(
"_rv_release: invalid or already-consumed lease {} for {handle}",
req.lease_id,
);
}
}
Ok(())
})
.build()
}