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
11 changes: 11 additions & 0 deletions nativelink-config/src/cas_server.rs
Original file line number Diff line number Diff line change
Expand Up @@ -449,6 +449,17 @@ pub struct HealthConfig {
/// Timeout on health checks. Default: 5s.
#[serde(default)]
pub timeout_seconds: u64,

/// Path of the readiness check, the stricter sibling of `path`: it
/// answers 503 while any component is still initializing, where `path`
/// answers 200. A worker's registration with its scheduler is such a
/// component, so a Kubernetes readiness probe on this path turns Ready
/// only once the worker can take work. Must differ from `path`; the
/// same path for both is refused at startup.
///
/// Default: "/ready"
#[serde(default)]
pub readiness_path: String,
}

#[derive(Deserialize, Serialize, Debug)]
Expand Down
46 changes: 46 additions & 0 deletions nativelink-service/src/health_server.rs
Original file line number Diff line number Diff line change
Expand Up @@ -25,6 +25,7 @@ use http_body_util::Full;
use hyper::header::{CONTENT_TYPE, HeaderValue};
use hyper::{Request, Response, StatusCode};
use nativelink_config::cas_server::HealthConfig;
use nativelink_error::{Error, make_input_err};
use nativelink_util::health_utils::{
HealthRegistry, HealthStatus, HealthStatusDescription, HealthStatusReporter,
};
Expand All @@ -42,6 +43,36 @@ const DEFAULT_HEALTH_CHECK_TIMEOUT_SECONDS: u64 = 5;
pub struct HealthServer {
health_registry: HealthRegistry,
timeout: Duration,
/// Readiness: a component still initializing makes the answer 503.
strict: bool,
}

/// Where the plain status check answers when `path` is unset.
pub const DEFAULT_STATUS_PATH: &str = "/status";
/// Where the readiness check answers when `readiness_path` is unset.
pub const DEFAULT_READINESS_PATH: &str = "/ready";

/// The two paths the health service serves, defaults filled in, as
/// `(status, readiness)`. They must differ: the router panics on a
/// duplicate route, so a configuration that names the same path for both
/// is refused here with a message instead.
pub fn health_paths(health_cfg: &HealthConfig) -> Result<(String, String), Error> {
let status = if health_cfg.path.is_empty() {
DEFAULT_STATUS_PATH
} else {
&health_cfg.path
};
let readiness = if health_cfg.readiness_path.is_empty() {
DEFAULT_READINESS_PATH
} else {
&health_cfg.readiness_path
};
if status == readiness {
return Err(make_input_err!(
"services.health.path and readiness_path are both {status}; the readiness check needs a path of its own"
));
}
Ok((status.to_string(), readiness.to_string()))
}

impl HealthServer {
Expand All @@ -54,8 +85,17 @@ impl HealthServer {
Self {
health_registry,
timeout,
strict: false,
}
}

/// The readiness form: unavailable while anything is initializing, not
/// only when something failed.
pub const fn readiness(health_registry: HealthRegistry, health_cfg: &HealthConfig) -> Self {
let mut server = Self::new(health_registry, health_cfg);
server.strict = true;
server
}
}

impl Service<Request<Body>> for HealthServer {
Expand All @@ -70,6 +110,7 @@ impl Service<Request<Body>> for HealthServer {
fn call(&mut self, _req: Request<Body>) -> Self::Future {
let health_registry = self.health_registry.clone();
let local_timeout = self.timeout;
let strict = self.strict;
Box::pin(error_span!("health_server_call").in_scope(|| async move {
let health_status_descriptions: Vec<HealthStatusDescription> = health_registry
.health_status_report(&local_timeout)
Expand All @@ -82,6 +123,11 @@ impl Service<Request<Body>> for HealthServer {
health_status_descriptions.iter().any(|description| {
matches!(description.status, HealthStatus::Failed { .. })
| matches!(description.status, HealthStatus::Timeout { .. })
| (strict
&& matches!(
description.status,
HealthStatus::Initializing { .. }
))
});
let status_code = if contains_failed_report {
StatusCode::SERVICE_UNAVAILABLE
Expand Down
82 changes: 81 additions & 1 deletion nativelink-service/tests/health_server_test.rs
Original file line number Diff line number Diff line change
Expand Up @@ -6,7 +6,7 @@ use axum::http::Request;
use hyper::StatusCode;
use nativelink_config::cas_server::HealthConfig;
use nativelink_macro::nativelink_test;
use nativelink_service::health_server::HealthServer;
use nativelink_service::health_server::{HealthServer, health_paths};
use nativelink_util::health_utils::{
HealthRegistry, HealthRegistryBuilder, HealthStatus, HealthStatusIndicator,
};
Expand Down Expand Up @@ -155,3 +155,83 @@ async fn health_test_with_sleep() -> Result<(), Box<dyn core::error::Error>> {
));
Ok(())
}

struct InitializingIndicator {}

#[async_trait]
impl HealthStatusIndicator for InitializingIndicator {
fn get_name(&self) -> &'static str {
"initializing_indicator"
}

async fn check_health(&self, _namespace: Cow<'static, str>) -> HealthStatus {
HealthStatus::Initializing {
struct_name: "InitializingIndicator",
message: "not yet".into(),
}
}

fn struct_name(&self) -> &'static str {
"InitializingIndicator"
}
}

async fn status_of(server: HealthServer) -> Result<StatusCode, Box<dyn core::error::Error>> {
let tonic_services = Routes::builder().routes();
let mut svc = tonic_services
.into_axum_router()
.route_service("/probe", server);
let request = Request::builder()
.method("GET")
.uri("/probe")
.body(Body::empty())?;
let response: hyper::Response<axum::body::Body> =
svc.as_service().ready().await?.call(request).await?;
Ok(response.status())
}

/// A component still initializing is healthy enough for `/status` and not
/// ready enough for `/ready`.
#[nativelink_test]
async fn readiness_is_unavailable_while_initializing() -> Result<(), Box<dyn core::error::Error>> {
let mut health_registry_builder = HealthRegistryBuilder::new("foo");
health_registry_builder.register_indicator(Arc::new(InitializingIndicator {}));
let health_registry = health_registry_builder.build();
let config = HealthConfig::default();

assert_eq!(
status_of(HealthServer::new(health_registry.clone(), &config)).await?,
StatusCode::OK
);
assert_eq!(
status_of(HealthServer::readiness(health_registry, &config)).await?,
StatusCode::SERVICE_UNAVAILABLE
);
Ok(())
}

/// Defaults fill the two paths; a configuration that gives both the same
/// path is refused with a message rather than left to panic the router.
#[nativelink_test]
async fn health_paths_are_distinct_or_refused() -> Result<(), Box<dyn core::error::Error>> {
assert_eq!(
health_paths(&HealthConfig::default())?,
("/status".to_string(), "/ready".to_string())
);
let custom = HealthConfig {
path: "/healthz".to_string(),
readiness_path: "/readyz".to_string(),
..HealthConfig::default()
};
assert_eq!(
health_paths(&custom)?,
("/healthz".to_string(), "/readyz".to_string())
);
let clash = HealthConfig {
path: "/ready".to_string(),
..HealthConfig::default()
};
let err = health_paths(&clash).expect_err("the same path for both must be refused");
assert!(err.to_string().contains("both /ready"), "{err}");
Ok(())
}
97 changes: 95 additions & 2 deletions nativelink-worker/src/local_worker.rs
Original file line number Diff line number Diff line change
Expand Up @@ -41,6 +41,7 @@ use nativelink_store::fast_slow_store::FastSlowStore;
use nativelink_util::action_messages::{ActionResult, ActionStage, OperationId};
use nativelink_util::common::fs;
use nativelink_util::digest_hasher::DigestHasherFunc;
use nativelink_util::health_utils::{HealthStatus, HealthStatusIndicator};
use nativelink_util::metrics_utils::{AsyncCounterWrapper, CounterWithTime};
use nativelink_util::shutdown_guard::ShutdownGuard;
use nativelink_util::store_trait::Store;
Expand All @@ -49,7 +50,7 @@ use opentelemetry::context::Context;
use tokio::sync::{broadcast, mpsc};
use tokio::{process, time};
use tokio_stream::wrappers::UnboundedReceiverStream;
use tonic::Streaming;
use tonic::{Streaming, async_trait};
use tracing::{Level, debug, error, event, info, info_span, instrument, trace, warn};

use crate::running_actions_manager::{
Expand Down Expand Up @@ -771,6 +772,80 @@ pub struct LocalWorker<T: WorkerApiClientTrait + 'static, U: RunningActionsManag
connection_factory: ConnectionFactory<T>,
sleep_fn: Option<Box<dyn Fn(Duration) -> BoxFuture<'static, ()> + Send + Sync>>,
metrics: Arc<Metrics>,
registration: Arc<WorkerRegistration>,
}

/// Whether the worker currently holds a registration with its scheduler.
/// A worker that has not registered, or lost its connection and is
/// reconnecting, cannot take work; as a health indicator this is what a
/// readiness probe reads.
#[derive(Debug)]
pub struct WorkerRegistration {
name: String,
registered: AtomicBool,
ever_registered: AtomicBool,
}

impl WorkerRegistration {
pub fn new(name: &str) -> Arc<Self> {
Arc::new(Self {
name: name.to_string(),
registered: AtomicBool::new(false),
ever_registered: AtomicBool::new(false),
})
}

pub fn is_registered(&self) -> bool {
self.registered.load(Ordering::Acquire)
}

fn set_registered(&self, registered: bool) {
self.registered.store(registered, Ordering::Release);
if registered {
self.ever_registered.store(true, Ordering::Release);
}
}
}

#[async_trait]
impl HealthStatusIndicator for WorkerRegistration {
fn get_name(&self) -> &'static str {
"WorkerRegistration"
}

async fn check_health(&self, _namespace: Cow<'static, str>) -> HealthStatus {
if self.is_registered() {
HealthStatus::Ok {
struct_name: "WorkerRegistration",
message: Cow::Owned(format!(
"worker '{}' registered with the scheduler",
self.name
)),
}
} else if self.ever_registered.load(Ordering::Acquire) {
// Lost after it was there: the worker is reconnecting, and until
// it does it holds no work. Initializing, not Failed: Failed
// turns the plain status check red too, and a liveness probe on
// it would restart every worker during a scheduler roll longer
// than its threshold, and redden a co-hosted CAS on one worker's
// blip. This belongs to readiness alone.
HealthStatus::Initializing {
struct_name: "WorkerRegistration",
message: Cow::Owned(format!(
"worker '{}' lost its scheduler connection, reconnecting",
self.name
)),
}
} else {
HealthStatus::Initializing {
struct_name: "WorkerRegistration",
message: Cow::Owned(format!(
"worker '{}' not yet registered with the scheduler",
self.name
)),
}
}
}
}

impl<
Expand Down Expand Up @@ -1031,15 +1106,30 @@ impl<T: WorkerApiClientTrait + 'static, U: RunningActionsManager> LocalWorker<T,
let metrics = Arc::new(Metrics::new(Arc::downgrade(
running_actions_manager.metrics(),
)));
let registration = WorkerRegistration::new(&config.name);
Self {
config,
running_actions_manager,
connection_factory,
sleep_fn: Some(sleep_fn),
metrics,
registration,
}
}

/// The registration flag this worker flips, for a health registry.
pub fn registration(&self) -> Arc<WorkerRegistration> {
self.registration.clone()
}

/// Flip a flag that was registered with a health registry before the
/// worker existed, as the server binary has to.
#[must_use]
pub fn with_registration(mut self, registration: Arc<WorkerRegistration>) -> Self {
self.registration = registration;
self
}

#[allow(
clippy::missing_const_for_fn,
reason = "False positive on stable, but not on nightly"
Expand Down Expand Up @@ -1189,9 +1279,12 @@ impl<T: WorkerApiClientTrait + 'static, U: RunningActionsManager> LocalWorker<T,
"Worker registered with scheduler"
);
attempts.store(0, Ordering::Release);
self.registration.set_registered(true);

// Now listen for connections and run all other services.
if let Err(err) = inner.run(update_for_worker_stream, &mut shutdown_rx).await {
let run_result = inner.run(update_for_worker_stream, &mut shutdown_rx).await;
self.registration.set_registered(false);
if let Err(err) = run_result {
// Give in-transit actions a chance to settle before we kill
// them, so their results still reach the scheduler.
const ITERATIONS: usize = 1_000;
Expand Down
Loading
Loading