use super::{
AttachError, ClusterError, CompiledProgramRef, EndpointLeaseId, EndpointLeaseRequest, 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]
pub(crate) fn storage_ref_ptr(&self) -> *const SessionStorage<'cfg, T> {
self.storage.get() as *const SessionStorage<'cfg, T>
}
#[inline]
fn storage_ref(&self) -> &SessionStorage<'cfg, T> {
unsafe { &*self.storage_ref_ptr() }
}
#[inline]
pub(in crate::session::cluster::core) fn locals(
&self,
) -> &crate::session::lease::core::RendezvousTable<'cfg, T> {
unsafe { &*self.storage_ref().locals.get() }
}
#[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(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,
role_descriptor,
} = request;
let resident_budget =
crate::rendezvous::core::EndpointResidentBudget::with_frontier_workspace(
crate::rendezvous::core::Rendezvous::<T>::frontier_workspace_guard_bytes(
role_descriptor.frontier_scratch_layout(),
),
);
let (slot, generation, offset, _len) = self
.locals()
.allocate_endpoint_lease_for_session_role(EndpointLeaseRequest {
rendezvous: rv_id,
session: sid,
role: ROLE,
program: role_descriptor.program(),
bytes: required_bytes,
align: required_align,
resident_budget,
})?;
let rv = self
.locals()
.get_checked(&rv_id)
.map_err(Self::map_rendezvous_access_error)?;
let (slab_ptr, slab_len) = rv.slab_ptr_and_len();
let storage_end = crate::invariant_some(offset.checked_add(required_bytes));
if storage_end > slab_len {
crate::invariant();
}
if let Err(resource) = rv.ensure_endpoint_resident_capacity() {
rv.abort_endpoint_lease_reservation(slot, generation);
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.abort_endpoint_lease_reservation(slot, generation);
return Err(ClusterError::resource_exhausted(resource));
}
Ok((
slot,
generation,
unsafe { slab_ptr.add(offset) }
.cast::<crate::endpoint::kernel::CursorEndpoint<'r, ROLE, T>>(),
))
}
pub(crate) fn public_endpoint_header_ptr(
&self,
rv_id: RendezvousId,
slot: EndpointLeaseId,
generation: u32,
) -> core::ptr::NonNull<crate::endpoint::carrier::KernelEndpointHeader<'cfg>> {
let rv = crate::invariant_some(
self.locals()
.published_endpoint_owner(rv_id, slot, generation),
);
let (slab_ptr, slab_len) = rv.slab_ptr_and_len();
let (offset, len) = crate::invariant_some(rv.endpoint_lease_storage(slot, generation));
let storage_end = crate::invariant_some(offset.checked_add(len));
if len == 0 || storage_end > slab_len {
crate::invariant();
}
let ptr = unsafe {
slab_ptr
.add(offset)
.cast::<crate::endpoint::carrier::KernelEndpointHeader<'cfg>>()
};
crate::invariant_some(core::ptr::NonNull::new(ptr))
}
pub(crate) fn try_public_endpoint_operation_lease(
&self,
rv_id: RendezvousId,
slot: EndpointLeaseId,
generation: u32,
) -> Option<crate::rendezvous::core::EndpointOperationLease<'_>> {
let rendezvous = self
.locals()
.published_endpoint_owner(rv_id, slot, generation)?;
rendezvous.try_endpoint_operation_lease()
}
#[inline]
pub(crate) fn session_fault(
&self,
rv_id: RendezvousId,
slot: EndpointLeaseId,
generation: u32,
sid: SessionId,
) -> Option<crate::rendezvous::SessionFaultKind> {
crate::invariant_some(
self.locals()
.published_endpoint_owner(rv_id, slot, generation),
)
.session_fault(sid)
}
#[inline]
pub(crate) fn poison_session<const ROLE: u8>(
&self,
rv_id: RendezvousId,
slot: EndpointLeaseId,
generation: u32,
sid: SessionId,
cause: crate::rendezvous::SessionFaultKind,
) -> crate::rendezvous::SessionFaultKind {
let rv = crate::invariant_some(
self.locals()
.published_endpoint_owner(rv_id, slot, generation),
);
rv.poison_session(sid, ROLE, cause)
}
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;
self.locals()
.get_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();
self.locals()
.get_checked(&rv_id)
.map_err(Self::map_rendezvous_access_error)
.map_err(AttachError::cluster)?;
Ok(RoleImageSlice::from_resident(compiled))
}
pub(crate) fn release_public_endpoint_slot_owned(
&self,
rv_id: RendezvousId,
slot: EndpointLeaseId,
generation: u32,
) {
crate::invariant_some(
self.locals()
.published_endpoint_owner(rv_id, slot, generation),
)
.release_endpoint_lease(slot, generation);
}
pub(crate) fn replace_public_endpoint_waiter(
&self,
rv_id: RendezvousId,
slot: EndpointLeaseId,
generation: u32,
replacement: core::task::Waker,
) -> Option<core::task::Waker> {
crate::invariant_some(
self.locals()
.published_endpoint_owner(rv_id, slot, generation),
)
.replace_endpoint_waiter(slot, generation, replacement)
}
pub(crate) fn take_public_endpoint_waiter(
&self,
rv_id: RendezvousId,
slot: EndpointLeaseId,
generation: u32,
) -> Option<core::task::Waker> {
crate::invariant_some(
self.locals()
.published_endpoint_owner(rv_id, slot, generation),
)
.take_endpoint_waiter(slot, generation)
}
pub(crate) fn publish_public_endpoint_slot(
&self,
rv_id: RendezvousId,
slot: EndpointLeaseId,
generation: u32,
) {
crate::invariant_some(self.get_local(&rv_id)).publish_endpoint_lease(slot, generation);
}
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(&'static 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> {
match self
.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>> {
self.locals().get(id)
}
pub(crate) fn lease_port<'lease>(
&'lease self,
rv_id: RendezvousId,
sid: SessionId,
lane: Lane,
role: u8,
) -> Result<LaneLease<'lease, 'cfg, T>, RendezvousError>
where
'cfg: 'lease,
{
let lease = match self.locals().lease(rv_id) {
Ok(lease) => lease,
Err(LeaseError::RendezvousUnregistered(_)) => {
return Err(RendezvousError::LaneOutOfRange { lane });
}
Err(LeaseError::AlreadyLeased(_)) => {
return Err(RendezvousError::LaneBusy { lane });
}
};
lease.with_rendezvous(|rv| rv.activate_lane_attachment(sid, lane))?;
let active = &self.storage_ref().active_leases;
let next = crate::invariant_some(active.get().checked_add(1));
active.set(next);
let brand = lease.brand();
Ok(LaneLease::new(lease, sid, lane, role, active, brand))
}
}