1#![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};
35pub 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 debug!(network_spec = ?network_spec,"Network spec to spawn");
184
185 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 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 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 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 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 let mut relaychain = network_spec
254 .relaychain
255 .take()
256 .expect("checked to be `Some` above; qed");
257
258 relaychain.chain_spec.build(&ns, &scoped_fs).await?;
260
261 debug!("relaychain spec built!");
262 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 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 ¶_to_register_in_genesis {
285 let genesis_config = para.get_genesis_config()?;
286 para_artifacts.push(genesis_config)
287 }
288
289 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 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 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 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 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 relaychain
342 .chain_spec
343 .build_raw(&ns, &scoped_fs, None)
344 .await?;
345
346 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 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 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 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 for para in &network_spec.parachains {
380 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 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 !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 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 network_spec.relaychain = Some(relaychain);
458
459 let (bootnodes, relaynodes) =
460 split_nodes_by_bootnodes(&network_spec.relaychain().nodes, false);
461
462 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 let mut node_ws_url: String = "".to_string();
498
499 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 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 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 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 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 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 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 for para in para_to_register_with_extrinsic {
591 let register_para_options: RegisterParachainOptions = RegisterParachainOptions {
592 id: para.id,
593 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, 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 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 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 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(¶.collators, para.no_default_bootnodes);
673
674 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 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, ¶chain, &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 for level in dependency_levels_among(&collators)? {
712 for node in self
713 .spawn_parachain_level(&level, ¶chain, &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 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 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 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 let jam_config = jam_config::generate(jam_spec)?;
776 scoped_fs
778 .write(
779 "jam_config.json",
780 serde_json::to_string_pretty(&jam_config)?,
781 )
782 .await?;
783 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 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 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 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 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 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 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 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
926fn 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 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
1013async 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
1078fn split_nodes_by_bootnodes(
1081 nodes: &[NodeSpec],
1082 no_default_bootnodes: bool,
1083) -> (Vec<&NodeSpec>, Vec<&NodeSpec>) {
1084 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
1102fn 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}
1116fn 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 if let Some(relaychain) = network_spec.try_relaychain() {
1126 if relaychain.default_image.is_none() {
1127 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 for para in &network_spec.parachains {
1139 if para.default_image.is_none() {
1140 let nodes = ¶.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 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 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 let path = std::env::var("PATH").unwrap_or_default(); 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 !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
1239fn 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 if has_tokens(&serde_json::to_string(spec)?) {
1270 (1, true)
1271 } else {
1272 (100, false)
1274 }
1275 };
1276
1277 Ok((spawn_concurrency, limited_by_tokens))
1278}
1279
1280fn 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 for &node in nodes {
1302 if let Ok(args_json) = serde_json::to_string(&node.args) {
1303 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 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 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 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#[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 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 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
1499pub 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 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 assert!(dependency_levels_among(&[]).unwrap().is_empty());
1731
1732 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 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 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 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 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}