Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

5 changes: 4 additions & 1 deletion opendut-edgar/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -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 }
Expand Down Expand Up @@ -74,7 +75,9 @@ shadow-rs = { workspace = true, default-features = true }

[features]
integration_testing = []
viper = []
viper = [
"opendut-viper-rt/run",
]

[lints]
workspace = true
3 changes: 3 additions & 0 deletions opendut-edgar/src/service/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down
10 changes: 8 additions & 2 deletions opendut-edgar/src/service/peer_configuration.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand Down Expand Up @@ -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 {
Expand All @@ -127,4 +134,3 @@ async fn apply_peer_configuration(params: ApplyPeerConfigurationParams) -> Colle
debug!("Peer configuration has been successfully applied.");
result
}

9 changes: 8 additions & 1 deletion opendut-edgar/src/service/peer_messaging_client.rs
Original file line number Diff line number Diff line change
Expand Up @@ -26,14 +26,15 @@ 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,
executor_manager: ExecutorManagerRef,
metrics_manager: NetworkMetricsManagerRef,
carl_disconnect_timeout: Duration,
tx_peer_configuration: mpsc::Sender<ApplyPeerConfigurationParams>,
#[cfg(feature = "viper")]
viper_run_manager: crate::service::viper_run_manager::ViperRunManagerRef,
}


Expand Down Expand Up @@ -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,
Expand All @@ -71,6 +74,8 @@ impl PeerMessagingClient {
metrics_manager,
carl_disconnect_timeout,
tx_peer_configuration,
#[cfg(feature = "viper")]
viper_run_manager,
})
}

Expand Down Expand Up @@ -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?;

Expand Down
9 changes: 7 additions & 2 deletions opendut-edgar/src/service/tasks/create_ethernet_bridge.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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;

Expand Down Expand Up @@ -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 {
Expand Down
54 changes: 54 additions & 0 deletions opendut-edgar/src/service/tasks/create_viper_runtime.rs
Original file line number Diff line number Diff line change
@@ -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<TaskStateFulfilled> {
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<Success> {
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<TaskStateFulfilled> {
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<Success> {
self.viper_run_manager.abort_test_run(
&self.parameter.run_id,
).await;

Ok(Success::default())
}
}
3 changes: 3 additions & 0 deletions opendut-edgar/src/service/tasks/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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;
18 changes: 16 additions & 2 deletions opendut-edgar/src/service/tasks/runner/task_resolver.rs
Original file line number Diff line number Diff line change
Expand Up @@ -14,18 +14,24 @@ 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 {
pub fn new(
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,
}
}
}
Expand Down Expand Up @@ -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
}
};
}
Expand Down
84 changes: 84 additions & 0 deletions opendut-edgar/src/service/viper_run_manager/mod.rs
Original file line number Diff line number Diff line change
@@ -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<ViperRunManager>;

pub struct ViperRunManager {
test_runs: RwLock<HashMap<ViperRunId, JoinHandle<Result<TestSuiteReport, StartTestRunError>>>>,
}

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<CompilationError>),

#[error(transparent)]
Run(Box<RunError>)
}
2 changes: 1 addition & 1 deletion opendut-viper/viper-rt/src/runtime/types/naming/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down
Loading