Skip to content

Commit 43dc9e5

Browse files
committed
refactor: use v2 owerref, labels, RoleGroupName, cleanup
1 parent f0e6cbb commit 43dc9e5

19 files changed

Lines changed: 1059 additions & 1043 deletions

File tree

rust/operator-binary/src/config/jvm.rs

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -99,7 +99,7 @@ mod tests {
9999
let rg = validated
100100
.role_group_configs
101101
.get(&KafkaRole::Broker)
102-
.and_then(|groups| groups.get("default"))
102+
.and_then(|groups| groups.get(&"default".parse().unwrap()))
103103
.expect("broker default role group should exist");
104104
(rg.config.clone(), rg.jvm_argument_overrides.clone())
105105
}

rust/operator-binary/src/config/node_id_hasher.rs

Lines changed: 3 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -1,19 +1,13 @@
1-
use stackable_operator::role_utils::RoleGroupRef;
2-
3-
use crate::crd::v1alpha1::KafkaCluster;
1+
use crate::crd::role::KafkaRole;
42

53
/// The Kafka node.id needs to be unique across the Kafka cluster.
64
/// This function generates an integer that is stable for a given role group
75
/// regardless if broker or controllers.
86
/// This integer is then added to the pod index to compute the final node.id
97
/// The node.id is only set and used in Kraft mode.
108
/// Warning: this is not safe from collisions.
11-
pub fn node_id_hash32_offset(rolegroup_ref: &RoleGroupRef<KafkaCluster>) -> u32 {
12-
let hash = fnv_hash32(&format!(
13-
"{role}-{rolegroup}",
14-
role = rolegroup_ref.role,
15-
rolegroup = rolegroup_ref.role_group
16-
));
9+
pub fn node_id_hash32_offset(role: &KafkaRole, role_group: &str) -> u32 {
10+
let hash = fnv_hash32(&format!("{role}-{role_group}"));
1711
let range = hash & 0x0000FFFF;
1812
// Kafka uses signed integer
1913
range * 0x00007FFF

rust/operator-binary/src/controller/build/config_map.rs

Lines changed: 53 additions & 38 deletions
Original file line numberDiff line numberDiff line change
@@ -3,65 +3,70 @@ use snafu::{ResultExt, Snafu};
33
use stackable_operator::{
44
builder::{configmap::ConfigMapBuilder, meta::ObjectMetaBuilder},
55
k8s_openapi::api::core::v1::ConfigMap,
6-
role_utils::RoleGroupRef,
7-
v2::config_file_writer::{PropertiesWriterError, to_java_properties_string},
6+
product_logging::framework::VECTOR_CONFIG_FILE,
7+
v2::{
8+
builder::meta::ownerreference_from_resource,
9+
config_file_writer::{PropertiesWriterError, to_java_properties_string},
10+
},
811
};
912

1013
use crate::{
1114
controller::{
12-
ValidatedCluster, ValidatedRoleGroupConfig,
15+
RoleGroupName, ValidatedCluster, ValidatedRoleGroupConfig,
1316
build::properties::logging::role_group_config_map_data,
1417
},
1518
crd::{
1619
ConfigFileName, STACKABLE_LISTENER_BOOTSTRAP_DIR, STACKABLE_LISTENER_BROKER_DIR,
1720
listener::{KafkaListenerConfig, node_address_cmd},
1821
role::AnyConfig,
19-
v1alpha1,
2022
},
21-
kafka_controller::{KAFKA_CONTROLLER_NAME, build_recommended_labels},
2223
};
2324

2425
#[derive(Snafu, Debug)]
2526
pub enum Error {
26-
#[snafu(display("failed to build ConfigMap for {}", rolegroup))]
27+
#[snafu(display("failed to build ConfigMap for role group {role_group}"))]
2728
BuildRoleGroupConfig {
2829
source: stackable_operator::builder::configmap::Error,
29-
rolegroup: RoleGroupRef<v1alpha1::KafkaCluster>,
30+
role_group: RoleGroupName,
3031
},
3132

32-
#[snafu(display("failed to serialize [{}] for {rolegroup}", ConfigFileName::Security))]
33+
#[snafu(display(
34+
"failed to serialize [{}] for role group {role_group}",
35+
ConfigFileName::Security
36+
))]
3337
JvmSecurityProperties {
3438
source: PropertiesWriterError,
35-
rolegroup: String,
36-
},
37-
38-
#[snafu(display("failed to build Metadata"))]
39-
MetadataBuild {
40-
source: stackable_operator::builder::meta::Error,
41-
},
42-
43-
#[snafu(display("object is missing metadata to build owner reference"))]
44-
ObjectMissingMetadataForOwnerRef {
45-
source: stackable_operator::builder::meta::Error,
39+
role_group: RoleGroupName,
4640
},
4741

48-
#[snafu(display("failed to serialize config for {rolegroup}"))]
42+
#[snafu(display("failed to serialize config for role group {role_group}"))]
4943
SerializeConfig {
5044
source: PropertiesWriterError,
51-
rolegroup: RoleGroupRef<v1alpha1::KafkaCluster>,
45+
role_group: RoleGroupName,
46+
},
47+
48+
#[snafu(display("failed to build pod descriptors"))]
49+
BuildPodDescriptors {
50+
source: crate::controller::PodDescriptorsError,
5251
},
5352

5453
#[snafu(display("no Kraft controllers found to build"))]
5554
NoKraftControllersFound,
5655
}
5756

58-
/// The rolegroup [`ConfigMap`] configures the rolegroup based on the configuration given by the administrator
57+
/// The rolegroup [`ConfigMap`] configures the rolegroup based on the configuration given by the administrator.
58+
///
59+
/// `vector_config` is the Vector agent config built by the caller (where a `RoleGroupRef` is
60+
/// available); it is `None` when the Vector agent is disabled. Resource naming and labels use the
61+
/// role (derived from `validated_rg.config`) and the typed `role_group_name`.
5962
pub fn build_rolegroup_config_map(
6063
validated_cluster: &ValidatedCluster,
61-
rolegroup: &RoleGroupRef<v1alpha1::KafkaCluster>,
64+
role_group_name: &RoleGroupName,
6265
validated_rg: &ValidatedRoleGroupConfig,
6366
listener_config: &KafkaListenerConfig,
67+
vector_config: Option<String>,
6468
) -> Result<ConfigMap, Error> {
69+
let role = validated_rg.config.kafka_role();
6570
let cluster_config = &validated_cluster.cluster_config;
6671
let kafka_security = &cluster_config.kafka_security;
6772
let resolved_product_image = &validated_cluster.image;
@@ -72,20 +77,26 @@ pub fn build_rolegroup_config_map(
7277
.overrides
7378
.clone();
7479

75-
if cluster_config.is_kraft_mode() && cluster_config.pod_descriptors.is_empty() {
80+
let pod_descriptors = validated_cluster
81+
.pod_descriptors(None)
82+
.context(BuildPodDescriptorsSnafu)?;
83+
84+
if cluster_config.is_kraft_mode() && pod_descriptors.is_empty() {
7685
return NoKraftControllersFoundSnafu.fail();
7786
}
7887

7988
let kafka_config = match &validated_rg.config {
8089
AnyConfig::Broker(_) => crate::controller::build::properties::broker_properties::build(
8190
cluster_config,
8291
listener_config,
92+
&pod_descriptors,
8393
config_overrides,
8494
),
8595
AnyConfig::Controller(_) => {
8696
crate::controller::build::properties::controller_properties::build(
8797
cluster_config,
8898
listener_config,
99+
&pod_descriptors,
89100
config_overrides,
90101
)
91102
}
@@ -101,32 +112,33 @@ pub fn build_rolegroup_config_map(
101112
.metadata(
102113
ObjectMetaBuilder::new()
103114
.name_and_namespace(validated_cluster)
104-
.name(rolegroup.object_name())
105-
.ownerreference_from_resource(validated_cluster, None, Some(true))
106-
.context(ObjectMissingMetadataForOwnerRefSnafu)?
107-
.with_recommended_labels(&build_recommended_labels(
115+
.name(
116+
validated_cluster
117+
.resource_names(&role, role_group_name)
118+
.role_group_config_map()
119+
.to_string(),
120+
)
121+
.ownerreference(ownerreference_from_resource(
108122
validated_cluster,
109-
KAFKA_CONTROLLER_NAME,
110-
&resolved_product_image.app_version_label_value,
111-
&rolegroup.role,
112-
&rolegroup.role_group,
123+
None,
124+
Some(true),
113125
))
114-
.context(MetadataBuildSnafu)?
126+
.with_labels(validated_cluster.recommended_labels(&role, role_group_name))
115127
.build(),
116128
)
117129
.add_data(
118130
kafka_config_file_name,
119131
to_java_properties_string(kafka_config.iter()).with_context(|_| {
120132
SerializeConfigSnafu {
121-
rolegroup: rolegroup.clone(),
133+
role_group: role_group_name.clone(),
122134
}
123135
})?,
124136
)
125137
.add_data(
126138
ConfigFileName::Security.to_string(),
127139
to_java_properties_string(jvm_sec_props.iter()).with_context(|_| {
128140
JvmSecurityPropertiesSnafu {
129-
rolegroup: rolegroup.role_group.clone(),
141+
role_group: role_group_name.clone(),
130142
}
131143
})?,
132144
)
@@ -139,7 +151,7 @@ pub fn build_rolegroup_config_map(
139151
.filter_map(|(k, v)| v.as_ref().map(|v| (k, v))),
140152
)
141153
.with_context(|_| JvmSecurityPropertiesSnafu {
142-
rolegroup: rolegroup.role_group.clone(),
154+
role_group: role_group_name.clone(),
143155
})?,
144156
)
145157
// This file contains the JAAS configuration for Kerberos authentication
@@ -156,7 +168,6 @@ pub fn build_rolegroup_config_map(
156168

157169
let config_data = role_group_config_map_data(
158170
&resolved_product_image.product_version,
159-
rolegroup,
160171
&validated_rg.config,
161172
);
162173
for (file_name, data) in config_data {
@@ -165,10 +176,14 @@ pub fn build_rolegroup_config_map(
165176
}
166177
}
167178

179+
if let Some(vector_config) = vector_config {
180+
cm_builder.add_data(VECTOR_CONFIG_FILE, vector_config);
181+
}
182+
168183
cm_builder
169184
.build()
170185
.with_context(|_| BuildRoleGroupConfigSnafu {
171-
rolegroup: rolegroup.clone(),
186+
role_group: role_group_name.clone(),
172187
})
173188
}
174189

rust/operator-binary/src/controller/build/discovery.rs

Lines changed: 14 additions & 33 deletions
Original file line numberDiff line numberDiff line change
@@ -1,27 +1,20 @@
1-
use std::num::TryFromIntError;
1+
use std::{num::TryFromIntError, str::FromStr};
22

33
use snafu::{OptionExt, ResultExt, Snafu};
44
use stackable_operator::{
55
builder::{configmap::ConfigMapBuilder, meta::ObjectMetaBuilder},
66
crd::listener,
77
k8s_openapi::api::core::v1::ConfigMap,
8-
kube::runtime::reflector::ObjectRef,
8+
v2::builder::meta::ownerreference_from_resource,
99
};
1010

1111
use crate::{
12-
controller::ValidatedCluster,
13-
crd::{role::KafkaRole, v1alpha1},
14-
kafka_controller::{KAFKA_CONTROLLER_NAME, build_recommended_labels},
12+
controller::{RoleGroupName, ValidatedCluster},
13+
crd::role::KafkaRole,
1514
};
1615

1716
#[derive(Snafu, Debug)]
1817
pub enum Error {
19-
#[snafu(display("object {} is missing metadata to build owner reference", kafka))]
20-
ObjectMissingMetadataForOwnerRef {
21-
source: stackable_operator::builder::meta::Error,
22-
kafka: ObjectRef<v1alpha1::KafkaCluster>,
23-
},
24-
2518
#[snafu(display("could not find service port with name {}", port_name))]
2619
NoServicePort { port_name: String },
2720

@@ -32,11 +25,6 @@ pub enum Error {
3225
BuildConfigMap {
3326
source: stackable_operator::builder::configmap::Error,
3427
},
35-
36-
#[snafu(display("failed to build metadata"))]
37-
MetadataBuild {
38-
source: stackable_operator::builder::meta::Error,
39-
},
4028
}
4129

4230
/// Build a discovery [`ConfigMap`] containing information about how to connect to a certain
@@ -46,7 +34,6 @@ pub fn build_discovery_configmap(
4634
listeners: &[listener::v1alpha1::Listener],
4735
) -> Result<ConfigMap, Error> {
4836
let kafka_security = &validated_cluster.cluster_config.kafka_security;
49-
let resolved_product_image = &validated_cluster.image;
5037

5138
let port_name = if kafka_security.has_kerberos_enabled() {
5239
kafka_security.bootstrap_port_name()
@@ -65,31 +52,25 @@ pub fn build_discovery_configmap(
6552
.metadata(
6653
ObjectMetaBuilder::new()
6754
.name_and_namespace(validated_cluster)
68-
.ownerreference_from_resource(validated_cluster, None, Some(true))
69-
.with_context(|_| ObjectMissingMetadataForOwnerRefSnafu {
70-
kafka: cluster_object_ref(validated_cluster),
71-
})?
72-
.with_recommended_labels(&build_recommended_labels(
55+
.ownerreference(ownerreference_from_resource(
7356
validated_cluster,
74-
KAFKA_CONTROLLER_NAME,
75-
&resolved_product_image.product_version,
76-
&KafkaRole::Broker.to_string(),
77-
"discovery",
57+
None,
58+
Some(true),
7859
))
79-
.context(MetadataBuildSnafu)?
60+
.with_labels(
61+
validated_cluster.recommended_labels(
62+
&KafkaRole::Broker,
63+
&RoleGroupName::from_str("discovery")
64+
.expect("'discovery' is a valid role group name"),
65+
),
66+
)
8067
.build(),
8168
)
8269
.add_data("KAFKA", bootstrap_servers)
8370
.build()
8471
.context(BuildConfigMapSnafu)
8572
}
8673

87-
/// An [`ObjectRef`] to the owning cluster, built from the validated identity — used only for
88-
/// error context.
89-
fn cluster_object_ref(cluster: &ValidatedCluster) -> ObjectRef<v1alpha1::KafkaCluster> {
90-
ObjectRef::new(cluster.name.as_ref()).within(cluster.namespace.as_ref())
91-
}
92-
9374
fn listener_hosts(
9475
listeners: &[listener::v1alpha1::Listener],
9576
port_name: &str,

rust/operator-binary/src/controller/build/properties/broker_properties.rs

Lines changed: 3 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -4,6 +4,7 @@ use super::kraft_controllers;
44
use crate::{
55
controller::ValidatedClusterConfig,
66
crd::{
7+
KafkaPodDescriptor,
78
listener::{KafkaListenerConfig, KafkaListenerName},
89
role::{
910
KAFKA_ADVERTISED_LISTENERS, KAFKA_BROKER_ID, KAFKA_CONTROLLER_QUORUM_BOOTSTRAP_SERVERS,
@@ -17,9 +18,10 @@ use crate::{
1718
pub fn build(
1819
cluster_config: &ValidatedClusterConfig,
1920
listener_config: &KafkaListenerConfig,
21+
pod_descriptors: &[KafkaPodDescriptor],
2022
overrides: BTreeMap<String, String>,
2123
) -> BTreeMap<String, String> {
22-
let kraft_controllers = kraft_controllers(&cluster_config.pod_descriptors);
24+
let kraft_controllers = kraft_controllers(pod_descriptors);
2325

2426
let mut result = BTreeMap::from([
2527
(

rust/operator-binary/src/controller/build/properties/controller_properties.rs

Lines changed: 3 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -4,6 +4,7 @@ use super::kraft_controllers;
44
use crate::{
55
controller::ValidatedClusterConfig,
66
crd::{
7+
KafkaPodDescriptor,
78
listener::{KafkaListenerConfig, KafkaListenerName},
89
role::{
910
KAFKA_CONTROLLER_QUORUM_BOOTSTRAP_SERVERS, KAFKA_LISTENER_SECURITY_PROTOCOL_MAP,
@@ -16,9 +17,10 @@ use crate::{
1617
pub fn build(
1718
cluster_config: &ValidatedClusterConfig,
1819
listener_config: &KafkaListenerConfig,
20+
pod_descriptors: &[KafkaPodDescriptor],
1921
overrides: BTreeMap<String, String>,
2022
) -> BTreeMap<String, String> {
21-
let kraft_controllers = kraft_controllers(&cluster_config.pod_descriptors).join(",");
23+
let kraft_controllers = kraft_controllers(pod_descriptors).join(",");
2224

2325
let mut result = BTreeMap::from([
2426
(

0 commit comments

Comments
 (0)