use sp_keyring::AccountKeyring;
use std::{
ffi::{
OsStr,
OsString,
},
io::{
BufRead,
BufReader,
Read,
},
process,
};
use subxt::{
Config,
OnlineClient,
};
pub struct TestNodeProcess<R: Config> {
proc: process::Child,
client: OnlineClient<R>,
url: String,
}
impl<R> Drop for TestNodeProcess<R>
where
R: Config,
{
fn drop(&mut self) {
let _ = self.kill();
}
}
impl<R> TestNodeProcess<R>
where
R: Config,
{
pub fn build<S>(program: S) -> TestNodeProcessBuilder<R>
where
S: AsRef<OsStr> + Clone,
{
TestNodeProcessBuilder::new(program)
}
pub fn kill(&mut self) -> Result<(), String> {
tracing::info!("Killing node process {}", self.proc.id());
if let Err(err) = self.proc.kill() {
let err = format!("Error killing node process {}: {}", self.proc.id(), err);
tracing::error!("{}", err);
return Err(err)
}
Ok(())
}
pub fn client(&self) -> OnlineClient<R> {
self.client.clone()
}
pub fn url(&self) -> &str {
&self.url
}
}
pub struct TestNodeProcessBuilder<R> {
node_path: OsString,
authority: Option<AccountKeyring>,
marker: std::marker::PhantomData<R>,
}
impl<R> TestNodeProcessBuilder<R>
where
R: Config,
{
pub fn new<P>(node_path: P) -> TestNodeProcessBuilder<R>
where
P: AsRef<OsStr>,
{
Self {
node_path: node_path.as_ref().into(),
authority: None,
marker: Default::default(),
}
}
pub fn with_authority(&mut self, account: AccountKeyring) -> &mut Self {
self.authority = Some(account);
self
}
pub async fn spawn(&self) -> Result<TestNodeProcess<R>, String> {
let mut cmd = process::Command::new(&self.node_path);
cmd.env("RUST_LOG", "info")
.arg("--dev")
.stdout(process::Stdio::piped())
.stderr(process::Stdio::piped())
.arg("--port=0")
.arg("--rpc-port=0");
if let Some(authority) = self.authority {
let authority = format!("{authority:?}");
let arg = format!("--{}", authority.as_str().to_lowercase());
cmd.arg(arg);
}
let mut proc = cmd.spawn().map_err(|e| {
format!(
"Error spawning substrate node '{}': {}",
self.node_path.to_string_lossy(),
e
)
})?;
let stderr = proc.stderr.take().unwrap();
let port = find_substrate_port_from_output(stderr);
let url = format!("ws://127.0.0.1:{port}");
let client = OnlineClient::from_url(url.clone()).await;
match client {
Ok(client) => {
Ok(TestNodeProcess {
proc,
client,
url: url.clone(),
})
}
Err(err) => {
let err = format!("Failed to connect to node rpc at {url}: {err}");
tracing::error!("{}", err);
proc.kill().map_err(|e| {
format!("Error killing substrate process '{}': {}", proc.id(), e)
})?;
Err(err)
}
}
}
}
fn find_substrate_port_from_output(r: impl Read + Send + 'static) -> u16 {
BufReader::new(r)
.lines()
.find_map(|line| {
let line =
line.expect("failed to obtain next line from stdout for port discovery");
let line_end = line
.rsplit_once("Listening for new connections on 127.0.0.1:")
.or_else(|| {
line.rsplit_once("Running JSON-RPC WS server: addr=127.0.0.1:")
})
.or_else(|| line.rsplit_once("Running JSON-RPC server: addr=127.0.0.1:"))
.map(|(_, port_str)| port_str)?;
let port_str = line_end.trim_end_matches(|b: char| !b.is_ascii_digit());
let port_num = port_str.parse().unwrap_or_else(|_| {
panic!("valid port expected for tracing line, got '{port_str}'")
});
Some(port_num)
})
.expect("We should find a port before the reader ends")
}
#[cfg(test)]
mod tests {
use super::*;
use subxt::PolkadotConfig as SubxtConfig;
#[tokio::test]
#[allow(unused_assignments)]
async fn spawning_and_killing_nodes_works() {
let mut client1: Option<OnlineClient<SubxtConfig>> = None;
let mut client2: Option<OnlineClient<SubxtConfig>> = None;
{
let node_proc1 =
TestNodeProcess::<SubxtConfig>::build("substrate-contracts-node")
.spawn()
.await
.unwrap();
client1 = Some(node_proc1.client());
let node_proc2 =
TestNodeProcess::<SubxtConfig>::build("substrate-contracts-node")
.spawn()
.await
.unwrap();
client2 = Some(node_proc2.client());
let res1 = node_proc1.client().rpc().block_hash(None).await;
let res2 = node_proc1.client().rpc().block_hash(None).await;
assert!(res1.is_ok());
assert!(res2.is_ok());
}
let res1 = client1.unwrap().rpc().block_hash(None).await;
let res2 = client2.unwrap().rpc().block_hash(None).await;
assert!(res1.is_err());
assert!(res2.is_err());
}
}