Skip to main content

zombienet_orchestrator/
lib.rs

1// TODO(Javier): Remove when we implement the logic in the orchestrator to spawn with the provider.
2#![allow(dead_code, clippy::expect_fun_call)]
3
4pub mod errors;
5pub mod generators;
6pub mod network;
7pub mod network_helper;
8pub mod observability;
9pub mod tx_helper;
10
11mod network_spec;
12pub mod shared;
13mod spawner;
14mod utils;
15
16use std::{
17    collections::{HashMap, HashSet, VecDeque},
18    env,
19    net::IpAddr,
20    path::{Path, PathBuf},
21    sync::Arc,
22    time::{Duration, SystemTime},
23};
24
25use anyhow::anyhow;
26use configuration::{types::JsonOverrides, NetworkConfig, RegistrationStrategy};
27use errors::OrchestratorError;
28use generators::{core_assignment, errors::GeneratorError, jam_config};
29use network::{
30    node::{NetworkNode, SpawnedNode},
31    parachain::Parachain,
32    relaychain::Relaychain,
33    Network,
34};
35// re-exported
36pub use network_spec::NetworkSpec;
37use network_spec::{jamchain::JamchainSpec, node::NodeSpec, parachain::ParachainSpec};
38use provider::{
39    types::{GenerateFileCommand, ProviderCapabilities, TransferedFile},
40    DynNamespace, DynProvider,
41};
42use serde_json::json;
43use support::{
44    constants::{
45        GRAPH_CONTAINS_DEP, GRAPH_CONTAINS_NAME, INDEGREE_CONTAINS_NAME, QUEUE_NOT_EMPTY,
46        THIS_IS_A_BUG,
47    },
48    fs::{FileSystem, FileSystemError},
49    replacer::{get_tokens_to_replace, has_tokens},
50};
51use tokio::time::timeout;
52use tracing::{debug, info, trace, warn};
53
54use crate::{
55    network::{
56        jamchain::{Jamchain, RawJamchain},
57        node::{jam::RawJamNetworkNode, JamNetworkNode, RawNetworkNode},
58        parachain::RawParachain,
59        relaychain::RawRelaychain,
60    },
61    shared::types::RegisterParachainOptions,
62    spawner::SpawnNodeCtx,
63    utils::write_zombie_json,
64};
65pub struct Orchestrator<T>
66where
67    T: FileSystem + Sync + Send,
68{
69    filesystem: T,
70    provider: DynProvider,
71}
72
73impl<T> Orchestrator<T>
74where
75    T: FileSystem + Sync + Send + Clone,
76{
77    pub fn new(filesystem: T, provider: DynProvider) -> Self {
78        Self {
79            filesystem,
80            provider,
81        }
82    }
83
84    pub async fn spawn(
85        &self,
86        network_config: NetworkConfig,
87    ) -> Result<Network<T>, OrchestratorError> {
88        let global_timeout = network_config.global_settings().network_spawn_timeout();
89        let network_spec = NetworkSpec::from_config(&network_config).await?;
90
91        let res = timeout(
92            Duration::from_secs(global_timeout.into()),
93            self.spawn_inner(network_spec),
94        )
95        .await
96        .map_err(|_| OrchestratorError::GlobalTimeOut(global_timeout));
97        res?
98    }
99
100    pub async fn spawn_from_spec(
101        &self,
102        network_spec: NetworkSpec,
103    ) -> Result<Network<T>, OrchestratorError> {
104        let global_timeout = network_spec.global_settings.network_spawn_timeout();
105        let res = timeout(
106            Duration::from_secs(global_timeout as u64),
107            self.spawn_inner(network_spec),
108        )
109        .await
110        .map_err(|_| OrchestratorError::GlobalTimeOut(global_timeout));
111        res?
112    }
113
114    pub async fn attach_to_live(
115        &self,
116        zombie_json_path: &Path,
117    ) -> Result<Network<T>, OrchestratorError> {
118        info!("attaching to live network...");
119        info!("reading zombie.json from {:?}", zombie_json_path);
120
121        let zombie_json_content = self.filesystem.read_to_string(zombie_json_path).await?;
122        let zombie_json: serde_json::Value = serde_json::from_str(&zombie_json_content)?;
123
124        info!("recreating namespace...");
125        let ns: DynNamespace = self
126            .provider
127            .create_namespace_from_json(&zombie_json)
128            .await?;
129
130        info!("recreating relaychain...");
131        let (relay, initial_spec) =
132            recreate_relaychain_from_json(&zombie_json, ns.clone(), self.provider.name()).await?;
133        let relay_nodes = relay.nodes.clone();
134
135        let mut network =
136            Network::new_with_relay(relay, ns.clone(), self.filesystem.clone(), initial_spec);
137
138        for node in relay_nodes {
139            if node.is_responsive().await {
140                node.set_is_running(true);
141            }
142            network.insert_node(node);
143        }
144
145        info!("recreating parachains...");
146        let parachains_map =
147            recreate_parachains_from_json(&zombie_json, ns.clone(), self.provider.name()).await?;
148        let para_nodes = parachains_map
149            .values()
150            .flat_map(|paras| paras.iter().flat_map(|para| para.collators.clone()))
151            .collect::<Vec<Arc<NetworkNode>>>();
152
153        network.set_parachains(parachains_map);
154        for node in para_nodes {
155            if node.is_responsive().await {
156                node.set_is_running(true);
157            }
158            network.insert_node(node);
159        }
160
161        if let Some(jamchain) =
162            recreate_jamchain_from_json(&zombie_json, ns.clone(), self.provider.name()).await?
163        {
164            info!("recreating jamchain...");
165            let jam_nodes = jamchain.nodes.clone();
166            network.set_jamchain(jamchain);
167            for node in jam_nodes {
168                if node.is_responsive().await {
169                    node.core().set_is_running(true);
170                }
171                network.insert_node(node);
172            }
173        }
174
175        Ok(network)
176    }
177
178    async fn spawn_inner(
179        &self,
180        mut network_spec: NetworkSpec,
181    ) -> Result<Network<T>, OrchestratorError> {
182        // main driver for spawn the network
183        debug!(network_spec = ?network_spec,"Network spec to spawn");
184
185        // TODO: move to Provider trait
186        validate_spec_with_provider_capabilities(&network_spec, self.provider.capabilities())
187            .map_err(|err| {
188                OrchestratorError::InvalidConfigForProvider(
189                    self.provider.name().into(),
190                    err.to_string(),
191                )
192            })?;
193
194        // create namespace
195        let ns = if let Some(base_dir) = network_spec.global_settings.base_dir() {
196            self.provider
197                .create_namespace_with_base_dir(base_dir)
198                .await?
199        } else {
200            self.provider.create_namespace().await?
201        };
202
203        // set the spawn_concurrency
204        let (spawn_concurrency, limited_by_tokens) = calculate_concurrency(&network_spec)?;
205
206        let start_time = SystemTime::now();
207        info!("🧰 ns: {}", ns.name());
208        info!("🧰 base_dir: {:?}", ns.base_dir());
209        info!("🕰 start time: {:?}", start_time);
210        info!("⚙️ spawn concurrency: {spawn_concurrency} (limited by tokens: {limited_by_tokens})");
211
212        network_spec
213            .populate_nodes_available_args(ns.clone())
214            .await?;
215
216        // A JAM chain takes the place of the relaychain and has its own, much simpler,
217        // spawn flow.
218        if network_spec.try_relaychain().is_none() {
219            let jam_spec =
220                network_spec
221                    .jamchain
222                    .clone()
223                    .ok_or(OrchestratorError::InvalidConfig(
224                        "A network needs either a relaychain or a jamchain.".to_string(),
225                    ))?;
226
227            return self
228                .spawn_jam(&jam_spec, &mut network_spec, ns, start_time)
229                .await;
230        }
231
232        // Resolve every node's `db_snapshot` AssetLocation into a local
233        // cache file once, serially, before any parallel spawn.
234        let all_nodes: Vec<&NodeSpec> = network_spec
235            .relaychain()
236            .nodes
237            .iter()
238            .chain(
239                network_spec
240                    .parachains_iter()
241                    .flat_map(|p| p.collators.iter()),
242            )
243            .collect();
244        let resolved_db_snapshots =
245            generators::resolve_db_snapshots(all_nodes, &ns, &self.filesystem).await?;
246
247        let base_dir = ns.base_dir().to_string_lossy();
248        let scoped_fs = ScopedFilesystem::new(&self.filesystem, &base_dir);
249
250        // The relaychain spec is moved out of the network spec while we set it up, so that
251        // it can be mutated independently of the parachains. It's put back before the nodes
252        // are spawned.
253        let mut relaychain = network_spec
254            .relaychain
255            .take()
256            .expect("checked to be `Some` above; qed");
257
258        // Create chain-spec for relaychain
259        relaychain.chain_spec.build(&ns, &scoped_fs).await?;
260
261        debug!("relaychain spec built!");
262        // Create parachain artifacts (chain-spec, wasm, state)
263        let relay_chain_id = relaychain.chain_spec.read_chain_id(&scoped_fs).await?;
264
265        let relay_chain_name = relaychain.chain.as_str().to_owned();
266        let base_dir_exists = network_spec.global_settings.base_dir().is_some();
267        network_spec
268            .build_parachain_artifacts(ns.clone(), &scoped_fs, &relay_chain_id, base_dir_exists)
269            .await?;
270
271        // Gather the parachains to register in genesis and the ones to register with extrinsic
272        let (para_to_register_in_genesis, para_to_register_with_extrinsic): (
273            Vec<&ParachainSpec>,
274            Vec<&ParachainSpec>,
275        ) = network_spec
276            .parachains
277            .iter()
278            .filter(|para| para.registration_strategy != RegistrationStrategy::Manual)
279            .partition(|para| {
280                matches!(para.registration_strategy, RegistrationStrategy::InGenesis)
281            });
282
283        let mut para_artifacts = vec![];
284        for para in &para_to_register_in_genesis {
285            let genesis_config = para.get_genesis_config()?;
286            para_artifacts.push(genesis_config)
287        }
288
289        // Customize relaychain
290        relaychain
291            .chain_spec
292            .customize_relay(
293                &relaychain,
294                &network_spec.hrmp_channels,
295                para_artifacts,
296                &scoped_fs,
297            )
298            .await?;
299
300        // Run post-process script if configured for the relaychain (run against plain spec before building raw)
301        if let Some(script_cmd) = relaychain.post_process_script.as_deref() {
302            relaychain
303                .chain_spec
304                .run_post_process_script(script_cmd, &scoped_fs)
305                .await?;
306        }
307
308        // Override cores if needed
309        let num_cores = network_spec.parachains_iter().fold(0u32, |mut acc, para| {
310            if let Some(cores) = para.num_cores {
311                acc += cores;
312            } else if &RegistrationStrategy::InGenesis == para.registration_strategy() {
313                // add 1 by default
314                acc += 1;
315            }
316
317            acc
318        });
319
320        if num_cores > para_to_register_in_genesis.len() as u32 {
321            let num_cores_to_set = num_cores - para_to_register_in_genesis.len() as u32;
322            // we should set the correct core config
323            let overrides = json!({
324                "configuration": {
325                    "config": {
326                        "scheduler_params": {
327                            "num_cores": num_cores_to_set,
328                            "max_validators_per_core": 1
329                        },
330                    }
331                }
332            });
333
334            relaychain
335                .chain_spec
336                .apply_genesis_override(&scoped_fs, &overrides)
337                .await?;
338        }
339
340        // Build raw version (after any post-processing of the plain spec)
341        relaychain
342            .chain_spec
343            .build_raw(&ns, &scoped_fs, None)
344            .await?;
345
346        // override wasm if needed
347        if let Some(ref wasm_override) = relaychain.wasm_override {
348            relaychain
349                .chain_spec
350                .override_code(&scoped_fs, wasm_override)
351                .await?;
352        }
353
354        // custom override raw spec if needed
355        if let Some(ref raw_spec_override) = relaychain.raw_spec_override {
356            relaychain
357                .chain_spec
358                .override_raw_spec(&scoped_fs, raw_spec_override)
359                .await?;
360        }
361
362        // assign extra cores if needed
363        if num_cores > para_to_register_in_genesis.len() as u32
364            || (relaychain.override_session_0 && !para_to_register_in_genesis.is_empty())
365        {
366            debug!("Raw overrides info: num_cores: {}, para_to_register_in_genesis_len: {:?}, override_session_0: {}", num_cores, para_to_register_in_genesis.len(), relaychain.override_session_0);
367            let mut core_index = 0u32;
368            // we should check with version the runtime is using
369            // could be ParaScheduler or CoretimeAssignmentProvider
370            let scheduler_key = core_assignment::get_parascheduler_storage_key();
371            let is_old = !relaychain
372                .chain_spec
373                .find_raw_key(&scoped_fs, &scheduler_key)
374                .await?;
375
376            let mut para_scheduler_value_parts: Vec<String> = vec![];
377            let mut raw_json_overrides = json!({});
378            // loop over para and assign cores from 0..
379            for para in &network_spec.parachains {
380                // cores we need to assign
381                let mut cores_for_para = 0_u32;
382                if let Some(cores) = para.num_cores {
383                    if &RegistrationStrategy::InGenesis == para.registration_strategy() {
384                        cores_for_para = cores;
385                    }
386                } else {
387                    // no num_cores set but we need to check if `override_session_0` is true
388                    // to assign the first core.
389                    if &RegistrationStrategy::InGenesis == para.registration_strategy()
390                        && relaychain.override_session_0
391                    {
392                        cores_for_para += 1;
393                    }
394                }
395
396                debug!(
397                    "Assigning {cores_for_para} cores in raw spec for para {}. Using pallet {}",
398                    para.id,
399                    if is_old {
400                        "CoretimeAssignmentProvider"
401                    } else {
402                        "ParaScheduler"
403                    }
404                );
405                for _core in 0..cores_for_para {
406                    if is_old {
407                        let (core_assign_key, core_assign_value) =
408                            core_assignment::generate_old(core_index, para.id);
409                        raw_json_overrides[core_assign_key] = json!(core_assign_value);
410                    } else {
411                        let part = core_assignment::generate(core_index, para.id);
412                        para_scheduler_value_parts.push(part);
413                    }
414                    core_index += 1;
415                }
416            }
417
418            // if not old we need to store the k/v to override
419            if !is_old {
420                let count_prefix = format!("{:02x}", para_scheduler_value_parts.len() * 4);
421                let core_assign_value =
422                    format!("{count_prefix}{}", para_scheduler_value_parts.join(""));
423                raw_json_overrides[scheduler_key] = json!(core_assign_value);
424            }
425
426            // extra check to ensure we need to override session 0
427            if relaychain.override_session_0 {
428                trace!("Overriding pallet ParaSessionInfo.session (0) to allow paras to produce blocks at first session.");
429                let raw_spec = relaychain.chain_spec.read_raw_spec(&scoped_fs).await?;
430                let overrides = generators::generate_session_0_overrides(&raw_spec, num_cores)?;
431
432                for (k, v) in overrides.as_object().ok_or(anyhow!(
433                    "'generate_session_0_overrides' should be a valid json Object."
434                ))? {
435                    raw_json_overrides[k] = v.clone();
436                }
437            }
438
439            debug!("Raw overrides keys: {:?}", raw_json_overrides);
440
441            relaychain
442                .chain_spec
443                .override_raw_spec(
444                    &scoped_fs,
445                    &JsonOverrides::Json(json!({
446                        "genesis": {
447                            "raw": {
448                                "top": raw_json_overrides
449                            }
450                        }
451                    })),
452                )
453                .await?;
454        }
455
456        // The relaychain is fully set up now, put it back into the network spec.
457        network_spec.relaychain = Some(relaychain);
458
459        let (bootnodes, relaynodes) =
460            split_nodes_by_bootnodes(&network_spec.relaychain().nodes, false);
461
462        // TODO: we want to still supporting spawn a dedicated bootnode??
463        let mut ctx = SpawnNodeCtx {
464            chain_id: &relay_chain_id,
465            parachain_id: None,
466            chain: relay_chain_name.as_str(),
467            role: ZombieRole::Node,
468            ns: &ns,
469            scoped_fs: &scoped_fs,
470            parachain: None,
471            bootnodes_addr: &vec![],
472            wait_ready: false,
473            nodes_by_name: json!({}),
474            global_settings: &network_spec.global_settings,
475            resolved_db_snapshots: &resolved_db_snapshots,
476        };
477
478        let global_files_to_inject = vec![TransferedFile::new(
479            PathBuf::from(format!(
480                "{}/{relay_chain_name}.json",
481                ns.base_dir().to_string_lossy()
482            )),
483            PathBuf::from(format!("/cfg/{relay_chain_name}.json")),
484        )];
485
486        let r = Relaychain::new(
487            relay_chain_name.to_string(),
488            relay_chain_id.clone(),
489            PathBuf::from(network_spec.relaychain().chain_spec.raw_path().ok_or(
490                OrchestratorError::InvariantError("chain-spec raw path should be set now"),
491            )?),
492        );
493        let mut network =
494            Network::new_with_relay(r, ns.clone(), self.filesystem.clone(), network_spec.clone());
495
496        // Initiate the node_ws_url which will be later used in the Parachain_with_extrinsic config
497        let mut node_ws_url: String = "".to_string();
498
499        // Calculate the bootnodes addr from the running nodes
500        let mut bootnodes_addr: Vec<String> = vec![];
501
502        for level in dependency_levels_among(&bootnodes)? {
503            let mut running_nodes_per_level = vec![];
504            for chunk in level.chunks(spawn_concurrency) {
505                let spawning_tasks = chunk
506                    .iter()
507                    .map(|node| spawner::spawn_node(node, global_files_to_inject.clone(), &ctx));
508
509                for node in futures::future::try_join_all(spawning_tasks).await? {
510                    let bootnode_multiaddr = node.multiaddr();
511
512                    bootnodes_addr.push(bootnode_multiaddr.to_string());
513
514                    // Is used in the register_para_options (We need to get this from the relay and not the collators)
515                    if node_ws_url.is_empty() {
516                        node_ws_url.clone_from(&node.ws_uri)
517                    }
518
519                    running_nodes_per_level.push(node);
520                }
521            }
522            info!(
523                "🕰  waiting for level: {:?} to be up...",
524                level.iter().map(|n| n.name.clone()).collect::<Vec<_>>()
525            );
526
527            // Wait for all nodes in the current level to be up
528            let waiting_tasks = running_nodes_per_level.iter().map(|node| {
529                node.wait_until_is_up(network_spec.global_settings.node_spawn_timeout())
530            });
531
532            let _ = futures::future::try_join_all(waiting_tasks).await?;
533
534            for node in running_nodes_per_level {
535                // Add the node to the  context and `Network` instance
536                ctx.nodes_by_name[node.name().to_owned()] = serde_json::to_value(&node)?;
537                network.add_running_node(node, None).await;
538            }
539        }
540
541        // Add the bootnodes to the relaychain spec file and ctx
542        network_spec
543            .relaychain()
544            .chain_spec
545            .add_bootnodes(&scoped_fs, &bootnodes_addr)
546            .await?;
547
548        ctx.bootnodes_addr = &bootnodes_addr;
549
550        for level in dependency_levels_among(&relaynodes)? {
551            let mut running_nodes_per_level = vec![];
552            for chunk in level.chunks(spawn_concurrency) {
553                let spawning_tasks = chunk
554                    .iter()
555                    .map(|node| spawner::spawn_node(node, global_files_to_inject.clone(), &ctx));
556
557                for node in futures::future::try_join_all(spawning_tasks).await? {
558                    running_nodes_per_level.push(node);
559                }
560            }
561            info!(
562                "🕰 waiting for level: {:?} to be up...",
563                level.iter().map(|n| n.name.clone()).collect::<Vec<_>>()
564            );
565
566            // Wait for all nodes in the current level to be up
567            let waiting_tasks = running_nodes_per_level.iter().map(|node| {
568                node.wait_until_is_up(network_spec.global_settings.network_spawn_timeout())
569            });
570
571            let _ = futures::future::try_join_all(waiting_tasks).await?;
572
573            for node in running_nodes_per_level {
574                ctx.nodes_by_name[node.name().to_owned()] = serde_json::to_value(&node)?;
575                network.add_running_node(node, None).await;
576            }
577        }
578
579        // spawn paras
580        self.spawn_parachains(
581            &network_spec.parachains,
582            &ctx,
583            &global_files_to_inject,
584            spawn_concurrency,
585            &mut network,
586        )
587        .await?;
588
589        // Now we need to register the paras with extrinsic from the Vec collected before;
590        for para in para_to_register_with_extrinsic {
591            let register_para_options: RegisterParachainOptions = RegisterParachainOptions {
592                id: para.id,
593                // This needs to resolve correctly
594                wasm_path: para
595                    .genesis_wasm
596                    .artifact_path()
597                    .ok_or(OrchestratorError::InvariantError(
598                        "artifact path for wasm must be set at this point",
599                    ))?
600                    .to_path_buf(),
601                state_path: para
602                    .genesis_state
603                    .artifact_path()
604                    .ok_or(OrchestratorError::InvariantError(
605                        "artifact path for state must be set at this point",
606                    ))?
607                    .to_path_buf(),
608                node_ws_url: node_ws_url.clone(),
609                onboard_as_para: para.onboard_as_parachain,
610                seed: None, // TODO: Seed is passed by?
611                finalization: false,
612            };
613
614            Parachain::register(register_para_options, &scoped_fs).await?;
615        }
616
617        if network_spec.global_settings.observability().enabled() {
618            match network
619                .start_observability(network_spec.global_settings.observability())
620                .await
621            {
622                Ok(obs) => {
623                    info!("📊 Prometheus URL: {}", obs.prometheus_url);
624                    info!("📊 Grafana URL: {}", obs.grafana_url);
625                },
626                Err(e) => {
627                    warn!("⚠️  Failed to spawn observability stack: {e}");
628                },
629            }
630        }
631
632        // start custom processes if needed
633        for cp in &network_spec.custom_processes {
634            if let Err(e) = spawner::spawn_process(cp, ns.clone()).await {
635                warn!("⚠️  Failed to spawn custom process {}, err: {e}", cp.name())
636            }
637        }
638
639        network.set_start_time_ts(start_time);
640
641        write_zombie_json(serde_json::to_value(&network)?, scoped_fs, ns.name()).await?;
642
643        if network_spec.global_settings.tear_down_on_failure() {
644            network.spawn_watching_task();
645        }
646
647        generators::cleanup_db_snapshot_cache(&resolved_db_snapshots).await;
648
649        Ok(network)
650    }
651
652    /// Spawn the collators of every parachain, in the context of an already running network.
653    ///
654    /// `ctx` is the context of the chain the parachains are anchored to (relaychain or JAM
655    /// chain); the per-parachain context is derived from it.
656    async fn spawn_parachains<'a>(
657        &self,
658        parachains: &[ParachainSpec],
659        ctx: &SpawnNodeCtx<'a, T>,
660        files_to_inject: &[TransferedFile],
661        spawn_concurrency: usize,
662        network: &mut Network<T>,
663    ) -> Result<(), OrchestratorError> {
664        let scoped_fs = ctx.scoped_fs;
665
666        for para in parachains {
667            // Create parachain (in the context of the running network)
668            let parachain = Parachain::from_spec(para, files_to_inject, scoped_fs).await?;
669            let parachain_id = parachain.chain_id.clone();
670
671            let (bootnodes, collators) =
672                split_nodes_by_bootnodes(&para.collators, para.no_default_bootnodes);
673
674            // Create `ctx` for spawn parachain nodes
675            let mut ctx_para = SpawnNodeCtx {
676                parachain: Some(para),
677                parachain_id: parachain_id.as_deref(),
678                role: if para.is_cumulus_based {
679                    ZombieRole::CumulusCollator
680                } else {
681                    ZombieRole::Collator
682                },
683                bootnodes_addr: &vec![],
684                ..ctx.clone()
685            };
686
687            // Calculate the bootnodes addr from the running nodes
688            let mut bootnodes_addr: Vec<String> = vec![];
689            let mut running_nodes: Vec<NetworkNode> = vec![];
690
691            for level in dependency_levels_among(&bootnodes)? {
692                for node in self
693                    .spawn_parachain_level(&level, &parachain, &ctx_para, spawn_concurrency)
694                    .await?
695                {
696                    bootnodes_addr.push(node.multiaddr().to_string());
697                    ctx_para.nodes_by_name[node.name().to_owned()] = serde_json::to_value(&node)?;
698                    running_nodes.push(node);
699                }
700            }
701
702            if let Some(para_chain_spec) = para.chain_spec.as_ref() {
703                para_chain_spec
704                    .add_bootnodes(scoped_fs, &bootnodes_addr)
705                    .await?;
706            }
707
708            ctx_para.bootnodes_addr = &bootnodes_addr;
709
710            // Spawn the rest of the nodes
711            for level in dependency_levels_among(&collators)? {
712                for node in self
713                    .spawn_parachain_level(&level, &parachain, &ctx_para, spawn_concurrency)
714                    .await?
715                {
716                    ctx_para.nodes_by_name[node.name().to_owned()] = serde_json::to_value(&node)?;
717                    running_nodes.push(node);
718                }
719            }
720
721            let running_para_id = parachain.para_id;
722            network.add_para(parachain);
723            for node in running_nodes {
724                network.add_running_node(node, Some(running_para_id)).await;
725            }
726        }
727
728        Ok(())
729    }
730
731    /// Spawn one dependency level of collators, concurrently, and wait until they are up.
732    async fn spawn_parachain_level<'a>(
733        &self,
734        level: &[&NodeSpec],
735        parachain: &Parachain,
736        ctx: &SpawnNodeCtx<'a, T>,
737        spawn_concurrency: usize,
738    ) -> Result<Vec<NetworkNode>, OrchestratorError> {
739        let mut running_nodes = vec![];
740        for chunk in level.chunks(spawn_concurrency) {
741            let spawning_tasks = chunk
742                .iter()
743                .map(|node| spawner::spawn_node(node, parachain.files_to_inject.clone(), ctx));
744
745            running_nodes.extend(futures::future::try_join_all(spawning_tasks).await?);
746        }
747
748        info!(
749            "🕰  waiting for level: {:?} to be up...",
750            level.iter().map(|n| n.name.clone()).collect::<Vec<_>>()
751        );
752
753        // Wait for all nodes in the current level to be up
754        let waiting_tasks = running_nodes
755            .iter()
756            .map(|node| node.wait_until_is_up(ctx.global_settings.network_spawn_timeout()));
757
758        let _ = futures::future::try_join_all(waiting_tasks).await?;
759
760        Ok(running_nodes)
761    }
762
763    /// Spawn a network whose root chain is a JAM chain instead of a relaychain.
764    async fn spawn_jam(
765        &self,
766        jam_spec: &JamchainSpec,
767        network_spec: &mut NetworkSpec,
768        ns: DynNamespace,
769        start_time: SystemTime,
770    ) -> Result<Network<T>, OrchestratorError> {
771        let base_dir = ns.base_dir().to_string_lossy().to_string();
772        let scoped_fs = ScopedFilesystem::new(&self.filesystem, &base_dir);
773
774        // generate config
775        let jam_config = jam_config::generate(jam_spec)?;
776        // store the config file
777        scoped_fs
778            .write(
779                "jam_config.json",
780                serde_json::to_string_pretty(&jam_config)?,
781            )
782            .await?;
783        // generate spec
784        let cmd_parts: Vec<&str> = jam_spec.chain_spec_command.split(" ").collect();
785        let cmd = cmd_parts
786            .first()
787            .expect("jam chain-spec generator cmd should be valid");
788        let jam_config_full_path = format!("{}/{}", base_dir, "jam_config.json");
789        let jam_spec_full_path = format!("{}/{}", base_dir, "jam_spec.json");
790        let args = vec![
791            "gen-spec",
792            jam_config_full_path.as_str(),
793            jam_spec_full_path.as_str(),
794        ];
795        let generate_command = GenerateFileCommand::new(cmd, jam_spec_full_path.clone()).args(args);
796        generators::chain_spec::build_locally(
797            generate_command,
798            &scoped_fs,
799            Some(&PathBuf::from(&jam_spec_full_path)),
800        )
801        .await?;
802
803        // Create parachain artifacts (chain-spec, wasm, state). The parachains are anchored
804        // to the JAM chain, so its id is the one written into their chain-spec.
805        let base_dir_exists = network_spec.global_settings.base_dir().is_some();
806        network_spec
807            .build_parachain_artifacts(
808                ns.clone(),
809                &scoped_fs,
810                jam_spec.id.as_str(),
811                base_dir_exists,
812            )
813            .await?;
814
815        // Resolve every collator's `db_snapshot` AssetLocation into a local cache file once,
816        // serially, before any parallel spawn.
817        let collators: Vec<&NodeSpec> = network_spec
818            .parachains_iter()
819            .flat_map(|para| para.collators.iter())
820            .collect();
821        let resolved_db_snapshots =
822            generators::resolve_db_snapshots(collators, &ns, &self.filesystem).await?;
823
824        let (spawn_concurrency, _) = calculate_concurrency(network_spec)?;
825        let global_settings = &network_spec.global_settings;
826
827        // context setup
828        let jam_ctx = SpawnNodeCtx {
829            chain_id: jam_spec.id.as_str(),
830            parachain_id: None,
831            chain: jam_spec.id.as_str(),
832            role: ZombieRole::Node,
833            ns: &ns,
834            scoped_fs: &scoped_fs,
835            parachain: None,
836            bootnodes_addr: &vec![],
837            wait_ready: false,
838            nodes_by_name: json!({}),
839            global_settings,
840            resolved_db_snapshots: &resolved_db_snapshots,
841        };
842
843        let jam_global_files_to_inject = vec![TransferedFile::new(
844            PathBuf::from(format!("{base_dir}/jam_spec.json")),
845            PathBuf::from("/cfg/jam_spec.json"),
846        )];
847
848        let mut network =
849            Network::new_without_relay(ns.clone(), self.filesystem.clone(), network_spec.clone());
850
851        network.set_jamchain(Jamchain::new(
852            jam_spec.id.as_str(),
853            PathBuf::from(&jam_spec_full_path),
854        ));
855
856        let mut bootnodes_addr = vec![];
857        // Nodes spawned so far, so the args of the next ones can reference them.
858        let mut nodes_by_name = json!({});
859        for jam_node in &jam_spec.nodes {
860            let jam_ctx = SpawnNodeCtx {
861                bootnodes_addr: &bootnodes_addr,
862                nodes_by_name: nodes_by_name.clone(),
863                ..jam_ctx.clone()
864            };
865
866            let running_node =
867                spawner::spawn_jam_node(jam_node, jam_global_files_to_inject.clone(), &jam_ctx)
868                    .await?;
869
870            running_node
871                .wait_until_is_up(global_settings.node_spawn_timeout() as u64)
872                .await
873                .map_err(|e| OrchestratorError::InvalidConfig(e.to_string()))?;
874
875            bootnodes_addr.push(running_node.peer_addr().to_string());
876            nodes_by_name[running_node.name().to_owned()] = serde_json::to_value(&running_node)?;
877            network.add_running_jam_node(running_node).await;
878        }
879
880        // spawn paras
881        for para in network_spec.parachains.iter() {
882            if para.registration_strategy() != &RegistrationStrategy::Manual {
883                warn!(
884                    "⚠️  Parachain {} can't be registered automatically on a JAM chain, it needs to be registered manually.",
885                    para.id
886                );
887            }
888        }
889
890        // the JAM nodes are already running, so the collators can reference them
891        let jam_ctx = SpawnNodeCtx {
892            nodes_by_name,
893            ..jam_ctx
894        };
895
896        self.spawn_parachains(
897            &network_spec.parachains,
898            &jam_ctx,
899            &jam_global_files_to_inject,
900            spawn_concurrency,
901            &mut network,
902        )
903        .await?;
904
905        // start custom processes if needed
906        for cp in &network_spec.custom_processes {
907            if let Err(e) = spawner::spawn_process(cp, ns.clone()).await {
908                warn!("⚠️  Failed to spawn custom process {}, err: {e}", cp.name())
909            }
910        }
911
912        network.set_start_time_ts(start_time);
913
914        write_zombie_json(serde_json::to_value(&network)?, scoped_fs, ns.name()).await?;
915
916        if network_spec.global_settings.tear_down_on_failure() {
917            network.spawn_watching_task();
918        }
919
920        generators::cleanup_db_snapshot_cache(&resolved_db_snapshots).await;
921
922        Ok(network)
923    }
924}
925
926// Helpers
927
928/// Make sure a node persisted in `zombie.json` was spawned by the provider we
929/// are attaching with.
930fn validate_provider_tag(
931    inner: &serde_json::Value,
932    node_name: &str,
933    provider_name: &str,
934) -> Result<(), OrchestratorError> {
935    let provider_tag = inner
936        .get("provider_tag")
937        .and_then(|v| v.as_str())
938        .ok_or_else(|| {
939            OrchestratorError::InvalidConfig(format!(
940                "Node '{node_name}' is missing `provider_tag` in inner node JSON"
941            ))
942        })?;
943
944    if provider_tag != provider_name {
945        return Err(OrchestratorError::InvalidConfigForProvider(
946            provider_name.to_string(),
947            provider_tag.to_string(),
948        ));
949    }
950
951    Ok(())
952}
953
954async fn recreate_network_nodes_from_json(
955    nodes_json: &serde_json::Value,
956    ns: DynNamespace,
957    provider_name: &str,
958) -> Result<Vec<Arc<NetworkNode>>, OrchestratorError> {
959    let raw_nodes: Vec<RawNetworkNode> = serde_json::from_value(nodes_json.clone())?;
960
961    let mut nodes = Vec::with_capacity(raw_nodes.len());
962    for raw in raw_nodes {
963        validate_provider_tag(&raw.inner, &raw.name, provider_name)?;
964
965        let inner = ns.spawn_node_from_json(&raw.inner).await?;
966        let relay_node = NetworkNode::new(
967            raw.name,
968            raw.ws_uri,
969            raw.prometheus_uri,
970            raw.multiaddr,
971            raw.spec,
972            inner,
973            raw.cmd_generator_opts,
974            raw.context,
975        );
976        nodes.push(Arc::new(relay_node));
977    }
978
979    Ok(nodes)
980}
981
982async fn recreate_relaychain_from_json(
983    zombie_json: &serde_json::Value,
984    ns: DynNamespace,
985    provider_name: &str,
986) -> Result<(Relaychain, NetworkSpec), OrchestratorError> {
987    let relay_json = zombie_json
988        .get("relay")
989        .ok_or(OrchestratorError::InvalidConfig(
990            "Missing `relay` field in zombie.json".into(),
991        ))?
992        .clone();
993
994    let mut relay_raw: RawRelaychain = serde_json::from_value(relay_json)?;
995
996    let initial_spec: NetworkSpec = serde_json::from_value(
997        zombie_json
998            .get("initial_spec")
999            .ok_or(OrchestratorError::InvalidConfig(
1000                "Missing `initial_spec` field in zombie.json".into(),
1001            ))?
1002            .clone(),
1003    )?;
1004
1005    // Populate relay nodes
1006    let nodes =
1007        recreate_network_nodes_from_json(&relay_raw.nodes, ns.clone(), provider_name).await?;
1008    relay_raw.inner.nodes = nodes;
1009
1010    Ok((relay_raw.inner, initial_spec))
1011}
1012
1013/// Rebuild the `jamchain` section of a `zombie.json`, if the network had one.
1014async fn recreate_jamchain_from_json(
1015    zombie_json: &serde_json::Value,
1016    ns: DynNamespace,
1017    provider_name: &str,
1018) -> Result<Option<Jamchain>, OrchestratorError> {
1019    let Some(jamchain_json) = zombie_json.get("jamchain") else {
1020        return Ok(None);
1021    };
1022
1023    let mut jamchain_raw: RawJamchain = serde_json::from_value(jamchain_json.clone())?;
1024    let raw_nodes: Vec<RawJamNetworkNode> = serde_json::from_value(jamchain_raw.nodes.clone())?;
1025
1026    let mut nodes = Vec::with_capacity(raw_nodes.len());
1027    for raw in raw_nodes {
1028        validate_provider_tag(&raw.inner, &raw.name, provider_name)?;
1029
1030        let inner = ns.spawn_node_from_json(&raw.inner).await?;
1031        nodes.push(Arc::new(JamNetworkNode::new(
1032            raw.name,
1033            inner,
1034            raw.spec,
1035            raw.ip,
1036            raw.cmd_generator_opts,
1037        )));
1038    }
1039
1040    jamchain_raw.inner.nodes = nodes;
1041
1042    Ok(Some(jamchain_raw.inner))
1043}
1044
1045async fn recreate_parachains_from_json(
1046    zombie_json: &serde_json::Value,
1047    ns: DynNamespace,
1048    provider_name: &str,
1049) -> Result<HashMap<u32, Vec<Parachain>>, OrchestratorError> {
1050    let paras_json = zombie_json
1051        .get("parachains")
1052        .ok_or(OrchestratorError::InvalidConfig(
1053            "Missing `parachains` field in zombie.json".into(),
1054        ))?
1055        .clone();
1056
1057    let raw_paras: HashMap<u32, Vec<RawParachain>> = serde_json::from_value(paras_json)?;
1058
1059    let mut parachains_map = HashMap::new();
1060
1061    for (id, parachain_entries) in raw_paras {
1062        let mut parsed_vec = Vec::with_capacity(parachain_entries.len());
1063
1064        for raw_para in parachain_entries {
1065            let mut para = raw_para.inner;
1066            para.collators =
1067                recreate_network_nodes_from_json(&raw_para.collators, ns.clone(), provider_name)
1068                    .await?;
1069            parsed_vec.push(para);
1070        }
1071
1072        parachains_map.insert(id, parsed_vec);
1073    }
1074
1075    Ok(parachains_map)
1076}
1077
1078// Split the node list depending if it's bootnode or not
1079// NOTE: if there isn't a bootnode declared we use the first one
1080fn split_nodes_by_bootnodes(
1081    nodes: &[NodeSpec],
1082    no_default_bootnodes: bool,
1083) -> (Vec<&NodeSpec>, Vec<&NodeSpec>) {
1084    // get the bootnodes to spawn first and calculate the bootnode string for use later
1085    let mut bootnodes = vec![];
1086    let mut other_nodes = vec![];
1087    nodes.iter().for_each(|node| {
1088        if node.is_bootnode {
1089            bootnodes.push(node)
1090        } else {
1091            other_nodes.push(node)
1092        }
1093    });
1094
1095    if bootnodes.is_empty() && !no_default_bootnodes {
1096        bootnodes.push(other_nodes.remove(0))
1097    }
1098
1099    (bootnodes, other_nodes)
1100}
1101
1102// Generate a bootnode multiaddress and return as string
1103fn generate_bootnode_addr(
1104    node: &NetworkNode,
1105    ip: &IpAddr,
1106    port: u16,
1107) -> Result<String, GeneratorError> {
1108    generators::generate_node_bootnode_addr(
1109        &node.spec.peer_id,
1110        ip,
1111        port,
1112        node.args().as_ref(),
1113        &node.spec.p2p_cert_hash,
1114    )
1115}
1116// Validate that the config fulfill all the requirements of the provider
1117fn validate_spec_with_provider_capabilities(
1118    network_spec: &NetworkSpec,
1119    capabilities: &ProviderCapabilities,
1120) -> Result<(), anyhow::Error> {
1121    let mut errs: Vec<String> = vec![];
1122
1123    if capabilities.requires_image {
1124        // Relaychain
1125        if let Some(relaychain) = network_spec.try_relaychain() {
1126            if relaychain.default_image.is_none() {
1127                // we should check if each node have an image
1128                let nodes = &relaychain.nodes;
1129                if nodes.iter().any(|node| node.image.is_none()) {
1130                    errs.push(String::from(
1131                        "Missing image for node, and not default is set at relaychain",
1132                    ));
1133                }
1134            }
1135        };
1136
1137        // Paras
1138        for para in &network_spec.parachains {
1139            if para.default_image.is_none() {
1140                let nodes = &para.collators;
1141                if nodes.iter().any(|node| node.image.is_none()) {
1142                    errs.push(format!(
1143                        "Missing image for node, and not default is set at parachain {}",
1144                        para.id
1145                    ));
1146                }
1147            }
1148        }
1149    } else {
1150        // native
1151        // We need to get all the `cmds` and verify if are part of the path
1152        let mut cmds: HashSet<&str> = Default::default();
1153        if let Some(relaychain) = network_spec.try_relaychain() {
1154            if let Some(cmd) = relaychain.default_command.as_ref() {
1155                cmds.insert(cmd.as_str());
1156            }
1157            for node in relaychain.nodes.iter() {
1158                cmds.insert(node.command());
1159            }
1160        }
1161
1162        if let Some(jamchain) = network_spec.jamchain.as_ref() {
1163            for node in jamchain.nodes.iter() {
1164                cmds.insert(node.command());
1165            }
1166        }
1167
1168        // Paras
1169        for para in &network_spec.parachains {
1170            if let Some(cmd) = para.default_command.as_ref() {
1171                cmds.insert(cmd.as_str());
1172            }
1173
1174            for node in para.collators.iter() {
1175                cmds.insert(node.command());
1176            }
1177        }
1178
1179        // now check the binaries
1180        let path = std::env::var("PATH").unwrap_or_default(); // path should always be set
1181        trace!("current PATH: {path}");
1182        let parts: Vec<_> = path.split(":").collect();
1183        for cmd in cmds {
1184            let missing = if cmd.contains('/') {
1185                trace!("checking {cmd}");
1186                if std::fs::metadata(cmd).is_err() {
1187                    true
1188                } else {
1189                    info!("🔎  We will use the full path {cmd} to spawn nodes.");
1190                    false
1191                }
1192            } else {
1193                // should be in the PATH
1194                !parts.iter().any(|part| {
1195                    let path_to = format!("{part}/{cmd}");
1196                    trace!("checking {path_to}");
1197                    let check_result = std::fs::metadata(&path_to);
1198                    trace!("result {:?}", check_result);
1199                    if check_result.is_ok() {
1200                        info!("🔎  We will use the cmd: '{cmd}' at path {path_to} to spawn nodes.");
1201                        true
1202                    } else {
1203                        false
1204                    }
1205                })
1206            };
1207
1208            if missing {
1209                errs.push(help_msg(cmd));
1210            }
1211        }
1212    }
1213
1214    if !errs.is_empty() {
1215        let msg = errs.join("\n");
1216        return Err(anyhow::anyhow!(format!("Invalid configuration: \n {msg}")));
1217    }
1218
1219    Ok(())
1220}
1221
1222fn help_msg(cmd: &str) -> String {
1223    match cmd {
1224        "parachain-template-node" | "solochain-template-node" | "minimal-template-node" => {
1225            format!("Missing binary {cmd}, compile by running: \n\tcargo build --package {cmd} --release")
1226        },
1227        "polkadot" => {
1228            format!("Missing binary {cmd}, compile by running (in the polkadot-sdk repo): \n\t cargo build --locked --release --features fast-runtime --bin {cmd} --bin polkadot-prepare-worker --bin polkadot-execute-worker")
1229        },
1230        "polkadot-parachain" => {
1231            format!("Missing binary {cmd}, compile by running (in the polkadot-sdk repo): \n\t cargo build --release --locked -p {cmd}-bin --bin {cmd}")
1232        },
1233        _ => {
1234            format!("Missing binary {cmd}, please compile it.")
1235        },
1236    }
1237}
1238
1239/// Allow to set the default concurrency through env var `ZOMBIE_SPAWN_CONCURRENCY`
1240fn spawn_concurrency_from_env() -> Option<usize> {
1241    if let Ok(concurrency) = env::var("ZOMBIE_SPAWN_CONCURRENCY") {
1242        concurrency.parse::<usize>().ok()
1243    } else {
1244        None
1245    }
1246}
1247
1248fn calculate_concurrency(spec: &NetworkSpec) -> Result<(usize, bool), anyhow::Error> {
1249    let desired_spawn_concurrency = match (
1250        spawn_concurrency_from_env(),
1251        spec.global_settings.spawn_concurrency(),
1252    ) {
1253        (Some(n), _) => Some(n),
1254        (None, Some(n)) => Some(n),
1255        _ => None,
1256    };
1257
1258    let (spawn_concurrency, limited_by_tokens) =
1259        if let Some(spawn_concurrency) = desired_spawn_concurrency {
1260            if spawn_concurrency == 1 {
1261                (1, false)
1262            } else if has_tokens(&serde_json::to_string(spec)?) {
1263                (1, true)
1264            } else {
1265                (spawn_concurrency, false)
1266            }
1267        } else {
1268            // not set
1269            if has_tokens(&serde_json::to_string(spec)?) {
1270                (1, true)
1271            } else {
1272                // use 100 as max concurrency, we can set a max by provider later
1273                (100, false)
1274            }
1275        };
1276
1277    Ok((spawn_concurrency, limited_by_tokens))
1278}
1279
1280/// Build deterministic dependency **levels** among the given nodes.
1281/// - Only dependencies **between nodes in `nodes`** are considered.
1282/// - Unknown/out-of-scope references are ignored.
1283/// - Self-dependencies are ignored.
1284fn dependency_levels_among<'a>(
1285    nodes: &'a [&'a NodeSpec],
1286) -> Result<Vec<Vec<&'a NodeSpec>>, OrchestratorError> {
1287    let by_name = nodes
1288        .iter()
1289        .map(|n| (n.name.as_str(), *n))
1290        .collect::<HashMap<_, _>>();
1291
1292    let mut graph = HashMap::with_capacity(nodes.len());
1293    let mut indegree = HashMap::with_capacity(nodes.len());
1294
1295    for node in nodes {
1296        graph.insert(node.name.as_str(), Vec::new());
1297        indegree.insert(node.name.as_str(), 0);
1298    }
1299
1300    // build dependency graph
1301    for &node in nodes {
1302        if let Ok(args_json) = serde_json::to_string(&node.args) {
1303            // collect dependencies
1304            let unique_deps = get_tokens_to_replace(&args_json)
1305                .into_iter()
1306                .filter(|dep| dep != &node.name)
1307                .filter_map(|dep| by_name.get(dep.as_str()))
1308                .map(|&dep_node| dep_node.name.as_str())
1309                .collect::<HashSet<_>>();
1310
1311            for dep_name in unique_deps {
1312                graph
1313                    .get_mut(dep_name)
1314                    .expect(&format!("{GRAPH_CONTAINS_DEP} {THIS_IS_A_BUG}"))
1315                    .push(node);
1316                *indegree
1317                    .get_mut(node.name.as_str())
1318                    .expect(&format!("{INDEGREE_CONTAINS_NAME} {THIS_IS_A_BUG}")) += 1;
1319            }
1320        }
1321    }
1322
1323    // find all nodes with no dependencies
1324    let mut queue = nodes
1325        .iter()
1326        .filter(|n| {
1327            *indegree
1328                .get(n.name.as_str())
1329                .expect(&format!("{INDEGREE_CONTAINS_NAME} {THIS_IS_A_BUG}"))
1330                == 0
1331        })
1332        .copied()
1333        .collect::<VecDeque<_>>();
1334
1335    let mut processed_count = 0;
1336    let mut levels = Vec::new();
1337
1338    // Kahn's algorithm
1339    while !queue.is_empty() {
1340        let level_size = queue.len();
1341        let mut current_level = Vec::with_capacity(level_size);
1342
1343        for _ in 0..level_size {
1344            let n = queue
1345                .pop_front()
1346                .expect(&format!("{QUEUE_NOT_EMPTY} {THIS_IS_A_BUG}"));
1347            current_level.push(n);
1348            processed_count += 1;
1349
1350            for &neighbour in graph
1351                .get(n.name.as_str())
1352                .expect(&format!("{GRAPH_CONTAINS_NAME} {THIS_IS_A_BUG}"))
1353            {
1354                let neighbour_indegree = indegree
1355                    .get_mut(neighbour.name.as_str())
1356                    .expect(&format!("{INDEGREE_CONTAINS_NAME} {THIS_IS_A_BUG}"));
1357                *neighbour_indegree -= 1;
1358
1359                if *neighbour_indegree == 0 {
1360                    queue.push_back(neighbour);
1361                }
1362            }
1363        }
1364
1365        current_level.sort_by_key(|n| &n.name);
1366        levels.push(current_level);
1367    }
1368
1369    // cycles detected, e.g A -> B -> A
1370    if processed_count != nodes.len() {
1371        return Err(OrchestratorError::InvalidConfig(
1372            "Tokens have cyclical dependencies".to_string(),
1373        ));
1374    }
1375
1376    Ok(levels)
1377}
1378
1379// TODO: get the fs from `DynNamespace` will make this not needed
1380// but the FileSystem trait isn't object-safe so we can't pass around
1381// as `dyn FileSystem`. We can refactor or using some `erase` techniques
1382// to resolve this and remove this struct
1383// TODO (Loris): Probably we could have a .scoped(base_dir) method on the
1384// filesystem itself (the trait), so it will return this and we can move this
1385// directly to the support crate, it can be useful in the future
1386#[derive(Clone, Debug)]
1387pub struct ScopedFilesystem<'a, FS: FileSystem> {
1388    fs: &'a FS,
1389    base_dir: &'a str,
1390}
1391
1392impl<'a, FS: FileSystem> ScopedFilesystem<'a, FS> {
1393    pub fn new(fs: &'a FS, base_dir: &'a str) -> Self {
1394        Self { fs, base_dir }
1395    }
1396
1397    async fn copy_files(&self, files: Vec<&TransferedFile>) -> Result<(), FileSystemError> {
1398        for file in files {
1399            let full_remote_path = PathBuf::from(format!(
1400                "{}/{}",
1401                self.base_dir,
1402                file.remote_path.to_string_lossy()
1403            ));
1404            trace!("coping file: {file}");
1405            self.fs
1406                .copy(file.local_path.as_path(), full_remote_path)
1407                .await?;
1408        }
1409        Ok(())
1410    }
1411
1412    async fn read(&self, file: impl AsRef<Path>) -> Result<Vec<u8>, FileSystemError> {
1413        let file = file.as_ref();
1414
1415        let full_path = if file.is_absolute() {
1416            file.to_owned()
1417        } else {
1418            PathBuf::from(format!("{}/{}", self.base_dir, file.to_string_lossy()))
1419        };
1420        let content = self.fs.read(full_path).await?;
1421        Ok(content)
1422    }
1423
1424    async fn read_to_string(&self, file: impl AsRef<Path>) -> Result<String, FileSystemError> {
1425        let file = file.as_ref();
1426
1427        let full_path = if file.is_absolute() {
1428            file.to_owned()
1429        } else {
1430            PathBuf::from(format!("{}/{}", self.base_dir, file.to_string_lossy()))
1431        };
1432        let content = self.fs.read_to_string(full_path).await?;
1433        Ok(content)
1434    }
1435
1436    async fn create_dir(&self, path: impl AsRef<Path>) -> Result<(), FileSystemError> {
1437        let path = PathBuf::from(format!(
1438            "{}/{}",
1439            self.base_dir,
1440            path.as_ref().to_string_lossy()
1441        ));
1442        self.fs.create_dir(path).await
1443    }
1444
1445    async fn create_dir_all(&self, path: impl AsRef<Path>) -> Result<(), FileSystemError> {
1446        let path = PathBuf::from(format!(
1447            "{}/{}",
1448            self.base_dir,
1449            path.as_ref().to_string_lossy()
1450        ));
1451        self.fs.create_dir_all(path).await
1452    }
1453
1454    async fn write(
1455        &self,
1456        path: impl AsRef<Path>,
1457        contents: impl AsRef<[u8]> + Send,
1458    ) -> Result<(), FileSystemError> {
1459        let path = path.as_ref();
1460
1461        let full_path = if path.is_absolute() {
1462            path.to_owned()
1463        } else {
1464            PathBuf::from(format!("{}/{}", self.base_dir, path.to_string_lossy()))
1465        };
1466
1467        self.fs.write(full_path, contents).await
1468    }
1469
1470    /// Get the full_path in the scoped FS
1471    fn full_path(&self, path: impl AsRef<Path>) -> PathBuf {
1472        let path = path.as_ref();
1473
1474        let full_path = if path.is_absolute() {
1475            path.to_owned()
1476        } else {
1477            PathBuf::from(format!("{}/{}", self.base_dir, path.to_string_lossy()))
1478        };
1479
1480        full_path
1481    }
1482
1483    /// Get the base_dir in the scoped FS
1484    fn base_dir(&self) -> &str {
1485        self.base_dir
1486    }
1487}
1488
1489#[derive(Clone, Debug)]
1490pub enum ZombieRole {
1491    Temp,
1492    Node,
1493    Bootnode,
1494    Collator,
1495    CumulusCollator,
1496    Companion,
1497}
1498
1499// re-exports
1500pub use network::{AddCollatorOptions, AddNodeOptions};
1501pub use network_helper::metrics;
1502pub use sc_chain_spec;
1503
1504#[cfg(test)]
1505mod tests {
1506    use configuration::{GlobalSettingsBuilder, NetworkConfigBuilder};
1507    use lazy_static::lazy_static;
1508    use tokio::sync::Mutex;
1509
1510    use super::*;
1511
1512    const ENV_KEY: &str = "ZOMBIE_SPAWN_CONCURRENCY";
1513    // mutex for test that use env
1514    lazy_static! {
1515        static ref ENV_MUTEX: Mutex<()> = Mutex::new(());
1516    }
1517
1518    fn set_env(concurrency: Option<u32>) {
1519        if let Some(value) = concurrency {
1520            env::set_var(ENV_KEY, value.to_string());
1521        } else {
1522            env::remove_var(ENV_KEY);
1523        }
1524    }
1525
1526    fn generate(
1527        with_image: bool,
1528        with_cmd: Option<&'static str>,
1529    ) -> Result<NetworkConfig, Vec<anyhow::Error>> {
1530        NetworkConfigBuilder::new()
1531            .with_relaychain(|r| {
1532                let mut relay = r
1533                    .with_chain("rococo-local")
1534                    .with_default_command(with_cmd.unwrap_or("polkadot"));
1535                if with_image {
1536                    relay = relay.with_default_image("docker.io/parity/polkadot")
1537                }
1538
1539                relay
1540                    .with_validator(|node| node.with_name("alice"))
1541                    .with_validator(|node| node.with_name("bob"))
1542            })
1543            .with_parachain(|p| {
1544                p.with_id(2000).cumulus_based(true).with_collator(|n| {
1545                    let node = n
1546                        .with_name("collator")
1547                        .with_command(with_cmd.unwrap_or("polkadot-parachain"));
1548                    if with_image {
1549                        node.with_image("docker.io/paritypr/test-parachain")
1550                    } else {
1551                        node
1552                    }
1553                })
1554            })
1555            .build()
1556    }
1557
1558    fn get_node_with_dependencies(name: &str, dependencies: Option<Vec<&NodeSpec>>) -> NodeSpec {
1559        let mut spec = NodeSpec {
1560            name: name.to_string(),
1561            ..Default::default()
1562        };
1563        if let Some(dependencies) = dependencies {
1564            for node in dependencies {
1565                spec.args.push(
1566                    format!("{{{{ZOMBIE:{}:someField}}}}", node.name)
1567                        .as_str()
1568                        .into(),
1569                );
1570            }
1571        }
1572        spec
1573    }
1574
1575    fn verify_levels(actual_levels: Vec<Vec<&NodeSpec>>, expected_levels: Vec<Vec<&str>>) {
1576        actual_levels
1577            .iter()
1578            .zip(expected_levels)
1579            .for_each(|(actual_level, expected_level)| {
1580                assert_eq!(actual_level.len(), expected_level.len());
1581                actual_level
1582                    .iter()
1583                    .zip(expected_level.iter())
1584                    .for_each(|(node, expected_name)| assert_eq!(node.name, *expected_name));
1585            });
1586    }
1587
1588    #[tokio::test]
1589    async fn valid_config_with_image() {
1590        let network_config = generate(true, None).unwrap();
1591        let spec = NetworkSpec::from_config(&network_config).await.unwrap();
1592        let caps = ProviderCapabilities {
1593            requires_image: true,
1594            has_resources: false,
1595            prefix_with_full_path: false,
1596            use_default_ports_in_cmd: false,
1597        };
1598
1599        let valid = validate_spec_with_provider_capabilities(&spec, &caps);
1600        assert!(valid.is_ok())
1601    }
1602
1603    #[tokio::test]
1604    async fn invalid_config_without_image() {
1605        let network_config = generate(false, None).unwrap();
1606        let spec = NetworkSpec::from_config(&network_config).await.unwrap();
1607        let caps = ProviderCapabilities {
1608            requires_image: true,
1609            has_resources: false,
1610            prefix_with_full_path: false,
1611            use_default_ports_in_cmd: false,
1612        };
1613
1614        let valid = validate_spec_with_provider_capabilities(&spec, &caps);
1615        assert!(valid.is_err())
1616    }
1617
1618    #[tokio::test]
1619    async fn invalid_config_missing_cmd() {
1620        let network_config = generate(false, Some("other")).unwrap();
1621        let spec = NetworkSpec::from_config(&network_config).await.unwrap();
1622        let caps = ProviderCapabilities {
1623            requires_image: false,
1624            has_resources: false,
1625            prefix_with_full_path: false,
1626            use_default_ports_in_cmd: false,
1627        };
1628
1629        let valid = validate_spec_with_provider_capabilities(&spec, &caps);
1630        assert!(valid.is_err())
1631    }
1632
1633    #[tokio::test]
1634    async fn valid_config_present_cmd() {
1635        let network_config = generate(false, Some("cargo")).unwrap();
1636        let spec = NetworkSpec::from_config(&network_config).await.unwrap();
1637        let caps = ProviderCapabilities {
1638            requires_image: false,
1639            has_resources: false,
1640            prefix_with_full_path: false,
1641            use_default_ports_in_cmd: false,
1642        };
1643
1644        let valid = validate_spec_with_provider_capabilities(&spec, &caps);
1645        println!("{valid:?}");
1646        assert!(valid.is_ok())
1647    }
1648
1649    #[tokio::test]
1650    async fn default_spawn_concurrency() {
1651        let _g = ENV_MUTEX.lock().await;
1652        set_env(None);
1653        let network_config = generate(false, Some("cargo")).unwrap();
1654        let spec = NetworkSpec::from_config(&network_config).await.unwrap();
1655        let (concurrency, _) = calculate_concurrency(&spec).unwrap();
1656        assert_eq!(concurrency, 100);
1657    }
1658
1659    #[tokio::test]
1660    async fn set_spawn_concurrency() {
1661        let _g = ENV_MUTEX.lock().await;
1662        set_env(None);
1663
1664        let network_config = generate(false, Some("cargo")).unwrap();
1665        let mut spec = NetworkSpec::from_config(&network_config).await.unwrap();
1666
1667        let global_settings = GlobalSettingsBuilder::new()
1668            .with_spawn_concurrency(4)
1669            .build()
1670            .unwrap();
1671
1672        spec.set_global_settings(global_settings);
1673        let (concurrency, limited) = calculate_concurrency(&spec).unwrap();
1674        assert_eq!(concurrency, 4);
1675        assert!(!limited);
1676    }
1677
1678    #[tokio::test]
1679    async fn set_spawn_concurrency_but_limited() {
1680        let _g = ENV_MUTEX.lock().await;
1681        set_env(None);
1682
1683        let network_config = generate(false, Some("cargo")).unwrap();
1684        let mut spec = NetworkSpec::from_config(&network_config).await.unwrap();
1685
1686        let global_settings = GlobalSettingsBuilder::new()
1687            .with_spawn_concurrency(4)
1688            .build()
1689            .unwrap();
1690
1691        spec.set_global_settings(global_settings);
1692        let node = spec.relaychain_mut().nodes.first_mut().unwrap();
1693        node.args
1694            .push("--bootnodes {{ZOMBIE:bob:multiAddress')}}".into());
1695        let (concurrency, limited) = calculate_concurrency(&spec).unwrap();
1696        assert_eq!(concurrency, 1);
1697        assert!(limited);
1698    }
1699
1700    #[tokio::test]
1701    async fn set_spawn_concurrency_from_env() {
1702        let _g = ENV_MUTEX.lock().await;
1703        set_env(Some(10));
1704
1705        let network_config = generate(false, Some("cargo")).unwrap();
1706        let spec = NetworkSpec::from_config(&network_config).await.unwrap();
1707        let (concurrency, limited) = calculate_concurrency(&spec).unwrap();
1708        assert_eq!(concurrency, 10);
1709        assert!(!limited);
1710    }
1711
1712    #[tokio::test]
1713    async fn set_spawn_concurrency_from_env_but_limited() {
1714        let _g = ENV_MUTEX.lock().await;
1715        set_env(Some(12));
1716
1717        let network_config = generate(false, Some("cargo")).unwrap();
1718        let mut spec = NetworkSpec::from_config(&network_config).await.unwrap();
1719        let node = spec.relaychain_mut().nodes.first_mut().unwrap();
1720        node.args
1721            .push("--bootnodes {{ZOMBIE:bob:multiAddress')}}".into());
1722        let (concurrency, limited) = calculate_concurrency(&spec).unwrap();
1723        assert_eq!(concurrency, 1);
1724        assert!(limited);
1725    }
1726
1727    #[test]
1728    fn dependency_levels_among_should_work() {
1729        // no nodes
1730        assert!(dependency_levels_among(&[]).unwrap().is_empty());
1731
1732        // one node
1733        let alice = get_node_with_dependencies("alice", None);
1734        let nodes = [&alice];
1735
1736        let levels = dependency_levels_among(&nodes).unwrap();
1737        let expected = vec![vec!["alice"]];
1738
1739        verify_levels(levels, expected);
1740
1741        // two independent nodes
1742        let alice = get_node_with_dependencies("alice", None);
1743        let bob = get_node_with_dependencies("bob", None);
1744        let nodes = [&alice, &bob];
1745
1746        let levels = dependency_levels_among(&nodes).unwrap();
1747        let expected = vec![vec!["alice", "bob"]];
1748
1749        verify_levels(levels, expected);
1750
1751        // alice -> bob -> charlie
1752        let alice = get_node_with_dependencies("alice", None);
1753        let bob = get_node_with_dependencies("bob", Some(vec![&alice]));
1754        let charlie = get_node_with_dependencies("charlie", Some(vec![&bob]));
1755        let nodes = [&alice, &bob, &charlie];
1756
1757        let levels = dependency_levels_among(&nodes).unwrap();
1758        let expected = vec![vec!["alice"], vec!["bob"], vec!["charlie"]];
1759
1760        verify_levels(levels, expected);
1761
1762        //         ┌─> bob
1763        // alice ──|
1764        //         └─> charlie
1765        let alice = get_node_with_dependencies("alice", None);
1766        let bob = get_node_with_dependencies("bob", Some(vec![&alice]));
1767        let charlie = get_node_with_dependencies("charlie", Some(vec![&alice]));
1768        let nodes = [&alice, &bob, &charlie];
1769
1770        let levels = dependency_levels_among(&nodes).unwrap();
1771        let expected = vec![vec!["alice"], vec!["bob", "charlie"]];
1772
1773        verify_levels(levels, expected);
1774
1775        //         ┌─>   bob  ──┐
1776        // alice ──|            ├─> dave
1777        //         └─> charlie  ┘
1778        let alice = get_node_with_dependencies("alice", None);
1779        let bob = get_node_with_dependencies("bob", Some(vec![&alice]));
1780        let charlie = get_node_with_dependencies("charlie", Some(vec![&alice]));
1781        let dave = get_node_with_dependencies("dave", Some(vec![&charlie, &bob]));
1782        let nodes = [&alice, &bob, &charlie, &dave];
1783
1784        let levels = dependency_levels_among(&nodes).unwrap();
1785        let expected = vec![vec!["alice"], vec!["bob", "charlie"], vec!["dave"]];
1786
1787        verify_levels(levels, expected);
1788    }
1789
1790    #[test]
1791    fn dependency_levels_among_should_detect_cycles() {
1792        let mut alice = get_node_with_dependencies("alice", None);
1793        let bob = get_node_with_dependencies("bob", Some(vec![&alice]));
1794        alice.args.push("{{ZOMBIE:bob:someField}}".into());
1795
1796        assert!(dependency_levels_among(&[&alice, &bob]).is_err())
1797    }
1798}