Skip to main content

hotshot_new_protocol/
coordinator.rs

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
62/// Views to retain in the VID reconstructor behind the decided view
63///
64/// A decide can land while an earlier view's payload is still being
65/// reconstructed, and GC at the decided view would abort that task.
66/// A decide proves a quorum reconstructed the payload
67/// so it can be fetched later assuming the quorum includes at
68/// least one query node serving catchup.
69/// The margin gives in flight reconstruction tasks time to finish, which is
70/// cheaper than fetching the payload through catchup.
71///
72/// Proposals are retained with the same margin: when a reconstruction
73/// finishes, `BlockPayloadReconstructed` is only emitted if the proposal
74/// (the block header) for that view is still available.
75pub(crate) const VID_RECONSTRUCT_GC_MARGIN: u64 = 5;
76
77/// Views to retain in storage GC behind the decided view, so in-flight
78/// storage writes for recent views aren't aborted before they persist.
79const STORAGE_GC_MARGIN: u64 = 5;
80
81/// Epoch changes claiming an epoch further ahead than this are dropped at
82/// intake. We could not verify them anyway: an epoch's stake table only
83/// materializes by walking the DRB chain, so parking such a message and
84/// driving catchup for an arbitrary claimed epoch just burns resources.
85/// Within the ceiling, deferred changes verify progressively as catchup
86/// advances ([`CertVerifiers::retry_pending`] runs on every DRB arrival).
87const 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    /// View of the last timer fire that counted towards participation and
135    /// timeout metrics, so a re-fire for the same (stuck) view doesn't
136    /// double-count it.
137    #[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 persisted on a prior run; restored so the lock survives restart.
171        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        // Seed `from_header` stubs for restored undecided proposals so a child
228        // proposal can be validated; anchor seeded last so its state wins.
229        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        // The anchor leaf and persisted proposals are blocks this node had
238        // reconstructed before it went down, so treat them as reconstructed on
239        // restart
240        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        // Seed every persisted proposal before `seed_parent` so its authoritative anchor wins.
253        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        // `seed_parent` sets the current epoch from the anchor proposal;
259        // `resume_from_restart` positions the view so the node never
260        // re-enters a view it may have voted or proposed in before it went
261        // down.
262        consensus.seed_parent(cert1, parent_proposal, reconstructed_blocks);
263        // Restore the persisted lock; it can be newer than the anchor QC, so
264        // this must run after `seed_parent`.
265        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    /// Emit `ViewChanged(current_view + 1)` and, if leader, a
352    /// `RequestBlockAndHeader`.
353    ///
354    /// A pre-cutover `seed` is applied first, so the coordinator starts
355    /// from the bridged legacy state instead of genesis. When the seed
356    /// carries no QC for the last legacy view, the coordinator parks on
357    /// that view instead of proposing; a bridged high QC or a timeout
358    /// advances it.
359    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            // Genesis DA never flows through the normal block-builder path.
375            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            // Emit `LeafDecided` for genesis so persistence sees the header.
386            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            // No parent proposal when restarting past the anchor view: the
405            // node cannot propose off the anchor for a later view; the
406            // timeout path takes over instead.
407            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                    // Re-arm for the same view: a node stuck exactly at TC2
441                    // threshold can lose its only timeout-vote broadcast, so
442                    // the vote is re-sent every timeout period until the
443                    // view advances.
444                    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                        // Refresh the network's peer set when a proposal is validated.
537                        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                        // Retain the payload and persist it when consensus proposes this
552                        // exact block (cf. SendProposal):
553                        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                        // We built this block; skip reconstructing it from our own loopback share.
566                        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                            // Only pair the payload with the header if the proposal commits to it
599                            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                        // Catchup/compute failed. The epoch manager clears
637                        // the pending guard; consensus's `maybe_propose`
638                        // will re-request the DRB when it next tries to
639                        // build a transition proposal and finds it missing.
640                        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                // `leaves` is ordered newest first.
709                //  Garbage collect the data for views < decided view
710                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                // Two blocks can be built for one view. Here we know which one
753                // wins and we persist just that one:
754                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                // Proactively fetch the DRB for the next epoch so
968                // late-starting nodes have it before they need it
969                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    /// Refresh the network's peer window for `epoch`.
1020    ///
1021    /// The coordinator does this itself whenever a proposal validates, but
1022    /// before its event loop is started callers can trigger this explicitly
1023    /// to keep the network up to date.
1024    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) => {}, // Still missing some fragments.
1088                            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                        // An epoch-root Vote1 MUST carry a state_vote.
1121                        // Reject otherwise.
1122                        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                    // The share belongs to the sender (`recipient_key == sender`,
1137                    // enforced by `handle_vid_share`); it is verified lazily
1138                    // against the pinned commitment at reconstruction time.
1139                    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 a peer times out in a view at or ahead of us we adopt its
1215                    // highest certificate. Evidence takes precedence over
1216                    // the vote being too far ahead, and a valid certificate proves
1217                    // the network reached that view.
1218                    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    /// The VID erasure parameters the committee fixes for `target_epoch`,
1429    /// matching what an honest disperser derives. Used to reject shares whose
1430    /// `common.param` is forged (the commitment binds `ns_commits`, not
1431    /// `param`). `None` if the committee cannot be resolved.
1432    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                // Rebroadcast so peer coordinators can aggregate too.
1585                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                // QC certifies the last legacy view; cutover view is the next.
1604                // Register idempotently so the smooth-start precondition holds
1605                // regardless of arrival order vs. the cutover seed.
1606                let qc_view = qc.view_number();
1607                let cutover_view = qc_view + 1;
1608                self.consensus.register_legacy_qc(&qc);
1609
1610                // Still parked on the last legacy view (seed landed without this
1611                // QC, waiting out the timer) and not yet skipped via TC2: propose
1612                // the cutover view on the real QC now. Self-idempotent — once
1613                // started, `cur_view` advances past `qc_view` and `maybe_propose`
1614                // dedups by `proposed_views`.
1615                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        // The consensus watermark already includes this batch, so the batch
1694        // advanced the decide frontier iff its newest view is the watermark;
1695        // a gap-fill decide of older views must not regress the gauges.
1696        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    /// Broadcast a signed proposal fetch request for `view` to all peers.
1743    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                // When we enter a new view, we do not want to GC certain data
1796                // for the previous view yet:
1797                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    /// Bridge legacy state into the coordinator before it starts.
1833    ///
1834    /// Returns `false` when the coordinator must park on the last seeded
1835    /// view because the cutover view cannot be proposed off it yet. A
1836    /// stale seed (the coordinator has already restarted past the
1837    /// cutover) is ignored and the normal start path proceeds.
1838    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        // State manager is owned by the coordinator, so the
1860        // validated-state map must be applied here before the
1861        // seed is consumed by consensus.
1862        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        // Refresh peers for the cutover epoch before kicking the
1884        // leader — the proposal-driven site can't fire yet.
1885        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    /// We ignore votes more than `MAX_VIEWS_AHEAD` ahead of ours.
1913    fn is_view_too_far_ahead(&self, v: ViewNumber) -> bool {
1914        v > self.consensus.current_view() + *MAX_VIEWS_AHEAD
1915    }
1916
1917    /// We ignore certificates more than `EPOCH_CHANGE_LOOKAHEAD` ahead of ours.
1918    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/// Garbage collection scope.
1955#[derive(Debug, Clone, Copy)]
1956pub enum GcScope {
1957    /// GC is invoked on local view changes.
1958    Local(ViewNumber),
1959    /// GC is invoked on local decided views.
1960    Decided(ViewNumber),
1961    /// GC is invoked on a view that advanced via timeout certificate.
1962    Timeout(ViewNumber),
1963}
1964
1965/// A payload built locally and awaiting DA persistence.
1966struct 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}