use std::error::Error as StdError;
use hashgraph_like_consensus::{storage::ConsensusStorage, types::ConsensusEvent};
use openmls_traits::signatures::Signer;
use openmls_traits::{OpenMlsProvider, storage::StorageProvider};
use prost::Message;
use tracing::{error, info};
use crate::{
ConsensusApplyResult, ConsensusPlugin, Conversation, ConversationError, ConversationEvent,
ConversationState, PeerScoringPlugin, ScoreOp, StewardListPlugin, apply_consensus_result,
emergency_score_ops,
protos::de_mls::messages::v1::{
ConversationUpdateRequest, StewardElectionProposal, conversation_update_request,
},
};
impl<C, Sc, St> Conversation<C, Sc, St>
where
C: ConsensusPlugin,
Sc: PeerScoringPlugin,
St: StewardListPlugin,
{
pub(crate) fn apply_consensus_outcome<Pr>(
&mut self,
provider: &Pr,
event: ConsensusEvent,
signer: &impl Signer,
) -> Result<(), ConversationError>
where
Pr: OpenMlsProvider,
<Pr::StorageProvider as StorageProvider<1>>::Error: StdError + Send + Sync + 'static,
{
let (proposal_id, approved, timestamp) = match &event {
ConsensusEvent::ConsensusReached {
proposal_id,
result,
timestamp,
} => (*proposal_id, *result, *timestamp),
ConsensusEvent::ConsensusFailed {
proposal_id,
timestamp,
} => (*proposal_id, false, *timestamp),
};
self.cancel_auto_vote(proposal_id);
self.unregister_consensus_timeout(proposal_id);
let already_applied = self.queues.is_consensus_outcome_applied(proposal_id);
if already_applied {
tracing::debug!(
conversation = %self.conversation_id,
proposal_id,
"duplicate consensus outcome dropped"
);
return Ok(());
}
self.emit_event(ConversationEvent::ConsensusReached {
proposal_id,
approved,
timestamp,
});
let scope = C::Scope::from(self.conversation_id.clone());
let proposal = self
.services
.consensus
.storage()
.get_proposal(&scope, proposal_id)?;
let request = ConversationUpdateRequest::decode(proposal.payload.as_slice())?;
info!(
conversation = %self.conversation_id,
proposal_id, approved, "consensus reached"
);
self.queues.mark_consensus_outcome_applied(proposal_id);
let consensus_apply =
apply_consensus_result(&mut self.queues, proposal_id, approved, &request)?;
self.replay_early_candidates()?;
match consensus_apply {
ConsensusApplyResult::NoAction => {}
ConsensusApplyResult::ElectionAccepted(election) => {
self.handle_election_accepted(provider, election, signer)?;
}
ConsensusApplyResult::ElectionRejected => {
self.handle_election_rejected(provider, signer)?;
}
ConsensusApplyResult::RecoveryModeOpened => {
self.enter_recovery_mode();
self.start_freezing_and_emit();
}
ConsensusApplyResult::UrgentRemoval { target } => {
self.start_freezing_and_emit();
self.refresh_stewards_after_removal(provider, &target, signer)?;
}
ConsensusApplyResult::QueuedRemoval { target } => {
self.refresh_stewards_after_removal(provider, &target, signer)?;
}
ConsensusApplyResult::RejectedMembership { target } => {
self.queues.remove_pending_update(&target);
}
}
let score_ops = emergency_score_ops(&request, approved);
if !score_ops.is_empty() {
self.handle_emergency_scored(provider, proposal_id, &request, &score_ops, signer)?;
}
Ok(())
}
fn emit_phase_change(&self, transition: Option<ConversationState>) {
if let Some(state) = transition {
self.emit_event(ConversationEvent::PhaseChange(state));
}
}
fn start_freezing_and_emit(&mut self) {
let transition = self.start_freezing();
self.emit_phase_change(transition);
}
fn refresh_stewards_after_removal<Pr>(
&mut self,
provider: &Pr,
target: &[u8],
signer: &impl Signer,
) -> Result<(), ConversationError>
where
Pr: OpenMlsProvider,
<Pr::StorageProvider as StorageProvider<1>>::Error: StdError + Send + Sync + 'static,
{
if !self.services.steward_list.is_steward(target) {
return Ok(());
}
if let Err(e) = self.initiate_steward_election(provider, true, signer) {
info!(
conversation = %self.conversation_id,
error = %e,
"post-removal steward-list refresh deferred"
);
}
Ok(())
}
fn handle_election_accepted<Pr>(
&mut self,
provider: &Pr,
election: StewardElectionProposal,
signer: &impl Signer,
) -> Result<(), ConversationError>
where
Pr: OpenMlsProvider,
<Pr::StorageProvider as StorageProvider<1>>::Error: StdError + Send + Sync + 'static,
{
let is_valid = self.services.steward_list.validate_proposed(
&election.proposed_stewards,
election.election_epoch,
&election.proposed_stewards,
election.retry_round,
)?;
if !is_valid {
info!(
conversation = %self.conversation_id,
"steward election rejected: invalid list"
);
return Ok(());
}
self.services.steward_list.install_list(
election.election_epoch,
&election.proposed_stewards,
election.proposed_stewards.len(),
election.retry_round,
)?;
self.exit_recovery_mode();
let resumed_from_reelection = if self.current_state() == ConversationState::Reelection {
Some(self.start_working())
} else {
None
};
self.emit_phase_change(resumed_from_reelection);
info!(
conversation = %self.conversation_id,
epoch = election.election_epoch,
stewards = election.proposed_stewards.len(),
retry_round = election.retry_round,
"steward election applied"
);
self.process_buffered_updates(provider, signer)
}
fn handle_election_rejected<Pr>(
&mut self,
provider: &Pr,
signer: &impl Signer,
) -> Result<(), ConversationError>
where
Pr: OpenMlsProvider,
<Pr::StorageProvider as StorageProvider<1>>::Error: StdError + Send + Sync + 'static,
{
self.services.steward_list.bump_retry();
let round = self.services.steward_list.next_retry_round();
let max = self.services.steward_list.max_retries();
if round > max {
info!(
conversation = %self.conversation_id,
round, max, "election retries exhausted; escalating to Layer 3"
);
if let Err(e) = self.initiate_deadlock_ecp(provider, signer) {
error!(conversation = %self.conversation_id, error = %e, "Deadlock ECP filing failed");
self.emit_event(ConversationEvent::Error {
operation: "Reelection stuck".to_string(),
message: e.to_string(),
});
}
return Ok(());
}
info!(
conversation = %self.conversation_id,
round, max, "steward election rejected, retrying"
);
if let Err(e) = self.initiate_steward_election(provider, true, signer) {
info!(conversation = %self.conversation_id, error = %e, "election retry deferred");
}
Ok(())
}
fn handle_emergency_scored<Pr>(
&mut self,
provider: &Pr,
proposal_id: u32,
request: &ConversationUpdateRequest,
score_ops: &[ScoreOp],
signer: &impl Signer,
) -> Result<(), ConversationError>
where
Pr: OpenMlsProvider,
<Pr::StorageProvider as StorageProvider<1>>::Error: StdError + Send + Sync + 'static,
{
let _ = self.services.scoring.apply_ops(score_ops);
if let Some(conversation_update_request::Payload::EmergencyCriteria(ec)) = &request.payload
&& let Some(ev) = &ec.evidence
{
self.queues.remove_pending_removal(&ev.target_member_id);
}
self.queues.remove_emergency(proposal_id);
let resumed_event = if self.current_state() == ConversationState::Reelection {
Some(self.start_working())
} else {
None
};
self.emit_phase_change(resumed_event);
if let Err(e) = self.check_and_initiate_score_removals(provider, signer) {
error!(conversation = %self.conversation_id, error = %e, "score-removal check failed");
}
Ok(())
}
}