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};
#[cfg(windows)]
use std::env;
@ -5,7 +6,7 @@ use std::env;
use std::error::Error;
#[cfg(windows)]
use std::error::Error;
#[cfg(unix)]
#[cfg(all(unix, any()))]
use std::fs;
#[cfg(windows)]
use std::fs;
@ -13,13 +14,23 @@ use std::fs;
use std::path::{Path, PathBuf};
#[cfg(windows)]
use std::process::Command;
#[cfg(unix)]
#[cfg(all(unix, any()))]
use std::process::{Command, Stdio};
#[cfg(windows)]
const DEFAULT_PG_VERSION: &str = "16.4-1";
#[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]
async fn main() -> Result<(), Box<dyn Error>> {
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.
ensure_administrator()?;
// Collect install settings and the database credentials to print at the end.
let host = prompt_visible_with_default("Enter Postgres host", "127.0.0.1").await;
let port = prompt_visible_with_default("Enter Postgres port", "5432").await;
let dbname =
prompt_visible_with_default("Enter new database name for blockchain", "contractless_db")
// A release installer can provide all values through process-local environment
// variables. Interactive use remains unchanged when those values are absent.
let host = env_or_prompt("CONTRACTLESS_PG_HOST", "Enter Postgres host", "127.0.0.1").await;
let port = env_or_prompt("CONTRACTLESS_PG_PORT", "Enter Postgres port", "5432").await;
let dbname = env_or_prompt(
"CONTRACTLESS_PG_DATABASE",
"Enter new database name for blockchain",
"contractless_db",
)
.await;
let user =
prompt_visible_with_default("Enter new username for blockchain database", "contractless")
let user = env_or_prompt(
"CONTRACTLESS_PG_USER",
"Enter new username for blockchain database",
"contractless",
)
.await;
let user_pass = prompt_hidden_nonempty(
let user_pass = env_or_hidden_prompt(
"CONTRACTLESS_PG_USER_PASSWORD",
"Enter password for new database user: ",
"Password cannot be empty. Please try again.",
)
.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: ",
"Password cannot be empty. Please try again.",
)
.await;
let install_dir = prompt_visible_with_default(
let install_dir = env_or_prompt(
"CONTRACTLESS_PG_INSTALL_DIR",
"Enter PostgreSQL install directory",
r"C:\Program Files\PostgreSQL\16",
)
.await;
let data_dir = prompt_visible_with_default(
let data_dir = env_or_prompt(
"CONTRACTLESS_PG_DATA_DIR",
"Enter PostgreSQL data directory",
r"C:\Program Files\PostgreSQL\16\data",
)
.await;
let service_name = prompt_visible_with_default(
let service_name = env_or_prompt(
"CONTRACTLESS_PG_SERVICE_NAME",
"Enter PostgreSQL Windows service name",
"postgresql-contractless",
)
@ -214,13 +235,29 @@ async fn main() -> Result<(), Box<dyn Error>> {
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)))]
fn main() {
eprintln!("postgres_installer is only supported on Unix-like systems and Windows.");
std::process::exit(1);
}
#[cfg(unix)]
#[cfg(all(unix, any()))]
fn find_pg_hba_path() -> Result<String, Box<dyn Error>> {
// Ask PostgreSQL for the active pg_hba.conf path instead of guessing distro paths.
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::log::error;
use crate::orphans::get_path_names::get_file_names;
use crate::orphans::save_blocks::save_new_blocks;
use crate::orphans::structs::UndoTransactions;
@ -228,6 +229,9 @@ pub async fn undo_transactions(
let final_height = true_start_height.saturating_sub(1);
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?;
// 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,
db: context.db,
map: context.map,
first: false,
run_startup_sync: false,
};
@ -218,7 +217,10 @@ async fn retry_dropped_outgoing(ip: String, port: u16) {
wallet: context.wallet.clone(),
db: context.db.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 {
@ -593,6 +595,7 @@ impl Connection {
}
info.ready = true;
info.catch_up_target = None;
spawn_monitor_update(
ip.clone(),
MONITOR_ACTION_ADD,
@ -679,6 +682,16 @@ impl Connection {
.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 {
self.connection_map
.values()
@ -728,6 +741,24 @@ impl Connection {
.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(
&self,
) -> 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(),
connections_key: format!("{ip}:{port}"),
};
let _ = if action == MONITOR_ACTION_ADD {
NodeInfo::add_monitor(params).await
if action == MONITOR_ACTION_ADD {
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 {
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 {
// Read the singleton connection manager and count live outgoing peers.
CONNECTIONS
@ -1013,6 +1096,15 @@ pub async fn peer_connection_count() -> usize {
.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 {
// Recovery logic uses raw miner sockets to avoid starting another
// 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)
}
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 {
if peer_is_operational(key).await {
return true;
@ -1185,6 +1326,34 @@ pub async fn peer_accepts_live_relay(key: &str) -> bool {
.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) {
let mut guard = CONNECTIONS.write().await;
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 crate::common::check_genesis::genesis_checkup;
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) {
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 {
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;
if let Some(node_info) = map.get_mut(address) {
// Counts are capped at u8-safe policy maximum used by node rules.
if node_info.blocks_mined < 250 {
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) {
@ -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> {
// Recompute node mined counts directly from saved block headers
// so startup and recovery can rebuild memory-only state.
@ -55,7 +149,7 @@ impl NodeInfo {
node_info.blocks_mined = 0;
}
drop(map);
Self::persist_recovery_snapshot("mined rebuild without genesis").await;
Self::persist_mined_counts(db).await?;
return Ok(());
}
@ -80,7 +174,38 @@ impl NodeInfo {
}
}
}
Self::persist_recovery_snapshot("mined rebuild").await;
Self::persist_mined_counts(db).await?;
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::log::{info, warn};
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::network_mapping::enums::NodeEditType;
use crate::records::memory::network_mapping::structs::{
@ -23,9 +21,27 @@ use crate::HashMap;
use crate::Mutex;
use crate::OnceLock;
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! {
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;
@ -93,6 +109,5 @@ mod add;
pub mod enums;
mod mined_counts;
pub(crate) mod monitor;
pub(crate) mod persistence;
mod queries;
pub mod structs;

View File

@ -14,7 +14,7 @@ pub const MONITOR_ACTION_REMOVE: u8 = 2;
lazy_static! {
// Keep the newest signed assertion for each monitor relationship. This is
// both replay protection and the authenticated state shared with joining peers.
static ref MONITOR_EVENT_STATE: Mutex<HashMap<String, SignedMonitorEdit>> =
pub(super) static ref MONITOR_EVENT_STATE: Mutex<HashMap<String, SignedMonitorEdit>> =
Mutex::new(HashMap::new());
}
@ -108,6 +108,39 @@ mod tests {
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 {
@ -226,31 +259,14 @@ impl NodeInfo {
)
}
async fn remember_newest_monitor_event(edit: &SignedMonitorEdit) -> bool {
let key = Self::monitor_relation_key(edit);
let mut state = MONITOR_EVENT_STATE.lock().await;
if state
.get(&key)
.map(|existing| Self::monitor_event_order(edit) <= Self::monitor_event_order(existing))
.unwrap_or(false)
{
return false;
}
state.insert(key, edit.clone());
true
}
async fn deletion_marker_for(monitored_address: &str) -> Option<(u64, u32)> {
MONITOR_EVENT_STATE
.lock()
.await
.values()
.filter(|event| {
event.monitored_address == monitored_address
&& event.action == MONITOR_ACTION_REMOVE
})
.map(|event| (event.modified_timestamp, event.modified_block))
.max()
pub(super) fn discard_old_target_ip_events_from(
state: &mut HashMap<String, SignedMonitorEdit>,
address: &str,
current_ip: &str,
) {
state.retain(|_, event| {
event.monitored_address != address || event.target_ip == current_ip
});
}
async fn broadcast_monitor_event(
@ -282,7 +298,7 @@ impl NodeInfo {
let connections_lock = CONNECTIONS.read().await;
connections_lock
.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()
};
@ -327,28 +343,55 @@ impl NodeInfo {
return Err("Could not validate monitor signature".to_string());
}
// Validate membership before recording the event so an invalid event
// cannot poison replay protection for a later valid assertion.
{
let address_map = ADDRESS_MAP.lock().await;
// Mapping snapshots lock these structures in this same order. Holding
// 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
{
return Err(
"monitor target must have a newer membership record before reconnection"
.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);
}
event_state.insert(relation_key, edit.clone());
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 {
None
};
let mut address_map = ADDRESS_MAP.lock().await;
let monitored = address_map
.get_mut(&edit.monitored_address)
.ok_or_else(|| "monitored address not found".to_string())?;
@ -384,6 +427,14 @@ impl NodeInfo {
}
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 {
mut edit,
remote_ip,
@ -419,7 +470,6 @@ impl NodeInfo {
Ok(true) => {}
}
Self::persist_recovery_snapshot("monitor update").await;
Self::broadcast_monitor_event(map, &edit, &remote_ip).await;
RpcResponse::Binary(b"Success".to_vec())
}
@ -432,9 +482,10 @@ impl NodeInfo {
Self::apply_verified_monitor_edit(&edit, db, local_short).await
}
pub(crate) async fn signed_monitor_state_bytes() -> Vec<u8> {
let mut events: Vec<SignedMonitorEdit> =
MONITOR_EVENT_STATE.lock().await.values().cloned().collect();
pub(super) fn signed_monitor_state_bytes_from(
state: &HashMap<String, SignedMonitorEdit>,
) -> Vec<u8> {
let mut events: Vec<SignedMonitorEdit> = state.values().cloned().collect();
events.sort_by(|left, right| {
Self::monitor_event_order(left)
.cmp(&Self::monitor_event_order(right))
@ -468,28 +519,51 @@ impl NodeInfo {
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 {
RpcResponse::Binary(Self::signed_monitor_state_bytes().await)
}
pub(crate) async fn load_monitor_recovery_state(bytes: &[u8]) -> usize {
let mut state = MONITOR_EVENT_STATE.lock().await;
state.clear();
let mut loaded = 0;
pub(crate) async fn validate_synced_monitor_state(
bytes: &[u8],
db: &Db,
expected_targets: &HashMap<String, String>,
) -> Result<HashMap<String, SignedMonitorEdit>, String> {
if bytes.len() % MONITOR_EVENT_BYTES != 0 {
return Err("monitor-state snapshot ended mid-record".to_string());
}
let mut replacement = HashMap::new();
for chunk in bytes.chunks_exact(MONITOR_EVENT_BYTES) {
let Some(edit) = Self::monitor_event_from_bytes(chunk) else {
continue;
};
// A restarted node no longer has the live sockets represented by
// old additions. Retain signed removals as historical deletion
// evidence; current additions are recreated by operational peers.
if edit.action != MONITOR_ACTION_REMOVE {
continue;
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());
}
state.insert(Self::monitor_relation_key(&edit), edit);
loaded += 1;
if !Self::verify_monitor_edit(&edit, db).await {
return Err("monitor-state snapshot contained an invalid signature".to_string());
}
loaded
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> {

View File

@ -3,10 +3,174 @@ use crate::common::governance::{
activation_vote_record_key, countable_governance_addresses, proposal_vote_record_key,
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 std::collections::HashSet;
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> {
let map = ADDRESS_MAP.lock().await;
map.get(address).map(|node| GovernanceNodeSnapshot {
@ -191,10 +355,13 @@ impl NodeInfo {
}
pub async fn request_valid_nodes() -> RpcResponse {
// Serialize the in-memory node map into the binary layout
// used by peer bootstrap and node-list synchronization.
// Serialize the complete in-memory node map used by peer bootstrap.
// 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 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() {
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();
let added_timestamp_bytes = node_info.added_timestamp.to_le_bytes();
// Network-map snapshots contain signed membership only. Liveness
// is transferred separately as signed monitor assertions.
let deleted_timestamp_bytes = 0_u64.to_le_bytes();
let deleted_block_bytes = 0_u32.to_le_bytes();
let monitor_count_bytes = 0_u16.to_le_bytes();
let monitor_bytes: Vec<u8> = node_info
.monitoring
.iter()
.filter_map(|monitor| Wallet::short_address_to_bytes(monitor))
.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();
// Field order here must match the parser used by node-list
// synchronization.
data.extend_from_slice(&address_bytes);
data.extend_from_slice(&ip_bytes);
data.extend_from_slice(&port_bytes);
data.push(blocks_mined);
mapping.extend_from_slice(&address_bytes);
mapping.extend_from_slice(&ip_bytes);
mapping.extend_from_slice(&port_bytes);
mapping.push(blocks_mined);
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 {
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);
data.extend_from_slice(&added_signature_bytes);
data.extend_from_slice(&deleted_timestamp_bytes);
data.extend_from_slice(&deleted_block_bytes);
data.extend_from_slice(&monitor_count_bytes);
mapping.extend_from_slice(&added_timestamp_bytes);
mapping.extend_from_slice(&added_signature_bytes);
mapping.extend_from_slice(&node_info.deleted_timestamp.to_le_bytes());
mapping.extend_from_slice(&node_info.deleted_block.to_le_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)
}
@ -253,3 +439,91 @@ impl NodeInfo {
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;
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.
#[derive(Debug, Clone)]
pub struct SignedNodeEdit {
@ -59,6 +63,14 @@ pub struct SignedMonitorEdit {
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.
#[derive(Clone)]
pub struct AddAddressParams {

View File

@ -29,6 +29,9 @@ pub struct ConnectionInfo {
pub local_setup_complete: bool,
pub remote_setup_complete: 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 health: ConnectionHealth,
}
@ -228,6 +231,7 @@ impl ConnectionInfo {
local_setup_complete: false,
remote_setup_complete: false,
local_setup_acknowledged: false,
catch_up_target: None,
ready: false,
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;
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 is_reorganizing_mode() && !allow_during_reorg {
return Err("Cannot save discovered block while reorganizing.".to_string());
@ -256,6 +261,17 @@ pub async fn save_block(params: SaveBlockParams) -> Result<(), String> {
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(
block_number: 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
// the persisted chain height forward.
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(())
@ -590,7 +608,9 @@ async fn save_binary_data(params: SaveBinaryDataParams<'_>) -> Result<(), String
// Only advance mined-count tracking when this save actually moved
// the persisted chain height forward.
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(())
@ -693,3 +713,20 @@ fn format_block_time(timestamp: u32) -> 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::sled::Db;
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::thread_rng;
use crate::wallets::structures::Wallet;
@ -43,6 +45,25 @@ use crate::Duration;
use crate::Mutex;
use crate::SliceRandom;
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;
#[derive(Clone)]
@ -52,19 +73,25 @@ pub struct BootstrapParams {
pub wallet: Arc<Wallet>,
pub db: Db,
pub map: Arc<Mutex<Command>>,
pub first: bool,
pub run_startup_sync: bool,
}
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 {
let _startup_guard = startup_guard;
if let Err(e) = bootstrap_peer_discovery(params).await {
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 {
Some(begin_chain_sync().await)
} else {
@ -72,29 +99,31 @@ pub async fn bootstrap_peer_discovery(mut params: BootstrapParams) -> Result<(),
};
let (_, _, local_endpoint) = get_ip_and_port().await;
let max = SETTINGS.outgoing_connections;
let mut current_key = params.connections_key.clone();
let mut stream = params.stream;
params.first = false;
let current_key = params.connections_key.clone();
let stream = params.stream;
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 {
let outgoing_connections = {
if params.run_startup_sync {
// 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;
connections
.as_ref()
.map(|connection| connection.count_ready_outgoing_connections())
.map(|connection| connection.count_outgoing_connections())
.unwrap_or(0)
};
if outgoing_connections >= max as usize {
break;
}
if addr_string == local_endpoint {
continue;
}
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()
@ -109,39 +138,53 @@ pub async fn bootstrap_peer_discovery(mut params: BootstrapParams) -> Result<(),
continue;
}
};
sleep(Duration::from_secs(2)).await;
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: params.wallet.clone(),
db: params.db.clone(),
map: params.map.clone(),
first: params.first,
wallet,
db,
map,
first: false,
};
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 !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 {
return Ok(());
}
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?;
request_remote_height(stream.clone(), params.map.clone(), current_key.clone())
.await?;
ensure_compatible_genesis(
stream.clone(),
params.map.clone(),
@ -180,8 +223,11 @@ pub async fn bootstrap_peer_discovery(mut params: BootstrapParams) -> Result<(),
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())
let imported_candidates = match hydrate_torrent_candidates(
stream.clone(),
params.map.clone(),
current_key.clone(),
)
.await
{
Ok(imported) => {
@ -238,9 +284,91 @@ pub async fn bootstrap_peer_discovery(mut params: BootstrapParams) -> Result<(),
if let Some(guard) = chain_sync_guard {
guard.finish();
}
}
// Outside startup, topology repair can use the ordinary sequential path.
// Startup itself already opened its torrent pool concurrently above.
if !params.run_startup_sync {
let mut candidate_endpoints = NodeInfo::active_node_endpoints().await;
candidate_endpoints.shuffle(&mut thread_rng());
for addr_string in candidate_endpoints {
let outgoing_connections = {
let connections = CONNECTIONS.read().await;
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}");
}
}
}
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(
response: Vec<u8>,
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}"
);
// Only the bootstrap owner imports a complete mapping snapshot.
// Discovered peers join the existing map and receive live updates.
announce_self_to_network(
broadcast_stream.clone(),
&wallet.saved.short_address.clone(),
params.map.clone(),
&params.db.clone(),
params.wallet.clone(),
&connections_key,
true,
params.first,
)
.await
.map_err(|err| {
@ -418,6 +547,46 @@ pub async fn process_handshake_response(
});
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;
if params.first {
@ -427,17 +596,19 @@ pub async fn process_handshake_response(
wallet: params.wallet.clone(),
db: params.db.clone(),
map: params.map.clone(),
first: params.first,
run_startup_sync: true,
};
spawn_bootstrap_peer_discovery(bsparams);
} else {
if is_normal_mode() {
} else if is_normal_mode() {
if !mark_peer_operational(&connections_key, params.map.clone()).await {
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(())
}

View File

@ -25,6 +25,31 @@ use std::collections::BTreeMap;
const SYNC_PREFETCH_WINDOW: usize = 100;
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 {
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 {
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);
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;
warn!(
"[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;
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_GOVERNANCE_STATE_SYNC: u8 = 56;
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 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::NodeInfo;
use crate::records::memory::enums::ClientType;
use crate::records::memory::response_channels::Command;
use crate::rpc::read_bytes_from_stream;
use crate::rpc::responses::RpcResponse;
@ -15,6 +16,8 @@ pub async fn add_network_node(
db: &Db,
wallet: Arc<Wallet>,
map: Arc<Mutex<Command>>,
client_type: ClientType,
reconciliation: bool,
) -> Result<(u32, RpcResponse), String> {
// Command 28 carries the signed node-add payload directly after the
// 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 {
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?;
// NodeInfo owns the signature checks, local re-signing rules, and
// in-memory/broadcast side effects for the actual add operation.
let result = NodeInfo::add_address(AddAddressParams {
map,
edit: SignedNodeEdit {
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
// in-memory/broadcast side effects for the actual add operation.
let result = NodeInfo::add_address(AddAddressParams {
map,
edit,
monitors: Vec::new(),
blocks_mined: 0_u8,
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 memory_by_signature;
pub mod network_info;
pub mod network_mapping_hash;
pub mod network_monitor_add;
pub mod network_monitor_remove;
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::sync_check::sync_checkup;
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::response_channels::Command;
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)
&& !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);
if !within_orphan_window(local_height, block_number) {
if is_syncing_mode()

View File

@ -1,10 +1,13 @@
use crate::io::ErrorKind;
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::rpc::command_maps::{
RPC_BLOCK_HEIGHT, RPC_BLOCK_PIECE, RPC_NETWORK_MONITOR_ADD, RPC_NETWORK_MONITOR_REMOVE,
RPC_REPLY, RPC_SUBMIT_TORRENT, RPC_SUBMIT_TRANSACTION, RPC_TORRENT_BY_HEIGHT,
RPC_BLOCK_HEIGHT, RPC_BLOCK_PIECE, RPC_NETWORK_MEMBERSHIP_RECONCILE,
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::flood_protection::check_request_frequency_with_client_type;
@ -80,6 +83,7 @@ fn requires_operational_miner(command: u8) -> bool {
| RPC_SUBMIT_TORRENT
| RPC_NETWORK_MONITOR_ADD
| RPC_NETWORK_MONITOR_REMOVE
| RPC_NETWORK_MEMBERSHIP_RECONCILE
)
}
@ -102,10 +106,17 @@ pub async fn next_incoming_command(
.await
.unwrap_or(ClientType::Miner);
if client_type == ClientType::Miner
&& requires_operational_miner(command)
&& !peer_accepts_live_relay(connections_key).await
{
let accepts_relay = if matches!(
command,
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!(
"[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::memory::enums::ClientType;
use crate::rpc::command_maps::{
RPC_ADD_NETWORK_NODE, RPC_NETWORK_MONITOR_ADD, RPC_NETWORK_MONITOR_REMOVE,
RPC_NETWORK_MONITOR_STATE, RPC_SETUP_COMPLETE,
RPC_ADD_NETWORK_NODE, RPC_NETWORK_MAPPING_HASH, RPC_NETWORK_MONITOR_ADD,
RPC_NETWORK_MEMBERSHIP_RECONCILE, RPC_NETWORK_MONITOR_REMOVE, RPC_NETWORK_MONITOR_STATE,
RPC_SETUP_COMPLETE,
};
use crate::rpc::server::structs::RpcFloodState;
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_REMOVE
| RPC_NETWORK_MONITOR_STATE
| RPC_NETWORK_MAPPING_HASH
| RPC_NETWORK_MEMBERSHIP_RECONCILE
| RPC_SETUP_COMPLETE
) {
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::records::block_height::get_block_height::get_height;
use crate::records::memory::connections::{
mark_peer_network_map_synced, mark_peer_operational, mark_peer_wallet_registry_synced,
spawn_peer_setup_retry,
mark_peer_catching_up, mark_peer_network_map_synced, mark_peer_operational,
mark_peer_wallet_registry_synced, spawn_peer_setup_retry,
};
use crate::records::memory::enums::ClientType;
use crate::records::memory::response_channels::generate_uid;
@ -40,6 +40,23 @@ use crate::Settings;
use crate::TcpStream;
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>>) {
// Failed handshakes are never stored in connection memory, but the
// accepted TCP socket should still be closed immediately.
@ -77,7 +94,7 @@ async fn sync_incoming_peer_before_operational(
wallet: Arc<Wallet>,
map: Arc<Mutex<Command>>,
connections_key: &str,
) -> Result<(bool, Option<ChainOperationGuard>), String> {
) -> Result<(bool, Option<ChainOperationGuard>, Option<u32>), String> {
let initial_local_height = get_height(db);
let initial_remote_height =
request_remote_height(stream.clone(), map.clone(), connections_key.to_string()).await?;
@ -96,7 +113,7 @@ async fn sync_incoming_peer_before_operational(
warn!(
"[startup] incoming peer is behind local chain; keeping connection passive until peer catches up: local_height={initial_local_height} remote_height={initial_remote_height}"
);
return Ok((false, None));
return Ok((false, None, Some(initial_local_height)));
}
// Wait for any existing startup, live catch-up, or orphan operation to
@ -118,7 +135,7 @@ async fn sync_incoming_peer_before_operational(
warn!(
"[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 {
@ -186,7 +203,7 @@ async fn sync_incoming_peer_before_operational(
.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(
@ -194,6 +211,7 @@ fn spawn_incoming_peer_promotion_watcher(
db: Db,
map: Arc<Mutex<Command>>,
connections_key: String,
catch_up_target: u32,
) {
tokio::spawn(async move {
loop {
@ -215,15 +233,46 @@ fn spawn_incoming_peer_promotion_watcher(
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 {
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(
stream: Arc<Mutex<TcpStream>>,
db: &Db,
@ -319,7 +368,6 @@ async fn complete_incoming_miner_setup(
&short_address,
map.clone(),
db,
wallet.clone(),
connections_key,
false,
)
@ -333,7 +381,6 @@ async fn complete_incoming_miner_setup(
stream.clone(),
map.clone(),
db,
wallet.clone(),
connections_key,
&short_address,
)
@ -346,7 +393,8 @@ async fn complete_incoming_miner_setup(
}
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) =
match sync_incoming_peer_before_operational(
stream.clone(),
db,
wallet.clone(),
@ -372,11 +420,26 @@ async fn complete_incoming_miner_setup(
guard.finish();
}
} 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(
stream.clone(),
db.clone(),
map.clone(),
connections_key.to_string(),
catch_up_target,
);
}
}

View File

@ -549,6 +549,8 @@ pub async fn start_loop(
&db,
wallet.clone(),
map.clone(),
client_type,
false,
)
.await?;
let should_drop_rejected_miner = client_type == ClientType::Miner
@ -561,6 +563,21 @@ pub async fn start_loop(
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 => {
// signed monitor-add event from a miner peer
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)
.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 => {
// Both sides announce setup completion independently. Receipt
// 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_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_NFT_DETAILS, RPC_NFT_LIST, RPC_REGISTER_WALLET, RPC_STORAGE_LOOKUP,
RPC_STORAGE_LOOKUP_COST, RPC_TIME, RPC_TOKEN_CATALOG, RPC_TOKEN_DETAILS, RPC_TOKEN_LIST,
RPC_TORRENT_BY_HEIGHT, RPC_TOTAL_CONFIRMED_TX, RPC_TRANSACTION_BY_TXID, RPC_UNBLOCK_IP,
RPC_VALIDATE_MESSAGE, RPC_VANITY_LOOKUP, RPC_VANITY_OWNER_LOOKUP,
RPC_WALLET_REGISTRATION_STATUS, RPC_WALLET_REGISTRY_SYNC,
RPC_NETWORK_MONITOR_STATE, RPC_NFT_DETAILS, RPC_NFT_LIST, RPC_REGISTER_WALLET,
RPC_REQUEST_NODE_LIST, RPC_STORAGE_LOOKUP, RPC_STORAGE_LOOKUP_COST, RPC_TIME,
RPC_TOKEN_CATALOG, RPC_TOKEN_DETAILS, RPC_TOKEN_LIST, RPC_TORRENT_BY_HEIGHT,
RPC_TOTAL_CONFIRMED_TX, RPC_TRANSACTION_BY_TXID, RPC_UNBLOCK_IP, RPC_VALIDATE_MESSAGE,
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::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(&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.
31 => {
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(&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.
57 => {
let command_number: u8 = RPC_GOVERNANCE_PROPOSAL;
@ -804,7 +815,10 @@ async fn build_request_bytes(
#[cfg(test)]
mod tests {
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;
#[tokio::test]
@ -843,4 +857,15 @@ mod tests {
assert_eq!(&bytes[1..4], &uid);
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::miner::flag::{
clear_mining_stop_request, is_normal_mode, set_mining_state, set_node_mode, MiningState,
NodeMode,
};
use crate::miner::flag::{is_mining_stop_requested, is_normal_mode};
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::rpc::client::handshake::connect_and_handshake;
use crate::rpc::client::structs::Connect;
@ -22,29 +21,38 @@ pub async fn handle_connections(
wallet: Arc<Wallet>,
map: Arc<Mutex<Command>>,
) -> Result<(), String> {
// A zero outgoing limit means this node should not open any bootstrap
// connection during startup.
// A zero outgoing limit means this node waits for an incoming sponsor. It
// does not bypass the independent-peer requirement or permit solo mining.
let outgoing_connections = crate::Settings::load()
.map(|settings| settings.outgoing_connections)
.unwrap_or(0);
if outgoing_connections == 0 {
info!("OUTGOING_CONNECTIONS is 0; skipping startup bootstrap.");
set_node_mode(NodeMode::Normal);
clear_mining_stop_request();
set_mining_state(MiningState::Idle);
info!("OUTGOING_CONNECTIONS is 0; waiting for an incoming sponsored peer.");
}
loop {
if outgoing_connections > 0
&& attempt_bootstrap_connections(db.clone(), wallet.clone(), map.clone(), "startup")
.await?
{
return Ok(());
}
let connected = attempt_bootstrap_connections(db, wallet, map, "startup").await?;
if connected {
info!("No existing network peer is available. Waiting for connections.");
for _ in 0..60 {
let peer_wallets = operational_peer_wallets().await;
if is_normal_mode()
&& !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(());
}
// Startup can continue as a standalone node even if no bootstrap peer is reachable.
set_node_mode(NodeMode::Normal);
clear_mining_stop_request();
set_mining_state(MiningState::Idle);
Ok(())
sleep(Duration::from_secs(1)).await;
}
warn!("No sponsored peer is operational; retrying configured bootstrap peers");
}
}
async fn attempt_bootstrap_connections(
@ -56,9 +64,15 @@ async fn attempt_bootstrap_connections(
// Try the configured bootstrap peers one by one until a
// handshake succeeds or the list is exhausted.
let filtered_servers = get_node_connections().await;
let (_, _, local_endpoint) = get_ip_and_port().await;
let mut last_error: Option<String> = None;
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
// shared state so each attempt can run independently
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 {
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::log::warn;
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_DELETED_TIMESTAMP_OFFSET, NODE_IP_OFFSET, NODE_MONITOR_COUNT_OFFSET, NODE_PORT_OFFSET,
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::response_channels::{reserve_entry, Command};
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::sled::Db;
@ -22,12 +24,42 @@ use crate::Mutex;
use crate::TcpStream;
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(
unlocked_stream: Arc<Mutex<TcpStream>>,
address: &str,
command_map: Arc<Mutex<Command>>,
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,
import_peer_map: bool,
) -> Result<(), String> {
@ -96,12 +128,12 @@ pub async fn announce_self_to_network(
}
if import_peer_map {
get_network_mapping(
get_network_mapping_response(
unlocked_stream,
command_map.clone(),
db,
wallet.clone(),
connections_key,
None,
)
.await?;
}
@ -112,44 +144,209 @@ pub async fn get_network_mapping(
unlocked_stream: Arc<Mutex<TcpStream>>,
command_map: Arc<Mutex<Command>>,
db: &Db,
wallet: Arc<Wallet>,
connections_key: &str,
) -> Result<(), String> {
get_network_mapping_inner(
unlocked_stream,
command_map,
db,
wallet,
connections_key,
None,
)
get_network_mapping_inner(unlocked_stream, command_map, db, connections_key, None).await
}
pub async fn compare_network_mapping_digest(
unlocked_stream: Arc<Mutex<TcpStream>>,
command_map: Arc<Mutex<Command>>,
connections_key: &str,
) -> Result<bool, String> {
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(
unlocked_stream: Arc<Mutex<TcpStream>>,
command_map: Arc<Mutex<Command>>,
db: &Db,
wallet: Arc<Wallet>,
connections_key: &str,
only_address: &str,
) -> Result<(), String> {
get_network_mapping_inner(
unlocked_stream,
command_map,
// During the first two-node bootstrap, the peer may acknowledge our
// reverse announcement while deferring its application until an in-flight
// full mapping snapshot finishes. Wait for that exact sponsored record to
// become visible instead of treating an absent record as a successful sync.
for attempt in 1..=40 {
match get_network_mapping_inner(
unlocked_stream.clone(),
command_map.clone(),
db,
wallet,
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(
unlocked_stream: Arc<Mutex<TcpStream>>,
command_map: Arc<Mutex<Command>>,
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,
only_address: Option<&str>,
) -> Result<(), String> {
@ -168,20 +365,42 @@ async fn get_network_mapping_inner(
RpcResponse::send_raw(&unlocked_stream, Some(connections_key), &message).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
.map_err(|_| "timed out waiting for network mapping response".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 {
let monitor_count = u16::from_le_bytes(
buffer[NODE_MONITOR_COUNT_OFFSET..NODE_RECORD_FIXED_BYTES]
.try_into()
.unwrap(),
) as usize;
if monitor_count != 0 {
return Err("network membership response contained unsigned monitor state".to_string());
}
let record_bytes =
NODE_RECORD_FIXED_BYTES + (monitor_count * Wallet::SHORT_ADDRESS_BYTES_LENGTH);
if buffer.len() < record_bytes {
@ -189,21 +408,21 @@ async fn get_network_mapping_inner(
}
let chunk: Vec<u8> = buffer.drain(0..record_bytes).collect();
// 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 {
continue;
};
let address = Wallet::bytes_to_short_address(&chunk[0..NODE_IP_OFFSET])
.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 port = u16::from_le_bytes(
chunk[NODE_PORT_OFFSET..NODE_BLOCKS_MINED_OFFSET]
.try_into()
.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 = if added_by_bytes.iter().all(|&byte| byte == 0) {
String::new()
} 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(
chunk[NODE_ADDED_TIMESTAMP_OFFSET..NODE_ADDED_SIGNATURE_OFFSET]
@ -222,10 +441,14 @@ async fn get_network_mapping_inner(
.try_into()
.unwrap(),
);
if deleted_timestamp != 0 || deleted_block != 0 {
return Err(
"network membership response contained unsigned deletion state".to_string(),
);
let mut monitoring = Vec::with_capacity(monitor_count);
for monitor_bytes in
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
@ -235,12 +458,32 @@ async fn get_network_mapping_inner(
continue;
}
// Import signed map records as remote state. This verifies the original
// signer but does not require the newly connecting node to be eligible
// to add nodes locally after the mature-network gate.
if let Err(err) = NodeInfo::import_signed_mapping_address(
requested_address_found = true;
if only_address.is_none() {
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,
SignedNodeEdit {
SyncedNodeState {
edit: SignedNodeEdit {
address: address.clone(),
ip: ip.clone(),
port,
@ -248,11 +491,15 @@ async fn get_network_mapping_inner(
modified_timestamp: added_timestamp,
modified_signature: added_signature,
},
blocks_mined,
deleted_timestamp,
deleted_block,
monitoring,
},
)
.await
{
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());
}
get_signed_monitor_state(
unlocked_stream,
command_map,
db,
&wallet.saved.short_address,
connections_key,
only_address,
)
.await?;
if !requested_address_found {
return Err(REQUESTED_MAPPING_RECORD_PENDING.to_string());
}
if only_address.is_none() {
let expected_targets = snapshot
.iter()
.map(|state| (state.edit.address.clone(), state.edit.ip.clone()))
.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(())
}
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::connections::initialize_node_runtime_context;
use crate::records::memory::mempool::{init_db, setup_mempool};
use crate::records::memory::network_mapping::NodeInfo;
use crate::records::memory::response_channels::Command;
use crate::records::record_chain::rewards_tx::rebuild_immature_reward_window;
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}");
}
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 map: Arc<Mutex<Command>> = Arc::new(Mutex::new(HashMap::new()));
// 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;
pub async fn is_within_30_days(timestamp: u32) -> bool {
// Convert the on-chain timestamp into UTC and reject values that fall
// 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);
is_within_30_days_at(timestamp, Utc::now().timestamp() as u32)
}
match transaction_time {
Some(transaction_time) => {
let current_time = Utc::now();
let duration = current_time.signed_duration_since(transaction_time);
pub fn is_within_30_days_at(timestamp: u32, as_of_timestamp: u32) -> bool {
// Convert the on-chain timestamp into UTC and reject values that fall
// more than 30 days before the supplied validation time or fail to decode.
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
}
None => false,
_ => false,
}
}

View File

@ -229,7 +229,7 @@ async fn verify_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) => {
reserve_verified_transaction(db, &transaction, balance_tracker, already_in_mempool)
.await?;

View File

@ -9,7 +9,7 @@ use crate::records::wallet_registry::{
};
use crate::sled::Db;
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::{
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 {
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;
// Both swap participants must provide valid wallet addresses and
// 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
// 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(
"Timestamp is to old. Transactions must be broadcast within 30 days of signing."
.to_string(),
);
}
let now = Utc::now().timestamp() as u32;
let offer_expiration = self.unsigned_swap.offer_expiration;
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());
}