Skip to content

Commit 662e9a1

Browse files
committed
thread through name, namespace, uid via ObjectMeta
1 parent 38f278c commit 662e9a1

10 files changed

Lines changed: 181 additions & 60 deletions

File tree

rust/operator-binary/src/controller.rs

Lines changed: 91 additions & 11 deletions
Original file line numberDiff line numberDiff line change
@@ -1,6 +1,6 @@
11
//! Ensures that `Pod`s are configured and running for each [`v1alpha1::KafkaCluster`].
22
3-
use std::{collections::BTreeMap, sync::Arc};
3+
use std::{borrow::Cow, collections::BTreeMap, sync::Arc};
44

55
use const_format::concatcp;
66
use snafu::{ResultExt, Snafu};
@@ -11,7 +11,7 @@ use stackable_operator::{
1111
crd::listener,
1212
kube::{
1313
Resource,
14-
api::DynamicObject,
14+
api::{DynamicObject, ObjectMeta},
1515
core::{DeserializeGuard, error_boundary},
1616
runtime::{controller::Action, reflector::ObjectRef},
1717
},
@@ -22,6 +22,10 @@ use stackable_operator::{
2222
compute_conditions, operations::ClusterOperationsConditionBuilder,
2323
statefulset::StatefulSetConditionBuilder,
2424
},
25+
v2::types::{
26+
kubernetes::{NamespaceName, Uid},
27+
operator::ClusterName,
28+
},
2529
};
2630
use strum::{EnumDiscriminants, IntoStaticStr};
2731

@@ -38,7 +42,6 @@ use crate::{
3842
security::KafkaTlsSecurity,
3943
v1alpha1,
4044
},
41-
discovery::{self, build_discovery_configmap},
4245
operations::pdb::add_pdbs,
4346
resource::{
4447
listener::build_broker_rolegroup_bootstrap_listener,
@@ -94,7 +97,7 @@ pub enum Error {
9497
},
9598

9699
#[snafu(display("failed to build discovery ConfigMap"))]
97-
BuildDiscoveryConfig { source: discovery::Error },
100+
BuildDiscoveryConfig { source: build::discovery::Error },
98101

99102
#[snafu(display("failed to apply discovery ConfigMap"))]
100103
ApplyDiscoveryConfig {
@@ -208,7 +211,17 @@ impl ReconcilerError for Error {
208211

209212
/// The validated cluster. Carries everything the build steps need, resolved once
210213
/// here so downstream code never re-derives it or touches the raw spec.
214+
///
215+
/// The cluster identity (`name`, `namespace`, `uid`) is captured here so that owner
216+
/// references for child objects can be built straight from this struct (via its
217+
/// [`Resource`] impl) without threading the raw [`v1alpha1::KafkaCluster`] around.
218+
/// This mirrors the hive-/opensearch-operator's `ValidatedCluster`.
211219
pub struct ValidatedKafkaCluster {
220+
/// `ObjectMeta` carrying `name`, `namespace` and `uid`, so this struct can act as the
221+
/// owner [`Resource`] for child objects.
222+
metadata: ObjectMeta,
223+
pub name: ClusterName,
224+
pub namespace: NamespaceName,
212225
pub image: ResolvedProductImage,
213226
pub kafka_security: KafkaTlsSecurity,
214227
// DESIGN DECISION: the dereferenced authorization config is folded into the
@@ -222,6 +235,69 @@ pub struct ValidatedKafkaCluster {
222235
pub metadata_manager: MetadataManager,
223236
}
224237

238+
impl ValidatedKafkaCluster {
239+
#[allow(clippy::too_many_arguments)]
240+
pub fn new(
241+
name: ClusterName,
242+
namespace: NamespaceName,
243+
uid: Uid,
244+
image: ResolvedProductImage,
245+
kafka_security: KafkaTlsSecurity,
246+
authorization_config: Option<KafkaAuthorizationConfig>,
247+
role_groups: BTreeMap<KafkaRole, BTreeMap<String, ValidatedRoleGroupConfig>>,
248+
pod_descriptors: Vec<KafkaPodDescriptor>,
249+
metadata_manager: MetadataManager,
250+
) -> Self {
251+
Self {
252+
metadata: ObjectMeta {
253+
name: Some(name.to_string()),
254+
namespace: Some(namespace.to_string()),
255+
uid: Some(uid.to_string()),
256+
..ObjectMeta::default()
257+
},
258+
name,
259+
namespace,
260+
image,
261+
kafka_security,
262+
authorization_config,
263+
role_groups,
264+
pod_descriptors,
265+
metadata_manager,
266+
}
267+
}
268+
}
269+
270+
/// Lets [`ValidatedKafkaCluster`] act as the owner [`Resource`] for child objects, so owner
271+
/// references are built from it (via the captured `metadata`) rather than the raw CR.
272+
impl Resource for ValidatedKafkaCluster {
273+
type DynamicType = <v1alpha1::KafkaCluster as Resource>::DynamicType;
274+
type Scope = <v1alpha1::KafkaCluster as Resource>::Scope;
275+
276+
fn kind(dt: &Self::DynamicType) -> Cow<'_, str> {
277+
v1alpha1::KafkaCluster::kind(dt)
278+
}
279+
280+
fn group(dt: &Self::DynamicType) -> Cow<'_, str> {
281+
v1alpha1::KafkaCluster::group(dt)
282+
}
283+
284+
fn version(dt: &Self::DynamicType) -> Cow<'_, str> {
285+
v1alpha1::KafkaCluster::version(dt)
286+
}
287+
288+
fn plural(dt: &Self::DynamicType) -> Cow<'_, str> {
289+
v1alpha1::KafkaCluster::plural(dt)
290+
}
291+
292+
fn meta(&self) -> &ObjectMeta {
293+
&self.metadata
294+
}
295+
296+
fn meta_mut(&mut self) -> &mut ObjectMeta {
297+
&mut self.metadata
298+
}
299+
}
300+
225301
pub struct ValidatedRoleGroupConfig {
226302
pub merged_config: AnyConfig,
227303
// DESIGN DECISION: overrides are resolved into flat maps HERE rather than stored
@@ -265,7 +341,7 @@ pub async fn reconcile_kafka(
265341
APP_NAME,
266342
OPERATOR_NAME,
267343
KAFKA_CONTROLLER_NAME,
268-
&kafka.object_ref(&()),
344+
&validated_cluster.object_ref(&()),
269345
ClusterResourceApplyStrategy::from(&kafka.spec.cluster_operation),
270346
&kafka.spec.object_overrides,
271347
)
@@ -306,16 +382,19 @@ pub async fn reconcile_kafka(
306382
let rolegroup_ref = kafka.rolegroup_ref(kafka_role, rolegroup_name);
307383

308384
let rg_headless_service = build_rolegroup_headless_service(
309-
kafka,
385+
&validated_cluster,
310386
&validated_cluster.image,
311387
&rolegroup_ref,
312388
&validated_cluster.kafka_security,
313389
)
314390
.context(BuildServiceSnafu)?;
315391

316-
let rg_metrics_service =
317-
build_rolegroup_metrics_service(kafka, &validated_cluster.image, &rolegroup_ref)
318-
.context(BuildServiceSnafu)?;
392+
let rg_metrics_service = build_rolegroup_metrics_service(
393+
&validated_cluster,
394+
&validated_cluster.image,
395+
&rolegroup_ref,
396+
)
397+
.context(BuildServiceSnafu)?;
319398

320399
let kafka_listeners = get_kafka_listener_config(
321400
kafka,
@@ -416,8 +495,9 @@ pub async fn reconcile_kafka(
416495
}
417496
}
418497

419-
let discovery_cm = build_discovery_configmap(kafka, validated_cluster, &bootstrap_listeners)
420-
.context(BuildDiscoveryConfigSnafu)?;
498+
let discovery_cm =
499+
build::discovery::build_discovery_configmap(&validated_cluster, &bootstrap_listeners)
500+
.context(BuildDiscoveryConfigSnafu)?;
421501

422502
cluster_resources
423503
.add(client, discovery_cm)

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

Lines changed: 3 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -121,12 +121,12 @@ pub fn build_rolegroup_config_map(
121121
cm_builder
122122
.metadata(
123123
ObjectMetaBuilder::new()
124-
.name_and_namespace(kafka)
124+
.name_and_namespace(validated_cluster)
125125
.name(rolegroup.object_name())
126-
.ownerreference_from_resource(kafka, None, Some(true))
126+
.ownerreference_from_resource(validated_cluster, None, Some(true))
127127
.context(ObjectMissingMetadataForOwnerRefSnafu)?
128128
.with_recommended_labels(&build_recommended_labels(
129-
kafka,
129+
validated_cluster,
130130
KAFKA_CONTROLLER_NAME,
131131
&resolved_product_image.app_version_label_value,
132132
&rolegroup.role,

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

Lines changed: 13 additions & 11 deletions
Original file line numberDiff line numberDiff line change
@@ -5,7 +5,7 @@ use stackable_operator::{
55
builder::{configmap::ConfigMapBuilder, meta::ObjectMetaBuilder},
66
crd::listener,
77
k8s_openapi::api::core::v1::ConfigMap,
8-
kube::{ResourceExt, runtime::reflector::ObjectRef},
8+
kube::runtime::reflector::ObjectRef,
99
};
1010

1111
use crate::{
@@ -22,9 +22,6 @@ pub enum Error {
2222
kafka: ObjectRef<v1alpha1::KafkaCluster>,
2323
},
2424

25-
#[snafu(display("object has no name associated"))]
26-
NoName,
27-
2825
#[snafu(display("could not find service port with name {}", port_name))]
2926
NoServicePort { port_name: String },
3027

@@ -45,8 +42,7 @@ pub enum Error {
4542
/// Build a discovery [`ConfigMap`] containing information about how to connect to a certain
4643
/// [`v1alpha1::KafkaCluster`].
4744
pub fn build_discovery_configmap(
48-
owner: &v1alpha1::KafkaCluster,
49-
validated_cluster: ValidatedKafkaCluster,
45+
validated_cluster: &ValidatedKafkaCluster,
5046
listeners: &[listener::v1alpha1::Listener],
5147
) -> Result<ConfigMap, Error> {
5248
let kafka_security = &validated_cluster.kafka_security;
@@ -68,14 +64,14 @@ pub fn build_discovery_configmap(
6864
ConfigMapBuilder::new()
6965
.metadata(
7066
ObjectMetaBuilder::new()
71-
.name_and_namespace(owner)
72-
.name(owner.name_unchecked())
73-
.ownerreference_from_resource(owner, None, Some(true))
67+
.name_and_namespace(validated_cluster)
68+
.name(validated_cluster.name.to_string())
69+
.ownerreference_from_resource(validated_cluster, None, Some(true))
7470
.with_context(|_| ObjectMissingMetadataForOwnerRefSnafu {
75-
kafka: ObjectRef::from_obj(owner),
71+
kafka: cluster_object_ref(validated_cluster),
7672
})?
7773
.with_recommended_labels(&build_recommended_labels(
78-
owner,
74+
validated_cluster,
7975
KAFKA_CONTROLLER_NAME,
8076
&resolved_product_image.product_version,
8177
&KafkaRole::Broker.to_string(),
@@ -89,6 +85,12 @@ pub fn build_discovery_configmap(
8985
.context(BuildConfigMapSnafu)
9086
}
9187

88+
/// An [`ObjectRef`] to the owning cluster, built from the validated identity — used only for
89+
/// error context.
90+
fn cluster_object_ref(cluster: &ValidatedKafkaCluster) -> ObjectRef<v1alpha1::KafkaCluster> {
91+
ObjectRef::new(cluster.name.as_ref()).within(cluster.namespace.as_ref())
92+
}
93+
9294
fn listener_hosts(
9395
listeners: &[listener::v1alpha1::Listener],
9496
port_name: &str,
Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1,4 +1,5 @@
11
//! Builders that assemble Kubernetes resources for kafka rolegroups.
22
33
pub mod config_map;
4+
pub mod discovery;
45
pub mod properties;

rust/operator-binary/src/controller/validate.rs

Lines changed: 42 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -3,14 +3,21 @@
33
//! Synchronously validates inputs that don't require a Kubernetes client. Produces
44
//! [`ValidatedKafkaCluster`], consumed by the rest of `reconcile_kafka`.
55
6-
use std::collections::BTreeMap;
6+
use std::{collections::BTreeMap, str::FromStr};
77

8-
use snafu::{ResultExt, Snafu};
8+
use snafu::{OptionExt, ResultExt, Snafu};
99
use stackable_operator::{
1010
cli::OperatorEnvironmentOptions,
1111
commons::product_image_selection,
1212
config::merge::{Merge, merge},
13-
v2::config_overrides::KeyValueConfigOverrides,
13+
kube::ResourceExt,
14+
v2::{
15+
config_overrides::KeyValueConfigOverrides,
16+
types::{
17+
kubernetes::{NamespaceName, Uid},
18+
operator::ClusterName,
19+
},
20+
},
1421
};
1522

1623
use crate::{
@@ -50,6 +57,27 @@ pub enum Error {
5057

5158
#[snafu(display("invalid metadata manager"))]
5259
InvalidMetadataManager { source: crate::crd::Error },
60+
61+
#[snafu(display("invalid cluster name"))]
62+
InvalidClusterName {
63+
source: stackable_operator::v2::macros::attributed_string_type::Error,
64+
},
65+
66+
#[snafu(display("object defines no namespace"))]
67+
ObjectHasNoNamespace,
68+
69+
#[snafu(display("invalid cluster namespace"))]
70+
InvalidNamespace {
71+
source: stackable_operator::v2::macros::attributed_string_type::Error,
72+
},
73+
74+
#[snafu(display("object has no uid"))]
75+
ObjectHasNoUid,
76+
77+
#[snafu(display("invalid cluster uid"))]
78+
InvalidUid {
79+
source: stackable_operator::v2::macros::attributed_string_type::Error,
80+
},
5381
}
5482

5583
type Result<T, E = Error> = std::result::Result<T, E>;
@@ -164,14 +192,22 @@ pub fn validate(
164192
.effective_metadata_manager()
165193
.context(InvalidMetadataManagerSnafu)?;
166194

167-
Ok(ValidatedKafkaCluster {
195+
let name = ClusterName::from_str(&kafka.name_any()).context(InvalidClusterNameSnafu)?;
196+
let namespace = NamespaceName::from_str(&kafka.namespace().context(ObjectHasNoNamespaceSnafu)?)
197+
.context(InvalidNamespaceSnafu)?;
198+
let uid = Uid::from_str(&kafka.uid().context(ObjectHasNoUidSnafu)?).context(InvalidUidSnafu)?;
199+
200+
Ok(ValidatedKafkaCluster::new(
201+
name,
202+
namespace,
203+
uid,
168204
image,
169205
kafka_security,
170-
authorization_config: dereferenced_objects.authorization_config,
206+
dereferenced_objects.authorization_config,
171207
role_groups,
172208
pod_descriptors,
173209
metadata_manager,
174-
})
210+
))
175211
}
176212

177213
/// Merge role-group overrides over the role-level overrides (role-group wins per key) via the

rust/operator-binary/src/main.rs

Lines changed: 0 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -43,7 +43,6 @@ use crate::{
4343
mod config;
4444
mod controller;
4545
mod crd;
46-
mod discovery;
4746
mod kerberos;
4847
mod operations;
4948
mod product_logging;

rust/operator-binary/src/resource/listener.rs

Lines changed: 3 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -36,12 +36,12 @@ pub fn build_broker_rolegroup_bootstrap_listener(
3636

3737
Ok(listener::v1alpha1::Listener {
3838
metadata: ObjectMetaBuilder::new()
39-
.name_and_namespace(kafka)
39+
.name_and_namespace(validated_cluster)
4040
.name(kafka.bootstrap_service_name(rolegroup))
41-
.ownerreference_from_resource(kafka, None, Some(true))
41+
.ownerreference_from_resource(validated_cluster, None, Some(true))
4242
.context(ObjectMissingMetadataForOwnerRefSnafu)?
4343
.with_recommended_labels(&build_recommended_labels(
44-
kafka,
44+
validated_cluster,
4545
KAFKA_CONTROLLER_NAME,
4646
&resolved_product_image.app_version_label_value,
4747
&rolegroup.role,

0 commit comments

Comments
 (0)