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 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 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 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 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 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 },
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 }
406}