Skip to content

Commit 09c8d1a

Browse files
committed
refactor: build step consumes ValidatedCluster; drop redundant builder params & raw druid
1 parent 86807e6 commit 09c8d1a

6 files changed

Lines changed: 91 additions & 109 deletions

File tree

rust/operator-binary/src/authentication/mod.rs

Lines changed: 3 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -8,12 +8,12 @@ use stackable_operator::{
88
};
99

1010
use crate::{
11+
controller::validate::ValidatedCluster,
1112
crd::{
1213
DruidRole,
1314
authentication::{AuthenticationClassResolved, AuthenticationClassesResolved},
1415
env_var_reference,
1516
security::INTERNAL_INITIAL_CLIENT_PASSWORD_ENV,
16-
v1alpha1,
1717
},
1818
internal_secret::{build_shared_internal_secret_name, env_var_from_secret},
1919
};
@@ -166,13 +166,9 @@ impl DruidAuthenticationConfig {
166166
command
167167
}
168168

169-
pub fn get_env_var_mounts(
170-
&self,
171-
druid: &v1alpha1::DruidCluster,
172-
role: &DruidRole,
173-
) -> Vec<EnvVar> {
169+
pub fn get_env_var_mounts(&self, cluster: &ValidatedCluster, role: &DruidRole) -> Vec<EnvVar> {
174170
let mut envs = vec![];
175-
let internal_secret_name = build_shared_internal_secret_name(druid);
171+
let internal_secret_name = build_shared_internal_secret_name(cluster);
176172
envs.push(env_var_from_secret(
177173
&internal_secret_name,
178174
None,

rust/operator-binary/src/controller.rs

Lines changed: 29 additions & 71 deletions
Original file line numberDiff line numberDiff line change
@@ -16,10 +16,8 @@ use stackable_operator::{
1616
},
1717
cli::OperatorEnvironmentOptions,
1818
cluster_resources::{ClusterResourceApplyStrategy, ClusterResources},
19-
commons::{product_image_selection::ResolvedProductImage, rbac::build_rbac_resources},
19+
commons::rbac::build_rbac_resources,
2020
constants::RESTART_CONTROLLER_ENABLED_LABEL,
21-
crd::s3,
22-
database_connections::drivers::jdbc::JdbcDatabaseConnection as _,
2321
k8s_openapi::{
2422
DeepMerge,
2523
api::{
@@ -54,7 +52,6 @@ use stackable_operator::{
5452
use strum::{EnumDiscriminants, IntoStaticStr};
5553

5654
use crate::{
57-
authentication::DruidAuthenticationConfig,
5855
controller::build::resource::{
5956
listener::{
6057
LISTENER_VOLUME_DIR, LISTENER_VOLUME_NAME, build_group_listener,
@@ -67,15 +64,15 @@ use crate::{
6764
APP_NAME, Container, DRUID_CONFIG_DIRECTORY, DeepStorageSpec, DruidClusterStatus,
6865
DruidRole, HDFS_CONFIG_DIRECTORY, LOG_CONFIG_DIRECTORY, METRICS_PORT, METRICS_PORT_NAME,
6966
OPERATOR_NAME, RW_CONFIG_DIRECTORY, STACKABLE_LOG_DIR, ValidatedDruidConfig,
70-
build_recommended_labels, security::DruidTlsSecurity, v1alpha1,
67+
build_recommended_labels, v1alpha1,
7168
},
7269
internal_secret::create_shared_internal_secret,
7370
operations::graceful_shutdown::add_graceful_shutdown_config,
7471
};
7572

7673
mod build;
7774
mod dereference;
78-
mod validate;
75+
pub(crate) mod validate;
7976

8077
use build::{
8178
properties::product_logging::MAX_DRUID_LOG_FILES_SIZE,
@@ -258,11 +255,6 @@ pub enum Error {
258255
BuildConfigMap {
259256
source: build::resource::config_map::Error,
260257
},
261-
262-
#[snafu(display("invalid metadata database connection"))]
263-
InvalidMetadataDatabaseConnection {
264-
source: stackable_operator::database_connections::Error,
265-
},
266258
}
267259

268260
type Result<T, E = Error> = std::result::Result<T, E>;
@@ -332,35 +324,12 @@ pub async fn reconcile_druid(
332324
.context(FailedInternalSecretCreationSnafu)?;
333325

334326
for (rolegroup_name, rg) in groups.iter() {
335-
let role_group_service_recommended_labels = build_recommended_labels(
336-
druid,
337-
DRUID_CONTROLLER_NAME,
338-
&validated_cluster.image.app_version_label_value,
339-
&role_name,
340-
rolegroup_name.as_ref(),
341-
);
342-
343-
let role_group_service_selector =
344-
Labels::role_group_selector(druid, APP_NAME, &role_name, rolegroup_name.as_ref())
345-
.context(LabelBuildSnafu)?;
346-
347-
let rg_headless_service = build_rolegroup_headless_service(
348-
&validated_cluster,
349-
&validated_cluster.cluster_config.druid_tls_security,
350-
druid_role,
351-
rolegroup_name,
352-
role_group_service_recommended_labels.clone(),
353-
role_group_service_selector.clone().into(),
354-
)
355-
.context(ServiceConfigurationSnafu)?;
356-
let rg_metrics_service = build_rolegroup_metrics_service(
357-
&validated_cluster,
358-
druid_role,
359-
rolegroup_name,
360-
role_group_service_recommended_labels,
361-
role_group_service_selector.into(),
362-
)
363-
.context(ServiceConfigurationSnafu)?;
327+
let rg_headless_service =
328+
build_rolegroup_headless_service(&validated_cluster, druid_role, rolegroup_name)
329+
.context(ServiceConfigurationSnafu)?;
330+
let rg_metrics_service =
331+
build_rolegroup_metrics_service(&validated_cluster, druid_role, rolegroup_name)
332+
.context(ServiceConfigurationSnafu)?;
364333

365334
let rg_configmap = build::resource::config_map::build_rolegroup_config_map(
366335
&validated_cluster,
@@ -370,15 +339,10 @@ pub async fn reconcile_druid(
370339
)
371340
.context(BuildConfigMapSnafu)?;
372341
let rg_statefulset = build_rolegroup_statefulset(
373-
druid,
374342
&validated_cluster,
375-
&validated_cluster.image,
376343
druid_role,
377344
rolegroup_name,
378345
rg,
379-
validated_cluster.cluster_config.s3_connection.as_ref(),
380-
&validated_cluster.cluster_config.druid_tls_security,
381-
&validated_cluster.cluster_config.druid_auth_config,
382346
&rbac_sa,
383347
)?;
384348

@@ -415,12 +379,12 @@ pub async fn reconcile_druid(
415379
}
416380

417381
if let Some(listener_class) = druid_role.listener_class_name(druid)
418-
&& let Some(listener_group_name) = group_listener_name(druid, druid_role)
382+
&& let Some(listener_group_name) = group_listener_name(&validated_cluster, druid_role)
419383
{
420384
let role_group_listener = build_group_listener(
421-
druid,
385+
&validated_cluster,
422386
build_recommended_labels(
423-
druid,
387+
&validated_cluster,
424388
DRUID_CONTROLLER_NAME,
425389
&validated_cluster.image.app_version_label_value,
426390
&role_name,
@@ -429,7 +393,6 @@ pub async fn reconcile_druid(
429393
listener_class.to_string(),
430394
listener_group_name,
431395
druid_role,
432-
&validated_cluster.cluster_config.druid_tls_security,
433396
)
434397
.context(ListenerConfigurationSnafu)?;
435398

@@ -484,26 +447,26 @@ pub async fn reconcile_druid(
484447
Ok(Action::await_change())
485448
}
486449

487-
#[allow(clippy::too_many_arguments)]
488450
/// The rolegroup [`StatefulSet`] runs the rolegroup, as configured by the administrator.
489451
///
490452
/// The [`Pod`](`stackable_operator::k8s_openapi::api::core::v1::Pod`)s are accessible through the
491453
/// corresponding [`stackable_operator::k8s_openapi::api::core::v1::Service`] (from [`build_rolegroup_headless_service`]).
492454
fn build_rolegroup_statefulset(
493-
druid: &v1alpha1::DruidCluster,
494455
cluster: &ValidatedCluster,
495-
resolved_product_image: &ResolvedProductImage,
496456
role: &DruidRole,
497457
role_group_name: &RoleGroupName,
498458
rg: &DruidRoleGroupConfig,
499-
s3_conn: Option<&s3::v1alpha1::ConnectionSpec>,
500-
druid_tls_security: &DruidTlsSecurity,
501-
druid_auth_config: &Option<DruidAuthenticationConfig>,
502459
service_account: &ServiceAccount,
503460
) -> Result<StatefulSet> {
504461
let merged_rolegroup_config = &rg.config;
505462
let role_name = role.to_string();
506463
let resource_names = cluster.resource_names(role, role_group_name);
464+
// Everything below used to be threaded in as separate parameters; it all lives on the
465+
// `ValidatedCluster` now.
466+
let resolved_product_image = &cluster.image;
467+
let s3_conn = cluster.cluster_config.s3_connection.as_ref();
468+
let druid_tls_security = &cluster.cluster_config.druid_tls_security;
469+
let druid_auth_config = &cluster.cluster_config.druid_auth_config;
507470
// prepare container builder
508471
let prepare_container_name = Container::Prepare.to_string();
509472
let mut cb_prepare = ContainerBuilder::new(&prepare_container_name).context(
@@ -530,12 +493,7 @@ fn build_rolegroup_statefulset(
530493
)
531494
.context(GracefulShutdownSnafu)?;
532495

533-
let metadata_database_connection_details = druid
534-
.spec
535-
.cluster_config
536-
.metadata_database
537-
.jdbc_connection_details("metadata")
538-
.context(InvalidMetadataDatabaseConnectionSnafu)?;
496+
let metadata_database_connection_details = &cluster.cluster_config.metadata_db_connection;
539497

540498
let mut main_container_commands = role.main_container_prepare_commands(s3_conn);
541499
let mut prepare_container_commands = vec![];
@@ -589,7 +547,7 @@ fn build_rolegroup_statefulset(
589547
)?;
590548
add_log_volume_and_volume_mounts(&mut cb_druid, &mut cb_prepare, &mut pb)?;
591549
add_hdfs_cm_volume_and_volume_mounts(
592-
&druid.spec.cluster_config.deep_storage,
550+
&cluster.cluster_config.deep_storage,
593551
&mut cb_druid,
594552
&mut pb,
595553
)?;
@@ -623,7 +581,7 @@ fn build_rolegroup_statefulset(
623581
let mut rest_env: Vec<EnvVar> = rg.env_overrides.clone().into();
624582

625583
if let Some(auth_config) = druid_auth_config {
626-
rest_env.extend(auth_config.get_env_var_mounts(druid, role))
584+
rest_env.extend(auth_config.get_env_var_mounts(cluster, role))
627585
}
628586

629587
// Needed for the `containerdebug` process to log it's tracing information to.
@@ -662,7 +620,7 @@ fn build_rolegroup_statefulset(
662620
// Known roles are MiddleManagers for ingestion and Historicals for deep storage (GCS plugin)
663621
// We may at some time in the future revisit this and limit it again to avoid needlessly
664622
// propagating potentially confidential files throughout the cluster
665-
for volume in &druid.spec.cluster_config.extra_volumes {
623+
for volume in &cluster.cluster_config.extra_volumes {
666624
// Extract values into vars so we make it impossible to log something other than
667625
// what we actually use to create the mounts - maybe paranoid, but hey ..
668626
let volume_name = &volume.name;
@@ -682,14 +640,14 @@ fn build_rolegroup_statefulset(
682640

683641
let mut pvcs: Option<Vec<PersistentVolumeClaim>> = None;
684642

685-
if let Some(group_listener_name) = group_listener_name(druid, role) {
643+
if let Some(group_listener_name) = group_listener_name(cluster, role) {
686644
cb_druid
687645
.add_volume_mount(LISTENER_VOLUME_NAME, LISTENER_VOLUME_DIR)
688646
.context(AddVolumeMountSnafu)?;
689647

690648
// Used for PVC templates that cannot be modified once they are deployed
691649
let unversioned_recommended_labels = Labels::recommended(&build_recommended_labels(
692-
druid,
650+
cluster,
693651
DRUID_CONTROLLER_NAME,
694652
// A version value is required, and we do want to use the "recommended" format for the other desired labels
695653
"none",
@@ -706,7 +664,7 @@ fn build_rolegroup_statefulset(
706664

707665
let metadata = ObjectMetaBuilder::new()
708666
.with_recommended_labels(&build_recommended_labels(
709-
druid,
667+
cluster,
710668
DRUID_CONTROLLER_NAME,
711669
&resolved_product_image.app_version_label_value,
712670
&role_name,
@@ -743,12 +701,12 @@ fn build_rolegroup_statefulset(
743701

744702
Ok(StatefulSet {
745703
metadata: ObjectMetaBuilder::new()
746-
.name_and_namespace(druid)
704+
.name_and_namespace(cluster)
747705
.name(resource_names.stateful_set_name().to_string())
748-
.ownerreference_from_resource(druid, None, Some(true))
706+
.ownerreference_from_resource(cluster, None, Some(true))
749707
.context(ObjectMissingMetadataForOwnerRefSnafu)?
750708
.with_recommended_labels(&build_recommended_labels(
751-
druid,
709+
cluster,
752710
DRUID_CONTROLLER_NAME,
753711
&resolved_product_image.app_version_label_value,
754712
&role_name,
@@ -763,7 +721,7 @@ fn build_rolegroup_statefulset(
763721
selector: LabelSelector {
764722
match_labels: Some(
765723
Labels::role_group_selector(
766-
druid,
724+
cluster,
767725
APP_NAME,
768726
&role_name,
769727
role_group_name.as_ref(),

rust/operator-binary/src/controller/build/resource/listener.rs

Lines changed: 18 additions & 16 deletions
Original file line numberDiff line numberDiff line change
@@ -6,14 +6,15 @@ use stackable_operator::{
66
},
77
crd::listener::{self, v1alpha1::Listener},
88
k8s_openapi::api::core::v1::PersistentVolumeClaim,
9-
kube::ResourceExt,
109
kvp::{Labels, ObjectLabels},
1110
};
1211

13-
use crate::crd::{
14-
DruidRole,
15-
security::{DruidTlsSecurity, PLAINTEXT_PORT_NAME, TLS_PORT_NAME},
16-
v1alpha1,
12+
use crate::{
13+
controller::validate::ValidatedCluster,
14+
crd::{
15+
DruidRole,
16+
security::{DruidTlsSecurity, PLAINTEXT_PORT_NAME, TLS_PORT_NAME},
17+
},
1718
};
1819

1920
pub const LISTENER_VOLUME_NAME: &str = "listener";
@@ -47,25 +48,29 @@ pub enum Error {
4748
}
4849

4950
pub fn build_group_listener(
50-
druid: &v1alpha1::DruidCluster,
51-
object_labels: ObjectLabels<v1alpha1::DruidCluster>,
51+
cluster: &ValidatedCluster,
52+
object_labels: ObjectLabels<ValidatedCluster>,
5253
listener_class: String,
5354
listener_group_name: String,
5455
druid_role: &DruidRole,
55-
druid_tls_security: &DruidTlsSecurity,
5656
) -> Result<Listener, Error> {
5757
Ok(Listener {
5858
metadata: ObjectMetaBuilder::new()
59-
.name_and_namespace(druid)
59+
.name_and_namespace(cluster)
6060
.name(listener_group_name)
61-
.ownerreference_from_resource(druid, None, Some(true))
61+
.ownerreference_from_resource(cluster, None, Some(true))
6262
.context(ObjectMissingMetadataForOwnerRefSnafu)?
6363
.with_recommended_labels(&object_labels)
6464
.context(BuildObjectMetaSnafu)?
6565
.build(),
6666
spec: listener::v1alpha1::ListenerSpec {
6767
class_name: Some(listener_class),
68-
ports: Some(druid_tls_security.listener_ports(druid_role)),
68+
ports: Some(
69+
cluster
70+
.cluster_config
71+
.druid_tls_security
72+
.listener_ports(druid_role),
73+
),
6974
..listener::v1alpha1::ListenerSpec::default()
7075
},
7176
status: None,
@@ -84,14 +89,11 @@ pub fn build_group_listener_pvc(
8489
.context(BuildListenerPersistentVolumeSnafu)
8590
}
8691

87-
pub fn group_listener_name(
88-
druid: &v1alpha1::DruidCluster,
89-
druid_role: &DruidRole,
90-
) -> Option<String> {
92+
pub fn group_listener_name(cluster: &ValidatedCluster, druid_role: &DruidRole) -> Option<String> {
9193
match druid_role {
9294
DruidRole::Coordinator | DruidRole::Broker | DruidRole::Router => Some(format!(
9395
"{cluster_name}-{druid_role}",
94-
cluster_name = druid.name_any(),
96+
cluster_name = cluster.name,
9597
)),
9698
DruidRole::Historical | DruidRole::MiddleManager => None,
9799
}

0 commit comments

Comments
 (0)