Skip to content

Commit ffefaf6

Browse files
committed
refactor: use RoleGroupName, use v2 infallible labels, ResourceNames and reduce usage of RoleGroupRef. N.B. compiled against local op-rs
1 parent 3334481 commit ffefaf6

8 files changed

Lines changed: 434 additions & 287 deletions

File tree

rust/operator-binary/src/airflow_controller.rs

Lines changed: 93 additions & 113 deletions
Original file line numberDiff line numberDiff line change
@@ -53,7 +53,7 @@ use stackable_operator::{
5353
core::{DeserializeGuard, error_boundary},
5454
runtime::{controller::Action, reflector::ObjectRef},
5555
},
56-
kvp::{Annotation, Label, LabelError, Labels, ObjectLabels},
56+
kvp::{Annotation, Label, LabelError},
5757
logging::controller::ReconcilerError,
5858
product_logging::{self, framework::LoggingError, spec::ContainerLogConfig},
5959
role_utils::RoleGroupRef,
@@ -63,7 +63,10 @@ use stackable_operator::{
6363
statefulset::StatefulSetConditionBuilder,
6464
},
6565
utils::COMMON_BASH_TRAP_FUNCTIONS,
66-
v2::builder::meta::ownerreference_from_resource,
66+
v2::{
67+
builder::meta::ownerreference_from_resource,
68+
types::operator::{RoleGroupName, RoleName},
69+
},
6770
};
6871
use strum::{EnumDiscriminants, IntoStaticStr};
6972

@@ -79,7 +82,6 @@ use crate::{
7982
authentication::{
8083
AirflowAuthenticationClassResolved, AirflowClientAuthenticationDetailsResolved,
8184
},
82-
build_recommended_labels,
8385
internal_secret::{
8486
FERNET_KEY_SECRET_KEY, INTERNAL_SECRET_SECRET_KEY, JWT_SECRET_SECRET_KEY,
8587
},
@@ -100,6 +102,33 @@ use crate::{
100102

101103
pub const AIRFLOW_CONTROLLER_NAME: &str = "airflowcluster";
102104
pub const CONTAINER_IMAGE_BASE_NAME: &str = "airflow";
105+
106+
/// Pseudo role/role-group names for the Kubernetes executor's resources (it is not a real
107+
/// [`AirflowRole`]). Used to derive its labels and ConfigMap name.
108+
pub const EXECUTOR_ROLE_NAME: &str = "executor";
109+
pub const EXECUTOR_ROLE_GROUP_NAME: &str = "kubernetes";
110+
111+
/// The executor pseudo-role name (`executor`) as a type-safe value.
112+
pub fn executor_role_name() -> RoleName {
113+
EXECUTOR_ROLE_NAME
114+
.parse()
115+
.expect("'executor' is a valid role name")
116+
}
117+
118+
/// The executor's role-group name (`kubernetes`), used for its role-group ConfigMap.
119+
pub fn executor_role_group_name() -> RoleGroupName {
120+
EXECUTOR_ROLE_GROUP_NAME
121+
.parse()
122+
.expect("'kubernetes' is a valid role group name")
123+
}
124+
125+
/// The executor *pod-template* role-group name (`executor-template`), used for the template
126+
/// ConfigMap/pod labels.
127+
pub fn executor_template_role_group_name() -> RoleGroupName {
128+
"executor-template"
129+
.parse()
130+
.expect("'executor-template' is a valid role group name")
131+
}
103132
pub const AIRFLOW_FULL_CONTROLLER_NAME: &str =
104133
concatcp!(AIRFLOW_CONTROLLER_NAME, '.', OPERATOR_NAME);
105134

@@ -265,9 +294,6 @@ pub enum Error {
265294
ApplyGroupListener {
266295
source: stackable_operator::cluster_resources::Error,
267296
},
268-
269-
#[snafu(display("failed to configure service"))]
270-
ServiceConfiguration { source: crate::service::Error },
271297
}
272298

273299
type Result<T, E = Error> = std::result::Result<T, E>;
@@ -439,16 +465,10 @@ pub async fn reconcile_airflow(
439465
{
440466
let rg_group_listener = build_group_listener(
441467
&validated_cluster,
442-
build_recommended_labels(
443-
airflow,
444-
AIRFLOW_CONTROLLER_NAME,
445-
&validated_cluster.image.app_version_label_value,
446-
&role_name,
447-
"none",
448-
),
468+
airflow_role,
449469
listener_class.to_string(),
450470
listener_group_name.clone(),
451-
)?;
471+
);
452472
cluster_resources
453473
.add(client, rg_group_listener)
454474
.await
@@ -476,29 +496,8 @@ pub async fn reconcile_airflow(
476496
)
477497
.context(InvalidGitSyncSpecSnafu)?;
478498

479-
let role_group_service_recommended_labels = build_recommended_labels(
480-
airflow,
481-
AIRFLOW_CONTROLLER_NAME,
482-
&validated_cluster.image.app_version_label_value,
483-
&rolegroup.role,
484-
&rolegroup.role_group,
485-
);
486-
487-
let role_group_service_selector = Labels::role_group_selector(
488-
airflow,
489-
APP_NAME,
490-
&rolegroup.role,
491-
&rolegroup.role_group,
492-
)
493-
.context(LabelBuildSnafu)?;
494-
495-
let rg_headless_service = build_rolegroup_headless_service(
496-
&validated_cluster,
497-
&rolegroup,
498-
role_group_service_recommended_labels.clone(),
499-
role_group_service_selector.clone().into(),
500-
)
501-
.context(ServiceConfigurationSnafu)?;
499+
let rg_headless_service =
500+
build_rolegroup_headless_service(&validated_cluster, airflow_role, rolegroup_name);
502501

503502
cluster_resources
504503
.add(client, rg_headless_service)
@@ -507,26 +506,25 @@ pub async fn reconcile_airflow(
507506
rolegroup: rolegroup.clone(),
508507
})?;
509508

510-
let rg_metrics_service = build_rolegroup_metrics_service(
511-
&validated_cluster,
512-
&rolegroup,
513-
role_group_service_recommended_labels,
514-
role_group_service_selector.into(),
515-
)
516-
.context(ServiceConfigurationSnafu)?;
509+
let rg_metrics_service =
510+
build_rolegroup_metrics_service(&validated_cluster, airflow_role, rolegroup_name);
517511
cluster_resources
518512
.add(client, rg_metrics_service)
519513
.await
520514
.context(ApplyRoleGroupServiceSnafu {
521515
rolegroup: rolegroup.clone(),
522516
})?;
523517

518+
let vector_config =
519+
config_map::build_vector_config(&rolegroup, &validated_rg_config.config.logging);
524520
let rg_configmap = config_map::build_rolegroup_config_map(
525521
&validated_cluster,
526-
&rolegroup,
522+
&airflow_role.role_name(),
523+
rolegroup_name,
527524
&validated_rg_config.config_overrides,
528525
&validated_rg_config.config.logging,
529526
&Container::Airflow,
527+
vector_config,
530528
)
531529
.context(BuildConfigMapSnafu)?;
532530
cluster_resources
@@ -593,18 +591,22 @@ async fn build_executor_template(
593591
.context(FailedToResolveConfigSnafu)?;
594592
let rolegroup = RoleGroupRef {
595593
cluster: ObjectRef::from_obj(airflow),
596-
role: "executor".into(),
597-
role_group: "kubernetes".into(),
594+
role: EXECUTOR_ROLE_NAME.into(),
595+
role_group: EXECUTOR_ROLE_GROUP_NAME.into(),
598596
};
599597

598+
let vector_config =
599+
config_map::build_vector_config(&rolegroup, &merged_executor_config.logging);
600600
let rg_configmap = config_map::build_rolegroup_config_map(
601601
validated_cluster,
602-
&rolegroup,
602+
&executor_role_name(),
603+
&executor_role_group_name(),
603604
// The kubernetes-executor pod template does not apply webserver_config.py overrides
604605
// (preserves prior behaviour, which passed an empty map here).
605606
&AirflowConfigOverrides::default(),
606607
&merged_executor_config.logging,
607608
&Container::Base,
609+
vector_config,
608610
)
609611
.context(BuildConfigMapSnafu)?;
610612
cluster_resources
@@ -646,48 +648,45 @@ async fn build_executor_template(
646648

647649
fn build_rolegroup_metadata(
648650
cluster: &ValidatedCluster,
649-
rolegroup: &&RoleGroupRef<v1alpha2::AirflowCluster>,
651+
role: &AirflowRole,
652+
role_group_name: &RoleGroupName,
650653
prometheus_label: Label,
651654
name: String,
652-
) -> Result<ObjectMeta, Error> {
653-
let metadata = ObjectMetaBuilder::new()
655+
) -> ObjectMeta {
656+
ObjectMetaBuilder::new()
654657
.name_and_namespace(cluster)
655658
.name(name)
656659
.ownerreference(ownerreference_from_resource(cluster, None, Some(true)))
657-
.with_recommended_labels(&build_recommended_labels(
658-
cluster,
659-
AIRFLOW_CONTROLLER_NAME,
660-
&cluster.image.app_version_label_value,
661-
&rolegroup.role,
662-
&rolegroup.role_group,
663-
))
664-
.context(ObjectMetaSnafu)?
660+
.with_labels(cluster.recommended_labels(role, role_group_name))
665661
.with_label(prometheus_label)
666-
.build();
667-
Ok(metadata)
662+
.build()
668663
}
669664

670665
pub fn build_group_listener(
671666
cluster: &ValidatedCluster,
672-
object_labels: ObjectLabels<v1alpha2::AirflowCluster>,
667+
role: &AirflowRole,
673668
listener_class: String,
674669
listener_group_name: String,
675-
) -> Result<listener::v1alpha1::Listener> {
676-
Ok(listener::v1alpha1::Listener {
670+
) -> listener::v1alpha1::Listener {
671+
listener::v1alpha1::Listener {
677672
metadata: ObjectMetaBuilder::new()
678673
.name_and_namespace(cluster)
679674
.name(listener_group_name)
680675
.ownerreference(ownerreference_from_resource(cluster, None, Some(true)))
681-
.with_recommended_labels(&object_labels)
682-
.context(ObjectMetaSnafu)?
676+
// The group listener is a role-level object, so a constant `none` role-group is used
677+
// as the role-group label value.
678+
.with_labels(cluster.recommended_labels_for(
679+
&role.role_name(),
680+
&"none".parse().expect("'none' is a valid role group name"),
681+
))
683682
.build(),
684683
spec: listener::v1alpha1::ListenerSpec {
685684
class_name: Some(listener_class),
686685
ports: Some(listener_ports()),
687686
..listener::v1alpha1::ListenerSpec::default()
688687
},
689688
status: None,
690-
})
689+
}
691690
}
692691

693692
/// We only use the http port here and intentionally omit
@@ -725,27 +724,21 @@ fn build_server_rolegroup_statefulset(
725724
let executor = &validated_cluster.cluster_config.executor;
726725

727726
let mut pb = PodBuilder::new();
728-
let recommended_object_labels = build_recommended_labels(
729-
airflow,
730-
AIRFLOW_CONTROLLER_NAME,
731-
&resolved_product_image.app_version_label_value,
732-
&rolegroup_ref.role,
733-
&rolegroup_ref.role_group,
734-
);
735-
// Used for PVC templates that cannot be modified once they are deployed
736-
let unversioned_recommended_labels = Labels::recommended(&build_recommended_labels(
737-
airflow,
738-
AIRFLOW_CONTROLLER_NAME,
739-
// A version value is required, and we do want to use the "recommended" format for the other desired labels
740-
"none",
741-
&rolegroup_ref.role,
742-
&rolegroup_ref.role_group,
743-
))
744-
.context(LabelBuildSnafu)?;
727+
let role_group_name: RoleGroupName = rolegroup_ref
728+
.role_group
729+
.parse()
730+
.expect("the role group name was validated during cluster validation");
731+
let resource_names = validated_cluster.resource_names(airflow_role, &role_group_name);
732+
733+
let recommended_object_labels =
734+
validated_cluster.recommended_labels(airflow_role, &role_group_name);
735+
// Used for PVC templates that cannot be modified once they are deployed (a constant "none"
736+
// version keeps the labels stable across version upgrades).
737+
let unversioned_recommended_labels =
738+
validated_cluster.unversioned_recommended_labels(airflow_role, &role_group_name);
745739

746740
let pb_metadata = ObjectMetaBuilder::new()
747-
.with_recommended_labels(&recommended_object_labels)
748-
.context(ObjectMetaSnafu)?
741+
.with_labels(recommended_object_labels)
749742
.with_annotation(
750743
Annotation::try_from((
751744
"kubectl.kubernetes.io/default-container",
@@ -925,7 +918,7 @@ fn build_server_rolegroup_statefulset(
925918
pb.add_volumes(airflow.volumes().clone())
926919
.context(AddVolumeSnafu)?;
927920
pb.add_volumes(controller_commons::create_volumes(
928-
&rolegroup_ref.object_name(),
921+
resource_names.role_group_config_map().as_ref(),
929922
merged_airflow_config
930923
.logging
931924
.containers
@@ -971,18 +964,14 @@ fn build_server_rolegroup_statefulset(
971964

972965
let metadata = build_rolegroup_metadata(
973966
validated_cluster,
974-
&rolegroup_ref,
967+
airflow_role,
968+
&role_group_name,
975969
restarter_label,
976-
rolegroup_ref.object_name(),
977-
)?;
970+
resource_names.stateful_set_name().to_string(),
971+
);
978972

979-
let statefulset_match_labels = Labels::role_group_selector(
980-
airflow,
981-
APP_NAME,
982-
&rolegroup_ref.role,
983-
&rolegroup_ref.role_group,
984-
)
985-
.context(BuildLabelSnafu)?;
973+
let statefulset_match_labels =
974+
validated_cluster.role_group_selector(airflow_role, &role_group_name);
986975

987976
let statefulset_spec = StatefulSetSpec {
988977
pod_management_policy: Some(
@@ -1002,7 +991,7 @@ fn build_server_rolegroup_statefulset(
1002991
match_labels: Some(statefulset_match_labels.into()),
1003992
..LabelSelector::default()
1004993
},
1005-
service_name: stateful_set_service_name(rolegroup_ref),
994+
service_name: stateful_set_service_name(validated_cluster, airflow_role, &role_group_name),
1006995
template: pod_template,
1007996
volume_claim_templates: pvcs,
1008997
..StatefulSetSpec::default()
@@ -1053,14 +1042,9 @@ fn build_executor_template_config_map(
10531042

10541043
let mut pb = PodBuilder::new();
10551044
let pb_metadata = ObjectMetaBuilder::new()
1056-
.with_recommended_labels(&build_recommended_labels(
1057-
airflow,
1058-
AIRFLOW_CONTROLLER_NAME,
1059-
&resolved_product_image.app_version_label_value,
1060-
"executor",
1061-
"executor-template",
1062-
))
1063-
.context(ObjectMetaSnafu)?
1045+
.with_labels(
1046+
cluster.recommended_labels_for(&executor_role_name(), &executor_template_role_group_name()),
1047+
)
10641048
.build();
10651049

10661050
pb.metadata(pb_metadata)
@@ -1161,14 +1145,10 @@ fn build_executor_template_config_map(
11611145
.name_and_namespace(airflow)
11621146
.name(airflow.executor_template_configmap_name())
11631147
.ownerreference(ownerreference_from_resource(cluster, None, Some(true)))
1164-
.with_recommended_labels(&build_recommended_labels(
1165-
airflow,
1166-
AIRFLOW_CONTROLLER_NAME,
1167-
&resolved_product_image.app_version_label_value,
1168-
"executor",
1169-
"executor-template",
1148+
.with_labels(cluster.recommended_labels_for(
1149+
&executor_role_name(),
1150+
&executor_template_role_group_name(),
11701151
))
1171-
.context(ObjectMetaSnafu)?
11721152
.with_label(restarter_label)
11731153
.build(),
11741154
)

0 commit comments

Comments
 (0)