diff --git a/settings.ini b/settings.ini index f0150b8..afed23a 100644 --- a/settings.ini +++ b/settings.ini @@ -1,40 +1,42 @@ ; settings.ini [Paths] -BLOCK_PATH = "./blocks" -TORRENT_PATH = "./torrents" -DB_PATH = "./state" -BALANCE_SHEET = "./balance_sheet" -LOG_PATH = "./logs" -WALLET_PATH = "./wallets" -WALLET_NAME = "contractless.wallet" +BLOCK_PATH = "./blocks" +TORRENT_PATH = "./torrents" +DB_PATH = "./state" +BALANCE_SHEET = "./balance_sheet" +LOG_PATH = "./logs" +WALLET_PATH = "./wallets" +WALLET_NAME = "contractless.wallet" [Settings] -LOG_LEVEL = "info" -PUBLIC_IP = "your_public_ip_address" -LISTEN_IP = "0.0.0.0" -RPC_PORT = "50055" -TESTNET_RPC_PORT = "50050" -INCOMING_CONNECTIONS = "100" -OUTGOING_CONNECTIONS = "10" -VALIDATOR = "false" -; THREADS must 1, 2, or a multple of 4 -THREADS = "8" +LOG_LEVEL = "info" +PUBLIC_IP = "your_public_ip_address" +LISTEN_IP = "0.0.0.0" +RPC_PORT = "50055" +TESTNET_RPC_PORT = "50050" +INCOMING_CONNECTIONS = "100" +OUTGOING_CONNECTIONS = "10" +VALIDATOR = "false" +; THREADS must 1, 2, +; or a multple of 4 +THREADS = "8" +STORAGE_LOOKUP_FEE_PER_BYTE = "0.0001" [Piggyback] -PIGGYBACK_1 = "contractless.dev:50050" +PIGGYBACK_1 = "contractless.dev:50050" [Postgres-Testnet] -host = 127.0.0.1 -port = 5432 -user = contractless2 -password = your_paddword -dbname = contractless_db2 +host = 127.0.0.1 +port = 5432 +user = contractless2 +password = your_paddword +dbname = contractless_db2 [Postgres] -host = 127.0.0.1 -port = 5432 -user = contractless -password = your_password -dbname = contractless_db +host = 127.0.0.1 +port = 5432 +user = contractless +password = your_password +dbname = contractless_db diff --git a/src/bin/create_burn_tx.rs b/src/bin/create_burn_tx.rs index d5df971..d5a9904 100644 --- a/src/bin/create_burn_tx.rs +++ b/src/bin/create_burn_tx.rs @@ -1,6 +1,8 @@ use contractless::blocks::burn::{BurnTransaction, UnsignedBurnTransaction}; use contractless::common::asset_names::padded_normalized_asset_input; -use contractless::common::cli_prompts::{prompt_hidden_nonempty, prompt_visible, prompt_wallet_path}; +use contractless::common::cli_prompts::{ + prompt_hidden_nonempty, prompt_visible, prompt_wallet_path, +}; use contractless::common::network_paths_and_settings::block_extension_and_paths; use contractless::common::types::BURN_FEE; diff --git a/src/bin/create_collateral_claim_tx.rs b/src/bin/create_collateral_claim_tx.rs index 5f146df..192871c 100644 --- a/src/bin/create_collateral_claim_tx.rs +++ b/src/bin/create_collateral_claim_tx.rs @@ -1,7 +1,9 @@ use contractless::blocks::collateral::{ CollateralClaimTransaction, UnsignedCollateralClaimTransaction, }; -use contractless::common::cli_prompts::{prompt_hidden_nonempty, prompt_visible, prompt_wallet_path}; +use contractless::common::cli_prompts::{ + prompt_hidden_nonempty, prompt_visible, prompt_wallet_path, +}; use contractless::common::types::COLLATERAL_FEE; use contractless::json; diff --git a/src/bin/create_issue_token_tx.rs b/src/bin/create_issue_token_tx.rs index 41490c8..78d0d92 100644 --- a/src/bin/create_issue_token_tx.rs +++ b/src/bin/create_issue_token_tx.rs @@ -1,6 +1,8 @@ use contractless::blocks::issue_token::{IssueTokenTransaction, UnsignedIssueTokenTransaction}; use contractless::common::asset_names::padded_normalized_asset_input; -use contractless::common::cli_prompts::{prompt_hidden_nonempty, prompt_visible, prompt_wallet_path}; +use contractless::common::cli_prompts::{ + prompt_hidden_nonempty, prompt_visible, prompt_wallet_path, +}; use contractless::common::types::ISSUE_TOKEN_FEE; use contractless::json; diff --git a/src/bin/create_loan_payment_tx.rs b/src/bin/create_loan_payment_tx.rs index 2168533..4841270 100644 --- a/src/bin/create_loan_payment_tx.rs +++ b/src/bin/create_loan_payment_tx.rs @@ -1,7 +1,9 @@ use contractless::blocks::loan_payment::{ ContractPaymentTransaction, UnsignedContractPaymentTransaction, }; -use contractless::common::cli_prompts::{prompt_hidden_nonempty, prompt_visible, prompt_wallet_path}; +use contractless::common::cli_prompts::{ + prompt_hidden_nonempty, prompt_visible, prompt_wallet_path, +}; use contractless::common::types::BORROWER_FEE; use contractless::json; diff --git a/src/bin/create_loan_tx.rs b/src/bin/create_loan_tx.rs index 5bb9c14..ef2d630 100644 --- a/src/bin/create_loan_tx.rs +++ b/src/bin/create_loan_tx.rs @@ -2,7 +2,9 @@ use contractless::blocks::loans::UnsignedLoanContractTransaction; use contractless::common::asset_names::{ padded_loan_collateral_input_or_base, padded_normalized_asset_input, }; -use contractless::common::cli_prompts::{prompt_hidden_nonempty, prompt_visible, prompt_wallet_path}; +use contractless::common::cli_prompts::{ + prompt_hidden_nonempty, prompt_visible, prompt_wallet_path, +}; use contractless::common::network_paths_and_settings::block_extension_and_paths; use contractless::common::types::LENDER_FEE; diff --git a/src/bin/create_marketing_tx.rs b/src/bin/create_marketing_tx.rs index 6b583bf..106e04a 100644 --- a/src/bin/create_marketing_tx.rs +++ b/src/bin/create_marketing_tx.rs @@ -1,5 +1,7 @@ use contractless::blocks::marketing::{MarketingTransaction, UnsignedMarketingTransaction}; -use contractless::common::cli_prompts::{prompt_hidden_nonempty, prompt_visible, prompt_wallet_path}; +use contractless::common::cli_prompts::{ + prompt_hidden_nonempty, prompt_visible, prompt_wallet_path, +}; use contractless::common::types::MARKETING_FEE; use contractless::json; diff --git a/src/bin/create_nft_tx.rs b/src/bin/create_nft_tx.rs index af0f6a8..f8f8a20 100644 --- a/src/bin/create_nft_tx.rs +++ b/src/bin/create_nft_tx.rs @@ -1,6 +1,8 @@ use contractless::blocks::nft::{CreateNftTransaction, UnsignedCreateNftTransaction}; use contractless::common::asset_names::padded_normalized_asset_input; -use contractless::common::cli_prompts::{prompt_hidden_nonempty, prompt_visible, prompt_wallet_path}; +use contractless::common::cli_prompts::{ + prompt_hidden_nonempty, prompt_visible, prompt_wallet_path, +}; use contractless::common::types::CREATE_NFT_FEE; use contractless::json; diff --git a/src/bin/create_swap_tx.rs b/src/bin/create_swap_tx.rs index 2f4bc42..2764e73 100644 --- a/src/bin/create_swap_tx.rs +++ b/src/bin/create_swap_tx.rs @@ -1,6 +1,8 @@ use contractless::blocks::swap::UnsignedSwapTransaction; use contractless::common::asset_names::padded_normalized_asset_input; -use contractless::common::cli_prompts::{prompt_hidden_nonempty, prompt_visible, prompt_wallet_path}; +use contractless::common::cli_prompts::{ + prompt_hidden_nonempty, prompt_visible, prompt_wallet_path, +}; use contractless::common::network_paths_and_settings::block_extension_and_paths; use contractless::common::types::SWAP_FEE; diff --git a/src/bin/create_tokens_tx.rs b/src/bin/create_tokens_tx.rs index 2e011cd..c8af80d 100644 --- a/src/bin/create_tokens_tx.rs +++ b/src/bin/create_tokens_tx.rs @@ -1,6 +1,8 @@ use contractless::blocks::token::{CreateTokenTransaction, UnsignedCreateTokenTransaction}; use contractless::common::asset_names::padded_normalized_asset_input; -use contractless::common::cli_prompts::{prompt_hidden_nonempty, prompt_visible, prompt_wallet_path}; +use contractless::common::cli_prompts::{ + prompt_hidden_nonempty, prompt_visible, prompt_wallet_path, +}; use contractless::common::types::CREATE_TOKEN_FEE; use contractless::json; diff --git a/src/bin/create_transfer_tx.rs b/src/bin/create_transfer_tx.rs index 40d4433..99e6d8e 100644 --- a/src/bin/create_transfer_tx.rs +++ b/src/bin/create_transfer_tx.rs @@ -1,6 +1,8 @@ use contractless::blocks::transfer::{TransferTransaction, UnsignedTransferTransaction}; use contractless::common::asset_names::padded_normalized_asset_input; -use contractless::common::cli_prompts::{prompt_hidden_nonempty, prompt_visible, prompt_wallet_path}; +use contractless::common::cli_prompts::{ + prompt_hidden_nonempty, prompt_visible, prompt_wallet_path, +}; use contractless::common::network_paths_and_settings::block_extension_and_paths; use contractless::common::types::{NON_BASE_TRANSFER_MIN_FEE, TRANSFER_FEE}; diff --git a/src/bin/create_vanity_tx.rs b/src/bin/create_vanity_tx.rs index c9e5c58..06d3a15 100644 --- a/src/bin/create_vanity_tx.rs +++ b/src/bin/create_vanity_tx.rs @@ -1,5 +1,7 @@ use contractless::blocks::vanity::{UnsignedVanityAddressTransaction, VanityAddressTransaction}; -use contractless::common::cli_prompts::{prompt_hidden_nonempty, prompt_visible, prompt_wallet_path}; +use contractless::common::cli_prompts::{ + prompt_hidden_nonempty, prompt_visible, prompt_wallet_path, +}; use contractless::common::types::{VANITY_ADDRESS_FEE, VANITY_ADDRESS_TYPE}; use contractless::env; diff --git a/src/bin/validate_torrent_and_block_headers.rs b/src/bin/validate_torrent_and_block_headers.rs index 25c316d..8fd7451 100644 --- a/src/bin/validate_torrent_and_block_headers.rs +++ b/src/bin/validate_torrent_and_block_headers.rs @@ -1,3 +1,4 @@ +use colored::*; use contractless::common::binary_conversions::hex_to_u64; use contractless::common::cli_prompts::{prompt_hidden_nonempty, prompt_wallet_path}; use contractless::common::network_paths_and_settings::block_extension_and_paths; @@ -9,7 +10,6 @@ use contractless::standalone_tools::wallet_registry_lookup::lookup_pubkey_from_l use contractless::torrent::structs::Torrent; use contractless::wallets::structures::Wallet; use contractless::{AsyncReadExt, File}; -use colored::*; use std::process; #[tokio::main] diff --git a/src/bin/verify_message.rs b/src/bin/verify_message.rs index 0b7f15b..3988b6d 100644 --- a/src/bin/verify_message.rs +++ b/src/bin/verify_message.rs @@ -1,4 +1,6 @@ -use contractless::common::cli_prompts::{prompt_hidden_nonempty, prompt_visible, prompt_wallet_path}; +use contractless::common::cli_prompts::{ + prompt_hidden_nonempty, prompt_visible, prompt_wallet_path, +}; use contractless::common::skein::skein_256_hash_data; use contractless::env; use contractless::standalone_tools::wallet_registry_lookup::lookup_pubkey_from_live_registry; diff --git a/src/miner/flag.rs b/src/miner/flag.rs index c14b419..b2478d9 100644 --- a/src/miner/flag.rs +++ b/src/miner/flag.rs @@ -3,7 +3,7 @@ use crate::Duration; use crate::{AtomicBool, AtomicOrdering}; use std::sync::atomic::AtomicU8; -static REORG_GATE: AtomicBool = AtomicBool::new(false); +static CHAIN_OPERATION_GATE: AtomicBool = AtomicBool::new(false); static MINING_STOP_REQUESTED: AtomicBool = AtomicBool::new(false); static NODE_MODE: AtomicU8 = AtomicU8::new(NodeMode::Startup as u8); static MINING_STATE: AtomicU8 = AtomicU8::new(MiningState::Idle as u8); @@ -108,35 +108,114 @@ pub async fn wait_for_mining_idle() { } } -fn try_acquire_reorg_gate() -> bool { - // Only one reorg path may own the gate at a time. - REORG_GATE +fn try_acquire_chain_operation_gate() -> bool { + CHAIN_OPERATION_GATE .compare_exchange(false, true, AtomicOrdering::SeqCst, AtomicOrdering::SeqCst) .is_ok() } -pub async fn begin_reorg_lock() { - // Reorg starts by asking mining to stop and waiting for the - // current worker round to leave the save path. - request_mining_stop(); - wait_for_mining_idle().await; - if is_reorganizing_mode() { - return; - } - // Wait for any competing reorg attempt to release the gate. - while !try_acquire_reorg_gate() { +async fn acquire_chain_operation_gate() { + while !try_acquire_chain_operation_gate() { sleep(Duration::from_millis(5)).await; } - // Reorganizing mode blocks new mining rounds until the lock ends. - set_node_mode(NodeMode::Reorganizing); } -pub fn end_reorg_lock() { - // Restore normal mode only if this task still owns the reorg state. - if is_reorganizing_mode() { - set_node_mode(NodeMode::Normal); - } - // Clear the stop flag last so mining cannot resume before mode is restored. +fn release_chain_operation(mode: NodeMode) { + set_node_mode(mode); clear_mining_stop_request(); - REORG_GATE.store(false, AtomicOrdering::SeqCst); + CHAIN_OPERATION_GATE.store(false, AtomicOrdering::SeqCst); +} + +// Owns this process's canonical chain mutation slot for one synchronization +// operation. Torrent metadata and pieces may still download concurrently; +// this guard only prevents overlapping save/replay and mode transitions. +pub struct ChainOperationGuard { + active: bool, +} + +impl ChainOperationGuard { + pub fn finish(mut self) { + if self.active { + release_chain_operation(NodeMode::Normal); + self.active = false; + } + } +} + +impl Drop for ChainOperationGuard { + fn drop(&mut self) { + if self.active { + release_chain_operation(NodeMode::Normal); + self.active = false; + } + } +} + +pub async fn begin_chain_sync() -> ChainOperationGuard { + acquire_chain_operation_gate().await; + request_mining_stop(); + wait_for_mining_idle().await; + set_node_mode(NodeMode::Syncing); + ChainOperationGuard { active: true } +} + +pub async fn begin_reorg_lock() -> ChainOperationGuard { + // Orphan correction shares the canonical chain gate with synchronization, + // so rollback/replay cannot overlap a startup or live catch-up save. + acquire_chain_operation_gate().await; + request_mining_stop(); + wait_for_mining_idle().await; + set_node_mode(NodeMode::Reorganizing); + ChainOperationGuard { active: true } +} + +#[cfg(test)] +mod tests { + use super::*; + use tokio::sync::oneshot; + + #[tokio::test] + async fn canonical_chain_operations_wait_and_release_on_drop() { + CHAIN_OPERATION_GATE.store(false, AtomicOrdering::SeqCst); + set_mining_state(MiningState::Idle); + clear_mining_stop_request(); + set_node_mode(NodeMode::Normal); + + let first = begin_chain_sync().await; + assert!(is_syncing_mode()); + + let (started_tx, mut started_rx) = oneshot::channel(); + let waiter = tokio::spawn(async move { + let second = begin_reorg_lock().await; + let _ = started_tx.send(()); + second.finish(); + }); + + sleep(Duration::from_millis(20)).await; + assert!(matches!( + started_rx.try_recv(), + Err(oneshot::error::TryRecvError::Empty) + )); + + drop(first); + crate::timeout(Duration::from_secs(1), &mut started_rx) + .await + .expect("waiting chain operation did not acquire the released gate") + .expect("waiting chain operation exited before acquiring the gate"); + waiter.await.expect("waiting chain operation task failed"); + + assert!(is_normal_mode()); + assert!(!is_mining_stop_requested()); + assert!(!CHAIN_OPERATION_GATE.load(AtomicOrdering::SeqCst)); + + async fn guarded_error() -> Result<(), &'static str> { + let _guard = begin_chain_sync().await; + Err("expected test error") + } + + assert_eq!(guarded_error().await, Err("expected test error")); + assert!(is_normal_mode()); + assert!(!is_mining_stop_requested()); + assert!(!CHAIN_OPERATION_GATE.load(AtomicOrdering::SeqCst)); + } } diff --git a/src/orphans/deep_sync_rollback.rs b/src/orphans/deep_sync_rollback.rs index 2ca01d2..4307b5a 100644 --- a/src/orphans/deep_sync_rollback.rs +++ b/src/orphans/deep_sync_rollback.rs @@ -1,5 +1,4 @@ use crate::log::warn; -use crate::miner::flag::begin_reorg_lock; use crate::orphans::structs::{CheckUp, UndoTransactions}; use crate::orphans::undo_block_transactions::undo_transactions; use crate::torrent::unpack_local_torrent::load_torrent; @@ -35,9 +34,6 @@ pub async fn deep_sync_rollback(mut params: CheckUp, wallet: Arc) { break; } - if !params.node_syncing { - begin_reorg_lock().await; - } // Undo one height, then repeat the comparison // against the next lower local height. let undo_transactions_params = UndoTransactions { diff --git a/src/orphans/orphan_checkup.rs b/src/orphans/orphan_checkup.rs index 8bfaa1f..4752b1f 100644 --- a/src/orphans/orphan_checkup.rs +++ b/src/orphans/orphan_checkup.rs @@ -1,6 +1,5 @@ use crate::common::skein::skein_128_hash_bytes; use crate::log::{info, warn}; -use crate::miner::flag::begin_reorg_lock; use crate::orphans::replay_errors::staged_candidate_status_for_error; use crate::orphans::structs::{OrphanCheckup, UndoTransactions}; use crate::orphans::undo_block_transactions::undo_transactions; @@ -223,9 +222,6 @@ pub async fn checkup(params: OrphanCheckup, wallet: Arc) -> Result<(), S }; set_torrent_status(height, &competing_info_hash, TorrentStatus::Valid).await; - if !params.node_syncing { - begin_reorg_lock().await; - } info!("[orphan] adopting proven staged chain from height {height}"); undo_transactions(undo_transactions_params, wallet.clone()).await?; return Ok(()); diff --git a/src/orphans/snapshot_check.rs b/src/orphans/snapshot_check.rs index 7cfb964..6ad637e 100644 --- a/src/orphans/snapshot_check.rs +++ b/src/orphans/snapshot_check.rs @@ -1,6 +1,6 @@ use crate::common::binary_conversions::binary_to_string; use crate::log::{error, info, warn}; -use crate::miner::flag::{begin_reorg_lock, is_syncing_mode}; +use crate::miner::flag::is_syncing_mode; use crate::orphans::structs::UndoTransactions; use crate::orphans::undo_block_transactions::undo_transactions; use crate::records::memory::connections::live_miner_peer_streams; @@ -183,9 +183,6 @@ pub async fn snapshot_verified(params: UndoTransactions, wallet: Arc) -> if local_hash != snap_hash { // Local state no longer matches the trusted checkpoint, // so rollback starts at the snapshot height. - if !params.node_syncing { - begin_reorg_lock().await; - } let undo_transactions_params = UndoTransactions { start_height: snap_height, replay_to_height: params.replay_to_height, diff --git a/src/orphans/sync_check.rs b/src/orphans/sync_check.rs index 1c08a25..0c30893 100644 --- a/src/orphans/sync_check.rs +++ b/src/orphans/sync_check.rs @@ -1,6 +1,6 @@ use crate::common::check_genesis::genesis_checkup; use crate::log::{error, info, warn}; -use crate::miner::flag::end_reorg_lock; +use crate::miner::flag::begin_reorg_lock; use crate::orphans::add_genesis::create_genesis_block; use crate::orphans::checkup_state::take_orphan_recheck_height; use crate::orphans::deep_sync_rollback::deep_sync_rollback; @@ -280,6 +280,12 @@ async fn sync_checkup_pass(params: &OrphanCheckup2, wallet: Arc) -> Resu } pub async fn sync_checkup(mut params: OrphanCheckup2, wallet: Arc) -> Result<(), String> { + let reorg_guard = if params.node_syncing { + None + } else { + Some(begin_reorg_lock().await) + }; + let result = loop { match sync_checkup_pass(¶ms, wallet.clone()).await { Ok(()) => {} @@ -315,8 +321,8 @@ pub async fn sync_checkup(mut params: OrphanCheckup2, wallet: Arc) -> Re ); }; - if !params.node_syncing { - end_reorg_lock(); + if let Some(guard) = reorg_guard { + guard.finish(); } if result.is_ok() { diff --git a/src/records/memory/connections.rs b/src/records/memory/connections.rs index 33ab01b..fee1415 100644 --- a/src/records/memory/connections.rs +++ b/src/records/memory/connections.rs @@ -1,7 +1,6 @@ use crate::common::binary_conversions::{binary_to_ip, ip_to_binary}; use crate::lazy_static; use crate::log::{info, warn}; -use crate::miner::flag::is_normal_mode; use crate::records::memory::enums::{ClientType, ConnectionType}; use crate::records::memory::network_mapping::monitor::{MONITOR_ACTION_ADD, MONITOR_ACTION_REMOVE}; use crate::records::memory::network_mapping::structs::{MonitorAddressParams, SignedMonitorEdit}; @@ -9,11 +8,11 @@ use crate::records::memory::network_mapping::NodeInfo; use crate::records::memory::response_channels::{ delete_entry, reserve_entry_with_context, Command, }; -use crate::records::memory::structs::{Connection, StoreConnectionParams}; +use crate::records::memory::structs::{Connection, ConnectionHealth, StoreConnectionParams}; use crate::rpc::client::handshake::connect_and_handshake; use crate::rpc::client::handshake_processing::{bootstrap_peer_discovery, BootstrapParams}; use crate::rpc::client::structs::Connect; -use crate::rpc::command_maps::RPC_BLOCK_HEIGHT; +use crate::rpc::command_maps::{RPC_BLOCK_HEIGHT, RPC_SETUP_COMPLETE}; use crate::rpc::responses::RpcResponse; use crate::sled::Db; use crate::sleep; @@ -22,13 +21,13 @@ use crate::timeout; use crate::wallets::structures::Wallet; use crate::Arc; use crate::AsyncWriteExt; -use crate::AtomicBool; -use crate::AtomicOrdering; use crate::Duration; use crate::IteratorRandom; use crate::Mutex; use crate::RwLock; use crate::TcpStream; +use std::collections::HashSet; +use std::sync::{Mutex as StdMutex, RwLock as StdRwLock}; fn split_ip_port_key(value: &str) -> Option<(String, u16)> { // Connection keys are stored as ip:port strings; IPv6 addresses may arrive @@ -46,47 +45,70 @@ fn split_ip_port_key(value: &str) -> Option<(String, u16)> { use crate::records::memory::structs::{ConnectionInfo, ConnectionKey}; #[derive(Clone)] -struct ReconnectContext { +struct NodeRuntimeContext { + // These dependencies belong to the running node, not to any individual + // outgoing handshake or reconnect attempt. db: Db, wallet: Arc, map: Arc>, } lazy_static! { - static ref RECONNECT_CONTEXT: Mutex> = Mutex::new(None); - static ref RECONNECT_IN_PROGRESS: AtomicBool = AtomicBool::new(false); + static ref RECOVERY_IN_PROGRESS: StdMutex> = StdMutex::new(HashSet::new()); + static ref NODE_RUNTIME_CONTEXT: StdRwLock> = StdRwLock::new(None); } -fn try_start_reconnect() -> bool { - // Only one reconnect path should run at a time, whether it came from - // liveness failure or bootstrap recovery. - RECONNECT_IN_PROGRESS - .compare_exchange(false, true, AtomicOrdering::SeqCst, AtomicOrdering::SeqCst) - .is_ok() +struct RecoveryGuard { + key: String, } -fn finish_reconnect() { - // Release the reconnect gate after the async reconnect attempt finishes. - RECONNECT_IN_PROGRESS.store(false, AtomicOrdering::SeqCst); +impl Drop for RecoveryGuard { + fn drop(&mut self) { + if let Ok(mut active) = RECOVERY_IN_PROGRESS.lock() { + active.remove(&self.key); + } + } } -pub async fn set_reconnect_context(db: Db, wallet: Arc, map: Arc>) { - let mut context = RECONNECT_CONTEXT.lock().await; - // Store enough state for later liveness checks to reconnect without - // needing the original startup stack. - *context = Some(ReconnectContext { db, wallet, map }); +fn try_start_recovery(key: String) -> Option { + // Recovery is deduplicated per peer. One dead endpoint can no longer + // suppress retries for every other endpoint for the duration of its backoff. + let mut active = RECOVERY_IN_PROGRESS.lock().ok()?; + if !active.insert(key.clone()) { + return None; + } + Some(RecoveryGuard { key }) +} + +pub fn initialize_node_runtime_context( + db: Db, + wallet: Arc, + map: Arc>, +) -> Result<(), String> { + // Replace the context when the Windows service relaunches an unlocked node + // task without restarting the service process. + let mut context = NODE_RUNTIME_CONTEXT + .write() + .map_err(|_| "node runtime context lock was poisoned".to_string())?; + *context = Some(NodeRuntimeContext { db, wallet, map }); + Ok(()) +} + +fn node_runtime_context() -> Option { + NODE_RUNTIME_CONTEXT + .read() + .ok() + .and_then(|context| context.clone()) } async fn reconnect_replacement_inner(excluded_ip: &str) { // When an outgoing peer disappears, try to replace it with another // active node that is not already connected and is not the failed IP. - let context = { - let guard = RECONNECT_CONTEXT.lock().await; - guard.clone() + let Some(_topology_guard) = try_start_recovery("topology".to_string()) else { + return; }; - - let Some(context) = context else { - warn!("[reconnect] no reconnect context configured"); + let Some(context) = node_runtime_context() else { + warn!("[reconnect] node runtime context is not initialized"); return; }; @@ -132,77 +154,75 @@ async fn reconnect_replacement_inner(excluded_ip: &str) { } async fn retry_dropped_outgoing(ip: String, port: u16) { - if !try_start_reconnect() { - warn!("[reconnect] reconnect attempt already in progress, skipping duplicate request"); + let recovery_key = format!("peer:{ip}"); + let Some(_recovery_guard) = try_start_recovery(recovery_key) else { + warn!("[reconnect] reconnect attempt already in progress for peer {ip}"); return; - } + }; - async { - let context = { - let guard = RECONNECT_CONTEXT.lock().await; - guard.clone() - }; + let Some(context) = node_runtime_context() else { + warn!("[reconnect] node runtime context is not initialized"); + return; + }; - let Some(context) = context else { - warn!("[reconnect] no reconnect context configured"); - return; - }; - - let addr_string = format!("{ip}:{port}"); - for attempt in 1..=3 { + let addr_string = format!("{ip}:{port}"); + for attempt in 1..=3 { + // A connection-manager drop is definitive, so the first reconnect is + // immediate. Backoff applies only after that direct attempt fails. + if attempt > 1 { sleep(Duration::from_secs(30)).await; - - if Connection::get_stream_from_memory(&addr_string).await.is_some() { - info!( - "[reconnect] dropped peer {addr_string} is already connected; stopping reconnect attempts" - ); - return; - } - - let socket_addr = match addr_string.parse() { - Ok(addr) => addr, - Err(err) => { - warn!("[reconnect] invalid dropped peer address {addr_string}: {err}"); - break; - } - }; - - let connect = Connect { - addr: socket_addr, - node_ip: addr_string.clone(), - wallet: context.wallet.clone(), - db: context.db.clone(), - map: context.map.clone(), - first: peer_connection_count().await == 0, - }; - - match connect_and_handshake(connect).await { - Ok(()) => { - info!("[reconnect] reconnected dropped peer {addr_string} on attempt {attempt}"); - return; - } - Err(err) => { - let err_string = err.to_string(); - if err_string.contains( - "The connection is already in the connection manager Please wait 10 minutes and try again", - ) { - warn!( - "[reconnect] dropped peer {addr_string} is still recorded by remote connection manager on attempt {attempt}/3; retrying" - ); - continue; - } - warn!( - "[reconnect] failed to reconnect dropped peer {addr_string} on attempt {attempt}/3: {err_string}" - ); - } - } } - reconnect_replacement_inner(&ip).await; - } - .await; + if Connection::get_stream_from_memory(&addr_string) + .await + .is_some() + { + info!( + "[reconnect] dropped peer {addr_string} is already connected; stopping reconnect attempts" + ); + return; + } - finish_reconnect(); + let socket_addr = match addr_string.parse() { + Ok(addr) => addr, + Err(err) => { + warn!("[reconnect] invalid dropped peer address {addr_string}: {err}"); + break; + } + }; + + let connect = Connect { + addr: socket_addr, + node_ip: addr_string.clone(), + wallet: context.wallet.clone(), + db: context.db.clone(), + map: context.map.clone(), + first: peer_connection_count().await == 0, + }; + + match connect_and_handshake(connect).await { + Ok(()) => { + info!("[reconnect] reconnected dropped peer {addr_string} on attempt {attempt}"); + return; + } + Err(err) => { + let err_string = err.to_string(); + if err_string.contains( + "The connection is already in the connection manager Please wait 10 minutes and try again", + ) { + warn!( + "[reconnect] dropped peer {addr_string} is still recorded by remote connection manager on attempt {attempt}/3; retrying" + ); + continue; + } + warn!( + "[reconnect] failed to reconnect dropped peer {addr_string} on attempt {attempt}/3: {err_string}" + ); + } + } + } + + reconnect_replacement_inner(&ip).await; } pub fn spawn_retry_dropped_outgoing(ip: String, port: u16) { @@ -212,21 +232,26 @@ pub fn spawn_retry_dropped_outgoing(ip: String, port: u16) { } pub fn spawn_reconnect_bootstrap(params: BootstrapParams) { - if !try_start_reconnect() { + let recovery_key = format!("bootstrap:{}", params.connections_key); + let Some(recovery_guard) = try_start_recovery(recovery_key) else { warn!("[reconnect] bootstrap recovery already in progress, skipping duplicate request"); return; - } + }; // Bootstrap discovery can perform network requests, so it runs detached // from the caller that noticed the connection problem. tokio::spawn(async move { + let _recovery_guard = recovery_guard; if let Err(err) = bootstrap_peer_discovery(params).await { warn!("[reconnect] bootstrap recovery failed: {err}"); } - finish_reconnect(); }); } +pub async fn refill_outgoing_connections_once() { + reconnect_replacement_inner("").await; +} + impl Connection { // Initialize the in-memory connection manager state. pub fn new() -> Self { @@ -312,7 +337,6 @@ impl Connection { if let Some(connection_info) = removed.as_ref() { if ClientType::from_bytes(&connection_info.client_type) == Some(ClientType::Miner) && connection_info.ready - && is_normal_mode() { spawn_monitor_update( ip.clone(), @@ -388,76 +412,82 @@ impl Connection { message.push(message_type); message.extend_from_slice(&checkup_key); - RpcResponse::send_raw(&stream, None, &message).await; - - let response_result = { + let sent = + RpcResponse::send_raw(&stream, Some(&format!("{ip}:{port}")), &message).await; + let reply_received = if sent { let mut checkup_rx = checkup_rx_mutex.lock().await; - timeout(Duration::from_secs(30), checkup_rx.recv()).await + matches!( + timeout(Duration::from_secs(30), checkup_rx.recv()).await, + Ok(Some(_)) + ) + } else { + false }; - match response_result { - Ok(Some(_reply)) => { - consecutive_failures = 0; - info!( - "[connection_manager] liveness check ok: type={} peer={}:{}", - connection_type.as_str(), - ip, - port - ); - } - _ => { - let still_monitoring_same_stream = { - let guard = CONNECTIONS.read().await; - guard - .as_ref() - .and_then(|conn| conn.connection_map.get(&connection_key)) - .map(|connection_info| { - Arc::ptr_eq(&connection_info.stream, &stream) - }) - .unwrap_or(false) - }; - - if !still_monitoring_same_stream { - delete_entry(command_map.clone(), checkup_key).await; - break; - } - - consecutive_failures = consecutive_failures.saturating_add(1); - warn!( - "[connection_manager] liveness check failed: type={} peer={}:{} attempt={}/3", - connection_type.as_str(), - ip, - port, - consecutive_failures - ); + if reply_received { + set_stream_health(&stream, ConnectionHealth::Connected).await; + consecutive_failures = 0; + info!( + "[connection_manager] liveness check ok: type={} peer={}:{}", + connection_type.as_str(), + ip, + port + ); + } else { + // A busy stream is still owned by active work. It is not + // equivalent to a peer that accepted a ping but never replied. + if stream_health(&stream).await == Some(ConnectionHealth::Busy) { delete_entry(command_map.clone(), checkup_key).await; + continue; + } + set_stream_health(&stream, ConnectionHealth::Unresponsive).await; + let still_monitoring_same_stream = { + let guard = CONNECTIONS.read().await; + guard + .as_ref() + .and_then(|conn| conn.connection_map.get(&connection_key)) + .map(|connection_info| Arc::ptr_eq(&connection_info.stream, &stream)) + .unwrap_or(false) + }; - if consecutive_failures < 3 { - continue; - } - - // Three consecutive timed-out or missing replies drop - // the connection, and outgoing peers trigger - // replacement discovery. - let mut guard = CONNECTIONS.write().await; - if let Some(conn) = guard.as_mut() { - let should_drop = conn - .connection_map - .get(&connection_key) - .map(|connection_info| { - Arc::ptr_eq(&connection_info.stream, &stream) - }) - .unwrap_or(false); - if should_drop { - conn.drop_connection(connection_type, ip.clone(), port); - } - } - drop(guard); - if connection_type == ConnectionType::Outgoing { - spawn_retry_dropped_outgoing(ip.clone(), port); - } + if !still_monitoring_same_stream { + delete_entry(command_map.clone(), checkup_key).await; break; } + + consecutive_failures = consecutive_failures.saturating_add(1); + warn!( + "[connection_manager] liveness check failed: type={} peer={}:{} attempt={}/3", + connection_type.as_str(), + ip, + port, + consecutive_failures + ); + delete_entry(command_map.clone(), checkup_key).await; + + if consecutive_failures < 3 { + continue; + } + + // Three consecutive timed-out or missing replies drop + // the connection, and outgoing peers trigger + // replacement discovery. + let mut guard = CONNECTIONS.write().await; + if let Some(conn) = guard.as_mut() { + let should_drop = conn + .connection_map + .get(&connection_key) + .map(|connection_info| Arc::ptr_eq(&connection_info.stream, &stream)) + .unwrap_or(false); + if should_drop { + conn.drop_connection(connection_type, ip.clone(), port); + } + } + drop(guard); + if connection_type == ConnectionType::Outgoing { + spawn_retry_dropped_outgoing(ip.clone(), port); + } + break; } } }); @@ -483,6 +513,17 @@ impl Connection { .count() } + pub fn count_ready_outgoing_connections(&self) -> usize { + self.connection_map + .values() + .filter(|info| { + ConnectionType::from_bytes(&info.connection_type) == Some(ConnectionType::Outgoing) + && ClientType::from_bytes(&info.client_type) == Some(ClientType::Miner) + && info.ready + }) + .count() + } + pub fn mark_wallet_registry_synced(&mut self, key: &str) -> bool { let Some((ip, port)) = split_ip_port_key(key) else { return false; @@ -511,38 +552,98 @@ impl Connection { false } - pub fn mark_operational(&mut self, key: &str, command_map: Arc>) -> bool { + fn finalize_setup_if_complete( + info: &mut ConnectionInfo, + ip: String, + port: u16, + command_map: Arc>, + ) -> bool { + if info.ready { + return true; + } + if ClientType::from_bytes(&info.client_type) != Some(ClientType::Miner) + || !info.wallet_registry_synced + || !info.network_map_synced + || !info.local_setup_complete + || !info.remote_setup_complete + || !info.local_setup_acknowledged + { + return false; + } + + info.ready = true; + spawn_monitor_update( + ip.clone(), + MONITOR_ACTION_ADD, + info.wallet_short_address.clone(), + port, + ); + Connection::client_checkup( + Arc::clone(&info.stream), + ConnectionType::from_bytes(&info.connection_type).unwrap_or(ConnectionType::Incoming), + ip, + port, + command_map, + ); + true + } + + pub fn mark_local_setup_complete( + &mut self, + key: &str, + command_map: Arc>, + ) -> Option<(Arc>, bool)> { + let Some((ip, port)) = split_ip_port_key(key) else { + return None; + }; + let ip_bytes = ip_to_binary(&ip); + for (connection_key, info) in self.connection_map.iter_mut() { + if connection_key.ip == ip_bytes && connection_key.port == port { + if ClientType::from_bytes(&info.client_type) != Some(ClientType::Miner) + || !info.wallet_registry_synced + || !info.network_map_synced + { + return None; + } + info.local_setup_complete = true; + Self::finalize_setup_if_complete(info, ip, port, command_map); + return Some((Arc::clone(&info.stream), info.local_setup_acknowledged)); + } + } + None + } + + pub fn mark_remote_setup_complete( + &mut self, + key: &str, + command_map: Arc>, + ) -> bool { let Some((ip, port)) = split_ip_port_key(key) else { return false; }; let ip_bytes = ip_to_binary(&ip); for (connection_key, info) in self.connection_map.iter_mut() { if connection_key.ip == ip_bytes && connection_key.port == port { - if info.ready { - return true; - } - if ClientType::from_bytes(&info.client_type) != Some(ClientType::Miner) - || !info.wallet_registry_synced - || !info.network_map_synced - { - return false; - } - info.ready = true; - spawn_monitor_update( - ip.clone(), - MONITOR_ACTION_ADD, - info.wallet_short_address.clone(), - port, - ); - Connection::client_checkup( - Arc::clone(&info.stream), - ConnectionType::from_bytes(&info.connection_type) - .unwrap_or(ConnectionType::Incoming), - ip, - port, - command_map, - ); - return true; + info.remote_setup_complete = true; + return Self::finalize_setup_if_complete(info, ip, port, command_map); + } + } + false + } + + pub fn mark_local_setup_acknowledged( + &mut self, + key: &str, + command_map: Arc>, + ) -> bool { + let Some((ip, port)) = split_ip_port_key(key) else { + return false; + }; + let ip_bytes = ip_to_binary(&ip); + for (connection_key, info) in self.connection_map.iter_mut() { + if connection_key.ip == ip_bytes && connection_key.port == port { + info.local_setup_acknowledged = true; + return Self::finalize_setup_if_complete(info, ip, port, command_map); } } false @@ -797,11 +898,8 @@ impl Connection { fn spawn_monitor_update(ip: String, action: u8, monitored_address: String, port: u16) { tokio::spawn(async move { - let context = { - let guard = RECONNECT_CONTEXT.lock().await; - guard.clone() - }; - let Some(context) = context else { + let Some(context) = node_runtime_context() else { + warn!("[network_map] node runtime context is not initialized"); return; }; if !Wallet::short_address_validation(&monitored_address) { @@ -812,12 +910,16 @@ fn spawn_monitor_update(ip: String, action: u8, monitored_address: String, port: return; } let timestamp = crate::Utc::now().timestamp_millis() as u64; + let modified_block = + crate::records::block_height::get_block_height::get_height(&context.db) + .saturating_add(1); let signature = NodeInfo::monitor_signature( action, &monitored_address, &monitoring_address, &ip, timestamp, + modified_block, &context.wallet, ) .await; @@ -827,6 +929,7 @@ fn spawn_monitor_update(ip: String, action: u8, monitored_address: String, port: monitoring_address, target_ip: ip.clone(), modified_timestamp: timestamp, + modified_block, modified_signature: signature, }; let params = MonitorAddressParams { @@ -868,6 +971,15 @@ pub async fn outgoing_connection_count() -> usize { .unwrap_or(0) } +pub async fn ready_outgoing_connection_count() -> usize { + CONNECTIONS + .read() + .await + .as_ref() + .map(|connection| connection.count_ready_outgoing_connections()) + .unwrap_or(0) +} + pub async fn peer_connection_count() -> usize { // Mining needs at least one fully initialized miner peer. A raw socket // is not enough because block validation depends on wallet registry and @@ -910,14 +1022,80 @@ pub async fn mark_peer_network_map_synced(key: &str) -> bool { } pub async fn mark_peer_operational(key: &str, map: Arc>) -> bool { + let setup = CONNECTIONS + .write() + .await + .as_mut() + .and_then(|connection| connection.mark_local_setup_complete(key, map.clone())); + let Some((stream, already_acknowledged)) = setup else { + return false; + }; + if already_acknowledged { + return peer_is_operational(key).await; + } + + for attempt in 1..=3 { + let (uid, _tx, rx) = reserve_entry_with_context( + map.clone(), + Some(RPC_SETUP_COMPLETE), + Some(key.to_string()), + ) + .await; + let mut message = Vec::with_capacity(4); + message.push(RPC_SETUP_COMPLETE); + message.extend_from_slice(&uid); + + if !RpcResponse::send_raw(&stream, Some(key), &message).await { + delete_entry(map.clone(), uid).await; + continue; + } + + let response = { + let mut rx = rx.lock().await; + timeout(Duration::from_secs(10), rx.recv()).await + }; + delete_entry(map.clone(), uid).await; + if matches!(response, Ok(Some(bytes)) if bytes == b"setup_complete_ack") { + return CONNECTIONS + .write() + .await + .as_mut() + .map(|connection| connection.mark_local_setup_acknowledged(key, map.clone())) + .unwrap_or(false); + } + warn!( + "[handshake_state] setup-complete acknowledgement failed: peer={key} attempt={attempt}/3" + ); + } + + false +} + +pub async fn mark_peer_remote_setup_complete(key: &str, map: Arc>) -> bool { CONNECTIONS .write() .await .as_mut() - .map(|connection| connection.mark_operational(key, map)) + .map(|connection| connection.mark_remote_setup_complete(key, map)) .unwrap_or(false) } +pub fn spawn_peer_setup_retry(key: String, map: Arc>) { + tokio::spawn(async move { + loop { + sleep(Duration::from_secs(10)).await; + if Connection::get_stream_from_memory(&key).await.is_none() + || peer_is_operational(&key).await + { + break; + } + if mark_peer_operational(&key, map.clone()).await { + break; + } + } + }); +} + pub async fn peer_is_operational(key: &str) -> bool { let Some((ip, port)) = split_ip_port_key(key) else { return false; @@ -974,6 +1152,30 @@ pub async fn peer_accepts_live_relay(key: &str) -> bool { .unwrap_or(false) } +pub async fn set_stream_health(stream: &Arc>, health: ConnectionHealth) { + let mut guard = CONNECTIONS.write().await; + let Some(connection) = guard.as_mut() else { + return; + }; + if let Some(info) = connection + .connection_map + .values_mut() + .find(|info| Arc::ptr_eq(&info.stream, stream)) + { + info.health = health; + } +} + +pub async fn stream_health(stream: &Arc>) -> Option { + CONNECTIONS.read().await.as_ref().and_then(|connection| { + connection + .connection_map + .values() + .find(|info| Arc::ptr_eq(&info.stream, stream)) + .map(|info| info.health) + }) +} + pub async fn live_miner_peer_streams() -> Vec<(String, Arc>)> { // Snapshot consensus and recovery checks vote only across currently // connected miner peers, regardless of incoming/outgoing direction. @@ -1026,3 +1228,43 @@ pub async fn get_client_type_from_memory(key: &str) -> Option { None } + +#[cfg(test)] +mod tests { + use super::*; + use crate::wallets::structures::SavedWallet; + + #[test] + fn runtime_context_can_be_initialized_before_any_handshake() { + let db = sled::Config::new().temporary(true).open().unwrap(); + let wallet = Arc::new(Wallet { + saved: SavedWallet { + short_address: "ab13318c26250b048db92920a80a86127c933b0c.cltc".to_string(), + vanity_address: None, + public_key: String::new(), + private_key: String::new(), + }, + encryption_key: String::new(), + }); + let map = Arc::new(Mutex::new(Command::new())); + + initialize_node_runtime_context(db.clone(), wallet.clone(), map.clone()).unwrap(); + let context = node_runtime_context().expect("runtime context should be available"); + + assert!(Arc::ptr_eq(&context.wallet, &wallet)); + assert!(Arc::ptr_eq(&context.map, &map)); + } + + #[test] + fn recovery_is_deduplicated_per_peer_instead_of_globally() { + let first = try_start_recovery("test-peer:first".to_string()) + .expect("first peer recovery should start"); + assert!(try_start_recovery("test-peer:first".to_string()).is_none()); + + let second = try_start_recovery("test-peer:second".to_string()) + .expect("another peer must recover independently"); + drop(first); + assert!(try_start_recovery("test-peer:first".to_string()).is_some()); + drop(second); + } +} diff --git a/src/records/memory/mempool/mod.rs b/src/records/memory/mempool/mod.rs index 1fefddf..905fa5c 100644 --- a/src/records/memory/mempool/mod.rs +++ b/src/records/memory/mempool/mod.rs @@ -232,8 +232,7 @@ pub use lookups::{ pub use processing::{ delete_by_signatures, delete_unprocessed_by_address, mark_processed_by_signatures, mark_selected_transactions_processed, restore_processed_by_signatures, - restore_selected_transactions_processed, - spawn_processed_cleanup, + restore_selected_transactions_processed, spawn_processed_cleanup, }; pub use schema::{ clear_mempool, db_client, ensure_db_connection, init_db, pg_execute, pg_query, setup_mempool, diff --git a/src/records/memory/network_mapping/add.rs b/src/records/memory/network_mapping/add.rs index fe91815..6355cdf 100644 --- a/src/records/memory/network_mapping/add.rs +++ b/src/records/memory/network_mapping/add.rs @@ -14,22 +14,10 @@ fn signature_is_empty(signature: &str) -> bool { .unwrap_or(false) } -fn merge_monitors(existing: &mut Vec, monitors: Vec) -> bool { - let mut changed = false; - for monitor in monitors { - if !existing.iter().any(|known| known == &monitor) { - existing.push(monitor); - changed = true; - } - } - changed -} - impl NodeInfo { pub async fn import_signed_mapping_address( db: &Db, edit: SignedNodeEdit, - monitors: Vec, _blocks_mined: u8, ) -> Result<(), String> { if !is_public_network_address(&edit.ip) { @@ -66,11 +54,8 @@ impl NodeInfo { existing_node.added_by = edit.modified_by; existing_node.added_timestamp = edit.modified_timestamp; existing_node.added_signature = edit.modified_signature; - existing_node.monitoring = monitors; - existing_node.deleted_timestamp = 0_u64; - existing_node.deleted_block = 0_u32; drop(address_map); - Self::persist_recovery_snapshot("import revive").await; + Self::persist_recovery_snapshot("import deleted membership").await; return Ok(()); } else { if existing_node.ip != edit.ip { @@ -88,7 +73,6 @@ impl NodeInfo { existing_node.added_timestamp = edit.modified_timestamp; existing_node.added_signature = edit.modified_signature; } - merge_monitors(&mut existing_node.monitoring, monitors); drop(address_map); Self::persist_recovery_snapshot("import existing").await; return Ok(()); @@ -102,16 +86,14 @@ impl NodeInfo { } address_map.insert(edit.address, { - let mut node = NodeInfo::new( + NodeInfo::new( edit.ip, edit.port, 0, edit.modified_by, edit.modified_timestamp, edit.modified_signature, - ); - node.monitoring = monitors; - node + ) }); drop(address_map); Self::persist_recovery_snapshot("import add").await; @@ -170,20 +152,9 @@ impl NodeInfo { }; let modified_timestamp_bytes = edit.modified_timestamp.to_le_bytes(); let modified_signature_bytes = decode(&edit.modified_signature).unwrap(); - let monitor_bytes = { - let address_map = ADDRESS_MAP.lock().await; - address_map - .get(&edit.address) - .map(|node| { - node.monitoring - .iter() - .filter_map(|monitor| Wallet::short_address_to_bytes(monitor)) - .collect::>() - }) - .unwrap_or_default() - }; - let monitor_count = monitor_bytes.len().min(u16::MAX as usize) as u16; - let monitor_count_bytes = monitor_count.to_le_bytes(); + // Membership broadcasts never carry liveness state. Monitor changes + // use their own signed commands and cannot be injected into node adds. + let monitor_count_bytes = 0_u16.to_le_bytes(); let streams = { let connections_lock = CONNECTIONS.read().await; connections_lock @@ -193,6 +164,7 @@ impl NodeInfo { }; match streams { streams if !streams.is_empty() => { + let mut messages = Vec::with_capacity(streams.len()); for (peer_key, unlocked_stream) in streams { let Some((peer_ip, _)) = peer_key.rsplit_once(':') else { continue; @@ -216,11 +188,9 @@ impl NodeInfo { message.extend_from_slice(&modified_timestamp_bytes); message.extend_from_slice(&modified_signature_bytes); message.extend_from_slice(&monitor_count_bytes); - for monitor in monitor_bytes.iter().take(monitor_count as usize) { - message.extend_from_slice(monitor); - } - RpcResponse::send_raw(&unlocked_stream, Some(&peer_key), &message).await; + messages.push((peer_key, unlocked_stream, message)); } + RpcResponse::broadcast_raw(messages).await; } _ => { warn!("No active connections found."); @@ -232,7 +202,7 @@ impl NodeInfo { let AddAddressParams { map, mut edit, - monitors, + monitors: _monitors, blocks_mined: _blocks_mined, remote_ip, db, @@ -384,8 +354,6 @@ impl NodeInfo { existing_node.added_by = edit.modified_by.clone(); existing_node.added_timestamp = edit.modified_timestamp; existing_node.added_signature = edit.modified_signature.clone(); - existing_node.deleted_timestamp = 0_u64; - existing_node.deleted_block = 0_u32; state_changed = true; } else { if existing_node.ip != edit.ip { @@ -403,9 +371,6 @@ impl NodeInfo { existing_node.added_signature = edit.modified_signature.clone(); state_changed = true; } - if merge_monitors(&mut existing_node.monitoring, monitors.clone()) { - state_changed = true; - } } } @@ -421,18 +386,17 @@ impl NodeInfo { if !penalize_duplicate_ip { // Persist the new node locally. Network-map entries are bare // IP membership records, separate from live socket keys. - address_map.insert(edit.address.clone(), { - let mut node = NodeInfo::new( + address_map.insert( + edit.address.clone(), + NodeInfo::new( edit.ip.clone(), edit.port, 0, edit.modified_by.clone(), edit.modified_timestamp, edit.modified_signature.clone(), - ); - node.monitoring = monitors.clone(); - node - }); + ), + ); state_changed = true; } } diff --git a/src/records/memory/network_mapping/mined_counts.rs b/src/records/memory/network_mapping/mined_counts.rs index d366c46..8052a83 100644 --- a/src/records/memory/network_mapping/mined_counts.rs +++ b/src/records/memory/network_mapping/mined_counts.rs @@ -43,30 +43,6 @@ impl NodeInfo { } } - pub async fn set_deleted_block_from_mapping(address: &str, deleted_block: u32) { - let mut map = ADDRESS_MAP.lock().await; - if let Some(node_info) = map.get_mut(address) { - // The deletion height is filled in once the chain knows the block - // where the delete action becomes active. - node_info.deleted_block = deleted_block; - } - } - - pub async fn set_deleted_metadata_from_mapping( - address: &str, - deleted_timestamp: u64, - deleted_block: u32, - ) { - { - let mut map = ADDRESS_MAP.lock().await; - if let Some(node_info) = map.get_mut(address) { - node_info.deleted_timestamp = deleted_timestamp; - node_info.deleted_block = deleted_block; - } - } - Self::persist_recovery_snapshot("deleted metadata import").await; - } - pub async fn rebuild_mined_counts_from_chain(db: &Db) -> Result<(), String> { // Recompute node mined counts directly from saved block headers // so startup and recovery can rebuild memory-only state. diff --git a/src/records/memory/network_mapping/monitor.rs b/src/records/memory/network_mapping/monitor.rs index d3bbfdf..b36f9ff 100644 --- a/src/records/memory/network_mapping/monitor.rs +++ b/src/records/memory/network_mapping/monitor.rs @@ -1,4 +1,9 @@ use super::*; +use crate::records::memory::network_mapping::structs::{ + MONITOR_EVENT_ACTION_OFFSET, MONITOR_EVENT_BLOCK_OFFSET, MONITOR_EVENT_BYTES, + MONITOR_EVENT_MONITORED_OFFSET, MONITOR_EVENT_MONITORING_OFFSET, + MONITOR_EVENT_SIGNATURE_OFFSET, MONITOR_EVENT_TARGET_IP_OFFSET, MONITOR_EVENT_TIMESTAMP_OFFSET, +}; use crate::records::memory::response_channels::reserve_transient_entry_with_context; use crate::records::wallet_registry::resolve_pubkey_from_short_address; use crate::rpc::command_maps::{RPC_NETWORK_MONITOR_ADD, RPC_NETWORK_MONITOR_REMOVE}; @@ -7,7 +12,102 @@ pub const MONITOR_ACTION_ADD: u8 = 1; pub const MONITOR_ACTION_REMOVE: u8 = 2; lazy_static! { - static ref MONITOR_EVENTS_SEEN: Mutex> = Mutex::new(HashMap::new()); + // Keep the newest signed assertion for each monitor relationship. This is + // both replay protection and the authenticated state shared with joining peers. + static ref MONITOR_EVENT_STATE: Mutex> = + Mutex::new(HashMap::new()); +} + +#[cfg(test)] +mod tests { + use super::*; + + fn node(monitors: &[&str]) -> NodeInfo { + let mut node = NodeInfo::new( + "1.2.3.4".to_string(), + 50050, + 0, + "sponsor.cltc".to_string(), + 1, + "00".repeat(Wallet::SIGNATURE_LENGTH), + ); + node.monitoring = monitors.iter().map(|monitor| monitor.to_string()).collect(); + node + } + + #[test] + fn deletion_cascade_preserves_nodes_with_another_monitor() { + let mut map = HashMap::new(); + map.insert("node-a".to_string(), node(&[])); + map.insert("node-b".to_string(), node(&["node-a", "node-c"])); + map.insert("node-c".to_string(), node(&["local"])); + + NodeInfo::mark_deleted_and_cascade(&mut map, "node-a", "local", 100, 50); + + assert_eq!(map["node-a"].deleted_block, 50); + assert_eq!(map["node-b"].monitoring, vec!["node-c".to_string()]); + assert_eq!(map["node-b"].deleted_timestamp, 0); + } + + #[test] + fn deletion_cascade_removes_offline_nodes_from_their_targets() { + let mut map = HashMap::new(); + map.insert("node-a".to_string(), node(&[])); + map.insert("node-b".to_string(), node(&["node-a"])); + map.insert("node-c".to_string(), node(&["node-b"])); + + NodeInfo::mark_deleted_and_cascade(&mut map, "node-a", "local", 100, 50); + + for address in ["node-a", "node-b", "node-c"] { + assert_eq!(map[address].deleted_timestamp, 100); + assert_eq!(map[address].deleted_block, 50); + assert!(map[address].monitoring.is_empty()); + } + } + + #[test] + fn remove_wins_an_equal_timestamp_monitor_race() { + let add = SignedMonitorEdit { + action: MONITOR_ACTION_ADD, + monitored_address: "target".to_string(), + monitoring_address: "monitor".to_string(), + target_ip: "1.2.3.4".to_string(), + modified_timestamp: 100, + modified_block: 50, + modified_signature: "a".to_string(), + }; + let mut remove = add.clone(); + remove.action = MONITOR_ACTION_REMOVE; + + assert!(NodeInfo::monitor_event_order(&remove) > NodeInfo::monitor_event_order(&add)); + } + + #[test] + fn signed_monitor_wire_record_round_trips() { + let monitored = "ab13318c26250b048db92920a80a86127c933b0c.cltc".to_string(); + let monitoring = "7ba97ce531937388ef8bca4d8a4ca54cefbefb87.cltc".to_string(); + let mut bytes = Vec::with_capacity(MONITOR_EVENT_BYTES); + bytes.push(MONITOR_ACTION_REMOVE); + bytes.extend_from_slice(&Wallet::short_address_to_bytes(&monitored).unwrap()); + bytes.extend_from_slice(&Wallet::short_address_to_bytes(&monitoring).unwrap()); + bytes.extend_from_slice(&ip_to_binary("1.2.3.4")); + bytes.extend_from_slice(&100_u64.to_le_bytes()); + bytes.extend_from_slice(&50_u32.to_le_bytes()); + bytes.extend_from_slice(&vec![7_u8; Wallet::SIGNATURE_LENGTH]); + + assert_eq!(bytes.len(), MONITOR_EVENT_BYTES); + let parsed = NodeInfo::monitor_event_from_bytes(&bytes).unwrap(); + assert_eq!(parsed.action, MONITOR_ACTION_REMOVE); + assert_eq!(parsed.monitored_address, monitored); + assert_eq!(parsed.monitoring_address, monitoring); + assert_eq!(parsed.target_ip, "1.2.3.4"); + assert_eq!(parsed.modified_timestamp, 100); + assert_eq!(parsed.modified_block, 50); + assert_eq!( + parsed.modified_signature, + crate::encode(vec![7_u8; Wallet::SIGNATURE_LENGTH]) + ); + } } impl NodeInfo { @@ -16,12 +116,13 @@ impl NodeInfo { monitored_address: &str, monitoring_address: &str, target_ip: &str, - current_timestamp: u64, + modified_timestamp: u64, + modified_block: u32, wallet: &Arc, ) -> String { let private_key = &wallet.saved.private_key; let data = format!( - "{action}{monitored_address}{monitoring_address}{target_ip}{current_timestamp}" + "{action}{monitored_address}{monitoring_address}{target_ip}{modified_timestamp}{modified_block}" ); let hashed_data = skein_256_hash_data(&data); Wallet::sign_transaction(&hashed_data, private_key).await @@ -31,7 +132,7 @@ impl NodeInfo { address_map: &mut HashMap, deleted_address: &str, local_address: &str, - current_timestamp: u64, + deleted_timestamp: u64, deleted_block: u32, ) { let mut stack = vec![deleted_address.to_string()]; @@ -40,7 +141,7 @@ impl NodeInfo { let should_cascade = match address_map.get_mut(&address) { Some(_) if address == local_address => false, Some(node) if node.deleted_timestamp == 0 && node.monitoring.is_empty() => { - node.deleted_timestamp = current_timestamp; + node.deleted_timestamp = deleted_timestamp; node.deleted_block = deleted_block; true } @@ -53,9 +154,12 @@ impl NodeInfo { info!( "[network_map] node marked deleted: address={} timestamp={} deleted_block={}", - address, current_timestamp, deleted_block + address, deleted_timestamp, deleted_block ); + // Once a node has no monitors, its own old monitoring claims can + // no longer keep other nodes active. Remove those claims locally; + // no signature from the now-offline node is required. let targets: Vec = address_map .iter_mut() .filter_map(|(target, node)| { @@ -81,12 +185,13 @@ impl NodeInfo { return false; }; let data = format!( - "{}{}{}{}{}", + "{}{}{}{}{}{}", edit.action, edit.monitored_address, edit.monitoring_address, edit.target_ip, - edit.modified_timestamp + edit.modified_timestamp, + edit.modified_block, ); let hashed_data = skein_256_hash_data(&data); Wallet::verify_transaction_with_public_key_bytes( @@ -105,24 +210,47 @@ impl NodeInfo { Self::apply_monitor(params, MONITOR_ACTION_REMOVE).await } - fn monitor_event_key(edit: &SignedMonitorEdit) -> String { + fn monitor_relation_key(edit: &SignedMonitorEdit) -> String { format!( - "{}:{}:{}:{}:{}:{}", - edit.action, - edit.monitored_address, - edit.monitoring_address, - edit.target_ip, - edit.modified_timestamp, - edit.modified_signature + "{}:{}:{}", + edit.monitored_address, edit.monitoring_address, edit.target_ip ) } - async fn remember_monitor_event(edit: &SignedMonitorEdit) -> bool { - let key = Self::monitor_event_key(edit); - let mut seen = MONITOR_EVENTS_SEEN.lock().await; - let now = Utc::now().timestamp_millis() as u64; - seen.retain(|_, timestamp| now.saturating_sub(*timestamp) < 3_600_000); - seen.insert(key, now).is_none() + fn monitor_event_order(edit: &SignedMonitorEdit) -> (u64, u32, u8, &str) { + ( + edit.modified_timestamp, + edit.modified_block, + edit.action, + &edit.modified_signature, + ) + } + + async fn remember_newest_monitor_event(edit: &SignedMonitorEdit) -> bool { + let key = Self::monitor_relation_key(edit); + let mut state = MONITOR_EVENT_STATE.lock().await; + if state + .get(&key) + .map(|existing| Self::monitor_event_order(edit) <= Self::monitor_event_order(existing)) + .unwrap_or(false) + { + return false; + } + state.insert(key, edit.clone()); + true + } + + async fn deletion_marker_for(monitored_address: &str) -> Option<(u64, u32)> { + MONITOR_EVENT_STATE + .lock() + .await + .values() + .filter(|event| { + event.monitored_address == monitored_address + && event.action == MONITOR_ACTION_REMOVE + }) + .map(|event| (event.modified_timestamp, event.modified_block)) + .max() } async fn broadcast_monitor_event( @@ -145,9 +273,10 @@ impl NodeInfo { }; let target_ip_bytes = ip_to_binary(&edit.target_ip); let timestamp_bytes = edit.modified_timestamp.to_le_bytes(); + let block_bytes = edit.modified_block.to_le_bytes(); let signature_bytes = match decode(&edit.modified_signature) { - Ok(bytes) => bytes, - Err(_) => return, + Ok(bytes) if bytes.len() == Wallet::SIGNATURE_LENGTH => bytes, + _ => return, }; let streams = { let connections_lock = CONNECTIONS.read().await; @@ -157,6 +286,7 @@ impl NodeInfo { .unwrap_or_default() }; + let mut messages = Vec::with_capacity(streams.len()); for (peer_key, unlocked_stream) in streams { let Some((peer_ip, _)) = peer_key.rsplit_once(':') else { continue; @@ -170,7 +300,7 @@ impl NodeInfo { Some(peer_key.clone()), ) .await; - let mut message: Vec = Vec::new(); + let mut message = Vec::with_capacity(1 + 3 + MONITOR_EVENT_BYTES); message.push(message_type); message.extend_from_slice(&hashmap_key); message.push(edit.action); @@ -178,9 +308,79 @@ impl NodeInfo { message.extend_from_slice(&monitoring_bytes); message.extend_from_slice(&target_ip_bytes); message.extend_from_slice(×tamp_bytes); + message.extend_from_slice(&block_bytes); message.extend_from_slice(&signature_bytes); - RpcResponse::send_raw(&unlocked_stream, Some(&peer_key), &message).await; + messages.push((peer_key, unlocked_stream, message)); } + RpcResponse::broadcast_raw(messages).await; + } + + async fn apply_verified_monitor_edit( + edit: &SignedMonitorEdit, + db: &Db, + local_short: &str, + ) -> Result { + if !matches!(edit.action, MONITOR_ACTION_ADD | MONITOR_ACTION_REMOVE) { + return Err("Invalid monitor action".to_string()); + } + if !Self::verify_monitor_edit(edit, db).await { + return Err("Could not validate monitor signature".to_string()); + } + + // Validate membership before recording the event so an invalid event + // cannot poison replay protection for a later valid assertion. + { + let address_map = ADDRESS_MAP.lock().await; + let monitored = address_map + .get(&edit.monitored_address) + .ok_or_else(|| "monitored address not found".to_string())?; + if monitored.ip != edit.target_ip { + return Err("monitor target IP mismatch".to_string()); + } + } + + if !Self::remember_newest_monitor_event(edit).await { + return Ok(false); + } + + let deletion_marker = if edit.action == MONITOR_ACTION_REMOVE { + Self::deletion_marker_for(&edit.monitored_address).await + } else { + None + }; + let mut address_map = ADDRESS_MAP.lock().await; + let monitored = address_map + .get_mut(&edit.monitored_address) + .ok_or_else(|| "monitored address not found".to_string())?; + + match edit.action { + MONITOR_ACTION_ADD => { + if !monitored.monitoring.contains(&edit.monitoring_address) { + monitored.monitoring.push(edit.monitoring_address.clone()); + } + monitored.deleted_timestamp = 0; + monitored.deleted_block = 0; + } + MONITOR_ACTION_REMOVE => { + monitored + .monitoring + .retain(|monitor| monitor != &edit.monitoring_address); + if monitored.monitoring.is_empty() { + let (deleted_timestamp, deleted_block) = + deletion_marker.unwrap_or((edit.modified_timestamp, edit.modified_block)); + Self::mark_deleted_and_cascade( + &mut address_map, + &edit.monitored_address, + local_short, + deleted_timestamp, + deleted_block, + ); + } + } + _ => unreachable!(), + } + + Ok(true) } async fn apply_monitor(params: MonitorAddressParams, action: u8) -> RpcResponse { @@ -192,7 +392,6 @@ impl NodeInfo { map, .. } = params; - let current_timestamp = Utc::now().timestamp_millis() as u64; if edit.action != action { return RpcResponse::Binary(b"Error: Invalid monitor action".to_vec()); @@ -200,74 +399,127 @@ impl NodeInfo { let local_short = wallet.saved.short_address.clone(); if remote_ip.is_empty() && edit.monitoring_address == local_short { - edit.modified_timestamp = current_timestamp; + edit.modified_timestamp = Utc::now().timestamp_millis() as u64; + edit.modified_block = get_height(&db).saturating_add(1); edit.modified_signature = Self::monitor_signature( action, &edit.monitored_address, &edit.monitoring_address, &edit.target_ip, - current_timestamp, + edit.modified_timestamp, + edit.modified_block, &wallet, ) .await; } - if !Self::verify_monitor_edit(&edit, &db).await { - return RpcResponse::Binary(b"Error: Could not validate monitor signature".to_vec()); - } - - let is_new_event = Self::remember_monitor_event(&edit).await; - if !is_new_event { - return RpcResponse::Binary(b"Success".to_vec()); - } - - { - let mut address_map = ADDRESS_MAP.lock().await; - let Some(monitored) = address_map.get_mut(&edit.monitored_address) else { - return RpcResponse::Binary(b"Error: monitored address not found".to_vec()); - }; - if monitored.ip != edit.target_ip { - return RpcResponse::Binary(b"Error: monitor target IP mismatch".to_vec()); - } - - match action { - MONITOR_ACTION_ADD => { - if !monitored.monitoring.contains(&edit.monitoring_address) { - monitored.monitoring.push(edit.monitoring_address.clone()); - } - monitored.deleted_timestamp = 0; - monitored.deleted_block = 0; - } - MONITOR_ACTION_REMOVE => { - let before = monitored.monitoring.len(); - monitored - .monitoring - .retain(|monitor| monitor != &edit.monitoring_address); - if before != monitored.monitoring.len() && monitored.monitoring.is_empty() { - let deleted_block = get_height(&db) + 1; - Self::mark_deleted_and_cascade( - &mut address_map, - &edit.monitored_address, - &local_short, - current_timestamp, - deleted_block, - ); - } - } - _ => return RpcResponse::Binary(b"Error: Invalid monitor action".to_vec()), - } + match Self::apply_verified_monitor_edit(&edit, &db, &local_short).await { + Ok(false) => return RpcResponse::Binary(b"Success".to_vec()), + Err(err) => return RpcResponse::Binary(format!("Error: {err}").into_bytes()), + Ok(true) => {} } Self::persist_recovery_snapshot("monitor update").await; Self::broadcast_monitor_event(map, &edit, &remote_ip).await; - RpcResponse::Binary(b"Success".to_vec()) } - pub async fn set_monitors_from_mapping(address: &str, monitors: Vec) { - let mut address_map = ADDRESS_MAP.lock().await; - if let Some(node) = address_map.get_mut(address) { - node.monitoring = monitors; + pub async fn import_signed_monitor_event( + edit: SignedMonitorEdit, + db: &Db, + local_short: &str, + ) -> Result { + Self::apply_verified_monitor_edit(&edit, db, local_short).await + } + + pub(crate) async fn signed_monitor_state_bytes() -> Vec { + let mut events: Vec = + MONITOR_EVENT_STATE.lock().await.values().cloned().collect(); + events.sort_by(|left, right| { + Self::monitor_event_order(left) + .cmp(&Self::monitor_event_order(right)) + .then_with(|| { + Self::monitor_relation_key(left).cmp(&Self::monitor_relation_key(right)) + }) + }); + + let mut data = Vec::with_capacity(events.len() * MONITOR_EVENT_BYTES); + for event in events { + let Some(monitored) = Wallet::short_address_to_bytes(&event.monitored_address) else { + continue; + }; + let Some(monitoring) = Wallet::short_address_to_bytes(&event.monitoring_address) else { + continue; + }; + let Ok(signature) = decode(&event.modified_signature) else { + continue; + }; + if signature.len() != Wallet::SIGNATURE_LENGTH { + continue; + } + data.push(event.action); + data.extend_from_slice(&monitored); + data.extend_from_slice(&monitoring); + data.extend_from_slice(&ip_to_binary(&event.target_ip)); + data.extend_from_slice(&event.modified_timestamp.to_le_bytes()); + data.extend_from_slice(&event.modified_block.to_le_bytes()); + data.extend_from_slice(&signature); } + data + } + + pub async fn signed_monitor_state() -> RpcResponse { + RpcResponse::Binary(Self::signed_monitor_state_bytes().await) + } + + pub(crate) async fn load_monitor_recovery_state(bytes: &[u8]) -> usize { + let mut state = MONITOR_EVENT_STATE.lock().await; + state.clear(); + let mut loaded = 0; + for chunk in bytes.chunks_exact(MONITOR_EVENT_BYTES) { + let Some(edit) = Self::monitor_event_from_bytes(chunk) else { + continue; + }; + // A restarted node no longer has the live sockets represented by + // old additions. Retain signed removals as historical deletion + // evidence; current additions are recreated by operational peers. + if edit.action != MONITOR_ACTION_REMOVE { + continue; + } + state.insert(Self::monitor_relation_key(&edit), edit); + loaded += 1; + } + loaded + } + + pub fn monitor_event_from_bytes(bytes: &[u8]) -> Option { + if bytes.len() != MONITOR_EVENT_BYTES { + return None; + } + let monitored_address = Wallet::bytes_to_short_address( + &bytes[MONITOR_EVENT_MONITORED_OFFSET..MONITOR_EVENT_MONITORING_OFFSET], + )?; + let monitoring_address = Wallet::bytes_to_short_address( + &bytes[MONITOR_EVENT_MONITORING_OFFSET..MONITOR_EVENT_TARGET_IP_OFFSET], + )?; + Some(SignedMonitorEdit { + action: bytes[MONITOR_EVENT_ACTION_OFFSET], + monitored_address, + monitoring_address, + target_ip: crate::common::binary_conversions::binary_to_ip( + bytes[MONITOR_EVENT_TARGET_IP_OFFSET..MONITOR_EVENT_TIMESTAMP_OFFSET].to_vec(), + ), + modified_timestamp: u64::from_le_bytes( + bytes[MONITOR_EVENT_TIMESTAMP_OFFSET..MONITOR_EVENT_BLOCK_OFFSET] + .try_into() + .ok()?, + ), + modified_block: u32::from_le_bytes( + bytes[MONITOR_EVENT_BLOCK_OFFSET..MONITOR_EVENT_SIGNATURE_OFFSET] + .try_into() + .ok()?, + ), + modified_signature: crate::encode(&bytes[MONITOR_EVENT_SIGNATURE_OFFSET..]), + }) } } diff --git a/src/records/memory/network_mapping/persistence.rs b/src/records/memory/network_mapping/persistence.rs index 3e4e51c..2f89d5b 100644 --- a/src/records/memory/network_mapping/persistence.rs +++ b/src/records/memory/network_mapping/persistence.rs @@ -13,6 +13,11 @@ fn recovery_snapshot_path() -> PathBuf { PathBuf::from(db_path).join("network_map_recovery.bin") } +fn monitor_recovery_snapshot_path() -> PathBuf { + let (_, _, _, _, _, _, db_path, _, _) = block_extension_and_paths(); + PathBuf::from(db_path).join("network_monitor_recovery.bin") +} + impl NodeInfo { pub async fn save_recovery_snapshot() -> Result<(), String> { let path = recovery_snapshot_path(); @@ -60,7 +65,13 @@ impl NodeInfo { crate::tokio::fs::write(&path, data) .await - .map_err(|err| format!("failed to write network-map recovery snapshot: {err}")) + .map_err(|err| format!("failed to write network-map recovery snapshot: {err}"))?; + crate::tokio::fs::write( + monitor_recovery_snapshot_path(), + Self::signed_monitor_state_bytes().await, + ) + .await + .map_err(|err| format!("failed to write monitor recovery snapshot: {err}")) } pub async fn load_recovery_snapshot() -> Result { @@ -140,6 +151,17 @@ impl NodeInfo { let mut map = ADDRESS_MAP.lock().await; *map = recovered; + drop(map); + + if let Ok(monitor_bytes) = read(monitor_recovery_snapshot_path()).await { + if monitor_bytes.len() + % crate::records::memory::network_mapping::structs::MONITOR_EVENT_BYTES + != 0 + { + warn!("[network_map] monitor recovery snapshot had trailing partial bytes"); + } + Self::load_monitor_recovery_state(&monitor_bytes).await; + } Ok(loaded) } diff --git a/src/records/memory/network_mapping/queries.rs b/src/records/memory/network_mapping/queries.rs index 995a11c..8b268f5 100644 --- a/src/records/memory/network_mapping/queries.rs +++ b/src/records/memory/network_mapping/queries.rs @@ -95,10 +95,11 @@ impl NodeInfo { Wallet::short_address_to_bytes(&node_info.added_by).unwrap_or_default(); let added_timestamp_bytes = node_info.added_timestamp.to_le_bytes(); - let deleted_timestamp_bytes = node_info.deleted_timestamp.to_le_bytes(); - let deleted_block_bytes = node_info.deleted_block.to_le_bytes(); - let monitor_count = node_info.monitoring.len().min(u16::MAX as usize) as u16; - let monitor_count_bytes = monitor_count.to_le_bytes(); + // Network-map snapshots contain signed membership only. Liveness + // is transferred separately as signed monitor assertions. + let deleted_timestamp_bytes = 0_u64.to_le_bytes(); + let deleted_block_bytes = 0_u32.to_le_bytes(); + let monitor_count_bytes = 0_u16.to_le_bytes(); let added_signature_bytes = decode(node_info.added_signature.clone()).unwrap(); @@ -118,13 +119,6 @@ impl NodeInfo { data.extend_from_slice(&deleted_timestamp_bytes); data.extend_from_slice(&deleted_block_bytes); data.extend_from_slice(&monitor_count_bytes); - for monitor in node_info.monitoring.iter().take(monitor_count as usize) { - if let Some(monitor_bytes) = Wallet::short_address_to_bytes(monitor) { - data.extend_from_slice(&monitor_bytes); - } else { - data.extend_from_slice(&[0u8; Wallet::SHORT_ADDRESS_BYTES_LENGTH]); - } - } } RpcResponse::Binary(data) } diff --git a/src/records/memory/network_mapping/structs.rs b/src/records/memory/network_mapping/structs.rs index 038df76..4d49ffb 100644 --- a/src/records/memory/network_mapping/structs.rs +++ b/src/records/memory/network_mapping/structs.rs @@ -25,6 +25,18 @@ pub const NODE_DELETED_BLOCK_OFFSET: usize = NODE_DELETED_TIMESTAMP_OFFSET + NOD pub const NODE_MONITOR_COUNT_OFFSET: usize = NODE_DELETED_BLOCK_OFFSET + NODE_DELETED_BLOCK_BYTES; pub const NODE_RECORD_FIXED_BYTES: usize = NODE_MONITOR_COUNT_OFFSET + NODE_MONITOR_COUNT_BYTES; +pub const MONITOR_EVENT_ACTION_OFFSET: usize = 0; +pub const MONITOR_EVENT_MONITORED_OFFSET: usize = MONITOR_EVENT_ACTION_OFFSET + 1; +pub const MONITOR_EVENT_MONITORING_OFFSET: usize = + MONITOR_EVENT_MONITORED_OFFSET + Wallet::SHORT_ADDRESS_BYTES_LENGTH; +pub const MONITOR_EVENT_TARGET_IP_OFFSET: usize = + MONITOR_EVENT_MONITORING_OFFSET + Wallet::SHORT_ADDRESS_BYTES_LENGTH; +pub const MONITOR_EVENT_TIMESTAMP_OFFSET: usize = MONITOR_EVENT_TARGET_IP_OFFSET + NODE_IP_BYTES; +pub const MONITOR_EVENT_BLOCK_OFFSET: usize = MONITOR_EVENT_TIMESTAMP_OFFSET + NODE_TIMESTAMP_BYTES; +pub const MONITOR_EVENT_SIGNATURE_OFFSET: usize = + MONITOR_EVENT_BLOCK_OFFSET + NODE_DELETED_BLOCK_BYTES; +pub const MONITOR_EVENT_BYTES: usize = MONITOR_EVENT_SIGNATURE_OFFSET + Wallet::SIGNATURE_LENGTH; + // SignedNodeEdit carries the signed node membership payload used by add/delete updates. #[derive(Debug, Clone)] pub struct SignedNodeEdit { @@ -43,6 +55,7 @@ pub struct SignedMonitorEdit { pub monitoring_address: String, pub target_ip: String, pub modified_timestamp: u64, + pub modified_block: u32, pub modified_signature: String, } diff --git a/src/records/memory/response_channels.rs b/src/records/memory/response_channels.rs index aeeb6bd..d95cf21 100644 --- a/src/records/memory/response_channels.rs +++ b/src/records/memory/response_channels.rs @@ -99,6 +99,21 @@ pub async fn reserve_entry_with_context( ); (tx, rx) }; + let cleanup_map = map.clone(); + tokio::spawn(async move { + // Normal replies retire their UID in route_reply. This fallback + // covers send failures and cancellation of the requesting task. + crate::sleep(Duration::from_secs(60)).await; + let still_active = { + let map = cleanup_map.lock().await; + map.get(&key) + .map(|entry| entry.expires_at.is_none()) + .unwrap_or(false) + }; + if still_active { + delete_entry(cleanup_map, key).await; + } + }); return (key, tx, rx); } } @@ -181,3 +196,20 @@ pub async fn delete_entry(map: Arc>, key: Byte3) { channel_pair.expires_at = Some(Instant::now() + Duration::from_secs(30)); } } + +#[cfg(test)] +mod tests { + use super::*; + + #[tokio::test(start_paused = true)] + async fn abandoned_response_uid_is_retired_automatically() { + let map = Arc::new(Mutex::new(Command::new())); + let (uid, _tx, _rx) = reserve_entry(map.clone()).await; + + tokio::task::yield_now().await; + tokio::time::advance(Duration::from_secs(61)).await; + tokio::task::yield_now().await; + + assert!(is_retired_entry(map, uid).await); + } +} diff --git a/src/records/memory/structs.rs b/src/records/memory/structs.rs index d569f8c..76ef0bb 100644 --- a/src/records/memory/structs.rs +++ b/src/records/memory/structs.rs @@ -7,6 +7,14 @@ use crate::Mutex; use crate::Ordering; use crate::TcpStream; +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub enum ConnectionHealth { + Connected, + Busy, + Unresponsive, + Closed, +} + // ConnectionInfo stores the live metadata associated with one tracked peer. #[derive(Debug, Clone)] pub struct ConnectionInfo { @@ -18,7 +26,11 @@ pub struct ConnectionInfo { pub wallet_short_address: String, pub wallet_registry_synced: bool, pub network_map_synced: bool, + pub local_setup_complete: bool, + pub remote_setup_complete: bool, + pub local_setup_acknowledged: bool, pub ready: bool, + pub health: ConnectionHealth, } // ConnectionKey is the stable lookup key used inside the in-memory connection map. @@ -213,7 +225,11 @@ impl ConnectionInfo { wallet_short_address, wallet_registry_synced: false, network_map_synced: false, + local_setup_complete: false, + remote_setup_complete: false, + local_setup_acknowledged: false, ready: false, + health: ConnectionHealth::Connected, } } } diff --git a/src/rpc/client/handshake.rs b/src/rpc/client/handshake.rs index 9aaf307..62f56c6 100644 --- a/src/rpc/client/handshake.rs +++ b/src/rpc/client/handshake.rs @@ -157,8 +157,15 @@ async fn perform_handshake(mut params: Handshake) -> io::Result<()> { } async fn send_and_receive_handshake(params: &mut Handshake, data: &[u8]) -> io::Result> { - params.stream.write_all(data).await?; - params.stream.flush().await?; + timeout( + Duration::from_secs(HANDSHAKE_RESPONSE_TIMEOUT_SECONDS), + async { + params.stream.write_all(data).await?; + params.stream.flush().await + }, + ) + .await + .map_err(|_| io::Error::new(io::ErrorKind::TimedOut, "Timed out sending handshake"))??; let mut buffer = vec![0u8; HANDSHAKE_RESPONSE_BYTES]; let mut total_read = 0usize; diff --git a/src/rpc/client/handshake_processing.rs b/src/rpc/client/handshake_processing.rs index adf3653..7030450 100644 --- a/src/rpc/client/handshake_processing.rs +++ b/src/rpc/client/handshake_processing.rs @@ -5,22 +5,18 @@ use crate::config::SETTINGS; use crate::encode; use crate::io; use crate::log::{info, warn}; -use crate::miner::flag::{ - clear_mining_stop_request, request_mining_stop, set_mining_state, set_node_mode, MiningState, - NodeMode, -}; +use crate::miner::flag::{begin_chain_sync, is_normal_mode}; use crate::orphans::structs::OrphanCheckup2; use crate::orphans::sync_check::sync_checkup; use crate::orphans::torrent_candidates::hydrate_torrent_candidates; use crate::records::block_height::get_block_height::get_height; use crate::records::memory::connections::{ mark_peer_network_map_synced, mark_peer_operational, mark_peer_wallet_registry_synced, - set_reconnect_context, startup_synced_peer_streams, CONNECTIONS, + spawn_peer_setup_retry, startup_synced_peer_streams, CONNECTIONS, }; -use crate::records::memory::enums::{ClientType, ConnectionType}; use crate::records::memory::network_mapping::NodeInfo; use crate::records::memory::response_channels::Command; -use crate::records::memory::structs::{Connection, StoreConnectionParams}; +use crate::records::memory::structs::Connection; use crate::rpc::client::genesis_compat::ensure_compatible_genesis; use crate::rpc::client::handshake::connect_and_handshake; use crate::rpc::client::register_wallet::register_connected_wallet; @@ -31,7 +27,9 @@ use crate::rpc::handshake_constants::{ HANDSHAKE_ADDRESS_OFFSET, HANDSHAKE_MESSAGE_BYTES, HANDSHAKE_RESPONSE_BYTES, HANDSHAKE_SIGNATURE_OFFSET, }; -use crate::rpc::server::connection_memory_manager::remove_key_from_memory; +use crate::rpc::server::connection_memory_manager::{ + remove_stream_from_memory, write_outgoing_miner_to_memory, +}; use crate::rpc::server::rpc_command_loop::start_loop; use crate::sled::Db; use crate::sleep; @@ -59,23 +57,18 @@ pub struct BootstrapParams { pub fn spawn_bootstrap_peer_discovery(params: BootstrapParams) { tokio::spawn(async move { - let run_startup_sync = params.run_startup_sync; if let Err(e) = bootstrap_peer_discovery(params).await { - if !run_startup_sync { - set_node_mode(NodeMode::Normal); - clear_mining_stop_request(); - set_mining_state(MiningState::Idle); - } eprintln!("[bootstrap] error: {e}"); } }); } pub async fn bootstrap_peer_discovery(mut params: BootstrapParams) -> Result<(), String> { - if params.run_startup_sync { - set_node_mode(NodeMode::Syncing); - request_mining_stop(); - } + let chain_sync_guard = if params.run_startup_sync { + Some(begin_chain_sync().await) + } else { + None + }; let (_, _, local_endpoint) = get_ip_and_port().await; let max = SETTINGS.outgoing_connections; let mut current_key = params.connections_key.clone(); @@ -89,7 +82,7 @@ pub async fn bootstrap_peer_discovery(mut params: BootstrapParams) -> Result<(), let connections = CONNECTIONS.read().await; connections .as_ref() - .map(|connection| connection.count_outgoing_connections()) + .map(|connection| connection.count_ready_outgoing_connections()) .unwrap_or(0) }; if outgoing_connections >= max as usize { @@ -221,13 +214,15 @@ pub async fn bootstrap_peer_discovery(mut params: BootstrapParams) -> Result<(), info!("[sync] post-sync checks complete, mining grace period started"); for (peer_key, _) in startup_synced_peer_streams().await { - mark_peer_operational(&peer_key, params.map.clone()).await; + if !mark_peer_operational(&peer_key, params.map.clone()).await { + spawn_peer_setup_retry(peer_key, params.map.clone()); + } } sleep(Duration::from_secs(15)).await; info!("[sync] mining grace period complete, mining resuming"); - set_node_mode(NodeMode::Normal); - clear_mining_stop_request(); - set_mining_state(MiningState::Idle); + if let Some(guard) = chain_sync_guard { + guard.finish(); + } Ok(()) } @@ -296,27 +291,17 @@ pub async fn process_handshake_response( // create and arc mutex of the stream let stream = Arc::new(Mutex::new(params.stream)); let connections_key = params.addr.clone(); - let socket_parts: Vec<&str> = params.addr.split(':').collect(); - if socket_parts.len() == 2 { - let ip = socket_parts[0]; - let port: u16 = socket_parts[1].parse().unwrap_or(0); - set_reconnect_context(params.db.clone(), params.wallet.clone(), params.map.clone()).await; - let mut conn = CONNECTIONS.write().await; - if let Some(manager) = conn.as_mut() { - if !manager.store_connection(StoreConnectionParams { - connection_type: ConnectionType::Outgoing, - ip: ip.to_string(), - port, - stream: Arc::clone(&stream), - client_type: ClientType::Miner, - wallet_short_address: returned_short_address.clone(), - command_map: params.map.clone(), - }) { - return Err(io::Error::other( - "The connection is already in the connection manager Please wait 10 minutes and try again", - )); - } - } + if !write_outgoing_miner_to_memory( + ¶ms.addr, + stream.clone(), + returned_short_address, + params.map.clone(), + ) + .await + { + return Err(io::Error::other( + "The connection is already in the connection manager Please wait 10 minutes and try again", + )); } let listener_stream = Arc::clone(&stream); tokio::spawn(start_loop( @@ -337,7 +322,7 @@ pub async fn process_handshake_response( ) .await { - remove_key_from_memory(&connections_key).await; + remove_stream_from_memory(&stream).await; return Err(io::Error::other(format!( "Wallet registry sync failed after handshake: {err}" ))); @@ -353,7 +338,7 @@ pub async fn process_handshake_response( ) .await { - remove_key_from_memory(&connections_key).await; + remove_stream_from_memory(&stream).await; return Err(io::Error::other(format!( "Wallet registration failed after handshake: {err}" ))); @@ -366,9 +351,9 @@ pub async fn process_handshake_response( ) .await .map_err(|err| { - let connections_key = connections_key.clone(); + let failed_stream = stream.clone(); tokio::spawn(async move { - remove_key_from_memory(&connections_key).await; + remove_stream_from_memory(&failed_stream).await; }); io::Error::other(format!( "Remote height lookup failed before network join: {err}" @@ -382,9 +367,9 @@ pub async fn process_handshake_response( ) .await .map_err(|err| { - let connections_key = connections_key.clone(); + let failed_stream = stream.clone(); tokio::spawn(async move { - remove_key_from_memory(&connections_key).await; + remove_stream_from_memory(&failed_stream).await; }); io::Error::other(format!("Genesis compatibility check failed: {err}")) })?; @@ -412,9 +397,9 @@ pub async fn process_handshake_response( ) .await .map_err(|err| { - let connections_key = connections_key.clone(); + let failed_stream = stream.clone(); tokio::spawn(async move { - remove_key_from_memory(&connections_key).await; + remove_stream_from_memory(&failed_stream).await; }); io::Error::other(format!("Network join/map sync failed: {err}")) })?; @@ -433,8 +418,10 @@ pub async fn process_handshake_response( spawn_bootstrap_peer_discovery(bsparams); } else { - if crate::miner::flag::is_normal_mode() { - mark_peer_operational(&connections_key, params.map.clone()).await; + if is_normal_mode() { + if !mark_peer_operational(&connections_key, params.map.clone()).await { + spawn_peer_setup_retry(connections_key.clone(), params.map.clone()); + } } } Ok(()) diff --git a/src/rpc/command_maps.rs b/src/rpc/command_maps.rs index 1364ea2..704fba1 100644 --- a/src/rpc/command_maps.rs +++ b/src/rpc/command_maps.rs @@ -67,6 +67,8 @@ pub const RPC_NETWORK_MONITOR_ADD: u8 = 50; pub const RPC_NETWORK_MONITOR_REMOVE: u8 = 51; pub const RPC_STORAGE_LOOKUP_COST: u8 = 52; pub const RPC_STORAGE_LOOKUP: u8 = 53; +pub const RPC_NETWORK_MONITOR_STATE: u8 = 54; +pub const RPC_SETUP_COMPLETE: u8 = 55; pub const RPC_REPLY: u8 = 255; pub const MAX_RPC_REPLY_BYTES: usize = 64 * 1024 * 1024; diff --git a/src/rpc/commands/add_network_node.rs b/src/rpc/commands/add_network_node.rs index 57b8555..52f2220 100644 --- a/src/rpc/commands/add_network_node.rs +++ b/src/rpc/commands/add_network_node.rs @@ -52,16 +52,8 @@ pub async fn add_network_node( let monitor_count = read_bytes_from_stream::read_u16_from_stream(connections_key, stream_locked.clone()).await? as usize; - let mut monitors = Vec::with_capacity(monitor_count); - for _ in 0..monitor_count { - let monitor_bytes = read_bytes_from_stream::read_short_address_from_stream( - connections_key, - stream_locked.clone(), - ) - .await?; - if let Some(monitor) = Wallet::bytes_to_short_address(&monitor_bytes) { - monitors.push(monitor); - } + if monitor_count != 0 { + return Err("error: Node membership updates cannot contain monitor state".to_string()); } let remote_ip = read_bytes_from_stream::read_caller_ip(stream_locked).await?; @@ -77,7 +69,7 @@ pub async fn add_network_node( modified_timestamp: added_timestamp, modified_signature: added_signature, }, - monitors, + monitors: Vec::new(), blocks_mined: 0_u8, remote_ip, db: db.clone(), diff --git a/src/rpc/commands/mod.rs b/src/rpc/commands/mod.rs index 56635af..7198791 100644 --- a/src/rpc/commands/mod.rs +++ b/src/rpc/commands/mod.rs @@ -21,6 +21,7 @@ pub mod memory_by_signature; pub mod network_info; pub mod network_monitor_add; pub mod network_monitor_remove; +pub mod network_monitor_state; pub mod nft_list; pub mod nft_lookup; pub mod receive_torrent; diff --git a/src/rpc/commands/network_monitor_add.rs b/src/rpc/commands/network_monitor_add.rs index 52e30f5..63b7b95 100644 --- a/src/rpc/commands/network_monitor_add.rs +++ b/src/rpc/commands/network_monitor_add.rs @@ -42,6 +42,9 @@ pub async fn network_monitor_add( let modified_timestamp = read_bytes_from_stream::read_u64_from_stream(connections_key, stream_locked.clone()) .await?; + let modified_block = + read_bytes_from_stream::read_u32_from_stream(connections_key, stream_locked.clone()) + .await?; let modified_signature = read_bytes_from_stream::read_signature_from_stream(connections_key, stream_locked.clone()) .await?; @@ -62,6 +65,7 @@ pub async fn network_monitor_add( monitoring_address, target_ip, modified_timestamp, + modified_block, modified_signature, }, remote_ip, diff --git a/src/rpc/commands/network_monitor_remove.rs b/src/rpc/commands/network_monitor_remove.rs index 4420026..b16a324 100644 --- a/src/rpc/commands/network_monitor_remove.rs +++ b/src/rpc/commands/network_monitor_remove.rs @@ -42,6 +42,9 @@ pub async fn network_monitor_remove( let modified_timestamp = read_bytes_from_stream::read_u64_from_stream(connections_key, stream_locked.clone()) .await?; + let modified_block = + read_bytes_from_stream::read_u32_from_stream(connections_key, stream_locked.clone()) + .await?; let modified_signature = read_bytes_from_stream::read_signature_from_stream(connections_key, stream_locked.clone()) .await?; @@ -62,6 +65,7 @@ pub async fn network_monitor_remove( monitoring_address, target_ip, modified_timestamp, + modified_block, modified_signature, }, remote_ip, diff --git a/src/rpc/commands/network_monitor_state.rs b/src/rpc/commands/network_monitor_state.rs new file mode 100644 index 0000000..7003958 --- /dev/null +++ b/src/rpc/commands/network_monitor_state.rs @@ -0,0 +1,6 @@ +use crate::records::memory::network_mapping::NodeInfo; +use crate::rpc::responses::RpcResponse; + +pub async fn network_monitor_state() -> RpcResponse { + NodeInfo::signed_monitor_state().await +} diff --git a/src/rpc/commands/receive_torrent.rs b/src/rpc/commands/receive_torrent.rs index 0a1da6f..ef9155c 100644 --- a/src/rpc/commands/receive_torrent.rs +++ b/src/rpc/commands/receive_torrent.rs @@ -1,11 +1,7 @@ use crate::common::check_genesis::genesis_checkup; use crate::common::skein::skein_128_hash_bytes; use crate::log::{error, info, warn}; -use crate::miner::flag::{ - clear_mining_stop_request, is_normal_mode, is_reorganizing_mode, is_syncing_mode, - request_mining_stop, set_mining_state, set_node_mode, wait_for_mining_idle, MiningState, - NodeMode, -}; +use crate::miner::flag::{begin_chain_sync, is_normal_mode, is_reorganizing_mode, is_syncing_mode}; use crate::orphans::checkup_state::{ finish_orphan_check, request_orphan_recheck, try_begin_orphan_check, }; @@ -79,9 +75,7 @@ async fn run_live_catchup_sync( info!( "[sync] live catch-up started: local_height={local_height} remote_height={remote_height}" ); - set_node_mode(NodeMode::Syncing); - request_mining_stop(); - wait_for_mining_idle().await; + let chain_sync_guard = begin_chain_sync().await; let sync_result = node_syncing( stream.clone(), @@ -124,9 +118,7 @@ async fn run_live_catchup_sync( } } - set_node_mode(NodeMode::Normal); - clear_mining_stop_request(); - set_mining_state(MiningState::Idle); + chain_sync_guard.finish(); if sync_result.is_ok() { info!("[sync] live catch-up complete, normal mode restored"); } diff --git a/src/rpc/commands/tx_submit.rs b/src/rpc/commands/tx_submit.rs index c98c6ad..afb9a4c 100644 --- a/src/rpc/commands/tx_submit.rs +++ b/src/rpc/commands/tx_submit.rs @@ -40,6 +40,7 @@ async fn broadcast_tx(tx_bytes: Vec, map: Arc>) { // Broadcast newly accepted mempool transactions only to miner peers, // since those are the nodes that need transaction fan-out most. let nodes = get_live_node_broadcast_peers().await; + let mut messages = Vec::with_capacity(nodes.len()); for (key, stream) in nodes { let client_type = get_client_type_from_memory(&key) @@ -61,8 +62,9 @@ async fn broadcast_tx(tx_bytes: Vec, map: Arc>) { message.push(RPC_SUBMIT_TRANSACTION); message.extend_from_slice(&uid); message.extend_from_slice(&tx_bytes); - RpcResponse::send_raw(&stream, Some(&key), &message).await; + messages.push((key, stream, message)); } + RpcResponse::broadcast_raw(messages).await; } macro_rules! submit_storage_value { diff --git a/src/rpc/commands/wallet_register.rs b/src/rpc/commands/wallet_register.rs index d3e4d19..ade8e18 100644 --- a/src/rpc/commands/wallet_register.rs +++ b/src/rpc/commands/wallet_register.rs @@ -33,6 +33,7 @@ async fn broadcast_wallet_registration( .unwrap_or_default() }; + let mut messages = Vec::with_capacity(streams.len()); for (peer_key, unlocked_stream) in streams { if let Some((peer_ip, _)) = peer_key.rsplit_once(':') { if !remote_ip.is_empty() && peer_ip == remote_ip { @@ -59,8 +60,9 @@ async fn broadcast_wallet_registration( message.extend_from_slice(&public_key_bytes); message.extend_from_slice(&signature_bytes); - RpcResponse::send_raw(&unlocked_stream, Some(&peer_key), &message).await; + messages.push((peer_key, unlocked_stream, message)); } + RpcResponse::broadcast_raw(messages).await; } pub async fn register( diff --git a/src/rpc/read_bytes_from_stream.rs b/src/rpc/read_bytes_from_stream.rs index 63e7608..06ad6f0 100644 --- a/src/rpc/read_bytes_from_stream.rs +++ b/src/rpc/read_bytes_from_stream.rs @@ -11,7 +11,8 @@ use crate::{timeout, AsyncReadExt, AsyncWriteExt, Duration}; const STREAM_READ_TIMEOUT_SECONDS: u64 = 30; async fn read_exact_from_stream( - key: &str, + _key: &str, + stream_locked: &Arc>, stream: &mut TcpStream, buffer: &mut [u8], ) -> Result<(), String> { @@ -28,14 +29,14 @@ async fn read_exact_from_stream( if let Err(e) = stream.shutdown().await { warn!("Error shutting down stream: {e}"); } - connection_memory_manager::remove_key_from_memory(key).await; + connection_memory_manager::remove_stream_from_memory(stream_locked).await; Err(err.to_string()) } Err(_) => { if let Err(e) = stream.shutdown().await { warn!("Error shutting down stream: {e}"); } - connection_memory_manager::remove_key_from_memory(key).await; + connection_memory_manager::remove_stream_from_memory(stream_locked).await; Err(format!( "Timed out reading from stream after {STREAM_READ_TIMEOUT_SECONDS} seconds" )) @@ -63,7 +64,7 @@ pub async fn read_first_byte( // if the peer disconnects mid-read. let mut stream = stream_locked.lock().await; let mut buffer = [0u8; 1]; - read_exact_from_stream(key, &mut stream, &mut buffer).await?; + read_exact_from_stream(key, &stream_locked, &mut stream, &mut buffer).await?; drop(stream); Ok(u8::from_le_bytes(buffer)) } @@ -76,7 +77,7 @@ pub async fn read_uid_from_stream( // can be matched through the shared hashmap. let mut stream = stream_locked.lock().await; let mut uid_bytes = [0u8; 3]; - read_exact_from_stream(key, &mut stream, &mut uid_bytes).await?; + read_exact_from_stream(key, &stream_locked, &mut stream, &mut uid_bytes).await?; let uid = u32::from_le_bytes([0, uid_bytes[0], uid_bytes[1], uid_bytes[2]]); drop(stream); Ok((uid, uid_bytes)) @@ -88,7 +89,7 @@ pub async fn read_u8_from_stream( ) -> Result { let mut stream = stream_locked.lock().await; let mut buffer = [0u8; 1]; - read_exact_from_stream(key, &mut stream, &mut buffer).await?; + read_exact_from_stream(key, &stream_locked, &mut stream, &mut buffer).await?; drop(stream); Ok(u8::from_le_bytes(buffer)) } @@ -99,7 +100,7 @@ pub async fn read_u16_from_stream( ) -> Result { let mut stream = stream_locked.lock().await; let mut buffer = [0u8; 2]; - read_exact_from_stream(key, &mut stream, &mut buffer).await?; + read_exact_from_stream(key, &stream_locked, &mut stream, &mut buffer).await?; drop(stream); Ok(u16::from_le_bytes(buffer)) } @@ -110,7 +111,7 @@ pub async fn read_u32_from_stream( ) -> Result { let mut stream = stream_locked.lock().await; let mut buffer = [0u8; 4]; - read_exact_from_stream(key, &mut stream, &mut buffer).await?; + read_exact_from_stream(key, &stream_locked, &mut stream, &mut buffer).await?; drop(stream); Ok(u32::from_le_bytes(buffer)) } @@ -121,7 +122,7 @@ pub async fn read_u64_from_stream( ) -> Result { let mut stream = stream_locked.lock().await; let mut buffer = [0u8; 8]; - read_exact_from_stream(key, &mut stream, &mut buffer).await?; + read_exact_from_stream(key, &stream_locked, &mut stream, &mut buffer).await?; drop(stream); Ok(u64::from_le_bytes(buffer)) } @@ -132,7 +133,7 @@ pub async fn read_u128_from_stream( ) -> Result { let mut stream = stream_locked.lock().await; let mut buffer = [0u8; 16]; - read_exact_from_stream(key, &mut stream, &mut buffer).await?; + read_exact_from_stream(key, &stream_locked, &mut stream, &mut buffer).await?; drop(stream); Ok(u128::from_le_bytes(buffer)) } @@ -146,7 +147,7 @@ pub async fn read_usize_from_stream( // leaving interpretation to the calling command handler. let mut stream = stream_locked.lock().await; let mut buffer = vec![0u8; bytes]; - read_exact_from_stream(key, &mut stream, &mut buffer).await?; + read_exact_from_stream(key, &stream_locked, &mut stream, &mut buffer).await?; drop(stream); Ok(buffer) } @@ -157,7 +158,7 @@ pub async fn read_public_key_from_stream( ) -> Result, String> { let mut stream = stream_locked.lock().await; let mut buffer = vec![0u8; Wallet::PUBLIC_KEY_LENGTH]; - read_exact_from_stream(key, &mut stream, &mut buffer).await?; + read_exact_from_stream(key, &stream_locked, &mut stream, &mut buffer).await?; drop(stream); Ok(buffer) } @@ -170,7 +171,7 @@ pub async fn read_short_address_from_stream( // wire and stay as raw bytes until the command handler validates them. let mut stream = stream_locked.lock().await; let mut buffer = vec![0u8; Wallet::SHORT_ADDRESS_BYTES_LENGTH]; - read_exact_from_stream(key, &mut stream, &mut buffer).await?; + read_exact_from_stream(key, &stream_locked, &mut stream, &mut buffer).await?; drop(stream); Ok(buffer) } @@ -192,7 +193,7 @@ pub async fn read_ip_from_stream( // the node list and peer-management RPC calls. let mut stream = stream_locked.lock().await; let mut buffer = vec![0u8; 16]; - read_exact_from_stream(key, &mut stream, &mut buffer).await?; + read_exact_from_stream(key, &stream_locked, &mut stream, &mut buffer).await?; drop(stream); let ip = binary_to_ip(buffer.to_vec()); Ok(ip) @@ -206,7 +207,7 @@ pub async fn read_signature_from_stream( // into hex strings for downstream verification logic. let mut stream = stream_locked.lock().await; let mut buffer = vec![0u8; Wallet::SIGNATURE_LENGTH]; - read_exact_from_stream(key, &mut stream, &mut buffer).await?; + read_exact_from_stream(key, &stream_locked, &mut stream, &mut buffer).await?; drop(stream); let signature = encode(buffer); Ok(signature) diff --git a/src/rpc/responses.rs b/src/rpc/responses.rs index 39d4e9b..6c0f665 100644 --- a/src/rpc/responses.rs +++ b/src/rpc/responses.rs @@ -1,12 +1,21 @@ use crate::log::warn; +use crate::records::memory::connections::set_stream_health; +use crate::records::memory::structs::ConnectionHealth; use crate::rpc::command_maps::{MAX_RPC_REPLY_BYTES, RPC_REPLY}; -use crate::rpc::server::connection_memory_manager::remove_key_from_memory; +use crate::rpc::server::connection_memory_manager::remove_stream_from_memory; +use crate::timeout; use crate::Arc; use crate::AsyncWriteExt; +use crate::Duration; use crate::Mutex; use crate::Serialize; use crate::TcpStream; use std::io::ErrorKind; +use tokio::sync::Semaphore; +use tokio::task::JoinSet; + +const NETWORK_WRITE_TIMEOUT: Duration = Duration::from_secs(10); +const MAX_CONCURRENT_BROADCAST_WRITES: usize = 16; #[derive(Debug, Serialize)] pub enum RpcResponse { @@ -25,27 +34,79 @@ impl RpcResponse { ) } + async fn write_bytes( + stream: &Arc>, + connections_key: Option<&str>, + bytes: &[u8], + ) -> bool { + let peer = connections_key.unwrap_or("unknown"); + let mut stream_guard = match timeout(NETWORK_WRITE_TIMEOUT, stream.lock()).await { + Ok(guard) => guard, + Err(_) => { + set_stream_health(stream, ConnectionHealth::Busy).await; + warn!("[rpc] timed out waiting for stream lock: peer={peer}"); + return false; + } + }; + + let write_result = timeout(NETWORK_WRITE_TIMEOUT, async { + stream_guard.write_all(bytes).await?; + stream_guard.flush().await + }) + .await; + + match write_result { + Ok(Ok(())) => { + drop(stream_guard); + set_stream_health(stream, ConnectionHealth::Connected).await; + true + } + Ok(Err(err)) => { + if !Self::is_expected_disconnect(&err) { + warn!("Error sending binary response to {peer}: {err:?}"); + } + let _ = stream_guard.shutdown().await; + drop(stream_guard); + set_stream_health(stream, ConnectionHealth::Closed).await; + remove_stream_from_memory(stream).await; + false + } + Err(_) => { + drop(stream_guard); + set_stream_health(stream, ConnectionHealth::Unresponsive).await; + warn!("[rpc] socket write timed out: peer={peer}"); + false + } + } + } + pub async fn send_raw( stream: &Arc>, connections_key: Option<&str>, bytes: &[u8], - ) { + ) -> bool { // Raw sends are used for messages that already include their // command framing and payload layout. - let mut stream_guard = stream.lock().await; - if let Err(err) = stream_guard.write_all(bytes).await { - if !Self::is_expected_disconnect(&err) { - warn!("Error sending binary response: {err:?}"); - } - if let Some(key) = connections_key { - remove_key_from_memory(key).await; - } - let _ = stream_guard.shutdown().await; - return; + Self::write_bytes(stream, connections_key, bytes).await + } + + pub async fn broadcast_raw(messages: Vec<(String, Arc>, Vec)>) { + let permits = Arc::new(Semaphore::new(MAX_CONCURRENT_BROADCAST_WRITES)); + let mut writes = JoinSet::new(); + + for (peer_key, stream, message) in messages { + let permits = permits.clone(); + writes.spawn(async move { + let Ok(_permit) = permits.acquire_owned().await else { + return; + }; + Self::send_raw(&stream, Some(&peer_key), &message).await; + }); } - if let Err(err) = stream_guard.flush().await { - if !Self::is_expected_disconnect(&err) { - warn!("Error flushing stream: {err:?}"); + + while let Some(result) = writes.join_next().await { + if let Err(err) = result { + warn!("[rpc] broadcast write task failed: {err}"); } } } @@ -55,11 +116,9 @@ impl RpcResponse { stream: &Arc>, connections_key: Option<&str>, uid: u32, - ) { + ) -> bool { // Standard RPC responses are framed with command 255 and the // request UID so the waiting caller can route the reply. - let mut stream_guard = stream.lock().await; - let uid_bytes = uid.to_le_bytes()[1..4].to_vec(); let command = RPC_REPLY; match self { @@ -82,21 +141,7 @@ impl RpcResponse { response_bytes.extend_from_slice(&uid_bytes); response_bytes.extend_from_slice(&bytes_as_slice); response_bytes.extend_from_slice(payload); - if let Err(err) = stream_guard.write_all(&response_bytes).await { - if !Self::is_expected_disconnect(&err) { - warn!("Error sending binary response: {err:?}"); - } - if let Some(key) = connections_key { - remove_key_from_memory(key).await; - } - let _ = stream_guard.shutdown().await; - return; - } - } - } - if let Err(err) = stream_guard.flush().await { - if !Self::is_expected_disconnect(&err) { - warn!("Error flushing stream: {err:?}"); + return Self::write_bytes(stream, connections_key, &response_bytes).await; } } } diff --git a/src/rpc/server/command_loop_state.rs b/src/rpc/server/command_loop_state.rs index 206badf..a6b88c7 100644 --- a/src/rpc/server/command_loop_state.rs +++ b/src/rpc/server/command_loop_state.rs @@ -107,10 +107,8 @@ pub async fn next_incoming_command( && !peer_accepts_live_relay(connections_key).await { warn!( - "[rpc] closing unsynced miner peer that sent relay command: peer={connections_key} cmd={command}" + "[rpc] setup-time relay received before peer became operational: peer={connections_key} cmd={command}" ); - remove_stream_from_memory(&stream_locked).await; - return Ok(None); } // Replies and miner sync traffic belong to expected node-to-node diff --git a/src/rpc/server/connection_memory_manager.rs b/src/rpc/server/connection_memory_manager.rs index da22642..3fb76e6 100644 --- a/src/rpc/server/connection_memory_manager.rs +++ b/src/rpc/server/connection_memory_manager.rs @@ -1,9 +1,18 @@ -use crate::records::memory::connections::{spawn_retry_dropped_outgoing, CONNECTIONS}; +use crate::common::binary_conversions::{binary_to_ip, ip_to_binary}; +use crate::log::warn; +use crate::records::memory::connections::{ + set_stream_health, spawn_retry_dropped_outgoing, CONNECTIONS, +}; use crate::records::memory::enums::{ClientType, ConnectionType}; -use crate::records::memory::response_channels::Command; -use crate::records::memory::structs::StoreConnectionParams; +use crate::records::memory::response_channels::{ + delete_entry, reserve_entry_with_context, Command, +}; +use crate::records::memory::structs::{ConnectionHealth, StoreConnectionParams}; +use crate::rpc::command_maps::RPC_BLOCK_HEIGHT; +use crate::timeout; use crate::Arc; use crate::AsyncWriteExt; +use crate::Duration; use crate::Mutex; use crate::TcpStream; @@ -20,6 +29,194 @@ fn split_ip_port_key(value: &str) -> Option<(String, u16)> { Some((ip, port)) } +async fn miner_stream_health( + key: &str, + stream: Arc>, + command_map: Arc>, +) -> ConnectionHealth { + let (checkup_key, _checkup_tx, checkup_rx_mutex) = reserve_entry_with_context( + command_map.clone(), + Some(RPC_BLOCK_HEIGHT), + Some(key.to_string()), + ) + .await; + + let mut message = Vec::with_capacity(4); + message.push(RPC_BLOCK_HEIGHT); + message.extend_from_slice(&checkup_key); + + let write_result = { + let Ok(mut stream_guard) = stream.try_lock() else { + delete_entry(command_map, checkup_key).await; + set_stream_health(&stream, ConnectionHealth::Busy).await; + return ConnectionHealth::Busy; + }; + timeout(Duration::from_secs(5), async { + stream_guard.write_all(&message).await?; + stream_guard.flush().await + }) + .await + }; + + match write_result { + Ok(Ok(())) => {} + Ok(Err(err)) => { + warn!("[connection_manager] stale duplicate probe write failed for {key}: {err}"); + delete_entry(command_map, checkup_key).await; + set_stream_health(&stream, ConnectionHealth::Closed).await; + return ConnectionHealth::Closed; + } + Err(_) => { + warn!("[connection_manager] stale duplicate probe write timed out for {key}"); + delete_entry(command_map, checkup_key).await; + set_stream_health(&stream, ConnectionHealth::Unresponsive).await; + return ConnectionHealth::Unresponsive; + } + } + + let response_result = { + let mut checkup_rx = checkup_rx_mutex.lock().await; + timeout(Duration::from_secs(5), checkup_rx.recv()).await + }; + delete_entry(command_map, checkup_key).await; + + let health = if matches!(response_result, Ok(Some(_))) { + ConnectionHealth::Connected + } else { + ConnectionHealth::Unresponsive + }; + set_stream_health(&stream, health).await; + health +} + +async fn purge_stale_duplicate_miner(ip: &str, command_map: Arc>) -> bool { + let ip_bytes = ip_to_binary(ip); + let duplicate = { + let guard = CONNECTIONS.read().await; + guard.as_ref().and_then(|connection| { + connection + .connection_map + .iter() + .find_map(|(connection_key, connection_info)| { + if connection_key.ip != ip_bytes + || ClientType::from_bytes(&connection_info.client_type) + != Some(ClientType::Miner) + { + return None; + } + let connection_type = + ConnectionType::from_bytes(&connection_key.connection_type)?; + let ip = binary_to_ip(connection_key.ip.clone()); + Some(( + connection_type, + ip, + connection_key.port, + Arc::clone(&connection_info.stream), + connection_info.health, + )) + }) + }) + }; + + let Some((connection_type, duplicate_ip, duplicate_port, duplicate_stream, recorded_health)) = + duplicate + else { + return false; + }; + + let duplicate_key = format!("{duplicate_ip}:{duplicate_port}"); + match recorded_health { + ConnectionHealth::Connected | ConnectionHealth::Busy => { + let health = + miner_stream_health(&duplicate_key, duplicate_stream.clone(), command_map).await; + if matches!(health, ConnectionHealth::Connected | ConnectionHealth::Busy) { + return false; + } + } + ConnectionHealth::Unresponsive | ConnectionHealth::Closed => {} + } + + let mut guard = CONNECTIONS.write().await; + let Some(connection) = guard.as_mut() else { + return false; + }; + let duplicate_ip_bytes = ip_to_binary(&duplicate_ip); + let still_same_stream = + connection + .connection_map + .iter() + .any(|(connection_key, connection_info)| { + connection_key.ip == duplicate_ip_bytes + && connection_key.port == duplicate_port + && ConnectionType::from_bytes(&connection_key.connection_type) + == Some(connection_type) + && Arc::ptr_eq(&connection_info.stream, &duplicate_stream) + }); + + if !still_same_stream { + return false; + } + + warn!( + "[connection_manager] removing stale duplicate miner stream before accepting reconnect: peer={duplicate_key}" + ); + connection.drop_connection(connection_type, duplicate_ip, duplicate_port); + true +} + +pub async fn write_outgoing_miner_to_memory( + received_ip_port: &str, + stream: Arc>, + wallet_short_address: String, + command_map: Arc>, +) -> bool { + let Some((ip, port)) = split_ip_port_key(received_ip_port) else { + return false; + }; + + let added = { + let mut guard = CONNECTIONS.write().await; + guard + .as_mut() + .map(|connection| { + connection.store_connection(StoreConnectionParams { + connection_type: ConnectionType::Outgoing, + ip: ip.clone(), + port, + stream: stream.clone(), + client_type: ClientType::Miner, + wallet_short_address: wallet_short_address.clone(), + command_map: command_map.clone(), + }) + }) + .unwrap_or(false) + }; + if added { + return true; + } + + if !purge_stale_duplicate_miner(&ip, command_map.clone()).await { + return false; + } + + CONNECTIONS + .write() + .await + .as_mut() + .map(|connection| { + connection.store_connection(StoreConnectionParams { + connection_type: ConnectionType::Outgoing, + ip, + port, + stream, + client_type: ClientType::Miner, + wallet_short_address, + command_map, + }) + }) + .unwrap_or(false) +} + // write an incoming connection to memory // client type defines if connection // is another node or a client connecting @@ -52,26 +249,51 @@ pub async fn write_to_memory( port, stream: stream.clone(), client_type, - wallet_short_address, - command_map, + wallet_short_address: wallet_short_address.clone(), + command_map: command_map.clone(), }); *connection_instance = Some(connection); + drop(connection_instance); + if !added + && client_type == ClientType::Miner + && purge_stale_duplicate_miner(&ip, command_map.clone()).await + { + let mut retry_connection_instance = CONNECTIONS.write().await; + if let Some(mut retry_connection) = retry_connection_instance.take() { + let retry_added = retry_connection.store_connection(StoreConnectionParams { + connection_type: ConnectionType::Incoming, + ip: ip.clone(), + port, + stream: stream.clone(), + client_type, + wallet_short_address, + command_map, + }); + *retry_connection_instance = Some(retry_connection); + if retry_added { + drop(retry_connection_instance); + return received_ip_port.to_string(); + } + } + } if added { // Return the original ip:port string as the connection key // used by stream readers and cleanup paths. - drop(connection_instance); received_ip_port.to_string() } else { // Duplicate streams are told why they are being closed before // the socket is shut down. - drop(connection_instance); let duplicate_message = "The connection is already in the connection manager Please wait 10 minutes and try again"; - let mut stream_guard = stream.lock().await; - let _ = stream_guard.write_all(duplicate_message.as_bytes()).await; - let _ = stream_guard.flush().await; - let _ = stream_guard.shutdown().await; + if let Ok(mut stream_guard) = timeout(Duration::from_secs(10), stream.lock()).await { + let _ = timeout(Duration::from_secs(10), async { + let _ = stream_guard.write_all(duplicate_message.as_bytes()).await; + let _ = stream_guard.flush().await; + stream_guard.shutdown().await + }) + .await; + } "false".to_string() } } else { @@ -79,28 +301,6 @@ pub async fn write_to_memory( } } -// delete a connection from memory -pub async fn remove_key_from_memory(key: &str) { - let mut connection_instance = CONNECTIONS.write().await; - // A key can represent either direction, so cleanup removes both - // incoming and outgoing entries with the same endpoint. - let Some((ip, port)) = split_ip_port_key(key) else { - return; - }; - if let Some(connection) = connection_instance.as_mut() { - connection.drop_connection(ConnectionType::Incoming, ip.clone(), port); - let dropped_outgoing = - connection.drop_connection(ConnectionType::Outgoing, ip.clone(), port); - if dropped_outgoing - .as_ref() - .and_then(|connection_info| ClientType::from_bytes(&connection_info.client_type)) - == Some(ClientType::Miner) - { - spawn_retry_dropped_outgoing(ip, port); - } - } -} - pub async fn remove_stream_from_memory(stream: &Arc>) { let mut connection_instance = CONNECTIONS.write().await; let Some(connection) = connection_instance.as_mut() else { diff --git a/src/rpc/server/handshake.rs b/src/rpc/server/handshake.rs index c7d6182..ddd66d9 100644 --- a/src/rpc/server/handshake.rs +++ b/src/rpc/server/handshake.rs @@ -1,15 +1,13 @@ use crate::common::check_genesis::genesis_checkup; use crate::log::{error, warn}; -use crate::miner::flag::{ - clear_mining_stop_request, request_mining_stop, set_mining_state, set_node_mode, MiningState, - NodeMode, -}; +use crate::miner::flag::{begin_chain_sync, ChainOperationGuard}; use crate::orphans::structs::OrphanCheckup2; use crate::orphans::sync_check::sync_checkup; use crate::orphans::torrent_candidates::hydrate_torrent_candidates; use crate::records::block_height::get_block_height::get_height; use crate::records::memory::connections::{ mark_peer_network_map_synced, mark_peer_operational, mark_peer_wallet_registry_synced, + spawn_peer_setup_retry, }; use crate::records::memory::response_channels::generate_uid; use crate::records::memory::response_channels::Command; @@ -19,7 +17,7 @@ use crate::rpc::client::register_wallet::register_connected_wallet; use crate::rpc::client::syncing::node_syncing; use crate::rpc::client::wallet_registry_sync::sync_wallet_registry_with_retries; use crate::rpc::responses::RpcResponse; -use crate::rpc::server::connection_memory_manager::{remove_key_from_memory, write_to_memory}; +use crate::rpc::server::connection_memory_manager::{remove_stream_from_memory, write_to_memory}; use crate::rpc::server::handshake_processing::{combine_and_send_data, parse_received_data}; use crate::rpc::server::handshake_verifications::{connection_count, perform_handshake_tests}; use crate::rpc::server::structs::{CombineAndSendDataParams, HandshakeTestParams}; @@ -42,9 +40,13 @@ use crate::Utc; async fn drop_failed_handshake(stream: &Arc>) { // Failed handshakes are never stored in connection memory, but the // accepted TCP socket should still be closed immediately. - let mut stream_guard = stream.lock().await; - let _ = stream_guard.flush().await; - let _ = stream_guard.shutdown().await; + if let Ok(mut stream_guard) = crate::timeout(Duration::from_secs(10), stream.lock()).await { + let _ = crate::timeout(Duration::from_secs(10), async { + let _ = stream_guard.flush().await; + stream_guard.shutdown().await + }) + .await; + } } async fn get_connection_counts() -> (u8, u8) { @@ -62,7 +64,29 @@ async fn sync_incoming_peer_before_operational( wallet: Arc, map: Arc>, connections_key: &str, -) -> Result { +) -> Result<(bool, Option), String> { + let initial_local_height = get_height(db); + let initial_remote_height = + request_remote_height(stream.clone(), map.clone(), connections_key.to_string()).await?; + ensure_compatible_genesis( + stream.clone(), + map.clone(), + connections_key.to_string(), + initial_remote_height, + ) + .await?; + let initial_local_genesis_exists = genesis_checkup().await; + + if initial_local_genesis_exists && initial_remote_height < initial_local_height { + warn!( + "[startup] incoming peer is behind local chain; keeping connection passive until peer catches up: local_height={initial_local_height} remote_height={initial_remote_height}" + ); + return Ok((false, None)); + } + + // Wait for any existing startup, live catch-up, or orphan operation to + // finish before this incoming peer can mutate the canonical chain. + let chain_sync_guard = begin_chain_sync().await; let local_height = get_height(db); let remote_height = request_remote_height(stream.clone(), map.clone(), connections_key.to_string()).await?; @@ -77,16 +101,12 @@ async fn sync_incoming_peer_before_operational( if local_genesis_exists && remote_height < local_height { warn!( - "[startup] incoming peer is behind local chain; keeping connection passive until peer catches up: local_height={local_height} remote_height={remote_height}" + "[startup] incoming peer fell behind while waiting for chain access; keeping connection passive: local_height={local_height} remote_height={remote_height}" ); - return Ok(false); + return Ok((false, None)); } if !local_genesis_exists || remote_height > local_height + 10 { - set_node_mode(NodeMode::Syncing); - request_mining_stop(); - set_mining_state(MiningState::Idle); - node_syncing( stream.clone(), db, @@ -145,7 +165,7 @@ async fn sync_incoming_peer_before_operational( .map_err(|err| format!("Incoming post-sync orphan check error: {err}"))?; } - Ok(true) + Ok((true, Some(chain_sync_guard))) } fn spawn_incoming_peer_promotion_watcher( @@ -175,8 +195,9 @@ fn spawn_incoming_peer_promotion_watcher( }; if remote_height >= local_height { - mark_peer_operational(&connections_key, map.clone()).await; - break; + if mark_peer_operational(&connections_key, map.clone()).await { + break; + } } } }); @@ -201,7 +222,7 @@ async fn complete_incoming_miner_setup( .await { error!("[startup] incoming peer wallet registry sync failed: {err}"); - remove_key_from_memory(connections_key).await; + remove_stream_from_memory(&stream).await; return; } mark_peer_wallet_registry_synced(connections_key).await; @@ -216,7 +237,7 @@ async fn complete_incoming_miner_setup( .await { error!("[startup] incoming peer wallet registration failed: {err}"); - remove_key_from_memory(connections_key).await; + remove_stream_from_memory(&stream).await; return; } @@ -228,7 +249,7 @@ async fn complete_incoming_miner_setup( error!( "[startup] incoming peer remote height lookup failed before network join: {err}" ); - remove_key_from_memory(connections_key).await; + remove_stream_from_memory(&stream).await; return; } }; @@ -241,7 +262,7 @@ async fn complete_incoming_miner_setup( .await { error!("[startup] incoming peer genesis compatibility check failed: {err}"); - remove_key_from_memory(connections_key).await; + remove_stream_from_memory(&stream).await; return; } let self_add_allowed = @@ -266,7 +287,7 @@ async fn complete_incoming_miner_setup( .await { error!("[startup] incoming peer network map sync failed: {err}"); - remove_key_from_memory(connections_key).await; + remove_stream_from_memory(&stream).await; return; } if let Err(err) = get_network_mapping_for_address( @@ -280,13 +301,13 @@ async fn complete_incoming_miner_setup( .await { error!("[startup] incoming peer local network record import failed: {err}"); - remove_key_from_memory(connections_key).await; + remove_stream_from_memory(&stream).await; return; } } mark_peer_network_map_synced(connections_key).await; - let operational = match sync_incoming_peer_before_operational( + let (operational, chain_sync_guard) = match sync_incoming_peer_before_operational( stream.clone(), db, wallet.clone(), @@ -298,17 +319,19 @@ async fn complete_incoming_miner_setup( Ok(operational) => operational, Err(err) => { error!("[startup] incoming peer chain sync failed: {err}"); - remove_key_from_memory(connections_key).await; + remove_stream_from_memory(&stream).await; return; } }; if operational { - mark_peer_operational(connections_key, map.clone()).await; + if !mark_peer_operational(connections_key, map.clone()).await { + spawn_peer_setup_retry(connections_key.to_string(), map.clone()); + } sleep(Duration::from_secs(15)).await; - set_node_mode(NodeMode::Normal); - clear_mining_stop_request(); - set_mining_state(MiningState::Idle); + if let Some(guard) = chain_sync_guard { + guard.finish(); + } } else { spawn_incoming_peer_promotion_watcher( stream.clone(), diff --git a/src/rpc/server/handshake_processing.rs b/src/rpc/server/handshake_processing.rs index ec5d53b..850af98 100644 --- a/src/rpc/server/handshake_processing.rs +++ b/src/rpc/server/handshake_processing.rs @@ -7,7 +7,7 @@ use crate::rpc::handshake_constants::{ HANDSHAKE_REQUEST_BYTES, HANDSHAKE_RESPONSE_BYTES, HANDSHAKE_SIGNATURE_OFFSET, HANDSHAKE_TIME_OFFSET, }; -use crate::rpc::server::connection_memory_manager::remove_key_from_memory; +use crate::rpc::server::connection_memory_manager::remove_stream_from_memory; use crate::rpc::server::rpc_command_loop::start_loop; use crate::rpc::server::structs::CombineAndSendDataParams; use crate::wallets::structures::Wallet; @@ -108,7 +108,7 @@ pub async fn return_handshake( public_key: &str, message: &str, signed_message: &str, - connections_key: &str, + _connections_key: &str, ) -> Result<(), String> { // convert to binary/bytes let public_key_bin = Wallet::normalize_saved_public_key_bytes(public_key) @@ -123,17 +123,24 @@ pub async fn return_handshake( data.extend_from_slice(&signed_bin); data.extend(&public_key_bin); - // get stream lock - let mut stream_guard = stream.lock().await; + // Handshake writes are bounded just like normal RPC writes so a peer + // that stops reading cannot hold this setup task indefinitely. + let mut stream_guard = crate::timeout(crate::Duration::from_secs(10), stream.lock()) + .await + .map_err(|_| "error: Timed out waiting for handshake stream".to_string())?; - // write to stream - if let Err(err) = stream_guard.write_all(&data).await { - let _ = remove_key_from_memory(connections_key).await; + let write_result = crate::timeout(crate::Duration::from_secs(10), async { + stream_guard.write_all(&data).await?; + stream_guard.flush().await + }) + .await; + if let Ok(Err(err)) = write_result { + remove_stream_from_memory(&stream).await; error!("Error writing to stream: {err:?}"); return Err("error: Error writing to stream".to_string()); - } else if let Err(err) = stream_guard.flush().await { - error!("Error flushing stream: {err:?}"); - return Err("error: Error flushing stream".to_string()); + } else if write_result.is_err() { + remove_stream_from_memory(&stream).await; + return Err("error: Timed out writing handshake response".to_string()); } // drop lock diff --git a/src/rpc/server/rpc_command_loop.rs b/src/rpc/server/rpc_command_loop.rs index de05f07..b79aca4 100644 --- a/src/rpc/server/rpc_command_loop.rs +++ b/src/rpc/server/rpc_command_loop.rs @@ -1,12 +1,11 @@ use crate::common::binary_conversions::binary_to_string; use crate::encode; use crate::log::warn; +use crate::records::memory::connections::mark_peer_remote_setup_complete; use crate::records::memory::enums::ClientType; use crate::records::memory::response_channels::Command; use crate::rpc::server::command_loop_state::next_incoming_command; -use crate::rpc::server::connection_memory_manager::{ - remove_key_from_memory, remove_stream_from_memory, -}; +use crate::rpc::server::connection_memory_manager::remove_stream_from_memory; use crate::rpc::*; use crate::sled::Db; use crate::wallets::structures::Wallet; @@ -558,7 +557,7 @@ pub async fn start_loop( .send(&stream_locked, Some(&connections_key), uid) .await; if should_drop_rejected_miner { - remove_key_from_memory(&connections_key).await; + remove_stream_from_memory(&stream_locked).await; break 'outer Ok(()); } } @@ -592,6 +591,32 @@ pub async fn start_loop( .send(&stream_locked, Some(&connections_key), uid) .await; } + 54 => { + // return the latest signed state for every known monitor relationship + let (uid, _) = read_bytes_from_stream::read_uid_from_stream( + &connections_key, + stream_locked.clone(), + ) + .await?; + let result = commands::network_monitor_state::network_monitor_state().await; + result + .send(&stream_locked, Some(&connections_key), uid) + .await; + } + 55 => { + // Both sides announce setup completion independently. Receipt + // records the remote half and the response acknowledges it. + let (uid, _) = read_bytes_from_stream::read_uid_from_stream( + &connections_key, + stream_locked.clone(), + ) + .await?; + mark_peer_remote_setup_complete(&connections_key, map.clone()).await; + let result = responses::RpcResponse::Binary(b"setup_complete_ack".to_vec()); + result + .send(&stream_locked, Some(&connections_key), uid) + .await; + } 30 => { // request node list let (uid, _) = read_bytes_from_stream::read_uid_from_stream( diff --git a/src/startup/connections.rs b/src/startup/connections.rs index 7bdd0ca..63135d0 100644 --- a/src/startup/connections.rs +++ b/src/startup/connections.rs @@ -4,7 +4,9 @@ use crate::miner::flag::{ clear_mining_stop_request, is_normal_mode, set_mining_state, set_node_mode, MiningState, NodeMode, }; -use crate::records::memory::connections::peer_connection_count; +use crate::records::memory::connections::{ + peer_connection_count, ready_outgoing_connection_count, refill_outgoing_connections_once, +}; use crate::records::memory::response_channels::Command; use crate::rpc::client::handshake::connect_and_handshake; use crate::rpc::client::structs::Connect; @@ -144,22 +146,32 @@ pub fn spawn_isolated_bootstrap_recovery(db: Db, wallet: Arc, map: Arc 0 { + if !is_normal_mode() { continue; } - info!("[reconnect] no operational peers remain; retrying bootstrap recovery"); - match attempt_bootstrap_connections( - db.clone(), - wallet.clone(), - map.clone(), - "reconnect", - ) - .await - { - Ok(true) => {} - Ok(false) => {} - Err(err) => warn!("[reconnect] bootstrap recovery aborted: {err}"), + if peer_connection_count().await == 0 { + info!("[reconnect] no operational peers remain; retrying bootstrap recovery"); + match attempt_bootstrap_connections( + db.clone(), + wallet.clone(), + map.clone(), + "reconnect", + ) + .await + { + Ok(true) => {} + Ok(false) => {} + Err(err) => warn!("[reconnect] bootstrap recovery aborted: {err}"), + } + continue; + } + + // A topology pass is bounded: it walks the current active map once, + // tries each eligible endpoint at most once, then sleeps until the + // next interval even when the configured limit exceeds network size. + if ready_outgoing_connection_count().await < outgoing_connections as usize { + refill_outgoing_connections_once().await; } } }); diff --git a/src/startup/network_broadcast.rs b/src/startup/network_broadcast.rs index fd63e1d..d2f0eae 100644 --- a/src/startup/network_broadcast.rs +++ b/src/startup/network_broadcast.rs @@ -2,13 +2,16 @@ use crate::common::binary_conversions::{binary_to_ip, binary_to_string, ip_to_bi use crate::common::network_startup::get_ip_and_port; use crate::log::warn; use crate::records::memory::network_mapping::structs::{ - SignedNodeEdit, NODE_ADDED_BY_OFFSET, NODE_ADDED_SIGNATURE_OFFSET, NODE_ADDED_TIMESTAMP_OFFSET, - NODE_BLOCKS_MINED_OFFSET, NODE_DELETED_BLOCK_OFFSET, NODE_DELETED_TIMESTAMP_OFFSET, - NODE_IP_OFFSET, NODE_MONITOR_COUNT_OFFSET, NODE_PORT_OFFSET, NODE_RECORD_FIXED_BYTES, + SignedNodeEdit, MONITOR_EVENT_BYTES, NODE_ADDED_BY_OFFSET, NODE_ADDED_SIGNATURE_OFFSET, + NODE_ADDED_TIMESTAMP_OFFSET, NODE_BLOCKS_MINED_OFFSET, NODE_DELETED_BLOCK_OFFSET, + NODE_DELETED_TIMESTAMP_OFFSET, NODE_IP_OFFSET, NODE_MONITOR_COUNT_OFFSET, NODE_PORT_OFFSET, + NODE_RECORD_FIXED_BYTES, }; use crate::records::memory::network_mapping::NodeInfo; use crate::records::memory::response_channels::{reserve_entry, Command}; -use crate::rpc::command_maps::{RPC_ADD_NETWORK_NODE, RPC_REQUEST_NODE_LIST}; +use crate::rpc::command_maps::{ + RPC_ADD_NETWORK_NODE, RPC_NETWORK_MONITOR_STATE, RPC_REQUEST_NODE_LIST, +}; use crate::rpc::responses::RpcResponse; use crate::sled::Db; use crate::timeout; @@ -109,17 +112,25 @@ pub async fn get_network_mapping( unlocked_stream: Arc>, command_map: Arc>, db: &Db, - _wallet: Arc, + wallet: Arc, connections_key: &str, ) -> Result<(), String> { - get_network_mapping_inner(unlocked_stream, command_map, db, connections_key, None).await + get_network_mapping_inner( + unlocked_stream, + command_map, + db, + wallet, + connections_key, + None, + ) + .await } pub async fn get_network_mapping_for_address( unlocked_stream: Arc>, command_map: Arc>, db: &Db, - _wallet: Arc, + wallet: Arc, connections_key: &str, only_address: &str, ) -> Result<(), String> { @@ -127,6 +138,7 @@ pub async fn get_network_mapping_for_address( unlocked_stream, command_map, db, + wallet, connections_key, Some(only_address), ) @@ -137,6 +149,7 @@ async fn get_network_mapping_inner( unlocked_stream: Arc>, command_map: Arc>, db: &Db, + wallet: Arc, connections_key: &str, only_address: Option<&str>, ) -> Result<(), String> { @@ -166,6 +179,9 @@ async fn get_network_mapping_inner( .try_into() .unwrap(), ) as usize; + if monitor_count != 0 { + return Err("network membership response contained unsigned monitor state".to_string()); + } let record_bytes = NODE_RECORD_FIXED_BYTES + (monitor_count * Wallet::SHORT_ADDRESS_BYTES_LENGTH); if buffer.len() < record_bytes { @@ -206,14 +222,10 @@ async fn get_network_mapping_inner( .try_into() .unwrap(), ); - let mut monitors = Vec::with_capacity(monitor_count); - for monitor_index in 0..monitor_count { - let start = - NODE_RECORD_FIXED_BYTES + monitor_index * Wallet::SHORT_ADDRESS_BYTES_LENGTH; - let end = start + Wallet::SHORT_ADDRESS_BYTES_LENGTH; - if let Some(monitor) = Wallet::bytes_to_short_address(&chunk[start..end]) { - monitors.push(monitor); - } + if deleted_timestamp != 0 || deleted_block != 0 { + return Err( + "network membership response contained unsigned deletion state".to_string(), + ); } if only_address @@ -236,23 +248,74 @@ async fn get_network_mapping_inner( modified_timestamp: added_timestamp, modified_signature: added_signature, }, - monitors, blocks_mined, ) .await { warn!("[network_map] skipped imported node record {address}: {err}"); } - - if deleted_timestamp > 0 { - NodeInfo::set_deleted_metadata_from_mapping(&address, deleted_timestamp, deleted_block) - .await; - } } if !buffer.is_empty() { return Err("network mapping response had trailing partial bytes".to_string()); } + get_signed_monitor_state( + unlocked_stream, + command_map, + db, + &wallet.saved.short_address, + connections_key, + only_address, + ) + .await?; + + Ok(()) +} + +async fn get_signed_monitor_state( + unlocked_stream: Arc>, + command_map: Arc>, + db: &Db, + local_short: &str, + connections_key: &str, + only_address: Option<&str>, +) -> Result<(), String> { + let (uid, _tx, rx) = reserve_entry(command_map).await; + let mut message = Vec::with_capacity(4); + message.push(RPC_NETWORK_MONITOR_STATE); + message.extend_from_slice(&uid); + RpcResponse::send_raw(&unlocked_stream, Some(connections_key), &message).await; + + let mut rx = rx.lock().await; + let buffer = timeout(Duration::from_secs(30), rx.recv()) + .await + .map_err(|_| "timed out waiting for signed monitor state".to_string())? + .ok_or_else(|| "signed monitor state response channel closed".to_string())?; + if buffer.len() % MONITOR_EVENT_BYTES != 0 { + return Err("signed monitor state response had trailing partial bytes".to_string()); + } + + let mut changed = false; + for chunk in buffer.chunks_exact(MONITOR_EVENT_BYTES) { + let Some(edit) = NodeInfo::monitor_event_from_bytes(chunk) else { + warn!("[network_map] skipped malformed signed monitor event"); + continue; + }; + if only_address + .map(|address| edit.monitored_address != address) + .unwrap_or(false) + { + continue; + } + match NodeInfo::import_signed_monitor_event(edit, db, local_short).await { + Ok(imported) => changed |= imported, + Err(err) => warn!("[network_map] skipped signed monitor event: {err}"), + } + } + + if changed { + NodeInfo::persist_recovery_snapshot("signed monitor state import").await; + } Ok(()) } diff --git a/src/startup/node_runtime.rs b/src/startup/node_runtime.rs index 4d6233b..f67ef48 100644 --- a/src/startup/node_runtime.rs +++ b/src/startup/node_runtime.rs @@ -10,6 +10,7 @@ use crate::miner::genesis::create_genesis_transaction; use crate::miner::mining::mine_block; use crate::panic; use crate::records::memory::chain_state::rebuild_chain_state_cache; +use crate::records::memory::connections::initialize_node_runtime_context; use crate::records::memory::mempool::{init_db, setup_mempool}; use crate::records::memory::network_mapping::NodeInfo; use crate::records::memory::response_channels::Command; @@ -143,6 +144,9 @@ pub async fn run_unlocked_node(wallet: Arc, install_shutdown: bool) -> R let wallet_for_server = wallet.clone(); let map: Arc> = Arc::new(Mutex::new(HashMap::new())); + // Install shared node dependencies before either incoming or outgoing + // handshakes can promote a connection and broadcast monitor state. + initialize_node_runtime_context(db.clone(), wallet.clone(), map.clone())?; let map_cloned = Arc::clone(&map); let db_server = db.clone(); let verification_service = Arc::new(initialize_global_verification_service()); diff --git a/src/torrent/create_metadata.rs b/src/torrent/create_metadata.rs index 3fefe03..e3fac4e 100644 --- a/src/torrent/create_metadata.rs +++ b/src/torrent/create_metadata.rs @@ -160,6 +160,7 @@ pub async fn broadcast_new_torrent_to_peers( // Send the torrent to the same live node-peer set used by transaction relay. let peers = get_live_node_broadcast_peers().await; + let mut messages = Vec::with_capacity(peers.len()); for (connections_key, stream) in peers { // Each peer gets a short-lived reply mapping so the ack can be drained. let uid_bytes = reserve_transient_entry_with_context( @@ -176,6 +177,7 @@ pub async fn broadcast_new_torrent_to_peers( message.extend_from_slice(&block_height.to_le_bytes()); // Block height message.extend_from_slice(torrent_bytes); // Torrent file contents - RpcResponse::send_raw(&stream, Some(&connections_key), &message).await; + messages.push((connections_key, stream, message)); } + RpcResponse::broadcast_raw(messages).await; }