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 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 #[rpc(meta, name = "exit")]
191 fn exit(&self, meta: Self::Metadata) -> Result<()>;
192
193 #[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 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 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 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 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 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 let (response_sender, response_receiver) = oneshot_channel();
424
425 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 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 let (response_sender, response_receiver) = oneshot_channel();
453
454 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 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 let (response_sender, response_receiver) = oneshot_channel();
481
482 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 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 let (response_sender, response_receiver) = oneshot_channel();
509
510 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 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 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
1025pub 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) .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 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
1090pub 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
1103pub fn runtime() -> Runtime {
1105 tokio::runtime::Builder::new_multi_thread()
1106 .thread_name("solAdminRpcRt")
1107 .enable_all()
1108 .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(), 0u16, ),
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 #[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, 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 #[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 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 assert!(result["result"].is_boolean());
1568 assert!(result["result"].as_bool().unwrap());
1570 }
1571
1572 #[test]
1573 fn test_is_generating_snapshots_no_controller() {
1574 let rpc = RpcHandler::start_with_config(TestConfig::default());
1576 let RpcHandler { io, .. } = rpc;
1577
1578 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 assert!(result["error"].is_object());
1606 assert_eq!(
1607 result["error"]["message"].as_str().unwrap(),
1608 "snapshot_controller unavailable"
1609 );
1610 }
1611}