use std::sync::Arc;
use crate::rendezvous::protocol::{
AcquireResponse, RvAcquireRequest, RvDetachRequest, RvMetadataRequest, RvPullRequest,
RvRefRequest, RvReleaseRequest,
};
use crate::rendezvous::store::{DEFAULT_CHUNK_SIZE, DataStore};
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"))?;
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()
}
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_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) {
Some(expected_local_id) if expected_local_id == local_id => {
store.release_read_lock(local_id);
store.remove_transfers_by_lease(req.lease_id);
}
Some(expected_local_id) => {
tracing::warn!(
"_rv_detach: lease {} maps to slot {}, not {}",
req.lease_id,
expected_local_id,
local_id,
);
}
None => {
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) {
Some(expected_local_id) if expected_local_id == local_id => {
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}");
}
}
Some(expected_local_id) => {
tracing::warn!(
"_rv_release: lease {} maps to slot {}, not {}",
req.lease_id,
expected_local_id,
local_id,
);
}
None => {
tracing::warn!(
"_rv_release: invalid or already-consumed lease {} for {handle}",
req.lease_id,
);
}
}
Ok(())
})
.build()
}