main.rs (53829B)
1 use std::{ 2 collections::BTreeMap, 3 net::SocketAddr, 4 path::{Path, PathBuf}, 5 sync::{ 6 Arc, 7 atomic::{AtomicBool, Ordering}, 8 }, 9 time::{Duration, Instant}, 10 }; 11 12 use anyhow::{Context, Result, bail}; 13 use iuna::{ 14 adapters::{ 15 chain_store::SqliteChainStore, config_store, http, p2p, stratum, 16 ui_data_store::SqliteUiDataStore, wallet_endpoint, wallet_store, 17 }, 18 app::{ 19 NodeCore, PeerBook, SharedNode, SharedPeerBook, StratumStatus, debug_logging_enabled, 20 now_ms, set_debug_logging, validate_network_genesis, 21 }, 22 domain::{ 23 Amount, ChainSnapshot, GenesisBurn, LaunchProfile, Ledger, MAX_VDF_ROUNDS, MICRO_IUNA, 24 VDF_TARGET_BLOCK_MS, VdfProgress, VdfProgressPhase, run_vdf, 25 run_vdf_cancellable_with_progress, 26 }, 27 }; 28 use secrecy::{ExposeSecret, SecretString}; 29 use tokio::sync::Mutex; 30 31 mod cli; 32 #[cfg(feature = "cli-updater")] 33 mod updater; 34 use cli::{ 35 ChainMode, CliOptions, apply_cli_p2p_config_overrides, apply_cli_stratum_config_overrides, 36 apply_cli_wallet_endpoint_config_overrides, configured_p2p_announce_addr, 37 configured_p2p_bind_addr, configured_stratum_addr, configured_wallet_endpoint_addr, 38 initial_burn_fee, initial_burn_per_block, validate_wallet_for_mode, 39 }; 40 #[cfg(test)] 41 use cli::{default_data_dir, help_text}; 42 43 const GENESIS_BOOTSTRAP_BALANCE: Amount = MICRO_IUNA; 44 const GENESIS_BOOTSTRAP_BURN_AMOUNT: Amount = MICRO_IUNA; 45 const GENESIS_INITIAL_BURN_PER_BLOCK: Amount = config_store::DEFAULT_BURN_AMOUNT; 46 const GENESIS_INITIAL_BURN_FEE: Amount = config_store::DEFAULT_BURN_FEE; 47 const VDF_MEASUREMENT_INITIAL_ROUNDS: u64 = 1_000; 48 const VDF_MEASUREMENT_MAX_ROUNDS: u64 = 10_000_000; 49 const VDF_MEASUREMENT_MIN_ELAPSED: Duration = Duration::from_millis(150); 50 const VDF_PROGRESS_LOG_INTERVAL: Duration = Duration::from_secs(10); 51 const SYNC_CHAIN_CHECKPOINT_INTERVAL: Duration = Duration::from_secs(30); 52 const AUTOMATIC_BURN_ENABLED_ENV: &str = "IUNA_AUTOMATIC_BURN_ENABLED"; 53 const POW_MINING_ENABLED_ENV: &str = "IUNA_POW_MINING_ENABLED"; 54 const POW_MINING_WORKERS_ENV: &str = "IUNA_POW_MINING_WORKERS"; 55 const LOCAL_TESTNET_ENV: &str = "IUNA_LOCAL_TESTNET"; 56 const SETUP_COMPLETE_ENV: &str = "IUNA_SETUP_COMPLETE"; 57 const WALLET_PASSWORD_ENV: &str = "IUNA_WALLET_PASSWORD"; 58 const WALLET_ENDPOINT_ENABLED_ENV: &str = "IUNA_WALLET_ENDPOINT_ENABLED"; 59 const WALLET_ENDPOINT_PORT_ENV: &str = "IUNA_WALLET_ENDPOINT_PORT"; 60 61 #[tokio::main] 62 async fn main() -> Result<()> { 63 #[cfg(feature = "cli-updater")] 64 if updater::handle_cli_command().await? { 65 return Ok(()); 66 } 67 let Some(opts) = CliOptions::parse()? else { 68 return Ok(()); 69 }; 70 let startup_local_testnet = startup_bool_from_env(LOCAL_TESTNET_ENV)?.unwrap_or(false); 71 #[cfg(feature = "e2e")] 72 if !startup_local_testnet { 73 bail!("the e2e build requires {LOCAL_TESTNET_ENV}=true"); 74 } 75 set_debug_logging(opts.debug); 76 let debug_logging = opts.debug; 77 let wallet_path = opts.wallet_path(); 78 let config_path = opts.config_path(); 79 let wallet_file_exists = wallet_path.exists(); 80 validate_wallet_for_mode(&opts, &wallet_path, wallet_file_exists)?; 81 let chain_db_path = opts.chain_db_path(); 82 let ui_data_db_path = ui_data_db_path(&chain_db_path); 83 let chain_store = SqliteChainStore::open(&chain_db_path)?; 84 let ui_data_store = SqliteUiDataStore::open(&ui_data_db_path)?; 85 let persisted_chain_exists = chain_store.contains_chain()?; 86 if opts.chain_mode == ChainMode::Genesis && persisted_chain_exists { 87 bail!( 88 "--genesis refuses to run because chain database already contains a blockchain at {}; start without --genesis to resume it", 89 chain_store.path().display() 90 ); 91 } 92 let mut ui_config = config_store::load_or_create(&config_path)?; 93 let startup_wallet_password = startup_wallet_password_from_env()?; 94 let startup_automatic_burn_enabled = startup_bool_from_env(AUTOMATIC_BURN_ENABLED_ENV)?; 95 let startup_pow_mining_enabled = startup_bool_from_env(POW_MINING_ENABLED_ENV)?; 96 let startup_pow_mining_workers = startup_pow_mining_workers_from_env()?; 97 let startup_setup_complete = startup_bool_from_env(SETUP_COMPLETE_ENV)?; 98 let startup_wallet_endpoint_enabled = startup_bool_from_env(WALLET_ENDPOINT_ENABLED_ENV)?; 99 let startup_wallet_endpoint_port = startup_port_from_env(WALLET_ENDPOINT_PORT_ENV)?; 100 let wallet_endpoint_env_dirty = apply_startup_wallet_endpoint_config_overrides( 101 &mut ui_config, 102 startup_wallet_endpoint_enabled, 103 startup_wallet_endpoint_port, 104 ); 105 let p2p_config_dirty = apply_cli_p2p_config_overrides(&opts, &mut ui_config); 106 let stratum_config_dirty = apply_cli_stratum_config_overrides(&opts, &mut ui_config); 107 let wallet_endpoint_cli_dirty = 108 apply_cli_wallet_endpoint_config_overrides(&opts, &mut ui_config); 109 let p2p_announce_addr = configured_p2p_announce_addr(&opts, &ui_config)?; 110 let configured_p2p_addr = configured_p2p_bind_addr(&opts, &ui_config); 111 let configured_stratum_addr = configured_stratum_addr(&opts, &ui_config); 112 let configured_wallet_endpoint_addr = configured_wallet_endpoint_addr(&opts, &ui_config); 113 if let Some(wallet_addr) = configured_wallet_endpoint_addr { 114 if !ui_config.p2p_accept_inbound { 115 bail!("wallet endpoint requires the node to be configured as Public"); 116 } 117 if wallet_addr.port() == 0 { 118 bail!("wallet endpoint port must be between 1 and 65535"); 119 } 120 if wallet_addr.port() == opts.http_addr.port() { 121 bail!("wallet endpoint port must differ from the management UI port"); 122 } 123 if ui_config.p2p_accept_inbound && wallet_addr.port() == configured_p2p_addr.port() { 124 bail!("wallet endpoint port must differ from the P2P listener port"); 125 } 126 if configured_stratum_addr.is_some_and(|addr| addr.port() == wallet_addr.port()) { 127 bail!("wallet endpoint port must differ from the Stratum listener port"); 128 } 129 } 130 let p2p_accept_inbound = ui_config.p2p_accept_inbound; 131 let advertised_p2p_addr = p2p_announce_addr.unwrap_or(configured_p2p_addr); 132 if opts.chain_mode == ChainMode::Genesis { 133 ui_config.setup_complete = false; 134 ui_config.mining_enabled = true; 135 ui_config.pow_mining_enabled = false; 136 ui_config.burn_per_block = GENESIS_INITIAL_BURN_PER_BLOCK; 137 ui_config.burn_fee = GENESIS_INITIAL_BURN_FEE; 138 } 139 let mining_config_dirty = apply_startup_mining_config_overrides( 140 &mut ui_config, 141 startup_automatic_burn_enabled, 142 startup_pow_mining_enabled, 143 startup_pow_mining_workers, 144 ); 145 let setup_config_dirty = apply_startup_setup_config_override( 146 &opts, 147 persisted_chain_exists, 148 &mut ui_config, 149 startup_setup_complete, 150 )?; 151 let auth_config_dirty = apply_startup_wallet_password_config( 152 &config_path, 153 &mut ui_config, 154 startup_wallet_password 155 .as_ref() 156 .map(ExposeSecret::expose_secret), 157 )?; 158 let wallet_load = load_startup_wallet( 159 &wallet_path, 160 startup_wallet_password 161 .as_ref() 162 .map(ExposeSecret::expose_secret), 163 )?; 164 let wallet_address = wallet_load.address().to_string(); 165 let ui_config_dirty = opts.chain_mode == ChainMode::Genesis 166 || p2p_config_dirty 167 || stratum_config_dirty 168 || wallet_endpoint_env_dirty 169 || wallet_endpoint_cli_dirty 170 || mining_config_dirty 171 || setup_config_dirty; 172 if ui_config_dirty || auth_config_dirty { 173 config_store::save(&config_path, &ui_config)?; 174 } 175 let initialized_ledger = initialize_ledger( 176 &opts, 177 &wallet_address, 178 &chain_store, 179 advertised_p2p_addr, 180 startup_local_testnet, 181 true, 182 ) 183 .await?; 184 let migration_from = initialized_ledger.migration_from.clone(); 185 let migration_required = migration_from.is_some(); 186 let ledger = initialized_ledger.ledger; 187 let has_chain = migration_from.is_none() && (opts.has_chain() || persisted_chain_exists); 188 let initial_burn_per_block = initial_burn_per_block(&opts, &ui_config); 189 let initial_burn_fee = initial_burn_fee(&opts, &ui_config); 190 191 let mut node_core = match wallet_load { 192 StartupWallet::Unlocked { wallet } => NodeCore::from_ledger_with_burn_fee_and_enabled( 193 wallet, 194 ledger, 195 ui_config.mining_enabled, 196 initial_burn_per_block, 197 initial_burn_fee, 198 ), 199 StartupWallet::Locked { address } => NodeCore::from_locked_wallet_address( 200 address, 201 ledger, 202 ui_config.mining_enabled, 203 initial_burn_per_block, 204 initial_burn_fee, 205 ), 206 }; 207 node_core.set_pow_mining_workers(ui_config.pow_mining_workers); 208 node_core.set_pow_mining_enabled(ui_config.pow_mining_enabled); 209 node_core.set_recovery_vdf_top_rank_percent(ui_config.recovery_vdf_top_rank_percent); 210 if let Some(from_network) = migration_from { 211 println!( 212 "network upgrade requires local chain reset: {from_network} -> {}", 213 iuna::app::NETWORK_ID 214 ); 215 node_core.require_network_migration(from_network); 216 } 217 if !migration_required { 218 let restored = restore_pending_transactions_v2(&mut node_core, &chain_store)?; 219 if restored > 0 { 220 println!("restored {restored} pending transaction-v2 envelope(s)"); 221 } 222 } 223 let node: SharedNode = Arc::new(Mutex::new(node_core)); 224 let ui_config = Arc::new(Mutex::new(ui_config)); 225 let mut peers = ui_config.lock().await.peers.clone(); 226 peers.extend(opts.peers); 227 let peers: SharedPeerBook = Arc::new(Mutex::new(PeerBook::from_addresses(peers))); 228 if has_chain { 229 let initial_snapshot = { node.lock().await.chain_snapshot() }; 230 let keep_metrics = ui_config.lock().await.keep_track_of_metrics; 231 persist_chain_snapshot(&chain_store, initial_snapshot.clone()).await?; 232 warm_ui_data_store(&ui_data_store, initial_snapshot, keep_metrics).await?; 233 } else if !migration_required { 234 clear_ui_data_store(&ui_data_store).await?; 235 } 236 237 println!("iuna wallet: {}", node.lock().await.wallet_address()); 238 if node.lock().await.wallet_is_locked() { 239 println!("wallet locked: unlock it in the management UI"); 240 } 241 println!("wallet file: {}", wallet_path.display()); 242 println!("config file: {}", config_path.display()); 243 println!("chain database: {}", chain_store.path().display()); 244 println!("UI data database: {}", ui_data_store.path().display()); 245 println!("management UI: http://{}", opts.http_addr); 246 if p2p_accept_inbound { 247 println!("p2p listener: {configured_p2p_addr}"); 248 } else { 249 println!("p2p listener: disabled (outbound-only)"); 250 } 251 if let Some(addr) = configured_wallet_endpoint_addr { 252 println!("public wallet endpoint configured on separate listener: {addr}"); 253 } else { 254 println!("public wallet endpoint: disabled"); 255 } 256 if p2p_accept_inbound { 257 if let Some(addr) = p2p_announce_addr { 258 println!("p2p announce address: {addr}"); 259 } 260 } 261 println!( 262 "automatic finalization: VDF-driven, burning {} IUNA per block with {} IUNA per byte fee rate", 263 format_iuna(initial_burn_per_block), 264 format_iuna(initial_burn_fee) 265 ); 266 267 let gossip = p2p::GossipNetwork::start( 268 Arc::clone(&node), 269 Arc::clone(&peers), 270 configured_p2p_addr, 271 p2p_announce_addr, 272 p2p_accept_inbound, 273 ) 274 .await?; 275 let mut stratum_status = StratumStatus { 276 enabled: false, 277 listen_addr: None, 278 }; 279 if let Some(stratum_addr) = configured_stratum_addr { 280 let stratum = 281 stratum::StratumServer::start(Arc::clone(&node), gossip.clone(), stratum_addr).await?; 282 println!("stratum listener: {}", stratum.listen_addr()); 283 stratum_status = StratumStatus { 284 enabled: true, 285 listen_addr: Some(stratum.listen_addr().to_string()), 286 }; 287 } 288 289 let persistence_node = Arc::clone(&node); 290 let persistence_store = chain_store.clone(); 291 let persistence_ui_data_store = ui_data_store.clone(); 292 let persistence_config = Arc::clone(&ui_config); 293 let persistence_gossip = gossip.clone(); 294 let persistence_initial_tip = { 295 let node = node.lock().await; 296 if node.has_real_chain() { 297 Some(node.chain_tip_hash()) 298 } else { 299 None 300 } 301 }; 302 let persistence_initial_keep_metrics = ui_config.lock().await.keep_track_of_metrics; 303 tokio::spawn(async move { 304 run_chain_persistence( 305 persistence_node, 306 persistence_store, 307 persistence_ui_data_store, 308 persistence_config, 309 persistence_gossip, 310 persistence_initial_tip, 311 persistence_initial_keep_metrics, 312 ) 313 .await; 314 }); 315 316 let finalizer_node = Arc::clone(&node); 317 let finalizer_gossip = gossip.clone(); 318 tokio::spawn(async move { 319 run_automatic_finalizer(finalizer_node, finalizer_gossip, debug_logging).await; 320 }); 321 322 let pow_miner_node = Arc::clone(&node); 323 let pow_miner_gossip = gossip.clone(); 324 tokio::spawn(async move { 325 run_automatic_pow_miner(pow_miner_node, pow_miner_gossip, debug_logging).await; 326 }); 327 328 let sync_node = Arc::clone(&node); 329 let sync_gossip = gossip.clone(); 330 tokio::spawn(async move { 331 run_peer_sync(sync_node, sync_gossip, debug_logging).await; 332 }); 333 334 if !has_chain { 335 println!("setup mode: waiting to join or create a chain"); 336 } 337 338 let wallet_endpoint_ui_data_store = ui_data_store.clone(); 339 let management = http::serve( 340 Arc::clone(&node), 341 peers, 342 gossip.clone(), 343 ui_config, 344 http::ServeOptions { 345 config_path, 346 chain_store, 347 ui_data_store, 348 wallet_path, 349 stratum: stratum_status, 350 wallet_endpoint_addr: configured_wallet_endpoint_addr, 351 addr: opts.http_addr, 352 }, 353 ); 354 if let Some(addr) = configured_wallet_endpoint_addr { 355 tokio::try_join!( 356 management, 357 wallet_endpoint::serve(node, gossip, wallet_endpoint_ui_data_store, addr) 358 )?; 359 Ok(()) 360 } else { 361 management.await 362 } 363 } 364 365 enum StartupWallet { 366 Unlocked { wallet: iuna::domain::Wallet }, 367 Locked { address: String }, 368 } 369 370 impl StartupWallet { 371 fn address(&self) -> &str { 372 match self { 373 Self::Unlocked { wallet, .. } => wallet.address(), 374 Self::Locked { address } => address, 375 } 376 } 377 } 378 379 fn load_startup_wallet( 380 wallet_path: &Path, 381 startup_wallet_password: Option<&str>, 382 ) -> Result<StartupWallet> { 383 if let Some(password) = startup_wallet_password { 384 if wallet_path.exists() { 385 wallet_store::encrypt_existing_with_password(wallet_path, password)?; 386 let wallet = wallet_store::load_with_password(wallet_path, password)?; 387 return Ok(StartupWallet::Unlocked { wallet }); 388 } 389 let (wallet, _) = 390 wallet_store::replace_with_generated_seed_phrase_encrypted(wallet_path, password)?; 391 return Ok(StartupWallet::Unlocked { wallet }); 392 } 393 394 match wallet_store::load_or_create(wallet_path) { 395 Ok(wallet) => Ok(StartupWallet::Unlocked { wallet }), 396 Err(error) => { 397 let Some(metadata) = wallet_store::metadata(wallet_path)? else { 398 return Err(error); 399 }; 400 if metadata.encrypted { 401 Ok(StartupWallet::Locked { 402 address: metadata.address, 403 }) 404 } else { 405 Err(error) 406 } 407 } 408 } 409 } 410 411 fn startup_wallet_password_from_env() -> Result<Option<SecretString>> { 412 let Some(password) = std::env::var_os(WALLET_PASSWORD_ENV) else { 413 return Ok(None); 414 }; 415 let password = password 416 .into_string() 417 .map_err(|_| anyhow::anyhow!("{WALLET_PASSWORD_ENV} must be valid UTF-8"))?; 418 http::validate_management_password(&password) 419 .with_context(|| format!("{WALLET_PASSWORD_ENV} is not a valid wallet password"))?; 420 Ok(Some(password.into())) 421 } 422 423 fn startup_bool_from_env(name: &str) -> Result<Option<bool>> { 424 let Some(value) = std::env::var_os(name) else { 425 return Ok(None); 426 }; 427 let value = value 428 .into_string() 429 .map_err(|_| anyhow::anyhow!("{name} must be valid UTF-8"))?; 430 let normalized = value.trim().to_ascii_lowercase(); 431 parse_startup_bool_env_value(name, &normalized).map(Some) 432 } 433 434 fn parse_startup_bool_env_value(name: &str, normalized: &str) -> Result<bool> { 435 match normalized { 436 "1" | "true" | "yes" | "on" => Ok(true), 437 "0" | "false" | "no" | "off" => Ok(false), 438 _ => bail!("{name} must be one of true, false, 1, 0, yes, no, on, or off"), 439 } 440 } 441 442 fn startup_pow_mining_workers_from_env() -> Result<Option<u8>> { 443 let Some(value) = std::env::var_os(POW_MINING_WORKERS_ENV) else { 444 return Ok(None); 445 }; 446 let value = value 447 .into_string() 448 .map_err(|_| anyhow::anyhow!("{POW_MINING_WORKERS_ENV} must be valid UTF-8"))?; 449 parse_startup_pow_mining_workers_env_value(value.trim()).map(Some) 450 } 451 452 fn startup_port_from_env(name: &str) -> Result<Option<u16>> { 453 let Some(value) = std::env::var_os(name) else { 454 return Ok(None); 455 }; 456 let value = value 457 .into_string() 458 .map_err(|_| anyhow::anyhow!("{name} must be valid UTF-8"))?; 459 let port = value 460 .trim() 461 .parse::<u16>() 462 .with_context(|| format!("{name} must be an integer between 1 and 65535"))?; 463 if port == 0 { 464 bail!("{name} must be between 1 and 65535"); 465 } 466 Ok(Some(port)) 467 } 468 469 fn apply_startup_wallet_endpoint_config_overrides( 470 ui_config: &mut config_store::UiConfig, 471 enabled: Option<bool>, 472 port: Option<u16>, 473 ) -> bool { 474 let mut dirty = false; 475 if let Some(enabled) = enabled { 476 dirty |= ui_config.wallet_endpoint_enabled != enabled; 477 ui_config.wallet_endpoint_enabled = enabled; 478 } 479 if let Some(port) = port { 480 dirty |= ui_config.wallet_endpoint_bind_port != port; 481 ui_config.wallet_endpoint_bind_port = port; 482 } 483 dirty 484 } 485 486 fn parse_startup_pow_mining_workers_env_value(value: &str) -> Result<u8> { 487 let workers = value 488 .parse::<u8>() 489 .with_context(|| format!("{POW_MINING_WORKERS_ENV} must be an integer"))?; 490 if !(1..=config_store::MAX_POW_MINING_WORKERS).contains(&workers) { 491 bail!( 492 "{POW_MINING_WORKERS_ENV} must be between 1 and {}", 493 config_store::MAX_POW_MINING_WORKERS 494 ); 495 } 496 Ok(workers) 497 } 498 499 fn apply_startup_mining_config_overrides( 500 ui_config: &mut config_store::UiConfig, 501 automatic_burn_enabled: Option<bool>, 502 pow_mining_enabled: Option<bool>, 503 pow_mining_workers: Option<u8>, 504 ) -> bool { 505 let mut dirty = false; 506 if let Some(enabled) = automatic_burn_enabled { 507 if ui_config.mining_enabled != enabled { 508 ui_config.mining_enabled = enabled; 509 dirty = true; 510 } 511 } 512 if let Some(enabled) = pow_mining_enabled { 513 if ui_config.pow_mining_enabled != enabled { 514 ui_config.pow_mining_enabled = enabled; 515 dirty = true; 516 } 517 } 518 if let Some(workers) = pow_mining_workers { 519 if ui_config.pow_mining_workers != workers { 520 ui_config.pow_mining_workers = workers; 521 dirty = true; 522 } 523 } 524 dirty 525 } 526 527 fn apply_startup_setup_config_override( 528 opts: &CliOptions, 529 persisted_chain_exists: bool, 530 ui_config: &mut config_store::UiConfig, 531 setup_complete: Option<bool>, 532 ) -> Result<bool> { 533 let Some(setup_complete) = setup_complete else { 534 return Ok(false); 535 }; 536 if setup_complete 537 && opts.chain_mode == ChainMode::Setup 538 && !persisted_chain_exists 539 && opts.join_peers.is_empty() 540 { 541 bail!( 542 "{SETUP_COMPLETE_ENV}=true requires --genesis, --join, or an existing chain database" 543 ); 544 } 545 if ui_config.setup_complete == setup_complete { 546 return Ok(false); 547 } 548 ui_config.setup_complete = setup_complete; 549 Ok(true) 550 } 551 552 fn apply_startup_wallet_password_config( 553 config_path: &Path, 554 ui_config: &mut config_store::UiConfig, 555 password: Option<&str>, 556 ) -> Result<bool> { 557 let Some(password) = password else { 558 return Ok(false); 559 }; 560 let Some(existing_hash) = ui_config.auth_password_hash.as_deref() else { 561 ui_config.auth_password_hash = Some(http::hash_management_password(password)?); 562 return Ok(true); 563 }; 564 if !http::verify_management_password(password, existing_hash)? { 565 bail!( 566 "{WALLET_PASSWORD_ENV} does not match the configured management UI and wallet password in {}", 567 config_path.display() 568 ); 569 } 570 Ok(false) 571 } 572 573 fn format_iuna(amount: Amount) -> String { 574 let whole = amount / MICRO_IUNA; 575 let fractional = amount % MICRO_IUNA; 576 if fractional == 0 { 577 whole.to_string() 578 } else { 579 let mut fractional = format!("{fractional:06}"); 580 while fractional.ends_with('0') { 581 fractional.pop(); 582 } 583 format!("{whole}.{fractional}") 584 } 585 } 586 587 fn ui_data_db_path(chain_db_path: &Path) -> PathBuf { 588 chain_db_path.with_file_name("ui_data.sqlite3") 589 } 590 591 async fn initialize_ledger( 592 opts: &CliOptions, 593 wallet_address: &str, 594 chain_store: &SqliteChainStore, 595 advertised_p2p_addr: SocketAddr, 596 local_testnet: bool, 597 enforce_pinned_genesis: bool, 598 ) -> Result<InitializedLedger> { 599 if let Some(loaded) = chain_store.load_with_verification_status()? { 600 let snapshot = loaded.snapshot; 601 if opts.chain_mode == ChainMode::Genesis { 602 bail!( 603 "--genesis refuses to run because chain database already contains a blockchain at {}; start without --genesis to resume it", 604 chain_store.path().display() 605 ); 606 } 607 let expected_profile = if local_testnet { 608 LaunchProfile::local_testnet() 609 } else { 610 LaunchProfile::default() 611 }; 612 if snapshot.launch_profile.profile_id != expected_profile.profile_id { 613 return Ok(InitializedLedger { 614 ledger: setup_ledger(local_testnet), 615 migration_from: Some(snapshot.launch_profile.profile_id), 616 }); 617 } 618 if enforce_pinned_genesis { 619 let stored_genesis = snapshot 620 .blocks 621 .first() 622 .context("persisted chain is missing its genesis block")?; 623 validate_network_genesis(&snapshot.launch_profile.profile_id, &stored_genesis.hash)?; 624 } 625 let height = snapshot_height(&snapshot); 626 match loaded.revalidation_from_height { 627 Some(from_height) => println!( 628 "validating local chain from height {from_height} for the current consensus ruleset..." 629 ), 630 None => println!( 631 "local chain is trusted under the current consensus ruleset; skipping historical validation" 632 ), 633 } 634 let ledger = Ledger::from_persisted_snapshot_revalidating_from( 635 snapshot, 636 loaded.revalidation_from_height, 637 ) 638 .with_context(|| { 639 format!( 640 "failed to load chain database {}", 641 chain_store.path().display() 642 ) 643 })?; 644 println!( 645 "resumed chain from {} at height {height}", 646 chain_store.path().display() 647 ); 648 Ok(InitializedLedger { 649 ledger, 650 migration_from: None, 651 }) 652 } else { 653 let ledger = match opts.chain_mode { 654 ChainMode::Setup => Ok(setup_ledger(local_testnet)), 655 ChainMode::Genesis => start_genesis_ledger(wallet_address, local_testnet), 656 ChainMode::Join => { 657 let expected_profile = if local_testnet { 658 LaunchProfile::local_testnet() 659 } else { 660 LaunchProfile::default() 661 }; 662 join_chain_ledger( 663 &opts.join_peers, 664 advertised_p2p_addr, 665 &expected_profile.profile_id, 666 ) 667 .await 668 } 669 }?; 670 Ok(InitializedLedger { 671 ledger, 672 migration_from: None, 673 }) 674 } 675 } 676 677 #[derive(Debug)] 678 struct InitializedLedger { 679 ledger: Ledger, 680 migration_from: Option<String>, 681 } 682 683 fn restore_pending_transactions_v2( 684 node: &mut NodeCore, 685 chain_store: &SqliteChainStore, 686 ) -> Result<usize> { 687 let mut restored = 0_usize; 688 for (transaction_id, envelope) in chain_store.load_pending_transactions_v2()? { 689 match node.receive_gossiped_transaction_v2(envelope) { 690 Ok(()) => restored = restored.saturating_add(1), 691 Err(error) if debug_logging_enabled() => { 692 eprintln!("dropping stale pending transaction v2 {transaction_id}: {error:#}"); 693 } 694 Err(_) => {} 695 } 696 } 697 chain_store.replace_pending_transactions_v2(&node.pending_transaction_v2_envelopes()?)?; 698 Ok(restored) 699 } 700 701 impl std::ops::Deref for InitializedLedger { 702 type Target = Ledger; 703 704 fn deref(&self) -> &Self::Target { 705 &self.ledger 706 } 707 } 708 709 impl std::ops::DerefMut for InitializedLedger { 710 fn deref_mut(&mut self) -> &mut Self::Target { 711 &mut self.ledger 712 } 713 } 714 715 fn snapshot_height(snapshot: &ChainSnapshot) -> u64 { 716 snapshot 717 .blocks 718 .last() 719 .map(|block| block.height) 720 .unwrap_or(0) 721 } 722 723 fn setup_ledger(local_testnet: bool) -> Ledger { 724 if local_testnet { 725 Ledger::new_with_genesis_burns_and_profile( 726 BTreeMap::new(), 727 Vec::new(), 728 1, 729 LaunchProfile::local_testnet(), 730 ) 731 .expect("empty local-testnet setup ledger is valid") 732 } else { 733 Ledger::new(BTreeMap::new(), 1) 734 } 735 } 736 737 fn start_genesis_ledger(wallet_address: &str, local_testnet: bool) -> Result<Ledger> { 738 if !local_testnet { 739 bail!( 740 "mainnet-candidate genesis is pinned; use --join instead of creating a new candidate chain" 741 ); 742 } 743 let vdf_rounds = measure_initial_vdf_rounds(); 744 let mut genesis = BTreeMap::new(); 745 genesis.insert(wallet_address.to_string(), GENESIS_BOOTSTRAP_BALANCE); 746 let launch_profile = if local_testnet { 747 LaunchProfile::local_testnet() 748 } else { 749 LaunchProfile::default() 750 }; 751 let ledger = Ledger::new_with_genesis_burns_and_profile( 752 genesis, 753 vec![GenesisBurn::new( 754 wallet_address, 755 GENESIS_BOOTSTRAP_BURN_AMOUNT, 756 )], 757 vdf_rounds, 758 launch_profile, 759 )?; 760 validate_network_genesis(&ledger.launch_profile().profile_id, ledger.genesis_hash())?; 761 Ok(ledger) 762 } 763 764 fn measure_initial_vdf_rounds() -> u64 { 765 let seed = "iuna-vdf-calibration"; 766 let (measured_rounds, elapsed) = measure_vdf_rounds( 767 seed, 768 VDF_MEASUREMENT_INITIAL_ROUNDS, 769 VDF_MEASUREMENT_MIN_ELAPSED, 770 VDF_MEASUREMENT_MAX_ROUNDS, 771 ); 772 let rounds = extrapolate_vdf_rounds( 773 measured_rounds, 774 elapsed, 775 Duration::from_millis(VDF_TARGET_BLOCK_MS), 776 ); 777 println!( 778 "measured {measured_rounds} VDF rounds in {:.3}ms; initial VDF rounds: {rounds}", 779 elapsed.as_secs_f64() * 1000.0 780 ); 781 rounds 782 } 783 784 fn measure_vdf_rounds( 785 seed: &str, 786 initial_rounds: u64, 787 min_elapsed: Duration, 788 max_rounds_per_attempt: u64, 789 ) -> (u64, Duration) { 790 let mut rounds = initial_rounds.max(1).min(max_rounds_per_attempt.max(1)); 791 let mut measured_rounds = 0_u64; 792 let mut measured_elapsed = Duration::ZERO; 793 794 loop { 795 let started = Instant::now(); 796 let _ = run_vdf(seed, rounds); 797 measured_elapsed += started.elapsed(); 798 measured_rounds = measured_rounds.saturating_add(rounds); 799 800 if measured_elapsed >= min_elapsed || rounds >= max_rounds_per_attempt { 801 return (measured_rounds, measured_elapsed); 802 } 803 rounds = rounds.saturating_mul(2).min(max_rounds_per_attempt); 804 } 805 } 806 807 fn extrapolate_vdf_rounds(measured_rounds: u64, elapsed: Duration, target: Duration) -> u64 { 808 let elapsed_ns = elapsed.as_nanos().max(1); 809 let target_ns = target.as_nanos().max(1); 810 let rounds = u128::from(measured_rounds) 811 .saturating_mul(target_ns) 812 .saturating_div(elapsed_ns) 813 .max(1); 814 rounds.min(u128::from(MAX_VDF_ROUNDS)) as u64 815 } 816 817 async fn join_chain_ledger( 818 join_peers: &[String], 819 advertised_addr: SocketAddr, 820 expected_profile_id: &str, 821 ) -> Result<Ledger> { 822 let mut errors = Vec::new(); 823 for peer in join_peers { 824 match p2p::fetch_snapshot_with_announcement( 825 peer, 826 Some(advertised_addr), 827 expected_profile_id, 828 ) 829 .await 830 { 831 Ok(snapshot) => { 832 let height = snapshot 833 .blocks 834 .last() 835 .map(|block| block.height) 836 .unwrap_or(0); 837 println!("joined chain from {peer} at height {height}"); 838 return p2p::validate_chain_snapshot(snapshot).await; 839 } 840 Err(error) => { 841 errors.push(format!("{peer}: {error:#}")); 842 } 843 } 844 } 845 846 bail!( 847 "could not join any requested peer; refusing to start a separate chain: {}", 848 errors.join("; ") 849 ) 850 } 851 852 async fn run_automatic_finalizer(node: SharedNode, gossip: p2p::GossipNetwork, debug: bool) { 853 let mut last_logged_skip: Option<(u64, String)> = None; 854 loop { 855 if !node.lock().await.has_real_chain() { 856 tokio::time::sleep(std::time::Duration::from_secs(1)).await; 857 continue; 858 } 859 let (height, plan, outbox) = { 860 let mut node = node.lock().await; 861 let height = node.chain_height(); 862 let plan = node.prepare_automatic_finalization(now_ms()); 863 let outbox = node.drain_outbox(); 864 (height, plan, outbox) 865 }; 866 867 if let Err(error) = gossip.broadcast(outbox).await { 868 if debug { 869 eprintln!("p2p broadcast failed after automatic burn: {error:#}"); 870 } 871 } 872 873 let Some(work) = plan.work else { 874 if let Some(reason) = &plan.skipped_reason { 875 if debug 876 && should_log_automatic_finalization_skip(&mut last_logged_skip, height, reason) 877 { 878 println!("auto-finalization skipped at height {height}: {reason}"); 879 } 880 } 881 tokio::time::sleep(std::time::Duration::from_secs(1)).await; 882 continue; 883 }; 884 885 let candidate_height = work.height(); 886 let candidate_parent = work.prev_hash().to_string(); 887 let seed = work.vdf_seed().to_string(); 888 let rounds = work.vdf_rounds(); 889 let publish_at_ms = work.timestamp_ms(); 890 let precheck_at_ms = automatic_finalization_precheck_time(now_ms(), publish_at_ms); 891 let precheck = { 892 let node = node.lock().await; 893 node.precheck_prepared_block_without_vdf_at(&work, precheck_at_ms) 894 }; 895 if let Err(error) = precheck { 896 let message = format!("skipped before VDF: {error:#}"); 897 if debug 898 && should_log_automatic_finalization_skip(&mut last_logged_skip, height, &message) 899 { 900 println!("auto-finalization {message}"); 901 } 902 node.lock() 903 .await 904 .record_automatic_finalization_status(message); 905 tokio::time::sleep(std::time::Duration::from_secs(1)).await; 906 continue; 907 } 908 last_logged_skip = None; 909 if debug { 910 println!( 911 "leader selected locally for candidate block {}; running VDF for {} rounds", 912 work.height(), 913 work.vdf_rounds() 914 ); 915 } 916 let (progress_tx, progress_rx) = std::sync::mpsc::channel(); 917 let cancellation = Arc::new(AtomicBool::new(false)); 918 let worker_cancellation = Arc::clone(&cancellation); 919 let mut vdf_worker = tokio::task::spawn_blocking(move || { 920 run_vdf_cancellable_with_progress( 921 &seed, 922 rounds, 923 VDF_PROGRESS_LOG_INTERVAL, 924 worker_cancellation.as_ref(), 925 |progress| { 926 let _ = progress_tx.send(progress); 927 }, 928 ) 929 }); 930 let mut cancelled_for_new_tip = false; 931 let mut cancelled_for_disabled = false; 932 let vdf_output = loop { 933 tokio::select! { 934 result = &mut vdf_worker => { 935 break match result { 936 Ok(output) => output, 937 Err(error) => { 938 if debug { 939 eprintln!("VDF worker failed: {error:#}"); 940 } 941 None 942 } 943 }; 944 } 945 _ = tokio::time::sleep(std::time::Duration::from_millis(500)) => { 946 while let Ok(progress) = progress_rx.try_recv() { 947 let message = format_vdf_progress(candidate_height, progress); 948 if debug { 949 println!("{message}"); 950 } 951 node.lock().await.record_automatic_finalization_status(message); 952 } 953 let (tip_changed, finalization_disabled) = { 954 let node = node.lock().await; 955 ( 956 node.ledger().tip_hash() != candidate_parent, 957 !node.automatic_mining_enabled(), 958 ) 959 }; 960 if tip_changed { 961 cancelled_for_new_tip = true; 962 cancellation.store(true, Ordering::Relaxed); 963 } else if finalization_disabled { 964 cancelled_for_disabled = true; 965 cancellation.store(true, Ordering::Relaxed); 966 } 967 } 968 } 969 }; 970 while let Ok(progress) = progress_rx.try_recv() { 971 let message = format_vdf_progress(candidate_height, progress); 972 if debug { 973 println!("{message}"); 974 } 975 node.lock() 976 .await 977 .record_automatic_finalization_status(message); 978 } 979 let Some(vdf_output) = vdf_output else { 980 let message = if cancelled_for_new_tip { 981 format!("cancelled stale VDF for candidate block {candidate_height}") 982 } else if cancelled_for_disabled { 983 format!( 984 "cancelled VDF for candidate block {candidate_height} because automatic finalization is disabled" 985 ) 986 } else { 987 format!("VDF worker failed for candidate block {candidate_height}") 988 }; 989 if debug { 990 println!("auto-finalization {message}"); 991 } 992 node.lock() 993 .await 994 .record_automatic_finalization_status(message); 995 continue; 996 }; 997 998 let completed_at_ms = now_ms(); 999 let mut stale_before_publish = false; 1000 let mut disabled_before_publish = false; 1001 if completed_at_ms < publish_at_ms { 1002 let wait_ms = publish_at_ms - completed_at_ms; 1003 let message = format!( 1004 "VDF complete for candidate block {}; waiting for finalizer rank time slot ({:.1}s)", 1005 work.height(), 1006 wait_ms as f64 / 1000.0 1007 ); 1008 if debug { 1009 println!("{message}"); 1010 } 1011 node.lock() 1012 .await 1013 .record_automatic_finalization_status(message); 1014 loop { 1015 let current_ms = now_ms(); 1016 if current_ms >= publish_at_ms { 1017 break; 1018 } 1019 let wait_ms = publish_at_ms - current_ms; 1020 tokio::time::sleep(std::time::Duration::from_millis(wait_ms.min(1_000))).await; 1021 let (tip_changed, finalization_disabled) = { 1022 let node = node.lock().await; 1023 ( 1024 node.ledger().tip_hash() != candidate_parent, 1025 !node.automatic_mining_enabled(), 1026 ) 1027 }; 1028 if tip_changed || finalization_disabled { 1029 stale_before_publish = tip_changed; 1030 disabled_before_publish = finalization_disabled; 1031 break; 1032 } 1033 } 1034 } 1035 if stale_before_publish || disabled_before_publish { 1036 let message = if stale_before_publish { 1037 format!("cancelled completed VDF for stale candidate block {candidate_height}") 1038 } else { 1039 format!( 1040 "cancelled completed VDF for candidate block {candidate_height} because automatic finalization is disabled" 1041 ) 1042 }; 1043 if debug { 1044 println!("auto-finalization {message}"); 1045 } 1046 node.lock() 1047 .await 1048 .record_automatic_finalization_status(message); 1049 continue; 1050 } 1051 let publish_timestamp_ms = now_ms().max(publish_at_ms); 1052 1053 let (finalized, outbox) = { 1054 let mut node = node.lock().await; 1055 if !node.automatic_mining_enabled() { 1056 node.record_automatic_finalization_status(format!( 1057 "cancelled completed VDF for candidate block {candidate_height} because automatic finalization is disabled" 1058 )); 1059 continue; 1060 } 1061 let finalized = node.complete_prepared_block_at(work, vdf_output, publish_timestamp_ms); 1062 match &finalized { 1063 Ok(block) => node.record_automatic_finalization_status(format!( 1064 "finalized block {} ({})", 1065 block.height, block.hash 1066 )), 1067 Err(error) => node 1068 .record_automatic_finalization_status(format!("skipped after VDF: {error:#}")), 1069 } 1070 let outbox = node.drain_outbox(); 1071 (finalized, outbox) 1072 }; 1073 1074 let failed_after_vdf = finalized.is_err(); 1075 match finalized { 1076 Ok(block) if debug => { 1077 println!("auto-finalized block {} ({})", block.height, block.hash); 1078 } 1079 Ok(_) => {} 1080 Err(error) if debug => println!("auto-finalization skipped after VDF: {error:#}"), 1081 Err(_) => {} 1082 } 1083 1084 if let Err(error) = gossip.broadcast(outbox).await { 1085 if debug { 1086 eprintln!("p2p broadcast failed after automatic block: {error:#}"); 1087 } 1088 } 1089 1090 if failed_after_vdf { 1091 tokio::time::sleep(std::time::Duration::from_secs(10)).await; 1092 } 1093 1094 tokio::task::yield_now().await; 1095 } 1096 } 1097 1098 fn automatic_finalization_precheck_time(now_ms: u64, publish_at_ms: u64) -> u64 { 1099 // A fallback ticket's rank slot can be well beyond the normal future-drift 1100 // allowance. Validate the candidate as it will stand when that known slot 1101 // opens, so its VDF can run in advance; final application still validates 1102 // the actual publication timestamp against the then-current clock. 1103 now_ms.max(publish_at_ms) 1104 } 1105 1106 fn should_log_automatic_finalization_skip( 1107 last_logged_skip: &mut Option<(u64, String)>, 1108 height: u64, 1109 reason: &str, 1110 ) -> bool { 1111 let skip = (height, reason.to_string()); 1112 if last_logged_skip.as_ref() == Some(&skip) { 1113 return false; 1114 } 1115 *last_logged_skip = Some(skip); 1116 true 1117 } 1118 1119 fn format_vdf_progress(candidate_height: u64, progress: VdfProgress) -> String { 1120 let phase = match progress.phase { 1121 VdfProgressPhase::Output => "output", 1122 VdfProgressPhase::Proof => "proof", 1123 }; 1124 let percent = if progress.total_steps == 0 { 1125 100.0 1126 } else { 1127 progress.completed_steps as f64 * 100.0 / progress.total_steps as f64 1128 }; 1129 format!( 1130 "running VDF for candidate block {candidate_height}: {phase} {}/{} rounds, total {}/{} steps ({percent:.1}%)", 1131 progress.completed_phase_rounds, 1132 progress.phase_rounds, 1133 progress.completed_steps, 1134 progress.total_steps 1135 ) 1136 } 1137 1138 async fn run_automatic_pow_miner(node: SharedNode, gossip: p2p::GossipNetwork, debug: bool) { 1139 loop { 1140 tokio::time::sleep(std::time::Duration::from_secs(1)).await; 1141 let (height, job) = { 1142 let mut node = node.lock().await; 1143 if !node.pow_mining_enabled() { 1144 continue; 1145 } 1146 if !node.has_real_chain() { 1147 continue; 1148 } 1149 let height = node.chain_height(); 1150 let job = match node.prepare_automatic_pow_mining_job() { 1151 Ok(job) => job, 1152 Err(error) => { 1153 node.record_automatic_pow_mining_error(format!( 1154 "automatic PoW mining failed: {error:#}" 1155 )); 1156 None 1157 } 1158 }; 1159 (height, job) 1160 }; 1161 let Some(job) = job else { 1162 continue; 1163 }; 1164 1165 let search = tokio::task::spawn_blocking(move || job.search()).await; 1166 let (pow_mined, outbox) = { 1167 let mut node = node.lock().await; 1168 let pow_mined = match search { 1169 Ok(Ok((job, outcome))) => { 1170 match node.finish_automatic_pow_mining_job(job, outcome) { 1171 Ok(tx) => tx, 1172 Err(error) => { 1173 node.record_automatic_pow_mining_error(format!( 1174 "automatic PoW mining failed: {error:#}" 1175 )); 1176 None 1177 } 1178 } 1179 } 1180 Ok(Err(error)) => { 1181 node.record_automatic_pow_mining_error(format!( 1182 "automatic PoW mining failed: {error:#}" 1183 )); 1184 None 1185 } 1186 Err(error) => { 1187 node.record_automatic_pow_mining_error(format!( 1188 "automatic PoW mining task failed: {error:#}" 1189 )); 1190 None 1191 } 1192 }; 1193 let outbox = node.drain_outbox(); 1194 (pow_mined, outbox) 1195 }; 1196 1197 if let Err(error) = gossip.broadcast(outbox).await { 1198 if debug { 1199 eprintln!("p2p broadcast failed after automatic PoW mining: {error:#}"); 1200 } 1201 } 1202 1203 if debug { 1204 if let Some(tx) = &pow_mined { 1205 println!( 1206 "auto-pow queued mine action for height {} ({})", 1207 height, 1208 tx.signature() 1209 ); 1210 } 1211 } 1212 } 1213 } 1214 1215 async fn run_peer_sync(node: SharedNode, gossip: p2p::GossipNetwork, debug: bool) { 1216 loop { 1217 tokio::time::sleep(std::time::Duration::from_secs(5)).await; 1218 let envelopes = { 1219 let mut node = node.lock().await; 1220 let mut envelopes = vec![node.peer_status()]; 1221 envelopes.extend(node.drain_outbox()); 1222 envelopes.extend(node.mempool_gossip()); 1223 envelopes 1224 }; 1225 let mut envelopes = envelopes; 1226 envelopes.push(gossip.peer_exchange().await); 1227 if let Err(error) = gossip.broadcast(envelopes).await { 1228 if debug { 1229 eprintln!("p2p sync gossip failed: {error:#}"); 1230 } 1231 } 1232 } 1233 } 1234 1235 async fn run_chain_persistence( 1236 node: SharedNode, 1237 store: SqliteChainStore, 1238 ui_data_store: SqliteUiDataStore, 1239 ui_config: Arc<Mutex<config_store::UiConfig>>, 1240 gossip: p2p::GossipNetwork, 1241 initial_saved_tip: Option<String>, 1242 initial_projected_keep_metrics: bool, 1243 ) { 1244 let initial_state = ChainPersistenceState { 1245 projected_tip: initial_saved_tip.clone(), 1246 saved_tip: initial_saved_tip, 1247 projected_keep_metrics: initial_projected_keep_metrics, 1248 sync_checkpoint_interval: SYNC_CHAIN_CHECKPOINT_INTERVAL, 1249 }; 1250 run_chain_persistence_loop( 1251 node, 1252 store, 1253 ui_data_store, 1254 ui_config, 1255 Duration::from_secs(2), 1256 Some(gossip), 1257 initial_state, 1258 ) 1259 .await; 1260 } 1261 1262 #[cfg(test)] 1263 async fn run_chain_persistence_with_interval( 1264 node: SharedNode, 1265 store: SqliteChainStore, 1266 ui_data_store: SqliteUiDataStore, 1267 ui_config: Arc<Mutex<config_store::UiConfig>>, 1268 interval: Duration, 1269 initial_saved_tip: Option<String>, 1270 initial_projected_keep_metrics: bool, 1271 ) { 1272 let initial_state = ChainPersistenceState { 1273 projected_tip: initial_saved_tip.clone(), 1274 saved_tip: initial_saved_tip, 1275 projected_keep_metrics: initial_projected_keep_metrics, 1276 sync_checkpoint_interval: SYNC_CHAIN_CHECKPOINT_INTERVAL, 1277 }; 1278 run_chain_persistence_loop( 1279 node, 1280 store, 1281 ui_data_store, 1282 ui_config, 1283 interval, 1284 None, 1285 initial_state, 1286 ) 1287 .await; 1288 } 1289 1290 struct ChainPersistenceState { 1291 saved_tip: Option<String>, 1292 projected_tip: Option<String>, 1293 projected_keep_metrics: bool, 1294 sync_checkpoint_interval: Duration, 1295 } 1296 1297 async fn run_chain_persistence_loop( 1298 node: SharedNode, 1299 store: SqliteChainStore, 1300 ui_data_store: SqliteUiDataStore, 1301 ui_config: Arc<Mutex<config_store::UiConfig>>, 1302 interval: Duration, 1303 gossip: Option<p2p::GossipNetwork>, 1304 initial_state: ChainPersistenceState, 1305 ) { 1306 let mut last_saved_tip = initial_state.saved_tip; 1307 let mut last_projected_tip = initial_state.projected_tip; 1308 let mut last_projected_keep_metrics = initial_state.projected_keep_metrics; 1309 let mut last_chain_checkpoint = Instant::now(); 1310 let mut last_saved_pending_v2_ids = None::<Vec<String>>; 1311 loop { 1312 tokio::time::sleep(interval).await; 1313 { 1314 let node = node.lock().await; 1315 match node.pending_transaction_v2_envelopes() { 1316 Ok(pending_v2) => { 1317 let pending_v2_ids = pending_v2 1318 .iter() 1319 .map(|(transaction_id, _)| transaction_id.clone()) 1320 .collect::<Vec<_>>(); 1321 if last_saved_pending_v2_ids.as_ref() != Some(&pending_v2_ids) { 1322 match persist_pending_transactions_v2(&store, pending_v2).await { 1323 Ok(()) => last_saved_pending_v2_ids = Some(pending_v2_ids), 1324 Err(error) if debug_logging_enabled() => { 1325 eprintln!("pending transaction-v2 persistence failed: {error:#}"); 1326 } 1327 Err(_) => {} 1328 } 1329 } 1330 } 1331 Err(error) if debug_logging_enabled() => { 1332 eprintln!("pending transaction-v2 encoding failed: {error:#}"); 1333 } 1334 Err(_) => {} 1335 } 1336 } 1337 let syncing = gossip 1338 .as_ref() 1339 .is_some_and(|network| network.chain_sync_active_or_recent(Duration::from_secs(5))); 1340 let defer_sync_checkpoint = should_defer_sync_checkpoint( 1341 syncing, 1342 last_chain_checkpoint.elapsed(), 1343 initial_state.sync_checkpoint_interval, 1344 ); 1345 let snapshot = { 1346 let node = node.lock().await; 1347 if !node.has_real_chain() { 1348 continue; 1349 } 1350 let tip_hash = node.ledger().tip_hash(); 1351 if syncing { 1352 let tip_changed = last_saved_tip.as_deref() != Some(tip_hash); 1353 if !tip_changed || defer_sync_checkpoint { 1354 continue; 1355 } 1356 } 1357 node.chain_snapshot() 1358 }; 1359 let Some(tip_hash) = snapshot.blocks.last().map(|block| block.hash.clone()) else { 1360 continue; 1361 }; 1362 let keep_metrics = ui_config.lock().await.keep_track_of_metrics; 1363 let tip_changed = last_saved_tip.as_deref() != Some(tip_hash.as_str()); 1364 let projected_tip_changed = last_projected_tip.as_deref() != Some(tip_hash.as_str()); 1365 let metrics_mode_changed = last_projected_keep_metrics != keep_metrics; 1366 if !tip_changed && !projected_tip_changed && !metrics_mode_changed { 1367 continue; 1368 } 1369 1370 let result = if syncing && tip_changed { 1371 persist_chain_snapshot(&store, snapshot).await 1372 } else if tip_changed { 1373 persist_chain_and_project_ui_data(&store, &ui_data_store, snapshot, keep_metrics).await 1374 } else { 1375 project_ui_data_store(&ui_data_store, snapshot, keep_metrics).await 1376 }; 1377 match result { 1378 Ok(()) => { 1379 if !syncing { 1380 last_projected_tip = Some(tip_hash.clone()); 1381 last_projected_keep_metrics = keep_metrics; 1382 } 1383 last_saved_tip = Some(tip_hash); 1384 if tip_changed { 1385 last_chain_checkpoint = Instant::now(); 1386 } 1387 } 1388 Err(error) if debug_logging_enabled() => { 1389 eprintln!("chain persistence failed: {error:#}") 1390 } 1391 Err(_) => {} 1392 } 1393 } 1394 } 1395 1396 fn should_defer_sync_checkpoint( 1397 syncing: bool, 1398 since_last_checkpoint: Duration, 1399 checkpoint_interval: Duration, 1400 ) -> bool { 1401 syncing && since_last_checkpoint < checkpoint_interval 1402 } 1403 1404 async fn persist_chain_and_project_ui_data( 1405 store: &SqliteChainStore, 1406 ui_data_store: &SqliteUiDataStore, 1407 snapshot: ChainSnapshot, 1408 keep_metrics: bool, 1409 ) -> Result<()> { 1410 persist_chain_snapshot(store, snapshot.clone()).await?; 1411 project_ui_data_store(ui_data_store, snapshot, keep_metrics).await 1412 } 1413 1414 async fn persist_chain_snapshot(store: &SqliteChainStore, snapshot: ChainSnapshot) -> Result<()> { 1415 let store = store.clone(); 1416 tokio::task::spawn_blocking(move || store.save_verified(&snapshot)) 1417 .await 1418 .context("chain persistence worker failed")??; 1419 Ok(()) 1420 } 1421 1422 async fn persist_pending_transactions_v2( 1423 store: &SqliteChainStore, 1424 pending: Vec<(String, String)>, 1425 ) -> Result<()> { 1426 let store = store.clone(); 1427 tokio::task::spawn_blocking(move || store.replace_pending_transactions_v2(&pending)) 1428 .await 1429 .context("pending transaction-v2 persistence worker failed")??; 1430 Ok(()) 1431 } 1432 1433 async fn warm_ui_data_store( 1434 store: &SqliteUiDataStore, 1435 snapshot: ChainSnapshot, 1436 keep_metrics: bool, 1437 ) -> Result<()> { 1438 println!("warming UI data database..."); 1439 let started = Instant::now(); 1440 let tip_hash = snapshot 1441 .blocks 1442 .last() 1443 .map(|block| block.hash.clone()) 1444 .context("cannot warm UI data database from an empty chain")?; 1445 let readiness_store = store.clone(); 1446 let readiness_tip = tip_hash.clone(); 1447 let (ui_ready, metrics_ready) = tokio::task::spawn_blocking(move || { 1448 Ok::<_, anyhow::Error>(( 1449 readiness_store.is_projected_to(&readiness_tip)?, 1450 readiness_store.metrics_are_projected_to(&readiness_tip)?, 1451 )) 1452 }) 1453 .await 1454 .context("UI data readiness worker failed")??; 1455 1456 if !ui_ready { 1457 project_ui_data_store(store, snapshot, keep_metrics).await?; 1458 } else if keep_metrics && !metrics_ready { 1459 let metrics_store = store.clone(); 1460 tokio::task::spawn_blocking(move || metrics_store.replace_metrics_for_snapshot(&snapshot)) 1461 .await 1462 .context("metrics warm-up worker failed")??; 1463 } else if !keep_metrics { 1464 let metrics_store = store.clone(); 1465 tokio::task::spawn_blocking(move || metrics_store.clear_metrics()) 1466 .await 1467 .context("metrics cleanup worker failed")??; 1468 } 1469 println!( 1470 "UI data database ready in {:.2}s", 1471 started.elapsed().as_secs_f64() 1472 ); 1473 Ok(()) 1474 } 1475 1476 async fn project_ui_data_store( 1477 store: &SqliteUiDataStore, 1478 snapshot: ChainSnapshot, 1479 keep_metrics: bool, 1480 ) -> Result<()> { 1481 let store = store.clone(); 1482 tokio::task::spawn_blocking(move || store.project_snapshot(&snapshot, keep_metrics)) 1483 .await 1484 .context("UI data projection worker failed")??; 1485 Ok(()) 1486 } 1487 1488 async fn clear_ui_data_store(store: &SqliteUiDataStore) -> Result<()> { 1489 println!("clearing UI data database..."); 1490 let started = Instant::now(); 1491 let store = store.clone(); 1492 tokio::task::spawn_blocking(move || store.clear_all()) 1493 .await 1494 .context("UI data cleanup worker failed")??; 1495 println!( 1496 "UI data database ready in {:.2}s", 1497 started.elapsed().as_secs_f64() 1498 ); 1499 Ok(()) 1500 } 1501 1502 #[cfg(test)] 1503 #[path = "main_tests.rs"] 1504 mod tests;