From ebf8f6312149e4721fc8435ea431e8fad7d64c58 Mon Sep 17 00:00:00 2001 From: sdairs Date: Thu, 1 Oct 2026 19:39:57 +0100 Subject: [PATCH 1/2] feat(cli): add cloud postgres backup list and cloud auth whoami (#1038) - `cloud postgres backup list [--cursor] [--limit 1..100]`: read-only page of retained base backups; human table (Key, Last modified) plus a "Next cursor:" line, and --json emits {result, limit, totalCount, nextCursor}. - `cloud auth whoami` (Beta): resolves the identity behind the active credentials. Unscoped (no organization lookup), allowed with OAuth, and auth failures keep the Auth kind / exit 4. Auth commands that need a client now take the standard client path; the local ones still return before credentials are resolved. - Permission declarations for both, clap/read-write tests, wiremock coverage for routes, pagination encoding, both whoami variants, OAuth, 401 and 404, and README notes. Co-Authored-By: Claude Opus 5.5 --- README.md | 9 +- crates/clickhousectl/src/cloud/auth.rs | 115 ++++++- crates/clickhousectl/src/cloud/cli.rs | 1 + crates/clickhousectl/src/cloud/mod.rs | 33 +- crates/clickhousectl/src/cloud/postgres.rs | 233 +++++++++++++- .../tests/cli_request_shape_test.rs | 291 ++++++++++++++++++ 6 files changed, 659 insertions(+), 23 deletions(-) diff --git a/README.md b/README.md index 215def96..f0505889 100644 --- a/README.md +++ b/README.md @@ -677,6 +677,7 @@ Learn how to [create API keys](https://clickhouse.com/docs/cloud/manage/openapi? ```bash clickhousectl cloud auth status # Show current auth state (including read-only/read-write labels) +clickhousectl cloud auth whoami # (Beta) Ask the Cloud API who the active credentials belong to clickhousectl cloud auth logout # Clear all saved credentials (credentials.json & tokens.json) clickhousectl cloud auth logout --oauth # Clear only OAuth tokens, keep API keys clickhousectl cloud auth logout --api-keys # Clear only API keys, keep OAuth tokens @@ -693,6 +694,8 @@ Credential resolution order: Supplying only `--api-key` or only `--api-secret` blocks fallback to other sources. `cloud auth status` still succeeds with no active authentication; `--debug` identifies the missing flag. +`cloud auth whoami` works with OAuth or API keys and needs no organization: a user shows its ID, email, +name and organizations; an API key shows its key ID, name and owning organization ID. When environment credentials are configured but a credentials file or explicit CLI flags take precedence, clickhousectl prints a one-line note to stderr. @@ -1361,7 +1364,11 @@ clickhousectl cloud postgres certs get --output ca.pem PGSSLMODE=verify-full PGSSLROOTCERT=ca.pem psql --host --port 5432 \ --username --dbname postgres -# Read replica and PITR restore +# Retained base backups, newest first; follow `Next cursor:` (JSON: nextCursor) with --cursor +clickhousectl cloud postgres backup list --limit 20 +clickhousectl cloud postgres backup list --cursor + +# Read replica and PITR restore (restore takes a point in time, not a backup key) clickhousectl cloud postgres read-replica create --source-name primary --name replica-1 clickhousectl cloud postgres read-replica create --name replica-2 \ --tag env=prod --pg-config-file ./pg.json diff --git a/crates/clickhousectl/src/cloud/auth.rs b/crates/clickhousectl/src/cloud/auth.rs index b1f7435d..9f5e042f 100644 --- a/crates/clickhousectl/src/cloud/auth.rs +++ b/crates/clickhousectl/src/cloud/auth.rs @@ -1,4 +1,5 @@ use super::permissions::Declaration as Permission; +use clickhouse_cloud_api::meta::operations as op; // Declare every API call made by these workflows, including optional lookups. pub(super) const PERMISSIONS: &[Permission] = &[ @@ -15,15 +16,18 @@ pub(super) const PERMISSIONS: &[Permission] = &[ "Reads local authentication state; no Cloud API call.", ), Permission::non_api("auth signup", "Opens account signup; no Cloud API call."), + Permission::api("auth whoami", &[&op::WHOAMI_GET]).unscoped(), ]; +use crate::cloud::client::{CloudClient, Result as CloudResult}; use crate::cloud::credentials; -use crate::cloud::output::eprint_line; +use crate::cloud::output::{eprint_line, print_human}; use crate::cloud::{ AuthSource, dotenv_env_provenance, env_cred_presence, resolve_active_auth_source, }; use crate::error::Error; use clap::Subcommand; +use clickhouse_cloud_api::models::Whoami; use serde::{Deserialize, Serialize}; use std::path::PathBuf; @@ -68,6 +72,12 @@ CONTEXT FOR AGENTS: }, /// Show current authentication status Status, + /// Show the identity behind the active credentials (Beta) + #[command(after_help = "\ +CONTEXT FOR AGENTS: + Works with OAuth or API key credentials and needs no --org-id. + `cloud auth status` shows which credential source is active.")] + Whoami, /// Create a ClickHouse Cloud account #[command(after_help = "\ CONTEXT FOR AGENTS: @@ -108,10 +118,60 @@ impl AuthCommands { AuthCommands::Logout { .. } => false, AuthCommands::Status => false, AuthCommands::Signup => false, + AuthCommands::Whoami => false, + } + } + + /// Whether this command calls the Cloud API, and so needs a `CloudClient` + /// built from the resolved credentials. The others manage local state. + pub fn needs_client(&self) -> bool { + match self { + AuthCommands::Login { .. } + | AuthCommands::Logout { .. } + | AuthCommands::Status + | AuthCommands::Signup => false, + AuthCommands::Whoami => true, + } + } +} + +/// Run an auth command that calls the Cloud API (see [`AuthCommands::needs_client`]). +pub async fn run_with_client( + client: &CloudClient, + command: AuthCommands, + json: bool, +) -> CloudResult<()> { + match command { + AuthCommands::Whoami => { + let identity = client.get_whoami().await?; + if json { + println!("{}", serde_json::to_string_pretty(&identity)?); + } else { + print_human(&identity)?; + } + Ok(()) + } + AuthCommands::Login { .. } + | AuthCommands::Logout { .. } + | AuthCommands::Status + | AuthCommands::Signup => { + unreachable!("local auth commands are handled before a client is built") } } } +impl CloudClient { + /// Resolve the caller behind the active credentials; not organization-scoped. + async fn get_whoami(&self) -> CloudResult { + let response = self + .api() + .whoami_get() + .await + .map_err(|error| self.convert_error(error))?; + Self::unwrap_response(response) + } +} + pub async fn run( command: AuthCommands, api_key: Option<&str>, @@ -158,6 +218,7 @@ pub async fn run( Ok(()) } } + AuthCommands::Whoami => unreachable!("whoami runs through run_with_client"), AuthCommands::Signup => { let api_url = api_url.unwrap_or("https://api.clickhouse.cloud"); let parsed = url::Url::parse(api_url) @@ -838,8 +899,60 @@ mod tests { AuthCli::try_parse_from(["clickhousectl", "signup"]) .unwrap() .command, + AuthCli::try_parse_from(["clickhousectl", "whoami"]) + .unwrap() + .command, ]; assert!(commands.iter().all(|command| !command.is_write())); + // Only whoami calls the Cloud API; the rest manage local state. + let remote: Vec<_> = commands + .iter() + .map(|command| matches!(command, AuthCommands::Whoami)) + .collect(); + assert_eq!( + commands + .iter() + .map(AuthCommands::needs_client) + .collect::>(), + remote + ); + assert_eq!(remote.iter().filter(|remote| **remote).count(), 1); + } + + #[test] + fn auth_whoami_takes_shared_credentials_and_rejects_arguments() { + let cli = Cli::try_parse_from([ + "clickhousectl", + "cloud", + "auth", + "whoami", + "--api-key", + "key", + "--api-secret", + "secret", + "--json", + ]) + .unwrap(); + let Commands::Cloud(args) = cli.command else { + panic!("expected cloud command"); + }; + assert_eq!(args.api_key.as_deref(), Some("key")); + assert_eq!(args.api_secret.as_deref(), Some("secret")); + assert!(args.json); + assert!(!args.command.is_write_command()); + assert!(matches!( + args.command, + crate::cloud::cli::CloudCommands::Auth { + command: AuthCommands::Whoami + } + )); + assert_eq!( + Cli::try_parse_from(["clickhousectl", "cloud", "auth", "whoami", "extra"]) + .err() + .unwrap() + .kind(), + clap::error::ErrorKind::UnknownArgument + ); } #[test] diff --git a/crates/clickhousectl/src/cloud/cli.rs b/crates/clickhousectl/src/cloud/cli.rs index e15adada..9562fce2 100644 --- a/crates/clickhousectl/src/cloud/cli.rs +++ b/crates/clickhousectl/src/cloud/cli.rs @@ -413,6 +413,7 @@ pub(crate) mod tests { &["activity", "get", "activity-1"], &["backup", "get", "svc-1", "backup-1"], &["postgres", "config", "get", "pg-1"], + &["postgres", "backup", "list", "pg-1"], &["clickstack", "dashboard", "get", "svc-1", "dashboard-1"], &[ "clickpipe", diff --git a/crates/clickhousectl/src/cloud/mod.rs b/crates/clickhousectl/src/cloud/mod.rs index f542038b..e90a075d 100644 --- a/crates/clickhousectl/src/cloud/mod.rs +++ b/crates/clickhousectl/src/cloud/mod.rs @@ -58,18 +58,21 @@ fn ignored_env_credentials_notice( } pub async fn run(args: CloudArgs, json: bool) -> Result<()> { - // Auth subcommands don't need a client. - if let CloudCommands::Auth { command } = args.command { - return auth::run( - command, - args.api_key.as_deref(), - args.api_secret.as_deref(), - args.url.as_deref(), - args.debug, - json, - ) - .await; - } + // Local auth subcommands don't need a client; `whoami` takes the normal path. + let command = match args.command { + CloudCommands::Auth { command } if !command.needs_client() => { + return auth::run( + command, + args.api_key.as_deref(), + args.api_secret.as_deref(), + args.url.as_deref(), + args.debug, + json, + ) + .await; + } + command => command, + }; // Refresh OAuth tokens if needed. Errors here are filesystem failures // (refresh-rpc failures are swallowed and tokens cleared), so this stays @@ -104,7 +107,7 @@ pub async fn run(args: CloudArgs, json: bool) -> Result<()> { // 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() { + if client.is_bearer_auth() && command.is_write_command() { return Err(Error::AuthRequired( "This command requires API key authentication. \ OAuth (browser login) provides read-only access.\n\n\ @@ -119,7 +122,7 @@ pub async fn run(args: CloudArgs, json: bool) -> Result<()> { )); } - dispatch(&client, args.command, json) + dispatch(&client, command, json) .await .map_err(cloud_error_to_top_level) } @@ -150,7 +153,7 @@ fn cloud_error_to_top_level(e: CloudError) -> Error { async fn dispatch(client: &CloudClient, command: CloudCommands, json: bool) -> client::Result<()> { match command { - CloudCommands::Auth { .. } => unreachable!("handled above"), + CloudCommands::Auth { command } => auth::run_with_client(client, command, json).await, CloudCommands::Org { command } => organizations::run_org(client, command, json).await, CloudCommands::Service { command } => services::run(client, command, json).await, CloudCommands::Member { command } => organizations::run_member(client, command, json).await, diff --git a/crates/clickhousectl/src/cloud/postgres.rs b/crates/clickhousectl/src/cloud/postgres.rs index 9060c9ad..059d61d0 100644 --- a/crates/clickhousectl/src/cloud/postgres.rs +++ b/crates/clickhousectl/src/cloud/postgres.rs @@ -60,6 +60,10 @@ pub(super) const PERMISSIONS: &[Permission] = &[ "postgres prometheus org", &[&op::POSTGRES_ORG_PROMETHEUS_GET], ), + Permission::api( + "postgres backup list", + &[&op::POSTGRES_SERVICE_BACKUP_GET_LIST], + ), Permission::api("postgres restore", &[&op::POSTGRES_INSTANCE_RESTORE]).when(&[ Conditional::flag("source-name", &[&op::POSTGRES_SERVICE_GET_LIST]), ]), @@ -88,12 +92,13 @@ use clickhouse_cloud_api::RunPostgresQueryRequest; use clickhouse_cloud_api::models::{ ApiResponse, PgBouncerConfig, PgConfig, PgConfigDefaultTransactionIsolation, PgConfigSslMinProtocolVersion, PgConfigWalCompression, PgHaType, PgIdProperty, PgProvider, - PgSize, PgVersion, PostgresInstanceConfig, PostgresLogEntry, PostgresLogsGetListSortorder, - PostgresMetrics, PostgresService, PostgresServiceListItem, PostgresServicePatchRequest, - PostgresServicePostRequest, PostgresServiceReadReplicaRequest, PostgresServiceRestoreRequest, - PostgresServiceSetPassword, PostgresServiceSetState, PostgresServiceSetStateCommand, - PostgresSlowQueryPattern, PostgresSlowQueryPatternDetail, ResourceTagsV1, - ResourceTagsV1Response, SlowQueryPatternsGetListSortby, SlowQueryPatternsGetListSortorder, + PgSize, PgVersion, PostgresBackup, PostgresInstanceConfig, PostgresLogEntry, + PostgresLogsGetListSortorder, PostgresMetrics, PostgresService, PostgresServiceListItem, + PostgresServicePatchRequest, PostgresServicePostRequest, PostgresServiceReadReplicaRequest, + PostgresServiceRestoreRequest, PostgresServiceSetPassword, PostgresServiceSetState, + PostgresServiceSetStateCommand, PostgresSlowQueryPattern, PostgresSlowQueryPatternDetail, + ResourceTagsV1, ResourceTagsV1Response, SlowQueryPatternsGetListSortby, + SlowQueryPatternsGetListSortorder, }; use futures_util::StreamExt; use serde::de::DeserializeOwned; @@ -331,6 +336,10 @@ CONTEXT FOR AGENTS: )] Prometheus(PrometheusCommands), + /// Inspect Postgres service backups + #[command(subcommand)] + Backup(PostgresBackupCommands), + /// Restore a Postgres service to a point in time #[command(after_help = "\ CONTEXT FOR AGENTS: @@ -481,6 +490,24 @@ CONTEXT FOR AGENTS: }, } +#[derive(Subcommand)] +pub enum PostgresBackupCommands { + /// List Postgres service backups + #[command(after_help = "\ +CONTEXT FOR AGENTS: + Restore targets a point in time, not a backup key: `cloud postgres restore --restore-target`.")] + List { + /// Postgres service ID (from `cloud postgres list`) + postgres_id: String, + /// Cursor from nextCursor + #[arg(long)] + cursor: Option, + /// Maximum records per page (1–100) + #[arg(long, allow_negative_numbers = true, value_parser = clap::value_parser!(i64).range(1..=100))] + limit: Option, + }, +} + #[derive(Subcommand)] pub enum PrometheusCommands { /// Get metrics for one Postgres service @@ -578,6 +605,7 @@ impl PostgresCommands { | PostgresCommands::Logs { .. } | PostgresCommands::SlowQueries(_) | PostgresCommands::Prometheus(_) => false, + PostgresCommands::Backup(PostgresBackupCommands::List { .. }) => false, PostgresCommands::Certs(CertsCommands::Get { .. }) => false, PostgresCommands::Config(ConfigCommands::Get { .. }) => false, @@ -852,6 +880,11 @@ pub async fn run(client: &CloudClient, command: PostgresCommands, json: bool) -> PostgresCommands::Prometheus(PrometheusCommands::Org) => { postgres_prometheus_org(client, json).await } + PostgresCommands::Backup(PostgresBackupCommands::List { + postgres_id, + cursor, + limit, + }) => postgres_backup_list(client, &postgres_id, cursor.as_deref(), limit, json).await, PostgresCommands::Restore { postgres_id, name, @@ -2740,11 +2773,94 @@ async fn postgres_prometheus_org(client: &CloudClient, json: bool) -> CloudResul print_prometheus(&metrics, json) } +/// One backup list page, keeping the envelope's continuation metadata for `--json`. +#[derive(serde::Serialize)] +struct PostgresBackupPage { + result: Vec, + #[serde(skip_serializing_if = "Option::is_none")] + limit: Option, + #[serde(rename = "totalCount", skip_serializing_if = "Option::is_none")] + total_count: Option, + #[serde(rename = "nextCursor", skip_serializing_if = "Option::is_none")] + next_cursor: Option, +} + +async fn postgres_backup_list( + client: &CloudClient, + postgres_id: &str, + cursor: Option<&str>, + limit: Option, + json: bool, +) -> CloudResult<()> { + let org_id = resolve_org_id(client).await?; + let page = client + .list_postgres_backups_page(&org_id, postgres_id, cursor, limit) + .await?; + if json { + println!("{}", serde_json::to_string_pretty(&page)?); + } else { + println!("{}", render_postgres_backup_page(page)); + } + Ok(()) +} + +fn render_postgres_backup_page(page: PostgresBackupPage) -> String { + #[derive(Tabled)] + struct Row { + #[tabled(rename = "Key")] + key: String, + #[tabled(rename = "Last modified")] + last_modified: String, + } + + let mut out = if page.result.is_empty() { + "No Postgres backups found".to_owned() + } else { + let rows = page.result.into_iter().map(|backup| Row { + key: or_absent(backup.key), + last_modified: or_absent( + backup + .last_modified + .map(|at| at.to_rfc3339_opts(chrono::SecondsFormat::AutoSi, true)), + ), + }); + Table::new(rows).with(Style::markdown()).to_string() + }; + if let Some(cursor) = page.next_cursor { + out.push_str(&format!("\nNext cursor: {cursor}")); + } + out +} + // --------------------------------------------------------------------------- // Role changes: promote / switchover // --------------------------------------------------------------------------- impl CloudClient { + /// Fetch one page of a Postgres service's retained base backups. + async fn list_postgres_backups_page( + &self, + org_id: &str, + postgres_id: &str, + cursor: Option<&str>, + limit: Option, + ) -> CloudResult { + let mut response = self + .api() + .postgres_service_backup_get_list(org_id, postgres_id, cursor, limit) + .await + .map_err(|error| self.convert_error_for_organization(error, org_id))?; + let limit = response.limit; + let total_count = response.total_count; + let next_cursor = response.next_cursor.take(); + Ok(PostgresBackupPage { + result: Self::unwrap_response(response)?, + limit, + total_count, + next_cursor, + }) + } + /// Fetch raw Prometheus exposition text for a Postgres service. pub async fn get_postgres_prometheus( &self, @@ -3022,6 +3138,111 @@ mod tests { command } + #[test] + fn parses_backup_list_pagination_as_read() { + let cmd = parse_postgres(&[ + "clickhousectl", + "cloud", + "postgres", + "backup", + "list", + "pg-1", + ]); + assert!(!cmd.is_write()); + let PostgresCommands::Backup(PostgresBackupCommands::List { + postgres_id, + cursor, + limit, + }) = cmd + else { + panic!("expected backup list"); + }; + assert_eq!(postgres_id, "pg-1"); + assert!(cursor.is_none()); + assert!(limit.is_none()); + + for limit_value in ["1", "100"] { + let cmd = parse_postgres(&[ + "clickhousectl", + "cloud", + "postgres", + "backup", + "list", + "pg-1", + "--cursor", + "next+/=", + "--limit", + limit_value, + "--org-id", + "org-1", + ]); + assert!(!cmd.is_write()); + let PostgresCommands::Backup(PostgresBackupCommands::List { cursor, limit, .. }) = cmd + else { + panic!("expected backup list"); + }; + assert_eq!(cursor.as_deref(), Some("next+/=")); + assert_eq!(limit, Some(limit_value.parse().unwrap())); + } + } + + #[test] + fn backup_list_requires_an_id_and_bounds_limit() { + let base = ["clickhousectl", "cloud", "postgres", "backup", "list"]; + assert_eq!( + Cli::try_parse_from(base).err().unwrap().kind(), + clap::error::ErrorKind::MissingRequiredArgument + ); + for value in ["-1", "0", "101", "many"] { + let mut args = base.to_vec(); + args.extend(["pg-1", "--limit", value]); + assert_eq!( + Cli::try_parse_from(args).err().unwrap().kind(), + clap::error::ErrorKind::ValueValidation, + "{value}" + ); + } + // The group itself is not executable without a leaf. + assert!(Cli::try_parse_from(["clickhousectl", "cloud", "postgres", "backup"]).is_err()); + } + + #[test] + fn backup_page_human_view_renders_absent_fields_and_next_cursor() { + let page = PostgresBackupPage { + result: vec![ + PostgresBackup { + key: Some("basebackups_005/0001_backup_stop_sentinel.json".into()), + last_modified: Some("2026-03-31T18:17:37Z".parse().unwrap()), + }, + PostgresBackup::default(), + ], + limit: Some(100), + total_count: Some(2), + next_cursor: Some("page-2".into()), + }; + let rendered = render_postgres_backup_page(page); + assert!(rendered.contains("| Key"), "{rendered}"); + assert!(rendered.contains("Last modified"), "{rendered}"); + assert!(rendered.contains("basebackups_005/0001_backup_stop_sentinel.json")); + assert!(rendered.contains("2026-03-31T18:17:37Z"), "{rendered}"); + assert!( + rendered.lines().any(|line| { + line.split('|').map(str::trim).collect::>() == ["", "-", "-", ""] + }), + "{rendered}" + ); + assert_eq!(rendered.lines().last(), Some("Next cursor: page-2")); + + let empty = render_postgres_backup_page(PostgresBackupPage { + result: vec![], + limit: None, + total_count: None, + next_cursor: None, + }); + assert!(!empty.contains("Next cursor")); + assert!(!empty.contains('|')); + } + #[test] fn metrics_human_view_preserves_series_points_labels_and_units() { let metrics = PostgresMetrics { diff --git a/crates/clickhousectl/tests/cli_request_shape_test.rs b/crates/clickhousectl/tests/cli_request_shape_test.rs index 5d691008..af1fd788 100644 --- a/crates/clickhousectl/tests/cli_request_shape_test.rs +++ b/crates/clickhousectl/tests/cli_request_shape_test.rs @@ -29863,3 +29863,294 @@ async fn saved_query_invalid_input_and_oauth_writes_fail_before_http() { } assert!(server.received_requests().await.unwrap().is_empty()); } + +const POSTGRES_BACKUPS_PATH: &str = "/v1/organizations/org-1/postgres/pg-1/backups"; +const WHOAMI_USER_ID: &str = "44444444-5555-4666-8777-888888888888"; +const WHOAMI_ORG_ID: &str = "55555555-6666-4777-8888-999999999999"; + +/// Run a cloud command with API-key flags, or with saved OAuth tokens only. +fn invoke_cloud_as( + server: &MockServer, + project: &Path, + oauth: bool, + json: bool, + args: &[&str], +) -> std::process::Output { + let home = project.join("home"); + let cloud_dir = home.join(".clickhouse"); + std::fs::create_dir_all(&cloud_dir).unwrap(); + if oauth { + write_oauth_tokens(&cloud_dir, &server.uri()); + } + let mut command = Command::new(clickhousectl_binary()); + clear_inherited_env(&mut command); + command + .env("HOME", home) + .env("DO_NOT_TRACK", "1") + .current_dir(project) + .args(["cloud", "--url", &server.uri()]); + if !oauth { + command.args(["--api-key", "drift-key", "--api-secret", "drift-secret"]); + } + if json { + command.arg("--json"); + } + command.args(args).stdin(Stdio::null()).output().unwrap() +} + +fn postgres_backup_page(result: Value, next_cursor: Option<&str>) -> ResponseTemplate { + ResponseTemplate::new(200).set_body_json(serde_json::json!({ + "status": 200, + "requestId": "postgres-backup-list", + "result": result, + "limit": 2, + "totalCount": 3, + "nextCursor": next_cursor + })) +} + +#[tokio::test] +async fn postgres_backup_list_sends_pagination_and_renders_json_and_table() { + let server = MockServer::start().await; + let backups = serde_json::json!([ + { + "key": "basebackups_005/000000010000000000000002_backup_stop_sentinel.json", + "lastModified": "2026-03-31T18:17:37Z" + }, + {} + ]); + Mock::given(method("GET")) + .and(path(POSTGRES_BACKUPS_PATH)) + .and(wiremock::matchers::basic_auth("drift-key", "drift-secret")) + .and(query_param("cursor", "page /+?")) + .and(query_param("limit", "2")) + .respond_with(postgres_backup_page(backups.clone(), Some("next /+?"))) + .with_priority(1) + .expect(2) + .mount(&server) + .await; + Mock::given(method("GET")) + .and(path(POSTGRES_BACKUPS_PATH)) + .respond_with(postgres_backup_page(serde_json::json!([]), None)) + .with_priority(10) + .expect(1) + .mount(&server) + .await; + let project = tempfile::tempdir().unwrap(); + let args = [ + "postgres", "backup", "list", "pg-1", "--cursor", "page /+?", "--limit", "2", "--org-id", + "org-1", + ]; + + let json = invoke_cloud_as(&server, project.path(), false, true, &args); + assert_success(&json); + assert_eq!( + serde_json::from_slice::(&json.stdout).unwrap(), + serde_json::json!({ + "result": backups, + "limit": 2, + "totalCount": 3, + "nextCursor": "next /+?" + }) + ); + + let human = invoke_cloud_as(&server, project.path(), false, false, &args); + assert_success(&human); + let stdout = String::from_utf8_lossy(&human.stdout); + for expected in [ + "Key", + "Last modified", + "basebackups_005/000000010000000000000002_backup_stop_sentinel.json", + "2026-03-31T18:17:37Z", + "Next cursor: next /+?", + ] { + assert!(stdout.contains(expected), "missing {expected}: {stdout}"); + } + + // Without pagination flags nothing is added to the query string. + let empty = invoke_cloud_as( + &server, + project.path(), + false, + false, + &["postgres", "backup", "list", "pg-1", "--org-id", "org-1"], + ); + assert_success(&empty); + assert!(!String::from_utf8_lossy(&empty.stdout).contains("Next cursor")); + let requests = server.received_requests().await.unwrap(); + assert_eq!(requests.len(), 3); + assert!(requests[2].url.query().is_none(), "{:?}", requests[2].url); +} + +#[tokio::test] +async fn postgres_backup_list_supports_oauth_and_reports_not_found() { + let server = MockServer::start().await; + Mock::given(method("GET")) + .and(path(POSTGRES_BACKUPS_PATH)) + .and(header("authorization", "Bearer test-bearer-token")) + .respond_with(postgres_backup_page(serde_json::json!([]), None)) + .expect(1) + .mount(&server) + .await; + let project = tempfile::tempdir().unwrap(); + let args = ["postgres", "backup", "list", "pg-1", "--org-id", "org-1"]; + let oauth = invoke_cloud_as(&server, project.path(), true, true, &args); + assert_success(&oauth); + assert_eq!( + serde_json::from_slice::(&oauth.stdout).unwrap()["result"], + serde_json::json!([]) + ); + + let missing = MockServer::start().await; + Mock::given(method("GET")) + .and(path(POSTGRES_BACKUPS_PATH)) + .respond_with(ResponseTemplate::new(404).set_body_json(serde_json::json!({ + "status": 404, + "requestId": "postgres-backup-missing", + "error": "Postgres service pg-1 was not found" + }))) + .expect(1) + .mount(&missing) + .await; + let project = tempfile::tempdir().unwrap(); + let output = invoke_cloud_as(&missing, project.path(), false, false, &args); + assert_eq!(output.status.code(), Some(1)); + assert!(output.stdout.is_empty()); + assert!( + String::from_utf8_lossy(&output.stderr).contains("Postgres service pg-1 was not found"), + "{}", + String::from_utf8_lossy(&output.stderr) + ); +} + +fn whoami_envelope(result: Value) -> ResponseTemplate { + ResponseTemplate::new(200).set_body_json(serde_json::json!({ + "status": 200, + "requestId": "whoami", + "result": result + })) +} + +#[tokio::test] +async fn auth_whoami_renders_both_identities_without_an_organization() { + let user = serde_json::json!({ + "actorType": "user", + "userId": WHOAMI_USER_ID, + "email": "ada@example.com", + "name": "Ada Lovelace", + "organizations": [ + {"organizationId": WHOAMI_ORG_ID, "organizationName": "Analytical Engines"} + ] + }); + let api_key = serde_json::json!({ + "actorType": "apiKey", + "keyId": "key-123", + "name": "ci deploy", + "organizationId": WHOAMI_ORG_ID + }); + let future = serde_json::json!({"actorType": "robot", "robotId": "r2"}); + + for (identity, oauth, human_expected) in [ + ( + user.clone(), + true, + vec![ + "actorType: user", + WHOAMI_USER_ID, + "ada@example.com", + "Ada Lovelace", + WHOAMI_ORG_ID, + "Analytical Engines", + ], + ), + ( + api_key.clone(), + false, + vec!["actorType: apiKey", "key-123", "ci deploy", WHOAMI_ORG_ID], + ), + ( + future.clone(), + false, + vec!["actorType: robot", "robotId: r2"], + ), + ] { + let server = MockServer::start().await; + let auth = if oauth { + header("authorization", "Bearer test-bearer-token") + } else { + header( + "authorization", + "Basic ZHJpZnQta2V5OmRyaWZ0LXNlY3JldA==", // drift-key:drift-secret + ) + }; + Mock::given(method("GET")) + .and(path("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/v1/whoami")) + .and(auth) + .respond_with(whoami_envelope(identity.clone())) + .expect(2) + .mount(&server) + .await; + let project = tempfile::tempdir().unwrap(); + + let json = invoke_cloud_as(&server, project.path(), oauth, true, &["auth", "whoami"]); + assert_success(&json); + assert_eq!( + serde_json::from_slice::(&json.stdout).unwrap(), + identity + ); + + let human = invoke_cloud_as(&server, project.path(), oauth, false, &["auth", "whoami"]); + assert_success(&human); + let stdout = String::from_utf8_lossy(&human.stdout); + for expected in human_expected { + assert!(stdout.contains(expected), "missing {expected}: {stdout}"); + } + + // No organization lookup or other call: exactly the two whoami requests. + let requests = server.received_requests().await.unwrap(); + assert_eq!(requests.len(), 2); + assert!( + requests + .iter() + .all(|request| request.url.path() == "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/v1/whoami" && request.url.query().is_none()) + ); + } +} + +#[tokio::test] +async fn auth_whoami_auth_failures_exit_4() { + let server = MockServer::start().await; + Mock::given(method("GET")) + .and(path("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/v1/whoami")) + .respond_with(ResponseTemplate::new(401).set_body_json(serde_json::json!({ + "status": 401, + "requestId": "whoami-401", + "error": "Invalid credentials" + }))) + .expect(2) + .mount(&server) + .await; + let project = tempfile::tempdir().unwrap(); + for json in [false, true] { + let output = invoke_cloud_as(&server, project.path(), false, json, &["auth", "whoami"]); + assert_eq!(output.status.code(), Some(4), "json={json}"); + assert!(output.stdout.is_empty(), "json={json}"); + let stderr = String::from_utf8_lossy(&output.stderr); + assert!(stderr.contains("Invalid credentials"), "{stderr}"); + if json { + let error: Value = serde_json::from_str(&stderr).unwrap(); + assert_eq!(error["error"]["code"], "auth_required"); + } + } + + // With no credentials at all, whoami fails as auth-required before any request. + let empty = MockServer::start().await; + let output = invoke_cli_without_cloud_credentials(&empty, &["auth".into(), "whoami".into()]); + assert_eq!( + output.status.code(), + Some(4), + "{}", + String::from_utf8_lossy(&output.stderr) + ); + assert!(empty.received_requests().await.unwrap().is_empty()); +} From bd56869c027b7234b84427edc24f1211c4a86575 Mon Sep 17 00:00:00 2001 From: sdairs Date: Thu, 1 Oct 2026 19:53:57 +0100 Subject: [PATCH 2/2] feat(cli): add Kafka Glue schema registry and backup encryption config (#1038) `cloud clickpipe create kafka` and `clickpipe schema-discover kafka` gain `--schema-registry-type confluent|glue` with `--glue-region`, `--glue-registry-name` and `--glue-role-arn`. Omitting the type keeps the existing Confluent request unchanged; Glue conflicts with the Confluent registry flags and needs a role from `--glue-role-arn` or the source's `--iam-role`, rejected as usage errors before any request. `cloud service create --backup-encryption-config ` (private preview, requires `--backup-id`) passes a backup's encryption_config.json object through unchanged as `backupEncryptionConfig`; unreadable or non-object input is a usage error before any request. Co-Authored-By: Claude Opus 5.5 --- README.md | 19 + crates/clickhousectl/src/cloud/clickpipes.rs | 535 ++++++++++++++++-- crates/clickhousectl/src/cloud/services.rs | 150 ++++- .../tests/cli_request_shape_test.rs | 330 +++++++++++ 4 files changed, 970 insertions(+), 64 deletions(-) diff --git a/README.md b/README.md index f0505889..df650c96 100644 --- a/README.md +++ b/README.md @@ -860,6 +860,14 @@ clickhousectl cloud service create --name restored-service \ --ip-allow /32 \ --backup-id +# Restore a TDE service's backup from your own bucket (private preview): pass the +# backup's encryption_config.json unchanged, or - to read it from stdin +clickhousectl cloud service create --name restored-service \ + --provider aws \ + --region us-east-1 \ + --ip-allow /32 \ + --backup-id --backup-encryption-config ./encryption_config.json + # Create with release channel clickhousectl cloud service create --name my-service \ --provider aws \ @@ -2120,6 +2128,17 @@ clickhousectl cloud clickpipe create kafka \ --database default --table events \ --column "event_id:Int64" +# Avro via the AWS Glue Schema Registry on MSK; --glue-role-arn defaults to +# --iam-role and is required when the source does not use --iam-role +clickhousectl cloud clickpipe create kafka \ + --name my-glue-pipe \ + --brokers 'broker:9098' --topics events --kafka-type msk \ + --format Avro --iam-role arn:aws:iam::123456789012:role/ClickPipes \ + --schema-registry-type glue \ + --glue-region us-east-1 --glue-registry-name my-registry \ + --database default --table events \ + --column "event_id:Int64" + # Protobuf schema from a file, with exactly-once delivery clickhousectl cloud clickpipe create kafka \ --name my-protobuf-pipe \ diff --git a/crates/clickhousectl/src/cloud/clickpipes.rs b/crates/clickhousectl/src/cloud/clickpipes.rs index 70da701e..03fb650f 100644 --- a/crates/clickhousectl/src/cloud/clickpipes.rs +++ b/crates/clickhousectl/src/cloud/clickpipes.rs @@ -551,7 +551,9 @@ impl ClickPipeCommands { ), ClickPipeCreateCommands::Kafka(args) => ( "kafka", - build_create_request_args(&args.request, ClickPipeSourceKind::Kafka).err(), + build_create_request_args(&args.request, ClickPipeSourceKind::Kafka) + .and_then(|_| validate_kafka_schema_registry_args(&args.source)) + .err(), ), ClickPipeCreateCommands::Kinesis(args) => ( "kinesis", @@ -583,15 +585,22 @@ impl ClickPipeCommands { pub(crate) fn clickpipe_schema_discover_validation_error( &self, ) -> Option<(&'static str, String)> { - let Self::SchemaDiscover { - command: ClickPipeSchemaDiscoverCommands::Kinesis(args), - } = self - else { + let Self::SchemaDiscover { command } = self else { return None; }; - validate_kinesis_source_args(&args.source) - .err() - .map(|error| ("kinesis", error.message)) + match command { + ClickPipeSchemaDiscoverCommands::Kinesis(args) => { + validate_kinesis_source_args(&args.source) + .err() + .map(|error| ("kinesis", error.message)) + } + ClickPipeSchemaDiscoverCommands::Kafka(args) => { + validate_kafka_schema_registry_args(&args.source) + .err() + .map(|error| ("kafka", error.message)) + } + _ => None, + } } pub(crate) fn reverse_private_endpoint_create_validation_error(&self) -> Option { @@ -1139,6 +1148,10 @@ pub struct ObjectStorageCreateArgs { pub destination_roles: DestinationRoleArgs, } +/// Glue-only schema registry flags; each conflicts with the Confluent flags. +const KAFKA_GLUE_SCHEMA_REGISTRY_ARGS: [&str; 3] = + ["glue_region", "glue_registry_name", "glue_role_arn"]; + /// Source-connection fields for a Kafka / Kafka-compatible ClickPipe source. /// Flattened into both `KafkaCreateArgs` (pipe creation) and the schema-discover /// Kafka subcommand so the source field set has a single definition. @@ -1226,27 +1239,63 @@ pub struct KafkaSourceFields { #[arg(long)] pub offset_timestamp: Option, - /// Schema registry URL (for Avro/Protobuf formats) - #[arg(long)] + /// Schema registry type; confluent when omitted + #[arg(long, value_name = "TYPE", value_parser = ["confluent", "glue"])] + pub schema_registry_type: Option, + + /// Schema registry URL (for Avro/Protobuf formats; confluent only) + #[arg( + long, + required_if_eq("schema_registry_type", "confluent"), + conflicts_with_all = KAFKA_GLUE_SCHEMA_REGISTRY_ARGS + )] pub schema_registry_url: Option, - /// Schema registry username - #[arg(long)] + /// Schema registry username (confluent only) + #[arg(long, conflicts_with_all = KAFKA_GLUE_SCHEMA_REGISTRY_ARGS)] pub schema_registry_username: Option, - /// Schema registry password - #[arg(long)] + /// Schema registry password (confluent only) + #[arg(long, conflicts_with_all = KAFKA_GLUE_SCHEMA_REGISTRY_ARGS)] pub schema_registry_password: Option, + /// Glue registry AWS region (only with --schema-registry-type glue) + #[arg( + long, + value_name = "REGION", + requires = "schema_registry_type", + required_if_eq("schema_registry_type", "glue") + )] + pub glue_region: Option, + + /// Glue registry name (only with --schema-registry-type glue) + #[arg( + long, + value_name = "NAME", + requires = "schema_registry_type", + required_if_eq("schema_registry_type", "glue") + )] + pub glue_registry_name: Option, + + /// IAM role ARN for Glue access (only with --schema-registry-type glue) + /// + /// Defaults to --iam-role, so it is required when --iam-role is not set. + #[arg(long, value_name = "ARN", requires = "schema_registry_type")] + pub glue_role_arn: Option, + /// Path to a .proto file or FileDescriptorSet, or - to read stdin #[arg( long, value_name = "PATH", conflicts_with_all = [ + "schema_registry_type", "schema_registry_url", "schema_registry_username", "schema_registry_password", - "schema_registry_ca_certificate" + "schema_registry_ca_certificate", + "glue_region", + "glue_registry_name", + "glue_role_arn" ] )] pub protobuf_schema_file: Option, @@ -1263,8 +1312,8 @@ pub struct KafkaSourceFields { #[arg(long, requires = "client_certificate")] pub client_key: Option, - /// Path to schema registry CA certificate file - #[arg(long)] + /// Path to schema registry CA certificate file (confluent only) + #[arg(long, conflicts_with_all = KAFKA_GLUE_SCHEMA_REGISTRY_ARGS)] pub schema_registry_ca_certificate: Option, /// Reverse private endpoint IDs (repeatable) @@ -2660,15 +2709,15 @@ fn validate_kafka_source_args(args: &KafkaSourceFields) -> CloudResult<()> { )); } if args.protobuf_schema_file.is_some() - && (args.schema_registry_url.is_some() - || args.schema_registry_username.is_some() - || args.schema_registry_password.is_some() - || args.schema_registry_ca_certificate.is_some()) + && (has_confluent_registry_flags(args) + || has_glue_registry_flags(args) + || args.schema_registry_type.is_some()) { return Err(CloudError::new( "--protobuf-schema-file cannot be combined with schema registry flags", )); } + validate_kafka_schema_registry_args(args)?; if args.event_hubs_connection_string.is_some() && source_type != ClickPipePostKafkaSourceType::Azureeventhub @@ -2781,6 +2830,121 @@ fn validate_kafka_source_args(args: &KafkaSourceFields) -> CloudResult<()> { Ok(()) } +fn has_confluent_registry_flags(args: &KafkaSourceFields) -> bool { + args.schema_registry_url.is_some() + || args.schema_registry_username.is_some() + || args.schema_registry_password.is_some() + || args.schema_registry_ca_certificate.is_some() +} + +fn has_glue_registry_flags(args: &KafkaSourceFields) -> bool { + args.glue_region.is_some() || args.glue_registry_name.is_some() || args.glue_role_arn.is_some() +} + +/// Check the schema-registry flag combination clap cannot express on its own. +/// Runs before credentials or network access, and again in the builder. +fn validate_kafka_schema_registry_args(args: &KafkaSourceFields) -> CloudResult<()> { + match args.schema_registry_type.as_deref() { + Some("glue") => { + if has_confluent_registry_flags(args) { + return Err(CloudError::usage( + "--schema-registry-type glue cannot be combined with --schema-registry-url, \ + --schema-registry-username, --schema-registry-password or \ + --schema-registry-ca-certificate", + )); + } + if args.glue_region.is_none() || args.glue_registry_name.is_none() { + return Err(CloudError::usage( + "--schema-registry-type glue requires --glue-region and --glue-registry-name", + )); + } + if args.glue_role_arn.is_none() && args.iam_role.is_none() { + return Err(CloudError::usage( + "--schema-registry-type glue requires --glue-role-arn when the source does not \ + authenticate with --iam-role", + )); + } + } + explicit => { + if has_glue_registry_flags(args) { + return Err(CloudError::usage( + "--glue-region, --glue-registry-name and --glue-role-arn require \ + --schema-registry-type glue", + )); + } + if explicit.is_some() && args.schema_registry_url.is_none() { + return Err(CloudError::usage( + "--schema-registry-type confluent requires --schema-registry-url", + )); + } + } + } + Ok(()) +} + +/// Build the Kafka source's schema registry from validated flags; `None` when +/// no registry was requested. +fn build_kafka_schema_registry( + args: &KafkaSourceFields, +) -> CloudResult> { + use clickhouse_cloud_api::models::{ + ClickPipeKafkaGlueSchemaRegistry, ClickPipeKafkaGlueSchemaRegistryType, + ClickPipeKafkaSchemaRegistryCredentials, ClickPipeMutateKafkaConfluentSchemaRegistry, + ClickPipeMutateKafkaConfluentSchemaRegistryType, + }; + + if args.schema_registry_type.as_deref() == Some("glue") { + let (Some(glue_region), Some(glue_registry_name)) = + (args.glue_region.as_ref(), args.glue_registry_name.as_ref()) + else { + return Err(CloudError::usage( + "--schema-registry-type glue requires --glue-region and --glue-registry-name", + )); + }; + return Ok(Some( + ClickPipeKafkaGlueSchemaRegistry { + r#type: ClickPipeKafkaGlueSchemaRegistryType::Glue, + glue_region: glue_region.clone(), + glue_registry_name: glue_registry_name.clone(), + glue_role_arn: args.glue_role_arn.clone(), + } + .into(), + )); + } + + let Some(url) = args.schema_registry_url.as_ref() else { + return Ok(None); + }; + let credentials = match ( + args.schema_registry_username.as_deref(), + args.schema_registry_password.as_deref(), + ) { + (Some(username), Some(password)) => ClickPipeKafkaSchemaRegistryCredentials { + username: username.to_string(), + password: password.to_string(), + }, + _ => ClickPipeKafkaSchemaRegistryCredentials::default(), + }; + let ca_certificate = match args.schema_registry_ca_certificate.as_deref() { + Some(path) => Some(std::fs::read_to_string(path)?), + None => None, + }; + Ok(Some( + ClickPipeMutateKafkaConfluentSchemaRegistry { + // Send `type` only when the user named it; the API defaults to confluent. + r#type: args + .schema_registry_type + .as_ref() + .map(|_| ClickPipeMutateKafkaConfluentSchemaRegistryType::Confluent), + url: url.clone(), + authentication: Default::default(), + credentials, + ca_certificate, + } + .into(), + )) +} + // Keep I/O failures typed; the source chooses how invalid schema content is // reported independently of failures to open or read its input. fn read_protobuf_schema_file( @@ -2822,9 +2986,7 @@ fn build_kafka_source_with_exactly_once( exactly_once: Option, ) -> CloudResult { use clickhouse_cloud_api::models::{ - ClickPipeKafkaOffset, ClickPipeKafkaSchemaRegistryCredentials, - ClickPipeMutateKafkaConfluentSchemaRegistry, ClickPipePostKafkaSource, - ClickPipePostKafkaSourceAuthentication, + ClickPipeKafkaOffset, ClickPipePostKafkaSource, ClickPipePostKafkaSourceAuthentication, }; validate_kafka_source_args(args)?; @@ -2855,34 +3017,7 @@ fn build_kafka_source_with_exactly_once( }; let credentials = build_kafka_credentials(authentication.as_ref(), args, mtls_cert_contents)?; - let schema_registry = args - .schema_registry_url - .as_ref() - .map(|url| -> CloudResult<_> { - let credentials = match ( - args.schema_registry_username.as_deref(), - args.schema_registry_password.as_deref(), - ) { - (Some(username), Some(password)) => ClickPipeKafkaSchemaRegistryCredentials { - username: username.to_string(), - password: password.to_string(), - }, - _ => ClickPipeKafkaSchemaRegistryCredentials::default(), - }; - let ca_certificate = match args.schema_registry_ca_certificate.as_deref() { - Some(path) => Some(std::fs::read_to_string(path)?), - None => None, - }; - Ok(ClickPipeMutateKafkaConfluentSchemaRegistry { - r#type: None, - url: url.clone(), - authentication: Default::default(), - credentials, - ca_certificate, - } - .into()) - }) - .transpose()?; + let schema_registry = build_kafka_schema_registry(args)?; let ca_certificate = match args.ca_certificate.as_deref() { Some(path) => Some(std::fs::read_to_string(path)?), @@ -13488,9 +13623,13 @@ mod tests { secret_key: None, offset: "from_beginning".into(), offset_timestamp: None, + schema_registry_type: None, schema_registry_url: None, schema_registry_username: None, schema_registry_password: None, + glue_region: None, + glue_registry_name: None, + glue_role_arn: None, protobuf_schema_file: None, ca_certificate: None, client_certificate: None, @@ -13836,6 +13975,298 @@ mod tests { assert_eq!(registry.ca_certificate.as_deref(), Some("REGISTRY_CA")); } + fn glue_registry( + source: clickhouse_cloud_api::models::ClickPipePostKafkaSource, + ) -> clickhouse_cloud_api::models::ClickPipeKafkaGlueSchemaRegistry { + let Some( + clickhouse_cloud_api::models::ClickPipeMutateKafkaSchemaRegistry::ClickPipeKafkaGlueSchemaRegistry( + registry, + ), + ) = source.schema_registry + else { + panic!("a Glue schema registry is populated"); + }; + registry + } + + #[test] + fn build_kafka_source_builds_a_minimal_glue_schema_registry() { + let mut args = kafka_args().source; + args.format = "Avro".into(); + args.kafka_type = "msk".into(); + args.iam_role = Some("arn:aws:iam::123456789012:role/Source".into()); + args.schema_registry_type = Some("glue".into()); + args.glue_region = Some("us-east-1".into()); + args.glue_registry_name = Some("my-registry".into()); + + let source = build_kafka_source(&args).unwrap(); + assert_eq!(kafka_auth(&source).as_deref(), Some("IAM_ROLE")); + let registry = glue_registry(source); + assert_eq!( + registry.r#type, + clickhouse_cloud_api::models::ClickPipeKafkaGlueSchemaRegistryType::Glue + ); + assert_eq!(registry.glue_region, "us-east-1"); + assert_eq!(registry.glue_registry_name, "my-registry"); + assert_eq!(registry.glue_role_arn, None); + } + + #[test] + fn build_kafka_source_builds_a_maximal_glue_schema_registry() { + let mut args = kafka_args().source; + args.format = "AvroConfluent".into(); + args.kafka_type = "msk".into(); + args.username = Some("user".into()); + args.password = Some("password".into()); + args.schema_registry_type = Some("glue".into()); + args.glue_region = Some("eu-west-1".into()); + args.glue_registry_name = Some("registry".into()); + args.glue_role_arn = Some("arn:aws:iam::123456789012:role/Glue".into()); + + let source = build_kafka_source(&args).unwrap(); + assert_eq!(kafka_auth(&source).as_deref(), Some("PLAIN")); + assert_eq!(source.iam_role, None); + let registry = glue_registry(source); + assert_eq!(registry.glue_region, "eu-west-1"); + assert_eq!(registry.glue_registry_name, "registry"); + assert_eq!( + registry.glue_role_arn.as_deref(), + Some("arn:aws:iam::123456789012:role/Glue") + ); + } + + #[test] + fn build_kafka_source_sends_confluent_type_only_when_named() { + let mut args = kafka_args().source; + args.schema_registry_url = Some("https://registry.example".into()); + args.schema_registry_type = Some("confluent".into()); + + let source = build_kafka_source(&args).unwrap(); + let Some( + clickhouse_cloud_api::models::ClickPipeMutateKafkaSchemaRegistry::ClickPipeMutateKafkaConfluentSchemaRegistry( + registry, + ), + ) = source.schema_registry + else { + panic!("a Confluent schema registry is populated"); + }; + assert_eq!( + registry.r#type, + Some(clickhouse_cloud_api::models::ClickPipeMutateKafkaConfluentSchemaRegistryType::Confluent) + ); + assert_eq!(registry.url, "https://registry.example"); + + args.schema_registry_type = None; + args.schema_registry_url = None; + assert_eq!(build_kafka_source(&args).unwrap().schema_registry, None); + } + + #[test] + fn kafka_schema_registry_misuse_is_a_usage_error() { + let glue = || { + let mut args = kafka_args().source; + args.schema_registry_type = Some("glue".into()); + args.glue_region = Some("us-east-1".into()); + args.glue_registry_name = Some("my-registry".into()); + args + }; + let mut cases = Vec::new(); + + // No role: neither --glue-role-arn nor a source --iam-role. + cases.push(glue()); + let mut confluent_flag = glue(); + confluent_flag.glue_role_arn = Some("arn:role".into()); + confluent_flag.schema_registry_url = Some("https://registry.example".into()); + cases.push(confluent_flag); + let mut missing_name = glue(); + missing_name.glue_role_arn = Some("arn:role".into()); + missing_name.glue_registry_name = None; + cases.push(missing_name); + let mut untyped = kafka_args().source; + untyped.glue_region = Some("us-east-1".into()); + cases.push(untyped); + let mut confluent_without_url = kafka_args().source; + confluent_without_url.schema_registry_type = Some("confluent".into()); + cases.push(confluent_without_url); + + for args in cases { + let error = build_kafka_source(&args).unwrap_err(); + assert_eq!( + error.kind, + crate::cloud::client::CloudErrorKind::Usage, + "{error:?}" + ); + } + + let mut protobuf = kafka_args().source; + protobuf.format = "Protobuf".into(); + protobuf.protobuf_schema_file = Some("/file/that/does/not/exist".into()); + protobuf.schema_registry_type = Some("glue".into()); + let error = build_kafka_source(&protobuf).unwrap_err(); + assert!(error.message.contains("schema registry flags"), "{error:?}"); + } + + #[test] + fn parses_kafka_glue_schema_registry_flags_on_both_surfaces() { + for base in [kafka_create_cli_args(), kafka_discover_cli_args()] { + let mut args = base; + args.extend([ + "--schema-registry-type", + "glue", + "--glue-region", + "us-east-1", + "--glue-registry-name", + "my-registry", + "--glue-role-arn", + "arn:role", + ]); + let source = match parse_clickpipe(&args) { + ClickPipeCommands::Create { + command: ClickPipeCreateCommands::Kafka(args), + } => args.source, + ClickPipeCommands::SchemaDiscover { + command: ClickPipeSchemaDiscoverCommands::Kafka(args), + } => args.source, + _ => panic!("expected a Kafka command"), + }; + assert_eq!(source.schema_registry_type.as_deref(), Some("glue")); + assert_eq!(source.glue_region.as_deref(), Some("us-east-1")); + assert_eq!(source.glue_registry_name.as_deref(), Some("my-registry")); + assert_eq!(source.glue_role_arn.as_deref(), Some("arn:role")); + assert_eq!(source.schema_registry_url, None); + } + + let ClickPipeCommands::Create { + command: ClickPipeCreateCommands::Kafka(args), + } = parse_clickpipe(&kafka_create_cli_args()) + else { + panic!("expected kafka create"); + }; + assert_eq!(args.source.schema_registry_type, None); + assert_eq!(args.source.glue_region, None); + assert_eq!(args.source.glue_registry_name, None); + assert_eq!(args.source.glue_role_arn, None); + } + + #[test] + fn kafka_glue_schema_registry_flag_constraints_are_clap_errors() { + use clap::error::ErrorKind; + let glue = [ + "--schema-registry-type", + "glue", + "--glue-region", + "us-east-1", + "--glue-registry-name", + "my-registry", + ]; + let cases: Vec<(Vec<&str>, ErrorKind)> = vec![ + ( + vec!["--schema-registry-type", "avro"], + ErrorKind::InvalidValue, + ), + ( + vec!["--schema-registry-type", "glue"], + ErrorKind::MissingRequiredArgument, + ), + ( + vec![ + "--schema-registry-type", + "glue", + "--glue-region", + "us-east-1", + ], + ErrorKind::MissingRequiredArgument, + ), + ( + vec!["--glue-region", "us-east-1"], + ErrorKind::MissingRequiredArgument, + ), + ( + vec!["--glue-role-arn", "arn:role"], + ErrorKind::MissingRequiredArgument, + ), + ( + vec!["--schema-registry-type", "confluent"], + ErrorKind::MissingRequiredArgument, + ), + ( + [ + &glue[..], + &["--schema-registry-url", "https://registry.example"], + ] + .concat(), + ErrorKind::ArgumentConflict, + ), + ( + [&glue[..], &["--schema-registry-username", "user"]].concat(), + ErrorKind::ArgumentConflict, + ), + ( + [&glue[..], &["--schema-registry-password", "password"]].concat(), + ErrorKind::ArgumentConflict, + ), + ( + [ + &glue[..], + &["--schema-registry-ca-certificate", "/tmp/ca.pem"], + ] + .concat(), + ErrorKind::ArgumentConflict, + ), + ( + [ + &[ + "--format", + "Protobuf", + "--protobuf-schema-file", + "/tmp/e.proto", + ][..], + &glue[..], + ] + .concat(), + ErrorKind::ArgumentConflict, + ), + ]; + for (extra, kind) in cases { + let mut args = kafka_create_cli_args(); + args.extend(extra.iter().copied()); + assert_eq!(clickpipe_parse_error(&args).kind(), kind, "{extra:?}"); + } + } + + #[test] + fn kafka_glue_schema_registry_runtime_rules_run_before_execution() { + let mut create = kafka_create_cli_args(); + create.extend([ + "--schema-registry-type", + "glue", + "--glue-region", + "us-east-1", + "--glue-registry-name", + "my-registry", + ]); + let message = clickpipe_validation_message(&parse_clickpipe(&create)) + .expect("Glue without a role is rejected before execution"); + assert!(message.contains("--glue-role-arn"), "{message}"); + + let mut discover = kafka_discover_cli_args(); + discover.extend(&create[create.len() - 6..]); + let message = parse_clickpipe(&discover) + .clickpipe_schema_discover_validation_error() + .map(|(source, message)| { + assert_eq!(source, "kafka"); + message + }) + .expect("Glue without a role is rejected before schema discovery"); + assert!(message.contains("--glue-role-arn"), "{message}"); + + create.extend(["--iam-role", "arn:role"]); + assert_eq!( + clickpipe_validation_message(&parse_clickpipe(&create)), + None + ); + } + #[test] fn kinesis_protobuf_parsing_and_builder_cover_create_and_discovery() { let dir = tempfile::tempdir().unwrap(); diff --git a/crates/clickhousectl/src/cloud/services.rs b/crates/clickhousectl/src/cloud/services.rs index 91664761..94c7e86a 100644 --- a/crates/clickhousectl/src/cloud/services.rs +++ b/crates/clickhousectl/src/cloud/services.rs @@ -168,17 +168,17 @@ use crate::failure::{self, ApiFailure, FailureKind, FailureStage, ProvisioningSt use clap::builder::PossibleValuesParser; use clap::{ArgGroup, Subcommand}; use clickhouse_cloud_api::models::{ - AutoscalingMode, InstancePrivateEndpointsPatch, InstanceServiceQueryApiEndpointsPostRequest, - InstanceTagsPatch, IpAccessListEntry, IpAccessListPatch, QueryEndpointRole, ScalingSchedule, - ScalingScheduleEntryRequest, ScalingSchedulePostRequest, ServicPrivateEndpointePostRequest, - Service, ServiceClickhouseSetting, ServiceClickhouseSettingsList, - ServiceClickhouseSettingsPatchRequest, ServiceClickhouseSettingsPatchResponse, - ServiceClickhouseSettingsSchema, ServiceEndpoint, ServiceEndpointChange, - ServiceEndpointChangeProtocol, ServiceEndpointProtocol, ServicePasswordPatchRequest, - ServicePatchRequest, ServicePatchRequestReleasechannel, ServicePostRequest, - ServicePostRequestCompliancetype, ServicePostRequestProfile, ServicePostRequestProvider, - ServicePostRequestRegion, ServicePostRequestReleasechannel, ServiceProfile, - ServiceQueryAPIEndpoint, ServiceReplicaScalingPatchRequest, ServiceState, + AutoscalingMode, BackupEncryptionConfig, InstancePrivateEndpointsPatch, + InstanceServiceQueryApiEndpointsPostRequest, InstanceTagsPatch, IpAccessListEntry, + IpAccessListPatch, QueryEndpointRole, ScalingSchedule, ScalingScheduleEntryRequest, + ScalingSchedulePostRequest, ServicPrivateEndpointePostRequest, Service, + ServiceClickhouseSetting, ServiceClickhouseSettingsList, ServiceClickhouseSettingsPatchRequest, + ServiceClickhouseSettingsPatchResponse, ServiceClickhouseSettingsSchema, ServiceEndpoint, + ServiceEndpointChange, ServiceEndpointChangeProtocol, ServiceEndpointProtocol, + ServicePasswordPatchRequest, ServicePatchRequest, ServicePatchRequestReleasechannel, + ServicePostRequest, ServicePostRequestCompliancetype, ServicePostRequestProfile, + ServicePostRequestProvider, ServicePostRequestRegion, ServicePostRequestReleasechannel, + ServiceProfile, ServiceQueryAPIEndpoint, ServiceReplicaScalingPatchRequest, ServiceState, ServiceStatePatchRequest, ServiceStatePatchRequestCommand, UpgradeWindow, UpgradeWindowPutRequest, UpgradeWindowStartHourUtc, }; @@ -330,6 +330,13 @@ CONTEXT FOR AGENTS: #[arg(long)] backup_id: Option, + /// Backup's encryption_config.json file, or - for stdin (private preview) + /// + /// Only with --backup-id, for a backup of a TDE-enabled service in your own + /// bucket. The restored service has TDE enabled. + #[arg(long, value_name = "PATH", requires = "backup_id")] + backup_encryption_config: Option, + /// Release channel: slow, default, fast (production services only) #[arg(long)] release_channel: Option, @@ -1141,6 +1148,7 @@ pub async fn run(client: &CloudClient, command: ServiceCommands, json: bool) -> idle_timeout_minutes, ip_allow, backup_id, + backup_encryption_config, release_channel, data_warehouse_id, readonly, @@ -1170,6 +1178,10 @@ pub async fn run(client: &CloudClient, command: ServiceCommands, json: bool) -> idle_timeout_minutes, ip_allow, backup_id, + backup_encryption_config: backup_encryption_config + .as_deref() + .map(read_backup_encryption_config) + .transpose()?, release_channel, data_warehouse_id, is_readonly: readonly, @@ -2111,6 +2123,7 @@ struct CreateServiceOptions { idle_timeout_minutes: Option, ip_allow: Vec, backup_id: Option, + backup_encryption_config: Option, release_channel: Option, data_warehouse_id: Option, is_readonly: bool, @@ -2185,6 +2198,41 @@ fn resolve_horizontal_autoscaling( }) } +/// Read a backup's `encryption_config.json` from a file or stdin (`-`), passing +/// the object through unchanged. Every failure is a usage error raised before +/// any request is sent. +fn read_backup_encryption_config(path: &str) -> CloudResult { + use std::io::Read as _; + + let (raw, source) = if path == "-" { + let mut raw = String::new(); + std::io::stdin().read_to_string(&mut raw).map_err(|error| { + CloudError::usage(format!( + "failed to read --backup-encryption-config from stdin: {error}" + )) + })?; + (raw, "stdin".to_string()) + } else { + ( + std::fs::read_to_string(path).map_err(|error| { + CloudError::usage(format!( + "failed to read --backup-encryption-config file '{path}': {error}" + )) + })?, + format!("file '{path}'"), + ) + }; + match serde_json::from_str::(&raw) { + Ok(serde_json::Value::Object(object)) => Ok(object.into_iter().collect()), + Ok(_) => Err(CloudError::usage(format!( + "--backup-encryption-config {source} must contain a JSON object" + ))), + Err(error) => Err(CloudError::usage(format!( + "--backup-encryption-config {source} is not valid JSON: {error}" + ))), + } +} + fn build_create_service_request(options: &CreateServiceOptions) -> CloudResult { let ip_access_list = if options.ip_allow.is_empty() { vec![IpAccessListEntry { @@ -2288,7 +2336,7 @@ fn build_create_service_request(options: &CreateServiceOptions) -> CloudResult Vec<&'static str> { + vec![ + "clickpipe", + "create", + "kafka", + "svc-id", + "--name", + "glue-pipe", + "--brokers", + "broker:9098", + "--topics", + "events", + "--format", + "Avro", + "--database", + "default", + "--table", + "events", + "--column", + "id:Int64", + "--kafka-type", + "msk", + "--org-id", + "org", + "--schema-registry-type", + "glue", + "--glue-region", + "us-east-1", + "--glue-registry-name", + "my-registry", + ] +} + +#[tokio::test] +async fn kafka_create_sends_the_exact_glue_schema_registry_body() { + let mock = start_mock_clickpipes_api().await; + let mut args = kafka_glue_args(); + args.extend([ + "--username", + "u", + "--password", + "p", + "--glue-role-arn", + "arn:aws:iam::123456789012:role/Glue", + ]); + let body = invoke_cli_capture_body(&mock, &args).await; + assert_eq!( + body["source"]["kafka"]["schemaRegistry"], + serde_json::json!({ + "type": "glue", + "glueRegion": "us-east-1", + "glueRegistryName": "my-registry", + "glueRoleArn": "arn:aws:iam::123456789012:role/Glue" + }) + ); + assert_eq!(body["source"]["kafka"]["authentication"], "PLAIN"); + + // With IAM role authentication the Glue role ARN may be omitted: the API + // falls back to the source's iamRole. + let mock = start_mock_clickpipes_api().await; + let mut args = kafka_glue_args(); + args.extend(["--iam-role", "arn:aws:iam::123456789012:role/Source"]); + let body = invoke_cli_capture_body(&mock, &args).await; + assert_eq!( + body["source"]["kafka"]["schemaRegistry"], + serde_json::json!({ + "type": "glue", + "glueRegion": "us-east-1", + "glueRegistryName": "my-registry" + }) + ); + assert_eq!( + body["source"]["kafka"]["iamRole"], + "arn:aws:iam::123456789012:role/Source" + ); +} + +#[tokio::test] +async fn kafka_create_sends_type_only_for_an_explicit_confluent_registry() { + let mock = start_mock_clickpipes_api().await; + let mut args = kafka_args_minimal(); + args.extend([ + "--schema-registry-type", + "confluent", + "--schema-registry-url", + "https://registry.example", + ]); + let body = invoke_cli_capture_body(&mock, &args).await; + assert_eq!( + body["source"]["kafka"]["schemaRegistry"], + serde_json::json!({ + "type": "confluent", + "url": "https://registry.example", + "authentication": "PLAIN", + "credentials": { "username": "", "password": "" } + }) + ); + + let mock = start_mock_clickpipes_api().await; + let mut args = kafka_args_minimal(); + args.extend(["--schema-registry-url", "https://registry.example"]); + let body = invoke_cli_capture_body(&mock, &args).await; + assert!( + body["source"]["kafka"]["schemaRegistry"] + .get("type") + .is_none(), + "{body}" + ); +} + +#[tokio::test] +async fn kafka_glue_schema_registry_misuse_is_a_usage_error_before_any_request() { + let mock = MockServer::start().await; + let cases: [(&[&str], &[&str]); 7] = [ + // Glue conflicts with every Confluent registry flag. + ( + &kafka_glue_args(), + &[ + "--iam-role", + "arn:role", + "--schema-registry-url", + "https://registry.example", + ], + ), + ( + &kafka_glue_args(), + &["--iam-role", "arn:role", "--schema-registry-username", "u"], + ), + // Glue needs a role: from --glue-role-arn or the source's --iam-role. + (&kafka_glue_args(), &["--username", "u", "--password", "p"]), + // Glue requires both region and registry name. + ( + &kafka_args_without_auth(), + &[ + "--schema-registry-type", + "glue", + "--glue-region", + "us-east-1", + ], + ), + // Glue flags require the glue registry type. + (&kafka_args_without_auth(), &["--glue-region", "us-east-1"]), + ( + &kafka_args_without_auth(), + &[ + "--schema-registry-type", + "confluent", + "--schema-registry-url", + "https://registry.example", + "--glue-registry-name", + "my-registry", + ], + ), + // An explicit Confluent registry needs its URL. + ( + &kafka_args_without_auth(), + &["--schema-registry-type", "confluent"], + ), + ]; + for (base, extra) in cases { + let mut args = base.to_vec(); + args.extend(extra); + let output = invoke_cli_with_cloud_credentials(&mock, &args); + assert_eq!( + output.status.code(), + Some(2), + "{extra:?}\nstderr:\n{}", + String::from_utf8_lossy(&output.stderr) + ); + } + assert!(mock.received_requests().await.unwrap().is_empty()); +} + #[tokio::test] async fn kafka_iam_user_credentials_shape() { let mock = start_mock_clickpipes_api().await; @@ -15463,6 +15636,46 @@ async fn mount_clickpipe_get(mock: &MockServer, source: Value) { .await; } +#[tokio::test] +async fn clickpipe_get_renders_a_glue_schema_registry_and_keeps_json_verbatim() { + let source = serde_json::json!({ + "kafka": { + "type": "msk", + "format": "Avro", + "brokers": "broker:9098", + "topics": "events", + "authentication": "IAM_ROLE", + "iamRole": "arn:aws:iam::123456789012:role/Source", + "schemaRegistry": { + "type": "glue", + "glueRegion": "us-east-1", + "glueRegistryName": "my-registry", + "glueRoleArn": "arn:aws:iam::123456789012:role/Glue" + } + } + }); + let mock = MockServer::start().await; + mount_clickpipe_get(&mock, source.clone()).await; + let args = ["clickpipe", "get", "svc-id", "pipe-id", "--org-id", "org"]; + + let output = invoke_cli_with_cloud_credentials(&mock, &args); + assert_success(&output); + let json: Value = serde_json::from_slice(&output.stdout).unwrap(); + assert_eq!(json["source"], source); + + let output = invoke_cli_with_cloud_credentials_human(&mock, &args); + assert_success(&output); + let stdout = String::from_utf8_lossy(&output.stdout); + for line in [ + "type: glue", + "glueRegion: us-east-1", + "glueRegistryName: my-registry", + "glueRoleArn: arn:aws:iam::123456789012:role/Glue", + ] { + assert!(stdout.contains(line), "{line}\n{stdout}"); + } +} + async fn mount_clickpipe_settings_put(mock: &MockServer, result: Value) { let updated_settings = serde_json::json!({ "result": result, @@ -24351,6 +24564,123 @@ async fn service_create_discovers_and_sends_a_dynamic_byoc_profile() { assert_success(&output); } +fn backup_encryption_config_document() -> Value { + serde_json::json!({ + "schema_version": 1, + "restore_key_pairs": [{ + "customer_managed_encryption_key": { + "aws_kms_key_arn": "arn:aws:kms:us-east-1:111122223333:key/1234abcd" + }, + "encrypted_dek": "d3JhcHBlZA==", + "future_field": { "nested": [1, null, true] } + }] + }) +} + +async fn mount_service_create_with_backup_encryption_config(mock: &MockServer) { + Mock::given(method("POST")) + .and(path("/v1/organizations/org-1/services")) + .and(body_json(serde_json::json!({ + "name": "restored", + "provider": "aws", + "region": "us-east-1", + "ipAccessList": [{ + "source": "0.0.0.0/0", + "description": "Allow all (created by clickhousectl)" + }], + "backupId": "a1a2a3a4-b1b2-c1c2-d1d2-e1e2e3e4e5e6", + "backupEncryptionConfig": backup_encryption_config_document() + }))) + .respond_with(ResponseTemplate::new(200).set_body_json(serde_json::json!({ + "result": { + "service": {"id": "22222222-3333-4444-5555-666666666666", "name": "restored"}, + "password": "generated-password" + }, + "status": 200, + "requestId": "stub-restore-service-create" + }))) + .expect(1) + .mount(mock) + .await; +} + +const BACKUP_RESTORE_ARGS: [&str; 7] = [ + "service", + "create", + "--name", + "restored", + "--backup-id", + "a1a2a3a4-b1b2-c1c2-d1d2-e1e2e3e4e5e6", + "--backup-encryption-config", +]; + +#[tokio::test] +async fn service_create_passes_backup_encryption_config_file_through_unchanged() { + let mock = MockServer::start().await; + mount_service_create_with_backup_encryption_config(&mock).await; + let directory = tempfile::tempdir().unwrap(); + let config = directory.path().join("encryption_config.json"); + std::fs::write( + &config, + serde_json::to_string_pretty(&backup_encryption_config_document()).unwrap(), + ) + .unwrap(); + + let mut args = BACKUP_RESTORE_ARGS.to_vec(); + args.extend([config.to_str().unwrap(), "--org-id", "org-1"]); + let output = invoke_cli_with_cloud_credentials(&mock, &args); + assert_success(&output); +} + +#[tokio::test] +async fn service_create_reads_backup_encryption_config_from_stdin() { + let mock = MockServer::start().await; + mount_service_create_with_backup_encryption_config(&mock).await; + + let mut args = BACKUP_RESTORE_ARGS.to_vec(); + args.extend(["-", "--org-id", "org-1"]); + let output = invoke_cli_with_cloud_credentials_and_stdin( + &mock, + &args, + &backup_encryption_config_document().to_string(), + ); + assert_success(&output); +} + +#[tokio::test] +async fn service_create_rejects_bad_backup_encryption_config_before_any_request() { + let mock = MockServer::start().await; + let directory = tempfile::tempdir().unwrap(); + let missing = directory.path().join("missing.json"); + for stdin in ["[1, 2, 3]", "\"just a string\"", "{not json"] { + let mut args = BACKUP_RESTORE_ARGS.to_vec(); + args.extend(["-", "--org-id", "org-1"]); + let output = invoke_cli_with_cloud_credentials_and_stdin(&mock, &args, stdin); + assert_eq!(output.status.code(), Some(2), "{stdin}"); + } + let mut args = BACKUP_RESTORE_ARGS.to_vec(); + args.extend([missing.to_str().unwrap(), "--org-id", "org-1"]); + let output = invoke_cli_with_cloud_credentials(&mock, &args); + assert_eq!(output.status.code(), Some(2)); + + // Clap requires --backup-id alongside the config. + let output = invoke_cli_with_cloud_credentials( + &mock, + &[ + "service", + "create", + "--name", + "restored", + "--backup-encryption-config", + "-", + "--org-id", + "org-1", + ], + ); + assert_eq!(output.status.code(), Some(2)); + assert!(mock.received_requests().await.unwrap().is_empty()); +} + #[tokio::test] async fn service_create_rejects_dynamic_profile_memory_mismatch_before_post() { let mock = MockServer::start().await;