Skip to content

Commit 585bcf0

Browse files
authored
(Fix) ha sandbox resilience with k8s (#3644)
* fix(server): retry internal store updates on version conflict - Re-read and reapply internal CAS updates (expected version 0) up to 5 times on conflict - Client-supplied versions still fail on conflict - Add concurrent-writer test Signed-off-by: divesh <dgude@nvidia.com> * fix(kubernetes): keep sandboxes running while the supervisor reconnects - Treat a running but not-Ready supervisor Pod as degraded, not unavailable - Report degraded sandboxes as not ready without suspending them - Suspend only when the supervisor Pod is missing, terminated, or deleting - Add availability test Signed-off-by: divesh <dgude@nvidia.com> * fix(kubernetes): skip fence generation check for suspended sandboxes - Check the fence generation only for running or bootstrapping sandboxes - Stop re-suspending stopped sandboxes and logging a warning every reconcile Signed-off-by: divesh <dgude@nvidia.com> * fix(kubernetes): rank degraded supervisor above unknown dependencies - Report a degraded supervisor as unavailable even when another dependency read is unknown - Extract dependency aggregation and readiness mapping into pure helpers - Add mixed degraded and unknown regression test Signed-off-by: divesh <dgude@nvidia.com> --------- Signed-off-by: divesh <dgude@nvidia.com>
1 parent 854b237 commit 585bcf0

3 files changed

Lines changed: 254 additions & 87 deletions

File tree

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

Lines changed: 120 additions & 21 deletions
Original file line numberDiff line numberDiff line change
@@ -3719,7 +3719,11 @@ impl KubernetesComputeDriver {
37193719
else {
37203720
continue;
37213721
};
3722-
if !sandbox_runtime_namespace_fence_generation_matches(&fence, &object) {
3722+
let workload_may_run = sandbox_runtime_should_run(&object)
3723+
|| sandbox_runtime_bootstrap_in_progress(&object);
3724+
if workload_may_run
3725+
&& !sandbox_runtime_namespace_fence_generation_matches(&fence, &object)
3726+
{
37233727
warn!(
37243728
sandbox_id,
37253729
"namespace workload fence generation changed; suspending stale boundary"
@@ -3751,7 +3755,8 @@ impl KubernetesComputeDriver {
37513755
)
37523756
.await
37533757
{
3754-
SandboxRuntimeControlAvailability::Available => {}
3758+
SandboxRuntimeControlAvailability::Available
3759+
| SandboxRuntimeControlAvailability::Degraded => {}
37553760
SandboxRuntimeControlAvailability::Unavailable => {
37563761
warn!(
37573762
sandbox_id,
@@ -3777,7 +3782,8 @@ impl KubernetesComputeDriver {
37773782
)
37783783
.await
37793784
{
3780-
SandboxRuntimeControlAvailability::Available => {}
3785+
SandboxRuntimeControlAvailability::Available
3786+
| SandboxRuntimeControlAvailability::Degraded => {}
37813787
SandboxRuntimeControlAvailability::Unavailable => {
37823788
warn!(
37833789
sandbox_id,
@@ -3934,10 +3940,8 @@ impl KubernetesComputeDriver {
39343940
object: &DynamicObject,
39353941
availability: SandboxRuntimeControlAvailability,
39363942
) {
3937-
let state = match availability {
3938-
SandboxRuntimeControlAvailability::Available => "ready",
3939-
SandboxRuntimeControlAvailability::Unavailable => "unavailable",
3940-
SandboxRuntimeControlAvailability::Unknown => return,
3943+
let Some(state) = sandbox_runtime_readiness_state(availability) else {
3944+
return;
39413945
};
39423946
if object
39433947
.metadata
@@ -5164,23 +5168,27 @@ fn sandbox_from_object(namespace: &str, obj: DynamicObject) -> Result<(String, S
51645168
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
51655169
enum SandboxRuntimeControlAvailability {
51665170
Available,
5171+
Degraded,
51675172
Unavailable,
51685173
Unknown,
51695174
}
51705175

51715176
fn sandbox_runtime_control_availability_from_pod(pod: &Pod) -> SandboxRuntimeControlAvailability {
5172-
if pod.metadata.deletion_timestamp.is_none()
5173-
&& pod
5174-
.status
5175-
.as_ref()
5176-
.and_then(|status| status.conditions.as_ref())
5177-
.is_some_and(|conditions| {
5178-
conditions
5179-
.iter()
5180-
.any(|condition| condition.type_ == "Ready" && condition.status == "True")
5181-
})
5182-
{
5177+
if pod.metadata.deletion_timestamp.is_some() {
5178+
return SandboxRuntimeControlAvailability::Unavailable;
5179+
}
5180+
let status = pod.status.as_ref();
5181+
let ready = status
5182+
.and_then(|status| status.conditions.as_ref())
5183+
.is_some_and(|conditions| {
5184+
conditions
5185+
.iter()
5186+
.any(|condition| condition.type_ == "Ready" && condition.status == "True")
5187+
});
5188+
if ready {
51835189
SandboxRuntimeControlAvailability::Available
5190+
} else if status.and_then(|status| status.phase.as_deref()) == Some("Running") {
5191+
SandboxRuntimeControlAvailability::Degraded
51845192
} else {
51855193
SandboxRuntimeControlAvailability::Unavailable
51865194
}
@@ -5277,16 +5285,34 @@ async fn sandbox_runtime_control_availability(
52775285
SandboxRuntimeControlAvailability::Unknown
52785286
}
52795287
};
5280-
let dependencies = [control, service, fence, supervisor_fence];
5288+
combine_sandbox_runtime_control_availability([control, service, fence, supervisor_fence])
5289+
}
5290+
5291+
fn combine_sandbox_runtime_control_availability(
5292+
dependencies: [SandboxRuntimeControlAvailability; 4],
5293+
) -> SandboxRuntimeControlAvailability {
52815294
if dependencies.contains(&SandboxRuntimeControlAvailability::Unavailable) {
52825295
SandboxRuntimeControlAvailability::Unavailable
5296+
} else if dependencies.contains(&SandboxRuntimeControlAvailability::Degraded) {
5297+
SandboxRuntimeControlAvailability::Degraded
52835298
} else if dependencies.contains(&SandboxRuntimeControlAvailability::Unknown) {
52845299
SandboxRuntimeControlAvailability::Unknown
52855300
} else {
52865301
SandboxRuntimeControlAvailability::Available
52875302
}
52885303
}
52895304

5305+
fn sandbox_runtime_readiness_state(
5306+
availability: SandboxRuntimeControlAvailability,
5307+
) -> Option<&'static str> {
5308+
match availability {
5309+
SandboxRuntimeControlAvailability::Available => Some("ready"),
5310+
SandboxRuntimeControlAvailability::Degraded
5311+
| SandboxRuntimeControlAvailability::Unavailable => Some("unavailable"),
5312+
SandboxRuntimeControlAvailability::Unknown => None,
5313+
}
5314+
}
5315+
52905316
async fn sandbox_runtime_workload_generation_matches(
52915317
client: &Client,
52925318
namespace: &str,
@@ -5395,8 +5421,11 @@ async fn sandbox_from_object_with_sandbox_runtime_readiness(
53955421
// A transient API read failure is not evidence that the running
53965422
// boundary disappeared. Preserve the CR's published readiness for
53975423
// Unknown and let periodic reconciliation retry.
5398-
if dependencies == SandboxRuntimeControlAvailability::Unavailable
5399-
|| workload_generation == SandboxRuntimeControlAvailability::Unavailable
5424+
if matches!(
5425+
dependencies,
5426+
SandboxRuntimeControlAvailability::Unavailable
5427+
| SandboxRuntimeControlAvailability::Degraded
5428+
) || workload_generation == SandboxRuntimeControlAvailability::Unavailable
54005429
|| supervisor_generation == SandboxRuntimeControlAvailability::Unavailable
54015430
{
54025431
mark_sandbox_runtime_control_unavailable(&mut sandbox);
@@ -11373,6 +11402,76 @@ mod tests {
1137311402
);
1137411403
}
1137511404

11405+
#[test]
11406+
fn sandbox_runtime_degraded_supervisor_is_not_masked_by_unknown_dependency() {
11407+
use SandboxRuntimeControlAvailability::{Available, Degraded, Unavailable, Unknown};
11408+
11409+
let combined =
11410+
combine_sandbox_runtime_control_availability([Degraded, Unknown, Available, Available]);
11411+
assert_eq!(combined, Degraded);
11412+
assert_eq!(
11413+
sandbox_runtime_readiness_state(combined),
11414+
Some("unavailable")
11415+
);
11416+
11417+
assert_eq!(
11418+
combine_sandbox_runtime_control_availability([
11419+
Degraded,
11420+
Unknown,
11421+
Unavailable,
11422+
Available
11423+
]),
11424+
Unavailable
11425+
);
11426+
let unknown = combine_sandbox_runtime_control_availability([
11427+
Available, Unknown, Available, Available,
11428+
]);
11429+
assert_eq!(unknown, Unknown);
11430+
assert_eq!(sandbox_runtime_readiness_state(unknown), None);
11431+
assert_eq!(
11432+
combine_sandbox_runtime_control_availability([Available; 4]),
11433+
Available
11434+
);
11435+
}
11436+
11437+
#[test]
11438+
fn sandbox_runtime_unready_running_supervisor_is_degraded_not_unavailable() {
11439+
let pod_with = |status: serde_json::Value| Pod {
11440+
status: Some(serde_json::from_value(status).expect("valid Pod status")),
11441+
..Default::default()
11442+
};
11443+
11444+
let reconnecting = pod_with(serde_json::json!({
11445+
"phase": "Running",
11446+
"conditions": [{"type": "Ready", "status": "False"}]
11447+
}));
11448+
assert_eq!(
11449+
sandbox_runtime_control_availability_from_pod(&reconnecting),
11450+
SandboxRuntimeControlAvailability::Degraded
11451+
);
11452+
11453+
for phase in ["Failed", "Succeeded", "Pending"] {
11454+
let pod = pod_with(serde_json::json!({
11455+
"phase": phase,
11456+
"conditions": [{"type": "Ready", "status": "False"}]
11457+
}));
11458+
assert_eq!(
11459+
sandbox_runtime_control_availability_from_pod(&pod),
11460+
SandboxRuntimeControlAvailability::Unavailable,
11461+
"phase {phase}"
11462+
);
11463+
}
11464+
11465+
let mut deleting = reconnecting;
11466+
deleting.metadata.deletion_timestamp = Some(
11467+
k8s_openapi::apimachinery::pkg::apis::meta::v1::Time(k8s_openapi::chrono::Utc::now()),
11468+
);
11469+
assert_eq!(
11470+
sandbox_runtime_control_availability_from_pod(&deleting),
11471+
SandboxRuntimeControlAvailability::Unavailable
11472+
);
11473+
}
11474+
1137611475
#[test]
1137711476
fn sandbox_runtime_readiness_transitions_bump_the_watched_cr() {
1137811477
let unavailable = sandbox_runtime_readiness_transition_patch("42", "unavailable");

‎crates/openshell-server/src/persistence/mod.rs‎

Lines changed: 79 additions & 66 deletions
Original file line numberDiff line numberDiff line change
@@ -1104,7 +1104,7 @@ impl Store {
11041104
/// Update a protobuf message using CAS (compare-and-swap).
11051105
///
11061106
/// Fetches the current object, validates the expected version, applies the
1107-
/// mutation function, and attempts a single CAS write. Returns Conflict on
1107+
/// mutation function, and attempts a CAS write. Returns Conflict on
11081108
/// version mismatch for caller-driven retry.
11091109
///
11101110
/// # Arguments
@@ -1137,79 +1137,92 @@ impl Store {
11371137
+ Clone,
11381138
F: FnMut(&mut T),
11391139
{
1140-
// Fetch current object with authoritative resource_version
1141-
let current = self
1142-
.get_message::<T>(id)
1143-
.await?
1144-
.ok_or_else(|| PersistenceError::Database(format!("object {id} not found")))?;
1140+
const INTERNAL_CAS_ATTEMPTS: u32 = 5;
1141+
let mut attempt = 0;
1142+
loop {
1143+
attempt += 1;
1144+
// Fetch current object with authoritative resource_version
1145+
let current = self
1146+
.get_message::<T>(id)
1147+
.await?
1148+
.ok_or_else(|| PersistenceError::Database(format!("object {id} not found")))?;
1149+
1150+
let current_version = current.get_resource_version();
1151+
1152+
// Determine the version to use for CAS:
1153+
// - If expected_version is 0, use current version (internal operations)
1154+
// - Otherwise, validate that expected matches current (client-facing operations)
1155+
let cas_version = if expected_version == 0 {
1156+
current_version
1157+
} else {
1158+
if expected_version != current_version {
1159+
return Err(PersistenceError::Conflict {
1160+
current_resource_version: Some(current_version),
1161+
});
1162+
}
1163+
expected_version
1164+
};
11451165

1146-
let current_version = current.get_resource_version();
1166+
// Apply mutation
1167+
let mut updated = current.clone();
1168+
mutate(&mut updated);
1169+
1170+
// Serialize labels
1171+
let labels_map = updated.object_labels();
1172+
let labels_json = if labels_map.as_ref().is_none_or(HashMap::is_empty) {
1173+
None
1174+
} else {
1175+
Some(serde_json::to_string(&labels_map).map_err(|e| {
1176+
PersistenceError::Encode(format!("failed to serialize labels: {e}"))
1177+
})?)
1178+
};
11471179

1148-
// Determine the version to use for CAS:
1149-
// - If expected_version is 0, use current version (internal operations)
1150-
// - Otherwise, validate that expected matches current (client-facing operations)
1151-
let cas_version = if expected_version == 0 {
1152-
current_version
1153-
} else {
1154-
if expected_version != current_version {
1155-
return Err(PersistenceError::Conflict {
1156-
current_resource_version: Some(current_version),
1157-
});
1180+
if T::requires_workspace() && updated.object_workspace().is_empty() {
1181+
return Err(PersistenceError::Encode(format!(
1182+
"{} requires a non-empty workspace",
1183+
T::object_type(),
1184+
)));
11581185
}
1159-
expected_version
1160-
};
1161-
1162-
// Apply mutation
1163-
let mut updated = current.clone();
1164-
mutate(&mut updated);
11651186

1166-
// Serialize labels
1167-
let labels_map = updated.object_labels();
1168-
let labels_json = if labels_map.as_ref().is_none_or(HashMap::is_empty) {
1169-
None
1170-
} else {
1171-
Some(serde_json::to_string(&labels_map).map_err(|e| {
1172-
PersistenceError::Encode(format!("failed to serialize labels: {e}"))
1173-
})?)
1174-
};
1187+
if updated.object_name() != current.object_name() {
1188+
return Err(PersistenceError::Encode(format!(
1189+
"{} name cannot be changed after creation",
1190+
T::object_type(),
1191+
)));
1192+
}
11751193

1176-
if T::requires_workspace() && updated.object_workspace().is_empty() {
1177-
return Err(PersistenceError::Encode(format!(
1178-
"{} requires a non-empty workspace",
1179-
T::object_type(),
1180-
)));
1181-
}
1194+
if updated.object_workspace() != current.object_workspace() {
1195+
return Err(PersistenceError::Encode(format!(
1196+
"{} workspace cannot be changed after creation",
1197+
T::object_type(),
1198+
)));
1199+
}
11821200

1183-
if updated.object_name() != current.object_name() {
1184-
return Err(PersistenceError::Encode(format!(
1185-
"{} name cannot be changed after creation",
1186-
T::object_type(),
1187-
)));
1188-
}
1201+
let result = match self
1202+
.put_if(
1203+
T::object_type(),
1204+
updated.object_id(),
1205+
updated.object_name(),
1206+
updated.object_workspace(),
1207+
&updated.encode_to_vec(),
1208+
labels_json.as_deref(),
1209+
WriteCondition::MatchResourceVersion(cas_version),
1210+
)
1211+
.await
1212+
{
1213+
Ok(result) => result,
1214+
Err(PersistenceError::Conflict { .. })
1215+
if expected_version == 0 && attempt < INTERNAL_CAS_ATTEMPTS =>
1216+
{
1217+
continue;
1218+
}
1219+
Err(error) => return Err(error),
1220+
};
11891221

1190-
if updated.object_workspace() != current.object_workspace() {
1191-
return Err(PersistenceError::Encode(format!(
1192-
"{} workspace cannot be changed after creation",
1193-
T::object_type(),
1194-
)));
1222+
// Success - hydrate the new resource_version and return
1223+
updated.set_resource_version(result.resource_version);
1224+
return Ok(updated);
11951225
}
1196-
1197-
// Single-attempt CAS write - fails with Conflict on version mismatch
1198-
let result = self
1199-
.put_if(
1200-
T::object_type(),
1201-
updated.object_id(),
1202-
updated.object_name(),
1203-
updated.object_workspace(),
1204-
&updated.encode_to_vec(),
1205-
labels_json.as_deref(),
1206-
WriteCondition::MatchResourceVersion(cas_version),
1207-
)
1208-
.await?;
1209-
1210-
// Success - hydrate the new resource_version and return
1211-
updated.set_resource_version(result.resource_version);
1212-
Ok(updated)
12131226
}
12141227
}
12151228

0 commit comments

Comments
 (0)