Skip to content

Commit 5a6f9dc

Browse files
committed
alter structure to attempt to being it mroe in line with e.g. hdfs. Added test for merge/extend difference
1 parent 6512d4d commit 5a6f9dc

4 files changed

Lines changed: 254 additions & 155 deletions

File tree

rust/operator-binary/src/airflow_controller.rs

Lines changed: 12 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -1,6 +1,6 @@
11
//! Ensures that `Pod`s are configured and running for each [`v1alpha2::AirflowCluster`]
22
use std::{
3-
collections::{BTreeMap, BTreeSet, HashMap},
3+
collections::{BTreeSet, HashMap},
44
sync::Arc,
55
};
66

@@ -73,10 +73,11 @@ use crate::{
7373
},
7474
controller_commons::{self, CONFIG_VOLUME_NAME, LOG_CONFIG_VOLUME_NAME, LOG_VOLUME_NAME},
7575
crd::{
76-
self, APP_NAME, AirflowClusterStatus, AirflowExecutor, AirflowExecutorCommonConfiguration,
77-
AirflowRole, CONFIG_PATH, Container, ExecutorConfig, HTTP_PORT, HTTP_PORT_NAME,
78-
LISTENER_VOLUME_DIR, LISTENER_VOLUME_NAME, LOG_CONFIG_DIR, METRICS_PORT, METRICS_PORT_NAME,
79-
OPERATOR_NAME, STACKABLE_LOG_DIR, TEMPLATE_LOCATION, TEMPLATE_NAME, TEMPLATE_VOLUME_NAME,
76+
self, APP_NAME, AirflowClusterStatus, AirflowConfigOverrides, AirflowExecutor,
77+
AirflowExecutorCommonConfiguration, AirflowRole, CONFIG_PATH, Container, ExecutorConfig,
78+
HTTP_PORT, HTTP_PORT_NAME, LISTENER_VOLUME_DIR, LISTENER_VOLUME_NAME, LOG_CONFIG_DIR,
79+
METRICS_PORT, METRICS_PORT_NAME, OPERATOR_NAME, STACKABLE_LOG_DIR, TEMPLATE_LOCATION,
80+
TEMPLATE_NAME, TEMPLATE_VOLUME_NAME,
8081
authentication::{
8182
AirflowAuthenticationClassResolved, AirflowClientAuthenticationDetailsResolved,
8283
},
@@ -475,7 +476,7 @@ pub async fn reconcile_airflow(
475476
let git_sync_resources = git_sync::v1alpha2::GitSyncResources::new(
476477
&airflow.spec.cluster_config.dags_git_sync,
477478
&validated_cluster.image,
478-
&env_vars_from_overrides(&validated_rg_config.overrides.env_overrides),
479+
&env_vars_from_overrides(&validated_rg_config.env_overrides),
479480
&airflow.volume_mounts(),
480481
LOG_VOLUME_NAME,
481482
&validated_rg_config
@@ -534,7 +535,7 @@ pub async fn reconcile_airflow(
534535
airflow,
535536
&validated_cluster,
536537
&rolegroup,
537-
&validated_rg_config.overrides.config_file_overrides,
538+
&validated_rg_config.config_overrides,
538539
&validated_rg_config.merged_config.logging,
539540
&Container::Airflow,
540541
)
@@ -611,7 +612,9 @@ async fn build_executor_template(
611612
airflow,
612613
validated_cluster,
613614
&rolegroup,
614-
&BTreeMap::new(),
615+
// The kubernetes-executor pod template does not apply webserver_config.py overrides
616+
// (preserves prior behaviour, which passed an empty map here).
617+
&AirflowConfigOverrides::default(),
615618
&merged_executor_config.logging,
616619
&Container::Base,
617620
)
@@ -730,7 +733,7 @@ fn build_server_rolegroup_statefulset(
730733
git_sync_resources: &git_sync::v1alpha2::GitSyncResources,
731734
) -> Result<StatefulSet> {
732735
let merged_airflow_config = &validated_rg_config.merged_config;
733-
let env_overrides = &validated_rg_config.overrides.env_overrides;
736+
let env_overrides = &validated_rg_config.env_overrides;
734737

735738
let resolved_product_image = &validated_cluster.image;
736739
let authentication_config = &validated_cluster.authentication_config;

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

Lines changed: 13 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -19,7 +19,8 @@ use crate::{
1919
config::webserver_config,
2020
controller::validate::ValidatedAirflowCluster,
2121
crd::{
22-
AIRFLOW_CONFIG_FILENAME, Container, STACKABLE_LOG_DIR, build_recommended_labels, v1alpha2,
22+
AIRFLOW_CONFIG_FILENAME, AirflowConfigOverrides, Container, STACKABLE_LOG_DIR,
23+
build_recommended_labels, v1alpha2,
2324
},
2425
product_logging::{LOG_CONFIG_FILE, create_airflow_config},
2526
};
@@ -60,15 +61,24 @@ pub fn build_rolegroup_config_map(
6061
airflow: &v1alpha2::AirflowCluster,
6162
validated_cluster: &ValidatedAirflowCluster,
6263
rolegroup: &RoleGroupRef<v1alpha2::AirflowCluster>,
63-
config_file_overrides: &BTreeMap<String, String>,
64+
config_overrides: &AirflowConfigOverrides,
6465
logging: &Logging<Container>,
6566
container: &Container,
6667
) -> Result<ConfigMap, Error> {
68+
// Flatten the typed `webserver_config.py` overrides into a plain map for the file writer,
69+
// dropping entries whose value is unset (`null`).
70+
let config_file_overrides: BTreeMap<String, String> = config_overrides
71+
.webserver_config_py
72+
.overrides
73+
.iter()
74+
.filter_map(|(key, value)| value.clone().map(|value| (key.clone(), value)))
75+
.collect();
76+
6777
let config_file = webserver_config::build(
6878
&validated_cluster.authentication_config,
6979
&validated_cluster.authorization_config,
7080
&validated_cluster.image.product_version,
71-
config_file_overrides,
81+
&config_file_overrides,
7282
)
7383
.with_context(|_| BuildWebserverConfigSnafu {
7484
rolegroup: rolegroup.clone(),

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

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

33
use snafu::{ResultExt, Snafu};
44
use stackable_operator::{
55
commons::product_image_selection::{self, ResolvedProductImage},
6+
config::merge::Merge,
67
role_utils::RoleGroupRef,
78
};
89
use strum::IntoEnumIterator;
@@ -11,7 +12,7 @@ use super::dereference::DereferencedObjects;
1112
use crate::{
1213
airflow_controller::CONTAINER_IMAGE_BASE_NAME,
1314
crd::{
14-
AirflowConfig, AirflowExecutor, AirflowRole, MergedOverrides,
15+
AirflowConfig, AirflowConfigOverrides, AirflowExecutor, AirflowRole, AirflowRoleType,
1516
authentication::AirflowClientAuthenticationDetailsResolved,
1617
authorization::AirflowAuthorizationResolved, v1alpha2,
1718
},
@@ -37,10 +38,15 @@ pub struct ValidatedRoleConfig {
3738
}
3839

3940
/// Per-rolegroup configuration: the merged CRD config plus overrides.
41+
///
42+
/// `config_overrides` is kept as the typed [`AirflowConfigOverrides`] (role-group merged over
43+
/// role); it is flattened into the rendered config file later, in the build step. This mirrors
44+
/// hdfs-operator. `env_overrides` is already a flat map.
4045
#[derive(Clone, Debug)]
4146
pub struct ValidatedRoleGroupConfig {
4247
pub merged_config: AirflowConfig,
43-
pub overrides: MergedOverrides,
48+
pub config_overrides: AirflowConfigOverrides,
49+
pub env_overrides: HashMap<String, String>,
4450
}
4551

4652
/// The validated cluster: proves that config merging succeeded for every role and
@@ -104,15 +110,15 @@ pub fn validate_cluster(
104110
.merged_config(&role, &rolegroup_ref)
105111
.context(FailedToResolveConfigSnafu)?;
106112

107-
let overrides = airflow
108-
.merged_overrides(&role, rolegroup_name)
109-
.context(FailedToResolveConfigSnafu)?;
113+
let (config_overrides, env_overrides) =
114+
merge_role_group_overrides(&resolved_role, rolegroup_name);
110115

111116
group_configs.insert(
112117
rolegroup_name.clone(),
113118
ValidatedRoleGroupConfig {
114119
merged_config,
115-
overrides,
120+
config_overrides,
121+
env_overrides,
116122
},
117123
);
118124
}
@@ -134,3 +140,217 @@ pub fn validate_cluster(
134140
authorization_config,
135141
})
136142
}
143+
144+
/// Merge a role group's config overrides over the role-level ones (role-group wins per key) via
145+
/// the `Merge` impl on [`AirflowConfigOverrides`], and combine env overrides (role first, then
146+
/// role-group on top). Mirrors hdfs-operator's `validate_role_group_config`.
147+
///
148+
/// The merged overrides are returned *typed*; flattening into the rendered `webserver_config.py`
149+
/// happens later, in the build step. Note the `Merge` semantics: a role-group `null` inherits the
150+
/// role-level value rather than unsetting it.
151+
fn merge_role_group_overrides(
152+
role: &AirflowRoleType,
153+
rolegroup_name: &str,
154+
) -> (AirflowConfigOverrides, HashMap<String, String>) {
155+
let rolegroup = role.role_groups.get(rolegroup_name);
156+
157+
let mut config_overrides = rolegroup
158+
.map(|rg| rg.config.config_overrides.clone())
159+
.unwrap_or_default();
160+
config_overrides.merge(&role.config.config_overrides);
161+
162+
let mut env_overrides = role.config.env_overrides.clone();
163+
if let Some(rg) = rolegroup {
164+
env_overrides.extend(rg.config.env_overrides.clone());
165+
}
166+
167+
(config_overrides, env_overrides)
168+
}
169+
170+
#[cfg(test)]
171+
mod tests {
172+
use std::collections::BTreeMap;
173+
174+
use super::merge_role_group_overrides;
175+
use crate::crd::{AirflowRole, v1alpha2};
176+
177+
fn test_cluster() -> v1alpha2::AirflowCluster {
178+
let cluster_yaml = r#"
179+
apiVersion: airflow.stackable.tech/v1alpha2
180+
kind: AirflowCluster
181+
metadata:
182+
name: airflow
183+
spec:
184+
image:
185+
productVersion: 3.1.6
186+
clusterConfig:
187+
loadExamples: false
188+
exposeConfig: false
189+
credentialsSecretName: airflow-admin-credentials
190+
metadataDatabase:
191+
postgresql:
192+
host: airflow-postgresql
193+
database: airflow
194+
credentialsSecretName: airflow-postgresql-credentials
195+
webservers:
196+
config: {}
197+
configOverrides:
198+
webserver_config.py:
199+
AUTH_TYPE: "AUTH_OID"
200+
ROLE_ONLY_KEY: "role-value"
201+
envOverrides:
202+
ROLE_ENV_VAR: "role-env-value"
203+
roleGroups:
204+
default:
205+
config: {}
206+
configOverrides:
207+
webserver_config.py:
208+
AUTH_TYPE: "AUTH_DB"
209+
GROUP_ONLY_KEY: "group-value"
210+
envOverrides:
211+
GROUP_ENV_VAR: "group-env-value"
212+
schedulers:
213+
config: {}
214+
roleGroups:
215+
default:
216+
config: {}
217+
kubernetesExecutors:
218+
config: {}
219+
"#;
220+
let deserializer = serde_yaml::Deserializer::from_str(cluster_yaml);
221+
serde_yaml::with::singleton_map_recursive::deserialize(deserializer).unwrap()
222+
}
223+
224+
#[test]
225+
fn role_group_overrides_merge_over_role_overrides() {
226+
let cluster = test_cluster();
227+
let role = cluster
228+
.get_role(&AirflowRole::Webserver)
229+
.expect("webserver role");
230+
231+
let (config_overrides, env_overrides) = merge_role_group_overrides(&role, "default");
232+
233+
// configOverrides are kept typed (values are `Option<String>`). The role-group AUTH_TYPE
234+
// overrides the role-level one; both role-only and group-only keys are kept.
235+
assert_eq!(
236+
config_overrides.webserver_config_py.overrides,
237+
BTreeMap::from([
238+
("AUTH_TYPE".to_string(), Some("AUTH_DB".to_string())),
239+
("ROLE_ONLY_KEY".to_string(), Some("role-value".to_string())),
240+
(
241+
"GROUP_ONLY_KEY".to_string(),
242+
Some("group-value".to_string())
243+
),
244+
])
245+
);
246+
247+
assert_eq!(env_overrides.len(), 2);
248+
assert_eq!(env_overrides.get("ROLE_ENV_VAR").unwrap(), "role-env-value");
249+
assert_eq!(
250+
env_overrides.get("GROUP_ENV_VAR").unwrap(),
251+
"group-env-value"
252+
);
253+
}
254+
255+
/// A role-group `null` override inherits the role-level value instead of unsetting it — the
256+
/// behavioural consequence of merging with `Merge` rather than `.extend()`. `main`'s
257+
/// product-config used `.extend()`, where the same input would have *removed* the key. This
258+
/// test pins that choice (mirrors the kafka-operator test).
259+
#[test]
260+
fn role_group_null_inherits_role_value_rather_than_unsetting_it() {
261+
let cluster_yaml = r#"
262+
apiVersion: airflow.stackable.tech/v1alpha2
263+
kind: AirflowCluster
264+
metadata:
265+
name: airflow
266+
spec:
267+
image:
268+
productVersion: 3.1.6
269+
clusterConfig:
270+
loadExamples: false
271+
exposeConfig: false
272+
credentialsSecretName: airflow-admin-credentials
273+
metadataDatabase:
274+
postgresql:
275+
host: airflow-postgresql
276+
database: airflow
277+
credentialsSecretName: airflow-postgresql-credentials
278+
webservers:
279+
config: {}
280+
configOverrides:
281+
webserver_config.py:
282+
AUTH_TYPE: "AUTH_OID"
283+
roleGroups:
284+
default:
285+
config: {}
286+
configOverrides:
287+
webserver_config.py:
288+
AUTH_TYPE: null
289+
schedulers:
290+
config: {}
291+
roleGroups:
292+
default:
293+
config: {}
294+
kubernetesExecutors:
295+
config: {}
296+
"#;
297+
let deserializer = serde_yaml::Deserializer::from_str(cluster_yaml);
298+
let cluster: v1alpha2::AirflowCluster =
299+
serde_yaml::with::singleton_map_recursive::deserialize(deserializer).unwrap();
300+
let role = cluster
301+
.get_role(&AirflowRole::Webserver)
302+
.expect("webserver role");
303+
304+
// For contrast: under the old `.extend()` layering the role-group `null` overwrites the
305+
// role value, and the key is dropped on flatten — i.e. AUTH_TYPE is unset entirely.
306+
let old_extend_behaviour: BTreeMap<String, String> = {
307+
let mut combined = role
308+
.config
309+
.config_overrides
310+
.webserver_config_py
311+
.overrides
312+
.clone();
313+
let rg = role.role_groups.get("default").expect("default role group");
314+
combined.extend(
315+
rg.config
316+
.config_overrides
317+
.webserver_config_py
318+
.overrides
319+
.clone(),
320+
);
321+
combined
322+
.into_iter()
323+
.filter_map(|(key, value)| value.map(|value| (key, value)))
324+
.collect()
325+
};
326+
assert!(
327+
!old_extend_behaviour.contains_key("AUTH_TYPE"),
328+
"under the old `.extend()` behaviour the role-group `null` unsets AUTH_TYPE"
329+
);
330+
331+
// What we do now (Merge): the role-group `null` inherits the role-level value, so
332+
// AUTH_TYPE survives as the role's "AUTH_OID".
333+
let (config_overrides, _env_overrides) = merge_role_group_overrides(&role, "default");
334+
assert_eq!(
335+
config_overrides
336+
.webserver_config_py
337+
.overrides
338+
.get("AUTH_TYPE"),
339+
Some(&Some("AUTH_OID".to_string())),
340+
"role-group `null` should inherit the role-level AUTH_TYPE under Merge semantics"
341+
);
342+
}
343+
344+
#[test]
345+
fn role_without_overrides_yields_empty() {
346+
let cluster = test_cluster();
347+
let role = cluster
348+
.get_role(&AirflowRole::Scheduler)
349+
.expect("scheduler role");
350+
351+
let (config_overrides, env_overrides) = merge_role_group_overrides(&role, "default");
352+
353+
assert!(config_overrides.webserver_config_py.overrides.is_empty());
354+
assert!(env_overrides.is_empty());
355+
}
356+
}

0 commit comments

Comments
 (0)