Skip to content

Commit dc1ea69

Browse files
committed
refactor: adapt pdb to build instead of mutate
1 parent bbaee25 commit dc1ea69

3 files changed

Lines changed: 58 additions & 86 deletions

File tree

rust/operator-binary/src/controller.rs

Lines changed: 22 additions & 24 deletions
Original file line numberDiff line numberDiff line change
@@ -16,7 +16,7 @@ use const_format::concatcp;
1616
use snafu::{ResultExt, Snafu};
1717
use stackable_operator::{
1818
cli::OperatorEnvironmentOptions,
19-
cluster_resources::{ClusterResourceApplyStrategy, ClusterResources},
19+
cluster_resources::ClusterResourceApplyStrategy,
2020
commons::{
2121
networking::DomainName, product_image_selection::ResolvedProductImage,
2222
rbac::build_rbac_resources,
@@ -38,6 +38,7 @@ use stackable_operator::{
3838
},
3939
v2::{
4040
HasName, HasUid, NameIsValidLabelValue,
41+
cluster_resources::cluster_resources_new,
4142
kvp::label::{recommended_labels, role_group_selector},
4243
role_group_utils::ResourceNames,
4344
types::{
@@ -63,7 +64,7 @@ use crate::{
6364
properties::listener::get_kafka_listener_config,
6465
resource::{
6566
listener::build_broker_rolegroup_bootstrap_listener,
66-
pdb::add_pdbs,
67+
pdb::build_pdb,
6768
service::{build_rolegroup_headless_service, build_rolegroup_metrics_service},
6869
statefulset::{
6970
build_broker_rolegroup_statefulset, build_controller_rolegroup_statefulset,
@@ -314,17 +315,17 @@ impl NameIsValidLabelValue for ValidatedCluster {
314315
}
315316

316317
/// The product name (`kafka`) as a type-safe label value.
317-
fn product_name() -> ProductName {
318+
pub(crate) fn product_name() -> ProductName {
318319
ProductName::from_str(APP_NAME).expect("'kafka' is a valid product name")
319320
}
320321

321322
/// The operator name as a type-safe label value.
322-
fn operator_name() -> OperatorName {
323+
pub(crate) fn operator_name() -> OperatorName {
323324
OperatorName::from_str(OPERATOR_NAME).expect("the operator name is a valid label value")
324325
}
325326

326327
/// The controller name as a type-safe label value.
327-
fn controller_name() -> ControllerName {
328+
pub(crate) fn controller_name() -> ControllerName {
328329
ControllerName::from_str(KAFKA_CONTROLLER_NAME)
329330
.expect("the controller name is a valid label value")
330331
}
@@ -461,11 +462,6 @@ pub enum Error {
461462
source: stackable_operator::cluster_resources::Error,
462463
},
463464

464-
#[snafu(display("failed to create cluster resources"))]
465-
CreateClusterResources {
466-
source: stackable_operator::cluster_resources::Error,
467-
},
468-
469465
#[snafu(display("failed to patch service account"))]
470466
ApplyServiceAccount {
471467
source: stackable_operator::cluster_resources::Error,
@@ -486,9 +482,9 @@ pub enum Error {
486482
source: stackable_operator::commons::rbac::Error,
487483
},
488484

489-
#[snafu(display("failed to create PodDisruptionBudget"))]
490-
FailedToCreatePdb {
491-
source: crate::controller::build::resource::pdb::Error,
485+
#[snafu(display("failed to apply PodDisruptionBudget"))]
486+
ApplyPdb {
487+
source: stackable_operator::cluster_resources::Error,
492488
},
493489

494490
#[snafu(display("failed to get required Labels"))]
@@ -530,12 +526,11 @@ impl ReconcilerError for Error {
530526
Error::BuildDiscoveryConfig { .. } => None,
531527
Error::ApplyDiscoveryConfig { .. } => None,
532528
Error::DeleteOrphans { .. } => None,
533-
Error::CreateClusterResources { .. } => None,
534529
Error::ApplyServiceAccount { .. } => None,
535530
Error::ApplyRoleBinding { .. } => None,
536531
Error::ApplyStatus { .. } => None,
537532
Error::BuildRbacResources { .. } => None,
538-
Error::FailedToCreatePdb { .. } => None,
533+
Error::ApplyPdb { .. } => None,
539534
Error::GetRequiredLabels { .. } => None,
540535
Error::InvalidKafkaCluster { .. } => None,
541536
Error::BuildStatefulset { .. } => None,
@@ -568,15 +563,16 @@ pub async fn reconcile_kafka(
568563
validate::validate(kafka, dereferenced_objects, &ctx.operator_environment)
569564
.context(ValidateClusterSnafu)?;
570565

571-
let mut cluster_resources = ClusterResources::new(
572-
APP_NAME,
573-
OPERATOR_NAME,
574-
KAFKA_CONTROLLER_NAME,
575-
&validated_cluster.object_ref(&()),
566+
let mut cluster_resources = cluster_resources_new(
567+
&product_name(),
568+
&operator_name(),
569+
&controller_name(),
570+
&validated_cluster.name,
571+
&validated_cluster.namespace,
572+
&validated_cluster.uid,
576573
ClusterResourceApplyStrategy::from(&kafka.spec.cluster_operation),
577574
&kafka.spec.object_overrides,
578-
)
579-
.context(CreateClusterResourcesSnafu)?;
575+
);
580576

581577
tracing::debug!(
582578
kerberos_enabled = validated_cluster.cluster_config.kafka_security.has_kerberos_enabled(),
@@ -718,10 +714,12 @@ pub async fn reconcile_kafka(
718714
if let Some(GenericRoleConfig {
719715
pod_disruption_budget: pdb,
720716
}) = role_cfg
717+
&& let Some(pdb) = build_pdb(pdb, &validated_cluster, kafka_role)
721718
{
722-
add_pdbs(pdb, kafka, kafka_role, client, &mut cluster_resources)
719+
cluster_resources
720+
.add(client, pdb)
723721
.await
724-
.context(FailedToCreatePdbSnafu)?;
722+
.context(ApplyPdbSnafu)?;
725723
}
726724
}
727725

rust/operator-binary/src/controller/build/resource/pdb.rs

Lines changed: 16 additions & 40 deletions
Original file line numberDiff line numberDiff line change
@@ -1,61 +1,37 @@
1-
use snafu::{ResultExt, Snafu};
21
use stackable_operator::{
3-
builder::pdb::PodDisruptionBudgetBuilder, client::Client, cluster_resources::ClusterResources,
4-
commons::pdb::PdbConfig, kube::ResourceExt,
2+
commons::pdb::PdbConfig, k8s_openapi::api::policy::v1::PodDisruptionBudget,
3+
v2::builder::pdb::pod_disruption_budget_builder_with_role,
54
};
65

76
use crate::{
8-
controller::KAFKA_CONTROLLER_NAME,
9-
crd::{APP_NAME, OPERATOR_NAME, role::KafkaRole, v1alpha1},
7+
controller::{ValidatedCluster, controller_name, operator_name, product_name},
8+
crd::role::KafkaRole,
109
};
1110

12-
#[derive(Snafu, Debug)]
13-
pub enum Error {
14-
#[snafu(display("Cannot create PodDisruptionBudget for role [{role}]"))]
15-
CreatePdb {
16-
source: stackable_operator::builder::pdb::Error,
17-
role: String,
18-
},
19-
#[snafu(display("Cannot apply PodDisruptionBudget [{name}]"))]
20-
ApplyPdb {
21-
source: stackable_operator::cluster_resources::Error,
22-
name: String,
23-
},
24-
}
25-
26-
pub async fn add_pdbs(
11+
/// Builds the [`PodDisruptionBudget`] for the given `role`, or `None` if PDBs are disabled.
12+
pub fn build_pdb(
2713
pdb: &PdbConfig,
28-
kafka: &v1alpha1::KafkaCluster,
14+
validated_cluster: &ValidatedCluster,
2915
role: &KafkaRole,
30-
client: &Client,
31-
cluster_resources: &mut ClusterResources<'_>,
32-
) -> Result<(), Error> {
16+
) -> Option<PodDisruptionBudget> {
3317
if !pdb.enabled {
34-
return Ok(());
18+
return None;
3519
}
3620
let max_unavailable = pdb.max_unavailable.unwrap_or(match role {
3721
KafkaRole::Broker => max_unavailable_brokers(),
3822
KafkaRole::Controller => max_unavailable_controllers(),
3923
});
40-
let pdb = PodDisruptionBudgetBuilder::new_with_role(
41-
kafka,
42-
APP_NAME,
43-
&role.to_string(),
44-
OPERATOR_NAME,
45-
KAFKA_CONTROLLER_NAME,
24+
let pdb = pod_disruption_budget_builder_with_role(
25+
validated_cluster,
26+
&product_name(),
27+
&ValidatedCluster::role_name(role),
28+
&operator_name(),
29+
&controller_name(),
4630
)
47-
.with_context(|_| CreatePdbSnafu {
48-
role: role.to_string(),
49-
})?
5031
.with_max_unavailable(max_unavailable)
5132
.build();
52-
let pdb_name = pdb.name_any();
53-
cluster_resources
54-
.add(client, pdb)
55-
.await
56-
.with_context(|_| ApplyPdbSnafu { name: pdb_name })?;
5733

58-
Ok(())
34+
Some(pdb)
5935
}
6036

6137
fn max_unavailable_brokers() -> u16 {

rust/operator-binary/src/controller/build/resource/statefulset.rs

Lines changed: 20 additions & 22 deletions
Original file line numberDiff line numberDiff line change
@@ -1,15 +1,12 @@
1-
use std::ops::Deref;
1+
use std::{ops::Deref, str::FromStr};
22

33
use snafu::{OptionExt, ResultExt, Snafu};
44
use stackable_operator::{
55
builder::{
66
meta::ObjectMetaBuilder,
77
pod::{
8-
PodBuilder,
9-
container::ContainerBuilder,
10-
resources::ResourceRequirementsBuilder,
11-
security::PodSecurityContextBuilder,
12-
volume::{ListenerOperatorVolumeSourceBuilder, ListenerReference, VolumeBuilder},
8+
PodBuilder, container::ContainerBuilder, resources::ResourceRequirementsBuilder,
9+
security::PodSecurityContextBuilder, volume::VolumeBuilder,
1310
},
1411
},
1512
commons::product_image_selection::ResolvedProductImage,
@@ -35,8 +32,13 @@ use stackable_operator::{
3532
},
3633
},
3734
v2::{
38-
builder::meta::ownerreference_from_resource, jvm_argument_overrides::JvmArgumentOverrides,
35+
builder::{
36+
meta::ownerreference_from_resource,
37+
pod::volume::{ListenerReference, listener_operator_volume_source_builder_build_pvc},
38+
},
39+
jvm_argument_overrides::JvmArgumentOverrides,
3940
role_group_utils::ResourceNames,
41+
types::kubernetes::{ListenerName, PersistentVolumeClaimName},
4042
},
4143
};
4244

@@ -97,11 +99,6 @@ pub enum Error {
9799
source: stackable_operator::builder::pod::Error,
98100
},
99101

100-
#[snafu(display("failed to build bootstrap listener pvc"))]
101-
BuildBootstrapListenerPvc {
102-
source: stackable_operator::builder::pod::volume::ListenerOperatorVolumeSourceBuilderError,
103-
},
104-
105102
#[snafu(display("failed to build pod descriptors"))]
106103
BuildPodDescriptors {
107104
source: crate::controller::PodDescriptorsError,
@@ -191,16 +188,17 @@ pub fn build_broker_rolegroup_statefulset(
191188

192189
// bootstrap listener should be persistent,
193190
// main broker listener is an ephemeral PVC instead
194-
pvcs.push(
195-
ListenerOperatorVolumeSourceBuilder::new(
196-
&ListenerReference::ListenerName(
197-
validated_cluster.bootstrap_listener_name(kafka_role, role_group_name),
198-
),
199-
&unversioned_recommended_labels,
200-
)
201-
.build_pvc(LISTENER_BOOTSTRAP_VOLUME_NAME)
202-
.context(BuildBootstrapListenerPvcSnafu)?,
203-
);
191+
let bootstrap_listener_name = ListenerName::from_str(
192+
&validated_cluster.bootstrap_listener_name(kafka_role, role_group_name),
193+
)
194+
.expect("the bootstrap listener name is a valid Listener name");
195+
let bootstrap_pvc_name = PersistentVolumeClaimName::from_str(LISTENER_BOOTSTRAP_VOLUME_NAME)
196+
.expect("the bootstrap listener volume name is a valid PVC name");
197+
pvcs.push(listener_operator_volume_source_builder_build_pvc(
198+
&ListenerReference::Listener(bootstrap_listener_name),
199+
&unversioned_recommended_labels,
200+
&bootstrap_pvc_name,
201+
));
204202

205203
if kafka_security.has_kerberos_enabled() {
206204
add_kerberos_pod_config(

0 commit comments

Comments
 (0)