Skip to content

Commit 8983642

Browse files
authored
fix(gateway): give image preparation its own deadline (#4038)
* feat(server): separate image preparation and admission deadlines Give sandbox image preparation its own deadline and start the admission deadline after preparation finishes. Persist both phases across gateway restarts and show recovery guidance for the phase that expired. Preserve the current service authorization schema and regenerate the Go bindings with the preparation timestamps. Fixes #3952 Related to #3955 Signed-off-by: Shiju <shiju@nvidia.com> * fix(cli): simplify preparation timeout fallback selection Use lazy Option fallbacks while preserving timeout messages and retained sandbox behavior. Signed-off-by: Shiju <shiju@nvidia.com> * docs(server): clarify admission timer prerequisites Signed-off-by: Shiju <shiju@nvidia.com> * fix(compute): enforce deadlines during initial sandbox create Release stalled create operations after preparation expires and preserve the timeout diagnosis across late driver results. Keep failed-create cleanup bound to its original attempt so another replica can retry safely. Signed-off-by: Shiju <shiju@nvidia.com> * test(server): satisfy deadline regression lints Drop the create-error mutex guard before matching its cloned value and use idiomatic iteration and timeout matching in the deadline fixtures. Signed-off-by: Shiju <shiju@nvidia.com> * fix(server): retain ownership of pending provisioning operations Keep submitted create and start requests alive after caller cancellation, monitor failure, or preparation timeout. Persist request ownership before dispatch and retain staged uploads while the driver response is pending. Require timeout cleanup to stop compute after driver settlement without discarding an active cleanup claim. Fence late result handling against newer operations, preserve failed-start recovery, and defer automatic restart while another request owns the sandbox. Expose pending ownership in CLI JSON. Add ordered multi-replica and cancellation regressions for the review findings. Signed-off-by: Shiju <shiju@nvidia.com> * fix(openshell): preserve tracing and accept ready create responses Carry the request span into the detached provisioning worker so compute driver calls remain attached to their parent trace after task handoff. Update the compensation regression to require no backend DELETE when the durable cleanup claim fails, matching the operation ownership requirement. Accept the gateway's current Ready snapshot when a sandbox becomes ready before CREATE returns. Do not require the client to observe an earlier provisioning phase. Cover the Ready-only watch and command attachment. Signed-off-by: Shiju <shiju@nvidia.com> --------- Signed-off-by: Shiju <shiju@nvidia.com>
1 parent a2429fc commit 8983642

19 files changed

Lines changed: 3569 additions & 604 deletions

File tree

‎crates/openshell-cli/src/run.rs‎

Lines changed: 152 additions & 30 deletions
Original file line numberDiff line numberDiff line change
@@ -800,11 +800,8 @@ pub async fn sandbox_create(
800800
// Non-interactive mode: track start time for timestamps.
801801
let provision_start = Instant::now();
802802

803-
// Don't use stop_on_terminal on the server — the Kubernetes CRD may
804-
// briefly report a stale Ready status before the controller reconciles
805-
// a newly created sandbox. Instead we handle termination client-side:
806-
// we wait until we have observed at least one non-Ready phase followed
807-
// by Ready (a genuine Provisioning → Ready transition).
803+
// Handle terminal states here so a provisional container exit can wait
804+
// for the supervisor's canonical-process result before cleanup.
808805
let sandbox_name = sandbox.object_name().to_string();
809806
let sandbox_workspace = sandbox.object_workspace().to_string();
810807
let mut stream = client
@@ -832,8 +829,6 @@ pub async fn sandbox_create(
832829
let mut last_sandbox = sandbox.clone();
833830
let mut last_error_reason = String::new();
834831
let mut last_condition_message = ready_false_condition_message(sandbox.status.as_ref());
835-
// Track whether we have seen a non-Ready phase during the watch.
836-
let mut saw_non_ready = SandboxPhase::try_from(sandbox.phase()) != Ok(SandboxPhase::Ready);
837832
let provision_timeout = Duration::from_secs(
838833
std::env::var("OPENSHELL_PROVISION_TIMEOUT")
839834
.ok()
@@ -911,10 +906,6 @@ pub async fn sandbox_create(
911906
last_condition_message = Some(message);
912907
}
913908

914-
if phase != SandboxPhase::Ready {
915-
saw_non_ready = true;
916-
}
917-
918909
let main_process_result = has_main_process_result(&s);
919910
if matches!(
920911
phase,
@@ -949,9 +940,10 @@ pub async fn sandbox_create(
949940
break;
950941
}
951942

952-
// Only accept Ready as terminal after we've observed a
953-
// non-Ready phase, proving the controller has reconciled.
954-
if saw_non_ready && phase == SandboxPhase::Ready {
943+
// The gateway owns readiness. Its initial watch snapshot may
944+
// already be Ready if provisioning finished before CREATE
945+
// returned; requiring an earlier phase would miss that state.
946+
if phase == SandboxPhase::Ready {
955947
if let Some(d) = display.as_interactive_mut() {
956948
d.clear();
957949
}
@@ -1218,25 +1210,31 @@ pub async fn sandbox_create(
12181210
SandboxPhase::Error => {
12191211
drop(stream);
12201212
drop(client);
1221-
let provisioning_timed_out = last_sandbox
1213+
let timed_out_provisioning = last_sandbox
12221214
.status
12231215
.as_ref()
12241216
.and_then(|status| status.provisioning.as_ref())
1225-
.is_some_and(|record| record.timeout_time.is_some());
1226-
let create_result = if provisioning_timed_out {
1227-
Err(miette::miette!(
1228-
"{last_error_reason}\nSandbox '{sandbox_name}' was retained. Inspect it with `openshell sandbox get {sandbox_name}`; repair its configuration, then run `openshell sandbox start {sandbox_name}` after cleanup completes."
1229-
))
1230-
} else if last_error_reason.is_empty() {
1231-
Err(miette::miette!(
1232-
"sandbox entered error phase while provisioning"
1233-
))
1234-
} else {
1235-
Err(miette::miette!(
1236-
"sandbox entered error phase while provisioning: {}",
1237-
last_error_reason
1238-
))
1239-
};
1217+
.filter(|record| record.timeout_time.is_some());
1218+
let create_result = timed_out_provisioning.map_or_else(
1219+
|| {
1220+
if last_error_reason.is_empty() {
1221+
Err(miette::miette!(
1222+
"sandbox entered error phase while provisioning"
1223+
))
1224+
} else {
1225+
Err(miette::miette!(
1226+
"sandbox entered error phase while provisioning: {}",
1227+
last_error_reason
1228+
))
1229+
}
1230+
},
1231+
|record| {
1232+
Err(miette::miette!(
1233+
"{}",
1234+
retained_sandbox_timeout_message(&sandbox_name, &last_error_reason, record)
1235+
))
1236+
},
1237+
);
12401238
finalize_sandbox_create_session(
12411239
&effective_server,
12421240
&sandbox_name,
@@ -1259,6 +1257,24 @@ pub async fn sandbox_create(
12591257
}
12601258
}
12611259

1260+
/// Use the persisted phase to select recovery guidance. Preparation may expire
1261+
/// before any policy is evaluated, so it must not tell the user to repair policy.
1262+
fn retained_sandbox_timeout_message(
1263+
sandbox_name: &str,
1264+
error_reason: &str,
1265+
record: &openshell_core::proto::SandboxProvisioning,
1266+
) -> String {
1267+
let recovery = if record.preparation_deadline.is_some() && record.admission_start_time.is_none()
1268+
{
1269+
"check image preparation and supervisor startup diagnostics and the gateway's `image_preparation_timeout_seconds` budget"
1270+
} else {
1271+
"repair its configuration"
1272+
};
1273+
format!(
1274+
"{error_reason}\nSandbox '{sandbox_name}' was retained. Inspect it with `openshell sandbox get {sandbox_name}`; {recovery}, then run `openshell sandbox start {sandbox_name}` after cleanup completes."
1275+
)
1276+
}
1277+
12621278
/// Resolved source for the `--from` flag on `sandbox create`.
12631279
#[derive(Debug)]
12641280
enum ResolvedSource {
@@ -2924,11 +2940,22 @@ fn sandbox_to_json(sandbox: &Sandbox) -> serde_json::Value {
29242940
"configuration_change_id": record.configuration_change_id,
29252941
"configuration_change_time": record.configuration_change_time.as_ref().map(ToString::to_string),
29262942
"first_rejection_time": record.first_rejection_time.as_ref().map(ToString::to_string),
2943+
"phase": if record.deadline.is_none() && record.timeout_time.is_none() {
2944+
"ready"
2945+
} else if record.preparation_deadline.is_some() && record.admission_start_time.is_none() {
2946+
"preparation"
2947+
} else {
2948+
"admission"
2949+
},
2950+
"preparation_deadline": record.preparation_deadline.as_ref().map(ToString::to_string),
2951+
"admission_start_time": record.admission_start_time.as_ref().map(ToString::to_string),
29272952
"deadline": record.deadline.as_ref().map(ToString::to_string),
29282953
"timeout_time": record.timeout_time.as_ref().map(ToString::to_string),
29292954
"cleanup_completed_time": record.cleanup_completed_time.as_ref().map(ToString::to_string),
29302955
"cleanup_error": record.cleanup_error,
29312956
"cleanup_retry_time": record.cleanup_retry_time.as_ref().map(ToString::to_string),
2957+
"driver_operation_pending": record.driver_operation_pending,
2958+
"driver_operation_id": record.driver_operation_id,
29322959
}));
29332960
serde_json::json!({
29342961
"id": sandbox.object_id(),
@@ -8131,6 +8158,101 @@ mod tests {
81318158
);
81328159
}
81338160

8161+
#[test]
8162+
fn retained_sandbox_timeout_message_matches_expired_phase() {
8163+
let mut record = openshell_core::proto::SandboxProvisioning {
8164+
preparation_deadline: openshell_core::time::timestamp_from_millis(1_800_000).ok(),
8165+
timeout_time: openshell_core::time::timestamp_from_millis(1_800_000).ok(),
8166+
..Default::default()
8167+
};
8168+
let message = super::retained_sandbox_timeout_message(
8169+
"cold-image",
8170+
"ImagePreparationTimedOut: preparation expired",
8171+
&record,
8172+
);
8173+
assert!(message.starts_with("ImagePreparationTimedOut: preparation expired\n"));
8174+
assert!(message.contains("Sandbox 'cold-image' was retained"));
8175+
assert!(message.contains("image preparation and supervisor startup diagnostics"));
8176+
assert!(message.contains("image_preparation_timeout_seconds"));
8177+
assert!(!message.contains("repair its configuration"));
8178+
assert!(message.contains("openshell sandbox get cold-image"));
8179+
assert!(message.contains("openshell sandbox start cold-image` after cleanup completes"));
8180+
8181+
// An admission timeout retains preparation timestamps. Its completed
8182+
// transition must select configuration repair rather than a larger budget.
8183+
record.admission_start_time = openshell_core::time::timestamp_from_millis(600_000).ok();
8184+
let admission_without_preparation = openshell_core::proto::SandboxProvisioning {
8185+
timeout_time: record.timeout_time,
8186+
..Default::default()
8187+
};
8188+
for admission_record in [&record, &admission_without_preparation] {
8189+
let message = super::retained_sandbox_timeout_message(
8190+
"invalid-policy",
8191+
"ProvisioningTimedOut: repair window expired",
8192+
admission_record,
8193+
);
8194+
assert!(message.contains("repair its configuration"));
8195+
assert!(!message.contains("image_preparation_timeout_seconds"));
8196+
assert!(
8197+
message.contains("openshell sandbox start invalid-policy` after cleanup completes")
8198+
);
8199+
}
8200+
}
8201+
8202+
#[test]
8203+
fn provisioning_json_exposes_pending_driver_operation() {
8204+
for pending in [true, false] {
8205+
let mut sandbox = Sandbox::default();
8206+
sandbox.set_phase(SandboxPhase::Provisioning.into());
8207+
sandbox.status.as_mut().unwrap().provisioning =
8208+
Some(openshell_core::proto::SandboxProvisioning {
8209+
driver_operation_pending: pending,
8210+
driver_operation_id: "operation-1".into(),
8211+
..Default::default()
8212+
});
8213+
assert_eq!(
8214+
super::sandbox_to_json(&sandbox)["provisioning"]["driver_operation_pending"],
8215+
pending
8216+
);
8217+
assert_eq!(
8218+
super::sandbox_to_json(&sandbox)["provisioning"]["driver_operation_id"],
8219+
"operation-1"
8220+
);
8221+
}
8222+
}
8223+
8224+
#[test]
8225+
fn provisioning_json_distinguishes_preparation_and_admission() {
8226+
let mut sandbox = Sandbox::default();
8227+
sandbox.set_phase(SandboxPhase::Provisioning.into());
8228+
let ceiling = openshell_core::time::timestamp_from_millis(1_800_000).ok();
8229+
sandbox.status.as_mut().unwrap().provisioning =
8230+
Some(openshell_core::proto::SandboxProvisioning {
8231+
preparation_deadline: ceiling,
8232+
deadline: ceiling,
8233+
..Default::default()
8234+
});
8235+
let json = super::sandbox_to_json(&sandbox);
8236+
assert_eq!(json["provisioning"]["phase"], "preparation");
8237+
assert_eq!(
8238+
json["provisioning"]["preparation_deadline"],
8239+
"1970-01-01T00:30:00Z"
8240+
);
8241+
assert!(json["provisioning"]["admission_start_time"].is_null());
8242+
sandbox
8243+
.status
8244+
.as_mut()
8245+
.unwrap()
8246+
.provisioning
8247+
.as_mut()
8248+
.unwrap()
8249+
.admission_start_time = openshell_core::time::timestamp_from_millis(600_000).ok();
8250+
assert_eq!(
8251+
super::sandbox_to_json(&sandbox)["provisioning"]["phase"],
8252+
"admission"
8253+
);
8254+
}
8255+
81348256
#[test]
81358257
fn sandbox_json_exposes_repair_diagnostic_and_accepted_generation() {
81368258
use openshell_core::proto::{ConfigurationAdmissionState, SandboxConfigurationAdmission};

‎crates/openshell-cli/tests/sandbox_create_lifecycle_integration.rs‎

Lines changed: 77 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -70,6 +70,7 @@ struct SandboxState {
7070
vm_error_with_observed_exit: Arc<AtomicBool>,
7171
vm_slow_progress_before_ready: Arc<AtomicBool>,
7272
vm_log_churn_before_ready: Arc<AtomicBool>,
73+
ready_before_create_returns: Arc<AtomicBool>,
7374
terminal_before_relay: Arc<AtomicBool>,
7475
terminal_after_provisional_container_exit: Arc<AtomicBool>,
7576
provisional_container_exit_without_result: Arc<AtomicBool>,
@@ -194,7 +195,17 @@ impl OpenShell for TestOpenShell {
194195
}),
195196
..Sandbox::default()
196197
};
197-
sandbox.set_phase(SandboxPhase::Provisioning as i32);
198+
sandbox.set_phase(
199+
if self
200+
.state
201+
.ready_before_create_returns
202+
.load(Ordering::SeqCst)
203+
{
204+
SandboxPhase::Ready as i32
205+
} else {
206+
SandboxPhase::Provisioning as i32
207+
},
208+
);
198209
Ok(Response::new(SandboxResponse {
199210
sandbox: Some(sandbox),
200211
service_urls,
@@ -677,6 +688,10 @@ impl OpenShell for TestOpenShell {
677688
.vm_slow_progress_before_ready
678689
.load(Ordering::SeqCst);
679690
let vm_log_churn_before_ready = self.state.vm_log_churn_before_ready.load(Ordering::SeqCst);
691+
let ready_before_create_returns = self
692+
.state
693+
.ready_before_create_returns
694+
.load(Ordering::SeqCst);
680695
let terminal_before_relay = self.state.terminal_before_relay.load(Ordering::SeqCst);
681696
let terminal_after_provisional_container_exit = self
682697
.state
@@ -723,6 +738,18 @@ impl OpenShell for TestOpenShell {
723738
}
724739
let mut ready = provisioning.clone();
725740
ready.set_phase(SandboxPhase::Ready as i32);
741+
if ready_before_create_returns {
742+
// A watch starts with the current snapshot. Keep it open so
743+
// stream closure cannot hide a client that ignores Ready.
744+
let _ = tx
745+
.send(Ok(SandboxStreamEvent {
746+
payload: Some(sandbox_stream_event::Payload::Sandbox(ready)),
747+
cursor: String::new(),
748+
}))
749+
.await;
750+
tx.closed().await;
751+
return;
752+
}
726753
let mut completed = provisioning.clone();
727754
completed.status = Some(SandboxStatus {
728755
exit_code: Some(0),
@@ -2372,6 +2399,55 @@ async fn sandbox_create_preserves_vm_error_when_exit_code_is_observed() {
23722399
assert!(rendered.contains("ProcessExited: VM process exited with status 0"));
23732400
}
23742401

2402+
#[tokio::test]
2403+
async fn sandbox_create_accepts_ready_before_create_returns() {
2404+
let server = run_server().await;
2405+
server
2406+
.openshell
2407+
.state
2408+
.ready_before_create_returns
2409+
.store(true, Ordering::SeqCst);
2410+
let fake_ssh_dir = tempfile::tempdir().unwrap();
2411+
let xdg_dir = tempfile::tempdir().unwrap();
2412+
let _env = test_env_with(
2413+
&fake_ssh_dir,
2414+
&xdg_dir,
2415+
&[("OPENSHELL_PROVISION_TIMEOUT", "1".to_string())],
2416+
);
2417+
let tls = test_tls(&server);
2418+
install_fake_ssh(&fake_ssh_dir);
2419+
2420+
let exit_code = tokio::time::timeout(
2421+
Duration::from_secs(10),
2422+
run::sandbox_create(
2423+
&server.endpoint,
2424+
"openshell",
2425+
run::SandboxCreateConfig {
2426+
name: Some("already-ready"),
2427+
command: &["echo".into(), "OK".into()],
2428+
..test_config()
2429+
},
2430+
"default",
2431+
&tls,
2432+
),
2433+
)
2434+
.await
2435+
.expect("creation must finish while the watch remains open")
2436+
.expect("an already-Ready sandbox must not wait for a new provisioning transition");
2437+
2438+
assert_eq!(exit_code, 0);
2439+
assert_eq!(create_requests(&server).await.len(), 1);
2440+
assert_eq!(
2441+
server
2442+
.openshell
2443+
.state
2444+
.ssh_session_requests
2445+
.load(Ordering::SeqCst),
2446+
1,
2447+
"the initial Ready snapshot must allow the command to attach"
2448+
);
2449+
}
2450+
23752451
#[tokio::test]
23762452
async fn sandbox_create_keeps_waiting_while_vm_progress_arrives() {
23772453
let server = run_server().await;

‎crates/openshell-core/src/config.rs‎

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -240,6 +240,10 @@ pub struct Config {
240240
/// TTL for SSH session tokens, in seconds. 0 disables expiry.
241241
pub ssh_session_ttl_secs: u64,
242242

243+
/// Absolute image preparation and initial supervisor startup budget for new
244+
/// sandbox attempts, in seconds. Must be between 1 and 86400, inclusive.
245+
pub image_preparation_timeout_seconds: u32,
246+
243247
/// Maximum gRPC requests allowed per rate-limit window.
244248
///
245249
/// When paired with [`Self::grpc_rate_limit_window_secs`], positive values
@@ -865,6 +869,7 @@ impl Config {
865869
credential_drivers: Vec::new(),
866870
default_credential_driver: None,
867871
ssh_session_ttl_secs: default_ssh_session_ttl_secs(),
872+
image_preparation_timeout_seconds: 1800,
868873
grpc_rate_limit_requests: None,
869874
grpc_rate_limit_window_secs: None,
870875
service_routing: ServiceRoutingConfig::default(),

0 commit comments

Comments
 (0)