use std::io::{IsTerminal, Write};
use std::path::{Path, PathBuf};
use std::process::Stdio;
use std::time::Duration;
use colored::Colorize;
use dialoguer::{theme::ColorfulTheme, Confirm, Select};
use serde_json::Value;
use tokio::process::Command;
use tokio::time::sleep;
use crate::utils::command_exists;
const KIND_NODE_IMAGE_DEFAULT: &str = "kindest/node:v1.36.1";
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum DoctorOutcome {
Ready,
Skipped,
StillDown,
}
#[derive(Debug, Clone)]
pub struct DoctorOptions {
pub yes: bool,
pub debug: bool,
pub prefer_kubeadm: bool,
}
impl Default for DoctorOptions {
fn default() -> Self {
Self {
yes: false,
debug: false,
prefer_kubeadm: true,
}
}
}
pub fn looks_like_docker_desktop_k8s_issue(text: &str) -> bool {
let lower = text.to_ascii_lowercase();
lower.contains("docker-desktop")
|| lower.contains("kubernetes.docker.internal")
|| lower.contains("docker desktop kubernetes")
|| lower.contains("kindest/node")
|| lower.contains("desktop-containerd")
|| (lower.contains("127.0.0.1:6443") && lower.contains("refused"))
|| (lower.contains("localhost:6443") && lower.contains("refused"))
|| (lower.contains("cluster not ready")
&& (lower.contains("6443") || lower.contains("docker")))
}
pub fn is_interactive_tty() -> bool {
std::io::stdin().is_terminal() && std::io::stdout().is_terminal()
}
pub async fn maybe_recover_docker_desktop_k8s(
failure_text: &str,
opts: DoctorOptions,
) -> Result<DoctorOutcome, String> {
if !looks_like_docker_desktop_k8s_issue(failure_text)
&& !looks_like_docker_desktop_k8s_issue(¤t_context_hint().await)
{
if !opts.yes {
return Ok(DoctorOutcome::Skipped);
}
}
if !command_exists("docker") {
println!(
"{} Docker CLI not found; cannot automate Docker Desktop recovery.",
"k8s doctor:".yellow().bold()
);
return Ok(DoctorOutcome::Skipped);
}
if !docker_desktop_cli_available().await {
println!(
"{} `docker desktop` CLI unavailable. Open Docker Desktop → Settings → Kubernetes and set Cluster to kubeadm, or install a Docker Desktop version that ships the desktop CLI.",
"k8s doctor:".yellow().bold()
);
print_manual_purge_steps();
return Ok(DoctorOutcome::Skipped);
}
run_doctor(opts).await
}
pub async fn run_doctor(opts: DoctorOptions) -> Result<DoctorOutcome, String> {
println!(
"{} Docker Desktop Kubernetes recovery",
"k8s doctor".bright_cyan().bold()
);
print_disk_hints().await;
let status = kubernetes_status(opts.debug).await;
match &status {
Ok(s) => {
println!(
" {} state={} mode={} version={}",
"status".bright_black(),
s.state.bright_white(),
s.mode.as_deref().unwrap_or("?").bright_yellow(),
s.version.as_deref().unwrap_or("?").bright_black()
);
if s.is_running() {
if kubectl_cluster_ready(opts.debug).await {
println!(
"{} cluster already reachable (`kubectl cluster-info` ok)",
"k8s doctor:".green().bold()
);
print_namespace_snapshot(opts.debug).await;
return Ok(DoctorOutcome::Ready);
}
println!(
" {} Docker reports running but kubectl cannot reach the API — will repair.",
"note".yellow()
);
} else if let Some(err) = &s.error {
if !err.is_empty() && err != "None" {
println!(" {} {}", "error".red(), err);
}
}
if let Some(msg) = &s.progress {
if !msg.is_empty() {
println!(" {} {}", "progress".bright_black(), msg);
}
}
}
Err(e) => {
println!(" {} could not read kubernetes status: {e}", "status".yellow());
}
}
let settings_mode = read_kubernetes_mode();
if let Some(mode) = &settings_mode {
println!(
" {} settings-store KubernetesMode={}",
"config".bright_black(),
mode.bright_yellow()
);
}
let action = if opts.yes {
if opts.prefer_kubeadm {
DoctorAction::SwitchKubeadm
} else {
DoctorAction::ResetCluster
}
} else if is_interactive_tty() {
prompt_action(settings_mode.as_deref())?
} else {
println!(
"{} non-interactive shell: pass --yes to auto-switch to kubeadm, or run `xbp kubernetes doctor` in a TTY.",
"k8s doctor:".yellow().bold()
);
return Ok(DoctorOutcome::Skipped);
};
match action {
DoctorAction::Cancel => {
println!("{} cancelled", "k8s doctor".dimmed());
return Ok(DoctorOutcome::Skipped);
}
DoctorAction::ManualPurge => {
print_manual_purge_steps();
return Ok(DoctorOutcome::Skipped);
}
DoctorAction::Recheck => {
return finalize_wait(opts.debug).await;
}
DoctorAction::SwitchKubeadm => {
switch_to_kubeadm(opts.debug).await?;
}
DoctorAction::KindRepullReset => {
kind_repull_and_reset(opts.debug).await?;
}
DoctorAction::ResetCluster => {
reset_cluster(opts.debug).await?;
}
}
finalize_wait(opts.debug).await
}
#[derive(Debug, Clone, Copy)]
enum DoctorAction {
SwitchKubeadm,
KindRepullReset,
ResetCluster,
Recheck,
ManualPurge,
Cancel,
}
fn prompt_action(current_mode: Option<&str>) -> Result<DoctorAction, String> {
let mode_note = current_mode.unwrap_or("unknown");
let items = [
format!(
"Switch to kubeadm + enable Kubernetes + restart Docker Desktop (recommended; current mode={mode_note})"
),
"Stay on kind: repull kindest/node + reset Kubernetes cluster".into(),
"Reset Kubernetes cluster only (keep current mode)".into(),
"Wait / recheck status only".into(),
"Show Clean / Purge data steps (manual; fixes desktop-containerd)".into(),
"Cancel".into(),
];
let idx = Select::with_theme(&ColorfulTheme::default())
.with_prompt("Docker Desktop Kubernetes recovery")
.items(&items)
.default(0)
.interact()
.map_err(|e| e.to_string())?;
Ok(match idx {
0 => DoctorAction::SwitchKubeadm,
1 => DoctorAction::KindRepullReset,
2 => DoctorAction::ResetCluster,
3 => DoctorAction::Recheck,
4 => DoctorAction::ManualPurge,
_ => DoctorAction::Cancel,
})
}
async fn switch_to_kubeadm(debug: bool) -> Result<(), String> {
println!(
"{} setting KubernetesMode=kubeadm, KubernetesEnabled=true",
"k8s doctor".bright_cyan()
);
write_kubernetes_settings("kubeadm", true)?;
if is_interactive_tty() {
let ok = Confirm::with_theme(&ColorfulTheme::default())
.with_prompt("Restart Docker Desktop now to apply kubeadm? (required)")
.default(true)
.interact()
.map_err(|e| e.to_string())?;
if !ok {
return Err(
"kubeadm settings written but Docker Desktop was not restarted — run `docker desktop restart` then `xbp kubernetes doctor`"
.into(),
);
}
}
restart_docker_desktop(debug).await?;
let _ = reset_cluster(debug).await;
Ok(())
}
async fn kind_repull_and_reset(debug: bool) -> Result<(), String> {
let image = resolve_kind_node_image().await;
println!(
"{} kind recovery: repull {image} then reset cluster",
"k8s doctor".bright_cyan()
);
let _ = run_docker(&["image", "rm", "-f", &image], debug).await;
let _ = run_docker(&["builder", "prune", "-af"], debug).await;
let _ = run_docker(&["image", "prune", "-af"], debug).await;
run_docker(&["pull", &image], debug).await?;
let _ = write_kubernetes_settings("kind", true);
reset_cluster(debug).await?;
Ok(())
}
async fn reset_cluster(debug: bool) -> Result<(), String> {
println!(
"{} `docker desktop kubernetes reset-cluster`",
"k8s doctor".bright_cyan()
);
run_docker(&["desktop", "kubernetes", "reset-cluster"], debug).await?;
Ok(())
}
async fn restart_docker_desktop(debug: bool) -> Result<(), String> {
println!(
"{} restarting Docker Desktop (this can take a minute)…",
"k8s doctor".bright_cyan()
);
run_docker(&["desktop", "restart"], debug).await?;
sleep(Duration::from_secs(5)).await;
Ok(())
}
async fn finalize_wait(debug: bool) -> Result<DoctorOutcome, String> {
println!(
"{} waiting for Kubernetes to become ready…",
"k8s doctor".bright_cyan()
);
let deadline = tokio::time::Instant::now() + Duration::from_secs(180);
let mut last = String::new();
while tokio::time::Instant::now() < deadline {
if let Ok(s) = kubernetes_status(debug).await {
let line = format!(
"state={} mode={} msg={}",
s.state,
s.mode.as_deref().unwrap_or("?"),
s.progress.as_deref().unwrap_or("")
);
if line != last {
println!(" {}", line.bright_black());
last = line;
}
if s.is_running() && kubectl_cluster_ready(debug).await {
println!(
"{} {}",
"k8s doctor:".green().bold(),
"cluster ready — kubectl can reach the API".bright_green()
);
print_namespace_snapshot(debug).await;
return Ok(DoctorOutcome::Ready);
}
if s.state.eq_ignore_ascii_case("error") || s.state.eq_ignore_ascii_case("failed") {
if let Some(err) = s.error {
println!(" {} {err}", "cluster error".red());
}
}
}
sleep(Duration::from_secs(4)).await;
}
println!(
"{} cluster still not ready after wait",
"k8s doctor:".red().bold()
);
print_manual_purge_steps();
Ok(DoctorOutcome::StillDown)
}
#[derive(Debug, Clone)]
struct K8sDesktopStatus {
state: String,
mode: Option<String>,
version: Option<String>,
progress: Option<String>,
error: Option<String>,
}
impl K8sDesktopStatus {
fn is_running(&self) -> bool {
self.state.eq_ignore_ascii_case("running")
}
}
async fn kubernetes_status(debug: bool) -> Result<K8sDesktopStatus, String> {
let out = run_docker_capture(
&["desktop", "kubernetes", "status", "--format", "json"],
debug,
)
.await?;
parse_kubernetes_status_json(&out)
}
fn parse_kubernetes_status_json(raw: &str) -> Result<K8sDesktopStatus, String> {
let v: Value = serde_json::from_str(raw).map_err(|e| format!("parse status json: {e}"))?;
let content = v.get("content").cloned().unwrap_or(v.clone());
let state = v
.get("status")
.and_then(|s| s.as_str())
.or_else(|| content.get("state").and_then(|s| s.as_str()))
.unwrap_or("unknown")
.to_string();
Ok(K8sDesktopStatus {
state,
mode: content
.get("mode")
.and_then(|m| m.as_str())
.map(str::to_string),
version: content
.get("version")
.and_then(|m| m.as_str())
.map(str::to_string),
progress: content
.get("progressMessage")
.and_then(|m| m.as_str())
.map(str::to_string),
error: content
.get("error")
.and_then(|m| m.as_str())
.map(str::to_string)
.or_else(|| {
v.get("error")
.and_then(|m| m.as_str())
.map(str::to_string)
}),
})
}
async fn resolve_kind_node_image() -> String {
if let Some(ver) = read_settings_string("KubernetesNodesVersion") {
let v = ver.trim().trim_start_matches('v');
if !v.is_empty() {
return format!("kindest/node:v{v}");
}
}
if let Ok(out) =
run_docker_capture(&["desktop", "kubernetes", "images", "--format", "json"], false).await
{
if let Ok(v) = serde_json::from_str::<Value>(&out) {
if let Some(tag) = v
.pointer("/kind")
.and_then(|k| k.as_array())
.and_then(|arr| {
arr.iter().find_map(|img| {
let name = img.get("name").and_then(|n| n.as_str()).unwrap_or("");
let repo = img.get("repo").and_then(|n| n.as_str()).unwrap_or("");
if name == "node" && repo.contains("kindest") {
img.get("tag").and_then(|t| t.as_str())
} else {
None
}
})
})
{
let t = tag.trim().trim_start_matches('v');
if !t.is_empty() {
return format!("kindest/node:v{t}");
}
}
}
}
KIND_NODE_IMAGE_DEFAULT.to_string()
}
fn settings_store_candidates() -> Vec<PathBuf> {
let mut paths = Vec::new();
if let Some(data) = dirs::data_dir() {
paths.push(data.join("Docker").join("settings-store.json"));
paths.push(data.join("Docker").join("settings.json"));
}
if let Some(home) = dirs::home_dir() {
paths.push(
home.join("Library/Group Containers/group.com.docker/settings-store.json"),
);
paths.push(home.join("Library/Group Containers/group.com.docker/settings.json"));
paths.push(home.join(".docker/desktop/settings-store.json"));
paths.push(home.join(".docker/desktop/settings.json"));
}
paths
}
fn find_settings_store() -> Option<PathBuf> {
settings_store_candidates().into_iter().find(|p| p.is_file())
}
fn read_kubernetes_mode() -> Option<String> {
read_settings_string("KubernetesMode")
}
fn read_settings_string(key: &str) -> Option<String> {
let path = find_settings_store()?;
let raw = std::fs::read_to_string(path).ok()?;
let v: Value = serde_json::from_str(&raw).ok()?;
v.get(key)
.and_then(|x| x.as_str())
.map(str::to_string)
.or_else(|| {
v.pointer(&format!("/{key}"))
.and_then(|x| x.as_str())
.map(str::to_string)
})
}
fn write_kubernetes_settings(mode: &str, enabled: bool) -> Result<(), String> {
let path = find_settings_store().ok_or_else(|| {
"Docker Desktop settings-store.json not found (is Docker Desktop installed?)".to_string()
})?;
write_kubernetes_settings_at(&path, mode, enabled)
}
fn write_kubernetes_settings_at(path: &Path, mode: &str, enabled: bool) -> Result<(), String> {
let raw = std::fs::read_to_string(path).map_err(|e| format!("read {}: {e}", path.display()))?;
let mut v: Value =
serde_json::from_str(&raw).map_err(|e| format!("parse {}: {e}", path.display()))?;
let obj = v
.as_object_mut()
.ok_or_else(|| "settings-store.json is not a JSON object".to_string())?;
obj.insert(
"KubernetesMode".into(),
Value::String(mode.to_string()),
);
obj.insert("KubernetesEnabled".into(), Value::Bool(enabled));
let pretty = serde_json::to_string_pretty(&v).map_err(|e| e.to_string())?;
let tmp = path.with_extension("json.tmp");
std::fs::write(&tmp, pretty.as_bytes()).map_err(|e| format!("write temp: {e}"))?;
std::fs::rename(&tmp, path).map_err(|e| format!("replace settings: {e}"))?;
println!(
" {} wrote {} (KubernetesMode={mode}, KubernetesEnabled={enabled})",
"ok".green(),
path.display()
);
Ok(())
}
async fn docker_desktop_cli_available() -> bool {
run_docker_capture(&["desktop", "version"], false)
.await
.map(|s| !s.trim().is_empty())
.unwrap_or(false)
}
async fn print_namespace_snapshot(debug: bool) {
let config = crate::commands::service::load_xbp_config_with_root()
.await
.ok()
.map(|(_, c)| c);
let inv = crate::commands::k8s_inventory::collect_namespace_inventory(
Some("docker-desktop"),
config.as_ref(),
debug,
)
.await;
crate::commands::k8s_inventory::print_namespace_inventory(&inv, false);
}
async fn kubectl_cluster_ready(debug: bool) -> bool {
if !command_exists("kubectl") {
return false;
}
let mut cmd = Command::new("kubectl");
cmd.args([
"--context",
"docker-desktop",
"cluster-info",
"--request-timeout=8s",
]);
cmd.stdout(Stdio::piped()).stderr(Stdio::piped());
match cmd.output().await {
Ok(out) if out.status.success() => {
if debug {
let s = String::from_utf8_lossy(&out.stdout);
println!(" {}", s.lines().next().unwrap_or("cluster-info ok").dimmed());
}
true
}
_ => false,
}
}
async fn current_context_hint() -> String {
if !command_exists("kubectl") {
return String::new();
}
let out = Command::new("kubectl")
.args(["config", "current-context"])
.stdout(Stdio::piped())
.stderr(Stdio::null())
.output()
.await;
match out {
Ok(o) => String::from_utf8_lossy(&o.stdout).trim().to_string(),
Err(_) => String::new(),
}
}
async fn print_disk_hints() {
println!(
" {} check free disk before retries (low disk causes incomplete containerd snapshots)",
"tip".bright_black()
);
if command_exists("docker") {
if let Ok(out) = run_docker_capture(&["system", "df"], false).await {
for line in out.lines().take(6) {
println!(" {}", line.bright_black());
}
}
}
}
fn print_manual_purge_steps() {
println!(
"{} if kubeadm/kind still fail, purge internal desktop-containerd:",
"k8s doctor".yellow().bold()
);
println!(" 1. Docker Desktop → Troubleshoot → Clean / Purge data");
println!(" 2. After restart: `docker pull kindest/node:<ver>` only if using kind");
println!(" 3. Settings → Kubernetes → Cluster = kubeadm (preferred) → start");
println!(" 4. `kubectl --context docker-desktop cluster-info`");
println!(
" {}",
"Note: docker system prune does not fix /var/lib/desktop-containerd/.".bright_black()
);
}
async fn run_docker(args: &[&str], debug: bool) -> Result<(), String> {
let out = run_docker_capture(args, debug).await?;
if debug && !out.trim().is_empty() {
for line in out.lines().take(20) {
println!(" {}", line.dimmed());
}
}
Ok(())
}
async fn run_docker_capture(args: &[&str], debug: bool) -> Result<String, String> {
if debug {
eprintln!("debug: docker {}", args.join(" "));
}
let output = Command::new("docker")
.args(args)
.stdout(Stdio::piped())
.stderr(Stdio::piped())
.output()
.await
.map_err(|e| format!("failed to start docker: {e}"))?;
let stdout = String::from_utf8_lossy(&output.stdout).trim().to_string();
let stderr = String::from_utf8_lossy(&output.stderr).trim().to_string();
if output.status.success() {
Ok(if stdout.is_empty() { stderr } else { stdout })
} else {
Err(if stderr.is_empty() {
stdout
} else {
stderr
})
}
}
pub async fn offer_deploy_recovery(failure_text: &str, yes: bool, debug: bool) -> bool {
if !looks_like_docker_desktop_k8s_issue(failure_text) {
let ctx = current_context_hint().await;
if !ctx.eq_ignore_ascii_case("docker-desktop")
&& !looks_like_docker_desktop_k8s_issue(&ctx)
{
return false;
}
if !failure_text.to_ascii_lowercase().contains("cluster not ready")
&& !failure_text.to_ascii_lowercase().contains("connection refused")
&& !failure_text.to_ascii_lowercase().contains("cannot reach")
&& !failure_text.to_ascii_lowercase().contains("not listening")
{
return false;
}
}
if !yes && !is_interactive_tty() {
println!(
"{} Docker Desktop Kubernetes looks unhealthy. Re-run in a TTY or: `xbp kubernetes doctor --yes`",
"k8s doctor:".yellow().bold()
);
return false;
}
if !yes {
let proceed = Confirm::with_theme(&ColorfulTheme::default())
.with_prompt("Local Docker Desktop Kubernetes looks down. Run interactive recovery now?")
.default(true)
.interact()
.unwrap_or(false);
if !proceed {
return false;
}
} else {
println!(
"{} auto-recovering Docker Desktop Kubernetes (prefer kubeadm)…",
"k8s doctor".bright_cyan().bold()
);
}
match maybe_recover_docker_desktop_k8s(
failure_text,
DoctorOptions {
yes,
debug,
prefer_kubeadm: true,
},
)
.await
{
Ok(DoctorOutcome::Ready) => {
let _ = std::io::stdout().flush();
true
}
Ok(_) => false,
Err(e) => {
println!("{} recovery error: {e}", "k8s doctor:".red().bold());
false
}
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn detects_docker_desktop_refused() {
let s = "cluster not ready: kubectl failed: Docker Desktop Kubernetes API is not listening at 127.0.0.1:6443 (connection refused)";
assert!(looks_like_docker_desktop_k8s_issue(s));
}
#[test]
fn detects_kindest_lchown() {
assert!(looks_like_docker_desktop_k8s_issue(
"failed to Lchown ... desktop-containerd ... kindest/node:v1.36.1"
));
}
#[test]
fn ignores_unrelated() {
assert!(!looks_like_docker_desktop_k8s_issue(
"oci: manifest unknown for ghcr.io/xylex-group/athena"
));
}
#[test]
fn parses_status_json() {
let raw = r#"{"content":{"mode":"kubeadm","progressMessage":"up","version":"v1.34.1"},"status":"running"}"#;
let s = parse_kubernetes_status_json(raw).unwrap();
assert!(s.is_running());
assert_eq!(s.mode.as_deref(), Some("kubeadm"));
}
#[test]
fn write_settings_roundtrip() {
let dir = tempfile_dir();
let path = dir.join("settings-store.json");
std::fs::write(
&path,
r#"{"KubernetesEnabled":false,"KubernetesMode":"kind","SettingsVersion":44}"#,
)
.unwrap();
write_kubernetes_settings_at(&path, "kubeadm", true).unwrap();
let raw = std::fs::read_to_string(&path).unwrap();
let v: Value = serde_json::from_str(&raw).unwrap();
assert_eq!(v["KubernetesMode"], "kubeadm");
assert_eq!(v["KubernetesEnabled"], true);
let _ = std::fs::remove_dir_all(dir);
}
fn tempfile_dir() -> PathBuf {
let mut p = std::env::temp_dir();
p.push(format!("xbp-dd-k8s-test-{}", std::process::id()));
let _ = std::fs::create_dir_all(&p);
p
}
}