diff --git a/Cargo.lock b/Cargo.lock index c351c0284..0912c7382 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -4644,6 +4644,7 @@ dependencies = [ "opendut-netbird-client-api", "opendut-telemetry", "opendut-util", + "opendut-viper-rt", "opentelemetry", "opentelemetry_sdk", "ping-rs", diff --git a/opendut-edgar/Cargo.toml b/opendut-edgar/Cargo.toml index c9cd05043..e9da9b6d7 100644 --- a/opendut-edgar/Cargo.toml +++ b/opendut-edgar/Cargo.toml @@ -13,6 +13,7 @@ opendut-netbird-client-api = { workspace = true } opendut-model = { workspace = true } opendut-telemetry = { workspace = true } opendut-util = { workspace = true, features = ["crypto", "pem", "settings", "serde"] } +opendut-viper-rt = { workspace = true, optional = true } anyhow = { workspace = true } async-trait = { workspace = true } @@ -74,7 +75,9 @@ shadow-rs = { workspace = true, default-features = true } [features] integration_testing = [] -viper = [] +viper = [ + "opendut-viper-rt/run", +] [lints] workspace = true diff --git a/opendut-edgar/src/service/mod.rs b/opendut-edgar/src/service/mod.rs index f13347190..2dd85a9b4 100644 --- a/opendut-edgar/src/service/mod.rs +++ b/opendut-edgar/src/service/mod.rs @@ -5,6 +5,9 @@ pub mod process_manager; pub mod start; pub mod vpn; +#[cfg(feature = "viper")] +pub mod viper_run_manager; + mod can; mod network_metrics; mod tasks; diff --git a/opendut-edgar/src/service/peer_configuration.rs b/opendut-edgar/src/service/peer_configuration.rs index 13e8be5bb..d06df5ae5 100644 --- a/opendut-edgar/src/service/peer_configuration.rs +++ b/opendut-edgar/src/service/peer_configuration.rs @@ -21,6 +21,8 @@ pub struct ApplyPeerConfigurationParams { pub network_interface_management: NetworkInterfaceManagement, pub executor_manager: ExecutorManagerRef, pub metrics_manager: NetworkMetricsManagerRef, + #[cfg(feature = "viper")] + pub viper_run_manager: crate::service::viper_run_manager::ViperRunManagerRef, } #[derive(Clone)] pub enum NetworkInterfaceManagement { @@ -102,13 +104,18 @@ async fn apply_peer_configuration(params: ApplyPeerConfigurationParams) -> Colle let ApplyPeerConfigurationParams { peer_configuration, network_interface_management, - executor_manager, metrics_manager + executor_manager, + metrics_manager, + #[cfg(feature = "viper")] + viper_run_manager, } = params; let resolver = runner::task_resolver::ServiceTaskResolver::new( peer_configuration.clone(), network_interface_management.clone(), Arc::clone(&metrics_manager), + #[cfg(feature = "viper")] + Arc::clone(&viper_run_manager), ); let result = runner::service_runner::run_tasks(peer_configuration.clone(), resolver).await; if result.success { @@ -127,4 +134,3 @@ async fn apply_peer_configuration(params: ApplyPeerConfigurationParams) -> Colle debug!("Peer configuration has been successfully applied."); result } - diff --git a/opendut-edgar/src/service/peer_messaging_client.rs b/opendut-edgar/src/service/peer_messaging_client.rs index fdd0a6a06..4e26d850e 100644 --- a/opendut-edgar/src/service/peer_messaging_client.rs +++ b/opendut-edgar/src/service/peer_messaging_client.rs @@ -26,7 +26,6 @@ use crate::service::network_metrics::manager::{NetworkMetricsManager, NetworkMet use crate::service::peer_configuration::{ApplyPeerConfigurationParams, NetworkInterfaceManagement}; use crate::service::test_execution::executor_manager::{ExecutorManager, ExecutorManagerRef}; - pub struct PeerMessagingClient { self_id: PeerId, network_interface_management: NetworkInterfaceManagement, @@ -34,6 +33,8 @@ pub struct PeerMessagingClient { metrics_manager: NetworkMetricsManagerRef, carl_disconnect_timeout: Duration, tx_peer_configuration: mpsc::Sender, + #[cfg(feature = "viper")] + viper_run_manager: crate::service::viper_run_manager::ViperRunManagerRef, } @@ -63,6 +64,8 @@ impl PeerMessagingClient { let metrics_manager: NetworkMetricsManagerRef = NetworkMetricsManager::load(settings)?; + #[cfg(feature = "viper")] + let viper_run_manager = crate::service::viper_run_manager::ViperRunManager::create(); Ok(PeerMessagingClient { self_id, @@ -71,6 +74,8 @@ impl PeerMessagingClient { metrics_manager, carl_disconnect_timeout, tx_peer_configuration, + #[cfg(feature = "viper")] + viper_run_manager, }) } @@ -266,6 +271,8 @@ impl PeerMessagingClient { network_interface_management: self.network_interface_management.clone(), executor_manager: Arc::clone(&self.executor_manager), metrics_manager: Arc::clone(&self.metrics_manager), + #[cfg(feature = "viper")] + viper_run_manager: Arc::clone(&self.viper_run_manager), }; peer_configuration_sender.send(apply_config_params).await?; diff --git a/opendut-edgar/src/service/tasks/create_ethernet_bridge.rs b/opendut-edgar/src/service/tasks/create_ethernet_bridge.rs index 7ed637ad5..605fbcd98 100644 --- a/opendut-edgar/src/service/tasks/create_ethernet_bridge.rs +++ b/opendut-edgar/src/service/tasks/create_ethernet_bridge.rs @@ -107,7 +107,6 @@ mod tests { use opendut_model::peer::configuration::{parameter, ParameterTarget, PeerConfiguration}; use opendut_model::util::net::NetworkInterfaceName; use std::sync::Arc; - use crate::service::tasks::runner; use crate::service::tasks::runner::task_resolver::ServiceTaskResolver; use crate::service::tasks::testing::NetworkInterfaceNameExt; @@ -201,10 +200,16 @@ mod tests { can_manager }; let metrics_manager = NetworkMetricsManager::new(NetworkMetricsOptions::default()); - let service_task_resolver = runner::task_resolver::ServiceTaskResolver::new( + + #[cfg(feature = "viper")] + let viper_run_manager = crate::service::viper_run_manager::ViperRunManager::create(); + + let service_task_resolver = ServiceTaskResolver::new( peer_configuration.clone(), network_interface_management.clone(), Arc::clone(&metrics_manager), + #[cfg(feature = "viper")] + Arc::clone(&viper_run_manager), ); Self { diff --git a/opendut-edgar/src/service/tasks/create_viper_runtime.rs b/opendut-edgar/src/service/tasks/create_viper_runtime.rs new file mode 100644 index 000000000..cba43575a --- /dev/null +++ b/opendut-edgar/src/service/tasks/create_viper_runtime.rs @@ -0,0 +1,54 @@ +use async_trait::async_trait; +use opendut_model::peer::configuration::parameter; +use crate::common::task::{Success, Task, TaskAbsent, TaskStateFulfilled}; +use crate::service::viper_run_manager::ViperRunManagerRef; + +pub struct CreateViperTestRun { + pub parameter: parameter::TestRunReport, + pub viper_run_manager: ViperRunManagerRef, +} + +#[async_trait] +impl Task for CreateViperTestRun { + fn description(&self) -> String { + format!("Create VIPER test run '{}'", self.parameter.run_id) + } + + async fn check_present(&self) -> anyhow::Result { + if self.viper_run_manager.contains_test_run(&self.parameter.run_id).await { + Ok(TaskStateFulfilled::Yes) + } else { + Ok(TaskStateFulfilled::No) + } + } + + async fn make_present(&self) -> anyhow::Result { + let run_id = self.parameter.run_id; + let source_code = Clone::clone(&self.parameter.source_code); + + self.viper_run_manager.start_test_run(run_id, source_code).await; + + Ok(Success::default()) + } +} + +#[async_trait] +impl TaskAbsent for CreateViperTestRun { + async fn check_absent( + &self, + ) -> anyhow::Result { + if self.viper_run_manager.contains_test_run(&self.parameter.run_id).await { + Ok(TaskStateFulfilled::No) + } else { + Ok(TaskStateFulfilled::Yes) + } + } + + async fn make_absent(&self) -> anyhow::Result { + self.viper_run_manager.abort_test_run( + &self.parameter.run_id, + ).await; + + Ok(Success::default()) + } +} diff --git a/opendut-edgar/src/service/tasks/mod.rs b/opendut-edgar/src/service/tasks/mod.rs index 8d49e1d8f..04f1a9c68 100644 --- a/opendut-edgar/src/service/tasks/mod.rs +++ b/opendut-edgar/src/service/tasks/mod.rs @@ -12,3 +12,6 @@ mod create_gre_interfaces; mod manage_joined_interfaces; mod require_interface_up; mod setup_cluster_metrics; + +#[cfg(feature = "viper")] +mod create_viper_runtime; diff --git a/opendut-edgar/src/service/tasks/runner/task_resolver.rs b/opendut-edgar/src/service/tasks/runner/task_resolver.rs index fb2712cef..936d97d51 100644 --- a/opendut-edgar/src/service/tasks/runner/task_resolver.rs +++ b/opendut-edgar/src/service/tasks/runner/task_resolver.rs @@ -14,6 +14,8 @@ pub struct ServiceTaskResolver { peer_configuration: PeerConfiguration, network_interface_management: NetworkInterfaceManagement, metrics_manager: NetworkMetricsManagerRef, + #[cfg(feature = "viper")] + viper_run_manager: crate::service::viper_run_manager::ViperRunManagerRef, } impl ServiceTaskResolver { @@ -21,11 +23,15 @@ impl ServiceTaskResolver { peer_configuration: PeerConfiguration, network_interface_management: NetworkInterfaceManagement, metrics_manager: NetworkMetricsManagerRef, + #[cfg(feature = "viper")] + viper_run_manager: crate::service::viper_run_manager::ViperRunManagerRef, ) -> Self { Self { peer_configuration, network_interface_management, metrics_manager, + #[cfg(feature = "viper")] + viper_run_manager, } } } @@ -65,8 +71,16 @@ impl TaskResolver for ServiceTaskResolver { tasks.push(Box::new(tasks::can_local_route::CanLocalRoute { parameter: parameter.value.clone(), network_interface_manager: network_interface_manager.clone(), can_fd: false })); tasks.push(Box::new(tasks::can_local_route::CanLocalRoute { parameter: parameter.value.clone(), network_interface_manager, can_fd: true })); } - ParameterVariant::TestRunReport(_parameter) => { - todo!("EDGAR does not yet have a Task for handling TestRunReport parameters."); + #[cfg(feature = "viper")] + ParameterVariant::TestRunReport(parameter) => { + tasks.push(Box::new(tasks::create_viper_runtime::CreateViperTestRun { + parameter: parameter.value.clone(), + viper_run_manager: self.viper_run_manager.clone(), + })); + } + #[cfg(not(feature = "viper"))] + ParameterVariant::TestRunReport(_) => { + // Ignore TestRunReport when VIPER is disabled } }; } diff --git a/opendut-edgar/src/service/viper_run_manager/mod.rs b/opendut-edgar/src/service/viper_run_manager/mod.rs new file mode 100644 index 000000000..2f4983f23 --- /dev/null +++ b/opendut-edgar/src/service/viper_run_manager/mod.rs @@ -0,0 +1,84 @@ +use std::collections::HashMap; +use std::fmt::{Debug, Formatter}; +use std::sync::Arc; +use thiserror::Error; +use tokio::sync::RwLock; +use tokio::task::JoinHandle; +use tracing::warn; +use opendut_model::viper::{TestRunSourceCode, ViperRunId}; +use opendut_viper_rt::ViperRuntime; +use opendut_viper_rt::compile::{CompilationError, IdentifierFilter}; +use opendut_viper_rt::events::emitter; +use opendut_viper_rt::run::{ParameterBindings, RunError, TestSuiteReport}; +use opendut_viper_rt::source::Source; + +pub type ViperRunManagerRef = Arc; + +pub struct ViperRunManager { + test_runs: RwLock>>>, +} + +impl Debug for ViperRunManager { + fn fmt(&self, f: &mut Formatter<'_>) -> std::fmt::Result { + f.debug_struct("ViperRunManager") + .finish() + } +} + +impl ViperRunManager { + pub fn create() -> ViperRunManagerRef { + Arc::new(Self { + test_runs: RwLock::new(HashMap::new()), + }) + } + + pub async fn start_test_run(&self, run_id: ViperRunId, source_code: TestRunSourceCode) { + let handle = tokio::task::spawn_blocking(move || { + tokio::runtime::Handle::current().block_on(async move { + #[allow(clippy::needless_update)] + let viper_runtime = ViperRuntime::default(); + + let source = Source::embedded(source_code.inner.code); + + let compilation = viper_runtime.compile(&source, &mut emitter::drain(), &IdentifierFilter::default()).await + .map_err(StartTestRunError::Compilation)?; + + let suite = compilation.into_suite(); + let bindings = ParameterBindings::new(); // todo: include parameter bindings from CARL + + let report = viper_runtime.run(suite, bindings, &mut emitter::drain()).await + .map_err(StartTestRunError::Run)?; + + Ok(report) + }) + }); + + let mut test_runs = self.test_runs.write().await; + test_runs.insert(run_id, handle); + } + + pub async fn contains_test_run(&self, run_id: &ViperRunId) -> bool { + let test_runs = self.test_runs.read().await; + test_runs.contains_key(run_id) + } + + pub async fn abort_test_run(&self, run_id: &ViperRunId) { + let mut test_runs = self.test_runs.write().await; + let removed_handle = test_runs.remove(run_id); + match removed_handle { + Some(handle) => handle.abort(), + None => { + warn!("Thread handle in ViperRunManager for test run {run_id} not found"); + } + } + } +} + +#[derive(Debug, Error)] +enum StartTestRunError { + #[error(transparent)] + Compilation(Box), + + #[error(transparent)] + Run(Box) +} diff --git a/opendut-viper/viper-rt/src/runtime/types/naming/mod.rs b/opendut-viper/viper-rt/src/runtime/types/naming/mod.rs index fe6a7fe47..566b1c83f 100644 --- a/opendut-viper/viper-rt/src/runtime/types/naming/mod.rs +++ b/opendut-viper/viper-rt/src/runtime/types/naming/mod.rs @@ -19,7 +19,7 @@ use crate::runtime::types::naming::error::InvalidIdentifierError; const SEPARATOR: &str = "::"; -pub trait Identifier : Debug + Display { +pub trait Identifier : Debug + Display + Send + Sync { /// Extracts a string slice containing the entire identifier. fn as_str(&self) -> &str;