Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
182 changes: 152 additions & 30 deletions crates/openshell-cli/src/run.rs
Original file line number Diff line number Diff line change
Expand Up @@ -800,11 +800,8 @@ pub async fn sandbox_create(
// Non-interactive mode: track start time for timestamps.
let provision_start = Instant::now();

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

if phase != SandboxPhase::Ready {
saw_non_ready = true;
}

let main_process_result = has_main_process_result(&s);
if matches!(
phase,
Expand Down Expand Up @@ -949,9 +940,10 @@ pub async fn sandbox_create(
break;
}

// Only accept Ready as terminal after we've observed a
// non-Ready phase, proving the controller has reconciled.
if saw_non_ready && phase == SandboxPhase::Ready {
// The gateway owns readiness. Its initial watch snapshot may
// already be Ready if provisioning finished before CREATE
// returned; requiring an earlier phase would miss that state.
if phase == SandboxPhase::Ready {
if let Some(d) = display.as_interactive_mut() {
d.clear();
}
Expand Down Expand Up @@ -1218,25 +1210,31 @@ pub async fn sandbox_create(
SandboxPhase::Error => {
drop(stream);
drop(client);
let provisioning_timed_out = last_sandbox
let timed_out_provisioning = last_sandbox
.status
.as_ref()
.and_then(|status| status.provisioning.as_ref())
.is_some_and(|record| record.timeout_time.is_some());
let create_result = if provisioning_timed_out {
Err(miette::miette!(
"{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."
))
} else if last_error_reason.is_empty() {
Err(miette::miette!(
"sandbox entered error phase while provisioning"
))
} else {
Err(miette::miette!(
"sandbox entered error phase while provisioning: {}",
last_error_reason
))
};
.filter(|record| record.timeout_time.is_some());
let create_result = timed_out_provisioning.map_or_else(
|| {
if last_error_reason.is_empty() {
Err(miette::miette!(
"sandbox entered error phase while provisioning"
))
} else {
Err(miette::miette!(
"sandbox entered error phase while provisioning: {}",
last_error_reason
))
}
},
|record| {
Err(miette::miette!(
"{}",
retained_sandbox_timeout_message(&sandbox_name, &last_error_reason, record)
))
},
);
finalize_sandbox_create_session(
&effective_server,
&sandbox_name,
Expand All @@ -1259,6 +1257,24 @@ pub async fn sandbox_create(
}
}

/// Use the persisted phase to select recovery guidance. Preparation may expire
/// before any policy is evaluated, so it must not tell the user to repair policy.
fn retained_sandbox_timeout_message(
sandbox_name: &str,
error_reason: &str,
record: &openshell_core::proto::SandboxProvisioning,
) -> String {
let recovery = if record.preparation_deadline.is_some() && record.admission_start_time.is_none()
{
"check image preparation and supervisor startup diagnostics and the gateway's `image_preparation_timeout_seconds` budget"
} else {
"repair its configuration"
};
format!(
"{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."
)
}

/// Resolved source for the `--from` flag on `sandbox create`.
#[derive(Debug)]
enum ResolvedSource {
Expand Down Expand Up @@ -2904,11 +2920,22 @@ fn sandbox_to_json(sandbox: &Sandbox) -> serde_json::Value {
"configuration_change_id": record.configuration_change_id,
"configuration_change_time": record.configuration_change_time.as_ref().map(ToString::to_string),
"first_rejection_time": record.first_rejection_time.as_ref().map(ToString::to_string),
"phase": if record.deadline.is_none() && record.timeout_time.is_none() {
"ready"
} else if record.preparation_deadline.is_some() && record.admission_start_time.is_none() {
"preparation"
} else {
"admission"
},
"preparation_deadline": record.preparation_deadline.as_ref().map(ToString::to_string),
"admission_start_time": record.admission_start_time.as_ref().map(ToString::to_string),
"deadline": record.deadline.as_ref().map(ToString::to_string),
"timeout_time": record.timeout_time.as_ref().map(ToString::to_string),
"cleanup_completed_time": record.cleanup_completed_time.as_ref().map(ToString::to_string),
"cleanup_error": record.cleanup_error,
"cleanup_retry_time": record.cleanup_retry_time.as_ref().map(ToString::to_string),
"driver_operation_pending": record.driver_operation_pending,
"driver_operation_id": record.driver_operation_id,
}));
serde_json::json!({
"id": sandbox.object_id(),
Expand Down Expand Up @@ -8110,6 +8137,101 @@ mod tests {
);
}

#[test]
fn retained_sandbox_timeout_message_matches_expired_phase() {
let mut record = openshell_core::proto::SandboxProvisioning {
preparation_deadline: openshell_core::time::timestamp_from_millis(1_800_000).ok(),
timeout_time: openshell_core::time::timestamp_from_millis(1_800_000).ok(),
..Default::default()
};
let message = super::retained_sandbox_timeout_message(
"cold-image",
"ImagePreparationTimedOut: preparation expired",
&record,
);
assert!(message.starts_with("ImagePreparationTimedOut: preparation expired\n"));
assert!(message.contains("Sandbox 'cold-image' was retained"));
assert!(message.contains("image preparation and supervisor startup diagnostics"));
assert!(message.contains("image_preparation_timeout_seconds"));
assert!(!message.contains("repair its configuration"));
assert!(message.contains("openshell sandbox get cold-image"));
assert!(message.contains("openshell sandbox start cold-image` after cleanup completes"));

// An admission timeout retains preparation timestamps. Its completed
// transition must select configuration repair rather than a larger budget.
record.admission_start_time = openshell_core::time::timestamp_from_millis(600_000).ok();
let admission_without_preparation = openshell_core::proto::SandboxProvisioning {
timeout_time: record.timeout_time,
..Default::default()
};
for admission_record in [&record, &admission_without_preparation] {
let message = super::retained_sandbox_timeout_message(
"invalid-policy",
"ProvisioningTimedOut: repair window expired",
admission_record,
);
assert!(message.contains("repair its configuration"));
assert!(!message.contains("image_preparation_timeout_seconds"));
assert!(
message.contains("openshell sandbox start invalid-policy` after cleanup completes")
);
}
}

#[test]
fn provisioning_json_exposes_pending_driver_operation() {
for pending in [true, false] {
let mut sandbox = Sandbox::default();
sandbox.set_phase(SandboxPhase::Provisioning.into());
sandbox.status.as_mut().unwrap().provisioning =
Some(openshell_core::proto::SandboxProvisioning {
driver_operation_pending: pending,
driver_operation_id: "operation-1".into(),
..Default::default()
});
assert_eq!(
super::sandbox_to_json(&sandbox)["provisioning"]["driver_operation_pending"],
pending
);
assert_eq!(
super::sandbox_to_json(&sandbox)["provisioning"]["driver_operation_id"],
"operation-1"
);
}
}

#[test]
fn provisioning_json_distinguishes_preparation_and_admission() {
let mut sandbox = Sandbox::default();
sandbox.set_phase(SandboxPhase::Provisioning.into());
let ceiling = openshell_core::time::timestamp_from_millis(1_800_000).ok();
sandbox.status.as_mut().unwrap().provisioning =
Some(openshell_core::proto::SandboxProvisioning {
preparation_deadline: ceiling,
deadline: ceiling,
..Default::default()
});
let json = super::sandbox_to_json(&sandbox);
assert_eq!(json["provisioning"]["phase"], "preparation");
assert_eq!(
json["provisioning"]["preparation_deadline"],
"1970-01-01T00:30:00Z"
);
assert!(json["provisioning"]["admission_start_time"].is_null());
sandbox
.status
.as_mut()
.unwrap()
.provisioning
.as_mut()
.unwrap()
.admission_start_time = openshell_core::time::timestamp_from_millis(600_000).ok();
assert_eq!(
super::sandbox_to_json(&sandbox)["provisioning"]["phase"],
"admission"
);
}

#[test]
fn sandbox_json_exposes_repair_diagnostic_and_accepted_generation() {
use openshell_core::proto::{ConfigurationAdmissionState, SandboxConfigurationAdmission};
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -70,6 +70,7 @@ struct SandboxState {
vm_error_with_observed_exit: Arc<AtomicBool>,
vm_slow_progress_before_ready: Arc<AtomicBool>,
vm_log_churn_before_ready: Arc<AtomicBool>,
ready_before_create_returns: Arc<AtomicBool>,
terminal_before_relay: Arc<AtomicBool>,
terminal_after_provisional_container_exit: Arc<AtomicBool>,
provisional_container_exit_without_result: Arc<AtomicBool>,
Expand Down Expand Up @@ -194,7 +195,17 @@ impl OpenShell for TestOpenShell {
}),
..Sandbox::default()
};
sandbox.set_phase(SandboxPhase::Provisioning as i32);
sandbox.set_phase(
if self
.state
.ready_before_create_returns
.load(Ordering::SeqCst)
{
SandboxPhase::Ready as i32
} else {
SandboxPhase::Provisioning as i32
},
);
Ok(Response::new(SandboxResponse {
sandbox: Some(sandbox),
service_urls,
Expand Down Expand Up @@ -677,6 +688,10 @@ impl OpenShell for TestOpenShell {
.vm_slow_progress_before_ready
.load(Ordering::SeqCst);
let vm_log_churn_before_ready = self.state.vm_log_churn_before_ready.load(Ordering::SeqCst);
let ready_before_create_returns = self
.state
.ready_before_create_returns
.load(Ordering::SeqCst);
let terminal_before_relay = self.state.terminal_before_relay.load(Ordering::SeqCst);
let terminal_after_provisional_container_exit = self
.state
Expand Down Expand Up @@ -723,6 +738,18 @@ impl OpenShell for TestOpenShell {
}
let mut ready = provisioning.clone();
ready.set_phase(SandboxPhase::Ready as i32);
if ready_before_create_returns {
// A watch starts with the current snapshot. Keep it open so
// stream closure cannot hide a client that ignores Ready.
let _ = tx
.send(Ok(SandboxStreamEvent {
payload: Some(sandbox_stream_event::Payload::Sandbox(ready)),
cursor: String::new(),
}))
.await;
tx.closed().await;
return;
}
let mut completed = provisioning.clone();
completed.status = Some(SandboxStatus {
exit_code: Some(0),
Expand Down Expand Up @@ -2372,6 +2399,55 @@ async fn sandbox_create_preserves_vm_error_when_exit_code_is_observed() {
assert!(rendered.contains("ProcessExited: VM process exited with status 0"));
}

#[tokio::test]
async fn sandbox_create_accepts_ready_before_create_returns() {
let server = run_server().await;
server
.openshell
.state
.ready_before_create_returns
.store(true, Ordering::SeqCst);
let fake_ssh_dir = tempfile::tempdir().unwrap();
let xdg_dir = tempfile::tempdir().unwrap();
let _env = test_env_with(
&fake_ssh_dir,
&xdg_dir,
&[("OPENSHELL_PROVISION_TIMEOUT", "1".to_string())],
);
let tls = test_tls(&server);
install_fake_ssh(&fake_ssh_dir);

let exit_code = tokio::time::timeout(
Duration::from_secs(10),
run::sandbox_create(
&server.endpoint,
"openshell",
run::SandboxCreateConfig {
name: Some("already-ready"),
command: &["echo".into(), "OK".into()],
..test_config()
},
"default",
&tls,
),
)
.await
.expect("creation must finish while the watch remains open")
.expect("an already-Ready sandbox must not wait for a new provisioning transition");

assert_eq!(exit_code, 0);
assert_eq!(create_requests(&server).await.len(), 1);
assert_eq!(
server
.openshell
.state
.ssh_session_requests
.load(Ordering::SeqCst),
1,
"the initial Ready snapshot must allow the command to attach"
);
}

#[tokio::test]
async fn sandbox_create_keeps_waiting_while_vm_progress_arrives() {
let server = run_server().await;
Expand Down
5 changes: 5 additions & 0 deletions crates/openshell-core/src/config.rs
Original file line number Diff line number Diff line change
Expand Up @@ -240,6 +240,10 @@ pub struct Config {
/// TTL for SSH session tokens, in seconds. 0 disables expiry.
pub ssh_session_ttl_secs: u64,

/// Absolute image preparation and initial supervisor startup budget for new
/// sandbox attempts, in seconds. Must be between 1 and 86400, inclusive.
pub image_preparation_timeout_seconds: u32,

/// Maximum gRPC requests allowed per rate-limit window.
///
/// When paired with [`Self::grpc_rate_limit_window_secs`], positive values
Expand Down Expand Up @@ -865,6 +869,7 @@ impl Config {
credential_drivers: Vec::new(),
default_credential_driver: None,
ssh_session_ttl_secs: default_ssh_session_ttl_secs(),
image_preparation_timeout_seconds: 1800,
grpc_rate_limit_requests: None,
grpc_rate_limit_window_secs: None,
service_routing: ServiceRoutingConfig::default(),
Expand Down
Loading
Loading