mirror of
https://github.com/NVIDIA/OpenShell.git
synced 2026-10-11 04:30:53 +08:00
3275 lines
107 KiB
Rust
3275 lines
107 KiB
Rust
// SPDX-FileCopyrightText: Copyright (c) 2025-2026 NVIDIA CORPORATION & AFFILIATES. All rights reserved.
|
|
// SPDX-License-Identifier: Apache-2.0
|
|
|
|
#![cfg(not(target_os = "windows"))]
|
|
|
|
mod helpers;
|
|
|
|
use helpers::{EnvVarGuard, build_ca, build_client_cert, build_server_cert};
|
|
use openshell_bootstrap::load_last_sandbox;
|
|
use openshell_cli::run;
|
|
use openshell_cli::tls::TlsOptions;
|
|
use openshell_core::proto::open_shell_server::{OpenShell, OpenShellServer};
|
|
use openshell_core::proto::{
|
|
AttachSandboxProviderRequest, AttachSandboxProviderResponse, CreateProviderRequest,
|
|
CreateSandboxRequest, CreateSandboxTemplateRequest, CreateSshSessionRequest,
|
|
CreateSshSessionResponse, DeleteProviderRequest, DeleteProviderResponse, DeleteSandboxRequest,
|
|
DeleteSandboxResponse, DeleteSandboxTemplateRequest, DetachSandboxProviderRequest,
|
|
DetachSandboxProviderResponse, ExchangeProviderSubjectTokenRequest,
|
|
ExchangeProviderSubjectTokenResponse, ExecSandboxEvent, ExecSandboxInput, ExecSandboxRequest,
|
|
GatewayMessage, GetGatewayConfigRequest, GetGatewayConfigResponse, GetProviderRequest,
|
|
GetSandboxConfigRequest, GetSandboxConfigResponse, GetSandboxProviderEnvironmentRequest,
|
|
GetSandboxProviderEnvironmentResponse, GetSandboxRequest, GetSandboxTemplateRequest,
|
|
GpuResourceRequirements, HealthRequest, HealthResponse, ListProvidersRequest,
|
|
ListProvidersResponse, ListSandboxProvidersRequest, ListSandboxProvidersResponse,
|
|
ListSandboxTemplatesRequest, ListSandboxTemplatesResponse, ListSandboxesRequest,
|
|
ListSandboxesResponse, PlatformEvent, Provider, ProviderResponse, ReportEndpointStatusRequest,
|
|
ReportEndpointStatusResponse, RevokeSshSessionRequest, RevokeSshSessionResponse, Sandbox,
|
|
SandboxCondition, SandboxLogLine, SandboxPhase, SandboxResponse, SandboxStatus,
|
|
SandboxStreamEvent, SandboxTemplateResponse, SandboxWorkloadTemplate, ServiceStatus,
|
|
SettingValue, SupervisorMessage, UpdateProviderRequest, WatchSandboxRequest,
|
|
sandbox_stream_event,
|
|
};
|
|
use std::collections::HashMap;
|
|
use std::fs;
|
|
use std::os::unix::fs::PermissionsExt;
|
|
use std::sync::Arc;
|
|
use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering};
|
|
use std::time::{Duration, Instant};
|
|
use tempfile::TempDir;
|
|
use tokio::io::{AsyncReadExt, AsyncWriteExt};
|
|
use tokio::net::TcpListener;
|
|
use tokio::sync::{Mutex, Notify, mpsc};
|
|
use tokio_stream::wrappers::TcpListenerStream;
|
|
use tonic::transport::{Certificate as TlsCertificate, Identity, Server, ServerTlsConfig};
|
|
use tonic::{Response, Status};
|
|
|
|
fn selected_workspace(
|
|
scope: &Option<openshell_core::proto::datamodel::v1::WorkspaceSelector>,
|
|
) -> Option<&str> {
|
|
match scope.as_ref()?.selection.as_ref()? {
|
|
openshell_core::proto::datamodel::v1::workspace_selector::Selection::Workspace(
|
|
workspace,
|
|
) => Some(workspace),
|
|
openshell_core::proto::datamodel::v1::workspace_selector::Selection::AllWorkspaces(_) => {
|
|
None
|
|
}
|
|
}
|
|
}
|
|
|
|
#[derive(Clone, Default)]
|
|
struct SandboxState {
|
|
deleted_names: Arc<Mutex<Vec<Vec<String>>>>,
|
|
delete_requests: Arc<Mutex<Vec<DeleteSandboxRequest>>>,
|
|
create_requests: Arc<Mutex<Vec<CreateSandboxRequest>>>,
|
|
fail_delete_sandbox_message: Arc<Mutex<Option<String>>>,
|
|
vm_error_after_started: Arc<AtomicBool>,
|
|
vm_error_with_observed_exit: Arc<AtomicBool>,
|
|
vm_slow_progress_before_ready: Arc<AtomicBool>,
|
|
vm_log_churn_before_ready: Arc<AtomicBool>,
|
|
terminal_before_relay: Arc<AtomicBool>,
|
|
terminal_after_provisional_container_exit: Arc<AtomicBool>,
|
|
provisional_container_exit_without_result: Arc<AtomicBool>,
|
|
provisional_container_exit_sent: Arc<Notify>,
|
|
ssh_session_failures_remaining: Arc<AtomicUsize>,
|
|
ssh_session_requests: Arc<AtomicUsize>,
|
|
global_settings: Arc<Mutex<HashMap<String, SettingValue>>>,
|
|
gateway_config_requests: Arc<AtomicUsize>,
|
|
providers: Arc<Mutex<Vec<Provider>>>,
|
|
template_create_requests: Arc<Mutex<Vec<CreateSandboxTemplateRequest>>>,
|
|
template_get_requests: Arc<Mutex<Vec<GetSandboxTemplateRequest>>>,
|
|
template_list_requests: Arc<Mutex<Vec<ListSandboxTemplatesRequest>>>,
|
|
template_delete_requests: Arc<Mutex<Vec<DeleteSandboxTemplateRequest>>>,
|
|
}
|
|
|
|
#[derive(Clone, Default)]
|
|
struct TestOpenShell {
|
|
state: SandboxState,
|
|
}
|
|
|
|
#[tonic::async_trait]
|
|
impl OpenShell for TestOpenShell {
|
|
async fn begin_rootfs_tar_staging(
|
|
&self,
|
|
_request: tonic::Request<openshell_core::proto::BeginRootfsTarStagingRequest>,
|
|
) -> Result<Response<openshell_core::proto::BeginRootfsTarStagingResponse>, Status> {
|
|
Err(Status::unimplemented("not used by this test server"))
|
|
}
|
|
|
|
async fn report_main_process_exit(
|
|
&self,
|
|
_request: tonic::Request<openshell_core::proto::ReportMainProcessExitRequest>,
|
|
) -> Result<Response<openshell_core::proto::ReportMainProcessExitResponse>, Status> {
|
|
Err(Status::unimplemented("not used by this test server"))
|
|
}
|
|
|
|
async fn finalize_main_process_exit(
|
|
&self,
|
|
_request: tonic::Request<openshell_core::proto::FinalizeMainProcessExitRequest>,
|
|
) -> Result<Response<openshell_core::proto::FinalizeMainProcessExitResponse>, Status> {
|
|
Err(Status::unimplemented("not used by this test server"))
|
|
}
|
|
|
|
async fn get_current_user(
|
|
&self,
|
|
_request: tonic::Request<openshell_core::proto::GetCurrentUserRequest>,
|
|
) -> Result<Response<openshell_core::proto::GetCurrentUserResponse>, Status> {
|
|
Err(Status::unimplemented("not used by this test server"))
|
|
}
|
|
|
|
async fn health(
|
|
&self,
|
|
_request: tonic::Request<HealthRequest>,
|
|
) -> Result<Response<HealthResponse>, Status> {
|
|
Ok(Response::new(HealthResponse {
|
|
status: ServiceStatus::Healthy.into(),
|
|
version: "test".to_string(),
|
|
}))
|
|
}
|
|
|
|
async fn get_gateway_info(
|
|
&self,
|
|
_request: tonic::Request<openshell_core::proto::GetGatewayInfoRequest>,
|
|
) -> Result<Response<openshell_core::proto::GetGatewayInfoResponse>, Status> {
|
|
Err(Status::unimplemented("unused"))
|
|
}
|
|
|
|
async fn create_sandbox(
|
|
&self,
|
|
request: tonic::Request<CreateSandboxRequest>,
|
|
) -> Result<Response<SandboxResponse>, Status> {
|
|
let request = request.into_inner();
|
|
let name = request.name.clone();
|
|
self.state.create_requests.lock().await.push(request);
|
|
let sandbox_name = if name.is_empty() {
|
|
"test-sandbox".to_string()
|
|
} else {
|
|
name
|
|
};
|
|
|
|
let mut sandbox = Sandbox {
|
|
metadata: Some(openshell_core::proto::datamodel::v1::ObjectMeta {
|
|
id: format!("id-{sandbox_name}"),
|
|
name: sandbox_name,
|
|
created_at_ms: 0,
|
|
labels: HashMap::new(),
|
|
resource_version: 0,
|
|
annotations: HashMap::new(),
|
|
workspace: String::new(),
|
|
deletion_timestamp_ms: 0,
|
|
}),
|
|
..Sandbox::default()
|
|
};
|
|
sandbox.set_phase(SandboxPhase::Provisioning as i32);
|
|
Ok(Response::new(SandboxResponse {
|
|
sandbox: Some(sandbox),
|
|
}))
|
|
}
|
|
|
|
async fn stop_sandbox(
|
|
&self,
|
|
_request: tonic::Request<openshell_core::proto::StopSandboxRequest>,
|
|
) -> Result<Response<SandboxResponse>, Status> {
|
|
Err(Status::unimplemented("unused"))
|
|
}
|
|
|
|
async fn start_sandbox(
|
|
&self,
|
|
_request: tonic::Request<openshell_core::proto::StartSandboxRequest>,
|
|
) -> Result<Response<SandboxResponse>, Status> {
|
|
Err(Status::unimplemented("unused"))
|
|
}
|
|
|
|
async fn get_sandbox(
|
|
&self,
|
|
request: tonic::Request<GetSandboxRequest>,
|
|
) -> Result<Response<SandboxResponse>, Status> {
|
|
let name = request.into_inner().name;
|
|
let mut sandbox = Sandbox {
|
|
metadata: Some(openshell_core::proto::datamodel::v1::ObjectMeta {
|
|
id: format!("id-{name}"),
|
|
name,
|
|
created_at_ms: 0,
|
|
labels: HashMap::new(),
|
|
resource_version: 0,
|
|
annotations: HashMap::new(),
|
|
workspace: String::new(),
|
|
deletion_timestamp_ms: 0,
|
|
}),
|
|
..Sandbox::default()
|
|
};
|
|
sandbox.set_phase(SandboxPhase::Ready as i32);
|
|
Ok(Response::new(SandboxResponse {
|
|
sandbox: Some(sandbox),
|
|
}))
|
|
}
|
|
|
|
async fn list_sandboxes(
|
|
&self,
|
|
_request: tonic::Request<ListSandboxesRequest>,
|
|
) -> Result<Response<ListSandboxesResponse>, Status> {
|
|
Ok(Response::new(ListSandboxesResponse::default()))
|
|
}
|
|
|
|
async fn create_sandbox_template(
|
|
&self,
|
|
request: tonic::Request<CreateSandboxTemplateRequest>,
|
|
) -> Result<Response<SandboxTemplateResponse>, Status> {
|
|
let request = request.into_inner();
|
|
let mut template = request.template.clone().unwrap_or_default();
|
|
let name = template
|
|
.metadata
|
|
.as_ref()
|
|
.map_or_else(|| "template".to_string(), |metadata| metadata.name.clone());
|
|
template.metadata = Some(openshell_core::proto::datamodel::v1::ObjectMeta {
|
|
id: format!("template-{name}"),
|
|
name,
|
|
created_at_ms: 0,
|
|
labels: template
|
|
.metadata
|
|
.as_ref()
|
|
.map(|metadata| metadata.labels.clone())
|
|
.unwrap_or_default(),
|
|
resource_version: 1,
|
|
annotations: HashMap::new(),
|
|
workspace: selected_workspace(&request.workspace_scope)
|
|
.unwrap_or("default")
|
|
.to_string(),
|
|
deletion_timestamp_ms: 0,
|
|
});
|
|
self.state
|
|
.template_create_requests
|
|
.lock()
|
|
.await
|
|
.push(request);
|
|
Ok(Response::new(SandboxTemplateResponse {
|
|
template: Some(template),
|
|
}))
|
|
}
|
|
|
|
async fn get_sandbox_template(
|
|
&self,
|
|
request: tonic::Request<GetSandboxTemplateRequest>,
|
|
) -> Result<Response<SandboxTemplateResponse>, Status> {
|
|
let request = request.into_inner();
|
|
self.state
|
|
.template_get_requests
|
|
.lock()
|
|
.await
|
|
.push(request.clone());
|
|
Ok(Response::new(SandboxTemplateResponse {
|
|
template: Some(SandboxWorkloadTemplate {
|
|
metadata: Some(openshell_core::proto::datamodel::v1::ObjectMeta {
|
|
id: format!("template-{}", request.name),
|
|
name: request.name,
|
|
created_at_ms: 0,
|
|
labels: HashMap::new(),
|
|
resource_version: 1,
|
|
annotations: HashMap::new(),
|
|
workspace: selected_workspace(&request.workspace_scope)
|
|
.unwrap_or("default")
|
|
.to_string(),
|
|
deletion_timestamp_ms: 0,
|
|
}),
|
|
spec: None,
|
|
}),
|
|
}))
|
|
}
|
|
|
|
async fn list_sandbox_templates(
|
|
&self,
|
|
request: tonic::Request<ListSandboxTemplatesRequest>,
|
|
) -> Result<Response<ListSandboxTemplatesResponse>, Status> {
|
|
self.state
|
|
.template_list_requests
|
|
.lock()
|
|
.await
|
|
.push(request.into_inner());
|
|
Ok(Response::new(ListSandboxTemplatesResponse {
|
|
templates: Vec::new(),
|
|
next_page_token: String::new(),
|
|
}))
|
|
}
|
|
|
|
async fn delete_sandbox_template(
|
|
&self,
|
|
request: tonic::Request<DeleteSandboxTemplateRequest>,
|
|
) -> Result<Response<openshell_core::proto::DeleteSandboxTemplateResponse>, Status> {
|
|
self.state
|
|
.template_delete_requests
|
|
.lock()
|
|
.await
|
|
.push(request.into_inner());
|
|
Ok(Response::new(
|
|
openshell_core::proto::DeleteSandboxTemplateResponse { deleted: true },
|
|
))
|
|
}
|
|
|
|
async fn list_sandbox_providers(
|
|
&self,
|
|
_request: tonic::Request<ListSandboxProvidersRequest>,
|
|
) -> Result<Response<ListSandboxProvidersResponse>, Status> {
|
|
Ok(Response::new(ListSandboxProvidersResponse::default()))
|
|
}
|
|
|
|
async fn attach_sandbox_provider(
|
|
&self,
|
|
_request: tonic::Request<AttachSandboxProviderRequest>,
|
|
) -> Result<Response<AttachSandboxProviderResponse>, Status> {
|
|
Ok(Response::new(AttachSandboxProviderResponse::default()))
|
|
}
|
|
|
|
async fn detach_sandbox_provider(
|
|
&self,
|
|
_request: tonic::Request<DetachSandboxProviderRequest>,
|
|
) -> Result<Response<DetachSandboxProviderResponse>, Status> {
|
|
Ok(Response::new(DetachSandboxProviderResponse::default()))
|
|
}
|
|
|
|
async fn delete_sandbox(
|
|
&self,
|
|
request: tonic::Request<DeleteSandboxRequest>,
|
|
) -> Result<Response<DeleteSandboxResponse>, Status> {
|
|
let request = request.into_inner();
|
|
self.state
|
|
.delete_requests
|
|
.lock()
|
|
.await
|
|
.push(request.clone());
|
|
self.state
|
|
.deleted_names
|
|
.lock()
|
|
.await
|
|
.push(vec![request.name.clone()]);
|
|
let delete_failure = self.state.fail_delete_sandbox_message.lock().await.take();
|
|
if let Some(message) = delete_failure {
|
|
return Err(Status::internal(message));
|
|
}
|
|
Ok(Response::new(DeleteSandboxResponse { deleted: true }))
|
|
}
|
|
|
|
async fn get_sandbox_config(
|
|
&self,
|
|
_request: tonic::Request<GetSandboxConfigRequest>,
|
|
) -> Result<Response<GetSandboxConfigResponse>, Status> {
|
|
Ok(Response::new(GetSandboxConfigResponse::default()))
|
|
}
|
|
|
|
async fn get_gateway_config(
|
|
&self,
|
|
_request: tonic::Request<GetGatewayConfigRequest>,
|
|
) -> Result<Response<GetGatewayConfigResponse>, Status> {
|
|
self.state
|
|
.gateway_config_requests
|
|
.fetch_add(1, Ordering::SeqCst);
|
|
Ok(Response::new(GetGatewayConfigResponse {
|
|
settings: self.state.global_settings.lock().await.clone(),
|
|
settings_revision: 1,
|
|
}))
|
|
}
|
|
|
|
async fn get_sandbox_provider_environment(
|
|
&self,
|
|
_request: tonic::Request<GetSandboxProviderEnvironmentRequest>,
|
|
) -> Result<Response<GetSandboxProviderEnvironmentResponse>, Status> {
|
|
Ok(Response::new(
|
|
GetSandboxProviderEnvironmentResponse::default(),
|
|
))
|
|
}
|
|
|
|
async fn create_ssh_session(
|
|
&self,
|
|
request: tonic::Request<CreateSshSessionRequest>,
|
|
) -> Result<Response<CreateSshSessionResponse>, Status> {
|
|
self.state
|
|
.ssh_session_requests
|
|
.fetch_add(1, Ordering::SeqCst);
|
|
if self
|
|
.state
|
|
.ssh_session_failures_remaining
|
|
.fetch_update(Ordering::SeqCst, Ordering::SeqCst, |remaining| {
|
|
remaining.checked_sub(1)
|
|
})
|
|
.is_ok()
|
|
{
|
|
return Err(Status::failed_precondition("sandbox is not ready"));
|
|
}
|
|
let sandbox_id = request.into_inner().sandbox_id;
|
|
Ok(Response::new(CreateSshSessionResponse {
|
|
sandbox_id,
|
|
token: "test-token".to_string(),
|
|
gateway_scheme: "https".to_string(),
|
|
gateway_host: "localhost".to_string(),
|
|
gateway_port: 443,
|
|
..CreateSshSessionResponse::default()
|
|
}))
|
|
}
|
|
|
|
async fn expose_service(
|
|
&self,
|
|
_request: tonic::Request<openshell_core::proto::ExposeServiceRequest>,
|
|
) -> Result<Response<openshell_core::proto::ServiceEndpointResponse>, Status> {
|
|
Ok(Response::new(
|
|
openshell_core::proto::ServiceEndpointResponse::default(),
|
|
))
|
|
}
|
|
|
|
async fn get_service(
|
|
&self,
|
|
_: tonic::Request<openshell_core::proto::GetServiceRequest>,
|
|
) -> Result<Response<openshell_core::proto::ServiceEndpointResponse>, Status> {
|
|
Err(Status::unimplemented("unused"))
|
|
}
|
|
|
|
async fn list_services(
|
|
&self,
|
|
_: tonic::Request<openshell_core::proto::ListServicesRequest>,
|
|
) -> Result<Response<openshell_core::proto::ListServicesResponse>, Status> {
|
|
Err(Status::unimplemented("unused"))
|
|
}
|
|
|
|
async fn delete_service(
|
|
&self,
|
|
_: tonic::Request<openshell_core::proto::DeleteServiceRequest>,
|
|
) -> Result<Response<openshell_core::proto::DeleteServiceResponse>, Status> {
|
|
Err(Status::unimplemented("unused"))
|
|
}
|
|
|
|
async fn revoke_ssh_session(
|
|
&self,
|
|
_request: tonic::Request<RevokeSshSessionRequest>,
|
|
) -> Result<Response<RevokeSshSessionResponse>, Status> {
|
|
Ok(Response::new(RevokeSshSessionResponse::default()))
|
|
}
|
|
|
|
async fn exchange_provider_subject_token(
|
|
&self,
|
|
_request: tonic::Request<ExchangeProviderSubjectTokenRequest>,
|
|
) -> Result<Response<ExchangeProviderSubjectTokenResponse>, Status> {
|
|
Err(Status::unimplemented("unused"))
|
|
}
|
|
|
|
async fn create_provider(
|
|
&self,
|
|
_request: tonic::Request<CreateProviderRequest>,
|
|
) -> Result<Response<ProviderResponse>, Status> {
|
|
Ok(Response::new(ProviderResponse::default()))
|
|
}
|
|
|
|
async fn get_provider(
|
|
&self,
|
|
_request: tonic::Request<GetProviderRequest>,
|
|
) -> Result<Response<ProviderResponse>, Status> {
|
|
Err(Status::not_found("provider not found"))
|
|
}
|
|
|
|
async fn list_providers(
|
|
&self,
|
|
_request: tonic::Request<ListProvidersRequest>,
|
|
) -> Result<Response<ListProvidersResponse>, Status> {
|
|
Ok(Response::new(ListProvidersResponse {
|
|
providers: self.state.providers.lock().await.clone(),
|
|
next_page_token: String::new(),
|
|
}))
|
|
}
|
|
|
|
async fn list_provider_profiles(
|
|
&self,
|
|
_request: tonic::Request<openshell_core::proto::ListProviderProfilesRequest>,
|
|
) -> Result<Response<openshell_core::proto::ListProviderProfilesResponse>, Status> {
|
|
let profiles = openshell_providers::builtin_profiles()
|
|
.iter()
|
|
.map(openshell_providers::ProviderTypeProfile::to_proto)
|
|
.collect();
|
|
Ok(Response::new(
|
|
openshell_core::proto::ListProviderProfilesResponse {
|
|
profiles,
|
|
next_page_token: String::new(),
|
|
},
|
|
))
|
|
}
|
|
|
|
async fn get_provider_profile(
|
|
&self,
|
|
request: tonic::Request<openshell_core::proto::GetProviderProfileRequest>,
|
|
) -> Result<Response<openshell_core::proto::ProviderProfileResponse>, Status> {
|
|
let id = request.into_inner().id;
|
|
let profile = openshell_providers::builtin_profiles()
|
|
.iter()
|
|
.find(|profile| profile.id == id)
|
|
.ok_or_else(|| Status::not_found("provider profile not found"))?
|
|
.to_proto();
|
|
Ok(Response::new(
|
|
openshell_core::proto::ProviderProfileResponse {
|
|
profile: Some(profile),
|
|
},
|
|
))
|
|
}
|
|
|
|
async fn import_provider_profiles(
|
|
&self,
|
|
_request: tonic::Request<openshell_core::proto::ImportProviderProfilesRequest>,
|
|
) -> Result<Response<openshell_core::proto::ImportProviderProfilesResponse>, Status> {
|
|
Err(Status::unimplemented("not implemented in test"))
|
|
}
|
|
|
|
async fn update_provider_profiles(
|
|
&self,
|
|
_request: tonic::Request<openshell_core::proto::UpdateProviderProfilesRequest>,
|
|
) -> Result<Response<openshell_core::proto::UpdateProviderProfilesResponse>, Status> {
|
|
Err(Status::unimplemented("not implemented in test"))
|
|
}
|
|
|
|
async fn lint_provider_profiles(
|
|
&self,
|
|
_request: tonic::Request<openshell_core::proto::LintProviderProfilesRequest>,
|
|
) -> Result<Response<openshell_core::proto::LintProviderProfilesResponse>, Status> {
|
|
Err(Status::unimplemented("not implemented in test"))
|
|
}
|
|
|
|
async fn delete_provider_profile(
|
|
&self,
|
|
_request: tonic::Request<openshell_core::proto::DeleteProviderProfileRequest>,
|
|
) -> Result<Response<openshell_core::proto::DeleteProviderProfileResponse>, Status> {
|
|
Err(Status::unimplemented("not implemented in test"))
|
|
}
|
|
|
|
async fn update_provider(
|
|
&self,
|
|
_request: tonic::Request<UpdateProviderRequest>,
|
|
) -> Result<Response<ProviderResponse>, Status> {
|
|
Ok(Response::new(ProviderResponse::default()))
|
|
}
|
|
async fn get_provider_refresh_status(
|
|
&self,
|
|
_: tonic::Request<openshell_core::proto::GetProviderRefreshStatusRequest>,
|
|
) -> Result<Response<openshell_core::proto::GetProviderRefreshStatusResponse>, Status> {
|
|
Err(Status::unimplemented("unused"))
|
|
}
|
|
|
|
async fn configure_provider_refresh(
|
|
&self,
|
|
_: tonic::Request<openshell_core::proto::ConfigureProviderRefreshRequest>,
|
|
) -> Result<Response<openshell_core::proto::ConfigureProviderRefreshResponse>, Status> {
|
|
Err(Status::unimplemented("unused"))
|
|
}
|
|
|
|
async fn rotate_provider_credential(
|
|
&self,
|
|
_: tonic::Request<openshell_core::proto::RotateProviderCredentialRequest>,
|
|
) -> Result<Response<openshell_core::proto::RotateProviderCredentialResponse>, Status> {
|
|
Err(Status::unimplemented("unused"))
|
|
}
|
|
|
|
async fn delete_provider_refresh(
|
|
&self,
|
|
_: tonic::Request<openshell_core::proto::DeleteProviderRefreshRequest>,
|
|
) -> Result<Response<openshell_core::proto::DeleteProviderRefreshResponse>, Status> {
|
|
Err(Status::unimplemented("unused"))
|
|
}
|
|
|
|
async fn delete_provider(
|
|
&self,
|
|
_request: tonic::Request<DeleteProviderRequest>,
|
|
) -> Result<Response<DeleteProviderResponse>, Status> {
|
|
Ok(Response::new(DeleteProviderResponse { deleted: true }))
|
|
}
|
|
|
|
type WatchSandboxStream =
|
|
tokio_stream::wrappers::ReceiverStream<Result<SandboxStreamEvent, Status>>;
|
|
type ExecSandboxStream =
|
|
tokio_stream::wrappers::ReceiverStream<Result<ExecSandboxEvent, Status>>;
|
|
type ConnectSupervisorStream =
|
|
tokio_stream::wrappers::ReceiverStream<Result<GatewayMessage, Status>>;
|
|
|
|
async fn watch_sandbox(
|
|
&self,
|
|
request: tonic::Request<WatchSandboxRequest>,
|
|
) -> Result<Response<Self::WatchSandboxStream>, Status> {
|
|
let sandbox_id = request.into_inner().id;
|
|
let (tx, rx) = mpsc::channel(4);
|
|
let vm_error_after_started = self.state.vm_error_after_started.load(Ordering::SeqCst);
|
|
let vm_error_with_observed_exit = self
|
|
.state
|
|
.vm_error_with_observed_exit
|
|
.load(Ordering::SeqCst);
|
|
let vm_slow_progress_before_ready = self
|
|
.state
|
|
.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 terminal_before_relay = self.state.terminal_before_relay.load(Ordering::SeqCst);
|
|
let terminal_after_provisional_container_exit = self
|
|
.state
|
|
.terminal_after_provisional_container_exit
|
|
.load(Ordering::SeqCst);
|
|
let provisional_container_exit_without_result = self
|
|
.state
|
|
.provisional_container_exit_without_result
|
|
.load(Ordering::SeqCst);
|
|
let provisional_container_exit_sent =
|
|
Arc::clone(&self.state.provisional_container_exit_sent);
|
|
|
|
tokio::spawn(async move {
|
|
let mut provisioning = Sandbox {
|
|
metadata: Some(openshell_core::proto::datamodel::v1::ObjectMeta {
|
|
id: sandbox_id.clone(),
|
|
name: sandbox_id.trim_start_matches("id-").to_string(),
|
|
created_at_ms: 0,
|
|
labels: HashMap::new(),
|
|
resource_version: 0,
|
|
annotations: HashMap::new(),
|
|
workspace: String::new(),
|
|
deletion_timestamp_ms: 0,
|
|
}),
|
|
..Sandbox::default()
|
|
};
|
|
provisioning.set_phase(SandboxPhase::Provisioning as i32);
|
|
let mut error = Sandbox {
|
|
status: Some(SandboxStatus {
|
|
sandbox_name: sandbox_id.trim_start_matches("id-").to_string(),
|
|
conditions: vec![SandboxCondition {
|
|
r#type: "Ready".to_string(),
|
|
status: "False".to_string(),
|
|
reason: "ProcessExited".to_string(),
|
|
message: "VM process exited with status 0".to_string(),
|
|
last_transition_time: String::new(),
|
|
}],
|
|
..Default::default()
|
|
}),
|
|
..provisioning.clone()
|
|
};
|
|
error.set_phase(SandboxPhase::Error as i32);
|
|
if vm_error_with_observed_exit {
|
|
error.status.as_mut().unwrap().exit_code = Some(137);
|
|
}
|
|
let mut ready = provisioning.clone();
|
|
ready.set_phase(SandboxPhase::Ready as i32);
|
|
let mut completed = provisioning.clone();
|
|
completed.status = Some(SandboxStatus {
|
|
exit_code: Some(0),
|
|
..SandboxStatus::default()
|
|
});
|
|
completed.set_phase(SandboxPhase::Completed as i32);
|
|
|
|
let mut provisional_container_exit = error.clone();
|
|
if let Some(ready) = provisional_container_exit
|
|
.status
|
|
.as_mut()
|
|
.and_then(|status| status.conditions.first_mut())
|
|
{
|
|
ready.reason = "ContainerExited".to_string();
|
|
ready.message = "Sandbox container exited".to_string();
|
|
}
|
|
|
|
let _ = tx
|
|
.send(Ok(SandboxStreamEvent {
|
|
payload: Some(sandbox_stream_event::Payload::Sandbox(provisioning)),
|
|
}))
|
|
.await;
|
|
if terminal_after_provisional_container_exit
|
|
|| provisional_container_exit_without_result
|
|
{
|
|
let _ = tx
|
|
.send(Ok(SandboxStreamEvent {
|
|
payload: Some(sandbox_stream_event::Payload::Sandbox(
|
|
provisional_container_exit,
|
|
)),
|
|
}))
|
|
.await;
|
|
provisional_container_exit_sent.notify_waiters();
|
|
if provisional_container_exit_without_result {
|
|
std::future::pending::<()>().await;
|
|
return;
|
|
}
|
|
tokio::task::yield_now().await;
|
|
let _ = tx
|
|
.send(Ok(SandboxStreamEvent {
|
|
payload: Some(sandbox_stream_event::Payload::Sandbox(completed)),
|
|
}))
|
|
.await;
|
|
return;
|
|
}
|
|
if vm_error_after_started {
|
|
let _ = tx
|
|
.send(Ok(SandboxStreamEvent {
|
|
payload: Some(sandbox_stream_event::Payload::Event(PlatformEvent {
|
|
source: "vm".to_string(),
|
|
reason: "Started".to_string(),
|
|
message: "Started VM launcher".to_string(),
|
|
..PlatformEvent::default()
|
|
})),
|
|
}))
|
|
.await;
|
|
let _ = tx
|
|
.send(Ok(SandboxStreamEvent {
|
|
payload: Some(sandbox_stream_event::Payload::Sandbox(error)),
|
|
}))
|
|
.await;
|
|
tokio::time::sleep(Duration::from_secs(5)).await;
|
|
return;
|
|
}
|
|
if vm_log_churn_before_ready {
|
|
for message in ["still booting", "still booting again"] {
|
|
tokio::time::sleep(Duration::from_millis(600)).await;
|
|
let _ = tx
|
|
.send(Ok(SandboxStreamEvent {
|
|
payload: Some(sandbox_stream_event::Payload::Log(SandboxLogLine {
|
|
sandbox_id: sandbox_id.clone(),
|
|
timestamp_ms: 0,
|
|
level: "INFO".to_string(),
|
|
target: "test".to_string(),
|
|
message: message.to_string(),
|
|
source: "gateway".to_string(),
|
|
fields: HashMap::new(),
|
|
})),
|
|
}))
|
|
.await;
|
|
}
|
|
let _ = tx
|
|
.send(Ok(SandboxStreamEvent {
|
|
payload: Some(sandbox_stream_event::Payload::Sandbox(ready)),
|
|
}))
|
|
.await;
|
|
return;
|
|
}
|
|
if terminal_before_relay {
|
|
let _ = tx
|
|
.send(Ok(SandboxStreamEvent {
|
|
payload: Some(sandbox_stream_event::Payload::Sandbox(completed)),
|
|
}))
|
|
.await;
|
|
return;
|
|
}
|
|
if vm_slow_progress_before_ready {
|
|
tokio::time::sleep(Duration::from_millis(600)).await;
|
|
let _ = tx
|
|
.send(Ok(SandboxStreamEvent {
|
|
payload: Some(sandbox_stream_event::Payload::Event(PlatformEvent {
|
|
source: "vm".to_string(),
|
|
reason: "PreparingRootfs".to_string(),
|
|
message: "Preparing rootfs".to_string(),
|
|
..PlatformEvent::default()
|
|
})),
|
|
}))
|
|
.await;
|
|
tokio::time::sleep(Duration::from_millis(600)).await;
|
|
let _ = tx
|
|
.send(Ok(SandboxStreamEvent {
|
|
payload: Some(sandbox_stream_event::Payload::Event(PlatformEvent {
|
|
source: "vm".to_string(),
|
|
reason: "CreatingRootDisk".to_string(),
|
|
message: "Formatting root disk".to_string(),
|
|
..PlatformEvent::default()
|
|
})),
|
|
}))
|
|
.await;
|
|
tokio::time::sleep(Duration::from_millis(600)).await;
|
|
let _ = tx
|
|
.send(Ok(SandboxStreamEvent {
|
|
payload: Some(sandbox_stream_event::Payload::Sandbox(ready)),
|
|
}))
|
|
.await;
|
|
return;
|
|
}
|
|
let _ = tx
|
|
.send(Ok(SandboxStreamEvent {
|
|
payload: Some(sandbox_stream_event::Payload::Event(PlatformEvent {
|
|
reason: "Scheduled".to_string(),
|
|
message: "Sandbox scheduled".to_string(),
|
|
..PlatformEvent::default()
|
|
})),
|
|
}))
|
|
.await;
|
|
let _ = tx
|
|
.send(Ok(SandboxStreamEvent {
|
|
payload: Some(sandbox_stream_event::Payload::Sandbox(ready)),
|
|
}))
|
|
.await;
|
|
});
|
|
|
|
Ok(Response::new(tokio_stream::wrappers::ReceiverStream::new(
|
|
rx,
|
|
)))
|
|
}
|
|
|
|
async fn exec_sandbox(
|
|
&self,
|
|
_request: tonic::Request<ExecSandboxRequest>,
|
|
) -> Result<Response<Self::ExecSandboxStream>, Status> {
|
|
let (_tx, rx) = mpsc::channel(1);
|
|
Ok(Response::new(tokio_stream::wrappers::ReceiverStream::new(
|
|
rx,
|
|
)))
|
|
}
|
|
|
|
type ExecSandboxInteractiveStream =
|
|
tokio_stream::wrappers::ReceiverStream<Result<ExecSandboxEvent, Status>>;
|
|
async fn exec_sandbox_interactive(
|
|
&self,
|
|
_request: tonic::Request<tonic::Streaming<ExecSandboxInput>>,
|
|
) -> Result<Response<Self::ExecSandboxInteractiveStream>, Status> {
|
|
Err(Status::unimplemented("not implemented in test"))
|
|
}
|
|
|
|
async fn update_config(
|
|
&self,
|
|
_request: tonic::Request<openshell_core::proto::UpdateConfigRequest>,
|
|
) -> Result<Response<openshell_core::proto::UpdateConfigResponse>, Status> {
|
|
Err(Status::unimplemented("not implemented in test"))
|
|
}
|
|
|
|
async fn get_sandbox_policy_status(
|
|
&self,
|
|
_request: tonic::Request<openshell_core::proto::GetSandboxPolicyStatusRequest>,
|
|
) -> Result<Response<openshell_core::proto::GetSandboxPolicyStatusResponse>, Status> {
|
|
Err(Status::unimplemented("not implemented in test"))
|
|
}
|
|
|
|
async fn list_sandbox_policies(
|
|
&self,
|
|
_request: tonic::Request<openshell_core::proto::ListSandboxPoliciesRequest>,
|
|
) -> Result<Response<openshell_core::proto::ListSandboxPoliciesResponse>, Status> {
|
|
Err(Status::unimplemented("not implemented in test"))
|
|
}
|
|
|
|
async fn report_policy_status(
|
|
&self,
|
|
_request: tonic::Request<openshell_core::proto::ReportPolicyStatusRequest>,
|
|
) -> Result<Response<openshell_core::proto::ReportPolicyStatusResponse>, Status> {
|
|
Err(Status::unimplemented("not implemented in test"))
|
|
}
|
|
|
|
async fn report_endpoint_status(
|
|
&self,
|
|
_request: tonic::Request<ReportEndpointStatusRequest>,
|
|
) -> Result<Response<ReportEndpointStatusResponse>, Status> {
|
|
Err(Status::unimplemented("not implemented in test"))
|
|
}
|
|
|
|
async fn get_sandbox_logs(
|
|
&self,
|
|
_request: tonic::Request<openshell_core::proto::GetSandboxLogsRequest>,
|
|
) -> Result<Response<openshell_core::proto::GetSandboxLogsResponse>, Status> {
|
|
Err(Status::unimplemented("not implemented in test"))
|
|
}
|
|
|
|
async fn push_sandbox_logs(
|
|
&self,
|
|
_request: tonic::Request<tonic::Streaming<openshell_core::proto::PushSandboxLogsRequest>>,
|
|
) -> Result<Response<openshell_core::proto::PushSandboxLogsResponse>, Status> {
|
|
Err(Status::unimplemented("not implemented in test"))
|
|
}
|
|
|
|
async fn submit_policy_analysis(
|
|
&self,
|
|
_request: tonic::Request<openshell_core::proto::SubmitPolicyAnalysisRequest>,
|
|
) -> Result<Response<openshell_core::proto::SubmitPolicyAnalysisResponse>, Status> {
|
|
Err(Status::unimplemented("not implemented in test"))
|
|
}
|
|
|
|
async fn get_draft_policy(
|
|
&self,
|
|
_request: tonic::Request<openshell_core::proto::GetDraftPolicyRequest>,
|
|
) -> Result<Response<openshell_core::proto::GetDraftPolicyResponse>, Status> {
|
|
Err(Status::unimplemented("not implemented in test"))
|
|
}
|
|
|
|
async fn approve_draft_chunk(
|
|
&self,
|
|
_request: tonic::Request<openshell_core::proto::ApproveDraftChunkRequest>,
|
|
) -> Result<Response<openshell_core::proto::ApproveDraftChunkResponse>, Status> {
|
|
Err(Status::unimplemented("not implemented in test"))
|
|
}
|
|
|
|
async fn reject_draft_chunk(
|
|
&self,
|
|
_request: tonic::Request<openshell_core::proto::RejectDraftChunkRequest>,
|
|
) -> Result<Response<openshell_core::proto::RejectDraftChunkResponse>, Status> {
|
|
Err(Status::unimplemented("not implemented in test"))
|
|
}
|
|
|
|
async fn approve_all_draft_chunks(
|
|
&self,
|
|
_request: tonic::Request<openshell_core::proto::ApproveAllDraftChunksRequest>,
|
|
) -> Result<Response<openshell_core::proto::ApproveAllDraftChunksResponse>, Status> {
|
|
Err(Status::unimplemented("not implemented in test"))
|
|
}
|
|
|
|
async fn edit_draft_chunk(
|
|
&self,
|
|
_request: tonic::Request<openshell_core::proto::EditDraftChunkRequest>,
|
|
) -> Result<Response<openshell_core::proto::EditDraftChunkResponse>, Status> {
|
|
Err(Status::unimplemented("not implemented in test"))
|
|
}
|
|
|
|
async fn undo_draft_chunk(
|
|
&self,
|
|
_request: tonic::Request<openshell_core::proto::UndoDraftChunkRequest>,
|
|
) -> Result<Response<openshell_core::proto::UndoDraftChunkResponse>, Status> {
|
|
Err(Status::unimplemented("not implemented in test"))
|
|
}
|
|
|
|
async fn clear_draft_chunks(
|
|
&self,
|
|
_request: tonic::Request<openshell_core::proto::ClearDraftChunksRequest>,
|
|
) -> Result<Response<openshell_core::proto::ClearDraftChunksResponse>, Status> {
|
|
Err(Status::unimplemented("not implemented in test"))
|
|
}
|
|
|
|
async fn get_draft_history(
|
|
&self,
|
|
_request: tonic::Request<openshell_core::proto::GetDraftHistoryRequest>,
|
|
) -> Result<Response<openshell_core::proto::GetDraftHistoryResponse>, Status> {
|
|
Err(Status::unimplemented("not implemented in test"))
|
|
}
|
|
|
|
async fn issue_sandbox_token(
|
|
&self,
|
|
_request: tonic::Request<openshell_core::proto::IssueSandboxTokenRequest>,
|
|
) -> Result<Response<openshell_core::proto::IssueSandboxTokenResponse>, Status> {
|
|
Err(Status::unimplemented("not implemented in test"))
|
|
}
|
|
|
|
async fn refresh_sandbox_token(
|
|
&self,
|
|
_request: tonic::Request<openshell_core::proto::RefreshSandboxTokenRequest>,
|
|
) -> Result<Response<openshell_core::proto::RefreshSandboxTokenResponse>, Status> {
|
|
Err(Status::unimplemented("not implemented in test"))
|
|
}
|
|
|
|
async fn connect_supervisor(
|
|
&self,
|
|
_request: tonic::Request<tonic::Streaming<SupervisorMessage>>,
|
|
) -> Result<Response<Self::ConnectSupervisorStream>, Status> {
|
|
Err(Status::unimplemented("not implemented in test"))
|
|
}
|
|
|
|
type RelayStreamStream =
|
|
tokio_stream::wrappers::ReceiverStream<Result<openshell_core::proto::RelayFrame, Status>>;
|
|
|
|
async fn relay_stream(
|
|
&self,
|
|
_request: tonic::Request<tonic::Streaming<openshell_core::proto::RelayFrame>>,
|
|
) -> Result<Response<Self::RelayStreamStream>, Status> {
|
|
Err(Status::unimplemented("not implemented in test"))
|
|
}
|
|
|
|
type ForwardTcpStream = tokio_stream::wrappers::ReceiverStream<
|
|
Result<openshell_core::proto::TcpForwardFrame, Status>,
|
|
>;
|
|
|
|
async fn forward_tcp(
|
|
&self,
|
|
_request: tonic::Request<tonic::Streaming<openshell_core::proto::TcpForwardFrame>>,
|
|
) -> Result<Response<Self::ForwardTcpStream>, Status> {
|
|
Err(Status::unimplemented("not implemented in test"))
|
|
}
|
|
|
|
async fn create_workspace(
|
|
&self,
|
|
_request: tonic::Request<openshell_core::proto::CreateWorkspaceRequest>,
|
|
) -> Result<Response<openshell_core::proto::CreateWorkspaceResponse>, Status> {
|
|
Err(Status::unimplemented("not implemented in test"))
|
|
}
|
|
|
|
async fn get_workspace(
|
|
&self,
|
|
_request: tonic::Request<openshell_core::proto::GetWorkspaceRequest>,
|
|
) -> Result<Response<openshell_core::proto::GetWorkspaceResponse>, Status> {
|
|
Err(Status::unimplemented("not implemented in test"))
|
|
}
|
|
|
|
async fn list_workspaces(
|
|
&self,
|
|
_request: tonic::Request<openshell_core::proto::ListWorkspacesRequest>,
|
|
) -> Result<Response<openshell_core::proto::ListWorkspacesResponse>, Status> {
|
|
Err(Status::unimplemented("not implemented in test"))
|
|
}
|
|
|
|
async fn delete_workspace(
|
|
&self,
|
|
_request: tonic::Request<openshell_core::proto::DeleteWorkspaceRequest>,
|
|
) -> Result<Response<openshell_core::proto::DeleteWorkspaceResponse>, Status> {
|
|
Err(Status::unimplemented("not implemented in test"))
|
|
}
|
|
|
|
async fn add_workspace_member(
|
|
&self,
|
|
_request: tonic::Request<openshell_core::proto::AddWorkspaceMemberRequest>,
|
|
) -> Result<Response<openshell_core::proto::AddWorkspaceMemberResponse>, Status> {
|
|
Err(Status::unimplemented("not implemented in test"))
|
|
}
|
|
|
|
async fn remove_workspace_member(
|
|
&self,
|
|
_request: tonic::Request<openshell_core::proto::RemoveWorkspaceMemberRequest>,
|
|
) -> Result<Response<openshell_core::proto::RemoveWorkspaceMemberResponse>, Status> {
|
|
Err(Status::unimplemented("not implemented in test"))
|
|
}
|
|
|
|
async fn list_workspace_members(
|
|
&self,
|
|
_request: tonic::Request<openshell_core::proto::ListWorkspaceMembersRequest>,
|
|
) -> Result<Response<openshell_core::proto::ListWorkspaceMembersResponse>, Status> {
|
|
Err(Status::unimplemented("not implemented in test"))
|
|
}
|
|
}
|
|
|
|
struct TestServer {
|
|
endpoint: String,
|
|
tls: TlsOptions,
|
|
openshell: TestOpenShell,
|
|
dir: TempDir,
|
|
}
|
|
|
|
async fn run_server() -> TestServer {
|
|
let (ca, ca_key) = build_ca();
|
|
let (server_cert, server_key) = build_server_cert(&ca, &ca_key);
|
|
let (client_cert, client_key) = build_client_cert(&ca, &ca_key);
|
|
let ca_cert = ca.pem();
|
|
|
|
let identity = Identity::from_pem(server_cert, server_key);
|
|
let client_ca = TlsCertificate::from_pem(ca_cert.clone());
|
|
let tls_config = ServerTlsConfig::new()
|
|
.identity(identity)
|
|
.client_ca_root(client_ca);
|
|
|
|
let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
|
|
let addr = listener.local_addr().unwrap();
|
|
let incoming = TcpListenerStream::new(listener);
|
|
|
|
let openshell = TestOpenShell::default();
|
|
let svc_openshell = openshell.clone();
|
|
|
|
tokio::spawn(async move {
|
|
Server::builder()
|
|
.tls_config(tls_config)
|
|
.unwrap()
|
|
.add_service(OpenShellServer::new(svc_openshell))
|
|
.serve_with_incoming(incoming)
|
|
.await
|
|
.unwrap();
|
|
});
|
|
|
|
let dir = tempfile::tempdir().unwrap();
|
|
let ca_path = dir.path().join("ca.crt");
|
|
let cert_path = dir.path().join("tls.crt");
|
|
let key_path = dir.path().join("tls.key");
|
|
fs::write(&ca_path, ca_cert).unwrap();
|
|
fs::write(&cert_path, client_cert).unwrap();
|
|
fs::write(&key_path, client_key).unwrap();
|
|
|
|
let tls = TlsOptions::new(Some(ca_path), Some(cert_path), Some(key_path));
|
|
let endpoint = format!("https://localhost:{}", addr.port());
|
|
|
|
TestServer {
|
|
endpoint,
|
|
tls,
|
|
openshell,
|
|
dir,
|
|
}
|
|
}
|
|
|
|
fn install_executable_script(
|
|
dir: &TempDir,
|
|
name: &str,
|
|
contents: impl AsRef<[u8]>,
|
|
) -> std::path::PathBuf {
|
|
let path = dir.path().join(name);
|
|
fs::write(&path, contents).unwrap();
|
|
let mut perms = fs::metadata(&path).unwrap().permissions();
|
|
perms.set_mode(0o755);
|
|
fs::set_permissions(&path, perms).unwrap();
|
|
path
|
|
}
|
|
|
|
fn install_fake_ssh(dir: &TempDir) -> std::path::PathBuf {
|
|
install_executable_script(dir, "ssh", "#!/bin/sh\nexit 0\n")
|
|
}
|
|
|
|
fn install_fake_pgrep_no_match(dir: &TempDir) -> std::path::PathBuf {
|
|
install_executable_script(dir, "pgrep", "#!/bin/sh\nexit 1\n")
|
|
}
|
|
|
|
fn install_fake_forward_process_helper(dir: &TempDir) -> std::path::PathBuf {
|
|
// Linux validation reads exact `/proc` argv, so the fake child must look
|
|
// like `ssh`, not Python or shell with appended tokens.
|
|
if let Some(path) = std::env::var_os("OPENSHELL_TEST_FAKE_FORWARD_PATH") {
|
|
return path.into();
|
|
}
|
|
|
|
let source_path = dir.path().join("fake-forward-process.rs");
|
|
let binary_path = dir.path().join("fake-forward-process");
|
|
fs::write(&source_path, include_bytes!("fixtures/fake_forward.rs")).unwrap();
|
|
let status = std::process::Command::new("rustc")
|
|
.arg("--edition=2021")
|
|
.arg(&source_path)
|
|
.arg("-o")
|
|
.arg(&binary_path)
|
|
.status()
|
|
.unwrap();
|
|
assert!(status.success(), "failed to compile fake forward process");
|
|
binary_path
|
|
}
|
|
|
|
fn install_fake_ps_for_pid_revalidation(
|
|
dir: &TempDir,
|
|
pid_path: &std::path::Path,
|
|
command_path: &std::path::Path,
|
|
) {
|
|
install_executable_script(
|
|
dir,
|
|
"ps",
|
|
format!(
|
|
r#"#!/bin/sh
|
|
set -eu
|
|
|
|
command_mode=0
|
|
requested_pid=""
|
|
previous=""
|
|
|
|
for arg in "$@"; do
|
|
if [ "$previous" = "-o" ]; then
|
|
if [ "$arg" = "command=" ]; then
|
|
command_mode=1
|
|
fi
|
|
previous=""
|
|
continue
|
|
fi
|
|
|
|
if [ "$previous" = "-p" ]; then
|
|
requested_pid="$arg"
|
|
previous=""
|
|
continue
|
|
fi
|
|
|
|
case "$arg" in
|
|
-o|-p)
|
|
previous="$arg"
|
|
;;
|
|
esac
|
|
done
|
|
|
|
expected_pid=""
|
|
if [ -s '{pid_path}' ]; then
|
|
expected_pid="$(cat '{pid_path}')"
|
|
fi
|
|
|
|
if [ "$command_mode" = "1" ] && [ -n "$expected_pid" ] && [ "$requested_pid" = "$expected_pid" ] && [ -s '{command_path}' ]; then
|
|
cat '{command_path}'
|
|
printf '\n'
|
|
exit 0
|
|
fi
|
|
|
|
exec /bin/ps "$@"
|
|
"#,
|
|
pid_path = pid_path.display(),
|
|
command_path = command_path.display(),
|
|
),
|
|
);
|
|
}
|
|
|
|
async fn wait_for_process_exit(pid: u32, timeout: Duration) -> bool {
|
|
let deadline = Instant::now() + timeout;
|
|
loop {
|
|
let output = std::process::Command::new("ps")
|
|
.arg("-o")
|
|
.arg("stat=")
|
|
.arg("-p")
|
|
.arg(pid.to_string())
|
|
.stderr(std::process::Stdio::null())
|
|
.output();
|
|
let alive = output.is_ok_and(|output| {
|
|
if !output.status.success() {
|
|
return false;
|
|
}
|
|
let stat = String::from_utf8_lossy(&output.stdout);
|
|
// Linux can leave the orphaned fake forward as a short-lived zombie
|
|
// until the container's init process reaps it. A zombie has already
|
|
// exited, so it satisfies this cleanup assertion.
|
|
!stat.trim_start().starts_with('Z')
|
|
});
|
|
if !alive {
|
|
return true;
|
|
}
|
|
if Instant::now() >= deadline {
|
|
return false;
|
|
}
|
|
tokio::time::sleep(Duration::from_millis(50)).await;
|
|
}
|
|
}
|
|
|
|
fn install_fake_forwarding_ssh(dir: &TempDir) -> std::path::PathBuf {
|
|
let pid_path = dir.path().join("fake-forward.pid");
|
|
let command_path = dir.path().join("fake-forward.command");
|
|
let helper_path = install_fake_forward_process_helper(dir);
|
|
let ssh_path = install_executable_script(
|
|
dir,
|
|
"ssh",
|
|
r#"#!/bin/sh
|
|
set -eu
|
|
|
|
forward=""
|
|
sandbox_id=""
|
|
saw_no_command=0
|
|
last_arg=""
|
|
previous=""
|
|
|
|
for arg in "$@"; do
|
|
if [ "$previous" = "-L" ]; then
|
|
forward="$arg"
|
|
previous=""
|
|
last_arg="$arg"
|
|
continue
|
|
fi
|
|
|
|
if [ "$previous" = "-o" ]; then
|
|
case "$arg" in
|
|
ProxyCommand=*)
|
|
sandbox_id="$(printf '%s\n' "$arg" | sed -n 's/.*--sandbox-id \([^ ]*\).*/\1/p')"
|
|
;;
|
|
esac
|
|
previous=""
|
|
last_arg="$arg"
|
|
continue
|
|
fi
|
|
|
|
case "$arg" in
|
|
-N)
|
|
saw_no_command=1
|
|
;;
|
|
-L|-o)
|
|
previous="$arg"
|
|
;;
|
|
esac
|
|
last_arg="$arg"
|
|
done
|
|
|
|
if [ -z "$forward" ]; then
|
|
exit 0
|
|
fi
|
|
|
|
if [ "$saw_no_command" != "1" ] || [ "$last_arg" != "sandbox" ]; then
|
|
exit 1
|
|
fi
|
|
|
|
first="${forward%%:*}"
|
|
rest="${forward#*:}"
|
|
case "$first" in
|
|
''|*[!0-9]*)
|
|
port="${rest%%:*}"
|
|
;;
|
|
*)
|
|
port="$first"
|
|
;;
|
|
esac
|
|
|
|
if [ -z "$port" ] || [ -z "$sandbox_id" ]; then
|
|
exit 1
|
|
fi
|
|
|
|
helper='@HELPER_PATH@'
|
|
echo "$$" > '@PID_PATH@'
|
|
printf '%s\n' "ssh -N -o ProxyCommand=/tmp/openshell ssh-proxy --gateway https://127.0.0.1:9443 --sandbox-id $sandbox_id --token test-token --gateway-name test-gateway -o ExitOnForwardFailure=yes -L $forward sandbox" > '@COMMAND_PATH@'
|
|
exec env OPENSHELL_FAKE_FORWARD_MODE=listen "$helper" -N -o "ProxyCommand=/tmp/openshell ssh-proxy --gateway https://127.0.0.1:9443 --sandbox-id $sandbox_id --token test-token --gateway-name test-gateway" -o ExitOnForwardFailure=yes -L "$forward" sandbox
|
|
"#
|
|
.replace("@PID_PATH@", &pid_path.display().to_string())
|
|
.replace("@COMMAND_PATH@", &command_path.display().to_string())
|
|
.replace("@HELPER_PATH@", &helper_path.display().to_string()),
|
|
);
|
|
|
|
install_fake_ps_for_pid_revalidation(dir, &pid_path, &command_path);
|
|
|
|
ssh_path
|
|
}
|
|
|
|
struct FakeUnreachableForward {
|
|
log_path: std::path::PathBuf,
|
|
pid_path: std::path::PathBuf,
|
|
}
|
|
|
|
fn install_fake_unreachable_forwarding_ssh(dir: &TempDir) -> FakeUnreachableForward {
|
|
let log_path = dir.path().join("fake-forward.log");
|
|
let pid_path = dir.path().join("fake-forward.pid");
|
|
let helper_path = install_fake_forward_process_helper(dir);
|
|
install_executable_script(
|
|
dir,
|
|
"ssh",
|
|
r#"#!/bin/sh
|
|
set -eu
|
|
|
|
forward=""
|
|
sandbox_id=""
|
|
saw_no_command=0
|
|
last_arg=""
|
|
previous=""
|
|
|
|
for arg in "$@"; do
|
|
if [ "$previous" = "-L" ]; then
|
|
forward="$arg"
|
|
previous=""
|
|
last_arg="$arg"
|
|
continue
|
|
fi
|
|
|
|
if [ "$previous" = "-o" ]; then
|
|
case "$arg" in
|
|
ProxyCommand=*)
|
|
sandbox_id="$(printf '%s\n' "$arg" | sed -n 's/.*--sandbox-id \([^ ]*\).*/\1/p')"
|
|
;;
|
|
esac
|
|
previous=""
|
|
last_arg="$arg"
|
|
continue
|
|
fi
|
|
|
|
case "$arg" in
|
|
-N)
|
|
saw_no_command=1
|
|
;;
|
|
-L|-o)
|
|
previous="$arg"
|
|
;;
|
|
esac
|
|
last_arg="$arg"
|
|
done
|
|
|
|
if [ -z "$forward" ] || [ -z "$sandbox_id" ]; then
|
|
exit 1
|
|
fi
|
|
|
|
if [ "$saw_no_command" != "1" ] || [ "$last_arg" != "sandbox" ]; then
|
|
exit 1
|
|
fi
|
|
|
|
helper='@HELPER_PATH@'
|
|
echo "$$" > '@PID_PATH@'
|
|
exec env OPENSHELL_FAKE_FORWARD_MODE=sleep "$helper" -N -o "ProxyCommand=/tmp/openshell ssh-proxy --gateway https://127.0.0.1:9443 --sandbox-id $sandbox_id --token test-token --gateway-name test-gateway" -o ExitOnForwardFailure=yes -L "$forward" sandbox >'@LOG_PATH@' 2>&1
|
|
"#
|
|
.replace("@LOG_PATH@", &log_path.display().to_string())
|
|
.replace("@PID_PATH@", &pid_path.display().to_string())
|
|
.replace("@HELPER_PATH@", &helper_path.display().to_string()),
|
|
);
|
|
|
|
FakeUnreachableForward { log_path, pid_path }
|
|
}
|
|
|
|
fn test_env(fake_ssh_dir: &TempDir, xdg_dir: &TempDir) -> EnvVarGuard {
|
|
test_env_with(fake_ssh_dir, xdg_dir, &[])
|
|
}
|
|
|
|
fn test_env_with(
|
|
fake_ssh_dir: &TempDir,
|
|
xdg_dir: &TempDir,
|
|
extra: &[(&'static str, String)],
|
|
) -> EnvVarGuard {
|
|
let path = format!(
|
|
"{}:{}",
|
|
fake_ssh_dir.path().display(),
|
|
std::env::var("PATH").unwrap_or_default()
|
|
);
|
|
let xdg = xdg_dir.path().to_str().unwrap().to_string();
|
|
|
|
let mut owned_pairs = vec![
|
|
("PATH", path),
|
|
("XDG_CONFIG_HOME", xdg.clone()),
|
|
("HOME", xdg),
|
|
];
|
|
owned_pairs.extend(extra.iter().cloned());
|
|
let pairs = owned_pairs
|
|
.iter()
|
|
.map(|(key, value)| (*key, value.as_str()))
|
|
.collect::<Vec<_>>();
|
|
|
|
EnvVarGuard::set(&pairs)
|
|
}
|
|
|
|
async fn deleted_names(server: &TestServer) -> Vec<Vec<String>> {
|
|
server.openshell.state.deleted_names.lock().await.clone()
|
|
}
|
|
|
|
async fn delete_requests(server: &TestServer) -> Vec<DeleteSandboxRequest> {
|
|
server.openshell.state.delete_requests.lock().await.clone()
|
|
}
|
|
|
|
async fn create_requests(server: &TestServer) -> Vec<CreateSandboxRequest> {
|
|
server.openshell.state.create_requests.lock().await.clone()
|
|
}
|
|
|
|
async fn template_create_requests(server: &TestServer) -> Vec<CreateSandboxTemplateRequest> {
|
|
server
|
|
.openshell
|
|
.state
|
|
.template_create_requests
|
|
.lock()
|
|
.await
|
|
.clone()
|
|
}
|
|
|
|
async fn template_list_requests(server: &TestServer) -> Vec<ListSandboxTemplatesRequest> {
|
|
server
|
|
.openshell
|
|
.state
|
|
.template_list_requests
|
|
.lock()
|
|
.await
|
|
.clone()
|
|
}
|
|
|
|
async fn template_delete_requests(server: &TestServer) -> Vec<DeleteSandboxTemplateRequest> {
|
|
server
|
|
.openshell
|
|
.state
|
|
.template_delete_requests
|
|
.lock()
|
|
.await
|
|
.clone()
|
|
}
|
|
|
|
async fn add_provider(server: &TestServer, name: &str, provider_type: &str) {
|
|
server
|
|
.openshell
|
|
.state
|
|
.providers
|
|
.lock()
|
|
.await
|
|
.push(Provider {
|
|
metadata: Some(openshell_core::proto::datamodel::v1::ObjectMeta {
|
|
id: format!("provider-{name}"),
|
|
name: name.to_string(),
|
|
created_at_ms: 0,
|
|
labels: HashMap::new(),
|
|
resource_version: 0,
|
|
annotations: HashMap::new(),
|
|
workspace: "default".to_string(),
|
|
deletion_timestamp_ms: 0,
|
|
}),
|
|
r#type: provider_type.to_string(),
|
|
credentials: HashMap::new(),
|
|
config: HashMap::new(),
|
|
credential_expires_at_ms: HashMap::new(),
|
|
profile_workspace: "default".to_string(),
|
|
credential_handles: HashMap::new(),
|
|
});
|
|
}
|
|
|
|
fn test_tls(server: &TestServer) -> TlsOptions {
|
|
server.tls.with_gateway_name("openshell")
|
|
}
|
|
|
|
fn gpu_requirements(count: Option<u32>) -> GpuResourceRequirements {
|
|
GpuResourceRequirements { count }
|
|
}
|
|
|
|
/// Shared defaults for integration tests. Note: `keep` is `true` here (most
|
|
/// tests expect persistent sandboxes) while `SandboxCreateConfig::default()`
|
|
/// sets `keep: false` (the safe production default). Tests that exercise
|
|
/// ephemeral behavior must explicitly override with `keep: false`.
|
|
fn test_config() -> run::SandboxCreateConfig<'static> {
|
|
run::SandboxCreateConfig {
|
|
keep: true,
|
|
tty_override: Some(false),
|
|
auto_providers_override: Some(false),
|
|
..Default::default()
|
|
}
|
|
}
|
|
|
|
#[tokio::test]
|
|
async fn sandbox_delete_continues_after_entry_failure() {
|
|
let server = run_server().await;
|
|
let tls = test_tls(&server);
|
|
*server
|
|
.openshell
|
|
.state
|
|
.fail_delete_sandbox_message
|
|
.lock()
|
|
.await = Some("simulated sandbox delete failure".to_string());
|
|
|
|
let err = run::sandbox_delete(
|
|
&server.endpoint,
|
|
&["failing-sandbox".to_string(), "later-sandbox".to_string()],
|
|
false,
|
|
None,
|
|
None,
|
|
"default",
|
|
&tls,
|
|
"openshell",
|
|
)
|
|
.await
|
|
.expect_err("sandbox delete should report aggregate failure");
|
|
|
|
let msg = err.to_string();
|
|
assert!(
|
|
msg.contains("failed to delete 1 sandbox: failing-sandbox"),
|
|
"unexpected error: {msg}"
|
|
);
|
|
assert_eq!(
|
|
deleted_names(&server).await,
|
|
vec![
|
|
vec!["failing-sandbox".to_string()],
|
|
vec!["later-sandbox".to_string()]
|
|
]
|
|
);
|
|
}
|
|
|
|
#[tokio::test]
|
|
async fn sandbox_delete_forwards_identity_preconditions() {
|
|
let server = run_server().await;
|
|
let tls = test_tls(&server);
|
|
|
|
run::sandbox_delete(
|
|
&server.endpoint,
|
|
&["guarded-sandbox".to_string()],
|
|
false,
|
|
Some("sb-123"),
|
|
Some(17),
|
|
"default",
|
|
&tls,
|
|
"openshell",
|
|
)
|
|
.await
|
|
.expect("identity-guarded delete should succeed");
|
|
|
|
let requests = delete_requests(&server).await;
|
|
assert_eq!(requests.len(), 1);
|
|
assert_eq!(requests[0].name, "guarded-sandbox");
|
|
assert_eq!(requests[0].expected_sandbox_id, "sb-123");
|
|
assert_eq!(requests[0].expected_resource_version, 17);
|
|
}
|
|
|
|
#[tokio::test]
|
|
async fn sandbox_delete_rejects_preconditions_for_multiple_names_before_rpc() {
|
|
let server = run_server().await;
|
|
let tls = test_tls(&server);
|
|
|
|
let error = run::sandbox_delete(
|
|
&server.endpoint,
|
|
&["sandbox-a".to_string(), "sandbox-b".to_string()],
|
|
false,
|
|
Some("sb-123"),
|
|
None,
|
|
"default",
|
|
&tls,
|
|
"openshell",
|
|
)
|
|
.await
|
|
.expect_err("identity preconditions must be single-sandbox only");
|
|
|
|
assert!(
|
|
error
|
|
.to_string()
|
|
.contains("require exactly one sandbox name")
|
|
);
|
|
assert!(delete_requests(&server).await.is_empty());
|
|
}
|
|
|
|
#[tokio::test]
|
|
async fn sandbox_delete_rejects_resource_version_without_immutable_identity() {
|
|
let server = run_server().await;
|
|
let tls = test_tls(&server);
|
|
|
|
let error = run::sandbox_delete(
|
|
&server.endpoint,
|
|
&["guarded-sandbox".to_string()],
|
|
false,
|
|
None,
|
|
Some(17),
|
|
"default",
|
|
&tls,
|
|
"openshell",
|
|
)
|
|
.await
|
|
.expect_err("resource version without immutable ID must fail closed");
|
|
|
|
assert!(error.to_string().contains("requires --expected-id"));
|
|
assert!(delete_requests(&server).await.is_empty());
|
|
}
|
|
|
|
#[tokio::test]
|
|
async fn sandbox_create_keeps_command_sessions_by_default() {
|
|
let server = run_server().await;
|
|
let fake_ssh_dir = tempfile::tempdir().unwrap();
|
|
let xdg_dir = tempfile::tempdir().unwrap();
|
|
let _env = test_env(&fake_ssh_dir, &xdg_dir);
|
|
let tls = test_tls(&server);
|
|
install_fake_ssh(&fake_ssh_dir);
|
|
|
|
run::sandbox_create(
|
|
&server.endpoint,
|
|
"openshell",
|
|
run::SandboxCreateConfig {
|
|
name: Some("default-command"),
|
|
command: &["echo".into(), "OK".into()],
|
|
..test_config()
|
|
},
|
|
"default",
|
|
&tls,
|
|
)
|
|
.await
|
|
.expect("sandbox create should succeed");
|
|
|
|
assert!(deleted_names(&server).await.is_empty());
|
|
assert_eq!(
|
|
load_last_sandbox("openshell", "default").as_deref(),
|
|
Some("default-command"),
|
|
"default sandboxes should be persisted as last-used"
|
|
);
|
|
}
|
|
|
|
#[tokio::test]
|
|
async fn sandbox_create_without_inferred_provider_skips_gateway_config() {
|
|
let server = run_server().await;
|
|
let fake_ssh_dir = tempfile::tempdir().unwrap();
|
|
let xdg_dir = tempfile::tempdir().unwrap();
|
|
let _env = test_env(&fake_ssh_dir, &xdg_dir);
|
|
let tls = test_tls(&server);
|
|
install_fake_ssh(&fake_ssh_dir);
|
|
|
|
run::sandbox_create(
|
|
&server.endpoint,
|
|
"openshell",
|
|
run::SandboxCreateConfig {
|
|
name: Some("no-provider-config"),
|
|
command: &["echo".into(), "OK".into()],
|
|
..test_config()
|
|
},
|
|
"default",
|
|
&tls,
|
|
)
|
|
.await
|
|
.expect("sandbox create should succeed without reading gateway config");
|
|
|
|
assert_eq!(
|
|
server
|
|
.openshell
|
|
.state
|
|
.gateway_config_requests
|
|
.load(Ordering::SeqCst),
|
|
0,
|
|
"commands without an inferred provider must not require global gateway settings"
|
|
);
|
|
}
|
|
|
|
#[tokio::test]
|
|
async fn sandbox_create_sends_cpu_and_memory_limits_only() {
|
|
let server = run_server().await;
|
|
let fake_ssh_dir = tempfile::tempdir().unwrap();
|
|
let xdg_dir = tempfile::tempdir().unwrap();
|
|
let _env = test_env(&fake_ssh_dir, &xdg_dir);
|
|
let tls = test_tls(&server);
|
|
install_fake_ssh(&fake_ssh_dir);
|
|
|
|
run::sandbox_create(
|
|
&server.endpoint,
|
|
"openshell",
|
|
run::SandboxCreateConfig {
|
|
name: Some("resources"),
|
|
cpu: Some("500m"),
|
|
memory: Some("2Gi"),
|
|
command: &["echo".into(), "OK".into()],
|
|
..test_config()
|
|
},
|
|
"default",
|
|
&tls,
|
|
)
|
|
.await
|
|
.expect("sandbox create should succeed");
|
|
|
|
let requests = create_requests(&server).await;
|
|
let resources = requests[0]
|
|
.spec
|
|
.as_ref()
|
|
.and_then(|spec| spec.template.as_ref())
|
|
.and_then(|template| template.resources.as_ref())
|
|
.expect("resource limits should be sent");
|
|
let limits = resources
|
|
.fields
|
|
.get("limits")
|
|
.and_then(|value| value.kind.as_ref())
|
|
.and_then(|kind| match kind {
|
|
prost_types::value::Kind::StructValue(inner) => Some(inner),
|
|
_ => None,
|
|
})
|
|
.expect("limits should be a struct");
|
|
|
|
assert_eq!(
|
|
limits
|
|
.fields
|
|
.get("cpu")
|
|
.and_then(|value| value.kind.as_ref())
|
|
.and_then(|kind| match kind {
|
|
prost_types::value::Kind::StringValue(value) => Some(value.as_str()),
|
|
_ => None,
|
|
}),
|
|
Some("500m")
|
|
);
|
|
assert_eq!(
|
|
limits
|
|
.fields
|
|
.get("memory")
|
|
.and_then(|value| value.kind.as_ref())
|
|
.and_then(|kind| match kind {
|
|
prost_types::value::Kind::StringValue(value) => Some(value.as_str()),
|
|
_ => None,
|
|
}),
|
|
Some("2Gi")
|
|
);
|
|
assert!(!resources.fields.contains_key("requests"));
|
|
}
|
|
|
|
#[tokio::test]
|
|
async fn sandbox_create_persists_exact_trailing_argv_as_main_process() {
|
|
let server = run_server().await;
|
|
let fake_ssh_dir = tempfile::tempdir().unwrap();
|
|
let xdg_dir = tempfile::tempdir().unwrap();
|
|
let _env = test_env(&fake_ssh_dir, &xdg_dir);
|
|
let tls = test_tls(&server);
|
|
install_fake_ssh(&fake_ssh_dir);
|
|
let command = vec![
|
|
"/opt/agent binary".to_string(),
|
|
"--prompt=keep spaces".to_string(),
|
|
"literal * $HOME".to_string(),
|
|
];
|
|
|
|
run::sandbox_create(
|
|
&server.endpoint,
|
|
"openshell",
|
|
run::SandboxCreateConfig {
|
|
name: Some("canonical-main"),
|
|
command: &command,
|
|
tty_override: Some(false),
|
|
..test_config()
|
|
},
|
|
"default",
|
|
&tls,
|
|
)
|
|
.await
|
|
.expect("sandbox create should succeed");
|
|
|
|
let requests = create_requests(&server).await;
|
|
let spec = requests[0]
|
|
.spec
|
|
.as_ref()
|
|
.expect("sandbox spec should be persisted at create time");
|
|
assert_eq!(spec.command, command);
|
|
assert!(!spec.tty);
|
|
assert!(requests[0].await_main_process_attachment);
|
|
}
|
|
|
|
#[tokio::test]
|
|
async fn detached_command_does_not_declare_main_process_attachment() {
|
|
let server = run_server().await;
|
|
let fake_ssh_dir = tempfile::tempdir().unwrap();
|
|
let xdg_dir = tempfile::tempdir().unwrap();
|
|
let _env = test_env(&fake_ssh_dir, &xdg_dir);
|
|
let tls = test_tls(&server);
|
|
install_fake_ssh(&fake_ssh_dir);
|
|
|
|
run::sandbox_create(
|
|
&server.endpoint,
|
|
"openshell",
|
|
run::SandboxCreateConfig {
|
|
name: Some("detached-main"),
|
|
command: &["echo".into(), "OK".into()],
|
|
detach: true,
|
|
..test_config()
|
|
},
|
|
"default",
|
|
&tls,
|
|
)
|
|
.await
|
|
.expect("detached sandbox create should succeed");
|
|
|
|
let requests = create_requests(&server).await;
|
|
assert!(!requests[0].await_main_process_attachment);
|
|
}
|
|
|
|
#[tokio::test]
|
|
async fn detached_ephemeral_command_delegates_cleanup_to_gateway() {
|
|
let server = run_server().await;
|
|
let fake_ssh_dir = tempfile::tempdir().unwrap();
|
|
let xdg_dir = tempfile::tempdir().unwrap();
|
|
let _env = test_env(&fake_ssh_dir, &xdg_dir);
|
|
let tls = test_tls(&server);
|
|
install_fake_ssh(&fake_ssh_dir);
|
|
|
|
run::sandbox_create(
|
|
&server.endpoint,
|
|
"openshell",
|
|
run::SandboxCreateConfig {
|
|
name: Some("detached-ephemeral-main"),
|
|
keep: false,
|
|
command: &["worker".into()],
|
|
detach: true,
|
|
..test_config()
|
|
},
|
|
"default",
|
|
&tls,
|
|
)
|
|
.await
|
|
.expect("detached ephemeral sandbox create should succeed");
|
|
|
|
let requests = create_requests(&server).await;
|
|
assert!(!requests[0].await_main_process_attachment);
|
|
assert_eq!(
|
|
requests[0]
|
|
.annotations
|
|
.get("openshell.nvidia.com/retention")
|
|
.map(String::as_str),
|
|
Some("ephemeral")
|
|
);
|
|
assert!(
|
|
deleted_names(&server).await.is_empty(),
|
|
"the gateway owns cleanup after a detached canonical process exits"
|
|
);
|
|
}
|
|
|
|
#[tokio::test]
|
|
async fn sandbox_create_sends_driver_config_json() {
|
|
let server = run_server().await;
|
|
let fake_ssh_dir = tempfile::tempdir().unwrap();
|
|
let xdg_dir = tempfile::tempdir().unwrap();
|
|
let _env = test_env(&fake_ssh_dir, &xdg_dir);
|
|
let tls = test_tls(&server);
|
|
install_fake_ssh(&fake_ssh_dir);
|
|
|
|
run::sandbox_create(
|
|
&server.endpoint,
|
|
"openshell",
|
|
run::SandboxCreateConfig {
|
|
name: Some("driver-config"),
|
|
driver_config_json: Some(
|
|
r#"{"kubernetes":{"pod":{"priority_class_name":"batch-low"}}}"#,
|
|
),
|
|
command: &["echo".into(), "OK".into()],
|
|
..test_config()
|
|
},
|
|
"default",
|
|
&tls,
|
|
)
|
|
.await
|
|
.expect("sandbox create should succeed");
|
|
|
|
let requests = create_requests(&server).await;
|
|
let driver_config = requests[0]
|
|
.spec
|
|
.as_ref()
|
|
.and_then(|spec| spec.template.as_ref())
|
|
.and_then(|template| template.driver_config.as_ref())
|
|
.expect("driver config should be sent");
|
|
let kubernetes = driver_config
|
|
.fields
|
|
.get("kubernetes")
|
|
.and_then(|value| value.kind.as_ref())
|
|
.and_then(|kind| match kind {
|
|
prost_types::value::Kind::StructValue(inner) => Some(inner),
|
|
_ => None,
|
|
})
|
|
.expect("kubernetes block should be a struct");
|
|
let pod = kubernetes
|
|
.fields
|
|
.get("pod")
|
|
.and_then(|value| value.kind.as_ref())
|
|
.and_then(|kind| match kind {
|
|
prost_types::value::Kind::StructValue(inner) => Some(inner),
|
|
_ => None,
|
|
})
|
|
.expect("pod block should be a struct");
|
|
|
|
assert_eq!(
|
|
pod.fields
|
|
.get("priority_class_name")
|
|
.and_then(|value| value.kind.as_ref())
|
|
.and_then(|kind| match kind {
|
|
prost_types::value::Kind::StringValue(value) => Some(value.as_str()),
|
|
_ => None,
|
|
}),
|
|
Some("batch-low")
|
|
);
|
|
}
|
|
|
|
#[tokio::test]
|
|
async fn sandbox_create_with_template_sends_workload_template_name() {
|
|
let server = run_server().await;
|
|
add_provider(&server, "github", "github").await;
|
|
let fake_ssh_dir = tempfile::tempdir().unwrap();
|
|
let xdg_dir = tempfile::tempdir().unwrap();
|
|
let _env = test_env(&fake_ssh_dir, &xdg_dir);
|
|
let tls = test_tls(&server);
|
|
install_fake_ssh(&fake_ssh_dir);
|
|
|
|
run::sandbox_create(
|
|
&server.endpoint,
|
|
"openshell",
|
|
run::SandboxCreateConfig {
|
|
name: Some("from-template"),
|
|
template: Some("gpu-kata"),
|
|
providers: &["github".to_string()],
|
|
command: &["echo".into(), "OK".into()],
|
|
..test_config()
|
|
},
|
|
"default",
|
|
&tls,
|
|
)
|
|
.await
|
|
.expect("sandbox create should succeed");
|
|
|
|
let requests = create_requests(&server).await;
|
|
let request = requests.first().expect("create request should be recorded");
|
|
assert_eq!(request.workload_template_name, "gpu-kata");
|
|
let spec = request
|
|
.spec
|
|
.as_ref()
|
|
.expect("governance spec should be sent");
|
|
assert_eq!(spec.providers, vec!["github".to_string()]);
|
|
assert!(spec.template.is_none());
|
|
assert!(spec.environment.is_empty());
|
|
assert!(spec.resource_requirements.is_none());
|
|
}
|
|
|
|
#[tokio::test]
|
|
async fn sandbox_template_create_sends_non_default_workspace_in_scope_and_metadata() {
|
|
let server = run_server().await;
|
|
let fake_ssh_dir = tempfile::tempdir().unwrap();
|
|
let xdg_dir = tempfile::tempdir().unwrap();
|
|
let _env = test_env(&fake_ssh_dir, &xdg_dir);
|
|
let tls = test_tls(&server);
|
|
|
|
run::sandbox_template_create(
|
|
&server.endpoint,
|
|
"gpu-kata",
|
|
Some("registry.example.com/agent:latest"),
|
|
Some("2"),
|
|
Some("4Gi"),
|
|
Some(GpuResourceRequirements { count: Some(1) }),
|
|
Some(r#"{"kubernetes":{"pod":{"node_selector":{"pool":"gpu"}}}}"#),
|
|
Some("5m"),
|
|
Some(3),
|
|
HashMap::from([("team".to_string(), "runtime".to_string())]),
|
|
HashMap::from([("owner".to_string(), "platform".to_string())]),
|
|
HashMap::from([("FEATURE_FLAG".to_string(), "on".to_string())]),
|
|
"table",
|
|
"team-a",
|
|
&tls,
|
|
)
|
|
.await
|
|
.expect("template create should succeed");
|
|
|
|
let requests = template_create_requests(&server).await;
|
|
let request = requests
|
|
.first()
|
|
.expect("template create request should be recorded");
|
|
assert_eq!(selected_workspace(&request.workspace_scope), Some("team-a"));
|
|
let template = request.template.as_ref().expect("template should be sent");
|
|
let metadata = template.metadata.as_ref().expect("metadata should be sent");
|
|
assert_eq!(metadata.name, "gpu-kata");
|
|
assert_eq!(metadata.workspace, "team-a");
|
|
assert_eq!(metadata.labels.get("team"), Some(&"runtime".to_string()));
|
|
assert_eq!(
|
|
metadata.annotations.get("owner"),
|
|
Some(&"platform".to_string())
|
|
);
|
|
|
|
let spec = template.spec.as_ref().expect("spec should be sent");
|
|
let workload = spec.workload.as_ref().expect("workload should be sent");
|
|
assert_eq!(workload.image, "registry.example.com/agent:latest");
|
|
assert_eq!(
|
|
workload.environment.get("FEATURE_FLAG"),
|
|
Some(&"on".to_string())
|
|
);
|
|
let resources = workload
|
|
.resources
|
|
.as_ref()
|
|
.expect("resources should be sent");
|
|
assert_eq!(resources.cpu, "2");
|
|
assert_eq!(resources.memory, "4Gi");
|
|
assert_eq!(resources.gpu.as_ref().and_then(|gpu| gpu.count), Some(1));
|
|
assert!(spec.driver_config.is_some());
|
|
let startup = spec
|
|
.desired_service_level
|
|
.as_ref()
|
|
.and_then(|service_level| service_level.startup.as_ref())
|
|
.expect("startup service level should be sent");
|
|
assert_eq!(startup.max_burst, 3);
|
|
assert_eq!(
|
|
startup
|
|
.ready_within
|
|
.as_ref()
|
|
.map(|duration| duration.seconds),
|
|
Some(300)
|
|
);
|
|
}
|
|
|
|
#[tokio::test]
|
|
async fn sandbox_template_list_and_delete_send_workspace_requests() {
|
|
let server = run_server().await;
|
|
let fake_ssh_dir = tempfile::tempdir().unwrap();
|
|
let xdg_dir = tempfile::tempdir().unwrap();
|
|
let _env = test_env(&fake_ssh_dir, &xdg_dir);
|
|
let tls = test_tls(&server);
|
|
|
|
run::sandbox_template_list(
|
|
&server.endpoint,
|
|
25,
|
|
"next-template-page",
|
|
Some("team=runtime"),
|
|
false,
|
|
"table",
|
|
"default",
|
|
false,
|
|
&tls,
|
|
)
|
|
.await
|
|
.expect("template list should succeed");
|
|
run::sandbox_template_delete(&server.endpoint, &["gpu-kata".to_string()], "default", &tls)
|
|
.await
|
|
.expect("template delete should succeed");
|
|
|
|
let list_requests = template_list_requests(&server).await;
|
|
let list_request = list_requests
|
|
.first()
|
|
.expect("template list request should be recorded");
|
|
assert_eq!(list_request.page_size, 25);
|
|
assert_eq!(list_request.page_token, "next-template-page");
|
|
assert_eq!(list_request.label_selector, "team=runtime");
|
|
assert_eq!(
|
|
selected_workspace(&list_request.workspace_scope),
|
|
Some("default")
|
|
);
|
|
|
|
let delete_requests = template_delete_requests(&server).await;
|
|
let delete_request = delete_requests
|
|
.first()
|
|
.expect("template delete request should be recorded");
|
|
assert_eq!(delete_request.name, "gpu-kata");
|
|
assert_eq!(
|
|
selected_workspace(&delete_request.workspace_scope),
|
|
Some("default")
|
|
);
|
|
}
|
|
|
|
#[tokio::test]
|
|
async fn sandbox_template_create_allows_omitted_image() {
|
|
let server = run_server().await;
|
|
let fake_ssh_dir = tempfile::tempdir().unwrap();
|
|
let xdg_dir = tempfile::tempdir().unwrap();
|
|
let _env = test_env(&fake_ssh_dir, &xdg_dir);
|
|
let tls = test_tls(&server);
|
|
|
|
run::sandbox_template_create(
|
|
&server.endpoint,
|
|
"base",
|
|
None,
|
|
None,
|
|
None,
|
|
None,
|
|
None,
|
|
None,
|
|
None,
|
|
HashMap::new(),
|
|
HashMap::new(),
|
|
HashMap::new(),
|
|
"table",
|
|
"default",
|
|
&tls,
|
|
)
|
|
.await
|
|
.expect("template create without image should succeed");
|
|
|
|
let requests = template_create_requests(&server).await;
|
|
let request = requests
|
|
.first()
|
|
.expect("template create request should be recorded");
|
|
let workload = request
|
|
.template
|
|
.as_ref()
|
|
.and_then(|template| template.spec.as_ref())
|
|
.and_then(|spec| spec.workload.as_ref())
|
|
.expect("workload should be sent");
|
|
assert_eq!(workload.image, "");
|
|
}
|
|
|
|
#[tokio::test]
|
|
async fn sandbox_create_sends_gpu_default_request() {
|
|
let server = run_server().await;
|
|
let fake_ssh_dir = tempfile::tempdir().unwrap();
|
|
let xdg_dir = tempfile::tempdir().unwrap();
|
|
let _env = test_env(&fake_ssh_dir, &xdg_dir);
|
|
let tls = test_tls(&server);
|
|
install_fake_ssh(&fake_ssh_dir);
|
|
|
|
run::sandbox_create(
|
|
&server.endpoint,
|
|
"openshell",
|
|
run::SandboxCreateConfig {
|
|
name: Some("gpu-default"),
|
|
gpu_requirements: Some(gpu_requirements(None)),
|
|
command: &["echo".into(), "OK".into()],
|
|
..test_config()
|
|
},
|
|
"default",
|
|
&tls,
|
|
)
|
|
.await
|
|
.expect("sandbox create should succeed");
|
|
|
|
let requests = create_requests(&server).await;
|
|
let gpu = requests[0]
|
|
.spec
|
|
.as_ref()
|
|
.and_then(|spec| spec.resource_requirements.as_ref())
|
|
.and_then(|requirements| requirements.gpu.as_ref())
|
|
.expect("GPU requirement should be sent");
|
|
|
|
assert_eq!(gpu.count, None);
|
|
}
|
|
|
|
#[tokio::test]
|
|
async fn sandbox_create_sends_gpu_count_request() {
|
|
let server = run_server().await;
|
|
let fake_ssh_dir = tempfile::tempdir().unwrap();
|
|
let xdg_dir = tempfile::tempdir().unwrap();
|
|
let _env = test_env(&fake_ssh_dir, &xdg_dir);
|
|
let tls = test_tls(&server);
|
|
install_fake_ssh(&fake_ssh_dir);
|
|
|
|
run::sandbox_create(
|
|
&server.endpoint,
|
|
"openshell",
|
|
run::SandboxCreateConfig {
|
|
name: Some("gpu-two"),
|
|
gpu_requirements: Some(gpu_requirements(Some(2))),
|
|
command: &["echo".into(), "OK".into()],
|
|
..test_config()
|
|
},
|
|
"default",
|
|
&tls,
|
|
)
|
|
.await
|
|
.expect("sandbox create should succeed");
|
|
|
|
let requests = create_requests(&server).await;
|
|
let gpu = requests[0]
|
|
.spec
|
|
.as_ref()
|
|
.and_then(|spec| spec.resource_requirements.as_ref())
|
|
.and_then(|requirements| requirements.gpu.as_ref())
|
|
.expect("GPU requirement should be sent");
|
|
|
|
assert_eq!(gpu.count, Some(2));
|
|
}
|
|
|
|
#[tokio::test]
|
|
async fn sandbox_create_skips_inferred_provider_without_local_credentials() {
|
|
let server = run_server().await;
|
|
let fake_ssh_dir = tempfile::tempdir().unwrap();
|
|
let xdg_dir = tempfile::tempdir().unwrap();
|
|
let _env = test_env(&fake_ssh_dir, &xdg_dir);
|
|
let tls = test_tls(&server);
|
|
install_fake_ssh(&fake_ssh_dir);
|
|
|
|
run::sandbox_create(
|
|
&server.endpoint,
|
|
"openshell",
|
|
run::SandboxCreateConfig {
|
|
name: Some("no-inferred-provider"),
|
|
command: &["claude".into(), "--version".into()],
|
|
tty_override: Some(true),
|
|
..test_config()
|
|
},
|
|
"default",
|
|
&tls,
|
|
)
|
|
.await
|
|
.expect("sandbox create should succeed without inferred provider");
|
|
|
|
let requests = create_requests(&server).await;
|
|
let providers = requests[0]
|
|
.spec
|
|
.as_ref()
|
|
.expect("sandbox spec should be sent")
|
|
.providers
|
|
.clone();
|
|
assert!(
|
|
providers.is_empty(),
|
|
"missing local credentials should skip inferred providers, got {providers:?}"
|
|
);
|
|
}
|
|
|
|
#[tokio::test]
|
|
async fn sandbox_create_returns_vm_error_without_waiting_for_timeout() {
|
|
let server = run_server().await;
|
|
server
|
|
.openshell
|
|
.state
|
|
.vm_error_after_started
|
|
.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 started_at = Instant::now();
|
|
let err = run::sandbox_create(
|
|
&server.endpoint,
|
|
"openshell",
|
|
run::SandboxCreateConfig {
|
|
name: Some("vm-error"),
|
|
command: &["echo".into(), "OK".into()],
|
|
..test_config()
|
|
},
|
|
"default",
|
|
&tls,
|
|
)
|
|
.await
|
|
.expect_err("sandbox create should fail on terminal VM error");
|
|
|
|
assert!(
|
|
started_at.elapsed() < Duration::from_secs(2),
|
|
"terminal VM errors should not wait for the provisioning timeout"
|
|
);
|
|
let rendered = err.to_string();
|
|
assert!(rendered.contains("sandbox entered error phase while provisioning"));
|
|
assert!(rendered.contains("ProcessExited: VM process exited with status 0"));
|
|
assert!(!rendered.contains("timed out"));
|
|
}
|
|
|
|
#[tokio::test]
|
|
async fn sandbox_create_preserves_vm_error_when_exit_code_is_observed() {
|
|
let server = run_server().await;
|
|
server
|
|
.openshell
|
|
.state
|
|
.vm_error_after_started
|
|
.store(true, Ordering::SeqCst);
|
|
server
|
|
.openshell
|
|
.state
|
|
.vm_error_with_observed_exit
|
|
.store(true, Ordering::SeqCst);
|
|
let fake_ssh_dir = tempfile::tempdir().unwrap();
|
|
let xdg_dir = tempfile::tempdir().unwrap();
|
|
let _env = test_env(&fake_ssh_dir, &xdg_dir);
|
|
let tls = test_tls(&server);
|
|
install_fake_ssh(&fake_ssh_dir);
|
|
|
|
let err = run::sandbox_create(
|
|
&server.endpoint,
|
|
"openshell",
|
|
run::SandboxCreateConfig {
|
|
name: Some("vm-error-with-exit"),
|
|
command: &["echo".into(), "OK".into()],
|
|
..test_config()
|
|
},
|
|
"default",
|
|
&tls,
|
|
)
|
|
.await
|
|
.expect_err("an observed process exit must not hide the infrastructure error");
|
|
|
|
let rendered = err.to_string();
|
|
assert!(rendered.contains("sandbox entered error phase while provisioning"));
|
|
assert!(rendered.contains("ProcessExited: VM process exited with status 0"));
|
|
}
|
|
|
|
#[tokio::test]
|
|
async fn sandbox_create_keeps_waiting_while_vm_progress_arrives() {
|
|
let server = run_server().await;
|
|
server
|
|
.openshell
|
|
.state
|
|
.vm_slow_progress_before_ready
|
|
.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);
|
|
|
|
run::sandbox_create(
|
|
&server.endpoint,
|
|
"openshell",
|
|
run::SandboxCreateConfig {
|
|
name: Some("vm-slow-progress"),
|
|
command: &["echo".into(), "OK".into()],
|
|
..test_config()
|
|
},
|
|
"default",
|
|
&tls,
|
|
)
|
|
.await
|
|
.expect("sandbox create should not time out while VM progress is active");
|
|
}
|
|
|
|
#[tokio::test]
|
|
async fn sandbox_create_times_out_when_only_logs_arrive() {
|
|
let server = run_server().await;
|
|
server
|
|
.openshell
|
|
.state
|
|
.vm_log_churn_before_ready
|
|
.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 started_at = Instant::now();
|
|
let err = run::sandbox_create(
|
|
&server.endpoint,
|
|
"openshell",
|
|
run::SandboxCreateConfig {
|
|
name: Some("vm-log-churn"),
|
|
command: &["echo".into(), "OK".into()],
|
|
..test_config()
|
|
},
|
|
"default",
|
|
&tls,
|
|
)
|
|
.await
|
|
.expect_err("sandbox create should time out when only logs arrive");
|
|
|
|
assert!(
|
|
started_at.elapsed() < Duration::from_secs(2),
|
|
"logs should not extend the provisioning timeout"
|
|
);
|
|
assert!(err.to_string().contains("sandbox provisioning timed out"));
|
|
}
|
|
|
|
#[tokio::test]
|
|
async fn sandbox_create_retries_terminal_attachment_until_relay_registers() {
|
|
let server = run_server().await;
|
|
server
|
|
.openshell
|
|
.state
|
|
.terminal_before_relay
|
|
.store(true, Ordering::SeqCst);
|
|
server
|
|
.openshell
|
|
.state
|
|
.ssh_session_failures_remaining
|
|
.store(1, Ordering::SeqCst);
|
|
let fake_ssh_dir = tempfile::tempdir().unwrap();
|
|
let xdg_dir = tempfile::tempdir().unwrap();
|
|
let _env = test_env(&fake_ssh_dir, &xdg_dir);
|
|
let tls = test_tls(&server);
|
|
install_fake_ssh(&fake_ssh_dir);
|
|
|
|
let exit_code = run::sandbox_create(
|
|
&server.endpoint,
|
|
"openshell",
|
|
run::SandboxCreateConfig {
|
|
name: Some("fast-command"),
|
|
command: &["echo".into(), "OK".into()],
|
|
..test_config()
|
|
},
|
|
"default",
|
|
&tls,
|
|
)
|
|
.await
|
|
.expect("sandbox create should wait for the declared terminal attachment relay");
|
|
|
|
assert_eq!(exit_code, 0);
|
|
assert_eq!(
|
|
server
|
|
.openshell
|
|
.state
|
|
.ssh_session_requests
|
|
.load(Ordering::SeqCst),
|
|
2
|
|
);
|
|
}
|
|
|
|
#[tokio::test]
|
|
async fn sandbox_create_waits_for_main_result_after_provisional_container_exit() {
|
|
let server = run_server().await;
|
|
server
|
|
.openshell
|
|
.state
|
|
.terminal_after_provisional_container_exit
|
|
.store(true, Ordering::SeqCst);
|
|
let fake_ssh_dir = tempfile::tempdir().unwrap();
|
|
let xdg_dir = tempfile::tempdir().unwrap();
|
|
let _env = test_env(&fake_ssh_dir, &xdg_dir);
|
|
let tls = test_tls(&server);
|
|
install_fake_ssh(&fake_ssh_dir);
|
|
|
|
let exit_code = run::sandbox_create(
|
|
&server.endpoint,
|
|
"openshell",
|
|
run::SandboxCreateConfig {
|
|
name: Some("fast-ephemeral-command"),
|
|
keep: false,
|
|
command: &["echo".into(), "OK".into()],
|
|
..test_config()
|
|
},
|
|
"default",
|
|
&tls,
|
|
)
|
|
.await
|
|
.expect("a provisional container exit must yield to the canonical main-process result");
|
|
|
|
assert_eq!(exit_code, 0);
|
|
assert_eq!(
|
|
deleted_names(&server).await,
|
|
vec![vec!["fast-ephemeral-command".to_string()]]
|
|
);
|
|
}
|
|
|
|
#[tokio::test]
|
|
async fn sandbox_create_bounds_provisional_container_exit_reconciliation() {
|
|
let server = run_server().await;
|
|
server
|
|
.openshell
|
|
.state
|
|
.provisional_container_exit_without_result
|
|
.store(true, Ordering::SeqCst);
|
|
let fake_ssh_dir = tempfile::tempdir().unwrap();
|
|
let xdg_dir = tempfile::tempdir().unwrap();
|
|
let _env = test_env(&fake_ssh_dir, &xdg_dir);
|
|
let tls = test_tls(&server);
|
|
install_fake_ssh(&fake_ssh_dir);
|
|
|
|
let provisional_container_exit_sent = server
|
|
.openshell
|
|
.state
|
|
.provisional_container_exit_sent
|
|
.notified();
|
|
tokio::pin!(provisional_container_exit_sent);
|
|
let command = ["echo".into(), "OK".into()];
|
|
let create = run::sandbox_create(
|
|
&server.endpoint,
|
|
"openshell",
|
|
run::SandboxCreateConfig {
|
|
name: Some("missing-main-result"),
|
|
command: &command,
|
|
..test_config()
|
|
},
|
|
"default",
|
|
&tls,
|
|
);
|
|
tokio::pin!(create);
|
|
|
|
tokio::select! {
|
|
() = &mut provisional_container_exit_sent => {}
|
|
result = &mut create => panic!("sandbox create returned before the provisional exit was observed: {result:?}"),
|
|
}
|
|
let reconciliation_started = Instant::now();
|
|
let result = tokio::time::timeout(Duration::from_secs(10), &mut create)
|
|
.await
|
|
.expect("provisional container exit reconciliation must remain bounded");
|
|
let reconciliation_elapsed = reconciliation_started.elapsed();
|
|
|
|
let err = result
|
|
.expect_err("a missing canonical main-process result must retain the container exit error");
|
|
|
|
let rendered = err.to_string();
|
|
assert!(
|
|
rendered.contains("sandbox entered error phase while provisioning"),
|
|
"unexpected error: {rendered}"
|
|
);
|
|
assert!(
|
|
rendered.contains("ContainerExited: Sandbox container exited"),
|
|
"unexpected error: {rendered}"
|
|
);
|
|
assert!(
|
|
!rendered.contains("timed out"),
|
|
"unexpected error: {rendered}"
|
|
);
|
|
assert!(
|
|
reconciliation_elapsed >= Duration::from_secs(5),
|
|
"provisional container exit returned before reconciliation: {reconciliation_elapsed:?}"
|
|
);
|
|
assert!(deleted_names(&server).await.is_empty());
|
|
}
|
|
|
|
#[tokio::test]
|
|
async fn sandbox_create_deletes_command_sessions_with_no_keep() {
|
|
let server = run_server().await;
|
|
let fake_ssh_dir = tempfile::tempdir().unwrap();
|
|
let xdg_dir = tempfile::tempdir().unwrap();
|
|
let _env = test_env(&fake_ssh_dir, &xdg_dir);
|
|
let tls = test_tls(&server);
|
|
install_fake_ssh(&fake_ssh_dir);
|
|
|
|
run::sandbox_create(
|
|
&server.endpoint,
|
|
"openshell",
|
|
run::SandboxCreateConfig {
|
|
name: Some("ephemeral-command"),
|
|
keep: false,
|
|
command: &["echo".into(), "OK".into()],
|
|
..test_config()
|
|
},
|
|
"default",
|
|
&tls,
|
|
)
|
|
.await
|
|
.expect("sandbox create should succeed");
|
|
|
|
assert_eq!(
|
|
deleted_names(&server).await,
|
|
vec![vec!["ephemeral-command".to_string()]]
|
|
);
|
|
let requests = create_requests(&server).await;
|
|
assert_eq!(
|
|
requests[0]
|
|
.annotations
|
|
.get("openshell.nvidia.com/retention")
|
|
.map(String::as_str),
|
|
Some("ephemeral")
|
|
);
|
|
assert_eq!(
|
|
load_last_sandbox("openshell", "default"),
|
|
None,
|
|
"no-keep sandboxes should not be persisted as last-used"
|
|
);
|
|
}
|
|
|
|
#[tokio::test]
|
|
async fn sandbox_create_returns_exact_main_status_after_no_keep_cleanup() {
|
|
let server = run_server().await;
|
|
let fake_ssh_dir = tempfile::tempdir().unwrap();
|
|
let xdg_dir = tempfile::tempdir().unwrap();
|
|
let _env = test_env(&fake_ssh_dir, &xdg_dir);
|
|
let tls = test_tls(&server);
|
|
install_executable_script(&fake_ssh_dir, "ssh", "#!/bin/sh\nexit 7\n");
|
|
|
|
let exit_code = run::sandbox_create(
|
|
&server.endpoint,
|
|
"openshell",
|
|
run::SandboxCreateConfig {
|
|
name: Some("ephemeral-failure"),
|
|
keep: false,
|
|
command: &["sh".into(), "-c".into(), "exit 7".into()],
|
|
..test_config()
|
|
},
|
|
"default",
|
|
&tls,
|
|
)
|
|
.await
|
|
.expect("a main-process failure is a command result, not a cleanup error");
|
|
|
|
assert_eq!(exit_code, 7);
|
|
assert_eq!(
|
|
deleted_names(&server).await,
|
|
vec![vec!["ephemeral-failure".to_string()]]
|
|
);
|
|
}
|
|
|
|
#[tokio::test]
|
|
async fn sandbox_create_deletes_shell_sessions_with_no_keep() {
|
|
let server = run_server().await;
|
|
let fake_ssh_dir = tempfile::tempdir().unwrap();
|
|
let xdg_dir = tempfile::tempdir().unwrap();
|
|
let _env = test_env(&fake_ssh_dir, &xdg_dir);
|
|
let tls = test_tls(&server);
|
|
install_fake_ssh(&fake_ssh_dir);
|
|
|
|
run::sandbox_create(
|
|
&server.endpoint,
|
|
"openshell",
|
|
run::SandboxCreateConfig {
|
|
name: Some("ephemeral-shell"),
|
|
keep: false,
|
|
tty_override: Some(true),
|
|
..test_config()
|
|
},
|
|
"default",
|
|
&tls,
|
|
)
|
|
.await
|
|
.expect("sandbox create shell should succeed");
|
|
|
|
assert_eq!(
|
|
deleted_names(&server).await,
|
|
vec![vec!["ephemeral-shell".to_string()]]
|
|
);
|
|
assert_eq!(
|
|
load_last_sandbox("openshell", "default"),
|
|
None,
|
|
"no-keep shell sessions should not be persisted as last-used"
|
|
);
|
|
}
|
|
|
|
#[tokio::test]
|
|
async fn sandbox_create_keeps_sandbox_with_hidden_keep_flag() {
|
|
let server = run_server().await;
|
|
let fake_ssh_dir = tempfile::tempdir().unwrap();
|
|
let xdg_dir = tempfile::tempdir().unwrap();
|
|
let _env = test_env(&fake_ssh_dir, &xdg_dir);
|
|
let tls = test_tls(&server);
|
|
install_fake_ssh(&fake_ssh_dir);
|
|
|
|
run::sandbox_create(
|
|
&server.endpoint,
|
|
"openshell",
|
|
run::SandboxCreateConfig {
|
|
name: Some("persistent-keep"),
|
|
command: &["echo".into(), "OK".into()],
|
|
..test_config()
|
|
},
|
|
"default",
|
|
&tls,
|
|
)
|
|
.await
|
|
.expect("sandbox create should succeed");
|
|
|
|
assert!(deleted_names(&server).await.is_empty());
|
|
assert_eq!(
|
|
load_last_sandbox("openshell", "default").as_deref(),
|
|
Some("persistent-keep"),
|
|
"persistent sandboxes should remain selectable as last-used"
|
|
);
|
|
}
|
|
|
|
#[tokio::test]
|
|
async fn sandbox_create_keeps_sandbox_with_forwarding() {
|
|
let server = run_server().await;
|
|
let fake_ssh_dir = tempfile::tempdir().unwrap();
|
|
let xdg_dir = tempfile::tempdir().unwrap();
|
|
let _env = test_env(&fake_ssh_dir, &xdg_dir);
|
|
let tls = test_tls(&server);
|
|
install_fake_forwarding_ssh(&fake_ssh_dir);
|
|
let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
|
|
let forward_port = listener.local_addr().unwrap().port();
|
|
drop(listener);
|
|
|
|
run::sandbox_create(
|
|
&server.endpoint,
|
|
"openshell",
|
|
run::SandboxCreateConfig {
|
|
name: Some("persistent-forward"),
|
|
keep: false,
|
|
forward: Some(openshell_core::forward::ForwardSpec::new(forward_port)),
|
|
command: &["echo".into(), "OK".into()],
|
|
..test_config()
|
|
},
|
|
"default",
|
|
&tls,
|
|
)
|
|
.await
|
|
.expect("sandbox create with forward should succeed");
|
|
|
|
assert!(deleted_names(&server).await.is_empty());
|
|
let record = openshell_core::forward::read_forward_pid("persistent-forward", forward_port)
|
|
.expect("fake forward should be tracked");
|
|
let _ = std::process::Command::new("kill")
|
|
.arg(record.pid.to_string())
|
|
.status();
|
|
}
|
|
|
|
#[tokio::test]
|
|
async fn sandbox_forward_background_tracks_owned_child_when_pid_discovery_fails() {
|
|
let server = run_server().await;
|
|
let fake_ssh_dir = tempfile::tempdir().unwrap();
|
|
let xdg_dir = tempfile::tempdir().unwrap();
|
|
let _env = test_env(&fake_ssh_dir, &xdg_dir);
|
|
let tls = test_tls(&server);
|
|
install_fake_forwarding_ssh(&fake_ssh_dir);
|
|
install_fake_pgrep_no_match(&fake_ssh_dir);
|
|
let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
|
|
let forward_port = listener.local_addr().unwrap().port();
|
|
drop(listener);
|
|
|
|
let spec = openshell_core::forward::ForwardSpec::new(forward_port);
|
|
run::sandbox_forward(
|
|
&server.endpoint,
|
|
"owned-forward",
|
|
&spec,
|
|
true,
|
|
&tls,
|
|
"default",
|
|
)
|
|
.await
|
|
.expect("background forward should track the owned SSH child without PID discovery");
|
|
let record = openshell_core::forward::read_forward_pid("owned-forward", forward_port)
|
|
.expect("owned background forward should write a PID file");
|
|
|
|
assert!(
|
|
openshell_core::forward::stop_forward("owned-forward", forward_port)
|
|
.expect("tracked fake forward should stop"),
|
|
"tracked fake forward should be recognized as alive and stopped",
|
|
);
|
|
assert!(
|
|
wait_for_process_exit(record.pid, Duration::from_secs(2)).await,
|
|
"tracked fake forward process should exit after stop"
|
|
);
|
|
}
|
|
|
|
#[tokio::test]
|
|
#[ignore = "flaky under concurrent test execution"]
|
|
async fn sandbox_forward_foreground_fails_when_ssh_exits_before_listener_opens() {
|
|
let server = run_server().await;
|
|
let fake_ssh_dir = tempfile::tempdir().unwrap();
|
|
let xdg_dir = tempfile::tempdir().unwrap();
|
|
let _env = test_env(&fake_ssh_dir, &xdg_dir);
|
|
let tls = test_tls(&server);
|
|
install_fake_ssh(&fake_ssh_dir);
|
|
let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
|
|
let forward_port = listener.local_addr().unwrap().port();
|
|
drop(listener);
|
|
|
|
let spec = openshell_core::forward::ForwardSpec::new(forward_port);
|
|
let err = run::sandbox_forward(
|
|
&server.endpoint,
|
|
"foreground-forward",
|
|
&spec,
|
|
false,
|
|
&tls,
|
|
"default",
|
|
)
|
|
.await
|
|
.expect_err("foreground forward should fail when ssh exits before listener readiness");
|
|
let msg = format!("{err}");
|
|
assert!(
|
|
msg.contains("ssh exited before local forward listener opened"),
|
|
"error should explain that ssh exited before listener readiness, got: {msg}",
|
|
);
|
|
}
|
|
|
|
#[tokio::test]
|
|
#[ignore = "flaky under concurrent test execution"]
|
|
async fn sandbox_forward_background_terminates_owned_child_when_listener_never_opens() {
|
|
let server = run_server().await;
|
|
let fake_ssh_dir = tempfile::tempdir().unwrap();
|
|
let xdg_dir = tempfile::tempdir().unwrap();
|
|
let _env = test_env(&fake_ssh_dir, &xdg_dir);
|
|
let tls = test_tls(&server);
|
|
let fake_forward = install_fake_unreachable_forwarding_ssh(&fake_ssh_dir);
|
|
let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
|
|
let forward_port = listener.local_addr().unwrap().port();
|
|
drop(listener);
|
|
|
|
let spec = openshell_core::forward::ForwardSpec::new(forward_port);
|
|
let err = run::sandbox_forward(
|
|
&server.endpoint,
|
|
"unreachable-forward",
|
|
&spec,
|
|
true,
|
|
&tls,
|
|
"default",
|
|
)
|
|
.await
|
|
.expect_err("background forward should fail when the listener never opens");
|
|
let msg = format!("{err}");
|
|
assert!(
|
|
msg.contains("ssh process started but local forward listener was not reachable"),
|
|
"error should preserve listener startup context, got: {msg}",
|
|
);
|
|
assert!(
|
|
openshell_core::forward::read_forward_pid("unreachable-forward", forward_port).is_none(),
|
|
"unreachable background forwards must not write a PID file",
|
|
);
|
|
let pid = fs::read_to_string(&fake_forward.pid_path)
|
|
.expect("fake forward should record a PID")
|
|
.trim()
|
|
.parse::<u32>()
|
|
.expect("fake forward PID should be numeric");
|
|
if !wait_for_process_exit(pid, Duration::from_secs(2)).await {
|
|
let log = fs::read_to_string(&fake_forward.log_path).unwrap_or_default();
|
|
let command = std::process::Command::new("ps")
|
|
.arg("-ww")
|
|
.arg("-o")
|
|
.arg("command=")
|
|
.arg("-p")
|
|
.arg(pid.to_string())
|
|
.output()
|
|
.ok()
|
|
.map(|output| String::from_utf8_lossy(&output.stdout).to_string())
|
|
.unwrap_or_default();
|
|
panic!(
|
|
"owned background SSH child should exit after listener failure cleanup; pid={}, command={}, log={}",
|
|
pid,
|
|
command.trim(),
|
|
log.trim(),
|
|
);
|
|
}
|
|
}
|
|
|
|
#[tokio::test]
|
|
async fn sandbox_create_sends_environment_variables() {
|
|
let server = run_server().await;
|
|
let fake_ssh_dir = tempfile::tempdir().unwrap();
|
|
let xdg_dir = tempfile::tempdir().unwrap();
|
|
let _env = test_env(&fake_ssh_dir, &xdg_dir);
|
|
let tls = test_tls(&server);
|
|
install_fake_ssh(&fake_ssh_dir);
|
|
|
|
run::sandbox_create(
|
|
&server.endpoint,
|
|
"openshell",
|
|
run::SandboxCreateConfig {
|
|
name: Some("env-test"),
|
|
command: &["echo".into(), "OK".into()],
|
|
environment: HashMap::from([
|
|
("FOO".into(), "bar".into()),
|
|
("BAZ".into(), "qux=with=equals".into()),
|
|
]),
|
|
..test_config()
|
|
},
|
|
"default",
|
|
&tls,
|
|
)
|
|
.await
|
|
.expect("sandbox create should succeed");
|
|
|
|
let requests = create_requests(&server).await;
|
|
let environment = &requests[0]
|
|
.spec
|
|
.as_ref()
|
|
.expect("spec should be present")
|
|
.environment;
|
|
assert_eq!(environment.get("FOO").map(String::as_str), Some("bar"));
|
|
assert_eq!(
|
|
environment.get("BAZ").map(String::as_str),
|
|
Some("qux=with=equals")
|
|
);
|
|
assert_eq!(environment.len(), 2);
|
|
}
|
|
|
|
#[tokio::test]
|
|
async fn sandbox_create_env_rejects_invalid_format() {
|
|
let err = run::parse_key_value_pairs(
|
|
&["VALID=ok".to_string(), "NOEQUALSSIGN".to_string()],
|
|
"--env",
|
|
)
|
|
.unwrap_err();
|
|
let msg = format!("{err}");
|
|
assert!(
|
|
msg.contains("--env") && msg.contains("NOEQUALSSIGN"),
|
|
"error should mention the flag and bad value, got: {msg}"
|
|
);
|
|
}
|
|
|
|
#[tokio::test]
|
|
async fn sandbox_create_env_rejects_reserved_prefix() {
|
|
let err = run::parse_env_pairs(&["VALID=ok".to_string(), "OPENSHELL_SECRET=bad".to_string()])
|
|
.unwrap_err();
|
|
let msg = format!("{err}");
|
|
assert!(
|
|
msg.contains("OPENSHELL_") && msg.contains("reserved"),
|
|
"error should mention reserved prefix, got: {msg}"
|
|
);
|
|
}
|
|
|
|
#[tokio::test]
|
|
async fn sandbox_create_env_rejects_invalid_key_name() {
|
|
let err = run::parse_env_pairs(&["1BAD=value".to_string()]).unwrap_err();
|
|
let msg = format!("{err}");
|
|
assert!(
|
|
msg.contains("1BAD"),
|
|
"error should mention invalid key, got: {msg}"
|
|
);
|
|
|
|
let err = run::parse_env_pairs(&["BAD-NAME=value".to_string()]).unwrap_err();
|
|
let msg = format!("{err}");
|
|
assert!(
|
|
msg.contains("BAD-NAME"),
|
|
"error should mention invalid key, got: {msg}"
|
|
);
|
|
}
|
|
|
|
fn prepare_cli_xdg(server: &TestServer, xdg_dir: &TempDir) {
|
|
let tls_dir = xdg_dir.path().join("openshell/gateways/openshell/mtls");
|
|
fs::create_dir_all(&tls_dir).unwrap();
|
|
for filename in ["ca.crt", "tls.crt", "tls.key"] {
|
|
fs::copy(server.dir.path().join(filename), tls_dir.join(filename)).unwrap();
|
|
}
|
|
}
|
|
|
|
async fn run_cli_sandbox_create_with_xdg(
|
|
server: &TestServer,
|
|
xdg_dir: &TempDir,
|
|
name: &str,
|
|
extra_args: &[&str],
|
|
) -> std::process::Output {
|
|
run_cli_sandbox_create_with_xdg_and_env(server, xdg_dir, name, extra_args, &[]).await
|
|
}
|
|
|
|
async fn run_cli_sandbox_create_with_xdg_and_env(
|
|
server: &TestServer,
|
|
xdg_dir: &TempDir,
|
|
name: &str,
|
|
extra_args: &[&str],
|
|
environment: &[(&str, &str)],
|
|
) -> std::process::Output {
|
|
let mut cmd = tokio::process::Command::new(env!("CARGO_BIN_EXE_openshell"));
|
|
for (key, _) in std::env::vars().filter(|(k, _)| k.starts_with("OPENSHELL_")) {
|
|
cmd.env_remove(&key);
|
|
}
|
|
cmd.args([
|
|
"--gateway",
|
|
"openshell",
|
|
"--gateway-endpoint",
|
|
&server.endpoint,
|
|
"sandbox",
|
|
"create",
|
|
"--name",
|
|
name,
|
|
"--no-tty",
|
|
"--no-auto-providers",
|
|
])
|
|
.args(extra_args)
|
|
.envs(environment.iter().copied())
|
|
.env("XDG_CONFIG_HOME", xdg_dir.path())
|
|
.env("HOME", xdg_dir.path())
|
|
.env("OPENSHELL_PROVISION_TIMEOUT", "5")
|
|
.output()
|
|
.await
|
|
.unwrap()
|
|
}
|
|
|
|
async fn run_cli_sandbox_create(
|
|
server: &TestServer,
|
|
name: &str,
|
|
extra_args: &[&str],
|
|
) -> std::process::Output {
|
|
let xdg_dir = tempfile::tempdir().unwrap();
|
|
prepare_cli_xdg(server, &xdg_dir);
|
|
run_cli_sandbox_create_with_xdg(server, &xdg_dir, name, extra_args).await
|
|
}
|
|
|
|
async fn run_cli_sandbox_create_with_env(
|
|
server: &TestServer,
|
|
name: &str,
|
|
extra_args: &[&str],
|
|
environment: &[(&str, &str)],
|
|
) -> std::process::Output {
|
|
let xdg_dir = tempfile::tempdir().unwrap();
|
|
prepare_cli_xdg(server, &xdg_dir);
|
|
run_cli_sandbox_create_with_xdg_and_env(server, &xdg_dir, name, extra_args, environment).await
|
|
}
|
|
|
|
#[tokio::test]
|
|
async fn sandbox_create_env_from_reaches_request() {
|
|
let server = run_server().await;
|
|
let value = "qualification-value-not-in-argv";
|
|
|
|
let output = run_cli_sandbox_create_with_env(
|
|
&server,
|
|
"env-from-test",
|
|
&["--env-from", "SANDBOX_VALUE=HOST_VALUE", "--output=json"],
|
|
&[("HOST_VALUE", value)],
|
|
)
|
|
.await;
|
|
assert!(
|
|
output.status.success(),
|
|
"sandbox create failed: {}",
|
|
String::from_utf8_lossy(&output.stderr)
|
|
);
|
|
|
|
let requests = create_requests(&server).await;
|
|
let environment = &requests[0]
|
|
.spec
|
|
.as_ref()
|
|
.expect("spec should be present")
|
|
.environment;
|
|
assert_eq!(
|
|
environment.get("SANDBOX_VALUE").map(String::as_str),
|
|
Some(value)
|
|
);
|
|
}
|
|
|
|
async fn run_cli_sandbox_template_create(
|
|
server: &TestServer,
|
|
name: &str,
|
|
extra_args: &[&str],
|
|
) -> std::process::Output {
|
|
let xdg_dir = tempfile::tempdir().unwrap();
|
|
let tls_dir = xdg_dir.path().join("openshell/gateways/openshell/mtls");
|
|
fs::create_dir_all(&tls_dir).unwrap();
|
|
for filename in ["ca.crt", "tls.crt", "tls.key"] {
|
|
fs::copy(server.dir.path().join(filename), tls_dir.join(filename)).unwrap();
|
|
}
|
|
|
|
let mut cmd = tokio::process::Command::new(env!("CARGO_BIN_EXE_openshell"));
|
|
for (key, _) in std::env::vars().filter(|(k, _)| k.starts_with("OPENSHELL_")) {
|
|
cmd.env_remove(&key);
|
|
}
|
|
cmd.args([
|
|
"--gateway",
|
|
"openshell",
|
|
"--gateway-endpoint",
|
|
&server.endpoint,
|
|
"sandbox",
|
|
"template",
|
|
"create",
|
|
name,
|
|
"--image",
|
|
"registry.example.com/agent:latest",
|
|
"--output=json",
|
|
])
|
|
.args(extra_args)
|
|
.env("XDG_CONFIG_HOME", xdg_dir.path())
|
|
.env("HOME", xdg_dir.path())
|
|
.env("OPENSHELL_PROVISION_TIMEOUT", "5")
|
|
.output()
|
|
.await
|
|
.unwrap()
|
|
}
|
|
|
|
async fn run_rejected_oidc_refresh_server() -> String {
|
|
let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
|
|
let issuer = format!("http://{}", listener.local_addr().unwrap());
|
|
let task_issuer = issuer.clone();
|
|
|
|
tokio::spawn(async move {
|
|
loop {
|
|
let Ok((mut stream, _)) = listener.accept().await else {
|
|
return;
|
|
};
|
|
let issuer = task_issuer.clone();
|
|
tokio::spawn(async move {
|
|
let mut request = vec![0; 8192];
|
|
let Ok(read) = stream.read(&mut request).await else {
|
|
return;
|
|
};
|
|
let request = String::from_utf8_lossy(&request[..read]);
|
|
let (status, body) =
|
|
if request.starts_with("GET /.well-known/openid-configuration ") {
|
|
(
|
|
"200 OK",
|
|
serde_json::json!({
|
|
"issuer": issuer,
|
|
"authorization_endpoint": format!("{issuer}/authorize"),
|
|
"token_endpoint": format!("{issuer}/token"),
|
|
})
|
|
.to_string(),
|
|
)
|
|
} else if request.starts_with("POST /token ") {
|
|
(
|
|
"400 Bad Request",
|
|
serde_json::json!({
|
|
"error": "invalid_grant",
|
|
"error_description": "Token is not active",
|
|
})
|
|
.to_string(),
|
|
)
|
|
} else {
|
|
("404 Not Found", "{}".to_string())
|
|
};
|
|
let response = format!(
|
|
"HTTP/1.1 {status}\r\nContent-Type: application/json\r\nContent-Length: {}\r\nConnection: close\r\n\r\n{body}",
|
|
body.len(),
|
|
);
|
|
let _ = stream.write_all(response.as_bytes()).await;
|
|
});
|
|
}
|
|
});
|
|
|
|
issuer
|
|
}
|
|
|
|
fn write_oidc_test_credentials(
|
|
server: &TestServer,
|
|
xdg_dir: &TempDir,
|
|
issuer: &str,
|
|
expires_at: u64,
|
|
) {
|
|
let gateway_dir = xdg_dir.path().join("openshell/gateways/openshell");
|
|
fs::write(
|
|
gateway_dir.join("metadata.json"),
|
|
serde_json::to_vec_pretty(&openshell_bootstrap::GatewayMetadata {
|
|
name: "openshell".to_string(),
|
|
gateway_endpoint: server.endpoint.clone(),
|
|
auth_mode: Some("oidc".to_string()),
|
|
oidc_issuer: Some(issuer.to_string()),
|
|
oidc_client_id: Some("openshell-cli".to_string()),
|
|
..Default::default()
|
|
})
|
|
.unwrap(),
|
|
)
|
|
.unwrap();
|
|
fs::write(
|
|
gateway_dir.join("oidc_token.json"),
|
|
serde_json::to_vec_pretty(&openshell_bootstrap::oidc_token::OidcTokenBundle {
|
|
access_token: "cached-access-token".to_string(),
|
|
refresh_token: Some("inactive-refresh-token".to_string()),
|
|
expires_at: Some(expires_at),
|
|
issuer: issuer.to_string(),
|
|
client_id: "openshell-cli".to_string(),
|
|
})
|
|
.unwrap(),
|
|
)
|
|
.unwrap();
|
|
}
|
|
|
|
/// Models the gateway's JWT-expiration leeway by using a test gateway that
|
|
/// still accepts the stale bearer. Before the CLI failed closed, the rejected
|
|
/// refresh was only logged and `CreateSandbox` reached this gateway anyway.
|
|
#[tokio::test]
|
|
async fn sandbox_create_fails_before_mutation_when_expired_oidc_refresh_is_rejected() {
|
|
let server = run_server().await;
|
|
let issuer = run_rejected_oidc_refresh_server().await;
|
|
let xdg_dir = tempfile::tempdir().unwrap();
|
|
prepare_cli_xdg(&server, &xdg_dir);
|
|
|
|
write_oidc_test_credentials(&server, &xdg_dir, &issuer, 0);
|
|
|
|
let result = run_cli_sandbox_create_with_xdg(
|
|
&server,
|
|
&xdg_dir,
|
|
"must-not-be-created",
|
|
&["--output=json"],
|
|
)
|
|
.await;
|
|
|
|
assert!(
|
|
!result.status.success(),
|
|
"refresh failure must abort create"
|
|
);
|
|
let stderr = String::from_utf8_lossy(&result.stderr);
|
|
assert!(
|
|
stderr.contains("invalid_grant"),
|
|
"unexpected error: {stderr}"
|
|
);
|
|
assert!(
|
|
stderr.contains("openshell gateway login openshell"),
|
|
"error should explain how to re-authenticate: {stderr}"
|
|
);
|
|
assert!(
|
|
create_requests(&server).await.is_empty(),
|
|
"CreateSandbox must not be called after OIDC refresh fails"
|
|
);
|
|
}
|
|
|
|
#[tokio::test]
|
|
async fn sandbox_create_continues_with_unexpired_cached_token_when_refresh_fails() {
|
|
let server = run_server().await;
|
|
let issuer = run_rejected_oidc_refresh_server().await;
|
|
let xdg_dir = tempfile::tempdir().unwrap();
|
|
prepare_cli_xdg(&server, &xdg_dir);
|
|
let expires_at = std::time::SystemTime::now()
|
|
.duration_since(std::time::UNIX_EPOCH)
|
|
.unwrap()
|
|
.as_secs()
|
|
+ 20;
|
|
write_oidc_test_credentials(&server, &xdg_dir, &issuer, expires_at);
|
|
|
|
let result =
|
|
run_cli_sandbox_create_with_xdg(&server, &xdg_dir, "uses-cached-token", &["--output=json"])
|
|
.await;
|
|
|
|
assert!(
|
|
result.status.success(),
|
|
"sandbox create should continue with an unexpired cached token:\n{}",
|
|
String::from_utf8_lossy(&result.stderr)
|
|
);
|
|
let stderr = String::from_utf8_lossy(&result.stderr);
|
|
assert!(
|
|
stderr.contains("continuing with the cached token"),
|
|
"refresh fallback warning should be visible: {stderr}"
|
|
);
|
|
assert_eq!(
|
|
create_requests(&server).await.len(),
|
|
1,
|
|
"CreateSandbox should be reached with the cached token"
|
|
);
|
|
}
|
|
|
|
#[tokio::test]
|
|
async fn sandbox_create_json_stdout_is_parseable() {
|
|
let server = run_server().await;
|
|
|
|
let result = run_cli_sandbox_create(&server, "json-clean", &["--output=json"]).await;
|
|
assert!(
|
|
result.status.success(),
|
|
"sandbox create failed:\n{}",
|
|
String::from_utf8_lossy(&result.stderr)
|
|
);
|
|
let stdout = String::from_utf8(result.stdout).expect("stdout should be UTF-8");
|
|
serde_json::from_str::<serde_json::Value>(&stdout)
|
|
.unwrap_or_else(|err| panic!("stdout should contain only JSON: {err}\n{stdout}"));
|
|
}
|
|
|
|
#[tokio::test]
|
|
async fn sandbox_create_yaml_stdout_is_parseable() {
|
|
let server = run_server().await;
|
|
|
|
let result = run_cli_sandbox_create(&server, "yaml-clean", &["--output=yaml"]).await;
|
|
assert!(
|
|
result.status.success(),
|
|
"sandbox create failed:\n{}",
|
|
String::from_utf8_lossy(&result.stderr)
|
|
);
|
|
let stdout = String::from_utf8(result.stdout).expect("stdout should be UTF-8");
|
|
serde_yml::from_str::<serde_yml::Value>(&stdout)
|
|
.unwrap_or_else(|err| panic!("stdout should contain only YAML: {err}\n{stdout}"));
|
|
}
|
|
|
|
#[tokio::test]
|
|
async fn sandbox_template_create_warns_for_credential_env_vars() {
|
|
let server = run_server().await;
|
|
|
|
let result = run_cli_sandbox_template_create(
|
|
&server,
|
|
"credential-env",
|
|
&["--env", "OPENAI_API_KEY=plain-secret"],
|
|
)
|
|
.await;
|
|
|
|
assert!(
|
|
result.status.success(),
|
|
"sandbox template create failed:\n{}",
|
|
String::from_utf8_lossy(&result.stderr)
|
|
);
|
|
let stderr = String::from_utf8_lossy(&result.stderr);
|
|
assert!(
|
|
stderr.contains("OPENAI_API_KEY looks like a credential"),
|
|
"template create should warn for credential-looking --env values: {stderr}"
|
|
);
|
|
assert!(
|
|
stderr.contains("To hide it from the agent, use a provider instead"),
|
|
"warning should point users toward providers: {stderr}"
|
|
);
|
|
|
|
let requests = template_create_requests(&server).await;
|
|
assert_eq!(requests.len(), 1);
|
|
}
|
|
|
|
#[tokio::test]
|
|
async fn sandbox_template_create_suppresses_credential_env_warnings() {
|
|
let server = run_server().await;
|
|
|
|
let result = run_cli_sandbox_template_create(
|
|
&server,
|
|
"credential-env-suppressed",
|
|
&[
|
|
"--env",
|
|
"OPENAI_API_KEY=plain-secret",
|
|
"--no-credential-warnings",
|
|
],
|
|
)
|
|
.await;
|
|
|
|
assert!(
|
|
result.status.success(),
|
|
"sandbox template create failed:\n{}",
|
|
String::from_utf8_lossy(&result.stderr)
|
|
);
|
|
let stderr = String::from_utf8_lossy(&result.stderr);
|
|
assert!(
|
|
!stderr.contains("OPENAI_API_KEY looks like a credential"),
|
|
"template create should suppress credential-looking --env warnings: {stderr}"
|
|
);
|
|
|
|
let requests = template_create_requests(&server).await;
|
|
assert_eq!(requests.len(), 1);
|
|
}
|