Skip to main content

agave_validator/
dashboard.rs

1use {
2    crate::{
3        ProgressBar, admin_rpc_service, format_name_value, new_spinner_progress_bar,
4        println_name_value,
5    },
6    console::style,
7    solana_clock::Slot,
8    solana_commitment_config::CommitmentConfig,
9    solana_core::validator::ValidatorStartProgress,
10    solana_native_token::Sol,
11    solana_pubkey::Pubkey,
12    solana_rpc_client::rpc_client::RpcClient,
13    solana_rpc_client_api::{client_error, request, response::RpcContactInfo},
14    solana_validator_exit::Exit,
15    std::{
16        net::SocketAddr,
17        path::{Path, PathBuf},
18        sync::{
19            Arc,
20            atomic::{AtomicBool, Ordering},
21        },
22        thread,
23        time::{Duration, SystemTime},
24    },
25};
26
27pub struct Dashboard {
28    progress_bar: ProgressBar,
29    ledger_path: PathBuf,
30    exit: Arc<AtomicBool>,
31}
32
33impl Dashboard {
34    pub fn new(
35        ledger_path: &Path,
36        log_path: Option<&Path>,
37        validator_exit: Option<&mut Exit>,
38    ) -> Self {
39        println_name_value("Ledger location:", &format!("{}", ledger_path.display()));
40        if let Some(log_path) = log_path {
41            println_name_value("Log:", &format!("{}", log_path.display()));
42        }
43
44        let progress_bar = new_spinner_progress_bar();
45        progress_bar.set_message("Initializing...");
46
47        let exit = Arc::new(AtomicBool::new(false));
48        if let Some(validator_exit) = validator_exit {
49            let exit = exit.clone();
50            validator_exit.register_exit(Box::new(move || exit.store(true, Ordering::Relaxed)));
51        }
52
53        Self {
54            exit,
55            ledger_path: ledger_path.to_path_buf(),
56            progress_bar,
57        }
58    }
59
60    pub fn run(self, refresh_interval: Duration) {
61        let Self {
62            exit,
63            ledger_path,
64            progress_bar,
65            ..
66        } = self;
67        drop(progress_bar);
68
69        let runtime = admin_rpc_service::runtime();
70        while !exit.load(Ordering::Relaxed) {
71            let progress_bar = new_spinner_progress_bar();
72            progress_bar.set_message("Connecting...");
73
74            let Some((rpc_addr, start_time, mut vat_status)) = runtime.block_on(
75                wait_for_validator_startup(&ledger_path, &exit, progress_bar, refresh_interval),
76            ) else {
77                continue;
78            };
79
80            let rpc_client = RpcClient::new_socket(rpc_addr);
81            let mut identity = match rpc_client.get_identity() {
82                Ok(identity) => identity,
83                Err(err) => {
84                    println!("Failed to get validator identity over RPC: {err}");
85                    continue;
86                }
87            };
88            println_name_value("Identity:", &identity.to_string());
89            if let Some(status) = vat_status.as_ref()
90                && status.voting_enabled
91            {
92                println_name_value("Vote Account:", &status.vote_account.to_string());
93            }
94
95            if let Ok(genesis_hash) = rpc_client.get_genesis_hash() {
96                println_name_value("Genesis Hash:", &genesis_hash.to_string());
97            }
98
99            if let Some(contact_info) = get_contact_info(&rpc_client, &identity) {
100                println_name_value(
101                    "Version:",
102                    &contact_info.version.unwrap_or_else(|| "?".to_string()),
103                );
104                if let Some(shred_version) = contact_info.shred_version {
105                    println_name_value("Shred Version:", &shred_version.to_string());
106                }
107                if let Some(gossip) = contact_info.gossip {
108                    println_name_value("Gossip Address:", &gossip.to_string());
109                }
110                if let Some(tpu) = contact_info.tpu_quic {
111                    println_name_value("TPU QUIC Address:", &tpu.to_string());
112                }
113                if let Some(rpc) = contact_info.rpc {
114                    println_name_value("JSON RPC URL:", &format!("http://{rpc}"));
115                }
116                if let Some(pubsub) = contact_info.pubsub {
117                    println_name_value("WebSocket PubSub URL:", &format!("ws://{pubsub}"));
118                }
119            }
120
121            let progress_bar = new_spinner_progress_bar();
122            let mut snapshot_slot_info = None;
123            let mut admin_client = None;
124            for i in 0.. {
125                if exit.load(Ordering::Relaxed) {
126                    break;
127                }
128                if i % 10 == 0 {
129                    snapshot_slot_info = rpc_client.get_highest_snapshot_slot().ok();
130                }
131
132                let new_identity = rpc_client.get_identity().unwrap_or(identity);
133                if identity != new_identity {
134                    identity = new_identity;
135                    progress_bar.println(format_name_value("Identity:", &identity.to_string()));
136                    if let Some(status) = vat_status.as_ref()
137                        && status.voting_enabled
138                    {
139                        progress_bar.println(format_name_value(
140                            "Vote Account:",
141                            &status.vote_account.to_string(),
142                        ));
143                    }
144                }
145
146                if i > 0 && i % 30 == 0 {
147                    if admin_client.is_none() {
148                        admin_client = runtime
149                            .block_on(admin_rpc_service::connect(&ledger_path))
150                            .ok();
151                    }
152
153                    let vat_status_result = admin_client
154                        .as_ref()
155                        .map(|admin_client| runtime.block_on(admin_client.vat_status()));
156                    vat_status = match vat_status_result {
157                        Some(Ok(status)) => Some(status),
158                        Some(Err(_err)) => {
159                            admin_client = None;
160                            None
161                        }
162                        None => None,
163                    };
164                }
165
166                match get_validator_stats(&rpc_client, &identity) {
167                    Ok((
168                        processed_slot,
169                        confirmed_slot,
170                        finalized_slot,
171                        transaction_count,
172                        identity_balance,
173                        health,
174                    )) => {
175                        let uptime = {
176                            let uptime =
177                                chrono::Duration::from_std(start_time.elapsed().unwrap()).unwrap();
178
179                            format!(
180                                "{:02}:{:02}:{:02} ",
181                                uptime.num_hours(),
182                                uptime.num_minutes() % 60,
183                                uptime.num_seconds() % 60
184                            )
185                        };
186
187                        let vat_status_formatted = format_vat_status(vat_status.as_ref());
188
189                        progress_bar.set_message(format!(
190                            "{}{}| Processed Slot: {} | Confirmed Slot: {} | Finalized Slot: {} | \
191                             Full Snapshot Slot: {} | Incremental Snapshot Slot: {} | \
192                             Transactions: {} | {}\n{}",
193                            uptime,
194                            if health == "ok" {
195                                "".to_string()
196                            } else {
197                                format!("| {} ", style(health).bold().red())
198                            },
199                            processed_slot,
200                            confirmed_slot,
201                            finalized_slot,
202                            snapshot_slot_info
203                                .as_ref()
204                                .map(|snapshot_slot_info| snapshot_slot_info.full.to_string())
205                                .unwrap_or_else(|| '-'.to_string()),
206                            snapshot_slot_info
207                                .as_ref()
208                                .and_then(|snapshot_slot_info| snapshot_slot_info
209                                    .incremental
210                                    .map(|incremental| incremental.to_string()))
211                                .unwrap_or_else(|| '-'.to_string()),
212                            transaction_count,
213                            identity_balance,
214                            vat_status_formatted,
215                        ));
216                        thread::sleep(refresh_interval);
217                    }
218                    Err(err) => {
219                        progress_bar.abandon_with_message(format!("RPC connection failure: {err}"));
220                        break;
221                    }
222                }
223            }
224        }
225    }
226}
227
228fn format_vat_status(
229    status: Option<&admin_rpc_service::AdminRpcValidatorAdmissionTicketStatus>,
230) -> String {
231    let Some(status) = status else {
232        return "VAT: failed to connect to admin RPC".to_string();
233    };
234
235    if !status.vat_active {
236        return "VAT: inactive".to_string();
237    }
238
239    format!(
240        "{}{}",
241        format_current_vat_status(status),
242        format_effective_epoch_vat_status(status)
243    )
244}
245
246fn format_current_vat_status(
247    status: &admin_rpc_service::AdminRpcValidatorAdmissionTicketStatus,
248) -> String {
249    if status.in_current_epoch_vat {
250        format!(
251            "VAT: epoch {} in (stake: {})",
252            status.current_epoch,
253            Sol(status.current_epoch_vote_account_stake)
254        )
255    } else {
256        format!("VAT: epoch {} out", status.current_epoch)
257    }
258}
259
260fn format_effective_epoch_vat_status(
261    status: &admin_rpc_service::AdminRpcValidatorAdmissionTicketStatus,
262) -> String {
263    let vat_effective_epoch = status.current_epoch.saturating_add(2);
264    if let Some(vat_failure_reason) = &status.next_epoch_vat_failure_reason {
265        format!(", epoch {vat_effective_epoch}: {vat_failure_reason}")
266    } else {
267        format!(", epoch {vat_effective_epoch}: eligible if staked")
268    }
269}
270
271async fn wait_for_validator_startup(
272    ledger_path: &Path,
273    exit: &AtomicBool,
274    progress_bar: ProgressBar,
275    refresh_interval: Duration,
276) -> Option<(
277    SocketAddr,
278    SystemTime,
279    Option<admin_rpc_service::AdminRpcValidatorAdmissionTicketStatus>,
280)> {
281    let mut admin_client = None;
282    loop {
283        if exit.load(Ordering::Relaxed) {
284            return None;
285        }
286
287        if admin_client.is_none() {
288            match admin_rpc_service::connect(ledger_path).await {
289                Ok(new_admin_client) => admin_client = Some(new_admin_client),
290                Err(err) => {
291                    progress_bar.set_message(format!("Unable to connect to validator: {err}"));
292                    thread::sleep(refresh_interval);
293                    continue;
294                }
295            }
296        }
297
298        match admin_client.as_ref().unwrap().start_progress().await {
299            Ok(start_progress) => {
300                if start_progress == ValidatorStartProgress::Running {
301                    let admin_client = admin_client.take().unwrap();
302
303                    let validator_info = async move {
304                        let rpc_addr = admin_client.rpc_addr().await?;
305                        let start_time = admin_client.start_time().await?;
306                        let vat_status = if rpc_addr.is_some() {
307                            admin_client.vat_status().await.ok()
308                        } else {
309                            None
310                        };
311                        Ok::<_, jsonrpc_core_client::RpcError>((rpc_addr, start_time, vat_status))
312                    }
313                    .await;
314                    match validator_info {
315                        Ok((None, _, _)) => progress_bar.set_message("RPC service not available"),
316                        Ok((Some(rpc_addr), start_time, vat_status)) => {
317                            return Some((rpc_addr, start_time, vat_status));
318                        }
319                        Err(err) => {
320                            progress_bar
321                                .set_message(format!("Failed to get validator info: {err}"));
322                        }
323                    }
324                } else {
325                    progress_bar.set_message(format!("Validator startup: {start_progress:?}..."));
326                }
327            }
328            Err(err) => {
329                admin_client = None;
330                progress_bar.set_message(format!("Failed to get validator start progress: {err}"));
331            }
332        }
333        thread::sleep(refresh_interval);
334    }
335}
336
337fn get_contact_info(rpc_client: &RpcClient, identity: &Pubkey) -> Option<RpcContactInfo> {
338    rpc_client
339        .get_cluster_nodes()
340        .ok()
341        .unwrap_or_default()
342        .into_iter()
343        .find(|node| node.pubkey == identity.to_string())
344}
345
346fn get_validator_stats(
347    rpc_client: &RpcClient,
348    identity: &Pubkey,
349) -> client_error::Result<(Slot, Slot, Slot, u64, Sol, String)> {
350    let finalized_slot = rpc_client.get_slot_with_commitment(CommitmentConfig::finalized())?;
351    let confirmed_slot = rpc_client.get_slot_with_commitment(CommitmentConfig::confirmed())?;
352    let processed_slot = rpc_client.get_slot_with_commitment(CommitmentConfig::processed())?;
353    let transaction_count =
354        rpc_client.get_transaction_count_with_commitment(CommitmentConfig::processed())?;
355    let identity_balance = rpc_client
356        .get_balance_with_commitment(identity, CommitmentConfig::confirmed())?
357        .value;
358
359    let health = match rpc_client.get_health() {
360        Ok(()) => "ok".to_string(),
361        Err(err) => {
362            if let client_error::ErrorKind::RpcError(request::RpcError::RpcResponseError {
363                code: _,
364                message: _,
365                data:
366                    request::RpcResponseErrorData::NodeUnhealthy {
367                        num_slots_behind: Some(num_slots_behind),
368                    },
369            }) = err.kind()
370            {
371                format!("{num_slots_behind} slots behind")
372            } else {
373                "health unknown".to_string()
374            }
375        }
376    };
377
378    Ok((
379        processed_slot,
380        confirmed_slot,
381        finalized_slot,
382        transaction_count,
383        Sol(identity_balance),
384        health,
385    ))
386}