use alloc::{boxed::Box, vec, vec::Vec};
use crate::{algebra::ResourceVector, wire::RecordAdmission};
use super::super::{
AttachCommit, AttachTransition, BindingState, ClaimFrontiers, ClosureAccounting, ClosureDebt,
CommittedDetachTransition, DebtCompletion, DetachCell, DetachedCredentialRecovery, Event,
FencedAttachCommit, FrontierBinding, FrontierParticipant, IdentityState,
InitialEnrollmentFrontierCommit, LeaveCommitError, LeaveCommitParameters, LiveMember,
MarkerAckCommit, NonzeroParticipantAckCommit, ObserverProgressProjection, OrderLedger,
ParticipantAckCommit, PendingFinalization, PendingLeaveCommitParameters,
PrepareLeaveAuthorityError, RetainedCausalRecord, RetainedCausalRecordKind, SequenceLedger,
StoredEdge, VerifiedLeaveRequest,
claim_frontier::{BindingTerminalOwner, FencedMarkerSourceRecord, LiveFrontierTransitionError},
commit_leave, commit_pending_leave,
};
use super::{
InitialEnrollmentOperationCommit, MarkerDeliveryProjection, MarkerDrainCommit,
RecordAdmissionPersistenceParts, RetainedRecordCharge, UnchangedRecordAdmission,
};
mod binding_fate_transition;
mod ledger;
mod state;
pub(super) use binding_fate_transition::BindingFateOwnerPlan;
use ledger::{
detach_order, detach_sequence, detached_attach_order, detached_attach_sequence,
enrollment_order, enrollment_sequence, superseding_attach_order, superseding_attach_sequence,
};
use state::{
accounting_after_fenced_attach, accounting_after_leave, accounting_after_marker_ack,
accounting_after_rows, retained_attached, retained_terminal,
};
#[derive(Debug, PartialEq, Eq)]
pub struct LiveFrontierOwner {
frontiers: ClaimFrontiers,
closure_accounting: ClosureAccounting,
retained_charges: Vec<RetainedRecordCharge>,
retained_record_limit: u64,
}
impl LiveFrontierOwner {
#[must_use]
pub fn from_initial_enrollment<F>(
initial: InitialEnrollmentFrontierCommit<F>,
retained_record_limit: u64,
) -> (InitialEnrollmentOperationCommit<F>, Self) {
let (operation, frontiers, closure_accounting, attached_charge) =
initial.into_conversation_parts();
let attached = operation.enrollment().attached;
let retained_charges = vec![RetainedRecordCharge::new(
attached.delivery_seq(),
attached.admission_order(),
attached_charge,
)];
(
operation,
Self {
frontiers,
closure_accounting,
retained_charges,
retained_record_limit,
},
)
}
#[cfg(any(test, feature = "test-support"))]
pub(in crate::lifecycle) const fn from_test_parts(
frontiers: ClaimFrontiers,
closure_accounting: ClosureAccounting,
retained_charges: Vec<RetainedRecordCharge>,
retained_record_limit: u64,
) -> Self {
Self {
frontiers,
closure_accounting,
retained_charges,
retained_record_limit,
}
}
#[cfg(any(test, feature = "test-support"))]
pub fn with_pending_finalizer_test_capacity(
mut self,
finalizer_rows: u64,
finalizer_charge: ResourceVector,
) -> Result<Self, &'static str> {
let retained_count = u64::try_from(self.retained_charges.len())
.map_err(|_| "retained record count exceeds u64")?;
self.retained_record_limit = retained_count
.checked_add(finalizer_rows)
.ok_or("pending finalizer retained-record capacity overflow")?;
let current = self.closure_accounting;
let configured = current.configured_cap();
let configured = ResourceVector::new(
configured
.entries
.checked_add(finalizer_charge.entries)
.ok_or("pending finalizer entry capacity overflow")?,
configured
.bytes
.checked_add(finalizer_charge.bytes)
.ok_or("pending finalizer byte capacity overflow")?,
);
self.closure_accounting = ClosureAccounting::try_new(
current.state(),
current.marker_capacity_credits(),
current.marker_anchors(),
current.edge_sequence_claims(),
current.edge_order_position_claims(),
current.edge_k_remaining(),
current.baseline(),
configured,
current.episode_churn_used(),
current.episode_churn_limit(),
)
.map_err(|_| "pending finalizer closure capacity extension refused")?;
Ok(self)
}
#[must_use]
pub const fn frontiers(&self) -> &ClaimFrontiers {
&self.frontiers
}
#[must_use]
pub const fn closure_accounting(&self) -> ClosureAccounting {
self.closure_accounting
}
#[must_use]
pub fn retained_charges(&self) -> &[RetainedRecordCharge] {
&self.retained_charges
}
#[must_use]
pub const fn retained_record_limit(&self) -> u64 {
self.retained_record_limit
}
#[must_use]
pub fn into_parts(
self,
) -> (
ClaimFrontiers,
ClosureAccounting,
Vec<RetainedRecordCharge>,
u64,
) {
(
self.frontiers,
self.closure_accounting,
self.retained_charges,
self.retained_record_limit,
)
}
#[must_use]
pub fn from_unchanged_record_admission<EF, V, LF>(
unchanged: UnchangedRecordAdmission<'_, EF, V, LF>,
retained_record_limit: u64,
) -> (Self, RecordAdmission, ResourceVector) {
let (prestate, encoded_record_charge) = unchanged.into_parts();
let (request, frontiers, closure_accounting, retained_charges) =
prestate.into_live_owner_parts();
(
Self {
frontiers,
closure_accounting,
retained_charges,
retained_record_limit,
},
request,
encoded_record_charge,
)
}
#[must_use]
pub fn from_record_admission_persistence(
persistence: RecordAdmissionPersistenceParts,
retained_record_limit: u64,
) -> Self {
Self {
frontiers: persistence.frontiers,
closure_accounting: persistence.accounting,
retained_charges: persistence.retained_charges,
retained_record_limit,
}
}
#[must_use]
pub fn from_marker_drain(
commit: MarkerDrainCommit,
retained_record_limit: u64,
) -> (Self, StoredEdge, MarkerDeliveryProjection) {
let (frontiers, closure_accounting, retained_charges, successor, projection) =
commit.into_parts();
(
Self {
frontiers,
closure_accounting,
retained_charges,
retained_record_limit,
},
successor,
projection,
)
}
pub(super) fn commit_binding_terminal_candidate(
self,
active_binding: super::super::ActiveBinding,
admission_order: super::super::AdmissionOrder,
delivery_seq: crate::wire::DeliverySeq,
charge: RetainedRecordCharge,
) -> Result<Self, Box<(Self, LiveFrontierError)>> {
let mut active = self.frontiers.active_identities().participants().to_vec();
let Some(participant) = active
.iter_mut()
.find(|participant| participant.participant_index() == active_binding.participant_id)
else {
return Err(Box::new((self, LiveFrontierError::Authority)));
};
if participant.binding() != FrontierBinding::Bound(active_binding.binding_epoch) {
return Err(Box::new((self, LiveFrontierError::Authority)));
}
*participant = FrontierParticipant::new(
participant.participant_index(),
participant.cursor(),
FrontierBinding::Detached(active_binding.binding_epoch),
);
let row = RetainedCausalRecord {
delivery_seq,
admission_order,
kind: RetainedCausalRecordKind::BindingTerminal(super::super::BindingTerminalOwner {
participant_index: active_binding.participant_id,
binding_epoch: active_binding.binding_epoch,
}),
};
let Some(sequence) = detach_sequence(self.frontiers.sequence().ledger(), delivery_seq)
else {
return Err(Box::new((self, LiveFrontierError::Frontier)));
};
let Some(order) = detach_order(
self.frontiers.order().ledger(),
admission_order.transaction_order(),
) else {
return Err(Box::new((self, LiveFrontierError::Frontier)));
};
match transition(self, (), active, &[row], vec![charge], sequence, order) {
Ok(committed) => {
let ((), owner) = committed.into_parts();
Ok(owner)
}
Err(failure) => {
let error = failure.error();
let ((), owner) = failure.into_parts();
Err(Box::new((owner, error)))
}
}
}
pub(super) fn pend_binding_terminal_candidate(
self,
active_binding: super::super::ActiveBinding,
admission_order: super::super::AdmissionOrder,
delivery_seq: crate::wire::DeliverySeq,
) -> Result<Self, Box<(Self, LiveFrontierError)>> {
let Some(order) = detach_order(
self.frontiers.order().ledger(),
admission_order.transaction_order(),
) else {
return Err(Box::new((self, LiveFrontierError::Frontier)));
};
let Self {
frontiers,
closure_accounting,
retained_charges,
retained_record_limit,
} = self;
match frontiers.apply_pending_binding_terminal(
active_binding.participant_id,
active_binding.binding_epoch,
delivery_seq,
admission_order,
order,
) {
Ok(frontiers) => Ok(Self {
frontiers,
closure_accounting,
retained_charges,
retained_record_limit,
}),
Err(failure) => {
let (frontiers, error) = *failure;
Err(Box::new((
Self {
frontiers,
closure_accounting,
retained_charges,
retained_record_limit,
},
map_frontier_error(error),
)))
}
}
}
pub fn drain_pending_terminal(
self,
pending: PendingFinalization,
terminal_charge: RetainedRecordCharge,
) -> Result<DrainedPendingTerminal, Box<PendingTerminalDrainRefused>> {
let Some(expected_sequence) = self
.frontiers
.sequence()
.ledger()
.high_watermark()
.checked_add(1)
else {
return drain_refusal(self, LiveFrontierError::Frontier);
};
if terminal_charge.delivery_seq() != expected_sequence
|| terminal_charge.admission_order() != pending.admission_order()
|| terminal_charge.encoded_charge().entries != 1
{
return drain_refusal(self, LiveFrontierError::RetainedCharge);
}
if self
.retained_charges
.len()
.checked_add(1)
.and_then(|len| u64::try_from(len).ok())
.is_none()
{
return drain_refusal(self, LiveFrontierError::RetainedRecordLimit);
}
let Some(accounting) = accounting_after_rows(self.closure_accounting, &[terminal_charge])
else {
return drain_refusal(self, LiveFrontierError::ClosureAccounting);
};
let Self {
frontiers,
closure_accounting,
mut retained_charges,
retained_record_limit,
} = self;
let expected_owner = BindingTerminalOwner {
participant_index: pending.participant_id(),
binding_epoch: pending.binding_epoch(),
};
let (frontiers, record) = match frontiers
.drain_first_binding_terminal(expected_owner, pending.admission_order())
{
Ok(drained) => drained,
Err(failure) => {
let (frontiers, error) = *failure;
return drain_refusal(
Self {
frontiers,
closure_accounting,
retained_charges,
retained_record_limit,
},
map_frontier_error(error),
);
}
};
retained_charges.push(terminal_charge);
let projection =
ObserverProgressProjection::new(pending.conversation_id(), record.delivery_seq);
Ok(DrainedPendingTerminal {
owner: Self {
frontiers,
closure_accounting: accounting,
retained_charges,
retained_record_limit,
},
projection,
})
}
pub fn retain_fenced_marker_source(
self,
recovery: DetachedCredentialRecovery,
) -> Result<RetainedFencedMarkerSource, Box<FencedMarkerSourceRetentionRefused>> {
let Some(source) = self.frontiers.fenced_marker_source(recovery) else {
return Err(Box::new(FencedMarkerSourceRetentionRefused {
owner: self,
recovery,
}));
};
Ok(RetainedFencedMarkerSource {
owner: self,
recovery,
expectation: FencedMarkerSourceExpectation { source },
})
}
#[must_use]
pub fn mint_fenced_attach(
mut self,
marker_source_sequence: u64,
recovery: DetachedCredentialRecovery,
debt: ClosureDebt,
event: Event,
successor: DebtCompletion,
) -> MintFencedAttachResult {
let Some(record) = self.frontiers.take_fenced_marker_record(recovery) else {
return MintFencedAttachResult::MintRefused(Box::new(MintFencedAttachRefused {
owner: self,
marker_source_sequence,
recovery,
debt,
event,
successor,
reason: FencedAttachMintRefusalReason::MarkerAuthority,
}));
};
match recovery.fenced_attach(record, debt, event, successor) {
Ok(proof) => MintFencedAttachResult::Minted(Box::new(MintedFencedAttach {
owner_without_marker_authority: self,
proof,
})),
Err(refusal) => {
self.frontiers
.reinstall_fenced_marker_record((*refusal).into_record());
MintFencedAttachResult::MintRefused(Box::new(MintFencedAttachRefused {
owner: self,
marker_source_sequence,
recovery,
debt,
event,
successor,
reason: FencedAttachMintRefusalReason::ProofInputs,
}))
}
}
}
}
#[derive(Debug, PartialEq, Eq)]
pub struct DrainedPendingTerminal {
owner: LiveFrontierOwner,
projection: ObserverProgressProjection,
}
impl DrainedPendingTerminal {
#[must_use]
pub fn into_parts(self) -> (LiveFrontierOwner, ObserverProgressProjection) {
(self.owner, self.projection)
}
}
#[derive(Debug, PartialEq, Eq)]
pub struct PendingTerminalDrainRefused {
owner: LiveFrontierOwner,
error: LiveFrontierError,
}
impl PendingTerminalDrainRefused {
#[must_use]
pub const fn error(&self) -> LiveFrontierError {
self.error
}
#[must_use]
pub fn into_owner(self) -> LiveFrontierOwner {
self.owner
}
}
fn drain_refusal(
owner: LiveFrontierOwner,
error: LiveFrontierError,
) -> Result<DrainedPendingTerminal, Box<PendingTerminalDrainRefused>> {
Err(Box::new(PendingTerminalDrainRefused { owner, error }))
}
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub struct FencedMarkerSourceExpectation {
source: FencedMarkerSourceRecord,
}
impl FencedMarkerSourceExpectation {
#[must_use]
pub const fn conversation_id(self) -> u64 {
self.source.conversation_id
}
#[must_use]
pub const fn marker_delivery_seq(self) -> u64 {
self.source.delivery_seq
}
#[must_use]
pub const fn admission_order(self) -> super::super::AdmissionOrder {
self.source.admission_order
}
#[must_use]
pub const fn participant_id(self) -> u64 {
self.source.participant_id
}
#[must_use]
pub const fn provenance(self) -> super::super::MarkerProvenance {
self.source.provenance
}
#[must_use]
pub const fn target_binding(self) -> FrontierBinding {
self.source.target_binding
}
}
#[derive(Debug, PartialEq, Eq)]
pub struct RetainedFencedMarkerSource {
owner: LiveFrontierOwner,
recovery: DetachedCredentialRecovery,
expectation: FencedMarkerSourceExpectation,
}
impl RetainedFencedMarkerSource {
#[must_use]
pub const fn expectation(&self) -> FencedMarkerSourceExpectation {
self.expectation
}
#[must_use]
pub fn into_parts(self) -> (LiveFrontierOwner, DetachedCredentialRecovery) {
(self.owner, self.recovery)
}
}
#[derive(Debug, PartialEq, Eq)]
pub struct FencedMarkerSourceRetentionRefused {
owner: LiveFrontierOwner,
recovery: DetachedCredentialRecovery,
}
impl FencedMarkerSourceRetentionRefused {
#[must_use]
pub fn into_parts(self) -> (LiveFrontierOwner, DetachedCredentialRecovery) {
(self.owner, self.recovery)
}
}
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub enum FencedAttachMintRefusalReason {
MarkerAuthority,
ProofInputs,
}
#[derive(Debug, PartialEq, Eq)]
pub struct MintedFencedAttach {
owner_without_marker_authority: LiveFrontierOwner,
proof: FencedAttachCommit,
}
impl MintedFencedAttach {
#[must_use]
pub fn into_parts(self) -> (LiveFrontierOwner, FencedAttachCommit) {
(self.owner_without_marker_authority, self.proof)
}
}
#[derive(Debug, PartialEq, Eq)]
pub struct MintFencedAttachRefused {
owner: LiveFrontierOwner,
marker_source_sequence: u64,
recovery: DetachedCredentialRecovery,
debt: ClosureDebt,
event: Event,
successor: DebtCompletion,
reason: FencedAttachMintRefusalReason,
}
impl MintFencedAttachRefused {
#[must_use]
pub const fn reason(&self) -> FencedAttachMintRefusalReason {
self.reason
}
#[must_use]
pub fn into_parts(
self,
) -> (
LiveFrontierOwner,
u64,
DetachedCredentialRecovery,
ClosureDebt,
Event,
DebtCompletion,
) {
(
self.owner,
self.marker_source_sequence,
self.recovery,
self.debt,
self.event,
self.successor,
)
}
}
#[derive(Debug, PartialEq, Eq)]
pub enum MintFencedAttachResult {
Minted(Box<MintedFencedAttach>),
MintRefused(Box<MintFencedAttachRefused>),
}
#[derive(Debug, PartialEq, Eq)]
pub struct LiveLeaveCommit<EF, V, LF> {
identity: IdentityState<EF, V, LF>,
owner: LiveFrontierOwner,
}
impl<EF, V, LF> LiveLeaveCommit<EF, V, LF> {
#[must_use]
pub const fn observer_progress_projection(&self) -> Option<ObserverProgressProjection> {
let IdentityState::Retired(retired) = &self.identity else {
return None;
};
let committed = retired.committed_result();
Some(ObserverProgressProjection::new(
committed.conversation_id(),
committed.left_delivery_seq(),
))
}
#[must_use]
pub fn into_parts(self) -> (IdentityState<EF, V, LF>, LiveFrontierOwner) {
(self.identity, self.owner)
}
}
#[derive(Clone, Debug, PartialEq, Eq)]
pub enum LiveLeaveError {
Prepare(PrepareLeaveAuthorityError),
Commit(LeaveCommitError),
RetainedCharge,
RetainedRecordLimit,
ClosureAccounting,
Identity,
}
pub fn commit_settled_leave_frontier<EF, V, LF, D>(
owner: LiveFrontierOwner,
member: LiveMember<EF>,
binding: BindingState,
detach_cell: DetachCell<D>,
verified: VerifiedLeaveRequest<V, LF>,
left_delivery_seq: u64,
left_charge: RetainedRecordCharge,
) -> Result<LiveLeaveCommit<EF, V, LF>, LiveLeaveError> {
let LiveFrontierOwner {
frontiers,
closure_accounting,
mut retained_charges,
retained_record_limit,
} = owner;
let retired_marker_charge =
retired_marker_charge(&frontiers, &retained_charges, member.participant_id())?;
let authority = frontiers
.prepare_settled_leave_authority(&member, binding)
.map_err(LiveLeaveError::Prepare)?;
let commit = commit_leave(
member,
binding,
detach_cell,
verified,
authority,
LeaveCommitParameters { left_delivery_seq },
)
.map_err(LiveLeaveError::Commit)?;
let (identity, frontiers) = commit.into_parts();
let IdentityState::Retired(retired) = &identity else {
return Err(LiveLeaveError::Identity);
};
if left_charge.delivery_seq() != retired.committed_result().left_delivery_seq()
|| left_charge.admission_order() != retired.left_admission_order()
|| left_charge.encoded_charge().entries != 1
{
return Err(LiveLeaveError::RetainedCharge);
}
retained_charges.push(left_charge);
retained_charges.sort_unstable_by_key(|charge| charge.delivery_seq());
let retained_len = u64::try_from(frontiers.retained_records().len())
.map_err(|_| LiveLeaveError::RetainedRecordLimit)?;
if retained_len > retained_record_limit
|| retained_charges.len() != frontiers.retained_records().len()
{
return Err(LiveLeaveError::RetainedRecordLimit);
}
let closure_accounting =
accounting_after_leave(closure_accounting, &[left_charge], retired_marker_charge)
.ok_or(LiveLeaveError::ClosureAccounting)?;
Ok(LiveLeaveCommit {
identity,
owner: LiveFrontierOwner {
frontiers,
closure_accounting,
retained_charges,
retained_record_limit,
},
})
}
pub fn commit_pending_leave_frontier<EF, V, LF, D>(
owner: LiveFrontierOwner,
member: LiveMember<EF>,
pending: PendingFinalization,
detach_cell: DetachCell<D>,
verified: VerifiedLeaveRequest<V, LF>,
parameters: PendingLeaveCommitParameters,
charges: [RetainedRecordCharge; 2],
) -> Result<LiveLeaveCommit<EF, V, LF>, LiveLeaveError> {
let [terminal_charge, left_charge] = charges;
let terminal_delivery_seq = parameters.terminal_delivery_seq;
let LiveFrontierOwner {
frontiers,
closure_accounting,
mut retained_charges,
retained_record_limit,
} = owner;
let retired_marker_charge =
retired_marker_charge(&frontiers, &retained_charges, member.participant_id())?;
let authority = frontiers
.prepare_pending_leave_authority(&member, pending)
.map_err(LiveLeaveError::Prepare)?;
let commit = commit_pending_leave(
member,
pending,
detach_cell,
verified,
authority,
parameters,
)
.map_err(LiveLeaveError::Commit)?;
let (identity, frontiers) = commit.into_parts();
let IdentityState::Retired(retired) = &identity else {
return Err(LiveLeaveError::Identity);
};
if retired.committed_result().prior_terminal_delivery_seq() != Some(terminal_delivery_seq)
|| terminal_charge.delivery_seq() != terminal_delivery_seq
|| terminal_charge.admission_order() != pending.admission_order()
|| terminal_charge.encoded_charge().entries != 1
|| left_charge.delivery_seq() != retired.committed_result().left_delivery_seq()
|| left_charge.admission_order() != retired.left_admission_order()
|| left_charge.encoded_charge().entries != 1
{
return Err(LiveLeaveError::RetainedCharge);
}
retained_charges.extend([terminal_charge, left_charge]);
retained_charges.sort_unstable_by_key(|charge| charge.delivery_seq());
let retained_len = u64::try_from(frontiers.retained_records().len())
.map_err(|_| LiveLeaveError::RetainedRecordLimit)?;
if retained_len > retained_record_limit {
return Err(LiveLeaveError::RetainedRecordLimit);
}
if retained_charges.len() != frontiers.retained_records().len() {
return Err(LiveLeaveError::RetainedCharge);
}
let closure_accounting = accounting_after_leave(
closure_accounting,
&[terminal_charge, left_charge],
retired_marker_charge,
)
.ok_or(LiveLeaveError::ClosureAccounting)?;
Ok(LiveLeaveCommit {
identity,
owner: LiveFrontierOwner {
frontiers,
closure_accounting,
retained_charges,
retained_record_limit,
},
})
}
fn retired_marker_charge(
frontiers: &ClaimFrontiers,
retained_charges: &[RetainedRecordCharge],
participant_id: crate::wire::ParticipantId,
) -> Result<Option<RetainedRecordCharge>, LiveLeaveError> {
let marker_sequence = frontiers
.retained_marker_records()
.iter()
.find_map(|record| {
matches!(
record.kind,
RetainedCausalRecordKind::CompactionMarker {
participant_index,
..
} if participant_index == participant_id
)
.then_some(record.delivery_seq)
});
let Some(marker_sequence) = marker_sequence else {
return Ok(None);
};
retained_charges
.iter()
.copied()
.find(|charge| charge.delivery_seq() == marker_sequence)
.map(Some)
.ok_or(LiveLeaveError::RetainedCharge)
}
#[derive(Debug, PartialEq, Eq)]
pub struct AttachFrontierCharges {
terminal: Option<RetainedRecordCharge>,
attached: RetainedRecordCharge,
seal: LiveTransitionInputSeal,
}
#[derive(Debug, PartialEq, Eq)]
enum LiveTransitionInputSeal {
Validated,
}
impl AttachFrontierCharges {
#[must_use]
pub const fn new(
terminal: Option<RetainedRecordCharge>,
attached: RetainedRecordCharge,
) -> Self {
Self {
terminal,
attached,
seal: LiveTransitionInputSeal::Validated,
}
}
const fn into_parts(self) -> (Option<RetainedRecordCharge>, RetainedRecordCharge) {
let Self {
terminal,
attached,
seal,
} = self;
match seal {
LiveTransitionInputSeal::Validated => (terminal, attached),
}
}
}
#[derive(Debug, PartialEq, Eq)]
pub struct LiveFrontierCommit<T> {
operation: T,
owner: LiveFrontierOwner,
}
impl<T> LiveFrontierCommit<T> {
#[must_use]
pub const fn operation(&self) -> &T {
&self.operation
}
#[must_use]
pub const fn owner(&self) -> &LiveFrontierOwner {
&self.owner
}
#[must_use]
pub fn into_parts(self) -> (T, LiveFrontierOwner) {
(self.operation, self.owner)
}
}
#[derive(Debug, PartialEq, Eq)]
pub struct LiveFrontierFailure<T> {
error: LiveFrontierError,
operation: T,
owner: LiveFrontierOwner,
}
impl<T> LiveFrontierFailure<T> {
#[must_use]
pub const fn error(&self) -> LiveFrontierError {
self.error
}
#[must_use]
pub fn into_parts(self) -> (T, LiveFrontierOwner) {
(self.operation, self.owner)
}
}
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub enum LiveFrontierError {
Authority,
Precedence,
RetainedCharge,
RetainedRecordLimit,
Frontier,
ClosureAccounting,
}
pub type LiveFrontierResult<T> = Result<LiveFrontierCommit<T>, Box<LiveFrontierFailure<T>>>;
pub fn apply_enrollment_frontier<F>(
owner: LiveFrontierOwner,
operation: super::super::EnrollmentCommit<F>,
attached_charge: RetainedRecordCharge,
) -> LiveFrontierResult<super::super::EnrollmentCommit<F>> {
let attached = operation.attached;
if attached.conversation_id() != owner.frontiers.conversation_id() {
return failure(owner, operation, LiveFrontierError::Authority);
}
let participant_id = attached.participant_id();
let mut active = owner.frontiers.active_identities().participants().to_vec();
if active
.iter()
.any(|participant| participant.participant_index() == participant_id)
{
return failure(owner, operation, LiveFrontierError::Authority);
}
active.push(FrontierParticipant::new(
participant_id,
operation.member.cursor(),
FrontierBinding::Bound(attached.binding_epoch()),
));
active.sort_unstable_by_key(|participant| participant.participant_index());
let rows = [retained_attached(attached)];
let Some(sequence) =
enrollment_sequence(owner.frontiers.sequence().ledger(), attached.delivery_seq())
else {
return failure(owner, operation, LiveFrontierError::Frontier);
};
let Some(order) = enrollment_order(
owner.frontiers.order().ledger(),
attached.admission_order().transaction_order(),
) else {
return failure(owner, operation, LiveFrontierError::Frontier);
};
transition(
owner,
operation,
active,
&rows,
vec![attached_charge],
sequence,
order,
)
}
pub fn apply_attach_frontier<F, V>(
owner: LiveFrontierOwner,
operation: AttachCommit<F, V>,
charges: AttachFrontierCharges,
) -> LiveFrontierResult<AttachCommit<F, V>> {
let (terminal_charge, attached_charge) = charges.into_parts();
let attached = operation.attached;
if attached.conversation_id() != owner.frontiers.conversation_id() {
return failure(owner, operation, LiveFrontierError::Authority);
}
let mut active = owner.frontiers.active_identities().participants().to_vec();
let Some(participant) = active
.iter_mut()
.find(|participant| participant.participant_index() == attached.participant_id())
else {
return failure(owner, operation, LiveFrontierError::Authority);
};
*participant = FrontierParticipant::new(
participant.participant_index(),
operation.member.cursor(),
FrontierBinding::Bound(attached.binding_epoch()),
);
let current_sequence = owner.frontiers.sequence().ledger();
let current_order = owner.frontiers.order().ledger();
let (rows, keyed_charges, sequence, order) = match operation.transition {
AttachTransition::Detached => {
if terminal_charge.is_some() {
return failure(owner, operation, LiveFrontierError::RetainedCharge);
}
let Some(sequence) =
detached_attach_sequence(current_sequence, attached.delivery_seq())
else {
return failure(owner, operation, LiveFrontierError::Frontier);
};
let Some(order) = detached_attach_order(
current_order,
attached.admission_order().transaction_order(),
) else {
return failure(owner, operation, LiveFrontierError::Frontier);
};
(
vec![retained_attached(attached)],
vec![attached_charge],
sequence,
order,
)
}
AttachTransition::Superseded { terminal } => {
let Some(terminal_charge) = terminal_charge else {
return failure(owner, operation, LiveFrontierError::RetainedCharge);
};
let rows = vec![
retained_terminal(terminal.into()),
retained_attached(attached),
];
let Some(sequence) = superseding_attach_sequence(current_sequence, &rows) else {
return failure(owner, operation, LiveFrontierError::Frontier);
};
let Some(order) = superseding_attach_order(
current_order,
attached.admission_order().transaction_order(),
) else {
return failure(owner, operation, LiveFrontierError::Frontier);
};
(
rows,
vec![terminal_charge, attached_charge],
sequence,
order,
)
}
AttachTransition::FencedRecovery {
prior_binding_epoch,
composed_terminal,
next_closure_state,
} => {
return apply_fenced_attach_frontier(
owner,
operation,
terminal_charge,
attached_charge,
prior_binding_epoch,
composed_terminal,
next_closure_state,
);
}
};
transition(
owner,
operation,
active,
&rows,
keyed_charges,
sequence,
order,
)
}
pub fn apply_detach_frontier<EF, V>(
owner: LiveFrontierOwner,
operation: CommittedDetachTransition<EF, V>,
terminal_charge: RetainedRecordCharge,
) -> LiveFrontierResult<CommittedDetachTransition<EF, V>> {
let terminal = operation.terminal();
if terminal.conversation_id() != owner.frontiers.conversation_id() {
return failure(owner, operation, LiveFrontierError::Authority);
}
let mut active = owner.frontiers.active_identities().participants().to_vec();
let Some(participant) = active
.iter_mut()
.find(|participant| participant.participant_index() == terminal.participant_id())
else {
return failure(owner, operation, LiveFrontierError::Authority);
};
*participant = FrontierParticipant::new(
participant.participant_index(),
operation.member().cursor(),
FrontierBinding::Detached(terminal.binding_epoch()),
);
let row = retained_terminal(terminal.into());
let Some(sequence) = detach_sequence(owner.frontiers.sequence().ledger(), row.delivery_seq)
else {
return failure(owner, operation, LiveFrontierError::Frontier);
};
let Some(order) = detach_order(
owner.frontiers.order().ledger(),
row.admission_order.transaction_order(),
) else {
return failure(owner, operation, LiveFrontierError::Frontier);
};
transition(
owner,
operation,
active,
&[row],
vec![terminal_charge],
sequence,
order,
)
}
pub fn apply_participant_ack_frontier(
mut owner: LiveFrontierOwner,
operation: ParticipantAckCommit,
) -> LiveFrontierResult<ParticipantAckCommit> {
let request = operation.outcome().request();
let Some(current) = owner
.frontiers
.active_identities()
.participants()
.iter()
.find(|participant| participant.participant_index() == request.participant_id)
.copied()
else {
return failure(owner, operation, LiveFrontierError::Authority);
};
let participant = FrontierParticipant::new(
request.participant_id,
request.through_seq,
current.binding(),
);
owner.frontiers = match owner.frontiers.apply_live_identity(participant) {
Ok(frontiers) => frontiers,
Err(frontier_failure) => {
let (frontiers, error) = *frontier_failure;
owner.frontiers = frontiers;
return failure(owner, operation, map_frontier_error(error));
}
};
Ok(LiveFrontierCommit { operation, owner })
}
pub fn apply_nonzero_participant_ack_frontier(
mut owner: LiveFrontierOwner,
operation: NonzeroParticipantAckCommit,
) -> LiveFrontierResult<NonzeroParticipantAckCommit> {
let request = operation.outcome().request();
let Some(current) = owner
.frontiers
.active_identities()
.participants()
.iter()
.find(|participant| participant.participant_index() == request.participant_id)
.copied()
else {
return failure(owner, operation, LiveFrontierError::Authority);
};
let participant = FrontierParticipant::new(
request.participant_id,
request.through_seq,
current.binding(),
);
owner.frontiers = match owner.frontiers.apply_live_identity(participant) {
Ok(frontiers) => frontiers,
Err(frontier_failure) => {
let (frontiers, error) = *frontier_failure;
owner.frontiers = frontiers;
return failure(owner, operation, map_frontier_error(error));
}
};
Ok(LiveFrontierCommit { operation, owner })
}
pub fn apply_marker_ack_frontier(
mut owner: LiveFrontierOwner,
operation: MarkerAckCommit,
) -> LiveFrontierResult<MarkerAckCommit> {
let request = operation.outcome().request();
if !owner
.frontiers
.retained_marker_records()
.iter()
.any(|record| {
record.delivery_seq == request.marker_delivery_seq
&& matches!(
record.kind,
RetainedCausalRecordKind::CompactionMarker { participant_index, .. }
if participant_index == request.participant_id
)
})
{
return failure(owner, operation, LiveFrontierError::Authority);
}
let Some(current) = owner
.frontiers
.active_identities()
.participants()
.iter()
.find(|participant| participant.participant_index() == request.participant_id)
.copied()
else {
return failure(owner, operation, LiveFrontierError::Authority);
};
let Some(accounting) = accounting_after_marker_ack(owner.closure_accounting) else {
return failure(owner, operation, LiveFrontierError::ClosureAccounting);
};
let participant = FrontierParticipant::new(
request.participant_id,
request.marker_delivery_seq,
current.binding(),
);
owner.frontiers = match owner.frontiers.apply_live_identity(participant) {
Ok(frontiers) => frontiers,
Err(frontier_failure) => {
let (frontiers, error) = *frontier_failure;
owner.frontiers = frontiers;
return failure(owner, operation, map_frontier_error(error));
}
};
owner.closure_accounting = accounting;
Ok(LiveFrontierCommit { operation, owner })
}
fn apply_fenced_attach_frontier<F, V>(
owner: LiveFrontierOwner,
operation: AttachCommit<F, V>,
terminal_charge: Option<RetainedRecordCharge>,
attached_charge: RetainedRecordCharge,
prior_binding_epoch: crate::wire::BindingEpoch,
composed_terminal: Option<super::super::CommittedBindingTerminal>,
next_closure_state: super::super::ClosureState,
) -> LiveFrontierResult<AttachCommit<F, V>> {
let attached = operation.attached;
let (rows, charges) = match (composed_terminal, terminal_charge) {
(None, None) => (vec![retained_attached(attached)], vec![attached_charge]),
(Some(terminal), Some(terminal_charge)) => (
vec![retained_terminal(terminal), retained_attached(attached)],
vec![terminal_charge, attached_charge],
),
(None, Some(_)) | (Some(_), None) => {
return failure(owner, operation, LiveFrontierError::RetainedCharge);
}
};
let participant = FrontierParticipant::new(
attached.participant_id(),
operation.member.cursor(),
FrontierBinding::Bound(attached.binding_epoch()),
);
fenced_attach_transition(
owner,
operation,
participant,
prior_binding_epoch,
next_closure_state,
&rows,
charges,
)
}
fn fenced_attach_transition<T>(
mut owner: LiveFrontierOwner,
operation: T,
participant: FrontierParticipant,
prior_binding_epoch: crate::wire::BindingEpoch,
next_closure_state: super::super::ClosureState,
rows: &[RetainedCausalRecord],
charges: Vec<RetainedRecordCharge>,
) -> LiveFrontierResult<T> {
if rows.len() != charges.len()
|| rows.iter().zip(&charges).any(|(row, charge)| {
row.delivery_seq != charge.delivery_seq()
|| row.admission_order != charge.admission_order()
|| charge.encoded_charge().entries != 1
})
{
return failure(owner, operation, LiveFrontierError::RetainedCharge);
}
let resulting_len = owner
.frontiers
.retained_records()
.len()
.checked_add(rows.len());
if resulting_len
.and_then(|len| u64::try_from(len).ok())
.is_none_or(|len| len > owner.retained_record_limit)
{
return failure(owner, operation, LiveFrontierError::RetainedRecordLimit);
}
let Some(accounting) =
accounting_after_fenced_attach(owner.closure_accounting, &charges, next_closure_state)
else {
return failure(owner, operation, LiveFrontierError::ClosureAccounting);
};
owner.frontiers =
match owner
.frontiers
.apply_live_fenced_attach(participant, prior_binding_epoch, rows)
{
Ok(frontiers) => frontiers,
Err(frontier_failure) => {
let (frontiers, error) = *frontier_failure;
owner.frontiers = frontiers;
return failure(owner, operation, map_frontier_error(error));
}
};
owner.retained_charges.extend(charges);
owner
.retained_charges
.sort_unstable_by_key(|charge| charge.delivery_seq());
owner.closure_accounting = accounting;
Ok(LiveFrontierCommit { operation, owner })
}
fn transition<T>(
mut owner: LiveFrontierOwner,
operation: T,
active: Vec<FrontierParticipant>,
rows: &[RetainedCausalRecord],
charges: Vec<RetainedRecordCharge>,
sequence: SequenceLedger,
order: OrderLedger,
) -> LiveFrontierResult<T> {
if rows.len() != charges.len()
|| rows.iter().zip(&charges).any(|(row, charge)| {
row.delivery_seq != charge.delivery_seq()
|| row.admission_order != charge.admission_order()
|| charge.encoded_charge().entries != 1
})
{
return failure(owner, operation, LiveFrontierError::RetainedCharge);
}
let resulting_len = owner
.frontiers
.retained_records()
.len()
.checked_add(rows.len());
if resulting_len
.and_then(|len| u64::try_from(len).ok())
.is_none_or(|len| len > owner.retained_record_limit)
{
return failure(owner, operation, LiveFrontierError::RetainedRecordLimit);
}
let Some(accounting) = accounting_after_rows(owner.closure_accounting, &charges) else {
return failure(owner, operation, LiveFrontierError::ClosureAccounting);
};
owner.frontiers = match owner
.frontiers
.apply_live_transition(active, rows, sequence, order)
{
Ok(frontiers) => frontiers,
Err(frontier_failure) => {
let (frontiers, error) = *frontier_failure;
owner.frontiers = frontiers;
return failure(owner, operation, map_frontier_error(error));
}
};
owner.retained_charges.extend(charges);
owner.closure_accounting = accounting;
Ok(LiveFrontierCommit { operation, owner })
}
const fn map_frontier_error(error: LiveFrontierTransitionError) -> LiveFrontierError {
match error {
LiveFrontierTransitionError::Authority => LiveFrontierError::Authority,
LiveFrontierTransitionError::Precedence => LiveFrontierError::Precedence,
LiveFrontierTransitionError::RecordPosition
| LiveFrontierTransitionError::Exhausted
| LiveFrontierTransitionError::ResultingFrontier => LiveFrontierError::Frontier,
}
}
fn failure<T, U>(
owner: LiveFrontierOwner,
operation: T,
error: LiveFrontierError,
) -> Result<U, Box<LiveFrontierFailure<T>>> {
Err(Box::new(LiveFrontierFailure {
error,
operation,
owner,
}))
}