diff --git a/.agents/skills/helm-dev-environment/SKILL.md b/.agents/skills/helm-dev-environment/SKILL.md index 932d0d905d..c0d1e68071 100644 --- a/.agents/skills/helm-dev-environment/SKILL.md +++ b/.agents/skills/helm-dev-environment/SKILL.md @@ -449,6 +449,7 @@ for dependencies still declared in `Chart.yaml`. | `deploy/helm/openshell/ci/values-cert-manager.yaml` | cert-manager PKI overlay (opt-in; disables pkiInitJob) | | `deploy/helm/openshell/ci/values-gateway.yaml` | Envoy Gateway GRPCRoute + Gateway overlay | | `deploy/helm/openshell/ci/values-high-availability.yaml` | HA test overlay (`replicaCount: 2` with external PostgreSQL Secret) | +| `deploy/helm/openshell/ci/values-autoscaling.yaml` | Render-only overlay for the optional gateway HorizontalPodAutoscaler (helm lint and helm-unittest) | | `deploy/helm/openshell/ci/values-keycloak.yaml` | Keycloak OIDC overlay | | `deploy/helm/openshell/ci/values-spire.yaml` | SPIFFE/SPIRE provider token grant overlay | | `deploy/helm/openshell/ci/values-spire-stack.yaml` | SPIRE hardened chart values for local dev | diff --git a/crates/openshell-server/src/gateway_metrics.rs b/crates/openshell-server/src/gateway_metrics.rs new file mode 100644 index 0000000000..28adda1712 --- /dev/null +++ b/crates/openshell-server/src/gateway_metrics.rs @@ -0,0 +1,754 @@ +// SPDX-FileCopyrightText: Copyright (c) 2025-2026 NVIDIA CORPORATION & AFFILIATES. All rights reserved. +// SPDX-License-Identifier: Apache-2.0 + +//! Gateway capacity and HA metrics. +//! +//! The metric names are an operator-facing contract, documented in +//! docs/observability/gateway-metrics.mdx. Labels are bounded enums only: never a sandbox, +//! channel, endpoint, token, or replica id. The scrape target already identifies the replica. +//! +//! Metric handles bind to whichever recorder is current when a macro runs. `run_server` +//! installs the Prometheus recorder after it builds `ServerState`, so never cache a handle in a +//! static or in state built before [`install_global_recorder`]. [`GaugeSlot`] acquires its +//! handle when the tracked object is created and releases it through the same handle, so a +//! slot can never drive a series negative. + +use std::time::{Duration, Instant}; + +use metrics::{ + Gauge, Unit, counter, describe_counter, describe_gauge, describe_histogram, gauge, histogram, +}; +use metrics_exporter_prometheus::{BuildError, Matcher, PrometheusBuilder, PrometheusHandle}; +use tonic::{Code, Status}; + +// Gauges +pub const SUPERVISOR_SESSIONS: &str = "openshell_server_supervisor_sessions"; +pub const RELAY_PENDING: &str = "openshell_server_relay_pending"; +pub const RELAY_PENDING_CAPACITY: &str = "openshell_server_relay_pending_capacity"; +// Counters +pub const RELAY_REJECTED_TOTAL: &str = "openshell_server_relay_rejected_total"; +pub const RELAY_EXPIRED_TOTAL: &str = "openshell_server_relay_expired_total"; +pub const ROUTED_REQUEST_ATTEMPTS_TOTAL: &str = "openshell_server_routed_request_attempts_total"; +// Histograms (explicit buckets, see BUCKETED_HISTOGRAMS) +pub const RELAY_CLAIM_DURATION_SECONDS: &str = "openshell_server_relay_claim_duration_seconds"; +pub const PEER_REQUEST_DURATION_SECONDS: &str = "openshell_server_peer_request_duration_seconds"; + +const LABEL_REASON: &str = "reason"; +const LABEL_OPERATION: &str = "operation"; +const LABEL_OUTCOME: &str = "outcome"; +const LABEL_GRPC_CODE: &str = "grpc_code"; +const LABEL_RELAY_KIND: &str = "relay_kind"; +const LABEL_ROUTE: &str = "route"; + +/// Buckets for the new latency histograms, 1 ms to 15 s. The top buckets cover the 10 s relay +/// claim timeout and the 15 s routed-relay wait. +const LATENCY_BUCKETS_SECONDS: [f64; 14] = [ + 0.001, 0.0025, 0.005, 0.01, 0.025, 0.05, 0.1, 0.25, 0.5, 1.0, 2.5, 5.0, 10.0, 15.0, +]; + +/// Only these names render as Prometheus histograms. Every existing `*_duration_seconds` metric +/// keeps its summary format, so current dashboards are unaffected. +const BUCKETED_HISTOGRAMS: [&str; 2] = + [RELAY_CLAIM_DURATION_SECONDS, PEER_REQUEST_DURATION_SECONDS]; + +/// Protocol the supervisor is asked to relay. Never label metrics with the target address. +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +pub enum RelayKind { + Ssh, + Tcp, +} + +impl RelayKind { + pub const ALL: [Self; 2] = [Self::Ssh, Self::Tcp]; + + pub const fn label(self) -> &'static str { + match self { + Self::Ssh => "ssh", + Self::Tcp => "tcp", + } + } +} + +/// Where the requesting replica tries to open a relay. +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +pub enum RelayRoute { + Local, + Peer, +} + +impl RelayRoute { + pub const ALL: [Self; 2] = [Self::Local, Self::Peer]; + + pub const fn label(self) -> &'static str { + match self { + Self::Local => "local", + Self::Peer => "peer", + } + } +} + +/// Which pending-relay cap rejected an open. +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +pub enum RelayRejection { + ReplicaCapacity, + SandboxCapacity, +} + +impl RelayRejection { + pub const ALL: [Self; 2] = [Self::ReplicaCapacity, Self::SandboxCapacity]; + + pub const fn label(self) -> &'static str { + match self { + Self::ReplicaCapacity => "replica_capacity", + Self::SandboxCapacity => "sandbox_capacity", + } + } +} + +/// Routed operation. For a peer request, the owning replica records the matching gRPC method in +/// `openshell_server_grpc_requests_total`, for example `PeerRelay` for `relay`. +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +pub enum PeerRpc { + Relay, + ReportProviderReadiness, + ReportEndpointStatus, + GetSandboxProviderStatus, +} + +impl PeerRpc { + pub const ALL: [Self; 4] = [ + Self::Relay, + Self::ReportProviderReadiness, + Self::ReportEndpointStatus, + Self::GetSandboxProviderStatus, + ]; + + const fn operation(self) -> &'static str { + match self { + Self::Relay => "relay", + Self::ReportProviderReadiness => "report_provider_readiness", + Self::ReportEndpointStatus => "report_endpoint_status", + Self::GetSandboxProviderStatus => "get_sandbox_provider_status", + } + } +} + +/// Where a routed attempt ended. A relay succeeds only when the supervisor claims it, on either +/// route, so the values mean the same thing for local and peer attempts. +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +enum AttemptOutcome { + Success, + /// Failed on this replica, including cancellation by the caller. + LocalError, + /// The owning replica returned an error, or the open peer connection failed. + RemoteError, +} + +impl AttemptOutcome { + const fn label(self) -> &'static str { + match self { + Self::Success => "success", + Self::LocalError => "local_error", + Self::RemoteError => "remote_error", + } + } +} + +/// Snake-case gRPC status name for the `grpc_code` label. It is not named `code` because +/// `openshell_server_grpc_requests_total` uses that name for the numeric status. Exhaustive on +/// purpose, so a new tonic variant fails to compile instead of producing an unbounded label. +const fn grpc_code_label(code: Code) -> &'static str { + match code { + Code::Ok => "ok", + Code::Cancelled => "cancelled", + Code::Unknown => "unknown", + Code::InvalidArgument => "invalid_argument", + Code::DeadlineExceeded => "deadline_exceeded", + Code::NotFound => "not_found", + Code::AlreadyExists => "already_exists", + Code::PermissionDenied => "permission_denied", + Code::ResourceExhausted => "resource_exhausted", + Code::FailedPrecondition => "failed_precondition", + Code::Aborted => "aborted", + Code::OutOfRange => "out_of_range", + Code::Unimplemented => "unimplemented", + Code::Internal => "internal", + Code::Unavailable => "unavailable", + Code::DataLoss => "data_loss", + Code::Unauthenticated => "unauthenticated", + } +} + +/// Relay cap published as `openshell_server_relay_pending_capacity`. The caller passes the value +/// that enforces the cap, so this module does not depend on the relay registry. +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +pub struct RelayCapacity { + /// Pending relays allowed on one replica. + pub per_replica: usize, +} + +/// Apply the bucket overrides. Tests build local recorders from the same builder. +pub fn configure_exporter(builder: PrometheusBuilder) -> Result { + BUCKETED_HISTOGRAMS + .iter() + .try_fold(builder, |builder, name| { + builder.set_buckets_for_metric( + Matcher::Full((*name).to_string()), + &LATENCY_BUCKETS_SECONDS, + ) + }) +} + +/// Install the process-wide recorder, then describe and zero-initialize the catalog. Call +/// once, from `run_server`. +pub fn install_global_recorder(relay: RelayCapacity) -> Result { + let handle = configure_exporter(PrometheusBuilder::new())?.install_recorder()?; + describe_and_initialize(relay); + Ok(handle) +} + +/// Emit HELP metadata, create every fixed-label series at 0, and publish the relay caps. An +/// idle replica then exports 0 instead of "no data", which HPA and `rate()` need. +pub fn describe_and_initialize(relay: RelayCapacity) { + describe_gauge!( + SUPERVISOR_SESSIONS, + Unit::Count, + "Supervisor control sessions registered on this gateway replica." + ); + describe_gauge!( + RELAY_PENDING, + Unit::Count, + "Relay channels on this replica waiting for the supervisor to connect back, including channels opened for peer replicas." + ); + describe_gauge!( + RELAY_PENDING_CAPACITY, + Unit::Count, + "Maximum pending relay channels on one gateway replica." + ); + describe_counter!( + ROUTED_REQUEST_ATTEMPTS_TOTAL, + Unit::Count, + "Local relay setup and outbound peer attempts completed or cancelled by this replica. Each retry counts separately." + ); + describe_counter!( + RELAY_REJECTED_TOTAL, + Unit::Count, + "Relay opens rejected because a pending relay cap was reached." + ); + describe_counter!( + RELAY_EXPIRED_TOTAL, + Unit::Count, + "Pending relay channels dropped because the supervisor did not connect back in time." + ); + describe_histogram!( + RELAY_CLAIM_DURATION_SECONDS, + Unit::Seconds, + "Time from opening a relay channel to the supervisor claiming it." + ); + describe_histogram!( + PEER_REQUEST_DURATION_SECONDS, + Unit::Seconds, + "Latency of outbound requests to the owning replica. For relays, until the owner's supervisor claimed the relay." + ); + + // `increment(0)` registers a series without overwriting a value recorded earlier. + gauge!(SUPERVISOR_SESSIONS).increment(0.0); + gauge!(RELAY_PENDING).increment(0.0); + gauge!(RELAY_PENDING_CAPACITY).set(count_as_f64(relay.per_replica)); + for kind in RelayKind::ALL { + for route in RelayRoute::ALL { + counter!( + ROUTED_REQUEST_ATTEMPTS_TOTAL, + LABEL_OPERATION => PeerRpc::Relay.operation(), + LABEL_ROUTE => route.label(), + LABEL_RELAY_KIND => kind.label(), + LABEL_OUTCOME => AttemptOutcome::Success.label(), + LABEL_GRPC_CODE => grpc_code_label(Code::Ok) + ) + .increment(0); + } + } + for reason in RelayRejection::ALL { + counter!(RELAY_REJECTED_TOTAL, LABEL_REASON => reason.label()).increment(0); + } + counter!(RELAY_EXPIRED_TOTAL).increment(0); + for rpc in PeerRpc::ALL { + if rpc == PeerRpc::Relay { + continue; + } + counter!( + ROUTED_REQUEST_ATTEMPTS_TOTAL, + LABEL_OPERATION => rpc.operation(), + LABEL_ROUTE => RelayRoute::Peer.label(), + LABEL_RELAY_KIND => "none", + LABEL_OUTCOME => AttemptOutcome::Success.label(), + LABEL_GRPC_CODE => grpc_code_label(Code::Ok) + ) + .increment(0); + } +} + +/// Counts in this module stay far below 2^53, so the conversion is exact. +#[allow(clippy::cast_precision_loss)] +pub fn count_as_f64(count: usize) -> f64 { + count as f64 +} + +/// One unit of an exact gauge, held for as long as the tracked object lives. Dropping it +/// decrements through the same handle it incremented, so every removal path is counted exactly +/// once, including paths added in the future. +#[must_use = "dropping a GaugeSlot immediately releases it"] +pub struct GaugeSlot(Gauge); + +impl GaugeSlot { + /// Share of `openshell_server_supervisor_sessions`. + pub fn supervisor_session() -> Self { + Self::acquire(SUPERVISOR_SESSIONS) + } + + /// Share of `openshell_server_relay_pending`. + pub fn relay_pending() -> Self { + Self::acquire(RELAY_PENDING) + } + + fn acquire(name: &'static str) -> Self { + let gauge = gauge!(name); + gauge.increment(1.0); + Self(gauge) + } +} + +impl Drop for GaugeSlot { + fn drop(&mut self) { + self.0.decrement(1.0); + } +} + +pub fn record_relay_rejected(reason: RelayRejection) { + counter!(RELAY_REJECTED_TOTAL, LABEL_REASON => reason.label()).increment(1); +} + +/// `count` pending relays were dropped unclaimed (late claim or reaper). +pub fn record_relay_expired(count: usize) { + if count > 0 { + counter!(RELAY_EXPIRED_TOTAL).increment(count as u64); + } +} + +pub fn record_relay_claimed(waited: Duration) { + histogram!(RELAY_CLAIM_DURATION_SECONDS).record(waited); +} + +/// Counts one local relay setup or outbound peer attempt exactly once, and times peer requests. +/// A relay succeeds when the supervisor claims it, on either route. Dropping an unfinished +/// timer (the caller gave up) records `local_error` / `cancelled`. +#[must_use = "finish the timer with local_error() or finish()"] +pub struct RoutedRequestTimer { + rpc: PeerRpc, + route: RelayRoute, + relay_kind: Option, + started: Instant, + recorded: bool, +} + +impl RoutedRequestTimer { + pub fn start(rpc: PeerRpc) -> Self { + Self { + rpc, + route: RelayRoute::Peer, + relay_kind: if rpc == PeerRpc::Relay { + Some(RelayKind::Ssh) + } else { + None + }, + started: Instant::now(), + recorded: false, + } + } + + pub fn relay(kind: RelayKind, route: RelayRoute) -> Self { + Self { + rpc: PeerRpc::Relay, + route, + relay_kind: Some(kind), + started: Instant::now(), + recorded: false, + } + } + + /// The attempt failed on this replica: local relay setup or claim (including the claim + /// window closing), or a peer request that failed before it reached the owner (token, + /// channel, headers, or stream setup). + pub fn local_error(&mut self, status: &Status) { + self.record(AttemptOutcome::LocalError, status.code()); + } + + /// Record the raw tonic result of the RPC itself. Call this BEFORE any remap to + /// `Unavailable`, so the owner's code (for example `resource_exhausted`) is kept. + pub fn finish(&mut self, result: &Result) { + match result { + Ok(_) => self.record(AttemptOutcome::Success, Code::Ok), + Err(status) if self.route == RelayRoute::Local => self.local_error(status), + Err(status) => self.record(AttemptOutcome::RemoteError, status.code()), + } + } + + fn record(&mut self, outcome: AttemptOutcome, code: Code) { + if self.recorded { + return; + } + self.recorded = true; + counter!( + ROUTED_REQUEST_ATTEMPTS_TOTAL, + LABEL_OPERATION => self.rpc.operation(), + LABEL_ROUTE => self.route.label(), + LABEL_RELAY_KIND => self.relay_kind.map_or("none", RelayKind::label), + LABEL_OUTCOME => outcome.label(), + LABEL_GRPC_CODE => grpc_code_label(code) + ) + .increment(1); + if self.route == RelayRoute::Peer { + histogram!( + PEER_REQUEST_DURATION_SECONDS, + LABEL_OPERATION => self.rpc.operation(), + LABEL_OUTCOME => outcome.label() + ) + .record(self.started.elapsed()); + } + } +} + +impl Drop for RoutedRequestTimer { + fn drop(&mut self) { + self.record(AttemptOutcome::LocalError, Code::Cancelled); + } +} + +/// Captures metrics recorded on the current thread through a configured Prometheus recorder. +/// +/// Works in `#[test]` and in the default current-thread `#[tokio::test]`, where tasks spawned on +/// the runtime share the thread. It does not work in `multi_thread` tests or inside +/// `spawn_blocking`. Never pass it into an `async fn` helper: it is `!Send`, and clippy +/// `future_not_send` (nursery) would fire. +#[cfg(test)] +pub struct MetricsCapture { + handle: PrometheusHandle, + _guard: metrics::LocalRecorderGuard<'static>, +} + +#[cfg(test)] +impl MetricsCapture { + pub fn install() -> Self { + // Leaked (test only, one small allocation per test) so the guard can borrow it for 'static. + let recorder: &'static metrics_exporter_prometheus::PrometheusRecorder = + Box::leak(Box::new( + configure_exporter(PrometheusBuilder::new()) + .expect("valid exporter config") + .build_recorder(), + )); + let handle = recorder.handle(); + let guard = metrics::set_default_local_recorder(recorder); + Self { + handle, + _guard: guard, + } + } + + pub fn render(&self) -> String { + self.handle.render() + } + + /// Integer value of one exact series, such as `name` or `name{a="b"}`. Parses as i64 to + /// avoid clippy `float_cmp` and so a negative gauge is visible. `None` if the series is absent. + pub fn value(&self, series: &str) -> Option { + series_value(&self.handle, series) + } + + /// Reads one series like [`Self::value`], from code that cannot hold `self`, such as a + /// waker that runs while the value is being produced. + pub fn value_reader(&self, series: &'static str) -> Box Option + Send + Sync> { + let handle = self.handle.clone(); + Box::new(move || series_value(&handle, series)) + } +} + +#[cfg(test)] +fn series_value(handle: &PrometheusHandle, series: &str) -> Option { + handle + .render() + .lines() + .find_map(|line| line.strip_prefix(series)?.strip_prefix(' ')?.parse().ok()) +} + +#[cfg(test)] +mod tests { + use super::*; + use std::collections::HashSet; + + #[test] + fn describe_and_initialize_exports_capacity_and_zero_series() { + let metrics = MetricsCapture::install(); + describe_and_initialize(RelayCapacity { per_replica: 256 }); + + for (series, expected) in [ + ("openshell_server_relay_pending_capacity", 256), + ("openshell_server_supervisor_sessions", 0), + ("openshell_server_relay_pending", 0), + ( + "openshell_server_relay_rejected_total{reason=\"replica_capacity\"}", + 0, + ), + ( + "openshell_server_relay_rejected_total{reason=\"sandbox_capacity\"}", + 0, + ), + ("openshell_server_relay_expired_total", 0), + ] { + assert_eq!(metrics.value(series), Some(expected), "{series}"); + } + for kind in ["ssh", "tcp"] { + for route in ["local", "peer"] { + let series = format!( + "openshell_server_routed_request_attempts_total{{operation=\"relay\",route=\"{route}\",relay_kind=\"{kind}\",outcome=\"success\",grpc_code=\"ok\"}}" + ); + assert_eq!(metrics.value(&series), Some(0), "{series}"); + } + } + for operation in [ + "report_provider_readiness", + "report_endpoint_status", + "get_sandbox_provider_status", + ] { + let series = format!( + "openshell_server_routed_request_attempts_total{{operation=\"{operation}\",route=\"peer\",relay_kind=\"none\",outcome=\"success\",grpc_code=\"ok\"}}" + ); + assert_eq!(metrics.value(&series), Some(0), "{series}"); + } + assert!( + metrics + .render() + .contains("# HELP openshell_server_supervisor_sessions ") + ); + } + + #[test] + fn routed_attempts_have_seven_bounded_success_series_and_keep_counts_on_initialize() { + let metrics = MetricsCapture::install(); + for kind in RelayKind::ALL { + for route in RelayRoute::ALL { + RoutedRequestTimer::relay(kind, route).finish(&Ok::<(), Status>(())); + } + } + RoutedRequestTimer::relay(RelayKind::Tcp, RelayRoute::Peer).finish(&Ok::<(), Status>(())); + describe_and_initialize(RelayCapacity { per_replica: 256 }); + + for kind in ["ssh", "tcp"] { + for route in ["local", "peer"] { + let series = format!( + "openshell_server_routed_request_attempts_total{{operation=\"relay\",route=\"{route}\",relay_kind=\"{kind}\",outcome=\"success\",grpc_code=\"ok\"}}" + ); + let expected = if kind == "tcp" && route == "peer" { + 2 + } else { + 1 + }; + assert_eq!(metrics.value(&series), Some(expected), "{series}"); + } + } + let rendered = metrics.render(); + assert!(rendered.contains("# TYPE openshell_server_routed_request_attempts_total counter")); + assert!(!rendered.contains("openshell_server_relay_setup_attempts_total")); + assert!(!rendered.contains("openshell_server_peer_requests_total")); + assert_eq!( + rendered + .lines() + .filter(|line| line.starts_with("openshell_server_routed_request_attempts_total{")) + .count(), + 7 + ); + } + + #[test] + fn configured_exporter_buckets_only_new_latency_histograms() { + let metrics = MetricsCapture::install(); + let sample = Duration::from_millis(3); + histogram!(RELAY_CLAIM_DURATION_SECONDS).record(sample); + histogram!( + PEER_REQUEST_DURATION_SECONDS, + LABEL_OPERATION => "relay", + LABEL_OUTCOME => "success" + ) + .record(sample); + histogram!( + "openshell_server_grpc_request_duration_seconds", + "method" => "ListSandboxes", + "code" => "0" + ) + .record(sample); + histogram!( + "openshell_server_http_request_duration_seconds", + "path" => "/healthz", + "status" => "200" + ) + .record(sample); + histogram!( + "openshell_server_readiness_database_probe_duration_seconds", + "outcome" => "success" + ) + .record(sample); + histogram!("openshell_gateway_interceptor_latency_seconds").record(sample); + + let rendered = metrics.render(); + for name in BUCKETED_HISTOGRAMS { + assert!( + rendered.contains(&format!("# TYPE {name} histogram")), + "{name} should render as a histogram" + ); + } + for name in [ + "openshell_server_grpc_request_duration_seconds", + "openshell_server_http_request_duration_seconds", + "openshell_server_readiness_database_probe_duration_seconds", + "openshell_gateway_interceptor_latency_seconds", + ] { + assert!( + rendered.contains(&format!("# TYPE {name} summary")), + "{name} should keep its summary format" + ); + } + assert!( + rendered.contains("openshell_server_relay_claim_duration_seconds_bucket{le=\"0.001\"}") + ); + assert!( + rendered.contains("openshell_server_relay_claim_duration_seconds_bucket{le=\"15\"}") + ); + } + + #[test] + fn gauge_slot_counts_until_dropped() { + let metrics = MetricsCapture::install(); + let first = GaugeSlot::relay_pending(); + let second = GaugeSlot::relay_pending(); + assert_eq!(metrics.value(RELAY_PENDING), Some(2)); + drop(first); + assert_eq!(metrics.value(RELAY_PENDING), Some(1)); + drop(second); + assert_eq!(metrics.value(RELAY_PENDING), Some(0)); + + let session = GaugeSlot::supervisor_session(); + assert_eq!(metrics.value(SUPERVISOR_SESSIONS), Some(1)); + drop(session); + assert_eq!(metrics.value(SUPERVISOR_SESSIONS), Some(0)); + } + + #[test] + fn gauge_slot_acquired_before_recorder_never_goes_negative() { + // No capture is installed yet, so this slot binds to the no-op recorder. + let slot = GaugeSlot::relay_pending(); + let metrics = MetricsCapture::install(); + drop(slot); + assert_eq!(metrics.value(RELAY_PENDING), None); + } + + #[test] + fn peer_request_timer_records_outcome_code_and_latency() { + let metrics = MetricsCapture::install(); + + let mut relay = RoutedRequestTimer::start(PeerRpc::Relay); + relay.finish(&Err::<(), _>(Status::resource_exhausted("x"))); + drop(relay); + assert_eq!( + metrics.value( + "openshell_server_routed_request_attempts_total{operation=\"relay\",route=\"peer\",relay_kind=\"ssh\",outcome=\"remote_error\",grpc_code=\"resource_exhausted\"}" + ), + Some(1) + ); + assert_eq!( + metrics.value( + "openshell_server_peer_request_duration_seconds_count{operation=\"relay\",outcome=\"remote_error\"}" + ), + Some(1) + ); + + let mut endpoint = RoutedRequestTimer::start(PeerRpc::ReportEndpointStatus); + endpoint.local_error(&Status::unavailable("x")); + drop(endpoint); + assert_eq!( + metrics.value( + "openshell_server_routed_request_attempts_total{operation=\"report_endpoint_status\",route=\"peer\",relay_kind=\"none\",outcome=\"local_error\",grpc_code=\"unavailable\"}" + ), + Some(1) + ); + + let mut provider_status = RoutedRequestTimer::start(PeerRpc::GetSandboxProviderStatus); + provider_status.finish(&Ok::<(), Status>(())); + drop(provider_status); + assert_eq!( + metrics.value( + "openshell_server_routed_request_attempts_total{operation=\"get_sandbox_provider_status\",route=\"peer\",relay_kind=\"none\",outcome=\"success\",grpc_code=\"ok\"}" + ), + Some(1) + ); + } + + #[test] + fn peer_request_timer_records_cancelled_when_dropped_unfinished() { + let metrics = MetricsCapture::install(); + drop(RoutedRequestTimer::start(PeerRpc::ReportProviderReadiness)); + assert_eq!( + metrics.value( + "openshell_server_routed_request_attempts_total{operation=\"report_provider_readiness\",route=\"peer\",relay_kind=\"none\",outcome=\"local_error\",grpc_code=\"cancelled\"}" + ), + Some(1) + ); + } + + #[test] + fn peer_request_timer_records_once() { + let metrics = MetricsCapture::install(); + let mut timer = RoutedRequestTimer::start(PeerRpc::Relay); + timer.finish(&Ok::<(), Status>(())); + timer.local_error(&Status::unavailable("x")); + drop(timer); + assert_eq!( + metrics.value( + "openshell_server_routed_request_attempts_total{operation=\"relay\",route=\"peer\",relay_kind=\"ssh\",outcome=\"success\",grpc_code=\"ok\"}" + ), + Some(1) + ); + assert!(!metrics.render().contains("outcome=\"local_error\"")); + } + + #[test] + fn local_relay_failure_records_status_without_peer_latency() { + let metrics = MetricsCapture::install(); + let mut timer = RoutedRequestTimer::relay(RelayKind::Tcp, RelayRoute::Local); + timer.finish(&Err::<(), _>(Status::resource_exhausted("capacity"))); + drop(timer); + assert_eq!( + metrics.value("openshell_server_routed_request_attempts_total{operation=\"relay\",route=\"local\",relay_kind=\"tcp\",outcome=\"local_error\",grpc_code=\"resource_exhausted\"}"), + Some(1) + ); + assert!(!metrics.render().contains(PEER_REQUEST_DURATION_SECONDS)); + } + + #[test] + fn grpc_code_labels_are_distinct_snake_case() { + let labels: HashSet<&str> = (0..=16) + .map(|code| grpc_code_label(Code::from(code))) + .collect(); + for label in &labels { + assert!( + label.chars().all(|c| c.is_ascii_lowercase() || c == '_'), + "{label} is not snake_case" + ); + } + assert_eq!(labels.len(), 17); + assert_eq!(grpc_code_label(Code::DeadlineExceeded), "deadline_exceeded"); + assert_eq!( + grpc_code_label(Code::ResourceExhausted), + "resource_exhausted" + ); + assert_eq!(grpc_code_label(Code::Ok), "ok"); + } +} diff --git a/crates/openshell-server/src/lib.rs b/crates/openshell-server/src/lib.rs index 72efa01fbf..1d8fbcfc2e 100644 --- a/crates/openshell-server/src/lib.rs +++ b/crates/openshell-server/src/lib.rs @@ -22,6 +22,7 @@ mod config_update_operation; mod credentials; mod defaults; mod gateway_listener; +mod gateway_metrics; mod gateway_ocsf; mod grpc; mod http; @@ -53,7 +54,6 @@ mod tracing_setup; mod watch_cursor; mod ws_tunnel; -use metrics_exporter_prometheus::PrometheusBuilder; use openshell_core::net::set_tcp_nodelay_best_effort; use openshell_core::telemetry::TelemetryComputeDriver; use openshell_core::{Config, Error, ObjectLabels, Result}; @@ -894,9 +894,9 @@ pub(crate) async fn run_server( // Bind the Prometheus metrics endpoint on a dedicated port when configured. if let Some(metrics_bind_address) = config.metrics_bind_address { - let prometheus_handle = PrometheusBuilder::new() - .install_recorder() - .map_err(|e| Error::config(format!("failed to install metrics recorder: {e}")))?; + let prometheus_handle = + gateway_metrics::install_global_recorder(supervisor_session::RELAY_CAPACITY) + .map_err(|e| Error::config(format!("failed to install metrics recorder: {e}")))?; let metrics_listener = TcpListener::bind(metrics_bind_address).await.map_err(|e| { Error::transport(format!( "failed to bind metrics port {metrics_bind_address}: {e}", diff --git a/crates/openshell-server/src/supervisor_session.rs b/crates/openshell-server/src/supervisor_session.rs index 7fc11594af..494e0730f0 100644 --- a/crates/openshell-server/src/supervisor_session.rs +++ b/crates/openshell-server/src/supervisor_session.rs @@ -27,6 +27,10 @@ use openshell_core::transport_errors::is_expected_transport_close_status; use crate::ServerState; use crate::auth::principal::Principal; +use crate::gateway_metrics::{ + self, GaugeSlot, PeerRpc, RelayCapacity, RelayKind, RelayRejection, RelayRoute, + RoutedRequestTimer, +}; use crate::grpc::provider_readiness::ProviderReadinessEvidence; use crate::persistence::ObjectId; use crate::supervisor_owner::{OWNER_TTL, OwnerError, OwnerGuard, SupervisorOwnerIndex}; @@ -48,6 +52,11 @@ const MAX_PENDING_RELAYS: usize = 256; /// consume the entire global budget. Sits above the SSH-tunnel per-sandbox /// cap (20) so tunnel-specific limits still fire first for that caller. const MAX_PENDING_RELAYS_PER_SANDBOX: usize = 32; +/// The replica relay cap above, published as a capacity gauge when the metrics recorder is +/// installed. +pub(crate) const RELAY_CAPACITY: RelayCapacity = RelayCapacity { + per_replica: MAX_PENDING_RELAYS, +}; const PEER_TLS_CA_FILE_ENV: &str = "OPENSHELL_PEER_TLS_CA_FILE"; const PEER_TLS_CERT_FILE_ENV: &str = "OPENSHELL_PEER_TLS_CERT_FILE"; const PEER_TLS_KEY_FILE_ENV: &str = "OPENSHELL_PEER_TLS_KEY_FILE"; @@ -297,6 +306,9 @@ struct LiveSession { provider_readiness: Option, #[allow(dead_code)] connected_at: Instant, + /// This session's share of `openshell_server_supervisor_sessions`, released when the entry + /// leaves the registry by any path (supersede, remove, disconnect, cleanup). + _gauge_slot: GaugeSlot, } /// Idempotency state for tool server endpoint-status reports from one live supervisor. @@ -338,6 +350,8 @@ struct PendingRelay { created_at: Instant, /// Last session whose outbound queue received this `RelayOpen`. delivered_session_id: Option, + /// This relay's share of `openshell_server_relay_pending`. + _gauge_slot: GaugeSlot, } #[derive(Debug)] @@ -418,6 +432,7 @@ impl SupervisorSessionRegistry { endpoint_report_cursor: None, provider_readiness: None, connected_at: Instant::now(), + _gauge_slot: GaugeSlot::supervisor_session(), }, ); match previous { @@ -857,6 +872,7 @@ impl SupervisorSessionRegistry { { let mut pending = self.pending_relays.lock().unwrap(); if pending.len() >= MAX_PENDING_RELAYS { + gateway_metrics::record_relay_rejected(RelayRejection::ReplicaCapacity); return Err(Status::resource_exhausted(format!( "gateway relay capacity reached ({MAX_PENDING_RELAYS} in flight)" ))); @@ -866,6 +882,7 @@ impl SupervisorSessionRegistry { .filter(|p| p.sandbox_id == sandbox_id) .count(); if per_sandbox >= MAX_PENDING_RELAYS_PER_SANDBOX { + gateway_metrics::record_relay_rejected(RelayRejection::SandboxCapacity); return Err(Status::resource_exhausted(format!( "per-sandbox relay limit reached ({MAX_PENDING_RELAYS_PER_SANDBOX} in flight for {sandbox_id})" ))); @@ -878,6 +895,7 @@ impl SupervisorSessionRegistry { relay_open: relay_open.clone(), created_at: Instant::now(), delivered_session_id: Some(session_id), + _gauge_slot: GaugeSlot::relay_pending(), }, ); // Insertion, delivery selection, and enqueueing are atomic with @@ -900,13 +918,19 @@ impl SupervisorSessionRegistry { } pub fn fail_pending_relay(&self, channel_id: &str, error: String) -> bool { - let pending = self.pending_relays.lock().unwrap().remove(channel_id); - if let Some(pending) = pending { - let _ = pending.sender.send(Err(Status::unavailable(error))); - true - } else { - false - } + // The rest of the entry, including its gauge slot, drops inside this statement while the + // lock is still held, so `relay_pending` never exceeds capacity. + let Some(sender) = self + .pending_relays + .lock() + .unwrap() + .remove(channel_id) + .map(|pending| pending.sender) + else { + return false; + }; + let _ = sender.send(Err(Status::unavailable(error))); + true } /// Claim a pending relay channel. Called by the `/relay/{channel_id}` HTTP handler @@ -920,7 +944,7 @@ impl SupervisorSessionRegistry { channel_id: &str, principal: Option<&Principal>, ) -> Result { - let pending = { + let (sender, sandbox_id) = { let mut map = self.pending_relays.lock().unwrap(); let pending = map .get(channel_id) @@ -940,13 +964,22 @@ impl SupervisorSessionRegistry { return Err(status); } - if pending.created_at.elapsed() > RELAY_PENDING_TIMEOUT { + let waited = pending.created_at.elapsed(); + if waited > RELAY_PENDING_TIMEOUT { map.remove(channel_id); + gateway_metrics::record_relay_expired(1); return Err(Status::deadline_exceeded("relay channel timed out")); } - - map.remove(channel_id) - .expect("pending relay existed before removal") + gateway_metrics::record_relay_claimed(waited); + + // The rest of the entry, including its gauge slot, drops at the end of this + // statement while the lock is still held, so `relay_pending` never exceeds capacity. + let PendingRelay { + sender, sandbox_id, .. + } = map + .remove(channel_id) + .expect("pending relay existed before removal"); + (sender, sandbox_id) }; // Create a duplex stream pair: one end for the gateway bridge, one for @@ -954,20 +987,25 @@ impl SupervisorSessionRegistry { let (gateway_stream, supervisor_stream) = tokio::io::duplex(64 * 1024); // Send the gateway-side stream to the waiter (exec handler or forward handler). - if pending.sender.send(Ok(gateway_stream)).is_err() { + if sender.send(Ok(gateway_stream)).is_err() { return Err(Status::internal("relay requester dropped")); } Ok(ClaimedRelay { stream: supervisor_stream, - sandbox_id: pending.sandbox_id, + sandbox_id, }) } /// Remove all pending relays that have exceeded the timeout. pub fn reap_expired_relays(&self) { - let mut map = self.pending_relays.lock().unwrap(); - map.retain(|_, pending| pending.created_at.elapsed() <= RELAY_PENDING_TIMEOUT); + let reaped = { + let mut map = self.pending_relays.lock().unwrap(); + let before = map.len(); + map.retain(|_, pending| pending.created_at.elapsed() <= RELAY_PENDING_TIMEOUT); + before - map.len() + }; + gateway_metrics::record_relay_expired(reaped); } /// Clean up all state for a sandbox (session + pending relays). @@ -1388,15 +1426,16 @@ pub(crate) async fn forward_provider_readiness_to_owner( request: ReportProviderReadinessRequest, ) -> Result { let sandbox_id = request.sandbox_id.clone(); - let mut client = peer_rpc_client(state, &owner.owner_peer_endpoint).await?; - client - .peer_report_provider_readiness(request) + let mut timer = RoutedRequestTimer::start(PeerRpc::ReportProviderReadiness); + let mut client = peer_rpc_client(state, &owner.owner_peer_endpoint) .await - .map(Response::into_inner) - .inspect_err(|_| { - state.peer_routes.evict_channel(&owner.owner_peer_endpoint); - state.peer_routes.evict_owner(&sandbox_id); - }) + .inspect_err(|status| timer.local_error(status))?; + let result = client.peer_report_provider_readiness(request).await; + timer.finish(&result); + result.map(Response::into_inner).inspect_err(|_| { + state.peer_routes.evict_channel(&owner.owner_peer_endpoint); + state.peer_routes.evict_owner(&sandbox_id); + }) } pub(crate) async fn forward_endpoint_status_to_owner( @@ -1405,15 +1444,16 @@ pub(crate) async fn forward_endpoint_status_to_owner( request: ReportEndpointStatusRequest, ) -> Result { let sandbox_id = request.sandbox_id.clone(); - let mut client = peer_rpc_client(state, &owner.owner_peer_endpoint).await?; - client - .peer_report_endpoint_status(request) + let mut timer = RoutedRequestTimer::start(PeerRpc::ReportEndpointStatus); + let mut client = peer_rpc_client(state, &owner.owner_peer_endpoint) .await - .map(Response::into_inner) - .inspect_err(|_| { - state.peer_routes.evict_channel(&owner.owner_peer_endpoint); - state.peer_routes.evict_owner(&sandbox_id); - }) + .inspect_err(|status| timer.local_error(status))?; + let result = client.peer_report_endpoint_status(request).await; + timer.finish(&result); + result.map(Response::into_inner).inspect_err(|_| { + state.peer_routes.evict_channel(&owner.owner_peer_endpoint); + state.peer_routes.evict_owner(&sandbox_id); + }) } pub(crate) async fn forward_provider_status_query_to_owner( @@ -1422,15 +1462,16 @@ pub(crate) async fn forward_provider_status_query_to_owner( sandbox_id: &str, request: GetSandboxProviderStatusRequest, ) -> Result { - let mut client = peer_rpc_client(state, &owner.owner_peer_endpoint).await?; - client - .peer_get_sandbox_provider_status(request) + let mut timer = RoutedRequestTimer::start(PeerRpc::GetSandboxProviderStatus); + let mut client = peer_rpc_client(state, &owner.owner_peer_endpoint) .await - .map(Response::into_inner) - .inspect_err(|_| { - state.peer_routes.evict_channel(&owner.owner_peer_endpoint); - state.peer_routes.evict_owner(sandbox_id); - }) + .inspect_err(|status| timer.local_error(status))?; + let result = client.peer_get_sandbox_provider_status(request).await; + timer.finish(&result); + result.map(Response::into_inner).inspect_err(|_| { + state.peer_routes.evict_channel(&owner.owner_peer_endpoint); + state.peer_routes.evict_owner(sandbox_id); + }) } pub async fn open_routed_relay_with_target( @@ -1455,6 +1496,68 @@ pub async fn open_routed_relay_with_target( open_routed_relay_with_message(state, sandbox_id, relay_open, session_wait_timeout).await } +fn relay_kind(relay_open: &RelayOpen) -> RelayKind { + // An absent target means SSH for compatibility with older callers. + match relay_open.target.as_ref() { + Some(relay_open::Target::Ssh(_)) | None => RelayKind::Ssh, + Some(relay_open::Target::Tcp(_)) => RelayKind::Tcp, + } +} + +/// Hand the caller a receiver that forwards the local relay's claim result unchanged, and count +/// the attempt when the supervisor claims the relay, the claim window closes, or the caller +/// gives up. The owner of a `PeerRelay` answers on the same events, so a relay outcome means +/// the same thing on the local and peer routes. +/// +/// The window is anchored before the caller can start its own wait, so a caller that times out +/// with the same 10 s budget is counted as an unclaimed relay, not as a cancellation. A claim +/// that lands just after the caller gave up still succeeds in the registry; the forwarder then +/// drops the stream and the supervisor sees it close, as when a caller drops right after a claim. +fn track_local_relay_claim( + mut timer: RoutedRequestTimer, + mut claimed: oneshot::Receiver>, +) -> oneshot::Receiver> { + let claim_deadline = tokio::time::Instant::now() + RELAY_PENDING_TIMEOUT; + let unclaimed = || Status::deadline_exceeded("relay was not claimed in time"); + let (mut forward_tx, forward_rx) = oneshot::channel(); + tokio::spawn(async move { + let claim_window = tokio::time::sleep_until(claim_deadline); + tokio::pin!(claim_window); + let mut window_open = true; + let result = loop { + tokio::select! { + biased; + result = &mut claimed => break result, + () = &mut claim_window, if window_open => { + window_open = false; + timer.local_error(&unclaimed()); + } + () = forward_tx.closed() => { + if tokio::time::Instant::now() >= claim_deadline { + timer.local_error(&unclaimed()); + } + // Otherwise dropping the timer records the cancellation. + return; + } + } + }; + match result { + Ok(claim) => { + timer.finish(&claim); + let _ = forward_tx.send(claim); + } + // The registry dropped the relay without an answer: it expired (reaper or a late + // claim), or the registry was torn down. Dropping `forward_tx` passes the same closed + // channel on to the caller. + Err(_) if tokio::time::Instant::now() >= claim_deadline => { + timer.local_error(&unclaimed()); + } + Err(_) => timer.local_error(&Status::unavailable("relay channel dropped")), + } + }); + forward_rx +} + pub async fn open_routed_relay_with_message( state: &Arc, sandbox_id: &str, @@ -1473,12 +1576,19 @@ pub async fn open_routed_relay_with_message( let owner_index = SupervisorOwnerIndex::new(state.store.clone(), OWNER_TTL); loop { if state.supervisor_sessions.has_session(sandbox_id) { - match state + let mut timer = + RoutedRequestTimer::relay(relay_kind(&relay_open), RelayRoute::Local); + let result = state .supervisor_sessions .open_relay_with_message_until(sandbox_id, relay_open.clone(), deadline, false) - .await - { - Ok(relay) => return Ok(relay), + .await; + if let Err(status) = &result { + timer.local_error(status); + } + match result { + Ok((channel_id, relay_rx)) => { + return Ok((channel_id, track_local_relay_claim(timer, relay_rx))); + } Err(status) if status.code() == tonic::Code::Unavailable => { // The session can migrate after `has_session` but before // RelayOpen reaches its sender. Fall through and reread the @@ -1605,15 +1715,30 @@ async fn open_peer_relay( Ok((channel_id, relay_rx)) } +/// Open a `PeerRelay` stream to the owner replica and bridge it to a local duplex stream. +/// +/// The peer request metrics count attempts, so the routed-relay retry loop spikes `unavailable` +/// during rollouts. `ok` means the owner's supervisor claimed the relay (response headers +/// arrived); later bridge failures are not counted. async fn connect_peer_relay( state: &Arc, owner_peer_endpoint: &str, sandbox_id: &str, relay_open: RelayOpen, ) -> Result { - let token = state.peer_routes.peer_token().await?; - let channel = state.peer_routes.channel(owner_peer_endpoint).await?; - let interceptor = PeerAuthInterceptor::new(&token, &state.replica_id)?; + let mut timer = RoutedRequestTimer::relay(relay_kind(&relay_open), RelayRoute::Peer); + let token = state + .peer_routes + .peer_token() + .await + .inspect_err(|s| timer.local_error(s))?; + let channel = state + .peer_routes + .channel(owner_peer_endpoint) + .await + .inspect_err(|s| timer.local_error(s))?; + let interceptor = PeerAuthInterceptor::new(&token, &state.replica_id) + .inspect_err(|s| timer.local_error(s))?; let mut client = open_shell_client::OpenShellClient::with_interceptor(channel, interceptor); let (out_tx, out_rx) = mpsc::channel::(16); @@ -1626,15 +1751,16 @@ async fn connect_peer_relay( })), }) .await - .map_err(|_| Status::internal("failed to initialize peer relay stream"))?; - - let response = client - .peer_relay(ReceiverStream::new(out_rx)) - .await - .map_err(|err| { - state.peer_routes.evict_channel(owner_peer_endpoint); - Status::unavailable(format!("gateway peer relay RPC failed: {err}")) - })?; + .map_err(|_| Status::internal("failed to initialize peer relay stream")) + .inspect_err(|s| timer.local_error(s))?; + + let result = client.peer_relay(ReceiverStream::new(out_rx)).await; + // Record the owner's code before the remap below hides it as `unavailable`. + timer.finish(&result); + let response = result.map_err(|err| { + state.peer_routes.evict_channel(owner_peer_endpoint); + Status::unavailable(format!("gateway peer relay RPC failed: {err}")) + })?; let inbound = response.into_inner(); let (gateway_stream, bridge_stream) = tokio::io::duplex(64 * 1024); spawn_peer_bridge(bridge_stream, inbound, out_tx, sandbox_id.to_string()); @@ -2359,7 +2485,12 @@ mod tests { use super::*; use crate::auth::identity::{Identity, IdentityProvider}; use crate::auth::principal::{SandboxIdentitySource, SandboxPrincipal, UserPrincipal}; + use crate::gateway_metrics::MetricsCapture; use crate::persistence::Store; + use bytes::Bytes; + use http_body::Frame; + use http_body_util::{BodyExt, Empty, StreamBody}; + use std::convert::Infallible; use tokio::io::{AsyncReadExt, AsyncWriteExt}; async fn test_store() -> Arc { @@ -2587,6 +2718,7 @@ mod tests { }, created_at, delivered_session_id: None, + _gauge_slot: GaugeSlot::relay_pending(), } } @@ -2788,10 +2920,46 @@ mod tests { assert_eq!(registry.remove_if_current("sbx", "s1"), Some(true)); } + #[test] + fn session_gauge_tracks_register_supersede_and_removal() { + let metrics = MetricsCapture::install(); + let registry = SupervisorSessionRegistry::new(); + let (tx, _rx) = mpsc::channel(1); + + registry.register( + "sbx-a".to_string(), + "s1".to_string(), + tx.clone(), + make_shutdown(), + ); + assert_eq!(metrics.value(gateway_metrics::SUPERVISOR_SESSIONS), Some(1)); + registry.register( + "sbx-a".to_string(), + "s2".to_string(), + tx.clone(), + make_shutdown(), + ); + assert_eq!( + metrics.value(gateway_metrics::SUPERVISOR_SESSIONS), + Some(1), + "a supersede on the same replica nets zero" + ); + registry.register("sbx-b".to_string(), "s3".to_string(), tx, make_shutdown()); + assert_eq!(metrics.value(gateway_metrics::SUPERVISOR_SESSIONS), Some(2)); + + assert_eq!(registry.remove_if_current("sbx-a", "s1"), None); + assert_eq!(metrics.value(gateway_metrics::SUPERVISOR_SESSIONS), Some(2)); + assert_eq!(registry.remove_if_current("sbx-a", "s2"), Some(false)); + assert_eq!(metrics.value(gateway_metrics::SUPERVISOR_SESSIONS), Some(1)); + assert!(registry.disconnect("sbx-b")); + assert_eq!(metrics.value(gateway_metrics::SUPERVISOR_SESSIONS), Some(0)); + } + // ---- open_relay: happy path and wait semantics ---- #[tokio::test] async fn open_relay_sends_relay_open_to_registered_session() { + let metrics = MetricsCapture::install(); let registry = SupervisorSessionRegistry::new(); let (tx, mut rx) = mpsc::channel(4); registry.register("sbx".to_string(), "s1".to_string(), tx, make_shutdown()); @@ -2809,6 +2977,12 @@ mod tests { } other => panic!("expected RelayOpen, got {other:?}"), } + assert!( + !metrics + .render() + .contains(gateway_metrics::ROUTED_REQUEST_ATTEMPTS_TOTAL), + "owner-side registry opens must not count another routing attempt" + ); } #[tokio::test] @@ -2849,6 +3023,7 @@ mod tests { #[tokio::test] async fn open_relay_fails_when_session_receiver_dropped() { + let metrics = MetricsCapture::install(); let registry = SupervisorSessionRegistry::new(); let (tx, rx) = mpsc::channel::(4); registry.register("sbx".to_string(), "s1".to_string(), tx, make_shutdown()); @@ -2864,10 +3039,14 @@ mod tests { assert_eq!(err.code(), tonic::Code::Unavailable); // The pending-relay entry must have been cleaned up on failure. assert!(registry.pending_relays.lock().unwrap().is_empty()); + // Queue capacity is reserved before insertion, so no gauge slot was ever taken. + assert_eq!(metrics.value(gateway_metrics::RELAY_PENDING), None); + assert_eq!(metrics.value(gateway_metrics::RELAY_EXPIRED_TOTAL), None); } #[tokio::test] async fn open_relay_rejects_when_global_cap_reached() { + let metrics = MetricsCapture::install(); let registry = SupervisorSessionRegistry::new(); let (tx, _rx) = mpsc::channel::(8); registry.register( @@ -2898,10 +3077,20 @@ mod tests { .expect_err("open_relay should reject once global cap is reached"); assert_eq!(err.code(), tonic::Code::ResourceExhausted); assert!(err.message().contains("gateway relay capacity")); + assert_eq!( + metrics.value("openshell_server_relay_rejected_total{reason=\"replica_capacity\"}"), + Some(1) + ); + assert_eq!(metrics.value(gateway_metrics::RELAY_PENDING), Some(256)); + assert_eq!( + metrics.value("openshell_server_relay_rejected_total{reason=\"sandbox_capacity\"}"), + None + ); } #[tokio::test] async fn open_relay_rejects_when_per_sandbox_cap_reached() { + let metrics = MetricsCapture::install(); let registry = SupervisorSessionRegistry::new(); let (tx, _rx) = mpsc::channel::(8); registry.register("sbx".to_string(), "s".to_string(), tx, make_shutdown()); @@ -2923,6 +3112,11 @@ mod tests { .expect_err("open_relay should reject when per-sandbox cap is reached"); assert_eq!(err.code(), tonic::Code::ResourceExhausted); assert!(err.message().contains("per-sandbox relay limit")); + assert_eq!( + metrics.value("openshell_server_relay_rejected_total{reason=\"sandbox_capacity\"}"), + Some(1) + ); + assert_eq!(metrics.value(gateway_metrics::RELAY_PENDING), Some(32)); // A different sandbox still has headroom. let (tx2, _rx2) = mpsc::channel::(8); @@ -2936,6 +3130,7 @@ mod tests { .open_relay("sbx-other", Duration::from_millis(50)) .await .expect("different sandbox should still accept new relays"); + assert_eq!(metrics.value(gateway_metrics::RELAY_PENDING), Some(33)); } #[tokio::test] @@ -3344,6 +3539,7 @@ mod tests { #[test] fn claim_relay_success() { + let metrics = MetricsCapture::install(); let registry = SupervisorSessionRegistry::new(); let (relay_tx, _relay_rx) = oneshot::channel(); registry.pending_relays.lock().unwrap().insert( @@ -3355,10 +3551,86 @@ mod tests { let result = registry.claim_relay("ch-1", Some(&principal)); assert!(result.is_ok()); assert!(!registry.pending_relays.lock().unwrap().contains_key("ch-1")); + assert_eq!(metrics.value(gateway_metrics::RELAY_PENDING), Some(0)); + assert_eq!( + metrics.value("openshell_server_relay_claim_duration_seconds_count"), + Some(1) + ); + assert_eq!(metrics.value(gateway_metrics::RELAY_EXPIRED_TOTAL), None); + } + + /// Waker that reads `relay_pending` each time it is woken. `oneshot::Sender::send` wakes a + /// registered receiver synchronously, so the probe sees the gauge exactly as a waiter on + /// another worker thread could at that instant. + struct PendingGaugeProbe { + read: Box Option + Send + Sync>, + seen: Mutex>>, + } + + impl PendingGaugeProbe { + fn register(metrics: &MetricsCapture, rx: &mut oneshot::Receiver) -> Arc { + let probe = Arc::new(Self { + read: metrics.value_reader(gateway_metrics::RELAY_PENDING), + seen: Mutex::new(Vec::new()), + }); + let waker = std::task::Waker::from(Arc::clone(&probe)); + let mut cx = std::task::Context::from_waker(&waker); + assert!(Pin::new(rx).poll(&mut cx).is_pending()); + probe + } + + fn seen(&self) -> Vec> { + self.seen.lock().unwrap().clone() + } + } + + impl std::task::Wake for PendingGaugeProbe { + fn wake(self: Arc) { + self.wake_by_ref(); + } + + fn wake_by_ref(self: &Arc) { + self.seen.lock().unwrap().push((self.read)()); + } + } + + #[test] + fn claim_relay_releases_pending_slot_before_waking_waiter() { + let metrics = MetricsCapture::install(); + let registry = SupervisorSessionRegistry::new(); + let (relay_tx, mut relay_rx) = oneshot::channel(); + registry.pending_relays.lock().unwrap().insert( + "ch-1".to_string(), + pending_relay("sbx-test", relay_tx, Instant::now()), + ); + let probe = PendingGaugeProbe::register(&metrics, &mut relay_rx); + + registry + .claim_relay("ch-1", Some(&sandbox_principal("sbx-test"))) + .expect("claim should succeed"); + // The slot is released under the pending lock, before the waiter is woken, so a + // concurrent open can never push `relay_pending` above capacity. + assert_eq!(probe.seen(), vec![Some(0)]); + } + + #[test] + fn fail_pending_relay_releases_pending_slot_before_waking_waiter() { + let metrics = MetricsCapture::install(); + let registry = SupervisorSessionRegistry::new(); + let (relay_tx, mut relay_rx) = oneshot::channel(); + registry.pending_relays.lock().unwrap().insert( + "ch-fail".to_string(), + pending_relay("sbx-test", relay_tx, Instant::now()), + ); + let probe = PendingGaugeProbe::register(&metrics, &mut relay_rx); + + assert!(registry.fail_pending_relay("ch-fail", "target refused".to_string())); + assert_eq!(probe.seen(), vec![Some(0)]); } #[test] fn claim_relay_rejects_cross_sandbox_principal_without_consuming_channel() { + let metrics = MetricsCapture::install(); let registry = SupervisorSessionRegistry::new(); let (relay_tx, _relay_rx) = oneshot::channel(); registry.pending_relays.lock().unwrap().insert( @@ -3379,6 +3651,11 @@ mod tests { .contains_key("ch-cross"), "failed cross-sandbox claim must not consume the channel" ); + assert_eq!(metrics.value(gateway_metrics::RELAY_PENDING), Some(1)); + assert_eq!( + metrics.value("openshell_server_relay_claim_duration_seconds_count"), + None + ); } #[test] @@ -3398,6 +3675,7 @@ mod tests { #[tokio::test] async fn relay_open_failure_completes_pending_waiter() { + let metrics = MetricsCapture::install(); let registry = SupervisorSessionRegistry::new(); let (relay_tx, relay_rx) = oneshot::channel(); registry.pending_relays.lock().unwrap().insert( @@ -3418,10 +3696,13 @@ mod tests { let status = result.expect_err("waiter should receive status failure"); assert_eq!(status.code(), tonic::Code::Unavailable); assert_eq!(status.message(), "target refused"); + assert_eq!(metrics.value(gateway_metrics::RELAY_PENDING), Some(0)); + assert_eq!(metrics.value(gateway_metrics::RELAY_EXPIRED_TOTAL), None); } #[test] fn claim_relay_expired_returns_deadline_exceeded() { + let metrics = MetricsCapture::install(); let registry = SupervisorSessionRegistry::new(); let (relay_tx, _relay_rx) = oneshot::channel(); registry.pending_relays.lock().unwrap().insert( @@ -3447,10 +3728,17 @@ mod tests { .unwrap() .contains_key("ch-old") ); + assert_eq!(metrics.value(gateway_metrics::RELAY_EXPIRED_TOTAL), Some(1)); + assert_eq!(metrics.value(gateway_metrics::RELAY_PENDING), Some(0)); + assert_eq!( + metrics.value("openshell_server_relay_claim_duration_seconds_count"), + None + ); } #[test] fn claim_relay_receiver_dropped_returns_internal() { + let metrics = MetricsCapture::install(); let registry = SupervisorSessionRegistry::new(); let (relay_tx, relay_rx) = oneshot::channel::>(); drop(relay_rx); // Gateway-side waiter has given up already. @@ -3463,6 +3751,11 @@ mod tests { .claim_relay("ch-1", Some(&sandbox_principal("sbx-test"))) .expect_err("should err when receiver is gone"); assert_eq!(err.code(), tonic::Code::Internal); + assert_eq!( + metrics.value("openshell_server_relay_claim_duration_seconds_count"), + Some(1) + ); + assert_eq!(metrics.value(gateway_metrics::RELAY_PENDING), Some(0)); } #[tokio::test] @@ -3500,6 +3793,7 @@ mod tests { #[test] fn reap_expired_relays_removes_old_entries() { + let metrics = MetricsCapture::install(); let registry = SupervisorSessionRegistry::new(); let (relay_tx, _relay_rx) = oneshot::channel(); registry.pending_relays.lock().unwrap().insert( @@ -3521,10 +3815,13 @@ mod tests { .unwrap() .contains_key("ch-old") ); + assert_eq!(metrics.value(gateway_metrics::RELAY_EXPIRED_TOTAL), Some(1)); + assert_eq!(metrics.value(gateway_metrics::RELAY_PENDING), Some(0)); } #[test] fn reap_expired_relays_keeps_fresh_entries() { + let metrics = MetricsCapture::install(); let registry = SupervisorSessionRegistry::new(); let (relay_tx, _relay_rx) = oneshot::channel(); registry.pending_relays.lock().unwrap().insert( @@ -3540,6 +3837,9 @@ mod tests { .unwrap() .contains_key("ch-fresh") ); + // Reaping nothing records nothing. + assert_eq!(metrics.value(gateway_metrics::RELAY_EXPIRED_TOTAL), None); + assert_eq!(metrics.value(gateway_metrics::RELAY_PENDING), Some(1)); } fn owner_record(replica: &str) -> crate::supervisor_owner::OwnerRecord { @@ -3618,6 +3918,630 @@ mod tests { assert!(OWNER_CACHE_TTL < OWNER_TTL); } + // ---- peer request metrics (requester side) ---- + + #[derive(Clone, Copy)] + enum FakePeerReply { + Status(tonic::Code), + EmptyOk, + } + + /// Minimal h2c server that answers every gRPC call the same way. It stands in for an owner + /// replica without implementing the full `OpenShell` service. + async fn spawn_fake_peer(reply: FakePeerReply) -> String { + let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap(); + let addr = listener.local_addr().unwrap(); + tokio::spawn(async move { + while let Ok((stream, _)) = listener.accept().await { + tokio::spawn(async move { + let service = hyper::service::service_fn( + move |_req: http::Request| async move { + Ok::<_, Infallible>(fake_peer_response(reply)) + }, + ); + let _ = hyper_util::server::conn::auto::Builder::new( + hyper_util::rt::TokioExecutor::new(), + ) + .serve_connection(hyper_util::rt::TokioIo::new(stream), service) + .await; + }); + } + }); + format!("http://{addr}") + } + + fn fake_peer_response( + reply: FakePeerReply, + ) -> http::Response> { + let builder = http::Response::builder() + .status(200) + .header("content-type", "application/grpc"); + match reply { + // Trailers-only error: tonic returns Err(status) for unary and streaming calls. + FakePeerReply::Status(code) => builder + .header("grpc-status", i32::from(code).to_string()) + .header("grpc-message", "fake peer") + .body(Empty::new().boxed_unsync()) + .unwrap(), + // One empty message (5-byte frame header, zero length), then grpc-status 0. This + // decodes as a default response for any unary RPC, and gives streaming calls an OK + // header. + FakePeerReply::EmptyOk => { + let mut trailers = http::HeaderMap::new(); + trailers.insert("grpc-status", http::HeaderValue::from_static("0")); + let frames = futures::stream::iter([ + Ok::<_, Infallible>(Frame::data(Bytes::from_static(&[0, 0, 0, 0, 0]))), + Ok(Frame::trailers(trailers)), + ]); + builder + .body(StreamBody::new(frames).boxed_unsync()) + .unwrap() + } + } + } + + fn seed_peer_token(state: &ServerState) { + *state.peer_routes.token.lock().unwrap() = Some(CachedPeerToken { + token: "test-peer-token".to_string(), + refresh_at: Instant::now() + Duration::from_mins(5), + }); + } + + fn owner_at(endpoint: &str) -> crate::supervisor_owner::OwnerRecord { + let mut owner = owner_record("replica-owner"); + owner.owner_peer_endpoint = endpoint.to_string(); + owner + } + + fn closed_local_endpoint() -> String { + let listener = std::net::TcpListener::bind("127.0.0.1:0").unwrap(); + let addr = listener.local_addr().unwrap(); + drop(listener); + format!("http://{addr}") + } + + fn peer_relay_open(channel_id: &str) -> RelayOpen { + RelayOpen { + channel_id: channel_id.to_string(), + target: Some(relay_open::Target::Ssh(SshRelayTarget {})), + service_id: String::new(), + } + } + + #[tokio::test] + async fn routed_relay_metrics_count_local_targets_including_legacy_ssh() { + let metrics = MetricsCapture::install(); + let state = crate::grpc::test_support::test_server_state().await; + let (tx, mut rx) = mpsc::channel(4); + state.supervisor_sessions.register( + "sbx-routing".into(), + "session-routing".into(), + tx, + make_shutdown(), + ); + + for (target, label, expected) in [ + (Some(relay_open::Target::Ssh(SshRelayTarget {})), "ssh", 1), + ( + Some(relay_open::Target::Tcp( + openshell_core::proto::TcpRelayTarget { + host: "127.0.0.1".into(), + port: 12345, + }, + )), + "tcp", + 1, + ), + (None, "ssh", 2), + ] { + let relay_open = RelayOpen { + target, + ..peer_relay_open(&Uuid::new_v4().to_string()) + }; + let (channel_id, relay_rx) = open_routed_relay_with_message( + &state, + "sbx-routing", + relay_open.clone(), + Duration::from_secs(1), + ) + .await + .unwrap(); + assert_eq!(channel_id, relay_open.channel_id); + assert_eq!( + rx.recv().await.unwrap().payload, + Some(gateway_message::Payload::RelayOpen(relay_open)) + ); + let series = format!( + "openshell_server_routed_request_attempts_total{{operation=\"relay\",route=\"local\",relay_kind=\"{label}\",outcome=\"success\",grpc_code=\"ok\"}}" + ); + // Enqueued but not claimed yet: success waits for the supervisor. + assert_eq!(metrics.value(&series).unwrap_or(0), expected - 1); + let _claimed = state + .supervisor_sessions + .claim_relay(&channel_id, None) + .unwrap(); + relay_rx.await.unwrap().unwrap(); + assert_eq!(metrics.value(&series), Some(expected)); + } + let rendered = metrics.render(); + assert!(!rendered.contains("route=\"peer\"")); + for identifier in ["sbx-routing", "session-routing", "127.0.0.1", "12345"] { + assert!( + !rendered.contains(identifier), + "must not label with {identifier}" + ); + } + } + + #[tokio::test] + async fn routed_relay_metrics_count_failed_local_setup() { + let metrics = MetricsCapture::install(); + let state = crate::grpc::test_support::test_server_state().await; + let (tx, rx) = mpsc::channel(1); + drop(rx); + state.supervisor_sessions.register( + "sbx-routing".into(), + "session-routing".into(), + tx, + make_shutdown(), + ); + + open_routed_relay_with_message( + &state, + "sbx-routing", + peer_relay_open("ch-routing"), + Duration::from_millis(50), + ) + .await + .expect_err("the local supervisor disconnected"); + assert_eq!( + metrics.value( + "openshell_server_routed_request_attempts_total{operation=\"relay\",route=\"local\",relay_kind=\"ssh\",outcome=\"local_error\",grpc_code=\"unavailable\"}" + ), + Some(1) + ); + assert!(!metrics.render().contains("route=\"peer\"")); + } + + #[tokio::test] + async fn routed_relay_metrics_count_local_to_peer_fallback() { + let metrics = MetricsCapture::install(); + let state = crate::grpc::test_support::test_server_state().await; + seed_peer_token(&state); + let endpoint = spawn_fake_peer(FakePeerReply::EmptyOk).await; + state + .peer_routes + .store_owner("sbx-routing", &owner_at(&endpoint)); + let (tx, rx) = mpsc::channel(1); + drop(rx); + state.supervisor_sessions.register( + "sbx-routing".into(), + "old-session".into(), + tx, + make_shutdown(), + ); + + open_routed_relay_with_message( + &state, + "sbx-routing", + peer_relay_open("ch-routing"), + Duration::from_secs(5), + ) + .await + .expect("the peer accepts after the local supervisor disconnects"); + for route in ["local", "peer"] { + let (outcome, code) = if route == "local" { + ("local_error", "unavailable") + } else { + ("success", "ok") + }; + assert_eq!( + metrics.value(&format!( + "openshell_server_routed_request_attempts_total{{operation=\"relay\",route=\"{route}\",relay_kind=\"ssh\",outcome=\"{outcome}\",grpc_code=\"{code}\"}}" + )), + Some(1) + ); + } + } + + #[tokio::test] + async fn routed_relay_metrics_count_peer_targets() { + let metrics = MetricsCapture::install(); + let state = crate::grpc::test_support::test_server_state().await; + seed_peer_token(&state); + let endpoint = spawn_fake_peer(FakePeerReply::EmptyOk).await; + state + .peer_routes + .store_owner("sbx-routing", &owner_at(&endpoint)); + + for (target, label) in [ + (relay_open::Target::Ssh(SshRelayTarget {}), "ssh"), + ( + relay_open::Target::Tcp(openshell_core::proto::TcpRelayTarget { + host: "127.0.0.1".into(), + port: 12345, + }), + "tcp", + ), + ] { + let (_channel_id, _relay_rx) = open_routed_relay_with_target( + &state, + "sbx-routing", + target, + String::new(), + Duration::from_secs(5), + ) + .await + .unwrap(); + assert_eq!( + metrics.value(&format!( + "openshell_server_routed_request_attempts_total{{operation=\"relay\",route=\"peer\",relay_kind=\"{label}\",outcome=\"success\",grpc_code=\"ok\"}}" + )), + Some(1) + ); + } + let rendered = metrics.render(); + assert!(!rendered.contains("route=\"local\"")); + assert!(!rendered.contains(endpoint.trim_start_matches("http://"))); + assert!(!rendered.contains("sbx-routing")); + } + + #[tokio::test] + async fn routed_relay_metrics_count_each_failed_peer_retry() { + let metrics = MetricsCapture::install(); + let state = crate::grpc::test_support::test_server_state().await; + seed_peer_token(&state); + let endpoint = spawn_fake_peer(FakePeerReply::Status(tonic::Code::Unavailable)).await; + SupervisorOwnerIndex::new(state.store.clone(), OWNER_TTL) + .publish( + "sbx-routing", + "session", + "instance", + 1, + "replica-owner", + &endpoint, + ) + .await + .unwrap(); + + open_routed_relay_with_message( + &state, + "sbx-routing", + peer_relay_open("ch-routing"), + Duration::from_secs(1), + ) + .await + .expect_err("the peer rejects every attempt"); + let attempts = metrics + .value("openshell_server_routed_request_attempts_total{operation=\"relay\",route=\"peer\",relay_kind=\"ssh\",outcome=\"remote_error\",grpc_code=\"unavailable\"}") + .unwrap(); + assert!(attempts > 1, "the routing loop must have retried"); + assert_eq!( + metrics.value( + "openshell_server_peer_request_duration_seconds_count{operation=\"relay\",outcome=\"remote_error\"}" + ), + Some(attempts) + ); + assert!(!metrics.render().contains("route=\"local\"")); + } + + #[tokio::test] + async fn routed_relay_metrics_do_not_count_waiting_for_an_owner() { + let metrics = MetricsCapture::install(); + let state = crate::grpc::test_support::test_server_state().await; + open_routed_relay_with_message( + &state, + "sbx-missing", + peer_relay_open("ch-routing"), + Duration::from_millis(250), + ) + .await + .expect_err("no supervisor or owner is available"); + assert!( + !metrics + .render() + .contains(gateway_metrics::ROUTED_REQUEST_ATTEMPTS_TOTAL) + ); + } + + #[tokio::test] + async fn routed_relay_metrics_count_cancelled_local_setup_once() { + let metrics = MetricsCapture::install(); + let state = crate::grpc::test_support::test_server_state().await; + let (tx, _rx) = mpsc::channel(1); + tx.try_send(GatewayMessage::default()).unwrap(); + state.supervisor_sessions.register( + "sbx-routing".into(), + "session-routing".into(), + tx, + make_shutdown(), + ); + let mut setup = Box::pin(open_routed_relay_with_message( + &state, + "sbx-routing", + peer_relay_open("ch-routing"), + Duration::from_secs(5), + )); + tokio::select! { + result = &mut setup => panic!("setup should wait for queue space: {result:?}"), + () = tokio::time::sleep(Duration::from_millis(10)) => {} + } + assert!( + !metrics + .render() + .contains(gateway_metrics::ROUTED_REQUEST_ATTEMPTS_TOTAL) + ); + drop(setup); + assert_eq!( + metrics.value("openshell_server_routed_request_attempts_total{operation=\"relay\",route=\"local\",relay_kind=\"ssh\",outcome=\"local_error\",grpc_code=\"cancelled\"}"), + Some(1) + ); + assert!( + !metrics + .render() + .contains(gateway_metrics::PEER_REQUEST_DURATION_SECONDS) + ); + } + + const LOCAL_RELAY_SERIES: &str = "openshell_server_routed_request_attempts_total{operation=\"relay\",route=\"local\",relay_kind=\"ssh\""; + + fn local_relay_outcome(metrics: &MetricsCapture, outcome: &str, code: &str) -> Option { + metrics.value(&format!( + "{LOCAL_RELAY_SERIES},outcome=\"{outcome}\",grpc_code=\"{code}\"}}" + )) + } + + type RelayClaim = Result; + + fn tracked_local_relay() -> (oneshot::Sender, oneshot::Receiver) { + let (claim_tx, claim_rx) = oneshot::channel(); + let timer = RoutedRequestTimer::relay(RelayKind::Ssh, RelayRoute::Local); + (claim_tx, track_local_relay_claim(timer, claim_rx)) + } + + #[tokio::test] + async fn local_relay_claim_records_success_once() { + let metrics = MetricsCapture::install(); + let (claim_tx, forwarded) = tracked_local_relay(); + let (stream, _peer) = tokio::io::duplex(64); + claim_tx.send(Ok(stream)).unwrap(); + forwarded.await.unwrap().unwrap(); + assert_eq!(local_relay_outcome(&metrics, "success", "ok"), Some(1)); + assert!(!metrics.render().contains("local_error")); + } + + #[tokio::test(start_paused = true)] + async fn local_relay_unclaimed_in_window_records_deadline_exceeded_once() { + let metrics = MetricsCapture::install(); + let (claim_tx, forwarded) = tracked_local_relay(); + tokio::time::sleep(RELAY_PENDING_TIMEOUT + Duration::from_millis(1)).await; + assert_eq!( + local_relay_outcome(&metrics, "local_error", "deadline_exceeded"), + Some(1) + ); + // A late answer still reaches the caller without a second count. + claim_tx + .send(Err(Status::unavailable("supervisor gone"))) + .unwrap(); + let answer = forwarded.await.unwrap(); + assert_eq!(answer.unwrap_err().code(), tonic::Code::Unavailable); + assert_eq!( + local_relay_outcome(&metrics, "local_error", "unavailable"), + None + ); + } + + #[tokio::test(start_paused = true)] + async fn local_relay_caller_timeout_counts_as_unclaimed_not_cancelled() { + let metrics = MetricsCapture::install(); + let (claim_tx, forwarded) = tracked_local_relay(); + // Callers wait with the same 10 s budget and drop the receiver when it runs out. + assert!( + tokio::time::timeout(RELAY_PENDING_TIMEOUT, forwarded) + .await + .is_err() + ); + while !claim_tx.is_closed() { + tokio::task::yield_now().await; + } + assert_eq!( + local_relay_outcome(&metrics, "local_error", "deadline_exceeded"), + Some(1) + ); + assert_eq!( + local_relay_outcome(&metrics, "local_error", "cancelled"), + None + ); + } + + #[tokio::test] + async fn local_relay_failed_by_registry_records_its_status() { + let metrics = MetricsCapture::install(); + let (claim_tx, forwarded) = tracked_local_relay(); + claim_tx + .send(Err(Status::unavailable("supervisor session disconnected"))) + .unwrap(); + assert_eq!( + forwarded.await.unwrap().unwrap_err().code(), + tonic::Code::Unavailable + ); + assert_eq!( + local_relay_outcome(&metrics, "local_error", "unavailable"), + Some(1) + ); + } + + #[tokio::test] + async fn local_relay_dropped_by_registry_closes_the_caller_channel() { + let metrics = MetricsCapture::install(); + let (claim_tx, forwarded) = tracked_local_relay(); + drop(claim_tx); + assert!( + forwarded.await.is_err(), + "the caller sees the same closed channel" + ); + assert_eq!( + local_relay_outcome(&metrics, "local_error", "unavailable"), + Some(1) + ); + } + + #[tokio::test] + async fn local_relay_abandoned_by_caller_records_cancelled() { + let metrics = MetricsCapture::install(); + let (claim_tx, forwarded) = tracked_local_relay(); + drop(forwarded); + while !claim_tx.is_closed() { + tokio::task::yield_now().await; + } + assert_eq!( + local_relay_outcome(&metrics, "local_error", "cancelled"), + Some(1) + ); + } + + #[tokio::test] + async fn peer_relay_metrics_keep_owner_code_before_unavailable_remap() { + let metrics = MetricsCapture::install(); + let state = crate::grpc::test_support::test_server_state().await; + seed_peer_token(&state); + let endpoint = spawn_fake_peer(FakePeerReply::Status(tonic::Code::ResourceExhausted)).await; + + let err = connect_peer_relay(&state, &endpoint, "sbx-peer", peer_relay_open("ch-peer")) + .await + .expect_err("the owner rejected the relay"); + assert_eq!(err.code(), tonic::Code::Unavailable); + assert_eq!( + metrics.value( + "openshell_server_routed_request_attempts_total{operation=\"relay\",route=\"peer\",relay_kind=\"ssh\",outcome=\"remote_error\",grpc_code=\"resource_exhausted\"}" + ), + Some(1) + ); + assert_eq!( + metrics.value( + "openshell_server_peer_request_duration_seconds_count{operation=\"relay\",outcome=\"remote_error\"}" + ), + Some(1) + ); + let rendered = metrics.render(); + assert!(rendered.contains( + "openshell_server_peer_request_duration_seconds_bucket{operation=\"relay\",outcome=\"remote_error\",le=\"0.001\"}" + )); + assert!( + !state + .peer_routes + .channels + .lock() + .unwrap() + .contains_key(&endpoint), + "a failed peer relay must evict the channel" + ); + let host_port = endpoint.trim_start_matches("http://"); + assert!( + !rendered.contains(host_port), + "metrics must not carry peer endpoints" + ); + } + + #[tokio::test] + async fn peer_relay_metrics_record_ok_when_owner_accepts() { + let metrics = MetricsCapture::install(); + let state = crate::grpc::test_support::test_server_state().await; + seed_peer_token(&state); + let endpoint = spawn_fake_peer(FakePeerReply::EmptyOk).await; + + connect_peer_relay(&state, &endpoint, "sbx-peer", peer_relay_open("ch-peer")) + .await + .expect("the owner accepted the relay"); + assert_eq!( + metrics.value( + "openshell_server_routed_request_attempts_total{operation=\"relay\",route=\"peer\",relay_kind=\"ssh\",outcome=\"success\",grpc_code=\"ok\"}" + ), + Some(1) + ); + } + + #[tokio::test] + async fn peer_forward_metrics_record_local_error_when_owner_unreachable() { + let metrics = MetricsCapture::install(); + let state = crate::grpc::test_support::test_server_state().await; + seed_peer_token(&state); + let endpoint = closed_local_endpoint(); + + let err = forward_provider_status_query_to_owner( + &state, + &owner_at(&endpoint), + "sbx-peer", + GetSandboxProviderStatusRequest::default(), + ) + .await + .expect_err("the owner is unreachable"); + assert_eq!(err.code(), tonic::Code::Unavailable); + assert_eq!( + metrics.value( + "openshell_server_routed_request_attempts_total{operation=\"get_sandbox_provider_status\",route=\"peer\",relay_kind=\"none\",outcome=\"local_error\",grpc_code=\"unavailable\"}" + ), + Some(1) + ); + assert!( + !metrics + .render() + .contains("operation=\"get_sandbox_provider_status\",outcome=\"remote_error\"") + ); + } + + #[tokio::test] + async fn peer_forward_metrics_record_owner_remote_error_code() { + let metrics = MetricsCapture::install(); + let state = crate::grpc::test_support::test_server_state().await; + seed_peer_token(&state); + let endpoint = spawn_fake_peer(FakePeerReply::Status(tonic::Code::PermissionDenied)).await; + + let err = forward_endpoint_status_to_owner( + &state, + &owner_at(&endpoint), + ReportEndpointStatusRequest { + sandbox_id: "sbx-peer".into(), + ..Default::default() + }, + ) + .await + .expect_err("the owner rejected the report"); + // Unary forwarders return the owner's status unchanged. + assert_eq!(err.code(), tonic::Code::PermissionDenied); + assert_eq!( + metrics.value( + "openshell_server_routed_request_attempts_total{operation=\"report_endpoint_status\",route=\"peer\",relay_kind=\"none\",outcome=\"remote_error\",grpc_code=\"permission_denied\"}" + ), + Some(1) + ); + } + + #[tokio::test] + async fn peer_forward_metrics_record_ok() { + let metrics = MetricsCapture::install(); + let state = crate::grpc::test_support::test_server_state().await; + seed_peer_token(&state); + let endpoint = spawn_fake_peer(FakePeerReply::EmptyOk).await; + + forward_provider_readiness_to_owner( + &state, + &owner_at(&endpoint), + ReportProviderReadinessRequest { + sandbox_id: "sbx-peer".into(), + ..Default::default() + }, + ) + .await + .expect("the owner accepted the report"); + assert_eq!( + metrics.value( + "openshell_server_routed_request_attempts_total{operation=\"report_provider_readiness\",route=\"peer\",relay_kind=\"none\",outcome=\"success\",grpc_code=\"ok\"}" + ), + Some(1) + ); + } + #[tokio::test] async fn evict_channel_removes_only_the_named_peer() { let cache = PeerRouteCache::default(); diff --git a/crates/openshell-server/tests/supervisor_relay_integration.rs b/crates/openshell-server/tests/supervisor_relay_integration.rs index 5856acbbda..de0cfee394 100644 --- a/crates/openshell-server/tests/supervisor_relay_integration.rs +++ b/crates/openshell-server/tests/supervisor_relay_integration.rs @@ -22,6 +22,8 @@ use hyper_util::{ rt::{TokioExecutor, TokioIo}, server::conn::auto::Builder, }; +use metrics::LocalRecorderGuard; +use metrics_exporter_prometheus::{PrometheusBuilder, PrometheusHandle, PrometheusRecorder}; use openshell_core::proto::{ GatewayMessage, PeerRelayFrame, RelayFrame, RelayInit, SupervisorMessage, TcpForwardFrame, open_shell_client::OpenShellClient, @@ -658,6 +660,38 @@ fn register_session_with_capacity( rx } +/// Captures metrics recorded on this test's thread. `#[tokio::test]` is current-thread, so the +/// in-process gateway tasks record here too. Do not pass it into async helper fns (it is !Send). +struct MetricsCapture { + handle: PrometheusHandle, + _guard: LocalRecorderGuard<'static>, +} + +impl MetricsCapture { + fn install() -> Self { + // Leaked (test only) so the guard can borrow the recorder for 'static. + let recorder: &'static PrometheusRecorder = + Box::leak(Box::new(PrometheusBuilder::new().build_recorder())); + let handle = recorder.handle(); + let guard = metrics::set_default_local_recorder(recorder); + Self { + handle, + _guard: guard, + } + } + + fn render(&self) -> String { + self.handle.render() + } + + /// Integer value of one exact series, or `None` when it was never emitted. + fn value(&self, series: &str) -> Option { + self.render() + .lines() + .find_map(|line| line.strip_prefix(series)?.strip_prefix(' ')?.parse().ok()) + } +} + /// Mock supervisor that opens a `RelayStream`, sends `Init`, then echoes every /// data frame it receives. Returns when the gateway drops the stream or when /// the supervisor's own outbound channel closes. @@ -859,6 +893,7 @@ async fn concurrent_relays_multiplex_independently() { /// rather than racing the pending map into an inconsistent state. #[tokio::test] async fn open_relay_enforces_per_sandbox_cap_under_concurrent_burst() { + let metrics = MetricsCapture::install(); let registry = Arc::new(SupervisorSessionRegistry::new()); let _channel = spawn_gateway(Arc::clone(®istry)).await; // Oversized mpsc so the session doesn't backpressure the burst — the cap, @@ -894,6 +929,23 @@ async fn open_relay_enforces_per_sandbox_cap_under_concurrent_burst() { } assert_eq!(ok, 32, "exactly per-sandbox cap should succeed"); assert_eq!(exhausted, 32, "overflow should be rejected, not dropped"); + // The successful receivers were dropped, so their entries stay pending + // until they are claimed or reaped. + assert_eq!(metrics.value("openshell_server_relay_pending"), Some(32)); + assert_eq!( + metrics.value("openshell_server_relay_rejected_total{reason=\"sandbox_capacity\"}"), + Some(32) + ); + assert_eq!( + metrics + .value("openshell_server_relay_rejected_total{reason=\"replica_capacity\"}") + .unwrap_or(0), + 0 + ); + assert_eq!( + metrics.value("openshell_server_supervisor_sessions"), + Some(1) + ); // A different sandbox still has headroom — the per-sandbox cap doesn't // leak onto unrelated tenants. @@ -902,6 +954,107 @@ async fn open_relay_enforces_per_sandbox_cap_under_concurrent_burst() { .open_relay("sbx-other", Duration::from_secs(1)) .await .expect("other sandbox should not be affected by sbx cap"); + assert_eq!(metrics.value("openshell_server_relay_pending"), Some(33)); + assert_eq!( + metrics.value("openshell_server_supervisor_sessions"), + Some(2) + ); +} + +/// Bursts more `open_relay` calls than the replica cap allows, spread so that +/// no sandbox reaches its own cap, and asserts the replica ceiling and its +/// metrics. +#[tokio::test] +async fn open_relay_enforces_replica_cap_under_concurrent_burst() { + let metrics = MetricsCapture::install(); + let registry = Arc::new(SupervisorSessionRegistry::new()); + let _channel = spawn_gateway(Arc::clone(®istry)).await; + let sandbox_ids: Vec = (0..9).map(|i| format!("sbx-{i}")).collect(); + let _session_rxs: Vec<_> = sandbox_ids + .iter() + .map(|id| register_session_with_capacity(®istry, id, 64)) + .collect(); + + // 9 x 32 = 288 opens against a replica cap of 256. No sandbox gets more + // than its cap of 32 attempts and the replica check runs first, so exactly + // 32 opens hit the replica cap. + let mut handles = Vec::with_capacity(288); + for id in &sandbox_ids { + for _ in 0..32 { + let r = Arc::clone(®istry); + let id = id.clone(); + handles.push(tokio::spawn(async move { + r.open_relay(&id, Duration::from_secs(1)).await + })); + } + } + + let mut ok = 0usize; + let mut exhausted = 0usize; + for h in handles { + match h.await.expect("task joined") { + Ok(_pair) => ok += 1, + Err(status) if status.code() == tonic::Code::ResourceExhausted => { + assert!( + status.message().contains("gateway relay capacity"), + "expected replica capacity error message, got: {}", + status.message() + ); + exhausted += 1; + } + Err(other) => panic!("unexpected open_relay error: {other:?}"), + } + } + assert_eq!(ok, 256, "exactly the replica cap should succeed"); + assert_eq!(exhausted, 32, "overflow should be rejected, not dropped"); + + assert_eq!(metrics.value("openshell_server_relay_pending"), Some(256)); + assert_eq!( + metrics.value("openshell_server_relay_rejected_total{reason=\"replica_capacity\"}"), + Some(32) + ); + assert_eq!( + metrics + .value("openshell_server_relay_rejected_total{reason=\"sandbox_capacity\"}") + .unwrap_or(0), + 0 + ); + assert_eq!( + metrics.value("openshell_server_supervisor_sessions"), + Some(9) + ); + let rendered = metrics.render(); + for forbidden in ["sandbox_id=", "channel_id=", "sbx-0", "endpoint="] { + assert!( + !rendered.contains(forbidden), + "metric labels must not carry identifiers ({forbidden})" + ); + } +} + +#[tokio::test] +async fn relay_claim_releases_pending_slot_and_records_claim_latency() { + let metrics = MetricsCapture::install(); + let registry = Arc::new(SupervisorSessionRegistry::new()); + let channel = spawn_gateway(Arc::clone(®istry)).await; + let _session_rx = register_session(®istry, "sbx"); + + let (channel_id, relay_rx) = registry + .open_relay("sbx", Duration::from_secs(2)) + .await + .expect("open_relay"); + assert_eq!(metrics.value("openshell_server_relay_pending"), Some(1)); + + tokio::spawn(run_echo_supervisor(channel, channel_id)); + let _relay = relay_rx.await.expect("relay result").expect("relay duplex"); + + // The claim releases the pending slot under the pending lock, before it + // wakes the waiter. + assert_eq!(metrics.value("openshell_server_relay_pending"), Some(0)); + assert_eq!( + metrics.value("openshell_server_relay_claim_duration_seconds_count"), + Some(1) + ); } /// Build an in-memory store sufficient for wiring `health_router` in tests diff --git a/deploy/helm/openshell/README.md b/deploy/helm/openshell/README.md index df3ea173ad..f617faded7 100644 --- a/deploy/helm/openshell/README.md +++ b/deploy/helm/openshell/README.md @@ -189,6 +189,7 @@ See [`values.yaml`](values.yaml) for source defaults. Selected overlays: - [`ci/values-cert-manager.yaml`](ci/values-cert-manager.yaml) - cert-manager integration - [`ci/values-keycloak.yaml`](ci/values-keycloak.yaml) - Keycloak OIDC integration - [`ci/values-high-availability.yaml`](ci/values-high-availability.yaml) - CI overlay for multi-replica external PostgreSQL testing +- [`ci/values-autoscaling.yaml`](ci/values-autoscaling.yaml) - CI overlay for rendering the optional gateway HorizontalPodAutoscaler - [`ci/values-spire.yaml`](ci/values-spire.yaml) - SPIFFE/SPIRE provider token grants - [`ci/values-spire-stack.yaml`](ci/values-spire-stack.yaml) - SPIRE hardened chart values for local development @@ -275,6 +276,32 @@ DNS name while connecting directly to the owning pod. Custom TLS Secrets must include that Service DNS name in the server certificate and provide the CA and client credentials configured by `server.tls`. +Set `autoscaling.enabled=true` to render an `autoscaling/v2` +HorizontalPodAutoscaler for the gateway workload. The Deployment or +StatefulSet then omits `spec.replicas`, and the chart applies its +multi-replica checks to `autoscaling.maxReplicas`. CPU and memory targets +need a matching `resources.requests` entry, or a `resources.limits` entry, +which Kubernetes copies into the request. Add custom metrics from a metrics +adapter with `autoscaling.metrics`. See the +[High Availability guide](https://docs.nvidia.com/openshell/latest/kubernetes/high-availability) +for the metrics to scale and alert on. + +Enabling autoscaling on an existing release removes `spec.replicas` from the +workload in that upgrade. Kubernetes can reset the workload to one replica +until the HPA scales it back to at least `autoscaling.minReplicas`, which +disconnects supervisor sessions from the other gateway pods; they reconnect to the remaining replica. Enable autoscaling +when you install the chart, or during a maintenance window. On an existing +release, upgrade with `--reset-then-reuse-values`, which needs Helm 3.14 or +later. With an older Helm, save the values with +`helm get values -o yaml > values.yaml` and upgrade with +`--reset-values -f values.yaml`. `--reuse-values` keeps the +previous chart version's defaults, which lack the autoscaling values, +including the scale-down `behavior`, when that release predates them. For background, refer to +[Migrating Deployments and StatefulSets to horizontal autoscaling](https://kubernetes.io/docs/tasks/run-application/horizontal-pod-autoscale/#migrating-deployments-and-statefulsets-to-horizontal-autoscaling). +Disabling autoscaling renders `spec.replicas` from `replicaCount` again, +which defaults to 1, so set `replicaCount` to the replica count you want in +that upgrade. + ## Secret bootstrap By default, a pre-install/pre-upgrade hook Job runs `openshell-gateway generate-certs` @@ -311,6 +338,13 @@ discovery endpoint or its TLS CA. |-----|------|---------|-------------| | affinity | object | `{}` | Affinity rules for the gateway pod. | | agentSandbox.preflight.enabled | bool | `true` | Check the live cluster for a supported Agent Sandbox API before rendering gateway resources. Disable only for offline rendering and linting. | +| autoscaling.behavior | object | `{"scaleDown":{"policies":[{"periodSeconds":120,"type":"Pods","value":1}],"stabilizationWindowSeconds":300}}` | HPA scaling behavior. Scale-down disconnects the removed pod's supervisor sessions; they reconnect to the remaining replicas, so the default removes at most one replica every two minutes after a five-minute stabilization window. Helm merges maps: set autoscaling.behavior.scaleDown to null to drop the default. | +| autoscaling.enabled | bool | `false` | Render a HorizontalPodAutoscaler and stop rendering spec.replicas. | +| autoscaling.maxReplicas | int | `4` | Maximum gateway replicas. Each replica opens its own PostgreSQL connection pool; size the database for rollouts at this count, as the High Availability guide describes. | +| autoscaling.metrics | list | `[]` | Additional autoscaling/v2 MetricSpec entries appended verbatim, such as Pods metrics served by prometheus-adapter. | +| autoscaling.minReplicas | int | `2` | Minimum gateway replicas. Use 2 or more to survive a pod failure. | +| autoscaling.targetCPUUtilizationPercentage | int | `80` | Target average CPU utilization, as a percentage of resources.requests.cpu. Set to null to disable. Requires resources.requests.cpu, or resources.limits.cpu, which Kubernetes copies into the request. | +| autoscaling.targetMemoryUtilizationPercentage | int | `nil` | Target average memory utilization, as a percentage of resources.requests.memory. Null disables it. Requires resources.requests.memory, or resources.limits.memory, which Kubernetes copies into the request. | | certManager.caSecretName | string | `"openshell-ca-tls"` | Secret created for the intermediate CA (Certificate with isCA: true). | | certManager.certificateDuration | string | `"8760h"` | Duration for cert-manager-issued certificates. | | certManager.certificateRenewBefore | string | `"720h"` | Renewal window for cert-manager-issued certificates. | diff --git a/deploy/helm/openshell/README.md.gotmpl b/deploy/helm/openshell/README.md.gotmpl index 1e86e0cbcf..59d6fb06b6 100644 --- a/deploy/helm/openshell/README.md.gotmpl +++ b/deploy/helm/openshell/README.md.gotmpl @@ -190,6 +190,7 @@ See [`values.yaml`](values.yaml) for source defaults. Selected overlays: - [`ci/values-cert-manager.yaml`](ci/values-cert-manager.yaml) - cert-manager integration - [`ci/values-keycloak.yaml`](ci/values-keycloak.yaml) - Keycloak OIDC integration - [`ci/values-high-availability.yaml`](ci/values-high-availability.yaml) - CI overlay for multi-replica external PostgreSQL testing +- [`ci/values-autoscaling.yaml`](ci/values-autoscaling.yaml) - CI overlay for rendering the optional gateway HorizontalPodAutoscaler - [`ci/values-spire.yaml`](ci/values-spire.yaml) - SPIFFE/SPIRE provider token grants - [`ci/values-spire-stack.yaml`](ci/values-spire-stack.yaml) - SPIRE hardened chart values for local development @@ -276,6 +277,32 @@ DNS name while connecting directly to the owning pod. Custom TLS Secrets must include that Service DNS name in the server certificate and provide the CA and client credentials configured by `server.tls`. +Set `autoscaling.enabled=true` to render an `autoscaling/v2` +HorizontalPodAutoscaler for the gateway workload. The Deployment or +StatefulSet then omits `spec.replicas`, and the chart applies its +multi-replica checks to `autoscaling.maxReplicas`. CPU and memory targets +need a matching `resources.requests` entry, or a `resources.limits` entry, +which Kubernetes copies into the request. Add custom metrics from a metrics +adapter with `autoscaling.metrics`. See the +[High Availability guide](https://docs.nvidia.com/openshell/latest/kubernetes/high-availability) +for the metrics to scale and alert on. + +Enabling autoscaling on an existing release removes `spec.replicas` from the +workload in that upgrade. Kubernetes can reset the workload to one replica +until the HPA scales it back to at least `autoscaling.minReplicas`, which +disconnects supervisor sessions from the other gateway pods; they reconnect to the remaining replica. Enable autoscaling +when you install the chart, or during a maintenance window. On an existing +release, upgrade with `--reset-then-reuse-values`, which needs Helm 3.14 or +later. With an older Helm, save the values with +`helm get values -o yaml > values.yaml` and upgrade with +`--reset-values -f values.yaml`. `--reuse-values` keeps the +previous chart version's defaults, which lack the autoscaling values, +including the scale-down `behavior`, when that release predates them. For background, refer to +[Migrating Deployments and StatefulSets to horizontal autoscaling](https://kubernetes.io/docs/tasks/run-application/horizontal-pod-autoscale/#migrating-deployments-and-statefulsets-to-horizontal-autoscaling). +Disabling autoscaling renders `spec.replicas` from `replicaCount` again, +which defaults to 1, so set `replicaCount` to the replica count you want in +that upgrade. + ## Secret bootstrap By default, a pre-install/pre-upgrade hook Job runs `openshell-gateway generate-certs` diff --git a/deploy/helm/openshell/ci/values-autoscaling.yaml b/deploy/helm/openshell/ci/values-autoscaling.yaml new file mode 100644 index 0000000000..0de3494ebc --- /dev/null +++ b/deploy/helm/openshell/ci/values-autoscaling.yaml @@ -0,0 +1,35 @@ +# SPDX-FileCopyrightText: Copyright (c) 2025-2026 NVIDIA CORPORATION & AFFILIATES. All rights reserved. +# SPDX-License-Identifier: Apache-2.0 + +# CI overlay for rendering the optional gateway HorizontalPodAutoscaler. Like +# values-high-availability.yaml, it expects a PostgreSQL Secret named +# openshell-ha-pg. It is exercised by helm lint and helm-unittest only; kind +# CI has no metrics-server, so HPA runtime behavior is not tested. +workload: + kind: deployment + +server: + externalDbSecret: openshell-ha-pg + +resources: + requests: + cpu: 250m + memory: 256Mi + +autoscaling: + enabled: true + minReplicas: 2 + maxReplicas: 4 + targetCPUUtilizationPercentage: 80 + # Render-only placeholder: the gateway exports no metric by this name. It + # checks that the chart passes a Pods metric to the HPA unchanged. An HPA + # skips scale-down while any of its metrics is unavailable, so drop it for a + # runtime test with metrics-server only: --set autoscaling.metrics=null. + metrics: + - type: Pods + pods: + metric: + name: example_adapter_metric + target: + type: AverageValue + averageValue: "100" diff --git a/deploy/helm/openshell/templates/_autoscaling.tpl b/deploy/helm/openshell/templates/_autoscaling.tpl new file mode 100644 index 0000000000..cceac448c7 --- /dev/null +++ b/deploy/helm/openshell/templates/_autoscaling.tpl @@ -0,0 +1,79 @@ +# SPDX-FileCopyrightText: Copyright (c) 2025-2026 NVIDIA CORPORATION & AFFILIATES. All rights reserved. +# SPDX-License-Identifier: Apache-2.0 + +{{/* +Whether the chart renders a HorizontalPodAutoscaler for the gateway workload. +A missing autoscaling map (for example `helm upgrade --reuse-values` from a +release that predates these values) means disabled. +*/}} +{{- define "openshell.autoscalingEnabled" -}} +{{- $autoscaling := .Values.autoscaling | default dict -}} +{{- if and (kindIs "map" $autoscaling) (get $autoscaling "enabled") -}}true{{- end -}} +{{- end }} + +{{/* +Largest replica count the chart can run: autoscaling.maxReplicas when the +HPA is enabled, otherwise replicaCount. +*/}} +{{- define "openshell.maxReplicas" -}} +{{- if eq (include "openshell.autoscalingEnabled" .) "true" -}} +{{- int (get .Values.autoscaling "maxReplicas" | default 1) -}} +{{- else -}} +{{- int (default 1 .Values.replicaCount) -}} +{{- end -}} +{{- end }} + +{{/* +Name of the value that sets the largest replica count, for error messages. +*/}} +{{- define "openshell.maxReplicasSource" -}} +{{- ternary "autoscaling.maxReplicas" "replicaCount" (eq (include "openshell.autoscalingEnabled" .) "true") -}} +{{- end }} + +{{/* +Validate autoscaling values. Called from openshell.validateValues. +*/}} +{{- define "openshell.validateAutoscaling" -}} +{{- if eq (include "openshell.autoscalingEnabled" .) "true" -}} +{{- $a := .Values.autoscaling -}} +{{- if not (and (hasKey $a "minReplicas") (hasKey $a "maxReplicas")) -}} +{{- fail "autoscaling.minReplicas and autoscaling.maxReplicas are not set. helm upgrade --reuse-values keeps the values of a release that predates the chart's autoscaling defaults; upgrade with --reset-then-reuse-values, or set the autoscaling values explicitly, including autoscaling.behavior." -}} +{{- end -}} +{{- $min := int (get $a "minReplicas" | default 0) -}} +{{- $max := int (get $a "maxReplicas" | default 0) -}} +{{- if lt $min 1 -}} +{{- fail "autoscaling.minReplicas must be at least 1." -}} +{{- end -}} +{{- if lt $max $min -}} +{{- fail "autoscaling.maxReplicas must be greater than or equal to autoscaling.minReplicas." -}} +{{- end -}} +{{- $cpu := get $a "targetCPUUtilizationPercentage" -}} +{{- $memory := get $a "targetMemoryUtilizationPercentage" -}} +{{- if not (or $cpu $memory (get $a "metrics")) -}} +{{- fail "autoscaling.enabled requires targetCPUUtilizationPercentage, targetMemoryUtilizationPercentage, or autoscaling.metrics." -}} +{{- end -}} +{{- if and $cpu (ne (include "openshell.hasPositiveRequest" (dict "resources" .Values.resources "name" "cpu")) "true") -}} +{{- fail "autoscaling.targetCPUUtilizationPercentage requires a positive resources.requests.cpu (or resources.limits.cpu, which Kubernetes copies into an unset request); Kubernetes computes utilization against the container request." -}} +{{- end -}} +{{- if and $memory (ne (include "openshell.hasPositiveRequest" (dict "resources" .Values.resources "name" "memory")) "true") -}} +{{- fail "autoscaling.targetMemoryUtilizationPercentage requires a positive resources.requests.memory (or resources.limits.memory, which Kubernetes copies into an unset request); Kubernetes computes utilization against the container request." -}} +{{- end -}} +{{- end -}} +{{- end }} + +{{/* +"true" when the container request Kubernetes uses for one resource is positive. +That request is resources.requests. when the key is present, even as +null, which the chart renders and Kubernetes stores as zero. Otherwise it is +resources.limits., which Kubernetes copies only into an absent request. +Helm cannot parse quantities, so the number before any suffix or exponent must +be unsigned or "+" and contain a nonzero digit: 0, 0m, 0Mi, 0.0 and 0e3 are +zero. +*/}} +{{- define "openshell.hasPositiveRequest" -}} +{{- $resources := .resources | default dict -}} +{{- $requests := $resources.requests | default dict -}} +{{- $limits := $resources.limits | default dict -}} +{{- $value := ternary (get $requests .name) (get $limits .name) (hasKey $requests .name) -}} +{{- if and (not (kindIs "invalid" $value)) (regexMatch "^[+]?[0-9.]*[1-9]" (toString $value)) -}}true{{- end -}} +{{- end }} diff --git a/deploy/helm/openshell/templates/_helpers.tpl b/deploy/helm/openshell/templates/_helpers.tpl index ab42458759..6b84d847f0 100644 --- a/deploy/helm/openshell/templates/_helpers.tpl +++ b/deploy/helm/openshell/templates/_helpers.tpl @@ -51,6 +51,18 @@ app.kubernetes.io/name: {{ include "openshell.name" . }} app.kubernetes.io/instance: {{ .Release.Name }} {{- end }} +{{/* +Pod labels for the certgen hook Jobs. They keep the release instance label but +do not match openshell.selectorLabels, so gateway selectors (the workload, +Services, HorizontalPodAutoscaler, anti-affinity, and PodDisruptionBudgets) +never select hook pods. +*/}} +{{- define "openshell.certgenPodLabels" -}} +app.kubernetes.io/name: {{ printf "%s-certgen" (include "openshell.name" . | trunc 55 | trimSuffix "-") }} +app.kubernetes.io/instance: {{ .Release.Name }} +app.kubernetes.io/component: certgen +{{- end }} + {{/* Create the name of the service account to use */}} @@ -385,7 +397,8 @@ Validate chart values that Helm would otherwise accept silently. {{- define "openshell.validateValues" -}} {{- $workloadKind := include "openshell.workloadKind" . -}} {{- $workload := .Values.workload | default dict -}} -{{- $replicaCount := int (default 1 .Values.replicaCount) -}} +{{- $maxReplicas := int (include "openshell.maxReplicas" .) -}} +{{- $maxReplicasSource := include "openshell.maxReplicasSource" . -}} {{- if and (hasKey .Values "postgres") (kindIs "map" .Values.postgres) (hasKey .Values.postgres "enabled") -}} {{- fail "postgres.enabled was removed; the OpenShell chart no longer deploys PostgreSQL. Provision PostgreSQL separately and set server.externalDbSecret to a Secret containing a PostgreSQL URI." -}} {{- end -}} @@ -395,11 +408,12 @@ Validate chart values that Helm would otherwise accept silently. {{- if and (eq $workloadKind "deployment") (not .Values.server.externalDbSecret) -}} {{- fail "workload.kind=deployment requires server.externalDbSecret; use workload.kind=statefulset for the default SQLite database." -}} {{- end -}} -{{- if and (gt $replicaCount 1) (not .Values.server.externalDbSecret) -}} -{{- fail "replicaCount > 1 requires server.externalDbSecret; multiple gateway replicas cannot share the default per-pod SQLite database." -}} +{{- include "openshell.validateAutoscaling" . -}} +{{- if and (gt $maxReplicas 1) (not .Values.server.externalDbSecret) -}} +{{- fail (printf "%s > 1 requires server.externalDbSecret; multiple gateway replicas cannot share the default per-pod SQLite database." $maxReplicasSource) -}} {{- end -}} -{{- if and (eq $workloadKind "statefulset") (gt $replicaCount 1) (not (get $workload "allowMultiReplicaStatefulSet" | default false)) -}} -{{- fail "replicaCount > 1 with workload.kind=statefulset requires workload.allowMultiReplicaStatefulSet=true; use workload.kind=deployment for external database-backed multi-replica gateways." -}} +{{- if and (eq $workloadKind "statefulset") (gt $maxReplicas 1) (not (get $workload "allowMultiReplicaStatefulSet" | default false)) -}} +{{- fail (printf "%s > 1 with workload.kind=statefulset requires workload.allowMultiReplicaStatefulSet=true; use workload.kind=deployment for external database-backed multi-replica gateways." $maxReplicasSource) -}} {{- end -}} {{- $workspaceMode := .Values.server.drivers.kubernetes.workspaceMode | default "shared" -}} {{- if not (has $workspaceMode (list "shared" "managed" "operator")) -}} diff --git a/deploy/helm/openshell/templates/certgen.yaml b/deploy/helm/openshell/templates/certgen.yaml index f7c9a751d3..d5610faf2f 100644 --- a/deploy/helm/openshell/templates/certgen.yaml +++ b/deploy/helm/openshell/templates/certgen.yaml @@ -75,7 +75,7 @@ spec: template: metadata: labels: - {{- include "openshell.selectorLabels" . | nindent 8 }} + {{- include "openshell.certgenPodLabels" . | nindent 8 }} spec: restartPolicy: OnFailure serviceAccountName: {{ $hookName }} @@ -146,7 +146,7 @@ spec: template: metadata: labels: - {{- include "openshell.selectorLabels" . | nindent 8 }} + {{- include "openshell.certgenPodLabels" . | nindent 8 }} spec: restartPolicy: OnFailure serviceAccountName: {{ $hookName }} diff --git a/deploy/helm/openshell/templates/deployment.yaml b/deploy/helm/openshell/templates/deployment.yaml index f94900b136..1c9f589bcf 100644 --- a/deploy/helm/openshell/templates/deployment.yaml +++ b/deploy/helm/openshell/templates/deployment.yaml @@ -9,7 +9,9 @@ metadata: labels: {{- include "openshell.labels" . | nindent 4 }} spec: + {{- if ne (include "openshell.autoscalingEnabled" .) "true" }} replicas: {{ .Values.replicaCount }} + {{- end }} selector: matchLabels: {{- include "openshell.selectorLabels" . | nindent 6 }} diff --git a/deploy/helm/openshell/templates/hpa.yaml b/deploy/helm/openshell/templates/hpa.yaml new file mode 100644 index 0000000000..d88c46290b --- /dev/null +++ b/deploy/helm/openshell/templates/hpa.yaml @@ -0,0 +1,42 @@ +# SPDX-FileCopyrightText: Copyright (c) 2025-2026 NVIDIA CORPORATION & AFFILIATES. All rights reserved. +# SPDX-License-Identifier: Apache-2.0 +{{- include "openshell.validateValues" . }} +{{- if eq (include "openshell.autoscalingEnabled" .) "true" }} +apiVersion: autoscaling/v2 +kind: HorizontalPodAutoscaler +metadata: + name: {{ include "openshell.fullname" . }} + labels: + {{- include "openshell.labels" . | nindent 4 }} +spec: + scaleTargetRef: + apiVersion: apps/v1 + kind: {{ ternary "Deployment" "StatefulSet" (eq (include "openshell.workloadKind" .) "deployment") }} + name: {{ include "openshell.fullname" . }} + minReplicas: {{ .Values.autoscaling.minReplicas }} + maxReplicas: {{ .Values.autoscaling.maxReplicas }} + metrics: + {{- with .Values.autoscaling.targetCPUUtilizationPercentage }} + - type: Resource + resource: + name: cpu + target: + type: Utilization + averageUtilization: {{ . }} + {{- end }} + {{- with .Values.autoscaling.targetMemoryUtilizationPercentage }} + - type: Resource + resource: + name: memory + target: + type: Utilization + averageUtilization: {{ . }} + {{- end }} + {{- with .Values.autoscaling.metrics }} + {{- toYaml . | nindent 4 }} + {{- end }} + {{- with .Values.autoscaling.behavior }} + behavior: + {{- toYaml . | nindent 4 }} + {{- end }} +{{- end }} diff --git a/deploy/helm/openshell/templates/statefulset.yaml b/deploy/helm/openshell/templates/statefulset.yaml index 10d0839f60..be65064c19 100644 --- a/deploy/helm/openshell/templates/statefulset.yaml +++ b/deploy/helm/openshell/templates/statefulset.yaml @@ -10,7 +10,9 @@ metadata: {{- include "openshell.labels" . | nindent 4 }} spec: serviceName: {{ include "openshell.peerServiceName" . }} + {{- if ne (include "openshell.autoscalingEnabled" .) "true" }} replicas: {{ .Values.replicaCount }} + {{- end }} selector: matchLabels: {{- include "openshell.selectorLabels" . | nindent 6 }} diff --git a/deploy/helm/openshell/tests/autoscaling_test.yaml b/deploy/helm/openshell/tests/autoscaling_test.yaml new file mode 100644 index 0000000000..ee3a2b7d79 --- /dev/null +++ b/deploy/helm/openshell/tests/autoscaling_test.yaml @@ -0,0 +1,386 @@ +# SPDX-FileCopyrightText: Copyright (c) 2025-2026 NVIDIA CORPORATION & AFFILIATES. All rights reserved. +# SPDX-License-Identifier: Apache-2.0 + +suite: gateway autoscaling +templates: + - templates/hpa.yaml + - templates/deployment.yaml + - templates/statefulset.yaml + - templates/gateway-config.yaml +release: + name: openshell + namespace: my-namespace + +tests: + - it: renders no HorizontalPodAutoscaler by default + template: templates/hpa.yaml + asserts: + - hasDocuments: + count: 0 + + - it: keeps spec.replicas on the default StatefulSet + template: templates/statefulset.yaml + asserts: + - equal: + path: spec.replicas + value: 1 + + - it: renders an autoscaling/v2 HPA targeting the Deployment from the CI overlay + template: templates/hpa.yaml + values: + - ../ci/values-autoscaling.yaml + asserts: + - isKind: + of: HorizontalPodAutoscaler + - isAPIVersion: + of: autoscaling/v2 + - equal: + path: metadata.name + value: openshell + - equal: + path: spec.scaleTargetRef + value: + apiVersion: apps/v1 + kind: Deployment + name: openshell + - equal: + path: spec.minReplicas + value: 2 + - equal: + path: spec.maxReplicas + value: 4 + - contains: + path: spec.metrics + content: + type: Resource + resource: + name: cpu + target: + type: Utilization + averageUtilization: 80 + - contains: + path: spec.metrics + content: + type: Pods + pods: + metric: + name: example_adapter_metric + target: + type: AverageValue + averageValue: "100" + - equal: + path: spec.behavior.scaleDown.stabilizationWindowSeconds + value: 300 + - equal: + path: spec.behavior.scaleDown.policies[0] + value: + type: Pods + value: 1 + periodSeconds: 120 + + - it: omits spec.replicas from the Deployment when autoscaling is enabled + template: templates/deployment.yaml + values: + - ../ci/values-autoscaling.yaml + set: + replicaCount: 3 + asserts: + - isKind: + of: Deployment + - notExists: + path: spec.replicas + + - it: renders memory utilization only when configured + template: templates/hpa.yaml + values: + - ../ci/values-autoscaling.yaml + set: + autoscaling.targetMemoryUtilizationPercentage: 70 + asserts: + - contains: + path: spec.metrics + content: + type: Resource + resource: + name: memory + target: + type: Utilization + averageUtilization: 70 + + - it: renders only custom metrics when utilization targets are disabled + template: templates/hpa.yaml + set: + workload.kind: deployment + server.externalDbSecret: my-pg-secret + autoscaling.enabled: true + autoscaling.targetCPUUtilizationPercentage: null + autoscaling.metrics: + - type: Pods + pods: + metric: + name: example_adapter_metric + target: + type: AverageValue + averageValue: "50" + asserts: + - lengthEqual: + path: spec.metrics + count: 1 + - equal: + path: spec.metrics[0].type + value: Pods + + - it: merges a scale-up policy with the default scale-down behavior + template: templates/hpa.yaml + values: + - ../ci/values-autoscaling.yaml + set: + autoscaling.behavior: + scaleUp: + stabilizationWindowSeconds: 0 + asserts: + - equal: + path: spec.behavior.scaleUp.stabilizationWindowSeconds + value: 0 + - equal: + path: spec.behavior.scaleDown.stabilizationWindowSeconds + value: 300 + + - it: omits behavior when set to null + template: templates/hpa.yaml + values: + - ../ci/values-autoscaling.yaml + set: + autoscaling.behavior: null + asserts: + - notExists: + path: spec.behavior + + - it: targets a multi-replica StatefulSet with the explicit override + template: templates/hpa.yaml + set: + server.externalDbSecret: my-pg-secret + workload.allowMultiReplicaStatefulSet: true + resources.requests.cpu: 100m + autoscaling.enabled: true + asserts: + - equal: + path: spec.scaleTargetRef.kind + value: StatefulSet + + - it: omits spec.replicas from the StatefulSet when autoscaling is enabled + template: templates/statefulset.yaml + set: + server.externalDbSecret: my-pg-secret + workload.allowMultiReplicaStatefulSet: true + resources.requests.cpu: 100m + autoscaling.enabled: true + asserts: + - notExists: + path: spec.replicas + + - it: treats a missing autoscaling map as disabled + template: templates/statefulset.yaml + set: + autoscaling: null + asserts: + - equal: + path: spec.replicas + value: 1 + + - it: fails when autoscaling can exceed one replica on SQLite + template: templates/statefulset.yaml + set: + autoscaling.enabled: true + resources.requests.cpu: 100m + asserts: + - failedTemplate: + errorPattern: "autoscaling.maxReplicas > 1 requires server.externalDbSecret" + + - it: fails when autoscaling a StatefulSet without the multi-replica override + template: templates/statefulset.yaml + set: + autoscaling.enabled: true + server.externalDbSecret: my-pg-secret + resources.requests.cpu: 100m + asserts: + - failedTemplate: + errorPattern: "autoscaling.maxReplicas > 1 with workload.kind=statefulset requires workload.allowMultiReplicaStatefulSet=true" + + - it: fails when maxReplicas is below minReplicas + template: templates/statefulset.yaml + values: + - ../ci/values-autoscaling.yaml + set: + autoscaling.maxReplicas: 1 + asserts: + - failedTemplate: + errorPattern: "autoscaling.maxReplicas must be greater than or equal to autoscaling.minReplicas" + + - it: fails when minReplicas is below one + template: templates/statefulset.yaml + values: + - ../ci/values-autoscaling.yaml + set: + autoscaling.minReplicas: 0 + asserts: + - failedTemplate: + errorPattern: "autoscaling.minReplicas must be at least 1" + + - it: fails when no autoscaling metric is configured + template: templates/statefulset.yaml + set: + workload.kind: deployment + server.externalDbSecret: my-pg-secret + autoscaling.enabled: true + autoscaling.targetCPUUtilizationPercentage: null + asserts: + - failedTemplate: + errorPattern: "autoscaling.enabled requires targetCPUUtilizationPercentage, targetMemoryUtilizationPercentage, or autoscaling.metrics" + + - it: fails when a CPU target has no CPU request + template: templates/statefulset.yaml + set: + workload.kind: deployment + server.externalDbSecret: my-pg-secret + autoscaling.enabled: true + asserts: + - failedTemplate: + errorPattern: "autoscaling.targetCPUUtilizationPercentage requires a positive resources.requests.cpu" + + - it: fails when a memory target has no memory request + template: templates/statefulset.yaml + values: + - ../ci/values-autoscaling.yaml + set: + autoscaling.targetMemoryUtilizationPercentage: 70 + resources.requests.memory: null + asserts: + - failedTemplate: + errorPattern: "autoscaling.targetMemoryUtilizationPercentage requires a positive resources.requests.memory" + + # Kubernetes copies a limit only into an unset request, so an explicit zero + # request stays zero and the HPA cannot compute utilization. + - it: fails when the CPU request is zero even with a CPU limit + template: templates/statefulset.yaml + values: + - ../ci/values-autoscaling.yaml + set: + resources.requests.cpu: 0 + resources.limits.cpu: "1" + asserts: + - failedTemplate: + errorPattern: "autoscaling.targetCPUUtilizationPercentage requires a positive resources.requests.cpu" + + - it: fails when the memory request is zero even with a memory limit + template: templates/statefulset.yaml + values: + - ../ci/values-autoscaling.yaml + set: + autoscaling.targetMemoryUtilizationPercentage: 70 + resources.requests.memory: 0Mi + resources.limits.memory: 1Gi + asserts: + - failedTemplate: + errorPattern: "autoscaling.targetMemoryUtilizationPercentage requires a positive resources.requests.memory" + + - it: fails when the CPU request is a zero quantity string + template: templates/statefulset.yaml + values: + - ../ci/values-autoscaling.yaml + set: + resources.requests.cpu: 0m + asserts: + - failedTemplate: + errorPattern: "autoscaling.targetCPUUtilizationPercentage requires a positive resources.requests.cpu" + + # The chart renders a null request, which Kubernetes stores as zero instead + # of copying the limit. + - it: fails when the CPU request is null even with a CPU limit + template: templates/statefulset.yaml + values: + - ../ci/values-autoscaling.yaml + set: + resources.requests.cpu: null + resources.limits.cpu: "1" + asserts: + - failedTemplate: + errorPattern: "autoscaling.targetCPUUtilizationPercentage requires a positive resources.requests.cpu" + + - it: fails when the CPU request is zero in exponent form + template: templates/statefulset.yaml + values: + - ../ci/values-autoscaling.yaml + set: + resources.requests.cpu: "0e3" + asserts: + - failedTemplate: + errorPattern: "autoscaling.targetCPUUtilizationPercentage requires a positive resources.requests.cpu" + + - it: fails when the memory request is zero in exponent form + template: templates/statefulset.yaml + values: + - ../ci/values-autoscaling.yaml + set: + autoscaling.targetMemoryUtilizationPercentage: 70 + resources.requests.memory: "0E+6" + asserts: + - failedTemplate: + errorPattern: "autoscaling.targetMemoryUtilizationPercentage requires a positive resources.requests.memory" + + - it: accepts CPU and memory limits in place of requests + template: templates/hpa.yaml + set: + workload.kind: deployment + server.externalDbSecret: my-pg-secret + autoscaling.enabled: true + autoscaling.targetMemoryUtilizationPercentage: 70 + resources.limits.cpu: "1" + resources.limits.memory: 1Gi + asserts: + - isKind: + of: HorizontalPodAutoscaler + - contains: + path: spec.metrics + content: + type: Resource + resource: + name: cpu + target: + type: Utilization + averageUtilization: 80 + - contains: + path: spec.metrics + content: + type: Resource + resource: + name: memory + target: + type: Utilization + averageUtilization: 70 + + - it: explains --reuse-values when the autoscaling defaults are missing + template: templates/statefulset.yaml + set: + workload.kind: deployment + server.externalDbSecret: my-pg-secret + resources.requests.cpu: 100m + # --reuse-values from a release without autoscaling values leaves only + # the enabled flag that the upgrade sets. + autoscaling.enabled: true + autoscaling.minReplicas: null + autoscaling.maxReplicas: null + autoscaling.targetCPUUtilizationPercentage: null + autoscaling.metrics: null + autoscaling.behavior: null + asserts: + - failedTemplate: + errorPattern: "autoscaling.minReplicas and autoscaling.maxReplicas are not set.*--reset-then-reuse-values" + + - it: keeps the replicaCount wording when autoscaling is disabled + template: templates/statefulset.yaml + set: + replicaCount: 2 + asserts: + - failedTemplate: + errorPattern: "replicaCount > 1 requires server.externalDbSecret" diff --git a/deploy/helm/openshell/tests/certgen_test.yaml b/deploy/helm/openshell/tests/certgen_test.yaml index cd88b60e97..7a478bbf76 100644 --- a/deploy/helm/openshell/tests/certgen_test.yaml +++ b/deploy/helm/openshell/tests/certgen_test.yaml @@ -114,3 +114,38 @@ tests: path: spec.template.spec.containers[0].args content: "--server-san=192.0.2.10" documentIndex: 3 + + - it: labels hook pods so that gateway selectors do not match them + template: templates/certgen.yaml + asserts: + - equal: + path: spec.template.metadata.labels + value: + app.kubernetes.io/name: openshell-certgen + app.kubernetes.io/instance: openshell + app.kubernetes.io/component: certgen + documentIndex: 3 + - notEqual: + path: spec.template.metadata.labels["app.kubernetes.io/name"] + value: openshell + documentIndex: 3 + + - it: labels the backend CA hook pod so that gateway selectors do not match it + template: templates/certgen.yaml + set: + certManager.enabled: true + grpcRoute.backendTLSPolicy.enabled: true + asserts: + - hasDocuments: + count: 5 + - equal: + path: metadata.name + value: openshell-certgen-backend-ca + documentIndex: 4 + - equal: + path: spec.template.metadata.labels + value: + app.kubernetes.io/name: openshell-certgen + app.kubernetes.io/instance: openshell + app.kubernetes.io/component: certgen + documentIndex: 4 diff --git a/deploy/helm/openshell/values.yaml b/deploy/helm/openshell/values.yaml index 4bf21b88a8..1c2b7738a9 100644 --- a/deploy/helm/openshell/values.yaml +++ b/deploy/helm/openshell/values.yaml @@ -217,6 +217,50 @@ podLifecycle: # covers supervisor session cleanup and draining queued OCSF records. terminationGracePeriodSeconds: 30 +# Optional HorizontalPodAutoscaler (autoscaling/v2) for the gateway workload. +# When enabled, the Deployment or StatefulSet omits spec.replicas and the HPA +# owns the replica count; replicaCount is ignored. Multi-replica limits +# (server.externalDbSecret, and workload.allowMultiReplicaStatefulSet for a +# StatefulSet) apply to maxReplicas. Scaling out does not move established +# supervisor sessions; read the High Availability guide before choosing metrics. +# The workload can drop to one replica when you enable this on an existing +# release, until the HPA scales it back up; read the chart README before you +# upgrade. +autoscaling: + # -- Render a HorizontalPodAutoscaler and stop rendering spec.replicas. + enabled: false + # -- Minimum gateway replicas. Use 2 or more to survive a pod failure. + minReplicas: 2 + # -- Maximum gateway replicas. Each replica opens its own PostgreSQL + # connection pool; size the database for rollouts at this count, as the High + # Availability guide describes. + maxReplicas: 4 + # -- Target average CPU utilization, as a percentage of + # resources.requests.cpu. Set to null to disable. Requires + # resources.requests.cpu, or resources.limits.cpu, which Kubernetes copies + # into the request. + targetCPUUtilizationPercentage: 80 + # -- (int) Target average memory utilization, as a percentage of + # resources.requests.memory. Null disables it. Requires + # resources.requests.memory, or resources.limits.memory, which Kubernetes + # copies into the request. + targetMemoryUtilizationPercentage: null + # -- Additional autoscaling/v2 MetricSpec entries appended verbatim, such as + # Pods metrics served by prometheus-adapter. + metrics: [] + # -- HPA scaling behavior. Scale-down disconnects the removed pod's + # supervisor sessions; they reconnect to the remaining replicas, so the + # default removes at most one replica every two minutes after a + # five-minute stabilization window. Helm merges maps: set + # autoscaling.behavior.scaleDown to null to drop the default. + behavior: + scaleDown: + stabilizationWindowSeconds: 300 + policies: + - type: Pods + value: 1 + periodSeconds: 120 + probes: startup: # -- Startup probe period, in seconds. diff --git a/docs/kubernetes/high-availability.mdx b/docs/kubernetes/high-availability.mdx index f83b181462..0153df050f 100644 --- a/docs/kubernetes/high-availability.mdx +++ b/docs/kubernetes/high-availability.mdx @@ -3,8 +3,8 @@ # SPDX-License-Identifier: Apache-2.0 title: "High Availability" sidebar-title: "High Availability" -description: "Run multiple OpenShell gateway replicas on Kubernetes with shared PostgreSQL and authenticated peer routing." -keywords: "Generative AI, Cybersecurity, Kubernetes, High Availability, HA, Gateway, PostgreSQL, Replicas, Failover" +description: "Run multiple OpenShell gateway replicas on Kubernetes with shared PostgreSQL, authenticated peer routing, capacity metrics, and optional autoscaling." +keywords: "Generative AI, Cybersecurity, Kubernetes, High Availability, HA, Gateway, PostgreSQL, Replicas, Failover, Autoscaling, HPA, Metrics" position: 3 --- @@ -32,6 +32,9 @@ An HA gateway deployment requires: database are intended for a single gateway replica. - An ingress or load balancer that routes clients to the gateway Service. Refer to [Ingress](/kubernetes/ingress) for a Gateway API configuration. +- PostgreSQL connection capacity for every replica, including pods that are + starting or terminating during a rollout. Refer to + [Size PostgreSQL Connections](#size-postgresql-connections). The Helm chart rejects `replicaCount` values above `1` unless `server.externalDbSecret` is set. It also rejects a multi-replica StatefulSet @@ -186,8 +189,68 @@ session instead of resuming the interrupted byte stream. Rolling updates can temporarily concentrate supervisor sessions on the replicas that stayed up; client requests remain routable through peer relay. +## Monitor Capacity + +Each gateway replica exposes Prometheus metrics on port `9090` +(`service.metricsPort`). Scrape every gateway pod, not the Service, because +these signals describe one replica. Refer to +[Gateway Metrics](/observability/gateway-metrics) for the full catalog, scrape +configuration, and access control. + +| Signal | Metrics | Use | +|---|---|---| +| Sessions per replica | `openshell_server_supervisor_sessions` | Placement and skew. Do not scale on it. | +| Relay utilization | `openshell_server_relay_pending`, `openshell_server_relay_pending_capacity` | Graph it. It is transient and usually near 0, so alert on rejections instead. | +| Relay rejections | `openshell_server_relay_rejected_total` | Alert on any increase. | +| Relay claims | `openshell_server_relay_claim_duration_seconds`, `openshell_server_relay_expired_total` | Supervisor connect-back time for exec, SSH, forwarding, and service traffic, including relays requested through peers. A high 99th percentile while rejections stay flat points at the supervisor or its node, not at relay capacity. Relays not claimed within 10 seconds are missing from the histogram and count as expired, so alert on any increase. | +| Routed requests | `openshell_server_routed_request_attempts_total`, `openshell_server_peer_request_duration_seconds` | Local/peer relay setup mix and outbound peer rate, failures, and latency. Counts completed or cancelled attempts, including retries. Spikes of `grpc_code="unavailable"` during rollouts are expected. Do not use as an HPA target. | + +These queries assume that Prometheus labels each series with `namespace` and +`pod`: + +```promql +# Relay capacity used on each replica. Usually near 0; graph it and alert +# on rejections instead. +max by (namespace, pod) (openshell_server_relay_pending / openshell_server_relay_pending_capacity) + +# Relays rejected for capacity. Alert on any increase. +sum by (namespace, pod, reason) (increase(openshell_server_relay_rejected_total[5m])) + +# 99th percentile time for supervisors to claim a relay on each replica. +histogram_quantile(0.99, sum by (namespace, pod, le) (rate(openshell_server_relay_claim_duration_seconds_bucket[5m]))) + +# Completed or cancelled relay setup attempts, by replica, relay kind, and route. +sum by (namespace, pod, relay_kind, route) (rate(openshell_server_routed_request_attempts_total{operation="relay"}[5m])) + +# Session skew across replicas. 1 means balanced. Starting and stopping pods +# raise it during rollouts. +max(openshell_server_supervisor_sessions) + / clamp_min(avg(openshell_server_supervisor_sessions), 1) + +# Failed peer requests as a share of attempts, by operation. +sum by (operation) (rate(openshell_server_routed_request_attempts_total{route="peer",outcome!="success"}[5m])) + / sum by (operation) (rate(openshell_server_routed_request_attempts_total{route="peer"}[5m])) + +# 99th percentile peer request latency across the fleet. +histogram_quantile(0.99, sum by (le, operation) (rate(openshell_server_peer_request_duration_seconds_bucket{outcome="success"}[5m]))) +``` + +Relay capacity is used on the replica that owns the sandbox's supervisor +session, including relays that other replicas request through peers, so +adding replicas does not relieve a busy owner. A relay whose client gave up +stays counted for up to about 40 seconds, the 10-second claim timeout plus the +30-second cleanup interval. Routed request counts include retries during +rollouts. + ## Scale the Gateway +Scale the gateway by changing a fixed replica count or by letting a +HorizontalPodAutoscaler choose one. In both cases, keep at least two ready +replicas when availability must survive one gateway pod failure, and size +PostgreSQL and the cluster nodes for the largest replica count. + +### Scale Manually + Change `replicaCount` in `values-ha.yaml`, then apply the release again: ```shell @@ -199,9 +262,155 @@ helm upgrade openshell \ --wait ``` -Keep at least two ready replicas when availability must survive one gateway pod -failure. Size PostgreSQL connection capacity and the cluster nodes for the -selected replica count. +### Autoscale with a HorizontalPodAutoscaler + +Set `autoscaling.enabled` to render an `autoscaling/v2` +HorizontalPodAutoscaler for the gateway Deployment, or for a StatefulSet with +`workload.allowMultiReplicaStatefulSet`. Add these values to +`values-ha.yaml`: + +```yaml +resources: + requests: + cpu: 500m + memory: 512Mi + +autoscaling: + enabled: true + minReplicas: 2 + maxReplicas: 6 + targetCPUUtilizationPercentage: 70 +``` + +With autoscaling enabled, the chart stops setting `spec.replicas` and ignores +`replicaCount`. Kubernetes computes CPU and memory utilization against the +container request, so those targets need a matching `resources.requests` +entry. A `resources.limits` entry also works, because Kubernetes copies a +limit into an unset request. The chart rejects an `autoscaling.maxReplicas` +above `1` without `server.externalDbSecret`. Each scale-down disconnects the +removed replica's supervisor sessions, and they reconnect to the remaining +replicas, so the default behavior removes at most one pod every two minutes +after a five-minute stabilization window. Override it with +`autoscaling.behavior`, or set `autoscaling.behavior.scaleDown` to `null` to +drop the default. New replicas do not take existing sessions. + + +A scale-down disconnects every supervisor session on the removed pod at once, +and requests through that pod must be retried. To choose when that happens, +set `autoscaling.behavior.scaleDown.selectPolicy` to `Disabled` so the +autoscaler only adds replicas, and remove replicas yourself with +`kubectl scale` at a quiet time. + + +Check the autoscaler after you apply the release: + +```shell +kubectl -n openshell get hpa openshell +kubectl -n openshell describe hpa openshell +``` + + +Enabling autoscaling on an existing release removes `spec.replicas` from the +workload in that upgrade. Kubernetes can reset the workload to one replica +until the HPA scales it back to at least `autoscaling.minReplicas`, which +disconnects supervisor sessions from the other gateway pods; they reconnect to +the remaining replica. Enable autoscaling +when you install the chart, or during a maintenance window. On an existing +release, upgrade with `--reset-then-reuse-values`, which needs Helm 3.14 or +later. With an older Helm, save the values with +`helm get values -o yaml > values.yaml` and upgrade with +`--reset-values -f values.yaml`. `--reuse-values` keeps the +previous chart version's defaults, which lack the autoscaling values, +including the scale-down `behavior`, when that release predates them. For background, refer to +[Migrating Deployments and StatefulSets to horizontal autoscaling](https://kubernetes.io/docs/tasks/run-application/horizontal-pod-autoscale/#migrating-deployments-and-statefulsets-to-horizontal-autoscaling). +Disabling autoscaling renders `spec.replicas` from `replicaCount` again, +which defaults to 1, so set `replicaCount` to the replica count you want in +that upgrade. + + +### Autoscale on Custom Metrics + +The HPA can also scale on a per-pod metric that a metrics adapter publishes. +The chart passes `autoscaling.metrics` to the HPA unchanged and does not +install Prometheus or an adapter. A custom metric needs these pieces: + +- Prometheus scrapes every gateway pod and labels each series with + `namespace` and `pod`. +- A metrics adapter, such as prometheus-adapter, maps those labels to + Kubernetes resources and publishes a per-pod metric through the + `custom.metrics.k8s.io` API. +- `autoscaling.metrics` references that metric as a `Pods` metric with an + `AverageValue` target. + +Scale on signals that fall when you add replicas, such as CPU utilization. Do +not scale on supervisor sessions, relay utilization, relay claim latency, or +peer request rate. New replicas do not take existing sessions, relay load stays +on the replica that owns each session, and peer traffic grows with the replica +count. Alert on those signals instead. + +```yaml +autoscaling: + enabled: true + minReplicas: 2 + maxReplicas: 6 + targetCPUUtilizationPercentage: 70 + metrics: + - type: Pods + pods: + metric: + name: + target: + type: AverageValue + averageValue: "100" +``` + +Replace `` with the name your adapter publishes, and set +`averageValue` to the per-pod value at which you want another replica. Check +that the adapter serves the metric for the gateway pods: + +```shell +kubectl get --raw "/apis/custom.metrics.k8s.io/v1beta1/namespaces/openshell/pods/*/" +``` + +### Size PostgreSQL Connections + +Each gateway pod opens up to 10 PostgreSQL connections on demand. The chart +uses the Kubernetes default rolling update strategy for each workload kind, so +the number of pods that run during a rollout depends on the kind. + +A Deployment rollout adds up to 25 percent of the replica count, rounded up, +as surge pods. Kubernetes does not count terminating pods against that surge, +and each replaced pod can keep its connections throughout its termination +grace period. A rollout can therefore run up to twice the replica count at +once, and an eviction or an autoscaler scale-down during the rollout adds more +terminating pods. Set PostgreSQL `max_connections` to at least +`(2 × replicas + surge) × 10`, which leaves one surge of margin for those +pods, plus headroom for your other clients and administration. Use +`autoscaling.maxReplicas` as the replica count when autoscaling is enabled. + +| Replicas | Surge | Minimum `max_connections` | +|---|---|---| +| 2 | 1 | `(4 + 1) × 10 = 50` | +| 3 | 1 | `(6 + 1) × 10 = 70` | +| 4 | 1 | `(8 + 1) × 10 = 90` | +| 6 | 2 | `(12 + 2) × 10 = 140` | + +With a Deployment, five or more replicas exceed what the PostgreSQL default of +100 leaves for the gateway, because PostgreSQL reserves 3 connections for +superusers. Raise `max_connections` or use a larger managed instance. Let one +rollout finish before you start another, because each overlapping rollout adds +its own terminating pods. + +If you run several replicas as a StatefulSet with +`workload.allowMultiReplicaStatefulSet`, a rollout adds no surge pods. It +replaces one pod at a time and creates each replacement only after the old pod +exits. Set `max_connections` to at least `replicas × 10`, using +`autoscaling.maxReplicas` when autoscaling is enabled, plus headroom for your +other clients and administration. + +If a connection pooler sits between the gateway and PostgreSQL, use session +pooling. The gateway holds session-level advisory locks, which transaction +pooling breaks. ## Next Steps @@ -211,3 +420,5 @@ selected replica count. [Managing Certificates](/kubernetes/managing-certificates). - To configure user authentication and authorization, refer to [Access Control](/kubernetes/access-control). +- To scrape and alert on gateway metrics, refer to + [Gateway Metrics](/observability/gateway-metrics). diff --git a/docs/kubernetes/setup.mdx b/docs/kubernetes/setup.mdx index 2b39ddf223..62bce316ad 100644 --- a/docs/kubernetes/setup.mdx +++ b/docs/kubernetes/setup.mdx @@ -239,6 +239,7 @@ The most commonly changed values are: | `replicaCount` | Number of gateway replicas. Values above `1` require shared PostgreSQL through `server.externalDbSecret`. | | `workload.kind` | Gateway workload controller. Use `statefulset` for SQLite or `deployment` with `server.externalDbSecret`. | | `workload.allowMultiReplicaStatefulSet` | Allow `replicaCount > 1` with `workload.kind=statefulset`. Prefer Deployment for external database-backed multi-replica gateways. | +| `autoscaling.enabled` / `autoscaling.minReplicas` / `autoscaling.maxReplicas` | Render a HorizontalPodAutoscaler instead of a fixed `replicaCount`. Refer to [High Availability](/kubernetes/high-availability#scale-the-gateway). | | `server.sandboxNamespace` | Namespace where sandbox pods are created. Defaults to the Helm release namespace when left empty. | | `workspaceResources.enabled` | Create namespace-scoped sandbox prerequisites from the gateway chart. Disable when installing the workspace chart separately. | | `server.externalDbSecret` | Secret containing a PostgreSQL connection URI in the `uri` key. Use when the database is managed outside the chart. | diff --git a/docs/observability/gateway-metrics.mdx b/docs/observability/gateway-metrics.mdx new file mode 100644 index 0000000000..aa849977d3 --- /dev/null +++ b/docs/observability/gateway-metrics.mdx @@ -0,0 +1,270 @@ +--- +# SPDX-FileCopyrightText: Copyright (c) 2025-2026 NVIDIA CORPORATION & AFFILIATES. All rights reserved. +# SPDX-License-Identifier: Apache-2.0 +title: "Gateway Metrics" +sidebar-title: "Gateway Metrics" +description: "Scrape Prometheus metrics from the OpenShell gateway, including supervisor session, relay capacity, and peer routing signals for multi-replica deployments." +keywords: "Generative AI, Cybersecurity, Observability, Metrics, Prometheus, Gateway, Kubernetes, High Availability" +--- + +The gateway exports Prometheus metrics for API requests, database readiness, +gateway interceptors, the OCSF log file, and multi-replica capacity. Use them +to alert on saturation, to follow rollouts, and to choose autoscaling signals. + +## Enable the Metrics Endpoint + +The gateway serves `GET /metrics` in the Prometheus text format on its own +listener, separate from the API and health ports. When the listener is +disabled, the gateway records no metrics. Enable it in one of these ways: + +- In `gateway.toml`, set `metrics_bind_address` under `[openshell.gateway]`, + for example `metrics_bind_address = "0.0.0.0:9090"`. Refer to + [Gateway Configuration](/how-it-works/gateways/configuration). +- On the command line, pass `--metrics-port` or set `OPENSHELL_METRICS_PORT`. + The standalone gateway binary defaults to `0`, which disables the endpoint. +- In the Helm chart, set `service.metricsPort`, which defaults to `9090`. The + chart renders the TOML key, a container port named `metrics`, and a + `metrics` port on the gateway Service. Set it to `0` to disable the + endpoint. + +A stopping gateway keeps serving metrics until the process exits. + + +The metrics endpoint uses plaintext HTTP with no authentication. Anyone who +can reach the port can read the metrics. The Helm chart publishes the port on +the gateway Service, so a `LoadBalancer` or `NodePort` Service exposes it +outside the cluster. Keep the Service internal, and restrict the port to your +monitoring system. + + +The following NetworkPolicy keeps the API and health ports open and accepts +metrics scrapes only from the `monitoring` namespace. Adjust the ports if you +changed `service.port`, `service.healthPort`, or `service.metricsPort`. + +```yaml +apiVersion: networking.k8s.io/v1 +kind: NetworkPolicy +metadata: + name: openshell-gateway-metrics + namespace: openshell +spec: + podSelector: + matchLabels: + app.kubernetes.io/name: openshell + app.kubernetes.io/instance: openshell + policyTypes: + - Ingress + ingress: + - ports: + - port: 8080 + - port: 8081 + - from: + - namespaceSelector: + matchLabels: + kubernetes.io/metadata.name: monitoring + ports: + - port: 9090 +``` + +## Scrape Each Replica + +Capacity metrics describe one gateway replica, so scrape every gateway pod +rather than the Service, and keep the `namespace` and `pod` labels on each +series. With the Prometheus Operator, a PodMonitor adds those labels +automatically: + +```yaml +apiVersion: monitoring.coreos.com/v1 +kind: PodMonitor +metadata: + name: openshell-gateway + namespace: openshell +spec: + selector: + matchLabels: + app.kubernetes.io/name: openshell + app.kubernetes.io/instance: openshell + podMetricsEndpoints: + - port: metrics + path: /metrics + interval: 30s +``` + +With a plain Prometheus configuration, use pod discovery and copy the labels +yourself: + +```yaml +scrape_configs: + - job_name: openshell-gateway + kubernetes_sd_configs: + - role: pod + namespaces: + names: [openshell] + relabel_configs: + - source_labels: [__meta_kubernetes_pod_label_app_kubernetes_io_instance] + regex: openshell + action: keep + - source_labels: [__meta_kubernetes_pod_container_port_name] + regex: metrics + action: keep + - source_labels: [__meta_kubernetes_namespace] + target_label: namespace + - source_labels: [__meta_kubernetes_pod_name] + target_label: pod +``` + +To check one pod without Prometheus, read its metrics through the Kubernetes +API server proxy: + +```shell +kubectl get --raw "/api/v1/namespaces/openshell/pods/:9090/proxy/metrics" \ + | grep '^openshell_server_supervisor_sessions' +``` + +The API server connects to the pod from the control plane. A NetworkPolicy +like the one above blocks that connection unless it also allows the control +plane on the metrics port. `kubectl port-forward` reaches the pod through the +kubelet instead, and NetworkPolicy does not block it. Start the port-forward: + +```shell +kubectl -n openshell port-forward pod/ 9090:9090 +``` + +Then read the metrics from another terminal: + +```shell +curl -s http://localhost:9090/metrics | grep '^openshell_server_supervisor_sessions' +``` + +## Histograms and Summaries + +The multi-replica latency metrics are Prometheus histograms with `_bucket`, +`_sum`, and `_count` series and buckets from 1 millisecond to 15 seconds, so +`histogram_quantile` can aggregate them across replicas. The older +`*_duration_seconds` metrics and the interceptor latency metric are summaries +with `quantile` labels that each replica computes on its own. Do not average +summary quantiles across pods. + +## Multi-Replica Capacity Metrics + +These metrics describe how much work one replica holds and how the replicas +interact. No metric carries a sandbox, channel, endpoint, or replica +identifier, because the scrape target already identifies the replica. + +Supervisor sessions: + +| Metric | Type | Labels | Description | +|---|---|---|---| +| `openshell_server_supervisor_sessions` | Gauge | None | Supervisor control sessions held by this replica. After a rollout settles, the sum across pods equals the number of connected sandboxes. | + +Relays: + +| Metric | Type | Labels | Description | +|---|---|---|---| +| `openshell_server_relay_pending` | Gauge | None | Relay channels on this replica waiting for the supervisor to connect back, including channels opened for peer replicas. A channel whose caller gave up stays counted until it is claimed or cleaned up, up to about 40 seconds. | +| `openshell_server_relay_pending_capacity` | Gauge | None | Pending relay limit per replica, `256`. | +| `openshell_server_relay_rejected_total` | Counter | `reason` | Relay opens rejected at a limit. Rejected relays fail the client request, so use this counter, not client error codes, to detect saturation. | +| `openshell_server_relay_expired_total` | Counter | None | Pending relays dropped because the supervisor did not connect back within 10 seconds. | +| `openshell_server_relay_claim_duration_seconds` | Histogram | None | Time from opening a relay to the supervisor claiming it. | + +Routed requests include local relay setup and outbound requests to the replica +that owns a sandbox's supervisor session: + +| Metric | Type | Labels | Description | +|---|---|---|---| +| `openshell_server_routed_request_attempts_total` | Counter | `operation`, `route`, `relay_kind`, `outcome`, `grpc_code` | Local relay setup and outbound peer attempts, counted once when they finish or are cancelled. Each retry counts separately. | +| `openshell_server_peer_request_duration_seconds` | Histogram | `operation`, `outcome` | Latency of outbound peer requests only (`route="peer"`). For relays, until the owner's supervisor claimed the relay. | + +The labels take these values: + +- `operation` is `relay`, `report_provider_readiness`, `report_endpoint_status`, + or `get_sandbox_provider_status`. Only relay setup includes the local route; + the other operations count outbound peer calls, not local API handling. On the + owning replica, `openshell_server_grpc_requests_total` records the matching + `method`, such as `PeerRelay` for `relay`. +- `relay_kind` is `ssh` or `tcp` for relays, and `none` for other operations. An + omitted relay target counts as `ssh`, matching the protocol's compatibility + behavior. TCP hostnames and ports are not labels. +- `route` is `local` when this replica tries its own supervisor session, or + `peer` when it tries another replica. A routed retry counts again, including a + change from local to peer routing. Waiting for an owner does not count. Incoming + `PeerRelay` requests do not increment this counter on the receiving replica; + they appear in that replica's gRPC request counter instead. This counter does + not measure active connections, unique client requests, or unfinished attempts. +- `reason` is `replica_capacity` when the gateway replica's pending-relay budget + is full, or `sandbox_capacity` when one sandbox's budget on that replica, 32 + pending relays, is full. +- `outcome` is `success`, `local_error`, or `remote_error`. A relay succeeds + when the supervisor claims it; this does not measure the data stream. + `local_error` means the attempt failed on this replica, and `remote_error` + means the owning replica returned an error or the open connection to it + failed. On `route="local"` this replica is the owner, so failures that appear + as `remote_error` on the peer route, such as an unclaimed relay + (`deadline_exceeded`) or a full relay budget (`resource_exhausted`), appear as + `local_error`. A cancelled attempt records `local_error` with + `grpc_code="cancelled"`, even when the owner was already handling it. +- `grpc_code` is the lowercase gRPC status name, such as `ok`, `unavailable`, + or `resource_exhausted`, from local setup, the local supervisor claim, or the + peer RPC. For a peer relay, the gateway then reports a failure to the client + as `UNAVAILABLE`. The other RPCs return the owner's status to the client + unchanged. The label is not named `code` because + `openshell_server_grpc_requests_total` uses that name for the numeric status + code. + +Gauges and counters exist from startup. They start at `0`, except +`openshell_server_relay_pending_capacity`, which starts at its limit. `openshell_server_routed_request_attempts_total` starts with seven +`outcome="success", grpc_code="ok"` series: four relay kind and route +combinations, and three other peer operations with `relay_kind="none"`. Other outcomes and codes, +and every histogram, appear after their first sample, so write alert expressions +that tolerate absent series. + +## Request and Readiness Metrics + +The gateway also records every API request and its background database check: + +| Metric | Type | Labels | Description | +|---|---|---|---| +| `openshell_server_grpc_requests_total` | Counter | `method`, `code` | gRPC requests by method name and numeric status code, including internal peer requests. Streaming RPCs record the status known when the response headers are sent. | +| `openshell_server_grpc_request_duration_seconds` | Summary | `method`, `code` | Time until the response headers are sent. For streaming RPCs this covers setup only. | +| `openshell_server_http_requests_total` | Counter | `path`, `status` | HTTP requests on the API port. `path` is `/_ws_tunnel`, `/auth`, or `unknown`. | +| `openshell_server_http_request_duration_seconds` | Summary | `path`, `status` | Duration of those HTTP requests. | +| `openshell_server_readiness_database_healthy` | Gauge | None | `1` when the last background database check succeeded, otherwise `0`. | +| `openshell_server_readiness_database_probe_duration_seconds` | Summary | `outcome` | Duration of the background database check. `outcome` is `success`, `db_error`, or `timeout`. | + +## Gateway Interceptor Metrics + +Gateways with [gateway interceptors](/extensibility/gateway-interceptors) +record their evaluations: + +| Metric | Type | Labels | Description | +|---|---|---|---| +| `openshell_gateway_interceptor_latency_seconds` | Summary | None | Duration of interceptor calls. | +| `openshell_gateway_interceptor_evaluations_total` | Counter | `decision`, `interceptor`, `binding_id` | Interceptor evaluations by decision. | +| `openshell_gateway_interceptor_patches_total` | Counter | `interceptor`, `binding_id` | Patches that interceptors applied. | +| `openshell_gateway_interceptor_fail_open_total` | Counter | None | Interceptor failures that let the request continue. | +| `openshell_gateway_interceptor_fail_closed_total` | Counter | None | Interceptor failures that rejected the request. | +| `openshell_gateway_interceptor_post_commit_observation_failures_total` | Counter | `stage` | Failed post-commit observations. `stage` is `response_body` or `evaluation`. | + +## OCSF Log Metrics + +Gateways that write an +[OCSF JSONL file](/how-it-works/gateways/configuration#ocsf-jsonl-output) +record its queue and writer: + +| Metric | Type | Labels | Description | +|---|---|---|---| +| `openshell_ocsf_log_queued_total` | Counter | None | Records queued for the file. | +| `openshell_ocsf_log_written_total` | Counter | None | Records written to the file. | +| `openshell_ocsf_log_dropped_total` | Counter | `reason` | Records discarded, by bounded reason. | +| `openshell_ocsf_log_writer_errors_total` | Counter | None | Failed writes, reopens, and retention cleanups. | +| `openshell_ocsf_log_recovery_discarded_bytes_total` | Counter | None | Bytes of an incomplete trailing line removed when the writer opens the file. | +| `openshell_ocsf_log_queue_records` | Gauge | None | Records waiting to be written. | +| `openshell_ocsf_log_queue_bytes` | Gauge | None | Encoded bytes waiting to be written. | + +## Next Steps + +- To monitor capacity and choose autoscaling signals for multiple gateway + replicas, refer to + [High Availability](/kubernetes/high-availability#monitor-capacity). +- To configure the metrics listener in `gateway.toml`, refer to + [Gateway Configuration](/how-it-works/gateways/configuration). diff --git a/skills/debug-openshell-cluster/SKILL.md b/skills/debug-openshell-cluster/SKILL.md index 46d5abc254..f454d9ea8f 100644 --- a/skills/debug-openshell-cluster/SKILL.md +++ b/skills/debug-openshell-cluster/SKILL.md @@ -511,6 +511,41 @@ name and load the chart CA plus client identity from the `peer-client-tls` volume exists, those files are readable, and the server certificate includes the name in `OPENSHELL_PEER_TLS_SERVER_NAME`. +To check per-replica capacity on a multi-replica gateway, read each gateway +pod's metrics and the autoscaler: + +```bash +for pod in $(kubectl -n openshell get pod \ + -l app.kubernetes.io/name=openshell,app.kubernetes.io/instance=openshell \ + -o jsonpath='{range .items[?(@.spec.containers[0].name=="openshell-gateway")]}{.metadata.name}{" "}{end}'); do + echo "${pod}" + kubectl get --raw "/api/v1/namespaces/openshell/pods/${pod}:9090/proxy/metrics" \ + | grep -E '^openshell_server_(supervisor_sessions|relay_pending|relay_rejected_total|routed_request_attempts_total)' +done +kubectl -n openshell get hpa +kubectl -n openshell describe hpa openshell +``` + +The JSONPath filter keeps only pods whose first container is the gateway +(`openshell-gateway`). It skips certificate hook Job pods, which older charts +labeled like gateway pods. The metrics port is `service.metricsPort` (default +`9090`). + +The API server proxy connects to each pod from the control plane. A +NetworkPolicy that accepts the metrics port only from a monitoring namespace +blocks it unless the policy also allows the control plane. In that case, +read one pod at a time through `kubectl port-forward`, which reaches the pod +through the kubelet and is not blocked by NetworkPolicy: + +```bash +kubectl -n openshell port-forward pod/ 9090:9090 >/dev/null & +pf_pid=$! +sleep 2 +curl -s http://localhost:9090/metrics \ + | grep -E '^openshell_server_(supervisor_sessions|relay_pending|relay_rejected_total|routed_request_attempts_total)' +kill "${pf_pid}" +``` + Check required Helm deployment secrets: ```bash @@ -957,6 +992,9 @@ credential failures. | Kubernetes gateway pod pending | PVC unbound, taint, selector, or insufficient resources | `kubectl -n openshell describe pod ` | | Kubernetes sandbox pod stuck pending, workspace PVC unbound | Cluster has no default `StorageClass` and OpenShell does not set `storageClassName` on the workspace PVC (clusters with a default `StorageClass` bind fine without it) | `kubectl -n openshell describe pvc`; set `server.workspaceStorageClass` (gateway config `workspace_storage_class`) to a valid `StorageClass` | | Kubernetes gateway pod crash loops | Missing secret, bad DB URL, bad TLS config | `kubectl -n openshell logs deployment/openshell -c openshell-gateway` or `kubectl -n openshell logs statefulset/openshell -c openshell-gateway` | +| `helm upgrade` fails with an `autoscaling.*` message | HPA values invalid: missing `resources.requests` (or `resources.limits`), `maxReplicas` above 1 without `server.externalDbSecret` (or on a StatefulSet without `workload.allowMultiReplicaStatefulSet`), no metric target, or min/max out of order. "`minReplicas` and `maxReplicas` are not set" means `--reuse-values` kept a release without the chart's autoscaling defaults | Fix the values named in the error; upgrade with `--reset-then-reuse-values` instead of `--reuse-values` | +| HPA shows `` targets | No metrics-server for CPU/memory, or the metrics adapter does not serve the custom metric | `kubectl -n openshell describe hpa openshell`, `kubectl get --raw /apis/custom.metrics.k8s.io/v1beta1` | +| One replica holds most sessions after a rollout | Expected: sessions stay where they reconnected | `openshell_server_supervisor_sessions` per pod; it fades as sandboxes are recreated | | OpenShift gateway pod fails to start with an SCC/`runAsUser` error (e.g. `unable to validate against any security context constraint`) | Chart's default `podSecurityContext`/`securityContext` hardcodes `runAsUser`/`fsGroup`, which the restricted-v2 SCC rejects; it must instead inject the namespace-assigned UID/GID range | `oc -n openshell describe pod `; deploy with `podSecurityContext: null` and clear `securityContext.runAsUser` (see `deploy/helm/openshell/ci/values-openshift-scc.yaml`) | | OpenShift sandbox pod fails to start (`unable to validate against any security context constraint`) | The `openshell-sandbox` service account lacks the privileged SCC it needs | `oc adm policy add-scc-to-user privileged -z openshell-sandbox -n openshell`; remove with `remove-scc-from-user` when done | | OpenShift self-hosted Vault/OpenBao credential store pod never schedules (waits time out with `no matching resources found`) | The store's Helm chart pins `runAsUser`/`fsGroup`/seccomp, which restricted-v2 rejects, so the StatefulSet controller never creates the pod | Deploy the store's chart in its OpenShift mode (`--set global.openshift=true` for the OpenBao/Vault chart) so the namespace SCC assigns a compliant security context — no manual SCC grant needed |