From 21620838f4dae6b8b87e685a257cab3cee178e7d Mon Sep 17 00:00:00 2001 From: Evie Howard Date: Fri, 18 Sep 2026 09:06:43 +0100 Subject: [PATCH 1/3] CARRY: test(odh): add gateway failover coverage Signed-off-by: Evie Howard --- e2e/rust/tests/odh/README.md | 80 +- e2e/rust/tests/odh/odh_harness/mod.rs | 1 + e2e/rust/tests/odh/odh_harness/oc.rs | 103 ++- e2e/rust/tests/odh/odh_harness/sandbox.rs | 34 + e2e/rust/tests/odh/tier3/gateway_failover.rs | 795 +++++++++++++++++++ e2e/rust/tests/odh/tier3/mod.rs | 2 + 6 files changed, 1002 insertions(+), 13 deletions(-) create mode 100644 e2e/rust/tests/odh/odh_harness/sandbox.rs create mode 100644 e2e/rust/tests/odh/tier3/gateway_failover.rs diff --git a/e2e/rust/tests/odh/README.md b/e2e/rust/tests/odh/README.md index 0382901ebd..a2e9cdec87 100644 --- a/e2e/rust/tests/odh/README.md +++ b/e2e/rust/tests/odh/README.md @@ -30,7 +30,8 @@ e2e/rust/tests/odh/ ├── tier2/ # Tier 2: medium/low priority positive tests │ └── mod.rs # empty — no scenarios yet ├── tier3/ # Tier 3: negative and destructive tests -│ └── mod.rs # empty — no scenarios yet +│ ├── mod.rs +│ └── gateway_failover.rs # external-PostgreSQL HA gateway failover ├── tiers.toml # tier → upstream test binaries + ODH module filter └── run-odh-test-tier.sh # runs one tier against a deployed gateway ``` @@ -55,8 +56,9 @@ plain module is enough — no new crate and no workspace change. `tokio::process::Command` for `oc`, injecting `--context` from `OPENSHELL_E2E_KUBE_CONTEXT_ACTIVE` (exported by `e2e/with-kube-gateway.sh`) when it is set, so every ODH test targets the same cluster the upstream - harness does. `oc_json()` runs a query and parses `-o json` output, panicking - with a descriptive message on any failure. + harness does. It also centralizes namespace/release resolution and Pod Ready + checks. `oc_json()` runs a query and parses `-o json` output, panicking with + a descriptive message on any failure. - `odh_harness::selinux::SelinuxAudit` — an opt-in scenario guard that detects OpenShift, requires every Ready worker to report `Enforcing`, records node-local audit cutoffs, and rejects OpenShell AVCs on completion. @@ -98,8 +100,9 @@ criticality, not just copied as-is. Smoke covers gateway reachability, sandbox lifecycle, and image provenance. Tier 1 checks the process-supervisor SELinux label and its mapped upstream -tests. Tier 2 and Tier 3 have no ODH scenarios yet; they run their mapped -upstream tests and the image provenance check. Add scenarios by creating a +tests. Tier 2 has no ODH scenarios yet; it runs its mapped upstream tests +and the image provenance check. Tier 3 includes an opt-in destructive +external-PostgreSQL HA gateway failover scenario. Add scenarios by creating a `.rs` file under the tier's directory and declaring it with a `mod` line in that tier's `mod.rs`. @@ -123,6 +126,28 @@ skips in the JUnit report. active gateway pointed at the deployed OpenShell instance (`openshell gateway add ...` / `openshell gateway select ...`) — the harness shells out to this binary and relies on its persisted config, not on any env var. +- For `tier3::gateway_failover`, a deployment installed with + `deploy/helm/openshell/ci/values-high-availability.yaml`, an external + PostgreSQL Secret, and a configured CLI gateway name. The test discovers two + ready gateway pods and opens its own loopback `oc port-forward` to direct the + initial create and session to one pod. It confirms the session disconnects + when that pod is deleted, then + reconnects through a second port-forward to the pre-existing surviving + replica. This pins the check to that replica and does not exercise the + configured Service or Route endpoint. The workload heartbeat includes a + UUID generated at process start; reconnect must report the same UUID, so a + restarted workload does not satisfy the session check. Set + `OPENSHELL_GATEWAY` to the configured gateway name; it also selects the + gateway used by the image-provenance check that runs with every tier. + `OPENSHELL_ODH_HA_GATEWAY_NAME` remains a compatibility fallback when + `OPENSHELL_GATEWAY` is not set. The test identity needs permission to get + the PostgreSQL Secret referenced by `OPENSHELL_DB_URL`, create pod + port-forwards, and delete gateway pods. The test + deletes a gateway pod only when `OPENSHELL_ODH_HA_FAILOVER=1` is set. Its + default pod selector is the Helm + `app.kubernetes.io/name=openshell,app.kubernetes.io/instance=` + selector; override it with `OPENSHELL_ODH_HA_GATEWAY_SELECTOR` for a renamed + chart deployment. ### The `KUBECONFIG` gotcha @@ -162,6 +187,51 @@ context you happen to have active elsewhere. This means: | `mise run e2e:odh:tier3` | Tier 3: mapped upstream tests + ODH `tier3::` + image provenance | | `cargo nextest run --manifest-path e2e/rust/Cargo.toml --features e2e-odh --test odh -E 'test(=module::test_name)'` | A single ODH test function | +The destructive HA failover scenario is disabled by default. Run it only +against a disposable or explicitly approved HA deployment: + +```bash +OPENSHELL_ODH_HA_FAILOVER=1 mise run e2e:odh:tier3 +``` + +The existing `e2e/with-kube-gateway.sh` wrapper can provision the HA fixture +in an ephemeral namespace: it installs the chart with the HA values overlay, +deploys the PostgreSQL fixture, configures the OpenShift Route and CLI, then +tears down its resources after the test. The Rust test itself assumes an +already-deployed gateway, as do the other ODH tier tests. To run only this +scenario through the wrapper, use a disposable OpenShift context and a built +`target/debug/openshell` CLI. The wrapper's existing-context mode uses one +registry and image tag for gateway, sandbox runtime, and supervisor; deployments +with different image sources for those components still need explicit Helm +overrides or a wrapper extension. + +```bash +OPENSHELL_E2E_KUBE_CONTEXT= \ +OPENSHELL_E2E_KUBE_EXTERNAL_POSTGRES_SECRET=openshell-ha-pg \ +OPENSHELL_E2E_KUBE_EXTRA_VALUES=deploy/helm/openshell/ci/values-high-availability.yaml \ +OPENSHELL_GATEWAY= \ +OPENSHELL_ODH_HA_FAILOVER=1 \ +e2e/with-kube-gateway.sh cargo test --manifest-path e2e/rust/Cargo.toml \ + --features e2e-odh --test odh -- \ + tier3::gateway_failover::gateway_pod_failover_preserves_sandbox_session_and_workspace \ + --exact --nocapture +``` + +Before running failover on an existing deployment, confirm that a normal +`smoke::sandbox::test_create_delete` passes. The failover test bounds sandbox +creation to five minutes and requires the initial session to remain attached +for more than two seconds before deleting any gateway pod. It retries that +preflight with a replacement pod port-forward when `oc port-forward` stalls +the SSH-based `sandbox connect` path. A failure after the retries means +failover has not been exercised. +After sandbox creation, the test awaits a sandbox deletion attempt even when a +later assertion fails. + +The failover scenario explicitly uses the shared E2E workload image, like the +smoke tests. It does not use the chart's default sandbox image: that image can +carry a discovered policy incompatible with the gateway under test, causing +the supervisor to exit before readiness and masking the failover behavior. + Example, running the Smoke tier against a real cluster: ```bash diff --git a/e2e/rust/tests/odh/odh_harness/mod.rs b/e2e/rust/tests/odh/odh_harness/mod.rs index 643594460d..411b0d80a0 100644 --- a/e2e/rust/tests/odh/odh_harness/mod.rs +++ b/e2e/rust/tests/odh/odh_harness/mod.rs @@ -10,4 +10,5 @@ //! collide with the upstream `harness` module. pub mod oc; +pub mod sandbox; pub mod selinux; diff --git a/e2e/rust/tests/odh/odh_harness/oc.rs b/e2e/rust/tests/odh/odh_harness/oc.rs index c8da725fdf..3d28dfcc86 100644 --- a/e2e/rust/tests/odh/odh_harness/oc.rs +++ b/e2e/rust/tests/odh/odh_harness/oc.rs @@ -8,7 +8,27 @@ //! command construction here keeps every test targeting the same cluster and //! reporting failures the same way. +use std::process::Stdio; + use serde_json::Value; +use tokio::io::AsyncWriteExt as _; + +/// Output from an `oc` invocation. +pub struct OcOutput { + pub success: bool, + pub stdout: String, + pub stderr: String, +} + +impl OcOutput { + pub fn contains(&self, needle: &str) -> bool { + self.stdout.contains(needle) || self.stderr.contains(needle) + } + + pub fn diagnostics(&self) -> String { + format!("stdout:\n{}\nstderr:\n{}", self.stdout, self.stderr) + } +} /// Builds an `oc` command targeting the active e2e cluster. /// @@ -27,22 +47,89 @@ pub fn oc_command() -> tokio::process::Command { cmd } +/// Builds a synchronous `oc` command targeting the active e2e cluster. +/// +/// This is intended for best-effort cleanup in `Drop` implementations, where +/// an async process cannot be awaited. +pub fn oc_std_command() -> std::process::Command { + let mut cmd = std::process::Command::new("oc"); + if let Ok(context) = std::env::var("OPENSHELL_E2E_KUBE_CONTEXT_ACTIVE") + && !context.is_empty() + { + cmd.args(["--context", &context]); + } + cmd +} + +/// Runs `oc `, optionally writing `input` to standard input. +/// +/// Returns stdout and stderr even when the command fails, allowing tests to +/// make assertions with the command's diagnostic output. +pub async fn oc(args: &[&str], input: Option<&str>) -> OcOutput { + let mut cmd = oc_command(); + cmd.args(args).stdout(Stdio::piped()).stderr(Stdio::piped()); + if input.is_some() { + cmd.stdin(Stdio::piped()); + } + let mut child = cmd.spawn().expect( + "failed to run `oc` — required for ODH cluster-state checks; ensure it is in PATH \\ + and KUBECONFIG targets the cluster", + ); + if let Some(input) = input { + child + .stdin + .take() + .expect("piped stdin") + .write_all(input.as_bytes()) + .await + .expect("write manifest to oc"); + } + let output = child.wait_with_output().await.expect("wait for oc"); + OcOutput { + success: output.status.success(), + stdout: String::from_utf8_lossy(&output.stdout).into_owned(), + stderr: String::from_utf8_lossy(&output.stderr).into_owned(), + } +} + +/// Returns whether the active cluster exposes the OpenShift Route API. +/// +/// ODH-only tests use this to skip cleanly on non-OpenShift clusters while +/// preserving the standard tier entry points. +pub async fn is_openshift() -> bool { + let output = oc_command() + .args([ + "api-resources", + "--api-group=route.openshift.io", + "--no-headers", + ]) + .output() + .await + .expect( + "failed to run `oc api-resources` — cannot decide whether the cluster is OpenShift; \\ + ensure `oc` is in PATH and KUBECONFIG targets the cluster", + ); + assert!( + output.status.success(), + "oc api-resources failed:\n{}", + String::from_utf8_lossy(&output.stderr) + ); + !output.stdout.is_empty() +} + /// Runs `oc ` and parses stdout as JSON. /// /// Panics with a descriptive message if `oc` cannot be launched, exits /// non-zero, or does not return valid JSON — use it for `-o json` queries /// whose failure should fail the test. pub async fn oc_json(args: &[&str]) -> Value { - let output = oc_command().args(args).output().await.expect( - "failed to run `oc` — required for ODH cluster-state checks; ensure it is in PATH \ - and KUBECONFIG targets the cluster", - ); + let output = oc(args, None).await; assert!( - output.status.success(), - "oc {args:?} failed: {}", - String::from_utf8_lossy(&output.stderr) + output.success, + "oc {args:?} failed:\n{}", + output.diagnostics() ); - serde_json::from_slice(&output.stdout) + serde_json::from_str(&output.stdout) .unwrap_or_else(|e| panic!("oc {args:?} did not return valid JSON: {e}")) } diff --git a/e2e/rust/tests/odh/odh_harness/sandbox.rs b/e2e/rust/tests/odh/odh_harness/sandbox.rs new file mode 100644 index 0000000000..cae71b657a --- /dev/null +++ b/e2e/rust/tests/odh/odh_harness/sandbox.rs @@ -0,0 +1,34 @@ +// SPDX-FileCopyrightText: Copyright (c) 2025-2026 NVIDIA CORPORATION & AFFILIATES. All rights reserved. +// SPDX-License-Identifier: Apache-2.0 + +//! Helpers for resolving OpenShell sandbox resources on an ODH cluster. + +use serde_json::Value; + +use super::oc::oc_json; + +/// Returns the pod selector reported by the Sandbox custom resource. +/// +/// The sandbox-agent controller owns pod creation and does not propagate the +/// OpenShell sandbox-name label to its pod. The custom resource's status is +/// therefore the stable way for downstream tests to discover that pod. +pub async fn sandbox_pod_selector(namespace: &str, sandbox_name: &str) -> Option { + let selector = format!("openshell.ai/sandbox-name={sandbox_name}"); + let sandboxes: Value = oc_json(&[ + "get", + "sandboxes.agents.x-k8s.io", + "-n", + namespace, + "-l", + &selector, + "-o", + "json", + ]) + .await; + sandboxes + .get("items") + .and_then(Value::as_array) + .and_then(|items| items.first()) + .and_then(|sandbox| sandbox["status"]["selector"].as_str()) + .map(str::to_string) +} diff --git a/e2e/rust/tests/odh/tier3/gateway_failover.rs b/e2e/rust/tests/odh/tier3/gateway_failover.rs new file mode 100644 index 0000000000..65135de674 --- /dev/null +++ b/e2e/rust/tests/odh/tier3/gateway_failover.rs @@ -0,0 +1,795 @@ +// SPDX-FileCopyrightText: Copyright (c) 2025-2026 NVIDIA CORPORATION & AFFILIATES. All rights reserved. +// SPDX-License-Identifier: Apache-2.0 + +//! Gateway failover coverage for an external-PostgreSQL HA deployment. +//! +//! This test deletes a gateway pod, so it runs only with +//! `OPENSHELL_ODH_HA_FAILOVER=1`. The test keeps the client's local endpoint +//! stable while switching its pod port-forward from the deleted replica to a +//! surviving replica, then verifies that the original client reconnects. + +use std::net::TcpListener; +use std::panic::AssertUnwindSafe; +use std::process::{Output, Stdio}; +use std::time::{Duration, SystemTime, UNIX_EPOCH}; + +use base64::Engine as _; +use futures_util::FutureExt as _; +use openshell_e2e::harness::binary::openshell_cmd; +use openshell_e2e::harness::sandbox::E2E_WORKLOAD_IMAGE; +use serde_json::Value; +use tokio::io::{AsyncBufReadExt, AsyncReadExt, BufReader}; +use tokio::process::{Child, ChildStdout}; +use tokio::time::{sleep, timeout}; + +use crate::odh_harness::oc::{ + is_openshift, namespace, oc, oc_command, oc_json, pod_is_ready, release, +}; +use crate::odh_harness::sandbox::sandbox_pod_selector; + +// constants +const ENABLE_ENV: &str = "OPENSHELL_ODH_HA_FAILOVER"; +const GATEWAY_SELECTOR_ENV: &str = "OPENSHELL_ODH_HA_GATEWAY_SELECTOR"; +const GATEWAY_NAME_ENV: &str = "OPENSHELL_ODH_HA_GATEWAY_NAME"; +const GATEWAY_PORT: u16 = 8080; +const POD_READY_TIMEOUT: Duration = Duration::from_secs(120); // sets timeouts +const SESSION_TIMEOUT: Duration = Duration::from_secs(60); +const CREATE_TIMEOUT: Duration = Duration::from_secs(300); +const CREATE_RECOVERY_TIMEOUT: Duration = Duration::from_secs(120); +const SESSION_START_ATTEMPTS: usize = 3; +// The CLI only opens its reconnect window after an attachment has survived +// two seconds. Keep a one-second margin so pod deletion cannot race that +// threshold on a slow runner. +const SESSION_ESTABLISHED_DURATION: Duration = Duration::from_secs(3); +const POLL_INTERVAL: Duration = Duration::from_secs(2); +const SENTINEL: &str = "odh-ha-workspace-sentinel"; // marker written to disk to test file persistence +const OUTPUT_MARKER: &str = "odh-ha-session-alive"; // prefix for a process-specific heartbeat used to detect workload restarts + +struct GatewayPod { + name: String, + ready: bool, +} // tracks pod readiness + +struct ActiveSession { + child: Child, + stdout: tokio::io::Lines>, +} // handles running terminal command streams + +struct PodPortForward { + endpoint: String, + port: u16, + child: Child, +} // tracks a local tunnel directly to a single pod + +struct SessionMarker { + process_id: String, + emitted_at: u64, +} + +struct ManagedSandbox { + name: String, + cleaned_up: bool, +} + +impl ManagedSandbox { + fn new(name: String) -> Self { + Self { + name, + cleaned_up: false, + } + } + + async fn cleanup(&mut self, endpoint: &str) -> Result<(), String> { + if self.cleaned_up { + return Ok(()); + } + + let mut command = direct_gateway_command(endpoint); + command + .args(["sandbox", "delete", &self.name]) + .stdout(Stdio::piped()) + .stderr(Stdio::piped()); + let output = command.output().await.map_err(|error| { + format!( + "failed to delete sandbox {} through {endpoint}: {error}", + self.name + ) + })?; + if !output.status.success() { + return Err(format!( + "delete sandbox {} through {endpoint} failed with exit {:?}:\n{}{}", + self.name, + output.status.code(), + String::from_utf8_lossy(&output.stdout), + String::from_utf8_lossy(&output.stderr), + )); + } + + self.cleaned_up = true; + Ok(()) + } +} + +fn gateway_selector(release: &str) -> String { + std::env::var(GATEWAY_SELECTOR_ENV).unwrap_or_else(|_| { + format!("app.kubernetes.io/name=openshell,app.kubernetes.io/instance={release}") + }) +} + +async fn gateway_pods(namespace: &str, selector: &str) -> Vec { + let pods = oc_json(&["get", "pods", "-n", namespace, "-l", selector, "-o", "json"]).await; + pods["items"] + .as_array() + .expect("gateway pod list should contain items") + .iter() + .filter_map(|pod| { + pod["metadata"]["name"].as_str().map(|name| GatewayPod { + name: name.to_string(), + ready: pod_is_ready(pod), + }) + }) + .collect() +} // uses pod_is_ready() on the array of gateway pods + +async fn wait_for_gateway_pod_ready(namespace: &str, selector: &str, pod_name: &str) { + let result = timeout(POD_READY_TIMEOUT, async { + loop { + if gateway_pods(namespace, selector) + .await + .iter() + .any(|pod| pod.name == pod_name && pod.ready) + { + return; + } + sleep(POLL_INTERVAL).await; + } + }) + .await; + assert!( + result.is_ok(), + "gateway pod {pod_name} did not become Ready within {POD_READY_TIMEOUT:?}" + ); +} // continuously polls OpenShift until a specific pre-existing gateway is ready + +fn direct_gateway_command(endpoint: &str) -> tokio::process::Command { + let gateway_name = std::env::var("OPENSHELL_GATEWAY") + .or_else(|_| std::env::var(GATEWAY_NAME_ENV)) + .unwrap_or_else(|_| { + panic!( + "set OPENSHELL_GATEWAY or {GATEWAY_NAME_ENV} to the configured mTLS gateway name" + ) + }); + let mut command = openshell_cmd(); + // The port-forward reaches the gateway directly, so use TLS rather than + // plaintext HTTP. Keep the configured gateway name so the CLI finds its + // mTLS client certificate instead of looking under a port-forward URL. + command + .args(["--gateway", &gateway_name]) + .arg("--gateway-endpoint") + .arg(endpoint); + command +} + +async fn assert_ha_deployment(namespace: &str, selector: &str) { + let deployments = oc_json(&[ + "get", + "deployments", + "-n", + namespace, + "-l", + selector, + "-o", + "json", + ]) + .await; + let items = deployments["items"] + .as_array() + .expect("gateway deployment list should contain items"); + assert_eq!( + items.len(), + 1, + "HA test requires exactly one gateway Deployment selected by '{selector}'" + ); + let deployment = &items[0]; + assert!( + deployment["spec"]["replicas"] + .as_u64() + .is_some_and(|n| n >= 2), + "HA test requires a gateway Deployment configured for at least two replicas" + ); + let containers = deployment["spec"]["template"]["spec"]["containers"] + .as_array() + .expect("gateway Deployment should contain containers"); + let secret_name = containers + .iter() + .find_map(|container| { + container["env"].as_array().and_then(|env| { + env.iter().find_map(|entry| { + (entry["name"].as_str() == Some("OPENSHELL_DB_URL") + && entry["valueFrom"]["secretKeyRef"]["key"].as_str() == Some("uri")) + .then(|| entry["valueFrom"]["secretKeyRef"]["name"].as_str()) + .flatten() + .filter(|name| !name.is_empty()) + }) + }) + }) + .expect("HA test requires OPENSHELL_DB_URL to reference a nonempty uri key in a Secret"); + let secret = oc_json(&["get", "secret", secret_name, "-n", namespace, "-o", "json"]).await; + let encoded_url = secret["data"]["uri"] + .as_str() + .expect("HA test requires the referenced Secret to contain data.uri"); + let url = base64::engine::general_purpose::STANDARD + .decode(encoded_url) + .expect("HA test requires the referenced Secret data.uri to be valid base64"); + let url = std::str::from_utf8(&url) + .expect("HA test requires the referenced Secret data.uri to be valid UTF-8"); + assert!( + url.starts_with("postgres://") || url.starts_with("postgresql://"), + "HA test requires the referenced Secret data.uri to contain a PostgreSQL URL" + ); +} + +fn reserve_loopback_port() -> u16 { + let listener = TcpListener::bind("127.0.0.1:0").expect("reserve a loopback port"); + listener.local_addr().expect("read reserved port").port() +} + +async fn port_forward_gateway_pod(namespace: &str, pod: &str, port: u16) -> PodPortForward { + let target = format!("{port}:{GATEWAY_PORT}"); + let mut command = oc_command(); + command + .args([ + "port-forward", + "-n", + namespace, + &format!("pod/{pod}"), + &target, + "--address=127.0.0.1", + ]) + .stdout(Stdio::null()) + .stderr(Stdio::null()) + .kill_on_drop(true); + let mut child = command.spawn().expect("start gateway pod port-forward"); + let result = timeout(SESSION_TIMEOUT, async { + loop { + if tokio::net::TcpStream::connect(("127.0.0.1", port)) + .await + .is_ok() + { + return; + } + if let Some(status) = child.try_wait().expect("poll gateway pod port-forward") { + panic!("gateway pod port-forward exited before becoming ready: {status}"); + } + sleep(POLL_INTERVAL).await; + } + }) + .await; + assert!( + result.is_ok(), + "gateway pod port-forward did not become reachable within {SESSION_TIMEOUT:?}" + ); + PodPortForward { + // The test deployment includes localhost in the server certificate + // SANs. `localhost` therefore keeps certificate validation intact + // while the port-forward still targets exactly one gateway pod. + endpoint: format!("https://localhost:{port}"), + port, + child, + } +} // bypasses the default load balancer and targets one gateway pod + +async fn stop_port_forward(forward: &mut PodPortForward) { + let _ = forward.child.kill().await; + let _ = forward.child.wait().await; +} + +async fn wait_for_created_sandbox(name: &str, endpoint: &str) -> Result<(), String> { + timeout(CREATE_RECOVERY_TIMEOUT, async { + loop { + let mut command = direct_gateway_command(endpoint); + command + .args(["sandbox", "get", name, "--output", "json"]) + .stdout(Stdio::piped()) + .stderr(Stdio::piped()); + let output = command + .output() + .await + .map_err(|error| format!("failed to check sandbox {name} after create: {error}"))?; + if output.status.success() { + let sandbox: Value = serde_json::from_slice(&output.stdout).map_err(|error| { + format!("sandbox get {name} returned invalid JSON after create: {error}") + })?; + match sandbox["phase"].as_str() { + Some("Ready") => return Ok(()), + Some("Error") => { + return Err(format!( + "sandbox {name} entered Error after the create connection dropped: {sandbox}" + )); + } + _ => {} + } + } + sleep(POLL_INTERVAL).await; + } + }) + .await + .map_err(|_| format!("sandbox {name} did not become Ready within {CREATE_RECOVERY_TIMEOUT:?}"))? +} + +async fn create_on_initial_gateway(endpoint: &str, script: &str) -> ManagedSandbox { + // Six random hex characters, together with the hexadecimal process ID, + // keep this below the 19-character sandbox-name limit while separating + // independent CI runs that happen to reuse a process ID. + let name = format!( + "ha-{:x}-{:06x}", + std::process::id(), + rand::random::() & 0x00ff_ffff + ); + let mut command = direct_gateway_command(endpoint); + // A failed create may still leave a retained sandbox behind. Register its + // cleanup before running the CLI so failure paths attempt to remove it. + let mut sandbox = ManagedSandbox::new(name.clone()); + command + .args([ + "sandbox", + "create", + "--detach", + "--name", + &name, + "--from", + E2E_WORKLOAD_IMAGE, + "--", + ]) + .args(["sh", "-lc", script]) + .stdout(Stdio::piped()) + .stderr(Stdio::piped()) + .kill_on_drop(true); + let mut child = command.spawn().expect("spawn sandbox create"); + let mut stdout = child + .stdout + .take() + .expect("sandbox create stdout must be piped"); + let mut stderr = child + .stderr + .take() + .expect("sandbox create stderr must be piped"); + let stdout_reader = tokio::spawn(async move { + let mut output = Vec::new(); + stdout.read_to_end(&mut output).await.map(|_| output) + }); + let stderr_reader = tokio::spawn(async move { + let mut output = Vec::new(); + stderr.read_to_end(&mut output).await.map(|_| output) + }); + let status = match timeout(CREATE_TIMEOUT, child.wait()).await { + Ok(Ok(status)) => status, + Ok(Err(error)) => { + let _ = stdout_reader.await; + let _ = stderr_reader.await; + if let Err(cleanup_error) = sandbox.cleanup(endpoint).await { + eprintln!("{cleanup_error}"); + } + panic!("failed to run sandbox create through the initial gateway: {error}"); + } + Err(_) => { + // Retain the child handle through timeout so it is explicitly + // killed and reaped before cleanup can race a still-running CLI. + let _ = child.kill().await; + let _ = child.wait().await; + let _ = stdout_reader.await; + let _ = stderr_reader.await; + if let Err(cleanup_error) = sandbox.cleanup(endpoint).await { + eprintln!("{cleanup_error}"); + } + panic!("sandbox create timed out after {CREATE_TIMEOUT:?}; failover was not exercised"); + } + }; + let output = Output { + status, + stdout: stdout_reader + .await + .expect("sandbox create stdout reader should not panic") + .expect("read sandbox create stdout"), + stderr: stderr_reader + .await + .expect("sandbox create stderr reader should not panic") + .expect("read sandbox create stderr"), + }; + if !output.status.success() { + let diagnostic = format!( + "create sandbox through initial gateway failed:\n{}{}", + String::from_utf8_lossy(&output.stdout), + String::from_utf8_lossy(&output.stderr), + ); + let stdout = String::from_utf8_lossy(&output.stdout); + let stderr = String::from_utf8_lossy(&output.stderr); + if stdout.contains(&format!("Created sandbox: {name}")) + && stderr.contains("peer closed connection without sending TLS close_notify") + { + match wait_for_created_sandbox(&name, endpoint).await { + Ok(()) => { + eprintln!( + "create watch through pod port-forward disconnected; sandbox {name} became Ready through the initial gateway pod" + ); + return sandbox; + } + Err(error) => { + if let Err(cleanup_error) = sandbox.cleanup(endpoint).await { + eprintln!("{cleanup_error}"); + } + panic!("{diagnostic}\nRecovery through initial gateway pod failed: {error}"); + } + } + } + if let Err(cleanup_error) = sandbox.cleanup(endpoint).await { + eprintln!("{cleanup_error}"); + } + panic!("{diagnostic}"); + } + sandbox +} // create a sandbox via replica-0 direct endpoint, execute a background script that writes sentinel text to /sandbox/.odh-ha-sentinel +// and enters an infinite loop printing OUTPUT_MARKER every second + +async fn start_session(sandbox_name: &str, gateway_endpoint: Option<&str>) -> ActiveSession { + let mut command = gateway_endpoint.map_or_else(openshell_cmd, direct_gateway_command); + command + .args(["sandbox", "connect", sandbox_name]) + .stdin(Stdio::null()) + .stdout(Stdio::piped()) + .stderr(Stdio::piped()) + .kill_on_drop(true); + let mut child = command.spawn().expect("spawn sandbox connect"); + let stdout = child + .stdout + .take() + .expect("sandbox connect stdout must be piped"); + ActiveSession { + child, + stdout: BufReader::new(stdout).lines(), + } +} // opens a streaming terminal session (openshell sandbox connect) + +async fn exec_on_gateway( + endpoint: &str, + sandbox_name: &str, + argv: &[&str], +) -> Result { + let mut command = direct_gateway_command(endpoint); + command + .args(["sandbox", "exec", "--name", sandbox_name, "--no-tty", "--"]) + .args(argv) + .stdout(Stdio::piped()) + .stderr(Stdio::piped()); + let output = command + .output() + .await + .map_err(|error| format!("failed to run sandbox exec through {endpoint}: {error}"))?; + let combined = format!( + "{}{}", + String::from_utf8_lossy(&output.stdout), + String::from_utf8_lossy(&output.stderr) + ); + if !output.status.success() { + return Err(format!( + "sandbox exec through {endpoint} failed with exit {:?}:\n{combined}", + output.status.code() + )); + } + Ok(combined) +} + +fn session_marker(line: &str) -> Option { + let (prefix_and_process_id, emitted_at) = line.rsplit_once(':')?; + let process_id = prefix_and_process_id.strip_prefix(&format!("{OUTPUT_MARKER}:"))?; + Some(SessionMarker { + process_id: process_id.to_string(), + emitted_at: emitted_at.parse().ok()?, + }) +} + +async fn session_marker_after( + session: &mut ActiveSession, + strictly_after: Option, +) -> Result { + timeout(SESSION_TIMEOUT, async { + loop { + let line = session + .stdout + .next_line() + .await + .map_err(|error| format!("cannot read sandbox connect output: {error}"))? + .ok_or_else(|| { + "sandbox connect exited before receiving session output".to_string() + })?; + if let Some(marker) = session_marker(&line) + && strictly_after.is_none_or(|minimum| marker.emitted_at > minimum) + { + return Ok(marker); + } + } + }) + .await + .map_err(|_| format!("did not receive '{OUTPUT_MARKER}' within {SESSION_TIMEOUT:?}"))? +} // returns a heartbeat emitted after the requested time, excluding buffered output + +async fn start_initial_session( + namespace: &str, + initial_pod: &str, + forward: &mut PodPortForward, + sandbox_name: &str, +) -> (ActiveSession, SessionMarker) { + let mut failures = Vec::new(); + for attempt in 1..=SESSION_START_ATTEMPTS { + let mut session = start_session(sandbox_name, Some(&forward.endpoint)).await; + match session_marker_after(&mut session, None).await { + Ok(_) => match session_marker_after( + &mut session, + Some(unix_timestamp() + SESSION_ESTABLISHED_DURATION.as_secs()), + ) + .await + { + Ok(marker) => return (session, marker), + Err(error) => { + failures.push(format!( + "attempt {attempt}: session did not remain attached for longer than {SESSION_ESTABLISHED_DURATION:?}: {error}" + )); + stop_session(&mut session).await; + if attempt < SESSION_START_ATTEMPTS { + stop_port_forward(forward).await; + *forward = + port_forward_gateway_pod(namespace, initial_pod, forward.port).await; + } + } + }, + Err(error) => { + failures.push(format!("attempt {attempt}: {error}")); + stop_session(&mut session).await; + if attempt < SESSION_START_ATTEMPTS { + // `oc port-forward` can accept TCP connections while its + // SSH stream is stalled. Recreate the tunnel before + // retrying so this remains a transport preflight, before + // the destructive failover action. + stop_port_forward(forward).await; + *forward = port_forward_gateway_pod(namespace, initial_pod, forward.port).await; + } + } + } + } + panic!( + "initial session through pod port-forward could not produce a heartbeat after {SESSION_START_ATTEMPTS} attempts; failover was not exercised:\n{}", + failures.join("\n") + ); +} + +fn unix_timestamp() -> u64 { + SystemTime::now() + .duration_since(UNIX_EPOCH) + .expect("system clock should be after the Unix epoch") + .as_secs() +} + +async fn stop_session(session: &mut ActiveSession) { + let _ = session.child.kill().await; + let _ = session.child.wait().await; +} + +async fn assert_sandbox_workload_is_live(namespace: &str, sandbox_name: &str) -> String { + let selector = sandbox_pod_selector(namespace, sandbox_name) + .await + .expect("sandbox CR should retain its workload selector after failover"); + let pods = oc_json(&[ + "get", "pods", "-n", namespace, "-l", &selector, "-o", "json", + ]) + .await; + let items = pods["items"] + .as_array() + .expect("sandbox pod list should contain items"); + assert!( + !items.is_empty() && items.iter().all(pod_is_ready), + "sandbox workload is orphaned or not ready after gateway failover: {pods}" + ); + selector +} // checks k8s to make sure the underlying container running the workload didn't crash or get orphaned during failover + +async fn assert_sandbox_deleted(endpoint: &str, namespace: &str, name: &str, pod_selector: &str) { + let cr_selector = format!("openshell.ai/sandbox-name={name}"); + let result = timeout(CREATE_RECOVERY_TIMEOUT, async { + loop { + let mut command = direct_gateway_command(endpoint); + command.args(["sandbox", "list", "--names"]); + let output = command + .output() + .await + .expect("list sandboxes after cleanup"); + assert!( + output.status.success(), + "could not verify sandbox deletion through the gateway: {}", + String::from_utf8_lossy(&output.stderr) + ); + let still_listed = String::from_utf8_lossy(&output.stdout) + .lines() + .any(|line| line.trim() == name); + let crs = oc_json(&[ + "get", + "sandboxes.agents.x-k8s.io", + "-n", + namespace, + "-l", + &cr_selector, + "-o", + "json", + ]) + .await; + let pods = oc_json(&[ + "get", + "pods", + "-n", + namespace, + "-l", + pod_selector, + "-o", + "json", + ]) + .await; + let cr_count = crs["items"] + .as_array() + .expect("Sandbox CR list items") + .len(); + let pod_count = pods["items"] + .as_array() + .expect("sandbox pod list items") + .len(); + if !still_listed && cr_count == 0 && pod_count == 0 { + return; + } + sleep(POLL_INTERVAL).await; + } + }) + .await; + assert!( + result.is_ok(), + "sandbox {name}, its custom resource, or its pod remained after cleanup for {CREATE_RECOVERY_TIMEOUT:?}" + ); +} + +async fn run_failover_scenario( + namespace: &str, + selector: &str, + initial_pod: &str, + surviving_pod: &str, + initial_forward: &mut PodPortForward, + sandbox: &ManagedSandbox, +) -> String { + let (mut initial_session, initial_marker) = + start_initial_session(namespace, initial_pod, initial_forward, &sandbox.name).await; + + let deleted = oc( + &["delete", "pod", initial_pod, "-n", namespace, "--wait=true"], + None, + ) + .await; + assert!( + deleted.success, + "delete gateway pod {initial_pod} failed:\n{}", + deleted.diagnostics() + ); + + // Keep the endpoint stable while replacing its direct backend. This lets + // the original client exercise its bounded reconnect behavior. + stop_port_forward(initial_forward).await; + wait_for_gateway_pod_ready(namespace, selector, surviving_pod).await; + let reconnect_started_at = unix_timestamp(); + let mut reconnect_forward = + port_forward_gateway_pod(namespace, surviving_pod, initial_forward.port).await; + let reconnected_marker = session_marker_after(&mut initial_session, Some(reconnect_started_at)) + .await + .unwrap_or_else(|error| { + panic!("session did not reconnect through the surviving gateway pod: {error}") + }); + assert_eq!( + reconnected_marker.process_id, initial_marker.process_id, + "sandbox workload process changed while the original client reconnected after gateway failover" + ); + stop_session(&mut initial_session).await; + + let sentinel = exec_on_gateway( + &reconnect_forward.endpoint, + &sandbox.name, + &["cat", "/sandbox/.odh-ha-sentinel"], + ) + .await + .expect("sandbox should remain executable after gateway failover"); + assert!( + sentinel.contains(SENTINEL), + "workspace sentinel was lost after gateway failover: {sentinel}" + ); + stop_port_forward(&mut reconnect_forward).await; + assert_sandbox_workload_is_live(namespace, &sandbox.name).await +} + +#[tokio::test] +async fn gateway_pod_failover_preserves_sandbox_session_and_workspace() { + // preflight checks + if std::env::var(ENABLE_ENV).as_deref() != Ok("1") { + eprintln!("skipping destructive HA failover test; set {ENABLE_ENV}=1 to enable it"); + return; + } + if !is_openshift().await { + eprintln!("skipping HA failover test; the active cluster is not OpenShift"); + return; + } + + // replica discovery + let namespace = namespace(); + let release = release(); + let selector = gateway_selector(&release); + assert_ha_deployment(&namespace, &selector).await; + let pods = gateway_pods(&namespace, &selector).await; // scan for running gateways + assert!( + pods.len() >= 2 && pods.iter().filter(|pod| pod.ready).count() >= 2, + "HA test requires two ready gateway pods selected by '{selector}', found: {:?}", + pods.iter() + .map(|pod| (&pod.name, pod.ready)) + .collect::>(), + ); // asserts at least 2 ready replicas exist + let initial_pod = pods + .iter() + .find(|pod| pod.ready) + .expect("one ready gateway pod") + .name + .clone(); // pick an initial_pod + let surviving_pod = pods + .iter() + .find(|pod| pod.ready && pod.name != initial_pod) + .expect("a second pre-existing ready gateway pod") + .name + .clone(); + let port = reserve_loopback_port(); + let mut initial_forward = port_forward_gateway_pod(&namespace, &initial_pod, port).await; + // Retain an independent connection to the replica that will survive the + // failure. It is deliberately separate from the stable client endpoint so + // cleanup and its verification never fall back to a Service or Route. + let cleanup_port = reserve_loopback_port(); + let mut surviving_forward = + port_forward_gateway_pod(&namespace, &surviving_pod, cleanup_port).await; + + let script = format!( + "printf '%s\\n' '{SENTINEL}' > /sandbox/.odh-ha-sentinel; session_id=$(cat /proc/sys/kernel/random/uuid); while :; do printf '%s:%s:%s\\n' '{OUTPUT_MARKER}' \"$session_id\" \"$(date +%s)\"; sleep 1; done" + ); // persist a workspace sentinel and stream a process-specific heartbeat + let mut sandbox = create_on_initial_gateway(&initial_forward.endpoint, &script).await; // make sandbox w/sentinel file + // Keep ownership of the sandbox here so every assertion failure still + // reaches awaited cleanup before the test unwinds. + let scenario = AssertUnwindSafe(run_failover_scenario( + &namespace, + &selector, + &initial_pod, + &surviving_pod, + &mut initial_forward, + &sandbox, + )) + .catch_unwind() + .await; + + let cleanup = sandbox.cleanup(&surviving_forward.endpoint).await; + match scenario { + Ok(pod_selector) => { + cleanup.expect("delete sandbox through the surviving gateway pod"); + assert_sandbox_deleted( + &surviving_forward.endpoint, + &namespace, + &sandbox.name, + &pod_selector, + ) + .await; + } + Err(failure) => { + if let Err(cleanup_error) = cleanup { + eprintln!("{cleanup_error}"); + } + stop_port_forward(&mut surviving_forward).await; + std::panic::resume_unwind(failure); + } + } + stop_port_forward(&mut surviving_forward).await; +} diff --git a/e2e/rust/tests/odh/tier3/mod.rs b/e2e/rust/tests/odh/tier3/mod.rs index 9b74f53bd2..c41c4433ef 100644 --- a/e2e/rust/tests/odh/tier3/mod.rs +++ b/e2e/rust/tests/odh/tier3/mod.rs @@ -2,3 +2,5 @@ // SPDX-License-Identifier: Apache-2.0 //! Tier 3: negative and destructive tests (no time limit). + +mod gateway_failover; From d7c7a018872cd114a45e2de5246f55ddb5ec5275 Mon Sep 17 00:00:00 2001 From: Evie Howard Date: Thu, 1 Oct 2026 12:10:32 +0100 Subject: [PATCH 2/3] CARRY: test(odh) add reusable OpenShift test helpers Signed-off-by: Evie Howard --- e2e/rust/tests/odh/odh_harness/oc.rs | 77 +++++++++++++++----- e2e/rust/tests/odh/smoke/image_provenance.rs | 6 +- 2 files changed, 61 insertions(+), 22 deletions(-) diff --git a/e2e/rust/tests/odh/odh_harness/oc.rs b/e2e/rust/tests/odh/odh_harness/oc.rs index 3d28dfcc86..b09c40befa 100644 --- a/e2e/rust/tests/odh/odh_harness/oc.rs +++ b/e2e/rust/tests/odh/odh_harness/oc.rs @@ -13,6 +13,40 @@ use std::process::Stdio; use serde_json::Value; use tokio::io::AsyncWriteExt as _; +const DEFAULT_DEPLOYMENT_NAME: &str = "openshell"; + +/// Resolves the namespace used by the ODH deployment under test. +/// +/// The explicit sandbox namespace takes precedence for compatibility with +/// existing ODH deployment environments. The remaining names match the +/// standard ODH and Kubernetes e2e configuration variables. +pub fn namespace() -> String { + std::env::var("SANDBOX_NAMESPACE") + .or_else(|_| std::env::var("NAMESPACE")) + .or_else(|_| std::env::var("OPENSHELL_E2E_KUBE_NAMESPACE")) + .unwrap_or_else(|_| DEFAULT_DEPLOYMENT_NAME.to_string()) +} + +/// Resolves the Helm release name used by the ODH deployment under test. +pub fn release() -> String { + std::env::var("RELEASE") + .or_else(|_| std::env::var("OPENSHELL_E2E_KUBE_RELEASE")) + .unwrap_or_else(|_| DEFAULT_DEPLOYMENT_NAME.to_string()) +} + +/// Returns whether a Kubernetes Pod JSON object is Running and Ready. +pub fn pod_is_ready(pod: &Value) -> bool { + pod["status"]["phase"].as_str() == Some("Running") + && pod["status"]["conditions"] + .as_array() + .is_some_and(|conditions| { + conditions.iter().any(|condition| { + condition["type"].as_str() == Some("Ready") + && condition["status"].as_str() == Some("True") + }) + }) +} + /// Output from an `oc` invocation. pub struct OcOutput { pub success: bool, @@ -21,10 +55,6 @@ pub struct OcOutput { } impl OcOutput { - pub fn contains(&self, needle: &str) -> bool { - self.stdout.contains(needle) || self.stderr.contains(needle) - } - pub fn diagnostics(&self) -> String { format!("stdout:\n{}\nstderr:\n{}", self.stdout, self.stderr) } @@ -47,20 +77,6 @@ pub fn oc_command() -> tokio::process::Command { cmd } -/// Builds a synchronous `oc` command targeting the active e2e cluster. -/// -/// This is intended for best-effort cleanup in `Drop` implementations, where -/// an async process cannot be awaited. -pub fn oc_std_command() -> std::process::Command { - let mut cmd = std::process::Command::new("oc"); - if let Ok(context) = std::env::var("OPENSHELL_E2E_KUBE_CONTEXT_ACTIVE") - && !context.is_empty() - { - cmd.args(["--context", &context]); - } - cmd -} - /// Runs `oc `, optionally writing `input` to standard input. /// /// Returns stdout and stderr even when the command fails, allowing tests to @@ -290,9 +306,32 @@ mod tests { use serde_json::json; use super::{ - pod_node_and_uid, pod_uid_cgroup_form, sandbox_id_from_json, supervisor_pod_from_json, + pod_is_ready, pod_node_and_uid, pod_uid_cgroup_form, sandbox_id_from_json, + supervisor_pod_from_json, }; + #[test] + fn recognizes_only_running_ready_pods() { + assert!(pod_is_ready(&json!({ + "status": { + "phase": "Running", + "conditions": [{"type": "Ready", "status": "True"}] + } + }))); + assert!(!pod_is_ready(&json!({ + "status": { + "phase": "Pending", + "conditions": [{"type": "Ready", "status": "True"}] + } + }))); + assert!(!pod_is_ready(&json!({ + "status": { + "phase": "Running", + "conditions": [{"type": "Ready", "status": "False"}] + } + }))); + } + #[test] fn resolves_sandbox_id_from_named_sandbox() { let sandbox = json!({ diff --git a/e2e/rust/tests/odh/smoke/image_provenance.rs b/e2e/rust/tests/odh/smoke/image_provenance.rs index 8ec46b1d5f..7de6ac9a53 100644 --- a/e2e/rust/tests/odh/smoke/image_provenance.rs +++ b/e2e/rust/tests/odh/smoke/image_provenance.rs @@ -18,7 +18,7 @@ use serde_json::Value; use openshell_e2e::harness::sandbox::SandboxGuard; -use crate::odh_harness::oc::{oc_command, oc_json}; +use crate::odh_harness::oc::{namespace, oc_command, oc_json, release}; fn allowed_prefixes() -> Vec { std::env::var("ALLOWED_IMAGE_REGISTRY_PREFIXES") @@ -108,8 +108,8 @@ fn parse_supervisor_image(gateway_toml: &str) -> Option { #[tokio::test] async fn test_sandbox_gateway_supervisor_images() { - let namespace = std::env::var("NAMESPACE").unwrap_or_else(|_| "openshell".to_string()); - let release = std::env::var("RELEASE").unwrap_or_else(|_| "openshell".to_string()); + let namespace = namespace(); + let release = release(); let allowed = allowed_prefixes(); let mut errors = Vec::new(); From d9bd8f1855f2fea3b76882d23282a6ffc02cd50a Mon Sep 17 00:00:00 2001 From: Evie Howard Date: Thu, 1 Oct 2026 14:38:21 +0100 Subject: [PATCH 3/3] CARRY: test(odh): add test coverage for when sandbox and gateway namespaces differ Signed-off-by: Evie Howard --- e2e/rust/tests/odh/README.md | 7 +- e2e/rust/tests/odh/odh_harness/oc.rs | 20 ++++-- e2e/rust/tests/odh/smoke/image_provenance.rs | 30 +++------ e2e/rust/tests/odh/tier1/selinux.rs | 4 +- e2e/rust/tests/odh/tier3/gateway_failover.rs | 67 ++++++++++++++------ 5 files changed, 78 insertions(+), 50 deletions(-) diff --git a/e2e/rust/tests/odh/README.md b/e2e/rust/tests/odh/README.md index a2e9cdec87..3fe392157c 100644 --- a/e2e/rust/tests/odh/README.md +++ b/e2e/rust/tests/odh/README.md @@ -248,7 +248,9 @@ see below. For RHOAI images, use too if your deployment doesn't use the defaults (`openshell`/`openshell`). The Quay deployment script sets `sandbox.image.pullPolicy=IfNotPresent`; other deployments must -configure it themselves. +configure it themselves. Set `SANDBOX_NAMESPACE` only when Sandbox +custom resources and workload Pods run in a namespace separate from the +gateway. ### SELinux-enforcing OCP validation @@ -317,6 +319,9 @@ ALLOWED_IMAGE_REGISTRY_PREFIXES="quay.io/opendatahub/,nvcr.io/nvidia/base/" \ otherwise match a lookalike host (e.g. `registry.redhat.io` would also match `registry.redhat.io.attacker.example/image`). - `NAMESPACE`/`RELEASE` env vars default to `openshell`/`openshell`. + `SANDBOX_NAMESPACE` defaults to the resolved gateway namespace and selects + Sandbox custom resources and workload Pods; `NAMESPACE` selects gateway + resources. - The check is a registry-prefix allowlist, not an exact image/digest match. It checks each observed image against the configured allowed prefixes. - The `imagePullPolicy: IfNotPresent` check applies to every container, diff --git a/e2e/rust/tests/odh/odh_harness/oc.rs b/e2e/rust/tests/odh/odh_harness/oc.rs index b09c40befa..9f4c5acbf0 100644 --- a/e2e/rust/tests/odh/odh_harness/oc.rs +++ b/e2e/rust/tests/odh/odh_harness/oc.rs @@ -15,18 +15,24 @@ use tokio::io::AsyncWriteExt as _; const DEFAULT_DEPLOYMENT_NAME: &str = "openshell"; -/// Resolves the namespace used by the ODH deployment under test. +/// Resolves the namespace containing gateway resources for the ODH deployment. /// -/// The explicit sandbox namespace takes precedence for compatibility with -/// existing ODH deployment environments. The remaining names match the -/// standard ODH and Kubernetes e2e configuration variables. -pub fn namespace() -> String { - std::env::var("SANDBOX_NAMESPACE") - .or_else(|_| std::env::var("NAMESPACE")) +/// The names match the standard ODH and Kubernetes e2e configuration +/// variables. +pub fn gateway_namespace() -> String { + std::env::var("NAMESPACE") .or_else(|_| std::env::var("OPENSHELL_E2E_KUBE_NAMESPACE")) .unwrap_or_else(|_| DEFAULT_DEPLOYMENT_NAME.to_string()) } +/// Resolves the namespace containing Sandbox custom resources and workload Pods. +/// +/// Sandbox resources normally share the gateway namespace. Set +/// `SANDBOX_NAMESPACE` when the compute driver uses a separate namespace. +pub fn sandbox_namespace() -> String { + std::env::var("SANDBOX_NAMESPACE").unwrap_or_else(|_| gateway_namespace()) +} + /// Resolves the Helm release name used by the ODH deployment under test. pub fn release() -> String { std::env::var("RELEASE") diff --git a/e2e/rust/tests/odh/smoke/image_provenance.rs b/e2e/rust/tests/odh/smoke/image_provenance.rs index 7de6ac9a53..46090e6d2a 100644 --- a/e2e/rust/tests/odh/smoke/image_provenance.rs +++ b/e2e/rust/tests/odh/smoke/image_provenance.rs @@ -18,7 +18,8 @@ use serde_json::Value; use openshell_e2e::harness::sandbox::SandboxGuard; -use crate::odh_harness::oc::{namespace, oc_command, oc_json, release}; +use crate::odh_harness::oc::{gateway_namespace, oc_command, oc_json, release, sandbox_namespace}; +use crate::odh_harness::sandbox::sandbox_pod_selector; fn allowed_prefixes() -> Vec { std::env::var("ALLOWED_IMAGE_REGISTRY_PREFIXES") @@ -108,7 +109,8 @@ fn parse_supervisor_image(gateway_toml: &str) -> Option { #[tokio::test] async fn test_sandbox_gateway_supervisor_images() { - let namespace = namespace(); + let gateway_namespace = gateway_namespace(); + let sandbox_namespace = sandbox_namespace(); let release = release(); let allowed = allowed_prefixes(); @@ -143,27 +145,11 @@ async fn test_sandbox_gateway_supervisor_images() { // agents.x-k8s.io/sandbox-name-hash label. So look up the CR first and // follow its reported `status.selector` to find the pod. let sandbox_selector = format!("openshell.ai/sandbox-name={}", sb.name); - let sandbox_crs = oc_json(&[ - "get", - "sandboxes.agents.x-k8s.io", - "-n", - &namespace, - "-l", - &sandbox_selector, - "-o", - "json", - ]) - .await; - let pod_selector = sandbox_crs - .get("items") - .and_then(Value::as_array) - .and_then(|items| items.first()) - .and_then(|cr| cr["status"]["selector"].as_str()) - .map(str::to_string); + let pod_selector = sandbox_pod_selector(&sandbox_namespace, &sb.name).await; match pod_selector { Some(pod_selector) => { - let sandbox_pods = oc_json(&["get", "pods", "-n", &namespace, "-l", &pod_selector, "-o", "json"]).await; + let sandbox_pods = oc_json(&["get", "pods", "-n", &sandbox_namespace, "-l", &pod_selector, "-o", "json"]).await; collect_pod_images(&sandbox_pods, &pod_selector, &mut errors, &mut images); } None => errors.push(format!( @@ -177,7 +163,7 @@ async fn test_sandbox_gateway_supervisor_images() { "get", "pods", "-n", - &namespace, + &gateway_namespace, "-l", &gateway_selector, "-o", @@ -196,7 +182,7 @@ async fn test_sandbox_gateway_supervisor_images() { "configmap", &cm_name, "-n", - &namespace, + &gateway_namespace, "-o", "jsonpath={.data.gateway\\.toml}", ]) diff --git a/e2e/rust/tests/odh/tier1/selinux.rs b/e2e/rust/tests/odh/tier1/selinux.rs index 6f5159b84b..eba372eff8 100644 --- a/e2e/rust/tests/odh/tier1/selinux.rs +++ b/e2e/rust/tests/odh/tier1/selinux.rs @@ -9,7 +9,7 @@ use openshell_e2e::harness::sandbox::SandboxGuard; use serial_test::serial; use tempfile::NamedTempFile; -use crate::odh_harness::oc::{paired_supervisor_pod, supervisor_selinux_label}; +use crate::odh_harness::oc::{paired_supervisor_pod, sandbox_namespace, supervisor_selinux_label}; use crate::odh_harness::selinux::SelinuxAudit; /// OCP's container `SELinux` type. The complete label includes per-pod MLS/MCS @@ -58,7 +58,7 @@ async fn selinux_is_enforcing_on_all_worker_nodes() { async fn supervisor_runs_in_container_selinux_domain() { run_audited("supervisor label", || async { let mut sandbox = SandboxGuard::create(&[]).await?; - let namespace = std::env::var("NAMESPACE").unwrap_or_else(|_| "openshell".to_string()); + let namespace = sandbox_namespace(); let supervisor_pod = paired_supervisor_pod(&namespace, &sandbox.name).await?; let label_result = supervisor_selinux_label(&namespace, &supervisor_pod).await; sandbox.cleanup().await; diff --git a/e2e/rust/tests/odh/tier3/gateway_failover.rs b/e2e/rust/tests/odh/tier3/gateway_failover.rs index 65135de674..f6f29687cc 100644 --- a/e2e/rust/tests/odh/tier3/gateway_failover.rs +++ b/e2e/rust/tests/odh/tier3/gateway_failover.rs @@ -23,7 +23,8 @@ use tokio::process::{Child, ChildStdout}; use tokio::time::{sleep, timeout}; use crate::odh_harness::oc::{ - is_openshift, namespace, oc, oc_command, oc_json, pod_is_ready, release, + gateway_namespace, is_openshift, oc, oc_command, oc_json, pod_is_ready, release, + sandbox_namespace, }; use crate::odh_harness::sandbox::sandbox_pod_selector; @@ -36,6 +37,7 @@ const POD_READY_TIMEOUT: Duration = Duration::from_secs(120); // sets timeouts const SESSION_TIMEOUT: Duration = Duration::from_secs(60); const CREATE_TIMEOUT: Duration = Duration::from_secs(300); const CREATE_RECOVERY_TIMEOUT: Duration = Duration::from_secs(120); +const SANDBOX_DELETION_TIMEOUT: Duration = Duration::from_secs(300); const SESSION_START_ATTEMPTS: usize = 3; // The CLI only opens its reconnect window after an attachment has survived // two seconds. Keep a one-second margin so pod deletion cannot race that @@ -594,7 +596,8 @@ async fn assert_sandbox_workload_is_live(namespace: &str, sandbox_name: &str) -> async fn assert_sandbox_deleted(endpoint: &str, namespace: &str, name: &str, pod_selector: &str) { let cr_selector = format!("openshell.ai/sandbox-name={name}"); - let result = timeout(CREATE_RECOVERY_TIMEOUT, async { + let mut last_observation = None; + let result = timeout(SANDBOX_DELETION_TIMEOUT, async { loop { let mut command = direct_gateway_command(endpoint); command.args(["sandbox", "list", "--names"]); @@ -640,6 +643,9 @@ async fn assert_sandbox_deleted(endpoint: &str, namespace: &str, name: &str, pod .as_array() .expect("sandbox pod list items") .len(); + last_observation = Some(format!( + "sandbox listed: {still_listed}, custom resources: {cr_count}, workload pods: {pod_count}" + )); if !still_listed && cr_count == 0 && pod_count == 0 { return; } @@ -649,23 +655,37 @@ async fn assert_sandbox_deleted(endpoint: &str, namespace: &str, name: &str, pod .await; assert!( result.is_ok(), - "sandbox {name}, its custom resource, or its pod remained after cleanup for {CREATE_RECOVERY_TIMEOUT:?}" + "sandbox {name}, its custom resource, or its pod remained after cleanup for {SANDBOX_DELETION_TIMEOUT:?}; last observation: {}", + last_observation.unwrap_or_else(|| "no deletion state observed".to_string()) ); } async fn run_failover_scenario( - namespace: &str, + gateway_namespace: &str, + sandbox_namespace: &str, selector: &str, initial_pod: &str, surviving_pod: &str, initial_forward: &mut PodPortForward, sandbox: &ManagedSandbox, ) -> String { - let (mut initial_session, initial_marker) = - start_initial_session(namespace, initial_pod, initial_forward, &sandbox.name).await; + let (mut initial_session, initial_marker) = start_initial_session( + gateway_namespace, + initial_pod, + initial_forward, + &sandbox.name, + ) + .await; let deleted = oc( - &["delete", "pod", initial_pod, "-n", namespace, "--wait=true"], + &[ + "delete", + "pod", + initial_pod, + "-n", + gateway_namespace, + "--wait=true", + ], None, ) .await; @@ -678,10 +698,18 @@ async fn run_failover_scenario( // Keep the endpoint stable while replacing its direct backend. This lets // the original client exercise its bounded reconnect behavior. stop_port_forward(initial_forward).await; - wait_for_gateway_pod_ready(namespace, selector, surviving_pod).await; - let reconnect_started_at = unix_timestamp(); + wait_for_gateway_pod_ready(gateway_namespace, selector, surviving_pod).await; let mut reconnect_forward = - port_forward_gateway_pod(namespace, surviving_pod, initial_forward.port).await; + port_forward_gateway_pod(gateway_namespace, surviving_pod, initial_forward.port).await; + // Read the baseline from the sandbox clock. Heartbeat timestamps use the + // same clock, so buffered pre-failover lines cannot satisfy the filter. + let reconnect_started_at: u64 = + exec_on_gateway(&reconnect_forward.endpoint, &sandbox.name, &["date", "+%s"]) + .await + .expect("read sandbox clock after failover") + .trim() + .parse() + .expect("sandbox date +%s should be an integer"); let reconnected_marker = session_marker_after(&mut initial_session, Some(reconnect_started_at)) .await .unwrap_or_else(|error| { @@ -705,7 +733,7 @@ async fn run_failover_scenario( "workspace sentinel was lost after gateway failover: {sentinel}" ); stop_port_forward(&mut reconnect_forward).await; - assert_sandbox_workload_is_live(namespace, &sandbox.name).await + assert_sandbox_workload_is_live(sandbox_namespace, &sandbox.name).await } #[tokio::test] @@ -721,11 +749,12 @@ async fn gateway_pod_failover_preserves_sandbox_session_and_workspace() { } // replica discovery - let namespace = namespace(); + let gateway_namespace = gateway_namespace(); + let sandbox_namespace = sandbox_namespace(); let release = release(); let selector = gateway_selector(&release); - assert_ha_deployment(&namespace, &selector).await; - let pods = gateway_pods(&namespace, &selector).await; // scan for running gateways + assert_ha_deployment(&gateway_namespace, &selector).await; + let pods = gateway_pods(&gateway_namespace, &selector).await; // scan for running gateways assert!( pods.len() >= 2 && pods.iter().filter(|pod| pod.ready).count() >= 2, "HA test requires two ready gateway pods selected by '{selector}', found: {:?}", @@ -746,13 +775,14 @@ async fn gateway_pod_failover_preserves_sandbox_session_and_workspace() { .name .clone(); let port = reserve_loopback_port(); - let mut initial_forward = port_forward_gateway_pod(&namespace, &initial_pod, port).await; + let mut initial_forward = + port_forward_gateway_pod(&gateway_namespace, &initial_pod, port).await; // Retain an independent connection to the replica that will survive the // failure. It is deliberately separate from the stable client endpoint so // cleanup and its verification never fall back to a Service or Route. let cleanup_port = reserve_loopback_port(); let mut surviving_forward = - port_forward_gateway_pod(&namespace, &surviving_pod, cleanup_port).await; + port_forward_gateway_pod(&gateway_namespace, &surviving_pod, cleanup_port).await; let script = format!( "printf '%s\\n' '{SENTINEL}' > /sandbox/.odh-ha-sentinel; session_id=$(cat /proc/sys/kernel/random/uuid); while :; do printf '%s:%s:%s\\n' '{OUTPUT_MARKER}' \"$session_id\" \"$(date +%s)\"; sleep 1; done" @@ -761,7 +791,8 @@ async fn gateway_pod_failover_preserves_sandbox_session_and_workspace() { // Keep ownership of the sandbox here so every assertion failure still // reaches awaited cleanup before the test unwinds. let scenario = AssertUnwindSafe(run_failover_scenario( - &namespace, + &gateway_namespace, + &sandbox_namespace, &selector, &initial_pod, &surviving_pod, @@ -777,7 +808,7 @@ async fn gateway_pod_failover_preserves_sandbox_session_and_workspace() { cleanup.expect("delete sandbox through the surviving gateway pod"); assert_sandbox_deleted( &surviving_forward.endpoint, - &namespace, + &sandbox_namespace, &sandbox.name, &pod_selector, )