#![forbid(unsafe_op_in_unsafe_fn)]
mod child_manager;
use child_manager::*;
#[cfg(windows)]
mod windows_kill_tree;
mod shutdown;
use shutdown::*;
mod utilities;
use utilities::*;
use anyhow::{Context as AnyhowContext, Result as AnyhowResult, anyhow, bail, ensure};
use chrono::Local;
use clap::Parser;
use dunce::canonicalize;
use itertools::Itertools;
use parking_lot::Mutex;
use std::{
env,
path::{Path, PathBuf},
process::{Child, Command as StdCommand, Stdio},
sync::Arc,
time::Duration,
};
use tracing::{self, Span};
#[cfg(unix)]
use {libc::setsid, std::os::unix::process::CommandExt};
use super::*;
use leo_ast::NetworkName;
const REST_RPS: &str = "999999999";
#[derive(Parser, Debug)]
pub struct LeoDevnet {
#[clap(long, help = "Number of validators", default_value = "4")]
pub(crate) num_validators: usize,
#[clap(long, help = "Number of clients", default_value = "2")]
pub(crate) num_clients: usize,
#[clap(short = 'n', long, help = "Network (mainnet=0, testnet=1, canary=2)", default_value = "testnet")]
pub(crate) network: NetworkName,
#[clap(short = 's', long, help = "Ledger / log root directory", default_value = "./")]
pub(crate) storage: PathBuf,
#[clap(long, help = "Path to snarkOS binary. If it does not exist, set `--install` to build it at this path.")]
pub(crate) snarkos: Option<PathBuf>,
#[clap(long, help = "Required features for snarkOS (e.g. `test_network`)", value_delimiter = ',')]
pub(crate) snarkos_features: Vec<String>,
#[clap(long, help = "Required version for snarkOS (e.g. `4.1.0`). Defaults to latest version on `crates.io`.")]
pub(crate) snarkos_version: Option<String>,
#[clap(long, help = "(Re)install snarkOS at the provided `--snarkos` path with the given `--snarkos-features`")]
pub(crate) install: bool,
#[clap(
long,
help = "Optional consensus heights to use. The `test_network` feature must be enabled for this to work.",
value_delimiter = ',',
env = "CONSENSUS_VERSION_HEIGHTS"
)]
pub(crate) consensus_heights: Option<Vec<u32>>,
#[clap(long, help = "Run nodes in tmux (only available on Unix)")]
pub(crate) tmux: bool,
#[clap(short = 'v', long, help = "snarkOS verbosity (0-4)", default_value = "1")]
pub(crate) verbosity: u8,
#[clap(long, short = 'y', help = "Skip confirmation prompts and proceed with the devnet startup")]
pub(crate) yes: bool,
#[clap(long, help = "Base REST port (each node uses base + dev_index)")]
pub(crate) rest_port: Option<u16>,
#[clap(long, help = "Base node port (each node uses base + dev_index)")]
pub(crate) node_port: Option<u16>,
#[clap(long, help = "Base BFT port (each node uses base + dev_index)")]
pub(crate) bft_port: Option<u16>,
#[clap(long, help = "Base metrics port (each validator uses base + dev_index)")]
pub(crate) metrics_port: Option<u16>,
#[clap(short = 'c', long, help = "Remove existing devnet storage before starting")]
pub(crate) clear_storage: bool,
#[clap(long, help = "Only clean devnet storage (ledgers, node data, logs) without starting")]
pub(crate) clean_only: bool,
}
impl Command for LeoDevnet {
type Input = ();
type Output = ();
fn log_span(&self) -> Span {
tracing::span!(tracing::Level::INFO, "LeoDevnet")
}
fn prelude(&self, _: Context) -> Result<Self::Input> {
Ok(())
}
fn apply(self, _cx: Context, _: Self::Input) -> Result<Self::Output> {
self.handle_apply().map_err(|e| crate::errors::custom(format!("Failed to run devnet command: {e}")).into())
}
}
impl LeoDevnet {
fn handle_apply(&self) -> AnyhowResult<()> {
if self.clean_only {
return self.handle_clean();
}
let snarkos_path = self.snarkos.as_ref().ok_or_else(|| {
anyhow!(
"The `--snarkos` flag is required when starting a devnet. Use `--clean-only` to only clean storage."
)
})?;
if cfg!(windows) && self.tmux {
bail!("tmux mode is not available on Windows โ remove `--tmux`.");
}
if self.tmux && std::env::var("TMUX").is_ok() {
bail!("Nested tmux session detected. Unset $TMUX and retry.");
}
if let Some(ref heights) = self.consensus_heights {
if !self.snarkos_features.contains(&"test_network".to_string()) {
bail!("The `test_network` feature must be enabled on snarkOS to use `--consensus-heights`.");
}
validate_consensus_heights(heights.as_slice())?;
}
if self.num_validators < 4 {
bail!("The number of validators must be at least 4.");
}
let total_nodes = self.num_validators + self.num_clients;
Self::validate_port_range("--rest-port", self.rest_port, total_nodes)?;
Self::validate_port_range("--node-port", self.node_port, total_nodes)?;
Self::validate_port_range("--bft-port", self.bft_port, total_nodes)?;
Self::validate_port_range("--metrics-port", self.metrics_port, self.num_validators)?;
if self.install {
if let Some(parent) = snarkos_path.parent()
&& !parent.exists()
{
std::fs::create_dir_all(parent)
.with_context(|| format!("Failed to create directory for binary: {}", parent.display()))?;
}
std::fs::write(snarkos_path, [0u8]).with_context(|| {
format!("Failed to write to path {} for snarkos installation", snarkos_path.display())
})?;
} else {
if !snarkos_path.exists() {
bail!(
"The snarkOS binary at `{}` does not exist. Please provide a valid path or use `--install`.",
snarkos_path.display()
);
}
};
let snarkos = canonicalize(snarkos_path)
.with_context(|| format!("Failed to resolve snarkOS path: {}", snarkos_path.display()))?;
println!("๐ง Starting devnet with the following options:");
println!(" โข Network: {}", self.network);
println!(" โข Validators: {}", self.num_validators);
println!(" โข Clients: {}", self.num_clients);
println!(" โข Storage: {}", self.storage.display());
if self.install {
println!(" โข Installing snarkOS at: {}", snarkos.display());
if let Some(ref version) = self.snarkos_version {
println!(" โข version: {version}");
}
if !self.snarkos_features.is_empty() {
println!(" โข features: {}", self.snarkos_features.iter().format(","));
}
} else {
println!(" โข Using snarkOS binary at: {}", snarkos.display());
}
if let Some(heights) = &self.consensus_heights {
println!(" โข Consensus heights: {}", heights.iter().format(","));
} else {
println!(" โข Consensus heights: default (based on your snarkOS binary)");
}
println!(" โข Clear storage: {}", if self.clear_storage { "yes" } else { "no" });
println!(" โข Verbosity: {}", self.verbosity);
println!(" โข tmux: {}", if self.tmux { "yes" } else { "no" });
let manager = Arc::new(Mutex::new(ChildManager::new()));
let (tx_shutdown, rx_shutdown) = crossbeam_channel::bounded::<()>(1);
let _signal_thread =
install_shutdown_listener(tx_shutdown.clone()).context("Failed to install shutdown listener")?;
let snarkos = if self.install {
if !confirm("\nProceed with snarkOS installation?", self.yes)? {
println!("โ Installation aborted.");
return Ok(());
}
install_snarkos(&snarkos, self.snarkos_version.as_deref(), &self.snarkos_features)?
} else {
snarkos
};
let version_output = StdCommand::new(&snarkos)
.arg("--version")
.output()
.context(format!("Failed to run `{}`", snarkos.display()))?;
if !version_output.status.success() {
bail!("Failed to run `{}`: {}", snarkos.display(), String::from_utf8_lossy(&version_output.stderr));
}
let version_str = String::from_utf8_lossy(&version_output.stdout);
println!("๐ Detected: {version_str}");
let features_str = version_str
.trim()
.split("features=[")
.nth(1)
.and_then(|s| s.split(']').next())
.ok_or_else(|| anyhow!("Failed to parse snarkOS features from version string: {version_str}"))?;
let found_features: Vec<String> = features_str.split(',').map(|s| s.trim().to_string()).collect();
for feature in &self.snarkos_features {
if !found_features.contains(feature) {
println!("โ ๏ธ Warning: snarkOS does not have the required feature `{feature}` enabled.");
}
}
if !confirm("\nProceed with devnet startup?", self.yes)? {
println!("โ Devnet aborted.");
return Ok(());
}
if !self.storage.exists() {
std::fs::create_dir_all(&self.storage)
.context(format!("Failed to create storage directory: {}", self.storage.display()))?;
} else if !self.storage.is_dir() {
bail!("The storage path `{}` is not a directory.", self.storage.display());
}
let storage = canonicalize(&self.storage)
.with_context(|| format!("Failed to resolve storage path: {}", self.storage.display()))?;
if self.clear_storage {
println!("๐งน Cleaning ledgers โฆ");
let mut cleaners = Vec::new();
for idx in 0..self.num_validators {
cleaners.push(clean_snarkos(
&snarkos,
self.network as usize,
idx,
&storage.join(format!("node-{idx}")),
&storage.join(format!("node-data-{idx}")),
)?);
}
for idx in 0..self.num_clients {
let dev_idx = idx + self.num_validators;
cleaners.push(clean_snarkos(
&snarkos,
self.network as usize,
dev_idx,
&storage.join(format!("node-{dev_idx}")),
&storage.join(format!("node-data-{dev_idx}")),
)?);
}
for mut c in cleaners {
c.wait()?;
}
}
let log_dir = {
let ts = Local::now().format(".logs-%Y-%m-%d-%H-%M-%S").to_string();
let p = storage.join(ts);
std::fs::create_dir_all(&p)?;
p
};
#[allow(clippy::too_many_arguments)]
fn build_args(
role: &str,
verbosity: u8,
network: usize,
num_validators: usize,
num_nodes: usize,
idx: usize,
log_file: &Path,
storage: &Path,
rest_port: Option<u16>,
node_port: Option<u16>,
bft_port: Option<u16>,
metrics_port: Option<u16>,
) -> Vec<String> {
let idx_u16 = idx as u16;
let mut base = vec![
"start".to_string(),
"--nodisplay".to_string(),
"--network".to_string(),
network.to_string(),
"--dev".to_string(),
idx.to_string(),
"--dev-num-validators".to_string(),
num_validators.to_string(),
"--rest-rps".to_string(),
REST_RPS.to_string(),
"--logfile".to_string(),
log_file.to_str().expect("log path is valid UTF-8").to_string(),
"--verbosity".to_string(),
verbosity.to_string(),
"--ledger-storage".to_string(),
storage.join(format!("node-{idx}")).to_str().expect("storage path is valid UTF-8").to_string(),
"--node-data-storage".to_string(),
storage.join(format!("node-data-{idx}")).to_str().expect("node-data path is valid UTF-8").to_string(),
];
if let Some(port) = rest_port {
base.extend(["--rest".into(), format!("0.0.0.0:{}", port + idx_u16)]);
}
if let Some(port) = node_port {
base.extend(["--node".into(), format!("0.0.0.0:{}", port + idx_u16)]);
let peers: String =
(0..num_nodes).map(|i| format!("127.0.0.1:{}", port + i as u16)).collect::<Vec<_>>().join(",");
base.extend(["--peers".into(), peers]);
}
if let Some(port) = bft_port {
base.extend(["--bft".into(), format!("0.0.0.0:{}", port + idx_u16)]);
let validators: String =
(0..num_validators).map(|i| format!("127.0.0.1:{}", port + i as u16)).collect::<Vec<_>>().join(",");
base.extend(["--validators".into(), validators]);
}
match role {
"validator" => {
base.extend(
["--allow-external-peers", "--validator", "--no-dev-txs"].into_iter().map(String::from),
);
if let Some(port) = metrics_port {
base.extend(["--metrics".into(), "--metrics-ip".into(), format!("0.0.0.0:{}", port + idx_u16)]);
}
}
"client" => base.push("--client".into()),
_ => unreachable!(),
}
base
}
if let Some(ref heights) = self.consensus_heights {
let heights = heights.iter().join(",");
println!("๐ง Setting consensus heights: {heights}");
#[allow(unsafe_code)]
unsafe {
env::set_var("CONSENSUS_VERSION_HEIGHTS", heights);
}
}
if self.tmux {
let mut args: Vec<String> =
vec!["new-session", "-d", "-s", "devnet", "-n", "validator-0"].into_iter().map(Into::into).collect();
if let Some(ref heights) = self.consensus_heights {
let heights = heights.iter().join(",");
args.push("-e".to_string());
args.push(format!("CONSENSUS_VERSION_HEIGHTS={heights}"));
}
ensure!(StdCommand::new("tmux").args(args).status()?.success(), "tmux failed to create session");
let num_nodes = self.num_validators + self.num_clients;
let base_index = {
let out = StdCommand::new("tmux").args(["show-option", "-gv", "base-index"]).output()?;
String::from_utf8_lossy(&out.stdout).trim().parse::<usize>().unwrap_or(0)
};
for idx in 0..self.num_validators {
let win_idx = idx + base_index;
let window_name = format!("validator-{idx}");
if idx != 0 {
StdCommand::new("tmux")
.args(["new-window", "-t", &format!("devnet:{win_idx}"), "-n", &window_name])
.status()?;
}
let log_file = log_dir.join(format!("{window_name}.log"));
let cmd = std::iter::once(snarkos.to_string_lossy().into_owned())
.chain(build_args(
"validator",
self.verbosity,
self.network as usize,
self.num_validators,
num_nodes,
idx,
log_file.as_path(),
&storage,
self.rest_port,
self.node_port,
self.bft_port,
self.metrics_port,
))
.collect::<Vec<_>>()
.join(" ");
StdCommand::new("tmux")
.args(["send-keys", "-t", &format!("devnet:{win_idx}"), &cmd, "C-m"])
.status()?;
}
for idx in 0..self.num_clients {
let dev_idx = idx + self.num_validators;
let win_idx = dev_idx + base_index;
let window_name = format!("client-{idx}");
StdCommand::new("tmux")
.args(["new-window", "-t", &format!("devnet:{win_idx}"), "-n", &window_name])
.status()?;
let log_file = log_dir.join(format!("{window_name}.log"));
let cmd = std::iter::once(snarkos.to_string_lossy().into_owned())
.chain(build_args(
"client",
self.verbosity,
self.network as usize,
self.num_validators,
num_nodes,
dev_idx,
log_file.as_path(),
&storage,
self.rest_port,
self.node_port,
self.bft_port,
None,
))
.collect::<Vec<_>>()
.join(" ");
StdCommand::new("tmux")
.args(["send-keys", "-t", &format!("devnet:{win_idx}"), &cmd, "C-m"])
.status()?;
}
println!("โ
tmux session \"devnet\" is ready โ attaching โฆ");
StdCommand::new("tmux").args(["attach-session", "-t", "devnet"]).status()?;
return Ok(()); }
println!("โ๏ธ Spawning nodes as background tasks โฆ");
let spawn_with_group = |mut cmd: StdCommand, log_file: &Path| -> AnyhowResult<Child> {
let log_handle = std::fs::OpenOptions::new().create(true).append(true).open(log_file)?;
cmd.stdout(Stdio::from(log_handle.try_clone()?));
cmd.stderr(Stdio::from(log_handle));
#[cfg(unix)]
#[allow(unsafe_code)]
unsafe {
cmd.pre_exec(|| {
setsid();
Ok(())
});
}
let child = cmd.spawn().map_err(|e| anyhow!("spawn {e}"))?;
#[cfg(windows)]
windows_kill_tree::attach_to_global_job(child.id())?;
Ok(child)
};
let num_nodes = self.num_validators + self.num_clients;
{
let mut guard = manager.lock();
for idx in 0..self.num_validators {
let log_file = log_dir.join(format!("validator-{idx}.log"));
let child = spawn_with_group(
{
let mut c = StdCommand::new(&snarkos);
c.args(build_args(
"validator",
self.verbosity,
self.network as usize,
self.num_validators,
num_nodes,
idx,
&log_file,
&storage,
self.rest_port,
self.node_port,
self.bft_port,
self.metrics_port,
));
c
},
&log_file,
)?;
println!(" โข validator {idx} (pid = {})", child.id());
guard.push(child);
}
for idx in 0..self.num_clients {
let dev_idx = idx + self.num_validators;
let log_file = log_dir.join(format!("client-{idx}.log"));
let child = spawn_with_group(
{
let mut c = StdCommand::new(&snarkos);
c.args(build_args(
"client",
self.verbosity,
self.network as usize,
self.num_validators,
num_nodes,
dev_idx,
&log_file,
&storage,
self.rest_port,
self.node_port,
self.bft_port,
None,
));
c
},
&log_file,
)?;
println!(" โข client {idx} (pid = {})", child.id());
guard.push(child);
}
}
println!("๐ Main process ID: {}", std::process::id());
println!("\nDevnet running โ Ctrl+C, SIGTERM, or terminal close to stop.");
let _ = rx_shutdown.recv();
manager.lock().shutdown_all(Duration::from_secs(30));
Ok(())
}
fn handle_clean(&self) -> AnyhowResult<()> {
let snarkos_path = self.snarkos.as_ref().ok_or_else(|| {
anyhow!("The `--snarkos` flag is required for `--clean-only`. Provide the path to the snarkOS binary.")
})?;
if !snarkos_path.exists() {
bail!("The snarkOS binary at `{}` does not exist.", snarkos_path.display());
}
let snarkos = canonicalize(snarkos_path)?;
if !self.storage.exists() {
println!("Storage path `{}` does not exist. Nothing to clean.", self.storage.display());
return Ok(());
}
if !self.storage.is_dir() {
bail!("Storage path `{}` is not a directory.", self.storage.display());
}
let storage = canonicalize(&self.storage)?;
let total_nodes = self.num_validators + self.num_clients;
if !confirm(
&format!("\nClean devnet storage for {total_nodes} nodes in `{}`?", self.storage.display()),
self.yes,
)? {
println!("Aborted.");
return Ok(());
}
println!("๐งน Cleaning ledgers โฆ");
let mut cleaners = Vec::new();
for idx in 0..self.num_validators {
cleaners.push(clean_snarkos(
&snarkos,
self.network as usize,
idx,
&storage.join(format!("node-{idx}")),
&storage.join(format!("node-data-{idx}")),
)?);
}
for idx in 0..self.num_clients {
let dev_idx = idx + self.num_validators;
cleaners.push(clean_snarkos(
&snarkos,
self.network as usize,
dev_idx,
&storage.join(format!("node-{dev_idx}")),
&storage.join(format!("node-data-{dev_idx}")),
)?);
}
for mut c in cleaners {
c.wait()?;
}
for entry in std::fs::read_dir(&storage)? {
let entry = entry?;
let name = entry.file_name();
let name = name.to_string_lossy();
if entry.file_type()?.is_dir() && name.starts_with(".logs-") {
std::fs::remove_dir_all(entry.path())?;
println!(" Removed {}", entry.path().display());
}
}
println!("Cleaned devnet storage.");
Ok(())
}
fn validate_port_range(flag: &str, base: Option<u16>, count: usize) -> AnyhowResult<()> {
if let Some(port) = base {
if count == 0 {
return Ok(());
}
let max_idx = (count - 1) as u16;
if port.checked_add(max_idx).is_none() {
bail!("{flag} {port} + {max_idx} nodes exceeds the maximum port number (65535).");
}
}
Ok(())
}
}