1use std::{
2 collections::{hash_map::Entry, HashMap},
3 sync::Arc,
4};
5
6use configuration::{CustomProcess, GlobalSettings, HrmpChannelConfig, NetworkConfig};
7use futures::future::try_join_all;
8use provider::{DynNamespace, ProviderError, ProviderNamespace};
9use serde::{Deserialize, Serialize};
10use support::{
11 constants::{RELAY_NOT_NONE, THIS_IS_A_BUG},
12 fs::FileSystem,
13};
14use tracing::{debug, trace};
15
16use crate::{
17 errors::{merge_errs, OrchestratorError},
18 network_spec::jamchain::JamchainSpec,
19 ScopedFilesystem,
20};
21
22pub mod jamchain;
23pub mod jamnode;
24pub mod node;
25pub mod parachain;
26pub mod relaychain;
27
28use self::{node::NodeSpec, parachain::ParachainSpec, relaychain::RelaychainSpec};
29
30#[derive(Debug, Clone, Serialize, Deserialize)]
31pub struct NetworkSpec {
32 pub(crate) relaychain: Option<RelaychainSpec>,
34
35 pub(crate) jamchain: Option<JamchainSpec>,
37
38 pub(crate) parachains: Vec<ParachainSpec>,
40
41 pub(crate) hrmp_channels: Vec<HrmpChannelConfig>,
43
44 pub(crate) global_settings: GlobalSettings,
46
47 #[serde(default)]
49 pub(crate) custom_processes: Vec<CustomProcess>,
50}
51
52impl NetworkSpec {
53 pub async fn from_config(
54 network_config: &NetworkConfig,
55 ) -> Result<NetworkSpec, OrchestratorError> {
56 let mut errs = vec![];
57 let relaychain = network_config
58 .try_relaychain()
59 .map(RelaychainSpec::from_config)
60 .transpose()?;
61 let jamchain = if let Some(jamchain_config) = network_config.jamchain() {
62 Some(JamchainSpec::from_config(jamchain_config)?)
63 } else {
64 None
65 };
66
67 let root_chain = relaychain
69 .as_ref()
70 .map(|relaychain| relaychain.chain.clone())
71 .or_else(|| jamchain.as_ref().map(|jamchain| jamchain.id.clone()));
72
73 let mut parachains = vec![];
74
75 if let Some(root_chain) = root_chain {
76 for para_config in network_config.parachains() {
78 match ParachainSpec::from_config(para_config, root_chain.clone()) {
79 Ok(para) => parachains.push(para),
80 Err(err) => errs.push(err),
81 }
82 }
83 }
84
85 if errs.is_empty() {
86 Ok(NetworkSpec {
87 relaychain,
88 jamchain,
89 parachains,
90 hrmp_channels: network_config
91 .hrmp_channels()
92 .into_iter()
93 .cloned()
94 .collect(),
95 global_settings: network_config.global_settings().clone(),
96 custom_processes: network_config
97 .custom_processes()
98 .into_iter()
99 .cloned()
100 .collect(),
101 })
102 } else {
103 Err(merge_errs(&errs))
104 }
105 }
106
107 pub async fn populate_nodes_available_args(
108 &mut self,
109 ns: Arc<dyn ProviderNamespace + Send + Sync>,
110 ) -> Result<(), OrchestratorError> {
111 let network_nodes = self.collect_network_nodes();
112
113 let mut image_command_to_nodes_mapping =
114 Self::create_image_command_to_nodes_mapping(network_nodes);
115
116 let available_args_outputs =
117 Self::retrieve_all_nodes_available_args_output(ns, &image_command_to_nodes_mapping)
118 .await?;
119
120 Self::update_nodes_available_args_output(
121 &mut image_command_to_nodes_mapping,
122 available_args_outputs,
123 );
124
125 Ok(())
126 }
127
128 pub async fn node_available_args_output(
130 &self,
131 node_spec: &NodeSpec,
132 ns: Arc<dyn ProviderNamespace + Send + Sync>,
133 ) -> Result<String, ProviderError> {
134 let cmp_fn = |ad_hoc: &&NodeSpec| -> bool {
136 ad_hoc.image == node_spec.image && ad_hoc.command == node_spec.command
137 };
138
139 let node = self
141 .relaychain
142 .iter()
143 .flat_map(|relaychain| relaychain.nodes.iter())
144 .find(cmp_fn);
145 let node = if let Some(node) = node {
146 Some(node)
147 } else {
148 let node = self
149 .parachains
150 .iter()
151 .find_map(|para| para.collators.iter().find(cmp_fn));
152
153 node
154 };
155
156 let output = if let Some(node) = node {
157 node.available_args_output.clone().expect(&format!(
158 "args_output should be set for running nodes {THIS_IS_A_BUG}"
159 ))
160 } else {
161 let image = node_spec
163 .image
164 .as_ref()
165 .map(|image| image.as_str().to_string());
166 let command = node_spec.command.as_str().to_string();
167
168 ns.get_node_available_args((command, image)).await?
169 };
170
171 Ok(output)
172 }
173
174 pub fn relaychain(&self) -> &RelaychainSpec {
175 self.relaychain
176 .as_ref()
177 .expect(&format!("{RELAY_NOT_NONE}, {THIS_IS_A_BUG}"))
178 }
179
180 pub fn relaychain_mut(&mut self) -> &mut RelaychainSpec {
181 self.relaychain
182 .as_mut()
183 .expect(&format!("{RELAY_NOT_NONE}, {THIS_IS_A_BUG}"))
184 }
185
186 pub fn try_relaychain(&self) -> Option<&RelaychainSpec> {
188 self.relaychain.as_ref()
189 }
190
191 pub fn parachains_iter(&self) -> impl Iterator<Item = &ParachainSpec> {
192 self.parachains.iter()
193 }
194
195 pub fn parachains_iter_mut(&mut self) -> impl Iterator<Item = &mut ParachainSpec> {
196 self.parachains.iter_mut()
197 }
198
199 pub fn set_global_settings(&mut self, global_settings: GlobalSettings) {
200 self.global_settings = global_settings;
201 }
202
203 pub async fn build_parachain_artifacts<'a, T: FileSystem>(
204 &mut self,
205 ns: DynNamespace,
206 scoped_fs: &ScopedFilesystem<'a, T>,
207 relaychain_id: &str,
208 base_dir_exists: bool,
209 ) -> Result<(), anyhow::Error> {
210 for para in self.parachains.iter_mut() {
211 let chain_spec_raw_path = para
212 .build_chain_spec(
213 relaychain_id,
214 &ns,
215 scoped_fs,
216 para.post_process_script.clone().as_deref(),
217 )
218 .await?;
219
220 trace!("creating dirs for {}", ¶.unique_id);
221 if base_dir_exists {
222 scoped_fs.create_dir_all(¶.unique_id).await?;
223 } else {
224 scoped_fs.create_dir(¶.unique_id).await?;
225 };
226 trace!("created dirs for {}", ¶.unique_id);
227
228 para.genesis_state
230 .build(
231 chain_spec_raw_path.clone(),
232 format!("{}/genesis-state", para.unique_id),
233 &ns,
234 scoped_fs,
235 None,
236 )
237 .await?;
238 debug!("parachain genesis state built!");
239 para.genesis_wasm
240 .build(
241 chain_spec_raw_path,
242 format!("{}/genesis-wasm", para.unique_id),
243 &ns,
244 scoped_fs,
245 None,
246 )
247 .await?;
248 debug!("parachain genesis wasm built!");
249 }
250
251 Ok(())
252 }
253
254 fn collect_network_nodes(&mut self) -> Vec<&mut NodeSpec> {
256 vec![
257 self.relaychain
258 .iter_mut()
259 .flat_map(|relaychain| relaychain.nodes.iter_mut())
260 .collect::<Vec<_>>(),
261 self.parachains
262 .iter_mut()
263 .flat_map(|para| para.collators.iter_mut())
264 .collect(),
265 ]
266 .into_iter()
267 .flatten()
268 .collect::<Vec<_>>()
269 }
270
271 fn create_image_command_to_nodes_mapping(
273 network_nodes: Vec<&mut NodeSpec>,
274 ) -> HashMap<(Option<String>, String), Vec<&mut NodeSpec>> {
275 network_nodes.into_iter().fold(
276 HashMap::new(),
277 |mut acc: HashMap<(Option<String>, String), Vec<&mut node::NodeSpec>>, node| {
278 let key = node
280 .image
281 .as_ref()
282 .map(|image| {
283 (
284 Some(image.as_str().to_string()),
285 node.command.as_str().to_string(),
286 )
287 })
288 .unwrap_or_else(|| (None, node.command.as_str().to_string()));
289
290 if let Entry::Vacant(entry) = acc.entry(key.clone()) {
292 entry.insert(vec![node]);
293 } else {
294 acc.get_mut(&key).unwrap().push(node);
295 }
296
297 acc
298 },
299 )
300 }
301
302 async fn retrieve_all_nodes_available_args_output(
303 ns: Arc<dyn ProviderNamespace + Send + Sync>,
304 image_command_to_nodes_mapping: &HashMap<(Option<String>, String), Vec<&mut NodeSpec>>,
305 ) -> Result<Vec<(Option<String>, String, String)>, OrchestratorError> {
306 try_join_all(
307 image_command_to_nodes_mapping
308 .keys()
309 .map(|(image, command)| async {
310 let image = image.clone();
311 let command = command.clone();
312 let available_args = ns
314 .get_node_available_args((command.clone(), image.clone()))
315 .await?;
316 debug!(
317 "retrieved available args for image: {:?}, command: {}",
318 image, command
319 );
320
321 Ok::<_, OrchestratorError>((image, command, available_args))
323 })
324 .collect::<Vec<_>>(),
325 )
326 .await
327 }
328
329 fn update_nodes_available_args_output(
330 image_command_to_nodes_mapping: &mut HashMap<(Option<String>, String), Vec<&mut NodeSpec>>,
331 available_args_outputs: Vec<(Option<String>, String, String)>,
332 ) {
333 for (image, command, available_args_output) in available_args_outputs {
334 let nodes = image_command_to_nodes_mapping
335 .get_mut(&(image, command))
336 .expect(&format!(
337 "node image/command key should exist {THIS_IS_A_BUG}"
338 ));
339
340 for node in nodes {
341 node.available_args_output = Some(available_args_output.clone());
342 }
343 }
344 }
345}
346
347#[cfg(test)]
348mod tests {
349 use crate::generators::generate_node_port;
350
351 #[tokio::test]
352 async fn small_network_config_get_spec() {
353 use configuration::NetworkConfigBuilder;
354
355 use super::*;
356
357 let config = NetworkConfigBuilder::new()
358 .with_relaychain(|r| {
359 r.with_chain("rococo-local")
360 .with_default_command("polkadot")
361 .with_validator(|node| node.with_name("alice"))
362 .with_fullnode(|node| node.with_name("bob").with_command("polkadot1"))
363 })
364 .with_parachain(|p| {
365 p.with_id(100)
366 .with_default_command("adder-collator")
367 .with_collator(|c| c.with_name("collator1"))
368 })
369 .with_custom_process(|c| c.with_name("eth-rpc").with_command("command"))
370 .build()
371 .unwrap();
372
373 let network_spec = NetworkSpec::from_config(&config).await.unwrap();
374 let alice = network_spec.relaychain().nodes.first().unwrap();
375 let bob = network_spec.relaychain().nodes.get(1).unwrap();
376 assert_eq!(alice.command.as_str(), "polkadot");
377 assert_eq!(bob.command.as_str(), "polkadot1");
378 assert!(alice.is_validator);
379 assert!(!bob.is_validator);
380 assert_eq!(network_spec.custom_processes.len(), 1);
381 let first_custom_process = network_spec.custom_processes.first().unwrap();
382 assert_eq!(first_custom_process.name(), "eth-rpc");
383
384 assert_eq!(network_spec.parachains.len(), 1);
386 let para_100 = network_spec.parachains.first().unwrap();
387 assert_eq!(para_100.id, 100);
388 }
389
390 #[tokio::test]
391 async fn jam_network_with_parachain_get_spec() {
392 use configuration::NetworkConfig;
393
394 use super::*;
395
396 let config = NetworkConfig::load_from_toml_string(
397 r#"
398[jamchain]
399id = "dev"
400default_command = "polkajam"
401
402[[jamchain.nodes]]
403name = "jam-or"
404mode = "ordinary"
405
406[[parachains]]
407id = 100
408default_command = "polkadot-omni-node"
409default_args = ["--jam-rpc-url http://{{ZOMBIE:jam-or:rpc_uri}}"]
410
411[[parachains.collators]]
412name = "collator"
413"#,
414 )
415 .unwrap();
416
417 let network_spec = NetworkSpec::from_config(&config).await.unwrap();
418
419 assert!(network_spec.try_relaychain().is_none());
420 assert_eq!(network_spec.jamchain.as_ref().unwrap().nodes.len(), 1);
421
422 assert_eq!(network_spec.parachains.len(), 1);
424 let para = network_spec.parachains.first().unwrap();
425 assert_eq!(para.id, 100);
426 let collator = para.collators.first().unwrap();
427 assert_eq!(collator.command.as_str(), "polkadot-omni-node");
428 }
429
430 #[tokio::test]
431 async fn used_port_in_node_reports_err() {
432 use configuration::NetworkConfigBuilder;
433
434 use super::*;
435 let used_port = generate_node_port(None).unwrap();
436
437 let config = NetworkConfigBuilder::new()
438 .with_relaychain(|r| {
439 r.with_chain("rococo-local")
440 .with_default_command("polkadot")
441 .with_validator(|node| node.with_name("alice"))
442 .with_fullnode(|node| node.with_name("bob").with_command("polkadot1"))
443 })
444 .with_parachain(|p| {
445 p.with_id(100)
446 .with_default_command("adder-collator")
447 .with_collator(|c| c.with_name("collator1").with_rpc_port(used_port.0))
448 })
449 .build()
450 .unwrap();
451
452 let network_spec_res = NetworkSpec::from_config(&config).await;
453 assert!(network_spec_res.is_err());
454 assert!(network_spec_res
455 .unwrap_err()
456 .to_string()
457 .contains("err Can\'t bind in socket"));
458 }
459
460 #[tokio::test]
461 async fn used_port_in_multiple_node_reports_err() {
462 use configuration::NetworkConfigBuilder;
463
464 use super::*;
465 let used_port_0 = generate_node_port(None).unwrap();
466 let used_port_1 = generate_node_port(None).unwrap();
467
468 let config = NetworkConfigBuilder::new()
469 .with_relaychain(|r| {
470 r.with_chain("rococo-local")
471 .with_default_command("polkadot")
472 .with_validator(|node| node.with_name("alice").with_rpc_port(used_port_0.0))
473 .with_fullnode(|node| node.with_name("bob").with_rpc_port(used_port_1.0))
474 })
475 .with_parachain(|p| {
476 p.with_id(100)
477 .with_default_command("adder-collator")
478 .with_collator(|c| c.with_name("collator1"))
479 })
480 .build()
481 .unwrap();
482
483 let network_spec_res = NetworkSpec::from_config(&config).await;
484 assert!(network_spec_res.is_err());
485 let err_str = network_spec_res.unwrap_err().to_string();
487 let socket_errs: Vec<&str> = err_str.matches("err Can\'t bind in socket").collect();
488 assert_eq!(socket_errs.len(), 2);
489 }
490}