pub mod consumer;
pub mod handle;
pub mod handlers;
pub mod protocol;
pub mod store;
pub mod transparent;
pub mod write;
pub use handle::DataHandle;
pub use protocol::DataMetadata;
pub use store::{RegisterOptions, StageMode};
pub use transparent::{RendezvousResolver, RendezvousStager};
pub use write::RendezvousWrite;
use std::sync::{Arc, OnceLock};
use std::time::Instant;
use crate::observability::{HandlerOutcome, RendezvousOp, VeloMetrics};
use anyhow::Result;
use bytes::Bytes;
use velo_ext::WorkerId;
pub struct RendezvousManager {
worker_id: WorkerId,
store: Arc<store::DataStore>,
messenger_lock: OnceLock<Arc<crate::messenger::Messenger>>,
metrics: Option<Arc<VeloMetrics>>,
}
impl RendezvousManager {
pub fn new(worker_id: WorkerId) -> Self {
Self {
worker_id,
store: Arc::new(store::DataStore::new()),
messenger_lock: OnceLock::new(),
metrics: None,
}
}
pub fn with_metrics(worker_id: WorkerId, metrics: Arc<VeloMetrics>) -> Self {
Self {
worker_id,
store: Arc::new(store::DataStore::new()),
messenger_lock: OnceLock::new(),
metrics: Some(metrics),
}
}
pub fn register_handlers(
self: &Arc<Self>,
messenger: Arc<crate::messenger::Messenger>,
) -> Result<()> {
use handlers::{
create_rv_acquire_handler, create_rv_detach_handler, create_rv_metadata_handler,
create_rv_pull_handler, create_rv_ref_handler, create_rv_release_handler,
};
messenger
.register_streaming_handler(create_rv_metadata_handler(Arc::clone(&self.store)))?;
messenger.register_streaming_handler(create_rv_acquire_handler(Arc::clone(&self.store)))?;
messenger.register_streaming_handler(create_rv_pull_handler(Arc::clone(&self.store)))?;
messenger.register_streaming_handler(create_rv_ref_handler(Arc::clone(&self.store)))?;
messenger.register_streaming_handler(create_rv_detach_handler(Arc::clone(&self.store)))?;
messenger.register_streaming_handler(create_rv_release_handler(Arc::clone(&self.store)))?;
self.messenger_lock
.set(messenger)
.map_err(|_| anyhow::anyhow!("register_handlers called twice"))?;
Ok(())
}
fn messenger(&self) -> &Arc<crate::messenger::Messenger> {
self.messenger_lock
.get()
.expect("RendezvousManager::register_handlers must be called before use")
}
pub fn register_data(&self, data: Bytes) -> DataHandle {
let started = Instant::now();
let data_len = data.len();
let local_id = self.store.register(data, None);
if let Some(m) = &self.metrics {
m.record_rendezvous_operation(
RendezvousOp::Register,
HandlerOutcome::Success,
started.elapsed(),
);
m.record_rendezvous_bytes(RendezvousOp::Register, data_len);
m.set_rendezvous_active_slots(self.store.slots.len());
}
DataHandle::pack(self.worker_id, local_id)
}
pub fn register_data_with(&self, data: Bytes, opts: RegisterOptions) -> DataHandle {
let started = Instant::now();
let data_len = data.len();
let local_id = self.store.register(data, Some(opts));
if let Some(m) = &self.metrics {
m.record_rendezvous_operation(
RendezvousOp::Register,
HandlerOutcome::Success,
started.elapsed(),
);
m.record_rendezvous_bytes(RendezvousOp::Register, data_len);
m.set_rendezvous_active_slots(self.store.slots.len());
}
DataHandle::pack(self.worker_id, local_id)
}
pub async fn metadata(&self, handle: DataHandle) -> Result<DataMetadata> {
let started = Instant::now();
let (target_worker, local_id) = handle.unpack();
let result = if target_worker == self.worker_id {
self.store
.metadata(local_id)
.ok_or_else(|| anyhow::anyhow!("rendezvous handle not found: {handle}"))
} else {
consumer::Consumer::metadata(self.messenger(), handle).await
};
if let Some(m) = &self.metrics {
let outcome = if result.is_ok() {
HandlerOutcome::Success
} else {
HandlerOutcome::Error
};
m.record_rendezvous_operation(RendezvousOp::Metadata, outcome, started.elapsed());
}
result
}
pub async fn get(&self, handle: DataHandle) -> Result<(Bytes, u64)> {
let started = Instant::now();
let (target_worker, local_id) = handle.unpack();
let result = if target_worker == self.worker_id {
let lease_id = self
.store
.acquire_read_lock(local_id)
.ok_or_else(|| anyhow::anyhow!("rendezvous handle not found: {handle}"))?;
let data = self
.store
.get_data(local_id)
.ok_or_else(|| anyhow::anyhow!("slot vanished after lock acquire"))?;
Ok((data, lease_id))
} else {
consumer::Consumer::get(self.messenger(), handle).await
};
if let Some(m) = &self.metrics {
let outcome = if result.is_ok() {
HandlerOutcome::Success
} else {
HandlerOutcome::Error
};
m.record_rendezvous_operation(RendezvousOp::Get, outcome, started.elapsed());
if let Ok((ref data, _)) = result {
m.record_rendezvous_bytes(RendezvousOp::Get, data.len());
}
}
result
}
pub async fn get_into(
&self,
handle: DataHandle,
dest: &mut impl RendezvousWrite,
) -> Result<u64> {
let started = Instant::now();
let (target_worker, local_id) = handle.unpack();
let result = if target_worker == self.worker_id {
let lease_id = self
.store
.acquire_read_lock(local_id)
.ok_or_else(|| anyhow::anyhow!("rendezvous handle not found: {handle}"))?;
let data = self
.store
.get_data(local_id)
.ok_or_else(|| anyhow::anyhow!("slot vanished after lock acquire"))?;
dest.write_chunk(0, &data)?;
Ok(lease_id)
} else {
consumer::Consumer::get_into(self.messenger(), handle, dest).await
};
if let Some(m) = &self.metrics {
let outcome = if result.is_ok() {
HandlerOutcome::Success
} else {
HandlerOutcome::Error
};
m.record_rendezvous_operation(RendezvousOp::Get, outcome, started.elapsed());
}
result
}
pub async fn ref_handle(&self, handle: DataHandle) -> Result<()> {
let started = Instant::now();
let (target_worker, local_id) = handle.unpack();
let result = if target_worker == self.worker_id {
if !self.store.ref_increment(local_id) {
anyhow::bail!("rendezvous handle not found: {handle}");
}
Ok(())
} else {
consumer::Consumer::ref_handle(self.messenger(), handle).await
};
if let Some(m) = &self.metrics {
let outcome = if result.is_ok() {
HandlerOutcome::Success
} else {
HandlerOutcome::Error
};
m.record_rendezvous_operation(RendezvousOp::Ref, outcome, started.elapsed());
}
result
}
pub async fn detach(&self, handle: DataHandle, lease_id: u64) -> Result<()> {
let started = Instant::now();
let (target_worker, local_id) = handle.unpack();
let result = if target_worker == self.worker_id {
match self.store.consume_lease(lease_id) {
Some(expected) if expected == local_id => {
self.store.release_read_lock(local_id);
self.store.remove_transfers_by_lease(lease_id);
Ok(())
}
Some(_) | None => {
anyhow::bail!("invalid or already-consumed lease {lease_id} for {handle}")
}
}
} else {
consumer::Consumer::detach(self.messenger(), handle, lease_id).await
};
if let Some(m) = &self.metrics {
let outcome = if result.is_ok() {
HandlerOutcome::Success
} else {
HandlerOutcome::Error
};
m.record_rendezvous_operation(RendezvousOp::Detach, outcome, started.elapsed());
}
result
}
pub async fn release(&self, handle: DataHandle, lease_id: u64) -> Result<()> {
let started = Instant::now();
let (target_worker, local_id) = handle.unpack();
let result = if target_worker == self.worker_id {
match self.store.consume_lease(lease_id) {
Some(expected) if expected == local_id => {
self.store.release_read_lock(local_id);
self.store.remove_transfers_by_lease(lease_id);
let should_free = self.store.ref_decrement(local_id);
if should_free {
self.store.try_free(local_id);
}
Ok(())
}
Some(_) | None => {
anyhow::bail!("invalid or already-consumed lease {lease_id} for {handle}")
}
}
} else {
consumer::Consumer::release(self.messenger(), handle, lease_id).await
};
if let Some(m) = &self.metrics {
let outcome = if result.is_ok() {
HandlerOutcome::Success
} else {
HandlerOutcome::Error
};
m.record_rendezvous_operation(RendezvousOp::Release, outcome, started.elapsed());
m.set_rendezvous_active_slots(self.store.slots.len());
}
result
}
pub fn worker_id(&self) -> WorkerId {
self.worker_id
}
pub fn data_store(&self) -> &Arc<store::DataStore> {
&self.store
}
}