diff --git a/CHANGELOG.md b/CHANGELOG.md index 3754173e7..8e7fd9f44 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -10,6 +10,7 @@ via `spec.webservers.roleConfig.trustedProxies` ([#835]). - Added airflow `3.3.1` ([#865]). - Add `/ready` endpoint to the operator Deployment, which reports the CRD installation status ([#868]). +- Webservers now have a default affinity to the OPA Pods when OPA authorization is configured ([#872]). ### Changed @@ -81,6 +82,7 @@ [#865]: https://github.com/stackabletech/airflow-operator/pull/865 [#868]: https://github.com/stackabletech/airflow-operator/pull/868 [#870]: https://github.com/stackabletech/airflow-operator/pull/870 +[#872]: https://github.com/stackabletech/airflow-operator/pull/872 ## [26.7.0] - 2026-07-21 diff --git a/docs/modules/airflow/pages/usage-guide/operations/pod-placement.adoc b/docs/modules/airflow/pages/usage-guide/operations/pod-placement.adoc index d21dd30b7..a0fa39e72 100644 --- a/docs/modules/airflow/pages/usage-guide/operations/pod-placement.adoc +++ b/docs/modules/airflow/pages/usage-guide/operations/pod-placement.adoc @@ -6,6 +6,8 @@ The default affinities created by the operator are: 1. Co-locate all the Airflow Pods (weight 20) 2. Distribute all Pods within the same role (worker, webserver, scheduler) (weight 70) +3. If OPA authorization is configured: co-locate the webservers with the OPA Pods (weight 50). + Only the webservers are configured with the OPA auth manager and thus have the affinity. == Kubernetes executors diff --git a/rust/operator-binary/src/controller/validate.rs b/rust/operator-binary/src/controller/validate.rs index e58b495d5..9147c5fdf 100644 --- a/rust/operator-binary/src/controller/validate.rs +++ b/rust/operator-binary/src/controller/validate.rs @@ -148,7 +148,8 @@ pub fn validate_cluster( }, ); - let default_config = AirflowConfig::default_config(&airflow.name_any(), &role); + let default_config = + AirflowConfig::default_config(&airflow.name_any(), &role, airflow.get_opa_config()); let mut group_configs = BTreeMap::new(); for (rolegroup_name, rolegroup) in &resolved_role.role_groups { @@ -431,7 +432,8 @@ mod tests { let role = cluster .get_role(&AirflowRole::Webserver) .expect("webserver role"); - let default_config = AirflowConfig::default_config("airflow", &AirflowRole::Webserver); + let default_config = + AirflowConfig::default_config("airflow", &AirflowRole::Webserver, None); let rolegroup = role.role_groups.get("default").expect("default role group"); let validated = validate_role_group( @@ -530,7 +532,8 @@ mod tests { let role = cluster .get_role(&AirflowRole::Scheduler) .expect("scheduler role"); - let default_config = AirflowConfig::default_config("airflow", &AirflowRole::Scheduler); + let default_config = + AirflowConfig::default_config("airflow", &AirflowRole::Scheduler, None); let rolegroup = role.role_groups.get("default").expect("default role group"); let validated = validate_role_group( @@ -607,7 +610,8 @@ mod tests { let role = cluster .get_role(&AirflowRole::Webserver) .expect("webserver role"); - let default_config = AirflowConfig::default_config("airflow", &AirflowRole::Webserver); + let default_config = + AirflowConfig::default_config("airflow", &AirflowRole::Webserver, None); let rolegroup = role.role_groups.get("default").expect("default role group"); let validated = validate_role_group( diff --git a/rust/operator-binary/src/crd/affinity.rs b/rust/operator-binary/src/crd/affinity.rs index aadcae904..bc5f1e3c7 100644 --- a/rust/operator-binary/src/crd/affinity.rs +++ b/rust/operator-binary/src/crd/affinity.rs @@ -1,6 +1,9 @@ use stackable_operator::{ - commons::affinity::{ - StackableAffinityFragment, affinity_between_cluster_pods, affinity_between_role_pods, + commons::{ + affinity::{ + StackableAffinityFragment, affinity_between_cluster_pods, affinity_between_role_pods, + }, + opa::OpaConfig, }, k8s_openapi::api::core::v1::{PodAffinity, PodAntiAffinity}, }; @@ -8,24 +11,51 @@ use stackable_operator::{ use crate::crd::{APP_NAME, AirflowRole}; /// Used for all [`AirflowRole`]s besides executors. -pub fn get_affinity(cluster_name: &str, role: &AirflowRole) -> StackableAffinityFragment { - get_affinity_for_role(cluster_name, &role.to_string()) +pub fn get_affinity( + cluster_name: &str, + role: &AirflowRole, + opa_config: Option<&OpaConfig>, +) -> StackableAffinityFragment { + let opa_config = match role { + // Only the webserver is configured with the OPA auth manager, so only the webserver is + // co-located with the OPA Pods. + AirflowRole::Webserver => opa_config, + AirflowRole::Scheduler + | AirflowRole::Worker + | AirflowRole::DagProcessor + | AirflowRole::Triggerer => None, + }; + get_affinity_for_role(cluster_name, &role.to_string(), opa_config) } /// There is no [`AirflowRole`] for executors (only for workers), so let's have a special case here. pub fn get_executor_affinity(cluster_name: &str) -> StackableAffinityFragment { - get_affinity_for_role(cluster_name, "executor") + get_affinity_for_role(cluster_name, "executor", None) } -fn get_affinity_for_role(cluster_name: &str, role: &str) -> StackableAffinityFragment { +fn get_affinity_for_role( + cluster_name: &str, + role: &str, + opa_config: Option<&OpaConfig>, +) -> StackableAffinityFragment { + // Built before the `let`s below, which shadow the helper functions with their results. + let affinity_to_opa_pods = opa_config.map(|opa_config| { + affinity_between_role_pods( + "opa", + &opa_config.config_map_name, // The discovery cm has the same name as the OpaCluster itself + "server", + 50, + ) + }); let affinity_between_cluster_pods = affinity_between_cluster_pods(APP_NAME, cluster_name, 20); let affinity_between_role_pods = affinity_between_role_pods(APP_NAME, cluster_name, role, 70); + let mut pod_affinities = vec![affinity_between_cluster_pods]; + pod_affinities.extend(affinity_to_opa_pods); + StackableAffinityFragment { pod_affinity: Some(PodAffinity { - preferred_during_scheduling_ignored_during_execution: Some(vec![ - affinity_between_cluster_pods, - ]), + preferred_during_scheduling_ignored_during_execution: Some(pod_affinities), required_during_scheduling_ignored_during_execution: None, }), pod_anti_affinity: Some(PodAntiAffinity { @@ -90,6 +120,10 @@ mod tests { redis: host: airflow-redis-master credentialsSecretName: airflow-redis-credentials + authorization: + opa: + configMapName: simple-opa + package: airflow webservers: roleGroups: default: @@ -111,36 +145,61 @@ mod tests { let resolved_role = airflow .get_role(&role) .expect("the role is defined in the test cluster"); - let default_config = AirflowConfig::default_config(&airflow.name_any(), &role); + let default_config = + AirflowConfig::default_config(&airflow.name_any(), &role, airflow.get_opa_config()); let rolegroup = resolved_role .role_groups .get("default") .expect("the 'default' role group is defined in the test cluster"); + let mut expected_pod_affinities = vec![WeightedPodAffinityTerm { + pod_affinity_term: PodAffinityTerm { + label_selector: Some(LabelSelector { + match_expressions: None, + match_labels: Some(BTreeMap::from([ + ("app.kubernetes.io/name".to_string(), "airflow".to_string()), + ( + "app.kubernetes.io/instance".to_string(), + "airflow".to_string(), + ), + ])), + }), + topology_key: "kubernetes.io/hostname".to_string(), + ..PodAffinityTerm::default() + }, + weight: 20, + }]; + // Only the webserver is configured with the OPA auth manager. + if role == AirflowRole::Webserver { + expected_pod_affinities.push(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(), + "simple-opa".to_string(), + ), + ( + "app.kubernetes.io/component".to_string(), + "server".to_string(), + ), + ])), + }), + topology_key: "kubernetes.io/hostname".to_string(), + ..PodAffinityTerm::default() + }, + weight: 50, + }); + } + let expected: StackableAffinity = StackableAffinity { node_affinity: None, node_selector: None, pod_affinity: Some(PodAffinity { required_during_scheduling_ignored_during_execution: None, - 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(), "airflow".to_string()), - ( - "app.kubernetes.io/instance".to_string(), - "airflow".to_string(), - ), - ])), - }), - topology_key: "kubernetes.io/hostname".to_string(), - ..PodAffinityTerm::default() - }, - weight: 20, - }, - ]), + preferred_during_scheduling_ignored_during_execution: Some(expected_pod_affinities), }), pod_anti_affinity: Some(PodAntiAffinity { required_during_scheduling_ignored_during_execution: None, diff --git a/rust/operator-binary/src/crd/mod.rs b/rust/operator-binary/src/crd/mod.rs index 781ddf10b..b4b60bb83 100644 --- a/rust/operator-binary/src/crd/mod.rs +++ b/rust/operator-binary/src/crd/mod.rs @@ -415,6 +415,16 @@ impl HasStatusCondition for v1alpha2::AirflowCluster { } impl v1alpha2::AirflowCluster { + /// The OPA config, if OPA authorization is configured. + pub fn get_opa_config(&self) -> Option<&OpaConfig> { + self.spec + .cluster_config + .authorization + .as_ref() + .and_then(|authorization| authorization.opa.as_ref()) + .map(|opa| &opa.opa) + } + /// The name of the group-listener provided for a specific role. /// Webservers will use this group listener so that only one load balancer /// is needed for that role. @@ -967,11 +977,15 @@ pub struct AirflowConfig { } impl AirflowConfig { - pub(crate) fn default_config(cluster_name: &str, role: &AirflowRole) -> AirflowConfigFragment { + pub(crate) fn default_config( + cluster_name: &str, + role: &AirflowRole, + opa_config: Option<&OpaConfig>, + ) -> AirflowConfigFragment { AirflowConfigFragment { resources: default_resources(role), logging: product_logging::spec::default_logging(), - affinity: get_affinity(cluster_name, role), + affinity: get_affinity(cluster_name, role, opa_config), graceful_shutdown_timeout: Some(match role { AirflowRole::Webserver | AirflowRole::Scheduler