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

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion .github/workflows/e2e-docker-test.yml
Original file line number Diff line number Diff line change
Expand Up @@ -33,7 +33,7 @@ on:
required: false
type: string
default: >-
[{"suite":"python","cmd":"mise run --no-deps --skip-deps e2e:python","apt_packages":"","python_proto":true,"mcp":false},{"suite":"oidc-python","cmd":"mise run --no-deps --skip-deps e2e:oidc-python:docker","apt_packages":"","python_proto":true,"mcp":false},{"suite":"oidc-pkce-docker","cmd":"mise run --no-deps --skip-deps e2e:oidc-pkce:docker","apt_packages":"openssh-client","python_proto":false,"mcp":false},{"suite":"rust-docker","cmd":"mise run --no-deps --skip-deps e2e:rust","apt_packages":"openssh-client","python_proto":false,"mcp":false},{"suite":"mcp","cmd":"mise run --no-deps --skip-deps e2e:mcp","apt_packages":"","python_proto":false,"mcp":true}]
[{"suite":"python","cmd":"mise run --no-deps --skip-deps e2e:python","apt_packages":"","python_proto":true,"mcp":false},{"suite":"oidc-python","cmd":"mise run --no-deps --skip-deps e2e:oidc-python:docker","apt_packages":"","python_proto":true,"mcp":false},{"suite":"oidc-pkce-docker","cmd":"mise run --no-deps --skip-deps e2e:oidc-pkce:docker","apt_packages":"openssh-client","python_proto":false,"mcp":false},{"suite":"rust-docker","cmd":"mise run --no-deps --skip-deps e2e:rust","apt_packages":"openssh-client","python_proto":false,"mcp":false},{"suite":"rust-docker-push","cmd":"mise run --no-deps --skip-deps e2e:rust:push","apt_packages":"openssh-client","python_proto":false,"mcp":false},{"suite":"mcp","cmd":"mise run --no-deps --skip-deps e2e:mcp","apt_packages":"","python_proto":false,"mcp":true}]

permissions:
actions: read
Expand Down
1 change: 1 addition & 0 deletions TESTING.md
Original file line number Diff line number Diff line change
Expand Up @@ -223,6 +223,7 @@ Suites:
the selected deployment.
- Docker suite (`--features e2e-docker`) - includes Docker-only coverage such as Dockerfile image builds, Docker preflight checks, and managed Docker gateway start.
- Docker GPU suite (`--features e2e-docker-gpu`) - Docker suite plus GPU sandbox smoke coverage.
- Configuration push suite (`--features e2e-config-push`) - CLI conformance and the policy tests against a Docker gateway with `config_delivery_mode = "push"`. Run it with `mise run e2e:rust:push`.
- VM suite (`--features e2e-vm`) - runs e2e tests on a VM.
- Kubernetes credential-driver suite (`--features e2e-kubernetes-credential-drivers`) - targeted Kubernetes Secrets and Vault provider credential storage coverage.

Expand Down
160 changes: 90 additions & 70 deletions crates/openshell-cli/src/commands/provider.rs
Original file line number Diff line number Diff line change
Expand Up @@ -129,6 +129,10 @@ pub async fn sandbox_provider_list(
Ok(())
}

/// Attempts for an attach or detach whose resource version went stale between
/// the command's read and its write.
const PROVIDER_ATTACHMENT_CAS_ATTEMPTS: u32 = 5;

/// Save a provider attachment and optionally wait for its installed authority.
pub async fn sandbox_provider_attach(
server: &str,
Expand All @@ -141,45 +145,53 @@ pub async fn sandbox_provider_attach(
readiness.validate()?;
let mut client = grpc_client(server, tls).await?;

// Fetch current sandbox to get resource_version for CAS
let sandbox = client
.get_sandbox(GetSandboxRequest {
name: (name).to_string(),
workspace_scope: Some(openshell_core::proto::workspace_selector(
(workspace).to_string(),
)),
})
.await
.map_err(|status| miette!("provider attachment lookup failed ({})", status.code()))?
.into_inner()
.sandbox
.ok_or_else(|| miette::miette!("sandbox not found"))?;
// The resource version only guards this command's own read, so a
// conflict from a concurrent status write is retried with a fresh read.
let mut attempt = 1;
let (sandbox, response) = loop {
let sandbox = client
.get_sandbox(GetSandboxRequest {
name: (name).to_string(),
workspace_scope: Some(openshell_core::proto::workspace_selector(
(workspace).to_string(),
)),
})
.await
.map_err(|status| miette!("provider attachment lookup failed ({})", status.code()))?
.into_inner()
.sandbox
.ok_or_else(|| miette::miette!("sandbox not found"))?;

let resource_version = sandbox.metadata.as_ref().map_or(0, |m| m.resource_version);
let resource_version = sandbox.metadata.as_ref().map_or(0, |m| m.resource_version);

let response = match client
.attach_sandbox_provider(AttachSandboxProviderRequest {
request_id: String::new(),
sandbox: (name).to_string(),
workspace_scope: Some(openshell_core::proto::workspace_selector(
(workspace).to_string(),
)),
provider: provider.to_string(),
expected_resource_version: resource_version,
})
.await
{
Ok(response) => response.into_inner(),
// Explicit post-save uncertainty takes precedence over a generic retry hint.
Err(status)
if status.code() == Code::Aborted && !provider_mutation_is_uncertain(&status) =>
match client
.attach_sandbox_provider(AttachSandboxProviderRequest {
request_id: String::new(),
sandbox: (name).to_string(),
workspace_scope: Some(openshell_core::proto::workspace_selector(
(workspace).to_string(),
)),
provider: provider.to_string(),
expected_resource_version: resource_version,
})
.await
{
return Err(miette::miette!(
"Failed to attach provider: sandbox was modified by another operation.\n\
Please retry the command."
));
Ok(response) => break (sandbox, response.into_inner()),
// Explicit post-save uncertainty takes precedence over a generic retry hint.
Err(status)
if status.code() == Code::Aborted && !provider_mutation_is_uncertain(&status) =>
{
if attempt < PROVIDER_ATTACHMENT_CAS_ATTEMPTS {
attempt += 1;
continue;
}
return Err(miette::miette!(
"Failed to attach provider: sandbox was modified by another operation.\n\
Please retry the command."
));
}
Err(error) => return Err(provider_mutation_error(&error, "attachment")),
}
Err(error) => return Err(provider_mutation_error(&error, "attachment")),
};

let receipt = response.receipt.ok_or_else(|| {
Expand Down Expand Up @@ -216,45 +228,53 @@ pub async fn sandbox_provider_detach(
readiness.validate()?;
let mut client = grpc_client(server, tls).await?;

// Fetch current sandbox to get resource_version for CAS
let sandbox = client
.get_sandbox(GetSandboxRequest {
name: (name).to_string(),
workspace_scope: Some(openshell_core::proto::workspace_selector(
(workspace).to_string(),
)),
})
.await
.map_err(|status| miette!("provider detachment lookup failed ({})", status.code()))?
.into_inner()
.sandbox
.ok_or_else(|| miette::miette!("sandbox not found"))?;
// The resource version only guards this command's own read, so a
// conflict from a concurrent status write is retried with a fresh read.
let mut attempt = 1;
let (sandbox, response) = loop {
let sandbox = client
.get_sandbox(GetSandboxRequest {
name: (name).to_string(),
workspace_scope: Some(openshell_core::proto::workspace_selector(
(workspace).to_string(),
)),
})
.await
.map_err(|status| miette!("provider detachment lookup failed ({})", status.code()))?
.into_inner()
.sandbox
.ok_or_else(|| miette::miette!("sandbox not found"))?;

let resource_version = sandbox.metadata.as_ref().map_or(0, |m| m.resource_version);
let resource_version = sandbox.metadata.as_ref().map_or(0, |m| m.resource_version);

let response = match client
.detach_sandbox_provider(DetachSandboxProviderRequest {
request_id: String::new(),
sandbox: (name).to_string(),
workspace_scope: Some(openshell_core::proto::workspace_selector(
(workspace).to_string(),
)),
provider: provider.to_string(),
expected_resource_version: resource_version,
})
.await
{
Ok(response) => response.into_inner(),
// Explicit post-save uncertainty takes precedence over a generic retry hint.
Err(status)
if status.code() == Code::Aborted && !provider_mutation_is_uncertain(&status) =>
match client
.detach_sandbox_provider(DetachSandboxProviderRequest {
request_id: String::new(),
sandbox: (name).to_string(),
workspace_scope: Some(openshell_core::proto::workspace_selector(
(workspace).to_string(),
)),
provider: provider.to_string(),
expected_resource_version: resource_version,
})
.await
{
return Err(miette::miette!(
"Failed to detach provider: sandbox was modified by another operation.\n\
Please retry the command."
));
Ok(response) => break (sandbox, response.into_inner()),
// Explicit post-save uncertainty takes precedence over a generic retry hint.
Err(status)
if status.code() == Code::Aborted && !provider_mutation_is_uncertain(&status) =>
{
if attempt < PROVIDER_ATTACHMENT_CAS_ATTEMPTS {
attempt += 1;
continue;
}
return Err(miette::miette!(
"Failed to detach provider: sandbox was modified by another operation.\n\
Please retry the command."
));
}
Err(error) => return Err(provider_mutation_error(&error, "detachment")),
}
Err(error) => return Err(provider_mutation_error(&error, "detachment")),
};

let receipt = response.receipt.ok_or_else(|| miette!("gateway did not return a provider receipt; saved detachment cannot establish revocation"))?;
Expand Down
7 changes: 7 additions & 0 deletions crates/openshell-cli/tests/ensure_providers_integration.rs
Original file line number Diff line number Diff line change
Expand Up @@ -79,6 +79,13 @@ impl TestOpenShell {

#[tonic::async_trait]
impl OpenShell for TestOpenShell {
async fn peer_notify_config_update(
&self,
_request: tonic::Request<openshell_core::proto::PeerNotifyConfigUpdateRequest>,
) -> Result<Response<openshell_core::proto::PeerNotifyConfigUpdateResponse>, Status> {
Err(Status::unimplemented("not used by this test server"))
}

async fn peer_report_provider_readiness(
&self,
_request: tonic::Request<openshell_core::proto::ReportProviderReadinessRequest>,
Expand Down
7 changes: 7 additions & 0 deletions crates/openshell-cli/tests/mtls_integration.rs
Original file line number Diff line number Diff line change
Expand Up @@ -32,6 +32,13 @@ struct TestOpenShell;

#[tonic::async_trait]
impl OpenShell for TestOpenShell {
async fn peer_notify_config_update(
&self,
_request: tonic::Request<openshell_core::proto::PeerNotifyConfigUpdateRequest>,
) -> Result<Response<openshell_core::proto::PeerNotifyConfigUpdateResponse>, Status> {
Err(Status::unimplemented("not used by this test server"))
}

async fn peer_report_provider_readiness(
&self,
_request: tonic::Request<openshell_core::proto::ReportProviderReadinessRequest>,
Expand Down
125 changes: 125 additions & 0 deletions crates/openshell-cli/tests/provider_commands_integration.rs
Original file line number Diff line number Diff line change
Expand Up @@ -110,6 +110,8 @@ struct ProviderState {
readiness_sequence: Arc<AtomicU64>,
corrupt_mutation_receipt: Arc<Mutex<Option<ReceiptCorruption>>>,
fail_mutation_after_save: Arc<Mutex<Option<Status>>>,
/// Plain resource-version conflicts to return before an attachment saves.
sandbox_provider_conflicts: Arc<AtomicU64>,
global_settings: Arc<Mutex<HashMap<String, SettingValue>>>,
}

Expand Down Expand Up @@ -232,6 +234,13 @@ impl TestOpenShell {

#[tonic::async_trait]
impl OpenShell for TestOpenShell {
async fn peer_notify_config_update(
&self,
_request: tonic::Request<openshell_core::proto::PeerNotifyConfigUpdateRequest>,
) -> Result<Response<openshell_core::proto::PeerNotifyConfigUpdateResponse>, Status> {
Err(Status::unimplemented("not used by this test server"))
}

async fn peer_report_provider_readiness(
&self,
_request: tonic::Request<openshell_core::proto::ReportProviderReadinessRequest>,
Expand Down Expand Up @@ -421,6 +430,18 @@ impl OpenShell for TestOpenShell {
sandbox_name: sandbox_name.clone(),
provider: request.provider.clone(),
});
if self
.state
.sandbox_provider_conflicts
.fetch_update(Ordering::SeqCst, Ordering::SeqCst, |left| {
left.checked_sub(1)
})
.is_ok()
{
return Err(Status::aborted(
"sandbox was modified concurrently (current resource_version: 2)",
));
}
if !self
.state
.providers
Expand Down Expand Up @@ -483,6 +504,18 @@ impl OpenShell for TestOpenShell {
sandbox_name: sandbox_name.clone(),
provider: request.provider.clone(),
});
if self
.state
.sandbox_provider_conflicts
.fetch_update(Ordering::SeqCst, Ordering::SeqCst, |left| {
left.checked_sub(1)
})
.is_ok()
{
return Err(Status::aborted(
"sandbox was modified concurrently (current resource_version: 2)",
));
}
let mut sandbox_providers = self.state.sandbox_providers.lock().await;
let providers = sandbox_providers.entry(sandbox_name.clone()).or_default();
let before_len = providers.len();
Expand Down Expand Up @@ -6028,3 +6061,95 @@ async fn provider_create_from_gcloud_adc_missing_client_secret() {
"no provider must be created when ADC validation fails"
);
}

#[tokio::test]
async fn provider_attachment_retries_resource_version_conflicts() {
let server = run_server().await;
seed_readiness_provider(&server).await;
for action in ["attach", "detach"] {
let sandbox_name = "conflicted";
server.state.sandbox_providers.lock().await.insert(
sandbox_name.to_string(),
if action == "attach" {
Vec::new()
} else {
vec![READINESS_PROVIDER.to_string()]
},
);
server.state.sandbox_provider_requests.lock().await.clear();
server
.state
.sandbox_provider_conflicts
.store(2, Ordering::SeqCst);
let output = run_readiness_cli(
&server,
&[
"sandbox",
"provider",
action,
sandbox_name,
READINESS_PROVIDER,
],
)
.await;
assert!(
output.status.success(),
"{action}: {}",
String::from_utf8_lossy(&output.stderr)
);
assert_eq!(
server.state.sandbox_provider_requests.lock().await.len(),
3,
"{action} retried each conflict once"
);
assert_eq!(
server
.state
.sandbox_providers
.lock()
.await
.get(sandbox_name)
.unwrap()
.contains(&READINESS_PROVIDER.to_string()),
action == "attach"
);
}
}

#[tokio::test]
async fn provider_attachment_reports_persistent_conflicts() {
let server = run_server().await;
seed_readiness_provider(&server).await;
server
.state
.sandbox_providers
.lock()
.await
.insert("contended".to_string(), Vec::new());
server
.state
.sandbox_provider_conflicts
.store(u64::MAX, Ordering::SeqCst);
let output = run_readiness_cli(
&server,
&[
"sandbox",
"provider",
"attach",
"contended",
READINESS_PROVIDER,
],
)
.await;
assert!(!output.status.success());
assert!(
String::from_utf8_lossy(&output.stderr).contains("Please retry the command"),
"{}",
String::from_utf8_lossy(&output.stderr)
);
assert_eq!(server.state.sandbox_provider_requests.lock().await.len(), 5);
server
.state
.sandbox_provider_conflicts
.store(0, Ordering::SeqCst);
}
Loading
Loading