Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
5 changes: 5 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -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

Expand Down Expand Up @@ -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

Expand Down
69 changes: 46 additions & 23 deletions rust/operator-binary/src/controller/build/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -27,7 +27,7 @@ use crate::{
statefulset::build_server_rolegroup_statefulset,
},
},
crd::{AirflowConfigOverrides, Container},
crd::{AirflowConfigOverrides, AirflowRole, Container},
};

pub mod graceful_shutdown;
Expand Down Expand Up @@ -93,23 +93,27 @@ pub fn build(cluster: &ValidatedCluster) -> Result<KubernetesResources<Prepared>
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;

Expand Down Expand Up @@ -149,6 +153,21 @@ pub fn build(cluster: &ValidatedCluster) -> Result<KubernetesResources<Prepared>
})?,
);
}

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 {
Expand Down Expand Up @@ -335,18 +354,22 @@ 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,
user_registration_role: "Public".to_string(),
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.
Expand Down
11 changes: 2 additions & 9 deletions rust/operator-binary/src/controller/build/properties/env_vars.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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::<Vec<_>>()
Expand All @@ -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
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -221,17 +221,13 @@ pub fn build_server_rolegroup_statefulset(

let mut pvcs: Option<Vec<PersistentVolumeClaim>> = 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,
);
Expand Down
Loading
Loading