From f7932ddcfaef5f877ae457e6e41b5676bba6d0e1 Mon Sep 17 00:00:00 2001 From: Lawrence Chen Date: Fri, 22 May 2026 18:26:15 -0700 Subject: [PATCH 1/9] Add Rust cloud CLI entrypoint --- .github/workflows/nightly.yml | 4 + .github/workflows/release.yml | 4 + CLI/cmux.swift | 103 ++- Native/CloudCLI/Cargo.lock | 107 +++ Native/CloudCLI/Cargo.toml | 19 + Native/CloudCLI/src/main.rs | 1130 ++++++++++++++++++++++++++++++++ cmux.xcodeproj/project.pbxproj | 24 + docs/cli-contract.md | 13 +- scripts/build-cloud-cli.sh | 114 ++++ 9 files changed, 1509 insertions(+), 9 deletions(-) create mode 100644 Native/CloudCLI/Cargo.lock create mode 100644 Native/CloudCLI/Cargo.toml create mode 100644 Native/CloudCLI/src/main.rs create mode 100755 scripts/build-cloud-cli.sh diff --git a/.github/workflows/nightly.yml b/.github/workflows/nightly.yml index 2d97259cf717..f7eeace222d4 100644 --- a/.github/workflows/nightly.yml +++ b/.github/workflows/nightly.yml @@ -212,15 +212,19 @@ jobs: set -euo pipefail APP_BINARY="build-universal/Build/Products/Release/cmux.app/Contents/MacOS/cmux" CLI_BINARY="build-universal/Build/Products/Release/cmux.app/Contents/Resources/bin/cmux" + CLOUD_CLI_BINARY="build-universal/Build/Products/Release/cmux.app/Contents/Resources/bin/cmux-cloud" HELPER_BINARY="build-universal/Build/Products/Release/cmux.app/Contents/Resources/bin/ghostty" APP_ARCHS="$(lipo -archs "$APP_BINARY")" CLI_ARCHS="$(lipo -archs "$CLI_BINARY")" + CLOUD_CLI_ARCHS="$(lipo -archs "$CLOUD_CLI_BINARY")" HELPER_ARCHS="$(lipo -archs "$HELPER_BINARY")" echo "App binary architectures: $APP_ARCHS" echo "CLI binary architectures: $CLI_ARCHS" + echo "Cloud CLI binary architectures: $CLOUD_CLI_ARCHS" echo "Ghostty helper architectures: $HELPER_ARCHS" [[ "$APP_ARCHS" == *arm64* && "$APP_ARCHS" == *x86_64* ]] [[ "$CLI_ARCHS" == *arm64* && "$CLI_ARCHS" == *x86_64* ]] + [[ "$CLOUD_CLI_ARCHS" == *arm64* && "$CLOUD_CLI_ARCHS" == *x86_64* ]] [[ "$HELPER_ARCHS" == *arm64* && "$HELPER_ARCHS" == *x86_64* ]] - name: Run CLI version memory guard regression diff --git a/.github/workflows/release.yml b/.github/workflows/release.yml index c842a3496972..2e8b32bada3d 100644 --- a/.github/workflows/release.yml +++ b/.github/workflows/release.yml @@ -176,15 +176,19 @@ jobs: set -euo pipefail APP_BINARY="build-universal/Build/Products/Release/cmux.app/Contents/MacOS/cmux" CLI_BINARY="build-universal/Build/Products/Release/cmux.app/Contents/Resources/bin/cmux" + CLOUD_CLI_BINARY="build-universal/Build/Products/Release/cmux.app/Contents/Resources/bin/cmux-cloud" HELPER_BINARY="build-universal/Build/Products/Release/cmux.app/Contents/Resources/bin/ghostty" APP_ARCHS="$(lipo -archs "$APP_BINARY")" CLI_ARCHS="$(lipo -archs "$CLI_BINARY")" + CLOUD_CLI_ARCHS="$(lipo -archs "$CLOUD_CLI_BINARY")" HELPER_ARCHS="$(lipo -archs "$HELPER_BINARY")" echo "App binary architectures: $APP_ARCHS" echo "CLI binary architectures: $CLI_ARCHS" + echo "Cloud CLI binary architectures: $CLOUD_CLI_ARCHS" echo "Ghostty helper architectures: $HELPER_ARCHS" [[ "$APP_ARCHS" == *arm64* && "$APP_ARCHS" == *x86_64* ]] [[ "$CLI_ARCHS" == *arm64* && "$CLI_ARCHS" == *x86_64* ]] + [[ "$CLOUD_CLI_ARCHS" == *arm64* && "$CLOUD_CLI_ARCHS" == *x86_64* ]] [[ "$HELPER_ARCHS" == *arm64* && "$HELPER_ARCHS" == *x86_64* ]] - name: Build remote daemon release assets and inject manifest diff --git a/CLI/cmux.swift b/CLI/cmux.swift index ef4a7071332e..a3663e3bc4b3 100644 --- a/CLI/cmux.swift +++ b/CLI/cmux.swift @@ -2629,6 +2629,83 @@ struct CMUXCLI { } } + private static func cloudCLIExecutablePath() -> String? { + let fileManager = FileManager.default + var candidates: [URL] = [] + if let override = normalizedEnvValue(ProcessInfo.processInfo.environment["CMUX_CLOUD_CLI_PATH"]) { + candidates.append(URL(fileURLWithPath: override)) + } + if let executableURL = CLIExecutableLocator.currentExecutableURL() { + candidates.append(executableURL.deletingLastPathComponent().appendingPathComponent("cmux-cloud")) + } + if let appBundle = CLIExecutableLocator.enclosingAppBundle() { + candidates.append( + appBundle.bundleURL + .appendingPathComponent("Contents", isDirectory: true) + .appendingPathComponent("Resources", isDirectory: true) + .appendingPathComponent("bin", isDirectory: true) + .appendingPathComponent("cmux-cloud", isDirectory: false) + ) + } + for candidate in candidates { + let path = candidate.path + if fileManager.isExecutableFile(atPath: path) { + return path + } + } + return nil + } + + private func runCloudRustCommand( + commandArgs: [String], + socketPath: String, + explicitPassword: String?, + jsonOutput: Bool, + idFormatArg: String?, + windowId: String? + ) throws -> Never { + guard let helperPath = Self.cloudCLIExecutablePath() else { + throw CLIError(message: "cmux cloud helper is missing. Reload cmux so Contents/Resources/bin/cmux-cloud is bundled.") + } + + let resolvedPassword = SocketPasswordResolver.resolve( + explicit: explicitPassword, + socketPath: socketPath + ) + setenv("CMUX_CLOUD_SOCKET_PATH", socketPath, 1) + setenv("CMUX_SOCKET_PATH", socketPath, 1) + setenv("CMUX_CLOUD_JSON", jsonOutput ? "1" : "0", 1) + if let resolvedPassword { + setenv("CMUX_CLOUD_SOCKET_PASSWORD", resolvedPassword, 1) + } else { + unsetenv("CMUX_CLOUD_SOCKET_PASSWORD") + } + if let idFormatArg { + setenv("CMUX_CLOUD_ID_FORMAT", idFormatArg, 1) + } else { + unsetenv("CMUX_CLOUD_ID_FORMAT") + } + if let windowId { + setenv("CMUX_CLOUD_WINDOW", windowId, 1) + } else { + unsetenv("CMUX_CLOUD_WINDOW") + } + if let currentExecutablePath = CLIExecutableLocator.currentExecutableURL()?.path { + setenv("CMUX_CLOUD_PARENT_CLI", currentExecutablePath, 1) + } + + var argv = ([helperPath] + commandArgs).map { strdup($0) } + defer { + for item in argv { + free(item) + } + } + argv.append(nil) + execv(helperPath, &argv) + let reason = String(cString: strerror(errno)) + throw CLIError(message: "Failed to launch cmux cloud helper at \(helperPath): \(reason)") + } + private static func shouldFocusWindowBeforeDispatch(command: String, commandArgs: [String]) -> Bool { let normalizedCommand = command.lowercased() if normalizedCommand == "surface-resume" { @@ -2993,6 +3070,17 @@ struct CMUXCLI { commandArgs: commandArgs ) + if command == "cloud" { + try runCloudRustCommand( + commandArgs: commandArgs, + socketPath: resolvedSocketPath, + explicitPassword: socketPasswordArg, + jsonOutput: jsonOutput, + idFormatArg: idFormatArg, + windowId: windowId + ) + } + let client = SocketClient(path: resolvedSocketPath) if resolvedSocketPath != socketPath { cliTelemetry.breadcrumb( @@ -3117,7 +3205,7 @@ struct CMUXCLI { throw CLIError(message: "Usage: cmux auth ") } - case "vm", "cloud": + case "vm": let sub = commandArgs.first?.lowercased() ?? "ls" let rest = Array(commandArgs.dropFirst()) switch sub { @@ -11576,7 +11664,7 @@ struct CMUXCLI { return """ Usage: cmux \(command) [args...] - Manage cloud VMs. `cloud` is an alias for `vm`. Requires `cmux auth login`. + Manage cloud VMs. Requires `cmux auth login`. Subcommands: ls List your cloud VMs. @@ -11601,10 +11689,10 @@ struct CMUXCLI { local testing from the web worktree. Example: - cmux vm new - cmux vm ls - cmux cloud exec -- echo hello - cmux vm rm + cmux \(command) new + cmux \(command) ls + cmux \(command) exec -- echo hello + cmux \(command) rm """ case "rpc": return """ @@ -29440,7 +29528,8 @@ export default function cmuxPiSessionExtension(pi: ExtensionAPI) { events [--after ] [--cursor-file ] [--name ] [--category ] [--reconnect] [--limit ] [--no-ack] [--no-heartbeat] auth login | logout (aliases for auth login/logout) - vm [args...] (alias: cloud) + cloud [args...] + vm [args...] (compatibility alias) rpc [json-params] identify [--workspace ] [--surface ] [--window ] [--no-caller] list-windows diff --git a/Native/CloudCLI/Cargo.lock b/Native/CloudCLI/Cargo.lock new file mode 100644 index 000000000000..c4d068d0245e --- /dev/null +++ b/Native/CloudCLI/Cargo.lock @@ -0,0 +1,107 @@ +# This file is automatically @generated by Cargo. +# It is not intended for manual editing. +version = 4 + +[[package]] +name = "cmux-cloud-cli" +version = "0.1.0" +dependencies = [ + "serde", + "serde_json", +] + +[[package]] +name = "itoa" +version = "1.0.18" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "8f42a60cbdf9a97f5d2305f08a87dc4e09308d1276d28c869c684d7777685682" + +[[package]] +name = "memchr" +version = "2.8.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "f8ca58f447f06ed17d5fc4043ce1b10dd205e060fb3ce5b979b8ed8e59ff3f79" + +[[package]] +name = "proc-macro2" +version = "1.0.106" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "8fd00f0bb2e90d81d1044c2b32617f68fcb9fa3bb7640c23e9c748e53fb30934" +dependencies = [ + "unicode-ident", +] + +[[package]] +name = "quote" +version = "1.0.45" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "41f2619966050689382d2b44f664f4bc593e129785a36d6ee376ddf37259b924" +dependencies = [ + "proc-macro2", +] + +[[package]] +name = "serde" +version = "1.0.228" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "9a8e94ea7f378bd32cbbd37198a4a91436180c5bb472411e48b5ec2e2124ae9e" +dependencies = [ + "serde_core", + "serde_derive", +] + +[[package]] +name = "serde_core" +version = "1.0.228" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "41d385c7d4ca58e59fc732af25c3983b67ac852c1a25000afe1175de458b67ad" +dependencies = [ + "serde_derive", +] + +[[package]] +name = "serde_derive" +version = "1.0.228" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "d540f220d3187173da220f885ab66608367b6574e925011a9353e4badda91d79" +dependencies = [ + "proc-macro2", + "quote", + "syn", +] + +[[package]] +name = "serde_json" +version = "1.0.150" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "e8014e44b4736ed0538adeecded0fce2a272f22dc9578a7eb6b2d9993c74cfb9" +dependencies = [ + "itoa", + "memchr", + "serde", + "serde_core", + "zmij", +] + +[[package]] +name = "syn" +version = "2.0.117" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "e665b8803e7b1d2a727f4023456bbbbe74da67099c585258af0ad9c5013b9b99" +dependencies = [ + "proc-macro2", + "quote", + "unicode-ident", +] + +[[package]] +name = "unicode-ident" +version = "1.0.24" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "e6e4313cd5fcd3dad5cafa179702e2b244f760991f45397d14d4ebf38247da75" + +[[package]] +name = "zmij" +version = "1.0.21" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "b8848ee67ecc8aedbaf3e4122217aff892639231befc6a1b58d29fff4c2cabaa" diff --git a/Native/CloudCLI/Cargo.toml b/Native/CloudCLI/Cargo.toml new file mode 100644 index 000000000000..e161342122a9 --- /dev/null +++ b/Native/CloudCLI/Cargo.toml @@ -0,0 +1,19 @@ +[package] +name = "cmux-cloud-cli" +version = "0.1.0" +edition = "2021" +publish = false + +[[bin]] +name = "cmux-cloud" +path = "src/main.rs" + +[dependencies] +serde = { version = "1", features = ["derive"] } +serde_json = "1" + +[profile.dev] +panic = "abort" + +[profile.release] +panic = "abort" diff --git a/Native/CloudCLI/src/main.rs b/Native/CloudCLI/src/main.rs new file mode 100644 index 000000000000..8c36a7f31a09 --- /dev/null +++ b/Native/CloudCLI/src/main.rs @@ -0,0 +1,1130 @@ +use serde::{Deserialize, Serialize}; +use serde_json::{json, Map, Value}; +use std::env; +use std::fs; +use std::io::{ErrorKind, Read, Write}; +use std::os::unix::fs::{FileTypeExt, MetadataExt}; +use std::os::unix::net::UnixStream; +use std::path::{Path, PathBuf}; +use std::process::{self, Command}; +use std::thread; +use std::time::{Duration, Instant, SystemTime, UNIX_EPOCH}; + +const DEFAULT_RESPONSE_TIMEOUT: Duration = Duration::from_secs(15); +const VM_CREATE_RESPONSE_TIMEOUT: Duration = Duration::from_secs(16 * 60); +const VM_CREATE_IDEMPOTENCY_TTL: u64 = 10 * 60; + +extern "C" { + fn getuid() -> u32; +} + +#[derive(Debug)] +struct CliError { + message: String, + exit_code: i32, +} + +impl CliError { + fn new(message: impl Into) -> Self { + Self { + message: message.into(), + exit_code: 1, + } + } + + fn exit(message: impl Into, exit_code: i32) -> Self { + Self { + message: message.into(), + exit_code, + } + } +} + +type CliResult = Result; + +#[derive(Debug, Clone)] +struct CloudContext { + socket_path: String, + socket_password: Option, + json_output: bool, + window_override: Option, + parent_cli: String, +} + +#[derive(Debug, Deserialize, Serialize)] +struct VMCreateIdempotencyStore { + #[serde(default)] + records: std::collections::BTreeMap, +} + +#[derive(Debug, Deserialize, Serialize)] +struct VMCreateIdempotencyRecord { + key: String, + #[serde(rename = "createdAt")] + created_at: f64, +} + +#[derive(Debug)] +struct ActiveVMCreateIdempotency { + signature: String, + key: String, +} + +struct SocketClient { + stream: UnixStream, +} + +impl SocketClient { + fn connect(path: String) -> CliResult { + let deadline = Instant::now() + Duration::from_millis(350); + loop { + match Self::connect_once(&path) { + Ok(stream) => return Ok(Self { stream }), + Err(error) if should_retry_connect(&error) && Instant::now() < deadline => { + thread::sleep(Duration::from_millis(25)); + } + Err(error) => return Err(error), + } + } + } + + fn connect_once(path: &str) -> CliResult { + let metadata = + fs::metadata(path).map_err(|_| CliError::new(format!("Socket not found at {path}")))?; + if !metadata.file_type().is_socket() { + return Err(CliError::new(format!( + "Path exists at {path} but is not a Unix socket" + ))); + } + let current_uid = unsafe { getuid() }; + if metadata.uid() != current_uid { + return Err(CliError::new(format!( + "Socket at {path} is not owned by the current user, refusing to connect" + ))); + } + + let stream = UnixStream::connect(path).map_err(|error| { + CliError::new(format!("Failed to connect to socket at {path} ({error})")) + })?; + stream + .set_read_timeout(Some(response_timeout())) + .map_err(|error| { + CliError::new(format!("Failed to configure socket timeout: {error}")) + })?; + stream + .set_write_timeout(Some(response_timeout())) + .map_err(|error| { + CliError::new(format!("Failed to configure socket timeout: {error}")) + })?; + Ok(stream) + } + + fn authenticate_if_needed(&mut self, password: Option<&str>) -> CliResult<()> { + let Some(password) = password else { + return Ok(()); + }; + let response = self.send_line(&format!("auth {password}"), response_timeout())?; + if response.starts_with("ERROR:") && !response.contains("Unknown command 'auth'") { + return Err(CliError::new(response)); + } + Ok(()) + } + + fn send_v2( + &mut self, + method: &str, + params: Value, + timeout: Duration, + ) -> CliResult> { + let request = json!({ + "id": request_id(), + "method": method, + "params": params, + }); + let raw = self.send_line(&request.to_string(), timeout)?; + if raw.starts_with("ERROR:") { + return Err(CliError::new(raw)); + } + + let response: Value = serde_json::from_str(&raw) + .map_err(|_| CliError::new(format!("Invalid v2 response: {raw}")))?; + if response.get("ok").and_then(Value::as_bool) == Some(true) { + return Ok(response + .get("result") + .and_then(Value::as_object) + .cloned() + .unwrap_or_default()); + } + if let Some(error) = response.get("error").and_then(Value::as_object) { + let code = error.get("code").and_then(Value::as_str).unwrap_or("error"); + let message = error + .get("message") + .and_then(Value::as_str) + .unwrap_or("Unknown v2 error"); + let action = error.get("action").and_then(Value::as_str); + let reason = error.get("reason").and_then(Value::as_str); + let details = safe_v2_details(error.get("details")); + return Err(CliError::new(format_v2_error( + code, + message, + action, + reason, + details.as_deref(), + ))); + } + Err(CliError::new("v2 request failed")) + } + + fn send_line(&mut self, line: &str, timeout: Duration) -> CliResult { + self.stream + .set_read_timeout(Some(timeout)) + .map_err(|error| { + CliError::new(format!("Failed to configure socket timeout: {error}")) + })?; + self.stream + .set_write_timeout(Some(timeout)) + .map_err(|error| { + CliError::new(format!("Failed to configure socket timeout: {error}")) + })?; + self.stream + .write_all(line.as_bytes()) + .and_then(|_| self.stream.write_all(b"\n")) + .map_err(|error| match error.kind() { + ErrorKind::WouldBlock | ErrorKind::TimedOut => CliError::new("Command timed out"), + _ => CliError::new(format!("Failed to write to socket ({error})")), + })?; + + let mut data = Vec::new(); + let mut byte = [0_u8; 1]; + loop { + match self.stream.read(&mut byte) { + Ok(0) if data.is_empty() => { + return Err(CliError::new("Socket closed before reply")); + } + Ok(0) => break, + Ok(_) if byte[0] == b'\n' => break, + Ok(_) => data.push(byte[0]), + Err(error) if matches!(error.kind(), ErrorKind::Interrupted) => {} + Err(error) + if matches!(error.kind(), ErrorKind::WouldBlock | ErrorKind::TimedOut) => + { + return Err(CliError::new("Command timed out")); + } + Err(error) => return Err(CliError::new(format!("Socket read error ({error})"))), + } + } + + String::from_utf8(data).map_err(|_| CliError::new("Invalid UTF-8 response")) + } +} + +fn main() { + if let Err(error) = run() { + eprintln!("Error: {}", error.message); + process::exit(error.exit_code); + } +} + +fn run() -> CliResult<()> { + let mut args = env::args().skip(1).collect::>(); + let json_from_args = take_flag_before_terminator(&mut args, "--json"); + if take_flag_before_terminator(&mut args, "--help") + || take_flag_before_terminator(&mut args, "-h") + { + println!("{}", usage()); + return Ok(()); + } + + let ctx = CloudContext::from_env(json_from_args)?; + let subcommand = args.first().map(|arg| arg.as_str()).unwrap_or("ls"); + let rest = if args.is_empty() { + Vec::new() + } else { + args[1..].to_vec() + }; + + match subcommand { + "ls" | "list" => run_list(&ctx), + "new" | "create" => run_new(&ctx, &rest), + "rm" | "destroy" | "delete" => run_destroy(&ctx, &rest), + "exec" => run_exec(&ctx, &rest), + "ssh-info" => run_ssh_info(&ctx, &rest), + "shell" | "attach" | "ssh" => run_delegated_interactive(&ctx, subcommand, &rest), + "ssh-attach" => { + let mut vm_args = vec![subcommand.to_string()]; + vm_args.extend(rest); + exec_parent_vm(&ctx, &vm_args) + } + "help" => { + println!("{}", usage()); + Ok(()) + } + _ => Err(CliError::exit( + format!( + "Usage: cmux cloud [args...]\n\nCommon commands:\n cmux cloud ls\n cmux cloud new\n cmux cloud ssh \n cmux cloud rm " + ), + 2, + )), + } +} + +impl CloudContext { + fn from_env(json_from_args: bool) -> CliResult { + let socket_path = normalized_env("CMUX_CLOUD_SOCKET_PATH") + .or_else(|| normalized_env("CMUX_SOCKET_PATH")) + .or_else(|| normalized_env("CMUX_SOCKET")) + .ok_or_else(|| { + CliError::new("cmux cloud needs CMUX_SOCKET_PATH from the cmux launcher") + })?; + let socket_password = normalized_env("CMUX_CLOUD_SOCKET_PASSWORD") + .or_else(|| normalized_env("CMUX_SOCKET_PASSWORD")); + let json_output = json_from_args || env_flag("CMUX_CLOUD_JSON"); + let window_override = normalized_env("CMUX_CLOUD_WINDOW"); + let parent_cli = + normalized_env("CMUX_CLOUD_PARENT_CLI").unwrap_or_else(|| "cmux".to_string()); + Ok(Self { + socket_path, + socket_password, + json_output, + window_override, + parent_cli, + }) + } + + fn connect(&self) -> CliResult { + let mut client = SocketClient::connect(self.socket_path.clone())?; + client.authenticate_if_needed(self.socket_password.as_deref())?; + if let Some(window_raw) = &self.window_override { + let normalized = normalize_window_handle(&mut client, window_raw)? + .unwrap_or_else(|| window_raw.to_string()); + client.send_v2( + "window.focus", + json!({ "window_id": normalized }), + response_timeout(), + )?; + } + Ok(client) + } +} + +fn run_list(ctx: &CloudContext) -> CliResult<()> { + let mut client = ctx.connect()?; + let response = client.send_v2("vm.list", json!({}), response_timeout())?; + if ctx.json_output { + println!("{}", Value::Object(response)); + return Ok(()); + } + + let vms = response + .get("vms") + .and_then(Value::as_array) + .cloned() + .unwrap_or_default(); + if vms.is_empty() { + println!("No cloud VMs. Try: cmux cloud new"); + return Ok(()); + } + + for vm in vms { + let id = value_str(vm.get("id")).unwrap_or("?"); + let provider = value_str(vm.get("provider")).unwrap_or("?"); + let image = value_str(vm.get("image")).unwrap_or("?"); + println!("{id} [{provider}] {image}"); + } + Ok(()) +} + +fn run_new(ctx: &CloudContext, args: &[String]) -> CliResult<()> { + let (image_opt, rem0) = parse_option(args, "--image"); + let (provider_opt, rem1) = parse_option(&rem0, "--provider"); + let (window_opt, rem2) = parse_option(&rem1, "--window"); + let detach = rem2.iter().any(|arg| arg == "--detach" || arg == "-d"); + let remaining = rem2 + .into_iter() + .filter(|arg| arg != "--detach" && arg != "-d") + .collect::>(); + + if let Some(unknown) = remaining + .iter() + .find(|arg| is_unknown_flag_token(arg, &["-d"])) + { + return Err(CliError::exit( + format!( + "cloud new: unknown flag '{unknown}'.\n\nKnown flags:\n --image \n --provider \n --detach, -d\n\nTry:\n cmux cloud new" + ), + 2, + )); + } + if let Some(extra) = remaining.iter().find(|arg| !is_flag_token(arg)) { + return Err(CliError::exit( + format!( + "cloud new: unexpected argument '{extra}'.\n\n`cmux cloud new` does not take a VM name or positional arguments.\n\nTry:\n cmux cloud new\n cmux cloud new --detach" + ), + 2, + )); + } + + let provider = normalized_vm_provider(provider_opt.as_deref())?; + let mut client = ctx.connect()?; + let effective_window = window_opt.as_ref().or(ctx.window_override.as_ref()); + let target_window = match effective_window { + Some(raw) => validate_window_handle(&mut client, raw)?, + None => None, + }; + + let idempotency = active_vm_create_idempotency(image_opt.as_deref(), provider.as_deref())?; + let mut params = Map::new(); + if let Some(image) = &image_opt { + params.insert("image".to_string(), Value::String(image.clone())); + } + if let Some(provider) = &provider { + params.insert("provider".to_string(), Value::String(provider.clone())); + } + params.insert( + "idempotency_key".to_string(), + Value::String(idempotency.key.clone()), + ); + + let response = client.send_v2( + "vm.create", + Value::Object(params), + VM_CREATE_RESPONSE_TIMEOUT, + )?; + if ctx.json_output { + clear_vm_create_idempotency(&idempotency)?; + println!("{}", Value::Object(response)); + return Ok(()); + } + + let id = response + .get("id") + .and_then(Value::as_str) + .unwrap_or("?") + .to_string(); + let provider_text = value_str(response.get("provider")).unwrap_or("?"); + let image_text = value_str(response.get("image")).unwrap_or("?"); + if detach { + clear_vm_create_idempotency(&idempotency)?; + println!("OK {id}"); + println!(" provider: {provider_text}"); + println!(" image: {image_text}"); + return Ok(()); + } + + println!("Created {id} [{provider_text}] {image_text}"); + drop(client); + clear_vm_create_idempotency(&idempotency)?; + + let mut vm_args = vec!["shell".to_string(), id]; + if let Some(window) = target_window { + vm_args.push("--window".to_string()); + vm_args.push(window); + } + exec_parent_vm(ctx, &vm_args) +} + +fn run_destroy(ctx: &CloudContext, args: &[String]) -> CliResult<()> { + let Some(vm_id) = args.first() else { + return Err(CliError::exit( + "Usage: cmux cloud rm \n\nFind an id:\n cmux cloud ls", + 2, + )); + }; + let mut client = ctx.connect()?; + client.send_v2( + "vm.destroy", + json!({ "id": vm_id }), + Duration::from_secs(60), + )?; + if ctx.json_output { + println!("{}", json!({ "ok": true, "id": vm_id })); + } else { + println!("OK {vm_id}"); + } + Ok(()) +} + +fn run_exec(ctx: &CloudContext, args: &[String]) -> CliResult<()> { + let Some(vm_id) = args.first() else { + return Err(CliError::exit( + "Usage: cmux cloud exec -- \n\nExamples:\n cmux cloud ls\n cmux cloud exec -- pwd", + 2, + )); + }; + let mut command_args = args[1..].to_vec(); + if command_args.first().map(String::as_str) == Some("--") { + command_args.remove(0); + } + if command_args.is_empty() { + return Err(CliError::exit( + format!( + "Usage: cmux cloud exec -- \n\nExample:\n cmux cloud exec {vm_id} -- uname -a" + ), + 2, + )); + } + + let command = command_args + .iter() + .map(|arg| shell_quote(arg)) + .collect::>() + .join(" "); + let mut client = ctx.connect()?; + let response = client.send_v2( + "vm.exec", + json!({ "id": vm_id, "command": command }), + Duration::from_secs(35), + )?; + let exit_code = response + .get("exit_code") + .and_then(Value::as_i64) + .unwrap_or(-1); + if ctx.json_output { + println!("{}", Value::Object(response)); + if exit_code != 0 { + return Err(CliError::exit(format!("exit {exit_code}"), 1)); + } + return Ok(()); + } + + if let Some(stdout) = response.get("stdout").and_then(Value::as_str) { + if !stdout.is_empty() { + print!("{stdout}"); + if !stdout.ends_with('\n') { + println!(); + } + } + } + if let Some(stderr) = response.get("stderr").and_then(Value::as_str) { + if !stderr.is_empty() { + eprint!("{stderr}"); + if !stderr.ends_with('\n') { + eprintln!(); + } + } + } + if exit_code != 0 { + return Err(CliError::exit(format!("exit {exit_code}"), 1)); + } + Ok(()) +} + +fn run_ssh_info(ctx: &CloudContext, args: &[String]) -> CliResult<()> { + let Some(vm_id) = args.first() else { + return Err(CliError::exit( + "Usage: cmux cloud ssh-info \n\nFind an id:\n cmux cloud ls", + 2, + )); + }; + let mut client = ctx.connect()?; + let response = client.send_v2( + "vm.ssh_info", + json!({ "id": vm_id }), + Duration::from_secs(60), + )?; + if ctx.json_output { + println!("{}", Value::Object(response)); + return Ok(()); + } + + let host = value_str(response.get("host")).unwrap_or("?"); + let port = response.get("port").and_then(Value::as_i64).unwrap_or(22); + let username = value_str(response.get("username")).unwrap_or("?"); + let credential = response.get("credential").and_then(Value::as_object); + let cred_kind = credential + .and_then(|cred| cred.get("kind")) + .and_then(Value::as_str) + .unwrap_or("?"); + let cred_value = credential + .and_then(|cred| cred.get("value")) + .and_then(Value::as_str) + .unwrap_or("?"); + if cred_kind == "password" { + println!("ssh {username}@{host} -p {port}"); + println!(); + println!(" host: {host}"); + println!(" port: {port}"); + println!(" username: {username}"); + println!(" password: {cred_value}"); + return Ok(()); + } + + println!("This Cloud VM does not support `cmux cloud ssh-info` in this cmux build."); + println!(); + println!("What to do:"); + println!(" Update cmux and retry."); + println!(" If this keeps happening, contact support with the VM id."); + Ok(()) +} + +fn run_delegated_interactive( + ctx: &CloudContext, + subcommand: &str, + args: &[String], +) -> CliResult<()> { + let (window_opt, vm_args) = parse_option(args, "--window"); + let Some(vm_id) = vm_args.first() else { + return Err(CliError::exit( + format!("Usage: cmux cloud {subcommand} \n\nFind an id:\n cmux cloud ls"), + 2, + )); + }; + + let mut delegated = vec![subcommand.to_string(), vm_id.clone()]; + if let Some(window) = window_opt.or_else(|| ctx.window_override.clone()) { + delegated.push("--window".to_string()); + delegated.push(window); + } + exec_parent_vm(ctx, &delegated) +} + +fn exec_parent_vm(ctx: &CloudContext, vm_args: &[String]) -> CliResult<()> { + let mut command = Command::new(&ctx.parent_cli); + command.arg("--socket").arg(&ctx.socket_path); + if let Some(password) = &ctx.socket_password { + command.arg("--password").arg(password); + } + if ctx.json_output { + command.arg("--json"); + } + command.arg("vm"); + command.args(vm_args); + let status = command + .status() + .map_err(|error| CliError::new(format!("Failed to launch cmux vm helper: {error}")))?; + process::exit(status.code().unwrap_or(1)); +} + +fn validate_window_handle(client: &mut SocketClient, raw: &str) -> CliResult> { + let Some(normalized) = normalize_window_handle(client, raw)? else { + return Ok(None); + }; + let listed = client.send_v2("window.list", json!({}), response_timeout())?; + let windows = listed + .get("windows") + .and_then(Value::as_array) + .cloned() + .unwrap_or_default(); + let found = windows + .iter() + .any(|item| window_handle_matches(&normalized, item)); + if !found { + return Err(CliError::new(format!("Window not found: {raw}"))); + } + Ok(Some(normalized)) +} + +fn normalize_window_handle(client: &mut SocketClient, raw: &str) -> CliResult> { + let trimmed = raw.trim(); + if trimmed.is_empty() { + return Ok(None); + } + if is_uuid(trimmed) { + return Ok(Some(trimmed.to_string())); + } + if is_handle_ref(trimmed) { + if let Some(matched) = matching_window_handle(client, trimmed)? { + return Ok(Some(matched)); + } + return Err(CliError::new(format!("Window not found: {trimmed}"))); + } + + let wanted_index = trimmed.parse::().map_err(|_| { + CliError::new(format!( + "Invalid window handle: {trimmed} (expected UUID, ref like window:1, or index)" + )) + })?; + let listed = client.send_v2("window.list", json!({}), response_timeout())?; + let windows = listed + .get("windows") + .and_then(Value::as_array) + .cloned() + .unwrap_or_default(); + for item in windows { + if item.get("index").and_then(Value::as_i64) == Some(wanted_index) { + if let Some(id) = item.get("id").and_then(Value::as_str) { + return Ok(Some(id.to_string())); + } + if let Some(reference) = item.get("ref").and_then(Value::as_str) { + return Ok(Some(reference.to_string())); + } + } + } + Err(CliError::new("Window index not found")) +} + +fn matching_window_handle(client: &mut SocketClient, handle: &str) -> CliResult> { + let listed = client.send_v2("window.list", json!({}), response_timeout())?; + let windows = listed + .get("windows") + .and_then(Value::as_array) + .cloned() + .unwrap_or_default(); + for item in windows { + if window_handle_matches(handle, &item) { + if let Some(id) = item.get("id").and_then(Value::as_str) { + return Ok(Some(id.to_string())); + } + if let Some(reference) = item.get("ref").and_then(Value::as_str) { + return Ok(Some(reference.to_string())); + } + return Ok(Some(handle.to_string())); + } + } + Ok(None) +} + +fn window_handle_matches(handle: &str, item: &Value) -> bool { + let Some(target) = normalized_handle_value(Some(handle)) else { + return false; + }; + for key in ["id", "ref"] { + let candidate = item.get(key).and_then(Value::as_str); + if let Some(candidate) = normalized_handle_value(candidate) { + if handles_match(&target, &candidate) { + return true; + } + } + } + false +} + +fn request_id() -> String { + let nanos = SystemTime::now() + .duration_since(UNIX_EPOCH) + .map(|duration| duration.as_nanos()) + .unwrap_or(0); + format!("cmux-cloud-{}-{nanos}", process::id()) +} + +fn response_timeout() -> Duration { + match normalized_env("CMUXTERM_CLI_RESPONSE_TIMEOUT_SEC") + .and_then(|raw| raw.parse::().ok()) + .filter(|seconds| seconds.is_finite() && *seconds > 0.0) + { + Some(seconds) => Duration::from_secs_f64(seconds), + None => DEFAULT_RESPONSE_TIMEOUT, + } +} + +fn should_retry_connect(error: &CliError) -> bool { + error.message.contains("Connection refused") + || error.message.contains("Resource temporarily unavailable") + || error.message.contains("would block") +} + +fn normalized_env(name: &str) -> Option { + env::var(name) + .ok() + .map(|value| value.trim().to_string()) + .filter(|value| !value.is_empty()) +} + +fn env_flag(name: &str) -> bool { + matches!( + normalized_env(name).as_deref(), + Some("1" | "true" | "TRUE" | "yes" | "YES") + ) +} + +fn take_flag_before_terminator(args: &mut Vec, flag: &str) -> bool { + let mut found = false; + let mut filtered = Vec::with_capacity(args.len()); + let mut past_terminator = false; + for arg in args.drain(..) { + if past_terminator { + filtered.push(arg); + continue; + } + if arg == "--" { + past_terminator = true; + filtered.push(arg); + continue; + } + if arg == flag { + found = true; + } else { + filtered.push(arg); + } + } + *args = filtered; + found +} + +fn parse_option(args: &[String], name: &str) -> (Option, Vec) { + let mut remaining = Vec::new(); + let mut value = None; + let mut index = 0; + let mut past_terminator = false; + while index < args.len() { + let arg = &args[index]; + if arg == "--" { + past_terminator = true; + remaining.push(arg.clone()); + index += 1; + continue; + } + if !past_terminator { + let equals_prefix = format!("{name}="); + if arg.starts_with(&equals_prefix) { + value = Some(arg[equals_prefix.len()..].to_string()); + index += 1; + continue; + } + if arg == name && index + 1 < args.len() { + value = Some(args[index + 1].clone()); + index += 2; + continue; + } + } + remaining.push(arg.clone()); + index += 1; + } + (value, remaining) +} + +fn normalized_vm_provider(provider: Option<&str>) -> CliResult> { + let Some(provider) = provider.map(str::trim).filter(|value| !value.is_empty()) else { + return Ok(None); + }; + let normalized = provider.to_ascii_lowercase(); + if normalized == "e2b" || normalized == "freestyle" { + return Ok(Some(normalized)); + } + Err(CliError::exit( + "cloud new: unsupported Cloud VM service override.\n\nTry:\n cmux cloud new", + 2, + )) +} + +fn is_flag_token(value: &str) -> bool { + value.starts_with('-') && value != "-" +} + +fn is_unknown_flag_token(value: &str, allowed_short_flags: &[&str]) -> bool { + is_flag_token(value) && !allowed_short_flags.contains(&value) +} + +fn idempotency_signature(image: Option<&str>, provider: Option<&str>) -> String { + format!( + "image={}\u{1f}provider={}", + image.unwrap_or("").trim(), + provider.unwrap_or("").trim().to_ascii_lowercase() + ) +} + +fn vm_create_idempotency_store_url() -> CliResult { + let home = env::var("HOME") + .map_err(|_| CliError::new("HOME is not set, cannot store VM idempotency state"))?; + Ok(Path::new(&home) + .join(".cmuxterm") + .join("vm-create-idempotency.json")) +} + +fn load_vm_create_idempotency_store(path: &Path) -> VMCreateIdempotencyStore { + fs::read_to_string(path) + .ok() + .and_then(|raw| serde_json::from_str(&raw).ok()) + .unwrap_or(VMCreateIdempotencyStore { + records: std::collections::BTreeMap::new(), + }) +} + +fn save_vm_create_idempotency_store( + store: &VMCreateIdempotencyStore, + path: &Path, +) -> CliResult<()> { + if let Some(parent) = path.parent() { + fs::create_dir_all(parent).map_err(|error| { + CliError::new(format!( + "Failed to create VM idempotency state directory: {error}" + )) + })?; + } + let data = serde_json::to_vec_pretty(store).map_err(|error| { + CliError::new(format!("Failed to encode VM idempotency state: {error}")) + })?; + fs::write(path, data) + .map_err(|error| CliError::new(format!("Failed to write VM idempotency state: {error}"))) +} + +fn active_vm_create_idempotency( + image: Option<&str>, + provider: Option<&str>, +) -> CliResult { + let path = vm_create_idempotency_store_url()?; + let signature = idempotency_signature(image, provider); + let now = unix_timestamp_secs(); + let mut store = load_vm_create_idempotency_store(&path); + store.records.retain(|_, record| { + !record.key.is_empty() + && now.saturating_sub(record.created_at as u64) < VM_CREATE_IDEMPOTENCY_TTL + }); + if let Some(existing) = store.records.get(&signature) { + save_vm_create_idempotency_store(&store, &path)?; + return Ok(ActiveVMCreateIdempotency { + signature, + key: existing.key.clone(), + }); + } + + let key = random_uuid_like(); + store.records.insert( + signature.clone(), + VMCreateIdempotencyRecord { + key: key.clone(), + created_at: now as f64, + }, + ); + save_vm_create_idempotency_store(&store, &path)?; + Ok(ActiveVMCreateIdempotency { signature, key }) +} + +fn clear_vm_create_idempotency(active: &ActiveVMCreateIdempotency) -> CliResult<()> { + let path = vm_create_idempotency_store_url()?; + let mut store = load_vm_create_idempotency_store(&path); + if store + .records + .get(&active.signature) + .map(|record| record.key.as_str()) + == Some(active.key.as_str()) + { + store.records.remove(&active.signature); + save_vm_create_idempotency_store(&store, &path)?; + } + Ok(()) +} + +fn unix_timestamp_secs() -> u64 { + SystemTime::now() + .duration_since(UNIX_EPOCH) + .map(|duration| duration.as_secs()) + .unwrap_or(0) +} + +fn random_uuid_like() -> String { + let mut bytes = [0_u8; 16]; + if fs::File::open("/dev/urandom") + .and_then(|mut file| file.read_exact(&mut bytes)) + .is_err() + { + let fallback = request_id(); + for (index, byte) in fallback.as_bytes().iter().take(16).enumerate() { + bytes[index] = *byte; + } + } + bytes[6] = (bytes[6] & 0x0f) | 0x40; + bytes[8] = (bytes[8] & 0x3f) | 0x80; + format!( + "{:02x}{:02x}{:02x}{:02x}-{:02x}{:02x}-{:02x}{:02x}-{:02x}{:02x}-{:02x}{:02x}{:02x}{:02x}{:02x}{:02x}", + bytes[0], + bytes[1], + bytes[2], + bytes[3], + bytes[4], + bytes[5], + bytes[6], + bytes[7], + bytes[8], + bytes[9], + bytes[10], + bytes[11], + bytes[12], + bytes[13], + bytes[14], + bytes[15] + ) +} + +fn shell_quote(value: &str) -> String { + if !value.is_empty() + && value + .chars() + .all(|ch| ch.is_ascii_alphanumeric() || "_@%+=:,./-".contains(ch)) + { + return value.to_string(); + } + format!("'{}'", value.replace('\'', "'\"'\"'")) +} + +fn value_str(value: Option<&Value>) -> Option<&str> { + value.and_then(Value::as_str) +} + +fn safe_v2_details(value: Option<&Value>) -> Option { + let Some(value) = value else { + return None; + }; + if let Some(text) = value + .as_str() + .map(str::trim) + .filter(|value| !value.is_empty()) + { + return Some(text.to_string()); + } + let object = value.as_object()?; + let allowed = [ + "amount", + "code", + "duration", + "durationMs", + "field", + "idempotencyKeySet", + "imageRequested", + "limit", + "operation", + "retryable", + "status", + "type", + "vmId", + ]; + let mut lines = Vec::new(); + for key in allowed { + if let Some(value) = object.get(key).filter(|value| !value.is_null()) { + lines.push(format!("{key}: {}", safe_v2_detail_value(value))); + } + } + if lines.is_empty() { + None + } else { + Some(lines.join("\n")) + } +} + +fn safe_v2_detail_value(value: &Value) -> String { + match value { + Value::String(value) => value.replace('\n', "\\n").replace('\r', "\\r"), + Value::Bool(value) => value.to_string(), + Value::Number(value) => value.to_string(), + _ => "".to_string(), + } +} + +fn format_v2_error( + code: &str, + message: &str, + action: Option<&str>, + reason: Option<&str>, + details: Option<&str>, +) -> String { + let header = if code == "vm_error" { + message.to_string() + } else if message.contains('\n') { + format!("{code}:\n{message}") + } else { + format!("{code}: {message}") + }; + let mut sections = vec![header]; + if let Some(reason) = trimmed_non_empty(reason) { + sections.push(format!("Reason:\n{}", indent_v2_error_lines(reason))); + } + if let Some(action) = trimmed_non_empty(action) { + sections.push(format!("What to do:\n{}", indent_v2_error_lines(action))); + } + if let Some(details) = trimmed_non_empty(details) { + sections.push(format!("Details:\n{}", indent_v2_error_lines(details))); + } + sections.join("\n\n") +} + +fn trimmed_non_empty(value: Option<&str>) -> Option<&str> { + value.map(str::trim).filter(|value| !value.is_empty()) +} + +fn indent_v2_error_lines(value: &str) -> String { + value + .lines() + .map(|line| format!(" {line}")) + .collect::>() + .join("\n") +} + +fn is_uuid(value: &str) -> bool { + let bytes = value.as_bytes(); + bytes.len() == 36 + && [8, 13, 18, 23].iter().all(|index| bytes[*index] == b'-') + && bytes + .iter() + .enumerate() + .all(|(index, byte)| [8, 13, 18, 23].contains(&index) || byte.is_ascii_hexdigit()) +} + +fn is_handle_ref(value: &str) -> bool { + let Some((kind, index)) = value.split_once(':') else { + return false; + }; + matches!( + kind.to_ascii_lowercase().as_str(), + "window" | "workspace" | "pane" | "surface" + ) && index.parse::().is_ok() +} + +fn normalized_handle_value(raw: Option<&str>) -> Option { + raw.map(str::trim) + .filter(|value| !value.is_empty()) + .map(ToString::to_string) +} + +fn handles_match(lhs: &str, rhs: &str) -> bool { + lhs.eq_ignore_ascii_case(rhs) +} + +fn usage() -> &'static str { + "Usage: cmux cloud [args...]\n\nManage cloud VMs. Requires `cmux auth login`.\n\nSubcommands:\n ls List your cloud VMs.\n new [--image