diff --git a/AGENTS.md b/AGENTS.md index 6bd32ff2..16d4872a 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -12,7 +12,7 @@ The user-facing CLI surface. Contains all logic for local commands, wraps `click - Cloud handlers go through the `CloudClient` wrapper (`src/cloud/client.rs`), not `clickhouse_cloud_api::Client` directly. The wrapper handles credential precedence, error conversion, and response unwrapping. - Cloud handlers always support `--json` output unless there is good reason not to. JSON is emitted automatically when `--json` is passed or a coding agent is detected (`is_ai_agent::detect()` via the `json_output()` helper in `main.rs`). -- `CloudError` carries a `kind: CloudErrorKind` (`Auth` for 401/403 and missing credentials, else `Generic`). It maps to `Error::AuthRequired` / `Error::Cloud` in `main.rs`. Dispatched commands exit with `0` on success or use `Error::exit_code()` for failures: `1` error, `3` cancelled, `4` auth required. Clap uses `2` for usage errors. +- `CloudError` carries a `kind: CloudErrorKind` (`Auth` for 401/403 and missing credentials, else `Generic`). It maps to `Error::AuthRequired` / `Error::Cloud` in `cloud::run`. Dispatched commands exit with `0` on success or use `Error::exit_code()` for failures: `1` error, `3` cancelled, `4` auth required. Clap uses `2` for usage errors. Use `--help` to learn the current command surface. @@ -22,7 +22,7 @@ The CLI does not need to have 100% coverage of endpoints exposed by the API libr #### Adding a command -For both local and cloud commands, define the clap variant in the appropriate `cli.rs`, then wire dispatch in `src/main.rs`. +For both local and cloud commands, define the clap variant in the appropriate `cli.rs`, then wire dispatch in the owning runtime module. **Local subcommand:** @@ -35,7 +35,7 @@ For both local and cloud commands, define the clap variant in the appropriate `c 1. Make sure `clickhouse-cloud-api` has already been updated to support necessary endpoints & models. 2. Add the variant to the relevant sub-enum in `src/cloud/cli.rs` (or `src/cloud/postgres.rs` for Postgres). Create a new sub-enum if the surface warrants its own grouping. 3. Classify the new variant in `CloudCommands::is_write_command()` in `src/cloud/cli.rs` (Postgres variants go in the equivalent `is_write()` on the Postgres enum). OAuth (Bearer) auth is read-only; write commands require API key auth and we fail fast on OAuth + write. The match has no wildcards, so the compiler will reject a missing arm — but you still need to make the read/write call deliberately, and add a case to both the `is_write_command_read_only_commands` and `is_write_command_destructive_commands` tests. -4. Add the match arm in `run_cloud()` in `src/main.rs`. +4. Add non-Postgres match arms to `cloud::dispatch()` in `src/cloud/mod.rs`. Postgres dispatch belongs in `src/cloud/postgres.rs::run`; if another domain later gains its own runtime dispatcher, have `cloud::dispatch()` delegate to it rather than keeping that domain's match arms centrally. 5. Add a thin wrapper method on `CloudClient` in `src/cloud/client.rs`. It should delegate to `self.api().()`, map errors via `self.convert_error(e)`, and unwrap with `Self::unwrap_response`. Use the library's request/response types here. 6. If the command sends a request body, extract a `build__request(...)` helper in `src/cloud/commands.rs` that returns the library's request struct. Cover the helper with minimal + maximal unit tests in the `mod tests` block at the bottom of `commands.rs`, asserting directly on library struct fields. 7. Implement the handler in `src/cloud/commands.rs`. For body-sending commands the handler calls the build helper, passes the result through the `CloudClient` wrapper, and prints with the `--json` output pattern. For detail/get views (rendering a single resource), drive human output through `print_human` so it shares serde's behaviour — including deprecated-field hiding — instead of hand-writing `println!` lines: diff --git a/crates/clickhousectl/src/cli.rs b/crates/clickhousectl/src/cli.rs index 8710f315..f020f77b 100644 --- a/crates/clickhousectl/src/cli.rs +++ b/crates/clickhousectl/src/cli.rs @@ -1,15 +1,6 @@ use clap::{Args, Parser, Subcommand}; -pub use crate::cloud::cli::{ - ActivityCommands, AuthCommands, BackupCommands, BackupConfigCommands, ClickPipeCommands, - ClickPipeCreateCommands, ClickPipeSettingsCommands, CloudArgs, CloudCommands, - InvitationCommands, KeyCommands, MemberCommands, OrgCommands, PrivateEndpointCommands, - QueryEndpointCommands, ServiceCommands, -}; -pub use crate::cloud::postgres::{ - CertsCommands as PostgresCertsCommands, ConfigCommands as PostgresConfigCommands, - PostgresCommands, ReadReplicaCommands as PostgresReadReplicaCommands, -}; +use crate::cloud::cli::CloudArgs; pub use crate::local::cli::LocalArgs; #[derive(Parser)] diff --git a/crates/clickhousectl/src/cloud/cli.rs b/crates/clickhousectl/src/cloud/cli.rs index 51698640..ef7a8f26 100644 --- a/crates/clickhousectl/src/cloud/cli.rs +++ b/crates/clickhousectl/src/cloud/cli.rs @@ -1,4 +1,4 @@ -use chrono::{DateTime, FixedOffset, NaiveDate, NaiveTime}; +use crate::cloud::shared::{parse_date_only, parse_datetime, parse_time_only}; use clap::builder::PossibleValuesParser; use clap::{Args, Subcommand}; @@ -74,32 +74,6 @@ const MONGODB_READ_PREFERENCES: &[&str] = &[ "nearest", ]; -fn parse_date_only(value: &str) -> Result { - if NaiveDate::parse_from_str(value, "%Y-%m-%d").is_err() { - return Err(format!("invalid date '{}': expected YYYY-MM-DD", value)); - } - - Ok(value.to_string()) -} - -pub(super) fn parse_datetime(value: &str) -> Result { - if DateTime::::parse_from_rfc3339(value).is_err() { - return Err(format!( - "invalid datetime '{}': expected ISO 8601 / RFC 3339", - value - )); - } - - Ok(value.to_string()) -} - -fn parse_time_only(value: &str) -> Result { - if NaiveTime::parse_from_str(value, "%H:%M").is_err() { - return Err(format!("invalid time '{}': expected HH:MM", value)); - } - - Ok(value.to_string()) -} #[derive(Subcommand)] pub enum AuthCommands { /// Log in to ClickHouse Cloud diff --git a/crates/clickhousectl/src/cloud/commands.rs b/crates/clickhousectl/src/cloud/commands.rs index 12ea4955..072dc99e 100644 --- a/crates/clickhousectl/src/cloud/commands.rs +++ b/crates/clickhousectl/src/cloud/commands.rs @@ -1,19 +1,19 @@ use crate::cloud::client::{CloudClient, CloudError}; use crate::cloud::credentials; use crate::cloud::output::{ABSENT, or_absent, print_human}; +use crate::cloud::shared::{parse_serde_enum, parse_tags, resolve_org_id}; use clickhouse_cloud_api::models::{ ApiKeyPatchRequest, ApiKeyPatchRequestState, ApiKeyPostRequest, ApiKeyPostRequestState, AutoscalingMode, BackupConfigurationPatchRequest, InstancePrivateEndpointsPatch, InstanceServiceQueryApiEndpointsPostRequest, InstanceTagsPatch, IpAccessListEntry, IpAccessListPatch, OrganizationPatchPrivateEndpoint, OrganizationPatchPrivateEndpointCloudprovider, OrganizationPatchPrivateEndpointRegion, - OrganizationPatchRequest, OrganizationPrivateEndpointsPatch, ResourceTagsV1, - ServicPrivateEndpointePostRequest, Service, ServiceEndpoint, ServiceEndpointChange, - ServiceEndpointChangeProtocol, ServicePasswordPatchRequest, ServicePatchRequest, - ServicePatchRequestReleasechannel, ServicePostRequest, ServicePostRequestCompliancetype, - ServicePostRequestProfile, ServicePostRequestProvider, ServicePostRequestRegion, - ServicePostRequestReleasechannel, ServiceReplicaScalingPatchRequest, ServiceState, - ServiceStatePatchRequestCommand, + OrganizationPatchRequest, OrganizationPrivateEndpointsPatch, ServicPrivateEndpointePostRequest, + Service, ServiceEndpoint, ServiceEndpointChange, ServiceEndpointChangeProtocol, + ServicePasswordPatchRequest, ServicePatchRequest, ServicePatchRequestReleasechannel, + ServicePostRequest, ServicePostRequestCompliancetype, ServicePostRequestProfile, + ServicePostRequestProvider, ServicePostRequestRegion, ServicePostRequestReleasechannel, + ServiceReplicaScalingPatchRequest, ServiceState, ServiceStatePatchRequestCommand, }; use std::io::{IsTerminal, Write}; use tabled::{Table, Tabled, settings::Style}; @@ -46,17 +46,6 @@ fn first_endpoint(endpoints: Option<&[ServiceEndpoint]>) -> String { .unwrap_or_else(|| ABSENT.to_string()) } -/// Resolve org ID from explicit arg or auto-detect -pub(super) async fn resolve_org_id( - client: &CloudClient, - org_id: Option<&str>, -) -> Result> { - match org_id { - Some(id) => Ok(id.to_string()), - None => Ok(client.get_default_org_id().await?), - } -} - /// Resolve a service by name or ID within the given org. /// Exactly one of `name` or `id` must be provided. async fn resolve_service( @@ -88,69 +77,6 @@ async fn resolve_service( } } -/// Parse a string into a library enum via serde deserialization, with client-side -/// validation against a known-values list. Library enums have an `Unknown(String)` -/// catch-all that prevents serde from ever failing, so we validate first. -pub(super) fn parse_serde_enum( - value: &str, - field: &str, - known_values: &[&str], -) -> Result> { - if !known_values.contains(&value) { - return Err(format!( - "invalid {}: unknown value '{}', expected one of: {}", - field, - value, - known_values.join(", ") - ) - .into()); - } - serde_json::from_value(serde_json::Value::String(value.to_string())) - .map_err(|e| format!("invalid {}: {}", field, e).into()) -} - -pub(super) fn parse_tag(value: &str) -> Result> { - match value.split_once('=') { - Some((key, tag_value)) => { - let key = key.trim(); - if key.is_empty() { - Err(format!("invalid tag '{}': tag key cannot be empty", value).into()) - } else { - Ok(ResourceTagsV1 { - key: key.to_string(), - value: Some(tag_value.to_string()), - }) - } - } - None => { - let key = value.trim(); - if key.is_empty() { - Err(format!("invalid tag '{}': tag key cannot be empty", value).into()) - } else { - Ok(ResourceTagsV1 { - key: key.to_string(), - value: None, - }) - } - } - } -} - -pub(super) fn parse_tags( - values: &[String], -) -> Result>, Box> { - if values.is_empty() { - Ok(None) - } else { - Ok(Some( - values - .iter() - .map(|value| parse_tag(value)) - .collect::, _>>()?, - )) - } -} - fn parse_ip_access_entries(values: &[String]) -> Option> { (!values.is_empty()).then(|| { values @@ -3869,21 +3795,6 @@ mod tests { ); } - #[test] - fn parse_tag_rejects_empty_keys() { - let err = parse_tag("=value").unwrap_err(); - assert_eq!( - err.to_string(), - "invalid tag '=value': tag key cannot be empty" - ); - - let err = parse_tag(" ").unwrap_err(); - assert_eq!( - err.to_string(), - "invalid tag ' ': tag key cannot be empty" - ); - } - #[test] fn build_create_service_request_supports_ga_optional_fields() { let opts = CreateServiceOptions { diff --git a/crates/clickhousectl/src/cloud/mod.rs b/crates/clickhousectl/src/cloud/mod.rs index 7f031890..d0e02fb7 100644 --- a/crates/clickhousectl/src/cloud/mod.rs +++ b/crates/clickhousectl/src/cloud/mod.rs @@ -6,6 +6,7 @@ pub mod credentials; pub mod output; pub mod postgres; pub mod service_query; +mod shared; pub mod types; #[cfg(test)] @@ -15,3 +16,1033 @@ pub use client::{ AuthSource, CloudClient, CloudError, CloudErrorKind, EnvCredPresence, dotenv_env_provenance, env_cred_presence, resolve_active_auth_source, }; + +use crate::error::{Error, Result}; +use cli::{ + ActivityCommands, AuthCommands, BackupCommands, BackupConfigCommands, ClickPipeCommands, + ClickPipeCreateCommands, ClickPipeSettingsCommands, CloudArgs, CloudCommands, + InvitationCommands, KeyCommands, MemberCommands, OrgCommands, PrivateEndpointCommands, + QueryEndpointCommands, ServiceCommands, +}; + +/// Explain when a configured environment credential cannot participate in +/// authentication because a higher-precedence source won. Keep this notice on +/// stderr so it does not contaminate JSON output. +fn ignored_env_credentials_notice( + active: AuthSource, + env_creds: EnvCredPresence, +) -> Option { + let configured = match (env_creds.key, env_creds.secret) { + (true, true) => { + "CLICKHOUSE_CLOUD_API_KEY and CLICKHOUSE_CLOUD_API_SECRET are set but ignored" + } + (true, false) => "CLICKHOUSE_CLOUD_API_KEY is set but ignored", + (false, true) => "CLICKHOUSE_CLOUD_API_SECRET is set but ignored", + (false, false) => return None, + }; + let winner = match active { + AuthSource::CliFlags => "CLI flags", + AuthSource::CredentialsFile => "credentials file", + AuthSource::EnvVars | AuthSource::OAuthTokens => return None, + }; + + Some(format!("note: {configured}; using {winner} — see --debug")) +} + +pub async fn run(args: CloudArgs, json: bool) -> Result<()> { + // Auth subcommands don't need a client. + if let CloudCommands::Auth { command } = args.command { + return match command { + AuthCommands::Login { + interactive, + api_key, + api_secret, + } => { + if interactive { + commands::auth_interactive().map_err(|e| Error::Cloud(e.to_string())) + } else if api_key.is_some() || api_secret.is_some() { + let key = api_key.ok_or_else(|| { + Error::AuthRequired( + "--api-key is required when --api-secret is provided".into(), + ) + })?; + let secret = api_secret.ok_or_else(|| { + Error::AuthRequired( + "--api-secret is required when --api-key is provided".into(), + ) + })?; + let mut creds = credentials::load_credentials().unwrap_or_default(); + creds.api_key = Some(key); + creds.api_secret = Some(secret); + credentials::save_credentials(&creds) + .map_err(|e| Error::Cloud(e.to_string()))?; + println!( + "Credentials saved to {}", + credentials::credentials_path().display() + ); + Ok(()) + } else { + let url = args + .url + .as_deref() + .unwrap_or("https://api.clickhouse.cloud"); + let tokens = auth::device_auth_login(url) + .await + .map_err(|e| Error::Cloud(e.to_string()))?; + auth::save_tokens(&tokens).map_err(|e| Error::Cloud(e.to_string()))?; + println!("Logged in successfully."); + let tokens_path = + auth::tokens_path().map_err(|e| Error::Cloud(e.to_string()))?; + println!("Tokens saved to {}", tokens_path.display()); + Ok(()) + } + } + AuthCommands::Signup => { + let api_url = args + .url + .as_deref() + .unwrap_or("https://api.clickhouse.cloud"); + let parsed = url::Url::parse(api_url) + .map_err(|e| Error::Cloud(format!("Invalid URL: {}", e)))?; + let host = parsed.host_str().unwrap_or("api.clickhouse.cloud"); + let base_host = host.strip_prefix("api.").unwrap_or(host); + let url = format!( + "https://console.{}/signUp?utm_source=clickhousectl", + base_host + ); + println!("Opening ClickHouse Cloud sign-up page..."); + if open::that(&url).is_err() { + println!("Could not open browser. Please visit: {}", url); + } + Ok(()) + } + AuthCommands::Logout { oauth, api_keys } => { + match (oauth, api_keys) { + (true, false) => { + auth::clear_tokens(); + println!("OAuth tokens cleared. API keys unchanged."); + } + (false, true) => { + credentials::clear_credentials(); + println!("API keys cleared. OAuth tokens unchanged."); + } + _ => { + auth::clear_tokens(); + credentials::clear_credentials(); + println!("Logged out. All saved credentials cleared."); + } + } + Ok(()) + } + AuthCommands::Status => { + use serde::Serialize; + use tabled::{Table, Tabled, settings::Style}; + + #[derive(Serialize, Tabled)] + struct AuthRow { + #[tabled(rename = "Type")] + #[serde(rename = "type")] + auth_type: String, + #[tabled(rename = "Status")] + status: String, + #[tabled(rename = "Scope")] + scope: String, + #[tabled(rename = "Active")] + active: String, + } + + // Determine which source would actually win precedence right now. + // CLI --api-key/--api-secret aren't relevant to `auth status` itself. + let active = resolve_active_auth_source(); + let mark = |src: AuthSource| -> String { + if active == Some(src) { + "yes".into() + } else { + "-".into() + } + }; + + let mut rows = Vec::new(); + + match auth::load_tokens() { + Some(tokens) if auth::is_token_valid(&tokens) => { + rows.push(AuthRow { + auth_type: "OAuth".into(), + status: "Active".into(), + scope: "read-only".into(), + active: mark(AuthSource::OAuthTokens), + }); + } + Some(_) => { + rows.push(AuthRow { + auth_type: "OAuth".into(), + status: "Expired".into(), + scope: "read-only".into(), + active: "-".into(), + }); + } + None => { + rows.push(AuthRow { + auth_type: "OAuth".into(), + status: "Not configured".into(), + scope: "-".into(), + active: "-".into(), + }); + } + } + + if credentials::load_credentials().is_some() { + rows.push(AuthRow { + auth_type: "API key".into(), + status: "Active".into(), + scope: "read/write".into(), + active: mark(AuthSource::CredentialsFile), + }); + } else { + rows.push(AuthRow { + auth_type: "API key".into(), + status: "Not configured".into(), + scope: "-".into(), + active: "-".into(), + }); + } + + // Presence is computed through the same `env_or_dotenv` merge + // the resolver uses (shell env with `.env` fallback, empties + // treated as absent) so this table can't disagree with which + // source actually wins. + let env_creds = env_cred_presence(); + let has_key = env_creds.key; + let has_secret = env_creds.secret; + + match (has_key, has_secret) { + (true, true) => { + // Only label the `.env` path when BOTH credentials + // come exclusively from it — otherwise the status + // would imply the file was the source even though + // one value is actually exported in the shell. Use + // the same rule as `dotenv_env_provenance()` so the + // table and `--debug describe()` stay consistent. + let provenance = dotenv_env_provenance() + .map(|path| format!(" (from {})", path.display())) + .unwrap_or_default(); + let status = match active { + Some(AuthSource::EnvVars) => format!("Active{provenance}"), + Some(AuthSource::CredentialsFile) => format!( + "Configured{provenance} (inactive, outranked by credentials file)" + ), + _ => format!("Configured{provenance} (inactive)"), + }; + rows.push(AuthRow { + auth_type: "Env vars".into(), + status, + scope: "read/write".into(), + active: mark(AuthSource::EnvVars), + }); + } + (true, false) => { + rows.push(AuthRow { + auth_type: "Env vars".into(), + status: "Incomplete (missing CLICKHOUSE_CLOUD_API_SECRET)".into(), + scope: "-".into(), + active: "-".into(), + }); + } + (false, true) => { + rows.push(AuthRow { + auth_type: "Env vars".into(), + status: "Incomplete (missing CLICKHOUSE_CLOUD_API_KEY)".into(), + scope: "-".into(), + active: "-".into(), + }); + } + (false, false) => { + rows.push(AuthRow { + auth_type: "Env vars".into(), + status: "Not configured".into(), + scope: "-".into(), + active: "-".into(), + }); + } + } + + if args.debug { + match active { + Some(src) => eprintln!("[debug] auth source: {}", src.describe()), + None => eprintln!("[debug] auth source: none (no credentials configured)"), + } + } + + if json { + println!("{}", serde_json::to_string_pretty(&rows)?); + } else { + println!("{}", Table::new(rows).with(Style::markdown())); + } + Ok(()) + } + }; + } + + // Refresh OAuth tokens if needed. Errors here are filesystem failures + // (refresh-rpc failures are swallowed and tokens cleared), so this stays + // a generic error rather than `AuthRequired`. + auth::ensure_fresh_tokens() + .await + .map_err(|e| Error::Cloud(e.to_string()))?; + + let client = CloudClient::new( + args.api_key.as_deref(), + args.api_secret.as_deref(), + args.url.as_deref(), + ) + .map_err(cloud_error_to_top_level)?; + + if let Some(notice) = ignored_env_credentials_notice(client.auth_source(), env_cred_presence()) + { + eprintln!("{notice}"); + } + + if args.debug { + eprintln!("[debug] auth source: {}", client.auth_source().describe()); + eprintln!("[debug] api url: {}", client.base_url()); + } + + // OAuth (Bearer) tokens are read-only. Block write commands early + // to avoid fail loops where agents repeatedly hit 403 errors. + if client.is_bearer_auth() && args.command.is_write_command() { + return Err(Error::AuthRequired( + "This command requires API key authentication. \ + OAuth (browser login) provides read-only access.\n\n\ + To authenticate with an API key:\n \ + clickhousectl cloud auth login --api-key YOUR_KEY --api-secret YOUR_SECRET\n\n\ + Or set environment variables:\n \ + export CLICKHOUSE_CLOUD_API_KEY=your-key\n \ + export CLICKHOUSE_CLOUD_API_SECRET=your-secret\n\n\ + Learn how to create API keys:\n \ + https://clickhouse.com/docs/cloud/manage/openapi?referrer=clickhousectl" + .into(), + )); + } + + dispatch(&client, args.command, json) + .await + .map_err(boxed_cloud_error_to_top_level) +} + +fn cloud_error_to_top_level(e: CloudError) -> Error { + match e.kind { + CloudErrorKind::Auth => Error::AuthRequired(e.message), + CloudErrorKind::Generic => Error::Cloud(e.message), + } +} + +// Cloud command fns return `Box`, so the `CloudError.kind` +// only survives via downcast — without it, auth-flagged errors silently fall back +// to `Error::Cloud` (exit 1) instead of `Error::AuthRequired` (exit 4). +fn boxed_cloud_error_to_top_level(e: Box) -> Error { + match e.downcast::() { + Ok(ce) => cloud_error_to_top_level(*ce), + Err(other) => Error::Cloud(other.to_string()), + } +} + +async fn dispatch( + client: &CloudClient, + command: CloudCommands, + json: bool, +) -> std::result::Result<(), Box> { + match command { + CloudCommands::Auth { .. } => unreachable!("handled above"), + CloudCommands::Org { command } => match command { + OrgCommands::List => commands::org_list(client, json).await, + OrgCommands::Get { org_id } => commands::org_get(client, &org_id, json).await, + OrgCommands::Update { + org_id, + name, + remove_private_endpoint, + enable_core_dumps, + } => { + let opts = commands::OrgUpdateOptions { + name, + remove_private_endpoints: remove_private_endpoint, + enable_core_dumps, + }; + commands::org_update(client, &org_id, opts, json).await + } + OrgCommands::Prometheus { + org_id, + legacy_org_id, + filtered_metrics, + } => { + let org_id = org_id.as_deref().or(legacy_org_id.as_deref()); + commands::org_prometheus(client, org_id, filtered_metrics, json).await + } + OrgCommands::Usage { + org_id, + legacy_org_id, + from_date, + to_date, + filter, + } => { + let org_id = org_id.as_deref().or(legacy_org_id.as_deref()); + commands::org_usage(client, org_id, &from_date, &to_date, &filter, json).await + } + }, + CloudCommands::Service { command } => match command { + ServiceCommands::List { org_id, filter } => { + commands::service_list(client, org_id.as_deref(), &filter, json).await + } + ServiceCommands::Get { service_id, org_id } => { + commands::service_get(client, &service_id, org_id.as_deref(), json).await + } + ServiceCommands::Create { + name, + provider, + region, + min_replica_memory_gb, + max_replica_memory_gb, + num_replicas, + min_replicas, + max_replicas, + autoscaling_mode, + idle_scaling, + idle_timeout_minutes, + ip_allow, + backup_id, + release_channel, + data_warehouse_id, + readonly, + encryption_key, + encryption_role, + enable_tde, + compliance_type, + profile, + tag, + enable_endpoint, + disable_endpoint, + private_preview_terms_checked, + enable_core_dumps, + org_id, + } => { + let opts = commands::CreateServiceOptions { + name, + provider, + region, + min_replica_memory_gb, + max_replica_memory_gb, + num_replicas, + min_replicas, + max_replicas, + autoscaling_mode, + idle_scaling, + idle_timeout_minutes, + ip_allow, + backup_id, + release_channel, + data_warehouse_id, + is_readonly: readonly, + encryption_key, + encryption_role, + enable_tde, + compliance_type, + profile, + tags: tag, + enable_endpoints: enable_endpoint, + disable_endpoints: disable_endpoint, + private_preview_terms_checked, + enable_core_dumps, + org_id, + }; + commands::service_create(client, opts, json).await + } + ServiceCommands::Delete { + service_id, + force, + org_id, + } => { + commands::service_delete(client, &service_id, force, org_id.as_deref(), json).await + } + ServiceCommands::Start { service_id, org_id } => { + commands::service_start(client, &service_id, org_id.as_deref(), json).await + } + ServiceCommands::Stop { service_id, org_id } => { + commands::service_stop(client, &service_id, org_id.as_deref(), json).await + } + ServiceCommands::Update { + service_id, + name, + add_ip_allow, + remove_ip_allow, + add_private_endpoint_id, + remove_private_endpoint_id, + release_channel, + enable_endpoint, + disable_endpoint, + transparent_data_encryption_key_id, + add_tag, + remove_tag, + enable_core_dumps, + org_id, + } => { + let opts = commands::ServiceUpdateOptions { + name, + add_ip_allow, + remove_ip_allow, + add_private_endpoint_ids: add_private_endpoint_id, + remove_private_endpoint_ids: remove_private_endpoint_id, + release_channel, + enable_endpoints: enable_endpoint, + disable_endpoints: disable_endpoint, + transparent_data_encryption_key_id, + add_tags: add_tag, + remove_tags: remove_tag, + enable_core_dumps, + org_id, + }; + commands::service_update(client, &service_id, opts, json).await + } + ServiceCommands::Scale { + service_id, + min_replica_memory_gb, + max_replica_memory_gb, + num_replicas, + min_replicas, + max_replicas, + autoscaling_mode, + idle_scaling, + idle_timeout_minutes, + org_id, + } => { + commands::service_scale( + client, + &service_id, + commands::ServiceScaleOptions { + min_replica_memory_gb, + max_replica_memory_gb, + num_replicas, + min_replicas, + max_replicas, + autoscaling_mode, + idle_scaling, + idle_timeout_minutes, + org_id, + }, + json, + ) + .await + } + ServiceCommands::ResetPassword { + service_id, + new_password_hash, + new_double_sha1_hash, + org_id, + } => { + let opts = commands::ServiceResetPasswordOptions { + new_password_hash, + new_double_sha1_hash, + org_id, + }; + commands::service_reset_password(client, &service_id, opts, json).await + } + ServiceCommands::QueryEndpoint { command } => match command { + QueryEndpointCommands::Get { service_id, org_id } => { + commands::query_endpoint_get(client, &service_id, org_id.as_deref(), json).await + } + QueryEndpointCommands::Create { + service_id, + role, + open_api_key, + allowed_origins, + org_id, + } => { + let opts = commands::QueryEndpointCreateOptions { + roles: role, + open_api_keys: open_api_key, + allowed_origins, + org_id, + }; + commands::query_endpoint_create(client, &service_id, opts, json).await + } + QueryEndpointCommands::Delete { service_id, org_id } => { + commands::query_endpoint_delete(client, &service_id, org_id.as_deref(), json) + .await + } + }, + ServiceCommands::PrivateEndpoint { command } => match command { + PrivateEndpointCommands::Create { + service_id, + endpoint_id, + description, + org_id, + } => { + commands::private_endpoint_create( + client, + &service_id, + &endpoint_id, + description.as_deref(), + org_id.as_deref(), + json, + ) + .await + } + PrivateEndpointCommands::GetConfig { service_id, org_id } => { + commands::private_endpoint_get_config( + client, + &service_id, + org_id.as_deref(), + json, + ) + .await + } + }, + ServiceCommands::BackupConfig { command } => match command { + BackupConfigCommands::Get { service_id, org_id } => { + commands::backup_config_get(client, &service_id, org_id.as_deref(), json).await + } + BackupConfigCommands::Update { + service_id, + backup_period_hours, + backup_retention_period_hours, + backup_start_time, + org_id, + } => { + let opts = commands::BackupConfigUpdateOptions { + backup_period_hours, + backup_retention_period_hours, + backup_start_time, + org_id, + }; + commands::backup_config_update(client, &service_id, opts, json).await + } + }, + ServiceCommands::Prometheus { + service_id, + org_id, + filtered_metrics, + } => { + commands::service_prometheus( + client, + &service_id, + org_id.as_deref(), + filtered_metrics, + ) + .await + } + ServiceCommands::Query { + name, + id, + query, + queries_file, + database, + format, + org_id, + no_auto_enable, + } => { + let opts = commands::ServiceQueryOptions { + name, + id, + query, + queries_file, + database, + format, + json, + org_id, + no_auto_enable, + }; + commands::service_query(client, opts).await + } + }, + CloudCommands::Member { command } => match command { + MemberCommands::List { org_id } => { + commands::member_list(client, org_id.as_deref(), json).await + } + MemberCommands::Get { user_id, org_id } => { + commands::member_get(client, &user_id, org_id.as_deref(), json).await + } + MemberCommands::Update { + user_id, + role_id, + org_id, + } => commands::member_update(client, &user_id, &role_id, org_id.as_deref(), json).await, + MemberCommands::Remove { user_id, org_id } => { + commands::member_remove(client, &user_id, org_id.as_deref(), json).await + } + }, + CloudCommands::Invitation { command } => match command { + InvitationCommands::List { org_id } => { + commands::invitation_list(client, org_id.as_deref(), json).await + } + InvitationCommands::Create { + email, + role_id, + org_id, + } => { + commands::invitation_create(client, &email, &role_id, org_id.as_deref(), json).await + } + InvitationCommands::Get { + invitation_id, + org_id, + } => commands::invitation_get(client, &invitation_id, org_id.as_deref(), json).await, + InvitationCommands::Delete { + invitation_id, + org_id, + } => commands::invitation_delete(client, &invitation_id, org_id.as_deref(), json).await, + }, + CloudCommands::Key { command } => match command { + KeyCommands::List { org_id } => { + commands::key_list(client, org_id.as_deref(), json).await + } + KeyCommands::Create { + name, + role_id, + expires_at, + state, + ip_allow, + hash_key_id, + hash_key_id_suffix, + hash_key_secret, + org_id, + } => { + let opts = commands::KeyCreateOptions { + name, + role_ids: role_id, + expires_at, + state, + ip_allow, + hash_key_id, + hash_key_id_suffix, + hash_key_secret, + org_id, + }; + commands::key_create(client, opts, json).await + } + KeyCommands::Get { key_id, org_id } => { + commands::key_get(client, &key_id, org_id.as_deref(), json).await + } + KeyCommands::Update { + key_id, + name, + role_id, + expires_at, + state, + ip_allow, + org_id, + } => { + let opts = commands::KeyUpdateOptions { + name, + role_ids: role_id, + expires_at, + state, + ip_allow, + org_id, + }; + commands::key_update(client, &key_id, opts, json).await + } + KeyCommands::Delete { key_id, org_id } => { + commands::key_delete(client, &key_id, org_id.as_deref(), json).await + } + }, + CloudCommands::Activity { command } => match command { + ActivityCommands::List { + org_id, + from_date, + to_date, + } => { + commands::activity_list( + client, + org_id.as_deref(), + from_date.as_deref(), + to_date.as_deref(), + json, + ) + .await + } + ActivityCommands::Get { + activity_id, + org_id, + } => commands::activity_get(client, &activity_id, org_id.as_deref(), json).await, + }, + CloudCommands::Backup { command } => match command { + BackupCommands::List { service_id, org_id } => { + commands::backup_list(client, &service_id, org_id.as_deref(), json).await + } + BackupCommands::Get { + service_id, + backup_id, + org_id, + } => { + commands::backup_get(client, &service_id, &backup_id, org_id.as_deref(), json).await + } + }, + CloudCommands::Postgres { command } => postgres::run(client, command, json).await, + CloudCommands::ClickPipe { command } => match *command { + ClickPipeCommands::List { service_id, org_id } => { + commands::clickpipe_list(client, &service_id, org_id.as_deref(), json).await + } + ClickPipeCommands::Get { + service_id, + clickpipe_id, + org_id, + } => { + commands::clickpipe_get(client, &service_id, &clickpipe_id, org_id.as_deref(), json) + .await + } + ClickPipeCommands::Delete { + service_id, + clickpipe_id, + org_id, + } => { + commands::clickpipe_delete( + client, + &service_id, + &clickpipe_id, + org_id.as_deref(), + json, + ) + .await + } + ClickPipeCommands::Start { + service_id, + clickpipe_id, + org_id, + } => { + commands::clickpipe_state( + client, + &service_id, + &clickpipe_id, + "start", + org_id.as_deref(), + json, + ) + .await + } + ClickPipeCommands::Stop { + service_id, + clickpipe_id, + org_id, + } => { + commands::clickpipe_state( + client, + &service_id, + &clickpipe_id, + "stop", + org_id.as_deref(), + json, + ) + .await + } + ClickPipeCommands::Resync { + service_id, + clickpipe_id, + org_id, + } => { + commands::clickpipe_state( + client, + &service_id, + &clickpipe_id, + "resync", + org_id.as_deref(), + json, + ) + .await + } + ClickPipeCommands::Scale { + service_id, + clickpipe_id, + replicas, + cpu_millicores, + memory_gb, + org_id, + } => { + commands::clickpipe_scale( + client, + &service_id, + &clickpipe_id, + replicas, + cpu_millicores, + memory_gb, + org_id.as_deref(), + json, + ) + .await + } + ClickPipeCommands::Settings { command } => match command { + ClickPipeSettingsCommands::Get { + service_id, + clickpipe_id, + org_id, + } => { + commands::clickpipe_settings_get( + client, + &service_id, + &clickpipe_id, + org_id.as_deref(), + json, + ) + .await + } + ClickPipeSettingsCommands::Update { + service_id, + clickpipe_id, + streaming_max_insert_wait_ms, + object_storage_concurrency, + object_storage_polling_interval_ms, + object_storage_max_insert_bytes, + object_storage_max_file_count, + clickhouse_max_threads, + clickhouse_max_insert_threads, + object_storage_use_cluster_function, + clickhouse_parallel_view_processing, + org_id, + } => { + commands::clickpipe_settings_update( + client, + &service_id, + &clickpipe_id, + streaming_max_insert_wait_ms, + object_storage_concurrency, + object_storage_polling_interval_ms, + object_storage_max_insert_bytes, + object_storage_max_file_count, + clickhouse_max_threads, + clickhouse_max_insert_threads, + object_storage_use_cluster_function, + clickhouse_parallel_view_processing, + org_id.as_deref(), + json, + ) + .await + } + }, + ClickPipeCommands::SchemaDiscover { + service_id, + command, + org_id, + } => { + commands::clickpipe_schema_discover( + client, + &service_id, + &command, + org_id.as_deref(), + json, + ) + .await + } + ClickPipeCommands::Create { command } => match command { + ClickPipeCreateCommands::ObjectStorage(args) => { + commands::clickpipe_create_s3(client, &args, json).await + } + ClickPipeCreateCommands::Kafka(args) => { + commands::clickpipe_create_kafka(client, &args, json).await + } + ClickPipeCreateCommands::Kinesis(args) => { + commands::clickpipe_create_kinesis(client, &args, json).await + } + ClickPipeCreateCommands::Postgres(args) => { + commands::clickpipe_create_postgres(client, &args, json).await + } + ClickPipeCreateCommands::MySQL(args) => { + commands::clickpipe_create_mysql(client, &args, json).await + } + ClickPipeCreateCommands::MongoDB(args) => { + commands::clickpipe_create_mongodb(client, &args, json).await + } + ClickPipeCreateCommands::BigQuery(args) => { + commands::clickpipe_create_bigquery(client, &args, json).await + } + }, + }, + } +} + +#[cfg(test)] +mod runtime_tests { + use super::*; + + #[test] + fn ignored_env_notice_names_the_overridden_variables_and_winner() { + assert_eq!( + ignored_env_credentials_notice( + AuthSource::CredentialsFile, + EnvCredPresence { + key: true, + secret: true, + }, + ) + .as_deref(), + Some( + "note: CLICKHOUSE_CLOUD_API_KEY and CLICKHOUSE_CLOUD_API_SECRET are set but \ + ignored; using credentials file — see --debug" + ) + ); + assert_eq!( + ignored_env_credentials_notice( + AuthSource::CliFlags, + EnvCredPresence { + key: true, + secret: false, + }, + ) + .as_deref(), + Some( + "note: CLICKHOUSE_CLOUD_API_KEY is set but ignored; using CLI flags — see --debug" + ) + ); + } + + #[test] + fn ignored_env_notice_is_absent_when_env_wins_or_is_unset() { + assert!( + ignored_env_credentials_notice( + AuthSource::EnvVars, + EnvCredPresence { + key: true, + secret: true, + }, + ) + .is_none() + ); + assert!( + ignored_env_credentials_notice( + AuthSource::CredentialsFile, + EnvCredPresence { + key: false, + secret: false, + }, + ) + .is_none() + ); + } + + #[test] + fn cloud_error_kind_routes_to_top_level() { + assert!(matches!( + cloud_error_to_top_level(CloudError::auth("nope")), + Error::AuthRequired(_) + )); + assert!(matches!( + cloud_error_to_top_level(CloudError::new("boom")), + Error::Cloud(_) + )); + assert_eq!(CloudError::new("x").kind, CloudErrorKind::Generic); + } + + #[test] + fn boxed_cloud_error_preserves_auth_kind_through_downcast() { + let boxed: Box = Box::new(CloudError::auth("nope")); + assert!(matches!( + boxed_cloud_error_to_top_level(boxed), + Error::AuthRequired(_) + )); + } + + #[test] + fn boxed_non_cloud_error_falls_back_to_generic() { + // Anything that isn't a CloudError must not downcast to AuthRequired. + let boxed: Box = "plain string error".into(); + assert!(matches!( + boxed_cloud_error_to_top_level(boxed), + Error::Cloud(_) + )); + } +} diff --git a/crates/clickhousectl/src/cloud/postgres.rs b/crates/clickhousectl/src/cloud/postgres.rs index a798525a..13b70062 100644 --- a/crates/clickhousectl/src/cloud/postgres.rs +++ b/crates/clickhousectl/src/cloud/postgres.rs @@ -1,7 +1,6 @@ -use crate::cloud::cli::parse_datetime; use crate::cloud::client::CloudClient; -use crate::cloud::commands::{parse_serde_enum, parse_tags, resolve_org_id}; use crate::cloud::output::{ABSENT, or_absent}; +use crate::cloud::shared::{parse_datetime, parse_serde_enum, parse_tags, resolve_org_id}; use clap::Subcommand; use clickhouse_cloud_api::models::{ ApiResponse, PgBouncerConfig, PgConfig, PgHaType, PgIdProperty, PgProvider, PgVersion, @@ -242,6 +241,199 @@ impl PostgresCommands { } } +pub async fn run( + client: &CloudClient, + command: PostgresCommands, + json: bool, +) -> Result<(), Box> { + match command { + PostgresCommands::List { org_id, filter } => { + postgres_list(client, org_id.as_deref(), &filter, json).await + } + PostgresCommands::Get { + postgres_id, + org_id, + } => postgres_get(client, &postgres_id, org_id.as_deref(), json).await, + PostgresCommands::Create { + name, + region, + size, + provider, + pg_version, + ha_type, + tag, + pg_config_file, + pg_bouncer_config_file, + org_id, + } => { + let opts = PostgresCreateOptions { + name: &name, + region: ®ion, + size: &size, + provider: &provider, + pg_version: pg_version.as_deref(), + ha_type: ha_type.as_deref(), + tags: &tag, + pg_config_file: pg_config_file.as_deref(), + pg_bouncer_config_file: pg_bouncer_config_file.as_deref(), + org_id: org_id.as_deref(), + }; + postgres_create(client, opts, json).await + } + PostgresCommands::Update { + postgres_id, + size, + ha_type, + add_tag, + remove_tag, + org_id, + } => { + let opts = PostgresUpdateOptions { + size: size.as_deref(), + ha_type: ha_type.as_deref(), + add_tag: &add_tag, + remove_tag: &remove_tag, + org_id: org_id.as_deref(), + }; + postgres_update(client, &postgres_id, opts, json).await + } + PostgresCommands::Delete { + postgres_id, + org_id, + } => postgres_delete(client, &postgres_id, org_id.as_deref(), json).await, + PostgresCommands::Certs(CertsCommands::Get { + postgres_id, + output, + org_id, + }) => { + postgres_certs_get( + client, + &postgres_id, + output.as_deref(), + org_id.as_deref(), + json, + ) + .await + } + PostgresCommands::Config(ConfigCommands::Get { + postgres_id, + org_id, + }) => postgres_config_get(client, &postgres_id, org_id.as_deref(), json).await, + PostgresCommands::Config(ConfigCommands::Replace { + postgres_id, + file, + org_id, + }) => postgres_config_replace(client, &postgres_id, &file, org_id.as_deref(), json).await, + PostgresCommands::Config(ConfigCommands::Patch { + postgres_id, + sets, + file, + org_id, + }) => { + postgres_config_patch( + client, + &postgres_id, + &sets, + file.as_deref(), + org_id.as_deref(), + json, + ) + .await + } + PostgresCommands::ResetPassword { + postgres_id, + password, + generate, + org_id, + } => { + postgres_reset_password( + client, + &postgres_id, + password.as_deref(), + generate, + org_id.as_deref(), + json, + ) + .await + } + PostgresCommands::ReadReplica(ReadReplicaCommands::Create { + postgres_id, + name, + tag, + pg_config_file, + pg_bouncer_config_file, + org_id, + }) => { + let opts = PostgresReadReplicaOptions { + name: &name, + tags: &tag, + pg_config_file: pg_config_file.as_deref(), + pg_bouncer_config_file: pg_bouncer_config_file.as_deref(), + org_id: org_id.as_deref(), + }; + postgres_read_replica_create(client, &postgres_id, opts, json).await + } + PostgresCommands::Restore { + postgres_id, + name, + restore_target, + tag, + pg_config_file, + pg_bouncer_config_file, + org_id, + } => { + let opts = PostgresRestoreOptions { + name: &name, + restore_target: &restore_target, + tags: &tag, + pg_config_file: pg_config_file.as_deref(), + pg_bouncer_config_file: pg_bouncer_config_file.as_deref(), + org_id: org_id.as_deref(), + }; + postgres_restore(client, &postgres_id, opts, json).await + } + PostgresCommands::Restart { + postgres_id, + org_id, + } => { + postgres_state_change( + client, + &postgres_id, + PostgresServiceSetStateCommand::Restart, + org_id.as_deref(), + json, + ) + .await + } + PostgresCommands::Promote { + postgres_id, + org_id, + } => { + postgres_state_change( + client, + &postgres_id, + PostgresServiceSetStateCommand::Promote, + org_id.as_deref(), + json, + ) + .await + } + PostgresCommands::Switchover { + postgres_id, + org_id, + } => { + postgres_state_change( + client, + &postgres_id, + PostgresServiceSetStateCommand::Switchover, + org_id.as_deref(), + json, + ) + .await + } + } +} + // --------------------------------------------------------------------------- // Helpers // --------------------------------------------------------------------------- @@ -1110,23 +1302,21 @@ pub async fn postgres_state_change( #[cfg(test)] mod tests { use super::*; - use crate::cli::{Cli, Commands}; - use crate::cloud::cli::CloudCommands; + use crate::cli::Cli; use clap::Parser; - fn parse_cloud(args: &[&str]) -> CloudCommands { - let cli = Cli::try_parse_from(args).expect("parse"); - match cli.command { - Commands::Cloud(a) => a.command, - _ => panic!("expected cloud command"), - } + #[derive(Parser)] + struct PostgresCli { + #[command(subcommand)] + command: PostgresCommands, } fn parse_postgres(args: &[&str]) -> PostgresCommands { - match parse_cloud(args) { - CloudCommands::Postgres { command } => command, - _ => panic!("expected postgres command"), - } + assert_eq!(args.get(1), Some(&"cloud")); + assert_eq!(args.get(2), Some(&"postgres")); + PostgresCli::try_parse_from(std::iter::once(args[0]).chain(args.iter().skip(3).copied())) + .expect("parse") + .command } #[test] diff --git a/crates/clickhousectl/src/cloud/shared.rs b/crates/clickhousectl/src/cloud/shared.rs new file mode 100644 index 00000000..eae5292e --- /dev/null +++ b/crates/clickhousectl/src/cloud/shared.rs @@ -0,0 +1,122 @@ +use crate::cloud::client::CloudClient; +use chrono::{DateTime, FixedOffset, NaiveDate, NaiveTime}; +use clickhouse_cloud_api::models::ResourceTagsV1; + +/// Resolve an organization ID from an explicit argument or auto-detection. +pub(super) async fn resolve_org_id( + client: &CloudClient, + org_id: Option<&str>, +) -> Result> { + match org_id { + Some(id) => Ok(id.to_string()), + None => Ok(client.get_default_org_id().await?), + } +} + +/// Parse a string into a library enum after validating its known wire values. +pub(super) fn parse_serde_enum( + value: &str, + field: &str, + known_values: &[&str], +) -> Result> { + if !known_values.contains(&value) { + return Err(format!( + "invalid {}: unknown value '{}', expected one of: {}", + field, + value, + known_values.join(", ") + ) + .into()); + } + serde_json::from_value(serde_json::Value::String(value.to_string())) + .map_err(|e| format!("invalid {}: {}", field, e).into()) +} + +pub(super) fn parse_tag(value: &str) -> Result> { + match value.split_once('=') { + Some((key, tag_value)) => { + let key = key.trim(); + if key.is_empty() { + Err(format!("invalid tag '{}': tag key cannot be empty", value).into()) + } else { + Ok(ResourceTagsV1 { + key: key.to_string(), + value: Some(tag_value.to_string()), + }) + } + } + None => { + let key = value.trim(); + if key.is_empty() { + Err(format!("invalid tag '{}': tag key cannot be empty", value).into()) + } else { + Ok(ResourceTagsV1 { + key: key.to_string(), + value: None, + }) + } + } + } +} + +pub(super) fn parse_tags( + values: &[String], +) -> Result>, Box> { + if values.is_empty() { + Ok(None) + } else { + Ok(Some( + values + .iter() + .map(|value| parse_tag(value)) + .collect::, _>>()?, + )) + } +} + +pub(super) fn parse_date_only(value: &str) -> Result { + if NaiveDate::parse_from_str(value, "%Y-%m-%d").is_err() { + return Err(format!("invalid date '{}': expected YYYY-MM-DD", value)); + } + + Ok(value.to_string()) +} + +pub(super) fn parse_datetime(value: &str) -> Result { + if DateTime::::parse_from_rfc3339(value).is_err() { + return Err(format!( + "invalid datetime '{}': expected ISO 8601 / RFC 3339", + value + )); + } + + Ok(value.to_string()) +} + +pub(super) fn parse_time_only(value: &str) -> Result { + if NaiveTime::parse_from_str(value, "%H:%M").is_err() { + return Err(format!("invalid time '{}': expected HH:MM", value)); + } + + Ok(value.to_string()) +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn parse_tag_rejects_empty_keys() { + let err = parse_tag("=value").unwrap_err(); + assert_eq!( + err.to_string(), + "invalid tag '=value': tag key cannot be empty" + ); + + let err = parse_tag(" ").unwrap_err(); + assert_eq!( + err.to_string(), + "invalid tag ' ': tag key cannot be empty" + ); + } +} diff --git a/crates/clickhousectl/src/main.rs b/crates/clickhousectl/src/main.rs index 944a57c2..6b96be35 100644 --- a/crates/clickhousectl/src/main.rs +++ b/crates/clickhousectl/src/main.rs @@ -15,15 +15,8 @@ mod version_manager; use clap::error::ErrorKind; use clap::{CommandFactory, FromArgMatches}; -use cli::{ - ActivityCommands, AuthCommands, BackupCommands, BackupConfigCommands, Cli, ClickPipeCommands, - ClickPipeCreateCommands, ClickPipeSettingsCommands, CloudArgs, CloudCommands, Commands, - InvitationCommands, KeyCommands, MemberCommands, OrgCommands, PostgresCertsCommands, - PostgresCommands, PostgresConfigCommands, PostgresReadReplicaCommands, PrivateEndpointCommands, - QueryEndpointCommands, ServiceCommands, SkillsArgs, UpdateArgs, -}; +use cli::{Cli, Commands, SkillsArgs, UpdateArgs}; -use cloud::CloudClient; use error::{Error, Result}; #[tokio::main] @@ -218,35 +211,14 @@ fn json_output(flag: bool) -> bool { flag || is_ai_agent::detect().is_some() } -/// Explain when a configured environment credential cannot participate in -/// authentication because a higher-precedence source won. Keep this notice on -/// stderr so it does not contaminate `--json` output. -fn ignored_env_credentials_notice( - active: cloud::AuthSource, - env_creds: cloud::EnvCredPresence, -) -> Option { - let configured = match (env_creds.key, env_creds.secret) { - (true, true) => { - "CLICKHOUSE_CLOUD_API_KEY and CLICKHOUSE_CLOUD_API_SECRET are set but ignored" - } - (true, false) => "CLICKHOUSE_CLOUD_API_KEY is set but ignored", - (false, true) => "CLICKHOUSE_CLOUD_API_SECRET is set but ignored", - (false, false) => return None, - }; - let winner = match active { - cloud::AuthSource::CliFlags => "CLI flags", - cloud::AuthSource::CredentialsFile => "credentials file", - cloud::AuthSource::EnvVars | cloud::AuthSource::OAuthTokens => return None, - }; - - Some(format!("note: {configured}; using {winner} — see --debug")) -} - async fn run(cmd: Commands) -> Result<()> { match cmd { Commands::Local(args) => local::run(args.command, json_output(args.json)).await, Commands::Skills(args) => run_skills(args).await, - Commands::Cloud(args) => run_cloud(*args).await, + Commands::Cloud(args) => { + let json = json_output(args.json); + cloud::run(*args, json).await + } Commands::Update(args) => run_update(args).await, #[cfg(feature = "telemetry")] Commands::Telemetry(args) => telemetry::run_command(args.command), @@ -274,1235 +246,16 @@ async fn run_skills(args: SkillsArgs) -> Result<()> { skills::install(args).await } -async fn run_cloud(args: CloudArgs) -> Result<()> { - // Auth subcommands don't need a client - if let CloudCommands::Auth { command } = args.command { - return match command { - AuthCommands::Login { - interactive, - api_key, - api_secret, - } => { - if interactive { - // Interactive prompt for API key/secret - cloud::commands::auth_interactive().map_err(|e| Error::Cloud(e.to_string())) - } else if api_key.is_some() || api_secret.is_some() { - // Non-interactive API key login - let key = api_key.ok_or_else(|| { - Error::AuthRequired( - "--api-key is required when --api-secret is provided".into(), - ) - })?; - let secret = api_secret.ok_or_else(|| { - Error::AuthRequired( - "--api-secret is required when --api-key is provided".into(), - ) - })?; - let mut creds = cloud::credentials::load_credentials().unwrap_or_default(); - creds.api_key = Some(key); - creds.api_secret = Some(secret); - cloud::credentials::save_credentials(&creds) - .map_err(|e| Error::Cloud(e.to_string()))?; - println!( - "Credentials saved to {}", - cloud::credentials::credentials_path().display() - ); - Ok(()) - } else { - // Default: OAuth device flow - let url = args - .url - .as_deref() - .unwrap_or("https://api.clickhouse.cloud"); - let tokens = cloud::auth::device_auth_login(url) - .await - .map_err(|e| Error::Cloud(e.to_string()))?; - cloud::auth::save_tokens(&tokens).map_err(|e| Error::Cloud(e.to_string()))?; - println!("Logged in successfully."); - let tokens_path = - cloud::auth::tokens_path().map_err(|e| Error::Cloud(e.to_string()))?; - println!("Tokens saved to {}", tokens_path.display()); - Ok(()) - } - } - AuthCommands::Signup => { - let api_url = args - .url - .as_deref() - .unwrap_or("https://api.clickhouse.cloud"); - let parsed = url::Url::parse(api_url) - .map_err(|e| Error::Cloud(format!("Invalid URL: {}", e)))?; - let host = parsed.host_str().unwrap_or("api.clickhouse.cloud"); - let base_host = host.strip_prefix("api.").unwrap_or(host); - let url = format!( - "https://console.{}/signUp?utm_source=clickhousectl", - base_host - ); - println!("Opening ClickHouse Cloud sign-up page..."); - if open::that(&url).is_err() { - println!("Could not open browser. Please visit: {}", url); - } - Ok(()) - } - AuthCommands::Logout { oauth, api_keys } => { - match (oauth, api_keys) { - (true, false) => { - cloud::auth::clear_tokens(); - println!("OAuth tokens cleared. API keys unchanged."); - } - (false, true) => { - cloud::credentials::clear_credentials(); - println!("API keys cleared. OAuth tokens unchanged."); - } - _ => { - cloud::auth::clear_tokens(); - cloud::credentials::clear_credentials(); - println!("Logged out. All saved credentials cleared."); - } - } - Ok(()) - } - AuthCommands::Status => { - use serde::Serialize; - use tabled::{Table, Tabled, settings::Style}; - - #[derive(Serialize, Tabled)] - struct AuthRow { - #[tabled(rename = "Type")] - #[serde(rename = "type")] - auth_type: String, - #[tabled(rename = "Status")] - status: String, - #[tabled(rename = "Scope")] - scope: String, - #[tabled(rename = "Active")] - active: String, - } - - // Determine which source would actually win precedence right now. - // CLI --api-key/--api-secret aren't relevant to `auth status` itself. - let active = cloud::resolve_active_auth_source(); - let mark = |src: cloud::AuthSource| -> String { - if active == Some(src) { - "yes".into() - } else { - "-".into() - } - }; - - let mut rows = Vec::new(); - - match cloud::auth::load_tokens() { - Some(tokens) if cloud::auth::is_token_valid(&tokens) => { - rows.push(AuthRow { - auth_type: "OAuth".into(), - status: "Active".into(), - scope: "read-only".into(), - active: mark(cloud::AuthSource::OAuthTokens), - }); - } - Some(_) => { - rows.push(AuthRow { - auth_type: "OAuth".into(), - status: "Expired".into(), - scope: "read-only".into(), - active: "-".into(), - }); - } - None => { - rows.push(AuthRow { - auth_type: "OAuth".into(), - status: "Not configured".into(), - scope: "-".into(), - active: "-".into(), - }); - } - } - - if cloud::credentials::load_credentials().is_some() { - rows.push(AuthRow { - auth_type: "API key".into(), - status: "Active".into(), - scope: "read/write".into(), - active: mark(cloud::AuthSource::CredentialsFile), - }); - } else { - rows.push(AuthRow { - auth_type: "API key".into(), - status: "Not configured".into(), - scope: "-".into(), - active: "-".into(), - }); - } - - // Presence is computed through the same `env_or_dotenv` merge - // the resolver uses (shell env with `.env` fallback, empties - // treated as absent) so this table can't disagree with which - // source actually wins. - let env_creds = cloud::env_cred_presence(); - let has_key = env_creds.key; - let has_secret = env_creds.secret; - - match (has_key, has_secret) { - (true, true) => { - // Only label the `.env` path when BOTH credentials - // come exclusively from it — otherwise the status - // would imply the file was the source even though - // one value is actually exported in the shell. Use - // the same rule as `dotenv_env_provenance()` so the - // table and `--debug describe()` stay consistent. - let provenance = cloud::dotenv_env_provenance() - .map(|path| format!(" (from {})", path.display())) - .unwrap_or_default(); - let status = match active { - Some(cloud::AuthSource::EnvVars) => format!("Active{provenance}"), - Some(cloud::AuthSource::CredentialsFile) => format!( - "Configured{provenance} (inactive, outranked by credentials file)" - ), - _ => format!("Configured{provenance} (inactive)"), - }; - rows.push(AuthRow { - auth_type: "Env vars".into(), - status, - scope: "read/write".into(), - active: mark(cloud::AuthSource::EnvVars), - }); - } - (true, false) => { - rows.push(AuthRow { - auth_type: "Env vars".into(), - status: "Incomplete (missing CLICKHOUSE_CLOUD_API_SECRET)".into(), - scope: "-".into(), - active: "-".into(), - }); - } - (false, true) => { - rows.push(AuthRow { - auth_type: "Env vars".into(), - status: "Incomplete (missing CLICKHOUSE_CLOUD_API_KEY)".into(), - scope: "-".into(), - active: "-".into(), - }); - } - (false, false) => { - rows.push(AuthRow { - auth_type: "Env vars".into(), - status: "Not configured".into(), - scope: "-".into(), - active: "-".into(), - }); - } - } - - if args.debug { - match active { - Some(src) => { - eprintln!("[debug] auth source: {}", src.describe()); - } - None => eprintln!("[debug] auth source: none (no credentials configured)"), - } - } - - if json_output(args.json) { - println!("{}", serde_json::to_string_pretty(&rows)?); - } else { - println!("{}", Table::new(rows).with(Style::markdown())); - } - Ok(()) - } - }; - } - - // Refresh OAuth tokens if needed. Errors here are filesystem failures - // (refresh-rpc failures are swallowed and tokens cleared), so this stays - // a generic error rather than `AuthRequired`. - cloud::auth::ensure_fresh_tokens() - .await - .map_err(|e| Error::Cloud(e.to_string()))?; - - let client = CloudClient::new( - args.api_key.as_deref(), - args.api_secret.as_deref(), - args.url.as_deref(), - ) - .map_err(cloud_error_to_top_level)?; - - if let Some(notice) = - ignored_env_credentials_notice(client.auth_source(), cloud::env_cred_presence()) - { - eprintln!("{notice}"); - } - - if args.debug { - eprintln!("[debug] auth source: {}", client.auth_source().describe()); - eprintln!("[debug] api url: {}", client.base_url()); - } - - // OAuth (Bearer) tokens are read-only. Block write commands early - // to avoid fail loops where agents repeatedly hit 403 errors. - if client.is_bearer_auth() && args.command.is_write_command() { - return Err(Error::AuthRequired( - "This command requires API key authentication. \ - OAuth (browser login) provides read-only access.\n\n\ - To authenticate with an API key:\n \ - clickhousectl cloud auth login --api-key YOUR_KEY --api-secret YOUR_SECRET\n\n\ - Or set environment variables:\n \ - export CLICKHOUSE_CLOUD_API_KEY=your-key\n \ - export CLICKHOUSE_CLOUD_API_SECRET=your-secret\n\n\ - Learn how to create API keys:\n \ - https://clickhouse.com/docs/cloud/manage/openapi?referrer=clickhousectl" - .into(), - )); - } - - let json = json_output(args.json); - - let result = match args.command { - CloudCommands::Auth { .. } => unreachable!("handled above"), - CloudCommands::Org { command } => match command { - OrgCommands::List => cloud::commands::org_list(&client, json).await, - OrgCommands::Get { org_id } => cloud::commands::org_get(&client, &org_id, json).await, - OrgCommands::Update { - org_id, - name, - remove_private_endpoint, - enable_core_dumps, - } => { - let opts = cloud::commands::OrgUpdateOptions { - name, - remove_private_endpoints: remove_private_endpoint, - enable_core_dumps, - }; - cloud::commands::org_update(&client, &org_id, opts, json).await - } - OrgCommands::Prometheus { - org_id, - legacy_org_id, - filtered_metrics, - } => { - let org_id = org_id.as_deref().or(legacy_org_id.as_deref()); - cloud::commands::org_prometheus(&client, org_id, filtered_metrics, json).await - } - OrgCommands::Usage { - org_id, - legacy_org_id, - from_date, - to_date, - filter, - } => { - let org_id = org_id.as_deref().or(legacy_org_id.as_deref()); - cloud::commands::org_usage(&client, org_id, &from_date, &to_date, &filter, json) - .await - } - }, - CloudCommands::Service { command } => match command { - ServiceCommands::List { org_id, filter } => { - cloud::commands::service_list(&client, org_id.as_deref(), &filter, json).await - } - ServiceCommands::Get { service_id, org_id } => { - cloud::commands::service_get(&client, &service_id, org_id.as_deref(), json).await - } - ServiceCommands::Create { - name, - provider, - region, - min_replica_memory_gb, - max_replica_memory_gb, - num_replicas, - min_replicas, - max_replicas, - autoscaling_mode, - idle_scaling, - idle_timeout_minutes, - ip_allow, - backup_id, - release_channel, - data_warehouse_id, - readonly, - encryption_key, - encryption_role, - enable_tde, - compliance_type, - profile, - tag, - enable_endpoint, - disable_endpoint, - private_preview_terms_checked, - enable_core_dumps, - org_id, - } => { - let opts = cloud::commands::CreateServiceOptions { - name, - provider, - region, - min_replica_memory_gb, - max_replica_memory_gb, - num_replicas, - min_replicas, - max_replicas, - autoscaling_mode, - idle_scaling, - idle_timeout_minutes, - ip_allow, - backup_id, - release_channel, - data_warehouse_id, - is_readonly: readonly, - encryption_key, - encryption_role, - enable_tde, - compliance_type, - profile, - tags: tag, - enable_endpoints: enable_endpoint, - disable_endpoints: disable_endpoint, - private_preview_terms_checked, - enable_core_dumps, - org_id, - }; - cloud::commands::service_create(&client, opts, json).await - } - ServiceCommands::Delete { - service_id, - force, - org_id, - } => { - cloud::commands::service_delete( - &client, - &service_id, - force, - org_id.as_deref(), - json, - ) - .await - } - ServiceCommands::Start { service_id, org_id } => { - cloud::commands::service_start(&client, &service_id, org_id.as_deref(), json).await - } - ServiceCommands::Stop { service_id, org_id } => { - cloud::commands::service_stop(&client, &service_id, org_id.as_deref(), json).await - } - ServiceCommands::Update { - service_id, - name, - add_ip_allow, - remove_ip_allow, - add_private_endpoint_id, - remove_private_endpoint_id, - release_channel, - enable_endpoint, - disable_endpoint, - transparent_data_encryption_key_id, - add_tag, - remove_tag, - enable_core_dumps, - org_id, - } => { - let opts = cloud::commands::ServiceUpdateOptions { - name, - add_ip_allow, - remove_ip_allow, - add_private_endpoint_ids: add_private_endpoint_id, - remove_private_endpoint_ids: remove_private_endpoint_id, - release_channel, - enable_endpoints: enable_endpoint, - disable_endpoints: disable_endpoint, - transparent_data_encryption_key_id, - add_tags: add_tag, - remove_tags: remove_tag, - enable_core_dumps, - org_id, - }; - cloud::commands::service_update(&client, &service_id, opts, json).await - } - ServiceCommands::Scale { - service_id, - min_replica_memory_gb, - max_replica_memory_gb, - num_replicas, - min_replicas, - max_replicas, - autoscaling_mode, - idle_scaling, - idle_timeout_minutes, - org_id, - } => { - cloud::commands::service_scale( - &client, - &service_id, - cloud::commands::ServiceScaleOptions { - min_replica_memory_gb, - max_replica_memory_gb, - num_replicas, - min_replicas, - max_replicas, - autoscaling_mode, - idle_scaling, - idle_timeout_minutes, - org_id, - }, - json, - ) - .await - } - ServiceCommands::ResetPassword { - service_id, - new_password_hash, - new_double_sha1_hash, - org_id, - } => { - let opts = cloud::commands::ServiceResetPasswordOptions { - new_password_hash, - new_double_sha1_hash, - org_id, - }; - cloud::commands::service_reset_password(&client, &service_id, opts, json).await - } - ServiceCommands::QueryEndpoint { command } => match command { - QueryEndpointCommands::Get { service_id, org_id } => { - cloud::commands::query_endpoint_get( - &client, - &service_id, - org_id.as_deref(), - json, - ) - .await - } - QueryEndpointCommands::Create { - service_id, - role, - open_api_key, - allowed_origins, - org_id, - } => { - let opts = cloud::commands::QueryEndpointCreateOptions { - roles: role, - open_api_keys: open_api_key, - allowed_origins, - org_id, - }; - cloud::commands::query_endpoint_create(&client, &service_id, opts, json).await - } - QueryEndpointCommands::Delete { service_id, org_id } => { - cloud::commands::query_endpoint_delete( - &client, - &service_id, - org_id.as_deref(), - json, - ) - .await - } - }, - ServiceCommands::PrivateEndpoint { command } => match command { - PrivateEndpointCommands::Create { - service_id, - endpoint_id, - description, - org_id, - } => { - cloud::commands::private_endpoint_create( - &client, - &service_id, - &endpoint_id, - description.as_deref(), - org_id.as_deref(), - json, - ) - .await - } - PrivateEndpointCommands::GetConfig { service_id, org_id } => { - cloud::commands::private_endpoint_get_config( - &client, - &service_id, - org_id.as_deref(), - json, - ) - .await - } - }, - ServiceCommands::BackupConfig { command } => match command { - BackupConfigCommands::Get { service_id, org_id } => { - cloud::commands::backup_config_get( - &client, - &service_id, - org_id.as_deref(), - json, - ) - .await - } - BackupConfigCommands::Update { - service_id, - backup_period_hours, - backup_retention_period_hours, - backup_start_time, - org_id, - } => { - let opts = cloud::commands::BackupConfigUpdateOptions { - backup_period_hours, - backup_retention_period_hours, - backup_start_time, - org_id, - }; - cloud::commands::backup_config_update(&client, &service_id, opts, json).await - } - }, - ServiceCommands::Prometheus { - service_id, - org_id, - filtered_metrics, - } => { - cloud::commands::service_prometheus( - &client, - &service_id, - org_id.as_deref(), - filtered_metrics, - ) - .await - } - ServiceCommands::Query { - name, - id, - query, - queries_file, - database, - format, - org_id, - no_auto_enable, - } => { - let opts = cloud::commands::ServiceQueryOptions { - name, - id, - query, - queries_file, - database, - format, - json, - org_id, - no_auto_enable, - }; - cloud::commands::service_query(&client, opts).await - } - }, - CloudCommands::Member { command } => match command { - MemberCommands::List { org_id } => { - cloud::commands::member_list(&client, org_id.as_deref(), json).await - } - MemberCommands::Get { user_id, org_id } => { - cloud::commands::member_get(&client, &user_id, org_id.as_deref(), json).await - } - MemberCommands::Update { - user_id, - role_id, - org_id, - } => { - cloud::commands::member_update(&client, &user_id, &role_id, org_id.as_deref(), json) - .await - } - MemberCommands::Remove { user_id, org_id } => { - cloud::commands::member_remove(&client, &user_id, org_id.as_deref(), json).await - } - }, - CloudCommands::Invitation { command } => match command { - InvitationCommands::List { org_id } => { - cloud::commands::invitation_list(&client, org_id.as_deref(), json).await - } - InvitationCommands::Create { - email, - role_id, - org_id, - } => { - cloud::commands::invitation_create( - &client, - &email, - &role_id, - org_id.as_deref(), - json, - ) - .await - } - InvitationCommands::Get { - invitation_id, - org_id, - } => { - cloud::commands::invitation_get(&client, &invitation_id, org_id.as_deref(), json) - .await - } - InvitationCommands::Delete { - invitation_id, - org_id, - } => { - cloud::commands::invitation_delete(&client, &invitation_id, org_id.as_deref(), json) - .await - } - }, - CloudCommands::Key { command } => match command { - KeyCommands::List { org_id } => { - cloud::commands::key_list(&client, org_id.as_deref(), json).await - } - KeyCommands::Create { - name, - role_id, - expires_at, - state, - ip_allow, - hash_key_id, - hash_key_id_suffix, - hash_key_secret, - org_id, - } => { - let opts = cloud::commands::KeyCreateOptions { - name, - role_ids: role_id, - expires_at, - state, - ip_allow, - hash_key_id, - hash_key_id_suffix, - hash_key_secret, - org_id, - }; - cloud::commands::key_create(&client, opts, json).await - } - KeyCommands::Get { key_id, org_id } => { - cloud::commands::key_get(&client, &key_id, org_id.as_deref(), json).await - } - KeyCommands::Update { - key_id, - name, - role_id, - expires_at, - state, - ip_allow, - org_id, - } => { - let opts = cloud::commands::KeyUpdateOptions { - name, - role_ids: role_id, - expires_at, - state, - ip_allow, - org_id, - }; - cloud::commands::key_update(&client, &key_id, opts, json).await - } - KeyCommands::Delete { key_id, org_id } => { - cloud::commands::key_delete(&client, &key_id, org_id.as_deref(), json).await - } - }, - CloudCommands::Activity { command } => match command { - ActivityCommands::List { - org_id, - from_date, - to_date, - } => { - cloud::commands::activity_list( - &client, - org_id.as_deref(), - from_date.as_deref(), - to_date.as_deref(), - json, - ) - .await - } - ActivityCommands::Get { - activity_id, - org_id, - } => { - cloud::commands::activity_get(&client, &activity_id, org_id.as_deref(), json).await - } - }, - CloudCommands::Backup { command } => match command { - BackupCommands::List { service_id, org_id } => { - cloud::commands::backup_list(&client, &service_id, org_id.as_deref(), json).await - } - BackupCommands::Get { - service_id, - backup_id, - org_id, - } => { - cloud::commands::backup_get( - &client, - &service_id, - &backup_id, - org_id.as_deref(), - json, - ) - .await - } - }, - CloudCommands::Postgres { command } => run_postgres(&client, command, json).await, - CloudCommands::ClickPipe { command } => match *command { - ClickPipeCommands::List { service_id, org_id } => { - cloud::commands::clickpipe_list(&client, &service_id, org_id.as_deref(), json).await - } - ClickPipeCommands::Get { - service_id, - clickpipe_id, - org_id, - } => { - cloud::commands::clickpipe_get( - &client, - &service_id, - &clickpipe_id, - org_id.as_deref(), - json, - ) - .await - } - ClickPipeCommands::Delete { - service_id, - clickpipe_id, - org_id, - } => { - cloud::commands::clickpipe_delete( - &client, - &service_id, - &clickpipe_id, - org_id.as_deref(), - json, - ) - .await - } - ClickPipeCommands::Start { - service_id, - clickpipe_id, - org_id, - } => { - cloud::commands::clickpipe_state( - &client, - &service_id, - &clickpipe_id, - "start", - org_id.as_deref(), - json, - ) - .await - } - ClickPipeCommands::Stop { - service_id, - clickpipe_id, - org_id, - } => { - cloud::commands::clickpipe_state( - &client, - &service_id, - &clickpipe_id, - "stop", - org_id.as_deref(), - json, - ) - .await - } - ClickPipeCommands::Resync { - service_id, - clickpipe_id, - org_id, - } => { - cloud::commands::clickpipe_state( - &client, - &service_id, - &clickpipe_id, - "resync", - org_id.as_deref(), - json, - ) - .await - } - ClickPipeCommands::Scale { - service_id, - clickpipe_id, - replicas, - cpu_millicores, - memory_gb, - org_id, - } => { - cloud::commands::clickpipe_scale( - &client, - &service_id, - &clickpipe_id, - replicas, - cpu_millicores, - memory_gb, - org_id.as_deref(), - json, - ) - .await - } - ClickPipeCommands::Settings { command } => match command { - ClickPipeSettingsCommands::Get { - service_id, - clickpipe_id, - org_id, - } => { - cloud::commands::clickpipe_settings_get( - &client, - &service_id, - &clickpipe_id, - org_id.as_deref(), - json, - ) - .await - } - ClickPipeSettingsCommands::Update { - service_id, - clickpipe_id, - streaming_max_insert_wait_ms, - object_storage_concurrency, - object_storage_polling_interval_ms, - object_storage_max_insert_bytes, - object_storage_max_file_count, - clickhouse_max_threads, - clickhouse_max_insert_threads, - object_storage_use_cluster_function, - clickhouse_parallel_view_processing, - org_id, - } => { - cloud::commands::clickpipe_settings_update( - &client, - &service_id, - &clickpipe_id, - streaming_max_insert_wait_ms, - object_storage_concurrency, - object_storage_polling_interval_ms, - object_storage_max_insert_bytes, - object_storage_max_file_count, - clickhouse_max_threads, - clickhouse_max_insert_threads, - object_storage_use_cluster_function, - clickhouse_parallel_view_processing, - org_id.as_deref(), - json, - ) - .await - } - }, - ClickPipeCommands::SchemaDiscover { - service_id, - command, - org_id, - } => { - cloud::commands::clickpipe_schema_discover( - &client, - &service_id, - &command, - org_id.as_deref(), - json, - ) - .await - } - ClickPipeCommands::Create { command } => match command { - ClickPipeCreateCommands::ObjectStorage(args) => { - cloud::commands::clickpipe_create_s3(&client, &args, json).await - } - ClickPipeCreateCommands::Kafka(args) => { - cloud::commands::clickpipe_create_kafka(&client, &args, json).await - } - ClickPipeCreateCommands::Kinesis(args) => { - cloud::commands::clickpipe_create_kinesis(&client, &args, json).await - } - ClickPipeCreateCommands::Postgres(args) => { - cloud::commands::clickpipe_create_postgres(&client, &args, json).await - } - ClickPipeCreateCommands::MySQL(args) => { - cloud::commands::clickpipe_create_mysql(&client, &args, json).await - } - ClickPipeCreateCommands::MongoDB(args) => { - cloud::commands::clickpipe_create_mongodb(&client, &args, json).await - } - ClickPipeCreateCommands::BigQuery(args) => { - cloud::commands::clickpipe_create_bigquery(&client, &args, json).await - } - }, - }, - }; - - result.map_err(boxed_cloud_error_to_top_level) -} - -fn cloud_error_to_top_level(e: cloud::CloudError) -> Error { - match e.kind { - cloud::CloudErrorKind::Auth => Error::AuthRequired(e.message), - cloud::CloudErrorKind::Generic => Error::Cloud(e.message), - } -} - -// Cloud command fns return `Box`, so the `CloudError.kind` -// only survives via downcast — without it, auth-flagged errors silently fall back -// to `Error::Cloud` (exit 1) instead of `Error::AuthRequired` (exit 4). -fn boxed_cloud_error_to_top_level(e: Box) -> Error { - match e.downcast::() { - Ok(ce) => cloud_error_to_top_level(*ce), - Err(other) => Error::Cloud(other.to_string()), - } -} - -async fn run_postgres( - client: &CloudClient, - command: PostgresCommands, - json: bool, -) -> std::result::Result<(), Box> { - use clickhouse_cloud_api::models::PostgresServiceSetStateCommand; - use cloud::postgres::{ - self as pg, PostgresCreateOptions, PostgresReadReplicaOptions, PostgresRestoreOptions, - PostgresUpdateOptions, - }; - - match command { - PostgresCommands::List { org_id, filter } => { - pg::postgres_list(client, org_id.as_deref(), &filter, json).await - } - PostgresCommands::Get { - postgres_id, - org_id, - } => pg::postgres_get(client, &postgres_id, org_id.as_deref(), json).await, - PostgresCommands::Create { - name, - region, - size, - provider, - pg_version, - ha_type, - tag, - pg_config_file, - pg_bouncer_config_file, - org_id, - } => { - let opts = PostgresCreateOptions { - name: &name, - region: ®ion, - size: &size, - provider: &provider, - pg_version: pg_version.as_deref(), - ha_type: ha_type.as_deref(), - tags: &tag, - pg_config_file: pg_config_file.as_deref(), - pg_bouncer_config_file: pg_bouncer_config_file.as_deref(), - org_id: org_id.as_deref(), - }; - pg::postgres_create(client, opts, json).await - } - PostgresCommands::Update { - postgres_id, - size, - ha_type, - add_tag, - remove_tag, - org_id, - } => { - let opts = PostgresUpdateOptions { - size: size.as_deref(), - ha_type: ha_type.as_deref(), - add_tag: &add_tag, - remove_tag: &remove_tag, - org_id: org_id.as_deref(), - }; - pg::postgres_update(client, &postgres_id, opts, json).await - } - PostgresCommands::Delete { - postgres_id, - org_id, - } => pg::postgres_delete(client, &postgres_id, org_id.as_deref(), json).await, - PostgresCommands::Certs(PostgresCertsCommands::Get { - postgres_id, - output, - org_id, - }) => { - pg::postgres_certs_get( - client, - &postgres_id, - output.as_deref(), - org_id.as_deref(), - json, - ) - .await - } - PostgresCommands::Config(PostgresConfigCommands::Get { - postgres_id, - org_id, - }) => pg::postgres_config_get(client, &postgres_id, org_id.as_deref(), json).await, - PostgresCommands::Config(PostgresConfigCommands::Replace { - postgres_id, - file, - org_id, - }) => { - pg::postgres_config_replace(client, &postgres_id, &file, org_id.as_deref(), json).await - } - PostgresCommands::Config(PostgresConfigCommands::Patch { - postgres_id, - sets, - file, - org_id, - }) => { - pg::postgres_config_patch( - client, - &postgres_id, - &sets, - file.as_deref(), - org_id.as_deref(), - json, - ) - .await - } - PostgresCommands::ResetPassword { - postgres_id, - password, - generate, - org_id, - } => { - pg::postgres_reset_password( - client, - &postgres_id, - password.as_deref(), - generate, - org_id.as_deref(), - json, - ) - .await - } - PostgresCommands::ReadReplica(PostgresReadReplicaCommands::Create { - postgres_id, - name, - tag, - pg_config_file, - pg_bouncer_config_file, - org_id, - }) => { - let opts = PostgresReadReplicaOptions { - name: &name, - tags: &tag, - pg_config_file: pg_config_file.as_deref(), - pg_bouncer_config_file: pg_bouncer_config_file.as_deref(), - org_id: org_id.as_deref(), - }; - pg::postgres_read_replica_create(client, &postgres_id, opts, json).await - } - PostgresCommands::Restore { - postgres_id, - name, - restore_target, - tag, - pg_config_file, - pg_bouncer_config_file, - org_id, - } => { - let opts = PostgresRestoreOptions { - name: &name, - restore_target: &restore_target, - tags: &tag, - pg_config_file: pg_config_file.as_deref(), - pg_bouncer_config_file: pg_bouncer_config_file.as_deref(), - org_id: org_id.as_deref(), - }; - pg::postgres_restore(client, &postgres_id, opts, json).await - } - PostgresCommands::Restart { - postgres_id, - org_id, - } => { - pg::postgres_state_change( - client, - &postgres_id, - PostgresServiceSetStateCommand::Restart, - org_id.as_deref(), - json, - ) - .await - } - PostgresCommands::Promote { - postgres_id, - org_id, - } => { - pg::postgres_state_change( - client, - &postgres_id, - PostgresServiceSetStateCommand::Promote, - org_id.as_deref(), - json, - ) - .await - } - PostgresCommands::Switchover { - postgres_id, - org_id, - } => { - pg::postgres_state_change( - client, - &postgres_id, - PostgresServiceSetStateCommand::Switchover, - org_id.as_deref(), - json, - ) - .await - } - } -} - #[cfg(test)] mod tests { use super::*; use clap::Parser; - use cloud::{CloudError, CloudErrorKind}; #[test] fn json_output_true_when_flag_set() { assert!(json_output(true)); } - #[test] - fn ignored_env_notice_names_the_overridden_variables_and_winner() { - assert_eq!( - ignored_env_credentials_notice( - cloud::AuthSource::CredentialsFile, - cloud::EnvCredPresence { - key: true, - secret: true, - }, - ) - .as_deref(), - Some( - "note: CLICKHOUSE_CLOUD_API_KEY and CLICKHOUSE_CLOUD_API_SECRET are set but \ - ignored; using credentials file — see --debug" - ) - ); - assert_eq!( - ignored_env_credentials_notice( - cloud::AuthSource::CliFlags, - cloud::EnvCredPresence { - key: true, - secret: false, - }, - ) - .as_deref(), - Some( - "note: CLICKHOUSE_CLOUD_API_KEY is set but ignored; using CLI flags — see --debug" - ) - ); - } - - #[test] - fn ignored_env_notice_is_absent_when_env_wins_or_is_unset() { - assert!( - ignored_env_credentials_notice( - cloud::AuthSource::EnvVars, - cloud::EnvCredPresence { - key: true, - secret: true, - }, - ) - .is_none() - ); - assert!( - ignored_env_credentials_notice( - cloud::AuthSource::CredentialsFile, - cloud::EnvCredPresence { - key: false, - secret: false, - }, - ) - .is_none() - ); - } - fn parse(args: &[&str]) -> Commands { Cli::try_parse_from(args).unwrap().command } @@ -1574,39 +327,4 @@ mod tests { "update" ]))); } - - #[test] - fn cloud_error_kind_routes_to_top_level() { - assert!(matches!( - cloud_error_to_top_level(CloudError::auth("nope")), - Error::AuthRequired(_) - )); - assert!(matches!( - cloud_error_to_top_level(CloudError::new("boom")), - Error::Cloud(_) - )); - // Default kind is Generic. - assert_eq!(CloudError::new("x").kind, CloudErrorKind::Generic); - } - - #[test] - fn boxed_cloud_error_preserves_auth_kind_through_downcast() { - let boxed: Box = Box::new(CloudError::auth("nope")); - assert!(matches!( - boxed_cloud_error_to_top_level(boxed), - Error::AuthRequired(_) - )); - } - - #[test] - fn boxed_non_cloud_error_falls_back_to_generic() { - // Anything that isn't a CloudError must not downcast to AuthRequired — - // it falls back to Error::Cloud (exit 1). This pins the contract that a - // handler stringifying a CloudError before boxing silently loses exit 4. - let boxed: Box = "plain string error".into(); - assert!(matches!( - boxed_cloud_error_to_top_level(boxed), - Error::Cloud(_) - )); - } }