Skip to content
Merged
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
637 changes: 338 additions & 299 deletions Cargo.lock

Large diffs are not rendered by default.

8 changes: 8 additions & 0 deletions crates/stackable-operator/CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -4,6 +4,14 @@ All notable changes to this project will be documented in this file.

## [Unreleased]

### Removed

- BREAKING: Removed `timeout_duration` parameter from `signal::crd_established`. It now waits indefinitely and
the timeout should be handled with a startup probe instead. As a consequence `DEFAULT_CRD_ESTABLISHED_TIMEOUT`
also got removed ([#1272]).

[#1272]: https://github.com/stackabletech/operator-rs/pull/1272

## [0.118.0] - 2026-09-14

### Added
Expand Down
53 changes: 24 additions & 29 deletions crates/stackable-operator/src/utils/signal.rs
Original file line number Diff line number Diff line change
@@ -1,7 +1,6 @@
use k8s_openapi::apiextensions_apiserver::pkg::apis::apiextensions::v1::CustomResourceDefinition;
use kube::runtime::wait;
use snafu::{ResultExt, Snafu};
use stackable_shared::time::Duration;
use tokio::{
signal::unix::{SignalKind, signal},
sync::watch,
Expand Down Expand Up @@ -77,45 +76,41 @@ impl SignalWatcher<()> {
}
}

pub const DEFAULT_CRD_ESTABLISHED_TIMEOUT: Duration = Duration::from_secs(5);

#[derive(Debug, Snafu)]
pub enum CrdEstablishedError {
#[snafu(display("failed to meet CRD established condition before the timeout elapsed"))]
TimeoutElapsed { source: tokio::time::error::Elapsed },

#[snafu(display("failed to await CRD established condition due to api error"))]
Api { source: kube::runtime::wait::Error },
#[snafu(display(
"failed to await CRD established condition for {crd_name:?} due to api error"
))]
Api {
source: kube::runtime::wait::Error,
crd_name: String,
},
}

/// Waits for a CRD named `crd_name` to be established before `timeout_duration` (or by default
/// [`DEFAULT_CRD_ESTABLISHED_TIMEOUT`]) is elapsed.
/// Waits for a CRD named `crd_name` to be established.
///
/// The same caveats from [`conditions::is_crd_established`](wait::conditions::is_crd_established)
/// apply here as well.
///
/// ### Errors
///
/// This function returns errors either if the timeout elapsed without the condition being met or
/// when the underlying API returned errors (CRD is unknown to the Kubernetes API server or due to
/// missing permissions).
pub async fn crd_established(
client: &Client,
crd_name: &str,
timeout_duration: impl Into<Option<Duration>>,
) -> Result<(), CrdEstablishedError> {
/// This function returns errors when the underlying API returned errors
/// (CRD is unknown to the Kubernetes API server or due to missing permissions).
pub async fn crd_established(client: &Client, crd_name: &str) -> Result<(), CrdEstablishedError> {
tracing::info!(
k8s.crd.name = crd_name,
"Waiting for the custom resource definition to be established"
);

let api: kube::Api<CustomResourceDefinition> = client.get_api(&());
let crd_established =
wait::await_condition(api, crd_name, wait::conditions::is_crd_established());
let _ = tokio::time::timeout(
*timeout_duration
.into()
.unwrap_or(DEFAULT_CRD_ESTABLISHED_TIMEOUT),
crd_established,
)
.await
.context(TimeoutElapsedSnafu)?
.context(ApiSnafu)?;
wait::await_condition(api, crd_name, wait::conditions::is_crd_established())
.await
.context(ApiSnafu { crd_name })?;

tracing::info!(
k8s.crd.name = crd_name,
"The custom resource definition is established"
);

Ok(())
}
8 changes: 8 additions & 0 deletions crates/stackable-webhook/CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -4,6 +4,14 @@ All notable changes to this project will be documented in this file.

## [Unreleased]

### Added

- Add `health` module containing `HealthCheck` and `HealthCheckRegistry` for health endpoints ([#1272]).
- BREAKING: The `WebhookServer` now serves a `/ready` endpoint for a startup probe.
For that, `WebhookServer::new` takes an additional `HealthCheckRegistry` argument ([#1272]).

[#1272]: https://github.com/stackabletech/operator-rs/pull/1272

## [0.9.2] - 2026-07-06

Note: There are only dependency bumps in this release.
Expand Down
164 changes: 164 additions & 0 deletions crates/stackable-webhook/src/health.rs
Original file line number Diff line number Diff line change
@@ -0,0 +1,164 @@
//! Health checks and health check registry for health endpoints used by probes
//!
//! The naming here follows the Kubernetes convention of a health check / health check registry and
//! their usage by the different probes.
//!
//! ## References
//!
//! - <https://github.com/kubernetes/kubernetes/blob/cc213a13bd4e8f0f253087187ffd3e1264d33f3b/staging/src/k8s.io/apiserver/pkg/server/healthz/healthz.go#L41>
//! - <https://github.com/kubernetes/kubernetes/blob/cc213a13bd4e8f0f253087187ffd3e1264d33f3b/staging/src/k8s.io/apiserver/pkg/server/healthz.go#L34>
//! - <https://github.com/kubernetes/kubernetes/blob/cc213a13bd4e8f0f253087187ffd3e1264d33f3b/staging/src/k8s.io/apiserver/pkg/server/genericapiserver.go#L205>
//! - <https://github.com/kubernetes/kubernetes/blob/cc213a13bd4e8f0f253087187ffd3e1264d33f3b/staging/src/k8s.io/apiserver/pkg/server/healthz/healthz.go#L330>
use std::{
fmt::Display,
sync::{
Arc,
atomic::{AtomicBool, Ordering},
},
};

use axum::{
http::StatusCode,
response::{IntoResponse, Response},
};

/// A single named check contributing to one health endpoint.
///
/// A check can be marked as passing with a call to [`HealthCheck::mark_passed`], and reset with
/// [`HealthCheck::mark_not_passed`], e.g. for a liveness check that can start failing again.
#[derive(Clone)]
pub struct HealthCheck {
name: String,
// This has to be an AtomicBool as we could otherwise not share references to it.
passed: Arc<AtomicBool>,
}
Comment thread
Techassi marked this conversation as resolved.

impl HealthCheck {
fn new(name: impl Into<String>) -> Self {
Self {
name: name.into(),
passed: Arc::new(AtomicBool::new(false)),
}
}

pub fn mark_passed(&self) {
self.passed.store(true, Ordering::Release);
}

pub fn mark_not_passed(&self) {
self.passed.store(false, Ordering::Release);
}

fn passed(&self) -> bool {
self.passed.load(Ordering::Acquire)
}
}
Comment thread
Techassi marked this conversation as resolved.

/// A set of checks to be used for a health endpoint a probe can call.
///
/// # Example
///
/// ```
/// use stackable_webhook::health::HealthCheckRegistry;
///
/// let mut startup_checks = HealthCheckRegistry::new();
/// let crds_established = startup_checks.register("crds-established");
///
/// assert!(!startup_checks.all_passed());
/// crds_established.mark_passed();
/// assert!(startup_checks.all_passed());
/// ```
#[derive(Default)]
pub struct HealthCheckRegistry {
checks: Vec<HealthCheck>,
}
Comment thread
Techassi marked this conversation as resolved.

impl HealthCheckRegistry {
/// Creates a new [`HealthCheckRegistry`] with no health checks registered.
pub fn new() -> Self {
Self { checks: Vec::new() }
}

Comment thread
Techassi marked this conversation as resolved.
/// Registers a new [`HealthCheck`] with the provided name and returns it.
///
/// The returned [`HealthCheck`] can be used to mark the check as passed.
pub fn register(&mut self, name: impl Into<String>) -> HealthCheck {
let check = HealthCheck::new(name);
self.checks.push(check.clone());
check
}

/// Returns `true` if all the registered health checks have passed or no health checks are
/// registered.
pub fn all_passed(&self) -> bool {
self.checks.iter().all(HealthCheck::passed)
}
}
Comment thread
Techassi marked this conversation as resolved.

impl IntoResponse for &HealthCheckRegistry {
fn into_response(self) -> Response {
let status = if self.all_passed() {
StatusCode::OK
} else {
StatusCode::SERVICE_UNAVAILABLE
};
// The response body carries check names and their status. Error causes etc. go to the
// log, never into a response to not leak internal information to the public endpoint.
(status, self.to_string()).into_response()
}
}

impl Display for HealthCheckRegistry {
/// Renders one line per check, with the check's name and status only. Anything else, error
/// causes in particular, must not end up in a response to an unauthenticated endpoint.
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
if self.checks.is_empty() {
return writeln!(f, "[ok] no checks registered");
}

for check in &self.checks {
let status = if check.passed() { "ok" } else { "pending" };
writeln!(f, "[{status}] {name}", name = check.name)?;
}
Ok(())
}
}
Comment thread
Techassi marked this conversation as resolved.

#[cfg(test)]
mod tests {
use super::*;

#[test]
fn passed_on_empty_registry() {
let registry = HealthCheckRegistry::new();

assert!(registry.all_passed());
}

#[test]
fn passed_only_once_every_check_passed() {
let mut registry = HealthCheckRegistry::new();
let crds = registry.register("crds-established");
let migration = registry.register("database-migrated");

assert!(!registry.all_passed());

crds.mark_passed();
assert!(!registry.all_passed());

migration.mark_passed();
assert!(registry.all_passed());
}

#[test]
fn not_passed_after_check_marked_not_passed() {
let mut registry = HealthCheckRegistry::new();
let crds = registry.register("crds-established");

crds.mark_passed();
assert!(registry.all_passed());

crds.mark_not_passed();
assert!(!registry.all_passed());
}
}
29 changes: 22 additions & 7 deletions crates/stackable-webhook/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -11,10 +11,13 @@
//!
//! For usage please look at the [`WebhookServer`] docs as well as the specific [`Webhook`] you are
//! using.
use std::net::{IpAddr, Ipv4Addr, SocketAddr};
use std::{
net::{IpAddr, Ipv4Addr, SocketAddr},
sync::Arc,
};

use ::x509_cert::Certificate;
use axum::{Router, routing::get};
use axum::{Router, response::IntoResponse, routing::get};
use futures_util::TryFutureExt;
use k8s_openapi::ByteString;
use snafu::{ResultExt, Snafu};
Expand All @@ -24,8 +27,9 @@ use tower::ServiceBuilder;
use webhooks::{Webhook, WebhookError};
use x509_cert::der::{EncodePem, pem::LineEnding};

use crate::tls::TlsServer;
use crate::{health::HealthCheckRegistry, tls::TlsServer};

pub mod health;
pub mod tls;
pub mod webhooks;

Expand Down Expand Up @@ -54,7 +58,9 @@ pub enum WebhookServerError {
/// ### Example usage
///
/// ```
/// use stackable_webhook::{WebhookServer, WebhookServerOptions, webhooks::Webhook};
/// use stackable_webhook::{
/// WebhookServer, WebhookServerOptions, health::HealthCheckRegistry, webhooks::Webhook,
/// };
/// use tokio::time::{Duration, sleep};
///
/// # async fn docs() {
Expand All @@ -65,7 +71,10 @@ pub enum WebhookServerError {
/// webhook_namespace: "my-namespace".to_owned(),
/// webhook_service_name: "my-operator".to_owned(),
/// };
/// let webhook_server = WebhookServer::new(webhooks, webhook_options).await.unwrap();
/// let readiness_checks = HealthCheckRegistry::new();
/// let webhook_server = WebhookServer::new(webhooks, webhook_options, readiness_checks)
/// .await
/// .unwrap();
/// let shutdown_signal = sleep(Duration::from_millis(100));
///
/// webhook_server.run(shutdown_signal).await.unwrap();
Expand Down Expand Up @@ -111,6 +120,7 @@ impl WebhookServer {
pub async fn new(
webhooks: Vec<Box<dyn Webhook>>,
options: WebhookServerOptions,
readiness_checks: HealthCheckRegistry,
Comment thread
Techassi marked this conversation as resolved.
) -> Result<Self> {
tracing::trace!("create new webhook server");

Expand All @@ -132,12 +142,17 @@ impl WebhookServer {
router = webhook.register_routes(router);
}

// Create the route handler for the startup probe.
let readiness_checks = Arc::new(readiness_checks);
let ready_route = move || async move { readiness_checks.as_ref().into_response() };

let router = router
// Enrich spans for routes added above.
// Routes defined below it will not be instrumented to reduce noise.
.layer(trace_service_builder)
// The health route is below the AxumTraceLayer so as not to be instrumented
.route("/health", get(|| async { "ok" }));
// The health and ready routes are below the AxumTraceLayer so as not to be instrumented
.route("/health", get(|| async { "ok" }))
.route("/ready", get(ready_route));

tracing::debug!("create TLS server");
let (tls_server, cert_rx) = TlsServer::new(router, &options)
Expand Down
12 changes: 9 additions & 3 deletions crates/stackable-webhook/src/webhooks/conversion_webhook.rs
Original file line number Diff line number Diff line change
Expand Up @@ -48,6 +48,7 @@ pub enum ConversionWebhookError {
/// };
/// use stackable_webhook::{
/// WebhookServer,
/// health::HealthCheckRegistry,
/// webhooks::{ConversionWebhook, ConversionWebhookOptions},
/// };
/// use tokio::time::{Duration, sleep};
Expand All @@ -72,9 +73,14 @@ pub enum ConversionWebhookError {
/// ConversionWebhook::new(crds_and_handlers, client, conversion_webhook_options);
///
/// let webhook_options = todo!();
/// let webhook_server = WebhookServer::new(vec![Box::new(conversion_webhook)], webhook_options)
/// .await
/// .unwrap();
/// let readiness_checks = HealthCheckRegistry::new();
/// let webhook_server = WebhookServer::new(
/// vec![Box::new(conversion_webhook)],
/// webhook_options,
/// readiness_checks,
/// )
/// .await
/// .unwrap();
/// let shutdown_signal = sleep(Duration::from_millis(100));
///
/// webhook_server.run(shutdown_signal).await.unwrap();
Expand Down
9 changes: 6 additions & 3 deletions crates/stackable-webhook/src/webhooks/mutating_webhook.rs
Original file line number Diff line number Diff line change
Expand Up @@ -48,6 +48,7 @@ pub enum MutatingWebhookError {
/// };
/// use stackable_webhook::{
/// WebhookServer,
/// health::HealthCheckRegistry,
/// webhooks::{MutatingWebhook, MutatingWebhookOptions},
/// };
/// use tokio::time::{Duration, sleep};
Expand All @@ -73,9 +74,11 @@ pub enum MutatingWebhookError {
/// ));
///
/// let webhook_options = todo!();
/// let webhook_server = WebhookServer::new(vec![mutating_webhook], webhook_options)
/// .await
/// .unwrap();
/// let readiness_checks = HealthCheckRegistry::new();
/// let webhook_server =
/// WebhookServer::new(vec![mutating_webhook], webhook_options, readiness_checks)
/// .await
/// .unwrap();
/// let shutdown_signal = sleep(Duration::from_millis(100));
///
/// webhook_server.run(shutdown_signal).await.unwrap();
Expand Down
Loading