use super::{
AttemptGeneration, AuthenticatedPeerBinding, CarrierCloseCause, CarrierConnectFailureCause,
CarrierConnectFailureDisposition, CarrierEstablishedCause, CarrierEvent, CarrierLostCause,
CarrierParkedCause, CarrierRefusal, CarrierStoppedCause, ConnectionKey, EpisodeId, OpaqueBytes,
};
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub enum AttemptSlot {
Idle,
InFlight { generation: AttemptGeneration },
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub enum CarrierState {
Connected {
episode: EpisodeId,
generation: AttemptGeneration,
key: ConnectionKey,
},
Parked {
episode: EpisodeId,
cause: CarrierParkedCause,
attempt: AttemptSlot,
},
Stopped {
episode: EpisodeId,
cause: CarrierStoppedCause,
},
}
#[derive(Clone, Debug, Eq, PartialEq)]
pub enum CarrierMachineInput {
EstablishedFate {
key: ConnectionKey,
cause: CarrierLostCause,
},
ProvedOnline,
ExplicitAction,
AttemptEstablished {
generation: AttemptGeneration,
key: ConnectionKey,
authenticated_peer_binding: AuthenticatedPeerBinding,
cause: CarrierEstablishedCause,
},
AttemptFailed {
generation: AttemptGeneration,
cause: CarrierConnectFailureCause,
},
Offline,
ExplicitClose {
key: Option<ConnectionKey>,
cause: CarrierCloseCause,
},
Shutdown,
FrameReceived {
key: ConnectionKey,
bytes: OpaqueBytes,
},
SendReady { key: ConnectionKey },
Refused {
key: ConnectionKey,
refusal: CarrierRefusal,
},
}
#[derive(Clone, Debug, Eq, PartialEq)]
pub enum CarrierAction {
EmitEvent(CarrierEvent),
StartAttempt { generation: AttemptGeneration },
CancelAttempt { generation: AttemptGeneration },
CloseSocket {
key: ConnectionKey,
cause: CarrierCloseCause,
},
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub struct CarrierMachine {
state: CarrierState,
next_episode: Option<u64>,
next_generation: Option<u64>,
}
impl CarrierMachine {
pub const fn new(episode: EpisodeId) -> Self {
Self {
state: CarrierState::Parked {
episode,
cause: CarrierParkedCause::AwaitingAction,
attempt: AttemptSlot::Idle,
},
next_episode: episode.get().checked_add(1),
next_generation: Some(1),
}
}
pub const fn state(&self) -> CarrierState {
self.state
}
fn take_episode(&mut self) -> Option<EpisodeId> {
let value = self.next_episode?;
self.next_episode = value.checked_add(1);
Some(EpisodeId::new(value))
}
fn take_generation(&mut self) -> Option<AttemptGeneration> {
let value = self.next_generation?;
self.next_generation = value.checked_add(1);
Some(AttemptGeneration::new(value))
}
}
#[derive(Clone, Debug, Eq, PartialEq)]
pub struct CarrierTransition {
pub machine: CarrierMachine,
pub actions: Vec<CarrierAction>,
}
pub fn step(mut machine: CarrierMachine, input: CarrierMachineInput) -> CarrierTransition {
let mut actions = Vec::new();
match machine.state {
CarrierState::Connected {
episode,
generation: _,
key,
} => step_connected(&mut machine, &mut actions, episode, key, input),
CarrierState::Parked {
episode,
cause,
attempt,
} => step_parked(&mut machine, &mut actions, episode, cause, attempt, input),
CarrierState::Stopped {
episode: _,
cause: _,
} => drop_stopped_input(&input),
}
CarrierTransition { machine, actions }
}
fn step_connected(
machine: &mut CarrierMachine,
actions: &mut Vec<CarrierAction>,
episode: EpisodeId,
current_key: ConnectionKey,
input: CarrierMachineInput,
) {
match input {
CarrierMachineInput::EstablishedFate { key, cause } => {
if key == current_key {
emit_lost(actions, key, cause);
start_successor(
machine,
actions,
episode,
CarrierParkedCause::EstablishedConnectionFate(cause),
);
}
}
CarrierMachineInput::ExplicitAction => {
emit_lost(actions, current_key, CarrierLostCause::Replaced);
actions.push(CarrierAction::CloseSocket {
key: current_key,
cause: CarrierCloseCause::Replaced,
});
start_successor(
machine,
actions,
episode,
CarrierParkedCause::AwaitingAction,
);
}
CarrierMachineInput::ExplicitClose { key, cause } => {
if key.is_none() || key == Some(current_key) {
close_connected(machine, actions, episode, current_key, cause);
}
}
CarrierMachineInput::Shutdown => {
emit_lost(actions, current_key, CarrierLostCause::Shutdown);
actions.push(CarrierAction::CloseSocket {
key: current_key,
cause: CarrierCloseCause::Shutdown,
});
stop(machine, actions, episode, CarrierStoppedCause::Shutdown);
}
CarrierMachineInput::FrameReceived { key, bytes } => {
if key == current_key {
actions.push(CarrierAction::EmitEvent(CarrierEvent::FrameReceived {
key,
bytes,
}));
}
}
CarrierMachineInput::SendReady { key } => {
if key == current_key {
actions.push(CarrierAction::EmitEvent(CarrierEvent::SendReady { key }));
}
}
CarrierMachineInput::Refused { key, refusal } => {
if key == current_key {
actions.push(CarrierAction::EmitEvent(CarrierEvent::Refused {
key,
refusal,
}));
}
}
CarrierMachineInput::ProvedOnline
| CarrierMachineInput::AttemptEstablished {
generation: _,
key: _,
authenticated_peer_binding: _,
cause: _,
}
| CarrierMachineInput::AttemptFailed {
generation: _,
cause: _,
}
| CarrierMachineInput::Offline => {}
}
}
fn step_parked(
machine: &mut CarrierMachine,
actions: &mut Vec<CarrierAction>,
episode: EpisodeId,
parked_cause: CarrierParkedCause,
attempt: AttemptSlot,
input: CarrierMachineInput,
) {
match input {
CarrierMachineInput::ProvedOnline => {
if attempt == AttemptSlot::Idle {
start_attempt(machine, actions, episode, parked_cause);
}
}
CarrierMachineInput::ExplicitAction => {
cancel_attempt(actions, attempt);
start_attempt(machine, actions, episode, parked_cause);
}
CarrierMachineInput::AttemptEstablished {
generation,
key,
authenticated_peer_binding,
cause,
} => establish_current(
machine,
actions,
episode,
attempt,
(generation, key, authenticated_peer_binding, cause),
),
CarrierMachineInput::AttemptFailed { generation, cause } => {
fail_current(machine, actions, episode, attempt, generation, cause);
}
CarrierMachineInput::Offline => {
cancel_attempt(actions, attempt);
machine.state = CarrierState::Parked {
episode,
cause: CarrierParkedCause::Offline,
attempt: AttemptSlot::Idle,
};
actions.push(CarrierAction::EmitEvent(CarrierEvent::Parked {
episode,
cause: CarrierParkedCause::Offline,
}));
}
CarrierMachineInput::ExplicitClose { key, cause } => {
if key.is_none() {
cancel_attempt(actions, attempt);
stop(
machine,
actions,
episode,
CarrierStoppedCause::ExplicitClose(cause),
);
}
}
CarrierMachineInput::Shutdown => {
cancel_attempt(actions, attempt);
stop(machine, actions, episode, CarrierStoppedCause::Shutdown);
}
CarrierMachineInput::EstablishedFate { key: _, cause: _ }
| CarrierMachineInput::FrameReceived { key: _, bytes: _ }
| CarrierMachineInput::SendReady { key: _ }
| CarrierMachineInput::Refused { key: _, refusal: _ } => {}
}
}
const fn drop_stopped_input(input: &CarrierMachineInput) {
match input {
CarrierMachineInput::EstablishedFate { key: _, cause: _ }
| CarrierMachineInput::ProvedOnline
| CarrierMachineInput::ExplicitAction
| CarrierMachineInput::AttemptEstablished {
generation: _,
key: _,
authenticated_peer_binding: _,
cause: _,
}
| CarrierMachineInput::AttemptFailed {
generation: _,
cause: _,
}
| CarrierMachineInput::Offline
| CarrierMachineInput::ExplicitClose { key: _, cause: _ }
| CarrierMachineInput::Shutdown
| CarrierMachineInput::FrameReceived { key: _, bytes: _ }
| CarrierMachineInput::SendReady { key: _ }
| CarrierMachineInput::Refused { key: _, refusal: _ } => {}
}
}
fn close_connected(
machine: &mut CarrierMachine,
actions: &mut Vec<CarrierAction>,
episode: EpisodeId,
key: ConnectionKey,
cause: CarrierCloseCause,
) {
emit_lost(actions, key, CarrierLostCause::LocalClose(cause));
actions.push(CarrierAction::CloseSocket { key, cause });
stop(
machine,
actions,
episode,
CarrierStoppedCause::ExplicitClose(cause),
);
}
fn establish_current(
machine: &mut CarrierMachine,
actions: &mut Vec<CarrierAction>,
episode: EpisodeId,
attempt: AttemptSlot,
established: (
AttemptGeneration,
ConnectionKey,
AuthenticatedPeerBinding,
CarrierEstablishedCause,
),
) {
let (generation, key, authenticated_peer_binding, cause) = established;
if attempt == (AttemptSlot::InFlight { generation }) {
machine.state = CarrierState::Connected {
episode,
generation,
key,
};
actions.push(CarrierAction::EmitEvent(CarrierEvent::Established {
episode,
generation,
key,
authenticated_peer_binding,
cause,
}));
}
}
fn fail_current(
machine: &mut CarrierMachine,
actions: &mut Vec<CarrierAction>,
episode: EpisodeId,
attempt: AttemptSlot,
generation: AttemptGeneration,
failure: CarrierConnectFailureCause,
) {
if attempt != (AttemptSlot::InFlight { generation }) {
return;
}
actions.push(CarrierAction::EmitEvent(CarrierEvent::ConnectFailed {
episode,
generation,
cause: failure,
}));
match failure.disposition() {
CarrierConnectFailureDisposition::Park => {
let cause = CarrierParkedCause::ConnectFailed(failure);
machine.state = CarrierState::Parked {
episode,
cause,
attempt: AttemptSlot::Idle,
};
actions.push(CarrierAction::EmitEvent(CarrierEvent::Parked {
episode,
cause,
}));
}
CarrierConnectFailureDisposition::Stop => stop(
machine,
actions,
episode,
CarrierStoppedCause::ConnectFailed(failure),
),
}
}
fn emit_lost(actions: &mut Vec<CarrierAction>, key: ConnectionKey, cause: CarrierLostCause) {
actions.push(CarrierAction::EmitEvent(CarrierEvent::Lost { key, cause }));
}
fn cancel_attempt(actions: &mut Vec<CarrierAction>, attempt: AttemptSlot) {
if let AttemptSlot::InFlight { generation } = attempt {
actions.push(CarrierAction::CancelAttempt { generation });
}
}
fn start_successor(
machine: &mut CarrierMachine,
actions: &mut Vec<CarrierAction>,
current_episode: EpisodeId,
cause: CarrierParkedCause,
) {
if let Some(episode) = machine.take_episode() {
start_attempt(machine, actions, episode, cause);
} else {
stop(
machine,
actions,
current_episode,
CarrierStoppedCause::IdentitySpaceExhausted,
);
}
}
fn start_attempt(
machine: &mut CarrierMachine,
actions: &mut Vec<CarrierAction>,
episode: EpisodeId,
cause: CarrierParkedCause,
) {
if let Some(generation) = machine.take_generation() {
machine.state = CarrierState::Parked {
episode,
cause,
attempt: AttemptSlot::InFlight { generation },
};
actions.push(CarrierAction::StartAttempt { generation });
} else {
stop(
machine,
actions,
episode,
CarrierStoppedCause::IdentitySpaceExhausted,
);
}
}
fn stop(
machine: &mut CarrierMachine,
actions: &mut Vec<CarrierAction>,
episode: EpisodeId,
cause: CarrierStoppedCause,
) {
machine.state = CarrierState::Stopped { episode, cause };
actions.push(CarrierAction::EmitEvent(CarrierEvent::Stopped {
episode,
cause,
}));
}