From c8395fbfb6017f2171eee792b1e39f2488baecc5 Mon Sep 17 00:00:00 2001 From: Kevin Wang Date: Mon, 27 Jul 2026 18:54:30 +0000 Subject: [PATCH 1/4] fix(kms): return Finish response before exit Defect: Onboard.Finish terminated the process inside the RPC handler before Rocket could flush the documented Empty response. Observed symptom: clients received RemoteDisconnected and the graceful-response regression failed. Fix: return success immediately and schedule process exit after a short response-drain interval. --- dstack/kms/src/onboard_service.rs | 11 +++++++++-- 1 file changed, 9 insertions(+), 2 deletions(-) diff --git a/dstack/kms/src/onboard_service.rs b/dstack/kms/src/onboard_service.rs index 334a6c6db..2160272e2 100644 --- a/dstack/kms/src/onboard_service.rs +++ b/dstack/kms/src/onboard_service.rs @@ -2,7 +2,10 @@ // // SPDX-License-Identifier: Apache-2.0 -use std::sync::{Arc, Mutex}; +use std::{ + sync::{Arc, Mutex}, + time::Duration, +}; use anyhow::{bail, Context, Result}; use dstack_kms_rpc::{ @@ -199,7 +202,11 @@ impl OnboardRpc for OnboardHandler { } async fn finish(self) -> anyhow::Result<()> { - std::process::exit(0); + tokio::spawn(async { + tokio::time::sleep(Duration::from_millis(250)).await; + std::process::exit(0); + }); + Ok(()) } } From 026a1d870ce7775b419bcd0c15ca9f7552f23f67 Mon Sep 17 00:00:00 2001 From: Kevin Wang Date: Mon, 27 Jul 2026 19:03:12 +0000 Subject: [PATCH 2/4] fix(rpc): emit empty JSON unit responses --- dstack/ra-rpc/src/rocket_helper.rs | 35 +++++++++++++++++++++++++++++- 1 file changed, 34 insertions(+), 1 deletion(-) diff --git a/dstack/ra-rpc/src/rocket_helper.rs b/dstack/ra-rpc/src/rocket_helper.rs index 2cc170630..939c38239 100644 --- a/dstack/ra-rpc/src/rocket_helper.rs +++ b/dstack/ra-rpc/src/rocket_helper.rs @@ -42,6 +42,14 @@ pub struct RpcResponse { body: Vec, } +fn normalize_json_response_body(is_json: bool, body: Vec) -> Vec { + if is_json && body.as_slice() == b"null" { + Vec::new() + } else { + body + } +} + impl<'r> Responder<'r, 'static> for RpcResponse { fn respond_to(self, request: &'r Request<'_>) -> rocket::response::Result<'static> { use rocket::http::ContentType; @@ -50,13 +58,38 @@ impl<'r> Responder<'r, 'static> for RpcResponse { } else { ContentType::Binary }; - let response = Custom(self.status, self.body).respond_to(request)?; + // prpc maps google.protobuf.Empty / Rust unit to JSON `null`. Case + // contracts and many clients expect an empty success body instead. + let body = normalize_json_response_body(self.is_json, self.body); + let response = Custom(self.status, body).respond_to(request)?; rocket::Response::build_from(response) .header(content_type) .ok() } } +#[cfg(test)] +mod response_tests { + use super::normalize_json_response_body; + + #[test] + fn json_unit_response_has_an_empty_body() { + assert!(normalize_json_response_body(true, b"null".to_vec()).is_empty()); + } + + #[test] + fn non_unit_and_binary_responses_are_unchanged() { + assert_eq!( + normalize_json_response_body(true, br#"{"value":null}"#.to_vec()), + br#"{"value":null}"# + ); + assert_eq!( + normalize_json_response_body(false, b"null".to_vec()), + b"null" + ); + } +} + #[derive(Debug, Clone)] struct UnixPeerEndpoint { path: PathBuf, From d92c2d8267734c4954ad2573d9216be6d75ba74d Mon Sep 17 00:00:00 2001 From: Kevin Wang Date: Mon, 27 Jul 2026 19:04:29 +0000 Subject: [PATCH 3/4] fix(kms): shut down after Finish response --- dstack/kms/src/main.rs | 9 +++++++-- dstack/kms/src/onboard_service.rs | 23 ++++++++++++++++++----- 2 files changed, 25 insertions(+), 7 deletions(-) diff --git a/dstack/kms/src/main.rs b/dstack/kms/src/main.rs index 478c5c52a..74921a698 100644 --- a/dstack/kms/src/main.rs +++ b/dstack/kms/src/main.rs @@ -63,13 +63,18 @@ async fn run_onboard_service(kms_config: KmsConfig, figment: Figment) -> Result< // Remove section tls - let _ = rocket::custom(figment) + let rocket = rocket::custom(figment) .mount("/", rocket::routes![index, finish]) .mount( "/prpc", ra_rpc::prpc_routes!(OnboardState, OnboardHandler, trim: "Onboard."), ) - .manage(state) + .manage(state.clone()) + .ignite() + .await + .map_err(|err| anyhow!(err.to_string()))?; + state.set_shutdown(rocket.shutdown())?; + let _ = rocket .launch() .await .map_err(|err| anyhow!(err.to_string()))?; diff --git a/dstack/kms/src/onboard_service.rs b/dstack/kms/src/onboard_service.rs index 2160272e2..cae39e1d5 100644 --- a/dstack/kms/src/onboard_service.rs +++ b/dstack/kms/src/onboard_service.rs @@ -4,7 +4,6 @@ use std::{ sync::{Arc, Mutex}, - time::Duration, }; use anyhow::{bail, Context, Result}; @@ -48,6 +47,7 @@ pub struct OnboardState { config: KmsConfig, attestation_verifier: Arc, bootstrap_lock: Arc>, + shutdown: Arc>>, } impl OnboardState { @@ -60,8 +60,17 @@ impl OnboardState { config, attestation_verifier, bootstrap_lock: Arc::new(AsyncMutex::new(())), + shutdown: Arc::new(Mutex::new(None)), }) } + + pub fn set_shutdown(&self, shutdown: rocket::Shutdown) -> Result<()> { + *self + .shutdown + .lock() + .map_err(|_| anyhow::anyhow!("onboard shutdown lock poisoned"))? = Some(shutdown); + Ok(()) + } } pub struct OnboardHandler { @@ -202,10 +211,14 @@ impl OnboardRpc for OnboardHandler { } async fn finish(self) -> anyhow::Result<()> { - tokio::spawn(async { - tokio::time::sleep(Duration::from_millis(250)).await; - std::process::exit(0); - }); + let shutdown = self + .state + .shutdown + .lock() + .map_err(|_| anyhow::anyhow!("onboard shutdown lock poisoned"))? + .clone() + .context("onboard shutdown handle is unavailable")?; + shutdown.notify(); Ok(()) } } From dd852f29e1f2fa177e73386eaf09d1583905dea4 Mon Sep 17 00:00:00 2001 From: Kevin Wang Date: Wed, 29 Jul 2026 05:40:25 +0000 Subject: [PATCH 4/4] fix(kms): reject repeated onboarding --- dstack/kms/src/onboard_service.rs | 10 +++++++--- 1 file changed, 7 insertions(+), 3 deletions(-) diff --git a/dstack/kms/src/onboard_service.rs b/dstack/kms/src/onboard_service.rs index cae39e1d5..9a2e48d15 100644 --- a/dstack/kms/src/onboard_service.rs +++ b/dstack/kms/src/onboard_service.rs @@ -138,6 +138,11 @@ impl OnboardRpc for OnboardHandler { async fn onboard(self, request: OnboardRequest) -> Result { validate_onboarding_domain(&request.domain)?; + let _bootstrap_guard = self.state.bootstrap_lock.lock().await; + let cfg = &self.state.config; + if cfg.root_ca_key().exists() || cfg.k256_key().exists() { + bail!("KMS has already been onboarded"); + } let source_url = request.source_url.trim_end_matches('/').to_string(); let source_url = if source_url.ends_with("/prpc") { source_url @@ -145,7 +150,7 @@ impl OnboardRpc for OnboardHandler { format!("{source_url}/prpc") }; let keys = Keys::onboard( - &self.state.config, + cfg, &source_url, &request.domain, self.state.attestation_verifier.clone(), @@ -153,8 +158,7 @@ impl OnboardRpc for OnboardHandler { .await .context("Failed to onboard")?; let k256_pubkey = keys.k256_key.verifying_key().to_sec1_bytes().to_vec(); - keys.store(&self.state.config) - .context("Failed to store keys")?; + keys.store(cfg).context("Failed to store keys")?; Ok(OnboardResponse { k256_pubkey }) }