use core::{convert::TryFrom, ops::ControlFlow, task::Poll};
use super::authority::{
Arm, DeferReason, DeferSource, LoopDecision, RouteDecisionSource, RouteDecisionToken,
RouteResolveStep, decision_policy_input_arg0, validate_route_decision_scope,
};
use super::evidence::{ScopeEvidence, ScopeFrameLabelMeta, ScopeLoopMeta};
use super::frontier::*;
use super::frontier_state::FrontierState;
use super::inbox::{BindingInbox, PackedIngressEvidence};
use super::lane_port;
use super::lane_slots::LaneSlotArray;
use super::layout::{EndpointArenaLayout, LeasedState};
use super::offer::*;
mod route_commit_helpers;
use super::decision_state::{RouteArmCommitProof, RouteCommitProofWorkspace, RouteState};
use crate::binding::{EndpointSlot, IngressEvidence, NoBinding};
use crate::eff::EffIndex;
use crate::global::ControlDesc;
#[cfg(test)]
use crate::global::LoopControlMeaning;
#[cfg(all(test, hibana_repo_tests))]
use crate::global::Message;
use crate::global::compiled::images::{ControlSemanticKind, ControlSemanticsTable};
use crate::global::const_dsl::{PolicyMode, ScopeId, ScopeKind};
use crate::global::role_program::LaneSetView;
use crate::global::typestate::{
ARM_SHARED, JumpReason, LocalAction, LoopRole, PassiveArmNavigation, PhaseCursor, RecvMeta,
SendMeta, StateIndex, state_index_to_usize,
};
use crate::{
control::types::{Lane, RendezvousId, SessionId},
control::{
cap::mint::{
CAP_HANDLE_LEN, CAP_TOKEN_LEN, CapHeader, CapShot, ControlOp, E0, EndpointEpoch,
EpochTable, EpochTbl, MintConfigMarker, Owner,
},
cap::resource_kinds::{LoopDecisionHandle, RouteArmHandle},
cluster::{
core::{DescriptorPublicationAuthority, DescriptorTerminal, DynamicPolicyResolution},
error::CpError,
},
},
endpoint::{
RecvError, RecvResult, SendError, SendResult, affine::LaneGuard, control::SessionControlCtx,
},
observe::core::{TapEvent, emit},
observe::scope::ScopeTrace,
observe::{events, ids},
policy_runtime::{self, PolicySlot},
rendezvous::SessionFaultKind,
rendezvous::{
capability::{CapEntry, CapReleaseCtx},
core::EndpointLeaseId,
port::Port,
},
runtime::consts::LabelUniverse,
transport::{
FrameLabelMask, Transport,
trace::TapFrameMeta,
wire::{CodecError, FrameFlags, Payload},
},
};
pub(in crate::endpoint::kernel) use route_commit_helpers::{
is_linger_route_from_cursor, preflight_route_arm_commit_after_clearing_other_lanes_from_parts,
preflight_route_arm_commit_from_parts, require_route_arm_commit_proof_from_parts,
scope_slot_for_route_from_cursor,
};
pub(in crate::endpoint::kernel::core) use route_commit_helpers::{
preview_selected_arm_for_scope_from_parts, route_scope_materialization_index_from_cursor,
};
#[derive(Clone, Copy)]
enum BindingLanePreference {
Any,
Arm(u8),
LabelMask(FrameLabelMask),
}
#[cfg(all(test, hibana_repo_tests))]
#[path = "test_support/core_offer_tests.rs"]
mod offer_regression_tests;
#[inline]
fn checked_state_index(idx: usize) -> Option<StateIndex> {
u16::try_from(idx).ok().map(StateIndex::new)
}
pub(crate) trait RecvKernelEndpoint<'r> {
fn prepare_recv_kernel_descriptor(
&mut self,
label: u8,
expects_control: bool,
accepts_empty_payload: bool,
) -> RecvResult<super::recv::PreparedRecv>;
fn poll_recv_kernel_payload_source(
&mut self,
desc: super::recv::RecvDescriptor,
accepts_empty_payload: bool,
state: &mut super::recv::RecvState,
cx: &mut core::task::Context<'_>,
) -> Poll<RecvResult<super::recv::RecvPayloadSource<'r>>>;
fn finish_recv_kernel_payload(
&mut self,
desc: super::recv::RecvDescriptor,
payload_source: super::recv::RecvPayloadSource<'r>,
erased: RecvRuntimeDesc,
control: Option<ControlDesc>,
validate: for<'a> fn(Payload<'a>) -> Result<(), CodecError>,
) -> RecvResult<Payload<'r>>;
}
pub(crate) trait DecodeKernelEndpoint<'r> {
fn prepare_decode_kernel_transport_wait(
&mut self,
desc: DecodeRuntimeDesc,
branch: &MaterializedRouteBranch<'r>,
) -> RecvResult<Option<RecvMeta>>;
fn poll_decode_kernel_transport_payload(
&mut self,
meta: RecvMeta,
pending_recv: &mut lane_port::PendingRecv,
cx: &mut core::task::Context<'_>,
) -> Poll<RecvResult<lane_port::ReceivedFrame<'r>>>;
fn finish_decode_kernel(
&mut self,
desc: DecodeRuntimeDesc,
control: Option<ControlDesc>,
prepared_meta: Option<RecvMeta>,
branch: &mut MaterializedRouteBranch<'r>,
) -> RecvResult<Payload<'r>>;
}
pub(crate) trait SendKernelEndpoint<'r> {
fn poll_send_init_kernel(
&mut self,
descriptor: SendRuntimeDesc,
meta: SendMeta,
preview_cursor_index: Option<StateIndex>,
payload: Option<lane_port::RawSendPayload>,
) -> SendInitOutcome<'r>;
fn poll_send_pending_kernel(
&mut self,
pending: &mut PendingSendIo<'r>,
cx: &mut core::task::Context<'_>,
) -> Poll<SendResult<SendCommitPlan<'r>>>;
fn finish_send_after_transport_kernel(
&mut self,
commit_plan: SendCommitPlan<'r>,
) -> SendCommitOutcome<'r>;
}
#[inline(never)]
pub(crate) fn kernel_recv<'r>(
endpoint: &mut dyn RecvKernelEndpoint<'r>,
logical_label: u8,
expects_control: bool,
control: Option<ControlDesc>,
accepts_empty_payload: bool,
validate: for<'a> fn(Payload<'a>) -> Result<(), CodecError>,
state: &mut super::recv::RecvState,
cx: &mut core::task::Context<'_>,
) -> Poll<RecvResult<Payload<'r>>> {
let prepared = match state.prepared() {
Some(prepared) => prepared,
None => {
let prepared = match endpoint.prepare_recv_kernel_descriptor(
logical_label,
expects_control,
accepts_empty_payload,
) {
Ok(prepared) => prepared,
Err(err) => return Poll::Ready(Err(err)),
};
state.set_prepared(prepared);
prepared
}
};
match endpoint.poll_recv_kernel_payload_source(
prepared.descriptor,
prepared.runtime.accepts_empty_payload(),
state,
cx,
) {
Poll::Pending => Poll::Pending,
Poll::Ready(Ok(payload_source)) => {
state.clear_prepared();
Poll::Ready(
endpoint
.finish_recv_kernel_payload(
prepared.descriptor,
payload_source,
prepared.runtime,
control,
validate,
)
.map(|payload| unsafe {
lane_port::endpoint_resident_payload(payload)
}),
)
}
Poll::Ready(Err(err)) => {
state.clear_prepared();
Poll::Ready(Err(err))
}
}
}
#[inline(never)]
pub(crate) fn kernel_decode<'r>(
endpoint: &mut dyn DecodeKernelEndpoint<'r>,
desc: DecodeRuntimeDesc,
control: Option<ControlDesc>,
state: &mut super::decode::DecodeState<'r>,
cx: &mut core::task::Context<'_>,
) -> Poll<RecvResult<Payload<'r>>> {
if state.branch().is_none() {
return Poll::Ready(Err(RecvError::PhaseInvariant));
}
if state.prepared_meta().is_none() {
let prepared = {
let branch = state.branch().expect("decode branch checked above");
match endpoint.prepare_decode_kernel_transport_wait(desc, branch) {
Ok(meta) => meta,
Err(err) => return Poll::Ready(Err(err)),
}
};
state.set_prepared_meta(prepared);
}
if let Some(meta) = state.prepared_meta() {
let needs_transport = {
let branch = state.branch().expect("decode branch checked above");
branch.staged_payload.is_none() && !branch.binding_evidence.is_present()
};
if needs_transport {
let frame = match endpoint.poll_decode_kernel_transport_payload(
meta,
state.pending_recv_mut(),
cx,
) {
Poll::Pending => return Poll::Pending,
Poll::Ready(Ok(payload)) => payload,
Poll::Ready(Err(err)) => {
state.set_prepared_meta(None);
return Poll::Ready(Err(err));
}
};
let branch = state.branch_mut().expect("decode branch checked above");
branch.staged_payload = Some(StagedPayload::Transport { frame });
}
}
let prepared_meta = state.prepared_meta();
let result = {
let branch = state.branch_mut().expect("decode branch checked above");
endpoint.finish_decode_kernel(desc, control, prepared_meta, branch)
};
match result {
Ok(payload) => {
let _ = state.take_branch();
state.restore_on_drop = false;
Poll::Ready(Ok(unsafe {
lane_port::endpoint_resident_payload(payload)
}))
}
Err(err) => Poll::Ready(Err(err)),
}
}
#[inline(never)]
pub(crate) fn kernel_send<'r>(
endpoint: &mut dyn SendKernelEndpoint<'r>,
state: &mut SendState<'r>,
payload: &mut Option<lane_port::RawSendPayload>,
cx: &mut core::task::Context<'_>,
) -> Poll<SendResult<SendCommitOutcome<'r>>> {
loop {
match state {
SendState::Init {
descriptor,
meta,
preview_cursor_index,
} => match endpoint.poll_send_init_kernel(
*descriptor,
*meta,
*preview_cursor_index,
payload.take(),
) {
SendInitOutcome::Ready(result) => {
*state = SendState::Done;
return Poll::Ready(result);
}
SendInitOutcome::Pending { pending } => {
*state = SendState::Sending { pending };
}
SendInitOutcome::Commit { commit_plan } => {
let result = endpoint.finish_send_after_transport_kernel(commit_plan);
*state = SendState::Done;
return Poll::Ready(Ok(result));
}
},
SendState::Sending { pending } => {
match endpoint.poll_send_pending_kernel(pending, cx) {
Poll::Pending => return Poll::Pending,
Poll::Ready(Ok(commit_plan)) => {
let result = endpoint.finish_send_after_transport_kernel(commit_plan);
*state = SendState::Done;
return Poll::Ready(Ok(result));
}
Poll::Ready(Err(err)) => {
*state = SendState::Done;
return Poll::Ready(Err(err));
}
}
}
SendState::Done => panic!("send future polled after completion"),
}
}
}
impl<'r, const ROLE: u8, T, U, C, E, const MAX_RV: usize, Mint, B> SendKernelEndpoint<'r>
for CursorEndpoint<'r, ROLE, T, U, C, E, MAX_RV, Mint, B>
where
T: Transport + 'r,
U: LabelUniverse,
C: crate::runtime::config::Clock,
E: EpochTable,
Mint: MintConfigMarker,
B: EndpointSlot + 'r,
<Mint as MintConfigMarker>::Policy: crate::control::cap::mint::AllowsEndpointMint,
{
#[inline]
fn poll_send_init_kernel(
&mut self,
descriptor: SendRuntimeDesc,
meta: SendMeta,
preview_cursor_index: Option<StateIndex>,
payload: Option<lane_port::RawSendPayload>,
) -> SendInitOutcome<'r> {
self.poll_send_init(descriptor, meta, preview_cursor_index, payload)
}
#[inline]
fn poll_send_pending_kernel(
&mut self,
pending: &mut PendingSendIo<'r>,
cx: &mut core::task::Context<'_>,
) -> Poll<SendResult<SendCommitPlan<'r>>> {
self.poll_send_pending(pending, cx)
}
#[inline]
fn finish_send_after_transport_kernel(
&mut self,
commit_plan: SendCommitPlan<'r>,
) -> SendCommitOutcome<'r> {
self.finish_send_after_transport_runtime(commit_plan)
}
}
#[inline]
fn controller_arm_label(cursor: &PhaseCursor, scope_id: ScopeId, arm: u8) -> Option<u8> {
cursor
.shared_controller_arm_entry_by_arm(scope_id, arm)
.map(|(_, label)| label)
}
#[inline]
fn controller_arm_semantic_kind(
cursor: &PhaseCursor,
_semantics: &ControlSemanticsTable,
scope_id: ScopeId,
arm: u8,
) -> Option<ControlSemanticKind> {
let (entry, _) = cursor.shared_controller_arm_entry_by_arm(scope_id, arm)?;
loop_control_semantic_kind(cursor.control_semantic_at(state_index_to_usize(entry)))
}
#[inline]
const fn loop_control_semantic_kind(kind: ControlSemanticKind) -> Option<ControlSemanticKind> {
if kind.is_loop() { Some(kind) } else { None }
}
#[inline]
const fn is_loop_control_semantic(kind: ControlSemanticKind) -> bool {
kind.is_loop()
}
#[inline]
const fn control_policy_is_validated_during_handle_preparation(op: ControlOp) -> bool {
matches!(op, ControlOp::TopologyBegin | ControlOp::TopologyAck)
}
#[inline]
fn next_preferred_lane_in_lane_set(
preferred_lane_idx: usize,
offer_lanes: LaneSetView,
lane_limit: usize,
scan_idx: &mut usize,
) -> Option<usize> {
if *scan_idx == 0 {
*scan_idx = 1;
if preferred_lane_idx < lane_limit && offer_lanes.contains(preferred_lane_idx) {
return Some(preferred_lane_idx);
}
}
let mut start = scan_idx.saturating_sub(1);
while let Some(lane_idx) = offer_lanes.next_set_from(start, lane_limit) {
*scan_idx = lane_idx.saturating_add(2);
start = lane_idx.saturating_add(1);
if lane_idx != preferred_lane_idx {
return Some(lane_idx);
}
}
None
}
#[inline]
#[cfg(test)]
const fn loop_control_meaning_from_semantic(
kind: ControlSemanticKind,
) -> Option<LoopControlMeaning> {
match kind {
ControlSemanticKind::LoopContinue => Some(LoopControlMeaning::Continue),
ControlSemanticKind::LoopBreak => Some(LoopControlMeaning::Break),
_ => None,
}
}
#[cfg(test)]
#[inline]
fn stage_transport_payload(scratch: &mut [u8], payload: &[u8]) -> RecvResult<usize> {
if payload.len() > scratch.len() {
return Err(RecvError::PhaseInvariant);
}
scratch[..payload.len()].copy_from_slice(payload);
Ok(payload.len())
}
#[cfg(test)]
fn endpoint_scope_frame_label_meta<'r, const ROLE: u8, T, U, C, E, const MAX_RV: usize, Mint, B>(
endpoint: &CursorEndpoint<'r, ROLE, T, U, C, E, MAX_RV, Mint, B>,
scope_id: ScopeId,
loop_meta: ScopeLoopMeta,
) -> ScopeFrameLabelMeta
where
T: Transport,
U: LabelUniverse,
C: crate::runtime::config::Clock,
E: EpochTable,
Mint: MintConfigMarker,
B: EndpointSlot + 'r,
{
CursorEndpoint::<ROLE, T, U, C, E, MAX_RV, Mint, B>::scope_frame_label_meta(
&endpoint.cursor,
&endpoint.control_semantics(),
scope_id,
loop_meta,
)
}
#[cfg(all(test, hibana_repo_tests))]
#[path = "core/decision_policy_tests.rs"]
mod decision_policy_tests;
#[cfg(all(test, hibana_repo_tests))]
#[path = "core/send_rollback_tests.rs"]
mod send_rollback_tests;
mod frontier_observation;
mod frontier_select;
mod offer_refresh;
mod scope_evidence_logic;
mod binding_ingress;
mod decision_policy;
mod frontier_helpers;
mod public_types;
mod route_preview;
mod route_preview_flow;
mod runtime_types;
mod scope_settlement;
mod send_control_commit;
mod send_control_ops;
mod send_descriptor_publication;
mod send_descriptor_terminal;
mod send_ops;
pub(crate) use public_types::*;
pub(crate) use runtime_types::*;
pub(crate) use send_descriptor_publication::*;
pub(crate) use send_descriptor_terminal::*;
impl<'r, const ROLE: u8, T, U, C, E, const MAX_RV: usize, Mint, B>
CursorEndpoint<'r, ROLE, T, U, C, E, MAX_RV, Mint, B>
where
T: Transport + 'r,
U: LabelUniverse,
C: crate::runtime::config::Clock,
E: EpochTable,
Mint: MintConfigMarker,
B: EndpointSlot,
{
pub(crate) fn matches_session(&self, sid: SessionId) -> bool {
self.sid == sid
}
pub(crate) fn for_each_physical_lane(&self, mut f: impl FnMut(Lane)) {
let logical_lane_count = self.cursor.logical_lane_count();
for slot in self.ports.iter().take(logical_lane_count) {
if let Some(port) = slot.as_ref() {
f(port.lane);
}
}
}
pub(crate) fn invalidate_public_owner(&mut self) {
self.public_header.invalidate();
self.public_generation = 0;
self.public_slot_owned = false;
}
pub(crate) fn prepare_public_owner_revocation(
&mut self,
terminal: &mut EndpointRevocationTerminal<'r>,
) {
terminal.set_waiter_lane(self.primary_physical_lane());
self.revoke_drain_public_send_terminal(terminal);
self.revoke_clear_public_recv_state();
self.revoke_clear_public_offer_state();
self.revoke_clear_public_decode_state();
if let Some(branch) = self.public_route_branch.take() {
branch.discard_terminal();
}
self.clear_public_op_terminal();
}
pub(crate) fn finish_public_owner_revocation(&mut self) {
self.invalidate_public_owner();
self.revoke_finish_public_send_state();
for port in self.ports.iter_mut() {
if let Some(port) = port.take() {
drop(port);
}
}
for guard in self.guards.iter_mut() {
if let Some(guard) = guard.as_mut() {
guard.detach_rendezvous();
}
if let Some(guard) = guard.take() {
drop(guard);
}
}
}
}
impl<'r, const ROLE: u8, T, U, C, E, const MAX_RV: usize, Mint, B> Drop
for CursorEndpoint<'r, ROLE, T, U, C, E, MAX_RV, Mint, B>
where
T: Transport + 'r,
U: LabelUniverse,
C: crate::runtime::config::Clock,
E: EpochTable,
Mint: MintConfigMarker,
B: EndpointSlot,
{
fn drop(&mut self) {
if self.public_generation != 0 && !self.cursor.is_terminal() {
let _ = self.poison_session(SessionFaultKind::EndpointDropped);
}
self.terminal_clear_public_send_state();
self.terminal_clear_public_recv_state();
self.terminal_clear_public_offer_state();
self.terminal_clear_public_decode_state();
if let Some(branch) = self.public_route_branch.take() {
branch.discard_terminal();
}
self.clear_public_op_terminal();
for port in self.ports.iter_mut() {
if let Some(p) = port.take() {
drop(p);
}
}
for guard in self.guards.iter_mut() {
if let Some(g) = guard.take() {
drop(g);
}
}
if self.public_generation != 0
&& let Some(cluster) = self.control.cluster()
{
if self.public_slot_owned {
cluster.release_public_endpoint_slot_owned(
self.public_rv,
self.public_slot,
self.public_generation,
);
}
self.public_header.invalidate();
self.public_generation = 0;
self.public_slot_owned = false;
}
}
}