use liminal_protocol::lifecycle::{
BindingState, ImmutableSequenceCandidate, LiveFrontierOwner, ObserverProgressProjection,
PendingDiedFinalization, PendingFinalization, RetainedRecordCharge,
};
use crate::server::participant::dispatch_impact::DispatchImpactAccumulator;
use super::binding_fate_completion::stored_died_cause;
use super::connection_fate::record_terminal_impact;
use super::fate_occurrence::PendingFinalizerRoute;
use super::frontier::terminal_charge;
use super::log::{
StoredBindingEpoch, StoredDied, StoredDrainedTerminal, StoredFinalizerPresentation,
StoredOperation, StoredOrdinaryTerminalSource, StoredPendingDiedFinalizer,
StoredTerminalDisposition,
};
use super::observer_progress::ObserverProgressSourceMetadata;
use super::outbox_projection::capture_projection_prestate;
use super::state::{ConversationAuthority, DurableAppend, StateError};
#[derive(Clone, Copy)]
struct DrainAuthority {
died: PendingDiedFinalization,
pending: PendingFinalization,
route: PendingFinalizerRoute,
}
impl ConversationAuthority {
pub(super) fn persist_drain_first(
&mut self,
candidate: ImmutableSequenceCandidate,
owner: LiveFrontierOwner,
appender: &dyn DurableAppend,
impact: &mut DispatchImpactAccumulator,
) -> Result<(), StateError> {
match candidate {
ImmutableSequenceCandidate::Marker(_) => {
self.persist_next_marker(candidate, owner, appender, impact)
}
ImmutableSequenceCandidate::BindingTerminal { .. } => {
self.persist_terminal_drain(candidate, owner, appender, impact)
}
}
}
fn persist_terminal_drain(
&mut self,
candidate: ImmutableSequenceCandidate,
owner: LiveFrontierOwner,
appender: &dyn DurableAppend,
impact: &mut DispatchImpactAccumulator,
) -> Result<(), StateError> {
let authority = match self.validate_drain_candidate(candidate) {
Ok(authority) => authority,
Err(error) => {
self.install_frontier(owner)?;
return Err(error);
}
};
let participant_id = authority.pending.participant_id();
let delivery_seq = candidate_delivery_seq(&candidate)?;
let source_log_sequence = self.next_log_sequence;
let (owner, projection, completes_ordinary) =
self.drain_validated_candidate(authority, delivery_seq, source_log_sequence, owner)?;
let row = StoredDied {
participant_id,
binding_epoch: StoredBindingEpoch::from(authority.pending.binding_epoch()),
cause: stored_died_cause(authority.died.cause()),
terminal_order: authority.pending.admission_order().transaction_order(),
disposition: StoredTerminalDisposition::Committed {
terminal_seq: delivery_seq,
},
connection_intent_sequence: None,
specific_fate_intent: None,
drained: Some(StoredDrainedTerminal {
pending_source_sequence: authority.route.pending_source_sequence,
finalizer_presentation: authority.route.presentation,
}),
};
let source = StoredOperation::Died { row };
let projection_facts = capture_projection_prestate(self, &source);
appender.append(&source, source_log_sequence)?;
self.install_frontier(owner)?;
self.release_drained_binding_slot(participant_id);
self.observe_replayed_position(
authority.pending.admission_order().transaction_order(),
delivery_seq,
)?;
self.advance_log_head()?;
self.record_drain_presentation(
authority.route.presentation,
source_log_sequence,
participant_id,
delivery_seq,
projection,
)?;
record_terminal_impact(
self,
source_log_sequence,
&source,
projection_facts,
participant_id,
impact,
)?;
if completes_ordinary {
self.complete_prepared_ordinary_finalizer(participant_id, appender)?;
self.record_episode_changed(impact);
}
Ok(())
}
pub(super) fn replay_died_row(
&mut self,
row: &StoredDied,
sequence: u64,
) -> Result<(), StateError> {
if row.drained.is_some() {
self.replay_terminal_drain(row, sequence)
} else {
self.replay_died_source(row, sequence)
}
}
fn replay_terminal_drain(&mut self, row: &StoredDied, sequence: u64) -> Result<(), StateError> {
let terminal_seq = validate_drain_row_shape(row)?;
if self.next_log_sequence != sequence {
return Err(StateError::invariant(
"terminal drain replay log head disagrees with durable sequence",
));
}
let owner = self.take_frontier()?;
let candidate = owner
.frontiers()
.sequence()
.immutable_candidates()
.first()
.copied()
.ok_or_else(|| StateError::invariant("durable terminal drain has no candidate"))?;
let authority = match self.validate_drain_candidate(candidate) {
Ok(authority) => authority,
Err(error) => {
self.install_frontier(owner)?;
return Err(error);
}
};
if let Err(error) =
validate_drain_row_authority(row, terminal_seq, authority, self.next_seq)
{
self.install_frontier(owner)?;
return Err(error);
}
let participant_id = authority.pending.participant_id();
let (owner, projection, _completes_ordinary) =
self.drain_validated_candidate(authority, terminal_seq, sequence, owner)?;
self.install_frontier(owner)?;
self.release_drained_binding_slot(participant_id);
self.observe_replayed_position(row.terminal_order, terminal_seq)?;
self.advance_log_head()?;
self.record_drain_presentation(
authority.route.presentation,
sequence,
participant_id,
terminal_seq,
projection,
)
}
fn drain_validated_candidate(
&mut self,
authority: DrainAuthority,
delivery_seq: u64,
finalizer_source_sequence: u64,
owner: LiveFrontierOwner,
) -> Result<(LiveFrontierOwner, ObserverProgressProjection, bool), StateError> {
let DrainAuthority {
died,
pending,
route,
} = authority;
let participant_id = pending.participant_id();
let (owner, completes_ordinary) = self.prepare_pending_died_finalizer(
participant_id,
route.pending_source_sequence,
died.commit(delivery_seq),
StoredOrdinaryTerminalSource::PendingDiedFinalized {
died_source_sequence: route.pending_source_sequence,
finalizer: StoredPendingDiedFinalizer::Drained {
source_sequence: finalizer_source_sequence,
},
},
owner,
)?;
let charge = match terminal_charge(
pending.conversation_id(),
participant_id,
pending.binding_epoch(),
pending.admission_order().transaction_order(),
delivery_seq,
) {
Ok(charge) => charge,
Err(error) => {
self.install_frontier(owner)?;
return Err(error);
}
};
let terminal_row_charge =
RetainedRecordCharge::new(delivery_seq, pending.admission_order(), charge);
let drained = match owner.drain_pending_terminal(pending, terminal_row_charge) {
Ok(drained) => drained,
Err(refused) => {
let error = refused.error();
self.install_frontier(refused.into_owner())?;
return Err(StateError::invariant(format!(
"candidate-lane terminal drain refused: {error:?}"
)));
}
};
let (owner, projection) = drained.into_parts();
Ok((owner, projection, completes_ordinary))
}
fn record_drain_presentation(
&mut self,
presentation: StoredFinalizerPresentation,
source_sequence: u64,
participant_id: u64,
terminal_seq: u64,
projection: ObserverProgressProjection,
) -> Result<(), StateError> {
if matches!(
presentation,
StoredFinalizerPresentation::ConsumeRecoveredReservation { .. }
) {
return Ok(());
}
let metadata = ObserverProgressSourceMetadata::died(
source_sequence,
self.conversation_id,
participant_id,
terminal_seq,
);
self.record_observer_progress_projection(projection, metadata)
}
fn release_drained_binding_slot(&mut self, participant_id: u64) {
self.slots.remove(&participant_id);
self.tokens.retain(|_, mapped| *mapped != participant_id);
}
fn validate_drain_candidate(
&mut self,
candidate: ImmutableSequenceCandidate,
) -> Result<DrainAuthority, StateError> {
let ImmutableSequenceCandidate::BindingTerminal {
delivery_seq,
admission_order,
owner: terminal_owner,
} = candidate
else {
return Err(StateError::invariant(
"terminal drain selected marker work instead of a binding terminal",
));
};
let participant_id = terminal_owner.participant_index;
let Some(slot) = self.slots.get(&participant_id) else {
return Err(StateError::invariant(
"terminal drain candidate names an absent participant slot",
));
};
let BindingState::PendingFinalization(pending) = slot.binding else {
return Err(StateError::invariant(
"terminal drain candidate does not rest in PendingFinalization",
));
};
let PendingFinalization::Died(died) = pending else {
return Err(StateError::invariant(
"terminal drain candidate is not a pending Died residence",
));
};
if pending.binding_epoch() != terminal_owner.binding_epoch
|| pending.admission_order() != admission_order
|| delivery_seq != self.next_seq
{
return Err(StateError::invariant(
"terminal drain candidate disagrees with its pending residence",
));
}
let route = self
.select_leave_finalizer(participant_id)?
.ok_or_else(|| {
StateError::invariant("pending terminal drain lost its finalizer route")
})?;
Ok(DrainAuthority {
died,
pending,
route,
})
}
}
fn candidate_delivery_seq(candidate: &ImmutableSequenceCandidate) -> Result<u64, StateError> {
match candidate {
ImmutableSequenceCandidate::BindingTerminal { delivery_seq, .. } => Ok(*delivery_seq),
ImmutableSequenceCandidate::Marker(_) => Err(StateError::invariant(
"terminal drain selected marker work instead of a binding terminal",
)),
}
}
fn validate_drain_row_shape(row: &StoredDied) -> Result<u64, StateError> {
if row.drained.is_none() {
return Err(StateError::invariant(
"terminal drain replay received a source Died row",
));
}
let StoredTerminalDisposition::Committed { terminal_seq } = row.disposition else {
return Err(StateError::invariant(
"durable terminal drain row is not committed",
));
};
if row.specific_fate_intent.is_some() || row.connection_intent_sequence.is_some() {
return Err(StateError::invariant(
"durable terminal drain row carries source-row authority",
));
}
Ok(terminal_seq)
}
fn validate_drain_row_authority(
row: &StoredDied,
terminal_seq: u64,
authority: DrainAuthority,
next_seq: u64,
) -> Result<(), StateError> {
let Some(drain) = row.drained else {
return Err(StateError::invariant(
"terminal drain replay received a source Died row",
));
};
if row.participant_id != authority.pending.participant_id()
|| row.binding_epoch.to_epoch()? != authority.pending.binding_epoch()
|| row.terminal_order != authority.pending.admission_order().transaction_order()
|| row.cause != stored_died_cause(authority.died.cause())
|| terminal_seq != next_seq
{
return Err(StateError::invariant(
"durable terminal drain row disagrees with its pending residence",
));
}
if authority.route.pending_source_sequence != drain.pending_source_sequence
|| authority.route.presentation != drain.finalizer_presentation
{
return Err(StateError::invariant(
"durable terminal drain finalizer route drifted",
));
}
Ok(())
}