Skip to main content

hotshot_task_impls/
stats.rs

1use std::{
2    collections::{BTreeMap, BTreeSet},
3    sync::Arc,
4};
5
6use async_broadcast::{Receiver, Sender};
7use async_trait::async_trait;
8use either::Either;
9use hotshot_task::task::TaskState;
10use hotshot_types::{
11    consensus::OuterConsensus,
12    data::{EpochNumber, ViewNumber},
13    epoch_membership::EpochMembershipCoordinator,
14    traits::{BlockPayload, block_contents::BlockHeader, node_implementation::NodeType},
15    vote::HasViewNumber,
16};
17use hotshot_utils::{anytrace::Result, warn};
18use serde::{Deserialize, Serialize};
19use time::OffsetDateTime;
20
21use crate::events::HotShotEvent;
22
23#[derive(Serialize, Deserialize)]
24pub struct LeaderViewStats {
25    pub view: ViewNumber,
26    pub prev_proposal_send: Option<i128>,
27    pub proposal_send: Option<i128>,
28    pub vote_recv: Option<i128>,
29    pub da_proposal_send: Option<i128>,
30    pub builder_start: Option<i128>,
31    pub block_built: Option<i128>,
32    pub vid_disperse_send: Option<i128>,
33    pub timeout_certificate_formed: Option<i128>,
34    pub qc_formed: Option<i128>,
35    pub da_cert_send: Option<i128>,
36}
37
38#[derive(Serialize, Deserialize)]
39pub struct ReplicaViewStats {
40    pub view: ViewNumber,
41    pub view_change: Option<i128>,
42    pub proposal_timestamp: Option<i128>,
43    pub proposal_recv: Option<i128>,
44    pub vote_send: Option<i128>,
45    pub timeout_vote_send: Option<i128>,
46    pub da_proposal_received: Option<i128>,
47    pub da_proposal_validated: Option<i128>,
48    pub da_certificate_recv: Option<i128>,
49    pub proposal_prelim_validated: Option<i128>,
50    pub proposal_validated: Option<i128>,
51    pub timeout_triggered: Option<i128>,
52    pub vid_share_validated: Option<i128>,
53    pub vid_share_recv: Option<i128>,
54}
55
56impl LeaderViewStats {
57    fn new(view: ViewNumber) -> Self {
58        Self {
59            view,
60            prev_proposal_send: None,
61            proposal_send: None,
62            vote_recv: None,
63            da_proposal_send: None,
64            builder_start: None,
65            block_built: None,
66            vid_disperse_send: None,
67            timeout_certificate_formed: None,
68            qc_formed: None,
69            da_cert_send: None,
70        }
71    }
72}
73
74impl ReplicaViewStats {
75    fn new(view: ViewNumber) -> Self {
76        Self {
77            view,
78            view_change: None,
79            proposal_timestamp: None,
80            proposal_recv: None,
81            vote_send: None,
82            timeout_vote_send: None,
83            da_proposal_received: None,
84            da_proposal_validated: None,
85            da_certificate_recv: None,
86            proposal_prelim_validated: None,
87            proposal_validated: None,
88            timeout_triggered: None,
89            vid_share_validated: None,
90            vid_share_recv: None,
91        }
92    }
93}
94
95pub struct StatsTaskState<TYPES: NodeType> {
96    view: ViewNumber,
97    epoch: Option<EpochNumber>,
98    public_key: TYPES::SignatureKey,
99    consensus: OuterConsensus<TYPES>,
100    membership_coordinator: EpochMembershipCoordinator<TYPES>,
101    leader_stats: BTreeMap<ViewNumber, LeaderViewStats>,
102    replica_stats: BTreeMap<ViewNumber, ReplicaViewStats>,
103    latencies_by_view: BTreeMap<ViewNumber, i128>,
104    sizes_by_view: BTreeMap<ViewNumber, i128>,
105    epoch_start_times: BTreeMap<EpochNumber, i128>,
106    timeouts: BTreeSet<ViewNumber>,
107}
108
109impl<TYPES: NodeType> StatsTaskState<TYPES> {
110    pub fn new(
111        view: ViewNumber,
112        epoch: Option<EpochNumber>,
113        public_key: TYPES::SignatureKey,
114        consensus: OuterConsensus<TYPES>,
115        membership_coordinator: EpochMembershipCoordinator<TYPES>,
116    ) -> Self {
117        Self {
118            view,
119            epoch,
120            public_key,
121            consensus,
122            membership_coordinator,
123            leader_stats: BTreeMap::new(),
124            replica_stats: BTreeMap::new(),
125            latencies_by_view: BTreeMap::new(),
126            sizes_by_view: BTreeMap::new(),
127            epoch_start_times: BTreeMap::new(),
128            timeouts: BTreeSet::new(),
129        }
130    }
131    fn leader_entry(&mut self, view: ViewNumber) -> &mut LeaderViewStats {
132        self.leader_stats
133            .entry(view)
134            .or_insert_with(|| LeaderViewStats::new(view))
135    }
136    fn replica_entry(&mut self, view: ViewNumber) -> &mut ReplicaViewStats {
137        self.replica_stats
138            .entry(view)
139            .or_insert_with(|| ReplicaViewStats::new(view))
140    }
141    fn garbage_collect(&mut self, view: ViewNumber) {
142        self.leader_stats = self.leader_stats.split_off(&view);
143        self.replica_stats = self.replica_stats.split_off(&view);
144        self.latencies_by_view = self.latencies_by_view.split_off(&view);
145        self.sizes_by_view = self.sizes_by_view.split_off(&view);
146        self.timeouts = BTreeSet::new();
147    }
148
149    fn dump_stats(&self) -> Result<()> {
150        let mut writer = csv::Writer::from_writer(vec![]);
151        for leader_stats in self.leader_stats.values() {
152            writer
153                .serialize(leader_stats)
154                .map_err(|e| warn!("Failed to serialize leader stats: {}", e))?;
155        }
156        let output = writer
157            .into_inner()
158            .map_err(|e| warn!("Failed to serialize replica stats: {}", e))?;
159        tracing::warn!(
160            "Leader stats: {}",
161            String::from_utf8(output)
162                .map_err(|e| warn!("Failed to convert leader stats to string: {}", e))?
163        );
164        let mut writer = csv::Writer::from_writer(vec![]);
165        for replica_stats in self.replica_stats.values() {
166            writer
167                .serialize(replica_stats)
168                .map_err(|e| warn!("Failed to serialize replica stats: {}", e))?;
169        }
170        let output = writer
171            .into_inner()
172            .map_err(|e| warn!("Failed to serialize replica stats: {}", e))?;
173        tracing::warn!(
174            "Replica stats: {}",
175            String::from_utf8(output)
176                .map_err(|e| warn!("Failed to convert replica stats to string: {}", e))?
177        );
178        Ok(())
179    }
180
181    fn log_basic_stats(&self, now: i128, epoch: &EpochNumber) {
182        let num_views = self.latencies_by_view.len();
183        let total_size = self.sizes_by_view.values().sum::<i128>();
184
185        // Either we have no views logged yet, no TXNs or we are not in the DA committee and don't know block sizes
186        if num_views == 0 || total_size == 0 {
187            return;
188        }
189
190        let total_latency = self.latencies_by_view.values().sum::<i128>();
191        let average_latency = total_latency / num_views as i128;
192        tracing::warn!("Average latency: {}ms", average_latency);
193        tracing::warn!(
194            "Number of timeouts in epoch: {}, is {}",
195            epoch,
196            self.timeouts.len()
197        );
198        if let Some(epoch_start_time) = self.epoch_start_times.get(epoch) {
199            let elapsed_time = now - epoch_start_time;
200            // multiply by 1000 to convert to seconds
201            let throughput = (total_size / elapsed_time) * 1000;
202            tracing::warn!("Throughput: {} bytes/s", throughput);
203        }
204    }
205}
206
207#[async_trait]
208impl<TYPES: NodeType> TaskState for StatsTaskState<TYPES> {
209    type Event = HotShotEvent<TYPES>;
210
211    async fn handle_event(
212        &mut self,
213        event: Arc<Self::Event>,
214        _sender: &Sender<Arc<Self::Event>>,
215        _receiver: &Receiver<Arc<Self::Event>>,
216    ) -> Result<()> {
217        let now = OffsetDateTime::now_utc().unix_timestamp_nanos();
218
219        match event.as_ref() {
220            HotShotEvent::BlockRecv(block_recv) => {
221                self.leader_entry(block_recv.view_number).block_built = Some(now);
222            },
223            HotShotEvent::QuorumProposalRecv(proposal, _) => {
224                self.replica_entry(proposal.data.view_number())
225                    .proposal_recv = Some(now);
226            },
227            HotShotEvent::QuorumVoteRecv(_vote) => {},
228            HotShotEvent::TimeoutVoteRecv(_vote) => {},
229            HotShotEvent::TimeoutVoteSend(vote) => {
230                self.replica_entry(vote.view_number()).timeout_vote_send = Some(now);
231            },
232            HotShotEvent::DaProposalRecv(proposal, _) => {
233                self.replica_entry(proposal.data.view_number())
234                    .da_proposal_received = Some(now);
235            },
236            HotShotEvent::DaProposalValidated(proposal, _) => {
237                self.replica_entry(proposal.data.view_number())
238                    .da_proposal_validated = Some(now);
239            },
240            HotShotEvent::DaVoteRecv(_simple_vote) => {},
241            HotShotEvent::DaCertificateRecv(simple_certificate) => {
242                self.replica_entry(simple_certificate.view_number())
243                    .da_certificate_recv = Some(now);
244            },
245            HotShotEvent::DaCertificateValidated(_simple_certificate) => {},
246            HotShotEvent::QuorumProposalSend(proposal, _) => {
247                self.leader_entry(proposal.data.view_number()).proposal_send = Some(now);
248
249                // If the last view succeeded, add the metric for time between proposals
250                if proposal.data.view_change_evidence().is_none()
251                    && let Some(previous_proposal_time) = self
252                        .replica_entry(proposal.data.view_number() - 1)
253                        .proposal_recv
254                {
255                    self.leader_entry(proposal.data.view_number())
256                        .prev_proposal_send = Some(previous_proposal_time);
257
258                    // calculate the elapsed time as milliseconds (from nanoseconds)
259                    let elapsed_time = (now - previous_proposal_time) / 1_000_000;
260                    if elapsed_time > 0 {
261                        self.consensus
262                            .read()
263                            .await
264                            .metrics
265                            .previous_proposal_to_proposal_time
266                            .add_point(elapsed_time as f64);
267                    } else {
268                        tracing::warn!("Previous proposal time is in the future");
269                    }
270                }
271            },
272            HotShotEvent::QuorumVoteSend(simple_vote) => {
273                self.replica_entry(simple_vote.view_number()).vote_send = Some(now);
274            },
275            HotShotEvent::ExtendedQuorumVoteSend(simple_vote) => {
276                self.replica_entry(simple_vote.view_number()).vote_send = Some(now);
277            },
278            HotShotEvent::QuorumProposalValidated(proposal, _) => {
279                self.replica_entry(proposal.data.view_number())
280                    .proposal_validated = Some(now);
281                self.replica_entry(proposal.data.view_number())
282                    .proposal_timestamp =
283                    Some(proposal.data.block_header().timestamp_millis() as i128);
284            },
285            HotShotEvent::DaProposalSend(proposal, _) => {
286                self.leader_entry(proposal.data.view_number())
287                    .da_proposal_send = Some(now);
288            },
289            HotShotEvent::DaVoteSend(simple_vote) => {
290                self.replica_entry(simple_vote.view_number()).vote_send = Some(now);
291            },
292            HotShotEvent::QcFormed(either) => {
293                match either {
294                    Either::Left(qc) => {
295                        self.leader_entry(qc.view_number() + 1).qc_formed = Some(now)
296                    },
297                    Either::Right(tc) => {
298                        self.leader_entry(tc.view_number())
299                            .timeout_certificate_formed = Some(now)
300                    },
301                };
302            },
303            HotShotEvent::Qc2Formed(either) => {
304                match either {
305                    Either::Left(qc) => {
306                        self.leader_entry(qc.view_number() + 1).qc_formed = Some(now)
307                    },
308                    Either::Right(tc) => {
309                        self.leader_entry(tc.view_number())
310                            .timeout_certificate_formed = Some(now)
311                    },
312                };
313            },
314            HotShotEvent::DacSend(simple_certificate, _) => {
315                self.leader_entry(simple_certificate.view_number())
316                    .da_cert_send = Some(now);
317            },
318            HotShotEvent::ViewChange(view, epoch) => {
319                // Record the timestamp of the first observed view change
320                // This can happen when transitioning to the next view, either due to voting
321                // or receiving a proposal, but we only store the first one
322                if self.replica_entry(*view + 1).view_change.is_none() {
323                    self.replica_entry(*view + 1).view_change = Some(now);
324                }
325
326                if *epoch <= self.epoch && *view <= self.view {
327                    return Ok(());
328                }
329                if self.view < *view {
330                    self.view = *view;
331                }
332                let prev_epoch = self.epoch;
333                let mut new_epoch = false;
334                if self.epoch < *epoch {
335                    self.epoch = *epoch;
336                    new_epoch = true;
337                }
338                if *view == ViewNumber::new(0) {
339                    return Ok(());
340                }
341
342                if new_epoch {
343                    if let Some(prev_epoch) = prev_epoch {
344                        self.log_basic_stats(now, &prev_epoch);
345                    }
346                    let _ = self.dump_stats();
347                    self.garbage_collect(*view - 1);
348                }
349
350                let leader = self
351                    .membership_coordinator
352                    .membership_for_epoch(*epoch)?
353                    .leader(*view)?;
354                if leader == self.public_key {
355                    self.leader_entry(*view).builder_start = Some(now);
356                }
357            },
358            HotShotEvent::Timeout(view, _) => {
359                self.replica_entry(*view).timeout_triggered = Some(now);
360                self.timeouts.insert(*view);
361            },
362            HotShotEvent::TransactionsRecv(_txns) => {
363                // TODO: Track transactions by time
364                // #3526 https://github.com/EspressoSystems/espresso-network/issues/3526
365            },
366            HotShotEvent::SendPayloadCommitmentAndMetadata(_, _, _, view, _) => {
367                self.leader_entry(*view).vid_disperse_send = Some(now);
368            },
369            HotShotEvent::VidShareRecv(_, proposal) => {
370                self.replica_entry(proposal.data.view_number())
371                    .vid_share_recv = Some(now);
372            },
373            HotShotEvent::VidShareValidated(proposal) => {
374                self.replica_entry(proposal.data.view_number())
375                    .vid_share_validated = Some(now);
376            },
377            HotShotEvent::QuorumProposalPreliminarilyValidated(proposal) => {
378                self.replica_entry(proposal.data.view_number())
379                    .proposal_prelim_validated = Some(now);
380            },
381            HotShotEvent::LeavesDecided(leaves) => {
382                for leaf in leaves {
383                    if leaf.view_number() == ViewNumber::genesis() {
384                        continue;
385                    }
386                    let view = leaf.view_number();
387                    let timestamp = leaf.block_header().timestamp_millis() as i128;
388                    let now_millis = now / 1_000_000;
389                    let latency = now_millis - timestamp;
390                    tracing::debug!("View {} Latency: {}ms", view, latency);
391                    self.latencies_by_view.insert(view, latency);
392                    self.sizes_by_view.insert(
393                        view,
394                        leaf.block_payload().map(|p| p.txn_bytes()).unwrap_or(0) as i128,
395                    );
396                }
397            },
398            _ => {},
399        }
400        Ok(())
401    }
402
403    fn cancel_subtasks(&mut self) {
404        // No subtasks to cancel
405    }
406}