Skip to content

Commit 9db51c0

Browse files
fix(kubernetes): serialize lifecycle cleanup with sandbox restart
Signed-off-by: Matthew Grossman <mgrossman@nvidia.com>
1 parent 76cfd0e commit 9db51c0

6 files changed

Lines changed: 407 additions & 1 deletion

File tree

‎crates/openshell-driver-kubernetes/README.md‎

Lines changed: 7 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -182,6 +182,13 @@ The workload Pod does not share host network, PID, IPC, or process namespaces.
182182
The driver uses a scheduling gate to inspect the admitted Pod and bind its UID
183183
into the bootstrap claims before kubelet starts it.
184184

185+
Lifecycle RPCs and runtime reconciliation share a per-sandbox mutation gate
186+
across clones of the driver. Reconciliation skips busy sandboxes and refreshes
187+
the Sandbox CR under that gate before cleanup, so a stopped or stopping LIST
188+
snapshot cannot delete a supervisor created by a concurrent restart in the same
189+
driver instance. The gate preserves concurrency across sandboxes; it does not
190+
provide distributed exclusion between separate gateway or driver processes.
191+
185192
## GPU Support
186193

187194
When a sandbox requests GPU support, the driver checks node allocatable capacity

‎crates/openshell-driver-kubernetes/src/driver.rs‎

Lines changed: 62 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -11,6 +11,7 @@ use crate::config::{
1111
use crate::isolation::{
1212
BOUNDARY_PAIR_LABEL, BOUNDARY_ROLE_LABEL, KubernetesSandboxRuntimeBoundarySpec,
1313
};
14+
use crate::lifecycle::LifecycleGates;
1415
use crate::sandbox_runtime::{
1516
BOUNDARY_CERTIFICATE_PATH, BOUNDARY_CONFIG_PATH, BOUNDARY_PRIVATE_KEY_PATH, ClientTlsMaterial,
1617
SUPERVISOR_TERMINATION_GRACE_PERIOD_SECONDS, SandboxRuntimeNames, SupervisorClientTls,
@@ -686,6 +687,7 @@ pub struct KubernetesComputeDriver {
686687
client: Client,
687688
watch_client: Client,
688689
sandbox_api_version: Arc<OnceCell<&'static str>>,
690+
lifecycle_gates: Arc<LifecycleGates>,
689691
config: KubernetesComputeConfig,
690692
operator_allowlist: Option<OperatorNamespaceAllowlist>,
691693
}
@@ -713,6 +715,7 @@ impl KubernetesComputeDriver {
713715
client: client.clone(),
714716
watch_client: client,
715717
sandbox_api_version: Arc::new(OnceCell::new()),
718+
lifecycle_gates: Arc::default(),
716719
config,
717720
operator_allowlist: None,
718721
}
@@ -794,6 +797,7 @@ impl KubernetesComputeDriver {
794797
client,
795798
watch_client,
796799
sandbox_api_version: Arc::new(OnceCell::new()),
800+
lifecycle_gates: Arc::default(),
797801
config,
798802
operator_allowlist,
799803
};
@@ -1764,6 +1768,11 @@ impl KubernetesComputeDriver {
17641768
)]
17651769
pub async fn create_sandbox(&self, sandbox: &Sandbox) -> Result<String, KubernetesDriverError> {
17661770
let span_status = openshell_otel::ErrorStatusGuard::current();
1771+
let _guard = self
1772+
.lifecycle_gates
1773+
.gate_for(&sandbox.id)
1774+
.lock_owned()
1775+
.await;
17671776
let result = Box::pin(self.create_sandbox_inner(sandbox)).await;
17681777
span_status.finish(result)
17691778
}
@@ -2880,6 +2889,7 @@ impl KubernetesComputeDriver {
28802889
)]
28812890
pub async fn stop_sandbox(&self, sandbox_id: &str) -> Result<(), KubernetesDriverError> {
28822891
let span_status = openshell_otel::ErrorStatusGuard::current();
2892+
let _guard = self.lifecycle_gates.gate_for(sandbox_id).lock_owned().await;
28832893
let result = Box::pin(self.stop_sandbox_inner(sandbox_id)).await;
28842894
span_status.finish(result)
28852895
}
@@ -2979,6 +2989,7 @@ impl KubernetesComputeDriver {
29792989
expected_runtime_identity: &str,
29802990
) -> Result<String, KubernetesDriverError> {
29812991
let span_status = openshell_otel::ErrorStatusGuard::current();
2992+
let _guard = self.lifecycle_gates.gate_for(sandbox_id).lock_owned().await;
29822993
let result = Box::pin(self.start_sandbox_runtime_generation(
29832994
sandbox_id,
29842995
generation_id,
@@ -3457,6 +3468,7 @@ impl KubernetesComputeDriver {
34573468
)]
34583469
pub async fn delete_sandbox(&self, sandbox_id: &str) -> Result<bool, String> {
34593470
let span_status = openshell_otel::ErrorStatusGuard::current();
3471+
let _guard = self.lifecycle_gates.gate_for(sandbox_id).lock_owned().await;
34603472
let result = self.delete_sandbox_inner(sandbox_id).await;
34613473
span_status.finish(result)
34623474
}
@@ -3657,6 +3669,44 @@ impl KubernetesComputeDriver {
36573669
let Ok(sandbox_id) = sandbox_id_from_object(&object) else {
36583670
continue;
36593671
};
3672+
// Lifecycle RPCs can replace the stable supervisor Pod name while
3673+
// this LIST snapshot still describes the previous stopped state.
3674+
// Skip in-flight mutations, then refresh under the shared gate so
3675+
// a snapshot taken before a completed restart cannot delete it.
3676+
let Ok(_guard) = self.lifecycle_gates.gate_for(&sandbox_id).try_lock_owned() else {
3677+
continue;
3678+
};
3679+
let Some(name) = object.metadata.name.as_deref() else {
3680+
continue;
3681+
};
3682+
let namespace = object
3683+
.metadata
3684+
.namespace
3685+
.as_deref()
3686+
.unwrap_or(&self.config.namespace);
3687+
let api = Self::agent_sandbox_api(
3688+
self.client.clone(),
3689+
&lookup_api.resource.version,
3690+
namespace,
3691+
);
3692+
let refreshed = match tokio::time::timeout(KUBE_API_TIMEOUT, api.api.get(name)).await {
3693+
Ok(Ok(refreshed)) => refreshed,
3694+
Ok(Err(KubeError::Api(error))) if error.code == 404 => continue,
3695+
Ok(Err(error)) => {
3696+
debug!(%sandbox_id, %error, "could not refresh Sandbox for runtime reconciliation");
3697+
continue;
3698+
}
3699+
Err(_) => {
3700+
warn!(%sandbox_id, "timed out refreshing Sandbox for runtime reconciliation");
3701+
continue;
3702+
}
3703+
};
3704+
if refreshed.metadata.uid != object.metadata.uid
3705+
|| sandbox_id_from_object(&refreshed).as_deref() != Ok(sandbox_id.as_str())
3706+
{
3707+
continue;
3708+
}
3709+
let object = refreshed;
36603710
if let Err(error) = self.admit_stored_resources(&object).await {
36613711
warn!(%sandbox_id, reason = %error.message(), "Sandbox resource admission revalidation failed");
36623712
if error.code() == tonic::Code::FailedPrecondition {
@@ -7449,10 +7499,15 @@ mod tests {
74497499
serde_json::json!({
74507500
"apiVersion": "agents.x-k8s.io/v1beta1",
74517501
"kind": "SandboxList",
7452-
"items": [sandbox]
7502+
"items": [sandbox.clone()]
74537503
}),
74547504
),
74557505
),
7506+
(
7507+
http::Method::GET,
7508+
"/apis/agents.x-k8s.io/v1beta1/namespaces/openshell/sandboxes/sandbox-cr",
7509+
kube_test_response(http::StatusCode::OK, sandbox),
7510+
),
74567511
(
74577512
http::Method::GET,
74587513
"/api/v1/namespaces/openshell/persistentvolumeclaims/team-data",
@@ -7488,6 +7543,7 @@ mod tests {
74887543
client: client.clone(),
74897544
watch_client: client,
74907545
sandbox_api_version: Arc::new(OnceCell::new()),
7546+
lifecycle_gates: Arc::default(),
74917547
config: KubernetesComputeConfig::default(),
74927548
operator_allowlist: None,
74937549
};
@@ -8462,6 +8518,7 @@ mod tests {
84628518
client: client.clone(),
84638519
watch_client: client,
84648520
sandbox_api_version: Arc::new(OnceCell::new()),
8521+
lifecycle_gates: Arc::default(),
84658522
config: KubernetesComputeConfig::default(),
84668523
operator_allowlist: None,
84678524
};
@@ -8552,6 +8609,7 @@ mod tests {
85528609
client: client.clone(),
85538610
watch_client: client,
85548611
sandbox_api_version: Arc::new(OnceCell::new()),
8612+
lifecycle_gates: Arc::default(),
85558613
config: KubernetesComputeConfig::default(),
85568614
operator_allowlist: None,
85578615
};
@@ -11013,6 +11071,7 @@ mod tests {
1101311071
client: client.clone(),
1101411072
watch_client: client,
1101511073
sandbox_api_version: Arc::new(OnceCell::new()),
11074+
lifecycle_gates: Arc::default(),
1101611075
config,
1101711076
operator_allowlist: None,
1101811077
};
@@ -11748,4 +11807,6 @@ mod tests {
1174811807
alpha.data = serde_json::json!({"spec": {"replicas": 1}});
1174911808
assert!(sandbox_runtime_should_run(&alpha));
1175011809
}
11810+
11811+
include!("lifecycle_tests.rs");
1175111812
}

‎crates/openshell-driver-kubernetes/src/lib.rs‎

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -5,6 +5,7 @@ pub mod config;
55
pub mod driver;
66
pub mod grpc;
77
pub mod isolation;
8+
mod lifecycle;
89
pub mod otel_tracing;
910
mod resource_admission;
1011
mod sandbox_runtime;
Lines changed: 26 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,26 @@
1+
// SPDX-FileCopyrightText: Copyright (c) 2025-2026 NVIDIA CORPORATION & AFFILIATES. All rights reserved.
2+
// SPDX-License-Identifier: Apache-2.0
3+
4+
//! Per-sandbox serialization shared by lifecycle RPCs and driver reconciliation.
5+
6+
use std::collections::HashMap;
7+
use std::sync::{Arc, Mutex, Weak};
8+
use tokio::sync::Mutex as AsyncMutex;
9+
10+
#[derive(Debug, Default)]
11+
pub struct LifecycleGates {
12+
gates: Mutex<HashMap<String, Weak<AsyncMutex<()>>>>,
13+
}
14+
15+
impl LifecycleGates {
16+
pub fn gate_for(&self, sandbox_id: &str) -> Arc<AsyncMutex<()>> {
17+
let mut gates = self.gates.lock().expect("lifecycle gate registry poisoned");
18+
gates.retain(|_, gate| gate.strong_count() > 0);
19+
if let Some(gate) = gates.get(sandbox_id).and_then(Weak::upgrade) {
20+
return gate;
21+
}
22+
let gate = Arc::new(AsyncMutex::new(()));
23+
gates.insert(sandbox_id.to_string(), Arc::downgrade(&gate));
24+
gate
25+
}
26+
}

0 commit comments

Comments
 (0)