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
65 changes: 28 additions & 37 deletions rust/operator-binary/src/controller/build/container.rs
Original file line number Diff line number Diff line change
Expand Up @@ -53,11 +53,9 @@ use stackable_operator::{
STACKABLE_LOG_DIR, ValidatedContainerLogConfigChoice, VectorContainerLogConfig,
vector_container,
},
role_utils::{JavaCommonConfig, RoleGroupConfig},
types::{
common::Port,
kubernetes::{ConfigMapName, ContainerName, VolumeName},
operator::RoleGroupName,
},
},
};
Expand All @@ -67,7 +65,7 @@ use crate::{
controller::{
ValidatedCluster,
build::{
self, ResolvedRoleGroup, RoleGroupResolver, RoleSpecificValues,
self, ResolvedRoleGroup, RoleGroupBuilder, RoleSpecificValues,
jvm::{self, construct_global_jvm_args, construct_role_specific_jvm_args},
kerberos::KERBEROS_CONTAINER_PATH,
properties::product_logging::{
Expand All @@ -92,7 +90,6 @@ use crate::{
SERVICE_PORT_NAME_RPC, STACKABLE_ROOT_DATA_DIR,
},
storage::DataNodeStorageConfig,
v1alpha1,
},
};

Expand Down Expand Up @@ -213,17 +210,19 @@ impl ContainerConfig {

/// Add all main, side and init containers as well as required volumes to the pod builder.
///
/// Every role-specific value is resolved by the caller into `resolved`; the role itself comes
/// from `C::ROLE`, the same `C` that produced it.
pub fn add_containers_and_volumes<C: RoleGroupResolver>(
/// Everything about the role group comes from `resolved`, the role and the merged overrides
/// included, so there is nothing here to pair with the wrong role group.
pub(crate) fn add_containers_and_volumes(
pb: &mut PodBuilder,
cluster: &ValidatedCluster,
cluster_info: &KubernetesClusterInfo,
role_group_name: &RoleGroupName,
rolegroup_config: &RoleGroupConfig<C, JavaCommonConfig, v1alpha1::HdfsConfigOverrides>,
resolved: &ResolvedRoleGroup<C>,
builder: &RoleGroupBuilder,
) -> Result<(), Error> {
let role = &C::ROLE;
let RoleGroupBuilder {
cluster,
cluster_info,
role_group_name,
resolved,
} = builder;
let role = &builder.role();
let namenode_podrefs = build::pod_refs(cluster, &HdfsNodeRole::Name);

// HDFS main container
Expand All @@ -241,7 +240,6 @@ impl ContainerConfig {
cluster,
cluster_info,
&resolved.logging.hdfs,
rolegroup_config,
resolved,
)?);

Expand Down Expand Up @@ -355,7 +353,6 @@ impl ContainerConfig {
cluster,
cluster_info,
zkfc,
rolegroup_config,
resolved,
)?);

Expand All @@ -371,7 +368,6 @@ impl ContainerConfig {
cluster,
cluster_info,
format_namenodes,
rolegroup_config,
resolved,
&namenode_podrefs,
)?);
Expand All @@ -388,7 +384,6 @@ impl ContainerConfig {
cluster,
cluster_info,
format_zookeeper,
rolegroup_config,
resolved,
&namenode_podrefs,
)?);
Expand All @@ -408,7 +403,6 @@ impl ContainerConfig {
cluster,
cluster_info,
wait_for_namenodes,
rolegroup_config,
resolved,
&namenode_podrefs,
)?);
Expand Down Expand Up @@ -491,23 +485,22 @@ impl ContainerConfig {
/// - Namenode ZooKeeper fail over controller (ZKFC)
/// - Datanode main process
/// - Journalnode main process
fn main_container<C: RoleGroupResolver>(
fn main_container(
&self,
cluster: &ValidatedCluster,
cluster_info: &KubernetesClusterInfo,
container_log_config: &ContainerLogConfig,
rolegroup_config: &RoleGroupConfig<C, JavaCommonConfig, v1alpha1::HdfsConfigOverrides>,
resolved: &ResolvedRoleGroup<C>,
resolved: &ResolvedRoleGroup,
) -> Result<Container, Error> {
let role = &C::ROLE;
let role = &resolved.role.node_role();
let mut cb = new_container_builder(self.container_name());

let resources = self.resources(&resolved.resources);

cb.image_from_product_image(&cluster.image)
.command(Self::command())
.args(self.args(cluster, cluster_info, role, container_log_config, &[])?)
.add_env_vars(self.env(cluster, role, rolegroup_config, resources.as_ref())?)
.add_env_vars(self.env(cluster, role, resolved, resources.as_ref())?)
.add_volume_mounts(self.volume_mounts(cluster, &resolved.volume_claim_templates))
.context(AddVolumeMountSnafu)?
.add_container_ports(self.container_ports(cluster));
Expand Down Expand Up @@ -539,16 +532,15 @@ impl ContainerConfig {
/// Creates respective init containers for:
/// - Namenode (format-namenodes, format-zookeeper)
/// - Datanode (wait-for-namenodes)
fn init_container<C: RoleGroupResolver>(
fn init_container(
&self,
cluster: &ValidatedCluster,
cluster_info: &KubernetesClusterInfo,
container_log_config: &ContainerLogConfig,
rolegroup_config: &RoleGroupConfig<C, JavaCommonConfig, v1alpha1::HdfsConfigOverrides>,
resolved: &ResolvedRoleGroup<C>,
resolved: &ResolvedRoleGroup,
namenode_podrefs: &[HdfsPodRef],
) -> Result<Container, Error> {
let role = &C::ROLE;
let role = &resolved.role.node_role();
let mut cb = new_container_builder(self.container_name());

cb.image_from_product_image(&cluster.image)
Expand All @@ -560,7 +552,7 @@ impl ContainerConfig {
container_log_config,
namenode_podrefs,
)?)
.add_env_vars(self.env(cluster, role, rolegroup_config, None)?)
.add_env_vars(self.env(cluster, role, resolved, None)?)
.add_volume_mounts(self.volume_mounts(cluster, &resolved.volume_claim_templates))
.context(AddVolumeMountSnafu)?;

Expand Down Expand Up @@ -881,11 +873,11 @@ impl ContainerConfig {
}

/// Returns the container env variables.
fn env<C>(
fn env(
&self,
cluster: &ValidatedCluster,
role: &HdfsNodeRole,
rolegroup_config: &RoleGroupConfig<C, JavaCommonConfig, v1alpha1::HdfsConfigOverrides>,
resolved: &ResolvedRoleGroup,
resources: Option<&ResourceRequirements>,
) -> Result<Vec<EnvVar>, Error> {
// Maps env var name to env var object. This allows env_overrides to work
Expand Down Expand Up @@ -916,7 +908,7 @@ impl ContainerConfig {
role_opts_name.clone(),
EnvVar {
name: role_opts_name,
value: Some(self.build_hadoop_opts(cluster, resources, rolegroup_config)?),
value: Some(self.build_hadoop_opts(cluster, resources, resolved)?),
..EnvVar::default()
},
);
Expand Down Expand Up @@ -980,7 +972,8 @@ impl ContainerConfig {
);

// Overrides need to come last
let mut env_override_vars: BTreeMap<String, EnvVar> = rolegroup_config
let mut env_override_vars: BTreeMap<String, EnvVar> = resolved
.merged
.env_overrides
.clone()
.into_iter()
Expand Down Expand Up @@ -1253,11 +1246,11 @@ impl ContainerConfig {
}

/// Build HADOOP_{*node}_OPTS for each namenode, datanodes and journalnodes.
fn build_hadoop_opts<C>(
fn build_hadoop_opts(
&self,
cluster: &ValidatedCluster,
resources: Option<&ResourceRequirements>,
rolegroup_config: &RoleGroupConfig<C, JavaCommonConfig, v1alpha1::HdfsConfigOverrides>,
resolved: &ResolvedRoleGroup,
) -> Result<String, Error> {
match self {
ContainerConfig::Hdfs {
Expand All @@ -1267,9 +1260,7 @@ impl ContainerConfig {
let config_dir = volume_mount_dirs.final_config();
construct_role_specific_jvm_args(
role,
&rolegroup_config
.product_specific_common_config
.jvm_argument_overrides,
&resolved.merged.jvm_argument_overrides,
cluster.has_kerberos_enabled(),
resources,
config_dir,
Expand Down
78 changes: 13 additions & 65 deletions rust/operator-binary/src/controller/build/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -59,6 +59,7 @@ pub mod opa;
pub mod properties;
pub mod resolve;
pub mod resource;
pub mod role_group_builder;

#[derive(Snafu, Debug)]
pub enum Error {
Expand Down Expand Up @@ -108,8 +109,10 @@ pub enum Error {
},
}

pub(crate) use resolve::RoleGroupResolver;
pub use resolve::{ResolvedRoleGroup, RoleGroupLogging, RoleSpecificValues};
pub(crate) use resolve::{
ResolvedRoleGroup, RoleGroupLogging, RoleGroupResolver, RoleSpecificValues,
};
pub(crate) use role_group_builder::RoleGroupBuilder;

/// The resources of every role, accumulated one role at a time by [`build_role`].
#[derive(Default)]
Expand All @@ -136,42 +139,15 @@ fn build_role<C: RoleGroupResolver>(
let role = &C::ROLE;

for (role_group_name, rg_config) in role_group_configs {
build_role_group_services(cluster, role, role_group_name, &mut rg_resources.services)?;

let selector_labels = rolegroup_selector_labels(cluster, role, role_group_name).context(
RoleGroupSelectorLabelsSnafu {
role: *role,
role_group: role_group_name.clone(),
},
)?;
let resolved = rg_config.config.resolve(role_group_name, selector_labels)?;
let builder = RoleGroupBuilder::new(cluster, cluster_info, role_group_name, rg_config)?;

rg_resources.config_maps.push(
resource::config_map::build_rolegroup_config_map(
cluster,
cluster_info,
role_group_name,
rg_config,
&resolved,
)
.context(ConfigMapSnafu {
role: *role,
role_group: role_group_name.clone(),
})?,
);
rg_resources.stateful_sets.entry(C::ROLE).or_default().push(
resource::statefulset::build_rolegroup_statefulset(
cluster,
cluster_info,
role_group_name,
rg_config,
&resolved,
)
.context(StatefulSetSnafu {
role: *role,
role_group: role_group_name.clone(),
})?,
);
rg_resources.services.extend(builder.build_services()?);
rg_resources.config_maps.push(builder.build_config_map()?);
rg_resources
.stateful_sets
.entry(C::ROLE)
.or_default()
.push(builder.build_stateful_set()?);
}

if let Some(pdb) = resource::pdb::build_pdb(cluster, role) {
Expand Down Expand Up @@ -250,34 +226,6 @@ pub fn build(
})
}

/// Builds the two Services for one role group. Role-agnostic: it reads nothing from the role
/// config.
fn build_role_group_services(
cluster: &ValidatedCluster,
role: &HdfsNodeRole,
role_group_name: &RoleGroupName,
services: &mut Vec<Service>,
) -> Result<(), Error> {
services.push(
resource::service::rolegroup_headless_service(cluster, role, role_group_name).context(
ServiceSnafu {
role: *role,
role_group: role_group_name.clone(),
},
)?,
);
services.push(
resource::service::rolegroup_metrics_service(cluster, role, role_group_name).context(
ServiceSnafu {
role: *role,
role_group: role_group_name.clone(),
},
)?,
);

Ok(())
}

/// The replica count a role group gets when it does not set one: Kubernetes runs a single pod for
/// a `StatefulSet` with `replicas: null`.
pub(crate) const DEFAULT_REPLICAS: u16 = 1;
Expand Down
Loading
Loading