use alloc::boxed::Box;
use crate::wire::{
AckCommitted, BindingEpoch, ConversationId, DeliverySeq, Generation, ParticipantAck,
ParticipantAckResponse, ParticipantId,
};
use super::super::{
BindingRequiredLookupResult, BindingState, CumulativeAckAuthorizationError,
CumulativeAckOutcome, LiveMember, NonzeroDebtCursorEpisode, ObserverProgressProjection,
ParticipantBindingRequest, PresentedIdentity, RecipientAckObligations,
RecipientAckObligationsContextError, SealedBindingFateToken, lookup_binding_required,
membership::{LiveMemberCursorUpdate, LiveMemberCursorUpdateError},
};
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub enum NonzeroAckEpisodePosition {
Before,
Resulting,
}
#[derive(Clone, Debug, PartialEq, Eq)]
pub struct NonzeroParticipantAckCommit {
outcome: AckCommitted,
from_cursor: DeliverySeq,
binding_epoch: BindingEpoch,
before_episode: NonzeroDebtCursorEpisode,
resulting_episode: NonzeroDebtCursorEpisode,
cursor_update: LiveMemberCursorUpdate,
}
impl NonzeroParticipantAckCommit {
#[must_use]
pub const fn outcome(&self) -> &AckCommitted {
&self.outcome
}
#[must_use]
pub const fn observer_progress_projection(&self) -> ObserverProgressProjection {
let request = self.outcome.request();
ObserverProgressProjection::new(request.conversation_id, request.through_seq)
}
#[must_use]
pub const fn resulting_episode(&self) -> &NonzeroDebtCursorEpisode {
&self.resulting_episode
}
pub fn progress_binding_fate_token(
&self,
token: SealedBindingFateToken,
) -> Result<SealedBindingFateToken, Box<SealedBindingFateToken>> {
let request = self.outcome.request();
token.participant_ack_progressed(
request.conversation_id,
request.participant_id,
self.binding_epoch,
self.from_cursor,
request.through_seq,
)
}
pub fn apply_to<F>(
self,
member: &mut LiveMember<F>,
episode: &mut NonzeroDebtCursorEpisode,
) -> Result<AckCommitted, NonzeroParticipantAckCommitError> {
let request = self.outcome.request();
if member.conversation_id() != request.conversation_id {
return Err(NonzeroParticipantAckCommitError::Conversation {
expected: request.conversation_id,
actual: member.conversation_id(),
});
}
if member.participant_id() != request.participant_id {
return Err(NonzeroParticipantAckCommitError::Participant {
expected: request.participant_id,
actual: member.participant_id(),
});
}
if member.generation() != request.capability_generation {
return Err(NonzeroParticipantAckCommitError::Generation {
expected: request.capability_generation,
actual: member.generation(),
});
}
let (position, expected_cursor) = if *episode == self.before_episode {
(NonzeroAckEpisodePosition::Before, self.from_cursor)
} else if *episode == self.resulting_episode {
(NonzeroAckEpisodePosition::Resulting, request.through_seq)
} else {
return Err(NonzeroParticipantAckCommitError::EpisodePrestate);
};
if member.cursor() != expected_cursor {
return Err(NonzeroParticipantAckCommitError::AggregateCursorPrestate {
episode_position: position,
expected_cursor,
actual_cursor: member.cursor(),
});
}
member
.apply_cursor_update(self.cursor_update)
.map_err(NonzeroParticipantAckCommitError::from_member_error)?;
if position == NonzeroAckEpisodePosition::Before {
*episode = self.resulting_episode;
}
Ok(self.outcome)
}
}
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub enum NonzeroParticipantAckCommitError {
Conversation {
expected: ConversationId,
actual: ConversationId,
},
Participant {
expected: ParticipantId,
actual: ParticipantId,
},
Generation {
expected: Generation,
actual: Generation,
},
EpisodePrestate,
AggregateCursorPrestate {
episode_position: NonzeroAckEpisodePosition,
expected_cursor: DeliverySeq,
actual_cursor: DeliverySeq,
},
NonAdvancing {
from_cursor: DeliverySeq,
resulting_cursor: DeliverySeq,
},
CursorPrestate {
expected_from_cursor: DeliverySeq,
resulting_cursor: DeliverySeq,
actual_cursor: DeliverySeq,
},
}
impl NonzeroParticipantAckCommitError {
const fn from_member_error(error: LiveMemberCursorUpdateError) -> Self {
match error {
LiveMemberCursorUpdateError::Conversation { expected, actual } => {
Self::Conversation { expected, actual }
}
LiveMemberCursorUpdateError::Participant { expected, actual } => {
Self::Participant { expected, actual }
}
LiveMemberCursorUpdateError::Generation { expected, actual } => {
Self::Generation { expected, actual }
}
LiveMemberCursorUpdateError::NonAdvancing {
from_cursor,
resulting_cursor,
} => Self::NonAdvancing {
from_cursor,
resulting_cursor,
},
LiveMemberCursorUpdateError::CursorPrestate {
expected_from_cursor,
resulting_cursor,
actual_cursor,
} => Self::CursorPrestate {
expected_from_cursor,
resulting_cursor,
actual_cursor,
},
}
}
}
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub enum NonzeroParticipantAckInvariantError {
Conversation {
member: ConversationId,
episode: ConversationId,
},
ParticipantMissing {
participant_id: ParticipantId,
},
Generation {
member: Generation,
episode: Generation,
},
Cursor {
member: DeliverySeq,
episode: DeliverySeq,
},
BindingEpoch {
active: BindingEpoch,
episode: BindingEpoch,
},
EpisodeTransition(CumulativeAckAuthorizationError),
}
#[derive(Clone, Debug, PartialEq, Eq)]
pub enum NonzeroParticipantAckDecision {
Respond(ParticipantAckResponse),
Invariant(NonzeroParticipantAckInvariantError),
Commit(Box<NonzeroParticipantAckCommit>),
}
#[must_use]
pub fn apply_nonzero_participant_ack<EF, V, LF>(
presented_identity: PresentedIdentity<'_, EF, V, LF>,
binding: &BindingState,
receiving_binding_epoch: BindingEpoch,
request: &ParticipantAck,
contiguously_available_through: DeliverySeq,
episode: &NonzeroDebtCursorEpisode,
) -> NonzeroParticipantAckDecision {
apply_nonzero_participant_ack_by_availability(
presented_identity,
binding,
receiving_binding_epoch,
request,
&AckAvailability::Contiguous(contiguously_available_through),
episode,
)
}
#[must_use]
pub fn apply_nonzero_participant_ack_with_obligations<EF, V, LF>(
presented_identity: PresentedIdentity<'_, EF, V, LF>,
binding: &BindingState,
receiving_binding_epoch: BindingEpoch,
request: &ParticipantAck,
obligations: &RecipientAckObligations,
episode: &NonzeroDebtCursorEpisode,
) -> NonzeroParticipantAckDecision {
apply_nonzero_participant_ack_by_availability(
presented_identity,
binding,
receiving_binding_epoch,
request,
&AckAvailability::Obligations(obligations),
episode,
)
}
pub fn scalar_audit_for_recipient_endpoint(
obligations: &RecipientAckObligations,
participant_id: ParticipantId,
acknowledged_through: DeliverySeq,
endpoint: DeliverySeq,
scalar_audit: DeliverySeq,
) -> Result<Option<DeliverySeq>, RecipientAckObligationsContextError> {
obligations
.contains_endpoint(participant_id, acknowledged_through, endpoint)
.map(|testified| testified.then_some(scalar_audit))
}
enum AckAvailability<'a> {
Contiguous(DeliverySeq),
Obligations(&'a RecipientAckObligations),
}
fn apply_nonzero_participant_ack_by_availability<EF, V, LF>(
presented_identity: PresentedIdentity<'_, EF, V, LF>,
binding: &BindingState,
receiving_binding_epoch: BindingEpoch,
request: &ParticipantAck,
availability: &AckAvailability<'_>,
episode: &NonzeroDebtCursorEpisode,
) -> NonzeroParticipantAckDecision {
let lookup_request = ParticipantBindingRequest::ParticipantAck(request.clone());
let (member, active_binding) = match lookup_binding_required(
presented_identity,
binding,
Some(receiving_binding_epoch),
&lookup_request,
) {
BindingRequiredLookupResult::Retired(outcome) => {
return NonzeroParticipantAckDecision::Respond(ParticipantAckResponse::from_retired(
outcome,
));
}
BindingRequiredLookupResult::ParticipantUnknown(outcome) => {
return NonzeroParticipantAckDecision::Respond(
ParticipantAckResponse::from_participant_unknown(outcome),
);
}
BindingRequiredLookupResult::StaleAuthority(outcome) => {
return NonzeroParticipantAckDecision::Respond(
ParticipantAckResponse::from_stale_authority(outcome),
);
}
BindingRequiredLookupResult::NoBinding(outcome) => {
return NonzeroParticipantAckDecision::Respond(
ParticipantAckResponse::from_no_binding(outcome),
);
}
BindingRequiredLookupResult::Authorized { member, binding } => (member, binding),
};
if let Err(error) = validate_aggregate(member, active_binding.binding_epoch, episode) {
return NonzeroParticipantAckDecision::Invariant(error);
}
let mut resulting_episode = episode.clone();
let selected = match availability {
AckAvailability::Contiguous(available_through) => resulting_episode.acknowledge(
member.participant_id(),
active_binding.binding_epoch,
request,
*available_through,
),
AckAvailability::Obligations(obligations) => resulting_episode
.acknowledge_with_obligations(
member.participant_id(),
active_binding.binding_epoch,
request,
obligations,
),
};
let outcome = match selected {
Ok(outcome) => outcome,
Err(error) => {
return NonzeroParticipantAckDecision::Invariant(
NonzeroParticipantAckInvariantError::EpisodeTransition(error),
);
}
};
match outcome {
CumulativeAckOutcome::Committed(outcome) => {
NonzeroParticipantAckDecision::Commit(Box::new(NonzeroParticipantAckCommit {
cursor_update: LiveMemberCursorUpdate::new(
request.conversation_id,
request.participant_id,
request.capability_generation,
member.cursor(),
request.through_seq,
),
outcome,
binding_epoch: active_binding.binding_epoch,
from_cursor: member.cursor(),
before_episode: episode.clone(),
resulting_episode,
}))
}
CumulativeAckOutcome::NoOp(outcome) => {
NonzeroParticipantAckDecision::Respond(ParticipantAckResponse::from_ack_no_op(outcome))
}
CumulativeAckOutcome::Gap(outcome) => {
NonzeroParticipantAckDecision::Respond(ParticipantAckResponse::ack_gap(outcome))
}
CumulativeAckOutcome::Regression(outcome) => {
NonzeroParticipantAckDecision::Respond(ParticipantAckResponse::ack_regression(outcome))
}
}
}
fn validate_aggregate<F>(
member: &LiveMember<F>,
active_binding_epoch: BindingEpoch,
episode: &NonzeroDebtCursorEpisode,
) -> Result<(), NonzeroParticipantAckInvariantError> {
if member.conversation_id() != episode.conversation_id() {
return Err(NonzeroParticipantAckInvariantError::Conversation {
member: member.conversation_id(),
episode: episode.conversation_id(),
});
}
let Some(episode_participant) = episode.participant(member.participant_id()) else {
return Err(NonzeroParticipantAckInvariantError::ParticipantMissing {
participant_id: member.participant_id(),
});
};
let episode_generation = episode_participant
.active_binding_epoch()
.capability_generation;
if episode_generation != member.generation() {
return Err(NonzeroParticipantAckInvariantError::Generation {
member: member.generation(),
episode: episode_generation,
});
}
if episode_participant.cursor() != member.cursor() {
return Err(NonzeroParticipantAckInvariantError::Cursor {
member: member.cursor(),
episode: episode_participant.cursor(),
});
}
if episode_participant.active_binding_epoch() != active_binding_epoch {
return Err(NonzeroParticipantAckInvariantError::BindingEpoch {
active: active_binding_epoch,
episode: episode_participant.active_binding_epoch(),
});
}
Ok(())
}