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..185a639 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -7,6 +7,10 @@ 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 - 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. @@ -18,6 +22,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/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..5b29fc8 100644 --- a/ant-cli/src/main.rs +++ b/ant-cli/src/main.rs @@ -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 @@ -454,39 +454,6 @@ fn resolve_evm_network( } } -/// 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." - ) -} - async fn create_client_node( bootstrap: &[SocketAddr], allow_loopback: bool, 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..0071de6 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::{DevnetManifest, MultiAddr}; use crate::error::{Error, Result}; /// Returns the platform-appropriate data directory for ant. @@ -79,6 +80,50 @@ 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) +} + #[derive(serde::Deserialize)] struct BootstrapConfig { peers: Vec, @@ -117,6 +162,53 @@ 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]); + } + + #[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..cfa212d 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,11 @@ 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("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();