Bootstrap fixes

This commit is contained in:
viraladmin 2026-07-14 15:47:04 -06:00
parent 9ed2dee62b
commit b999feee94
51 changed files with 1670 additions and 681 deletions

View File

@ -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

View File

@ -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;

View File

@ -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;

View File

@ -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;

View File

@ -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;

View File

@ -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;

View File

@ -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;

View File

@ -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;

View File

@ -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;

View File

@ -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;

View File

@ -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};

View File

@ -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;

View File

@ -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]

View File

@ -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;

View File

@ -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));
}
}

View File

@ -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<Wallet>) {
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 {

View File

@ -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<Wallet>) -> 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(());

View File

@ -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<Wallet>) ->
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,

View File

@ -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<Wallet>) -> Resu
}
pub async fn sync_checkup(mut params: OrphanCheckup2, wallet: Arc<Wallet>) -> Result<(), String> {
let reorg_guard = if params.node_syncing {
None
} else {
Some(begin_reorg_lock().await)
};
let result = loop {
match sync_checkup_pass(&params, wallet.clone()).await {
Ok(()) => {}
@ -315,8 +321,8 @@ pub async fn sync_checkup(mut params: OrphanCheckup2, wallet: Arc<Wallet>) -> Re
);
};
if !params.node_syncing {
end_reorg_lock();
if let Some(guard) = reorg_guard {
guard.finish();
}
if result.is_ok() {

View File

@ -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<Wallet>,
map: Arc<Mutex<Command>>,
}
lazy_static! {
static ref RECONNECT_CONTEXT: Mutex<Option<ReconnectContext>> = Mutex::new(None);
static ref RECONNECT_IN_PROGRESS: AtomicBool = AtomicBool::new(false);
static ref RECOVERY_IN_PROGRESS: StdMutex<HashSet<String>> = StdMutex::new(HashSet::new());
static ref NODE_RUNTIME_CONTEXT: StdRwLock<Option<NodeRuntimeContext>> = 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<Wallet>, map: Arc<Mutex<Command>>) {
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<RecoveryGuard> {
// 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<Wallet>,
map: Arc<Mutex<Command>>,
) -> 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<NodeRuntimeContext> {
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<Mutex<Command>>) -> bool {
fn finalize_setup_if_complete(
info: &mut ConnectionInfo,
ip: String,
port: u16,
command_map: Arc<Mutex<Command>>,
) -> 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<Mutex<Command>>,
) -> Option<(Arc<Mutex<TcpStream>>, 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<Mutex<Command>>,
) -> 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<Mutex<Command>>,
) -> 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<Mutex<Command>>) -> 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<Mutex<Command>>) -> 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<Mutex<Command>>) {
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<Mutex<TcpStream>>, 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<Mutex<TcpStream>>) -> Option<ConnectionHealth> {
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<Mutex<TcpStream>>)> {
// 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<ClientType> {
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);
}
}

View File

@ -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,

View File

@ -14,22 +14,10 @@ fn signature_is_empty(signature: &str) -> bool {
.unwrap_or(false)
}
fn merge_monitors(existing: &mut Vec<String>, monitors: Vec<String>) -> 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<String>,
_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::<Vec<_>>()
})
.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;
}
}

View File

@ -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.

View File

@ -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<HashMap<String, u64>> = 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<HashMap<String, SignedMonitorEdit>> =
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<Wallet>,
) -> 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<String, NodeInfo>,
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<String> = 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<u8> = 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(&timestamp_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<bool, String> {
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<String>) {
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<bool, String> {
Self::apply_verified_monitor_edit(&edit, db, local_short).await
}
pub(crate) async fn signed_monitor_state_bytes() -> Vec<u8> {
let mut events: Vec<SignedMonitorEdit> =
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<SignedMonitorEdit> {
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..]),
})
}
}

View File

@ -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<usize, String> {
@ -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)
}

View File

@ -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)
}

View File

@ -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,
}

View File

@ -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<Mutex<Command>>, 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);
}
}

View File

@ -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,
}
}
}

View File

@ -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<Vec<u8>> {
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;

View File

@ -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(
&params.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(())

View File

@ -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;

View File

@ -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(),

View File

@ -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;

View File

@ -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,

View File

@ -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,

View File

@ -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
}

View File

@ -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");
}

View File

@ -40,6 +40,7 @@ async fn broadcast_tx(tx_bytes: Vec<u8>, map: Arc<Mutex<Command>>) {
// 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<u8>, map: Arc<Mutex<Command>>) {
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 {

View File

@ -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(

View File

@ -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<Mutex<TcpStream>>,
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<u8, String> {
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<u16, String> {
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<u32, String> {
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<u64, String> {
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<u128, String> {
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<Vec<u8>, 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)

View File

@ -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<Mutex<TcpStream>>,
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<Mutex<TcpStream>>,
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<Mutex<TcpStream>>, Vec<u8>)>) {
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<Mutex<TcpStream>>,
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;
}
}
}

View File

@ -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

View File

@ -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<Mutex<TcpStream>>,
command_map: Arc<Mutex<Command>>,
) -> 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<Mutex<Command>>) -> 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<Mutex<TcpStream>>,
wallet_short_address: String,
command_map: Arc<Mutex<Command>>,
) -> 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<Mutex<TcpStream>>) {
let mut connection_instance = CONNECTIONS.write().await;
let Some(connection) = connection_instance.as_mut() else {

View File

@ -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<Mutex<TcpStream>>) {
// 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<Wallet>,
map: Arc<Mutex<Command>>,
connections_key: &str,
) -> Result<bool, String> {
) -> Result<(bool, Option<ChainOperationGuard>), 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(),

View File

@ -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

View File

@ -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(

View File

@ -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<Wallet>, map: Arc<M
loop {
sleep(Duration::from_secs(60)).await;
if !is_normal_mode() || peer_connection_count().await > 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;
}
}
});

View File

@ -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<Mutex<TcpStream>>,
command_map: Arc<Mutex<Command>>,
db: &Db,
_wallet: Arc<Wallet>,
wallet: Arc<Wallet>,
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<Mutex<TcpStream>>,
command_map: Arc<Mutex<Command>>,
db: &Db,
_wallet: Arc<Wallet>,
wallet: Arc<Wallet>,
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<Mutex<TcpStream>>,
command_map: Arc<Mutex<Command>>,
db: &Db,
wallet: Arc<Wallet>,
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<Mutex<TcpStream>>,
command_map: Arc<Mutex<Command>>,
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(())
}

View File

@ -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<Wallet>, install_shutdown: bool) -> R
let wallet_for_server = wallet.clone();
let map: Arc<Mutex<Command>> = 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());

View File

@ -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;
}