From 92aea5edb0855432cafa3ecea05eecc691c58c6f Mon Sep 17 00:00:00 2001 From: Marc-Merino <35137847+Marc-Merino@users.noreply.github.com> Date: Mon, 28 Sep 2026 11:49:30 -0500 Subject: [PATCH 1/2] feat: Add default affinity to OPA server Pods Knit-Group: kg_20260928_53b541 Knit-Bundle: opa-client-pod-affinity --- CHANGELOG.md | 3 + .../usage-guide/operations/pod-placement.adoc | 10 +++ .../src/controller/validate.rs | 6 +- rust/operator-binary/src/crd/affinity.rs | 83 ++++++++++++++++--- rust/operator-binary/src/crd/role/broker.rs | 17 ++-- rust/operator-binary/src/crd/role/commons.rs | 10 ++- .../src/crd/role/controller.rs | 2 +- 7 files changed, 111 insertions(+), 20 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 9cdeb3b9c..fe83796d3 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -6,6 +6,8 @@ All notable changes to this project will be documented in this file. ### Added +- Prefer scheduling Kafka broker Pods on nodes with OPA server Pods when OPA authorization is + configured ([#XXX]). - Support floating tags for product images via the new `spec.image.stackableVersionPolicy` field ([#1021]). - Support Kerberos (GSSAPI) authentication on KRaft controllers, covering both @@ -79,6 +81,7 @@ All notable changes to this project will be documented in this file. [#1021]: https://github.com/stackabletech/kafka-operator/pull/1021 [#1022]: https://github.com/stackabletech/kafka-operator/pull/1022 [#1024]: https://github.com/stackabletech/kafka-operator/pull/1024 +[#XXX]: https://github.com/stackabletech/kafka-operator/pull/XXX ## [26.7.0] - 2026-07-21 diff --git a/docs/modules/kafka/pages/usage-guide/operations/pod-placement.adoc b/docs/modules/kafka/pages/usage-guide/operations/pod-placement.adoc index 97a0d19c0..4202c1980 100644 --- a/docs/modules/kafka/pages/usage-guide/operations/pod-placement.adoc +++ b/docs/modules/kafka/pages/usage-guide/operations/pod-placement.adoc @@ -20,3 +20,13 @@ affinity: ---- In the example above `cluster-name` is the name of the Kafka custom resource that owns this Pod. + +When OPA authorization is configured, Kafka brokers also prefer to run on the same node as the +OPA server Pods. This preferred Pod affinity has weight `50` and uses +`kubernetes.io/hostname` as its topology key. It matches the OPA server Pods by +`app.kubernetes.io/name=opa`, `app.kubernetes.io/instance=`, and +`app.kubernetes.io/component=server`. +The configured OPA discovery ConfigMap is assumed to have the same name as the OpaCluster and +to be in the KafkaCluster's namespace. +This is a scheduling preference, not a guarantee of co-location, and it does not change how +Kafka brokers route requests to the OPA service. diff --git a/rust/operator-binary/src/controller/validate.rs b/rust/operator-binary/src/controller/validate.rs index dd3bafb53..87bb7a067 100644 --- a/rust/operator-binary/src/controller/validate.rs +++ b/rust/operator-binary/src/controller/validate.rs @@ -260,7 +260,11 @@ pub fn validate( let broker_role = &kafka.spec.brokers; let broker_groups = validate_role_group_configs( broker_role, - BrokerConfig::default_config(&kafka.name_any(), &KafkaRole::Broker.to_string()), + BrokerConfig::default_config( + &kafka.name_any(), + &KafkaRole::Broker.to_string(), + kafka.spec.cluster_config.authorization.opa.as_ref(), + ), cluster_id, AnyConfig::Broker, AnyConfigOverrides::Broker, diff --git a/rust/operator-binary/src/crd/affinity.rs b/rust/operator-binary/src/crd/affinity.rs index 152c8e587..f6f4862f2 100644 --- a/rust/operator-binary/src/crd/affinity.rs +++ b/rust/operator-binary/src/crd/affinity.rs @@ -1,13 +1,31 @@ use stackable_operator::{ - commons::affinity::{StackableAffinityFragment, affinity_between_role_pods}, - k8s_openapi::api::core::v1::PodAntiAffinity, + commons::{ + affinity::{StackableAffinityFragment, affinity_between_role_pods}, + opa::OpaConfig, + }, + k8s_openapi::api::core::v1::{PodAffinity, PodAntiAffinity}, }; use crate::crd::APP_NAME; -pub fn get_affinity(cluster_name: &str, role: &str) -> StackableAffinityFragment { +pub fn get_affinity( + cluster_name: &str, + role: &str, + opa_config: Option<&OpaConfig>, +) -> StackableAffinityFragment { + // Only brokers use the OPA authorizer. The discovery ConfigMap has the same name as the + // OpaCluster, and Kafka resolves it in its own namespace. + let pod_affinity = opa_config + .filter(|_| role == "broker") + .map(|opa_config| PodAffinity { + preferred_during_scheduling_ignored_during_execution: Some(vec![ + affinity_between_role_pods("opa", &opa_config.config_map_name, "server", 50), + ]), + required_during_scheduling_ignored_during_execution: None, + }); + StackableAffinityFragment { - pod_affinity: None, + pod_affinity, pod_anti_affinity: Some(PodAntiAffinity { preferred_during_scheduling_ignored_during_execution: Some(vec![ affinity_between_role_pods(APP_NAME, cluster_name, role, 70), @@ -27,7 +45,9 @@ mod tests { use stackable_operator::{ commons::affinity::StackableAffinity, k8s_openapi::{ - api::core::v1::{PodAffinityTerm, PodAntiAffinity, WeightedPodAffinityTerm}, + api::core::v1::{ + PodAffinity, PodAffinityTerm, PodAntiAffinity, WeightedPodAffinityTerm, + }, apimachinery::pkg::apis::meta::v1::LabelSelector, }, }; @@ -38,8 +58,9 @@ mod tests { }; #[rstest] - #[case(KafkaRole::Broker)] - fn test_affinity_defaults(#[case] role: KafkaRole) { + #[case(false)] + #[case(true)] + fn test_affinity_defaults(#[case] with_opa: bool) { let input = r#" apiVersion: kafka.stackable.tech/v1alpha1 kind: KafkaCluster @@ -58,11 +79,19 @@ mod tests { replicas: 1 "#; - let kafka = minimal_kafka(input); + let input = if with_opa { + input.replace( + " zookeeperConfigMapName: xyz", + " zookeeperConfigMapName: xyz\n authorization:\n opa:\n configMapName: test-opa", + ) + } else { + input.to_string() + }; + let kafka = minimal_kafka(&input); let validated = validated_cluster(&kafka); let merged_config = validated .role_group_configs - .get(&role) + .get(&KafkaRole::Broker) .and_then(|groups| groups.get(&"default".parse().unwrap())) .map(|rg| &rg.config.config) .expect("role group should exist"); @@ -70,7 +99,32 @@ mod tests { assert_eq!( merged_config.affinity, StackableAffinity { - pod_affinity: None, + pod_affinity: with_opa.then_some(PodAffinity { + preferred_during_scheduling_ignored_during_execution: Some(vec![ + WeightedPodAffinityTerm { + pod_affinity_term: PodAffinityTerm { + label_selector: Some(LabelSelector { + match_expressions: None, + match_labels: Some(BTreeMap::from([ + ("app.kubernetes.io/name".to_string(), "opa".to_string()), + ( + "app.kubernetes.io/instance".to_string(), + "test-opa".to_string(), + ), + ( + "app.kubernetes.io/component".to_string(), + "server".to_string(), + ), + ])), + }), + topology_key: "kubernetes.io/hostname".to_string(), + ..PodAffinityTerm::default() + }, + weight: 50, + }, + ]), + required_during_scheduling_ignored_during_execution: None, + }), pod_anti_affinity: Some(PodAntiAffinity { preferred_during_scheduling_ignored_during_execution: Some(vec![ WeightedPodAffinityTerm { @@ -102,4 +156,13 @@ mod tests { } ); } + + #[test] + fn controller_has_no_opa_affinity() { + let opa_config = serde_yaml::from_str("configMapName: test-opa").unwrap(); + assert_eq!( + super::get_affinity("simple-kafka", "controller", Some(&opa_config)).pod_affinity, + None + ); + } } diff --git a/rust/operator-binary/src/crd/role/broker.rs b/rust/operator-binary/src/crd/role/broker.rs index 6b50930e6..d9c499398 100644 --- a/rust/operator-binary/src/crd/role/broker.rs +++ b/rust/operator-binary/src/crd/role/broker.rs @@ -2,9 +2,12 @@ use std::str::FromStr; use serde::{Deserialize, Serialize}; use stackable_operator::{ - commons::resources::{ - CpuLimitsFragment, MemoryLimitsFragment, NoRuntimeLimits, NoRuntimeLimitsFragment, - PvcConfigFragment, Resources, ResourcesFragment, + commons::{ + opa::OpaConfig, + resources::{ + CpuLimitsFragment, MemoryLimitsFragment, NoRuntimeLimits, NoRuntimeLimitsFragment, + PvcConfigFragment, Resources, ResourcesFragment, + }, }, config::{fragment::Fragment, merge::Merge}, constant, @@ -87,9 +90,13 @@ pub struct BrokerConfig { } impl BrokerConfig { - pub fn default_config(cluster_name: &str, role: &str) -> BrokerConfigFragment { + pub fn default_config( + cluster_name: &str, + role: &str, + opa_config: Option<&OpaConfig>, + ) -> BrokerConfigFragment { BrokerConfigFragment { - common_config: CommonConfig::default_config(cluster_name, role), + common_config: CommonConfig::default_config(cluster_name, role, opa_config), bootstrap_listener_class: Some(DEFAULT_LISTENER_CLASS.clone()), broker_listener_class: Some(DEFAULT_LISTENER_CLASS.clone()), logging: product_logging::spec::default_logging(), diff --git a/rust/operator-binary/src/crd/role/commons.rs b/rust/operator-binary/src/crd/role/commons.rs index 32ac469f0..2610fe7d3 100644 --- a/rust/operator-binary/src/crd/role/commons.rs +++ b/rust/operator-binary/src/crd/role/commons.rs @@ -1,6 +1,6 @@ use serde::{Deserialize, Serialize}; use stackable_operator::{ - commons::{affinity::StackableAffinity, resources::PvcConfig}, + commons::{affinity::StackableAffinity, opa::OpaConfig, resources::PvcConfig}, config::{fragment::Fragment, merge::Merge}, k8s_openapi::api::core::v1::PersistentVolumeClaim, schemars::{self, JsonSchema}, @@ -70,9 +70,13 @@ impl CommonConfig { // Auto TLS certificate lifetime const DEFAULT_SECRET_LIFETIME: Duration = Duration::from_days_unchecked(1); - pub fn default_config(cluster_name: &str, role: &str) -> CommonConfigFragment { + pub fn default_config( + cluster_name: &str, + role: &str, + opa_config: Option<&OpaConfig>, + ) -> CommonConfigFragment { CommonConfigFragment { - affinity: get_affinity(cluster_name, role), + affinity: get_affinity(cluster_name, role, opa_config), graceful_shutdown_timeout: Some(Self::DEFAULT_GRACEFUL_SHUTDOWN_TIMEOUT), requested_secret_lifetime: Some(Self::DEFAULT_SECRET_LIFETIME), } diff --git a/rust/operator-binary/src/crd/role/controller.rs b/rust/operator-binary/src/crd/role/controller.rs index 39c98e1ef..2afb9e375 100644 --- a/rust/operator-binary/src/crd/role/controller.rs +++ b/rust/operator-binary/src/crd/role/controller.rs @@ -80,7 +80,7 @@ pub struct ControllerConfig { impl ControllerConfig { pub fn default_config(cluster_name: &str, role: &str) -> ControllerConfigFragment { ControllerConfigFragment { - common_config: CommonConfig::default_config(cluster_name, role), + common_config: CommonConfig::default_config(cluster_name, role, None), logging: product_logging::spec::default_logging(), resources: ResourcesFragment { cpu: CpuLimitsFragment { From f5966cc4f9f683488f10d7d51064fadcd2fb0f47 Mon Sep 17 00:00:00 2001 From: Marc-Merino <35137847+Marc-Merino@users.noreply.github.com> Date: Mon, 28 Sep 2026 12:06:34 -0500 Subject: [PATCH 2/2] docs: Link the pull request in the changelog Knit-Group: kg_20260928_ef2936 Knit-Bundle: opa-client-pod-affinity --- CHANGELOG.md | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index fe83796d3..f32553685 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -7,7 +7,7 @@ All notable changes to this project will be documented in this file. ### Added - Prefer scheduling Kafka broker Pods on nodes with OPA server Pods when OPA authorization is - configured ([#XXX]). + configured ([#1031]). - Support floating tags for product images via the new `spec.image.stackableVersionPolicy` field ([#1021]). - Support Kerberos (GSSAPI) authentication on KRaft controllers, covering both @@ -81,7 +81,7 @@ All notable changes to this project will be documented in this file. [#1021]: https://github.com/stackabletech/kafka-operator/pull/1021 [#1022]: https://github.com/stackabletech/kafka-operator/pull/1022 [#1024]: https://github.com/stackabletech/kafka-operator/pull/1024 -[#XXX]: https://github.com/stackabletech/kafka-operator/pull/XXX +[#1031]: https://github.com/stackabletech/kafka-operator/pull/1031 ## [26.7.0] - 2026-07-21