use core::{ptr, ptr::NonNull};
use super::{EndpointLeaseId, EndpointResidentBudget, RendezvousEntry, RendezvousTable};
use crate::{
rendezvous::core::Rendezvous,
session::{
cluster::error::ClusterError,
types::{RendezvousId, SessionId},
},
transport::Transport,
};
impl<'cfg, T> RendezvousTable<'cfg, T>
where
T: Transport,
{
fn contains_key(&self, key: &RendezvousId) -> bool {
self.entry_ref(key).is_some()
}
fn next_available_rendezvous_id(&self) -> Option<RendezvousId> {
let mut raw = 1u16;
while raw != 0 {
let id = RendezvousId::new(raw);
if !self.contains_key(&id) {
return Some(id);
}
raw = raw.wrapping_add(1);
}
None
}
pub(crate) fn register_local_from_resources_auto(
&mut self,
resources: crate::runtime_core::resources::RuntimeResources<'cfg>,
transport: T,
) -> Result<RendezvousId, RegisterRendezvousError> {
let id = self
.next_available_rendezvous_id()
.ok_or(RegisterRendezvousError::CapacityExceeded)?;
let rendezvous = unsafe {
Rendezvous::init_in_slab_auto(id, resources, transport)
.ok_or(RegisterRendezvousError::StorageExhausted)?
};
let entry_ptr = unsafe {
let rv = &mut *rendezvous;
match rv.allocate_external_persistent_sidecar_bytes(
core::mem::size_of::<RendezvousEntry<'cfg, T>>(),
core::mem::align_of::<RendezvousEntry<'cfg, T>>(),
) {
Some(sidecar) => sidecar.ptr().cast::<RendezvousEntry<'cfg, T>>(),
None => {
ptr::drop_in_place(rendezvous);
return Err(RegisterRendezvousError::StorageExhausted);
}
}
};
unsafe {
RendezvousEntry::init_from_parts(
entry_ptr,
id,
NonNull::new_unchecked(rendezvous),
self.head,
);
}
self.head = NonNull::new(entry_ptr);
self.len = self
.len
.checked_add(1)
.ok_or(RegisterRendezvousError::CapacityExceeded)?;
Ok(id)
}
pub(crate) fn allocate_endpoint_lease_for_session_role(
&mut self,
rv: RendezvousId,
sid: SessionId,
role: u8,
bytes: usize,
align: usize,
resident_budget: EndpointResidentBudget,
) -> Result<(EndpointLeaseId, u32, usize, usize), ClusterError> {
if role >= crate::g::ROLE_DOMAIN_SIZE {
crate::invariant();
}
let mut target = None;
let mut current = self.head;
while let Some(mut entry_ptr) = current {
let entry = unsafe {
entry_ptr.as_mut()
};
if entry.id == rv {
target = Some(entry_ptr);
}
if entry.is_active() {
return Err(ClusterError::RendezvousBusy { id: entry.id.raw() });
}
if crate::invariant_some(entry.rendezvous_ref())
.has_live_endpoint_session_role(sid, role)
{
if entry.id != rv {
return Err(ClusterError::RendezvousMismatch {
expected: entry.id.raw(),
actual: rv.raw(),
});
}
return Err(ClusterError::RendezvousBusy { id: rv.raw() });
}
current = entry.next;
}
let Some(mut target) = target else {
return Err(ClusterError::RendezvousUnregistered { id: rv.raw() });
};
let entry = unsafe {
target.as_mut()
};
let Some(rendezvous) = entry.rendezvous_mut() else {
return Err(ClusterError::RendezvousBusy { id: rv.raw() });
};
unsafe { rendezvous.allocate_endpoint_lease(sid, role, bytes, align, resident_budget) }
.map_err(ClusterError::resource_exhausted)
}
pub(crate) fn ensure_dynamic_resolver_capacity(
&mut self,
rv_id: RendezvousId,
additional_entries: usize,
) -> Result<(), crate::session::cluster::error::ClusterError> {
if additional_entries == 0 {
return Ok(());
}
let entry = self.entry_mut(&rv_id).ok_or(
crate::session::cluster::error::ClusterError::RendezvousMismatch {
expected: rv_id.raw(),
actual: 0,
},
)?;
let rv = entry.rendezvous_mut().ok_or(
crate::session::cluster::error::ClusterError::RendezvousMismatch {
expected: rv_id.raw(),
actual: 0,
},
)?;
let rv_ptr = core::ptr::from_mut(rv);
entry.resolver_bucket.ensure_capacity(
additional_entries,
|bytes, align| {
unsafe { (&mut *rv_ptr).allocate_external_persistent_sidecar_bytes(bytes, align) }
},
|sidecar| {
unsafe { (&mut *rv_ptr).release_external_persistent_sidecar(sidecar) }
},
)
}
pub(crate) fn insert_dynamic_resolver(
&mut self,
key: crate::session::cluster::core::DynamicResolverKey,
entry: crate::session::cluster::core::DynamicResolverEntry<'cfg>,
) -> Result<(), crate::session::cluster::error::ClusterError> {
self.entry_mut(&key.rv)
.ok_or(
crate::session::cluster::error::ClusterError::RendezvousMismatch {
expected: key.rv.raw(),
actual: 0,
},
)?
.resolver_bucket
.insert(key.scope, entry)
}
pub(crate) fn dynamic_resolver(
&self,
key: crate::session::cluster::core::DynamicResolverKey,
) -> Option<&crate::session::cluster::core::DynamicResolverEntry<'cfg>> {
self.entry_ref(&key.rv)?.resolver_bucket.get(key.scope)
}
}
impl<'cfg, T> Drop for RendezvousTable<'cfg, T>
where
T: Transport,
{
fn drop(&mut self) {
let mut current = self.head;
while let Some(entry_ptr) = current {
let entry = entry_ptr.as_ptr();
current = unsafe {
(*entry).next
};
unsafe {
ptr::drop_in_place(entry);
}
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub(crate) enum RegisterRendezvousError {
CapacityExceeded,
StorageExhausted,
}