refactor(tf): clean up example integration seams

This commit is contained in:
andrussal 2026-04-10 11:06:03 +02:00
parent e04b08c004
commit 36d7f3a0cf
18 changed files with 332 additions and 134 deletions

View File

@ -6,8 +6,7 @@ exclude-dev = true
no-default-features = true no-default-features = true
[advisories] [advisories]
ignore = [ ignore = []
]
yanked = "deny" yanked = "deny"
[bans] [bans]

View File

@ -72,6 +72,17 @@ pub trait StaticNodeConfigProvider: Application + Sized {
} }
} }
pub fn serialize_yaml_config<T>(config: &T) -> Result<String, serde_yaml::Error>
where
T: Serialize,
{
serde_yaml::to_string(config)
}
pub fn serialize_plain_text_config(config: &str) -> Result<String, std::convert::Infallible> {
Ok(config.to_owned())
}
impl<T> StaticArtifactRenderer for T impl<T> StaticArtifactRenderer for T
where where
T: StaticNodeConfigProvider, T: StaticNodeConfigProvider,

View File

@ -1,6 +1,7 @@
use std::{collections::HashMap, error::Error, io}; use std::{collections::HashMap, error::Error, io};
use cfgsync_artifacts::{ArtifactFile, ArtifactSet}; use cfgsync_artifacts::{ArtifactFile, ArtifactSet};
use serde::Serialize;
use crate::{ use crate::{
cfgsync::StaticNodeConfigProvider, cfgsync::StaticNodeConfigProvider,
@ -119,6 +120,13 @@ pub trait ClusterNodeConfigApplication: Application {
) -> Result<String, Self::ConfigError>; ) -> Result<String, Self::ConfigError>;
} }
pub fn serialize_cluster_yaml_config<T>(config: &T) -> Result<String, serde_yaml::Error>
where
T: Serialize,
{
serde_yaml::to_string(config)
}
impl<T> StaticNodeConfigProvider for T impl<T> StaticNodeConfigProvider for T
where where
T: ClusterNodeConfigApplication, T: ClusterNodeConfigApplication,

View File

@ -28,7 +28,9 @@ pub use capabilities::{
}; };
pub use client::NodeAccess; pub use client::NodeAccess;
pub use common_builder_ext::CoreBuilderExt; pub use common_builder_ext::CoreBuilderExt;
pub use config::{ClusterNodeConfigApplication, ClusterNodeView, ClusterPeerView}; pub use config::{
ClusterNodeConfigApplication, ClusterNodeView, ClusterPeerView, serialize_cluster_yaml_config,
};
pub use control::{ClusterWaitHandle, NodeControlHandle}; pub use control::{ClusterWaitHandle, NodeControlHandle};
pub use definition::{Scenario, ScenarioBuildError, ScenarioBuilder}; pub use definition::{Scenario, ScenarioBuildError, ScenarioBuilder};
pub use deployment_policy::{CleanupPolicy, DeploymentPolicy, RetryPolicy}; pub use deployment_policy::{CleanupPolicy, DeploymentPolicy, RetryPolicy};

View File

@ -3,7 +3,10 @@ services:
{{ node.name }}: {{ node.name }}:
image: {{ node.image }} image: {{ node.image }}
{% if node.platform %} platform: {{ node.platform }} {% if node.platform %} platform: {{ node.platform }}
{% endif %} entrypoint: {{ node.entrypoint }} {% endif %} entrypoint:
{% for arg in node.entrypoint %}
- {{ arg }}
{% endfor %}
volumes: volumes:
{% for volume in node.volumes %} {% for volume in node.volumes %}
- {{ volume }} - {{ volume }}

View File

@ -20,7 +20,9 @@ impl<E: ComposeDeployEnv> PortManager<E> {
environment: &mut StackEnvironment, environment: &mut StackEnvironment,
descriptors: &E::Deployment, descriptors: &E::Deployment,
) -> Result<HostPortMapping, ComposeRunnerError> { ) -> Result<HostPortMapping, ComposeRunnerError> {
let nodes = E::node_container_ports(descriptors); let nodes = E::node_container_ports(descriptors).map_err(|source| {
ComposeRunnerError::Config(crate::errors::ConfigError::Descriptor { source })
})?;
debug!( debug!(
nodes = nodes.len(), nodes = nodes.len(),
"resolving host ports for compose services" "resolving host ports for compose services"

View File

@ -9,7 +9,7 @@ use crate::infrastructure::ports::node_identifier;
pub struct NodeDescriptor { pub struct NodeDescriptor {
name: String, name: String,
image: String, image: String,
entrypoint: String, entrypoint: Vec<String>,
volumes: Vec<String>, volumes: Vec<String>,
extra_hosts: Vec<String>, extra_hosts: Vec<String>,
ports: Vec<String>, ports: Vec<String>,
@ -51,7 +51,7 @@ impl NodeDescriptor {
pub fn new( pub fn new(
name: impl Into<String>, name: impl Into<String>,
image: impl Into<String>, image: impl Into<String>,
entrypoint: impl Into<String>, entrypoint: Vec<String>,
volumes: Vec<String>, volumes: Vec<String>,
extra_hosts: Vec<String>, extra_hosts: Vec<String>,
ports: Vec<String>, ports: Vec<String>,
@ -62,7 +62,7 @@ impl NodeDescriptor {
Self { Self {
name: name.into(), name: name.into(),
image: image.into(), image: image.into(),
entrypoint: entrypoint.into(), entrypoint,
volumes, volumes,
extra_hosts, extra_hosts,
ports, ports,
@ -76,7 +76,7 @@ impl NodeDescriptor {
pub fn with_loopback_ports( pub fn with_loopback_ports(
name: impl Into<String>, name: impl Into<String>,
image: impl Into<String>, image: impl Into<String>,
entrypoint: impl Into<String>, entrypoint: Vec<String>,
volumes: Vec<String>, volumes: Vec<String>,
extra_hosts: Vec<String>, extra_hosts: Vec<String>,
container_ports: Vec<u16>, container_ports: Vec<u16>,
@ -135,7 +135,7 @@ impl NodeDescriptor {
#[derive(Clone, Debug)] #[derive(Clone, Debug)]
pub struct LoopbackNodeRuntimeSpec { pub struct LoopbackNodeRuntimeSpec {
pub image: String, pub image: String,
pub entrypoint: String, pub entrypoint: Vec<String>,
pub volumes: Vec<String>, pub volumes: Vec<String>,
pub extra_hosts: Vec<String>, pub extra_hosts: Vec<String>,
pub container_ports: Vec<u16>, pub container_ports: Vec<u16>,
@ -229,10 +229,11 @@ pub fn binary_config_node_runtime_spec(
) -> LoopbackNodeRuntimeSpec { ) -> LoopbackNodeRuntimeSpec {
let image = env::var(&spec.image_env_var).unwrap_or_else(|_| spec.default_image.clone()); let image = env::var(&spec.image_env_var).unwrap_or_else(|_| spec.default_image.clone());
let platform = env::var(&spec.platform_env_var).ok(); let platform = env::var(&spec.platform_env_var).ok();
let entrypoint = format!( let entrypoint = vec![
"/bin/sh -c '{} --config {}'", spec.binary_path.clone(),
spec.binary_path, spec.config_container_path "--config".to_owned(),
); spec.config_container_path.clone(),
];
LoopbackNodeRuntimeSpec { LoopbackNodeRuntimeSpec {
image, image,

View File

@ -10,11 +10,15 @@ use testing_framework_core::{
}, },
topology::DeploymentDescriptor, topology::DeploymentDescriptor,
}; };
use tokio::{
net::TcpStream,
time::{Instant, sleep},
};
use crate::{ use crate::{
descriptor::{ descriptor::{
BinaryConfigNodeSpec, ComposeDescriptor, LoopbackNodeRuntimeSpec, NodeDescriptor, BinaryConfigNodeSpec, ComposeDescriptor, LoopbackNodeRuntimeSpec, NodeDescriptor,
binary_config_node_runtime_spec, build_loopback_node_descriptors, binary_config_node_runtime_spec,
}, },
docker::config_server::DockerConfigServerSpec, docker::config_server::DockerConfigServerSpec,
infrastructure::ports::{ infrastructure::ports::{
@ -31,6 +35,12 @@ pub trait ConfigServerHandle: Send + Sync {
} }
} }
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub enum ComposeConfigServerMode {
Disabled,
Docker,
}
/// Compose-specific topology surface needed by the runner. /// Compose-specific topology surface needed by the runner.
#[async_trait] #[async_trait]
pub trait ComposeDeployEnv: Application { pub trait ComposeDeployEnv: Application {
@ -70,27 +80,39 @@ pub trait ComposeDeployEnv: Application {
fn compose_descriptor( fn compose_descriptor(
topology: &<Self as Application>::Deployment, topology: &<Self as Application>::Deployment,
_cfgsync_port: u16, _cfgsync_port: u16,
) -> ComposeDescriptor { ) -> Result<ComposeDescriptor, DynError> {
let nodes = build_loopback_node_descriptors(topology.node_count(), |index| { let mut nodes = Vec::with_capacity(topology.node_count());
Self::loopback_node_runtime_spec(topology, index) for index in 0..topology.node_count() {
.unwrap_or_else(|| panic!("compose_descriptor is not implemented for this app")) let spec = Self::loopback_node_runtime_spec(topology, index).ok_or_else(|| {
}); std::io::Error::other("compose_descriptor is not implemented for this app")
ComposeDescriptor::new(nodes) })?;
nodes.push(NodeDescriptor::with_loopback_ports(
crate::infrastructure::ports::node_identifier(index),
spec.image,
spec.entrypoint,
spec.volumes,
spec.extra_hosts,
spec.container_ports,
spec.environment,
spec.platform,
));
}
Ok(ComposeDescriptor::new(nodes))
} }
/// Container ports (API/testing) per node, used for docker-compose port /// Container ports (API/testing) per node, used for docker-compose port
/// discovery. /// discovery.
fn node_container_ports( fn node_container_ports(
topology: &<Self as Application>::Deployment, topology: &<Self as Application>::Deployment,
) -> Vec<NodeContainerPorts> { ) -> Result<Vec<NodeContainerPorts>, DynError> {
let descriptor = Self::compose_descriptor(topology, 0); let descriptor = Self::compose_descriptor(topology, 0)?;
descriptor Ok(descriptor
.nodes() .nodes()
.iter() .iter()
.enumerate() .enumerate()
.take(topology.node_count()) .take(topology.node_count())
.filter_map(|(index, node)| parse_node_container_ports(index, node)) .filter_map(|(index, node)| parse_node_container_ports(index, node))
.collect() .collect())
} }
/// Hostnames used when rewriting node configs for cfgsync delivery. /// Hostnames used when rewriting node configs for cfgsync delivery.
@ -126,10 +148,10 @@ pub trait ComposeDeployEnv: Application {
/// Build the config server container specification. /// Build the config server container specification.
fn cfgsync_container_spec( fn cfgsync_container_spec(
_cfgsync_path: &Path, _cfgsync_path: &Path,
port: u16, _port: u16,
network: &str, _network: &str,
) -> Result<DockerConfigServerSpec, DynError> { ) -> Result<DockerConfigServerSpec, DynError> {
Ok(dummy_cfgsync_spec(port, network)) Err(std::io::Error::other("cfgsync_container_spec is not implemented for this app").into())
} }
/// Timeout used when launching the config server container. /// Timeout used when launching the config server container.
@ -137,6 +159,10 @@ pub trait ComposeDeployEnv: Application {
Duration::from_secs(180) Duration::from_secs(180)
} }
fn cfgsync_server_mode() -> ComposeConfigServerMode {
ComposeConfigServerMode::Disabled
}
/// Build node clients from discovered host ports. /// Build node clients from discovered host ports.
fn node_client_from_ports( fn node_client_from_ports(
ports: &NodeHostPorts, ports: &NodeHostPorts,
@ -167,6 +193,12 @@ pub trait ComposeDeployEnv: Application {
<Self as Application>::node_readiness_path() <Self as Application>::node_readiness_path()
} }
fn node_readiness_probe() -> ComposeReadinessProbe {
ComposeReadinessProbe::Http {
path: <Self as ComposeDeployEnv>::node_readiness_path(),
}
}
/// Host used by default remote readiness checks. /// Host used by default remote readiness checks.
fn compose_runner_host() -> String { fn compose_runner_host() -> String {
compose_runner_host() compose_runner_host()
@ -178,14 +210,15 @@ pub trait ComposeDeployEnv: Application {
mapping: &HostPortMapping, mapping: &HostPortMapping,
requirement: HttpReadinessRequirement, requirement: HttpReadinessRequirement,
) -> Result<(), DynError> { ) -> Result<(), DynError> {
let host = Self::compose_runner_host(); match <Self as ComposeDeployEnv>::node_readiness_probe() {
let urls = readiness_urls( ComposeReadinessProbe::Http { path } => {
&host, let host = Self::compose_runner_host();
mapping, let urls = readiness_urls(&host, mapping, path)?;
<Self as ComposeDeployEnv>::node_readiness_path(), wait_http_readiness(&urls, requirement).await?;
)?; Ok(())
wait_http_readiness(&urls, requirement).await?; }
Ok(()) ComposeReadinessProbe::Tcp => wait_for_tcp_readiness(&mapping.nodes, requirement).await,
}
} }
/// Wait for HTTP readiness on node ports. /// Wait for HTTP readiness on node ports.
@ -194,14 +227,24 @@ pub trait ComposeDeployEnv: Application {
host: &str, host: &str,
requirement: HttpReadinessRequirement, requirement: HttpReadinessRequirement,
) -> Result<(), DynError> { ) -> Result<(), DynError> {
wait_for_http_ports_with_host_and_requirement( match <Self as ComposeDeployEnv>::node_readiness_probe() {
ports, ComposeReadinessProbe::Http { path } => {
host, wait_for_http_ports_with_host_and_requirement(ports, host, path, requirement)
<Self as ComposeDeployEnv>::node_readiness_path(), .await?;
requirement, Ok(())
) }
.await?; ComposeReadinessProbe::Tcp => {
Ok(()) let ports = ports
.iter()
.copied()
.map(|port| NodeHostPorts {
api: port,
testing: port,
})
.collect::<Vec<_>>();
wait_for_tcp_readiness(&ports, requirement).await
}
}
} }
} }
@ -257,27 +300,10 @@ fn write_dummy_cfgsync_config(path: &Path, port: u16) -> Result<(), DynError> {
Ok(()) Ok(())
} }
fn dummy_cfgsync_spec(port: u16, network: &str) -> DockerConfigServerSpec {
use crate::docker::config_server::DockerPortBinding;
DockerConfigServerSpec::new(
"cfgsync".to_owned(),
network.to_owned(),
"sh".to_owned(),
"busybox:1.36".to_owned(),
)
.with_network_alias("cfgsync".to_owned())
.with_args(vec![
"-c".to_owned(),
format!("while true; do nc -l -p {port} >/dev/null 2>&1; done"),
])
.with_ports(vec![DockerPortBinding::tcp(port, port)])
}
fn parse_node_container_ports(index: usize, node: &NodeDescriptor) -> Option<NodeContainerPorts> { fn parse_node_container_ports(index: usize, node: &NodeDescriptor) -> Option<NodeContainerPorts> {
let mut ports = node.container_ports().iter().copied(); let mut ports = node.container_ports().iter().copied();
let api = ports.next()?; let api = ports.next()?;
let testing = ports.next()?; let testing = ports.next().unwrap_or(api);
Some(NodeContainerPorts { Some(NodeContainerPorts {
index, index,
@ -286,6 +312,52 @@ fn parse_node_container_ports(index: usize, node: &NodeDescriptor) -> Option<Nod
}) })
} }
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub enum ComposeReadinessProbe {
Http { path: &'static str },
Tcp,
}
async fn wait_for_tcp_readiness(
ports: &[NodeHostPorts],
requirement: HttpReadinessRequirement,
) -> Result<(), DynError> {
let timeout = Duration::from_secs(60);
let deadline = Instant::now() + timeout;
loop {
let mut ready = 0;
for node in ports {
if TcpStream::connect(("127.0.0.1", node.testing))
.await
.is_ok()
{
ready += 1;
}
}
let total = ports.len();
let satisfied = match requirement {
HttpReadinessRequirement::AllNodesReady => ready == total,
HttpReadinessRequirement::AnyNodeReady => ready >= 1,
HttpReadinessRequirement::AtLeast(min_ready) => ready >= min_ready,
};
if satisfied {
return Ok(());
}
if Instant::now() >= deadline {
return Err(format!(
"tcp readiness timed out: ready={ready}, total={total}, requirement={requirement:?}"
)
.into());
}
sleep(Duration::from_millis(200)).await;
}
}
pub fn discovered_node_access(host: &str, ports: &NodeHostPorts) -> NodeAccess { pub fn discovered_node_access(host: &str, ports: &NodeHostPorts) -> NodeAccess {
NodeAccess::new(host, ports.api).with_testing_port(ports.testing) NodeAccess::new(host, ports.api).with_testing_port(ports.testing)
} }

View File

@ -73,6 +73,11 @@ impl WorkspaceError {
#[derive(Debug, thiserror::Error)] #[derive(Debug, thiserror::Error)]
/// Configuration-related failures while preparing compose runs. /// Configuration-related failures while preparing compose runs.
pub enum ConfigError { pub enum ConfigError {
#[error("failed to build compose descriptor: {source}")]
Descriptor {
#[source]
source: DynError,
},
#[error("failed to update cfgsync configuration at {path:?}: {source}")] #[error("failed to update cfgsync configuration at {path:?}: {source}")]
Cfgsync { Cfgsync {
path: PathBuf, path: PathBuf,

View File

@ -25,7 +25,7 @@ use crate::{
ensure_image_present, ensure_image_present,
workspace::ComposeWorkspace, workspace::ComposeWorkspace,
}, },
env::{ComposeCfgsyncEnv, ComposeDeployEnv, ConfigServerHandle}, env::{ComposeCfgsyncEnv, ComposeConfigServerMode, ComposeDeployEnv, ConfigServerHandle},
errors::{ComposeRunnerError, ConfigError, WorkspaceError}, errors::{ComposeRunnerError, ConfigError, WorkspaceError},
infrastructure::template::write_compose_file, infrastructure::template::write_compose_file,
lifecycle::cleanup::RunnerCleanup, lifecycle::cleanup::RunnerCleanup,
@ -216,7 +216,11 @@ pub async fn start_cfgsync_stage<E: ComposeDeployEnv>(
workspace: &WorkspaceState, workspace: &WorkspaceState,
cfgsync_port: u16, cfgsync_port: u16,
project_name: &str, project_name: &str,
) -> Result<Box<dyn ConfigServerHandle>, ComposeRunnerError> { ) -> Result<Option<Box<dyn ConfigServerHandle>>, ComposeRunnerError> {
if matches!(E::cfgsync_server_mode(), ComposeConfigServerMode::Disabled) {
return Ok(None);
}
info!(cfgsync_port = cfgsync_port, "launching cfgsync server"); info!(cfgsync_port = cfgsync_port, "launching cfgsync server");
let network = compose_network_name(project_name); let network = compose_network_name(project_name);
@ -244,7 +248,7 @@ pub async fn start_cfgsync_stage<E: ComposeDeployEnv>(
wait_for_cfgsync_ready(cfgsync_port, Some(&handle)).await?; wait_for_cfgsync_ready(cfgsync_port, Some(&handle)).await?;
log_cfgsync_started(&handle); log_cfgsync_started(&handle);
Ok(Box::new(handle)) Ok(Some(Box::new(handle)))
} }
/// Write cfgsync YAML from topology data. /// Write cfgsync YAML from topology data.
@ -299,7 +303,8 @@ pub fn write_compose_artifacts<E: ComposeDeployEnv>(
workspace_root = %workspace.root.display(), workspace_root = %workspace.root.display(),
"building compose descriptor" "building compose descriptor"
); );
let descriptor = E::compose_descriptor(descriptors, cfgsync_port); let descriptor = E::compose_descriptor(descriptors, cfgsync_port)
.map_err(|source| ConfigError::Descriptor { source })?;
let compose_path = workspace.root.join("compose.generated.yml"); let compose_path = workspace.root.join("compose.generated.yml");
write_compose_file(&descriptor, &compose_path) write_compose_file(&descriptor, &compose_path)
@ -326,10 +331,12 @@ pub async fn bring_up_stack(
compose_path: &Path, compose_path: &Path,
project_name: &str, project_name: &str,
workspace_root: &Path, workspace_root: &Path,
cfgsync_handle: &mut dyn ConfigServerHandle, cfgsync_handle: &mut Option<Box<dyn ConfigServerHandle>>,
) -> Result<(), ComposeRunnerError> { ) -> Result<(), ComposeRunnerError> {
if let Err(err) = compose_up(compose_path, project_name, workspace_root).await { if let Err(err) = compose_up(compose_path, project_name, workspace_root).await {
cfgsync_handle.shutdown(); if let Some(cfgsync_handle) = cfgsync_handle.as_deref_mut() {
cfgsync_handle.shutdown();
}
return Err(ComposeRunnerError::Compose(err)); return Err(ComposeRunnerError::Compose(err));
} }
debug!(project = %project_name, "docker compose up completed"); debug!(project = %project_name, "docker compose up completed");
@ -341,7 +348,7 @@ pub async fn bring_up_stack_logged(
compose_path: &Path, compose_path: &Path,
project_name: &str, project_name: &str,
workspace_root: &Path, workspace_root: &Path,
cfgsync_handle: &mut dyn ConfigServerHandle, cfgsync_handle: &mut Option<Box<dyn ConfigServerHandle>>,
) -> Result<(), ComposeRunnerError> { ) -> Result<(), ComposeRunnerError> {
info!(project = %project_name, "bringing up docker compose stack"); info!(project = %project_name, "bringing up docker compose stack");
bring_up_stack(compose_path, project_name, workspace_root, cfgsync_handle).await bring_up_stack(compose_path, project_name, workspace_root, cfgsync_handle).await
@ -357,13 +364,10 @@ where
{ {
let prepared = prepare_stack_artifacts::<E>(descriptors, metrics_otlp_ingest_url).await?; let prepared = prepare_stack_artifacts::<E>(descriptors, metrics_otlp_ingest_url).await?;
let mut cfgsync_handle = start_cfgsync_for_prepared::<E>(&prepared).await?; let mut cfgsync_handle = start_cfgsync_for_prepared::<E>(&prepared).await?;
start_compose_stack(&prepared, cfgsync_handle.as_mut()).await?; start_compose_stack(&prepared, &mut cfgsync_handle).await?;
log_compose_environment_ready(&prepared, "compose stack is up"); log_compose_environment_ready(&prepared, "compose stack is up");
Ok(stack_environment_from_prepared( Ok(stack_environment_from_prepared(prepared, cfgsync_handle))
prepared,
Some(cfgsync_handle),
))
} }
/// Prepare workspace, cfgsync, and compose artifacts without starting services. /// Prepare workspace, cfgsync, and compose artifacts without starting services.
@ -379,10 +383,7 @@ where
log_compose_environment_ready(&prepared, "compose manual environment prepared"); log_compose_environment_ready(&prepared, "compose manual environment prepared");
Ok(stack_environment_from_prepared( Ok(stack_environment_from_prepared(prepared, cfgsync_handle))
prepared,
Some(cfgsync_handle),
))
} }
async fn prepare_stack_artifacts<E>( async fn prepare_stack_artifacts<E>(
@ -418,25 +419,28 @@ async fn ensure_compose_images_present<E: ComposeDeployEnv>(
descriptors: &E::Deployment, descriptors: &E::Deployment,
cfgsync_port: u16, cfgsync_port: u16,
) -> Result<(), ComposeRunnerError> { ) -> Result<(), ComposeRunnerError> {
let descriptor = E::compose_descriptor(descriptors, 0); let descriptor = E::compose_descriptor(descriptors, 0)
.map_err(|source| ComposeRunnerError::Config(ConfigError::Descriptor { source }))?;
let mut images = descriptor let mut images = descriptor
.nodes() .nodes()
.iter() .iter()
.map(|node| (node.image().to_owned(), node.platform().map(str::to_owned))) .map(|node| (node.image().to_owned(), node.platform().map(str::to_owned)))
.collect::<BTreeSet<_>>(); .collect::<BTreeSet<_>>();
let cfgsync_spec = E::cfgsync_container_spec( if matches!(E::cfgsync_server_mode(), ComposeConfigServerMode::Docker) {
&workspace.cfgsync_path, let cfgsync_spec = E::cfgsync_container_spec(
cfgsync_port, &workspace.cfgsync_path,
&compose_network_name("compose-image-check"), cfgsync_port,
) &compose_network_name("compose-image-check"),
.map_err(|source| { )
ComposeRunnerError::Config(ConfigError::CfgsyncStart { .map_err(|source| {
port: cfgsync_port, ComposeRunnerError::Config(ConfigError::CfgsyncStart {
source, port: cfgsync_port,
}) source,
})?; })
})?;
images.insert((cfgsync_spec.image, cfgsync_spec.platform)); images.insert((cfgsync_spec.image, cfgsync_spec.platform));
}
for (image, platform) in images { for (image, platform) in images {
ensure_image_present(&image, platform.as_deref()).await?; ensure_image_present(&image, platform.as_deref()).await?;
@ -451,7 +455,7 @@ fn create_project_name() -> String {
async fn start_cfgsync_for_prepared<E: ComposeDeployEnv>( async fn start_cfgsync_for_prepared<E: ComposeDeployEnv>(
prepared: &PreparedEnvironment, prepared: &PreparedEnvironment,
) -> Result<Box<dyn ConfigServerHandle>, ComposeRunnerError> { ) -> Result<Option<Box<dyn ConfigServerHandle>>, ComposeRunnerError> {
start_cfgsync_stage::<E>( start_cfgsync_stage::<E>(
&prepared.workspace, &prepared.workspace,
prepared.cfgsync_port, prepared.cfgsync_port,
@ -462,7 +466,7 @@ async fn start_cfgsync_for_prepared<E: ComposeDeployEnv>(
async fn handle_compose_start_failure( async fn handle_compose_start_failure(
prepared: &PreparedEnvironment, prepared: &PreparedEnvironment,
cfgsync_handle: &mut dyn ConfigServerHandle, cfgsync_handle: &mut Option<Box<dyn ConfigServerHandle>>,
) { ) {
dump_compose_logs( dump_compose_logs(
&prepared.compose_path, &prepared.compose_path,
@ -470,7 +474,9 @@ async fn handle_compose_start_failure(
&prepared.workspace.root, &prepared.workspace.root,
) )
.await; .await;
cfgsync_handle.shutdown(); if let Some(cfgsync_handle) = cfgsync_handle.as_deref_mut() {
cfgsync_handle.shutdown();
}
} }
fn stack_environment_from_prepared( fn stack_environment_from_prepared(
@ -487,7 +493,7 @@ fn stack_environment_from_prepared(
async fn start_compose_stack( async fn start_compose_stack(
prepared: &PreparedEnvironment, prepared: &PreparedEnvironment,
cfgsync_handle: &mut dyn ConfigServerHandle, cfgsync_handle: &mut Option<Box<dyn ConfigServerHandle>>,
) -> Result<(), ComposeRunnerError> { ) -> Result<(), ComposeRunnerError> {
if let Err(error) = bring_up_stack_logged( if let Err(error) = bring_up_stack_logged(
&prepared.compose_path, &prepared.compose_path,

View File

@ -20,7 +20,10 @@ pub use docker::{
}, },
platform::host_gateway_entry, platform::host_gateway_entry,
}; };
pub use env::{ComposeDeployEnv, ConfigServerHandle, discovered_node_access}; pub use env::{
ComposeConfigServerMode, ComposeDeployEnv, ComposeReadinessProbe, ConfigServerHandle,
discovered_node_access,
};
pub use errors::ComposeRunnerError; pub use errors::ComposeRunnerError;
pub use infrastructure::{ pub use infrastructure::{
ports::{HostPortMapping, NodeHostPorts, compose_runner_host, node_identifier}, ports::{HostPortMapping, NodeHostPorts, compose_runner_host, node_identifier},

View File

@ -9,6 +9,7 @@ use async_trait::async_trait;
use cfgsync_artifacts::ArtifactSet; use cfgsync_artifacts::ArtifactSet;
use kube::Client; use kube::Client;
use reqwest::Url; use reqwest::Url;
use serde::Serialize;
use tempfile::TempDir; use tempfile::TempDir;
use testing_framework_core::{ use testing_framework_core::{
cfgsync::StaticNodeConfigProvider, cfgsync::StaticNodeConfigProvider,
@ -35,6 +36,11 @@ pub struct RenderedHelmChartAssets {
_tempdir: TempDir, _tempdir: TempDir,
} }
#[derive(Clone, Debug, Default)]
pub struct HelmManifest {
documents: Vec<String>,
}
#[derive(Clone, Debug)] #[derive(Clone, Debug)]
pub struct BinaryConfigK8sSpec { pub struct BinaryConfigK8sSpec {
pub chart_name: String, pub chart_name: String,
@ -55,6 +61,38 @@ impl HelmReleaseAssets for RenderedHelmChartAssets {
} }
} }
impl HelmManifest {
#[must_use]
pub fn new() -> Self {
Self::default()
}
pub fn push_yaml<T>(&mut self, value: &T) -> Result<(), DynError>
where
T: Serialize,
{
self.documents
.push(normalize_yaml_document(&serde_yaml::to_string(value)?));
Ok(())
}
pub fn push_raw_yaml(&mut self, yaml: &str) {
let yaml = yaml.trim();
if !yaml.is_empty() {
self.documents.push(yaml.to_owned());
}
}
pub fn extend(&mut self, other: Self) {
self.documents.extend(other.documents);
}
#[must_use]
pub fn render(&self) -> String {
self.documents.join("\n---\n")
}
}
pub fn standard_port_specs(node_count: usize, api: u16, auxiliary: u16) -> PortSpecs { pub fn standard_port_specs(node_count: usize, api: u16, auxiliary: u16) -> PortSpecs {
PortSpecs { PortSpecs {
nodes: (0..node_count) nodes: (0..node_count)
@ -201,6 +239,14 @@ pub fn render_single_template_chart_assets(
}) })
} }
pub fn render_manifest_chart_assets(
chart_name: &str,
template_name: &str,
manifest: &HelmManifest,
) -> Result<RenderedHelmChartAssets, DynError> {
render_single_template_chart_assets(chart_name, template_name, &manifest.render())
}
pub fn discovered_node_access(host: &str, api_port: u16, auxiliary_port: u16) -> NodeAccess { pub fn discovered_node_access(host: &str, api_port: u16, auxiliary_port: u16) -> NodeAccess {
NodeAccess::new(host, api_port).with_testing_port(auxiliary_port) NodeAccess::new(host, api_port).with_testing_port(auxiliary_port)
} }
@ -209,6 +255,10 @@ fn render_chart_yaml(chart_name: &str) -> String {
format!("apiVersion: v2\nname: {chart_name}\nversion: 0.1.0\n") format!("apiVersion: v2\nname: {chart_name}\nversion: 0.1.0\n")
} }
fn normalize_yaml_document(yaml: &str) -> String {
yaml.trim_start_matches("---\n").trim().to_owned()
}
pub async fn install_helm_release_with_cleanup<A: HelmReleaseAssets>( pub async fn install_helm_release_with_cleanup<A: HelmReleaseAssets>(
client: &Client, client: &Client,
assets: &A, assets: &A,

View File

@ -7,6 +7,8 @@ mod manual;
mod workspace; mod workspace;
use std::sync::Once; use std::sync::Once;
pub use k8s_openapi;
pub mod wait { pub mod wait {
pub use crate::lifecycle::wait::*; pub use crate::lifecycle::wait::*;
} }
@ -21,10 +23,10 @@ pub(crate) fn ensure_rustls_provider_installed() {
pub use deployer::{K8sDeployer, K8sDeploymentMetadata, K8sRunnerError}; pub use deployer::{K8sDeployer, K8sDeploymentMetadata, K8sRunnerError};
pub use env::{ pub use env::{
BinaryConfigK8sSpec, HelmReleaseAssets, K8sDeployEnv, RenderedHelmChartAssets, BinaryConfigK8sSpec, HelmManifest, HelmReleaseAssets, K8sDeployEnv, RenderedHelmChartAssets,
discovered_node_access, install_helm_release_with_cleanup, discovered_node_access, install_helm_release_with_cleanup,
render_binary_config_node_chart_assets, render_binary_config_node_manifest, render_binary_config_node_chart_assets, render_binary_config_node_manifest,
render_single_template_chart_assets, standard_port_specs, render_manifest_chart_assets, render_single_template_chart_assets, standard_port_specs,
}; };
pub use infrastructure::{ pub use infrastructure::{
chart_values::{ chart_values::{

View File

@ -1,5 +1,6 @@
use std::{ use std::{
fmt, io, fmt, io,
io::Read,
net::{Ipv4Addr, TcpListener, TcpStream}, net::{Ipv4Addr, TcpListener, TcpStream},
process::{Child, Command as StdCommand, ExitStatus, Stdio}, process::{Child, Command as StdCommand, ExitStatus, Stdio},
thread, thread,
@ -86,7 +87,7 @@ fn spawn_kubectl_port_forward(
.arg(format!("svc/{service}")) .arg(format!("svc/{service}"))
.arg(format!("{local_port}:{remote_port}")) .arg(format!("{local_port}:{remote_port}"))
.stdout(Stdio::null()) .stdout(Stdio::null())
.stderr(Stdio::null()) .stderr(Stdio::piped())
.spawn() .spawn()
} }
@ -107,7 +108,13 @@ fn wait_until_port_forward_ready(
} }
let _ = child.kill(); let _ = child.kill();
Err(port_forward_ready_timeout_error(service, remote_port)) let _ = child.wait();
let details = read_port_forward_stderr(child);
Err(port_forward_ready_timeout_error(
service,
remote_port,
details.as_deref(),
))
} }
fn ensure_port_forward_running( fn ensure_port_forward_running(
@ -122,7 +129,7 @@ fn ensure_port_forward_running(
Err(port_forward_error( Err(port_forward_error(
service, service,
remote_port, remote_port,
anyhow!("kubectl exited with {status}"), port_forward_process_error(status, read_port_forward_stderr(child)),
)) ))
} }
@ -134,11 +141,18 @@ fn port_forward_error(service: &str, remote_port: u16, source: anyhow::Error) ->
} }
} }
fn port_forward_ready_timeout_error(service: &str, remote_port: u16) -> ClusterWaitError { fn port_forward_ready_timeout_error(
service: &str,
remote_port: u16,
details: Option<&str>,
) -> ClusterWaitError {
port_forward_error( port_forward_error(
service, service,
remote_port, remote_port,
anyhow!("port-forward did not become ready"), anyhow!(
"port-forward did not become ready{}",
format_port_forward_details(details)
),
) )
} }
@ -146,6 +160,30 @@ fn exited_status(child: &mut Child) -> Option<ExitStatus> {
child.try_wait().ok().flatten() child.try_wait().ok().flatten()
} }
fn read_port_forward_stderr(child: &mut Child) -> Option<String> {
let mut stderr = child.stderr.take()?;
let mut output = String::new();
if stderr.read_to_string(&mut output).is_err() {
return None;
}
let trimmed = output.trim();
(!trimmed.is_empty()).then(|| trimmed.to_owned())
}
fn port_forward_process_error(status: ExitStatus, details: Option<String>) -> anyhow::Error {
anyhow!(
"kubectl exited with {status}{}",
format_port_forward_details(details.as_deref())
)
}
fn format_port_forward_details(details: Option<&str>) -> String {
match details {
Some(details) => format!(": {details}"),
None => String::new(),
}
}
fn local_port_reachable(local_port: u16) -> bool { fn local_port_reachable(local_port: u16) -> bool {
TcpStream::connect(localhost_addr(local_port)).is_ok() TcpStream::connect(localhost_addr(local_port)).is_ok()
} }

View File

@ -24,7 +24,6 @@ use crate::{
K8sDeployer, K8sDeployer,
env::{HelmReleaseAssets, K8sDeployEnv, discovered_node_access}, env::{HelmReleaseAssets, K8sDeployEnv, discovered_node_access},
host::node_host, host::node_host,
infrastructure::helm::install_release,
lifecycle::{ lifecycle::{
cleanup::RunnerCleanup, cleanup::RunnerCleanup,
wait::{ wait::{
@ -140,20 +139,9 @@ where
let assets = E::prepare_assets(&topology, None) let assets = E::prepare_assets(&topology, None)
.map_err(|source| ManualClusterError::Assets { source })?; .map_err(|source| ManualClusterError::Assets { source })?;
let (namespace, release) = E::cluster_identifiers(); let (namespace, release) = E::cluster_identifiers();
let install_spec = assets let cleanup = E::install_stack(&client, &assets, &namespace, &release, nodes)
.release_bundle() .await
.install_spec(release.clone(), namespace.clone()); .map_err(|source| ManualClusterError::InstallStack { source })?;
install_release(&install_spec).await.map_err(|source| {
ManualClusterError::InstallStack {
source: source.into(),
}
})?;
let cleanup = RunnerCleanup::new(
client.clone(),
namespace.clone(),
release.clone(),
std::env::var("K8S_RUNNER_PRESERVE").is_ok(),
);
let node_ports = E::collect_port_specs(&topology).nodes; let node_ports = E::collect_port_specs(&topology).nodes;
let node_allocations = let node_allocations =

View File

@ -511,15 +511,17 @@ where
None None
} }
fn node_endpoints(config: &<Self as Application>::NodeConfig) -> NodeEndpoints { fn node_endpoints(
config: &<Self as Application>::NodeConfig,
) -> Result<NodeEndpoints, DynError> {
if let Some(port) = Self::http_api_port(config) { if let Some(port) = Self::http_api_port(config) {
return NodeEndpoints { return Ok(NodeEndpoints {
api: SocketAddr::from((Ipv4Addr::LOCALHOST, port)), api: SocketAddr::from((Ipv4Addr::LOCALHOST, port)),
extra_ports: HashMap::new(), extra_ports: HashMap::new(),
}; });
} }
panic!("node_endpoints is not implemented for this app"); Err(std::io::Error::other("node_endpoints is not implemented for this app").into())
} }
fn node_peer_port(node: &Node<Self>) -> u16 { fn node_peer_port(node: &Node<Self>) -> u16 {
@ -530,18 +532,18 @@ where
None None
} }
fn node_client(endpoints: &NodeEndpoints) -> Self::NodeClient { fn node_client(endpoints: &NodeEndpoints) -> Result<Self::NodeClient, DynError> {
if let Ok(client) = if let Ok(client) =
<Self as Application>::build_node_client(&discovered_node_access(endpoints)) <Self as Application>::build_node_client(&discovered_node_access(endpoints))
{ {
return client; return Ok(client);
} }
if let Some(client) = Self::node_client_from_api_endpoint(endpoints.api) { if let Some(client) = Self::node_client_from_api_endpoint(endpoints.api) {
return client; return Ok(client);
} }
panic!("node_client is not implemented for this app"); Err(std::io::Error::other("node_client is not implemented for this app").into())
} }
fn readiness_endpoint_path() -> &'static str { fn readiness_endpoint_path() -> &'static str {
@ -680,11 +682,15 @@ mod tests {
Ok(LaunchSpec::default()) Ok(LaunchSpec::default())
} }
fn node_endpoints(_config: &<Self as Application>::NodeConfig) -> NodeEndpoints { fn node_endpoints(
NodeEndpoints::default() _config: &<Self as Application>::NodeConfig,
) -> Result<NodeEndpoints, DynError> {
Ok(NodeEndpoints::default())
} }
fn node_client(_endpoints: &NodeEndpoints) -> Self::NodeClient {} fn node_client(_endpoints: &NodeEndpoints) -> Result<Self::NodeClient, DynError> {
Ok(())
}
async fn wait_readiness_stable(_nodes: &[Node<Self>]) -> Result<(), DynError> { async fn wait_readiness_stable(_nodes: &[Node<Self>]) -> Result<(), DynError> {
STABLE_CALLS.fetch_add(1, Ordering::SeqCst); STABLE_CALLS.fetch_add(1, Ordering::SeqCst);

View File

@ -31,7 +31,7 @@ pub fn build_external_client<E: LocalDeployerEnv>(
let api = resolve_api_socket(source)?; let api = resolve_api_socket(source)?;
let mut endpoints = NodeEndpoints::default(); let mut endpoints = NodeEndpoints::default();
endpoints.api = api; endpoints.api = api;
Ok(E::node_client(&endpoints)) E::node_client(&endpoints)
} }
fn resolve_api_socket(source: &ExternalNodeSource) -> Result<std::net::SocketAddr, DynError> { fn resolve_api_socket(source: &ExternalNodeSource) -> Result<std::net::SocketAddr, DynError> {

View File

@ -203,11 +203,11 @@ impl<Config: Clone + Send + Sync + 'static, Client: Clone + Send + Sync + 'stati
label: &str, label: &str,
config: Config, config: Config,
build_launch_spec: impl FnOnce(&Config, &Path, &str) -> Result<LaunchSpec, DynError>, build_launch_spec: impl FnOnce(&Config, &Path, &str) -> Result<LaunchSpec, DynError>,
endpoints_from_config: impl FnOnce(&Config) -> NodeEndpoints, endpoints_from_config: impl FnOnce(&Config) -> Result<NodeEndpoints, DynError>,
keep_tempdir: bool, keep_tempdir: bool,
persist_dir: Option<&Path>, persist_dir: Option<&Path>,
snapshot_dir: Option<&Path>, snapshot_dir: Option<&Path>,
client_from_endpoints: impl FnOnce(&NodeEndpoints) -> Client, client_from_endpoints: impl FnOnce(&NodeEndpoints) -> Result<Client, DynError>,
) -> Result<Self, ProcessSpawnError> { ) -> Result<Self, ProcessSpawnError> {
let tempdir = create_tempdir(persist_dir)?; let tempdir = create_tempdir(persist_dir)?;
if let Some(snapshot_dir) = snapshot_dir { if let Some(snapshot_dir) = snapshot_dir {
@ -217,8 +217,10 @@ impl<Config: Clone + Send + Sync + 'static, Client: Clone + Send + Sync + 'stati
let launch = build_launch_spec(&config, tempdir.path(), label) let launch = build_launch_spec(&config, tempdir.path(), label)
.map_err(|source| ProcessSpawnError::Config { source })?; .map_err(|source| ProcessSpawnError::Config { source })?;
let endpoints = endpoints_from_config(&config); let endpoints = endpoints_from_config(&config)
let client = client_from_endpoints(&endpoints); .map_err(|source| ProcessSpawnError::Config { source })?;
let client = client_from_endpoints(&endpoints)
.map_err(|source| ProcessSpawnError::Config { source })?;
let child = spawn_child_for_launch(tempdir.path(), &launch).await?; let child = spawn_child_for_launch(tempdir.path(), &launch).await?;