Compare commits

...

2 Commits

Author SHA1 Message Date
contractless d2a13cb134 revised network mapping and install process 2026-08-14 14:27:22 -06:00
contractless c40fb0ce4d fixed swap validations 2026-08-02 10:40:38 -06:00
37 changed files with 3101 additions and 929 deletions

View File

@ -1,3 +1,4 @@
#[cfg(windows)]
use contractless::common::cli_prompts::{prompt_hidden_nonempty, prompt_visible_with_default}; use contractless::common::cli_prompts::{prompt_hidden_nonempty, prompt_visible_with_default};
#[cfg(windows)] #[cfg(windows)]
use std::env; use std::env;
@ -5,7 +6,7 @@ use std::env;
use std::error::Error; use std::error::Error;
#[cfg(windows)] #[cfg(windows)]
use std::error::Error; use std::error::Error;
#[cfg(unix)] #[cfg(all(unix, any()))]
use std::fs; use std::fs;
#[cfg(windows)] #[cfg(windows)]
use std::fs; use std::fs;
@ -13,13 +14,23 @@ use std::fs;
use std::path::{Path, PathBuf}; use std::path::{Path, PathBuf};
#[cfg(windows)] #[cfg(windows)]
use std::process::Command; use std::process::Command;
#[cfg(unix)] #[cfg(all(unix, any()))]
use std::process::{Command, Stdio}; use std::process::{Command, Stdio};
#[cfg(windows)] #[cfg(windows)]
const DEFAULT_PG_VERSION: &str = "16.4-1"; const DEFAULT_PG_VERSION: &str = "16.4-1";
#[cfg(unix)] #[cfg(unix)]
#[path = "postgres_installer/unix.rs"]
mod postgres_installer_unix;
#[cfg(unix)]
#[tokio::main]
async fn main() -> Result<(), Box<dyn Error>> {
postgres_installer_unix::run().await
}
#[cfg(all(unix, any()))]
#[tokio::main] #[tokio::main]
async fn main() -> Result<(), Box<dyn Error>> { async fn main() -> Result<(), Box<dyn Error>> {
println!("\n=== Contractless Postgres Installer ===\n"); println!("\n=== Contractless Postgres Installer ===\n");
@ -126,36 +137,46 @@ async fn main() -> Result<(), Box<dyn Error>> {
// Windows installation writes into Program Files and manages a service. // Windows installation writes into Program Files and manages a service.
ensure_administrator()?; ensure_administrator()?;
// Collect install settings and the database credentials to print at the end. // A release installer can provide all values through process-local environment
let host = prompt_visible_with_default("Enter Postgres host", "127.0.0.1").await; // variables. Interactive use remains unchanged when those values are absent.
let port = prompt_visible_with_default("Enter Postgres port", "5432").await; let host = env_or_prompt("CONTRACTLESS_PG_HOST", "Enter Postgres host", "127.0.0.1").await;
let dbname = let port = env_or_prompt("CONTRACTLESS_PG_PORT", "Enter Postgres port", "5432").await;
prompt_visible_with_default("Enter new database name for blockchain", "contractless_db") let dbname = env_or_prompt(
.await; "CONTRACTLESS_PG_DATABASE",
let user = "Enter new database name for blockchain",
prompt_visible_with_default("Enter new username for blockchain database", "contractless") "contractless_db",
.await; )
let user_pass = prompt_hidden_nonempty( .await;
let user = env_or_prompt(
"CONTRACTLESS_PG_USER",
"Enter new username for blockchain database",
"contractless",
)
.await;
let user_pass = env_or_hidden_prompt(
"CONTRACTLESS_PG_USER_PASSWORD",
"Enter password for new database user: ", "Enter password for new database user: ",
"Password cannot be empty. Please try again.",
) )
.await; .await;
let postgres_pass = prompt_hidden_nonempty( let postgres_pass = env_or_hidden_prompt(
"CONTRACTLESS_PG_ADMIN_PASSWORD",
"Enter password for the postgres superuser/service account: ", "Enter password for the postgres superuser/service account: ",
"Password cannot be empty. Please try again.",
) )
.await; .await;
let install_dir = prompt_visible_with_default( let install_dir = env_or_prompt(
"CONTRACTLESS_PG_INSTALL_DIR",
"Enter PostgreSQL install directory", "Enter PostgreSQL install directory",
r"C:\Program Files\PostgreSQL\16", r"C:\Program Files\PostgreSQL\16",
) )
.await; .await;
let data_dir = prompt_visible_with_default( let data_dir = env_or_prompt(
"CONTRACTLESS_PG_DATA_DIR",
"Enter PostgreSQL data directory", "Enter PostgreSQL data directory",
r"C:\Program Files\PostgreSQL\16\data", r"C:\Program Files\PostgreSQL\16\data",
) )
.await; .await;
let service_name = prompt_visible_with_default( let service_name = env_or_prompt(
"CONTRACTLESS_PG_SERVICE_NAME",
"Enter PostgreSQL Windows service name", "Enter PostgreSQL Windows service name",
"postgresql-contractless", "postgresql-contractless",
) )
@ -214,13 +235,29 @@ async fn main() -> Result<(), Box<dyn Error>> {
Ok(()) Ok(())
} }
#[cfg(windows)]
async fn env_or_prompt(variable: &str, prompt: &str, default: &str) -> String {
match env::var(variable) {
Ok(value) if !value.trim().is_empty() => value,
_ => prompt_visible_with_default(prompt, default).await,
}
}
#[cfg(windows)]
async fn env_or_hidden_prompt(variable: &str, prompt: &str) -> String {
match env::var(variable) {
Ok(value) if !value.is_empty() => value,
_ => prompt_hidden_nonempty(prompt, "Password cannot be empty. Please try again.").await,
}
}
#[cfg(not(any(unix, windows)))] #[cfg(not(any(unix, windows)))]
fn main() { fn main() {
eprintln!("postgres_installer is only supported on Unix-like systems and Windows."); eprintln!("postgres_installer is only supported on Unix-like systems and Windows.");
std::process::exit(1); std::process::exit(1);
} }
#[cfg(unix)] #[cfg(all(unix, any()))]
fn find_pg_hba_path() -> Result<String, Box<dyn Error>> { fn find_pg_hba_path() -> Result<String, Box<dyn Error>> {
// Ask PostgreSQL for the active pg_hba.conf path instead of guessing distro paths. // Ask PostgreSQL for the active pg_hba.conf path instead of guessing distro paths.
let output = Command::new("sudo") let output = Command::new("sudo")

View File

@ -0,0 +1,388 @@
// Linux-specific PostgreSQL installation and configuration.
use contractless::common::cli_prompts::{prompt_hidden_nonempty, prompt_visible_with_default};
use std::env;
use std::error::Error;
use std::fs;
use std::path::Path;
use std::process::{Command, Output, Stdio};
use std::thread::sleep;
use std::time::Duration;
#[derive(Clone, Copy)]
enum PackageManager {
Apt,
Dnf,
Yum,
Zypper,
Pacman,
Apk,
}
impl PackageManager {
fn command(self) -> &'static str {
match self {
Self::Apt => "apt-get",
Self::Dnf => "dnf",
Self::Yum => "yum",
Self::Zypper => "zypper",
Self::Pacman => "pacman",
Self::Apk => "apk",
}
}
}
pub async fn run() -> Result<(), Box<dyn Error>> {
println!("\n=== Contractless PostgreSQL Installer (Linux) ===\n");
if !nix::unistd::Uid::effective().is_root() {
return Err("This installer requires root privileges. Run it with sudo.".into());
}
let host = env_or_prompt("CONTRACTLESS_PG_HOST", "Enter PostgreSQL host", "127.0.0.1").await;
let port = env_or_prompt("CONTRACTLESS_PG_PORT", "Enter PostgreSQL port", "5432").await;
let database = env_or_prompt(
"CONTRACTLESS_PG_DATABASE",
"Enter database name for Contractless",
"contractless_db",
)
.await;
let user = env_or_prompt(
"CONTRACTLESS_PG_USER",
"Enter database username for Contractless",
"contractless",
)
.await;
let password = env_or_hidden_prompt(
"CONTRACTLESS_PG_USER_PASSWORD",
"Enter password for the Contractless database user: ",
)
.await;
validate_identifier("database name", &database)?;
validate_identifier("database user", &user)?;
let installed_with = if command_exists("psql") {
println!("[OK] PostgreSQL client already installed.");
None
} else {
let manager = detect_package_manager()?;
println!(
"[+] PostgreSQL was not found. Installing with {}...",
manager.command()
);
install_postgres(manager)?;
Some(manager)
};
start_postgres(installed_with)?;
wait_for_postgres()?;
ensure_database_user(&user, &password)?;
ensure_database(&database, &user)?;
ensure_password_authentication(&database, &user)?;
test_login(&host, &port, &database, &user, &password)?;
println!("\n[OK] PostgreSQL is ready for Contractless.");
print_settings("Postgres", &host, &port, &user, &password, &database);
print_settings(
"Postgres-Testnet",
&host,
&port,
&user,
&password,
&database,
);
Ok(())
}
async fn env_or_prompt(variable: &str, prompt: &str, default: &str) -> String {
match env::var(variable) {
Ok(value) if !value.trim().is_empty() => value,
_ => prompt_visible_with_default(prompt, default).await,
}
}
async fn env_or_hidden_prompt(variable: &str, prompt: &str) -> String {
match env::var(variable) {
Ok(value) if !value.is_empty() => value,
_ => prompt_hidden_nonempty(prompt, "Password cannot be empty. Please try again.").await,
}
}
fn command_exists(command: &str) -> bool {
env::var_os("PATH")
.is_some_and(|paths| env::split_paths(&paths).any(|path| path.join(command).is_file()))
}
fn detect_package_manager() -> Result<PackageManager, Box<dyn Error>> {
for (command, manager) in [
("apt-get", PackageManager::Apt),
("dnf", PackageManager::Dnf),
("yum", PackageManager::Yum),
("zypper", PackageManager::Zypper),
("pacman", PackageManager::Pacman),
("apk", PackageManager::Apk),
] {
if command_exists(command) {
return Ok(manager);
}
}
Err("No supported package manager was found. Install PostgreSQL manually and rerun this installer. Supported package managers: apt-get, dnf, yum, zypper, pacman, apk.".into())
}
fn run_checked(program: &str, args: &[&str]) -> Result<(), Box<dyn Error>> {
let status = Command::new(program).args(args).status()?;
if status.success() {
Ok(())
} else {
Err(format!("Command failed: {program} {}", args.join(" ")).into())
}
}
fn install_postgres(manager: PackageManager) -> Result<(), Box<dyn Error>> {
match manager {
PackageManager::Apt => {
run_checked("apt-get", &["update"])?;
run_checked("apt-get", &["install", "-y", "postgresql"])
}
PackageManager::Dnf => run_checked("dnf", &["install", "-y", "postgresql-server"]),
PackageManager::Yum => run_checked("yum", &["install", "-y", "postgresql-server"]),
PackageManager::Zypper => run_checked(
"zypper",
&["--non-interactive", "install", "postgresql-server"],
),
PackageManager::Pacman => run_checked("pacman", &["-Sy", "--noconfirm", "postgresql"]),
PackageManager::Apk => run_checked("apk", &["add", "postgresql", "postgresql-client"]),
}
}
fn start_postgres(manager: Option<PackageManager>) -> Result<(), Box<dyn Error>> {
if matches!(manager, Some(PackageManager::Dnf | PackageManager::Yum))
&& !Path::new("/var/lib/pgsql/data/PG_VERSION").exists()
&& command_exists("postgresql-setup")
{
run_checked("postgresql-setup", &["--initdb"])?;
}
if matches!(manager, Some(PackageManager::Pacman))
&& !Path::new("/var/lib/postgres/data/PG_VERSION").exists()
{
fs::create_dir_all("/var/lib/postgres/data")?;
run_as_postgres_status("initdb", &["-D", "/var/lib/postgres/data"])?;
}
if matches!(manager, Some(PackageManager::Apk)) {
if !Path::new("/var/lib/postgresql/data/PG_VERSION").exists() {
let _ = run_checked("rc-service", &["postgresql", "setup"]);
}
run_checked("rc-service", &["postgresql", "start"])?;
let _ = run_checked("rc-update", &["add", "postgresql", "default"]);
return Ok(());
}
if command_exists("systemctl") {
return run_checked("systemctl", &["enable", "--now", "postgresql"]);
}
if command_exists("service") {
return run_checked("service", &["postgresql", "start"]);
}
Err("PostgreSQL is installed, but no supported service manager was found. Start PostgreSQL manually and rerun this installer.".into())
}
fn run_as_postgres_output(program: &str, args: &[&str]) -> Result<Output, Box<dyn Error>> {
let output = if command_exists("runuser") {
Command::new("runuser")
.args(["-u", "postgres", "--", program])
.args(args)
.stdout(Stdio::piped())
.stderr(Stdio::piped())
.output()?
} else if command_exists("sudo") {
Command::new("sudo")
.args(["-u", "postgres", program])
.args(args)
.stdout(Stdio::piped())
.stderr(Stdio::piped())
.output()?
} else if command_exists("su") {
let command = std::iter::once(program)
.chain(args.iter().copied())
.map(shell_quote)
.collect::<Vec<_>>()
.join(" ");
Command::new("su")
.args(["postgres", "-s", "/bin/sh", "-c", &command])
.stdout(Stdio::piped())
.stderr(Stdio::piped())
.output()?
} else {
return Err("No supported tool is available to run commands as the postgres account. Install runuser, sudo, or su.".into());
};
Ok(output)
}
fn run_as_postgres_status(program: &str, args: &[&str]) -> Result<(), Box<dyn Error>> {
let output = run_as_postgres_output(program, args)?;
if output.status.success() {
Ok(())
} else {
Err(format!(
"PostgreSQL command failed: {}",
String::from_utf8_lossy(&output.stderr).trim()
)
.into())
}
}
fn postgres_sql(sql: &str) -> Result<Output, Box<dyn Error>> {
run_as_postgres_output(
"psql",
&[
"-d",
"postgres",
"-v",
"ON_ERROR_STOP=1",
"-t",
"-A",
"-c",
sql,
],
)
}
fn wait_for_postgres() -> Result<(), Box<dyn Error>> {
for _ in 0..30 {
if postgres_sql("SELECT 1;").is_ok_and(|output| output.status.success()) {
return Ok(());
}
sleep(Duration::from_secs(1));
}
Err("PostgreSQL did not become ready within 30 seconds.".into())
}
fn ensure_database_user(user: &str, password: &str) -> Result<(), Box<dyn Error>> {
let escaped_password = password.replace('\'', "''");
let sql = format!(
"DO $$ BEGIN IF NOT EXISTS (SELECT 1 FROM pg_roles WHERE rolname = '{user}') THEN CREATE ROLE {user} LOGIN PASSWORD '{escaped_password}'; ELSE ALTER ROLE {user} WITH LOGIN PASSWORD '{escaped_password}'; END IF; END $$;"
);
require_success(postgres_sql(&sql)?, "Failed to create the database user")
}
fn ensure_database(database: &str, user: &str) -> Result<(), Box<dyn Error>> {
let check = postgres_sql(&format!(
"SELECT 1 FROM pg_database WHERE datname = '{database}';"
))?;
if !check.status.success() {
return Err(format!(
"Failed to check the database: {}",
String::from_utf8_lossy(&check.stderr).trim()
)
.into());
}
if String::from_utf8_lossy(&check.stdout).trim() == "1" {
return Ok(());
}
run_as_postgres_status("createdb", &["--owner", user, database])
}
fn ensure_password_authentication(database: &str, user: &str) -> Result<(), Box<dyn Error>> {
let output = postgres_sql("SHOW hba_file;")?;
if !output.status.success() {
return Err(format!(
"Could not determine pg_hba.conf path: {}",
String::from_utf8_lossy(&output.stderr).trim()
)
.into());
}
let path = String::from_utf8_lossy(&output.stdout).trim().to_string();
if path.is_empty() {
return Err("Could not determine pg_hba.conf path.".into());
}
let rule = format!("host {database} {user} 127.0.0.1/32 md5");
let contents = fs::read_to_string(&path)?;
if !contents.lines().any(|line| line.trim() == rule) {
fs::write(&path, format!("{rule}\n{contents}"))?;
require_success(
postgres_sql("SELECT pg_reload_conf();")?,
"Failed to reload PostgreSQL configuration",
)?;
}
Ok(())
}
fn test_login(
host: &str,
port: &str,
database: &str,
user: &str,
password: &str,
) -> Result<(), Box<dyn Error>> {
let output = Command::new("psql")
.env("PGPASSWORD", password)
.args([
"-h",
host,
"-p",
port,
"-U",
user,
"-d",
database,
"-v",
"ON_ERROR_STOP=1",
"-t",
"-A",
"-c",
"SELECT 1;",
])
.output()?;
require_success(output, "Failed to verify the Contractless database login")
}
fn require_success(output: Output, context: &str) -> Result<(), Box<dyn Error>> {
if output.status.success() {
Ok(())
} else {
Err(format!(
"{context}: {}",
String::from_utf8_lossy(&output.stderr).trim()
)
.into())
}
}
fn shell_quote(value: &str) -> String {
format!("'{}'", value.replace('\'', "'\"'\"'"))
}
fn validate_identifier(label: &str, value: &str) -> Result<(), Box<dyn Error>> {
if value.is_empty()
|| !value
.chars()
.all(|character| character.is_ascii_alphanumeric() || character == '_')
{
return Err(format!(
"Invalid {label}. Only letters, numbers, and underscores are allowed."
)
.into());
}
Ok(())
}
fn print_settings(
section: &str,
host: &str,
port: &str,
user: &str,
password: &str,
database: &str,
) {
println!("\n[{section}]");
println!("host = {host}");
println!("port = {port}");
println!("user = {user}");
println!("password = {password}");
println!("dbname = {database}");
}

View File

@ -1,4 +1,5 @@
use crate::common::types::Transaction; use crate::common::types::Transaction;
use crate::log::error;
use crate::orphans::get_path_names::get_file_names; use crate::orphans::get_path_names::get_file_names;
use crate::orphans::save_blocks::save_new_blocks; use crate::orphans::save_blocks::save_new_blocks;
use crate::orphans::structs::UndoTransactions; use crate::orphans::structs::UndoTransactions;
@ -228,6 +229,9 @@ pub async fn undo_transactions(
let final_height = true_start_height.saturating_sub(1); let final_height = true_start_height.saturating_sub(1);
crate::orphans::undo_block::finalize_undo_height(final_height, &params.db).await; crate::orphans::undo_block::finalize_undo_height(final_height, &params.db).await;
if let Err(err) = NodeInfo::persist_mined_counts(&params.db).await {
error!("Failed to checkpoint mined counts after rollback: {err}");
}
rebuild_chain_state_cache(&params.db).await?; rebuild_chain_state_cache(&params.db).await?;
// Only now that every rolled-back block has been unwound do we test // Only now that every rolled-back block has been unwound do we test

View File

@ -165,7 +165,6 @@ async fn reconnect_replacement_inner(excluded_ip: &str) {
wallet: context.wallet, wallet: context.wallet,
db: context.db, db: context.db,
map: context.map, map: context.map,
first: false,
run_startup_sync: false, run_startup_sync: false,
}; };
@ -218,7 +217,10 @@ async fn retry_dropped_outgoing(ip: String, port: u16) {
wallet: context.wallet.clone(), wallet: context.wallet.clone(),
db: context.db.clone(), db: context.db.clone(),
map: context.map.clone(), map: context.map.clone(),
first: peer_connection_count().await == 0, // A miner stream that is still completing setup already has a
// startup owner. Counting only ready peers could launch a second
// full startup synchronization while the first one is active.
first: miner_connection_count().await == 0,
}; };
match connect_and_handshake(connect).await { match connect_and_handshake(connect).await {
@ -593,6 +595,7 @@ impl Connection {
} }
info.ready = true; info.ready = true;
info.catch_up_target = None;
spawn_monitor_update( spawn_monitor_update(
ip.clone(), ip.clone(),
MONITOR_ACTION_ADD, MONITOR_ACTION_ADD,
@ -679,6 +682,16 @@ impl Connection {
.count() .count()
} }
pub fn ready_miner_wallets(&self) -> Vec<String> {
self.connection_map
.values()
.filter(|info| {
ClientType::from_bytes(&info.client_type) == Some(ClientType::Miner) && info.ready
})
.map(|info| info.wallet_short_address.clone())
.collect()
}
pub fn count_miner_connections(&self) -> usize { pub fn count_miner_connections(&self) -> usize {
self.connection_map self.connection_map
.values() .values()
@ -728,6 +741,24 @@ impl Connection {
.collect() .collect()
} }
pub fn get_mapping_relay_peer_streams_with_keys(
&self,
) -> Vec<(String, Arc<Mutex<TcpStream>>)> {
self.connection_map
.iter()
.filter_map(|(key, connection_info)| {
if ClientType::from_bytes(&connection_info.client_type) != Some(ClientType::Miner)
|| (!connection_info.ready && !connection_info.wallet_registry_synced)
{
return None;
}
let ip = binary_to_ip(key.ip.clone());
let connections_key = format!("{}:{}", ip, key.port);
Some((connections_key, Arc::clone(&connection_info.stream)))
})
.collect()
}
pub fn get_startup_synced_peer_streams_with_keys( pub fn get_startup_synced_peer_streams_with_keys(
&self, &self,
) -> Vec<(String, Arc<Mutex<TcpStream>>)> { ) -> Vec<(String, Arc<Mutex<TcpStream>>)> {
@ -961,11 +992,34 @@ fn spawn_monitor_update(ip: String, action: u8, monitored_address: String, port:
wallet: context.wallet.clone(), wallet: context.wallet.clone(),
connections_key: format!("{ip}:{port}"), connections_key: format!("{ip}:{port}"),
}; };
let _ = if action == MONITOR_ACTION_ADD { if action == MONITOR_ACTION_ADD {
NodeInfo::add_monitor(params).await for attempt in 1..=40 {
let RpcResponse::Binary(response) = NodeInfo::add_monitor(params.clone()).await;
if response.as_slice() == b"Success" {
return;
}
if attempt == 40 {
warn!(
"[network_map] monitor-add failed after {attempt} attempts: monitored={} monitoring={} error={}",
monitored_address,
params.edit.monitoring_address,
String::from_utf8_lossy(&response)
);
return;
}
sleep(Duration::from_millis(250)).await;
}
} else { } else {
NodeInfo::remove_monitor(params).await let RpcResponse::Binary(response) = NodeInfo::remove_monitor(params.clone()).await;
}; if response.as_slice() != b"Success" {
warn!(
"[network_map] monitor-remove failed: monitored={} monitoring={} error={}",
monitored_address,
params.edit.monitoring_address,
String::from_utf8_lossy(&response)
);
}
}
}); });
} }
@ -982,6 +1036,35 @@ pub async fn initialize_connection() {
} }
} }
pub async fn drop_miner_connections_for_wallet(wallet_short_address: &str) {
let removed = {
let mut connections = CONNECTIONS.write().await;
let Some(connection) = connections.as_mut() else {
return;
};
let keys: Vec<ConnectionKey> = connection
.connection_map
.iter()
.filter_map(|(key, info)| {
(ClientType::from_bytes(&info.client_type) == Some(ClientType::Miner)
&& info.wallet_short_address == wallet_short_address)
.then_some(key.clone())
})
.collect();
keys.into_iter()
.filter_map(|key| connection.connection_map.remove(&key))
.collect::<Vec<_>>()
};
for info in removed {
let stream = Arc::clone(&info.stream);
tokio::spawn(async move {
let mut stream = stream.lock().await;
let _ = stream.shutdown().await;
});
}
}
pub async fn outgoing_connection_count() -> usize { pub async fn outgoing_connection_count() -> usize {
// Read the singleton connection manager and count live outgoing peers. // Read the singleton connection manager and count live outgoing peers.
CONNECTIONS CONNECTIONS
@ -1013,6 +1096,15 @@ pub async fn peer_connection_count() -> usize {
.unwrap_or(0) .unwrap_or(0)
} }
pub async fn operational_peer_wallets() -> Vec<String> {
CONNECTIONS
.read()
.await
.as_ref()
.map(|connection| connection.ready_miner_wallets())
.unwrap_or_default()
}
pub async fn miner_connection_count() -> usize { pub async fn miner_connection_count() -> usize {
// Recovery logic uses raw miner sockets to avoid starting another // Recovery logic uses raw miner sockets to avoid starting another
// bootstrap attempt while an accepted peer is still finishing setup. // bootstrap attempt while an accepted peer is still finishing setup.
@ -1042,6 +1134,55 @@ pub async fn mark_peer_network_map_synced(key: &str) -> bool {
.unwrap_or(false) .unwrap_or(false)
} }
pub async fn mark_peer_catching_up(key: &str, target_height: u32) -> bool {
let Some((ip, port)) = split_ip_port_key(key) else {
return false;
};
let ip_bytes = ip_to_binary(&ip);
CONNECTIONS
.write()
.await
.as_mut()
.map(|connection| {
for (connection_key, info) in connection.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
|| info.ready
{
return false;
}
info.catch_up_target = Some(target_height);
return true;
}
}
false
})
.unwrap_or(false)
}
pub async fn peer_catch_up_target(key: &str) -> Option<u32> {
let (ip, port) = split_ip_port_key(key)?;
let ip_bytes = ip_to_binary(&ip);
CONNECTIONS.read().await.as_ref().and_then(|connection| {
connection
.connection_map
.iter()
.find(|(connection_key, info)| {
connection_key.ip == ip_bytes
&& connection_key.port == port
&& ClientType::from_bytes(&info.client_type) == Some(ClientType::Miner)
&& !info.ready
&& info.wallet_registry_synced
&& info.network_map_synced
})
.and_then(|(_, info)| info.catch_up_target)
})
}
pub async fn mark_peer_operational(key: &str, map: Arc<Mutex<Command>>) -> bool { pub async fn mark_peer_operational(key: &str, map: Arc<Mutex<Command>>) -> bool {
if peer_is_operational(key).await { if peer_is_operational(key).await {
return true; return true;
@ -1185,6 +1326,34 @@ pub async fn peer_accepts_live_relay(key: &str) -> bool {
.unwrap_or(false) .unwrap_or(false)
} }
pub async fn peer_accepts_mapping_relay(key: &str) -> bool {
let Some((ip, port)) = split_ip_port_key(key) else {
return false;
};
let ip_bytes = ip_to_binary(&ip);
CONNECTIONS
.read()
.await
.as_ref()
.and_then(|connection| {
connection
.connection_map
.iter()
.find_map(|(connection_key, info)| {
if connection_key.ip == ip_bytes
&& connection_key.port == port
&& ClientType::from_bytes(&info.client_type) == Some(ClientType::Miner)
{
Some(info.wallet_registry_synced)
} else {
None
}
})
})
.unwrap_or(false)
}
pub async fn set_stream_health(stream: &Arc<Mutex<TcpStream>>, health: ConnectionHealth) { pub async fn set_stream_health(stream: &Arc<Mutex<TcpStream>>, health: ConnectionHealth) {
let mut guard = CONNECTIONS.write().await; let mut guard = CONNECTIONS.write().await;
let Some(connection) = guard.as_mut() else { let Some(connection) = guard.as_mut() else {

File diff suppressed because it is too large Load Diff

View File

@ -1,8 +0,0 @@
use super::*;
impl NodeInfo {
pub async fn delete_address(params: DeleteAddressParams) -> RpcResponse {
let _ = params;
RpcResponse::Binary(b"Error: explicit node deletion is disabled".to_vec())
}
}

View File

@ -1,6 +1,10 @@
use super::*; use super::*;
use crate::common::check_genesis::genesis_checkup; use crate::common::check_genesis::genesis_checkup;
use crate::records::unpack_block::unpack_header::load_block_header; use crate::records::unpack_block::unpack_header::load_block_header;
use crate::sled::Batch;
const MINED_COUNTS_TREE: &str = "network_mapping_mined_counts";
const MINED_COUNTS_HEIGHT_KEY: &[u8] = b"\0chain_height";
fn increment_count(counts: &mut HashMap<String, u8>, address: &str) { fn increment_count(counts: &mut HashMap<String, u8>, address: &str) {
let entry = counts.entry(address.to_string()).or_insert(0); let entry = counts.entry(address.to_string()).or_insert(0);
@ -9,17 +13,68 @@ fn increment_count(counts: &mut HashMap<String, u8>, address: &str) {
} }
} }
fn read_cached_mined_counts(db: &Db) -> Result<Option<HashMap<String, u8>>, String> {
let tree = db
.open_tree(MINED_COUNTS_TREE)
.map_err(|err| format!("Failed to open mined-count cache: {err}"))?;
let Some(stored_height) = tree
.get(MINED_COUNTS_HEIGHT_KEY)
.map_err(|err| format!("Failed to read mined-count cache height: {err}"))?
else {
return Ok(None);
};
let stored_height: [u8; 4] = stored_height
.as_ref()
.try_into()
.map_err(|_| "Invalid mined-count cache height".to_string())?;
if u32::from_le_bytes(stored_height) != get_height(db) {
return Ok(None);
}
let mut restored_counts: HashMap<String, u8> = HashMap::new();
for entry in tree.iter() {
let (key, value) =
entry.map_err(|err| format!("Failed to read mined-count cache: {err}"))?;
if key.as_ref() == MINED_COUNTS_HEIGHT_KEY {
continue;
}
if value.len() != 1 {
return Err("Invalid mined-count cache entry".to_string());
}
let address = String::from_utf8(key.to_vec())
.map_err(|_| "Invalid mined-count cache address".to_string())?;
restored_counts.insert(address, value[0]);
}
Ok(Some(restored_counts))
}
impl NodeInfo { impl NodeInfo {
pub async fn increment_mined(address: &str) { pub async fn increment_mined(db: &Db, address: &str) -> Result<(), String> {
{ let count = {
let mut map = ADDRESS_MAP.lock().await; let mut map = ADDRESS_MAP.lock().await;
if let Some(node_info) = map.get_mut(address) { if let Some(node_info) = map.get_mut(address) {
// Counts are capped at u8-safe policy maximum used by node rules. // Counts are capped at u8-safe policy maximum used by node rules.
if node_info.blocks_mined < 250 { if node_info.blocks_mined < 250 {
node_info.blocks_mined += 1; node_info.blocks_mined += 1;
} }
Some(node_info.blocks_mined)
} else {
None
} }
};
if let Some(count) = count {
let mut batch = Batch::default();
batch.insert(address.as_bytes(), &[count]);
batch.insert(MINED_COUNTS_HEIGHT_KEY, &get_height(db).to_le_bytes());
db.open_tree(MINED_COUNTS_TREE)
.map_err(|err| format!("Failed to open mined-count cache: {err}"))?
.apply_batch(batch)
.map_err(|err| format!("Failed to update mined-count cache: {err}"))?;
} }
Ok(())
} }
pub async fn decrement_mined(address: &str) { pub async fn decrement_mined(address: &str) {
@ -43,6 +98,45 @@ impl NodeInfo {
} }
} }
pub async fn restore_mined_counts(db: &Db) -> Result<bool, String> {
let Some(restored_counts) = read_cached_mined_counts(db)? else {
return Ok(false);
};
let mut map = ADDRESS_MAP.lock().await;
for (address, node_info) in map.iter_mut() {
node_info.blocks_mined = restored_counts.get(address).copied().unwrap_or(0);
}
Ok(true)
}
pub async fn restore_or_rebuild_mined_counts(db: &Db) -> Result<(), String> {
match Self::restore_mined_counts(db).await {
Ok(true) => Ok(()),
Ok(false) => Self::rebuild_mined_counts_from_chain(db).await,
Err(_) => Self::rebuild_mined_counts_from_chain(db).await,
}
}
pub async fn persist_mined_counts(db: &Db) -> Result<(), String> {
let counts: Vec<(String, u8)> = {
let map = ADDRESS_MAP.lock().await;
map.iter()
.map(|(address, node_info)| (address.clone(), node_info.blocks_mined))
.collect()
};
let mut batch = Batch::default();
for (address, count) in counts {
batch.insert(address.as_bytes(), &[count]);
}
batch.insert(MINED_COUNTS_HEIGHT_KEY, &get_height(db).to_le_bytes());
db.open_tree(MINED_COUNTS_TREE)
.map_err(|err| format!("Failed to open mined-count cache: {err}"))?
.apply_batch(batch)
.map_err(|err| format!("Failed to persist mined-count cache: {err}"))
}
pub async fn rebuild_mined_counts_from_chain(db: &Db) -> Result<(), String> { pub async fn rebuild_mined_counts_from_chain(db: &Db) -> Result<(), String> {
// Recompute node mined counts directly from saved block headers // Recompute node mined counts directly from saved block headers
// so startup and recovery can rebuild memory-only state. // so startup and recovery can rebuild memory-only state.
@ -55,7 +149,7 @@ impl NodeInfo {
node_info.blocks_mined = 0; node_info.blocks_mined = 0;
} }
drop(map); drop(map);
Self::persist_recovery_snapshot("mined rebuild without genesis").await; Self::persist_mined_counts(db).await?;
return Ok(()); return Ok(());
} }
@ -80,7 +174,38 @@ impl NodeInfo {
} }
} }
} }
Self::persist_recovery_snapshot("mined rebuild").await; Self::persist_mined_counts(db).await?;
Ok(()) Ok(())
} }
} }
#[cfg(test)]
mod tests {
use super::*;
fn cache_at_height(db: &Db, chain_height: u32, cache_height: u32) {
db.insert(b"height", &chain_height.to_le_bytes()).unwrap();
let tree = db.open_tree(MINED_COUNTS_TREE).unwrap();
let mut batch = Batch::default();
batch.insert(MINED_COUNTS_HEIGHT_KEY, &cache_height.to_le_bytes());
batch.insert(b"miner.cltc", &[100]);
tree.apply_batch(batch).unwrap();
}
#[test]
fn accepts_cache_at_the_active_chain_height() {
let db = crate::sled::Config::new().temporary(true).open().unwrap();
cache_at_height(&db, 42, 42);
let counts = read_cached_mined_counts(&db).unwrap().unwrap();
assert_eq!(counts.get("miner.cltc"), Some(&100));
}
#[test]
fn rejects_cache_from_a_different_chain_height() {
let db = crate::sled::Config::new().temporary(true).open().unwrap();
cache_at_height(&db, 43, 42);
assert!(read_cached_mined_counts(&db).unwrap().is_none());
}
}

View File

@ -6,8 +6,6 @@ use crate::decode;
use crate::lazy_static; use crate::lazy_static;
use crate::log::{info, warn}; use crate::log::{info, warn};
use crate::records::block_height::get_block_height::get_height; use crate::records::block_height::get_block_height::get_height;
use crate::records::ip_score::enums::InfractionType;
use crate::records::ip_score::score::update_ip_score;
use crate::records::memory::connections::CONNECTIONS; use crate::records::memory::connections::CONNECTIONS;
use crate::records::memory::network_mapping::enums::NodeEditType; use crate::records::memory::network_mapping::enums::NodeEditType;
use crate::records::memory::network_mapping::structs::{ use crate::records::memory::network_mapping::structs::{
@ -23,9 +21,27 @@ use crate::HashMap;
use crate::Mutex; use crate::Mutex;
use crate::OnceLock; use crate::OnceLock;
use crate::Utc; use crate::Utc;
use std::collections::VecDeque;
pub(crate) enum DeferredMappingUpdate {
Add(Db, SignedNodeEdit),
Monitor(MonitorAddressParams, u8),
SnapshotAddress(Db, structs::SyncedNodeState),
ReconciledMembership(Db, structs::SyncedNodeState),
}
struct MappingSnapshotState {
active: bool,
updates: VecDeque<DeferredMappingUpdate>,
}
lazy_static! { lazy_static! {
static ref ADDRESS_MAP: Mutex<HashMap<String, NodeInfo>> = Mutex::new(HashMap::new()); static ref ADDRESS_MAP: Mutex<HashMap<String, NodeInfo>> = Mutex::new(HashMap::new());
static ref MAPPING_SNAPSHOT_STATE: Mutex<MappingSnapshotState> =
Mutex::new(MappingSnapshotState {
active: false,
updates: VecDeque::new(),
});
} }
pub const SELF_ADD_BLOCK: u32 = 0; pub const SELF_ADD_BLOCK: u32 = 0;
@ -93,6 +109,5 @@ mod add;
pub mod enums; pub mod enums;
mod mined_counts; mod mined_counts;
pub(crate) mod monitor; pub(crate) mod monitor;
pub(crate) mod persistence;
mod queries; mod queries;
pub mod structs; pub mod structs;

View File

@ -14,7 +14,7 @@ pub const MONITOR_ACTION_REMOVE: u8 = 2;
lazy_static! { lazy_static! {
// Keep the newest signed assertion for each monitor relationship. This is // Keep the newest signed assertion for each monitor relationship. This is
// both replay protection and the authenticated state shared with joining peers. // both replay protection and the authenticated state shared with joining peers.
static ref MONITOR_EVENT_STATE: Mutex<HashMap<String, SignedMonitorEdit>> = pub(super) static ref MONITOR_EVENT_STATE: Mutex<HashMap<String, SignedMonitorEdit>> =
Mutex::new(HashMap::new()); Mutex::new(HashMap::new());
} }
@ -108,6 +108,39 @@ mod tests {
crate::encode(vec![7_u8; Wallet::SIGNATURE_LENGTH]) crate::encode(vec![7_u8; Wallet::SIGNATURE_LENGTH])
); );
} }
#[test]
fn ip_change_discards_only_old_target_ip_events() {
let mut state = HashMap::new();
let target = "target.cltc";
let monitor = "monitor.cltc";
let other = "other.cltc";
for (key, monitored_address, target_ip) in [
("old", target, "1.2.3.4"),
("current", target, "5.6.7.8"),
("other", other, "1.2.3.4"),
] {
state.insert(
key.to_string(),
SignedMonitorEdit {
action: MONITOR_ACTION_ADD,
monitored_address: monitored_address.to_string(),
monitoring_address: monitor.to_string(),
target_ip: target_ip.to_string(),
modified_timestamp: 100,
modified_block: 50,
modified_signature: "signature".to_string(),
},
);
}
NodeInfo::discard_old_target_ip_events_from(&mut state, target, "5.6.7.8");
assert!(!state.contains_key("old"));
assert!(state.contains_key("current"));
assert!(state.contains_key("other"));
}
} }
impl NodeInfo { impl NodeInfo {
@ -226,31 +259,14 @@ impl NodeInfo {
) )
} }
async fn remember_newest_monitor_event(edit: &SignedMonitorEdit) -> bool { pub(super) fn discard_old_target_ip_events_from(
let key = Self::monitor_relation_key(edit); state: &mut HashMap<String, SignedMonitorEdit>,
let mut state = MONITOR_EVENT_STATE.lock().await; address: &str,
if state current_ip: &str,
.get(&key) ) {
.map(|existing| Self::monitor_event_order(edit) <= Self::monitor_event_order(existing)) state.retain(|_, event| {
.unwrap_or(false) event.monitored_address != address || event.target_ip == current_ip
{ });
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( async fn broadcast_monitor_event(
@ -282,7 +298,7 @@ impl NodeInfo {
let connections_lock = CONNECTIONS.read().await; let connections_lock = CONNECTIONS.read().await;
connections_lock connections_lock
.as_ref() .as_ref()
.map(|connection| connection.get_all_ready_peer_streams_with_keys()) .map(|connection| connection.get_mapping_relay_peer_streams_with_keys())
.unwrap_or_default() .unwrap_or_default()
}; };
@ -327,28 +343,55 @@ impl NodeInfo {
return Err("Could not validate monitor signature".to_string()); return Err("Could not validate monitor signature".to_string());
} }
// Validate membership before recording the event so an invalid event // Mapping snapshots lock these structures in this same order. Holding
// cannot poison replay protection for a later valid assertion. // both here makes the signed event and its derived map state atomic.
let mut address_map = ADDRESS_MAP.lock().await;
let mut event_state = MONITOR_EVENT_STATE.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 monitored.deleted_timestamp > 0
&& monitored.added_timestamp <= monitored.deleted_timestamp
{ {
let address_map = ADDRESS_MAP.lock().await; return Err(
let monitored = address_map "monitor target must have a newer membership record before reconnection"
.get(&edit.monitored_address) .to_string(),
.ok_or_else(|| "monitored address not found".to_string())?; );
if monitored.ip != edit.target_ip { }
return Err("monitor target IP mismatch".to_string()); let monitoring = address_map
} .get(&edit.monitoring_address)
.ok_or_else(|| "monitoring address not found".to_string())?;
if monitoring.deleted_timestamp > 0 {
return Err("deleted node cannot update monitor relationships".to_string());
} }
if !Self::remember_newest_monitor_event(edit).await { let relation_key = Self::monitor_relation_key(edit);
if event_state
.get(&relation_key)
.map(|existing| Self::monitor_event_order(edit) <= Self::monitor_event_order(existing))
.unwrap_or(false)
{
return Ok(false); return Ok(false);
} }
event_state.insert(relation_key, edit.clone());
let deletion_marker = if edit.action == MONITOR_ACTION_REMOVE { let deletion_marker = if edit.action == MONITOR_ACTION_REMOVE {
Self::deletion_marker_for(&edit.monitored_address).await event_state
.values()
.filter(|event| {
event.monitored_address == edit.monitored_address
&& event.action == MONITOR_ACTION_REMOVE
})
.map(|event| (event.modified_timestamp, event.modified_block))
.max()
} else { } else {
None None
}; };
let mut address_map = ADDRESS_MAP.lock().await;
let monitored = address_map let monitored = address_map
.get_mut(&edit.monitored_address) .get_mut(&edit.monitored_address)
.ok_or_else(|| "monitored address not found".to_string())?; .ok_or_else(|| "monitored address not found".to_string())?;
@ -384,6 +427,14 @@ impl NodeInfo {
} }
async fn apply_monitor(params: MonitorAddressParams, action: u8) -> RpcResponse { async fn apply_monitor(params: MonitorAddressParams, action: u8) -> RpcResponse {
if Self::defer_mapping_update(DeferredMappingUpdate::Monitor(params.clone(), action)).await
{
return RpcResponse::Binary(b"Success".to_vec());
}
Self::apply_monitor_now(params, action).await
}
pub(crate) async fn apply_monitor_now(params: MonitorAddressParams, action: u8) -> RpcResponse {
let MonitorAddressParams { let MonitorAddressParams {
mut edit, mut edit,
remote_ip, remote_ip,
@ -419,7 +470,6 @@ impl NodeInfo {
Ok(true) => {} Ok(true) => {}
} }
Self::persist_recovery_snapshot("monitor update").await;
Self::broadcast_monitor_event(map, &edit, &remote_ip).await; Self::broadcast_monitor_event(map, &edit, &remote_ip).await;
RpcResponse::Binary(b"Success".to_vec()) RpcResponse::Binary(b"Success".to_vec())
} }
@ -432,9 +482,10 @@ impl NodeInfo {
Self::apply_verified_monitor_edit(&edit, db, local_short).await Self::apply_verified_monitor_edit(&edit, db, local_short).await
} }
pub(crate) async fn signed_monitor_state_bytes() -> Vec<u8> { pub(super) fn signed_monitor_state_bytes_from(
let mut events: Vec<SignedMonitorEdit> = state: &HashMap<String, SignedMonitorEdit>,
MONITOR_EVENT_STATE.lock().await.values().cloned().collect(); ) -> Vec<u8> {
let mut events: Vec<SignedMonitorEdit> = state.values().cloned().collect();
events.sort_by(|left, right| { events.sort_by(|left, right| {
Self::monitor_event_order(left) Self::monitor_event_order(left)
.cmp(&Self::monitor_event_order(right)) .cmp(&Self::monitor_event_order(right))
@ -468,28 +519,51 @@ impl NodeInfo {
data data
} }
pub(crate) async fn signed_monitor_state_bytes() -> Vec<u8> {
let state = MONITOR_EVENT_STATE.lock().await;
Self::signed_monitor_state_bytes_from(&state)
}
pub async fn signed_monitor_state() -> RpcResponse { pub async fn signed_monitor_state() -> RpcResponse {
RpcResponse::Binary(Self::signed_monitor_state_bytes().await) RpcResponse::Binary(Self::signed_monitor_state_bytes().await)
} }
pub(crate) async fn load_monitor_recovery_state(bytes: &[u8]) -> usize { pub(crate) async fn validate_synced_monitor_state(
let mut state = MONITOR_EVENT_STATE.lock().await; bytes: &[u8],
state.clear(); db: &Db,
let mut loaded = 0; expected_targets: &HashMap<String, String>,
for chunk in bytes.chunks_exact(MONITOR_EVENT_BYTES) { ) -> Result<HashMap<String, SignedMonitorEdit>, String> {
let Some(edit) = Self::monitor_event_from_bytes(chunk) else { if bytes.len() % MONITOR_EVENT_BYTES != 0 {
continue; return Err("monitor-state snapshot ended mid-record".to_string());
};
// 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
let mut replacement = HashMap::new();
for chunk in bytes.chunks_exact(MONITOR_EVENT_BYTES) {
let edit = Self::monitor_event_from_bytes(chunk)
.ok_or_else(|| "monitor-state snapshot contained an invalid event".to_string())?;
if !matches!(edit.action, MONITOR_ACTION_ADD | MONITOR_ACTION_REMOVE) {
return Err("monitor-state snapshot contained an invalid action".to_string());
}
if !Self::verify_monitor_edit(&edit, db).await {
return Err("monitor-state snapshot contained an invalid signature".to_string());
}
let target_ip = expected_targets
.get(&edit.monitored_address)
.ok_or_else(|| {
"monitor-state snapshot referenced an unknown monitored address".to_string()
})?;
if target_ip != &edit.target_ip {
return Err("monitor-state snapshot target IP did not match mapping".to_string());
}
let key = Self::monitor_relation_key(&edit);
if replacement.insert(key, edit).is_some() {
return Err("monitor-state snapshot contained a duplicate relationship".to_string());
}
}
Ok(replacement)
} }
pub fn monitor_event_from_bytes(bytes: &[u8]) -> Option<SignedMonitorEdit> { pub fn monitor_event_from_bytes(bytes: &[u8]) -> Option<SignedMonitorEdit> {

View File

@ -3,10 +3,174 @@ use crate::common::governance::{
activation_vote_record_key, countable_governance_addresses, proposal_vote_record_key, activation_vote_record_key, countable_governance_addresses, proposal_vote_record_key,
GovernanceNodeSnapshot, ACTIVATION_VOTES_TREE, PROPOSAL_VOTES_TREE, GovernanceNodeSnapshot, ACTIVATION_VOTES_TREE, PROPOSAL_VOTES_TREE,
}; };
use crate::common::skein::skein_256_hash_bytes;
use crate::records::memory::network_mapping::monitor::MONITOR_EVENT_STATE;
use crate::records::memory::network_mapping::structs::{
SyncedNodeState, NETWORK_SNAPSHOT_HEADER_BYTES, NETWORK_SNAPSHOT_MAGIC,
NETWORK_SNAPSHOT_VERSION,
};
use crate::sled::Db; use crate::sled::Db;
use std::collections::HashSet; use std::collections::HashSet;
impl NodeInfo { impl NodeInfo {
pub async fn has_reciprocal_sponsorship(
local_address: &str,
operational_peer_wallets: &[String],
) -> bool {
let map = ADDRESS_MAP.lock().await;
let Some(local) = map.get(local_address) else {
return false;
};
if local.deleted_timestamp > 0
|| local.added_by == local_address
|| !operational_peer_wallets.contains(&local.added_by)
{
return false;
}
map.get(&local.added_by)
.map(|peer| peer.deleted_timestamp == 0 && peer.added_by == local_address)
.unwrap_or(false)
}
pub(crate) async fn signed_mapping_state(
) -> (Vec<SyncedNodeState>, Vec<SignedMonitorEdit>) {
let map = ADDRESS_MAP.lock().await;
let monitor_events = MONITOR_EVENT_STATE.lock().await;
let mut memberships: Vec<SyncedNodeState> = map
.iter()
.map(|(address, node)| SyncedNodeState {
edit: SignedNodeEdit {
address: address.clone(),
ip: node.ip.clone(),
port: node.port,
modified_by: node.added_by.clone(),
modified_timestamp: node.added_timestamp,
modified_signature: node.added_signature.clone(),
},
deleted_timestamp: node.deleted_timestamp,
deleted_block: node.deleted_block,
monitoring: Vec::new(),
})
.collect();
memberships.sort_by(|left, right| left.edit.address.cmp(&right.edit.address));
let mut monitors: Vec<SignedMonitorEdit> = monitor_events.values().cloned().collect();
monitors.sort_by(|left, right| {
(
left.modified_timestamp,
left.modified_block,
left.action,
&left.modified_signature,
&left.monitored_address,
&left.monitoring_address,
&left.target_ip,
)
.cmp(&(
right.modified_timestamp,
right.modified_block,
right.action,
&right.modified_signature,
&right.monitored_address,
&right.monitoring_address,
&right.target_ip,
))
});
(memberships, monitors)
}
pub(super) fn canonical_mapping_records(
map: &HashMap<String, NodeInfo>,
) -> Result<Vec<u8>, String> {
let mut records = map
.iter()
.map(|(address, node)| {
let mut monitoring = node.monitoring.clone();
monitoring.sort();
monitoring.dedup();
(
address.clone(),
node.ip.clone(),
node.port,
node.added_by.clone(),
node.added_timestamp,
node.added_signature.clone(),
node.deleted_timestamp,
node.deleted_block,
monitoring,
)
})
.collect::<Vec<_>>();
records.sort_by(|left, right| left.0.cmp(&right.0));
let mut canonical = b"contractless-network-map-v1".to_vec();
canonical.extend_from_slice(&(records.len() as u32).to_le_bytes());
for (
address,
ip,
port,
added_by,
added_timestamp,
added_signature,
deleted_timestamp,
deleted_block,
monitoring,
) in records
{
let address_bytes = Wallet::short_address_to_bytes(&address)
.ok_or_else(|| "mapping contained an invalid address".to_string())?;
let added_by_bytes = if added_by.is_empty() {
vec![0; Wallet::SHORT_ADDRESS_BYTES_LENGTH]
} else {
Wallet::short_address_to_bytes(&added_by)
.ok_or_else(|| "mapping contained an invalid sponsor".to_string())?
};
let signature_bytes = if added_signature.is_empty() {
vec![0; Wallet::SIGNATURE_LENGTH]
} else {
let bytes = decode(&added_signature)
.map_err(|_| "mapping contained an invalid add signature".to_string())?;
if bytes.len() != Wallet::SIGNATURE_LENGTH {
return Err("mapping contained an invalid add signature length".to_string());
}
bytes
};
let monitor_count = u16::try_from(monitoring.len())
.map_err(|_| "mapping contained too many monitors".to_string())?;
canonical.extend_from_slice(&address_bytes);
canonical.extend_from_slice(&ip_to_binary(&ip));
canonical.extend_from_slice(&port.to_le_bytes());
canonical.extend_from_slice(&added_by_bytes);
canonical.extend_from_slice(&added_timestamp.to_le_bytes());
canonical.extend_from_slice(&signature_bytes);
canonical.extend_from_slice(&deleted_timestamp.to_le_bytes());
canonical.extend_from_slice(&deleted_block.to_le_bytes());
canonical.extend_from_slice(&monitor_count.to_le_bytes());
for monitor in monitoring {
canonical.extend_from_slice(
&Wallet::short_address_to_bytes(&monitor)
.ok_or_else(|| "mapping contained an invalid monitor".to_string())?,
);
}
}
Ok(canonical)
}
pub async fn canonical_mapping_digest() -> Result<Vec<u8>, String> {
let map = ADDRESS_MAP.lock().await;
let monitor_events = MONITOR_EVENT_STATE.lock().await;
let mut canonical = Self::canonical_mapping_records(&map)?;
let monitor_state = Self::signed_monitor_state_bytes_from(&monitor_events);
canonical.extend_from_slice(&(monitor_state.len() as u32).to_le_bytes());
canonical.extend_from_slice(&monitor_state);
decode(skein_256_hash_bytes(&canonical))
.map_err(|_| "failed to encode mapping digest".to_string())
}
pub async fn governance_node_snapshot(address: &str) -> Option<GovernanceNodeSnapshot> { pub async fn governance_node_snapshot(address: &str) -> Option<GovernanceNodeSnapshot> {
let map = ADDRESS_MAP.lock().await; let map = ADDRESS_MAP.lock().await;
map.get(address).map(|node| GovernanceNodeSnapshot { map.get(address).map(|node| GovernanceNodeSnapshot {
@ -191,10 +355,13 @@ impl NodeInfo {
} }
pub async fn request_valid_nodes() -> RpcResponse { pub async fn request_valid_nodes() -> RpcResponse {
// Serialize the in-memory node map into the binary layout // Serialize the complete in-memory node map used by peer bootstrap.
// used by peer bootstrap and node-list synchronization. // Signed membership proves how each node entered the map, while the
// remaining fields carry the current derived state needed by a syncing
// node to validate blocks and avoid dialing deleted peers.
let map = ADDRESS_MAP.lock().await; let map = ADDRESS_MAP.lock().await;
let mut data: Vec<u8> = Vec::with_capacity(map.len() * NODE_RECORD_FIXED_BYTES); let monitor_events = MONITOR_EVENT_STATE.lock().await;
let mut mapping: Vec<u8> = Vec::with_capacity(map.len() * NODE_RECORD_FIXED_BYTES);
for (address, node_info) in map.iter() { for (address, node_info) in map.iter() {
let address_bytes = match Wallet::short_address_to_bytes(address) { let address_bytes = match Wallet::short_address_to_bytes(address) {
@ -208,31 +375,50 @@ impl NodeInfo {
Wallet::short_address_to_bytes(&node_info.added_by).unwrap_or_default(); Wallet::short_address_to_bytes(&node_info.added_by).unwrap_or_default();
let added_timestamp_bytes = node_info.added_timestamp.to_le_bytes(); let added_timestamp_bytes = node_info.added_timestamp.to_le_bytes();
// Network-map snapshots contain signed membership only. Liveness let monitor_bytes: Vec<u8> = node_info
// is transferred separately as signed monitor assertions. .monitoring
let deleted_timestamp_bytes = 0_u64.to_le_bytes(); .iter()
let deleted_block_bytes = 0_u32.to_le_bytes(); .filter_map(|monitor| Wallet::short_address_to_bytes(monitor))
let monitor_count_bytes = 0_u16.to_le_bytes(); .flatten()
.collect();
let monitor_count = monitor_bytes.len() / Wallet::SHORT_ADDRESS_BYTES_LENGTH;
let Ok(monitor_count) = u16::try_from(monitor_count) else {
warn!(
"[network_map] too many monitors for node {address}; skipping snapshot record"
);
continue;
};
let added_signature_bytes = decode(node_info.added_signature.clone()).unwrap(); let added_signature_bytes = decode(node_info.added_signature.clone()).unwrap();
// Field order here must match the parser used by node-list // Field order here must match the parser used by node-list
// synchronization. // synchronization.
data.extend_from_slice(&address_bytes); mapping.extend_from_slice(&address_bytes);
data.extend_from_slice(&ip_bytes); mapping.extend_from_slice(&ip_bytes);
data.extend_from_slice(&port_bytes); mapping.extend_from_slice(&port_bytes);
data.push(blocks_mined); mapping.push(blocks_mined);
if added_by_bytes.len() == Wallet::SHORT_ADDRESS_BYTES_LENGTH { if added_by_bytes.len() == Wallet::SHORT_ADDRESS_BYTES_LENGTH {
data.extend_from_slice(&added_by_bytes); mapping.extend_from_slice(&added_by_bytes);
} else { } else {
data.extend_from_slice(&[0u8; Wallet::SHORT_ADDRESS_BYTES_LENGTH]); mapping.extend_from_slice(&[0u8; Wallet::SHORT_ADDRESS_BYTES_LENGTH]);
} }
data.extend_from_slice(&added_timestamp_bytes); mapping.extend_from_slice(&added_timestamp_bytes);
data.extend_from_slice(&added_signature_bytes); mapping.extend_from_slice(&added_signature_bytes);
data.extend_from_slice(&deleted_timestamp_bytes); mapping.extend_from_slice(&node_info.deleted_timestamp.to_le_bytes());
data.extend_from_slice(&deleted_block_bytes); mapping.extend_from_slice(&node_info.deleted_block.to_le_bytes());
data.extend_from_slice(&monitor_count_bytes); mapping.extend_from_slice(&monitor_count.to_le_bytes());
mapping.extend_from_slice(&monitor_bytes);
} }
let monitor_state = Self::signed_monitor_state_bytes_from(&monitor_events);
let mut data =
Vec::with_capacity(NETWORK_SNAPSHOT_HEADER_BYTES + mapping.len() + monitor_state.len());
data.extend_from_slice(NETWORK_SNAPSHOT_MAGIC);
data.push(NETWORK_SNAPSHOT_VERSION);
data.extend_from_slice(&(mapping.len() as u32).to_le_bytes());
data.extend_from_slice(&(monitor_state.len() as u32).to_le_bytes());
data.extend_from_slice(&mapping);
data.extend_from_slice(&monitor_state);
RpcResponse::Binary(data) RpcResponse::Binary(data)
} }
@ -253,3 +439,91 @@ impl NodeInfo {
Wallet::sign_transaction(&hashed_data, private_key).await Wallet::sign_transaction(&hashed_data, private_key).await
} }
} }
#[cfg(test)]
mod tests {
use super::*;
use crate::records::memory::network_mapping::structs::{
NODE_DELETED_BLOCK_OFFSET, NODE_DELETED_TIMESTAMP_OFFSET, NODE_MONITOR_COUNT_OFFSET,
NODE_RECORD_FIXED_BYTES,
};
#[tokio::test]
async fn bootstrap_snapshot_keeps_deletion_and_monitor_state() {
let address = "1111111111111111111111111111111111111111.cltc";
let monitor = "2222222222222222222222222222222222222222.cltc";
let mut node = NodeInfo::new(
"198.51.100.10".to_string(),
50050,
73,
monitor.to_string(),
1234,
"00".repeat(Wallet::SIGNATURE_LENGTH),
);
node.deleted_timestamp = 5678;
node.deleted_block = 9012;
node.monitoring = vec![monitor.to_string()];
{
let mut map = ADDRESS_MAP.lock().await;
map.clear();
map.insert(address.to_string(), node);
}
MONITOR_EVENT_STATE.lock().await.clear();
let RpcResponse::Binary(bytes) = NodeInfo::request_valid_nodes().await;
let record_start = NETWORK_SNAPSHOT_HEADER_BYTES;
let record_end =
record_start + NODE_RECORD_FIXED_BYTES + Wallet::SHORT_ADDRESS_BYTES_LENGTH;
assert_eq!(
bytes.len(),
NETWORK_SNAPSHOT_HEADER_BYTES
+ NODE_RECORD_FIXED_BYTES
+ Wallet::SHORT_ADDRESS_BYTES_LENGTH
);
assert_eq!(&bytes[0..4], NETWORK_SNAPSHOT_MAGIC);
assert_eq!(bytes[4], NETWORK_SNAPSHOT_VERSION);
assert_eq!(
u32::from_le_bytes(bytes[5..9].try_into().unwrap()) as usize,
NODE_RECORD_FIXED_BYTES + Wallet::SHORT_ADDRESS_BYTES_LENGTH
);
assert_eq!(u32::from_le_bytes(bytes[9..13].try_into().unwrap()), 0);
assert_eq!(
u64::from_le_bytes(
bytes[record_start + NODE_DELETED_TIMESTAMP_OFFSET
..record_start + NODE_DELETED_BLOCK_OFFSET]
.try_into()
.unwrap()
),
5678
);
assert_eq!(
u32::from_le_bytes(
bytes[record_start + NODE_DELETED_BLOCK_OFFSET
..record_start + NODE_MONITOR_COUNT_OFFSET]
.try_into()
.unwrap()
),
9012
);
assert_eq!(
u16::from_le_bytes(
bytes[record_start + NODE_MONITOR_COUNT_OFFSET
..record_start + NODE_RECORD_FIXED_BYTES]
.try_into()
.unwrap()
),
1
);
assert_eq!(
Wallet::bytes_to_short_address(
&bytes[record_start + NODE_RECORD_FIXED_BYTES..record_end]
)
.as_deref(),
Some(monitor)
);
ADDRESS_MAP.lock().await.clear();
MONITOR_EVENT_STATE.lock().await.clear();
}
}

View File

@ -37,6 +37,10 @@ pub const MONITOR_EVENT_SIGNATURE_OFFSET: usize =
MONITOR_EVENT_BLOCK_OFFSET + NODE_DELETED_BLOCK_BYTES; MONITOR_EVENT_BLOCK_OFFSET + NODE_DELETED_BLOCK_BYTES;
pub const MONITOR_EVENT_BYTES: usize = MONITOR_EVENT_SIGNATURE_OFFSET + Wallet::SIGNATURE_LENGTH; pub const MONITOR_EVENT_BYTES: usize = MONITOR_EVENT_SIGNATURE_OFFSET + Wallet::SIGNATURE_LENGTH;
pub const NETWORK_SNAPSHOT_MAGIC: &[u8; 4] = b"CNMS";
pub const NETWORK_SNAPSHOT_VERSION: u8 = 1;
pub const NETWORK_SNAPSHOT_HEADER_BYTES: usize = 13;
// SignedNodeEdit carries the signed node membership payload used by add/delete updates. // SignedNodeEdit carries the signed node membership payload used by add/delete updates.
#[derive(Debug, Clone)] #[derive(Debug, Clone)]
pub struct SignedNodeEdit { pub struct SignedNodeEdit {
@ -59,6 +63,14 @@ pub struct SignedMonitorEdit {
pub modified_signature: String, pub modified_signature: String,
} }
#[derive(Debug, Clone)]
pub(crate) struct SyncedNodeState {
pub edit: SignedNodeEdit,
pub deleted_timestamp: u64,
pub deleted_block: u32,
pub monitoring: Vec<String>,
}
// AddAddressParams groups the shared context needed to add a node to the network map. // AddAddressParams groups the shared context needed to add a node to the network map.
#[derive(Clone)] #[derive(Clone)]
pub struct AddAddressParams { pub struct AddAddressParams {

View File

@ -29,6 +29,9 @@ pub struct ConnectionInfo {
pub local_setup_complete: bool, pub local_setup_complete: bool,
pub remote_setup_complete: bool, pub remote_setup_complete: bool,
pub local_setup_acknowledged: bool, pub local_setup_acknowledged: bool,
// Incoming peers that connect behind the tip retain the height they
// originally need to reach. This avoids chasing a moving local tip.
pub catch_up_target: Option<u32>,
pub ready: bool, pub ready: bool,
pub health: ConnectionHealth, pub health: ConnectionHealth,
} }
@ -228,6 +231,7 @@ impl ConnectionInfo {
local_setup_complete: false, local_setup_complete: false,
remote_setup_complete: false, remote_setup_complete: false,
local_setup_acknowledged: false, local_setup_acknowledged: false,
catch_up_target: None,
ready: false, ready: false,
health: ConnectionHealth::Connected, health: ConnectionHealth::Connected,
} }

View File

@ -1,248 +0,0 @@
use crate::blocks::loans::LoanContractTransaction;
use crate::common::types::Transaction;
use crate::decode;
use crate::records::memory::mempool::BASECOIN;
use crate::rpc::commands::transaction_by_txid::request_transaction_by_txid;
use crate::rpc::responses::RpcResponse;
use crate::sled::Db;
use std::collections::BTreeMap;
pub const MINER_EARNING_FEE: u8 = 0;
pub const MINER_EARNING_TIP: u8 = 1;
pub const MINER_EARNING_ASSET_BYTES: usize = 15;
pub const MINER_EARNING_ENTRY_BYTES: usize = 1 + MINER_EARNING_ASSET_BYTES + 4 + 8;
const FINALIZED_MINER_EARNINGS_TREE: &str = "finalized_miner_earnings";
#[derive(Clone, Debug, Default, PartialEq, Eq)]
pub struct MinerEarnings {
entries: BTreeMap<(u8, String, u32), u64>,
}
impl MinerEarnings {
pub fn add_fee(&mut self, amount: u64) -> Result<(), String> {
self.add(MINER_EARNING_FEE, BASECOIN.clone(), 0, amount)
}
pub fn add_tip(&mut self, asset: &str, nft_series: u32, amount: u64) -> Result<(), String> {
self.add(MINER_EARNING_TIP, asset.to_string(), nft_series, amount)
}
fn add(
&mut self,
earning_type: u8,
asset: String,
nft_series: u32,
amount: u64,
) -> Result<(), String> {
if amount == 0 {
return Ok(());
}
let asset = fixed_asset_name(&asset)?;
let entry = self
.entries
.entry((earning_type, asset, nft_series))
.or_default();
*entry = entry
.checked_add(amount)
.ok_or_else(|| "Miner earnings amount overflowed u64.".to_string())?;
Ok(())
}
pub fn merge(&mut self, other: MinerEarnings) -> Result<(), String> {
for ((earning_type, asset, nft_series), amount) in other.entries {
self.add(earning_type, asset, nft_series, amount)?;
}
Ok(())
}
pub fn encode(&self) -> Result<Vec<u8>, String> {
let count = u32::try_from(self.entries.len())
.map_err(|_| "Miner earnings entry count overflowed u32.".to_string())?;
let mut bytes =
Vec::with_capacity(4 + self.entries.len().saturating_mul(MINER_EARNING_ENTRY_BYTES));
bytes.extend_from_slice(&count.to_le_bytes());
for ((earning_type, asset, nft_series), amount) in &self.entries {
bytes.push(*earning_type);
bytes.extend_from_slice(asset.as_bytes());
bytes.extend_from_slice(&nft_series.to_le_bytes());
bytes.extend_from_slice(&amount.to_le_bytes());
}
Ok(bytes)
}
}
fn fixed_asset_name(asset: &str) -> Result<String, String> {
let bytes = asset.as_bytes();
if bytes.len() != MINER_EARNING_ASSET_BYTES {
return Err(format!(
"Miner earning asset must be exactly {MINER_EARNING_ASSET_BYTES} bytes: {asset:?}"
));
}
Ok(asset.to_string())
}
async fn loan_payment_tip_asset(db: &Db, contract_hash: &str) -> Result<String, String> {
let contract_hash = decode(contract_hash)
.map_err(|err| format!("Failed to decode loan contract hash for miner tip: {err}"))?;
let RpcResponse::Binary(bytes) = request_transaction_by_txid(db, contract_hash).await;
if bytes.is_empty() || bytes[0] != 7 {
return Err("Loan payment miner tip references an invalid loan contract.".to_string());
}
let loan = LoanContractTransaction::from_bytes(7, &bytes[1..])
.await
.map_err(|err| format!("Failed to decode loan contract for miner tip: {err}"))?;
fixed_asset_name(&loan.unsigned_loan_contract.loan_coin)
}
pub async fn earnings_from_transactions(
db: &Db,
transactions: &[Transaction],
) -> Result<MinerEarnings, String> {
let mut earnings = MinerEarnings::default();
for transaction in transactions {
match transaction {
Transaction::Genesis(_) | Transaction::Rewards(_) => {}
Transaction::Transfer(tx) => earnings.add_fee(tx.unsigned_transfer.txfee)?,
Transaction::Token(tx) => earnings.add_fee(tx.unsigned_create_token.txfee)?,
Transaction::IssueToken(tx) => earnings.add_fee(tx.unsigned_issue_token.txfee)?,
Transaction::Burn(tx) => earnings.add_fee(tx.unsigned_burn.txfee)?,
Transaction::Nft(tx) => earnings.add_fee(tx.unsigned_create_nft.txfee)?,
Transaction::Marketing(tx) => earnings.add_fee(tx.unsigned_marketing.txfee)?,
Transaction::Swap(tx) => {
earnings.add_fee(tx.unsigned_swap.txfee1)?;
earnings.add_fee(tx.unsigned_swap.txfee2)?;
earnings.add_tip(
&tx.unsigned_swap.ticker1,
tx.unsigned_swap.nft_series1,
tx.unsigned_swap.tip1,
)?;
earnings.add_tip(
&tx.unsigned_swap.ticker2,
tx.unsigned_swap.nft_series2,
tx.unsigned_swap.tip2,
)?;
}
Transaction::Lender(tx) => earnings.add_fee(tx.unsigned_loan_contract.txfee)?,
Transaction::Borrower(tx) => {
earnings.add_fee(tx.unsigned_contract_payment.txfee)?;
if tx.unsigned_contract_payment.tip > 0 {
let asset =
loan_payment_tip_asset(db, &tx.unsigned_contract_payment.contract_hash)
.await?;
earnings.add_tip(&asset, 0, tx.unsigned_contract_payment.tip)?;
}
}
Transaction::Collateral(tx) => earnings.add_fee(tx.unsigned_collateral_claim.txfee)?,
Transaction::Vanity(tx) => earnings.add_fee(tx.unsigned_vanity_address.txfee)?,
}
}
Ok(earnings)
}
pub fn store_finalized_miner_earnings(
db: &Db,
reward_txid: &[u8],
earnings: &MinerEarnings,
) -> Result<(), String> {
let encoded = earnings.encode()?;
db.open_tree(FINALIZED_MINER_EARNINGS_TREE)
.map_err(|err| format!("Failed to open finalized miner earnings tree: {err}"))?
.insert(reward_txid, encoded)
.map_err(|err| format!("Failed to store finalized miner earnings: {err}"))?;
Ok(())
}
pub fn finalized_miner_earnings_bytes(db: &Db, reward_txid: &[u8]) -> Result<Vec<u8>, String> {
let tree = db
.open_tree(FINALIZED_MINER_EARNINGS_TREE)
.map_err(|err| format!("Failed to open finalized miner earnings tree: {err}"))?;
let encoded = tree
.get(reward_txid)
.map_err(|err| format!("Failed to read finalized miner earnings: {err}"))?
.map(|value| value.to_vec())
.unwrap_or_else(|| 0u32.to_le_bytes().to_vec());
validate_encoded_miner_earnings(&encoded)?;
Ok(encoded)
}
pub fn remove_finalized_miner_earnings(db: &Db, reward_txid: &[u8]) {
if let Ok(tree) = db.open_tree(FINALIZED_MINER_EARNINGS_TREE) {
let _ = tree.remove(reward_txid);
}
}
fn validate_encoded_miner_earnings(encoded: &[u8]) -> Result<(), String> {
let count_bytes = encoded
.get(0..4)
.ok_or_else(|| "Finalized miner earnings record is missing its count.".to_string())?;
let count = u32::from_le_bytes(
count_bytes
.try_into()
.map_err(|_| "Finalized miner earnings count was invalid.".to_string())?,
) as usize;
let expected = 4usize
.checked_add(
count
.checked_mul(MINER_EARNING_ENTRY_BYTES)
.ok_or_else(|| "Finalized miner earnings length overflowed.".to_string())?,
)
.ok_or_else(|| "Finalized miner earnings length overflowed.".to_string())?;
if encoded.len() != expected {
return Err(format!(
"Finalized miner earnings record has invalid length: expected {expected}, found {}.",
encoded.len()
));
}
Ok(())
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn combines_matching_fees_and_tips() {
let mut earnings = MinerEarnings::default();
earnings.add_fee(10).unwrap();
earnings.add_fee(20).unwrap();
earnings.add_tip("example ", 0, 7).unwrap();
earnings.add_tip("example ", 0, 8).unwrap();
let encoded = earnings.encode().unwrap();
assert_eq!(u32::from_le_bytes(encoded[0..4].try_into().unwrap()), 2);
assert_eq!(encoded.len(), 4 + (2 * MINER_EARNING_ENTRY_BYTES));
assert!(encoded.windows(8).any(|bytes| bytes == 30u64.to_le_bytes()));
assert!(encoded.windows(8).any(|bytes| bytes == 15u64.to_le_bytes()));
}
#[test]
fn separates_nft_tip_series() {
let mut earnings = MinerEarnings::default();
earnings.add_tip("example ", 1, 1).unwrap();
earnings.add_tip("example ", 2, 1).unwrap();
let encoded = earnings.encode().unwrap();
assert_eq!(u32::from_le_bytes(encoded[0..4].try_into().unwrap()), 2);
}
#[test]
fn persists_and_removes_finalized_earnings() {
let db = crate::sled::Config::new().temporary(true).open().unwrap();
let reward_txid = [7u8; 32];
let mut earnings = MinerEarnings::default();
earnings.add_fee(42).unwrap();
store_finalized_miner_earnings(&db, &reward_txid, &earnings).unwrap();
let stored = finalized_miner_earnings_bytes(&db, &reward_txid).unwrap();
assert_eq!(u32::from_le_bytes(stored[0..4].try_into().unwrap()), 1);
remove_finalized_miner_earnings(&db, &reward_txid);
assert_eq!(
finalized_miner_earnings_bytes(&db, &reward_txid).unwrap(),
0u32.to_le_bytes()
);
}
}

View File

@ -56,6 +56,11 @@ pub async fn save_block(params: SaveBlockParams) -> Result<(), String> {
} = params; } = params;
let genesis_missing = !genesis_checkup().await; let genesis_missing = !genesis_checkup().await;
let contains_genesis = block
.transactions
.iter()
.any(|transaction| matches!(transaction, Transaction::Genesis(_)));
validate_genesis_save_state(genesis_missing, contains_genesis)?;
if save_type.is_updating() { if save_type.is_updating() {
if is_reorganizing_mode() && !allow_during_reorg { if is_reorganizing_mode() && !allow_during_reorg {
return Err("Cannot save discovered block while reorganizing.".to_string()); return Err("Cannot save discovered block while reorganizing.".to_string());
@ -256,6 +261,17 @@ pub async fn save_block(params: SaveBlockParams) -> Result<(), String> {
Ok(()) Ok(())
} }
fn validate_genesis_save_state(
genesis_missing: bool,
contains_genesis: bool,
) -> Result<(), String> {
match (genesis_missing, contains_genesis) {
(true, false) => Err("Cannot save a non-genesis block before genesis exists".to_string()),
(false, true) => Err("Genesis block already exists".to_string()),
_ => Ok(()),
}
}
async fn log_saved_block_difficulty( async fn log_saved_block_difficulty(
block_number: u32, block_number: u32,
timestamp: u32, timestamp: u32,
@ -429,7 +445,9 @@ async fn save_binary_data_with_mempool_stream(
// Only advance mined-count tracking when this save actually moved // Only advance mined-count tracking when this save actually moved
// the persisted chain height forward. // the persisted chain height forward.
if get_height(db) > previous_height { if get_height(db) > previous_height {
NodeInfo::increment_mined(&miner).await; if let Err(err) = NodeInfo::increment_mined(db, &miner).await {
error!("Failed to update mined-count cache: {err}");
}
} }
Ok(()) Ok(())
@ -590,7 +608,9 @@ async fn save_binary_data(params: SaveBinaryDataParams<'_>) -> Result<(), String
// Only advance mined-count tracking when this save actually moved // Only advance mined-count tracking when this save actually moved
// the persisted chain height forward. // the persisted chain height forward.
if get_height(db) > previous_height { if get_height(db) > previous_height {
NodeInfo::increment_mined(&miner).await; if let Err(err) = NodeInfo::increment_mined(db, &miner).await {
error!("Failed to update mined-count cache: {err}");
}
} }
Ok(()) Ok(())
@ -693,3 +713,20 @@ fn format_block_time(timestamp: u32) -> String {
None => "invalid-time".to_string(), None => "invalid-time".to_string(),
} }
} }
#[cfg(test)]
mod tests {
use super::validate_genesis_save_state;
#[test]
fn only_genesis_can_be_saved_before_genesis_exists() {
assert!(validate_genesis_save_state(true, true).is_ok());
assert!(validate_genesis_save_state(true, false).is_err());
}
#[test]
fn genesis_cannot_be_saved_twice() {
assert!(validate_genesis_save_state(false, false).is_ok());
assert!(validate_genesis_save_state(false, true).is_err());
}
}

View File

@ -34,7 +34,9 @@ use crate::rpc::server::connection_memory_manager::{
use crate::rpc::server::rpc_command_loop::start_loop; use crate::rpc::server::rpc_command_loop::start_loop;
use crate::sled::Db; use crate::sled::Db;
use crate::sleep; use crate::sleep;
use crate::startup::network_broadcast::announce_self_to_network; use crate::startup::network_broadcast::{
announce_self_to_network, compare_network_mapping_digest, reconcile_network_mapping_with_peer,
};
use crate::startup::remote_height::request_remote_height; use crate::startup::remote_height::request_remote_height;
use crate::thread_rng; use crate::thread_rng;
use crate::wallets::structures::Wallet; use crate::wallets::structures::Wallet;
@ -43,6 +45,25 @@ use crate::Duration;
use crate::Mutex; use crate::Mutex;
use crate::SliceRandom; use crate::SliceRandom;
use crate::SocketAddr; use crate::SocketAddr;
use std::sync::atomic::{AtomicBool, Ordering};
use tokio::time::Instant as TokioInstant;
static STARTUP_DISCOVERY_ACTIVE: AtomicBool = AtomicBool::new(false);
struct StartupDiscoveryGuard;
impl Drop for StartupDiscoveryGuard {
fn drop(&mut self) {
STARTUP_DISCOVERY_ACTIVE.store(false, Ordering::SeqCst);
}
}
fn try_start_startup_discovery() -> Option<StartupDiscoveryGuard> {
STARTUP_DISCOVERY_ACTIVE
.compare_exchange(false, true, Ordering::SeqCst, Ordering::SeqCst)
.ok()
.map(|_| StartupDiscoveryGuard)
}
use crate::TcpStream; use crate::TcpStream;
#[derive(Clone)] #[derive(Clone)]
@ -52,19 +73,25 @@ pub struct BootstrapParams {
pub wallet: Arc<Wallet>, pub wallet: Arc<Wallet>,
pub db: Db, pub db: Db,
pub map: Arc<Mutex<Command>>, pub map: Arc<Mutex<Command>>,
pub first: bool,
pub run_startup_sync: bool, pub run_startup_sync: bool,
} }
pub fn spawn_bootstrap_peer_discovery(params: BootstrapParams) { pub fn spawn_bootstrap_peer_discovery(params: BootstrapParams) {
let Some(startup_guard) = try_start_startup_discovery() else {
info!("[startup] startup synchronization is already active; coalescing duplicate request");
spawn_peer_setup_retry(params.connections_key, params.map);
return;
};
tokio::spawn(async move { tokio::spawn(async move {
let _startup_guard = startup_guard;
if let Err(e) = bootstrap_peer_discovery(params).await { if let Err(e) = bootstrap_peer_discovery(params).await {
eprintln!("[bootstrap] error: {e}"); eprintln!("[bootstrap] error: {e}");
} }
}); });
} }
pub async fn bootstrap_peer_discovery(mut params: BootstrapParams) -> Result<(), String> { pub async fn bootstrap_peer_discovery(params: BootstrapParams) -> Result<(), String> {
let chain_sync_guard = if params.run_startup_sync { let chain_sync_guard = if params.run_startup_sync {
Some(begin_chain_sync().await) Some(begin_chain_sync().await)
} else { } else {
@ -72,117 +99,136 @@ pub async fn bootstrap_peer_discovery(mut params: BootstrapParams) -> Result<(),
}; };
let (_, _, local_endpoint) = get_ip_and_port().await; let (_, _, local_endpoint) = get_ip_and_port().await;
let max = SETTINGS.outgoing_connections; let max = SETTINGS.outgoing_connections;
let mut current_key = params.connections_key.clone(); let current_key = params.connections_key.clone();
let mut stream = params.stream; let stream = params.stream;
params.first = false;
let mut imported_chain = false; let mut imported_chain = false;
let mut candidate_endpoints = NodeInfo::active_node_endpoints().await;
candidate_endpoints.shuffle(&mut thread_rng());
for addr_string in candidate_endpoints { if params.run_startup_sync {
let outgoing_connections = { // Build the multi-peer torrent pool before chain synchronization. Peer
// setup runs concurrently so one unresponsive registry endpoint cannot
// delay every healthy peer or serialize startup for several minutes.
let connected_outgoing = {
let connections = CONNECTIONS.read().await; let connections = CONNECTIONS.read().await;
connections connections
.as_ref() .as_ref()
.map(|connection| connection.count_ready_outgoing_connections()) .map(|connection| connection.count_outgoing_connections())
.unwrap_or(0) .unwrap_or(0)
}; };
if outgoing_connections >= max as usize { let available_slots = (max as usize).saturating_sub(connected_outgoing);
let mut candidate_endpoints = NodeInfo::active_node_endpoints().await;
candidate_endpoints.shuffle(&mut thread_rng());
let mut setup_tasks = Vec::new();
for addr_string in candidate_endpoints
.into_iter()
.filter(|endpoint| endpoint != &local_endpoint && endpoint != &current_key)
.take(available_slots)
{
if Connection::get_stream_from_memory(&addr_string)
.await
.is_some()
{
continue;
}
let socket_addr: SocketAddr = match addr_string.parse() {
Ok(addr) => addr,
Err(_) => {
warn!("Invalid mapped peer endpoint: {addr_string}");
continue;
}
};
let wallet = params.wallet.clone();
let db = params.db.clone();
let map = params.map.clone();
setup_tasks.push(tokio::spawn(async move {
if Connection::get_stream_from_memory(&addr_string)
.await
.is_some()
{
return;
}
let connect = Connect {
addr: socket_addr,
node_ip: addr_string.clone(),
wallet,
db,
map,
first: false,
};
if let Err(err) = connect_and_handshake(connect).await {
warn!("Failed to connect to discovered peer {addr_string}: {err}");
}
}));
}
if !setup_tasks.is_empty() {
let warmup_deadline = TokioInstant::now() + Duration::from_secs(10);
while TokioInstant::now() < warmup_deadline
&& setup_tasks.iter().any(|task| !task.is_finished())
{
sleep(Duration::from_millis(100)).await;
}
let pool_size = startup_synced_peer_streams().await.len();
info!(
"[sync] torrent peer pool warmed with {pool_size} synchronized peer(s); starting chain synchronization"
);
// Dropping JoinHandle values detaches unfinished setup tasks. They
// may finish later, but cannot hold the canonical sync path open.
drop(setup_tasks);
}
}
if params.run_startup_sync {
loop {
let local_height = get_height(&params.db);
let remote_height =
request_remote_height(stream.clone(), params.map.clone(), current_key.clone())
.await?;
ensure_compatible_genesis(
stream.clone(),
params.map.clone(),
current_key.clone(),
remote_height,
)
.await?;
let local_genesis_exists = genesis_checkup().await;
if !local_genesis_exists || remote_height > local_height + 10 {
imported_chain = true;
info!("[sync] Starting sync from {local_height} to {remote_height}");
node_syncing(
stream.clone(),
&params.db,
remote_height,
params.map.clone(),
true,
params.wallet.clone(),
current_key.clone(),
)
.await
.map_err(|e| format!("Sync error: {e}"))?;
if !local_genesis_exists && !genesis_checkup().await && remote_height > 0 {
return Err("Sync completed without obtaining remote genesis".to_string());
}
if !local_genesis_exists && !genesis_checkup().await {
break;
}
continue;
}
break; break;
} }
if addr_string == local_endpoint { let post_sync_local_height = get_height(&params.db);
continue; let post_sync_remote_height =
}
if Connection::get_stream_from_memory(&addr_string)
.await
.is_some()
{
continue;
}
let socket_addr: SocketAddr = match addr_string.parse() {
Ok(addr) => addr,
Err(_) => {
warn!("Invalid mapped peer endpoint: {addr_string}");
continue;
}
};
sleep(Duration::from_secs(2)).await;
let connect = Connect {
addr: socket_addr,
node_ip: addr_string.clone(),
wallet: params.wallet.clone(),
db: params.db.clone(),
map: params.map.clone(),
first: params.first,
};
if let Err(err) = connect_and_handshake(connect).await {
warn!("Failed to connect to discovered peer {addr_string}: {err}");
continue;
}
if let Some(new_stream) = Connection::get_stream_from_memory(&addr_string).await {
stream = new_stream;
current_key = addr_string;
} else {
warn!("Failed to retrieve new stream for: {addr_string}");
}
}
if !params.run_startup_sync {
return Ok(());
}
loop {
let local_height = get_height(&params.db);
let remote_height =
request_remote_height(stream.clone(), params.map.clone(), current_key.clone()).await?; request_remote_height(stream.clone(), params.map.clone(), current_key.clone()).await?;
ensure_compatible_genesis(
let imported_candidates = match hydrate_torrent_candidates(
stream.clone(), stream.clone(),
params.map.clone(), params.map.clone(),
current_key.clone(), current_key.clone(),
remote_height,
) )
.await?; .await
let local_genesis_exists = genesis_checkup().await;
if !local_genesis_exists || remote_height > local_height + 10 {
imported_chain = true;
info!("[sync] Starting sync from {local_height} to {remote_height}");
node_syncing(
stream.clone(),
&params.db,
remote_height,
params.map.clone(),
true,
params.wallet.clone(),
current_key.clone(),
)
.await
.map_err(|e| format!("Sync error: {e}"))?;
if !local_genesis_exists && !genesis_checkup().await && remote_height > 0 {
return Err("Sync completed without obtaining remote genesis".to_string());
}
if !local_genesis_exists && !genesis_checkup().await {
break;
}
continue;
}
break;
}
let post_sync_local_height = get_height(&params.db);
let post_sync_remote_height =
request_remote_height(stream.clone(), params.map.clone(), current_key.clone()).await?;
let imported_candidates =
match hydrate_torrent_candidates(stream.clone(), params.map.clone(), current_key.clone())
.await
{ {
Ok(imported) => { Ok(imported) => {
if imported > 0 { if imported > 0 {
@ -198,49 +244,131 @@ pub async fn bootstrap_peer_discovery(mut params: BootstrapParams) -> Result<(),
} }
}; };
if post_sync_remote_height != post_sync_local_height || imported_candidates > 0 { if post_sync_remote_height != post_sync_local_height || imported_candidates > 0 {
let orphan_checkup_params = OrphanCheckup2 { let orphan_checkup_params = OrphanCheckup2 {
stream: stream.clone(), stream: stream.clone(),
db: params.db.clone(), db: params.db.clone(),
local_height: post_sync_local_height, local_height: post_sync_local_height,
remote_height: post_sync_remote_height, remote_height: post_sync_remote_height,
recheck_from_height: Some(post_sync_local_height.min(post_sync_remote_height)), recheck_from_height: Some(post_sync_local_height.min(post_sync_remote_height)),
map: params.map.clone(), map: params.map.clone(),
node_syncing: true, node_syncing: true,
connections_key: current_key.clone(), connections_key: current_key.clone(),
}; };
match sync_checkup(orphan_checkup_params, params.wallet.clone()).await { match sync_checkup(orphan_checkup_params, params.wallet.clone()).await {
Ok(()) => {} Ok(()) => {}
Err(err) => return Err(format!("Post-sync orphan check error: {err}")), Err(err) => return Err(format!("Post-sync orphan check error: {err}")),
}
}
if imported_chain {
sync_governance_state(
stream.clone(),
&params.db,
params.map.clone(),
current_key.clone(),
)
.await
.map_err(|err| format!("Governance state sync error: {err}"))?;
info!("[governance] finalized state sync completed: peer={current_key}");
}
info!("[sync] post-sync checks complete, mining grace period started");
for (peer_key, _) in startup_synced_peer_streams().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");
if let Some(guard) = chain_sync_guard {
guard.finish();
} }
} }
if imported_chain { // Outside startup, topology repair can use the ordinary sequential path.
sync_governance_state( // Startup itself already opened its torrent pool concurrently above.
stream.clone(), if !params.run_startup_sync {
&params.db, let mut candidate_endpoints = NodeInfo::active_node_endpoints().await;
params.map.clone(), candidate_endpoints.shuffle(&mut thread_rng());
current_key.clone(),
)
.await
.map_err(|err| format!("Governance state sync error: {err}"))?;
info!("[governance] finalized state sync completed: peer={current_key}");
}
info!("[sync] post-sync checks complete, mining grace period started"); for addr_string in candidate_endpoints {
for (peer_key, _) in startup_synced_peer_streams().await { let outgoing_connections = {
if !mark_peer_operational(&peer_key, params.map.clone()).await { let connections = CONNECTIONS.read().await;
spawn_peer_setup_retry(peer_key, params.map.clone()); connections
.as_ref()
.map(|connection| connection.count_ready_outgoing_connections())
.unwrap_or(0)
};
if outgoing_connections >= max as usize {
break;
}
if addr_string == local_endpoint || addr_string == current_key {
continue;
}
if Connection::get_stream_from_memory(&addr_string)
.await
.is_some()
{
continue;
}
let socket_addr: SocketAddr = match addr_string.parse() {
Ok(addr) => addr,
Err(_) => {
warn!("Invalid mapped peer endpoint: {addr_string}");
continue;
}
};
sleep(Duration::from_secs(2)).await;
// The peer may have opened an incoming connection during the delay.
// Recheck before dialing so simultaneous topology discovery does not
// create an avoidable incoming/outgoing duplicate pair.
if Connection::get_stream_from_memory(&addr_string)
.await
.is_some()
{
continue;
}
let connect = Connect {
addr: socket_addr,
node_ip: addr_string.clone(),
wallet: params.wallet.clone(),
db: params.db.clone(),
map: params.map.clone(),
first: false,
};
if let Err(err) = connect_and_handshake(connect).await {
warn!("Failed to connect to discovered peer {addr_string}: {err}");
}
} }
} }
sleep(Duration::from_secs(15)).await;
info!("[sync] mining grace period complete, mining resuming");
if let Some(guard) = chain_sync_guard {
guard.finish();
}
Ok(()) Ok(())
} }
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn duplicate_startup_discovery_is_coalesced_until_owner_finishes() {
STARTUP_DISCOVERY_ACTIVE.store(false, Ordering::SeqCst);
let owner = try_start_startup_discovery().expect("first startup owner should acquire");
assert!(try_start_startup_discovery().is_none());
drop(owner);
assert!(try_start_startup_discovery().is_some());
}
}
pub async fn process_handshake_response( pub async fn process_handshake_response(
response: Vec<u8>, response: Vec<u8>,
wallet: &Wallet, wallet: &Wallet,
@ -401,14 +529,15 @@ pub async fn process_handshake_response(
"[handshake_state] peer={connections_key} remote_height={remote_height} self_add_block={self_add_block} self_add_limit={self_add_limit} join_mode={join_mode}" "[handshake_state] peer={connections_key} remote_height={remote_height} self_add_block={self_add_block} self_add_limit={self_add_limit} join_mode={join_mode}"
); );
// Only the bootstrap owner imports a complete mapping snapshot.
// Discovered peers join the existing map and receive live updates.
announce_self_to_network( announce_self_to_network(
broadcast_stream.clone(), broadcast_stream.clone(),
&wallet.saved.short_address.clone(), &wallet.saved.short_address.clone(),
params.map.clone(), params.map.clone(),
&params.db.clone(), &params.db.clone(),
params.wallet.clone(),
&connections_key, &connections_key,
true, params.first,
) )
.await .await
.map_err(|err| { .map_err(|err| {
@ -418,6 +547,46 @@ pub async fn process_handshake_response(
}); });
io::Error::other(format!("Network join/map sync failed: {err}")) io::Error::other(format!("Network join/map sync failed: {err}"))
})?; })?;
if !params.first {
let matched = compare_network_mapping_digest(
broadcast_stream.clone(),
params.map.clone(),
&connections_key,
)
.await
.map_err(|err| io::Error::other(format!("Network-map digest comparison failed: {err}")))?;
if !matched {
info!(
"[network_map] reconciling signed mapping state with peer {connections_key}"
);
reconcile_network_mapping_with_peer(
broadcast_stream.clone(),
params.map.clone(),
&params.db,
&connections_key,
)
.await
.map_err(|err| io::Error::other(format!("Network-map reconciliation failed: {err}")))?;
let reconciled = compare_network_mapping_digest(
broadcast_stream.clone(),
params.map.clone(),
&connections_key,
)
.await
.map_err(|err| {
io::Error::other(format!("Reconciled network-map comparison failed: {err}"))
})?;
if !reconciled {
remove_stream_from_memory(&stream).await;
return Err(io::Error::other(format!(
"Network-map reconciliation did not converge with peer {connections_key}"
)));
}
}
}
mark_peer_network_map_synced(&connections_key).await; mark_peer_network_map_synced(&connections_key).await;
if params.first { if params.first {
@ -427,17 +596,19 @@ pub async fn process_handshake_response(
wallet: params.wallet.clone(), wallet: params.wallet.clone(),
db: params.db.clone(), db: params.db.clone(),
map: params.map.clone(), map: params.map.clone(),
first: params.first,
run_startup_sync: true, run_startup_sync: true,
}; };
spawn_bootstrap_peer_discovery(bsparams); spawn_bootstrap_peer_discovery(bsparams);
} else { } else if is_normal_mode() {
if is_normal_mode() { if !mark_peer_operational(&connections_key, params.map.clone()).await {
if !mark_peer_operational(&connections_key, params.map.clone()).await { spawn_peer_setup_retry(connections_key.clone(), params.map.clone());
spawn_peer_setup_retry(connections_key.clone(), params.map.clone());
}
} }
} else {
// A concurrently discovered peer may finish setup after the startup
// promotion snapshot. Retry until the canonical sync owner releases
// rather than leaving that valid peer passive indefinitely.
spawn_peer_setup_retry(connections_key.clone(), params.map.clone());
} }
Ok(()) Ok(())
} }

View File

@ -25,6 +25,31 @@ use std::collections::BTreeMap;
const SYNC_PREFETCH_WINDOW: usize = 100; const SYNC_PREFETCH_WINDOW: usize = 100;
const SYNC_TRANSIENT_RETRY_LIMIT: u8 = 3; const SYNC_TRANSIENT_RETRY_LIMIT: u8 = 3;
const SYNC_MAPPING_RETRY_LIMIT: u8 = 8;
const SYNC_MAPPING_RETRY_DELAY: Duration = Duration::from_millis(250);
fn is_mapping_sync_race(err: &str) -> bool {
err.contains("This address is not eligable to mine")
|| err.contains("This miner address is not registered")
|| err.contains("Miner wallet address is not registered")
}
#[cfg(test)]
mod tests {
use super::is_mapping_sync_race;
#[test]
fn missing_or_not_yet_eligible_miner_is_a_mapping_race() {
assert!(is_mapping_sync_race("This miner address is not registered"));
assert!(is_mapping_sync_race("This address is not eligable to mine"));
}
#[test]
fn permanent_block_errors_are_not_mapping_races() {
assert!(!is_mapping_sync_race("Invalid miner proof."));
assert!(!is_mapping_sync_race("Incorrect previous_block_hash."));
}
}
fn is_transient_sync_error(err: &str) -> bool { fn is_transient_sync_error(err: &str) -> bool {
err.contains("No available peer could provide remaining pieces") err.contains("No available peer could provide remaining pieces")
@ -260,14 +285,22 @@ pub async fn node_syncing(
if let Err(err) = verify_and_save_downloaded_block(download, wallet.clone()).await { if let Err(err) = verify_and_save_downloaded_block(download, wallet.clone()).await {
abort_pending_downloads(&mut pending); abort_pending_downloads(&mut pending);
if is_transient_sync_error(&err) { if is_transient_sync_error(&err) || is_mapping_sync_race(&err) {
let retries = retry_counts.entry(next_to_save).or_insert(0); let retries = retry_counts.entry(next_to_save).or_insert(0);
if *retries < SYNC_TRANSIENT_RETRY_LIMIT { let retry_limit = if is_mapping_sync_race(&err) {
SYNC_MAPPING_RETRY_LIMIT
} else {
SYNC_TRANSIENT_RETRY_LIMIT
};
if *retries < retry_limit {
*retries += 1; *retries += 1;
warn!( warn!(
"[sync] transient save/download error at height={next_to_save}; retry {}/{} before orphan recovery: {err}", "[sync] transient save/download error at height={next_to_save}; retry {}/{} before orphan recovery: {err}",
*retries, SYNC_TRANSIENT_RETRY_LIMIT *retries, retry_limit
); );
if is_mapping_sync_race(&err) {
crate::sleep(SYNC_MAPPING_RETRY_DELAY).await;
}
next_to_schedule = next_to_save; next_to_schedule = next_to_save;
continue; continue;
} }

View File

@ -75,6 +75,8 @@ pub const RPC_NETWORK_MONITOR_STATE: u8 = 54;
pub const RPC_SETUP_COMPLETE: u8 = 55; pub const RPC_SETUP_COMPLETE: u8 = 55;
pub const RPC_GOVERNANCE_STATE_SYNC: u8 = 56; pub const RPC_GOVERNANCE_STATE_SYNC: u8 = 56;
pub const RPC_GOVERNANCE_PROPOSAL: u8 = 57; pub const RPC_GOVERNANCE_PROPOSAL: u8 = 57;
pub const RPC_NETWORK_MAPPING_HASH: u8 = 58;
pub const RPC_NETWORK_MEMBERSHIP_RECONCILE: u8 = 59;
pub const RPC_REPLY: u8 = 255; pub const RPC_REPLY: u8 = 255;
pub const MAX_RPC_REPLY_BYTES: usize = 64 * 1024 * 1024; pub const MAX_RPC_REPLY_BYTES: usize = 64 * 1024 * 1024;

View File

@ -1,5 +1,6 @@
use crate::records::memory::network_mapping::structs::{AddAddressParams, SignedNodeEdit}; use crate::records::memory::network_mapping::structs::{AddAddressParams, SignedNodeEdit};
use crate::records::memory::network_mapping::NodeInfo; use crate::records::memory::network_mapping::NodeInfo;
use crate::records::memory::enums::ClientType;
use crate::records::memory::response_channels::Command; use crate::records::memory::response_channels::Command;
use crate::rpc::read_bytes_from_stream; use crate::rpc::read_bytes_from_stream;
use crate::rpc::responses::RpcResponse; use crate::rpc::responses::RpcResponse;
@ -15,6 +16,8 @@ pub async fn add_network_node(
db: &Db, db: &Db,
wallet: Arc<Wallet>, wallet: Arc<Wallet>,
map: Arc<Mutex<Command>>, map: Arc<Mutex<Command>>,
client_type: ClientType,
reconciliation: bool,
) -> Result<(u32, RpcResponse), String> { ) -> Result<(u32, RpcResponse), String> {
// Command 28 carries the signed node-add payload directly after the // Command 28 carries the signed node-add payload directly after the
// UID: node address, node IP, node port, signer, timestamp, and signature. // UID: node address, node IP, node port, signer, timestamp, and signature.
@ -55,20 +58,52 @@ pub async fn add_network_node(
if monitor_count != 0 { if monitor_count != 0 {
return Err("error: Node membership updates cannot contain monitor state".to_string()); 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?; let remote_ip = read_bytes_from_stream::read_caller_ip(stream_locked.clone()).await?;
let edit = SignedNodeEdit {
address,
ip,
port,
modified_by: added_by,
modified_timestamp: added_timestamp,
modified_signature: added_signature,
};
if reconciliation {
if client_type != ClientType::Miner {
return Ok((
uid,
RpcResponse::Binary(b"Error: mapping reconciliation is miner-peer only".to_vec()),
));
}
let deleted_timestamp =
read_bytes_from_stream::read_u64_from_stream(connections_key, stream_locked.clone())
.await?;
let deleted_block =
read_bytes_from_stream::read_u32_from_stream(connections_key, stream_locked.clone())
.await?;
let result = match NodeInfo::import_reconciled_membership(
db,
crate::records::memory::network_mapping::structs::SyncedNodeState {
edit,
deleted_timestamp,
deleted_block,
monitoring: Vec::new(),
},
)
.await
{
Ok(()) => RpcResponse::Binary(b"Success".to_vec()),
Err(err) => RpcResponse::Binary(format!("Error: {err}").into_bytes()),
};
return Ok((uid, result));
}
// NodeInfo owns the signature checks, local re-signing rules, and // NodeInfo owns the signature checks, local re-signing rules, and
// in-memory/broadcast side effects for the actual add operation. // in-memory/broadcast side effects for the actual add operation.
let result = NodeInfo::add_address(AddAddressParams { let result = NodeInfo::add_address(AddAddressParams {
map, map,
edit: SignedNodeEdit { edit,
address,
ip,
port,
modified_by: added_by,
modified_timestamp: added_timestamp,
modified_signature: added_signature,
},
monitors: Vec::new(), monitors: Vec::new(),
blocks_mined: 0_u8, blocks_mined: 0_u8,
remote_ip, remote_ip,

View File

@ -1,67 +0,0 @@
use crate::records::memory::network_mapping::structs::{DeleteAddressParams, SignedNodeEdit};
use crate::records::memory::network_mapping::NodeInfo;
use crate::records::memory::response_channels::Command;
use crate::rpc::read_bytes_from_stream;
use crate::rpc::responses::RpcResponse;
use crate::sled::Db;
use crate::wallets::structures::Wallet;
use crate::Arc;
use crate::Mutex;
use crate::TcpStream;
pub async fn delete_network_node(
connections_key: &str,
stream_locked: Arc<Mutex<TcpStream>>,
db: &Db,
wallet: Arc<Wallet>,
map: Arc<Mutex<Command>>,
) -> Result<(u32, RpcResponse), String> {
// Command 29 uses the same signed node-edit shape as add, but the
// payload represents a delete marker for an existing network node.
let (uid, _) =
read_bytes_from_stream::read_uid_from_stream(connections_key, stream_locked.clone())
.await?;
let address_bytes = read_bytes_from_stream::read_short_address_from_stream(
connections_key,
stream_locked.clone(),
)
.await?;
let address = Wallet::bytes_to_short_address(&address_bytes)
.ok_or_else(|| "error: Invalid short address bytes".to_string())?;
let ip =
read_bytes_from_stream::read_ip_from_stream(connections_key, stream_locked.clone()).await?;
let deleted_by_bytes = read_bytes_from_stream::read_short_address_from_stream(
connections_key,
stream_locked.clone(),
)
.await?;
let deleted_by = Wallet::bytes_to_short_address(&deleted_by_bytes)
.ok_or_else(|| "error: Invalid short address bytes".to_string())?;
let deleted_timestamp =
read_bytes_from_stream::read_u64_from_stream(connections_key, stream_locked.clone())
.await?;
let deleted_signature =
read_bytes_from_stream::read_signature_from_stream(connections_key, stream_locked.clone())
.await?;
let remote_ip = read_bytes_from_stream::read_caller_ip(stream_locked).await?;
// NodeInfo handles validating the signer and applying the delete
// metadata to the network map.
let result = NodeInfo::delete_address(DeleteAddressParams {
map,
edit: SignedNodeEdit {
address,
ip,
modified_by: deleted_by,
modified_timestamp: deleted_timestamp,
modified_signature: deleted_signature,
},
remote_ip,
db: db.clone(),
wallet,
connections_key: connections_key.to_string(),
})
.await;
Ok((uid, result))
}

View File

@ -21,6 +21,7 @@ pub mod latest_block;
pub mod marketing_campaign_history; pub mod marketing_campaign_history;
pub mod memory_by_signature; pub mod memory_by_signature;
pub mod network_info; pub mod network_info;
pub mod network_mapping_hash;
pub mod network_monitor_add; pub mod network_monitor_add;
pub mod network_monitor_remove; pub mod network_monitor_remove;
pub mod network_monitor_state; pub mod network_monitor_state;

View File

@ -0,0 +1,9 @@
use crate::records::memory::network_mapping::NodeInfo;
use crate::rpc::responses::RpcResponse;
pub async fn network_mapping_hash() -> RpcResponse {
match NodeInfo::canonical_mapping_digest().await {
Ok(hash) => RpcResponse::Binary(hash),
Err(err) => RpcResponse::Binary(format!("Error: {err}").into_bytes()),
}
}

View File

@ -1,52 +0,0 @@
use crate::common::binary_conversions::binary_to_ip;
use crate::records::block_height::get_block_height::get_height;
use crate::records::memory::connections::CONNECTIONS;
use crate::records::memory::network_mapping::NodeInfo;
use crate::rpc::responses::RpcResponse;
use crate::sled::Db;
pub async fn eligible_sponsor_endpoint(db: &Db, connections_key: &str) -> Option<String> {
// Before the mature-network gate, node adds are intentionally permissive.
// After height 10,000, only mature sponsors can be returned.
let eligible_ips = if get_height(db) > 10000 {
NodeInfo::eligible_sponsor_ips().await
} else {
NodeInfo::active_node_ips().await
};
if eligible_ips.is_empty() {
return None;
}
{
let connections_lock = CONNECTIONS.read().await;
if let Some(connection) = connections_lock.as_ref() {
connection.get_random_connection(Some(connections_key), Some(&eligible_ips))
} else {
None
}
}
.map(|(ip, port)| format!("{}:{}", binary_to_ip(ip), port))
}
pub async fn request_node(db: &Db, connections_key: &str) -> RpcResponse {
match eligible_sponsor_endpoint(db, connections_key).await {
Some(endpoint) => {
let Some((ip, port)) = endpoint.rsplit_once(':') else {
return RpcResponse::Binary(b"Error: No Eligible Sponsor Found".to_vec());
};
let Ok(port) = port.parse::<u16>() else {
return RpcResponse::Binary(b"Error: No Eligible Sponsor Found".to_vec());
};
let mut ip = crate::common::binary_conversions::ip_to_binary(ip);
ip.extend(port.to_le_bytes());
RpcResponse::Binary(ip)
}
None => {
let msg = "Error: No Eligible Sponsor Found"
.to_string()
.as_bytes()
.to_vec();
RpcResponse::Binary(msg)
}
}
}

View File

@ -8,7 +8,9 @@ use crate::orphans::checkup_state::{
use crate::orphans::structs::OrphanCheckup2; use crate::orphans::structs::OrphanCheckup2;
use crate::orphans::sync_check::sync_checkup; use crate::orphans::sync_check::sync_checkup;
use crate::records::block_height::get_block_height::get_height; use crate::records::block_height::get_block_height::get_height;
use crate::records::memory::connections::{get_client_type_from_memory, peer_is_operational}; use crate::records::memory::connections::{
get_client_type_from_memory, peer_catch_up_target, peer_is_operational,
};
use crate::records::memory::enums::ClientType; use crate::records::memory::enums::ClientType;
use crate::records::memory::response_channels::Command; use crate::records::memory::response_channels::Command;
use crate::rpc::client::syncing::node_syncing; use crate::rpc::client::syncing::node_syncing;
@ -422,6 +424,19 @@ pub async fn receive_torrent(
if get_client_type_from_memory(connections_key).await == Some(ClientType::Miner) if get_client_type_from_memory(connections_key).await == Some(ClientType::Miner)
&& !peer_is_operational(connections_key).await && !peer_is_operational(connections_key).await
{ {
if peer_catch_up_target(connections_key).await.is_none() {
warn!(
"[broadcast] ignored torrent from non-operational peer outside authenticated catch-up state: peer={connections_key} height={block_number}"
);
return Ok((
uid,
RpcResponse::Binary(
"Torrent ignored from incomplete peer setup."
.as_bytes()
.to_vec(),
),
));
}
let local_height = get_height(db); let local_height = get_height(db);
if !within_orphan_window(local_height, block_number) { if !within_orphan_window(local_height, block_number) {
if is_syncing_mode() if is_syncing_mode()

View File

@ -1,10 +1,13 @@
use crate::io::ErrorKind; use crate::io::ErrorKind;
use crate::log::warn; use crate::log::warn;
use crate::records::memory::connections::{get_client_type_from_memory, peer_accepts_live_relay}; use crate::records::memory::connections::{
get_client_type_from_memory, peer_accepts_live_relay, peer_accepts_mapping_relay,
};
use crate::records::memory::enums::ClientType; use crate::records::memory::enums::ClientType;
use crate::rpc::command_maps::{ use crate::rpc::command_maps::{
RPC_BLOCK_HEIGHT, RPC_BLOCK_PIECE, RPC_NETWORK_MONITOR_ADD, RPC_NETWORK_MONITOR_REMOVE, RPC_BLOCK_HEIGHT, RPC_BLOCK_PIECE, RPC_NETWORK_MEMBERSHIP_RECONCILE,
RPC_REPLY, RPC_SUBMIT_TORRENT, RPC_SUBMIT_TRANSACTION, RPC_TORRENT_BY_HEIGHT, RPC_NETWORK_MONITOR_ADD, RPC_NETWORK_MONITOR_REMOVE, RPC_REPLY, RPC_SUBMIT_TORRENT,
RPC_SUBMIT_TRANSACTION, RPC_TORRENT_BY_HEIGHT,
}; };
use crate::rpc::server::connection_memory_manager::remove_stream_from_memory; use crate::rpc::server::connection_memory_manager::remove_stream_from_memory;
use crate::rpc::server::flood_protection::check_request_frequency_with_client_type; use crate::rpc::server::flood_protection::check_request_frequency_with_client_type;
@ -80,6 +83,7 @@ fn requires_operational_miner(command: u8) -> bool {
| RPC_SUBMIT_TORRENT | RPC_SUBMIT_TORRENT
| RPC_NETWORK_MONITOR_ADD | RPC_NETWORK_MONITOR_ADD
| RPC_NETWORK_MONITOR_REMOVE | RPC_NETWORK_MONITOR_REMOVE
| RPC_NETWORK_MEMBERSHIP_RECONCILE
) )
} }
@ -102,10 +106,17 @@ pub async fn next_incoming_command(
.await .await
.unwrap_or(ClientType::Miner); .unwrap_or(ClientType::Miner);
if client_type == ClientType::Miner let accepts_relay = if matches!(
&& requires_operational_miner(command) command,
&& !peer_accepts_live_relay(connections_key).await RPC_NETWORK_MONITOR_ADD
{ | RPC_NETWORK_MONITOR_REMOVE
| RPC_NETWORK_MEMBERSHIP_RECONCILE
) {
peer_accepts_mapping_relay(connections_key).await
} else {
peer_accepts_live_relay(connections_key).await
};
if client_type == ClientType::Miner && requires_operational_miner(command) && !accepts_relay {
warn!( warn!(
"[rpc] setup-time relay received before peer became operational: peer={connections_key} cmd={command}" "[rpc] setup-time relay received before peer became operational: peer={connections_key} cmd={command}"
); );

View File

@ -2,8 +2,9 @@ use crate::records::ip_score::enums::InfractionType;
use crate::records::ip_score::score::update_ip_score; use crate::records::ip_score::score::update_ip_score;
use crate::records::memory::enums::ClientType; use crate::records::memory::enums::ClientType;
use crate::rpc::command_maps::{ use crate::rpc::command_maps::{
RPC_ADD_NETWORK_NODE, RPC_NETWORK_MONITOR_ADD, RPC_NETWORK_MONITOR_REMOVE, RPC_ADD_NETWORK_NODE, RPC_NETWORK_MAPPING_HASH, RPC_NETWORK_MONITOR_ADD,
RPC_NETWORK_MONITOR_STATE, RPC_SETUP_COMPLETE, RPC_NETWORK_MEMBERSHIP_RECONCILE, RPC_NETWORK_MONITOR_REMOVE, RPC_NETWORK_MONITOR_STATE,
RPC_SETUP_COMPLETE,
}; };
use crate::rpc::server::structs::RpcFloodState; use crate::rpc::server::structs::RpcFloodState;
use crate::sled::Db; use crate::sled::Db;
@ -33,6 +34,8 @@ fn request_subject(ip: &str, client_type: ClientType, command: u8) -> String {
| RPC_NETWORK_MONITOR_ADD | RPC_NETWORK_MONITOR_ADD
| RPC_NETWORK_MONITOR_REMOVE | RPC_NETWORK_MONITOR_REMOVE
| RPC_NETWORK_MONITOR_STATE | RPC_NETWORK_MONITOR_STATE
| RPC_NETWORK_MAPPING_HASH
| RPC_NETWORK_MEMBERSHIP_RECONCILE
| RPC_SETUP_COMPLETE | RPC_SETUP_COMPLETE
) { ) {
return format!("miner:{ip}:control:{command}"); return format!("miner:{ip}:control:{command}");

View File

@ -7,8 +7,8 @@ use crate::orphans::sync_check::sync_checkup;
use crate::orphans::torrent_candidates::hydrate_torrent_candidates; use crate::orphans::torrent_candidates::hydrate_torrent_candidates;
use crate::records::block_height::get_block_height::get_height; use crate::records::block_height::get_block_height::get_height;
use crate::records::memory::connections::{ use crate::records::memory::connections::{
mark_peer_network_map_synced, mark_peer_operational, mark_peer_wallet_registry_synced, mark_peer_catching_up, mark_peer_network_map_synced, mark_peer_operational,
spawn_peer_setup_retry, mark_peer_wallet_registry_synced, spawn_peer_setup_retry,
}; };
use crate::records::memory::enums::ClientType; use crate::records::memory::enums::ClientType;
use crate::records::memory::response_channels::generate_uid; use crate::records::memory::response_channels::generate_uid;
@ -40,6 +40,23 @@ use crate::Settings;
use crate::TcpStream; use crate::TcpStream;
use crate::Utc; use crate::Utc;
const LIVE_ORPHAN_WINDOW: u32 = 10;
const PASSIVE_CATCH_UP_BLOCK_LIMIT: u32 = 20;
fn peer_reached_catch_up_target(
local_height: u32,
remote_height: u32,
catch_up_target: u32,
) -> bool {
remote_height >= catch_up_target
&& local_height.saturating_sub(remote_height) <= LIVE_ORPHAN_WINDOW
}
fn passive_catch_up_expired(local_height: u32, remote_height: u32, catch_up_target: u32) -> bool {
local_height.saturating_sub(catch_up_target) > PASSIVE_CATCH_UP_BLOCK_LIMIT
&& local_height.saturating_sub(remote_height) > LIVE_ORPHAN_WINDOW
}
async fn drop_failed_handshake(stream: &Arc<Mutex<TcpStream>>) { async fn drop_failed_handshake(stream: &Arc<Mutex<TcpStream>>) {
// Failed handshakes are never stored in connection memory, but the // Failed handshakes are never stored in connection memory, but the
// accepted TCP socket should still be closed immediately. // accepted TCP socket should still be closed immediately.
@ -77,7 +94,7 @@ async fn sync_incoming_peer_before_operational(
wallet: Arc<Wallet>, wallet: Arc<Wallet>,
map: Arc<Mutex<Command>>, map: Arc<Mutex<Command>>,
connections_key: &str, connections_key: &str,
) -> Result<(bool, Option<ChainOperationGuard>), String> { ) -> Result<(bool, Option<ChainOperationGuard>, Option<u32>), String> {
let initial_local_height = get_height(db); let initial_local_height = get_height(db);
let initial_remote_height = let initial_remote_height =
request_remote_height(stream.clone(), map.clone(), connections_key.to_string()).await?; request_remote_height(stream.clone(), map.clone(), connections_key.to_string()).await?;
@ -96,7 +113,7 @@ async fn sync_incoming_peer_before_operational(
warn!( 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}" "[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)); return Ok((false, None, Some(initial_local_height)));
} }
// Wait for any existing startup, live catch-up, or orphan operation to // Wait for any existing startup, live catch-up, or orphan operation to
@ -118,7 +135,7 @@ async fn sync_incoming_peer_before_operational(
warn!( warn!(
"[startup] incoming peer fell behind while waiting for chain access; keeping connection passive: 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, None)); return Ok((false, None, Some(local_height)));
} }
if !local_genesis_exists || remote_height > local_height + 10 { if !local_genesis_exists || remote_height > local_height + 10 {
@ -186,7 +203,7 @@ async fn sync_incoming_peer_before_operational(
.map_err(|err| format!("Incoming governance state sync error: {err}"))?; .map_err(|err| format!("Incoming governance state sync error: {err}"))?;
} }
Ok((true, Some(chain_sync_guard))) Ok((true, Some(chain_sync_guard), None))
} }
fn spawn_incoming_peer_promotion_watcher( fn spawn_incoming_peer_promotion_watcher(
@ -194,6 +211,7 @@ fn spawn_incoming_peer_promotion_watcher(
db: Db, db: Db,
map: Arc<Mutex<Command>>, map: Arc<Mutex<Command>>,
connections_key: String, connections_key: String,
catch_up_target: u32,
) { ) {
tokio::spawn(async move { tokio::spawn(async move {
loop { loop {
@ -215,15 +233,46 @@ fn spawn_incoming_peer_promotion_watcher(
Err(_) => continue, Err(_) => continue,
}; };
if remote_height >= local_height { if peer_reached_catch_up_target(local_height, remote_height, catch_up_target) {
if mark_peer_operational(&connections_key, map.clone()).await { if mark_peer_operational(&connections_key, map.clone()).await {
break; break;
} }
continue;
}
if passive_catch_up_expired(local_height, remote_height, catch_up_target) {
warn!(
"[startup] incoming peer failed to catch up before passive limit; disconnecting: peer={connections_key} target_height={catch_up_target} local_height={local_height} remote_height={remote_height}"
);
remove_stream_from_memory(&stream).await;
break;
} }
} }
}); });
} }
#[cfg(test)]
mod tests {
use super::{passive_catch_up_expired, peer_reached_catch_up_target};
#[test]
fn moving_tip_does_not_move_the_original_catch_up_target() {
assert!(peer_reached_catch_up_target(66_092, 66_091, 66_091));
}
#[test]
fn peer_must_reach_original_target_and_remain_near_tip() {
assert!(!peer_reached_catch_up_target(66_091, 66_090, 66_091));
assert!(!peer_reached_catch_up_target(66_120, 66_091, 66_091));
}
#[test]
fn stale_passive_peer_expires_only_after_falling_outside_live_window() {
assert!(!passive_catch_up_expired(66_112, 66_102, 66_091));
assert!(passive_catch_up_expired(66_112, 66_100, 66_091));
}
}
async fn complete_incoming_miner_setup( async fn complete_incoming_miner_setup(
stream: Arc<Mutex<TcpStream>>, stream: Arc<Mutex<TcpStream>>,
db: &Db, db: &Db,
@ -319,7 +368,6 @@ async fn complete_incoming_miner_setup(
&short_address, &short_address,
map.clone(), map.clone(),
db, db,
wallet.clone(),
connections_key, connections_key,
false, false,
) )
@ -333,7 +381,6 @@ async fn complete_incoming_miner_setup(
stream.clone(), stream.clone(),
map.clone(), map.clone(),
db, db,
wallet.clone(),
connections_key, connections_key,
&short_address, &short_address,
) )
@ -346,22 +393,23 @@ async fn complete_incoming_miner_setup(
} }
mark_peer_network_map_synced(connections_key).await; mark_peer_network_map_synced(connections_key).await;
let (operational, chain_sync_guard) = match sync_incoming_peer_before_operational( let (operational, chain_sync_guard, catch_up_target) =
stream.clone(), match sync_incoming_peer_before_operational(
db, stream.clone(),
wallet.clone(), db,
map.clone(), wallet.clone(),
connections_key, map.clone(),
) connections_key,
.await )
{ .await
Ok(operational) => operational, {
Err(err) => { Ok(operational) => operational,
error!("[startup] incoming peer chain sync failed: {err}"); Err(err) => {
remove_stream_from_memory(&stream).await; error!("[startup] incoming peer chain sync failed: {err}");
return; remove_stream_from_memory(&stream).await;
} return;
}; }
};
if operational { if operational {
if !mark_peer_operational(connections_key, map.clone()).await { if !mark_peer_operational(connections_key, map.clone()).await {
@ -372,11 +420,26 @@ async fn complete_incoming_miner_setup(
guard.finish(); guard.finish();
} }
} else { } else {
let Some(catch_up_target) = catch_up_target else {
warn!(
"[startup] incoming peer entered passive setup without a catch-up target; disconnecting: peer={connections_key}"
);
remove_stream_from_memory(&stream).await;
return;
};
if !mark_peer_catching_up(connections_key, catch_up_target).await {
warn!(
"[startup] incoming peer could not enter catch-up state; disconnecting: peer={connections_key} target_height={catch_up_target}"
);
remove_stream_from_memory(&stream).await;
return;
}
spawn_incoming_peer_promotion_watcher( spawn_incoming_peer_promotion_watcher(
stream.clone(), stream.clone(),
db.clone(), db.clone(),
map.clone(), map.clone(),
connections_key.to_string(), connections_key.to_string(),
catch_up_target,
); );
} }
} }

View File

@ -549,6 +549,8 @@ pub async fn start_loop(
&db, &db,
wallet.clone(), wallet.clone(),
map.clone(), map.clone(),
client_type,
false,
) )
.await?; .await?;
let should_drop_rejected_miner = client_type == ClientType::Miner let should_drop_rejected_miner = client_type == ClientType::Miner
@ -561,6 +563,21 @@ pub async fn start_loop(
break 'outer Ok(()); break 'outer Ok(());
} }
} }
59 => {
let (uid, result) = commands::add_network_node::add_network_node(
&connections_key,
stream_locked.clone(),
&db,
wallet.clone(),
map.clone(),
client_type,
true,
)
.await?;
result
.send(&stream_locked, Some(&connections_key), uid)
.await;
}
50 => { 50 => {
// signed monitor-add event from a miner peer // signed monitor-add event from a miner peer
let (uid, result) = commands::network_monitor_add::network_monitor_add( let (uid, result) = commands::network_monitor_add::network_monitor_add(
@ -603,6 +620,18 @@ pub async fn start_loop(
.send(&stream_locked, Some(&connections_key), uid) .send(&stream_locked, Some(&connections_key), uid)
.await; .await;
} }
58 => {
// Return a deterministic digest of shared network-map state.
let (uid, _) = read_bytes_from_stream::read_uid_from_stream(
&connections_key,
stream_locked.clone(),
)
.await?;
let result = commands::network_mapping_hash::network_mapping_hash().await;
result
.send(&stream_locked, Some(&connections_key), uid)
.await;
}
55 => { 55 => {
// Both sides announce setup completion independently. Receipt // Both sides announce setup completion independently. Receipt
// records the remote half and the response acknowledges it. // records the remote half and the response acknowledges it.

View File

@ -9,11 +9,12 @@ use crate::rpc::command_maps::{
RPC_CONTRACT_BY_ADDRESS, RPC_DIFFICULTY, RPC_GOVERNANCE_PROPOSAL, RPC_LARGEST_TX_FEE, RPC_CONTRACT_BY_ADDRESS, RPC_DIFFICULTY, RPC_GOVERNANCE_PROPOSAL, RPC_LARGEST_TX_FEE,
RPC_LATEST_ADDRESS_TRANSACTIONS, RPC_LOAN_CONTRACT, RPC_MARKETING_CAMPAIGN_HISTORY, RPC_LATEST_ADDRESS_TRANSACTIONS, RPC_LOAN_CONTRACT, RPC_MARKETING_CAMPAIGN_HISTORY,
RPC_MEMPOOL_TX_BY_ADDRESS, RPC_MEMPOOL_TX_BY_SIGNATURE, RPC_MEMPOOL_TX_COUNT, RPC_NETWORK_INFO, RPC_MEMPOOL_TX_BY_ADDRESS, RPC_MEMPOOL_TX_BY_SIGNATURE, RPC_MEMPOOL_TX_COUNT, RPC_NETWORK_INFO,
RPC_NFT_DETAILS, RPC_NFT_LIST, RPC_REGISTER_WALLET, RPC_STORAGE_LOOKUP, RPC_NETWORK_MONITOR_STATE, RPC_NFT_DETAILS, RPC_NFT_LIST, RPC_REGISTER_WALLET,
RPC_STORAGE_LOOKUP_COST, RPC_TIME, RPC_TOKEN_CATALOG, RPC_TOKEN_DETAILS, RPC_TOKEN_LIST, RPC_REQUEST_NODE_LIST, RPC_STORAGE_LOOKUP, RPC_STORAGE_LOOKUP_COST, RPC_TIME,
RPC_TORRENT_BY_HEIGHT, RPC_TOTAL_CONFIRMED_TX, RPC_TRANSACTION_BY_TXID, RPC_UNBLOCK_IP, RPC_TOKEN_CATALOG, RPC_TOKEN_DETAILS, RPC_TOKEN_LIST, RPC_TORRENT_BY_HEIGHT,
RPC_VALIDATE_MESSAGE, RPC_VANITY_LOOKUP, RPC_VANITY_OWNER_LOOKUP, RPC_TOTAL_CONFIRMED_TX, RPC_TRANSACTION_BY_TXID, RPC_UNBLOCK_IP, RPC_VALIDATE_MESSAGE,
RPC_WALLET_REGISTRATION_STATUS, RPC_WALLET_REGISTRY_SYNC, RPC_VANITY_LOOKUP, RPC_VANITY_OWNER_LOOKUP, RPC_WALLET_REGISTRATION_STATUS,
RPC_WALLET_REGISTRY_SYNC,
}; };
use crate::standalone_tools::transaction_creator::create_transaction_request; use crate::standalone_tools::transaction_creator::create_transaction_request;
use crate::standalone_tools::vanity_resolver::canonical_vanity_address_input; use crate::standalone_tools::vanity_resolver::canonical_vanity_address_input;
@ -357,6 +358,11 @@ async fn build_request_bytes(
bin_msg.extend_from_slice(&ip_bytes); bin_msg.extend_from_slice(&ip_bytes);
bin_msg.extend_from_slice(&signature_bytes); bin_msg.extend_from_slice(&signature_bytes);
} }
// Request the signed network membership snapshot.
30 => {
bin_msg.push(RPC_REQUEST_NODE_LIST);
bin_msg.extend_from_slice(&hashmap_key);
}
// Token-list request. // Token-list request.
31 => { 31 => {
let command_number: u8 = RPC_TOKEN_LIST; let command_number: u8 = RPC_TOKEN_LIST;
@ -771,6 +777,11 @@ async fn build_request_bytes(
bin_msg.extend_from_slice(&address_bytes); bin_msg.extend_from_slice(&address_bytes);
bin_msg.extend_from_slice(&transfer_tx); bin_msg.extend_from_slice(&transfer_tx);
} }
// Request the latest signed monitor relationship state.
54 => {
bin_msg.push(RPC_NETWORK_MONITOR_STATE);
bin_msg.extend_from_slice(&hashmap_key);
}
// Look up complete governance state by its 32-byte proposal key. // Look up complete governance state by its 32-byte proposal key.
57 => { 57 => {
let command_number: u8 = RPC_GOVERNANCE_PROPOSAL; let command_number: u8 = RPC_GOVERNANCE_PROPOSAL;
@ -804,7 +815,10 @@ async fn build_request_bytes(
#[cfg(test)] #[cfg(test)]
mod tests { mod tests {
use super::build_request_bytes; use super::build_request_bytes;
use crate::rpc::command_maps::{RPC_GOVERNANCE_PROPOSAL, RPC_VALIDATE_MESSAGE}; use crate::rpc::command_maps::{
RPC_GOVERNANCE_PROPOSAL, RPC_NETWORK_MONITOR_STATE, RPC_REQUEST_NODE_LIST,
RPC_VALIDATE_MESSAGE,
};
use crate::wallets::structures::Wallet; use crate::wallets::structures::Wallet;
#[tokio::test] #[tokio::test]
@ -843,4 +857,15 @@ mod tests {
assert_eq!(&bytes[1..4], &uid); assert_eq!(&bytes[1..4], &uid);
assert_eq!(&bytes[4..], proposal_key); assert_eq!(&bytes[4..], proposal_key);
} }
#[tokio::test]
async fn network_map_snapshot_requests_use_uid_only_wire_layout() {
let uid = [7, 8, 9];
let node_list = build_request_bytes(String::new(), 30, uid).await.unwrap();
assert_eq!(node_list, [RPC_REQUEST_NODE_LIST, 7, 8, 9]);
let monitor_state = build_request_bytes(String::new(), 54, uid).await.unwrap();
assert_eq!(monitor_state, [RPC_NETWORK_MONITOR_STATE, 7, 8, 9]);
}
} }

View File

@ -1,12 +1,11 @@
use crate::common::network_startup::get_node_connections; use crate::common::network_startup::{get_ip_and_port, get_node_connections};
use crate::log::{error, info, warn}; use crate::log::{error, info, warn};
use crate::miner::flag::{ use crate::miner::flag::{is_mining_stop_requested, is_normal_mode};
clear_mining_stop_request, is_normal_mode, set_mining_state, set_node_mode, MiningState,
NodeMode,
};
use crate::records::memory::connections::{ use crate::records::memory::connections::{
peer_connection_count, ready_outgoing_connection_count, refill_outgoing_connections_once, operational_peer_wallets, peer_connection_count, ready_outgoing_connection_count,
refill_outgoing_connections_once,
}; };
use crate::records::memory::network_mapping::NodeInfo;
use crate::records::memory::response_channels::Command; use crate::records::memory::response_channels::Command;
use crate::rpc::client::handshake::connect_and_handshake; use crate::rpc::client::handshake::connect_and_handshake;
use crate::rpc::client::structs::Connect; use crate::rpc::client::structs::Connect;
@ -22,29 +21,38 @@ pub async fn handle_connections(
wallet: Arc<Wallet>, wallet: Arc<Wallet>,
map: Arc<Mutex<Command>>, map: Arc<Mutex<Command>>,
) -> Result<(), String> { ) -> Result<(), String> {
// A zero outgoing limit means this node should not open any bootstrap // A zero outgoing limit means this node waits for an incoming sponsor. It
// connection during startup. // does not bypass the independent-peer requirement or permit solo mining.
let outgoing_connections = crate::Settings::load() let outgoing_connections = crate::Settings::load()
.map(|settings| settings.outgoing_connections) .map(|settings| settings.outgoing_connections)
.unwrap_or(0); .unwrap_or(0);
if outgoing_connections == 0 { if outgoing_connections == 0 {
info!("OUTGOING_CONNECTIONS is 0; skipping startup bootstrap."); info!("OUTGOING_CONNECTIONS is 0; waiting for an incoming sponsored peer.");
set_node_mode(NodeMode::Normal);
clear_mining_stop_request();
set_mining_state(MiningState::Idle);
return Ok(());
} }
let connected = attempt_bootstrap_connections(db, wallet, map, "startup").await?; loop {
if connected { if outgoing_connections > 0
return Ok(()); && attempt_bootstrap_connections(db.clone(), wallet.clone(), map.clone(), "startup")
} .await?
{
return Ok(());
}
// Startup can continue as a standalone node even if no bootstrap peer is reachable. info!("No existing network peer is available. Waiting for connections.");
set_node_mode(NodeMode::Normal); for _ in 0..60 {
clear_mining_stop_request(); let peer_wallets = operational_peer_wallets().await;
set_mining_state(MiningState::Idle); if is_normal_mode()
Ok(()) && !is_mining_stop_requested()
&& NodeInfo::has_reciprocal_sponsorship(&wallet.saved.short_address, &peer_wallets)
.await
{
info!("Reciprocal node sponsorship completed; startup may continue.");
return Ok(());
}
sleep(Duration::from_secs(1)).await;
}
warn!("No sponsored peer is operational; retrying configured bootstrap peers");
}
} }
async fn attempt_bootstrap_connections( async fn attempt_bootstrap_connections(
@ -56,9 +64,15 @@ async fn attempt_bootstrap_connections(
// Try the configured bootstrap peers one by one until a // Try the configured bootstrap peers one by one until a
// handshake succeeds or the list is exhausted. // handshake succeeds or the list is exhausted.
let filtered_servers = get_node_connections().await; let filtered_servers = get_node_connections().await;
let (_, _, local_endpoint) = get_ip_and_port().await;
let mut last_error: Option<String> = None; let mut last_error: Option<String> = None;
for server in filtered_servers { for server in filtered_servers {
// A node can never sponsor itself. Ignore its own public endpoint and
// wait for a genuinely independent incoming or outgoing peer.
if server == local_endpoint {
continue;
}
// build the outbound handshake request using cloned // build the outbound handshake request using cloned
// shared state so each attempt can run independently // shared state so each attempt can run independently
let db_clone = db.clone(); let db_clone = db.clone();

View File

@ -209,9 +209,6 @@ pub fn install_shutdown_cleanup(db: Db) {
} }
} }
crate::records::memory::network_mapping::NodeInfo::persist_recovery_snapshot("shutdown")
.await;
if let Err(err) = db.flush_async().await { if let Err(err) = db.flush_async().await {
error!("Failed to flush sled during shutdown: {err}"); error!("Failed to flush sled during shutdown: {err}");
} }

View File

@ -2,7 +2,8 @@ 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::common::network_startup::get_ip_and_port;
use crate::log::warn; use crate::log::warn;
use crate::records::memory::network_mapping::structs::{ use crate::records::memory::network_mapping::structs::{
SignedNodeEdit, MONITOR_EVENT_BYTES, NODE_ADDED_BY_OFFSET, NODE_ADDED_SIGNATURE_OFFSET, SignedNodeEdit, SyncedNodeState, NETWORK_SNAPSHOT_HEADER_BYTES, NETWORK_SNAPSHOT_MAGIC,
NETWORK_SNAPSHOT_VERSION, NODE_ADDED_BY_OFFSET, NODE_ADDED_SIGNATURE_OFFSET,
NODE_ADDED_TIMESTAMP_OFFSET, NODE_BLOCKS_MINED_OFFSET, NODE_DELETED_BLOCK_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_DELETED_TIMESTAMP_OFFSET, NODE_IP_OFFSET, NODE_MONITOR_COUNT_OFFSET, NODE_PORT_OFFSET,
NODE_RECORD_FIXED_BYTES, NODE_RECORD_FIXED_BYTES,
@ -10,7 +11,8 @@ use crate::records::memory::network_mapping::structs::{
use crate::records::memory::network_mapping::NodeInfo; use crate::records::memory::network_mapping::NodeInfo;
use crate::records::memory::response_channels::{reserve_entry, Command}; use crate::records::memory::response_channels::{reserve_entry, Command};
use crate::rpc::command_maps::{ use crate::rpc::command_maps::{
RPC_ADD_NETWORK_NODE, RPC_NETWORK_MONITOR_STATE, RPC_REQUEST_NODE_LIST, RPC_ADD_NETWORK_NODE, RPC_NETWORK_MAPPING_HASH, RPC_NETWORK_MEMBERSHIP_RECONCILE,
RPC_NETWORK_MONITOR_ADD, RPC_NETWORK_MONITOR_REMOVE, RPC_REQUEST_NODE_LIST,
}; };
use crate::rpc::responses::RpcResponse; use crate::rpc::responses::RpcResponse;
use crate::sled::Db; use crate::sled::Db;
@ -22,12 +24,42 @@ use crate::Mutex;
use crate::TcpStream; use crate::TcpStream;
use crate::Utc; use crate::Utc;
const REQUESTED_MAPPING_RECORD_PENDING: &str =
"requested network mapping address was not present in the peer snapshot";
pub async fn announce_self_to_network( pub async fn announce_self_to_network(
unlocked_stream: Arc<Mutex<TcpStream>>, unlocked_stream: Arc<Mutex<TcpStream>>,
address: &str, address: &str,
command_map: Arc<Mutex<Command>>, command_map: Arc<Mutex<Command>>,
db: &Db, db: &Db,
wallet: Arc<Wallet>, connections_key: &str,
import_peer_map: bool,
) -> Result<(), String> {
if import_peer_map {
NodeInfo::begin_mapping_snapshot_sync().await;
}
let result = announce_self_to_network_inner(
unlocked_stream,
address,
command_map,
db,
connections_key,
import_peer_map,
)
.await;
if import_peer_map {
NodeInfo::finish_mapping_snapshot_sync().await;
}
result
}
async fn announce_self_to_network_inner(
unlocked_stream: Arc<Mutex<TcpStream>>,
address: &str,
command_map: Arc<Mutex<Command>>,
db: &Db,
connections_key: &str, connections_key: &str,
import_peer_map: bool, import_peer_map: bool,
) -> Result<(), String> { ) -> Result<(), String> {
@ -96,12 +128,12 @@ pub async fn announce_self_to_network(
} }
if import_peer_map { if import_peer_map {
get_network_mapping( get_network_mapping_response(
unlocked_stream, unlocked_stream,
command_map.clone(), command_map.clone(),
db, db,
wallet.clone(),
connections_key, connections_key,
None,
) )
.await?; .await?;
} }
@ -112,44 +144,209 @@ pub async fn get_network_mapping(
unlocked_stream: Arc<Mutex<TcpStream>>, unlocked_stream: Arc<Mutex<TcpStream>>,
command_map: Arc<Mutex<Command>>, command_map: Arc<Mutex<Command>>,
db: &Db, db: &Db,
wallet: Arc<Wallet>,
connections_key: &str, connections_key: &str,
) -> Result<(), String> { ) -> Result<(), String> {
get_network_mapping_inner( get_network_mapping_inner(unlocked_stream, command_map, db, connections_key, None).await
unlocked_stream, }
command_map,
db, pub async fn compare_network_mapping_digest(
wallet, unlocked_stream: Arc<Mutex<TcpStream>>,
connections_key, command_map: Arc<Mutex<Command>>,
None, connections_key: &str,
) ) -> Result<bool, String> {
.await let (uid, _tx, rx) = reserve_entry(command_map).await;
let mut request = Vec::with_capacity(4);
request.push(RPC_NETWORK_MAPPING_HASH);
request.extend_from_slice(&uid);
RpcResponse::send_raw(&unlocked_stream, Some(connections_key), &request).await;
let mut receiver = rx.lock().await;
let remote = timeout(Duration::from_secs(30), receiver.recv())
.await
.map_err(|_| "timed out waiting for network-map digest".to_string())?
.ok_or_else(|| "network-map digest response channel closed".to_string())?;
if remote.len() != 32 {
return Err(format!(
"peer returned an invalid network-map digest length: {}",
remote.len()
));
}
let local = NodeInfo::canonical_mapping_digest().await?;
Ok(local == remote)
}
async fn send_mapping_record(
unlocked_stream: &Arc<Mutex<TcpStream>>,
command_map: Arc<Mutex<Command>>,
connections_key: &str,
command: u8,
body: &[u8],
) -> Result<(), String> {
let (uid, _tx, rx) = reserve_entry(command_map).await;
let mut request = Vec::with_capacity(4 + body.len());
request.push(command);
request.extend_from_slice(&uid);
request.extend_from_slice(body);
RpcResponse::send_raw(unlocked_stream, Some(connections_key), &request).await;
let mut receiver = rx.lock().await;
let response = timeout(Duration::from_secs(30), receiver.recv())
.await
.map_err(|_| "timed out while reconciling network-map record".to_string())?
.ok_or_else(|| "network-map reconciliation response channel closed".to_string())?;
let response = binary_to_string(response);
if response != "Success" {
return Err(format!("peer rejected reconciled mapping record: {response}"));
}
Ok(())
}
pub async fn reconcile_network_mapping_with_peer(
unlocked_stream: Arc<Mutex<TcpStream>>,
command_map: Arc<Mutex<Command>>,
db: &Db,
connections_key: &str,
) -> Result<(), String> {
let (memberships, monitor_events) = NodeInfo::signed_mapping_state().await;
for state in memberships {
let edit = state.edit;
let address = Wallet::short_address_to_bytes(&edit.address)
.ok_or_else(|| "local mapping contained an invalid address".to_string())?;
let modified_by = Wallet::short_address_to_bytes(&edit.modified_by)
.ok_or_else(|| "local mapping contained an invalid signer".to_string())?;
let signature = crate::decode(&edit.modified_signature)
.map_err(|_| "local mapping contained an invalid membership signature".to_string())?;
if signature.len() != Wallet::SIGNATURE_LENGTH {
return Err("local mapping contained an invalid membership signature length".to_string());
}
let mut body = Vec::new();
body.extend_from_slice(&address);
body.extend_from_slice(&ip_to_binary(&edit.ip));
body.extend_from_slice(&edit.port.to_le_bytes());
body.extend_from_slice(&modified_by);
body.extend_from_slice(&edit.modified_timestamp.to_le_bytes());
body.extend_from_slice(&signature);
body.extend_from_slice(&0_u16.to_le_bytes());
body.extend_from_slice(&state.deleted_timestamp.to_le_bytes());
body.extend_from_slice(&state.deleted_block.to_le_bytes());
send_mapping_record(
&unlocked_stream,
command_map.clone(),
connections_key,
RPC_NETWORK_MEMBERSHIP_RECONCILE,
&body,
)
.await?;
}
for edit in monitor_events {
let monitored = Wallet::short_address_to_bytes(&edit.monitored_address)
.ok_or_else(|| "local monitor state contained an invalid target".to_string())?;
let monitoring = Wallet::short_address_to_bytes(&edit.monitoring_address)
.ok_or_else(|| "local monitor state contained an invalid signer".to_string())?;
let signature = crate::decode(&edit.modified_signature)
.map_err(|_| "local monitor state contained an invalid signature".to_string())?;
if signature.len() != Wallet::SIGNATURE_LENGTH {
return Err("local monitor state contained an invalid signature length".to_string());
}
let command = match edit.action {
crate::records::memory::network_mapping::monitor::MONITOR_ACTION_ADD => {
RPC_NETWORK_MONITOR_ADD
}
crate::records::memory::network_mapping::monitor::MONITOR_ACTION_REMOVE => {
RPC_NETWORK_MONITOR_REMOVE
}
_ => return Err("local monitor state contained an invalid action".to_string()),
};
let mut body = Vec::new();
body.push(edit.action);
body.extend_from_slice(&monitored);
body.extend_from_slice(&monitoring);
body.extend_from_slice(&ip_to_binary(&edit.target_ip));
body.extend_from_slice(&edit.modified_timestamp.to_le_bytes());
body.extend_from_slice(&edit.modified_block.to_le_bytes());
body.extend_from_slice(&signature);
send_mapping_record(
&unlocked_stream,
command_map.clone(),
connections_key,
command,
&body,
)
.await?;
}
// The peer has now merged every valid current signed record we hold. Pull
// its resulting complete state so records that only it possessed are also
// installed locally and derived deletion cascades are identical.
get_network_mapping(unlocked_stream, command_map, db, connections_key).await
} }
pub async fn get_network_mapping_for_address( pub async fn get_network_mapping_for_address(
unlocked_stream: Arc<Mutex<TcpStream>>, unlocked_stream: Arc<Mutex<TcpStream>>,
command_map: Arc<Mutex<Command>>, command_map: Arc<Mutex<Command>>,
db: &Db, db: &Db,
wallet: Arc<Wallet>,
connections_key: &str, connections_key: &str,
only_address: &str, only_address: &str,
) -> Result<(), String> { ) -> Result<(), String> {
get_network_mapping_inner( // During the first two-node bootstrap, the peer may acknowledge our
unlocked_stream, // reverse announcement while deferring its application until an in-flight
command_map, // full mapping snapshot finishes. Wait for that exact sponsored record to
db, // become visible instead of treating an absent record as a successful sync.
wallet, for attempt in 1..=40 {
connections_key, match get_network_mapping_inner(
Some(only_address), unlocked_stream.clone(),
) command_map.clone(),
.await db,
connections_key,
Some(only_address),
)
.await
{
Ok(()) => return Ok(()),
Err(err) if err == REQUESTED_MAPPING_RECORD_PENDING && attempt < 40 => {
crate::sleep(crate::Duration::from_millis(250)).await;
}
Err(err) => return Err(err),
}
}
Err(REQUESTED_MAPPING_RECORD_PENDING.to_string())
} }
async fn get_network_mapping_inner( async fn get_network_mapping_inner(
unlocked_stream: Arc<Mutex<TcpStream>>, unlocked_stream: Arc<Mutex<TcpStream>>,
command_map: Arc<Mutex<Command>>, command_map: Arc<Mutex<Command>>,
db: &Db, db: &Db,
wallet: Arc<Wallet>, connections_key: &str,
only_address: Option<&str>,
) -> Result<(), String> {
if only_address.is_none() {
NodeInfo::begin_mapping_snapshot_sync().await;
}
let result = get_network_mapping_response(
unlocked_stream,
command_map,
db,
connections_key,
only_address,
)
.await;
if only_address.is_none() {
NodeInfo::finish_mapping_snapshot_sync().await;
}
result
}
async fn get_network_mapping_response(
unlocked_stream: Arc<Mutex<TcpStream>>,
command_map: Arc<Mutex<Command>>,
db: &Db,
connections_key: &str, connections_key: &str,
only_address: Option<&str>, only_address: Option<&str>,
) -> Result<(), String> { ) -> Result<(), String> {
@ -168,20 +365,42 @@ async fn get_network_mapping_inner(
RpcResponse::send_raw(&unlocked_stream, Some(connections_key), &message).await; RpcResponse::send_raw(&unlocked_stream, Some(connections_key), &message).await;
let mut rx = download_hashmap_rx.lock().await; let mut rx = download_hashmap_rx.lock().await;
let mut buffer = timeout(Duration::from_secs(30), rx.recv()) let response = timeout(Duration::from_secs(30), rx.recv())
.await .await
.map_err(|_| "timed out waiting for network mapping response".to_string())? .map_err(|_| "timed out waiting for network mapping response".to_string())?
.ok_or_else(|| "network mapping response channel closed".to_string())?; .ok_or_else(|| "network mapping response channel closed".to_string())?;
if response.len() < NETWORK_SNAPSHOT_HEADER_BYTES {
return Err("network mapping response was missing its snapshot header".to_string());
}
if &response[0..4] != NETWORK_SNAPSHOT_MAGIC {
return Err("network mapping response had invalid snapshot magic".to_string());
}
if response[4] != NETWORK_SNAPSHOT_VERSION {
return Err("network mapping response used an unsupported snapshot version".to_string());
}
let mapping_len = u32::from_le_bytes(response[5..9].try_into().unwrap()) as usize;
let monitor_len = u32::from_le_bytes(response[9..13].try_into().unwrap()) as usize;
let expected_len = NETWORK_SNAPSHOT_HEADER_BYTES
.checked_add(mapping_len)
.and_then(|len| len.checked_add(monitor_len))
.ok_or_else(|| "network mapping snapshot lengths overflowed".to_string())?;
if response.len() != expected_len {
return Err("network mapping snapshot lengths did not match its payload".to_string());
}
let mapping_end = NETWORK_SNAPSHOT_HEADER_BYTES + mapping_len;
let mut buffer = response[NETWORK_SNAPSHOT_HEADER_BYTES..mapping_end].to_vec();
let monitor_state = &response[mapping_end..];
let mut snapshot = Vec::new();
let mut requested_address_found = only_address.is_none();
while buffer.len() >= NODE_RECORD_FIXED_BYTES { while buffer.len() >= NODE_RECORD_FIXED_BYTES {
let monitor_count = u16::from_le_bytes( let monitor_count = u16::from_le_bytes(
buffer[NODE_MONITOR_COUNT_OFFSET..NODE_RECORD_FIXED_BYTES] buffer[NODE_MONITOR_COUNT_OFFSET..NODE_RECORD_FIXED_BYTES]
.try_into() .try_into()
.unwrap(), .unwrap(),
) as usize; ) as usize;
if monitor_count != 0 {
return Err("network membership response contained unsigned monitor state".to_string());
}
let record_bytes = let record_bytes =
NODE_RECORD_FIXED_BYTES + (monitor_count * Wallet::SHORT_ADDRESS_BYTES_LENGTH); NODE_RECORD_FIXED_BYTES + (monitor_count * Wallet::SHORT_ADDRESS_BYTES_LENGTH);
if buffer.len() < record_bytes { if buffer.len() < record_bytes {
@ -189,21 +408,21 @@ async fn get_network_mapping_inner(
} }
let chunk: Vec<u8> = buffer.drain(0..record_bytes).collect(); let chunk: Vec<u8> = buffer.drain(0..record_bytes).collect();
// The first part of each record describes the advertised node address and IP. // The first part of each record describes the advertised node address and IP.
let Some(address) = Wallet::bytes_to_short_address(&chunk[0..NODE_IP_OFFSET]) else { let address = Wallet::bytes_to_short_address(&chunk[0..NODE_IP_OFFSET])
continue; .ok_or_else(|| "network mapping response contained an invalid address".to_string())?;
};
let ip = binary_to_ip(chunk[NODE_IP_OFFSET..NODE_PORT_OFFSET].to_vec()); let ip = binary_to_ip(chunk[NODE_IP_OFFSET..NODE_PORT_OFFSET].to_vec());
let port = u16::from_le_bytes( let port = u16::from_le_bytes(
chunk[NODE_PORT_OFFSET..NODE_BLOCKS_MINED_OFFSET] chunk[NODE_PORT_OFFSET..NODE_BLOCKS_MINED_OFFSET]
.try_into() .try_into()
.unwrap(), .unwrap(),
); );
let blocks_mined = chunk[NODE_BLOCKS_MINED_OFFSET]; let _blocks_mined = chunk[NODE_BLOCKS_MINED_OFFSET];
let added_by_bytes = &chunk[NODE_ADDED_BY_OFFSET..NODE_ADDED_TIMESTAMP_OFFSET]; let added_by_bytes = &chunk[NODE_ADDED_BY_OFFSET..NODE_ADDED_TIMESTAMP_OFFSET];
let added_by = if added_by_bytes.iter().all(|&byte| byte == 0) { let added_by = if added_by_bytes.iter().all(|&byte| byte == 0) {
String::new() String::new()
} else { } else {
Wallet::bytes_to_short_address(added_by_bytes).unwrap_or_default() Wallet::bytes_to_short_address(added_by_bytes)
.ok_or_else(|| "network mapping response contained an invalid signer".to_string())?
}; };
let added_timestamp = u64::from_le_bytes( let added_timestamp = u64::from_le_bytes(
chunk[NODE_ADDED_TIMESTAMP_OFFSET..NODE_ADDED_SIGNATURE_OFFSET] chunk[NODE_ADDED_TIMESTAMP_OFFSET..NODE_ADDED_SIGNATURE_OFFSET]
@ -222,10 +441,14 @@ async fn get_network_mapping_inner(
.try_into() .try_into()
.unwrap(), .unwrap(),
); );
if deleted_timestamp != 0 || deleted_block != 0 { let mut monitoring = Vec::with_capacity(monitor_count);
return Err( for monitor_bytes in
"network membership response contained unsigned deletion state".to_string(), chunk[NODE_RECORD_FIXED_BYTES..].chunks_exact(Wallet::SHORT_ADDRESS_BYTES_LENGTH)
); {
let Some(monitor) = Wallet::bytes_to_short_address(monitor_bytes) else {
return Err("network mapping response contained an invalid monitor".to_string());
};
monitoring.push(monitor);
} }
if only_address if only_address
@ -235,24 +458,48 @@ async fn get_network_mapping_inner(
continue; continue;
} }
// Import signed map records as remote state. This verifies the original requested_address_found = true;
// signer but does not require the newly connecting node to be eligible
// to add nodes locally after the mature-network gate. if only_address.is_none() {
if let Err(err) = NodeInfo::import_signed_mapping_address( snapshot.push(SyncedNodeState {
edit: SignedNodeEdit {
address,
ip,
port,
modified_by: added_by,
modified_timestamp: added_timestamp,
modified_signature: added_signature,
},
deleted_timestamp,
deleted_block,
monitoring,
});
continue;
}
// A reverse bootstrap announcement imports only the local node's
// newly sponsored record and must not replace the complete map. Keep
// its membership and derived state together if a full map is syncing.
if let Err(err) = NodeInfo::import_synced_mapping_address(
db, db,
SignedNodeEdit { SyncedNodeState {
address: address.clone(), edit: SignedNodeEdit {
ip: ip.clone(), address: address.clone(),
port, ip: ip.clone(),
modified_by: added_by, port,
modified_timestamp: added_timestamp, modified_by: added_by,
modified_signature: added_signature, modified_timestamp: added_timestamp,
modified_signature: added_signature,
},
deleted_timestamp,
deleted_block,
monitoring,
}, },
blocks_mined,
) )
.await .await
{ {
warn!("[network_map] skipped imported node record {address}: {err}"); warn!("[network_map] skipped imported node record {address}: {err}");
continue;
} }
} }
@ -260,62 +507,20 @@ async fn get_network_mapping_inner(
return Err("network mapping response had trailing partial bytes".to_string()); return Err("network mapping response had trailing partial bytes".to_string());
} }
get_signed_monitor_state( if !requested_address_found {
unlocked_stream, return Err(REQUESTED_MAPPING_RECORD_PENDING.to_string());
command_map, }
db,
&wallet.saved.short_address, if only_address.is_none() {
connections_key, let expected_targets = snapshot
only_address, .iter()
) .map(|state| (state.edit.address.clone(), state.edit.ip.clone()))
.await?; .collect();
let validated_monitor_state =
NodeInfo::validate_synced_monitor_state(monitor_state, db, &expected_targets).await?;
NodeInfo::install_synced_snapshot(db, snapshot, validated_monitor_state).await?;
NodeInfo::restore_or_rebuild_mined_counts(db).await?;
}
Ok(()) 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

@ -12,7 +12,6 @@ use crate::panic;
use crate::records::memory::chain_state::rebuild_chain_state_cache; use crate::records::memory::chain_state::rebuild_chain_state_cache;
use crate::records::memory::connections::initialize_node_runtime_context; use crate::records::memory::connections::initialize_node_runtime_context;
use crate::records::memory::mempool::{init_db, setup_mempool}; use crate::records::memory::mempool::{init_db, setup_mempool};
use crate::records::memory::network_mapping::NodeInfo;
use crate::records::memory::response_channels::Command; use crate::records::memory::response_channels::Command;
use crate::records::record_chain::rewards_tx::rebuild_immature_reward_window; use crate::records::record_chain::rewards_tx::rebuild_immature_reward_window;
use crate::rpc::server::start_rpc::start_rpc; use crate::rpc::server::start_rpc::start_rpc;
@ -131,17 +130,6 @@ pub async fn run_unlocked_node(wallet: Arc<Wallet>, install_shutdown: bool) -> R
error!("Failed to clear IP scores: {e}"); error!("Failed to clear IP scores: {e}");
} }
match NodeInfo::load_recovery_snapshot().await {
Ok(loaded) if loaded > 0 => {
info!("[network_map] loaded {loaded} recovered node records with monitors cleared");
if let Err(err) = NodeInfo::rebuild_mined_counts_from_chain(&db).await {
error!("[network_map] failed to rebuild recovered mined counts: {err}");
}
}
Ok(_) => {}
Err(err) => error!("[network_map] failed to load recovery snapshot: {err}"),
}
let wallet_for_server = wallet.clone(); let wallet_for_server = wallet.clone();
let map: Arc<Mutex<Command>> = Arc::new(Mutex::new(HashMap::new())); let map: Arc<Mutex<Command>> = Arc::new(Mutex::new(HashMap::new()));
// Install shared node dependencies before either incoming or outgoing // Install shared node dependencies before either incoming or outgoing

View File

@ -1,7 +0,0 @@
use crate::records::memory::mempool::signature_exists;
pub async fn memcheck(signature: &str, hash: &str) -> bool {
// Mempool uniqueness is checked by signature and hash so duplicate
// pending transactions can be rejected before deeper verification.
signature_exists(signature, hash).await.unwrap_or(false)
}

View File

@ -2,16 +2,20 @@ use crate::DateTime;
use crate::Utc; use crate::Utc;
pub async fn is_within_30_days(timestamp: u32) -> bool { pub async fn is_within_30_days(timestamp: u32) -> bool {
// Convert the on-chain timestamp into UTC and reject values that fall is_within_30_days_at(timestamp, Utc::now().timestamp() as u32)
// more than 30 days in the past or fail to decode as valid datetimes. }
let transaction_time = DateTime::<Utc>::from_timestamp(timestamp as i64, 0);
match transaction_time { pub fn is_within_30_days_at(timestamp: u32, as_of_timestamp: u32) -> bool {
Some(transaction_time) => { // Convert the on-chain timestamp into UTC and reject values that fall
let current_time = Utc::now(); // more than 30 days before the supplied validation time or fail to decode.
let duration = current_time.signed_duration_since(transaction_time); let transaction_time = DateTime::<Utc>::from_timestamp(timestamp as i64, 0);
let as_of_time = DateTime::<Utc>::from_timestamp(as_of_timestamp as i64, 0);
match (transaction_time, as_of_time) {
(Some(transaction_time), Some(as_of_time)) => {
let duration = as_of_time.signed_duration_since(transaction_time);
duration.num_days() <= 30 duration.num_days() <= 30
} }
None => false, _ => false,
} }
} }

View File

@ -229,7 +229,7 @@ async fn verify_transaction(
} }
} }
if let Transaction::Swap(swap_tx) = &transaction { if let Transaction::Swap(swap_tx) = &transaction {
match swap_tx.verify(db).await { match swap_tx.verify_at(db, block_timestamp).await {
Ok(value) => { Ok(value) => {
reserve_verified_transaction(db, &transaction, balance_tracker, already_in_mempool) reserve_verified_transaction(db, &transaction, balance_tracker, already_in_mempool)
.await?; .await?;

View File

@ -9,7 +9,7 @@ use crate::records::wallet_registry::{
}; };
use crate::sled::Db; use crate::sled::Db;
use crate::verifications::async_funcs::checks::balance_check::balance_checkup; use crate::verifications::async_funcs::checks::balance_check::balance_checkup;
use crate::verifications::async_funcs::checks::time_checks::is_within_30_days; use crate::verifications::async_funcs::checks::time_checks::is_within_30_days_at;
use crate::verifications::async_funcs::checks::verify_db::{ use crate::verifications::async_funcs::checks::verify_db::{
db_bytes_verification, db_hex_verification, db_bytes_verification, db_hex_verification,
}; };
@ -35,6 +35,10 @@ fn validate_swap_tip(value: u64, tip: u64, is_nft: bool, sender: &str) -> Result
impl SwapTransaction { impl SwapTransaction {
pub async fn verify(&self, db: &Db) -> Result<String, String> { pub async fn verify(&self, db: &Db) -> Result<String, String> {
self.verify_at(db, Utc::now().timestamp() as u32).await
}
pub async fn verify_at(&self, db: &Db, as_of_timestamp: u32) -> Result<String, String> {
let hash = self.unsigned_swap.hash().await; let hash = self.unsigned_swap.hash().await;
// Both swap participants must provide valid wallet addresses and // Both swap participants must provide valid wallet addresses and
// matching signatures over the shared unsigned swap payload. // matching signatures over the shared unsigned swap payload.
@ -203,14 +207,13 @@ impl SwapTransaction {
// Swap offers are bounded both by signature age and by the // Swap offers are bounded both by signature age and by the
// explicit offer-expiration window carried in the transaction. // explicit offer-expiration window carried in the transaction.
if !is_within_30_days(self.unsigned_swap.timestamp).await { if !is_within_30_days_at(self.unsigned_swap.timestamp, as_of_timestamp) {
return Err( return Err(
"Timestamp is to old. Transactions must be broadcast within 30 days of signing." "Timestamp is to old. Transactions must be broadcast within 30 days of signing."
.to_string(), .to_string(),
); );
} }
let now = Utc::now().timestamp() as u32;
let offer_expiration = self.unsigned_swap.offer_expiration; let offer_expiration = self.unsigned_swap.offer_expiration;
let timestamp = self.unsigned_swap.timestamp; let timestamp = self.unsigned_swap.timestamp;
@ -227,7 +230,7 @@ impl SwapTransaction {
); );
} }
if now > offer_expiration { if as_of_timestamp > offer_expiration {
return Err("This swap offer has expired.".to_string()); return Err("This swap offer has expired.".to_string());
} }