diff --git a/.gitignore b/.gitignore index fc94dea..8f84e46 100644 --- a/.gitignore +++ b/.gitignore @@ -1,5 +1,7 @@ /target .cargo/config.toml +.cache/ +*.orig .claude/plans/ .claude/scheduled_tasks.lock .claude/settings.local.json diff --git a/CHANGELOG.md b/CHANGELOG.md index cd8eb0e..08e0c27 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -7,7 +7,12 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 ## [Unreleased] +### Fixed +- `ant node start`/`ant node stop` with `--service-name` now resolve the node ID through the daemon API instead of reading `node_registry.json` directly, eliminating a race against concurrent registry mutations by the daemon. +- `ant node add --json` no longer interleaves binary-download progress with the JSON result; progress output now goes to stderr (was stdout) and is suppressed entirely in JSON mode. + ### Changed +- `--evm-network` still defaults to `arbitrum-one`, **except** when a devnet manifest carrying an EVM block is loaded: that combination now errors and asks for an explicit choice (`local` to use the manifest, or a preset to override it). The old behavior silently overrode the manifest's EVM config, producing no-op mainnet-vault transactions on other chains that spent gas and set a useless ANT allowance before every chunk PUT failed payment verification. An explicit preset selected alongside a manifest EVM block now prints a warning that the manifest's EVM config is ignored. No change for mainnet users or read-only operations. - Default network binding changed from IPv4-only to IPv6 dual-stack. Hosts without a working IPv6 stack should pass `--ipv4-only` to avoid advertising unreachable v6 addresses to the DHT (which causes slow connects and junk address records). - `ant file upload` now writes datamaps as `..datamap` instead of stripping the extension. Uploading `photo.jpg` produces `photo.jpg.datamap` (was `photo.datamap`). Existing datamaps remain readable. - `ant file upload` no longer silently overwrites an existing datamap. Repeated uploads of the same source path produce `name-2.datamap`, `name-3.datamap`, … capped at 100 attempts. Pass `--overwrite` to restore the previous behaviour. @@ -18,6 +23,7 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 - `ant file download --datamap` now reads both msgpack (canonical) and legacy JSON datamaps, so datamaps produced by older versions of the GUI download cleanly via the CLI. ### Internal +- CLI-audit thinning (V2-189): `PortRange` parsing (`FromStr`), env `KEY=VALUE` parsing (`AddNodeOpts::parse_env_vars`), and bootstrap-peer resolution (`config::resolve_bootstrap_peers`) moved from ant-cli into ant-core; `node add`/`node reset` daemon calls now go through `ant_core::node::daemon::client` (new `add_node`/`reset`/`resolve_node_id_by_name` functions) instead of hand-rolled HTTP; the two CLI `ProgressReporter` impls collapsed into one. - New `ant_core::datamap_file` module owns the on-disk datamap format (msgpack canonical, JSON legacy auto-detect on read) and naming convention. `ant-cli` and consumers like `ant-gui` route through this single helper instead of reimplementing serialization. ## [0.1.1] - 2026-03-28 diff --git a/Cargo.lock b/Cargo.lock index 20ed608..6c2eded 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -824,7 +824,6 @@ dependencies = [ "colored", "hex", "indicatif", - "reqwest 0.12.28", "rmp-serde", "serde", "serde_json", diff --git a/ant-cli/Cargo.toml b/ant-cli/Cargo.toml index 39b3cd8..34303c9 100644 --- a/ant-cli/Cargo.toml +++ b/ant-cli/Cargo.toml @@ -19,7 +19,6 @@ clap = { version = "4", features = ["derive", "env"] } colored = "3.1.1" hex = "0.4" indicatif = "0.17" -reqwest = { version = "0.12", features = ["json"] } rmp-serde = "1" serde = { version = "1", features = ["derive"] } serde_json = "1" diff --git a/ant-cli/src/cli.rs b/ant-cli/src/cli.rs index 109ddb6..39cfb5d 100644 --- a/ant-cli/src/cli.rs +++ b/ant-cli/src/cli.rs @@ -76,9 +76,13 @@ pub struct Cli { #[arg(short, long, action = ArgAction::Count)] pub verbose: u8, - /// EVM network for payment processing (arbitrum-one, arbitrum-sepolia, local). - #[arg(long, default_value = "arbitrum-one")] - pub evm_network: String, + /// EVM network for payment processing: arbitrum-one, arbitrum-sepolia, + /// or local (reads the devnet manifest's EVM config). Defaults to + /// arbitrum-one — except when a devnet manifest with EVM config is + /// loaded, where an explicit choice is required so an irreversible + /// on-chain payment never silently targets the wrong network. + #[arg(long)] + pub evm_network: Option, #[command(subcommand)] pub command: Commands, diff --git a/ant-cli/src/commands/node/add.rs b/ant-cli/src/commands/node/add.rs index d9fea24..5dc61e7 100644 --- a/ant-cli/src/commands/node/add.rs +++ b/ant-cli/src/commands/node/add.rs @@ -3,7 +3,7 @@ use std::path::PathBuf; use clap::Args; use colored::Colorize; -use ant_core::node::binary::ProgressReporter; +use ant_core::node::binary::{NoopProgress, ProgressReporter}; use ant_core::node::daemon::client; use ant_core::node::types::DaemonConfig; use ant_core::node::types::{ @@ -103,8 +103,8 @@ impl AddArgs { // Check if daemon is running; if so, POST to API; otherwise call directly let config = DaemonConfig::default(); let result = match client::status(&config).await { - Ok(status) if status.running => self.add_via_daemon(&config, &opts).await?, - _ => self.add_directly(&config, &opts).await?, + Ok(status) if status.running => client::add_node(&config, &opts).await?, + _ => self.add_directly(&config, &opts, json_output).await?, }; if json_output { @@ -150,7 +150,11 @@ impl AddArgs { } fn to_add_node_opts(&self) -> anyhow::Result { - let node_port = self.parse_port_range(&self.node_port)?; + let node_port = self + .node_port + .as_deref() + .map(str::parse::) + .transpose()?; let binary_source = if let Some(ref path) = self.path { BinarySource::LocalPath(path.clone()) @@ -162,18 +166,7 @@ impl AddArgs { BinarySource::Latest }; - let env_variables: Vec<(String, String)> = self - .env - .iter() - .map(|e| { - let parts: Vec<&str> = e.splitn(2, '=').collect(); - if parts.len() == 2 { - Ok((parts[0].to_string(), parts[1].to_string())) - } else { - anyhow::bail!("Invalid env variable format: '{e}'. Expected KEY=VALUE") - } - }) - .collect::>>()?; + let env_variables = AddNodeOpts::parse_env_vars(&self.env)?; Ok(AddNodeOpts { count: self.count, @@ -189,92 +182,21 @@ impl AddArgs { }) } - fn parse_port_range(&self, input: &Option) -> anyhow::Result> { - match input { - None => Ok(None), - Some(s) => { - if let Some((start, end)) = s.split_once('-') { - let start: u16 = start - .parse() - .map_err(|_| anyhow::anyhow!("Invalid port range start: '{start}'"))?; - let end: u16 = end - .parse() - .map_err(|_| anyhow::anyhow!("Invalid port range end: '{end}'"))?; - if end < start { - anyhow::bail!("Port range end ({end}) must be >= start ({start})"); - } - Ok(Some(PortRange::Range(start, end))) - } else { - let port: u16 = s - .parse() - .map_err(|_| anyhow::anyhow!("Invalid port: '{s}'"))?; - Ok(Some(PortRange::Single(port))) - } - } - } - } - - async fn add_via_daemon( - &self, - config: &DaemonConfig, - opts: &AddNodeOpts, - ) -> anyhow::Result { - let info = client::info(config); - let api_base = info - .api_base - .ok_or_else(|| anyhow::anyhow!("Daemon is running but API base URL not available"))?; - - let client = reqwest::Client::new(); - let resp = client - .post(format!("{api_base}/nodes")) - .json(opts) - .send() - .await?; - - if resp.status().is_success() { - Ok(resp.json().await?) - } else { - let body = resp.text().await?; - anyhow::bail!("Daemon returned error: {body}"); - } - } - async fn add_directly( &self, config: &DaemonConfig, opts: &AddNodeOpts, + json_output: bool, ) -> anyhow::Result { - let progress = CliProgress; + // Suppress progress in JSON mode so stdout stays parseable. + let progress: Box = if json_output { + Box::new(NoopProgress) + } else { + Box::new(crate::progress::CliProgress) + }; let result = - ant_core::node::add_nodes(opts.clone(), &config.registry_path, &progress).await?; + ant_core::node::add_nodes(opts.clone(), &config.registry_path, progress.as_ref()) + .await?; Ok(result) } } - -/// CLI progress reporter that prints to the terminal. -struct CliProgress; - -impl ProgressReporter for CliProgress { - fn report_started(&self, message: &str) { - println!("{} {message}", "⟳".cyan()); - } - - fn report_progress(&self, bytes: u64, total: u64) { - if total > 0 { - let pct = (bytes as f64 / total as f64 * 100.0) as u32; - let bar_width = 30; - let filled = (pct as usize * bar_width) / 100; - let empty = bar_width - filled; - let bar = format!( - "{}{}", - "█".repeat(filled).cyan(), - "░".repeat(empty).dimmed() - ); - print!("\r {} {bar} {pct:>3}%", "Downloading".dimmed()); - } - } - - fn report_complete(&self, message: &str) { - println!("\r{} {message}", "✓".green().bold()); - } -} diff --git a/ant-cli/src/commands/node/reset.rs b/ant-cli/src/commands/node/reset.rs index f6f2ad8..2d17d90 100644 --- a/ant-cli/src/commands/node/reset.rs +++ b/ant-cli/src/commands/node/reset.rs @@ -48,7 +48,7 @@ impl ResetArgs { } let result = if daemon_running { - self.reset_via_daemon(&config).await? + client::reset(&config).await? } else { self.reset_directly(&config)? }; @@ -85,23 +85,6 @@ impl ResetArgs { Ok(()) } - async fn reset_via_daemon(&self, config: &DaemonConfig) -> anyhow::Result { - let info = client::info(config); - let api_base = info - .api_base - .ok_or_else(|| anyhow::anyhow!("Daemon is running but API base URL not available"))?; - - let client = reqwest::Client::new(); - let resp = client.post(format!("{api_base}/reset")).send().await?; - - if resp.status().is_success() { - Ok(resp.json().await?) - } else { - let body = resp.text().await?; - anyhow::bail!("Daemon returned error: {body}"); - } - } - fn reset_directly(&self, config: &DaemonConfig) -> anyhow::Result { let result = ant_core::node::reset(&config.registry_path)?; Ok(result) diff --git a/ant-cli/src/commands/node/start.rs b/ant-cli/src/commands/node/start.rs index 506e209..75cb51e 100644 --- a/ant-cli/src/commands/node/start.rs +++ b/ant-cli/src/commands/node/start.rs @@ -33,12 +33,9 @@ impl StartArgs { service_name: &str, json_output: bool, ) -> anyhow::Result<()> { - // Look up node ID by service name via the registry - let registry = ant_core::node::registry::NodeRegistry::load(&config.registry_path)?; - let node = registry - .find_by_service_name(service_name) - .ok_or_else(|| anyhow::anyhow!("No node found with service name '{service_name}'"))?; - let node_id = node.id; + // Resolve the node ID through the daemon API — the daemon owns the + // registry, so the CLI must not read node_registry.json directly. + let node_id = client::resolve_node_id_by_name(config, service_name).await?; let result = client::start_node(config, node_id).await?; diff --git a/ant-cli/src/commands/node/stop.rs b/ant-cli/src/commands/node/stop.rs index b522734..8c2bfc0 100644 --- a/ant-cli/src/commands/node/stop.rs +++ b/ant-cli/src/commands/node/stop.rs @@ -33,11 +33,9 @@ impl StopArgs { service_name: &str, json_output: bool, ) -> anyhow::Result<()> { - let registry = ant_core::node::registry::NodeRegistry::load(&config.registry_path)?; - let node = registry - .find_by_service_name(service_name) - .ok_or_else(|| anyhow::anyhow!("No node found with service name '{service_name}'"))?; - let node_id = node.id; + // Resolve the node ID through the daemon API — the daemon owns the + // registry, so the CLI must not read node_registry.json directly. + let node_id = client::resolve_node_id_by_name(config, service_name).await?; let result = client::stop_node(config, node_id).await?; diff --git a/ant-cli/src/commands/update.rs b/ant-cli/src/commands/update.rs index 45de66c..3b65070 100644 --- a/ant-cli/src/commands/update.rs +++ b/ant-cli/src/commands/update.rs @@ -4,25 +4,7 @@ use colored::Colorize; use ant_core::node::binary::NoopProgress; use ant_core::update; -/// Progress reporter that prints to the terminal. -struct CliUpdateProgress; - -impl ant_core::node::binary::ProgressReporter for CliUpdateProgress { - fn report_started(&self, message: &str) { - eprintln!("{}", message.dimmed()); - } - - fn report_progress(&self, bytes: u64, total: u64) { - if total > 0 { - let pct = (bytes as f64 / total as f64 * 100.0) as u64; - eprint!("\r{}", format!(" Downloading... {pct}%").dimmed()); - } - } - - fn report_complete(&self, message: &str) { - eprintln!("\r{}", message.green()); - } -} +use crate::progress::CliProgress; #[derive(Args)] pub struct UpdateArgs { @@ -72,7 +54,7 @@ impl UpdateArgs { let progress: Box = if json_output { Box::new(NoopProgress) } else { - Box::new(CliUpdateProgress) + Box::new(CliProgress) }; let result = update::perform_update(&check, progress.as_ref()).await?; diff --git a/ant-cli/src/main.rs b/ant-cli/src/main.rs index b8523a3..fb21be2 100644 --- a/ant-cli/src/main.rs +++ b/ant-cli/src/main.rs @@ -13,8 +13,8 @@ use tracing_subscriber::{fmt, layer::SubscriberExt, util::SubscriberInitExt, Env use ant_core::data::{ peer_cache::{self, BootstrapAddressFilter}, - Client, ClientConfig, CoreNodeConfig, CustomNetwork, DevnetManifest, EvmAddress, EvmNetwork, - IPDiversityConfig, MultiAddr, NodeMode, P2PNode, Wallet, MAX_WIRE_MESSAGE_SIZE, + Client, ClientConfig, CoreNodeConfig, DevnetManifest, EvmNetwork, IPDiversityConfig, MultiAddr, + NodeMode, P2PNode, Wallet, MAX_WIRE_MESSAGE_SIZE, }; use cli::{Cli, Commands}; @@ -180,7 +180,7 @@ struct DataCliContext { quote_timeout_secs: u64, store_timeout_secs: Option, chunk_get_timeout_secs: Option, - evm_network: String, + evm_network: Option, quote_concurrency: Option, store_concurrency: Option, } @@ -209,7 +209,7 @@ async fn build_data_client( } let manifest = load_manifest(ctx)?; - let bootstrap = resolve_bootstrap_from(ctx, manifest.as_ref())?; + let bootstrap = ant_core::config::resolve_bootstrap_peers(&ctx.bootstrap, manifest.as_ref())?; // Explicit network selectors should be isolated from the general client // peer cache. `--bootstrap` and `--devnet-manifest` both mean "use exactly // this network entrypoint", so cached public-network peers must not be @@ -357,7 +357,7 @@ async fn build_data_client( let key = private_key .as_ref() .ok_or_else(|| anyhow::anyhow!("SECRET_KEY environment variable required"))?; - let network = resolve_evm_network(&ctx.evm_network, manifest.as_ref())?; + let network = resolve_evm_network(ctx.evm_network.as_deref(), manifest.as_ref())?; let wallet = create_wallet(key, network)?; info!("Wallet configured for EVM payments"); client = client.with_wallet(wallet); @@ -411,80 +411,30 @@ fn resolve_evm_network_and_manifest( ctx: &DataCliContext, ) -> anyhow::Result<(EvmNetwork, Option)> { let manifest = load_manifest(ctx)?; - let network = resolve_evm_network(&ctx.evm_network, manifest.as_ref())?; + let network = resolve_evm_network(ctx.evm_network.as_deref(), manifest.as_ref())?; Ok((network, manifest)) } +/// Resolve the EVM network through ant-core, with a CLI-side warning for +/// the one remaining silent-discard path: a devnet manifest that carries +/// an EVM block while a preset network is selected (V2-471). fn resolve_evm_network( - evm_network: &str, + evm_network: Option<&str>, manifest: Option<&DevnetManifest>, ) -> anyhow::Result { - match evm_network { - "arbitrum-one" => Ok(EvmNetwork::ArbitrumOne), - "arbitrum-sepolia" => Ok(EvmNetwork::ArbitrumSepoliaTest), - "local" => { - if let Some(m) = manifest { - if let Some(ref evm) = m.evm { - let rpc_url: reqwest::Url = evm - .rpc_url - .parse() - .map_err(|e| anyhow::anyhow!("Invalid RPC URL: {e}"))?; - let token_addr: EvmAddress = evm - .payment_token_address - .parse() - .map_err(|e| anyhow::anyhow!("Invalid token address: {e}"))?; - let vault_addr: EvmAddress = evm - .payment_vault_address - .parse() - .map_err(|e| anyhow::anyhow!("Invalid payment vault address: {e}"))?; - return Ok(EvmNetwork::Custom(CustomNetwork { - rpc_url_http: rpc_url, - payment_token_address: token_addr, - payment_vault_address: vault_addr, - })); - } - } - anyhow::bail!("EVM network 'local' requires --devnet-manifest with EVM info") - } - other => { - anyhow::bail!( - "Unsupported EVM network: {other}. Use 'arbitrum-one', 'arbitrum-sepolia', or 'local'." - ) + if let (Some(m), Some(name)) = (manifest, evm_network) { + if m.evm.is_some() && name != "local" { + eprintln!( + "warning: the devnet manifest contains an EVM block, but \ + --evm-network={name} selects a preset; the manifest's EVM \ + config is ignored. Pass --evm-network local to use it." + ); } } -} - -/// Resolve bootstrap peers from a pre-loaded manifest. -/// -/// Priority: CLI `--bootstrap` > devnet manifest > `bootstrap_peers.toml` config file. -fn resolve_bootstrap_from( - ctx: &DataCliContext, - manifest: Option<&DevnetManifest>, -) -> anyhow::Result> { - if !ctx.bootstrap.is_empty() { - return Ok(ctx.bootstrap.clone()); - } - - if let Some(m) = manifest { - let bootstrap: Vec = m - .bootstrap - .iter() - .filter_map(MultiAddr::socket_addr) - .collect(); - return Ok(bootstrap); - } - - if let Some(peers) = ant_core::config::load_bootstrap_peers() - .map_err(|e| anyhow::anyhow!("Failed to load bootstrap config: {e}"))? - { - info!("Loaded {} bootstrap peer(s) from config file", peers.len()); - return Ok(peers); - } - - anyhow::bail!( - "No bootstrap peers provided. Use --bootstrap, --devnet-manifest, \ - or install bootstrap_peers.toml to your config directory." - ) + Ok(ant_core::config::resolve_evm_network( + evm_network, + manifest, + )?) } async fn create_client_node( diff --git a/ant-cli/src/progress.rs b/ant-cli/src/progress.rs index 1711266..d8a0177 100644 --- a/ant-cli/src/progress.rs +++ b/ant-cli/src/progress.rs @@ -13,9 +13,12 @@ use std::io::{self, IsTerminal, Write}; use std::sync::OnceLock; use std::time::Duration; +use colored::Colorize; use indicatif::{MultiProgress, ProgressBar, ProgressStyle}; use tracing_subscriber::fmt::MakeWriter; +use ant_core::node::binary::ProgressReporter; + static MULTI: OnceLock = OnceLock::new(); /// The shared `MultiProgress` instance. Created on first access. @@ -51,6 +54,36 @@ pub fn attach(pb: ProgressBar) -> ProgressBar { } } +/// Terminal implementation of ant-core's `ProgressReporter` (binary +/// downloads during `node add` and self-update). Writes to stderr so +/// stdout stays clean for command output. +pub struct CliProgress; + +impl ProgressReporter for CliProgress { + fn report_started(&self, message: &str) { + eprintln!("{} {message}", "⟳".cyan()); + } + + fn report_progress(&self, bytes: u64, total: u64) { + if total > 0 { + let pct = (bytes as f64 / total as f64 * 100.0) as u32; + let bar_width = 30; + let filled = (pct as usize * bar_width) / 100; + let empty = bar_width - filled; + let bar = format!( + "{}{}", + "█".repeat(filled).cyan(), + "░".repeat(empty).dimmed() + ); + eprint!("\r {} {bar} {pct:>3}%", "Downloading".dimmed()); + } + } + + fn report_complete(&self, message: &str) { + eprintln!("\r{} {message}", "✓".green().bold()); + } +} + /// Tracing writer that suspends active progress bars while writing log lines. #[derive(Clone)] pub struct ProgressAwareWriter; diff --git a/ant-core/src/config.rs b/ant-core/src/config.rs index 806f15a..a2baccb 100644 --- a/ant-core/src/config.rs +++ b/ant-core/src/config.rs @@ -1,6 +1,7 @@ use std::net::SocketAddr; use std::path::PathBuf; +use crate::data::{CustomNetwork, DevnetManifest, EvmAddress, EvmNetwork, MultiAddr}; use crate::error::{Error, Result}; /// Returns the platform-appropriate data directory for ant. @@ -79,6 +80,104 @@ pub fn load_bootstrap_peers() -> Result>> { Ok(Some(addrs)) } +/// Resolve the bootstrap peers for a client connection. +/// +/// Priority: explicitly supplied peers (e.g. a frontend's `--bootstrap` +/// flag) > devnet manifest peers > the platform `bootstrap_peers.toml` +/// config file. Manifest peers without a resolvable socket address are +/// filtered out. A selected manifest is authoritative: if it yields no +/// usable peers, resolution fails rather than falling back to the +/// config file. +/// +/// # Errors +/// +/// Returns [`Error::NoBootstrapPeers`] when the selected source yields +/// no peers, and propagates config-file read/parse failures. +pub fn resolve_bootstrap_peers( + explicit: &[SocketAddr], + manifest: Option<&DevnetManifest>, +) -> Result> { + if !explicit.is_empty() { + return Ok(explicit.to_vec()); + } + + if let Some(m) = manifest { + let peers: Vec = m + .bootstrap + .iter() + .filter_map(MultiAddr::socket_addr) + .collect(); + // An explicitly selected manifest never falls back to the public + // config: an empty (or fully filtered) manifest is an error here, + // not later when the first data operation fails. + if peers.is_empty() { + return Err(Error::NoBootstrapPeers); + } + return Ok(peers); + } + + if let Some(peers) = load_bootstrap_peers()? { + tracing::info!("Loaded {} bootstrap peer(s) from config file", peers.len()); + return Ok(peers); + } + + Err(Error::NoBootstrapPeers) +} + +/// Resolve the EVM network for payment operations. +/// +/// `name` is the frontend's network selector (e.g. the CLI's +/// `--evm-network` flag): `arbitrum-one`, `arbitrum-sepolia`, or `local`, +/// which reads the RPC URL and contract addresses from the devnet +/// manifest's `evm` block. +/// +/// With no selector, the default is Arbitrum One (mainnet) — **unless** +/// the devnet manifest carries an `evm` block. Defaulting to mainnet +/// used to silently discard that config: the mainnet vault address +/// doesn't exist on other chains, so payments "succeeded" as no-op +/// transactions and every subsequent chunk PUT failed verification +/// (V2-471). In that one ambiguous case the caller must choose +/// explicitly or get [`Error::EvmNetworkAmbiguous`]. +pub fn resolve_evm_network( + name: Option<&str>, + manifest: Option<&DevnetManifest>, +) -> Result { + match name { + None => { + if manifest.is_some_and(|m| m.evm.is_some()) { + Err(Error::EvmNetworkAmbiguous) + } else { + Ok(EvmNetwork::ArbitrumOne) + } + } + Some("arbitrum-one") => Ok(EvmNetwork::ArbitrumOne), + Some("arbitrum-sepolia") => Ok(EvmNetwork::ArbitrumSepoliaTest), + Some("local") => { + let evm = manifest + .and_then(|m| m.evm.as_ref()) + .ok_or(Error::EvmManifestRequired)?; + let rpc_url: reqwest::Url = evm + .rpc_url + .parse() + .map_err(|e| Error::InvalidEvmManifest(format!("invalid RPC URL: {e}")))?; + let payment_token_address: EvmAddress = + evm.payment_token_address.parse().map_err(|e| { + Error::InvalidEvmManifest(format!("invalid payment token address: {e}")) + })?; + let payment_vault_address: EvmAddress = + evm.payment_vault_address.parse().map_err(|e| { + Error::InvalidEvmManifest(format!("invalid payment vault address: {e}")) + })?; + Ok(EvmNetwork::Custom(CustomNetwork { + rpc_url_http: rpc_url, + payment_token_address, + payment_vault_address, + })) + } + Some(other) => Err(Error::UnsupportedEvmNetwork(other.to_string())), + } +} + #[derive(serde::Deserialize)] struct BootstrapConfig { peers: Vec, @@ -117,6 +216,144 @@ mod tests { ); } + fn test_manifest(addrs: Vec) -> DevnetManifest { + DevnetManifest { + base_port: 10000, + node_count: addrs.len(), + bootstrap: addrs.into_iter().map(MultiAddr::quic).collect(), + data_dir: PathBuf::new(), + created_at: String::new(), + evm: None, + } + } + + #[test] + fn resolve_bootstrap_prefers_explicit_peers() { + let explicit: Vec = vec!["10.0.0.1:10000".parse().unwrap()]; + let manifest = test_manifest(vec!["10.0.0.2:10000".parse().unwrap()]); + let peers = resolve_bootstrap_peers(&explicit, Some(&manifest)).unwrap(); + assert_eq!(peers, explicit); + } + + #[test] + fn resolve_bootstrap_uses_manifest_when_no_explicit_peers() { + let addr: SocketAddr = "10.0.0.2:10000".parse().unwrap(); + let manifest = test_manifest(vec![addr]); + let peers = resolve_bootstrap_peers(&[], Some(&manifest)).unwrap(); + assert_eq!(peers, vec![addr]); + } + + fn manifest_with_evm() -> DevnetManifest { + let mut m = test_manifest(vec!["10.0.0.2:10000".parse().unwrap()]); + m.evm = Some(ant_protocol::DevnetEvmInfo { + rpc_url: "http://127.0.0.1:8545".to_string(), + wallet_private_key: + "0xac0974bec39a17e36ba4a6b4d238ff944bacb478cbed5efcae784d7bf4f2ff80".to_string(), + payment_token_address: "0x5FbDB2315678afecb367f032d93F642f64180aa3".to_string(), + payment_vault_address: "0xe7f1725E7734CE288F8367e1Bb143E90bb3F0512".to_string(), + }); + m + } + + #[test] + fn resolve_evm_network_defaults_to_mainnet_without_manifest_evm() { + assert!(matches!( + resolve_evm_network(None, None), + Ok(EvmNetwork::ArbitrumOne) + )); + let no_evm = test_manifest(vec![]); + assert!(matches!( + resolve_evm_network(None, Some(&no_evm)), + Ok(EvmNetwork::ArbitrumOne) + )); + } + + #[test] + fn resolve_evm_network_requires_choice_when_manifest_has_evm() { + // A manifest with an EVM block plus no explicit selection is the + // V2-471 trap: refuse rather than silently pay against mainnet. + assert!(matches!( + resolve_evm_network(None, Some(&manifest_with_evm())), + Err(Error::EvmNetworkAmbiguous) + )); + } + + #[test] + fn resolve_evm_network_maps_presets() { + assert!(matches!( + resolve_evm_network(Some("arbitrum-one"), None), + Ok(EvmNetwork::ArbitrumOne) + )); + assert!(matches!( + resolve_evm_network(Some("arbitrum-sepolia"), None), + Ok(EvmNetwork::ArbitrumSepoliaTest) + )); + } + + #[test] + fn resolve_evm_network_local_reads_manifest() { + let manifest = manifest_with_evm(); + let network = resolve_evm_network(Some("local"), Some(&manifest)).unwrap(); + match network { + EvmNetwork::Custom(custom) => { + assert_eq!(custom.rpc_url_http.as_str(), "http://127.0.0.1:8545/"); + assert_eq!( + format!("{:?}", custom.payment_token_address).to_lowercase(), + "0x5fbdb2315678afecb367f032d93f642f64180aa3" + ); + } + other => panic!("expected Custom network, got {other:?}"), + } + } + + #[test] + fn resolve_evm_network_local_requires_manifest_evm_block() { + // No manifest at all, and a manifest without an evm block. + assert!(matches!( + resolve_evm_network(Some("local"), None), + Err(Error::EvmManifestRequired) + )); + let no_evm = test_manifest(vec![]); + assert!(matches!( + resolve_evm_network(Some("local"), Some(&no_evm)), + Err(Error::EvmManifestRequired) + )); + } + + #[test] + fn resolve_evm_network_rejects_unknown_and_bad_manifest_values() { + assert!(matches!( + resolve_evm_network(Some("mainnet"), None), + Err(Error::UnsupportedEvmNetwork(_)) + )); + let mut bad = manifest_with_evm(); + bad.evm.as_mut().unwrap().payment_vault_address = "not-an-address".to_string(); + assert!(matches!( + resolve_evm_network(Some("local"), Some(&bad)), + Err(Error::InvalidEvmManifest(_)) + )); + } + + #[test] + fn resolve_bootstrap_errors_on_empty_manifest() { + let manifest = test_manifest(vec![]); + let err = resolve_bootstrap_peers(&[], Some(&manifest)).unwrap_err(); + assert!(matches!(err, Error::NoBootstrapPeers)); + } + + #[test] + fn resolve_bootstrap_errors_when_all_manifest_peers_filtered() { + // A non-IP transport has no socket address, so the peer is + // filtered out and the manifest yields nothing usable. + let bt: MultiAddr = "/bt/00:11:22:33:44:55/rfcomm/1".parse().unwrap(); + assert!(bt.socket_addr().is_none()); + let mut manifest = test_manifest(vec![]); + manifest.bootstrap = vec![bt]; + manifest.node_count = 1; + let err = resolve_bootstrap_peers(&[], Some(&manifest)).unwrap_err(); + assert!(matches!(err, Error::NoBootstrapPeers)); + } + #[test] fn load_bootstrap_peers_returns_none_when_no_file() { // Set config dir to a temp location where no file exists diff --git a/ant-core/src/error.rs b/ant-core/src/error.rs index 2b62460..d9365b9 100644 --- a/ant-core/src/error.rs +++ b/ant-core/src/error.rs @@ -11,6 +11,9 @@ pub enum Error { #[error("Node not found: {0}")] NodeNotFound(u32), + #[error("No node found with service name '{0}'")] + NodeNotFoundByName(String), + #[error("Node already running: {0}")] NodeAlreadyRunning(u32), @@ -44,6 +47,12 @@ pub enum Error { #[error("Port range length ({range_len}) does not match node count ({count})")] PortRangeMismatch { range_len: u16, count: u16 }, + #[error("Invalid port range: {0}")] + InvalidPortRange(String), + + #[error("Invalid env variable format: '{0}'. Expected KEY=VALUE")] + InvalidEnvVar(String), + #[error("Binary not found at path: {0}")] BinaryNotFound(PathBuf), @@ -65,6 +74,25 @@ pub enum Error { #[error("Failed to parse bootstrap_peers.toml: {0}")] BootstrapConfigParse(String), + #[error( + "No bootstrap peers available: pass peers explicitly, use a devnet manifest, or install bootstrap_peers.toml in the config directory" + )] + NoBootstrapPeers, + + #[error( + "Devnet manifest contains EVM config but no EVM network was selected: pass 'local' to use the manifest's EVM config, or an explicit preset ('arbitrum-one', 'arbitrum-sepolia') to override it" + )] + EvmNetworkAmbiguous, + + #[error("Unsupported EVM network: {0}. Use 'arbitrum-one', 'arbitrum-sepolia', or 'local'.")] + UnsupportedEvmNetwork(String), + + #[error("EVM network 'local' requires a devnet manifest with EVM info")] + EvmManifestRequired, + + #[error("Invalid EVM info in devnet manifest: {0}")] + InvalidEvmManifest(String), + #[error("Node count {count} exceeds maximum of {max} per call")] InvalidNodeCount { count: u16, max: u16 }, diff --git a/ant-core/src/node/daemon/client.rs b/ant-core/src/node/daemon/client.rs index eeb81a4..90704b3 100644 --- a/ant-core/src/node/daemon/client.rs +++ b/ant-core/src/node/daemon/client.rs @@ -5,8 +5,9 @@ use crate::error::{Error, Result}; use crate::node::daemon::health::FleetHealth; use crate::node::process::detach; use crate::node::types::{ - DaemonConfig, DaemonInfo, DaemonStartResult, DaemonStatus, DaemonStopResult, NodeStarted, - NodeStatusResult, NodeStopped, RemoveNodeResult, StartNodeResult, StopNodeResult, + AddNodeOpts, AddNodeResult, DaemonConfig, DaemonInfo, DaemonStartResult, DaemonStatus, + DaemonStopResult, NodeStarted, NodeStatusResult, NodeStopped, RemoveNodeResult, ResetResult, + StartNodeResult, StopNodeResult, }; /// Get the daemon's current status by querying its REST API. @@ -205,6 +206,50 @@ pub async fn stop_node(config: &DaemonConfig, node_id: u32) -> Result Result { + let port = read_port_file(&config.port_file_path).ok_or(Error::DaemonNotRunning)?; + + let url = format!("http://127.0.0.1:{port}/api/v1/nodes"); + let resp = reqwest::Client::new() + .post(&url) + .json(opts) + .send() + .await + .map_err(|e| Error::HttpRequest(e.to_string()))?; + + if resp.status().is_success() { + resp.json::() + .await + .map_err(|e| Error::HttpRequest(e.to_string())) + } else { + let body = resp.text().await.unwrap_or_default(); + Err(Error::HttpRequest(body)) + } +} + +/// Reset all node state — clear the registry and remove node data/log +/// directories — via the daemon REST API. +pub async fn reset(config: &DaemonConfig) -> Result { + let port = read_port_file(&config.port_file_path).ok_or(Error::DaemonNotRunning)?; + + let url = format!("http://127.0.0.1:{port}/api/v1/reset"); + let resp = reqwest::Client::new() + .post(&url) + .send() + .await + .map_err(|e| Error::HttpRequest(e.to_string()))?; + + if resp.status().is_success() { + resp.json::() + .await + .map_err(|e| Error::HttpRequest(e.to_string())) + } else { + let body = resp.text().await.unwrap_or_default(); + Err(Error::HttpRequest(body)) + } +} + /// Dismiss a node — remove it from the registry — via the daemon REST API. /// /// Intended for evicted nodes (whose data directory has already been deleted), but the daemon will @@ -267,6 +312,21 @@ pub async fn node_status(config: &DaemonConfig) -> Result { } } +/// Resolve a node's ID from its service name via the daemon REST API. +/// +/// Goes through the daemon — the registry's owner — rather than reading +/// node_registry.json directly, so the lookup cannot observe a +/// half-written registry or race a concurrent mutation by the daemon. +pub async fn resolve_node_id_by_name(config: &DaemonConfig, service_name: &str) -> Result { + let status = node_status(config).await?; + status + .nodes + .iter() + .find(|n| n.name == service_name) + .map(|n| n.node_id) + .ok_or_else(|| Error::NodeNotFoundByName(service_name.to_string())) +} + /// Stop all running nodes via the daemon REST API. pub async fn stop_all_nodes(config: &DaemonConfig) -> Result { let port = read_port_file(&config.port_file_path).ok_or(Error::DaemonNotRunning)?; diff --git a/ant-core/src/node/types.rs b/ant-core/src/node/types.rs index 70481e5..1661969 100644 --- a/ant-core/src/node/types.rs +++ b/ant-core/src/node/types.rs @@ -279,6 +279,36 @@ impl PortRange { } } +impl std::str::FromStr for PortRange { + type Err = crate::error::Error; + + /// Parse `"12000"` into [`PortRange::Single`] or `"12000-12004"` into + /// [`PortRange::Range`]. + fn from_str(s: &str) -> std::result::Result { + use crate::error::Error; + + if let Some((start, end)) = s.split_once('-') { + let start: u16 = start + .parse() + .map_err(|_| Error::InvalidPortRange(format!("invalid start port '{start}'")))?; + let end: u16 = end + .parse() + .map_err(|_| Error::InvalidPortRange(format!("invalid end port '{end}'")))?; + if end < start { + return Err(Error::InvalidPortRange(format!( + "end ({end}) must be >= start ({start})" + ))); + } + Ok(Self::Range(start, end)) + } else { + let port: u16 = s + .parse() + .map_err(|_| Error::InvalidPortRange(format!("invalid port '{s}'")))?; + Ok(Self::Single(port)) + } + } +} + /// Options for adding one or more nodes to the registry. #[derive(Debug, Clone, Serialize, Deserialize, utoipa::ToSchema)] pub struct AddNodeOpts { @@ -310,6 +340,20 @@ pub struct AddNodeOpts { pub evm_network: EvmNetwork, } +impl AddNodeOpts { + /// Parse `KEY=VALUE` strings (the format frontends accept for node env + /// vars) into the pair list `env_variables` expects. + pub fn parse_env_vars(vars: &[String]) -> crate::error::Result> { + vars.iter() + .map(|e| { + e.split_once('=') + .map(|(k, v)| (k.to_string(), v.to_string())) + .ok_or_else(|| crate::error::Error::InvalidEnvVar(e.clone())) + }) + .collect() + } +} + impl Default for AddNodeOpts { fn default() -> Self { Self { @@ -586,6 +630,33 @@ mod tests { assert_eq!(pr.port_at(3), None); } + #[test] + fn port_range_parses_single() { + let pr: PortRange = "12000".parse().unwrap(); + assert!(matches!(pr, PortRange::Single(12000))); + } + + #[test] + fn port_range_parses_range() { + let pr: PortRange = "12000-12004".parse().unwrap(); + assert!(matches!(pr, PortRange::Range(12000, 12004))); + } + + #[test] + fn port_range_rejects_inverted_range() { + let err = "12004-12000".parse::().unwrap_err(); + assert!(err.to_string().contains("must be >= start")); + } + + #[test] + fn port_range_rejects_garbage() { + assert!("abc".parse::().is_err()); + assert!("".parse::().is_err()); + assert!("12000-abc".parse::().is_err()); + assert!("-12000".parse::().is_err()); + assert!("70000".parse::().is_err()); + } + #[test] fn binary_source_serializes_with_tag() { let src = BinarySource::Latest; @@ -598,6 +669,32 @@ mod tests { assert!(json.contains("1.0.0")); } + #[test] + fn parse_env_vars_splits_on_first_equals() { + let parsed = AddNodeOpts::parse_env_vars(&[ + "KEY=VALUE".to_string(), + "RUST_LOG=info,ant_node=debug".to_string(), + "EMPTY=".to_string(), + "URL=http://host?a=b".to_string(), + ]) + .unwrap(); + assert_eq!( + parsed, + vec![ + ("KEY".to_string(), "VALUE".to_string()), + ("RUST_LOG".to_string(), "info,ant_node=debug".to_string()), + ("EMPTY".to_string(), String::new()), + ("URL".to_string(), "http://host?a=b".to_string()), + ] + ); + } + + #[test] + fn parse_env_vars_rejects_missing_equals() { + let err = AddNodeOpts::parse_env_vars(&["NOVALUE".to_string()]).unwrap_err(); + assert!(err.to_string().contains("Expected KEY=VALUE")); + } + #[test] fn add_node_opts_default() { let opts = AddNodeOpts::default();