use super::{
AttachError, ClusterError, CompiledProgramRef, EndpointLeaseId, Lane, LaneLease, LeaseError,
PublicEndpointStorageLayout, PublicEndpointStorageRequest, RegisterRendezvousError, Rendezvous,
RendezvousError, RendezvousId, ResourceScope, RoleImageSlice, SessionId, SessionStorage,
};
pub(crate) struct SessionCluster<'cfg, T>
where
T: crate::transport::Transport + 'cfg,
{
pub(crate) storage: core::cell::UnsafeCell<SessionStorage<'cfg, T>>,
_local_only: crate::local::LocalOnly,
}
impl<'cfg, T> SessionCluster<'cfg, T>
where
T: crate::transport::Transport + 'cfg,
{
pub(crate) unsafe fn init_empty(dst: *mut Self) {
unsafe {
SessionStorage::<T>::init_empty(core::ptr::addr_of_mut!((*dst).storage).cast());
core::ptr::addr_of_mut!((*dst)._local_only).write(crate::local::LocalOnly::new());
}
}
#[inline]
fn storage_ptr(&self) -> *mut SessionStorage<'cfg, T> {
self.storage.get()
}
#[inline]
pub(crate) fn storage_ref_ptr(&self) -> *const SessionStorage<'cfg, T> {
self.storage.get() as *const SessionStorage<'cfg, T>
}
pub(crate) fn with_storage_mut<F, R>(&self, f: F) -> R
where
F: FnOnce(&mut SessionStorage<'cfg, T>) -> R,
{
unsafe { f(&mut *self.storage_ptr()) }
}
#[inline]
pub(crate) fn map_rendezvous_access_error(error: LeaseError) -> ClusterError {
match error {
LeaseError::RendezvousUnregistered(id) => {
ClusterError::RendezvousUnregistered { id: id.raw() }
}
LeaseError::AlreadyLeased(id) => ClusterError::RendezvousBusy { id: id.raw() },
}
}
#[inline]
pub(in crate::session::cluster::core) fn public_endpoint_storage_requirement<const ROLE: u8>(
role_image: RoleImageSlice<ROLE>,
) -> PublicEndpointStorageLayout {
let arena_layout = role_image.endpoint_arena_layout();
let storage_layout = crate::endpoint::kernel::cursor_endpoint_storage_layout::<0, T>(
&arena_layout,
role_image.endpoint_lane_slot_count(),
);
PublicEndpointStorageLayout {
total_bytes: storage_layout.total_bytes(),
total_align: storage_layout.total_align(),
arena_offset: storage_layout.arena_offset(),
}
}
#[inline]
pub(crate) fn public_endpoint_resident_budget<const ROLE: u8>(
compiled_role: RoleImageSlice<ROLE>,
) -> crate::rendezvous::core::EndpointResidentBudget {
crate::rendezvous::core::EndpointResidentBudget::with_route_storage(
compiled_role.route_table_frame_slots(),
compiled_role.route_table_lane_slots(),
crate::rendezvous::core::Rendezvous::<T>::frontier_workspace_guard_bytes(
compiled_role.descriptor().frontier_scratch_layout(),
),
)
}
#[inline(never)]
pub(in crate::session::cluster::core) fn allocate_public_endpoint_storage_for_rv<
'r,
const ROLE: u8,
>(
&self,
request: PublicEndpointStorageRequest,
) -> Result<
(
EndpointLeaseId,
u32,
*mut crate::endpoint::kernel::CursorEndpoint<'r, ROLE, T>,
),
ClusterError,
>
where
'cfg: 'r,
{
let PublicEndpointStorageRequest {
rv_id,
sid,
required_bytes,
required_align,
logical_lane_count,
required_assoc_slots,
resident_budget,
} = request;
let core = unsafe { &mut *self.storage_ptr() };
let (slot, generation, offset, _len) =
core.locals.allocate_endpoint_lease_for_session_role(
rv_id,
sid,
ROLE,
required_bytes,
required_align,
resident_budget,
)?;
let rv = core
.locals
.get_mut_checked(&rv_id)
.map_err(Self::map_rendezvous_access_error)?;
let (slab_ptr, slab_len) = rv.slab_ptr_and_len();
let Some(storage_end) = offset.checked_add(required_bytes) else {
rv.release_endpoint_lease(slot, generation)
.map_err(ClusterError::resource_exhausted)?;
return Err(ClusterError::resource_exhausted(
ResourceScope::EndpointBounds,
));
};
if storage_end > slab_len {
rv.release_endpoint_lease(slot, generation)
.map_err(ClusterError::resource_exhausted)?;
return Err(ClusterError::resource_exhausted(
ResourceScope::EndpointBounds,
));
}
if let Err(resource) = rv.ensure_endpoint_lease_live(slot, generation) {
rv.release_endpoint_lease(slot, generation)
.map_err(ClusterError::resource_exhausted)?;
return Err(ClusterError::resource_exhausted(resource));
}
if let Err(resource) = rv.ensure_endpoint_resident_budget(resident_budget) {
rv.release_endpoint_lease(slot, generation)
.map_err(ClusterError::resource_exhausted)?;
return Err(ClusterError::resource_exhausted(resource));
}
if let Err(resource) =
rv.ensure_core_lane_storage_for_assoc_entries(logical_lane_count, required_assoc_slots)
{
rv.release_endpoint_lease(slot, generation)
.map_err(ClusterError::resource_exhausted)?;
return Err(ClusterError::resource_exhausted(resource));
}
Ok((
slot,
generation,
unsafe { slab_ptr.add(offset) }
.cast::<crate::endpoint::kernel::CursorEndpoint<'r, ROLE, T>>(),
))
}
fn public_endpoint_storage_raw_ptr(
&self,
rv_id: RendezvousId,
slot: EndpointLeaseId,
generation: u32,
) -> Option<*mut ()> {
let rv = self.get_local(&rv_id)?;
let (slab_ptr, slab_len) = rv.slab_ptr_and_len();
let (offset, len) = rv.endpoint_lease_storage(slot, generation)?;
let storage_end = offset.checked_add(len)?;
if len == 0 || storage_end > slab_len {
return None;
}
Some(
unsafe { slab_ptr.add(offset).cast() },
)
}
pub(crate) fn public_endpoint_header_ptr(
&self,
rv_id: RendezvousId,
slot: EndpointLeaseId,
generation: u32,
) -> Option<core::ptr::NonNull<crate::endpoint::carrier::KernelEndpointHeader<'cfg>>> {
core::ptr::NonNull::new(
self.public_endpoint_storage_raw_ptr(rv_id, slot, generation)?
.cast::<crate::endpoint::carrier::KernelEndpointHeader<'cfg>>(),
)
}
#[inline]
pub(crate) fn session_fault(
&self,
rv_id: RendezvousId,
sid: SessionId,
) -> Option<crate::rendezvous::SessionFaultKind> {
self.with_storage_mut(|core| {
core.locals
.get_mut(&rv_id)
.and_then(|rv| rv.session_fault(sid))
})
}
#[inline]
pub(crate) fn poison_session(
&self,
rv_id: RendezvousId,
sid: SessionId,
cause: crate::rendezvous::SessionFaultKind,
) -> crate::rendezvous::SessionFaultKind {
self.with_storage_mut(|core| {
let Some(rv) = core.locals.get_mut(&rv_id) else {
crate::invariant();
};
rv.poison_session(sid, cause)
})
}
#[inline]
pub(crate) fn register_session_waiter(
&self,
rv_id: RendezvousId,
sid: SessionId,
lane: Lane,
waker: &core::task::Waker,
) {
self.with_storage_mut(|core| {
if let Some(rv) = core.locals.get_mut(&rv_id) {
rv.register_session_waiter(sid, lane, waker);
}
});
}
#[inline]
pub(crate) fn clear_session_waiter(&self, rv_id: RendezvousId, sid: SessionId, lane: Lane) {
self.with_storage_mut(|core| {
if let Some(rv) = core.locals.get_mut(&rv_id) {
rv.clear_session_waiter(sid, lane);
}
});
}
fn ensure_compiled_program_ref<const ROLE: u8, P>(
&self,
rv_id: RendezvousId,
program: &P,
) -> Result<&'static CompiledProgramRef, AttachError>
where
P: crate::global::RoleProgramView<ROLE>,
{
let compiled = program.role_image_ref().program;
let core = unsafe { &mut *self.storage_ptr() };
core.locals
.get_mut_checked(&rv_id)
.map_err(Self::map_rendezvous_access_error)
.map_err(AttachError::cluster)?;
Ok(compiled)
}
#[inline(never)]
pub(crate) fn ensure_role_image_slice<const ROLE: u8, P>(
&self,
rv_id: RendezvousId,
program: &P,
) -> Result<RoleImageSlice<ROLE>, AttachError>
where
P: crate::global::RoleProgramView<ROLE>,
{
let compiled = program.role_image_ref();
let core = unsafe { &mut *self.storage_ptr() };
core.locals
.get_mut_checked(&rv_id)
.map_err(Self::map_rendezvous_access_error)
.map_err(AttachError::cluster)?;
Ok(RoleImageSlice::from_resident(compiled))
}
pub(crate) unsafe fn release_public_endpoint_slot(
&self,
rv_id: RendezvousId,
slot: EndpointLeaseId,
generation: u32,
) -> Result<(), ResourceScope> {
self.with_storage_mut(|core| {
if let Some(rv) = core.locals.get_mut(&rv_id) {
rv.release_endpoint_lease(slot, generation)?;
}
Ok(())
})
}
#[inline]
pub(crate) fn release_public_endpoint_slot_owned(
&self,
rv_id: RendezvousId,
slot: EndpointLeaseId,
generation: u32,
) {
unsafe {
if self
.release_public_endpoint_slot(rv_id, slot, generation)
.is_err()
{
crate::invariant();
}
}
}
pub(crate) fn with_resident_program_ref<const ROLE: u8, P, F, R, E>(
&self,
rv_id: RendezvousId,
program: &P,
f: F,
) -> Result<R, E>
where
E: From<ClusterError>,
P: crate::global::RoleProgramView<ROLE>,
F: FnOnce(&CompiledProgramRef) -> Result<R, E>,
{
let compiled = self
.ensure_compiled_program_ref(rv_id, program)
.map_err(|err| E::from(crate::invariant_some(err.cluster_cause())))?;
f(compiled)
}
pub(crate) fn register_rendezvous(
&self,
resources: crate::runtime_core::resources::RuntimeResources<'cfg>,
transport: T,
) -> Result<RendezvousId, ClusterError> {
self.with_storage_mut(|core| {
match core
.locals
.register_local_from_resources_auto(resources, transport)
{
Ok(id) => Ok(id),
Err(
RegisterRendezvousError::CapacityExceeded
| RegisterRendezvousError::StorageExhausted,
) => Err(ClusterError::resource_exhausted(
ResourceScope::RendezvousTable,
)),
}
})
}
pub(crate) fn get_local(&self, id: &RendezvousId) -> Option<&Rendezvous<'cfg, 'cfg, T>> {
unsafe { (*self.storage_ref_ptr()).locals.get(id) }
}
pub(crate) fn lease_port<'lease>(
&'lease self,
rv_id: RendezvousId,
sid: SessionId,
lane: Lane,
role: u8,
role_count: u8,
) -> Result<LaneLease<'lease, 'cfg, T>, RendezvousError>
where
'cfg: 'lease,
{
let core = unsafe { &mut *self.storage_ptr() };
let mut lease = match core.locals.lease(rv_id) {
Ok(lease) => lease,
Err(LeaseError::RendezvousUnregistered(_)) => {
return Err(RendezvousError::LaneOutOfRange { lane });
}
Err(LeaseError::AlreadyLeased(_)) => {
return Err(RendezvousError::LaneBusy { lane });
}
};
let active = &core.active_leases;
let current = active.get();
active.set(current + 1);
let brand = lease.brand();
Ok(LaneLease::new(
lease, sid, lane, role, role_count, active, brand,
))
}
}