Skip to content
3 changes: 3 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -11,10 +11,13 @@ All notable changes to this project will be documented in this file.
- The RBAC ServiceAccount and RoleBinding are now built with the operator-rs `v2::rbac`
functions and carry the full set of recommended labels ([#846]).
- Bump stackable-operator to 0.114.0 ([#855]).
- The reconciler now applies resources and derives the cluster status in discrete
apply and update_status steps ([#856]).

[#841]: https://github.com/stackabletech/druid-operator/pull/841
[#846]: https://github.com/stackabletech/druid-operator/pull/846
[#855]: https://github.com/stackabletech/druid-operator/pull/855
[#856]: https://github.com/stackabletech/druid-operator/pull/856

## [26.7.0] - 2026-07-21

Expand Down
260 changes: 113 additions & 147 deletions rust/operator-binary/src/controller.rs
Original file line number Diff line number Diff line change
@@ -1,7 +1,7 @@
//! Ensures that `Pod`s are configured and running for each [`DruidCluster`][v1alpha1]
//!
//! [v1alpha1]: v1alpha1::DruidCluster
use std::{str::FromStr, sync::Arc};
use std::{marker::PhantomData, str::FromStr, sync::Arc};

use const_format::concatcp;
use snafu::{ResultExt, Snafu};
Expand All @@ -21,25 +21,23 @@ use stackable_operator::{
},
logging::controller::ReconcilerError,
shared::time::Duration,
status::condition::{
compute_conditions, operations::ClusterOperationsConditionBuilder,
statefulset::StatefulSetConditionBuilder,
},
v2::{
cluster_resources::cluster_resources_new,
types::operator::{ControllerName, OperatorName, ProductName},
},
v2::types::operator::{ControllerName, OperatorName, ProductName},
};
use strum::{EnumDiscriminants, IntoStaticStr};

use crate::{
controller::build::resource::listener::{build_group_listener, group_listener_name},
crd::{APP_NAME, DruidClusterStatus, DruidRole, OPERATOR_NAME, v1alpha1},
internal_secret::create_shared_internal_secret,
controller::{
apply::{Applier, ensure_internal_secret},
build::resource::listener::group_listener_name,
update_status::update_status,
},
crd::{APP_NAME, DruidRole, OPERATOR_NAME, v1alpha1},
};

mod apply;
mod build;
mod dereference;
mod update_status;
pub(crate) mod validate;

use build::resource::discovery::{self, build_discovery_configmaps};
Expand Down Expand Up @@ -70,68 +68,78 @@ pub struct Ctx {
pub operator_environment: OperatorEnvironmentOptions,
}

/// Every Kubernetes resource produced by the client-free [`build`](build::build) step.
/// Marker for prepared Kubernetes resources which are not applied yet.
pub struct Prepared;

/// Marker for applied Kubernetes resources.
pub struct Applied;

/// Every Kubernetes resource produced by the build step.
///
/// The discovery `ConfigMap`s are intentionally *not* part of this bundle: they derive from the
/// *applied* Router listener's ingress address, so they are built after the first apply phase
/// and applied through the same [`apply::Applier`] before its orphan deletion runs.
///
/// The Router group `Listener` and the discovery `ConfigMap`s are intentionally *not* yet part of this
/// bundle: the discovery `ConfigMap` derives from the *applied* Router listener's ingress address,
/// so both are built and applied in the reconcile step instead.
pub struct KubernetesResources {
/// `T` is a marker that indicates if these resources are only [`Prepared`] or already [`Applied`].
/// The marker is useful e.g. to ensure that the cluster status is updated based on the applied
/// resources.
pub struct KubernetesResources<T> {
pub stateful_sets: Vec<StatefulSet>,
pub services: Vec<Service>,
pub listeners: Vec<Listener>,
pub config_maps: Vec<ConfigMap>,
pub pod_disruption_budgets: Vec<PodDisruptionBudget>,
pub service_accounts: Vec<ServiceAccount>,
pub role_bindings: Vec<RoleBinding>,
pub status: PhantomData<T>,
}

impl KubernetesResources<Applied> {
/// The applied group [`Listener`] of the given role, if the role exposes one.
pub fn group_listener(
&self,
cluster: &validate::ValidatedCluster,
role: &DruidRole,
) -> Option<&Listener> {
let listener_name = group_listener_name(cluster, role)?;
self.listeners
.iter()
.find(|listener| listener.metadata.name.as_deref() == Some(listener_name.as_ref()))
}
}

#[derive(Snafu, Debug, EnumDiscriminants)]
#[strum_discriminants(derive(IntoStaticStr))]
pub enum Error {
#[snafu(display("failed to apply the Kubernetes resources"))]
ApplyResources { source: apply::Error },

#[snafu(display("failed to build the Kubernetes resources"))]
BuildResources { source: build::Error },

#[snafu(display("failed to apply Kubernetes resource"))]
ApplyResource {
source: stackable_operator::cluster_resources::Error,
},

#[snafu(display("failed to dereference cluster objects"))]
Dereference { source: dereference::Error },

#[snafu(display("failed to build discovery ConfigMap"))]
BuildDiscoveryConfig { source: discovery::Error },

#[snafu(display("failed to apply discovery ConfigMap"))]
ApplyDiscoveryConfig {
source: stackable_operator::cluster_resources::Error,
},
ApplyDiscoveryConfig { source: apply::Error },

#[snafu(display("failed to apply cluster status"))]
ApplyStatus {
source: stackable_operator::client::Error,
},
#[snafu(display("failed to update the cluster status"))]
UpdateStatus { source: update_status::Error },

#[snafu(display("failed to delete orphaned resources"))]
DeleteOrphanedResources {
source: stackable_operator::cluster_resources::Error,
},
DeleteOrphanedResources { source: apply::Error },

#[snafu(display("failed to retrieve secret for internal communications"))]
FailedInternalSecretCreation {
source: crate::internal_secret::Error,
},
FailedInternalSecretCreation { source: apply::Error },

#[snafu(display("DruidCluster object is invalid"))]
InvalidDruidCluster {
source: error_boundary::InvalidObject,
},

#[snafu(display("failed to apply group listener"))]
ApplyGroupListener {
source: stackable_operator::cluster_resources::Error,
},

#[snafu(display("failed to validate cluster"))]
ValidateCluster { source: validate::Error },
}
Expand Down Expand Up @@ -165,124 +173,46 @@ pub async fn reconcile_druid(
validate::validate(druid, &dereferenced_objects, &ctx.operator_environment)
.context(ValidateClusterSnafu)?;

let mut cluster_resources = cluster_resources_new(
&product_name(),
&operator_name(),
&controller_name(),
&validated_cluster.name,
&validated_cluster.namespace,
&validated_cluster.uid,
ClusterResourceApplyStrategy::from(&druid.spec.cluster_operation),
&druid.spec.object_overrides,
);
let resources = build::build(&validated_cluster).context(BuildResourcesSnafu)?;

// The internal secret is shared across all roles and role groups, so it only needs to be
// created once per reconcile rather than inside the role loop below.
create_shared_internal_secret(&validated_cluster, client, DRUID_CONTROLLER_NAME)
ensure_internal_secret(client, &validated_cluster)
.await
.context(FailedInternalSecretCreationSnafu)?;

let resources = build::build(&validated_cluster).context(BuildResourcesSnafu)?;

let mut ss_cond_builder = StatefulSetConditionBuilder::default();

for service_account in resources.service_accounts {
cluster_resources
.add(client, service_account)
.await
.context(ApplyResourceSnafu)?;
}

for role_binding in resources.role_bindings {
cluster_resources
.add(client, role_binding)
.await
.context(ApplyResourceSnafu)?;
}

// Apply order: everything a Pod mounts (ConfigMaps) must exist before the StatefulSets, so the
// StatefulSets are applied last to prevent unnecessary Pod restarts.
// See https://github.com/stackabletech/commons-operator/issues/111 for details.
for service in resources.services {
cluster_resources
.add(client, service)
.await
.context(ApplyResourceSnafu)?;
}
for listener in resources.listeners {
cluster_resources
.add(client, listener)
.await
.context(ApplyResourceSnafu)?;
}
for config_map in resources.config_maps {
cluster_resources
.add(client, config_map)
.await
.context(ApplyResourceSnafu)?;
}
for pdb in resources.pod_disruption_budgets {
cluster_resources
.add(client, pdb)
.await
.context(ApplyResourceSnafu)?;
}
for stateful_set in resources.stateful_sets {
ss_cond_builder.add(
cluster_resources
.add(client, stateful_set)
.await
.context(ApplyResourceSnafu)?,
);
}

// The Router group Listener and its discovery ConfigMaps are applied here rather than in the
// build step: the discovery ConfigMap derives from the *applied* Router listener's ingress
// address, which is only known after the Listener has been applied.
if let Some(listener_class) = &validated_cluster
.role_config(&DruidRole::Router)
.listener_class
&& let Some(listener_group_name) =
group_listener_name(&validated_cluster, &DruidRole::Router)
{
let router_listener = build_group_listener(
&validated_cluster,
listener_class,
listener_group_name,
&DruidRole::Router,
);

let listener = cluster_resources
.add(client, router_listener)
.await
.context(ApplyGroupListenerSnafu)?;
let mut applier = Applier::new(
client,
&validated_cluster,
ClusterResourceApplyStrategy::from(&druid.spec.cluster_operation),
&druid.spec.object_overrides,
);

for discovery_cm in build_discovery_configmaps(&validated_cluster, listener)
let applied = applier
.apply(resources)
.await
.context(ApplyResourcesSnafu)?;

// Second apply phase: the discovery ConfigMaps derive from the *applied* Router listener's
// ingress address, which is only known after the Listener has been applied. The Router
// listener itself is built in `build()` and applied with all the other listeners in the
// first apply phase above; here it is only read back. The discovery ConfigMaps must go
// through the same Applier, so that the orphan deletion in `finish` sees them.

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

The discovery ConfigMap should be built in the build step, where it belongs.

The current comment is also inaccurate. It claims the discovery ConfigMap can be built from the applied Listener object, but what actually happens is:

  1. The build step runs without building the discovery ConfigMap.
  2. The Listener is applied, but the returned Listener object does not yet carry a status.
  3. Building the discovery ConfigMap fails with:
    ▎ failed to build discovery ConfigMap, failed to configure listener discovery configmap, router listener has no adress [sic]
  4. Reconciliation aborts: orphaned resources are not cleaned up and the cluster status is not updated.
  5. The Listener status is eventually set, which triggers another reconciliation.
  6. The build step runs again, still without building the discovery ConfigMap.
  7. …and so on.

Instead, we can simply:

  1. Dereference the Listener.
  2. Build the discovery ConfigMap in the build step if the Listener status is set; otherwise, skip it. (As a bonus, we no longer emit a misleading ERROR log for what is actually an expected, transient state.)
  3. Apply everything, update the status, and wait for the next reconciliation (once the Listener status is finally populated).

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

…and so on.

It must eventually converge otherwise the smoke test would not work (but, yes, we emit a few unecessary ERRORs until that happens). And the fact that we error out of the reconcile means it gets triggered again, whereas if we just miss out the ConfigMap when there is no status, there is no guarantee it gets retriggered unless we watch listeners, too. So should we add a watch on listeners (with an RBAC change, too), or return a short requeue if the status is not yet available?

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I would watch the Listeners. When the discovery ConfigMap is properly integrated into the controller steps then a lot of code can be removed that handled this special case.

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

See 5e6d481 and 173a9d0.
Re-tested:

--- PASS: kuttl (520.13s)
    --- PASS: kuttl/harness (0.00s)
        --- PASS: kuttl/harness/smoke_druid-37.0.0_zookeeper-3.9.5_hadoop-3.5.0_openshift-false (520.12s)
PASS

if let Some(router_listener) = applied.group_listener(&validated_cluster, &DruidRole::Router) {
let discovery_config_maps = build_discovery_configmaps(&validated_cluster, router_listener)
.context(BuildDiscoveryConfigSnafu)?;
applier
.apply_config_maps(discovery_config_maps)
.await
.context(BuildDiscoveryConfigSnafu)?
{
cluster_resources
.add(client, discovery_cm)
.await
.context(ApplyDiscoveryConfigSnafu)?;
}
.context(ApplyDiscoveryConfigSnafu)?;
}

let cluster_operation_cond_builder =
ClusterOperationsConditionBuilder::new(&druid.spec.cluster_operation);

let status = DruidClusterStatus {
conditions: compute_conditions(druid, &[&ss_cond_builder, &cluster_operation_cond_builder]),
};

cluster_resources
.delete_orphaned_resources(client)
applier
.finish()
.await
.context(DeleteOrphanedResourcesSnafu)?;
client
.apply_patch_status(OPERATOR_NAME, druid, &status)

update_status(client, druid, &applied)
.await
.context(ApplyStatusSnafu)?;
.context(UpdateStatusSnafu)?;

Ok(Action::await_change())
}
Expand Down Expand Up @@ -313,6 +243,42 @@ mod test {
crd::PROP_SEGMENT_CACHE_LOCATIONS,
};

/// `group_listener` finds the applied group Listener of a role by its name, and yields
/// `None` for roles that expose no Listener (here: Historical).
#[test]
fn group_listener_finds_the_role_listener_by_name() {
let druid = crate::controller::validate::test_support::druid_from_yaml(
crate::controller::validate::test_support::MINIMAL_DRUID_YAML,
);
let cluster = crate::controller::validate::test_support::validated_cluster(&druid);
let built = build::build(&cluster).expect("the test cluster builds");
// `KubernetesResources<Applied>` normally only exists after an apply; constructing it
// from the built resources is fine here because the lookup only reads names.
let applied = KubernetesResources::<Applied> {
stateful_sets: built.stateful_sets,
services: built.services,
listeners: built.listeners,
config_maps: built.config_maps,
pod_disruption_budgets: built.pod_disruption_budgets,
service_accounts: built.service_accounts,
role_bindings: built.role_bindings,
status: std::marker::PhantomData,
};

let router_listener = applied
.group_listener(&cluster, &DruidRole::Router)
.expect("the Router exposes a group Listener");
assert_eq!(
router_listener.metadata.name.as_deref(),
Some("simple-druid-router")
);
assert!(
applied
.group_listener(&cluster, &DruidRole::Historical)
.is_none()
);
}

#[rstest]
#[case(
"segment_cache.yaml",
Expand Down
Loading
Loading