1pub mod error;
2pub(crate) mod metrics;
3pub mod timer;
4
5use std::{
6 collections::{BTreeMap, HashMap, HashSet},
7 sync::Arc,
8 time::{Duration, Instant},
9};
10
11use bon::{Builder, bon};
12use committable::Commitment;
13use hotshot::{HotShotInitializer, traits::BlockPayload, types::SignatureKey};
14use hotshot_types::{
15 consensus::{ConsensusMetricsValue, ParticipationTracker},
16 data::{
17 EpochNumber, Leaf2, VidCommitment, VidCommitment2, ViewNumber,
18 vid_disperse::vid_total_weight,
19 },
20 epoch_membership::EpochMembershipCoordinator,
21 message::{Proposal as SignedProposal, UpgradeLock},
22 simple_certificate::{QuorumCertificate2, TimeoutCertificate2},
23 simple_vote::{HasEpoch, QuorumVote2, TimeoutVote2},
24 traits::{
25 block_contents::BlockHeader, metrics::Metrics, node_implementation::NodeType,
26 signature_key::StateSignatureKey,
27 },
28 utils::{epoch_from_block_number, is_epoch_root},
29 vid::avidm_gf2::{AvidmGf2Param, init_avidm_gf2_param},
30 vote::{HasViewNumber, Vote},
31};
32use time::OffsetDateTime;
33use tokio::{select, sync::oneshot};
34use tracing::{debug, error, info, warn};
35
36use crate::{
37 block::{BlockAndHeaderRequest, BlockBuilder, BlockBuilderConfig},
38 cert_verifier::CertVerifiers,
39 client::{ClientApi, ClientRequest, CoordinatorClient, QueryError},
40 consensus::{Consensus, ConsensusInput, ConsensusOutput, PreCutoverSeed},
41 coordinator::{
42 error::{CoordinatorError, ErrorSource, Severity},
43 timer::Timer,
44 },
45 epoch::{EpochManager, EpochRootResult},
46 helpers::proposal_commitment,
47 logging::KeyPrefix,
48 message::{
49 self, BlockMessage, CatchupEvidence, Certificate1, Certificate2, ConsensusMessage, Message,
50 MessageType, OpaqueMessage, Proposal, ProposalFetchMessage, ProposalMessage,
51 TimeoutOneHonest, TransactionMessage, Unchecked, Validated, Vote2,
52 },
53 network::Cliquenet,
54 outbox::Outbox,
55 proposal::{ProposalValidator, VidShareValidator},
56 state::{HeaderRequest, StateEntry, StateManager, StateManagerOutput},
57 storage::{NewProtocolStorage, Storage},
58 vid::{VidDisperseRequest, VidDisperser, VidFragmentAccumulator, VidReconstructor},
59 vote::{EpochRootTally, SimpleTally, VoteCollector},
60};
61
62pub(crate) const VID_RECONSTRUCT_GC_MARGIN: u64 = 5;
76
77const STORAGE_GC_MARGIN: u64 = 5;
80
81const EPOCH_CHANGE_LOOKAHEAD: u64 = 3;
88
89pub(crate) const MAX_VIEWS_AHEAD: ViewNumber = ViewNumber::new(30);
90
91#[derive(Builder)]
92pub struct Coordinator<T: NodeType, S> {
93 membership_coordinator: EpochMembershipCoordinator<T>,
94 consensus: Consensus<T>,
95 network: Cliquenet<T>,
96 state_manager: StateManager<T>,
97 #[builder(default)]
98 client: CoordinatorClient<T>,
99 vid_disperser: VidDisperser<T>,
100 vid_reconstructor: VidReconstructor<T>,
101 #[builder(default)]
102 vid_fragment_accumulator: VidFragmentAccumulator<T>,
103 vote1_collector: VoteCollector<T, SimpleTally<T, QuorumVote2<T>, QuorumCertificate2<T>>>,
104 vote2_collector: VoteCollector<T, SimpleTally<T, Vote2<T>, Certificate2<T>>>,
105 timeout_collector: VoteCollector<T, SimpleTally<T, TimeoutVote2<T>, TimeoutCertificate2<T>>>,
106 timeout_one_honest_collector:
107 VoteCollector<T, SimpleTally<T, TimeoutVote2<T>, TimeoutOneHonest<T>>>,
108 epoch_root_collector: VoteCollector<T, EpochRootTally<T>>,
109 cert_verifiers: CertVerifiers<T>,
110 epoch_manager: EpochManager<T>,
111 block_builder: BlockBuilder<T>,
112 proposal_validator: ProposalValidator<T>,
113 share_validator: VidShareValidator<T>,
114 storage: Storage<T, S>,
115 #[builder(default)]
116 outbox: Outbox<ConsensusOutput<T>>,
117 #[builder(default)]
118 coordinator_outbox: Outbox<OpaqueMessage<T::SignatureKey>>,
119 public_key: T::SignatureKey,
120 #[builder(default = KeyPrefix::from(&public_key))]
121 node_id: KeyPrefix,
122 timer: Timer,
123 #[builder(skip)]
124 pending_proposal_fetches: PendingProposalFetches<T>,
125 #[builder(skip)]
126 requested_missing_proposals: HashSet<ProposalFetchKey<T>>,
127 #[builder(skip)]
128 da_payloads: BTreeMap<(ViewNumber, VidCommitment2), PendingDa<T>>,
129 metrics: Option<metrics::Metrics>,
130 #[builder(default)]
131 participation: ParticipationTracker<T>,
132 #[builder(skip)]
133 voted_view: Option<ViewNumber>,
134 #[builder(skip)]
138 last_timeout_view: Option<ViewNumber>,
139 #[builder(skip)]
140 view_started: Option<(ViewNumber, EpochNumber, Instant)>,
141 #[builder(skip)]
142 proposal_received_at: Option<(ViewNumber, Instant)>,
143 #[builder(skip)]
144 invalid_certs_at_decide: u64,
145 #[builder(skip)]
146 payload_txn_bytes: BTreeMap<ViewNumber, usize>,
147}
148
149#[bon]
150impl<T, S> Coordinator<T, S>
151where
152 T: NodeType,
153 S: NewProtocolStorage<T>,
154{
155 #[builder(builder_type = CoordinatorMaker, finish_fn = make)]
156 #[allow(clippy::too_many_arguments)]
157 pub fn maker(
158 membership_coordinator: EpochMembershipCoordinator<T>,
159 network: Cliquenet<T>,
160 initializer: &HotShotInitializer<T>,
161 upgrade_lock: UpgradeLock<T>,
162 public_key: T::SignatureKey,
163 private_key: <T::SignatureKey as SignatureKey>::PrivateKey,
164 state_private_key: <T::StateSignatureKey as StateSignatureKey>::StatePrivateKey,
165 stake_table_capacity: usize,
166 timeout_duration: Duration,
167 storage: S,
168 metrics: &dyn Metrics,
169 consensus_metrics: ConsensusMetricsValue,
170 locked_qc: Option<Certificate1<T>>,
172 ) -> Self {
173 let mut consensus = Consensus::new(
174 membership_coordinator.clone(),
175 public_key.clone(),
176 private_key.clone(),
177 state_private_key,
178 stake_table_capacity,
179 upgrade_lock.clone(),
180 initializer.anchor_leaf.clone(),
181 initializer.epoch_height,
182 );
183
184 let anchor_leaf = &initializer.anchor_leaf;
185 let anchor_view = anchor_leaf.view_number();
186 let anchor_epoch = anchor_leaf
187 .epoch(initializer.epoch_height)
188 .unwrap_or(EpochNumber::genesis());
189 let cert1 = initializer.high_qc.clone();
190 let parent_proposal = message::Proposal {
191 block_header: anchor_leaf.block_header().clone(),
192 view_number: anchor_view,
193 epoch: anchor_epoch,
194 justify_qc: anchor_leaf.justify_qc(),
195 next_epoch_justify_qc: None,
196 upgrade_certificate: anchor_leaf.upgrade_certificate(),
197 view_change_evidence: anchor_leaf
198 .view_change_evidence
199 .clone()
200 .and_then(|e| match e {
201 hotshot_types::data::ViewChangeEvidence2::Timeout(tc) => Some(tc),
202 hotshot_types::data::ViewChangeEvidence2::ViewSync(_) => None,
203 }),
204 next_drb_result: anchor_leaf.next_drb_result,
205 state_cert: None,
206 };
207
208 let coordinator_metrics = metrics
209 .is_recording()
210 .then(|| metrics::Metrics::new(consensus_metrics));
211
212 let mut state_manager = StateManager::new(
213 Arc::new(initializer.instance_state.clone()),
214 upgrade_lock.clone(),
215 )
216 .with_metrics(
217 coordinator_metrics.as_ref().map(|m| {
218 m.consensus
219 .validate_and_apply_header_duration
220 .clone()
221 .into()
222 }),
223 coordinator_metrics
224 .as_ref()
225 .map(|m| m.consensus.update_leaf_duration.clone().into()),
226 );
227 for p in initializer.saved_proposals.values() {
230 state_manager.seed_from_header(message::Proposal::from(p.data.clone()));
231 }
232 state_manager.seed_state(
233 anchor_view,
234 initializer.anchor_state.clone(),
235 anchor_leaf.clone(),
236 );
237 let reconstructed_blocks =
241 std::iter::once((anchor_view, anchor_leaf.block_header().clone()))
242 .chain(
243 initializer
244 .saved_proposals
245 .iter()
246 .map(|(view, p)| (*view, p.data.block_header().clone())),
247 )
248 .filter_map(|(view, header)| match header.payload_commitment() {
249 VidCommitment::V2(commitment) => Some((view, commitment)),
250 _ => None,
251 });
252 let saved_proposals = initializer
254 .saved_proposals
255 .values()
256 .map(|p| message::Proposal::from(p.data.clone()));
257 consensus.seed_proposals(saved_proposals);
258 consensus.seed_parent(cert1, parent_proposal, reconstructed_blocks);
263 if let Some(locked_qc) = locked_qc {
266 consensus.seed_locked_cert(locked_qc);
267 }
268 consensus.resume_from_restart(
269 anchor_view,
270 initializer.start_view,
271 initializer.last_actioned_view,
272 );
273 if let Some(state_cert) = initializer.state_cert.clone() {
274 consensus.seed_state_cert(state_cert);
275 }
276
277 let participation = ParticipationTracker::new(&membership_coordinator, anchor_epoch);
278
279 let vid_disperser = VidDisperser::new(
280 membership_coordinator.clone(),
281 network.sender().clone(),
282 public_key.clone(),
283 private_key.clone(),
284 )
285 .with_metrics(
286 coordinator_metrics
287 .as_ref()
288 .map(|m| m.consensus.vid_disperse_duration.clone().into()),
289 );
290
291 let lock = upgrade_lock.clone();
292 Self::builder()
293 .consensus(consensus)
294 .network(network)
295 .state_manager(state_manager)
296 .vid_disperser(vid_disperser)
297 .vid_reconstructor(VidReconstructor::new())
298 .vote1_collector(VoteCollector::new(
299 membership_coordinator.clone(),
300 lock.clone(),
301 ))
302 .vote2_collector(VoteCollector::new(
303 membership_coordinator.clone(),
304 lock.clone(),
305 ))
306 .timeout_collector(VoteCollector::new(
307 membership_coordinator.clone(),
308 lock.clone(),
309 ))
310 .timeout_one_honest_collector(VoteCollector::new(
311 membership_coordinator.clone(),
312 lock.clone(),
313 ))
314 .epoch_root_collector(VoteCollector::new(
315 membership_coordinator.clone(),
316 lock.clone(),
317 ))
318 .cert_verifiers(CertVerifiers::new(
319 membership_coordinator.clone(),
320 lock.clone(),
321 ))
322 .epoch_manager(EpochManager::new(
323 initializer.epoch_height,
324 membership_coordinator.clone(),
325 ))
326 .block_builder(BlockBuilder::new(
327 Arc::new(initializer.instance_state.clone()),
328 membership_coordinator.clone(),
329 BlockBuilderConfig::default(),
330 upgrade_lock.clone(),
331 ))
332 .proposal_validator(ProposalValidator::new(
333 membership_coordinator.clone(),
334 initializer.epoch_height,
335 upgrade_lock.clone(),
336 ))
337 .share_validator(VidShareValidator::new(
338 membership_coordinator.clone(),
339 initializer.epoch_height,
340 upgrade_lock,
341 ))
342 .storage(Storage::new(storage, private_key).with_metrics(metrics))
343 .membership_coordinator(membership_coordinator)
344 .timer(Timer::new(timeout_duration, anchor_view, anchor_epoch))
345 .public_key(public_key)
346 .maybe_metrics(coordinator_metrics)
347 .participation(participation)
348 .build()
349 }
350
351 pub fn start(&mut self, seed: Option<PreCutoverSeed<T>>) {
360 if let Some(seed) = seed
361 && !self.apply_cutover_seed(seed)
362 {
363 return;
364 }
365
366 let cur_view = self.consensus.current_view();
367 let next_view = cur_view + 1;
368 let epoch = self
369 .consensus
370 .current_epoch()
371 .unwrap_or(EpochNumber::genesis());
372
373 if self.consensus.last_decided_leaf().view_number() == ViewNumber::genesis() {
374 let genesis_leaf = self.consensus.last_decided_leaf().clone();
376 let (payload, metadata) = T::BlockPayload::empty();
377 self.storage.append_da(
378 ViewNumber::genesis(),
379 EpochNumber::genesis(),
380 payload,
381 metadata,
382 genesis_leaf.payload_commitment(),
383 );
384
385 self.outbox.push_back(ConsensusOutput::LeafDecided {
387 leaves: vec![genesis_leaf],
388 cert1: self
389 .consensus
390 .cert1_at(ViewNumber::genesis())
391 .cloned()
392 .expect("genesis cert1 must be seeded"),
393 cert2: None,
394 vid_shares: vec![None],
395 });
396 }
397
398 self.outbox
399 .push_back(ConsensusOutput::ViewChanged(next_view, epoch));
400
401 if let Some(leader) = self.leader(next_view, epoch)
402 && leader == self.public_key
403 {
404 if let Some(parent_proposal) = self.consensus.proposal_at(cur_view).cloned() {
408 self.outbox
409 .push_back(ConsensusOutput::RequestBlockAndHeader(
410 BlockAndHeaderRequest {
411 view: next_view,
412 epoch,
413 parent_proposal,
414 },
415 ));
416 }
417 }
418 }
419
420 pub async fn stop(mut self) {
421 futures::join!(self.network.shutdown(), self.storage.flush());
422 }
423
424 pub async fn next_consensus_input(&mut self) -> Result<ConsensusInput<T>, CoordinatorError> {
425 loop {
426 select! {
427 message = self.network.receive() => match message {
428 Ok(m) => {
429 if let Some(input) = self.on_network_message(m) {
430 return Ok(input)
431 }
432 }
433 Err(e) => {
434 return Err(CoordinatorError::from(e).context("network receive"))
435 }
436 },
437 () = &mut self.timer => {
438 let view = self.timer.view();
439 let epoch = self.timer.epoch();
440 self.timer.reset();
445 if let Some(stats) = self.vote1_collector.stats(view, epoch) {
446 warn!(
447 %view, %epoch,
448 stake = %stats.stake,
449 threshold = %stats.threshold,
450 "timeout: vote1 stake observed (deduped by signer)"
451 );
452 } else {
453 warn!(%view, %epoch, "timeout: no vote1 received for this view");
454 }
455 let input = ConsensusInput::Timeout(view, epoch);
456 if self.last_timeout_view != Some(view) {
457 self.last_timeout_view = Some(view);
458 let leader = self.leader(view, epoch);
459 if let Some(leader) = leader.clone() {
460 self.participation.leader_missed(leader, epoch);
461 }
462 if let Some(m) = &self.metrics {
463 m.consensus.number_of_timeouts.add(1);
464 if leader.as_ref() == Some(&self.public_key) {
465 m.consensus.number_of_timeouts_as_leader.add(1);
466 }
467 }
468 }
469 return Ok(input)
470 }
471 Some(output) = self.state_manager.next() => {
472 if let Some(input) = self.on_state_manager_output(output) {
473 return Ok(input)
474 }
475 }
476 Some(request) = self.client.next_request() => {
477 if let Err(err) = self.on_client_request(request) {
478 error!(%err, "error while handling client request");
479 }
480 }
481 Some(tcert) = self.timeout_collector.next() => {
482 self.cert_verifiers.timeout.mark_completed(tcert.view_number());
483 return Ok(ConsensusInput::TimeoutCertificate(tcert))
484 }
485 Some(out) = self.timeout_one_honest_collector.next() => {
486 let Some(epoch) = out.data.epoch else {
487 let msg = format!("missing epoch in view {}", out.view_number());
488 return Err(CoordinatorError::regular(msg).context("gc timeout one honest"))
489 };
490 return Ok(ConsensusInput::TimeoutOneHonest(out.view_number(), epoch))
491 }
492 Some(cert1) = self.vote1_collector.next() => {
493 self.cert_verifiers.cert1.mark_completed(cert1.view_number());
494 return Ok(ConsensusInput::Certificate1(cert1))
495 }
496 Some(cert2) = self.vote2_collector.next() => {
497 self.cert_verifiers.cert2.mark_completed(cert2.view_number());
498 return Ok(ConsensusInput::Certificate2(cert2))
499 }
500 Some(cert1) = self.cert_verifiers.cert1.next() => {
501 return Ok(ConsensusInput::Certificate1(cert1))
502 }
503 Some(cert2) = self.cert_verifiers.cert2.next() => {
504 return Ok(ConsensusInput::Certificate2(cert2))
505 }
506 Some(tc) = self.cert_verifiers.timeout.next() => {
507 return Ok(ConsensusInput::TimeoutCertificate(tc))
508 }
509 Some(cert1) = self.cert_verifiers.advance.next() => {
510 return Ok(ConsensusInput::AdvanceView(cert1))
511 }
512 Some(epoch_change) = self.cert_verifiers.epoch_change.next() => {
513 return Ok(ConsensusInput::EpochChange(epoch_change.into_cert()))
514 }
515 Some((cert1, state_cert)) = self.epoch_root_collector.next() => {
516 self.cert_verifiers.cert1.mark_completed(cert1.view_number());
517 self.storage.append_state_cert(
518 ViewNumber::new(state_cert.light_client_state.view_number),
519 state_cert.clone(),
520 );
521 return Ok(ConsensusInput::EpochRootCertificates { cert1, state_cert })
522 }
523 Some(item) = self.share_validator.next() => match item {
524 Ok(vid_share) => {
525 return Ok(ConsensusInput::VidShare(vid_share))
526 },
527 Err(e) => {
528 return Err(CoordinatorError::regular(e).context("vid share validation"))
529 }
530 },
531 Some(item) = self.proposal_validator.next() => match item {
532 Ok(validated) if validated.fetched => {
533 return Ok(ConsensusInput::FetchedProposal(validated.message))
534 }
535 Ok(validated) => {
536 let epoch = validated.message.proposal.data.epoch;
538 self.bump_network_epoch(epoch);
539 return Ok(ConsensusInput::Proposal(validated.sender, validated.message))
540 }
541 Err(e) => {
542 return Err(CoordinatorError::regular(e).context("proposal validation"))
543 }
544 },
545 Some(item) = self.block_builder.next() => match item {
546 Ok(block) => {
547 self.state_manager.request_header(HeaderRequest::from(&block));
548 let next_view = block.view + 1;
549 let epoch = block.epoch;
550 let manifest = block.manifest.clone();
551 if let VidCommitment::V2(commit) = block.payload_commitment {
554 self.da_payloads.insert(
555 (block.view, commit),
556 PendingDa {
557 epoch: block.epoch,
558 payload: block.payload.payload.clone(),
559 metadata: block.payload.metadata.clone(),
560 },
561 );
562 } else {
563 warn!(view = %block.view, "block payload commitment is not V2");
564 }
565 self.vid_reconstructor.retire_view(block.view);
567 self.unicast_to_leader(
568 next_view,
569 epoch,
570 BlockMessage::DedupManifest(manifest),
571 )?;
572 return Ok(block.into())
573 }
574 Err(err) => {
575 return Err(CoordinatorError::regular(err).context("block building"))
576 }
577 },
578 Some(item) = self.vid_disperser.next() => match item {
579 Ok(out) => {
580 return Ok(ConsensusInput::VidDisperseCreated(out.view, out.payload_commitment))
581 }
582 Err(err) => {
583 return Err(CoordinatorError::from(err).context("vid disperse"))
584 }
585 },
586 Some(item) = self.vid_reconstructor.next() => match item {
587 Ok(out) => {
588 self.payload_txn_bytes.insert(out.view, out.payload.txn_bytes());
589 self.block_builder.on_block_reconstructed(out.tx_commitments);
590 self.storage.append_da(
591 out.view,
592 out.epoch,
593 out.payload.clone(),
594 out.metadata.clone(),
595 VidCommitment::V2(out.payload_commitment),
596 );
597 if let Some(proposal) = self.consensus.proposal_at(out.view) {
598 if proposal.block_header.payload_commitment()
600 == VidCommitment::V2(out.payload_commitment)
601 {
602 self.outbox.push_back(ConsensusOutput::BlockPayloadReconstructed {
603 view: out.view,
604 header: proposal.block_header.clone(),
605 payload: out.payload,
606 });
607 } else {
608 warn!(
609 view = %out.view,
610 header = %proposal.block_header.payload_commitment(),
611 reconstructed = %out.payload_commitment,
612 "reconstructed payload commitment does not match proposal header"
613 );
614 }
615 }
616 return Ok(ConsensusInput::BlockReconstructed(out.view, out.payload_commitment))
617 }
618 Err(err) => {
619 return Err(CoordinatorError::regular(err).context("vid reconstruction"))
620 }
621 },
622 Some(stored) = self.storage.next() => {
623 return Ok(ConsensusInput::Stored(stored))
624 },
625 Some(result) = self.epoch_manager.next() => match result {
626 Ok(EpochRootResult::DrbResult(epoch, drb_result)) => {
627 self.vote1_collector.retry_pending_votes();
628 self.vote2_collector.retry_pending_votes();
629 self.timeout_collector.retry_pending_votes();
630 self.timeout_one_honest_collector.retry_pending_votes();
631 self.epoch_root_collector.retry_pending_votes();
632 self.cert_verifiers.retry_pending(|e| self.epoch_manager.request_drb_result(e));
633 return Ok(ConsensusInput::DrbResult(epoch, drb_result))
634 }
635 Err(failure) => {
636 warn!(%failure.error, epoch = %failure.epoch, "DRB request failed");
641 continue;
642 }
643 },
644 else => {
645 return Err(CoordinatorError::critical(ErrorSource::NoInput))
646 }
647 }
648 }
649 }
650
651 pub fn apply_consensus(&mut self, input: ConsensusInput<T>) {
652 self.consensus.apply(input, &mut self.outbox)
653 }
654
655 pub fn process_consensus_output(
656 &mut self,
657 output: ConsensusOutput<T>,
658 ) -> Result<(), CoordinatorError> {
659 let node = self.node_id;
660 match output {
661 ConsensusOutput::RequestState(state_request) => {
662 debug!(
663 %node,
664 view = %state_request.view,
665 epoch = %state_request.epoch,
666 block = %state_request.block,
667 "request state validation"
668 );
669 self.state_manager.request_state(state_request);
670 },
671 ConsensusOutput::RequestVidDisperse {
672 view,
673 epoch,
674 payload,
675 metadata,
676 payload_commitment,
677 } => {
678 debug!(%node, %view, %epoch, "request vid disperse");
679 self.vid_disperser.request_vid_disperse(VidDisperseRequest {
680 view,
681 epoch,
682 block: payload,
683 metadata,
684 payload_commitment,
685 });
686 },
687 ConsensusOutput::RequestDrbResult(epoch) => {
688 debug!(%node, %epoch, "request drb result");
689 self.epoch_manager.request_drb_result(epoch);
690 },
691 ConsensusOutput::LeafDecided {
692 leaves,
693 cert1,
694 cert2,
695 ..
696 } => {
697 info!(
698 %node,
699 view = %cert1.view_number(),
700 epoch = ?cert1.epoch().map(|e| *e),
701 leaves = leaves.len(),
702 "leaves decided"
703 );
704 self.on_decide_metrics(&leaves);
705 if let Some(cert2) = cert2 {
706 self.storage.append_cert2(cert2.view_number, cert2.clone());
707 }
708 if let Some(newest) = leaves.first() {
711 let gc_view = newest.view_number();
712 let gc_epoch = newest.justify_qc().epoch().unwrap_or_default();
713 self.gc(gc_epoch, GcScope::Decided(gc_view))?;
714 }
715 for leaf in leaves.into_iter().rev() {
716 self.participation
717 .on_leaf_decided(&leaf, &self.membership_coordinator);
718 self.epoch_manager.handle_leaf_decided(leaf);
719 }
720 },
721 ConsensusOutput::LockUpdated(cert) => {
722 debug!(
723 %node,
724 view = %cert.view_number(),
725 epoch = ?cert.epoch().map(|e| *e),
726 "lock updated"
727 );
728 },
729 ConsensusOutput::RequestMissingProposal { view, leaf_commit } => {
730 debug!(%node, %view, "request missing proposal");
731 if let Err(err) = self.request_missing_proposal(view, leaf_commit) {
732 warn!(%node, %view, %err, "failed to request missing proposal");
733 }
734 },
735 ConsensusOutput::RequestBlockAndHeader(request) => {
736 debug!(
737 %node,
738 view = %request.view,
739 epoch = %request.epoch,
740 "request block and header"
741 );
742 self.block_builder.request_block(request);
743 },
744 ConsensusOutput::RecordAction(view, epoch, kind) => {
745 debug!(%node, %view, ?kind, "record action");
746 self.storage.record_action(view, epoch, kind);
747 },
748 ConsensusOutput::PersistProposal(proposal) => {
749 let view = proposal.data.view_number;
750 debug!(%node, %view, "persist proposal");
751 self.storage.append_proposal(proposal.data.clone());
752 if let VidCommitment::V2(commit) = proposal.data.block_header.payload_commitment() {
755 if let Some(da) = self.da_payloads.remove(&(view, commit)) {
756 self.payload_txn_bytes.insert(view, da.payload.txn_bytes());
757 if let Some(m) = &self.metrics
758 && da.payload.transactions(&da.metadata).next().is_none()
759 {
760 m.consensus.number_of_empty_blocks_proposed.add(1);
761 }
762 self.storage.append_da(
763 view,
764 da.epoch,
765 da.payload,
766 da.metadata,
767 VidCommitment::V2(commit),
768 );
769 } else {
770 warn!(%node, %view, "no payload for proposed block");
771 }
772 }
773 },
774 ConsensusOutput::ProposalPaired {
775 proposal,
776 vid_share,
777 } => {
778 let view = proposal.data.view_number;
779 debug!(%node, %view, "proposal paired with vid share");
780 self.storage.append_vid(vid_share.clone());
781 self.storage.append_proposal(proposal.data.clone());
782 if let Some(state_cert) = &proposal.data.state_cert {
783 self.storage.append_state_cert(
784 ViewNumber::new(state_cert.light_client_state.view_number),
785 state_cert.clone(),
786 );
787 }
788 let expected_param = self.expected_vid_param(vid_share.target_epoch);
789 self.vid_reconstructor.handle_proposal(
790 view,
791 vid_share.payload_commitment,
792 proposal.data.block_header.metadata().clone(),
793 proposal.data.epoch,
794 expected_param,
795 );
796 self.vid_reconstructor
797 .handle_vid_share(self.public_key.clone(), vid_share);
798 },
799 ConsensusOutput::SendProposal(proposal) => {
800 let view = proposal.data.view_number;
801 let epoch = proposal.data.epoch;
802 let block = proposal.data.block_header.block_number();
803 info!(%node, %view, %epoch, %block, "send proposal");
804 if let Some(m) = &self.metrics
805 && proposal.data.view_change_evidence.is_none()
806 && let Some((prev_view, received_at)) = self.proposal_received_at
807 && (*prev_view).checked_add(1) == Some(*view)
808 {
809 m.consensus
810 .previous_proposal_to_proposal_time
811 .add_point(received_at.elapsed().as_millis() as f64);
812 }
813 let message = Message {
814 sender: self.public_key.clone(),
815 message_type: MessageType::Consensus(ConsensusMessage::Proposal(
816 ProposalMessage::validated(proposal.clone()),
817 )),
818 };
819 if let Err(err) = self
820 .network
821 .sender()
822 .broadcast(self.consensus.current_view(), &message)
823 {
824 let err = CoordinatorError::from(err).context("proposal broadcast");
825 if err.severity == Severity::Critical {
826 return Err(err);
827 } else {
828 warn!(%node, %err, "network error while broadcasting proposal")
829 }
830 }
831 },
832 ConsensusOutput::SendTimeoutVote(vote, evidence) => {
833 let view = vote.view_number();
834 debug!(
835 %node, %view,
836 has_evidence = evidence.is_some(),
837 "send timeout vote"
838 );
839 self.broadcast(
840 ConsensusMessage::TimeoutVote(message::TimeoutVoteMessage { vote, evidence }),
841 "broadcast timeout vote",
842 )?
843 },
844 ConsensusOutput::SendTimeoutCertificate(tc, view, epoch) => {
845 debug!(
846 %node, %view, %epoch,
847 cert_view = %tc.view_number(),
848 "send timeout certificate"
849 );
850 if let Some(leader) = self.leader(view, epoch) {
851 let message = Message {
852 sender: self.public_key.clone(),
853 message_type: MessageType::Consensus(ConsensusMessage::TimeoutCertificate(
854 tc,
855 )),
856 };
857 self.network
858 .sender()
859 .unicast(self.consensus.current_view(), &leader, &message)
860 .map_err(|e| CoordinatorError::from(e).context("timeout certificate"))?;
861 }
862 },
863 ConsensusOutput::SendVote1(vote1) => {
864 let view = vote1.vote.view_number();
865 debug!(
866 %node, %view,
867 epoch_root = vote1.state_vote.is_some(),
868 "send vote1"
869 );
870 self.record_voted_view(view);
871 if let Some(epoch) = vote1.vote.data.epoch
872 && let Some(leader) = self.leader(view, epoch)
873 {
874 self.participation.leader_proposed(leader, epoch);
875 }
876 self.broadcast(ConsensusMessage::Vote1(vote1), "broadcast vote1")?
877 },
878 ConsensusOutput::BroadcastVidShare(share) => {
879 debug!(%node, view = %share.view_number(), "send vid share");
880 self.broadcast(
881 ConsensusMessage::VidShareBroadcast(share),
882 "broadcast vid share",
883 )?
884 },
885 ConsensusOutput::SendVote2(vote2) => {
886 let view = vote2.view_number();
887 debug!(%node, %view, "send vote2");
888 self.record_voted_view(view);
889 self.broadcast(ConsensusMessage::Vote2(vote2), "broadcast vote2")?
890 },
891 ConsensusOutput::PersistHighQc(high_qc) => {
892 debug!(%node, view = %high_qc.view_number(), "persist high qc");
893 self.storage.append_high_qc2(high_qc);
894 },
895 ConsensusOutput::SendEpochChange(epoch_change) => {
896 info!(
897 %node,
898 view = %epoch_change.cert1.view_number(),
899 epoch = ?epoch_change.cert1.epoch().map(|e| *e),
900 "send epoch change"
901 );
902 self.broadcast(
903 ConsensusMessage::EpochChange(epoch_change),
904 "broadcast epoch change",
905 )?
906 },
907 ConsensusOutput::SendCertificate1(cert1) => {
908 debug!(
909 %node,
910 view = %cert1.view_number(),
911 epoch = ?cert1.epoch().map(|e| *e),
912 "send certificate1"
913 );
914 self.broadcast(
915 ConsensusMessage::Certificate1(cert1, self.public_key.clone()),
916 "broadcast certificate1",
917 )?
918 },
919 ConsensusOutput::SendCertificate2(cert2) => {
920 debug!(
921 %node,
922 view = %cert2.view_number(),
923 epoch = ?cert2.epoch().map(|e| *e),
924 "send certificate2"
925 );
926 self.broadcast(
927 ConsensusMessage::Certificate2(cert2, self.public_key.clone()),
928 "broadcast certificate2",
929 )?
930 },
931 ConsensusOutput::ProposalValidated { proposal, sender } => {
932 debug!(
933 %node,
934 view = %proposal.data.view_number,
935 sender = %KeyPrefix::from(&sender),
936 "proposal validated"
937 );
938 },
939 ConsensusOutput::ViewChanged(view, epoch) => {
940 let current_view = self.consensus.current_view();
941 if view < current_view {
942 warn!(
943 %node, %view, %epoch, %current_view,
944 "ignoring view change to stale view"
945 );
946 return Ok(());
947 }
948 info!(%node, %view, %epoch, "view changed");
949 self.timer.reset_with_epoch(view, epoch);
950 self.gc(epoch, GcScope::Local(view))?;
951 let txns = self.block_builder.on_view_changed(view);
952 self.participation.on_view_changed(epoch);
953 self.on_view_changed_metrics(view, epoch);
954 if !txns.is_empty() {
955 let next_view = view + 1;
956 self.unicast_to_leader(
957 next_view,
958 epoch,
959 BlockMessage::Transactions(TransactionMessage {
960 view: next_view,
961 transactions: txns,
962 }),
963 )
964 .map_err(|e| e.context("unicast transactions"))?;
965 }
966
967 let next_epoch = epoch + 1;
970 if next_epoch > EpochNumber::genesis() + 1 {
971 self.epoch_manager.request_drb_result(next_epoch);
972 }
973 },
974 ConsensusOutput::ViewTimedOut(view) => {
975 debug!(%node, %view, "view timed out");
976 let epoch = self
977 .consensus
978 .current_epoch()
979 .unwrap_or_else(EpochNumber::genesis);
980 self.gc(epoch, GcScope::Timeout(view))?;
981 },
982 ConsensusOutput::BlockPayloadReconstructed { .. } => {},
983 }
984 Ok(())
985 }
986
987 pub fn node_id(&self) -> &KeyPrefix {
988 &self.node_id
989 }
990
991 pub fn outbox(&self) -> &Outbox<ConsensusOutput<T>> {
992 &self.outbox
993 }
994
995 pub fn outbox_mut(&mut self) -> &mut Outbox<ConsensusOutput<T>> {
996 &mut self.outbox
997 }
998
999 pub fn coordinator_outbox(&self) -> &Outbox<OpaqueMessage<T::SignatureKey>> {
1000 &self.coordinator_outbox
1001 }
1002
1003 pub fn coordinator_outbox_mut(&mut self) -> &mut Outbox<OpaqueMessage<T::SignatureKey>> {
1004 &mut self.coordinator_outbox
1005 }
1006
1007 pub fn current_view(&self) -> ViewNumber {
1008 self.consensus.current_view()
1009 }
1010
1011 pub fn state(&self, v: ViewNumber) -> Option<&StateEntry<T>> {
1012 self.state_manager.get_state(v)
1013 }
1014
1015 pub fn client_api(&self) -> &ClientApi<T> {
1016 self.client.handle()
1017 }
1018
1019 pub fn bump_network_epoch(&mut self, epoch: EpochNumber) {
1025 if let Err(err) = self
1026 .network
1027 .apply_epoch(epoch, &self.membership_coordinator)
1028 {
1029 error!(%epoch, %err, "network apply_epoch failed");
1030 }
1031 }
1032
1033 pub(crate) fn on_network_message(
1034 &mut self,
1035 message: Message<T, Unchecked>,
1036 ) -> Option<ConsensusInput<T>> {
1037 let sender = KeyPrefix::from(&message.sender);
1038 let node = self.node_id;
1039 match message.message_type {
1040 MessageType::Consensus(msg) => match msg {
1041 ConsensusMessage::Proposal(p) => {
1042 let view = p.view_number();
1043 let epoch = p.proposal.data.epoch;
1044 let block = p.proposal.data.block_header.block_number();
1045 debug!(%node, %sender, %view, %epoch, %block, "recv proposal");
1046 if !self.is_view_too_far_ahead(view)
1047 && self.proposal_received_at.is_none_or(|(v, _)| v < view)
1048 {
1049 self.proposal_received_at = Some((view, Instant::now()));
1050 }
1051 if self.consensus.wants_proposal_for_view(&view) {
1052 self.proposal_validator.validate(p);
1053 }
1054 None
1055 },
1056 ConsensusMessage::VidShareFragment(fragment) => {
1057 let view = fragment.data.view_number();
1058 debug!(%node, %sender, %view, "received vid share fragment");
1059 if fragment.data.recipient_key != self.public_key {
1060 warn!(
1061 %node,
1062 %sender,
1063 %view,
1064 "ignoring vid share fragment not addressed to this node"
1065 );
1066 return None;
1067 }
1068 let leader = fragment
1069 .data
1070 .epoch
1071 .and_then(|epoch| self.leader(view, epoch));
1072 if leader.as_ref() != Some(&message.sender) {
1073 warn!(
1074 %node,
1075 %sender,
1076 %view,
1077 "ignoring vid share fragment not from the view leader"
1078 );
1079 return None;
1080 }
1081 if self.consensus.wants_proposal_for_view(&view) {
1082 let signature = fragment.signature.clone();
1083 match self.vid_fragment_accumulator.accept(fragment.data) {
1084 Ok(Some(share)) => self
1085 .share_validator
1086 .validate(SignedProposal::new(share, signature)),
1087 Ok(None) => {}, Err(err) => {
1089 warn!(
1090 %node,
1091 %sender,
1092 %view,
1093 %err, "rejecting malformed vid share fragment"
1094 );
1095 },
1096 }
1097 }
1098 None
1099 },
1100 ConsensusMessage::Vote1(vote1) => {
1101 let view = vote1.vote.view_number();
1102 if self.is_view_too_far_ahead(view) {
1103 warn!(%node, %sender, %view, "vote1 is too far ahead");
1104 return None;
1105 }
1106 if vote1.vote.signing_key() != message.sender {
1107 warn!(%node, %sender, %view, "vote1 signing key != sender");
1108 return None;
1109 }
1110 let bn = vote1.vote.data.block_number.unwrap_or(0);
1111 let epoch_height = *self.consensus.epoch_height;
1112 let is_epoch_root_vote = is_epoch_root(bn, epoch_height);
1113 debug!(
1114 %node, %sender, %view,
1115 epoch_root = is_epoch_root_vote,
1116 has_state_vote = vote1.state_vote.is_some(),
1117 "recv vote1"
1118 );
1119 if is_epoch_root_vote {
1120 vote1.state_vote.as_ref()?;
1123 self.epoch_root_collector.accumulate_vote(vote1);
1124 } else {
1125 self.vote1_collector.accumulate_vote(vote1.vote);
1126 }
1127 None
1128 },
1129 ConsensusMessage::VidShareBroadcast(share) => {
1130 let view = share.view_number();
1131 if self.is_view_too_far_ahead(view) {
1132 warn!(%node, %sender, %view, "vid share broadcast is too far ahead");
1133 return None;
1134 }
1135 debug!(%node, %sender, %view, "recv vid share broadcast");
1136 self.vid_reconstructor
1140 .handle_vid_share(message.sender.clone(), share);
1141 None
1142 },
1143 ConsensusMessage::Vote2(vote2) => {
1144 let view = vote2.view_number();
1145 if self.is_view_too_far_ahead(view) {
1146 warn!(%node, %sender, %view, "vote2 is too far ahead");
1147 return None;
1148 }
1149 if vote2.signing_key() != message.sender {
1150 warn!(%node, %sender, %view, "vote2 signing key != sender");
1151 return None;
1152 }
1153 debug!(%node, %sender, %view, "recv vote2");
1154 self.vote2_collector.accumulate_vote(vote2);
1155 None
1156 },
1157 ConsensusMessage::Certificate1(certificate1, _key) => {
1158 let view = certificate1.view_number();
1159 debug!(
1160 %node, %sender, %view,
1161 epoch = ?certificate1.epoch().map(|e| *e),
1162 "recv certificate1"
1163 );
1164 if self.is_view_too_far_ahead(view) {
1165 warn!(%node, %sender, %view, "certificate1 is too far ahead");
1166 return None;
1167 }
1168 if self.is_epoch_too_far_ahead(certificate1.epoch()) {
1169 warn!(%node, %sender, %view, "certificate1 epoch is too far ahead");
1170 return None;
1171 }
1172 if let Some(epoch) = self
1173 .cert_verifiers
1174 .cert1
1175 .verify(message.sender, certificate1)
1176 {
1177 self.epoch_manager.request_drb_result(epoch);
1178 }
1179 None
1180 },
1181 ConsensusMessage::Certificate2(certificate2, _key) => {
1182 let view = certificate2.view_number();
1183 debug!(
1184 %node, %sender, %view,
1185 epoch = ?certificate2.epoch().map(|e| *e),
1186 "recv certificate2"
1187 );
1188 if self.is_view_too_far_ahead(view) {
1189 warn!(%node, %sender, %view, "certificate2 is too far ahead");
1190 return None;
1191 }
1192 if self.is_epoch_too_far_ahead(certificate2.epoch()) {
1193 warn!(%node, %sender, %view, "certificate2 epoch is too far ahead");
1194 return None;
1195 }
1196 if let Some(epoch) = self
1197 .cert_verifiers
1198 .cert2
1199 .verify(message.sender, certificate2)
1200 {
1201 self.epoch_manager.request_drb_result(epoch);
1202 }
1203 None
1204 },
1205 ConsensusMessage::TimeoutVote(timeout_msg) => {
1206 let view = timeout_msg.vote.view_number();
1207 if timeout_msg.vote.signing_key() != message.sender {
1208 warn!(%node, %sender, %view, "timeout vote signing key != sender");
1209 return None;
1210 }
1211 let current_view = self.consensus.current_view();
1212 let has_evidence = timeout_msg.evidence.is_some();
1213
1214 if let Some(e) = timeout_msg
1219 .evidence
1220 .filter(|e| e.view_number() >= current_view)
1221 {
1222 match e {
1223 CatchupEvidence::Qc(qc) => {
1224 if let Some(epoch) = self
1225 .cert_verifiers
1226 .advance
1227 .verify(message.sender.clone(), qc)
1228 {
1229 self.epoch_manager.request_drb_result(epoch);
1230 }
1231 },
1232 CatchupEvidence::Tc(tc) => {
1233 if let Some(epoch) = self
1234 .cert_verifiers
1235 .timeout
1236 .verify(message.sender.clone(), tc)
1237 {
1238 self.epoch_manager.request_drb_result(epoch);
1239 }
1240 },
1241 }
1242 }
1243
1244 if self.is_view_too_far_ahead(view) {
1245 warn!(%node, %sender, %view, "timeout vote is too far ahead");
1246 return None;
1247 }
1248
1249 if view < current_view {
1250 debug!(
1251 %node, %sender, %view, %current_view,
1252 "timeout vote for stale view; replying with catchup evidence"
1253 );
1254 self.send_catchup_evidence(&message.sender, view);
1255 return None;
1256 }
1257
1258 debug!(%node, %sender, %view, has_evidence, "recv timeout vote");
1259
1260 self.timeout_collector
1261 .accumulate_vote(timeout_msg.vote.clone());
1262 self.timeout_one_honest_collector
1263 .accumulate_vote(timeout_msg.vote);
1264
1265 None
1266 },
1267 ConsensusMessage::TimeoutCertificate(tc) => {
1268 debug!(
1269 %node, %sender,
1270 view = %tc.view_number(),
1271 epoch = ?tc.epoch().map(|e| *e),
1272 "recv timeout certificate"
1273 );
1274 if let Some(epoch) = self
1275 .cert_verifiers
1276 .timeout
1277 .verify(message.sender.clone(), tc)
1278 {
1279 self.epoch_manager.request_drb_result(epoch);
1280 }
1281 None
1282 },
1283 ConsensusMessage::HighQc(qc) => {
1284 debug!(
1285 %node, %sender,
1286 view = %qc.view_number(),
1287 epoch = ?qc.epoch().map(|e| *e),
1288 "recv high qc"
1289 );
1290 if let Some(epoch) = self
1291 .cert_verifiers
1292 .advance
1293 .verify(message.sender.clone(), qc)
1294 {
1295 self.epoch_manager.request_drb_result(epoch);
1296 }
1297 None
1298 },
1299 ConsensusMessage::EpochChange(epoch_change) => {
1300 let view = epoch_change.cert1.view_number();
1301 let epoch = epoch_change.cert1.epoch();
1302 debug!(%node, %sender, %view, epoch = ?epoch.map(|e| *e), "recv epoch change");
1303 if self.is_epoch_too_far_ahead(epoch) {
1304 warn!(%node, %sender, %view, ?epoch, "epoch change is too far ahead");
1305 return None;
1306 }
1307 if let Some(epoch) = self
1308 .cert_verifiers
1309 .epoch_change
1310 .verify(message.sender, epoch_change)
1311 {
1312 self.epoch_manager.request_drb_result(epoch);
1313 }
1314 None
1315 },
1316 },
1317 MessageType::Block(msg) => {
1318 match msg {
1319 BlockMessage::Transactions(msg) => {
1320 debug!(
1321 %node, %sender,
1322 view = %msg.view,
1323 count = msg.transactions.len(),
1324 "recv transactions"
1325 );
1326 self.block_builder.on_transactions(msg)
1327 },
1328 BlockMessage::DedupManifest(manifest) => {
1329 debug!(
1330 %node, %sender,
1331 view = %manifest.view,
1332 epoch = %manifest.epoch,
1333 hashes = manifest.hashes.len(),
1334 "recv dedup manifest"
1335 );
1336 if !self.is_view_too_far_ahead(manifest.view)
1337 && let Some(view_leader) = self.leader(manifest.view, manifest.epoch)
1338 && view_leader == message.sender
1339 {
1340 self.block_builder.on_dedup_manifest(manifest)
1341 }
1342 },
1343 }
1344 None
1345 },
1346 MessageType::ProposalFetch(ProposalFetchMessage::Request(request)) => {
1347 let view = request.view_number();
1348 debug!(%node, %sender, %view, "recv proposal fetch request");
1349 if !request.validate_sender(&message.sender) {
1350 warn!(
1351 %node,
1352 sender = %message.sender,
1353 %view,
1354 "ignoring invalid proposal fetch request signature"
1355 );
1356 return None;
1357 }
1358 if let Some(proposal) = self.consensus.signed_proposal(&view).cloned() {
1359 let response = Message {
1360 sender: self.public_key.clone(),
1361 message_type: MessageType::ProposalFetch(ProposalFetchMessage::Response(
1362 Box::new(proposal),
1363 )),
1364 };
1365 if let Err(err) = self.network.sender().unicast(
1366 self.consensus.current_view(),
1367 &message.sender,
1368 &response,
1369 ) {
1370 let err = CoordinatorError::from(err).context("proposal response");
1371 warn!(%node, %err, "network error while sending proposal response");
1372 }
1373 }
1374 None
1375 },
1376 MessageType::ProposalFetch(ProposalFetchMessage::Response(proposal)) => {
1377 debug!(
1378 %node, %sender,
1379 view = %proposal.data.view_number,
1380 "recv proposal fetch response"
1381 );
1382 self.pending_proposal_fetches.resolve(&proposal);
1383 self.maybe_validate_fetched_proposal(*proposal);
1384 None
1385 },
1386 MessageType::External(data) => {
1387 debug!(%node, %sender, bytes = data.len(), "recv external message");
1388 self.coordinator_outbox.push_back(OpaqueMessage {
1389 sender: message.sender,
1390 data,
1391 });
1392 None
1393 },
1394 }
1395 }
1396
1397 fn on_state_manager_output(
1398 &mut self,
1399 output: StateManagerOutput<T>,
1400 ) -> Option<ConsensusInput<T>> {
1401 match output {
1402 StateManagerOutput::State {
1403 response,
1404 validated: true,
1405 } => Some(ConsensusInput::StateValidated(response)),
1406 StateManagerOutput::State {
1407 response,
1408 validated: false,
1409 } => Some(ConsensusInput::StateValidationFailed(response)),
1410 StateManagerOutput::Header {
1411 response,
1412 header: Some(hdr),
1413 } => Some(ConsensusInput::HeaderCreated(
1414 response.view,
1415 proposal_commitment(&response.parent_proposal),
1416 hdr,
1417 )),
1418 StateManagerOutput::Header {
1419 response,
1420 header: None,
1421 } => {
1422 warn!(view = %response.view, "header creation failed");
1423 None
1424 },
1425 }
1426 }
1427
1428 fn expected_vid_param(&self, target_epoch: Option<EpochNumber>) -> Option<AvidmGf2Param> {
1433 let membership = self
1434 .membership_coordinator
1435 .stake_table_for_epoch(target_epoch)
1436 .ok()?;
1437 let total_weight = vid_total_weight::<T, _>(membership.stake_table(), target_epoch);
1438 init_avidm_gf2_param(total_weight).ok()
1439 }
1440
1441 fn broadcast(
1442 &self,
1443 message_type: ConsensusMessage<T, Validated>,
1444 ctx: &'static str,
1445 ) -> Result<(), CoordinatorError> {
1446 let message = Message {
1447 sender: self.public_key.clone(),
1448 message_type: MessageType::Consensus(message_type),
1449 };
1450 self.network
1451 .sender()
1452 .broadcast(self.consensus.current_view(), &message)
1453 .map_err(|e| CoordinatorError::from(e).context(ctx))
1454 }
1455
1456 fn unicast_to_leader(
1457 &mut self,
1458 view: ViewNumber,
1459 epoch: EpochNumber,
1460 msg: BlockMessage<T>,
1461 ) -> Result<(), CoordinatorError> {
1462 let Some(leader) = self.leader(view, epoch) else {
1463 warn!(%view, %epoch, "failed to resolve leader for unicast");
1464 return Ok(());
1465 };
1466 let message = Message {
1467 sender: self.public_key.clone(),
1468 message_type: MessageType::Block(msg),
1469 };
1470 self.network
1471 .sender()
1472 .unicast(self.consensus.current_view(), &leader, &message)
1473 .map_err(|e| CoordinatorError::from(e).context("leader unicast"))
1474 }
1475
1476 fn leader(&mut self, view: ViewNumber, epoch: EpochNumber) -> Option<T::SignatureKey> {
1477 let membership = self
1478 .membership_coordinator
1479 .membership_for_epoch(Some(epoch))
1480 .ok()?;
1481 membership.leader(view).ok()
1482 }
1483
1484 fn on_client_request(&mut self, request: ClientRequest<T>) -> Result<(), CoordinatorError> {
1485 match request {
1486 ClientRequest::CurrentView(tx) => {
1487 let _ = tx.send(self.consensus.current_view());
1488 },
1489 ClientRequest::CurrentEpoch(tx) => {
1490 let _ = tx.send(self.consensus.current_epoch());
1491 },
1492 ClientRequest::DecidedLeaf(tx) => {
1493 let _ = tx.send(self.consensus.last_decided_leaf().clone());
1494 },
1495 ClientRequest::DecidedState(tx) => {
1496 let view = self.consensus.last_decided_leaf().view_number();
1497 let _ = tx.send(self.state(view).map(|s| s.state.clone()));
1498 },
1499 ClientRequest::UndecidedLeaves(tx) => {
1500 let _ = tx.send(self.consensus.undecided_leaves().cloned().collect());
1501 },
1502 ClientRequest::GetState { view, respond } => {
1503 let _ = respond.send(self.state(view).map(|s| s.state.clone()));
1504 },
1505 ClientRequest::GetStateAndDelta { view, respond } => {
1506 let _ = respond.send(match self.state(view) {
1507 Some(s) => (Some(s.state.clone()), s.delta.clone()),
1508 None => (None, None),
1509 });
1510 },
1511 ClientRequest::ProposalParticipation { epoch, respond } => {
1512 let _ = respond.send(match epoch {
1513 Some(epoch) => self.participation.proposal_participation(epoch),
1514 None => self.participation.current_proposal_participation(),
1515 });
1516 },
1517 ClientRequest::VoteParticipation { epoch, respond } => {
1518 let _ = respond.send(match epoch {
1519 Some(epoch) => self.participation.vote_participation(epoch),
1520 None => self.participation.current_vote_participation(),
1521 });
1522 },
1523 ClientRequest::SubmitTransaction { tx, respond } => {
1524 self.block_builder.on_submit_transaction(tx);
1525 let _ = respond.send(());
1526 },
1527 ClientRequest::UpdateLeaf { update, respond } => {
1528 self.state_manager.update_state(update);
1529 let _ = respond.send(());
1530 },
1531 ClientRequest::RequestProposal {
1532 view,
1533 leaf_commitment,
1534 respond,
1535 } => {
1536 if let Some(proposal) = self.consensus.signed_proposal(&view)
1537 && proposal_commitment(&proposal.data) == leaf_commitment
1538 {
1539 let _ = respond.send(Ok(proposal.clone()));
1540 return Ok(());
1541 }
1542 if !self
1543 .pending_proposal_fetches
1544 .contains_request(view, leaf_commitment)
1545 {
1546 self.broadcast_proposal_fetch(view)?;
1547 }
1548 self.pending_proposal_fetches
1549 .push(view, leaf_commitment, respond);
1550 },
1551 ClientRequest::SendExternalMessage {
1552 payload,
1553 recipient,
1554 respond,
1555 } => {
1556 let message = Message {
1557 sender: self.public_key.clone(),
1558 message_type: MessageType::External(payload),
1559 };
1560 let result = self
1561 .network
1562 .sender()
1563 .unicast(self.consensus.current_view(), &recipient, &message)
1564 .map_err(|err| {
1565 CoordinatorError::from(err)
1566 .context("send external message")
1567 .into()
1568 });
1569 let _ = respond.send(result);
1570 },
1571 ClientRequest::SubmitTimeoutVote { vote } => {
1572 let view = vote.view_number();
1573 let current_view = self.consensus.current_view();
1574 if view < current_view {
1575 debug!(
1576 %view, %current_view,
1577 "ignoring bridged timeout vote for stale view"
1578 );
1579 return Ok(());
1580 }
1581 self.timeout_collector.accumulate_vote(vote.clone());
1582 self.timeout_one_honest_collector
1583 .accumulate_vote(vote.clone());
1584 let message = Message {
1586 sender: self.public_key.clone(),
1587 message_type: MessageType::Consensus(ConsensusMessage::TimeoutVote(
1588 message::TimeoutVoteMessage {
1589 vote,
1590 evidence: None,
1591 },
1592 )),
1593 };
1594 if let Err(err) = self
1595 .network
1596 .sender()
1597 .broadcast(self.consensus.current_view(), &message)
1598 {
1599 warn!(%err, "failed to rebroadcast bridged timeout vote");
1600 }
1601 },
1602 ClientRequest::SubmitLegacyHighQc { qc } => {
1603 let qc_view = qc.view_number();
1607 let cutover_view = qc_view + 1;
1608 self.consensus.register_legacy_qc(&qc);
1609
1610 let cur_view = self.consensus.current_view();
1616 if cur_view == qc_view
1617 && self.consensus.timeout_cert_at(cutover_view).is_none()
1618 && self.consensus.cert1_at(qc_view).is_some()
1619 && self.consensus.proposal_at(qc_view).is_some()
1620 {
1621 info!(
1622 %cutover_view,
1623 "bridged late legacy high QC; proposing cutover view on it (no timeout)"
1624 );
1625 self.start(None);
1626 while let Some(output) = self.outbox.pop_front() {
1627 if let Err(err) = self.process_consensus_output(output) {
1628 warn!(
1629 %err,
1630 "error processing bridged-high-qc bootstrap output"
1631 );
1632 }
1633 }
1634 }
1635 },
1636 }
1637
1638 Ok(())
1639 }
1640
1641 fn record_voted_view(&mut self, view: ViewNumber) {
1642 if self.voted_view.is_some_and(|v| v >= view) {
1643 return;
1644 }
1645 self.voted_view = Some(view);
1646 if let Some(m) = &self.metrics {
1647 m.consensus.last_voted_view.set(*view as usize);
1648 }
1649 }
1650
1651 fn on_view_changed_metrics(&mut self, view: ViewNumber, epoch: EpochNumber) {
1652 if let Some((started_view, started_epoch, started_at)) = self.view_started
1653 && started_view == view
1654 {
1655 if epoch > started_epoch {
1656 self.view_started = Some((view, epoch, started_at));
1657 }
1658 return;
1659 }
1660 let prev = self.view_started.replace((view, epoch, Instant::now()));
1661 let duration_as_leader = prev.and_then(|(prev_view, prev_epoch, entered)| {
1662 (self.leader(prev_view, prev_epoch).as_ref() == Some(&self.public_key))
1663 .then(|| entered.elapsed())
1664 });
1665 let Some(m) = &self.metrics else { return };
1666 let consensus = &m.consensus;
1667 consensus.current_view.set(*view as usize);
1668 let invalid_certs = self
1669 .cert_verifiers
1670 .num_invalid_certs()
1671 .saturating_sub(self.invalid_certs_at_decide);
1672 consensus.invalid_qc.set(invalid_certs as usize);
1673 let last_decided_view = self.consensus.last_decided_view();
1674 if view > last_decided_view {
1675 consensus
1676 .number_of_views_since_last_decide
1677 .set((*view - *last_decided_view) as usize);
1678 }
1679 if let Some(duration) = duration_as_leader {
1680 consensus
1681 .view_duration_as_leader
1682 .add_point(duration.as_secs_f64());
1683 }
1684 let (outstanding_txns, outstanding_bytes) = self.block_builder.outstanding_transactions();
1685 consensus.outstanding_transactions.set(outstanding_txns);
1686 consensus
1687 .outstanding_transactions_memory_size
1688 .set(outstanding_bytes);
1689 }
1690
1691 fn on_decide_metrics(&mut self, leaves: &[Leaf2<T>]) {
1692 let Some(newest) = leaves.first() else { return };
1693 let advanced = newest.view_number() == self.consensus.last_decided_view();
1697 if advanced {
1698 self.invalid_certs_at_decide = self.cert_verifiers.num_invalid_certs();
1699 }
1700 let Some(m) = &self.metrics else { return };
1701 let consensus = &m.consensus;
1702 let now = OffsetDateTime::now_utc().unix_timestamp();
1703 for leaf in leaves {
1704 let txn_bytes = self
1705 .payload_txn_bytes
1706 .get(&leaf.view_number())
1707 .copied()
1708 .or_else(|| leaf.block_payload_ref().map(|p| p.txn_bytes()));
1709 if let Some(txn_bytes) = txn_bytes {
1710 consensus.finalized_bytes.add_point(txn_bytes as f64);
1711 }
1712 if advanced {
1713 match (now as u64).checked_sub(leaf.block_header().timestamp()) {
1714 Some(age) => consensus.proposal_to_decide_time.add_point(age as f64),
1715 None => error!(
1716 timestamp = leaf.block_header().timestamp(),
1717 "failed to calculate proposal to decide time: timestamp in the future"
1718 ),
1719 }
1720 }
1721 }
1722 if !advanced {
1723 return;
1724 }
1725 consensus.last_decided_time.set(now as usize);
1726 consensus.invalid_qc.set(0);
1727 consensus
1728 .last_decided_view
1729 .set(*newest.view_number() as usize);
1730 consensus
1731 .last_synced_block_height
1732 .set(newest.block_header().block_number() as usize);
1733 if let Some(views_in_flight) =
1734 (*self.consensus.current_view()).checked_sub(*newest.view_number())
1735 {
1736 consensus
1737 .number_of_views_per_decide_event
1738 .add_point(views_in_flight as f64);
1739 }
1740 }
1741
1742 fn broadcast_proposal_fetch(&mut self, view: ViewNumber) -> Result<(), CoordinatorError> {
1744 let request = self
1745 .consensus
1746 .signed_proposal_fetch_request(view)
1747 .map_err(|err| {
1748 let err = format!("failed to sign proposal request: {err}");
1749 CoordinatorError::regular(err).context("sign proposal request")
1750 })?;
1751 let message = Message {
1752 sender: self.public_key.clone(),
1753 message_type: MessageType::ProposalFetch(ProposalFetchMessage::Request(request)),
1754 };
1755 self.network
1756 .sender()
1757 .broadcast(self.consensus.current_view(), &message)
1758 .map_err(|err| CoordinatorError::from(err).context("broadcast proposal request"))
1759 }
1760
1761 fn request_missing_proposal(
1762 &mut self,
1763 view: ViewNumber,
1764 leaf_commit: Commitment<Leaf2<T>>,
1765 ) -> Result<(), CoordinatorError> {
1766 if !self
1767 .requested_missing_proposals
1768 .insert(ProposalFetchKey::new(view, leaf_commit))
1769 {
1770 return Ok(());
1771 }
1772 self.broadcast_proposal_fetch(view)
1773 }
1774
1775 fn maybe_validate_fetched_proposal(&mut self, proposal: SignedProposal<T, Proposal<T>>) {
1776 let view = proposal.data.view_number;
1777 let key = ProposalFetchKey::new(view, proposal_commitment(&proposal.data));
1778 if !self.requested_missing_proposals.remove(&key) {
1779 return;
1780 }
1781 if self.consensus.proposal_at(view).is_some() {
1782 return;
1783 }
1784 self.proposal_validator
1785 .validate_fetched(ProposalMessage::unchecked(proposal));
1786 }
1787
1788 fn gc(&mut self, epoch: EpochNumber, scope: GcScope) -> Result<(), CoordinatorError> {
1789 self.consensus.gc(scope);
1790 match scope {
1791 GcScope::Local(view) => {
1792 self.block_builder.gc(view);
1793 self.vid_disperser.gc(view);
1794 self.vid_fragment_accumulator.gc(view);
1795 let view = view.saturating_sub(1).into();
1798 self.network.gc(view)?;
1799 self.timeout_collector.gc(view);
1800 self.timeout_one_honest_collector.gc(view);
1801 self.vote1_collector.gc(view);
1802 self.vote2_collector.gc(view);
1803 },
1804 GcScope::Decided(view) => {
1805 let decide_floor = self.consensus.decide_floor();
1806 self.epoch_manager.gc(epoch);
1807 self.epoch_root_collector.gc(view);
1808 self.cert_verifiers.gc(decide_floor, epoch);
1809 self.pending_proposal_fetches.gc(view);
1810 self.requested_missing_proposals
1811 .retain(|key| key.view > view);
1812 self.state_manager.gc(view);
1813 self.storage
1814 .gc(view.saturating_sub(STORAGE_GC_MARGIN).into());
1815 self.vid_reconstructor
1816 .gc(view.saturating_sub(VID_RECONSTRUCT_GC_MARGIN).into());
1817 let vc = VidCommitment2::default();
1818 self.da_payloads = self.da_payloads.split_off(&(view, vc));
1819 self.payload_txn_bytes = self.payload_txn_bytes.split_off(&(decide_floor + 1));
1820 },
1821 GcScope::Timeout(view) => {
1822 self.vid_reconstructor.retire_view(view);
1823 let vc = VidCommitment2::default();
1824 self.da_payloads
1825 .extract_if((view, vc)..(view + 1, vc), |_, _| true)
1826 .for_each(drop);
1827 },
1828 }
1829 Ok(())
1830 }
1831
1832 fn apply_cutover_seed(&mut self, seed: PreCutoverSeed<T>) -> bool {
1839 let current_view = self.consensus.current_view();
1840 if seed.cutover_view > ViewNumber::genesis() && current_view >= seed.cutover_view {
1841 info!(
1842 node = %self.node_id,
1843 %current_view,
1844 cutover_view = *seed.cutover_view,
1845 "ignoring pre-cutover seed; already past the cutover",
1846 );
1847 return true;
1848 }
1849 info!(
1850 node = %self.node_id,
1851 undecided = seed.undecided.len(),
1852 anchor_view = *seed.decided_anchor.view_number(),
1853 high_qc_view = seed.high_qc.as_ref().map(|qc| *qc.view_number()),
1854 cutover_view = *seed.cutover_view,
1855 states = seed.validated_states.len(),
1856 "applying legacy -> new-protocol seed",
1857 );
1858
1859 let anchor_view = seed.decided_anchor.view_number();
1863 if let Some(state) = seed.validated_states.get(&anchor_view).cloned() {
1864 self.state_manager
1865 .seed_state(anchor_view, state, seed.decided_anchor.clone());
1866 }
1867 for leaf in &seed.undecided {
1868 let view = leaf.view_number();
1869 if let Some(state) = seed.validated_states.get(&view).cloned() {
1870 self.state_manager.seed_state(view, state, leaf.clone());
1871 }
1872 }
1873
1874 let highest_seeded_leaf = seed.undecided.last().unwrap_or(&seed.decided_anchor);
1875 let cutover_epoch = EpochNumber::new(epoch_from_block_number(
1876 highest_seeded_leaf.block_header().block_number(),
1877 *self.consensus.epoch_height,
1878 ));
1879 let cutover_view = seed.cutover_view;
1880
1881 self.consensus.apply_pre_cutover_seed(seed);
1882
1883 if let Err(err) = self
1886 .network
1887 .apply_epoch(cutover_epoch, &self.membership_coordinator)
1888 {
1889 error!(
1890 %cutover_epoch,
1891 %err,
1892 "network on_epoch_change failed while applying the cutover seed",
1893 );
1894 }
1895
1896 let cur_view = self.consensus.current_view();
1897 if cur_view + 1 == cutover_view
1898 && self.consensus.cert1_at(cur_view).is_some()
1899 && self.consensus.proposal_at(cur_view).is_some()
1900 {
1901 return true;
1902 }
1903 let epoch = self
1904 .consensus
1905 .current_epoch()
1906 .unwrap_or(EpochNumber::genesis());
1907 self.outbox
1908 .push_back(ConsensusOutput::ViewChanged(cur_view, epoch));
1909 false
1910 }
1911
1912 fn is_view_too_far_ahead(&self, v: ViewNumber) -> bool {
1914 v > self.consensus.current_view() + *MAX_VIEWS_AHEAD
1915 }
1916
1917 fn is_epoch_too_far_ahead(&self, epoch: Option<EpochNumber>) -> bool {
1919 let current = self
1920 .consensus
1921 .current_epoch()
1922 .unwrap_or(EpochNumber::genesis());
1923 epoch.is_some_and(|e| e > current + EPOCH_CHANGE_LOOKAHEAD)
1924 }
1925
1926 pub(crate) fn catchup_evidence(&self) -> Option<ConsensusMessage<T, Validated>> {
1927 Some(match self.consensus.catchup_evidence()? {
1928 CatchupEvidence::Qc(qc) => ConsensusMessage::HighQc(qc),
1929 CatchupEvidence::Tc(tc) => ConsensusMessage::TimeoutCertificate(tc),
1930 })
1931 }
1932
1933 fn send_catchup_evidence(&mut self, peer: &T::SignatureKey, stale_view: ViewNumber) {
1934 let Some(evidence) = self.catchup_evidence() else {
1935 return;
1936 };
1937 if evidence.view_number() < stale_view {
1938 return;
1939 }
1940 let message = Message {
1941 sender: self.public_key.clone(),
1942 message_type: MessageType::Consensus(evidence),
1943 };
1944 if let Err(err) =
1945 self.network
1946 .sender()
1947 .unicast(self.consensus.current_view(), peer, &message)
1948 {
1949 warn!(%stale_view, %err, "failed to send catchup evidence");
1950 }
1951 }
1952}
1953
1954#[derive(Debug, Clone, Copy)]
1956pub enum GcScope {
1957 Local(ViewNumber),
1959 Decided(ViewNumber),
1961 Timeout(ViewNumber),
1963}
1964
1965struct PendingDa<T: NodeType> {
1967 epoch: EpochNumber,
1968 payload: T::BlockPayload,
1969 metadata: <T::BlockPayload as BlockPayload<T>>::Metadata,
1970}
1971
1972type ProposalFetchResponseSender<T> =
1973 oneshot::Sender<Result<SignedProposal<T, Proposal<T>>, QueryError>>;
1974
1975#[derive(Clone, Copy, Debug, Eq, Hash, PartialEq)]
1976struct ProposalFetchKey<T: NodeType> {
1977 view: ViewNumber,
1978 leaf_commitment: Commitment<Leaf2<T>>,
1979}
1980
1981impl<T: NodeType> ProposalFetchKey<T> {
1982 fn new(view: ViewNumber, leaf_commitment: Commitment<Leaf2<T>>) -> Self {
1983 Self {
1984 view,
1985 leaf_commitment,
1986 }
1987 }
1988}
1989
1990#[derive(Default)]
1991struct PendingProposalFetches<T: NodeType> {
1992 pending: HashMap<ProposalFetchKey<T>, Vec<ProposalFetchResponseSender<T>>>,
1993}
1994
1995impl<T: NodeType> PendingProposalFetches<T> {
1996 fn prune_closed(&mut self) {
1997 self.pending.retain(|_, responders| {
1998 responders.retain(|respond| !respond.is_closed());
1999 !responders.is_empty()
2000 });
2001 }
2002
2003 fn contains_request(
2004 &mut self,
2005 view: ViewNumber,
2006 leaf_commitment: Commitment<Leaf2<T>>,
2007 ) -> bool {
2008 self.prune_closed();
2009 self.pending
2010 .contains_key(&ProposalFetchKey::new(view, leaf_commitment))
2011 }
2012
2013 fn push(
2014 &mut self,
2015 view: ViewNumber,
2016 leaf_commitment: Commitment<Leaf2<T>>,
2017 respond: ProposalFetchResponseSender<T>,
2018 ) {
2019 self.pending
2020 .entry(ProposalFetchKey::new(view, leaf_commitment))
2021 .or_default()
2022 .push(respond);
2023 }
2024
2025 fn gc(&mut self, view: ViewNumber) {
2026 self.pending.retain(|key, responders| {
2027 responders.retain(|respond| !respond.is_closed());
2028 key.view >= view && !responders.is_empty()
2029 });
2030 }
2031
2032 fn resolve(&mut self, proposal: &SignedProposal<T, Proposal<T>>) {
2033 self.prune_closed();
2034 let view = proposal.data.view_number;
2035 let leaf_commitment = proposal_commitment(&proposal.data);
2036 let key = ProposalFetchKey::new(view, leaf_commitment);
2037
2038 if let Some(responders) = self.pending.remove(&key) {
2039 for respond in responders {
2040 let _ = respond.send(Ok(proposal.clone()));
2041 }
2042 }
2043 }
2044}