Skip to content

Commit 3215d08

Browse files
adwk67claude
andcommitted
feat: remove product-config; merge config/env overrides directly in validate
Replaces the product-config validation path with a ValidatedKafkaCluster that carries, per role group, the merged config plus the config-file, jvm-security, and env overrides resolved directly from the CRD (role <- role-group). Overrides now use stackable_operator::v2::config_overrides::KeyValueConfigOverrides (matching trino/hdfs); the v1 KeyValueOverridesProvider impls and the per-role Configuration impls are removed, and KAFKA_CLUSTER_ID injection moves into the override merge (collect_*_role_group_overrides). The dereferenced authorization config is folded into the validated cluster. Drops the product-config crate dependency (it remains transitive via stackable-operator). The CRD gains `nullable: true` on configOverrides values (v2 allows null to delete a key). Rendered .properties and env vars are unchanged (18 tests pass; byte parity to be confirmed via kuttl). Regenerated extra/crds.yaml and Cargo.nix. Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
1 parent cb850dc commit 3215d08

14 files changed

Lines changed: 314 additions & 290 deletions

File tree

Cargo.lock

Lines changed: 0 additions & 1 deletion
Some generated files are not rendered by default. Learn more about customizing how changed files appear on GitHub.

Cargo.nix

Lines changed: 0 additions & 4 deletions
Some generated files are not rendered by default. Learn more about customizing how changed files appear on GitHub.

Cargo.toml

Lines changed: 0 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -10,7 +10,6 @@ edition = "2021"
1010
repository = "https://github.com/stackabletech/kafka-operator"
1111

1212
[workspace.dependencies]
13-
product-config = { git = "https://github.com/stackabletech/product-config.git", tag = "0.8.0" }
1413
stackable-operator = { git = "https://github.com/stackabletech/operator-rs.git", tag = "stackable-operator-0.111.1", features = ["crds", "webhook"] }
1514

1615
anyhow = "1.0"

extra/crds.yaml

Lines changed: 8 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -541,6 +541,7 @@ spec:
541541
properties:
542542
broker.properties:
543543
additionalProperties:
544+
nullable: true
544545
type: string
545546
description: |-
546547
Flat key-value overrides for `*.properties`, Hadoop XML, etc.
@@ -551,6 +552,7 @@ spec:
551552
type: object
552553
security.properties:
553554
additionalProperties:
555+
nullable: true
554556
type: string
555557
description: |-
556558
Flat key-value overrides for `*.properties`, Hadoop XML, etc.
@@ -1158,6 +1160,7 @@ spec:
11581160
properties:
11591161
broker.properties:
11601162
additionalProperties:
1163+
nullable: true
11611164
type: string
11621165
description: |-
11631166
Flat key-value overrides for `*.properties`, Hadoop XML, etc.
@@ -1168,6 +1171,7 @@ spec:
11681171
type: object
11691172
security.properties:
11701173
additionalProperties:
1174+
nullable: true
11711175
type: string
11721176
description: |-
11731177
Flat key-value overrides for `*.properties`, Hadoop XML, etc.
@@ -1786,6 +1790,7 @@ spec:
17861790
properties:
17871791
controller.properties:
17881792
additionalProperties:
1793+
nullable: true
17891794
type: string
17901795
description: |-
17911796
Flat key-value overrides for `*.properties`, Hadoop XML, etc.
@@ -1796,6 +1801,7 @@ spec:
17961801
type: object
17971802
security.properties:
17981803
additionalProperties:
1804+
nullable: true
17991805
type: string
18001806
description: |-
18011807
Flat key-value overrides for `*.properties`, Hadoop XML, etc.
@@ -2235,6 +2241,7 @@ spec:
22352241
properties:
22362242
controller.properties:
22372243
additionalProperties:
2244+
nullable: true
22382245
type: string
22392246
description: |-
22402247
Flat key-value overrides for `*.properties`, Hadoop XML, etc.
@@ -2245,6 +2252,7 @@ spec:
22452252
type: object
22462253
security.properties:
22472254
additionalProperties:
2255+
nullable: true
22482256
type: string
22492257
description: |-
22502258
Flat key-value overrides for `*.properties`, Hadoop XML, etc.

rust/operator-binary/Cargo.toml

Lines changed: 0 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -9,7 +9,6 @@ repository.workspace = true
99
publish = false
1010

1111
[dependencies]
12-
product-config.workspace = true
1312
stackable-operator.workspace = true
1413

1514
indoc.workspace = true

rust/operator-binary/src/controller.rs

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

55
use const_format::concatcp;
6-
use product_config::ProductConfigManager;
76
use snafu::{ResultExt, Snafu};
87
use stackable_operator::{
98
cli::OperatorEnvironmentOptions,
@@ -51,7 +50,6 @@ pub const KAFKA_FULL_CONTROLLER_NAME: &str = concatcp!(KAFKA_CONTROLLER_NAME, '.
5150

5251
pub struct Ctx {
5352
pub client: stackable_operator::client::Client,
54-
pub product_config: ProductConfigManager,
5553
pub operator_environment: OperatorEnvironmentOptions,
5654
}
5755

@@ -114,9 +112,6 @@ pub enum Error {
114112
source: stackable_operator::cluster_resources::Error,
115113
},
116114

117-
#[snafu(display("failed to resolve and merge config for role and role group"))]
118-
FailedToResolveConfig { source: crate::crd::role::Error },
119-
120115
#[snafu(display("failed to patch service account"))]
121116
ApplyServiceAccount {
122117
source: stackable_operator::cluster_resources::Error,
@@ -156,9 +151,6 @@ pub enum Error {
156151
#[snafu(display("KafkaCluster object is misconfigured"))]
157152
MisconfiguredKafkaCluster { source: crd::Error },
158153

159-
#[snafu(display("failed to parse role: {source}"))]
160-
ParseRole { source: strum::ParseError },
161-
162154
#[snafu(display("failed to build statefulset"))]
163155
BuildStatefulset {
164156
source: crate::resource::statefulset::Error,
@@ -198,7 +190,6 @@ impl ReconcilerError for Error {
198190
Error::ApplyDiscoveryConfig { .. } => None,
199191
Error::DeleteOrphans { .. } => None,
200192
Error::CreateClusterResources { .. } => None,
201-
Error::FailedToResolveConfig { .. } => None,
202193
Error::ApplyServiceAccount { .. } => None,
203194
Error::ApplyRoleBinding { .. } => None,
204195
Error::ApplyStatus { .. } => None,
@@ -207,7 +198,6 @@ impl ReconcilerError for Error {
207198
Error::GetRequiredLabels { .. } => None,
208199
Error::InvalidKafkaCluster { .. } => None,
209200
Error::MisconfiguredKafkaCluster { .. } => None,
210-
Error::ParseRole { .. } => None,
211201
Error::BuildStatefulset { .. } => None,
212202
Error::BuildConfigMap { .. } => None,
213203
Error::BuildService { .. } => None,
@@ -238,18 +228,13 @@ pub async fn reconcile_kafka(
238228
.context(DereferenceSnafu)?;
239229

240230
// validate (no client required)
241-
let validate::ValidatedInputs {
231+
let validate::ValidatedKafkaCluster {
242232
authorization_config,
243233
image,
244234
kafka_security,
245-
role_config: validated_config,
246-
} = validate::validate(
247-
kafka,
248-
dereferenced_objects,
249-
&ctx.operator_environment,
250-
&ctx.product_config,
251-
)
252-
.context(ValidateClusterSnafu)?;
235+
role_groups,
236+
} = validate::validate(kafka, dereferenced_objects, &ctx.operator_environment)
237+
.context(ValidateClusterSnafu)?;
253238

254239
let opa_connect = authorization_config
255240
.as_ref()
@@ -295,15 +280,9 @@ pub async fn reconcile_kafka(
295280

296281
let mut bootstrap_listeners = Vec::<listener::v1alpha1::Listener>::new();
297282

298-
for (kafka_role_str, role_config) in &validated_config {
299-
let kafka_role = KafkaRole::from_str(kafka_role_str).context(ParseRoleSnafu)?;
300-
301-
for (rolegroup_name, rolegroup_config) in role_config.iter() {
302-
let rolegroup_ref = kafka.rolegroup_ref(&kafka_role, rolegroup_name);
303-
304-
let merged_config = kafka_role
305-
.merged_config(kafka, &rolegroup_ref.role_group)
306-
.context(FailedToResolveConfigSnafu)?;
283+
for (kafka_role, rg_map) in &role_groups {
284+
for (rolegroup_name, validated_rg) in rg_map {
285+
let rolegroup_ref = kafka.rolegroup_ref(kafka_role, rolegroup_name);
307286

308287
let rg_headless_service =
309288
build_rolegroup_headless_service(kafka, &image, &rolegroup_ref, &kafka_security)
@@ -333,8 +312,9 @@ pub async fn reconcile_kafka(
333312
&image,
334313
&kafka_security,
335314
&rolegroup_ref,
336-
rolegroup_config,
337-
&merged_config,
315+
validated_rg.config_file_overrides.clone(),
316+
validated_rg.jvm_security_overrides.clone(),
317+
&validated_rg.merged_config,
338318
&kafka_listeners,
339319
&pod_descriptors,
340320
opa_connect.as_deref(),
@@ -344,37 +324,37 @@ pub async fn reconcile_kafka(
344324
let rg_statefulset = match kafka_role {
345325
KafkaRole::Broker => build_broker_rolegroup_statefulset(
346326
kafka,
347-
&kafka_role,
327+
kafka_role,
348328
&image,
349329
&rolegroup_ref,
350-
rolegroup_config,
330+
&validated_rg.env_overrides,
351331
&kafka_security,
352-
&merged_config,
332+
&validated_rg.merged_config,
353333
&rbac_sa,
354334
&client.kubernetes_cluster_info,
355335
)
356336
.context(BuildStatefulsetSnafu)?,
357337
KafkaRole::Controller => build_controller_rolegroup_statefulset(
358338
kafka,
359-
&kafka_role,
339+
kafka_role,
360340
&image,
361341
&rolegroup_ref,
362-
rolegroup_config,
342+
&validated_rg.env_overrides,
363343
&kafka_security,
364-
&merged_config,
344+
&validated_rg.merged_config,
365345
&rbac_sa,
366346
&client.kubernetes_cluster_info,
367347
)
368348
.context(BuildStatefulsetSnafu)?,
369349
};
370350

371-
if let AnyConfig::Broker(broker_config) = merged_config {
351+
if let AnyConfig::Broker(broker_config) = &validated_rg.merged_config {
372352
let rg_bootstrap_listener = build_broker_rolegroup_bootstrap_listener(
373353
kafka,
374354
&image,
375355
&kafka_security,
376356
&rolegroup_ref,
377-
&broker_config,
357+
broker_config,
378358
)
379359
.context(BuildListenerSnafu)?;
380360
bootstrap_listeners.push(
@@ -417,12 +397,12 @@ pub async fn reconcile_kafka(
417397
);
418398
}
419399

420-
let role_config = kafka.role_config(&kafka_role);
400+
let role_cfg = kafka.role_config(kafka_role);
421401
if let Some(GenericRoleConfig {
422402
pod_disruption_budget: pdb,
423-
}) = role_config
403+
}) = role_cfg
424404
{
425-
add_pdbs(pdb, kafka, &kafka_role, client, &mut cluster_resources)
405+
add_pdbs(pdb, kafka, kafka_role, client, &mut cluster_resources)
426406
.await
427407
.context(FailedToCreatePdbSnafu)?;
428408
}

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

Lines changed: 6 additions & 16 deletions
Original file line numberDiff line numberDiff line change
@@ -1,7 +1,6 @@
1-
use std::collections::{BTreeMap, HashMap};
1+
use std::collections::BTreeMap;
22

33
use indoc::formatdoc;
4-
use product_config::types::PropertyNameKind;
54
use snafu::{ResultExt, Snafu};
65
use stackable_operator::{
76
builder::{configmap::ConfigMapBuilder, meta::ObjectMetaBuilder},
@@ -78,7 +77,8 @@ pub fn build_rolegroup_config_map(
7877
resolved_product_image: &ResolvedProductImage,
7978
kafka_security: &KafkaTlsSecurity,
8079
rolegroup: &RoleGroupRef<v1alpha1::KafkaCluster>,
81-
rolegroup_config: &HashMap<PropertyNameKind, BTreeMap<String, String>>,
80+
config_file_overrides: BTreeMap<String, String>,
81+
jvm_security_overrides: BTreeMap<String, String>,
8282
merged_config: &AnyConfig,
8383
listener_config: &KafkaListenerConfig,
8484
pod_descriptors: &[KafkaPodDescriptor],
@@ -90,11 +90,6 @@ pub fn build_rolegroup_config_map(
9090
.effective_metadata_manager()
9191
.context(InvalidMetadataManagerSnafu)?;
9292

93-
let overrides = rolegroup_config
94-
.get(&PropertyNameKind::File(kafka_config_file_name.to_string()))
95-
.cloned()
96-
.unwrap_or_default();
97-
9893
let kafka_config = match merged_config {
9994
AnyConfig::Broker(_) => crate::controller::build::properties::broker_properties::build(
10095
kafka_security,
@@ -107,15 +102,15 @@ pub fn build_rolegroup_config_map(
107102
.cluster_config
108103
.broker_id_pod_config_map_name
109104
.is_some(),
110-
overrides,
105+
config_file_overrides,
111106
),
112107
AnyConfig::Controller(_) => {
113108
crate::controller::build::properties::controller_properties::build(
114109
kafka_security,
115110
listener_config,
116111
pod_descriptors,
117112
metadata_manager == MetadataManager::KRaft,
118-
overrides,
113+
config_file_overrides,
119114
)
120115
}
121116
}
@@ -128,12 +123,7 @@ pub fn build_rolegroup_config_map(
128123
.map(|(k, v)| (k, Some(v)))
129124
.collect::<Vec<_>>();
130125

131-
let jvm_sec_props: BTreeMap<String, Option<String>> = rolegroup_config
132-
.get(&PropertyNameKind::File(
133-
JVM_SECURITY_PROPERTIES_FILE.to_string(),
134-
))
135-
.cloned()
136-
.unwrap_or_default()
126+
let jvm_sec_props: BTreeMap<String, Option<String>> = jvm_security_overrides
137127
.into_iter()
138128
.map(|(k, v)| (k, Some(v)))
139129
.collect();

0 commit comments

Comments
 (0)