Skip to main content

zombienet_orchestrator/
network_spec.rs

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    /// Relaychain configuration, `None` for a JAM network.
33    pub(crate) relaychain: Option<RelaychainSpec>,
34
35    /// Jamchain configuration, takes the place of the relaychain
36    pub(crate) jamchain: Option<JamchainSpec>,
37
38    /// Parachains configurations.
39    pub(crate) parachains: Vec<ParachainSpec>,
40
41    /// HRMP channels configurations.
42    pub(crate) hrmp_channels: Vec<HrmpChannelConfig>,
43
44    /// Global settings
45    pub(crate) global_settings: GlobalSettings,
46
47    /// Custom processes
48    #[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        // The chain the parachains are anchored to, either the relaychain or the JAM chain.
68        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            // TODO: move to `fold` or map+fold
77            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    //
129    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        // try to find a node that use the same combination of image/cmd
135        let cmp_fn = |ad_hoc: &&NodeSpec| -> bool {
136            ad_hoc.image == node_spec.image && ad_hoc.command == node_spec.command
137        };
138
139        // check if we already had computed the args output for this cmd/[image]
140        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            // we need to compute the args output
162            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    /// The relaychain spec, `None` for a JAM network.
187    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 {}", &para.unique_id);
221            if base_dir_exists {
222                scoped_fs.create_dir_all(&para.unique_id).await?;
223            } else {
224                scoped_fs.create_dir(&para.unique_id).await?;
225            };
226            trace!("created dirs for {}", &para.unique_id);
227
228            // create wasm/state
229            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    // collect mutable references to all nodes from relaychain and parachains
255    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    // initialize the mapping of all possible node image/commands to corresponding nodes
272    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                // build mapping key using image and command if image is present or command only
279                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                // append the node to the vector of nodes for this image/command tuple
291                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                    // get node available args output from image/command
313                    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                    // map the result to include image and command
322                    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        // paras
385        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        // the parachain is anchored to the JAM chain
423        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        // ensure we report 2 errors
486        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}