Skip to main content

agave_validator/
admin_rpc_service.rs

1use {
2    agave_votor::{
3        event::VotorEvent, vote_history::VoteHistory, vote_history_storage::VoteHistoryStorage,
4    },
5    crossbeam_channel::Sender,
6    jsonrpc_core::{BoxFuture, ErrorCode, MetaIoHandler, Metadata, Result},
7    jsonrpc_core_client::{RpcError, transports::ipc},
8    jsonrpc_derive::rpc,
9    jsonrpc_ipc_server::{
10        RequestContext, ServerBuilder, tokio::sync::oneshot::channel as oneshot_channel,
11    },
12    log::*,
13    serde::{Deserialize, Serialize, de::Deserializer},
14    solana_clock::{Epoch, Slot},
15    solana_core::{
16        admin_rpc_post_init::AdminRpcRequestMetadataPostInit,
17        banking_stage::{
18            BankingControlMsg, BankingStage,
19            transaction_scheduler::scheduler_controller::SchedulerConfig,
20        },
21        consensus::{Tower, tower_storage::TowerStorage},
22        repair::repair_service,
23        validator::{
24            BlockProductionMethod, SchedulerPacing, TransactionStructure, ValidatorStartProgress,
25            should_require_vote_history_file,
26        },
27    },
28    solana_geyser_plugin_manager::GeyserPluginManagerRequest,
29    solana_gossip::contact_info::{ContactInfo, Protocol, SOCKET_ADDR_UNSPECIFIED},
30    solana_keypair::{Keypair, read_keypair_file},
31    solana_metrics::{datapoint_info, datapoint_warn},
32    solana_pubkey::Pubkey,
33    solana_runtime::{bank::VATHealthError, snapshot_controller::SnapshotController},
34    solana_signer::Signer,
35    solana_validator_exit::Exit,
36    std::{
37        collections::{HashMap, HashSet},
38        env, error,
39        fmt::{self, Display},
40        net::{IpAddr, SocketAddr},
41        num::NonZeroUsize,
42        path::{Path, PathBuf},
43        sync::{
44            Arc, RwLock,
45            atomic::{AtomicBool, Ordering},
46        },
47        thread::{self, Builder},
48        time::{Duration, Instant, SystemTime},
49    },
50    tokio::runtime::Runtime,
51};
52
53#[derive(Clone)]
54pub struct AdminRpcRequestMetadata {
55    pub rpc_addr: Option<SocketAddr>,
56    pub start_time: SystemTime,
57    pub start_progress: Arc<RwLock<ValidatorStartProgress>>,
58    pub validator_exit: Arc<RwLock<Exit>>,
59    pub validator_exit_backpressure: HashMap<String, Arc<AtomicBool>>,
60    pub authorized_voter_keypairs: Arc<RwLock<Vec<Arc<Keypair>>>>,
61    pub tower_storage: Arc<dyn TowerStorage>,
62    pub vote_history_storage: Arc<dyn VoteHistoryStorage>,
63    pub staked_nodes_overrides: Arc<RwLock<HashMap<Pubkey, u64>>>,
64    pub post_init: Arc<RwLock<Option<AdminRpcRequestMetadataPostInit>>>,
65    pub rpc_to_plugin_manager_sender: Option<Sender<GeyserPluginManagerRequest>>,
66}
67
68impl Metadata for AdminRpcRequestMetadata {}
69
70impl AdminRpcRequestMetadata {
71    fn with_post_init<F, R>(&self, func: F) -> Result<R>
72    where
73        F: FnOnce(&AdminRpcRequestMetadataPostInit) -> Result<R>,
74    {
75        if let Some(post_init) = self.post_init.read().unwrap().as_ref() {
76            func(post_init)
77        } else {
78            Err(jsonrpc_core::error::Error::invalid_params(
79                "Retry once validator start up is complete",
80            ))
81        }
82    }
83
84    fn snapshot_controller(&self) -> Option<Arc<SnapshotController>> {
85        self.with_post_init(|post_init| Ok(post_init.snapshot_controller.clone()))
86            .map_err(|_| {
87                // The error from with_post_init is not relevant, as it is meant for RPC callers
88                warn!("snapshot_controller unavailable, shutting down without taking snapshot");
89            })
90            .ok()
91    }
92}
93
94#[derive(Debug, Deserialize, Serialize)]
95pub struct AdminRpcContactInfo {
96    pub id: String,
97    pub gossip: SocketAddr,
98    pub tvu: SocketAddr,
99    pub serve_repair_quic: SocketAddr,
100    pub tpu: SocketAddr,
101    pub tpu_quic: Option<SocketAddr>,
102    pub tpu_forwards: SocketAddr,
103    pub tpu_vote: SocketAddr,
104    pub rpc: SocketAddr,
105    pub rpc_pubsub: SocketAddr,
106    pub serve_repair: SocketAddr,
107    pub last_updated_timestamp: u64,
108    pub shred_version: u16,
109}
110
111#[derive(Debug, Deserialize, Serialize)]
112pub struct AdminRpcRepairWhitelist {
113    pub whitelist: Vec<Pubkey>,
114}
115
116#[derive(Debug, Deserialize, Serialize)]
117pub struct AdminRpcValidatorAdmissionTicketStatus {
118    pub vote_account: Pubkey,
119    pub voting_enabled: bool,
120    pub current_epoch: Epoch,
121    pub in_current_epoch_vat: bool,
122    pub current_epoch_vote_account_stake: u64,
123    pub current_epoch_vote_accounts: usize,
124    pub next_epoch_vat_failure_reason: Option<VATHealthError>,
125}
126
127impl From<ContactInfo> for AdminRpcContactInfo {
128    fn from(node: ContactInfo) -> Self {
129        macro_rules! unwrap_socket {
130            ($name:ident) => {
131                node.$name().unwrap_or(SOCKET_ADDR_UNSPECIFIED)
132            };
133            ($name:ident, $protocol:expr) => {
134                node.$name($protocol).unwrap_or(SOCKET_ADDR_UNSPECIFIED)
135            };
136        }
137        Self {
138            id: node.pubkey().to_string(),
139            last_updated_timestamp: node.wallclock(),
140            gossip: unwrap_socket!(gossip),
141            tvu: unwrap_socket!(tvu, Protocol::UDP),
142            serve_repair_quic: unwrap_socket!(serve_repair, Protocol::QUIC),
143            tpu: unwrap_socket!(tpu, Protocol::UDP),
144            tpu_quic: node.tpu(Protocol::QUIC),
145            tpu_forwards: unwrap_socket!(tpu_forwards, Protocol::UDP),
146            tpu_vote: unwrap_socket!(tpu_vote, Protocol::UDP),
147            rpc: unwrap_socket!(rpc),
148            rpc_pubsub: unwrap_socket!(rpc_pubsub),
149            serve_repair: unwrap_socket!(serve_repair, Protocol::UDP),
150            shred_version: node.shred_version(),
151        }
152    }
153}
154
155impl Display for AdminRpcContactInfo {
156    fn fmt(&self, f: &mut fmt::Formatter) -> fmt::Result {
157        writeln!(f, "Identity: {}", self.id)?;
158        writeln!(f, "Gossip: {}", self.gossip)?;
159        writeln!(f, "TVU: {}", self.tvu)?;
160        writeln!(f, "TPU: {}", self.tpu)?;
161        if let Some(tpu_quic) = self.tpu_quic {
162            writeln!(f, "TPU QUIC: {tpu_quic}")?;
163        }
164        writeln!(f, "TPU Forwards: {}", self.tpu_forwards)?;
165        writeln!(f, "TPU Votes: {}", self.tpu_vote)?;
166        writeln!(f, "RPC: {}", self.rpc)?;
167        writeln!(f, "RPC Pubsub: {}", self.rpc_pubsub)?;
168        writeln!(f, "Serve Repair: {}", self.serve_repair)?;
169        writeln!(f, "Last Updated Timestamp: {}", self.last_updated_timestamp)?;
170        writeln!(f, "Shred Version: {}", self.shred_version)
171    }
172}
173impl solana_cli_output::VerboseDisplay for AdminRpcContactInfo {}
174impl solana_cli_output::QuietDisplay for AdminRpcContactInfo {}
175
176impl Display for AdminRpcRepairWhitelist {
177    fn fmt(&self, f: &mut fmt::Formatter) -> fmt::Result {
178        writeln!(f, "Repair whitelist: {:?}", self.whitelist)
179    }
180}
181impl solana_cli_output::VerboseDisplay for AdminRpcRepairWhitelist {}
182impl solana_cli_output::QuietDisplay for AdminRpcRepairWhitelist {}
183
184#[rpc]
185pub trait AdminRpc {
186    type Metadata;
187
188    /// Initiates validator exit; exit is asynchronous so the validator
189    /// will almost certainly still be running when this method returns
190    #[rpc(meta, name = "exit")]
191    fn exit(&self, meta: Self::Metadata) -> Result<()>;
192
193    /// Return the process id (pid)
194    #[rpc(meta, name = "pid")]
195    fn pid(&self, meta: Self::Metadata) -> Result<u32>;
196
197    #[rpc(meta, name = "reloadPlugin")]
198    fn reload_plugin(
199        &self,
200        meta: Self::Metadata,
201        name: String,
202        config_file: String,
203    ) -> BoxFuture<Result<()>>;
204
205    #[rpc(meta, name = "unloadPlugin")]
206    fn unload_plugin(&self, meta: Self::Metadata, name: String) -> BoxFuture<Result<()>>;
207
208    #[rpc(meta, name = "loadPlugin")]
209    fn load_plugin(&self, meta: Self::Metadata, config_file: String) -> BoxFuture<Result<String>>;
210
211    #[rpc(meta, name = "listPlugins")]
212    fn list_plugins(&self, meta: Self::Metadata) -> BoxFuture<Result<Vec<String>>>;
213
214    #[rpc(meta, name = "rpcAddress")]
215    fn rpc_addr(&self, meta: Self::Metadata) -> Result<Option<SocketAddr>>;
216
217    #[rpc(name = "setLogFilter")]
218    fn set_log_filter(&self, filter: String) -> Result<()>;
219
220    #[rpc(meta, name = "startTime")]
221    fn start_time(&self, meta: Self::Metadata) -> Result<SystemTime>;
222
223    #[rpc(meta, name = "startProgress")]
224    fn start_progress(&self, meta: Self::Metadata) -> Result<ValidatorStartProgress>;
225
226    #[rpc(meta, name = "addAuthorizedVoter")]
227    fn add_authorized_voter(&self, meta: Self::Metadata, keypair_file: String) -> Result<()>;
228
229    #[rpc(meta, name = "addAuthorizedVoterFromBytes")]
230    fn add_authorized_voter_from_bytes(&self, meta: Self::Metadata, keypair: Vec<u8>)
231    -> Result<()>;
232
233    #[rpc(meta, name = "removeAllAuthorizedVoters")]
234    fn remove_all_authorized_voters(&self, meta: Self::Metadata) -> Result<()>;
235
236    #[rpc(meta, name = "setIdentity")]
237    fn set_identity(
238        &self,
239        meta: Self::Metadata,
240        keypair_file: String,
241        require_tower: bool,
242        require_vote_history: bool,
243    ) -> Result<()>;
244
245    #[rpc(meta, name = "setIdentityFromBytes")]
246    fn set_identity_from_bytes(
247        &self,
248        meta: Self::Metadata,
249        identity_keypair: Vec<u8>,
250        require_tower: bool,
251        require_vote_history: bool,
252    ) -> Result<()>;
253
254    #[rpc(meta, name = "setStakedNodesOverrides")]
255    fn set_staked_nodes_overrides(&self, meta: Self::Metadata, path: String) -> Result<()>;
256
257    #[rpc(meta, name = "contactInfo")]
258    fn contact_info(&self, meta: Self::Metadata) -> Result<AdminRpcContactInfo>;
259
260    #[rpc(meta, name = "validatorAdmissionTicketStatus")]
261    fn vat_status(&self, meta: Self::Metadata) -> Result<AdminRpcValidatorAdmissionTicketStatus>;
262
263    #[rpc(meta, name = "selectActiveInterface")]
264    fn select_active_interface(&self, meta: Self::Metadata, interface: IpAddr) -> Result<()>;
265
266    #[rpc(meta, name = "repairShredFromPeer")]
267    fn repair_shred_from_peer(
268        &self,
269        meta: Self::Metadata,
270        pubkey: Option<Pubkey>,
271        slot: u64,
272        shred_index: u64,
273    ) -> Result<()>;
274
275    #[rpc(meta, name = "repairWhitelist")]
276    fn repair_whitelist(&self, meta: Self::Metadata) -> Result<AdminRpcRepairWhitelist>;
277
278    #[rpc(meta, name = "setRepairWhitelist")]
279    fn set_repair_whitelist(&self, meta: Self::Metadata, whitelist: Vec<Pubkey>) -> Result<()>;
280
281    #[rpc(meta, name = "setPublicTpuAddress")]
282    fn set_public_tpu_address(
283        &self,
284        meta: Self::Metadata,
285        public_tpu_addr: SocketAddr,
286    ) -> Result<()>;
287
288    #[rpc(meta, name = "setPublicTpuForwardsAddress")]
289    fn set_public_tpu_forwards_address(
290        &self,
291        meta: Self::Metadata,
292        public_tpu_forwards_addr: SocketAddr,
293    ) -> Result<()>;
294
295    #[rpc(meta, name = "setPublicTvuAddress")]
296    fn set_public_tvu_address(
297        &self,
298        meta: Self::Metadata,
299        public_tvu_addr: SocketAddr,
300    ) -> Result<()>;
301
302    #[rpc(meta, name = "manageBlockProduction")]
303    fn manage_block_production(
304        &self,
305        meta: Self::Metadata,
306        block_production_method: BlockProductionMethod,
307        transaction_struct: TransactionStructure,
308        num_workers: NonZeroUsize,
309        scheduler_pacing: SchedulerPacing,
310    ) -> Result<()>;
311
312    #[rpc(meta, name = "isGeneratingSnapshots")]
313    fn is_generating_snapshots(&self, meta: Self::Metadata) -> Result<bool>;
314
315    #[rpc(meta, name = "blockstorePurge")]
316    fn blockstore_purge(&self, meta: Self::Metadata, maximum_purge_slot: Slot) -> Result<()>;
317}
318
319pub struct AdminRpcImpl;
320impl AdminRpc for AdminRpcImpl {
321    type Metadata = AdminRpcRequestMetadata;
322
323    fn exit(&self, meta: Self::Metadata) -> Result<()> {
324        debug!("exit admin rpc request received");
325
326        thread::Builder::new()
327            .name("solProcessExit".into())
328            .spawn(move || {
329                let start_time = Instant::now();
330
331                // Trigger a fastboot snapshot before exiting
332                if let Some(snapshot_controller) = meta.snapshot_controller() {
333                    let latest_snapshot_slot = snapshot_controller.latest_bank_snapshot_slot();
334
335                    info!("Requesting fastboot snapshot before exit");
336                    snapshot_controller.request_fastboot_snapshot();
337
338                    // Wait up to 5s for a snapshot to finish. This should allow time for the
339                    // fastboot snapshot to complete without stalling exit indefinitely.
340                    // The timeout will be hit in the event new roots are not being created.
341                    let timeout = Duration::from_secs(5);
342                    while snapshot_controller.latest_bank_snapshot_slot() == latest_snapshot_slot {
343                        if start_time.elapsed() > timeout {
344                            warn!("Timeout waiting for snapshot to complete");
345                            datapoint_warn!(
346                                "admin-rpc-snapshot-timeout",
347                                ("timeout_us", start_time.elapsed().as_micros(), i64)
348                            );
349                            break;
350                        }
351                        thread::sleep(Duration::from_millis(100));
352                    }
353                    info!(
354                        "Requesting fastboot snapshot before exit... Done in {:?}",
355                        start_time.elapsed()
356                    );
357                }
358
359                // Delay exit signal until this RPC request completes, otherwise the caller of `exit` might
360                // receive a confusing error as the validator shuts down before a response is sent back.
361                // If elapsed time has already taken 100ms, there is no need for further delay
362                if start_time.elapsed().as_millis() < 100 {
363                    thread::sleep(Duration::from_millis(100));
364                }
365
366                info!("validator exit requested");
367                meta.validator_exit.write().unwrap().exit();
368
369                if !meta.validator_exit_backpressure.is_empty() {
370                    let service_names = meta.validator_exit_backpressure.keys();
371                    info!("Wait for these services to complete: {service_names:?}");
372                    loop {
373                        // The initial sleep is a grace period to allow the services to raise their
374                        // backpressure flags.
375                        // Subsequent sleeps are to throttle how often we check and log.
376                        thread::sleep(Duration::from_secs(1));
377
378                        let mut any_flags_raised = false;
379                        for (name, flag) in meta.validator_exit_backpressure.iter() {
380                            let is_flag_raised = flag.load(Ordering::Relaxed);
381                            if is_flag_raised {
382                                info!("{name}'s exit backpressure flag is raised");
383                                any_flags_raised = true;
384                            }
385                        }
386                        if !any_flags_raised {
387                            break;
388                        }
389                    }
390                    info!("All services have completed");
391                }
392
393                // TODO: Debug why Exit doesn't always cause the validator to fully exit
394                // (rocksdb background processing or some other stuck thread perhaps?).
395                //
396                // If the process is still alive after five seconds, exit harder
397                thread::sleep(Duration::from_secs(
398                    env::var("SOLANA_VALIDATOR_EXIT_TIMEOUT")
399                        .ok()
400                        .and_then(|x| x.parse().ok())
401                        .unwrap_or(5),
402                ));
403                warn!("validator exit timeout");
404                std::process::exit(0);
405            })
406            .unwrap();
407
408        Ok(())
409    }
410
411    fn pid(&self, _meta: Self::Metadata) -> Result<u32> {
412        Ok(std::process::id())
413    }
414
415    fn reload_plugin(
416        &self,
417        meta: Self::Metadata,
418        name: String,
419        config_file: String,
420    ) -> BoxFuture<Result<()>> {
421        Box::pin(async move {
422            // Construct channel for plugin to respond to this particular rpc request instance
423            let (response_sender, response_receiver) = oneshot_channel();
424
425            // Send request to plugin manager if there is a geyser service
426            if let Some(ref rpc_to_manager_sender) = meta.rpc_to_plugin_manager_sender {
427                rpc_to_manager_sender
428                    .send(GeyserPluginManagerRequest::ReloadPlugin {
429                        name,
430                        config_file,
431                        response_sender,
432                    })
433                    .expect("GeyerPluginService should never drop request receiver");
434            } else {
435                return Err(jsonrpc_core::Error {
436                    code: ErrorCode::InvalidRequest,
437                    message: "No geyser plugin service".to_string(),
438                    data: None,
439                });
440            }
441
442            // Await response from plugin manager
443            response_receiver
444                .await
445                .expect("GeyerPluginService's oneshot sender shouldn't drop early")
446        })
447    }
448
449    fn load_plugin(&self, meta: Self::Metadata, config_file: String) -> BoxFuture<Result<String>> {
450        Box::pin(async move {
451            // Construct channel for plugin to respond to this particular rpc request instance
452            let (response_sender, response_receiver) = oneshot_channel();
453
454            // Send request to plugin manager if there is a geyser service
455            if let Some(ref rpc_to_manager_sender) = meta.rpc_to_plugin_manager_sender {
456                rpc_to_manager_sender
457                    .send(GeyserPluginManagerRequest::LoadPlugin {
458                        config_file,
459                        response_sender,
460                    })
461                    .expect("GeyerPluginService should never drop request receiver");
462            } else {
463                return Err(jsonrpc_core::Error {
464                    code: ErrorCode::InvalidRequest,
465                    message: "No geyser plugin service".to_string(),
466                    data: None,
467                });
468            }
469
470            // Await response from plugin manager
471            response_receiver
472                .await
473                .expect("GeyerPluginService's oneshot sender shouldn't drop early")
474        })
475    }
476
477    fn unload_plugin(&self, meta: Self::Metadata, name: String) -> BoxFuture<Result<()>> {
478        Box::pin(async move {
479            // Construct channel for plugin to respond to this particular rpc request instance
480            let (response_sender, response_receiver) = oneshot_channel();
481
482            // Send request to plugin manager if there is a geyser service
483            if let Some(ref rpc_to_manager_sender) = meta.rpc_to_plugin_manager_sender {
484                rpc_to_manager_sender
485                    .send(GeyserPluginManagerRequest::UnloadPlugin {
486                        name,
487                        response_sender,
488                    })
489                    .expect("GeyerPluginService should never drop request receiver");
490            } else {
491                return Err(jsonrpc_core::Error {
492                    code: ErrorCode::InvalidRequest,
493                    message: "No geyser plugin service".to_string(),
494                    data: None,
495                });
496            }
497
498            // Await response from plugin manager
499            response_receiver
500                .await
501                .expect("GeyerPluginService's oneshot sender shouldn't drop early")
502        })
503    }
504
505    fn list_plugins(&self, meta: Self::Metadata) -> BoxFuture<Result<Vec<String>>> {
506        Box::pin(async move {
507            // Construct channel for plugin to respond to this particular rpc request instance
508            let (response_sender, response_receiver) = oneshot_channel();
509
510            // Send request to plugin manager
511            if let Some(ref rpc_to_manager_sender) = meta.rpc_to_plugin_manager_sender {
512                rpc_to_manager_sender
513                    .send(GeyserPluginManagerRequest::ListPlugins { response_sender })
514                    .expect("GeyerPluginService should never drop request receiver");
515            } else {
516                return Err(jsonrpc_core::Error {
517                    code: ErrorCode::InvalidRequest,
518                    message: "No geyser plugin service".to_string(),
519                    data: None,
520                });
521            }
522
523            // Await response from plugin manager
524            response_receiver
525                .await
526                .expect("GeyerPluginService's oneshot sender shouldn't drop early")
527        })
528    }
529
530    fn rpc_addr(&self, meta: Self::Metadata) -> Result<Option<SocketAddr>> {
531        debug!("rpc_addr admin rpc request received");
532        Ok(meta.rpc_addr)
533    }
534
535    fn set_log_filter(&self, filter: String) -> Result<()> {
536        debug!("set_log_filter admin rpc request received");
537        agave_logger::setup_with(&filter);
538        Ok(())
539    }
540
541    fn start_time(&self, meta: Self::Metadata) -> Result<SystemTime> {
542        debug!("start_time admin rpc request received");
543        Ok(meta.start_time)
544    }
545
546    fn start_progress(&self, meta: Self::Metadata) -> Result<ValidatorStartProgress> {
547        debug!("start_progress admin rpc request received");
548        Ok(*meta.start_progress.read().unwrap())
549    }
550
551    fn add_authorized_voter(&self, meta: Self::Metadata, keypair_file: String) -> Result<()> {
552        debug!("add_authorized_voter request received");
553
554        let authorized_voter = read_keypair_file(keypair_file)
555            .map_err(|err| jsonrpc_core::error::Error::invalid_params(format!("{err}")))?;
556
557        AdminRpcImpl::add_authorized_voter_keypair(meta, authorized_voter)
558    }
559
560    fn add_authorized_voter_from_bytes(
561        &self,
562        meta: Self::Metadata,
563        keypair: Vec<u8>,
564    ) -> Result<()> {
565        debug!("add_authorized_voter_from_bytes request received");
566
567        let authorized_voter = Keypair::try_from(keypair.as_ref()).map_err(|err| {
568            jsonrpc_core::error::Error::invalid_params(format!(
569                "Failed to read authorized voter keypair from provided byte array: {err}"
570            ))
571        })?;
572
573        AdminRpcImpl::add_authorized_voter_keypair(meta, authorized_voter)
574    }
575
576    fn remove_all_authorized_voters(&self, meta: Self::Metadata) -> Result<()> {
577        debug!("remove_all_authorized_voters received");
578        meta.authorized_voter_keypairs.write().unwrap().clear();
579        Ok(())
580    }
581
582    fn set_identity(
583        &self,
584        meta: Self::Metadata,
585        keypair_file: String,
586        require_tower: bool,
587        require_vote_history: bool,
588    ) -> Result<()> {
589        debug!("set_identity request received");
590
591        let identity_keypair = read_keypair_file(&keypair_file).map_err(|err| {
592            jsonrpc_core::error::Error::invalid_params(format!(
593                "Failed to read identity keypair from {keypair_file}: {err}"
594            ))
595        })?;
596
597        AdminRpcImpl::set_identity_keypair(
598            meta,
599            identity_keypair,
600            require_tower,
601            require_vote_history,
602        )
603    }
604
605    fn set_identity_from_bytes(
606        &self,
607        meta: Self::Metadata,
608        identity_keypair: Vec<u8>,
609        require_tower: bool,
610        require_vote_history: bool,
611    ) -> Result<()> {
612        debug!("set_identity_from_bytes request received");
613
614        let identity_keypair = Keypair::try_from(identity_keypair.as_ref()).map_err(|err| {
615            jsonrpc_core::error::Error::invalid_params(format!(
616                "Failed to read identity keypair from provided byte array: {err}"
617            ))
618        })?;
619
620        AdminRpcImpl::set_identity_keypair(
621            meta,
622            identity_keypair,
623            require_tower,
624            require_vote_history,
625        )
626    }
627
628    fn set_staked_nodes_overrides(&self, meta: Self::Metadata, path: String) -> Result<()> {
629        let loaded_config = load_staked_nodes_overrides(&path)
630            .map_err(|err| {
631                error!("Failed to load staked nodes overrides from {path}: {err}");
632                jsonrpc_core::error::Error::internal_error()
633            })?
634            .staked_map_id;
635        let mut write_staked_nodes = meta.staked_nodes_overrides.write().unwrap();
636        write_staked_nodes.clear();
637        write_staked_nodes.extend(loaded_config);
638        info!("Staked nodes overrides loaded from {path}");
639        debug!("overrides map: {write_staked_nodes:?}");
640        Ok(())
641    }
642
643    fn contact_info(&self, meta: Self::Metadata) -> Result<AdminRpcContactInfo> {
644        meta.with_post_init(|post_init| Ok(post_init.cluster_info.my_contact_info().into()))
645    }
646
647    fn vat_status(&self, meta: Self::Metadata) -> Result<AdminRpcValidatorAdmissionTicketStatus> {
648        debug!("vat_status request received");
649        meta.with_post_init(|post_init| {
650            let bank = post_init.bank_forks.read().unwrap().root_bank();
651            let current_epoch = bank.epoch();
652            let epoch_vote_accounts = bank
653                .epoch_vote_accounts(current_epoch)
654                .ok_or_else(jsonrpc_core::Error::internal_error)?;
655            let current_epoch_vote_account_stake = epoch_vote_accounts
656                .get(&post_init.vote_account)
657                .map(|(stake, _vote_account)| *stake)
658                .unwrap_or_default();
659            let in_current_epoch_vote_accounts =
660                epoch_vote_accounts.contains_key(&post_init.vote_account);
661
662            let next_epoch_vat_failure_reason = bank
663                .get_vat_health_for_next_epoch(&post_init.vote_account)
664                .err();
665
666            let voting_enabled = !meta.authorized_voter_keypairs.read().unwrap().is_empty();
667            Ok(AdminRpcValidatorAdmissionTicketStatus {
668                vote_account: post_init.vote_account,
669                voting_enabled,
670                current_epoch,
671                in_current_epoch_vat: in_current_epoch_vote_accounts,
672                current_epoch_vote_account_stake,
673                current_epoch_vote_accounts: epoch_vote_accounts.len(),
674                next_epoch_vat_failure_reason,
675            })
676        })
677    }
678
679    fn select_active_interface(&self, meta: Self::Metadata, interface: IpAddr) -> Result<()> {
680        debug!("select_active_interface received: {interface}");
681        meta.with_post_init(|post_init| {
682            let node = post_init.node.as_ref().ok_or_else(|| {
683                jsonrpc_core::Error::invalid_params("`Node` not initialized in post_init")
684            })?;
685
686            node.switch_active_interface(interface, &post_init.cluster_info)
687                .map_err(|e| {
688                    jsonrpc_core::Error::invalid_params(format!(
689                        "Switching failed due to error {e}"
690                    ))
691                })?;
692            info!("Switched primary interface to {interface}");
693            Ok(())
694        })
695    }
696
697    fn repair_shred_from_peer(
698        &self,
699        meta: Self::Metadata,
700        pubkey: Option<Pubkey>,
701        slot: u64,
702        shred_index: u64,
703    ) -> Result<()> {
704        debug!("repair_shred_from_peer request received");
705
706        meta.with_post_init(|post_init| {
707            repair_service::RepairService::request_repair_for_shred_from_peer(
708                post_init.cluster_info.clone(),
709                post_init.cluster_slots.clone(),
710                pubkey,
711                slot,
712                shred_index,
713                &post_init.repair_socket,
714                post_init.outstanding_repair_requests.clone(),
715            );
716            Ok(())
717        })
718    }
719
720    fn repair_whitelist(&self, meta: Self::Metadata) -> Result<AdminRpcRepairWhitelist> {
721        debug!("repair_whitelist request received");
722
723        meta.with_post_init(|post_init| {
724            let whitelist: Vec<_> = post_init
725                .repair_whitelist
726                .read()
727                .unwrap()
728                .iter()
729                .copied()
730                .collect();
731            Ok(AdminRpcRepairWhitelist { whitelist })
732        })
733    }
734
735    fn set_repair_whitelist(&self, meta: Self::Metadata, whitelist: Vec<Pubkey>) -> Result<()> {
736        debug!("set_repair_whitelist request received");
737
738        let whitelist: HashSet<Pubkey> = whitelist.into_iter().collect();
739        meta.with_post_init(|post_init| {
740            *post_init.repair_whitelist.write().unwrap() = whitelist;
741            warn!(
742                "Repair whitelist set to {:?}",
743                post_init.repair_whitelist.read().unwrap()
744            );
745            Ok(())
746        })
747    }
748
749    fn set_public_tpu_address(
750        &self,
751        meta: Self::Metadata,
752        public_tpu_addr: SocketAddr,
753    ) -> Result<()> {
754        debug!("set_public_tpu_address rpc request received: {public_tpu_addr}");
755
756        meta.with_post_init(|post_init| {
757            post_init
758                .cluster_info
759                .my_contact_info()
760                .tpu(Protocol::QUIC)
761                .ok_or_else(|| {
762                    error!(
763                        "The public TPU QUIC address isn't being published. The node is likely in \
764                         repair mode. See help for --restricted-repair-only-mode for more \
765                         information."
766                    );
767                    jsonrpc_core::error::Error::internal_error()
768                })?;
769            post_init
770                .cluster_info
771                .set_tpu_quic(public_tpu_addr)
772                .map_err(|err| {
773                    error!("Failed to set public TPU QUIC address to {public_tpu_addr}: {err}");
774                    jsonrpc_core::error::Error::internal_error()
775                })?;
776            let my_contact_info = post_init.cluster_info.my_contact_info();
777            warn!(
778                "Public TPU addresses set to {:?} (quic)",
779                my_contact_info.tpu(Protocol::QUIC),
780            );
781            Ok(())
782        })
783    }
784
785    fn set_public_tpu_forwards_address(
786        &self,
787        meta: Self::Metadata,
788        public_tpu_forwards_addr: SocketAddr,
789    ) -> Result<()> {
790        debug!("set_public_tpu_forwards_address rpc request received: {public_tpu_forwards_addr}");
791
792        meta.with_post_init(|post_init| {
793            post_init
794                .cluster_info
795                .my_contact_info()
796                .tpu_forwards(Protocol::QUIC)
797                .ok_or_else(|| {
798                    error!(
799                        "The public TPU Forwards address isn't being published. The node is \
800                         likely in repair mode. See help for --restricted-repair-only-mode for \
801                         more information."
802                    );
803                    jsonrpc_core::error::Error::internal_error()
804                })?;
805            post_init
806                .cluster_info
807                .set_tpu_forwards_quic(public_tpu_forwards_addr)
808                .map_err(|err| {
809                    error!(
810                        "Failed to set public TPU QUIC address to {public_tpu_forwards_addr}: \
811                         {err}"
812                    );
813                    jsonrpc_core::error::Error::internal_error()
814                })?;
815            let my_contact_info = post_init.cluster_info.my_contact_info();
816            warn!(
817                "Public TPU Forwards address set to {:?} (quic)",
818                my_contact_info.tpu_forwards(Protocol::QUIC),
819            );
820            Ok(())
821        })
822    }
823
824    fn set_public_tvu_address(
825        &self,
826        meta: Self::Metadata,
827        public_tvu_addr: SocketAddr,
828    ) -> Result<()> {
829        debug!("set_public_tvu_address rpc request received: {public_tvu_addr}");
830
831        meta.with_post_init(|post_init| {
832            post_init
833                .cluster_info
834                .my_contact_info()
835                .tvu(Protocol::UDP)
836                .ok_or_else(|| {
837                    error!(
838                        "The public TVU address isn't being published. The node is likely in \
839                         repair mode. See help for --restricted-repair-only-mode for more \
840                         information."
841                    );
842                    jsonrpc_core::error::Error::internal_error()
843                })?;
844            post_init
845                .cluster_info
846                .set_tvu_socket(public_tvu_addr)
847                .map_err(|err| {
848                    error!("Failed to set public TVU address to {public_tvu_addr}: {err}");
849                    jsonrpc_core::error::Error::internal_error()
850                })?;
851            let my_contact_info = post_init.cluster_info.my_contact_info();
852            warn!(
853                "Public TVU addresses set to {:?}",
854                my_contact_info.tvu(Protocol::UDP),
855            );
856            Ok(())
857        })
858    }
859
860    fn manage_block_production(
861        &self,
862        meta: Self::Metadata,
863        block_production_method: BlockProductionMethod,
864        transaction_struct: TransactionStructure,
865        num_workers: NonZeroUsize,
866        scheduler_pacing: SchedulerPacing,
867    ) -> Result<()> {
868        debug!("manage_block_production rpc request received");
869
870        if num_workers > BankingStage::max_num_workers() {
871            return Err(jsonrpc_core::error::Error::invalid_params(format!(
872                "Number of workers ({}) exceeds maximum allowed ({})",
873                num_workers,
874                BankingStage::max_num_workers()
875            )));
876        }
877
878        if transaction_struct != TransactionStructure::View {
879            warn!("TransactionStructure::Sdk has no effect on block production");
880        }
881
882        meta.with_post_init(|post_init| {
883            if post_init
884                .banking_control_sender
885                .try_send(BankingControlMsg::Internal {
886                    block_production_method,
887                    num_workers,
888                    config: SchedulerConfig { scheduler_pacing },
889                })
890                .is_err()
891            {
892                error!("Banking stage already switching schedulers");
893
894                return Err(jsonrpc_core::error::Error::internal_error());
895            }
896
897            Ok(())
898        })
899    }
900
901    fn is_generating_snapshots(&self, meta: Self::Metadata) -> Result<bool> {
902        if let Some(snapshot_controller) = meta.snapshot_controller() {
903            Ok(snapshot_controller.is_generating_snapshots())
904        } else {
905            Err(jsonrpc_core::error::Error::invalid_params(
906                "snapshot_controller unavailable",
907            ))
908        }
909    }
910
911    fn blockstore_purge(&self, meta: Self::Metadata, maximum_purge_slot: Slot) -> Result<()> {
912        meta.with_post_init(|post_init| {
913            post_init
914                .blockstore
915                .send_manual_purge_request(maximum_purge_slot)
916                .map_err(|err| jsonrpc_core::Error {
917                    code: ErrorCode::InvalidRequest,
918                    message: format!("{err}"),
919                    data: None,
920                })
921        })
922    }
923}
924
925impl AdminRpcImpl {
926    fn add_authorized_voter_keypair(
927        meta: AdminRpcRequestMetadata,
928        authorized_voter: Keypair,
929    ) -> Result<()> {
930        let mut authorized_voter_keypairs = meta.authorized_voter_keypairs.write().unwrap();
931
932        if authorized_voter_keypairs
933            .iter()
934            .any(|x| x.pubkey() == authorized_voter.pubkey())
935        {
936            Err(jsonrpc_core::error::Error::invalid_params(
937                "Authorized voter already present",
938            ))
939        } else {
940            authorized_voter_keypairs.push(Arc::new(authorized_voter));
941            Ok(())
942        }
943    }
944
945    fn set_identity_keypair(
946        meta: AdminRpcRequestMetadata,
947        identity_keypair: Keypair,
948        require_tower: bool,
949        require_vote_history: bool,
950    ) -> Result<()> {
951        meta.with_post_init(|post_init| {
952            if require_tower {
953                let _ = Tower::restore(meta.tower_storage.as_ref(), &identity_keypair.pubkey())
954                    .map_err(|err| {
955                        jsonrpc_core::error::Error::invalid_params(format!(
956                            "Unable to load tower file for identity {}: {}",
957                            identity_keypair.pubkey(),
958                            err
959                        ))
960                    })?;
961            }
962
963            if require_vote_history {
964                let should_require_vote_history = {
965                    let bank_forks = post_init.bank_forks.read().unwrap();
966                    should_require_vote_history_file(
967                        &bank_forks.working_bank(),
968                        &post_init.vote_account,
969                        &identity_keypair.pubkey(),
970                    )
971                };
972                if should_require_vote_history {
973                    let _ = VoteHistory::restore(
974                        meta.vote_history_storage.as_ref(),
975                        &identity_keypair.pubkey(),
976                    )
977                    .map_err(|err| {
978                        jsonrpc_core::error::Error::invalid_params(format!(
979                            "Unable to load vote history file for identity {}: {}. The vote \
980                             account {} has prior Alpenglow votes. Ensure the vote history file \
981                             is present or (dangerous) use --do-not-require-vote-history if you \
982                             know what you're doing",
983                            identity_keypair.pubkey(),
984                            err,
985                            post_init.vote_account
986                        ))
987                    })?;
988                }
989            }
990
991            for (key, notifier) in &*post_init.notifies.read().unwrap() {
992                if let Err(err) = notifier.update_key(&identity_keypair) {
993                    error!("Error updating network layer keypair: {err} on {key:?}");
994                }
995            }
996
997            let old_identity = post_init.cluster_info.id();
998            let new_identity = identity_keypair.pubkey();
999            solana_metrics::set_host_id(new_identity.to_string());
1000            // Emit the datapoint after updating metrics to emit the new pubkey
1001            datapoint_info!(
1002                "validator-set_identity",
1003                ("old_id", old_identity.to_string(), String),
1004                ("new_id", new_identity.to_string(), String),
1005                ("version", solana_version::version!(), String),
1006            );
1007            post_init
1008                .cluster_info
1009                .set_keypair(Arc::new(identity_keypair));
1010            post_init
1011                .votor_event_sender
1012                .send(VotorEvent::SetIdentity)
1013                .map_err(|err| jsonrpc_core::error::Error {
1014                    code: ErrorCode::InternalError,
1015                    message: format!("Failed to send SetIdentity event: {err}").to_string(),
1016                    data: None,
1017                })?;
1018
1019            warn!("Identity set to {new_identity}");
1020            Ok(())
1021        })
1022    }
1023}
1024
1025// Start the Admin RPC interface
1026pub fn run(ledger_path: &Path, metadata: AdminRpcRequestMetadata) {
1027    let admin_rpc_path = admin_rpc_path(ledger_path);
1028
1029    let event_loop = tokio::runtime::Builder::new_multi_thread()
1030        .thread_name("solAdminRpcEl")
1031        .worker_threads(3) // Three still seems like a lot, and better than the default of available core count
1032        .enable_all()
1033        .build()
1034        .unwrap();
1035
1036    Builder::new()
1037        .name("solAdminRpc".to_string())
1038        .spawn(move || {
1039            let mut io = MetaIoHandler::default();
1040            io.extend_with(AdminRpcImpl.to_delegate());
1041
1042            let validator_exit = metadata.validator_exit.clone();
1043            let server = ServerBuilder::with_meta_extractor(io, move |_req: &RequestContext| {
1044                metadata.clone()
1045            })
1046            .event_loop_executor(event_loop.handle().clone())
1047            .start(&format!("{}", admin_rpc_path.display()));
1048
1049            match server {
1050                Err(err) => {
1051                    warn!("Unable to start admin rpc service: {err:?}");
1052                }
1053                Ok(server) => {
1054                    info!("started admin rpc service!");
1055                    let close_handle = server.close_handle();
1056                    validator_exit
1057                        .write()
1058                        .unwrap()
1059                        .register_exit(Box::new(move || {
1060                            close_handle.close();
1061                        }));
1062
1063                    server.wait();
1064                }
1065            }
1066        })
1067        .unwrap();
1068}
1069
1070fn admin_rpc_path(ledger_path: &Path) -> PathBuf {
1071    #[cfg(target_family = "windows")]
1072    {
1073        // More information about the wackiness of pipe names over at
1074        // https://docs.microsoft.com/en-us/windows/win32/ipc/pipe-names
1075        if let Some(ledger_filename) = ledger_path.file_name() {
1076            PathBuf::from(format!(
1077                "\\\\.\\pipe\\{}-admin.rpc",
1078                ledger_filename.to_string_lossy()
1079            ))
1080        } else {
1081            PathBuf::from("\\\\.\\pipe\\admin.rpc")
1082        }
1083    }
1084    #[cfg(not(target_family = "windows"))]
1085    {
1086        ledger_path.join("admin.rpc")
1087    }
1088}
1089
1090// Connect to the Admin RPC interface
1091pub async fn connect(ledger_path: &Path) -> std::result::Result<gen_client::Client, RpcError> {
1092    let admin_rpc_path = admin_rpc_path(ledger_path);
1093    if !admin_rpc_path.exists() {
1094        Err(RpcError::Client(format!(
1095            "{} does not exist",
1096            admin_rpc_path.display()
1097        )))
1098    } else {
1099        ipc::connect::<_, gen_client::Client>(&format!("{}", admin_rpc_path.display())).await
1100    }
1101}
1102
1103// Create a runtime for use by client side admin RPC interface calls
1104pub fn runtime() -> Runtime {
1105    tokio::runtime::Builder::new_multi_thread()
1106        .thread_name("solAdminRpcRt")
1107        .enable_all()
1108        // The agave-validator subcommands make few admin RPC calls and block
1109        // on the results so two workers is plenty
1110        .worker_threads(2)
1111        .build()
1112        .expect("new tokio runtime")
1113}
1114
1115#[derive(Default, Deserialize, Clone)]
1116pub struct StakedNodesOverrides {
1117    #[serde(deserialize_with = "deserialize_pubkey_map")]
1118    pub staked_map_id: HashMap<Pubkey, u64>,
1119}
1120
1121pub fn deserialize_pubkey_map<'de, D>(des: D) -> std::result::Result<HashMap<Pubkey, u64>, D::Error>
1122where
1123    D: Deserializer<'de>,
1124{
1125    let container: HashMap<String, u64> = serde::Deserialize::deserialize(des)?;
1126    let mut container_typed: HashMap<Pubkey, u64> = HashMap::new();
1127    for (key, value) in container.iter() {
1128        let typed_key = Pubkey::try_from(key.as_str())
1129            .map_err(|_| serde::de::Error::invalid_type(serde::de::Unexpected::Map, &"PubKey"))?;
1130        container_typed.insert(typed_key, *value);
1131    }
1132    Ok(container_typed)
1133}
1134
1135pub fn load_staked_nodes_overrides(
1136    path: &String,
1137) -> std::result::Result<StakedNodesOverrides, Box<dyn error::Error>> {
1138    debug!("Loading staked nodes overrides configuration from {path}");
1139    if Path::new(&path).exists() {
1140        let file = std::fs::File::open(path)?;
1141        Ok(serde_yaml::from_reader(file)?)
1142    } else {
1143        Err(format!("Staked nodes overrides provided '{path}' a non-existing file path.").into())
1144    }
1145}
1146
1147#[cfg(test)]
1148mod tests {
1149    use {
1150        super::*,
1151        agave_snapshots::snapshot_config::SnapshotConfig,
1152        agave_votor::event::VotorEventSender,
1153        assert_matches::assert_matches,
1154        crossbeam_channel::bounded,
1155        serde_json::Value,
1156        solana_accounts_db::{
1157            accounts_db::{ACCOUNTS_DB_CONFIG_FOR_TESTING, AccountsDbConfig},
1158            accounts_index::AccountSecondaryIndexes,
1159        },
1160        solana_core::{
1161            admin_rpc_post_init::{KeyUpdaterType, KeyUpdaters},
1162            consensus::tower_storage::NullTowerStorage,
1163            validator::{Validator, ValidatorConfig, ValidatorTpuConfig},
1164        },
1165        solana_gossip::{cluster_info::ClusterInfo, node::Node},
1166        solana_ledger::{
1167            blockstore::Blockstore,
1168            create_new_tmp_ledger,
1169            genesis_utils::{
1170                GenesisConfigInfo, create_genesis_config, create_genesis_config_with_leader,
1171            },
1172            get_tmp_ledger_path_auto_delete,
1173        },
1174        solana_net_utils::{SocketAddrSpace, sockets::bind_to_localhost_unique},
1175        solana_rpc::rpc::create_validator_exit,
1176        solana_runtime::{
1177            bank::{Bank, BankTestConfig},
1178            bank_forks::BankForks,
1179        },
1180        std::{collections::HashSet, fs::remove_dir_all, sync::atomic::AtomicBool},
1181        tokio::sync::mpsc,
1182    };
1183
1184    #[derive(Default)]
1185    struct TestConfig {
1186        account_indexes: AccountSecondaryIndexes,
1187        votor_event_sender: Option<VotorEventSender>,
1188    }
1189
1190    struct RpcHandler {
1191        io: MetaIoHandler<AdminRpcRequestMetadata>,
1192        meta: AdminRpcRequestMetadata,
1193    }
1194
1195    impl RpcHandler {
1196        fn _start() -> Self {
1197            Self::start_with_config(TestConfig::default())
1198        }
1199
1200        fn start_with_config(config: TestConfig) -> Self {
1201            let keypair = Arc::new(Keypair::new());
1202            let cluster_info = Arc::new(ClusterInfo::new(
1203                ContactInfo::new(
1204                    keypair.pubkey(),
1205                    solana_time_utils::timestamp(), // wallclock
1206                    0u16,                           // shred_version
1207                ),
1208                keypair,
1209                SocketAddrSpace::Unspecified,
1210            ));
1211            let exit = Arc::new(AtomicBool::new(false));
1212            let validator_exit = create_validator_exit(exit);
1213            let (bank_forks, vote_keypair) = new_bank_forks_with_config(BankTestConfig {
1214                accounts_db_config: AccountsDbConfig {
1215                    account_indexes: Some(config.account_indexes),
1216                    ..ACCOUNTS_DB_CONFIG_FOR_TESTING
1217                },
1218            });
1219
1220            let (snapshot_request_sender, _) = bounded(1024);
1221            let snapshot_controller = Arc::new(SnapshotController::new(
1222                snapshot_request_sender.clone(),
1223                SnapshotConfig::default(),
1224                bank_forks.read().unwrap().root(),
1225            ));
1226
1227            let ledger_path = get_tmp_ledger_path_auto_delete!();
1228            let blockstore = Arc::new(Blockstore::open(ledger_path.path()).unwrap());
1229
1230            let vote_account = vote_keypair.pubkey();
1231            let start_progress = Arc::new(RwLock::new(ValidatorStartProgress::default()));
1232            let repair_whitelist = Arc::new(RwLock::new(HashSet::new()));
1233            let votor_event_sender = config.votor_event_sender.unwrap_or_else(|| {
1234                let (votor_event_sender, _) = bounded(1024);
1235                votor_event_sender
1236            });
1237            let meta = AdminRpcRequestMetadata {
1238                rpc_addr: None,
1239                start_time: SystemTime::now(),
1240                start_progress,
1241                validator_exit,
1242                validator_exit_backpressure: HashMap::default(),
1243                authorized_voter_keypairs: Arc::new(RwLock::new(vec![vote_keypair])),
1244                tower_storage: Arc::new(NullTowerStorage {}),
1245                vote_history_storage: Arc::new(
1246                    agave_votor::vote_history_storage::NullVoteHistoryStorage::default(),
1247                ),
1248                post_init: Arc::new(RwLock::new(Some(AdminRpcRequestMetadataPostInit {
1249                    cluster_info,
1250                    bank_forks,
1251                    vote_account,
1252                    repair_whitelist,
1253                    notifies: Arc::new(RwLock::new(KeyUpdaters::default())),
1254                    repair_socket: Arc::new(bind_to_localhost_unique().expect("should bind")),
1255                    outstanding_repair_requests: Arc::<
1256                        RwLock<repair_service::OutstandingShredRepairs>,
1257                    >::default(),
1258                    cluster_slots: Arc::new(
1259                        solana_core::cluster_slots_service::cluster_slots::ClusterSlots::default_for_tests(),
1260                    ),
1261                    node: None,
1262                    banking_control_sender: mpsc::channel(1).0,
1263                    snapshot_controller,
1264                    blockstore,
1265                    votor_event_sender,
1266                }))),
1267                staked_nodes_overrides: Arc::new(RwLock::new(HashMap::new())),
1268                rpc_to_plugin_manager_sender: None,
1269            };
1270            let mut io = MetaIoHandler::default();
1271            io.extend_with(AdminRpcImpl.to_delegate());
1272
1273            Self { io, meta }
1274        }
1275    }
1276
1277    fn new_bank_forks_with_config(
1278        config: BankTestConfig,
1279    ) -> (Arc<RwLock<BankForks>>, Arc<Keypair>) {
1280        let GenesisConfigInfo {
1281            genesis_config,
1282            voting_keypair,
1283            ..
1284        } = create_genesis_config(1_000_000_000);
1285
1286        let bank = Bank::new_with_paths_for_tests(&genesis_config, Some(config), vec![], None);
1287        (BankForks::new_rw_arc(bank), Arc::new(voting_keypair))
1288    }
1289
1290    // This test checks that the rpc call to `set_identity` works a expected with
1291    // Bank but without validator.
1292    #[test]
1293    fn test_set_identity() {
1294        let (votor_event_sender, votor_event_receiver) = bounded(1024);
1295        let rpc = RpcHandler::start_with_config(TestConfig {
1296            account_indexes: AccountSecondaryIndexes::default(),
1297            votor_event_sender: Some(votor_event_sender),
1298        });
1299
1300        let RpcHandler { io, meta, .. } = rpc;
1301
1302        let expected_validator_id = Keypair::new();
1303        let validator_id_bytes = format!("{:?}", expected_validator_id.to_bytes());
1304
1305        let set_id_request = format!(
1306            r#"{{"jsonrpc":"2.0","id":1,"method":"setIdentityFromBytes","params":[{validator_id_bytes}, false, false]}}"#,
1307        );
1308        let response = io.handle_request_sync(&set_id_request, meta.clone());
1309        let actual_parsed_response: Value =
1310            serde_json::from_str(&response.expect("actual response"))
1311                .expect("actual response deserialization");
1312
1313        let expected_parsed_response: Value = serde_json::from_str(
1314            r#"{
1315                "id": 1,
1316                "jsonrpc": "2.0",
1317                "result": null
1318            }"#,
1319        )
1320        .expect("Failed to parse expected response");
1321        assert_eq!(actual_parsed_response, expected_parsed_response);
1322
1323        let contact_info_request =
1324            r#"{"jsonrpc":"2.0","id":1,"method":"contactInfo","params":[]}"#.to_string();
1325        let response = io.handle_request_sync(&contact_info_request, meta.clone());
1326        let parsed_response: Value = serde_json::from_str(&response.expect("actual response"))
1327            .expect("actual response deserialization");
1328        let actual_validator_id = parsed_response["result"]["id"]
1329            .as_str()
1330            .expect("Expected a string");
1331        assert_eq!(
1332            actual_validator_id,
1333            expected_validator_id.pubkey().to_string()
1334        );
1335        let event = votor_event_receiver
1336            .recv()
1337            .expect("Failed to receive SetIdentity event");
1338        assert_matches!(event, VotorEvent::SetIdentity);
1339    }
1340
1341    #[test]
1342    fn test_vat_status() {
1343        let rpc = RpcHandler::_start();
1344        let RpcHandler { io, meta, .. } = rpc;
1345        let post_init = meta.post_init.read().unwrap().clone().unwrap();
1346        let bank = post_init.bank_forks.read().unwrap().root_bank();
1347        let current_epoch = bank.epoch();
1348        let epoch_vote_accounts = bank.epoch_vote_accounts(current_epoch).unwrap();
1349        let expected_in_vat = epoch_vote_accounts.contains_key(&post_init.vote_account);
1350        let expected_stake = epoch_vote_accounts
1351            .get(&post_init.vote_account)
1352            .map(|(stake, _vote_account)| *stake)
1353            .unwrap_or_default();
1354
1355        let request =
1356            r#"{"jsonrpc":"2.0","id":1,"method":"validatorAdmissionTicketStatus","params":[]}"#;
1357        let response = io.handle_request_sync(request, meta.clone());
1358        let result: Value = serde_json::from_str(&response.expect("actual response"))
1359            .expect("actual response deserialization");
1360        let status: AdminRpcValidatorAdmissionTicketStatus =
1361            serde_json::from_value(result["result"].clone()).unwrap();
1362
1363        assert_eq!(status.vote_account, post_init.vote_account);
1364        assert!(status.voting_enabled);
1365        assert_eq!(status.current_epoch, current_epoch);
1366        assert_eq!(status.in_current_epoch_vat, expected_in_vat);
1367        assert_eq!(status.current_epoch_vote_account_stake, expected_stake);
1368        assert_eq!(
1369            status.current_epoch_vote_accounts,
1370            epoch_vote_accounts.len()
1371        );
1372    }
1373
1374    struct TestValidatorWithAdminRpc {
1375        meta: AdminRpcRequestMetadata,
1376        io: MetaIoHandler<AdminRpcRequestMetadata>,
1377        validator_ledger_path: PathBuf,
1378    }
1379
1380    impl TestValidatorWithAdminRpc {
1381        fn new() -> Self {
1382            let leader_keypair = Keypair::new();
1383            let leader_node = Node::new_localhost_with_pubkey(&leader_keypair.pubkey());
1384
1385            let validator_keypair = Keypair::new();
1386            let validator_node = Node::new_localhost_with_pubkey(&validator_keypair.pubkey());
1387            let genesis_config =
1388                create_genesis_config_with_leader(10_000, &leader_keypair.pubkey(), 1000)
1389                    .genesis_config;
1390            let (validator_ledger_path, _blockhash) = create_new_tmp_ledger!(&genesis_config);
1391
1392            let voting_keypair = Arc::new(Keypair::new());
1393            let voting_pubkey = voting_keypair.pubkey();
1394            let authorized_voter_keypairs = Arc::new(RwLock::new(vec![voting_keypair]));
1395            let validator_config = ValidatorConfig {
1396                rpc_addrs: Some((
1397                    validator_node.info.rpc().unwrap(),
1398                    validator_node.info.rpc_pubsub().unwrap(),
1399                )),
1400                ..ValidatorConfig::default_for_test()
1401            };
1402            let start_progress = Arc::new(RwLock::new(ValidatorStartProgress::default()));
1403
1404            let post_init = Arc::new(RwLock::new(None));
1405            let meta = AdminRpcRequestMetadata {
1406                rpc_addr: validator_config.rpc_addrs.map(|(rpc_addr, _)| rpc_addr),
1407                start_time: SystemTime::now(),
1408                start_progress: start_progress.clone(),
1409                validator_exit: validator_config.validator_exit.clone(),
1410                validator_exit_backpressure: HashMap::default(),
1411                authorized_voter_keypairs: authorized_voter_keypairs.clone(),
1412                tower_storage: Arc::new(NullTowerStorage {}),
1413                vote_history_storage: Arc::new(
1414                    agave_votor::vote_history_storage::NullVoteHistoryStorage::default(),
1415                ),
1416                post_init: post_init.clone(),
1417                staked_nodes_overrides: Arc::new(RwLock::new(HashMap::new())),
1418                rpc_to_plugin_manager_sender: None,
1419            };
1420
1421            let _validator = Validator::new(
1422                validator_node,
1423                Arc::new(validator_keypair),
1424                &validator_ledger_path,
1425                &voting_pubkey,
1426                authorized_voter_keypairs,
1427                vec![leader_node.info],
1428                &validator_config,
1429                None, // rpc_to_plugin_manager_receiver
1430                start_progress.clone(),
1431                SocketAddrSpace::Unspecified,
1432                ValidatorTpuConfig::new_for_tests(),
1433                post_init.clone(),
1434                None,
1435            )
1436            .expect("assume successful validator start");
1437            assert_eq!(
1438                *start_progress.read().unwrap(),
1439                ValidatorStartProgress::Running
1440            );
1441            let post_init = post_init.read().unwrap();
1442
1443            assert!(post_init.is_some());
1444            let post_init = post_init.as_ref().unwrap();
1445            let notifies = post_init.notifies.read().unwrap();
1446            let updater_keys: HashSet<KeyUpdaterType> =
1447                notifies.into_iter().map(|(key, _)| key.clone()).collect();
1448            assert_eq!(
1449                updater_keys,
1450                HashSet::from_iter(vec![
1451                    KeyUpdaterType::Tpu,
1452                    KeyUpdaterType::TpuForwards,
1453                    KeyUpdaterType::TpuVote,
1454                    KeyUpdaterType::Forward,
1455                    KeyUpdaterType::RpcService,
1456                    KeyUpdaterType::Votor,
1457                    KeyUpdaterType::VotorPeerListService,
1458                ])
1459            );
1460            let mut io = MetaIoHandler::default();
1461            io.extend_with(AdminRpcImpl.to_delegate());
1462            Self {
1463                meta,
1464                io,
1465                validator_ledger_path,
1466            }
1467        }
1468
1469        fn handle_request(&self, request: &str) -> Option<String> {
1470            self.io.handle_request_sync(request, self.meta.clone())
1471        }
1472    }
1473
1474    impl Drop for TestValidatorWithAdminRpc {
1475        fn drop(&mut self) {
1476            remove_dir_all(self.validator_ledger_path.clone()).unwrap();
1477        }
1478    }
1479
1480    #[test]
1481    fn test_no_post_init_no_snapshot_controller() {
1482        let validator_exit = create_validator_exit(Arc::new(AtomicBool::new(false)));
1483        let voting_keypair = Arc::new(Keypair::new());
1484        let authorized_voter_keypairs = Arc::new(RwLock::new(vec![voting_keypair]));
1485        let start_progress = Arc::new(RwLock::new(ValidatorStartProgress::default()));
1486
1487        let post_init = Arc::new(RwLock::new(None));
1488        let meta = AdminRpcRequestMetadata {
1489            rpc_addr: None,
1490            start_time: SystemTime::now(),
1491            start_progress: start_progress.clone(),
1492            validator_exit,
1493            validator_exit_backpressure: HashMap::default(),
1494            authorized_voter_keypairs: authorized_voter_keypairs.clone(),
1495            tower_storage: Arc::new(NullTowerStorage {}),
1496            vote_history_storage: Arc::new(
1497                agave_votor::vote_history_storage::NullVoteHistoryStorage::default(),
1498            ),
1499            post_init: post_init.clone(),
1500            staked_nodes_overrides: Arc::new(RwLock::new(HashMap::new())),
1501            rpc_to_plugin_manager_sender: None,
1502        };
1503
1504        let snapshot_controller = meta.snapshot_controller();
1505        assert!(snapshot_controller.is_none());
1506    }
1507
1508    // This test checks that `set_identity` call works with working validator and client.
1509    #[test]
1510    fn test_set_identity_with_validator() {
1511        let test_validator = TestValidatorWithAdminRpc::new();
1512        let expected_validator_id = Keypair::new();
1513        let validator_id_bytes = format!("{:?}", expected_validator_id.to_bytes());
1514
1515        let set_id_request = format!(
1516            r#"{{"jsonrpc":"2.0","id":1,"method":"setIdentityFromBytes","params":[{validator_id_bytes}, false, false]}}"#,
1517        );
1518        let response = test_validator.handle_request(&set_id_request);
1519        let actual_parsed_response: Value =
1520            serde_json::from_str(&response.expect("actual response"))
1521                .expect("actual response deserialization");
1522
1523        let expected_parsed_response: Value = serde_json::from_str(
1524            r#"{
1525                "id": 1,
1526                "jsonrpc": "2.0",
1527                "result": null
1528            }"#,
1529        )
1530        .expect("Failed to parse expected response");
1531        assert_eq!(actual_parsed_response, expected_parsed_response);
1532
1533        let contact_info_request =
1534            r#"{"jsonrpc":"2.0","id":1,"method":"contactInfo","params":[]}"#.to_string();
1535        let response = test_validator.handle_request(&contact_info_request);
1536        let parsed_response: Value = serde_json::from_str(&response.expect("actual response"))
1537            .expect("actual response deserialization");
1538        let actual_validator_id = parsed_response["result"]["id"]
1539            .as_str()
1540            .expect("Expected a string");
1541        assert_eq!(
1542            actual_validator_id,
1543            expected_validator_id.pubkey().to_string()
1544        );
1545
1546        let contact_info_request =
1547            r#"{"jsonrpc":"2.0","id":1,"method":"exit","params":[]}"#.to_string();
1548        let exit_response = test_validator.handle_request(&contact_info_request);
1549        let actual_parsed_response: Value =
1550            serde_json::from_str(&exit_response.expect("actual response"))
1551                .expect("actual response deserialization");
1552        assert_eq!(actual_parsed_response, expected_parsed_response);
1553    }
1554
1555    #[test]
1556    fn test_is_generating_snapshots() {
1557        // Test with snapshots enabled
1558        let rpc = RpcHandler::start_with_config(TestConfig::default());
1559        let RpcHandler { io, meta, .. } = rpc;
1560
1561        let request = r#"{"jsonrpc":"2.0","id":1,"method":"isGeneratingSnapshots","params":[]}"#;
1562        let response = io.handle_request_sync(request, meta.clone());
1563        let result: Value = serde_json::from_str(&response.expect("actual response"))
1564            .expect("actual response deserialization");
1565
1566        // Should return a boolean result indicating if snapshots are being generated
1567        assert!(result["result"].is_boolean());
1568        // Verify that snapshots are being generated since the test setup includes a snapshot controller
1569        assert!(result["result"].as_bool().unwrap());
1570    }
1571
1572    #[test]
1573    fn test_is_generating_snapshots_no_controller() {
1574        // Test with snapshots enabled
1575        let rpc = RpcHandler::start_with_config(TestConfig::default());
1576        let RpcHandler { io, .. } = rpc;
1577
1578        // Test with no post_init (snapshot_controller unavailable)
1579        let request = r#"{"jsonrpc":"2.0","id":1,"method":"isGeneratingSnapshots","params":[]}"#;
1580        let validator_exit = create_validator_exit(Arc::new(AtomicBool::new(false)));
1581        let authorized_voter_keypairs = Arc::new(RwLock::new(vec![Arc::new(Keypair::new())]));
1582        let start_progress = Arc::new(RwLock::new(ValidatorStartProgress::default()));
1583
1584        let meta_no_post_init = AdminRpcRequestMetadata {
1585            rpc_addr: None,
1586            start_time: SystemTime::now(),
1587            start_progress,
1588            validator_exit,
1589            validator_exit_backpressure: HashMap::default(),
1590            authorized_voter_keypairs,
1591            tower_storage: Arc::new(NullTowerStorage {}),
1592            vote_history_storage: Arc::new(
1593                agave_votor::vote_history_storage::NullVoteHistoryStorage::default(),
1594            ),
1595            post_init: Arc::new(RwLock::new(None)),
1596            staked_nodes_overrides: Arc::new(RwLock::new(HashMap::new())),
1597            rpc_to_plugin_manager_sender: None,
1598        };
1599
1600        let response = io.handle_request_sync(request, meta_no_post_init);
1601        let result: Value = serde_json::from_str(&response.expect("actual response"))
1602            .expect("actual response deserialization");
1603
1604        // Should return an error when snapshot_controller is unavailable
1605        assert!(result["error"].is_object());
1606        assert_eq!(
1607            result["error"]["message"].as_str().unwrap(),
1608            "snapshot_controller unavailable"
1609        );
1610    }
1611}