use core::{marker::PhantomData, ptr, ptr::NonNull};
use crate::rendezvous::core::{EndpointLeaseId, EndpointResidentBudget, LaneRelease, Rendezvous};
use crate::session::types::{Lane, RendezvousId, SessionId};
use crate::transport::Transport;
mod registry_ops;
pub(crate) use registry_ops::RegisterRendezvousError;
pub(crate) struct RendezvousTable<'cfg, T: Transport> {
head: Option<NonNull<RendezvousEntry<'cfg, T>>>,
len: u16,
}
impl<'cfg, T> RendezvousTable<'cfg, T>
where
T: Transport,
{
pub(crate) unsafe fn init_empty(dst: *mut Self) {
unsafe {
core::ptr::addr_of_mut!((*dst).head).write(None);
core::ptr::addr_of_mut!((*dst).len).write(0);
}
}
fn entry_ref(&self, id: &RendezvousId) -> Option<&RendezvousEntry<'cfg, T>> {
let mut current = self.head;
while let Some(entry_ptr) = current {
let entry = unsafe {
entry_ptr.as_ref()
};
if entry.id == *id {
return Some(entry);
}
current = entry.next;
}
None
}
fn entry_mut(&mut self, id: &RendezvousId) -> Option<&mut RendezvousEntry<'cfg, T>> {
let mut current = self.head;
while let Some(mut entry_ptr) = current {
let entry = unsafe {
entry_ptr.as_mut()
};
if entry.id == *id {
return Some(entry);
}
current = entry.next;
}
None
}
pub(crate) fn get(&self, id: &RendezvousId) -> Option<&Rendezvous<'cfg, 'cfg, T>> {
self.entry_ref(id).and_then(|entry| entry.rendezvous_ref())
}
pub(crate) fn get_mut(&mut self, id: &RendezvousId) -> Option<&mut Rendezvous<'cfg, 'cfg, T>> {
self.entry_mut(id).and_then(|entry| entry.rendezvous_mut())
}
pub(crate) fn get_mut_checked(
&mut self,
id: &RendezvousId,
) -> Result<&mut Rendezvous<'cfg, 'cfg, T>, LeaseError> {
let slot = self
.entry_mut(id)
.ok_or(LeaseError::RendezvousUnregistered(*id))?;
if slot.is_active() {
return Err(LeaseError::AlreadyLeased(*id));
}
Ok(slot.rendezvous())
}
pub(crate) fn lease<'lease>(
&'lease mut self,
rv_id: RendezvousId,
) -> Result<RendezvousLease<'lease, 'cfg, T>, LeaseError>
where
'cfg: 'lease,
{
let slot = self
.entry_mut(&rv_id)
.ok_or(LeaseError::RendezvousUnregistered(rv_id))?;
if slot.is_active() {
return Err(LeaseError::AlreadyLeased(rv_id));
}
slot.mark_active();
Ok(RendezvousLease::new(slot))
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub(crate) enum LeaseError {
RendezvousUnregistered(RendezvousId),
AlreadyLeased(RendezvousId),
}
struct RendezvousEntry<'cfg, T>
where
T: Transport,
{
id: RendezvousId,
rendezvous: NonNull<Rendezvous<'cfg, 'cfg, T>>,
state: RendezvousEntryState,
resolver_bucket: crate::session::cluster::core::ResolverBucket<'cfg>,
next: Option<NonNull<RendezvousEntry<'cfg, T>>>,
_marker: PhantomData<&'cfg mut Rendezvous<'cfg, 'cfg, T>>,
}
#[derive(Clone, Copy, Eq, PartialEq)]
enum RendezvousEntryState {
Available,
Leased,
}
impl<'cfg, T> RendezvousEntry<'cfg, T>
where
T: Transport,
{
fn is_active(&self) -> bool {
self.state == RendezvousEntryState::Leased
}
fn mark_active(&mut self) {
self.state = RendezvousEntryState::Leased;
}
fn clear_active(&mut self) {
self.state = RendezvousEntryState::Available;
}
fn rendezvous_ref(&self) -> Option<&Rendezvous<'cfg, 'cfg, T>> {
match self.state {
RendezvousEntryState::Available => Some(
unsafe { self.rendezvous.as_ref() },
),
RendezvousEntryState::Leased => None,
}
}
fn rendezvous_mut(&mut self) -> Option<&mut Rendezvous<'cfg, 'cfg, T>> {
match self.state {
RendezvousEntryState::Available => Some(
unsafe { self.rendezvous.as_mut() },
),
RendezvousEntryState::Leased => None,
}
}
fn rendezvous(&mut self) -> &mut Rendezvous<'cfg, 'cfg, T> {
unsafe { self.rendezvous.as_mut() }
}
}
impl<'cfg, T> RendezvousEntry<'cfg, T>
where
T: Transport,
{
unsafe fn init_from_parts(
dst: *mut Self,
rv_id: RendezvousId,
rendezvous: NonNull<Rendezvous<'cfg, 'cfg, T>>,
next: Option<NonNull<RendezvousEntry<'cfg, T>>>,
) {
unsafe {
core::ptr::addr_of_mut!((*dst).id).write(rv_id);
core::ptr::addr_of_mut!((*dst).rendezvous).write(rendezvous);
core::ptr::addr_of_mut!((*dst).state).write(RendezvousEntryState::Available);
crate::session::cluster::core::ResolverBucket::init_empty(core::ptr::addr_of_mut!(
(*dst).resolver_bucket
));
core::ptr::addr_of_mut!((*dst).next).write(next);
core::ptr::addr_of_mut!((*dst)._marker).write(PhantomData);
}
}
}
impl<'cfg, T> Drop for RendezvousEntry<'cfg, T>
where
T: Transport,
{
fn drop(&mut self) {
unsafe {
ptr::drop_in_place(self.rendezvous.as_ptr());
}
}
}
pub(crate) struct RendezvousLease<'lease, 'cfg, T: Transport>
where
'cfg: 'lease,
{
slot: Option<&'lease mut RendezvousEntry<'cfg, T>>,
}
impl<'lease, 'cfg, T> RendezvousLease<'lease, 'cfg, T>
where
T: Transport,
'cfg: 'lease,
{
fn new(slot: &'lease mut RendezvousEntry<'cfg, T>) -> Self {
Self { slot: Some(slot) }
}
#[inline]
fn entry_mut(&mut self) -> &mut RendezvousEntry<'cfg, T> {
match self.slot.as_mut() {
Some(slot) => slot,
None => crate::invariant(),
}
}
#[inline]
pub(crate) fn with_rendezvous<R>(
&mut self,
f: impl FnOnce(&mut Rendezvous<'cfg, 'cfg, T>) -> R,
) -> R {
let entry = self.entry_mut();
f(entry.rendezvous())
}
#[inline]
pub(crate) fn brand(&mut self) -> crate::session::brand::Guard<'cfg> {
self.with_rendezvous(|rv| rv.brand())
}
#[inline]
pub(crate) fn release_lane_with_tap(&mut self, sid: SessionId, lane: Lane) {
self.with_rendezvous(|rv| match rv.release_lane(sid, lane) {
LaneRelease::Released => {
rv.emit_lane_release(sid, lane);
}
LaneRelease::StillHeld => {}
});
}
}
impl<'lease, 'cfg, T> Drop for RendezvousLease<'lease, 'cfg, T>
where
T: Transport,
'cfg: 'lease,
{
fn drop(&mut self) {
if let Some(slot) = self.slot.take() {
slot.clear_active();
}
}
}