From 11240ab222778e99c31721a7a8e28d5c26fa0f72 Mon Sep 17 00:00:00 2001 From: Andrew Kenworthy Date: Wed, 23 Sep 2026 16:39:08 +0200 Subject: [PATCH 1/2] refactor: Flatten role configs --- .../src/controller/build/mod.rs | 69 +++-- .../controller/build/properties/env_vars.rs | 11 +- .../controller/build/resource/statefulset.rs | 8 +- rust/operator-binary/src/controller/mod.rs | 243 ++++++++++++++++-- .../src/controller/validate.rs | 239 ++++++++++++----- rust/operator-binary/src/crd/mod.rs | 162 +----------- .../src/crd/trusted_proxies.rs | 21 ++ 7 files changed, 475 insertions(+), 278 deletions(-) diff --git a/rust/operator-binary/src/controller/build/mod.rs b/rust/operator-binary/src/controller/build/mod.rs index e8b3410e..e5aae8d9 100644 --- a/rust/operator-binary/src/controller/build/mod.rs +++ b/rust/operator-binary/src/controller/build/mod.rs @@ -27,7 +27,7 @@ use crate::{ statefulset::build_server_rolegroup_statefulset, }, }, - crd::{AirflowConfigOverrides, Container}, + crd::{AirflowConfigOverrides, AirflowRole, Container}, }; pub mod graceful_shutdown; @@ -93,23 +93,27 @@ pub fn build(cluster: &ValidatedCluster) -> Result config_maps.push(executor_template_config_map); } - for (role, role_group_configs) in &cluster.role_groups { - if let Some(role_config) = cluster.role_configs.get(role) { - if let Some(pdb_config) = &role_config.pdb { - pod_disruption_budgets.extend(build_pdb(pdb_config, cluster, role)); - } - if let Some(listener_class) = &role_config.listener_class - && let Some(group_listener_name) = &role_config.group_listener_name - { - listeners.push(build_group_listener( - cluster, - role, - listener_class.clone(), - group_listener_name.clone(), - )); - } - } - + // One entry per role, in `AirflowRole` declaration order. Each role's groups come from its own + // field, so the role and its groups cannot be paired up wrongly here. + for (role, role_group_configs) in [ + ( + &AirflowRole::Webserver, + &cluster.webserver_role_group_configs, + ), + ( + &AirflowRole::Scheduler, + &cluster.scheduler_role_group_configs, + ), + (&AirflowRole::Worker, &cluster.worker_role_group_configs), + ( + &AirflowRole::DagProcessor, + &cluster.dagprocessor_role_group_configs, + ), + ( + &AirflowRole::Triggerer, + &cluster.triggerer_role_group_configs, + ), + ] { for (role_group_name, rg_config) in role_group_configs { let logging = &rg_config.config.logging; @@ -149,6 +153,21 @@ pub fn build(cluster: &ValidatedCluster) -> Result })?, ); } + + if let Some(pdb) = cluster.pdb(role) { + pod_disruption_budgets.extend(build_pdb(pdb, cluster, role)); + } + } + + // Only the webserver serves the web UI, so it is the only role with a group listener; it is + // built once here rather than inside the role loop. + if let Some(webserver) = &cluster.webserver_config { + listeners.push(build_group_listener( + cluster, + &AirflowRole::Webserver, + webserver.listener_class.clone(), + webserver.group_listener_name.clone(), + )); } Ok(KubernetesResources { @@ -335,7 +354,14 @@ pub(crate) mod test_support { serde_yaml::with::singleton_map_recursive::deserialize(cluster_value) .expect("the test CR deserialises"); - let dereferenced = DereferencedObjects { + validate_cluster(&cluster, "oci.stackable.tech/sdp", dereferenced_objects()) + .expect("test cluster validates") + } + + /// The resolved external objects a test cluster is validated against: no AuthenticationClasses + /// and no OPA authorization, so validation depends on the CR alone. + pub fn dereferenced_objects() -> DereferencedObjects { + DereferencedObjects { authentication_config: AirflowClientAuthenticationDetailsResolved { authentication_classes_resolved: vec![], user_registration: true, @@ -343,10 +369,7 @@ pub(crate) mod test_support { sync_roles_at: FlaskRolesSyncMoment::default(), }, authorization_config: AirflowAuthorizationResolved { opa: None }, - }; - - validate_cluster(&cluster, "oci.stackable.tech/sdp", dereferenced) - .expect("test cluster validates") + } } /// A Celery-executor cluster whose webserver trusts the given reverse proxies. diff --git a/rust/operator-binary/src/controller/build/properties/env_vars.rs b/rust/operator-binary/src/controller/build/properties/env_vars.rs index 02361145..7b7151d9 100644 --- a/rust/operator-binary/src/controller/build/properties/env_vars.rs +++ b/rust/operator-binary/src/controller/build/properties/env_vars.rs @@ -320,10 +320,7 @@ fn add_version_specific_env_vars( // behind a reverse proxy. // This covers the uvicorn backend, the only one the SDP image can run. let trusted_proxies = cluster - .role_configs - .get(airflow_role) - .map(|role_config| role_config.trusted_proxies.as_slice()) - .unwrap_or_default() + .trusted_proxies(airflow_role) .iter() .map(TrustedProxy::to_string) .collect::>() @@ -342,11 +339,7 @@ fn add_version_specific_env_vars( // The 2.x uses Werkzeug's `ProxyFix` to allow forwarded-headers and it does so regardless // of the peer source. The only valid value for `spec.webservers.roleConfig.trustedProxies` is `["*"]`. if airflow_role == &AirflowRole::Webserver { - let trusted_proxies = cluster - .role_configs - .get(airflow_role) - .map(|role_config| role_config.trusted_proxies.as_slice()) - .unwrap_or_default(); + let trusted_proxies = cluster.trusted_proxies(airflow_role); if !trusted_proxies.is_empty() { env_vars = env_vars diff --git a/rust/operator-binary/src/controller/build/resource/statefulset.rs b/rust/operator-binary/src/controller/build/resource/statefulset.rs index 9e7f4488..d25d61ef 100644 --- a/rust/operator-binary/src/controller/build/resource/statefulset.rs +++ b/rust/operator-binary/src/controller/build/resource/statefulset.rs @@ -221,17 +221,13 @@ pub fn build_server_rolegroup_statefulset( let mut pvcs: Option> = None; - if let Some(listener_group_name) = validated_cluster - .role_configs - .get(airflow_role) - .and_then(|role_config| role_config.group_listener_name.clone()) - { + if let Some(listener_group_name) = validated_cluster.group_listener_name(airflow_role) { // Listener endpoints for the Webserver role will use persistent volumes // so that load balancers can hard-code the target addresses. This will // be the case even when no class is set (and the value defaults to // cluster-internal) as the address should still be consistent. let pvc = listener_operator_volume_source_builder_build_pvc( - &ListenerReference::Listener(listener_group_name), + &ListenerReference::Listener(listener_group_name.clone()), &unversioned_recommended_labels, &LISTENER_PVC_NAME, ); diff --git a/rust/operator-binary/src/controller/mod.rs b/rust/operator-binary/src/controller/mod.rs index 17448784..ceb99241 100644 --- a/rust/operator-binary/src/controller/mod.rs +++ b/rust/operator-binary/src/controller/mod.rs @@ -3,6 +3,7 @@ use std::{collections::BTreeMap, marker::PhantomData, str::FromStr}; use stackable_operator::{ commons::{ affinity::StackableAffinity, + pdb::PdbConfig, product_image_selection::ResolvedProductImage, resources::{NoRuntimeLimits, Resources}, }, @@ -87,14 +88,19 @@ pub struct KubernetesResources { pub role_bindings: Vec, pub status: PhantomData, } +// Webserver role only — all non-Option +#[derive(Clone, Debug)] +pub struct ValidatedWebserverRoleConfig { + pub pdb: PdbConfig, + pub listener_class: ListenerClassName, + pub group_listener_name: ListenerName, + pub trusted_proxies: Vec, +} -/// Per-role configuration extracted during validation. +// Other roles: scheduler, worker, dagprocessor, triggerer #[derive(Clone, Debug)] pub struct ValidatedRoleConfig { - pub pdb: Option, - pub listener_class: Option, - pub group_listener_name: Option, - pub trusted_proxies: Vec, + pub pdb: PdbConfig, } /// Per-rolegroup configuration: the merged CRD config plus overrides. @@ -218,20 +224,60 @@ pub struct ValidatedCluster { pub product_version: ProductVersion, pub image: ResolvedProductImage, pub cluster_config: ValidatedClusterConfig, - pub role_groups: BTreeMap>, - pub role_configs: BTreeMap, + pub webserver_config: Option, + pub webserver_role_group_configs: BTreeMap, + pub scheduler_config: Option, + pub scheduler_role_group_configs: BTreeMap, + pub dagprocessor_config: Option, + pub dagprocessor_role_group_configs: BTreeMap, + pub triggerer_config: Option, + pub triggerer_role_group_configs: BTreeMap, + pub worker_config: Option, + pub worker_role_group_configs: BTreeMap, +} + +/// The non-derived inputs to [`ValidatedCluster::new`]. +/// +/// Named fields, so the five same-typed role-group maps — and the four +/// `Option` — cannot be swapped silently. +pub struct ValidatedClusterParams { + pub name: ClusterName, + pub namespace: NamespaceName, + pub uid: Uid, + pub image: ResolvedProductImage, + pub cluster_config: ValidatedClusterConfig, + pub webserver_config: Option, + pub webserver_role_group_configs: BTreeMap, + pub scheduler_config: Option, + pub scheduler_role_group_configs: BTreeMap, + pub dagprocessor_config: Option, + pub dagprocessor_role_group_configs: BTreeMap, + pub triggerer_config: Option, + pub triggerer_role_group_configs: BTreeMap, + pub worker_config: Option, + pub worker_role_group_configs: BTreeMap, } impl ValidatedCluster { - pub fn new( - name: ClusterName, - namespace: NamespaceName, - uid: Uid, - image: ResolvedProductImage, - cluster_config: ValidatedClusterConfig, - role_groups: BTreeMap>, - role_configs: BTreeMap, - ) -> Self { + pub fn new(params: ValidatedClusterParams) -> Self { + let ValidatedClusterParams { + name, + namespace, + uid, + image, + cluster_config, + webserver_config, + webserver_role_group_configs, + scheduler_config, + scheduler_role_group_configs, + dagprocessor_config, + dagprocessor_role_group_configs, + triggerer_config, + triggerer_role_group_configs, + worker_config, + worker_role_group_configs, + } = params; + // `app_version_label_value` is constructed to be a valid label value, so it is also a valid // `ProductVersion`. let product_version = ProductVersion::from_str(&image.app_version_label_value) @@ -250,14 +296,73 @@ impl ValidatedCluster { product_version, image, cluster_config, - role_groups, - role_configs, + webserver_config, + webserver_role_group_configs, + scheduler_config, + scheduler_role_group_configs, + dagprocessor_config, + dagprocessor_role_group_configs, + triggerer_config, + triggerer_role_group_configs, + worker_config, + worker_role_group_configs, } } - /// Whether the cluster has the given role configured (i.e. it has role groups for it). + /// Whether the cluster declares the given role. pub fn has_role(&self, role: &AirflowRole) -> bool { - self.role_groups.contains_key(role) + match role { + AirflowRole::Webserver => self.webserver_config.is_some(), + AirflowRole::Scheduler => self.scheduler_config.is_some(), + AirflowRole::Worker => self.worker_config.is_some(), + AirflowRole::DagProcessor => self.dagprocessor_config.is_some(), + AirflowRole::Triggerer => self.triggerer_config.is_some(), + } + } + + /// The PodDisruptionBudget config of `role`, or `None` if the cluster does not declare it. + pub(crate) fn pdb(&self, role: &AirflowRole) -> Option<&PdbConfig> { + match role { + AirflowRole::Webserver => self.webserver_config.as_ref().map(|config| &config.pdb), + AirflowRole::Scheduler => self.scheduler_config.as_ref().map(|config| &config.pdb), + AirflowRole::Worker => self.worker_config.as_ref().map(|config| &config.pdb), + AirflowRole::DagProcessor => { + self.dagprocessor_config.as_ref().map(|config| &config.pdb) + } + AirflowRole::Triggerer => self.triggerer_config.as_ref().map(|config| &config.pdb), + } + } + + /// The name of the group Listener provided for `role`, if the role serves the web UI. + pub(crate) fn group_listener_name(&self, role: &AirflowRole) -> Option<&ListenerName> { + match role { + AirflowRole::Webserver => self + .webserver_config + .as_ref() + .map(|config| &config.group_listener_name), + AirflowRole::Scheduler + | AirflowRole::Worker + | AirflowRole::DagProcessor + | AirflowRole::Triggerer => None, + } + } + + /// The reverse proxies `role` trusts `X-Forwarded-*` headers from. + /// + /// Empty for every role but the webserver, which alone serves the web UI — and empty for the + /// webserver too when the cluster declares no webserver role, or it trusts no proxies. + pub(crate) fn trusted_proxies(&self, role: &AirflowRole) -> &[TrustedProxy] { + match role { + AirflowRole::Webserver => self + .webserver_config + .as_ref() + .map(|config| config.trusted_proxies.as_slice()) + .unwrap_or_default(), + AirflowRole::Scheduler + | AirflowRole::Worker + | AirflowRole::DagProcessor + | AirflowRole::Triggerer => &[], + } } /// The Secret holding the shared internal secret (`-internal-secret`). @@ -450,7 +555,12 @@ impl HasUid for ValidatedCluster { #[cfg(test)] mod tests { + use indoc::formatdoc; + use super::*; + use crate::controller::{ + build::test_support::dereferenced_objects, validate::validate_cluster, + }; #[test] fn test_constants() { @@ -462,4 +572,97 @@ mod tests { let _ = *EXECUTOR_ROLE_GROUP_NAME; let _ = *EXECUTOR_TEMPLATE_ROLE_GROUP_NAME; } + + #[test] + fn webserver_trusted_proxies_are_parsed() { + let cluster = validated_cluster_with_webserver_role_config( + " trustedProxies:\n - 10.244.0.0/16\n - 192.168.1.1", + ); + + let trusted_proxies = cluster.trusted_proxies(&AirflowRole::Webserver); + + let rendered: Vec = trusted_proxies + .iter() + .map(TrustedProxy::to_string) + .collect(); + assert_eq!(rendered, ["10.244.0.0/16", "192.168.1.1"]); + } + + /// Only the webserver serves HTTP, so no other role may pick the setting up even if a + /// webserver configured it. + #[test] + fn non_webserver_roles_have_no_trusted_proxies() { + let cluster = validated_cluster_with_webserver_role_config( + " trustedProxies:\n - 10.244.0.0/16", + ); + + for role in [ + AirflowRole::Scheduler, + AirflowRole::Worker, + AirflowRole::DagProcessor, + AirflowRole::Triggerer, + ] { + assert!( + cluster.trusted_proxies(&role).is_empty(), + "role {role:?} must not have trusted proxies" + ); + } + } + + #[test] + fn a_webserver_without_trusted_proxies_yields_an_empty_list() { + let cluster = + validated_cluster_with_webserver_role_config(" listenerClass: external-stable"); + + assert!(cluster.trusted_proxies(&AirflowRole::Webserver).is_empty()); + } + + /// The validated cluster for a CR with the given `webservers.roleConfig` block spliced in. + /// + /// The `roleConfig` must be one the webserver accepts: the trusted proxies are parsed and + /// checked by `validate_cluster`, not by [`ValidatedCluster::trusted_proxies`], which only + /// hands back what validation already accepted. The rejection cases live in + /// [`crate::crd::trusted_proxies`], next to the parsing they exercise. + fn validated_cluster_with_webserver_role_config(role_config: &str) -> ValidatedCluster { + validate_cluster( + &test_cluster_with_webserver_role_config(role_config), + "oci.stackable.tech/sdp", + dereferenced_objects(), + ) + .expect("test cluster validates") + } + + /// A cluster CR with the given `webservers.roleConfig` block spliced in. + fn test_cluster_with_webserver_role_config(role_config: &str) -> v1alpha2::AirflowCluster { + let cluster = formatdoc! {" + apiVersion: airflow.stackable.tech/v1alpha2 + kind: AirflowCluster + metadata: + name: airflow + namespace: default + uid: e6ac237d-a6d4-43a1-8135-f36506110912 + spec: + image: + productVersion: 3.2.2 + clusterConfig: + credentialsSecretName: airflow-admin-credentials + metadataDatabase: + postgresql: + host: airflow-postgresql + database: airflow + credentialsSecretName: airflow-postgresql-credentials + webservers: + roleConfig: + {role_config} + roleGroups: + default: + config: {{}} + kubernetesExecutors: + config: {{}} + "}; + + let deserializer = serde_yaml::Deserializer::from_str(&cluster); + serde_yaml::with::singleton_map_recursive::deserialize(deserializer) + .expect("the test CR deserialises") + } } diff --git a/rust/operator-binary/src/controller/validate.rs b/rust/operator-binary/src/controller/validate.rs index 2d092c89..5361d11d 100644 --- a/rust/operator-binary/src/controller/validate.rs +++ b/rust/operator-binary/src/controller/validate.rs @@ -2,7 +2,7 @@ use std::{collections::BTreeMap, str::FromStr}; use snafu::{OptionExt, ResultExt, Snafu}; use stackable_operator::{ - commons::product_image_selection, + commons::product_image_selection::{self, ResolvedProductImage}, config::fragment, crd::git_sync, k8s_openapi::api::core::v1::VolumeMount, @@ -22,18 +22,20 @@ use stackable_operator::{ }, }, }; -use strum::IntoEnumIterator; use super::{ AirflowRoleGroupConfig, ValidatedAirflowConfig, ValidatedCluster, ValidatedClusterConfig, - ValidatedExecutorTemplate, ValidatedLogging, ValidatedRoleConfig, - build::volumes::LOG_VOLUME_NAME, dereference::DereferencedObjects, + ValidatedClusterParams, ValidatedExecutorTemplate, ValidatedLogging, ValidatedRoleConfig, + ValidatedWebserverRoleConfig, build::volumes::LOG_VOLUME_NAME, + dereference::DereferencedObjects, }; use crate::{ airflow_controller::CONTAINER_IMAGE_BASE_NAME, crd::{ AirflowConfig, AirflowConfigFragment, AirflowConfigOverrides, AirflowExecutor, AirflowRole, - AirflowRoleType, Container, v1alpha2, + AirflowRoleType, Container, + trusted_proxies::{self, TrustedProxy}, + v1alpha2, }, }; @@ -124,55 +126,82 @@ pub fn validate_cluster( .vector_aggregator_config_map_name .clone(); - let mut role_groups = BTreeMap::new(); - let mut role_configs = BTreeMap::new(); - - // if the kubernetes executor is specified there will be no worker role as the pods - // are provisioned by airflow as defined by the task (default: one pod per task) - for role in AirflowRole::iter() { - let Some(resolved_role) = airflow.get_role(&role) else { - continue; - }; - - role_configs.insert( - role.clone(), - ValidatedRoleConfig { - pdb: airflow - .role_config(&role) - .map(|rc| rc.pod_disruption_budget), - listener_class: role.listener_class_name(airflow), - group_listener_name: airflow.group_listener_name(&role), - trusted_proxies: role - .trusted_proxies(airflow) - .context(ParseTrustedProxiesSnafu)?, - }, - ); + // Only the webserver serves the web UI, so only it has a listener class, a group listener and + // trusted proxies. The other roles carry nothing but their Pod disruption budget. + let webserver_config = airflow + .spec + .webservers + .as_ref() + .map(|webservers| { + let trusted_proxies = webservers + .role_config + .trusted_proxies + .iter() + .map(|trusted_proxy| TrustedProxy::from_str(trusted_proxy)) + .collect::, _>>() + .and_then(|entries| { + trusted_proxies::ensure_wildcard_is_sole_entry(&entries)?; + Ok(entries) + }) + .context(ParseTrustedProxiesSnafu)?; + + Ok(ValidatedWebserverRoleConfig { + pdb: webservers.role_config.common.pod_disruption_budget.clone(), + listener_class: webservers.role_config.listener_class.clone(), + group_listener_name: airflow + .group_listener_name(&AirflowRole::Webserver) + .expect("the webserver role always has a group listener"), + trusted_proxies, + }) + }) + .transpose()?; + let webserver_role_group_configs = validate_role_groups( + airflow, + &AirflowRole::Webserver, + &vector_aggregator_config_map_name, + &resolved_product_image, + )?; - let default_config = AirflowConfig::default_config(&airflow.name_any(), &role); + let scheduler_config = airflow.spec.schedulers.as_ref().map(pdb_only_role_config); + let scheduler_role_group_configs = validate_role_groups( + airflow, + &AirflowRole::Scheduler, + &vector_aggregator_config_map_name, + &resolved_product_image, + )?; - let mut group_configs = BTreeMap::new(); - for (rolegroup_name, rolegroup) in &resolved_role.role_groups { - let role_group_name = RoleGroupName::from_str(rolegroup_name).with_context(|_| { - ParseRoleGroupNameSnafu { - role_group: rolegroup_name.clone(), - } - })?; - let config = validate_role_group( - &resolved_role, - &role_group_name, - rolegroup, - &default_config, - &vector_aggregator_config_map_name, - &resolved_product_image, - &airflow.spec.cluster_config.dags_git_sync, - &airflow.spec.cluster_config.volume_mounts, - )?; + let dagprocessor_config = airflow + .spec + .dag_processors + .as_ref() + .map(pdb_only_role_config); + let dagprocessor_role_group_configs = validate_role_groups( + airflow, + &AirflowRole::DagProcessor, + &vector_aggregator_config_map_name, + &resolved_product_image, + )?; - group_configs.insert(role_group_name, config); - } + let triggerer_config = airflow.spec.triggerers.as_ref().map(pdb_only_role_config); + let triggerer_role_group_configs = validate_role_groups( + airflow, + &AirflowRole::Triggerer, + &vector_aggregator_config_map_name, + &resolved_product_image, + )?; - role_groups.insert(role, group_configs); - } + // If the Kubernetes executor is specified there is no worker role, as the Pods are provisioned + // by Airflow as defined by the task (default: one Pod per task). + let worker_config = match &airflow.spec.executor { + AirflowExecutor::CeleryExecutors { config } => Some(pdb_only_role_config(config.as_ref())), + AirflowExecutor::KubernetesExecutors { .. } => None, + }; + let worker_role_group_configs = validate_role_groups( + airflow, + &AirflowRole::Worker, + &vector_aggregator_config_map_name, + &resolved_product_image, + )?; let DereferencedObjects { authentication_config, @@ -236,12 +265,12 @@ pub fn validate_cluster( AirflowExecutor::CeleryExecutors { .. } => None, }; - Ok(ValidatedCluster::new( - cluster_name, + Ok(ValidatedCluster::new(ValidatedClusterParams { + name: cluster_name, namespace, uid, - resolved_product_image, - ValidatedClusterConfig { + image: resolved_product_image, + cluster_config: ValidatedClusterConfig { executor: airflow.spec.executor.clone(), executor_template, authentication_config, @@ -260,9 +289,60 @@ pub fn validate_cluster( volumes: airflow.spec.cluster_config.volumes.clone(), volume_mounts: airflow.spec.cluster_config.volume_mounts.clone(), }, - role_groups, - role_configs, - )) + webserver_config, + webserver_role_group_configs, + scheduler_config, + scheduler_role_group_configs, + dagprocessor_config, + dagprocessor_role_group_configs, + triggerer_config, + triggerer_role_group_configs, + worker_config, + worker_role_group_configs, + })) +} + +/// The validated role-level config of a role that carries nothing but its Pod disruption budget. +fn pdb_only_role_config(role: &AirflowRoleType) -> ValidatedRoleConfig { + ValidatedRoleConfig { + pdb: role.role_config.pod_disruption_budget.clone(), + } +} + +/// The validated config of every role group of `role`, or an empty map if the role is absent. +fn validate_role_groups( + airflow: &v1alpha2::AirflowCluster, + role: &AirflowRole, + vector_aggregator_config_map_name: &Option, + resolved_product_image: &ResolvedProductImage, +) -> Result, Error> { + let Some(resolved_role) = airflow.get_role(role) else { + return Ok(BTreeMap::new()); + }; + let default_config = AirflowConfig::default_config(&airflow.name_any(), role); + + resolved_role + .role_groups + .iter() + .map(|(rolegroup_name, rolegroup)| { + let role_group_name = RoleGroupName::from_str(rolegroup_name).with_context(|_| { + ParseRoleGroupNameSnafu { + role_group: rolegroup_name.clone(), + } + })?; + let config = validate_role_group( + &resolved_role, + &role_group_name, + rolegroup, + &default_config, + vector_aggregator_config_map_name, + resolved_product_image, + &airflow.spec.cluster_config.dags_git_sync, + &airflow.spec.cluster_config.volume_mounts, + )?; + Ok((role_group_name, config)) + }) + .collect() } /// Validate and merge one role group against its role. @@ -359,10 +439,14 @@ pub(crate) fn validate_logging( mod tests { use std::collections::BTreeMap; + use rstest::rstest; use stackable_operator::v2::builder::pod::container::{EnvVarName, EnvVarSet}; - use super::validate_role_group; - use crate::crd::{AirflowConfig, AirflowRole, v1alpha2}; + use super::{Error, validate_cluster, validate_role_group}; + use crate::{ + controller::build::test_support::dereferenced_objects, + crd::{AirflowConfig, AirflowRole, v1alpha2}, + }; /// A minimal resolved product image for tests that exercise `validate_role_group` (which needs /// one to resolve git-sync resources; the test role groups configure no git-sync, so the value @@ -637,4 +721,41 @@ mod tests { assert_eq!(labels.get("rg-label"), Some(&"rg".to_string())); assert_eq!(labels.get("shared"), Some(&"rg".to_string())); } + + /// A `trustedProxies` list the webserver cannot accept must be rejected here, where it is + /// parsed and checked. `ValidatedCluster::trusted_proxies` is infallible and only hands back + /// what validation already accepted, so nothing downstream can catch this. + /// + /// The entries themselves are covered by `crd::trusted_proxies`; what this adds is that a + /// cluster carrying them fails to validate rather than silently losing the setting. + #[rstest] + #[case::wildcard_combined_with_another_entry(&["*", "10.0.0.0/8"])] + #[case::not_an_ip_address(&["airflow.example.com"])] + fn an_invalid_webserver_trusted_proxy_list_is_rejected(#[case] trusted_proxies: &[&str]) { + let mut cluster = test_cluster(); + // The namespace and uid are resolved before the trusted proxies are, so the fixture needs + // both for this to fail on the list rather than on the metadata. + cluster.metadata.namespace = Some("default".to_owned()); + cluster.metadata.uid = Some("e6ac237d-a6d4-43a1-8135-f36506110912".to_owned()); + cluster + .spec + .webservers + .as_mut() + .expect("the test CR declares a webserver role") + .role_config + .trusted_proxies = trusted_proxies + .iter() + .map(|proxy| proxy.to_string()) + .collect(); + + let error = validate_cluster(&cluster, "oci.stackable.tech/sdp", dereferenced_objects()) + // `ValidatedCluster` is not `Debug`, so map the success case away before `expect_err`. + .map(|_| ()) + .expect_err("the webserver must reject this trusted proxy list"); + + assert!( + matches!(error, Error::ParseTrustedProxies { .. }), + "error was: {error:?}" + ); + } } diff --git a/rust/operator-binary/src/crd/mod.rs b/rust/operator-binary/src/crd/mod.rs index 846c643d..d43f2087 100644 --- a/rust/operator-binary/src/crd/mod.rs +++ b/rust/operator-binary/src/crd/mod.rs @@ -66,7 +66,6 @@ use crate::{ databases::{ CeleryBrokerConnection, CeleryResultBackendConnection, MetadataDatabaseConnection, }, - trusted_proxies::TrustedProxy, }, util::role_service_name, }; @@ -453,10 +452,6 @@ impl v1alpha2::AirflowCluster { } } - pub fn role_config(&self, role: &AirflowRole) -> Option { - self.get_role(role).map(|r| r.role_config) - } - /// Retrieve and merge resource configs for the executor template pub fn merged_executor_config( &self, @@ -708,10 +703,7 @@ impl AirflowRole { return ""; } - let has_trusted_proxies = cluster - .role_configs - .get(&AirflowRole::Webserver) - .is_some_and(|role_config| !role_config.trusted_proxies.is_empty()); + let has_trusted_proxies = !cluster.trusted_proxies(self).is_empty(); if has_trusted_proxies { " --proxy-headers" @@ -763,44 +755,6 @@ impl AirflowRole { AirflowRole::Triggerer => None, } } - - pub fn listener_class_name( - &self, - airflow: &v1alpha2::AirflowCluster, - ) -> Option { - match self { - Self::Webserver => airflow - .spec - .webservers - .to_owned() - .map(|webserver| webserver.role_config.listener_class), - Self::Worker | Self::Scheduler | Self::DagProcessor | Self::Triggerer => None, - } - } - - /// The reverse proxies this role trusts `X-Forwarded-*` headers from. - /// - /// Only the webserver serves HTTP, so every other role returns an empty list regardless of - /// what the webserver role configured. - pub fn trusted_proxies( - &self, - airflow: &v1alpha2::AirflowCluster, - ) -> Result, trusted_proxies::Error> { - match self { - Self::Webserver => { - let entries: Vec = airflow - .spec - .webservers - .iter() - .flat_map(|webserver| &webserver.role_config.trusted_proxies) - .map(|trusted_proxy| TrustedProxy::from_str(trusted_proxy)) - .collect::>()?; - trusted_proxies::ensure_wildcard_is_sole_entry(&entries)?; - Ok(entries) - } - Self::Worker | Self::Scheduler | Self::DagProcessor | Self::Triggerer => Ok(Vec::new()), - } - } } impl Deref for AirflowRole { @@ -1170,120 +1124,6 @@ mod tests { ); } - /// A cluster CR with the given `webservers.roleConfig` block spliced in. - fn test_cluster_with_webserver_role_config(role_config: &str) -> v1alpha2::AirflowCluster { - let cluster = formatdoc! {" - apiVersion: airflow.stackable.tech/v1alpha2 - kind: AirflowCluster - metadata: - name: airflow - spec: - image: - productVersion: 3.2.2 - clusterConfig: - credentialsSecretName: airflow-admin-credentials - metadataDatabase: - postgresql: - host: airflow-postgresql - database: airflow - credentialsSecretName: airflow-postgresql-credentials - webservers: - roleConfig: - {role_config} - roleGroups: - default: - config: {{}} - kubernetesExecutors: - config: {{}} - "}; - - let deserializer = serde_yaml::Deserializer::from_str(&cluster); - serde_yaml::with::singleton_map_recursive::deserialize(deserializer) - .expect("the test CR deserialises") - } - - #[test] - fn webserver_trusted_proxies_are_parsed() { - let cluster = test_cluster_with_webserver_role_config( - " trustedProxies:\n - 10.244.0.0/16\n - 192.168.1.1", - ); - - let trusted_proxies = AirflowRole::Webserver - .trusted_proxies(&cluster) - .expect("the trusted proxies are valid"); - - let rendered: Vec = trusted_proxies - .iter() - .map(TrustedProxy::to_string) - .collect(); - assert_eq!(rendered, ["10.244.0.0/16", "192.168.1.1"]); - } - - #[test] - fn wildcard_combined_with_another_entry_is_rejected() { - let cluster = test_cluster_with_webserver_role_config( - " trustedProxies:\n - \"*\"\n - 10.0.0.0/8", - ); - - let error = AirflowRole::Webserver - .trusted_proxies(&cluster) - .expect_err("* combined with another entry must be rejected"); - assert!( - matches!( - error, - crate::crd::trusted_proxies::Error::WildcardMustBeSoleEntry - ), - "error was: {error:?}" - ); - } - - #[test] - fn an_invalid_trusted_proxy_is_rejected() { - let cluster = test_cluster_with_webserver_role_config( - " trustedProxies:\n - airflow.example.com", - ); - - AirflowRole::Webserver - .trusted_proxies(&cluster) - .expect_err("a hostname is not a valid trusted proxy"); - } - - /// Only the webserver serves HTTP, so no other role may pick the setting up even if a - /// webserver configured it. - #[test] - fn non_webserver_roles_have_no_trusted_proxies() { - let cluster = test_cluster_with_webserver_role_config( - " trustedProxies:\n - 10.244.0.0/16", - ); - - for role in [ - AirflowRole::Scheduler, - AirflowRole::Worker, - AirflowRole::DagProcessor, - AirflowRole::Triggerer, - ] { - assert!( - role.trusted_proxies(&cluster) - .expect("no proxies to parse") - .is_empty(), - "role {role:?} must not have trusted proxies" - ); - } - } - - #[test] - fn a_webserver_without_trusted_proxies_yields_an_empty_list() { - let cluster = - test_cluster_with_webserver_role_config(" listenerClass: external-stable"); - - assert!( - AirflowRole::Webserver - .trusted_proxies(&cluster) - .expect("nothing to parse") - .is_empty() - ); - } - /// The commands that background a process, i.e. the candidates for `$!`. fn backgrounded_commands(role: &AirflowRole, cluster: &ValidatedCluster) -> Vec { role.get_commands(cluster) diff --git a/rust/operator-binary/src/crd/trusted_proxies.rs b/rust/operator-binary/src/crd/trusted_proxies.rs index d44f1dfd..35b9a82f 100644 --- a/rust/operator-binary/src/crd/trusted_proxies.rs +++ b/rust/operator-binary/src/crd/trusted_proxies.rs @@ -207,6 +207,12 @@ mod tests { )); } + #[test] + fn an_invalid_trusted_proxy_is_rejected() { + TrustedProxy::from_str("airflow.example.com") + .expect_err("a hostname is not a valid trusted proxy"); + } + #[test] fn rejects_a_non_numeric_prefix_length() { assert!(matches!( @@ -310,6 +316,21 @@ mod tests { )); } + #[test] + fn a_wildcard_with_another_entry_reports_wildcard_must_be_sole_entry() { + let entries = [ + TrustedProxy::from_str("*").expect("must be accepted"), + TrustedProxy::from_str("10.0.0.0/8").expect("must be accepted"), + ]; + + let error = ensure_wildcard_is_sole_entry(&entries) + .expect_err("* combined with another entry must be rejected"); + assert!( + matches!(error, Error::WildcardMustBeSoleEntry), + "error was: {error:?}" + ); + } + #[test] fn a_list_without_a_wildcard_is_accepted() { let entries = [TrustedProxy::from_str("10.0.0.0/8").expect("must be accepted")]; From 32b420e1eaf2641a31033f6932c607feca664aaa Mon Sep 17 00:00:00 2001 From: Andrew Kenworthy Date: Wed, 23 Sep 2026 17:08:53 +0200 Subject: [PATCH 2/2] changelog --- CHANGELOG.md | 5 +++++ 1 file changed, 5 insertions(+) diff --git a/CHANGELOG.md b/CHANGELOG.md index 439f212a..5211e4cc 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -44,6 +44,10 @@ overridden, whereas previously the operator's value always took precedence ([#838]). - Make operations infallible where appropriate ([#852], [#860]). - Deprecated airflow `3.2.2` ([#865]). +- Internal operator refactoring: the validated cluster carries each role's configuration in its own + typed fields instead of maps keyed by role, and only the webserver carries a listener class, a + group listener and trusted proxies, so those are no longer optional fields that four of the five + roles leave unset ([#867]). ### Fixed @@ -76,6 +80,7 @@ [#860]: https://github.com/stackabletech/airflow-operator/pull/860 [#862]: https://github.com/stackabletech/airflow-operator/pull/862 [#865]: https://github.com/stackabletech/airflow-operator/pull/865 +[#867]: https://github.com/stackabletech/airflow-operator/pull/867 ## [26.7.0] - 2026-07-21