Skip to content
Merged
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
4 changes: 4 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -36,6 +36,9 @@ All notable changes to this project will be documented in this file.
- Internal operator refactoring: the validated cluster carries each role's configuration in typed
per-role fields instead of maps keyed by role, so role-specific values are resolved once and the
shared resource builders can no longer be handed another role's configuration ([#830]).
- Internal operator refactoring: each role now has its own module that gathers its role group's
inputs and role-specific containers, and a shared role-builder then builds the Services,
ConfigMap and StatefulSet from them ([#833]).
- Bump stackable-operator to 0.119.0 ([#835]).

### Fixed
Expand Down Expand Up @@ -63,6 +66,7 @@ All notable changes to this project will be documented in this file.
[#829]: https://github.com/stackabletech/hdfs-operator/pull/829
[#830]: https://github.com/stackabletech/hdfs-operator/pull/830
[#831]: https://github.com/stackabletech/hdfs-operator/pull/831
[#833]: https://github.com/stackabletech/hdfs-operator/pull/833
[#835]: https://github.com/stackabletech/hdfs-operator/pull/835

## [26.7.0] - 2026-07-21
Expand Down
383 changes: 115 additions & 268 deletions rust/operator-binary/src/controller/build/container.rs

Large diffs are not rendered by default.

221 changes: 68 additions & 153 deletions rust/operator-binary/src/controller/build/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -23,12 +23,19 @@ use stackable_operator::{
},
},
};
use strum::IntoEnumIterator;

use crate::{
controller::{
CONTROLLER_NAME, KubernetesResources, OPERATOR_NAME, PRODUCT_NAME, Prepared,
ValidatedCluster,
build::resource::rbac::{build_role_binding, build_service_account},
build::{
resource::rbac::{build_role_binding, build_service_account},
role_group::{
RoleGroupBuilder, build_datanode_role_group, build_journalnode_role_group,
build_namenode_role_group,
},
},
},
crd::{
HdfsNodeRole, HdfsPodRef,
Expand Down Expand Up @@ -57,61 +64,19 @@ pub mod jvm;
pub mod kerberos;
pub mod opa;
pub mod properties;
pub mod resolve;
pub mod resource;
pub mod role_group;

#[derive(Snafu, Debug)]
pub enum Error {
#[snafu(display("failed to build Service for role {role} role group {role_group}", role = role.as_ref()))]
Service {
source: resource::service::Error,
role: HdfsNodeRole,
role_group: RoleGroupName,
},

#[snafu(display("failed to build ConfigMap for role {role} role group {role_group}", role = role.as_ref()))]
ConfigMap {
source: resource::config_map::Error,
role: HdfsNodeRole,
role_group: RoleGroupName,
},

#[snafu(display("failed to build StatefulSet for role {role} role group {role_group}", role = role.as_ref()))]
StatefulSet {
source: resource::statefulset::Error,
role: HdfsNodeRole,
role_group: RoleGroupName,
},
#[snafu(display("failed to build the resources of a role group"))]
RoleGroup { source: role_group::Error },

#[snafu(display("failed to build the discovery ConfigMap"))]
DiscoveryConfigMap { source: resource::discovery::Error },

#[snafu(display("failed to build selector labels for role {role} role group {role_group}", role = role.as_ref()))]
RoleGroupSelectorLabels {
source: LabelError,
role: HdfsNodeRole,
role_group: RoleGroupName,
},

#[snafu(display("failed to build volume claim templates for role {role} role group {role_group}", role = role.as_ref()))]
VolumeClaimTemplates {
source: container::Error,
role: HdfsNodeRole,
role_group: RoleGroupName,
},

#[snafu(display("failed to build listener volume for role {role} role group {role_group}", role = role.as_ref()))]
ListenerVolume {
source: container::Error,
role: HdfsNodeRole,
role_group: RoleGroupName,
},
}

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

/// The resources of every role, accumulated one role at a time by [`build_role`].
/// The resources of every role, accumulated one role group at a time by [`build`].
#[derive(Default)]
struct RoleGroupResources {
services: Vec<Service>,
Expand All @@ -122,63 +87,18 @@ struct RoleGroupResources {
pod_disruption_budgets: Vec<PodDisruptionBudget>,
}

/// Builds every resource of every role group of one role, plus that role's PDB, appending them to
/// `rg_resources`.
fn build_role<C: RoleGroupResolver>(
cluster: &ValidatedCluster,
cluster_info: &KubernetesClusterInfo,
role_group_configs: &BTreeMap<
RoleGroupName,
RoleGroupConfig<C, JavaCommonConfig, v1alpha1::HdfsConfigOverrides>,
>,
rg_resources: &mut RoleGroupResources,
) -> Result<(), Error> {
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)?;

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(),
})?,
);
impl RoleGroupResources {
/// Builds the Services, ConfigMap and StatefulSet of one role group and adds them to the
/// collections.
fn add(&mut self, builder: &RoleGroupBuilder) -> Result<(), role_group::Error> {
self.services.extend(builder.build_services()?);
self.config_maps.push(builder.build_config_map()?);
self.stateful_sets
.entry(builder.role)
.or_default()
.push(builder.build_statefulset()?);
Ok(())
}

if let Some(pdb) = resource::pdb::build_pdb(cluster, role) {
rg_resources.pod_disruption_budgets.push(pdb);
}

Ok(())
}

/// Builds every Kubernetes resource for the given validated cluster.
Expand All @@ -188,6 +108,11 @@ fn build_role<C: RoleGroupResolver>(
/// `cluster_info` carries static cluster information resolved at operator startup (e.g. the
/// cluster domain used to build Kerberos principals), not a live client.
///
/// Each of the three loops hands its role group's typed config to that role's builder, which is
/// where everything specific to the role lives. The loops are free to be reordered: the
/// StatefulSets are keyed by role, and the apply step does not depend on the order of the other
/// three collections.
///
/// The resources are returned as flat collections. `stateful_sets` comes out in [`HdfsNodeRole`]
/// order, which the apply step depends on; that is structural, from a [`BTreeMap`] flattened in
/// key order, not from the order the roles are built in.
Expand All @@ -200,26 +125,33 @@ pub fn build(
) -> Result<KubernetesResources<Prepared>, Error> {
let mut built = RoleGroupResources::default();

// These three calls are free to be reordered: the StatefulSets are keyed by role, and the
// apply step does not depend on the order of the other three collections.
build_role(
cluster,
cluster_info,
&cluster.journalnode_role_group_configs,
&mut built,
)?;
build_role(
cluster,
cluster_info,
&cluster.namenode_role_group_configs,
&mut built,
)?;
build_role(
cluster,
cluster_info,
&cluster.datanode_role_group_configs,
&mut built,
)?;
for (role_group_name, rg_config) in &cluster.journalnode_role_group_configs {
let builder =
build_journalnode_role_group(cluster, cluster_info, role_group_name, rg_config)
.context(RoleGroupSnafu)?;

built.add(&builder).context(RoleGroupSnafu)?;
}

for (role_group_name, rg_config) in &cluster.namenode_role_group_configs {
let builder = build_namenode_role_group(cluster, cluster_info, role_group_name, rg_config)
.context(RoleGroupSnafu)?;

built.add(&builder).context(RoleGroupSnafu)?;
}

for (role_group_name, rg_config) in &cluster.datanode_role_group_configs {
let builder = build_datanode_role_group(cluster, cluster_info, role_group_name, rg_config)
.context(RoleGroupSnafu)?;

built.add(&builder).context(RoleGroupSnafu)?;
}

for role in HdfsNodeRole::iter() {
if let Some(pdb) = resource::pdb::build_pdb(cluster, &role) {
built.pod_disruption_budgets.push(pdb);
}
}

let RoleGroupResources {
services,
Expand Down Expand Up @@ -250,34 +182,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 Expand Up @@ -456,8 +360,19 @@ pub(crate) fn native_metrics_port(cluster: &ValidatedCluster, role: &HdfsNodeRol
}
}

/// The deprecated JMX exporter metrics port for the given `role`.
fn jmx_metrics_port(role: &HdfsNodeRole) -> Port {
/// The name of the port the given `role` serves IPC/RPC on, which its readiness probe checks.
///
/// The datanodes call theirs `ipc`, the other two `rpc`; the same names [`role_data_ports`]
/// exposes them under.
pub(crate) fn ipc_port_name(role: &HdfsNodeRole) -> &'static str {
match role {
HdfsNodeRole::Name | HdfsNodeRole::Journal => SERVICE_PORT_NAME_RPC,
HdfsNodeRole::Data => SERVICE_PORT_NAME_IPC,
}
}

/// The deprecated JMX Exporter metrics port for the given `role`.
pub(crate) fn jmx_metrics_port(role: &HdfsNodeRole) -> Port {
match role {
HdfsNodeRole::Name => DEFAULT_NAME_NODE_METRICS_PORT,
HdfsNodeRole::Data => DEFAULT_DATA_NODE_METRICS_PORT,
Expand Down
Loading
Loading