use crate::batch::{ActiveRegionSet, BatchError, ExpectedBatch, TransferBatch};
use crate::control::{ControlError, ControlFrame};
use core::cell::Cell;
use core::marker::PhantomData;
#[cfg(any(target_os = "macos", target_os = "windows"))]
use core::sync::atomic::{AtomicBool, Ordering};
#[cfg(target_os = "linux")]
use core::sync::atomic::{AtomicI32, Ordering};
use std::ffi::OsString;
use std::num::NonZeroU32;
#[cfg(target_os = "linux")]
use std::os::fd::{FromRawFd, OwnedFd};
#[cfg(any(target_os = "linux", target_os = "macos", target_os = "windows"))]
use std::path::Path;
use std::path::PathBuf;
use std::time::{Duration, Instant};
pub use crate::liveness::{ActiveLeaseFacts, LeaseFactsConsistency};
#[cfg(target_os = "linux")]
const RECEIVER_BOOTSTRAP_ENV_PREFIX: &[u8] = b"NATIVE_IPC_VNEXT_BOOTSTRAP_FD=";
#[cfg(target_os = "linux")]
const RECEIVER_PUBLIC_BOOTSTRAP_ENV_ENTRY: &[u8] = b"NATIVE_IPC_VNEXT_PUBLIC_BOOTSTRAP=1";
#[cfg(target_os = "linux")]
const PR_GET_MDWE: libc::c_int = 66;
#[cfg(target_os = "linux")]
const PR_MDWE_REFUSE_EXEC_GAIN: libc::c_ulong = 1;
#[cfg(target_os = "linux")]
const BOOTSTRAP_ABSENT: i32 = -2;
#[cfg(target_os = "linux")]
const BOOTSTRAP_INVALID: i32 = -1;
#[cfg(target_os = "linux")]
const BOOTSTRAP_TAKEN: i32 = -3;
#[cfg(target_os = "linux")]
static RECEIVER_BOOTSTRAP_FD: AtomicI32 = AtomicI32::new(BOOTSTRAP_ABSENT);
#[cfg(any(target_os = "macos", target_os = "windows"))]
static RECEIVER_BOOTSTRAP_TAKEN: AtomicBool = AtomicBool::new(false);
#[doc(hidden)]
pub unsafe extern "C" fn __receiver_bootstrap_preinit(
_argument_count: core::ffi::c_int,
_arguments: *mut *mut core::ffi::c_char,
environment: *mut *mut core::ffi::c_char,
) {
#[cfg(target_os = "linux")]
unsafe {
receiver_bootstrap_preinit_linux(environment);
}
#[cfg(not(target_os = "linux"))]
let _ = environment;
}
#[cfg(target_os = "linux")]
unsafe fn receiver_bootstrap_preinit_linux(environment: *mut *mut libc::c_char) {
let mut entry = environment;
let public_bootstrap = loop {
if entry.is_null() {
return;
}
let candidate = unsafe { *entry };
if candidate.is_null() {
return;
}
let mut matches = true;
for (offset, expected) in RECEIVER_PUBLIC_BOOTSTRAP_ENV_ENTRY.iter().enumerate() {
let actual = unsafe { *candidate.add(offset) }.to_ne_bytes()[0];
if actual == 0 || actual != *expected {
matches = false;
break;
}
}
if matches
&& unsafe { *candidate.add(RECEIVER_PUBLIC_BOOTSTRAP_ENV_ENTRY.len()) } == 0
{
break candidate;
}
entry = unsafe { entry.add(1) };
};
unsafe { *public_bootstrap = 0 };
let mut entry = environment;
let value = loop {
if entry.is_null() {
return;
}
let candidate = unsafe { *entry };
if candidate.is_null() {
return;
}
let mut matches = true;
for (offset, expected) in RECEIVER_BOOTSTRAP_ENV_PREFIX.iter().enumerate() {
let actual = unsafe { *candidate.add(offset) }.to_ne_bytes()[0];
if actual == 0 || actual != *expected {
matches = false;
break;
}
}
if matches {
let value = unsafe { candidate.add(RECEIVER_BOOTSTRAP_ENV_PREFIX.len()) };
unsafe { *candidate = 0 };
break value;
}
entry = unsafe { entry.add(1) };
};
let mut raw = 0_i32;
let mut length = 0_usize;
loop {
if length == 10 {
RECEIVER_BOOTSTRAP_FD.store(BOOTSTRAP_INVALID, Ordering::Release);
return;
}
let byte = unsafe { *value.add(length) }.to_ne_bytes()[0];
if byte == 0 {
break;
}
if !byte.is_ascii_digit() || (length == 0 && byte == b'0') {
RECEIVER_BOOTSTRAP_FD.store(BOOTSTRAP_INVALID, Ordering::Release);
return;
}
let Some(next) = raw
.checked_mul(10)
.and_then(|current| current.checked_add(i32::from(byte - b'0')))
else {
RECEIVER_BOOTSTRAP_FD.store(BOOTSTRAP_INVALID, Ordering::Release);
return;
};
raw = next;
length += 1;
}
if length == 0 || raw < 3 {
RECEIVER_BOOTSTRAP_FD.store(BOOTSTRAP_INVALID, Ordering::Release);
return;
}
let descriptor_flags = unsafe { libc::fcntl(raw, libc::F_GETFD) };
let descriptor_status = unsafe { libc::fcntl(raw, libc::F_GETFL) };
let mdwe = unsafe { libc::prctl(PR_GET_MDWE, 0, 0, 0, 0) } as libc::c_ulong;
let pid = unsafe { libc::getpid() };
let sid = unsafe { libc::getsid(0) };
let process_group = unsafe { libc::getpgrp() };
let mut socket_type = 0_i32;
let mut socket_type_len = core::mem::size_of::<i32>() as libc::socklen_t;
let socket_result = unsafe {
libc::getsockopt(
raw,
libc::SOL_SOCKET,
libc::SO_TYPE,
(&mut socket_type as *mut i32).cast(),
&mut socket_type_len,
)
};
if descriptor_flags != 0
|| descriptor_status < 0
|| descriptor_status & libc::O_NONBLOCK == 0
|| mdwe != PR_MDWE_REFUSE_EXEC_GAIN
|| pid <= 0
|| sid != pid
|| process_group != pid
|| socket_result != 0
|| socket_type_len as usize != core::mem::size_of::<i32>()
|| socket_type != libc::SOCK_SEQPACKET
{
let _ = unsafe { libc::close(raw) };
RECEIVER_BOOTSTRAP_FD.store(BOOTSTRAP_INVALID, Ordering::Release);
return;
}
if unsafe { libc::fcntl(raw, libc::F_SETFD, libc::FD_CLOEXEC) } != 0 {
let _ = unsafe { libc::close(raw) };
RECEIVER_BOOTSTRAP_FD.store(BOOTSTRAP_INVALID, Ordering::Release);
return;
}
RECEIVER_BOOTSTRAP_FD.store(raw, Ordering::Release);
}
pub const HARD_MAX_REGIONS_PER_BATCH: u16 = 16;
pub const HARD_MAX_BOOTSTRAP_PAYLOAD_BYTES: u32 = 16 * 1024 * 1024;
pub const HARD_MAX_CONTROL_PAYLOAD_BYTES: u32 = 16 * 1024 * 1024;
pub const HARD_MAX_REGION_BYTES: u64 = 1 << 40;
pub const HARD_MAX_BATCH_BYTES: u64 = 1 << 42;
pub const HARD_MAX_ACTIVE_REGIONS: u32 = 1 << 20;
pub const HARD_MAX_ACTIVE_BYTES: u64 = 1 << 44;
pub const HARD_MAX_TRANSACTIONS: u64 = 1 << 48;
pub struct Coordinator;
pub struct Receiver;
pub struct Negotiating;
pub struct Ready;
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub enum BackendStatus {
Available,
Unavailable,
}
pub const fn backend_status() -> BackendStatus {
BackendStatus::Available
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub struct ProtocolVersion {
major: u16,
minor: u16,
}
impl ProtocolVersion {
#[allow(dead_code, reason = "wired into accepted session facts below")]
pub(crate) const fn new(major: u16, minor: u16) -> Self {
Self { major, minor }
}
pub const fn major(self) -> u16 {
self.major
}
pub const fn minor(self) -> u16 {
self.minor
}
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub enum SessionState {
Ready,
Poisoned,
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub enum PeerStatus {
Connected,
Disconnected,
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub enum ChildExitStatus {
Exited(i32),
Signaled {
signal: i32,
dumped_core: bool,
},
AlreadyReaped,
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub enum DescendantCleanupStatus {
NotEstablished,
FreshGroupUnverified,
FreshGroupTerminated,
ContainedProcessTreeComplete,
OwnedContainmentUnverified,
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub struct ChildCleanupFacts {
direct_child: Option<ChildExitStatus>,
descendants: DescendantCleanupStatus,
native_error: Option<i32>,
}
impl ChildCleanupFacts {
#[allow(dead_code, reason = "wired into coordinator lifecycle facts below")]
pub(crate) const fn new(
direct_child: Option<ChildExitStatus>,
descendants: DescendantCleanupStatus,
native_error: Option<i32>,
) -> Self {
Self {
direct_child,
descendants,
native_error,
}
}
pub const fn direct_child(self) -> Option<ChildExitStatus> {
self.direct_child
}
pub const fn descendants(self) -> DescendantCleanupStatus {
self.descendants
}
pub const fn native_error(self) -> Option<i32> {
self.native_error
}
pub const fn direct_child_complete(self) -> bool {
self.direct_child.is_some()
}
}
pub enum CoordinatorCloseOutcome {
Closed(ChildCleanupFacts),
ActiveLeases {
session: CoordinatorSession<Ready>,
facts: ActiveLeaseFacts,
},
CleanupPending {
session: CoordinatorSession<Ready>,
facts: ChildCleanupFacts,
failure: SessionFailure,
},
Failed {
session: CoordinatorSession<Ready>,
error: SessionFailure,
},
}
pub enum ReceiverCloseOutcome {
Closed,
ActiveLeases {
session: ReceiverSession<Ready>,
facts: ActiveLeaseFacts,
},
Failed {
session: ReceiverSession<Ready>,
error: SessionFailure,
},
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub struct CoordinatorAbortOutcome {
cleanup: ChildCleanupFacts,
failure: Option<SessionFailure>,
}
impl CoordinatorAbortOutcome {
pub const fn cleanup(self) -> ChildCleanupFacts {
self.cleanup
}
pub const fn failure(self) -> Option<SessionFailure> {
self.failure
}
}
pub struct Session<Role, State> {
inner: SessionInner,
role: PhantomData<Role>,
state: PhantomData<State>,
not_sync: PhantomData<Cell<()>>,
}
pub type CoordinatorSession<State> = Session<Coordinator, State>;
pub type ReceiverSession<State> = Session<Receiver, State>;
pub struct ReceiverBootstrap {
#[cfg(target_os = "linux")]
inherited: OwnedFd,
not_sync: PhantomData<Cell<()>>,
}
#[macro_export]
macro_rules! receiver_main {
($entry:expr) => {
#[cfg(target_os = "linux")]
#[used]
#[unsafe(link_section = ".preinit_array")]
static NATIVE_IPC_RECEIVER_BOOTSTRAP_PREINIT: unsafe extern "C" fn(
::core::ffi::c_int,
*mut *mut ::core::ffi::c_char,
*mut *mut ::core::ffi::c_char,
) = $crate::session::__receiver_bootstrap_preinit;
fn main() {
let bootstrap = $crate::session::__take_receiver_bootstrap();
($entry)(bootstrap);
}
};
}
#[doc(hidden)]
pub fn __take_receiver_bootstrap() -> Result<ReceiverBootstrap, SessionFailure> {
#[cfg(target_os = "linux")]
{
let raw = RECEIVER_BOOTSTRAP_FD.swap(BOOTSTRAP_TAKEN, Ordering::AcqRel);
if raw < 3 {
return Err(SessionFailure::new(
SessionOperation::Bootstrap,
SessionTransactionState::NotEstablished,
SessionError::InvalidInput,
));
}
let inherited = unsafe { OwnedFd::from_raw_fd(raw) };
Ok(ReceiverBootstrap {
inherited,
not_sync: PhantomData,
})
}
#[cfg(not(any(target_os = "linux", target_os = "macos", target_os = "windows")))]
{
Err(SessionFailure::new(
SessionOperation::Bootstrap,
SessionTransactionState::NotEstablished,
SessionError::BackendUnavailable,
))
}
#[cfg(target_os = "windows")]
{
if RECEIVER_BOOTSTRAP_TAKEN.swap(true, Ordering::AcqRel)
|| std::env::var_os("NATIVE_IPC_VNEXT_PUBLIC_BOOTSTRAP").as_deref()
!= Some(std::ffi::OsStr::new("1"))
{
return Err(SessionFailure::new(
SessionOperation::Bootstrap,
SessionTransactionState::NotEstablished,
SessionError::InvalidInput,
));
}
Ok(ReceiverBootstrap {
not_sync: PhantomData,
})
}
#[cfg(target_os = "macos")]
{
if RECEIVER_BOOTSTRAP_TAKEN.swap(true, Ordering::AcqRel)
|| std::env::var_os("NATIVE_IPC_VNEXT_PUBLIC_BOOTSTRAP").as_deref()
!= Some(std::ffi::OsStr::new("1"))
{
return Err(SessionFailure::new(
SessionOperation::Bootstrap,
SessionTransactionState::NotEstablished,
SessionError::InvalidInput,
));
}
unsafe { std::env::remove_var("NATIVE_IPC_VNEXT_PUBLIC_BOOTSTRAP") };
Ok(ReceiverBootstrap {
not_sync: PhantomData,
})
}
}
enum SessionInner {
#[cfg(target_os = "linux")]
CoordinatorNegotiating(crate::backend::linux_vnext::spawn::LinuxCoordinatorNegotiatingSession),
#[cfg(target_os = "linux")]
ReceiverNegotiating(crate::backend::linux_vnext::spawn::LinuxReceiverNegotiatingSession),
#[cfg(target_os = "linux")]
CoordinatorReady(crate::backend::linux_vnext::spawn::LinuxCoordinatorReadySession),
#[cfg(target_os = "linux")]
ReceiverReady(crate::backend::linux_vnext::spawn::LinuxReceiverReadySession),
#[cfg(target_os = "macos")]
CoordinatorNegotiating(crate::backend::macos::vnext_session::MacCoordinatorNegotiatingSession),
#[cfg(target_os = "macos")]
ReceiverNegotiating(crate::backend::macos::vnext_session::MacReceiverNegotiatingSession),
#[cfg(target_os = "macos")]
CoordinatorReady(crate::backend::macos::vnext_session::MacCoordinatorReadySession),
#[cfg(target_os = "macos")]
ReceiverReady(crate::backend::macos::vnext_session::MacReceiverReadySession),
#[cfg(target_os = "windows")]
CoordinatorNegotiating(
Box<crate::backend::windows::vnext_session::WindowsCoordinatorNegotiatingSession>,
),
#[cfg(target_os = "windows")]
ReceiverNegotiating(
Box<crate::backend::windows::vnext_session::WindowsReceiverNegotiatingSession>,
),
#[cfg(target_os = "windows")]
CoordinatorReady(Box<crate::backend::windows::vnext_session::WindowsCoordinatorReadySession>),
#[cfg(target_os = "windows")]
ReceiverReady(Box<crate::backend::windows::vnext_session::WindowsReceiverReadySession>),
#[allow(dead_code)]
Unavailable,
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub enum ExecutableIdentityPolicy {
ExactOpenedFile,
}
#[derive(Clone, Debug, Eq, PartialEq)]
pub struct SessionCommand {
executable: PathBuf,
arguments: Vec<OsString>,
environment: Vec<(OsString, OsString)>,
}
impl SessionCommand {
pub fn new(executable: impl Into<PathBuf>) -> Self {
let executable = executable.into();
Self {
arguments: vec![executable.as_os_str().to_owned()],
executable,
environment: Vec::new(),
}
}
pub fn arg0(mut self, argument: impl Into<OsString>) -> Self {
self.arguments[0] = argument.into();
self
}
pub fn arg(mut self, argument: impl Into<OsString>) -> Self {
self.arguments.push(argument.into());
self
}
pub fn env(mut self, key: impl Into<OsString>, value: impl Into<OsString>) -> Self {
let key = key.into();
let value = value.into();
if let Some((_, existing)) = self
.environment
.iter_mut()
.find(|(existing, _)| *existing == key)
{
*existing = value;
} else {
self.environment.push((key, value));
}
self
}
fn has_reserved_environment(&self) -> bool {
const RESERVED: [&str; 6] = [
"NATIVE_IPC_VNEXT_BOOTSTRAP_FD",
"NATIVE_IPC_VNEXT_PUBLIC_BOOTSTRAP",
"NATIVE_IPC_MACH_NONCE",
"NATIVE_IPC_PARENT_PID",
"NATIVE_IPC_WINDOWS_PIPE",
"NATIVE_IPC_WINDOWS_NONCE",
];
self.environment.iter().any(|(key, _)| {
RESERVED
.iter()
.any(|name| key.as_os_str() == std::ffi::OsStr::new(name))
})
}
#[cfg(any(target_os = "linux", target_os = "macos", target_os = "windows"))]
pub(crate) fn executable(&self) -> &Path {
&self.executable
}
#[cfg(any(target_os = "linux", target_os = "macos", target_os = "windows"))]
pub(crate) fn arguments(&self) -> &[OsString] {
&self.arguments
}
#[cfg(any(target_os = "linux", target_os = "macos", target_os = "windows"))]
pub(crate) fn environment(&self) -> &[(OsString, OsString)] {
&self.environment
}
}
#[derive(Clone, Debug, Eq, PartialEq)]
pub struct SessionOptions {
deadline: AbsoluteDeadline,
limits: SessionLimits,
application_payload: Vec<u8>,
executable_identity: ExecutableIdentityPolicy,
require_atomic_u32: bool,
require_atomic_u64: bool,
}
impl SessionOptions {
pub fn new(deadline: AbsoluteDeadline, executable_identity: ExecutableIdentityPolicy) -> Self {
Self {
deadline,
limits: SessionLimits::default(),
application_payload: Vec::new(),
executable_identity,
require_atomic_u32: false,
require_atomic_u64: false,
}
}
pub fn with_limits(mut self, limits: SessionLimits) -> Self {
self.limits = limits;
self
}
pub fn with_application_payload(mut self, payload: Vec<u8>) -> Self {
self.application_payload = payload;
self
}
pub fn require_atomic_u32(mut self) -> Self {
self.require_atomic_u32 = true;
self
}
pub fn require_atomic_u64(mut self) -> Self {
self.require_atomic_u64 = true;
self
}
#[cfg(any(target_os = "linux", target_os = "macos", target_os = "windows"))]
pub(crate) const fn limits(&self) -> SessionLimits {
self.limits
}
#[cfg(any(target_os = "linux", target_os = "macos", target_os = "windows"))]
pub(crate) fn application_payload(&self) -> &[u8] {
&self.application_payload
}
#[cfg(any(target_os = "linux", target_os = "macos", target_os = "windows"))]
pub(crate) const fn requires_atomic_u32(&self) -> bool {
self.require_atomic_u32
}
#[cfg(any(target_os = "linux", target_os = "macos", target_os = "windows"))]
pub(crate) const fn requires_atomic_u64(&self) -> bool {
self.require_atomic_u64
}
#[cfg(any(target_os = "linux", target_os = "macos", target_os = "windows"))]
pub(crate) const fn deadline(&self) -> AbsoluteDeadline {
self.deadline
}
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub enum SessionEndpoint {
Coordinator,
Receiver,
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub struct RejectionReason(NonZeroU32);
impl RejectionReason {
pub const APPLICATION_DECLINED: Self = Self(NonZeroU32::MIN);
pub const INCOMPATIBLE_APPLICATION_PROTOCOL: Self =
Self(NonZeroU32::new(2).expect("two is nonzero"));
pub const APPLICATION_POLICY: Self = Self(NonZeroU32::new(3).expect("three is nonzero"));
pub const fn application_specific(value: u32) -> Option<Self> {
if value < 0x8000_0000 {
return None;
}
match NonZeroU32::new(value) {
Some(value) => Some(Self(value)),
None => None,
}
}
pub const fn get(self) -> u32 {
self.0.get()
}
#[cfg(any(target_os = "linux", target_os = "macos", target_os = "windows"))]
fn from_wire(value: NonZeroU32) -> Option<Self> {
match value.get() {
1 => Some(Self::APPLICATION_DECLINED),
2 => Some(Self::INCOMPATIBLE_APPLICATION_PROTOCOL),
3 => Some(Self::APPLICATION_POLICY),
value => Self::application_specific(value),
}
}
#[cfg(any(target_os = "linux", target_os = "macos", target_os = "windows"))]
const fn as_nonzero(self) -> NonZeroU32 {
self.0
}
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub enum NegotiationDecision {
Accept,
Reject(RejectionReason),
}
pub enum NegotiationOutcome<T> {
Accepted(T),
Rejected {
by: SessionEndpoint,
reason: RejectionReason,
cleanup: Option<ChildCleanupFacts>,
},
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub enum SessionError {
BackendUnavailable,
InvalidInput,
DeadlineExpired,
PeerDisconnected,
IdentityMismatch,
MalformedPeer,
Ambiguous,
NegotiationFailed,
NativeNegotiation(NegotiationError),
Control(ControlError),
Batch(BatchError),
ActiveLimit,
PeerPreparationFailed,
ActivationFailed,
Poisoned,
Native,
}
impl core::fmt::Display for SessionError {
fn fmt(&self, formatter: &mut core::fmt::Formatter<'_>) -> core::fmt::Result {
write!(formatter, "session operation failed: {self:?}")
}
}
impl std::error::Error for SessionError {}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub enum SessionOperation {
Bootstrap,
Spawn,
Negotiate,
PollPeer,
WaitForExit,
Close,
Abort,
TransferBatch,
ReceiveBatch,
SendControl,
ReceiveControl,
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub enum SessionTransactionState {
NotEstablished,
Spawned,
Negotiating,
Ready,
TransactionOpen,
Poisoned,
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub struct SessionFailure {
operation: SessionOperation,
transaction_state: SessionTransactionState,
reason: SessionError,
native_code: Option<i32>,
poisoned: bool,
peer: Option<PeerStatus>,
cleanup: Option<ChildCleanupFacts>,
}
impl SessionFailure {
const fn new(
operation: SessionOperation,
transaction_state: SessionTransactionState,
reason: SessionError,
) -> Self {
Self {
operation,
transaction_state,
reason,
native_code: None,
poisoned: matches!(transaction_state, SessionTransactionState::Poisoned),
peer: if matches!(reason, SessionError::PeerDisconnected) {
Some(PeerStatus::Disconnected)
} else {
None
},
cleanup: None,
}
}
const fn with_cleanup(mut self, cleanup: ChildCleanupFacts) -> Self {
self.cleanup = Some(cleanup);
self
}
#[cfg(any(target_os = "linux", target_os = "macos", target_os = "windows"))]
const fn with_optional_cleanup(mut self, cleanup: Option<ChildCleanupFacts>) -> Self {
self.cleanup = cleanup;
self
}
const fn with_native_code(mut self, native_code: Option<i32>) -> Self {
self.native_code = native_code;
self
}
const fn with_poisoned(mut self, poisoned: bool) -> Self {
self.poisoned = poisoned;
self
}
pub const fn operation(self) -> SessionOperation {
self.operation
}
pub const fn transaction_state(self) -> SessionTransactionState {
self.transaction_state
}
pub const fn reason(self) -> SessionError {
self.reason
}
pub const fn native_code(self) -> Option<i32> {
self.native_code
}
pub const fn is_poisoned(self) -> bool {
self.poisoned
}
pub const fn peer(self) -> Option<PeerStatus> {
self.peer
}
pub const fn cleanup(self) -> Option<ChildCleanupFacts> {
self.cleanup
}
}
impl core::fmt::Display for SessionFailure {
fn fmt(&self, formatter: &mut core::fmt::Formatter<'_>) -> core::fmt::Result {
write!(
formatter,
"session {:?} failed in {:?}: {:?}",
self.operation, self.transaction_state, self.reason
)
}
}
impl std::error::Error for SessionFailure {}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub struct SessionLimits {
pub max_regions_per_batch: u16,
pub max_region_bytes: u64,
pub max_batch_bytes: u64,
pub max_active_regions: u32,
pub max_active_bytes: u64,
pub max_transactions: u64,
pub max_bootstrap_payload_bytes: u32,
pub max_control_payload_bytes: u32,
}
impl Default for SessionLimits {
fn default() -> Self {
Self {
max_regions_per_batch: 16,
max_region_bytes: 256 * 1024 * 1024,
max_batch_bytes: 1024 * 1024 * 1024,
max_active_regions: 4096,
max_active_bytes: 8 * 1024 * 1024 * 1024,
max_transactions: 1 << 32,
max_bootstrap_payload_bytes: 1024 * 1024,
max_control_payload_bytes: 1024 * 1024,
}
}
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub enum NegotiationError {
ZeroLimit,
AboveHardMaximum,
NativeSizeNarrowing,
AtomicUnsupported,
InvalidDeadline,
}
impl core::fmt::Display for NegotiationError {
fn fmt(&self, formatter: &mut core::fmt::Formatter<'_>) -> core::fmt::Result {
write!(formatter, "session negotiation failed: {self:?}")
}
}
impl std::error::Error for NegotiationError {}
impl SessionLimits {
pub fn validate(self) -> Result<Self, NegotiationError> {
self.validate_for_native_max(usize::MAX as u64)
}
fn validate_for_native_max(self, native_usize_max: u64) -> Result<Self, NegotiationError> {
if self.max_regions_per_batch == 0
|| self.max_region_bytes == 0
|| self.max_batch_bytes == 0
|| self.max_active_regions == 0
|| self.max_active_bytes == 0
|| self.max_transactions == 0
|| self.max_bootstrap_payload_bytes == 0
|| self.max_control_payload_bytes == 0
{
return Err(NegotiationError::ZeroLimit);
}
if self.max_regions_per_batch > HARD_MAX_REGIONS_PER_BATCH
|| self.max_region_bytes > HARD_MAX_REGION_BYTES
|| self.max_batch_bytes > HARD_MAX_BATCH_BYTES
|| self.max_active_regions > HARD_MAX_ACTIVE_REGIONS
|| self.max_active_bytes > HARD_MAX_ACTIVE_BYTES
|| self.max_transactions > HARD_MAX_TRANSACTIONS
|| self.max_bootstrap_payload_bytes > HARD_MAX_BOOTSTRAP_PAYLOAD_BYTES
|| self.max_control_payload_bytes > HARD_MAX_CONTROL_PAYLOAD_BYTES
{
return Err(NegotiationError::AboveHardMaximum);
}
if self.max_region_bytes > native_usize_max
|| self.max_batch_bytes > native_usize_max
|| self.max_active_bytes > native_usize_max
|| u64::from(self.max_bootstrap_payload_bytes) > native_usize_max
|| u64::from(self.max_control_payload_bytes) > native_usize_max
{
return Err(NegotiationError::NativeSizeNarrowing);
}
Ok(self)
}
pub fn negotiate(local: Self, peer: Self) -> Result<Self, NegotiationError> {
let local = local.validate()?;
let peer = peer.validate()?;
Self {
max_regions_per_batch: local.max_regions_per_batch.min(peer.max_regions_per_batch),
max_region_bytes: local.max_region_bytes.min(peer.max_region_bytes),
max_batch_bytes: local.max_batch_bytes.min(peer.max_batch_bytes),
max_active_regions: local.max_active_regions.min(peer.max_active_regions),
max_active_bytes: local.max_active_bytes.min(peer.max_active_bytes),
max_transactions: local.max_transactions.min(peer.max_transactions),
max_bootstrap_payload_bytes: local
.max_bootstrap_payload_bytes
.min(peer.max_bootstrap_payload_bytes),
max_control_payload_bytes: local
.max_control_payload_bytes
.min(peer.max_control_payload_bytes),
}
.validate()
}
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub struct AtomicCapabilities {
atomic_u32_lock_free: bool,
atomic_u32_alignment: usize,
atomic_u64_lock_free: bool,
atomic_u64_alignment: usize,
page_alignment: usize,
cache_line_alignment: usize,
}
impl AtomicCapabilities {
pub(crate) const fn from_accepted_offer(value: crate::negotiation::AtomicOffer) -> Self {
Self {
atomic_u32_lock_free: value.u32_lock_free,
atomic_u32_alignment: value.u32_alignment as usize,
atomic_u64_lock_free: value.u64_lock_free,
atomic_u64_alignment: value.u64_alignment as usize,
page_alignment: value.page_alignment as usize,
cache_line_alignment: value.cache_line_alignment as usize,
}
}
#[allow(dead_code, reason = "wired into native HELLO discovery in phase 4b")]
pub(crate) fn from_verified_native(
page_alignment: usize,
cache_line_alignment: usize,
atomic_u32_lock_free: bool,
atomic_u64_lock_free: bool,
) -> Result<Self, NegotiationError> {
let atomic_u32_alignment = core::mem::align_of::<core::sync::atomic::AtomicU32>();
let atomic_u64_alignment = core::mem::align_of::<core::sync::atomic::AtomicU64>();
if !page_alignment.is_power_of_two()
|| !cache_line_alignment.is_power_of_two()
|| page_alignment < atomic_u32_alignment.max(atomic_u64_alignment)
|| cache_line_alignment < atomic_u32_alignment.max(atomic_u64_alignment)
{
return Err(NegotiationError::AtomicUnsupported);
}
Ok(Self {
atomic_u32_lock_free,
atomic_u32_alignment,
atomic_u64_lock_free,
atomic_u64_alignment,
page_alignment,
cache_line_alignment,
})
}
pub fn atomic_u32_lock_free(self) -> bool {
self.atomic_u32_lock_free
}
pub fn atomic_u32_alignment(self) -> usize {
self.atomic_u32_alignment
}
pub fn atomic_u64_lock_free(self) -> bool {
self.atomic_u64_lock_free
}
pub fn atomic_u64_alignment(self) -> usize {
self.atomic_u64_alignment
}
pub fn page_alignment(self) -> usize {
self.page_alignment
}
pub fn cache_line_alignment(self) -> usize {
self.cache_line_alignment
}
#[allow(dead_code, reason = "wired into native HELLO negotiation in phase 4b")]
pub(crate) fn require(
self,
u32_required: bool,
u64_required: bool,
) -> Result<Self, NegotiationError> {
if (u32_required && !self.atomic_u32_lock_free)
|| (u64_required && !self.atomic_u64_lock_free)
{
return Err(NegotiationError::AtomicUnsupported);
}
Ok(self)
}
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub struct AbsoluteDeadline(Instant);
impl AbsoluteDeadline {
pub fn after(duration: Duration) -> Result<Self, NegotiationError> {
if duration.is_zero() {
return Err(NegotiationError::InvalidDeadline);
}
Instant::now()
.checked_add(duration)
.map(Self)
.ok_or(NegotiationError::InvalidDeadline)
}
pub fn remaining(self) -> Duration {
self.0.saturating_duration_since(Instant::now())
}
pub fn is_expired(self) -> bool {
self.remaining().is_zero()
}
}
#[cfg(any(target_os = "linux", target_os = "macos", target_os = "windows"))]
impl<Role, State> Session<Role, State> {
fn from_inner(inner: SessionInner) -> Self {
Self {
inner,
role: PhantomData,
state: PhantomData,
not_sync: PhantomData,
}
}
}
impl Session<Coordinator, Negotiating> {
pub fn spawn(command: SessionCommand, options: SessionOptions) -> Result<Self, SessionFailure> {
validate_public_options(&options).map_err(|reason| {
SessionFailure::new(
SessionOperation::Spawn,
SessionTransactionState::NotEstablished,
reason,
)
})?;
if command.has_reserved_environment() {
return Err(SessionFailure::new(
SessionOperation::Spawn,
SessionTransactionState::NotEstablished,
SessionError::InvalidInput,
));
}
#[cfg(target_os = "linux")]
{
let inner =
crate::backend::linux_vnext::spawn::LinuxCoordinatorNegotiatingSession::spawn(
&command, &options,
)
.map_err(|failure| {
let native_code = linux_public_native_code(failure.error);
let transaction_state = match failure.state {
crate::backend::linux_vnext::spawn::LinuxCoordinatorFailureState::NotEstablished => {
SessionTransactionState::NotEstablished
}
crate::backend::linux_vnext::spawn::LinuxCoordinatorFailureState::Spawned => {
SessionTransactionState::Spawned
}
crate::backend::linux_vnext::spawn::LinuxCoordinatorFailureState::Negotiating => {
SessionTransactionState::Negotiating
}
};
SessionFailure::new(
SessionOperation::Spawn,
transaction_state,
failure.error.into(),
)
.with_native_code(native_code)
.with_poisoned(failure.poisoned)
.with_optional_cleanup(failure.cleanup)
})?;
Ok(Self::from_inner(SessionInner::CoordinatorNegotiating(
inner,
)))
}
#[cfg(target_os = "macos")]
{
let inner =
crate::backend::macos::vnext_session::MacCoordinatorNegotiatingSession::spawn(
&command, &options,
)
.map_err(|failure| {
let transaction_state = match failure.state {
crate::backend::macos::vnext_session::MacCoordinatorFailureState::NotEstablished => {
SessionTransactionState::NotEstablished
}
crate::backend::macos::vnext_session::MacCoordinatorFailureState::Spawned => {
SessionTransactionState::Spawned
}
crate::backend::macos::vnext_session::MacCoordinatorFailureState::Negotiating => {
SessionTransactionState::Negotiating
}
};
mac_session_failure(
SessionOperation::Spawn,
transaction_state,
failure.error,
failure.poisoned,
)
.with_optional_cleanup(failure.cleanup)
})?;
Ok(Self::from_inner(SessionInner::CoordinatorNegotiating(
inner,
)))
}
#[cfg(target_os = "windows")]
{
let inner = crate::backend::windows::vnext_session::WindowsCoordinatorNegotiatingSession::spawn(
&command,
&options,
)
.map_err(|failure| {
let transaction_state = match failure.state {
crate::backend::windows::vnext_session::WindowsCoordinatorFailureState::NotEstablished => SessionTransactionState::NotEstablished,
crate::backend::windows::vnext_session::WindowsCoordinatorFailureState::Spawned => SessionTransactionState::Spawned,
crate::backend::windows::vnext_session::WindowsCoordinatorFailureState::Negotiating => SessionTransactionState::Negotiating,
};
windows_session_failure(
SessionOperation::Spawn,
transaction_state,
failure.error,
failure.poisoned,
)
.with_optional_cleanup(failure.cleanup)
})?;
Ok(Self::from_inner(SessionInner::CoordinatorNegotiating(
Box::new(inner),
)))
}
#[cfg(not(any(target_os = "linux", target_os = "macos", target_os = "windows")))]
{
let _ = command;
Err(SessionFailure::new(
SessionOperation::Spawn,
SessionTransactionState::NotEstablished,
SessionError::BackendUnavailable,
))
}
}
pub fn peer_application_payload(&self) -> &[u8] {
match &self.inner {
#[cfg(target_os = "linux")]
SessionInner::CoordinatorNegotiating(inner) => inner.peer_application_payload(),
#[cfg(target_os = "macos")]
SessionInner::CoordinatorNegotiating(inner) => inner.peer_application_payload(),
#[cfg(target_os = "windows")]
SessionInner::CoordinatorNegotiating(inner) => inner.peer_application_payload(),
#[cfg(target_os = "linux")]
_ => unreachable!("coordinator negotiating typestate owns its exact backend state"),
#[cfg(target_os = "macos")]
_ => unreachable!("coordinator negotiating typestate owns its exact backend state"),
#[cfg(target_os = "windows")]
_ => unreachable!("coordinator negotiating typestate owns its exact backend state"),
#[cfg(not(any(target_os = "linux", target_os = "macos", target_os = "windows")))]
SessionInner::Unavailable => {
unreachable!("unavailable backend cannot construct a session")
}
}
}
pub fn decide(
self,
decision: NegotiationDecision,
) -> Result<NegotiationOutcome<Session<Coordinator, Ready>>, SessionFailure> {
#[cfg(not(any(target_os = "linux", target_os = "macos", target_os = "windows")))]
let _ = decision;
match self.inner {
#[cfg(target_os = "linux")]
SessionInner::CoordinatorNegotiating(inner) => {
let outcome = inner
.decide(decision_rejection(decision))
.map_err(|failure| {
let native_code = linux_public_native_code(failure.error);
SessionFailure::new(
SessionOperation::Negotiate,
SessionTransactionState::Negotiating,
failure.error.into(),
)
.with_native_code(native_code)
.with_poisoned(failure.poisoned)
.with_optional_cleanup(failure.cleanup)
})?;
map_linux_coordinator_outcome(outcome)
}
#[cfg(target_os = "macos")]
SessionInner::CoordinatorNegotiating(inner) => {
let outcome = inner
.decide(decision_rejection(decision))
.map_err(|failure| {
mac_session_failure(
SessionOperation::Negotiate,
SessionTransactionState::Negotiating,
failure.error,
failure.poisoned,
)
.with_optional_cleanup(failure.cleanup)
})?;
map_mac_coordinator_outcome(outcome)
}
#[cfg(target_os = "windows")]
SessionInner::CoordinatorNegotiating(inner) => {
let outcome = (*inner)
.decide(decision_rejection(decision))
.map_err(|failure| {
windows_session_failure(
SessionOperation::Negotiate,
SessionTransactionState::Negotiating,
failure.error,
failure.poisoned,
)
.with_optional_cleanup(failure.cleanup)
})?;
map_windows_coordinator_outcome(outcome)
}
#[cfg(target_os = "linux")]
_ => unreachable!("coordinator negotiating typestate owns its exact backend state"),
#[cfg(target_os = "macos")]
_ => unreachable!("coordinator negotiating typestate owns its exact backend state"),
#[cfg(target_os = "windows")]
_ => unreachable!("coordinator negotiating typestate owns its exact backend state"),
#[cfg(not(any(target_os = "linux", target_os = "macos", target_os = "windows")))]
SessionInner::Unavailable => {
unreachable!("unavailable backend cannot construct a session")
}
}
}
}
impl Session<Receiver, Negotiating> {
pub fn from_bootstrap(
bootstrap: ReceiverBootstrap,
options: SessionOptions,
) -> Result<Self, SessionFailure> {
validate_public_options(&options).map_err(|reason| {
SessionFailure::new(
SessionOperation::Bootstrap,
SessionTransactionState::NotEstablished,
reason,
)
})?;
#[cfg(target_os = "linux")]
{
let inner = crate::backend::linux_vnext::spawn::LinuxReceiverNegotiatingSession::from_inherited_bootstrap(
bootstrap.inherited,
options.limits,
options.application_payload,
options.require_atomic_u32,
options.require_atomic_u64,
options.deadline,
)
.map_err(|error| {
SessionFailure::new(
SessionOperation::Bootstrap,
SessionTransactionState::Negotiating,
error.into(),
)
.with_native_code(linux_public_native_code(error))
.with_poisoned(true)
})?;
Ok(Self::from_inner(SessionInner::ReceiverNegotiating(inner)))
}
#[cfg(target_os = "macos")]
{
let _bootstrap = bootstrap;
let inner =
crate::backend::macos::vnext_session::MacReceiverNegotiatingSession::from_environment(
options.limits,
options.application_payload,
options.require_atomic_u32,
options.require_atomic_u64,
options.deadline,
)
.map_err(|error| {
let invalid_input = matches!(
error,
crate::backend::macos::vnext_session::MacPublicSessionError::InvalidInput
);
let state = if invalid_input {
SessionTransactionState::NotEstablished
} else {
SessionTransactionState::Negotiating
};
mac_session_failure(
SessionOperation::Bootstrap,
state,
error,
!invalid_input,
)
})?;
Ok(Self::from_inner(SessionInner::ReceiverNegotiating(inner)))
}
#[cfg(target_os = "windows")]
{
let _bootstrap = bootstrap;
let inner = crate::backend::windows::vnext_session::WindowsReceiverNegotiatingSession::from_environment(&options)
.map_err(|error| {
let invalid_input = matches!(
error,
crate::backend::windows::vnext_session::WindowsPublicSessionError::InvalidInput
);
let state = if invalid_input {
SessionTransactionState::NotEstablished
} else {
SessionTransactionState::Negotiating
};
windows_session_failure(
SessionOperation::Bootstrap,
state,
error,
!invalid_input,
)
})?;
Ok(Self::from_inner(SessionInner::ReceiverNegotiating(
Box::new(inner),
)))
}
#[cfg(not(any(target_os = "linux", target_os = "macos", target_os = "windows")))]
{
let _ = bootstrap;
Err(SessionFailure::new(
SessionOperation::Bootstrap,
SessionTransactionState::NotEstablished,
SessionError::BackendUnavailable,
))
}
}
pub fn peer_application_payload(&self) -> &[u8] {
match &self.inner {
#[cfg(target_os = "linux")]
SessionInner::ReceiverNegotiating(inner) => inner.peer_application_payload(),
#[cfg(target_os = "macos")]
SessionInner::ReceiverNegotiating(inner) => inner.peer_application_payload(),
#[cfg(target_os = "windows")]
SessionInner::ReceiverNegotiating(inner) => inner.peer_application_payload(),
#[cfg(target_os = "linux")]
_ => unreachable!("receiver negotiating typestate owns its exact backend state"),
#[cfg(target_os = "macos")]
_ => unreachable!("receiver negotiating typestate owns its exact backend state"),
#[cfg(target_os = "windows")]
_ => unreachable!("receiver negotiating typestate owns its exact backend state"),
#[cfg(not(any(target_os = "linux", target_os = "macos", target_os = "windows")))]
SessionInner::Unavailable => {
unreachable!("unavailable backend cannot construct a session")
}
}
}
pub fn decide_after_coordinator(
self,
decide: impl FnOnce(&[u8]) -> NegotiationDecision,
) -> Result<NegotiationOutcome<Session<Receiver, Ready>>, SessionFailure> {
#[cfg(not(any(target_os = "linux", target_os = "macos", target_os = "windows")))]
let _ = decide;
match self.inner {
#[cfg(target_os = "linux")]
SessionInner::ReceiverNegotiating(inner) => {
let outcome = inner
.decide_after_coordinator(|payload| decision_rejection(decide(payload)))
.map_err(|error| {
SessionFailure::new(
SessionOperation::Negotiate,
SessionTransactionState::Negotiating,
error.into(),
)
.with_native_code(linux_public_native_code(error))
.with_poisoned(true)
})?;
map_linux_receiver_outcome(outcome)
}
#[cfg(target_os = "macos")]
SessionInner::ReceiverNegotiating(inner) => {
let outcome = inner
.decide_after_coordinator(|payload| decision_rejection(decide(payload)))
.map_err(|error| {
mac_session_failure(
SessionOperation::Negotiate,
SessionTransactionState::Negotiating,
error,
true,
)
})?;
map_mac_receiver_outcome(outcome)
}
#[cfg(target_os = "windows")]
SessionInner::ReceiverNegotiating(inner) => {
let outcome = (*inner)
.decide_after_coordinator(|payload| decision_rejection(decide(payload)))
.map_err(|error| {
windows_session_failure(
SessionOperation::Negotiate,
SessionTransactionState::Negotiating,
error,
true,
)
})?;
map_windows_receiver_outcome(outcome)
}
#[cfg(target_os = "linux")]
_ => unreachable!("receiver negotiating typestate owns its exact backend state"),
#[cfg(target_os = "macos")]
_ => unreachable!("receiver negotiating typestate owns its exact backend state"),
#[cfg(target_os = "windows")]
_ => unreachable!("receiver negotiating typestate owns its exact backend state"),
#[cfg(not(any(target_os = "linux", target_os = "macos", target_os = "windows")))]
SessionInner::Unavailable => {
unreachable!("unavailable backend cannot construct a session")
}
}
}
}
impl Session<Coordinator, Ready> {
#[cfg(all(test, target_os = "linux"))]
pub(crate) fn fail_next_cleanup_signal_for_test(&self, code: i32) {
match &self.inner {
SessionInner::CoordinatorReady(inner) => {
inner.fail_next_cleanup_signal_for_test(code);
}
_ => unreachable!("coordinator ready typestate owns its exact backend state"),
}
}
pub fn negotiated_limits(&self) -> SessionLimits {
match &self.inner {
#[cfg(any(target_os = "linux", target_os = "macos", target_os = "windows"))]
SessionInner::CoordinatorReady(inner) => inner.limits(),
#[cfg(any(target_os = "linux", target_os = "macos", target_os = "windows"))]
_ => unreachable!("coordinator ready typestate owns its exact backend state"),
#[cfg(not(any(target_os = "linux", target_os = "macos", target_os = "windows")))]
SessionInner::Unavailable => {
unreachable!("unavailable backend cannot construct a session")
}
}
}
pub fn atomic_capabilities(&self) -> AtomicCapabilities {
match &self.inner {
#[cfg(any(target_os = "linux", target_os = "macos", target_os = "windows"))]
SessionInner::CoordinatorReady(inner) => inner.atomics(),
#[cfg(any(target_os = "linux", target_os = "macos", target_os = "windows"))]
_ => unreachable!("coordinator ready typestate owns its exact backend state"),
#[cfg(not(any(target_os = "linux", target_os = "macos", target_os = "windows")))]
SessionInner::Unavailable => {
unreachable!("unavailable backend cannot construct a session")
}
}
}
pub fn protocol_version(&self) -> ProtocolVersion {
match &self.inner {
#[cfg(any(target_os = "linux", target_os = "macos", target_os = "windows"))]
SessionInner::CoordinatorReady(inner) => inner.protocol_version(),
#[cfg(any(target_os = "linux", target_os = "macos", target_os = "windows"))]
_ => unreachable!("coordinator ready typestate owns its exact backend state"),
#[cfg(not(any(target_os = "linux", target_os = "macos", target_os = "windows")))]
SessionInner::Unavailable => {
unreachable!("unavailable backend cannot construct a session")
}
}
}
pub fn state(&self) -> SessionState {
match &self.inner {
#[cfg(any(target_os = "linux", target_os = "macos", target_os = "windows"))]
SessionInner::CoordinatorReady(inner) => inner.state(),
#[cfg(any(target_os = "linux", target_os = "macos", target_os = "windows"))]
_ => unreachable!("coordinator ready typestate owns its exact backend state"),
#[cfg(not(any(target_os = "linux", target_os = "macos", target_os = "windows")))]
SessionInner::Unavailable => {
unreachable!("unavailable backend cannot construct a session")
}
}
}
pub fn active_leases(&self) -> ActiveLeaseFacts {
match &self.inner {
#[cfg(any(target_os = "linux", target_os = "macos", target_os = "windows"))]
SessionInner::CoordinatorReady(inner) => inner.active_leases(),
#[cfg(any(target_os = "linux", target_os = "macos", target_os = "windows"))]
_ => unreachable!("coordinator ready typestate owns its exact backend state"),
#[cfg(not(any(target_os = "linux", target_os = "macos", target_os = "windows")))]
SessionInner::Unavailable => {
unreachable!("unavailable backend cannot construct a session")
}
}
}
pub fn poll_peer(&mut self) -> Result<PeerStatus, SessionFailure> {
match &mut self.inner {
#[cfg(target_os = "linux")]
SessionInner::CoordinatorReady(inner) => {
let result = inner.poll_peer();
let state = inner.state();
result
.map_err(|error| linux_ready_failure(SessionOperation::PollPeer, state, error))
}
#[cfg(target_os = "macos")]
SessionInner::CoordinatorReady(inner) => {
let result = inner.poll_peer();
let state = inner.state();
result.map_err(|error| mac_ready_failure(SessionOperation::PollPeer, state, error))
}
#[cfg(target_os = "windows")]
SessionInner::CoordinatorReady(inner) => {
let result = inner.poll_peer();
let state = inner.state();
result.map_err(|error| {
windows_ready_failure(SessionOperation::PollPeer, state, error)
})
}
#[cfg(target_os = "linux")]
_ => unreachable!("coordinator ready typestate owns its exact backend state"),
#[cfg(target_os = "macos")]
_ => unreachable!("coordinator ready typestate owns its exact backend state"),
#[cfg(target_os = "windows")]
_ => unreachable!("coordinator ready typestate owns its exact backend state"),
#[cfg(not(any(target_os = "linux", target_os = "macos", target_os = "windows")))]
SessionInner::Unavailable => Err(SessionFailure::new(
SessionOperation::PollPeer,
SessionTransactionState::NotEstablished,
SessionError::BackendUnavailable,
)),
}
}
pub fn wait_for_exit(&mut self, deadline: AbsoluteDeadline) -> ChildCleanupFacts {
#[cfg(not(any(target_os = "linux", target_os = "macos", target_os = "windows")))]
let _ = deadline;
match &mut self.inner {
#[cfg(any(target_os = "linux", target_os = "macos", target_os = "windows"))]
SessionInner::CoordinatorReady(inner) => inner.wait_for_exit(deadline),
#[cfg(any(target_os = "linux", target_os = "macos", target_os = "windows"))]
_ => unreachable!("coordinator ready typestate owns its exact backend state"),
#[cfg(not(any(target_os = "linux", target_os = "macos", target_os = "windows")))]
SessionInner::Unavailable => {
ChildCleanupFacts::new(None, DescendantCleanupStatus::NotEstablished, None)
}
}
}
pub fn try_close(mut self, deadline: AbsoluteDeadline) -> CoordinatorCloseOutcome {
let facts = self.active_leases();
if !facts.is_empty() {
return CoordinatorCloseOutcome::ActiveLeases {
session: self,
facts,
};
}
let cleanup = self.wait_for_exit(deadline);
if !cleanup.direct_child_complete() {
let reason = if cleanup.native_error().is_some() {
SessionError::Native
} else {
SessionError::DeadlineExpired
};
let failure = SessionFailure::new(
SessionOperation::Close,
if self.state() == SessionState::Poisoned {
SessionTransactionState::Poisoned
} else {
SessionTransactionState::Ready
},
reason,
)
.with_native_code(cleanup.native_error())
.with_poisoned(self.state() == SessionState::Poisoned)
.with_cleanup(cleanup);
return CoordinatorCloseOutcome::CleanupPending {
session: self,
facts: cleanup,
failure,
};
}
let close: Result<(), SessionFailure> = match &mut self.inner {
#[cfg(target_os = "linux")]
SessionInner::CoordinatorReady(inner) => {
let result = inner.close_resources();
let state = inner.state();
result.map_err(|error| {
linux_ready_failure(SessionOperation::Close, state, error).with_cleanup(cleanup)
})
}
#[cfg(target_os = "macos")]
SessionInner::CoordinatorReady(inner) => {
let result = inner.close_resources();
let state = inner.state();
result.map_err(|error| {
mac_ready_failure(SessionOperation::Close, state, error).with_cleanup(cleanup)
})
}
#[cfg(target_os = "windows")]
SessionInner::CoordinatorReady(inner) => {
let result = inner.close_resources();
let state = inner.state();
result.map_err(|error| {
windows_ready_failure(SessionOperation::Close, state, error)
.with_cleanup(cleanup)
})
}
#[cfg(target_os = "linux")]
_ => unreachable!("coordinator ready typestate owns its exact backend state"),
#[cfg(target_os = "macos")]
_ => unreachable!("coordinator ready typestate owns its exact backend state"),
#[cfg(target_os = "windows")]
_ => unreachable!("coordinator ready typestate owns its exact backend state"),
#[cfg(not(any(target_os = "linux", target_os = "macos", target_os = "windows")))]
SessionInner::Unavailable => Err(SessionFailure::new(
SessionOperation::Close,
SessionTransactionState::NotEstablished,
SessionError::BackendUnavailable,
)
.with_cleanup(cleanup)),
};
if let Err(error) = close {
return CoordinatorCloseOutcome::Failed {
session: self,
error,
};
}
CoordinatorCloseOutcome::Closed(cleanup)
}
pub fn abort(mut self, deadline: AbsoluteDeadline) -> CoordinatorAbortOutcome {
#[cfg(not(any(target_os = "linux", target_os = "macos", target_os = "windows")))]
let _ = deadline;
let cleanup = match &mut self.inner {
#[cfg(any(target_os = "linux", target_os = "macos", target_os = "windows"))]
SessionInner::CoordinatorReady(inner) => inner.abort(deadline),
#[cfg(any(target_os = "linux", target_os = "macos", target_os = "windows"))]
_ => unreachable!("coordinator ready typestate owns its exact backend state"),
#[cfg(not(any(target_os = "linux", target_os = "macos", target_os = "windows")))]
SessionInner::Unavailable => {
ChildCleanupFacts::new(None, DescendantCleanupStatus::NotEstablished, None)
}
};
let failure = if cleanup.direct_child_complete() {
None
} else {
let reason = if cleanup.native_error().is_some() {
SessionError::Native
} else {
SessionError::DeadlineExpired
};
Some(
SessionFailure::new(
SessionOperation::Abort,
SessionTransactionState::Poisoned,
reason,
)
.with_native_code(cleanup.native_error())
.with_poisoned(true)
.with_cleanup(cleanup),
)
};
CoordinatorAbortOutcome { cleanup, failure }
}
pub fn new_transfer_batch(&self) -> Result<TransferBatch, BatchError> {
let limits = self.negotiated_limits();
TransferBatch::new(
limits.max_regions_per_batch,
limits.max_region_bytes,
limits.max_batch_bytes,
)
}
pub fn transfer_batch(
&mut self,
batch: TransferBatch,
deadline: AbsoluteDeadline,
) -> Result<ActiveRegionSet, SessionFailure> {
#[cfg(not(any(target_os = "linux", target_os = "macos", target_os = "windows")))]
let _ = (batch, deadline);
match &mut self.inner {
#[cfg(target_os = "linux")]
SessionInner::CoordinatorReady(inner) => {
let result = inner.transfer_batch(batch, deadline);
let state = inner.state();
result.map_err(|error| {
linux_ready_batch_failure(SessionOperation::TransferBatch, state, error)
})
}
#[cfg(target_os = "macos")]
SessionInner::CoordinatorReady(inner) => {
let result = inner.transfer_batch(batch, deadline);
let state = inner.state();
result.map_err(|error| {
mac_ready_failure(SessionOperation::TransferBatch, state, error)
})
}
#[cfg(target_os = "windows")]
SessionInner::CoordinatorReady(inner) => {
let result = inner.transfer_batch(batch, deadline);
let state = inner.state();
result.map_err(|error| {
windows_ready_failure(SessionOperation::TransferBatch, state, error)
})
}
#[cfg(target_os = "linux")]
_ => unreachable!("coordinator ready typestate owns its exact backend state"),
#[cfg(target_os = "macos")]
_ => unreachable!("coordinator ready typestate owns its exact backend state"),
#[cfg(target_os = "windows")]
_ => unreachable!("coordinator ready typestate owns its exact backend state"),
#[cfg(not(any(target_os = "linux", target_os = "macos", target_os = "windows")))]
SessionInner::Unavailable => {
unreachable!("unavailable backend cannot construct a session")
}
}
}
pub fn send_control(
&mut self,
kind: u32,
payload: &[u8],
deadline: AbsoluteDeadline,
) -> Result<(), SessionFailure> {
#[cfg(not(any(target_os = "linux", target_os = "macos", target_os = "windows")))]
let _ = (kind, payload, deadline);
match &mut self.inner {
#[cfg(target_os = "linux")]
SessionInner::CoordinatorReady(inner) => {
let result = inner.send_control(kind, payload, deadline);
let state = inner.state();
result.map_err(|error| {
linux_ready_failure(SessionOperation::SendControl, state, error)
})
}
#[cfg(target_os = "macos")]
SessionInner::CoordinatorReady(inner) => {
let result = inner.send_control(kind, payload, deadline);
let state = inner.state();
result
.map_err(|error| mac_ready_failure(SessionOperation::SendControl, state, error))
}
#[cfg(target_os = "windows")]
SessionInner::CoordinatorReady(inner) => {
let result = inner.send_control(kind, payload, deadline);
let state = inner.state();
result.map_err(|error| {
windows_ready_failure(SessionOperation::SendControl, state, error)
})
}
#[cfg(target_os = "linux")]
_ => unreachable!("coordinator ready typestate owns its exact backend state"),
#[cfg(target_os = "macos")]
_ => unreachable!("coordinator ready typestate owns its exact backend state"),
#[cfg(target_os = "windows")]
_ => unreachable!("coordinator ready typestate owns its exact backend state"),
#[cfg(not(any(target_os = "linux", target_os = "macos", target_os = "windows")))]
SessionInner::Unavailable => {
unreachable!("unavailable backend cannot construct a session")
}
}
}
pub fn receive_control(
&mut self,
deadline: AbsoluteDeadline,
) -> Result<ControlFrame, SessionFailure> {
#[cfg(not(any(target_os = "linux", target_os = "macos", target_os = "windows")))]
let _ = deadline;
match &mut self.inner {
#[cfg(target_os = "linux")]
SessionInner::CoordinatorReady(inner) => {
let result = inner.receive_control(deadline);
let state = inner.state();
result.map_err(|error| {
linux_ready_failure(SessionOperation::ReceiveControl, state, error)
})
}
#[cfg(target_os = "macos")]
SessionInner::CoordinatorReady(inner) => {
let result = inner.receive_control(deadline);
let state = inner.state();
result.map_err(|error| {
mac_ready_failure(SessionOperation::ReceiveControl, state, error)
})
}
#[cfg(target_os = "windows")]
SessionInner::CoordinatorReady(inner) => {
let result = inner.receive_control(deadline);
let state = inner.state();
result.map_err(|error| {
windows_ready_failure(SessionOperation::ReceiveControl, state, error)
})
}
#[cfg(target_os = "linux")]
_ => unreachable!("coordinator ready typestate owns its exact backend state"),
#[cfg(target_os = "macos")]
_ => unreachable!("coordinator ready typestate owns its exact backend state"),
#[cfg(target_os = "windows")]
_ => unreachable!("coordinator ready typestate owns its exact backend state"),
#[cfg(not(any(target_os = "linux", target_os = "macos", target_os = "windows")))]
SessionInner::Unavailable => {
unreachable!("unavailable backend cannot construct a session")
}
}
}
}
impl Session<Receiver, Ready> {
pub fn negotiated_limits(&self) -> SessionLimits {
match &self.inner {
#[cfg(any(target_os = "linux", target_os = "macos", target_os = "windows"))]
SessionInner::ReceiverReady(inner) => inner.limits(),
#[cfg(any(target_os = "linux", target_os = "macos", target_os = "windows"))]
_ => unreachable!("receiver ready typestate owns its exact backend state"),
#[cfg(not(any(target_os = "linux", target_os = "macos", target_os = "windows")))]
SessionInner::Unavailable => {
unreachable!("unavailable backend cannot construct a session")
}
}
}
pub fn atomic_capabilities(&self) -> AtomicCapabilities {
match &self.inner {
#[cfg(any(target_os = "linux", target_os = "macos", target_os = "windows"))]
SessionInner::ReceiverReady(inner) => inner.atomics(),
#[cfg(any(target_os = "linux", target_os = "macos", target_os = "windows"))]
_ => unreachable!("receiver ready typestate owns its exact backend state"),
#[cfg(not(any(target_os = "linux", target_os = "macos", target_os = "windows")))]
SessionInner::Unavailable => {
unreachable!("unavailable backend cannot construct a session")
}
}
}
pub fn protocol_version(&self) -> ProtocolVersion {
match &self.inner {
#[cfg(any(target_os = "linux", target_os = "macos", target_os = "windows"))]
SessionInner::ReceiverReady(inner) => inner.protocol_version(),
#[cfg(any(target_os = "linux", target_os = "macos", target_os = "windows"))]
_ => unreachable!("receiver ready typestate owns its exact backend state"),
#[cfg(not(any(target_os = "linux", target_os = "macos", target_os = "windows")))]
SessionInner::Unavailable => {
unreachable!("unavailable backend cannot construct a session")
}
}
}
pub fn state(&self) -> SessionState {
match &self.inner {
#[cfg(any(target_os = "linux", target_os = "macos", target_os = "windows"))]
SessionInner::ReceiverReady(inner) => inner.state(),
#[cfg(any(target_os = "linux", target_os = "macos", target_os = "windows"))]
_ => unreachable!("receiver ready typestate owns its exact backend state"),
#[cfg(not(any(target_os = "linux", target_os = "macos", target_os = "windows")))]
SessionInner::Unavailable => {
unreachable!("unavailable backend cannot construct a session")
}
}
}
pub fn active_leases(&self) -> ActiveLeaseFacts {
match &self.inner {
#[cfg(any(target_os = "linux", target_os = "macos", target_os = "windows"))]
SessionInner::ReceiverReady(inner) => inner.active_leases(),
#[cfg(any(target_os = "linux", target_os = "macos", target_os = "windows"))]
_ => unreachable!("receiver ready typestate owns its exact backend state"),
#[cfg(not(any(target_os = "linux", target_os = "macos", target_os = "windows")))]
SessionInner::Unavailable => {
unreachable!("unavailable backend cannot construct a session")
}
}
}
pub fn poll_peer(&mut self) -> Result<PeerStatus, SessionFailure> {
match &mut self.inner {
#[cfg(target_os = "linux")]
SessionInner::ReceiverReady(inner) => {
let result = inner.poll_peer();
let state = inner.state();
result
.map_err(|error| linux_ready_failure(SessionOperation::PollPeer, state, error))
}
#[cfg(target_os = "macos")]
SessionInner::ReceiverReady(inner) => {
let result = inner.poll_peer();
let state = inner.state();
result.map_err(|error| mac_ready_failure(SessionOperation::PollPeer, state, error))
}
#[cfg(target_os = "windows")]
SessionInner::ReceiverReady(inner) => {
let result = inner.poll_peer();
let state = inner.state();
result.map_err(|error| {
windows_ready_failure(SessionOperation::PollPeer, state, error)
})
}
#[cfg(target_os = "linux")]
_ => unreachable!("receiver ready typestate owns its exact backend state"),
#[cfg(target_os = "macos")]
_ => unreachable!("receiver ready typestate owns its exact backend state"),
#[cfg(target_os = "windows")]
_ => unreachable!("receiver ready typestate owns its exact backend state"),
#[cfg(not(any(target_os = "linux", target_os = "macos", target_os = "windows")))]
SessionInner::Unavailable => Err(SessionFailure::new(
SessionOperation::PollPeer,
SessionTransactionState::NotEstablished,
SessionError::BackendUnavailable,
)),
}
}
pub fn wait_for_exit(
&mut self,
deadline: AbsoluteDeadline,
) -> Result<PeerStatus, SessionFailure> {
#[cfg(not(any(target_os = "linux", target_os = "macos", target_os = "windows")))]
let _ = deadline;
match &mut self.inner {
#[cfg(target_os = "linux")]
SessionInner::ReceiverReady(inner) => {
let result = inner.wait_for_exit(deadline);
let state = inner.state();
result.map_err(|error| {
linux_ready_failure(SessionOperation::WaitForExit, state, error)
})
}
#[cfg(target_os = "macos")]
SessionInner::ReceiverReady(inner) => {
let result = inner.wait_for_exit(deadline);
let state = inner.state();
result
.map_err(|error| mac_ready_failure(SessionOperation::WaitForExit, state, error))
}
#[cfg(target_os = "windows")]
SessionInner::ReceiverReady(inner) => {
let result = inner.wait_for_exit(deadline);
let state = inner.state();
result.map_err(|error| {
windows_ready_failure(SessionOperation::WaitForExit, state, error)
})
}
#[cfg(target_os = "linux")]
_ => unreachable!("receiver ready typestate owns its exact backend state"),
#[cfg(target_os = "macos")]
_ => unreachable!("receiver ready typestate owns its exact backend state"),
#[cfg(target_os = "windows")]
_ => unreachable!("receiver ready typestate owns its exact backend state"),
#[cfg(not(any(target_os = "linux", target_os = "macos", target_os = "windows")))]
SessionInner::Unavailable => Err(SessionFailure::new(
SessionOperation::WaitForExit,
SessionTransactionState::NotEstablished,
SessionError::BackendUnavailable,
)),
}
}
pub fn try_close(mut self) -> ReceiverCloseOutcome {
let facts = self.active_leases();
if !facts.is_empty() {
return ReceiverCloseOutcome::ActiveLeases {
session: self,
facts,
};
}
let close: Result<(), SessionFailure> = match &mut self.inner {
#[cfg(target_os = "linux")]
SessionInner::ReceiverReady(inner) => {
let result = inner.close_resources();
let state = inner.state();
result.map_err(|error| linux_ready_failure(SessionOperation::Close, state, error))
}
#[cfg(target_os = "macos")]
SessionInner::ReceiverReady(inner) => {
let result = inner.close_resources();
let state = inner.state();
result.map_err(|error| mac_ready_failure(SessionOperation::Close, state, error))
}
#[cfg(target_os = "windows")]
SessionInner::ReceiverReady(inner) => {
let result = inner.close_resources();
let state = inner.state();
result.map_err(|error| windows_ready_failure(SessionOperation::Close, state, error))
}
#[cfg(target_os = "linux")]
_ => unreachable!("receiver ready typestate owns its exact backend state"),
#[cfg(target_os = "macos")]
_ => unreachable!("receiver ready typestate owns its exact backend state"),
#[cfg(target_os = "windows")]
_ => unreachable!("receiver ready typestate owns its exact backend state"),
#[cfg(not(any(target_os = "linux", target_os = "macos", target_os = "windows")))]
SessionInner::Unavailable => Err(SessionFailure::new(
SessionOperation::Close,
SessionTransactionState::NotEstablished,
SessionError::BackendUnavailable,
)),
};
if let Err(error) = close {
return ReceiverCloseOutcome::Failed {
session: self,
error,
};
}
ReceiverCloseOutcome::Closed
}
pub fn abort(mut self) {
match &mut self.inner {
#[cfg(target_os = "linux")]
SessionInner::ReceiverReady(inner) => inner.abort(),
#[cfg(target_os = "macos")]
SessionInner::ReceiverReady(inner) => inner.abort(),
#[cfg(target_os = "windows")]
SessionInner::ReceiverReady(inner) => inner.abort(),
#[cfg(target_os = "linux")]
_ => unreachable!("receiver ready typestate owns its exact backend state"),
#[cfg(target_os = "macos")]
_ => unreachable!("receiver ready typestate owns its exact backend state"),
#[cfg(target_os = "windows")]
_ => unreachable!("receiver ready typestate owns its exact backend state"),
#[cfg(not(any(target_os = "linux", target_os = "macos", target_os = "windows")))]
SessionInner::Unavailable => {}
}
}
pub fn receive_batch(
&mut self,
expected: ExpectedBatch,
deadline: AbsoluteDeadline,
) -> Result<ActiveRegionSet, SessionFailure> {
#[cfg(not(any(target_os = "linux", target_os = "macos", target_os = "windows")))]
let _ = (expected, deadline);
match &mut self.inner {
#[cfg(target_os = "linux")]
SessionInner::ReceiverReady(inner) => {
let result = inner.receive_batch(expected, deadline);
let state = inner.state();
result.map_err(|error| {
linux_ready_batch_failure(SessionOperation::ReceiveBatch, state, error)
})
}
#[cfg(target_os = "macos")]
SessionInner::ReceiverReady(inner) => {
let result = inner.receive_batch(expected, deadline);
let state = inner.state();
result.map_err(|error| {
mac_ready_failure(SessionOperation::ReceiveBatch, state, error)
})
}
#[cfg(target_os = "windows")]
SessionInner::ReceiverReady(inner) => {
let result = inner.receive_batch(expected, deadline);
let state = inner.state();
result.map_err(|error| {
windows_ready_failure(SessionOperation::ReceiveBatch, state, error)
})
}
#[cfg(target_os = "linux")]
_ => unreachable!("receiver ready typestate owns its exact backend state"),
#[cfg(target_os = "macos")]
_ => unreachable!("receiver ready typestate owns its exact backend state"),
#[cfg(target_os = "windows")]
_ => unreachable!("receiver ready typestate owns its exact backend state"),
#[cfg(not(any(target_os = "linux", target_os = "macos", target_os = "windows")))]
SessionInner::Unavailable => {
unreachable!("unavailable backend cannot construct a session")
}
}
}
pub fn send_control(
&mut self,
kind: u32,
payload: &[u8],
deadline: AbsoluteDeadline,
) -> Result<(), SessionFailure> {
#[cfg(not(any(target_os = "linux", target_os = "macos", target_os = "windows")))]
let _ = (kind, payload, deadline);
match &mut self.inner {
#[cfg(target_os = "linux")]
SessionInner::ReceiverReady(inner) => {
let result = inner.send_control(kind, payload, deadline);
let state = inner.state();
result.map_err(|error| {
linux_ready_failure(SessionOperation::SendControl, state, error)
})
}
#[cfg(target_os = "macos")]
SessionInner::ReceiverReady(inner) => {
let result = inner.send_control(kind, payload, deadline);
let state = inner.state();
result
.map_err(|error| mac_ready_failure(SessionOperation::SendControl, state, error))
}
#[cfg(target_os = "windows")]
SessionInner::ReceiverReady(inner) => {
let result = inner.send_control(kind, payload, deadline);
let state = inner.state();
result.map_err(|error| {
windows_ready_failure(SessionOperation::SendControl, state, error)
})
}
#[cfg(target_os = "linux")]
_ => unreachable!("receiver ready typestate owns its exact backend state"),
#[cfg(target_os = "macos")]
_ => unreachable!("receiver ready typestate owns its exact backend state"),
#[cfg(target_os = "windows")]
_ => unreachable!("receiver ready typestate owns its exact backend state"),
#[cfg(not(any(target_os = "linux", target_os = "macos", target_os = "windows")))]
SessionInner::Unavailable => {
unreachable!("unavailable backend cannot construct a session")
}
}
}
pub fn receive_control(
&mut self,
deadline: AbsoluteDeadline,
) -> Result<ControlFrame, SessionFailure> {
#[cfg(not(any(target_os = "linux", target_os = "macos", target_os = "windows")))]
let _ = deadline;
match &mut self.inner {
#[cfg(target_os = "linux")]
SessionInner::ReceiverReady(inner) => {
let result = inner.receive_control(deadline);
let state = inner.state();
result.map_err(|error| {
linux_ready_failure(SessionOperation::ReceiveControl, state, error)
})
}
#[cfg(target_os = "macos")]
SessionInner::ReceiverReady(inner) => {
let result = inner.receive_control(deadline);
let state = inner.state();
result.map_err(|error| {
mac_ready_failure(SessionOperation::ReceiveControl, state, error)
})
}
#[cfg(target_os = "windows")]
SessionInner::ReceiverReady(inner) => {
let result = inner.receive_control(deadline);
let state = inner.state();
result.map_err(|error| {
windows_ready_failure(SessionOperation::ReceiveControl, state, error)
})
}
#[cfg(target_os = "linux")]
_ => unreachable!("receiver ready typestate owns its exact backend state"),
#[cfg(target_os = "macos")]
_ => unreachable!("receiver ready typestate owns its exact backend state"),
#[cfg(target_os = "windows")]
_ => unreachable!("receiver ready typestate owns its exact backend state"),
#[cfg(not(any(target_os = "linux", target_os = "macos", target_os = "windows")))]
SessionInner::Unavailable => {
unreachable!("unavailable backend cannot construct a session")
}
}
}
}
fn validate_public_options(options: &SessionOptions) -> Result<(), SessionError> {
if options.deadline.is_expired()
|| options.application_payload.len() > options.limits.max_bootstrap_payload_bytes as usize
{
return Err(if options.deadline.is_expired() {
SessionError::DeadlineExpired
} else {
SessionError::InvalidInput
});
}
match options.executable_identity {
ExecutableIdentityPolicy::ExactOpenedFile => {}
}
options
.limits
.validate()
.map(|_| ())
.map_err(SessionError::NativeNegotiation)
}
#[cfg(any(target_os = "linux", target_os = "macos", target_os = "windows"))]
const fn decision_rejection(decision: NegotiationDecision) -> Option<NonZeroU32> {
match decision {
NegotiationDecision::Accept => None,
NegotiationDecision::Reject(reason) => Some(reason.as_nonzero()),
}
}
#[cfg(target_os = "macos")]
fn map_mac_role(role: crate::backend::macos::vnext_session::MacNegotiationRole) -> SessionEndpoint {
match role {
crate::backend::macos::vnext_session::MacNegotiationRole::Coordinator => {
SessionEndpoint::Coordinator
}
crate::backend::macos::vnext_session::MacNegotiationRole::Receiver => {
SessionEndpoint::Receiver
}
}
}
#[cfg(target_os = "macos")]
fn map_mac_coordinator_outcome(
outcome: crate::backend::macos::vnext_session::MacNegotiationOutcome<
crate::backend::macos::vnext_session::MacCoordinatorReadySession,
>,
) -> Result<NegotiationOutcome<Session<Coordinator, Ready>>, SessionFailure> {
match outcome {
crate::backend::macos::vnext_session::MacNegotiationOutcome::Accepted(inner) => {
Ok(NegotiationOutcome::Accepted(Session::from_inner(
SessionInner::CoordinatorReady(inner),
)))
}
crate::backend::macos::vnext_session::MacNegotiationOutcome::Rejected {
by,
reason,
cleanup,
} => {
let reason = RejectionReason::from_wire(reason).ok_or_else(|| {
let failure = SessionFailure::new(
SessionOperation::Negotiate,
SessionTransactionState::Poisoned,
SessionError::MalformedPeer,
);
cleanup.map_or(failure, |facts| failure.with_cleanup(facts))
})?;
Ok(NegotiationOutcome::Rejected {
by: map_mac_role(by),
reason,
cleanup,
})
}
}
}
#[cfg(target_os = "macos")]
fn map_mac_receiver_outcome(
outcome: crate::backend::macos::vnext_session::MacNegotiationOutcome<
crate::backend::macos::vnext_session::MacReceiverReadySession,
>,
) -> Result<NegotiationOutcome<Session<Receiver, Ready>>, SessionFailure> {
match outcome {
crate::backend::macos::vnext_session::MacNegotiationOutcome::Accepted(inner) => Ok(
NegotiationOutcome::Accepted(Session::from_inner(SessionInner::ReceiverReady(inner))),
),
crate::backend::macos::vnext_session::MacNegotiationOutcome::Rejected {
by,
reason,
cleanup,
} => {
let reason = RejectionReason::from_wire(reason).ok_or_else(|| {
SessionFailure::new(
SessionOperation::Negotiate,
SessionTransactionState::Poisoned,
SessionError::MalformedPeer,
)
})?;
Ok(NegotiationOutcome::Rejected {
by: map_mac_role(by),
reason,
cleanup,
})
}
}
}
#[cfg(target_os = "windows")]
fn map_windows_role(
role: crate::backend::windows::vnext_session::WindowsNegotiationRole,
) -> SessionEndpoint {
match role {
crate::backend::windows::vnext_session::WindowsNegotiationRole::Coordinator => {
SessionEndpoint::Coordinator
}
crate::backend::windows::vnext_session::WindowsNegotiationRole::Receiver => {
SessionEndpoint::Receiver
}
}
}
#[cfg(target_os = "windows")]
fn map_windows_coordinator_outcome(
outcome: crate::backend::windows::vnext_session::WindowsNegotiationOutcome<
crate::backend::windows::vnext_session::WindowsCoordinatorReadySession,
>,
) -> Result<NegotiationOutcome<Session<Coordinator, Ready>>, SessionFailure> {
match outcome {
crate::backend::windows::vnext_session::WindowsNegotiationOutcome::Accepted(inner) => {
Ok(NegotiationOutcome::Accepted(Session::from_inner(
SessionInner::CoordinatorReady(Box::new(inner)),
)))
}
crate::backend::windows::vnext_session::WindowsNegotiationOutcome::Rejected {
by,
reason,
cleanup,
} => {
let reason = RejectionReason::from_wire(reason).ok_or_else(|| {
let failure = SessionFailure::new(
SessionOperation::Negotiate,
SessionTransactionState::Poisoned,
SessionError::MalformedPeer,
);
cleanup.map_or(failure, |facts| failure.with_cleanup(facts))
})?;
Ok(NegotiationOutcome::Rejected {
by: map_windows_role(by),
reason,
cleanup,
})
}
}
}
#[cfg(target_os = "windows")]
fn map_windows_receiver_outcome(
outcome: crate::backend::windows::vnext_session::WindowsNegotiationOutcome<
crate::backend::windows::vnext_session::WindowsReceiverReadySession,
>,
) -> Result<NegotiationOutcome<Session<Receiver, Ready>>, SessionFailure> {
match outcome {
crate::backend::windows::vnext_session::WindowsNegotiationOutcome::Accepted(inner) => {
Ok(NegotiationOutcome::Accepted(Session::from_inner(
SessionInner::ReceiverReady(Box::new(inner)),
)))
}
crate::backend::windows::vnext_session::WindowsNegotiationOutcome::Rejected {
by,
reason,
cleanup,
} => {
let reason = RejectionReason::from_wire(reason).ok_or_else(|| {
SessionFailure::new(
SessionOperation::Negotiate,
SessionTransactionState::Poisoned,
SessionError::MalformedPeer,
)
})?;
Ok(NegotiationOutcome::Rejected {
by: map_windows_role(by),
reason,
cleanup,
})
}
}
}
#[cfg(target_os = "linux")]
fn map_linux_role(
role: crate::backend::linux_vnext::spawn::LinuxNegotiationRole,
) -> SessionEndpoint {
match role {
crate::backend::linux_vnext::spawn::LinuxNegotiationRole::Coordinator => {
SessionEndpoint::Coordinator
}
crate::backend::linux_vnext::spawn::LinuxNegotiationRole::Receiver => {
SessionEndpoint::Receiver
}
}
}
#[cfg(target_os = "linux")]
fn map_linux_coordinator_outcome(
outcome: crate::backend::linux_vnext::spawn::LinuxNegotiationOutcome<
crate::backend::linux_vnext::spawn::LinuxCoordinatorReadySession,
>,
) -> Result<NegotiationOutcome<Session<Coordinator, Ready>>, SessionFailure> {
match outcome {
crate::backend::linux_vnext::spawn::LinuxNegotiationOutcome::Accepted(inner) => {
Ok(NegotiationOutcome::Accepted(Session::from_inner(
SessionInner::CoordinatorReady(inner),
)))
}
crate::backend::linux_vnext::spawn::LinuxNegotiationOutcome::Rejected {
by,
reason,
cleanup,
} => {
let reason = RejectionReason::from_wire(reason).ok_or_else(|| {
let failure = SessionFailure::new(
SessionOperation::Negotiate,
SessionTransactionState::Poisoned,
SessionError::MalformedPeer,
);
cleanup.map_or(failure, |facts| failure.with_cleanup(facts))
})?;
Ok(NegotiationOutcome::Rejected {
by: map_linux_role(by),
reason,
cleanup,
})
}
}
}
#[cfg(target_os = "linux")]
fn map_linux_receiver_outcome(
outcome: crate::backend::linux_vnext::spawn::LinuxNegotiationOutcome<
crate::backend::linux_vnext::spawn::LinuxReceiverReadySession,
>,
) -> Result<NegotiationOutcome<Session<Receiver, Ready>>, SessionFailure> {
match outcome {
crate::backend::linux_vnext::spawn::LinuxNegotiationOutcome::Accepted(inner) => Ok(
NegotiationOutcome::Accepted(Session::from_inner(SessionInner::ReceiverReady(inner))),
),
crate::backend::linux_vnext::spawn::LinuxNegotiationOutcome::Rejected {
by,
reason,
cleanup,
} => {
let reason = RejectionReason::from_wire(reason).ok_or_else(|| {
SessionFailure::new(
SessionOperation::Negotiate,
SessionTransactionState::Poisoned,
SessionError::MalformedPeer,
)
})?;
Ok(NegotiationOutcome::Rejected {
by: map_linux_role(by),
reason,
cleanup,
})
}
}
}
#[cfg(target_os = "linux")]
fn linux_ready_failure(
operation: SessionOperation,
state: SessionState,
error: crate::backend::linux_vnext::spawn::LinuxPublicSessionError,
) -> SessionFailure {
let native_code = linux_public_native_code(error);
let poisoned = state == SessionState::Poisoned;
SessionFailure::new(
operation,
if poisoned {
SessionTransactionState::Poisoned
} else {
SessionTransactionState::Ready
},
error.into(),
)
.with_native_code(native_code)
.with_poisoned(poisoned)
}
#[cfg(target_os = "linux")]
fn linux_ready_batch_failure(
operation: SessionOperation,
state: SessionState,
failure: crate::backend::linux_vnext::spawn::LinuxPublicReadyFailure,
) -> SessionFailure {
let native_code = linux_public_native_code(failure.error);
let poisoned = state == SessionState::Poisoned;
SessionFailure::new(
operation,
if failure.transaction_open_on_failure {
SessionTransactionState::TransactionOpen
} else if poisoned {
SessionTransactionState::Poisoned
} else {
SessionTransactionState::Ready
},
failure.error.into(),
)
.with_native_code(native_code)
.with_poisoned(poisoned)
}
#[cfg(target_os = "linux")]
const fn linux_public_native_code(
error: crate::backend::linux_vnext::spawn::LinuxPublicSessionError,
) -> Option<i32> {
match error {
crate::backend::linux_vnext::spawn::LinuxPublicSessionError::Native(code) => code,
crate::backend::linux_vnext::spawn::LinuxPublicSessionError::ActivationFailed(code) => code,
_ => None,
}
}
#[cfg(target_os = "linux")]
impl From<crate::backend::linux_vnext::spawn::LinuxPublicSessionError> for SessionError {
fn from(error: crate::backend::linux_vnext::spawn::LinuxPublicSessionError) -> Self {
match error {
crate::backend::linux_vnext::spawn::LinuxPublicSessionError::InvalidInput => {
Self::InvalidInput
}
crate::backend::linux_vnext::spawn::LinuxPublicSessionError::DeadlineExpired => {
Self::DeadlineExpired
}
crate::backend::linux_vnext::spawn::LinuxPublicSessionError::PeerExited => {
Self::PeerDisconnected
}
crate::backend::linux_vnext::spawn::LinuxPublicSessionError::IdentityMismatch => {
Self::IdentityMismatch
}
crate::backend::linux_vnext::spawn::LinuxPublicSessionError::MalformedPeer => {
Self::MalformedPeer
}
crate::backend::linux_vnext::spawn::LinuxPublicSessionError::Ambiguous => {
Self::Ambiguous
}
crate::backend::linux_vnext::spawn::LinuxPublicSessionError::NegotiationFailed => {
Self::NegotiationFailed
}
crate::backend::linux_vnext::spawn::LinuxPublicSessionError::NativeNegotiation(
error,
) => Self::NativeNegotiation(error),
crate::backend::linux_vnext::spawn::LinuxPublicSessionError::Control(error) => {
Self::Control(error)
}
crate::backend::linux_vnext::spawn::LinuxPublicSessionError::Batch(error) => {
Self::Batch(error)
}
crate::backend::linux_vnext::spawn::LinuxPublicSessionError::ActiveLimit => {
Self::ActiveLimit
}
crate::backend::linux_vnext::spawn::LinuxPublicSessionError::PeerPreparationFailed => {
Self::PeerPreparationFailed
}
crate::backend::linux_vnext::spawn::LinuxPublicSessionError::ActivationFailed(_) => {
Self::ActivationFailed
}
crate::backend::linux_vnext::spawn::LinuxPublicSessionError::Poisoned => Self::Poisoned,
crate::backend::linux_vnext::spawn::LinuxPublicSessionError::Native(_) => Self::Native,
}
}
}
#[cfg(target_os = "macos")]
fn mac_session_failure(
operation: SessionOperation,
transaction_state: SessionTransactionState,
error: crate::backend::macos::vnext_session::MacPublicSessionError,
poisoned: bool,
) -> SessionFailure {
SessionFailure::new(operation, transaction_state, error.into())
.with_native_code(mac_public_native_code(error))
.with_poisoned(poisoned)
}
#[cfg(target_os = "macos")]
fn mac_ready_failure(
operation: SessionOperation,
state: SessionState,
error: crate::backend::macos::vnext_session::MacPublicSessionError,
) -> SessionFailure {
let poisoned = state == SessionState::Poisoned;
mac_session_failure(
operation,
if poisoned {
SessionTransactionState::Poisoned
} else {
SessionTransactionState::Ready
},
error,
poisoned,
)
}
#[cfg(target_os = "macos")]
const fn mac_public_native_code(
error: crate::backend::macos::vnext_session::MacPublicSessionError,
) -> Option<i32> {
match error {
crate::backend::macos::vnext_session::MacPublicSessionError::Native(code) => code,
_ => None,
}
}
#[cfg(target_os = "macos")]
impl From<crate::backend::macos::vnext_session::MacPublicSessionError> for SessionError {
fn from(error: crate::backend::macos::vnext_session::MacPublicSessionError) -> Self {
use crate::backend::macos::vnext_session::MacPublicSessionError as MacError;
match error {
MacError::InvalidInput => Self::InvalidInput,
MacError::DeadlineExpired => Self::DeadlineExpired,
MacError::PeerExited => Self::PeerDisconnected,
MacError::IdentityMismatch => Self::IdentityMismatch,
MacError::MalformedPeer => Self::MalformedPeer,
MacError::Ambiguous => Self::Ambiguous,
MacError::NegotiationFailed => Self::NegotiationFailed,
MacError::NativeNegotiation(error) => Self::NativeNegotiation(error),
MacError::Control(error) => Self::Control(error),
MacError::Batch(error) => Self::Batch(error),
MacError::ActiveLimit => Self::ActiveLimit,
MacError::PeerPreparationFailed => Self::PeerPreparationFailed,
MacError::ActivationFailed => Self::ActivationFailed,
MacError::Poisoned => Self::Poisoned,
MacError::Native(_) => Self::Native,
}
}
}
#[cfg(target_os = "windows")]
fn windows_session_failure(
operation: SessionOperation,
transaction_state: SessionTransactionState,
error: crate::backend::windows::vnext_session::WindowsPublicSessionError,
poisoned: bool,
) -> SessionFailure {
let native_code = match &error {
crate::backend::windows::vnext_session::WindowsPublicSessionError::Native(code) => *code,
_ => None,
};
SessionFailure::new(operation, transaction_state, error.into())
.with_native_code(native_code)
.with_poisoned(poisoned)
}
#[cfg(target_os = "windows")]
fn windows_ready_failure(
operation: SessionOperation,
state: SessionState,
error: crate::backend::windows::vnext_session::WindowsPublicSessionError,
) -> SessionFailure {
let poisoned = state == SessionState::Poisoned;
windows_session_failure(
operation,
if poisoned {
SessionTransactionState::Poisoned
} else {
SessionTransactionState::Ready
},
error,
poisoned,
)
}
#[cfg(target_os = "windows")]
impl From<crate::backend::windows::vnext_session::WindowsPublicSessionError> for SessionError {
fn from(error: crate::backend::windows::vnext_session::WindowsPublicSessionError) -> Self {
use crate::backend::windows::vnext_session::WindowsPublicSessionError as WindowsError;
match error {
WindowsError::InvalidInput => Self::InvalidInput,
WindowsError::DeadlineExpired => Self::DeadlineExpired,
WindowsError::PeerExited => Self::PeerDisconnected,
WindowsError::IdentityMismatch => Self::IdentityMismatch,
WindowsError::MalformedPeer => Self::MalformedPeer,
WindowsError::Ambiguous => Self::Ambiguous,
WindowsError::NegotiationFailed => Self::NegotiationFailed,
WindowsError::NativeNegotiation(error) => Self::NativeNegotiation(error),
WindowsError::Control(error) => Self::Control(error),
WindowsError::Batch(error) => Self::Batch(error),
WindowsError::ActiveLimit => Self::ActiveLimit,
WindowsError::PeerPreparationFailed => Self::PeerPreparationFailed,
WindowsError::ActivationFailed => Self::ActivationFailed,
WindowsError::Poisoned => Self::Poisoned,
WindowsError::Native(_) => Self::Native,
}
}
}
const _: () = assert!(cfg!(target_has_atomic = "32"));
const _: () = assert!(cfg!(target_has_atomic = "64"));
#[cfg(test)]
#[path = "session_test.rs"]
mod tests;