From 9db51c01936c6eeafdc59c74612e9839e1623e1d Mon Sep 17 00:00:00 2001 From: Matthew Grossman Date: Thu, 1 Oct 2026 14:15:57 -0700 Subject: [PATCH] fix(kubernetes): serialize lifecycle cleanup with sandbox restart Signed-off-by: Matthew Grossman --- crates/openshell-driver-kubernetes/README.md | 7 + .../openshell-driver-kubernetes/src/driver.rs | 63 +++- crates/openshell-driver-kubernetes/src/lib.rs | 1 + .../src/lifecycle.rs | 26 ++ .../src/lifecycle_tests.rs | 305 ++++++++++++++++++ skills/debug-openshell-cluster/SKILL.md | 6 + 6 files changed, 407 insertions(+), 1 deletion(-) create mode 100644 crates/openshell-driver-kubernetes/src/lifecycle.rs create mode 100644 crates/openshell-driver-kubernetes/src/lifecycle_tests.rs diff --git a/crates/openshell-driver-kubernetes/README.md b/crates/openshell-driver-kubernetes/README.md index 6cb277fbbb..6dd2f764b8 100644 --- a/crates/openshell-driver-kubernetes/README.md +++ b/crates/openshell-driver-kubernetes/README.md @@ -182,6 +182,13 @@ The workload Pod does not share host network, PID, IPC, or process namespaces. The driver uses a scheduling gate to inspect the admitted Pod and bind its UID into the bootstrap claims before kubelet starts it. +Lifecycle RPCs and runtime reconciliation share a per-sandbox mutation gate +across clones of the driver. Reconciliation skips busy sandboxes and refreshes +the Sandbox CR under that gate before cleanup, so a stopped or stopping LIST +snapshot cannot delete a supervisor created by a concurrent restart in the same +driver instance. The gate preserves concurrency across sandboxes; it does not +provide distributed exclusion between separate gateway or driver processes. + ## GPU Support When a sandbox requests GPU support, the driver checks node allocatable capacity diff --git a/crates/openshell-driver-kubernetes/src/driver.rs b/crates/openshell-driver-kubernetes/src/driver.rs index 64ffb0ca97..093a16de80 100644 --- a/crates/openshell-driver-kubernetes/src/driver.rs +++ b/crates/openshell-driver-kubernetes/src/driver.rs @@ -11,6 +11,7 @@ use crate::config::{ use crate::isolation::{ BOUNDARY_PAIR_LABEL, BOUNDARY_ROLE_LABEL, KubernetesSandboxRuntimeBoundarySpec, }; +use crate::lifecycle::LifecycleGates; use crate::sandbox_runtime::{ BOUNDARY_CERTIFICATE_PATH, BOUNDARY_CONFIG_PATH, BOUNDARY_PRIVATE_KEY_PATH, ClientTlsMaterial, SUPERVISOR_TERMINATION_GRACE_PERIOD_SECONDS, SandboxRuntimeNames, SupervisorClientTls, @@ -686,6 +687,7 @@ pub struct KubernetesComputeDriver { client: Client, watch_client: Client, sandbox_api_version: Arc>, + lifecycle_gates: Arc, config: KubernetesComputeConfig, operator_allowlist: Option, } @@ -713,6 +715,7 @@ impl KubernetesComputeDriver { client: client.clone(), watch_client: client, sandbox_api_version: Arc::new(OnceCell::new()), + lifecycle_gates: Arc::default(), config, operator_allowlist: None, } @@ -794,6 +797,7 @@ impl KubernetesComputeDriver { client, watch_client, sandbox_api_version: Arc::new(OnceCell::new()), + lifecycle_gates: Arc::default(), config, operator_allowlist, }; @@ -1764,6 +1768,11 @@ impl KubernetesComputeDriver { )] pub async fn create_sandbox(&self, sandbox: &Sandbox) -> Result { let span_status = openshell_otel::ErrorStatusGuard::current(); + let _guard = self + .lifecycle_gates + .gate_for(&sandbox.id) + .lock_owned() + .await; let result = Box::pin(self.create_sandbox_inner(sandbox)).await; span_status.finish(result) } @@ -2880,6 +2889,7 @@ impl KubernetesComputeDriver { )] pub async fn stop_sandbox(&self, sandbox_id: &str) -> Result<(), KubernetesDriverError> { let span_status = openshell_otel::ErrorStatusGuard::current(); + let _guard = self.lifecycle_gates.gate_for(sandbox_id).lock_owned().await; let result = Box::pin(self.stop_sandbox_inner(sandbox_id)).await; span_status.finish(result) } @@ -2979,6 +2989,7 @@ impl KubernetesComputeDriver { expected_runtime_identity: &str, ) -> Result { let span_status = openshell_otel::ErrorStatusGuard::current(); + let _guard = self.lifecycle_gates.gate_for(sandbox_id).lock_owned().await; let result = Box::pin(self.start_sandbox_runtime_generation( sandbox_id, generation_id, @@ -3457,6 +3468,7 @@ impl KubernetesComputeDriver { )] pub async fn delete_sandbox(&self, sandbox_id: &str) -> Result { let span_status = openshell_otel::ErrorStatusGuard::current(); + let _guard = self.lifecycle_gates.gate_for(sandbox_id).lock_owned().await; let result = self.delete_sandbox_inner(sandbox_id).await; span_status.finish(result) } @@ -3657,6 +3669,44 @@ impl KubernetesComputeDriver { let Ok(sandbox_id) = sandbox_id_from_object(&object) else { continue; }; + // Lifecycle RPCs can replace the stable supervisor Pod name while + // this LIST snapshot still describes the previous stopped state. + // Skip in-flight mutations, then refresh under the shared gate so + // a snapshot taken before a completed restart cannot delete it. + let Ok(_guard) = self.lifecycle_gates.gate_for(&sandbox_id).try_lock_owned() else { + continue; + }; + let Some(name) = object.metadata.name.as_deref() else { + continue; + }; + let namespace = object + .metadata + .namespace + .as_deref() + .unwrap_or(&self.config.namespace); + let api = Self::agent_sandbox_api( + self.client.clone(), + &lookup_api.resource.version, + namespace, + ); + let refreshed = match tokio::time::timeout(KUBE_API_TIMEOUT, api.api.get(name)).await { + Ok(Ok(refreshed)) => refreshed, + Ok(Err(KubeError::Api(error))) if error.code == 404 => continue, + Ok(Err(error)) => { + debug!(%sandbox_id, %error, "could not refresh Sandbox for runtime reconciliation"); + continue; + } + Err(_) => { + warn!(%sandbox_id, "timed out refreshing Sandbox for runtime reconciliation"); + continue; + } + }; + if refreshed.metadata.uid != object.metadata.uid + || sandbox_id_from_object(&refreshed).as_deref() != Ok(sandbox_id.as_str()) + { + continue; + } + let object = refreshed; if let Err(error) = self.admit_stored_resources(&object).await { warn!(%sandbox_id, reason = %error.message(), "Sandbox resource admission revalidation failed"); if error.code() == tonic::Code::FailedPrecondition { @@ -7449,10 +7499,15 @@ mod tests { serde_json::json!({ "apiVersion": "agents.x-k8s.io/v1beta1", "kind": "SandboxList", - "items": [sandbox] + "items": [sandbox.clone()] }), ), ), + ( + http::Method::GET, + "/apis/agents.x-k8s.io/v1beta1/namespaces/openshell/sandboxes/sandbox-cr", + kube_test_response(http::StatusCode::OK, sandbox), + ), ( http::Method::GET, "/api/v1/namespaces/openshell/persistentvolumeclaims/team-data", @@ -7488,6 +7543,7 @@ mod tests { client: client.clone(), watch_client: client, sandbox_api_version: Arc::new(OnceCell::new()), + lifecycle_gates: Arc::default(), config: KubernetesComputeConfig::default(), operator_allowlist: None, }; @@ -8462,6 +8518,7 @@ mod tests { client: client.clone(), watch_client: client, sandbox_api_version: Arc::new(OnceCell::new()), + lifecycle_gates: Arc::default(), config: KubernetesComputeConfig::default(), operator_allowlist: None, }; @@ -8552,6 +8609,7 @@ mod tests { client: client.clone(), watch_client: client, sandbox_api_version: Arc::new(OnceCell::new()), + lifecycle_gates: Arc::default(), config: KubernetesComputeConfig::default(), operator_allowlist: None, }; @@ -11013,6 +11071,7 @@ mod tests { client: client.clone(), watch_client: client, sandbox_api_version: Arc::new(OnceCell::new()), + lifecycle_gates: Arc::default(), config, operator_allowlist: None, }; @@ -11748,4 +11807,6 @@ mod tests { alpha.data = serde_json::json!({"spec": {"replicas": 1}}); assert!(sandbox_runtime_should_run(&alpha)); } + + include!("lifecycle_tests.rs"); } diff --git a/crates/openshell-driver-kubernetes/src/lib.rs b/crates/openshell-driver-kubernetes/src/lib.rs index 607f95d382..aa980db44f 100644 --- a/crates/openshell-driver-kubernetes/src/lib.rs +++ b/crates/openshell-driver-kubernetes/src/lib.rs @@ -5,6 +5,7 @@ pub mod config; pub mod driver; pub mod grpc; pub mod isolation; +mod lifecycle; pub mod otel_tracing; mod resource_admission; mod sandbox_runtime; diff --git a/crates/openshell-driver-kubernetes/src/lifecycle.rs b/crates/openshell-driver-kubernetes/src/lifecycle.rs new file mode 100644 index 0000000000..74d2be3841 --- /dev/null +++ b/crates/openshell-driver-kubernetes/src/lifecycle.rs @@ -0,0 +1,26 @@ +// SPDX-FileCopyrightText: Copyright (c) 2025-2026 NVIDIA CORPORATION & AFFILIATES. All rights reserved. +// SPDX-License-Identifier: Apache-2.0 + +//! Per-sandbox serialization shared by lifecycle RPCs and driver reconciliation. + +use std::collections::HashMap; +use std::sync::{Arc, Mutex, Weak}; +use tokio::sync::Mutex as AsyncMutex; + +#[derive(Debug, Default)] +pub struct LifecycleGates { + gates: Mutex>>>, +} + +impl LifecycleGates { + pub fn gate_for(&self, sandbox_id: &str) -> Arc> { + let mut gates = self.gates.lock().expect("lifecycle gate registry poisoned"); + gates.retain(|_, gate| gate.strong_count() > 0); + if let Some(gate) = gates.get(sandbox_id).and_then(Weak::upgrade) { + return gate; + } + let gate = Arc::new(AsyncMutex::new(())); + gates.insert(sandbox_id.to_string(), Arc::downgrade(&gate)); + gate + } +} diff --git a/crates/openshell-driver-kubernetes/src/lifecycle_tests.rs b/crates/openshell-driver-kubernetes/src/lifecycle_tests.rs new file mode 100644 index 0000000000..9d70b4340e --- /dev/null +++ b/crates/openshell-driver-kubernetes/src/lifecycle_tests.rs @@ -0,0 +1,305 @@ +// SPDX-FileCopyrightText: Copyright (c) 2025-2026 NVIDIA CORPORATION & AFFILIATES. All rights reserved. +// SPDX-License-Identifier: Apache-2.0 + +mod lifecycle_reconciliation { + use super::*; + + // Availability reads run concurrently. Match pending responses by path + // rather than assuming tokio::join! polls its branches in a fixed order. + fn read_only_driver(steps: Vec) -> ScriptedDriver { + let steps = Arc::new(std::sync::Mutex::new(VecDeque::from(steps))); + let pending = steps.clone(); + let service = tower::service_fn(move |request: http::Request| { + let pending = pending.clone(); + async move { + assert_eq!( + request.method(), + http::Method::GET, + "reconciliation must preserve the replacement" + ); + let mut pending = pending.lock().unwrap(); + let index = pending + .iter() + .position(|(method, path, _)| { + request.method() == method && request.uri().path() == *path + }) + .unwrap_or_else(|| panic!("unexpected read {}", request.uri().path())); + let (_, _, response) = pending.remove(index).unwrap(); + Ok::<_, std::convert::Infallible>(response) + } + }); + let client = Client::new(service, "openshell"); + let driver = KubernetesComputeDriver { + client: client.clone(), + watch_client: client, + sandbox_api_version: Arc::new(OnceCell::new()), + lifecycle_gates: Arc::default(), + config: KubernetesComputeConfig::default(), + operator_allowlist: None, + }; + (driver, steps, Arc::default()) + } + + fn snapshot(version: &str, phase: Option) -> serde_json::Value { + let mut object = serde_json::json!({ + "apiVersion": format!("{SANDBOX_GROUP}/{version}"), "kind": "Sandbox", + "metadata": { + "name": "sandbox-cr", "namespace": "openshell", "uid": "cr-uid", + "resourceVersion": "42", + "labels": {LABEL_SANDBOX_ID: "sandbox-1", LABEL_SANDBOX_WORKSPACE: "team-a"}, + "annotations": { + crate::resource_admission::CONFIG_USED: "false", + crate::resource_admission::IDENTITIES: "{}", + ANNOTATION_SANDBOX_RUNTIME_READINESS: "unavailable" + } + }, + "spec": {"podTemplate": {"spec": { + "automountServiceAccountToken": false, + "volumes": [{"name": SANDBOX_BOOTSTRAP_VOLUME_NAME, + "secret": {"secretName": "os-sandbox-sandbox-1-gen2"}}] + }}} + }); + object["spec"][if version == SANDBOX_VERSION_V1ALPHA1 { + "replicas" + } else { + "operatingMode" + }] = if version == SANDBOX_VERSION_V1ALPHA1 { + serde_json::json!(0) + } else { + serde_json::json!("Suspended") + }; + if let Some(phase) = phase { + object["metadata"]["annotations"][ANNOTATION_SANDBOX_RUNTIME_BOOTSTRAPPING] = + serde_json::json!("true"); + object["metadata"]["annotations"][ANNOTATION_SANDBOX_RUNTIME_BOOTSTRAP_OPERATION] = + serde_json::json!("stop"); + object["metadata"]["annotations"][ANNOTATION_SANDBOX_RUNTIME_BOOTSTRAP_PHASE] = + serde_json::json!(phase.as_str()); + object["metadata"]["annotations"][ANNOTATION_SANDBOX_RUNTIME_SUPERVISOR_UID] = + serde_json::json!("old-supervisor-uid"); + } + object + } + + fn restarted_snapshot(version: &str) -> serde_json::Value { + let mut object = snapshot(version, Some(SandboxRuntimeBootstrapPhase::Preparing)); + object["metadata"]["resourceVersion"] = serde_json::json!("43"); + let annotations = &mut object["metadata"]["annotations"]; + annotations[ANNOTATION_SANDBOX_RUNTIME_BOOTSTRAP_OPERATION] = serde_json::json!("restart"); + annotations[ANNOTATION_SANDBOX_RUNTIME_BOOTSTRAP_STARTED_AT] = + serde_json::json!(openshell_core::time::now_ms().to_string()); + annotations[ANNOTATION_SANDBOX_RUNTIME_SUPERVISOR_UID] = + serde_json::json!("new-supervisor-uid"); + annotations[ANNOTATION_SANDBOX_RUNTIME_NETWORK_POLICY_UID] = + serde_json::json!("sandbox-workload-fence-uid"); + annotations[ANNOTATION_SANDBOX_RUNTIME_NETWORK_POLICY_GENERATION] = serde_json::json!("1"); + object["spec"][if version == SANDBOX_VERSION_V1ALPHA1 { + "replicas" + } else { + "operatingMode" + }] = if version == SANDBOX_VERSION_V1ALPHA1 { + serde_json::json!(1) + } else { + serde_json::json!("Running") + }; + object + } + + fn listed(version: &str, object: serde_json::Value) -> KubeTestStep { + ( + http::Method::GET, + if version == SANDBOX_VERSION_V1ALPHA1 { + "/apis/agents.x-k8s.io/v1alpha1/namespaces/openshell/sandboxes" + } else { + "/apis/agents.x-k8s.io/v1beta1/namespaces/openshell/sandboxes" + }, + kube_test_response( + http::StatusCode::OK, + serde_json::json!({ + "apiVersion": format!("{SANDBOX_GROUP}/{version}"), + "kind": "SandboxList", "items": [object] + }), + ), + ) + } + + fn refreshed(version: &str, object: serde_json::Value) -> KubeTestStep { + ( + http::Method::GET, + if version == SANDBOX_VERSION_V1ALPHA1 { + "/apis/agents.x-k8s.io/v1alpha1/namespaces/openshell/sandboxes/sandbox-cr" + } else { + "/apis/agents.x-k8s.io/v1beta1/namespaces/openshell/sandboxes/sandbox-cr" + }, + kube_test_response(http::StatusCode::OK, object), + ) + } + + fn preparing_dependencies() -> Vec { + let mut fence = workload_fence("openshell", &SandboxRuntimeNames::new("sandbox-1"), 5500); + for (policy, component) in [ + (&mut fence.workload_policy, "sandbox-workload-fence"), + (&mut fence.supervisor_policy, "sandbox-supervisor-egress"), + ] { + policy.metadata.labels.get_or_insert_default().extend([ + ( + LABEL_MANAGED_BY.to_string(), + LABEL_MANAGED_BY_VALUE.to_string(), + ), + ("openshell.ai/component".to_string(), component.to_string()), + ]); + policy.metadata.uid = Some(format!("{component}-uid")); + policy.metadata.generation = Some(1); + } + let workload = || { + ( + http::Method::GET, + "/apis/networking.k8s.io/v1/namespaces/openshell/networkpolicies/openshell-sandbox-workloads", + kube_test_response( + http::StatusCode::OK, + serde_json::to_value(&fence.workload_policy).unwrap(), + ), + ) + }; + let supervisor = || { + ( + http::Method::GET, + "/apis/networking.k8s.io/v1/namespaces/openshell/networkpolicies/openshell-sandbox-supervisors", + kube_test_response( + http::StatusCode::OK, + serde_json::to_value(&fence.supervisor_policy).unwrap(), + ), + ) + }; + vec![ + workload(), + supervisor(), + workload(), + ( + http::Method::GET, + "/api/v1/namespaces/openshell/pods/os-supervisor-sandbox-1", + kube_test_response( + http::StatusCode::OK, + serde_json::json!({ + "apiVersion": "v1", "kind": "Pod", + "metadata": {"name": "os-supervisor-sandbox-1", "uid": "new-supervisor-uid"}, + "spec": {"containers": []}, "status": {"phase": "Pending"} + }), + ), + ), + ( + http::Method::GET, + "/api/v1/namespaces/openshell/services/os-boundary-sandbox-1", + kube_test_response( + http::StatusCode::OK, + serde_json::json!({ + "apiVersion": "v1", "kind": "Service", "metadata": {"name": "os-boundary-sandbox-1"} + }), + ), + ), + workload(), + supervisor(), + ] + } + + #[tokio::test] + async fn stale_stop_snapshots_cannot_delete_a_restarted_supervisor() { + for version in [SANDBOX_VERSION_V1ALPHA1, SANDBOX_VERSION_V1BETA1] { + for phase in [ + None, + Some(SandboxRuntimeBootstrapPhase::Releasing), + Some(SandboxRuntimeBootstrapPhase::Suspending), + ] { + // Restart completed its critical section after LIST. The old + // implementation either deleted the replacement or attempted + // cleanup using the old stop transition. Every allowed request + // here is a read; unexpected DELETE/PATCH/Secret access fails. + let mut steps = vec![ + listed(version, snapshot(version, phase)), + refreshed(version, restarted_snapshot(version)), + ]; + steps.extend(preparing_dependencies()); + let (driver, steps, _) = read_only_driver(steps); + driver.sandbox_api_version.set(version).unwrap(); + driver.reconcile_sandbox_runtime_resources().await; + assert!(steps.lock().unwrap().is_empty()); + } + } + } + + #[tokio::test] + async fn reconciliation_skips_a_busy_restart_and_retries_after_release() { + let version = SANDBOX_VERSION_V1BETA1; + let mut steps = vec![ + listed(version, snapshot(version, None)), + listed(version, snapshot(version, None)), + refreshed(version, restarted_snapshot(version)), + ]; + steps.extend(preparing_dependencies()); + let (driver, steps, _) = read_only_driver(steps); + driver.sandbox_api_version.set(version).unwrap(); + let clone = driver.clone(); + let guard = clone + .lifecycle_gates + .gate_for("sandbox-1") + .lock_owned() + .await; + // Restart may already have created its Pod while the CR is still + // stopped. Reconciliation must skip even reading that companion. + driver.reconcile_sandbox_runtime_resources().await; + drop(guard); + driver.reconcile_sandbox_runtime_resources().await; + assert!(steps.lock().unwrap().is_empty()); + } + + #[tokio::test] + async fn refresh_cannot_adopt_a_replaced_sandbox_cr() { + let version = SANDBOX_VERSION_V1BETA1; + let mut replacement = restarted_snapshot(version); + replacement["metadata"]["uid"] = serde_json::json!("another-cr-uid"); + let (driver, steps, _) = read_only_driver(vec![ + listed(version, snapshot(version, None)), + refreshed(version, replacement), + ]); + driver.sandbox_api_version.set(version).unwrap(); + driver.reconcile_sandbox_runtime_resources().await; + assert!(steps.lock().unwrap().is_empty()); + } + + #[tokio::test] + async fn lifecycle_requests_share_gates_across_clones_without_blocking_other_sandboxes() { + let driver = KubernetesComputeDriver::new_for_test(KubernetesComputeConfig::default()); + let clone = driver.clone(); + let guard = clone + .lifecycle_gates + .gate_for("sandbox-1") + .lock_owned() + .await; + let sandbox = Sandbox { + id: "sandbox-1".to_string(), + ..Default::default() + }; + let mut create = Box::pin(driver.create_sandbox(&sandbox)); + let mut stop = Box::pin(driver.stop_sandbox("sandbox-1")); + let mut delete = Box::pin(driver.delete_sandbox("sandbox-1")); + let mut start = Box::pin(driver.start_sandbox("sandbox-1", "", &[], "")); + assert!(futures::poll!(&mut create).is_pending()); + assert!(futures::poll!(&mut stop).is_pending()); + assert!(futures::poll!(&mut delete).is_pending()); + assert!(futures::poll!(&mut start).is_pending()); + driver + .start_sandbox("sandbox-2", "", &[], "") + .await + .unwrap_err(); + // Cancel queued mutations, then let start obtain the gate. Its normal + // generation validation proves release did not leave it deadlocked. + drop(create); + drop(stop); + drop(delete); + drop(guard); + assert!(matches!( + start.await, + Err(KubernetesDriverError::InvalidArgument(_)) + )); + } +} diff --git a/skills/debug-openshell-cluster/SKILL.md b/skills/debug-openshell-cluster/SKILL.md index 46d5abc254..ab774e3a88 100644 --- a/skills/debug-openshell-cluster/SKILL.md +++ b/skills/debug-openshell-cluster/SKILL.md @@ -747,6 +747,12 @@ remain unchanged. A generation-bound session-token rejection usually means the supervisor is presenting credentials from a runtime that was replaced; inspect the persisted generation before retrying bootstrap. +The Kubernetes driver serializes lifecycle mutations and runtime reconciliation +per sandbox within one driver instance. A busy sandbox is checked again on the +next reconciliation pass. If restart still loses its supervisor, compare the +Sandbox and Pod UIDs and identify which gateway or external driver process +performed cleanup; the local mutation gate does not coordinate separate processes. + ```bash helm -n openshell get values openshell | grep -A3 sandboxServiceAccount kubectl -n get serviceaccount openshell-sandbox