use std::path::PathBuf;
use std::process::Command;
use serde_json::Value;
use wiremock::matchers::{header, method, path, path_regex};
use wiremock::{Mock, MockServer, ResponseTemplate};
fn clickhousectl_binary() -> PathBuf {
PathBuf::from(env!("CARGO_BIN_EXE_clickhousectl"))
}
async fn start_mock_clickpipes_api() -> MockServer {
let mock = MockServer::start().await;
let stub_pipe = serde_json::json!({
"result": {
"id": "00000000-0000-0000-0000-000000000000",
"name": "stub",
"state": "Stopped",
"scaling": { "replicas": 1 },
"source": {},
"destination": { "database": "default" },
"metrics": {},
},
"status": 200,
"requestId": "stub-request-id"
});
Mock::given(method("POST"))
.and(path_regex(
r"^/v1/organizations/[^/]+/services/[^/]+/clickpipes$",
))
.respond_with(ResponseTemplate::new(200).set_body_json(stub_pipe))
.mount(&mock)
.await;
mock
}
fn assert_success(output: &std::process::Output) {
assert!(
output.status.success(),
"clickhousectl exited {}\nstderr:\n{}\nstdout:\n{}",
output.status.code().unwrap_or(-1),
String::from_utf8_lossy(&output.stderr),
String::from_utf8_lossy(&output.stdout),
);
}
async fn invoke_cli_capture_body(mock: &MockServer, cli_args: &[&str]) -> Value {
let mut full_args: Vec<&str> = vec!["cloud", "--url"];
let url = mock.uri();
full_args.push(&url);
full_args.push("--json");
full_args.extend(cli_args);
let output = Command::new(clickhousectl_binary())
.env("DO_NOT_TRACK", "1")
.args(&full_args)
.env("CLICKHOUSE_CLOUD_API_KEY", "fake-key-for-tests")
.env("CLICKHOUSE_CLOUD_API_SECRET", "fake-secret-for-tests")
.output()
.expect("failed to spawn clickhousectl");
assert!(
output.status.success(),
"clickhousectl exited {} for args {:?}\nstderr:\n{}\nstdout:\n{}",
output.status.code().unwrap_or(-1),
full_args,
String::from_utf8_lossy(&output.stderr),
String::from_utf8_lossy(&output.stdout),
);
let requests = mock
.received_requests()
.await
.expect("mock requests log unavailable");
let post = requests
.iter()
.find(|r| r.method == wiremock::http::Method::POST)
.expect("no POST request recorded by mock");
serde_json::from_slice(&post.body).expect("POST body wasn't valid JSON")
}
#[tokio::test]
async fn postgres_cdc_omits_publication_name_and_slot_when_not_passed() {
let mock = start_mock_clickpipes_api().await;
let body = invoke_cli_capture_body(
&mock,
&[
"clickpipe",
"create",
"postgres",
"svc-id",
"--name",
"test-pipe",
"--host",
"pg.example.com",
"--port",
"5432",
"--pg-database",
"test",
"--username",
"u",
"--password",
"p",
"--table-mapping",
"public.t:t",
"--replication-mode",
"cdc",
"--org-id",
"11dfa1ec-767d-43cb-bfad-618ce2aaf959",
],
)
.await;
let settings = &body["source"]["postgres"]["settings"];
assert!(
settings.get("publicationName").is_none(),
"publicationName leaked into wire body: {settings}",
);
assert!(
settings.get("replicationSlotName").is_none(),
"replicationSlotName leaked into wire body: {settings}",
);
}
#[tokio::test]
async fn postgres_destination_omits_table_columns_managed_table_definition() {
let mock = start_mock_clickpipes_api().await;
let body = invoke_cli_capture_body(
&mock,
&[
"clickpipe",
"create",
"postgres",
"svc-id",
"--name",
"test-pipe",
"--host",
"pg.example.com",
"--port",
"5432",
"--pg-database",
"test",
"--username",
"u",
"--password",
"p",
"--table-mapping",
"public.t:t",
"--replication-mode",
"cdc",
"--org-id",
"11dfa1ec-767d-43cb-bfad-618ce2aaf959",
],
)
.await;
let dest = &body["destination"];
assert_eq!(
dest["database"], "default",
"database should default to 'default' for postgres CDC, got {dest}"
);
for field in ["table", "columns", "managedTable", "tableDefinition"] {
assert!(
dest.get(field).is_none(),
"{field} leaked into destination body — Al's Bug 2 regression: {dest}",
);
}
}
#[tokio::test]
async fn mysql_destination_omits_table_columns_managed_table_definition() {
let mock = start_mock_clickpipes_api().await;
let body = invoke_cli_capture_body(
&mock,
&[
"clickpipe",
"create",
"mysql",
"svc-id",
"--name",
"test-pipe",
"--host",
"mysql.example.com",
"--port",
"3306",
"--username",
"u",
"--password",
"p",
"--table-mapping",
"mydb.t:t",
"--replication-mode",
"cdc",
"--org-id",
"11dfa1ec-767d-43cb-bfad-618ce2aaf959",
],
)
.await;
let dest = &body["destination"];
assert_eq!(dest["database"], "default");
for field in ["table", "columns", "managedTable", "tableDefinition"] {
assert!(
dest.get(field).is_none(),
"{field} leaked into MySQL destination body: {dest}",
);
}
}
#[tokio::test]
async fn mongodb_destination_omits_table_columns_managed_table_definition() {
let mock = start_mock_clickpipes_api().await;
let body = invoke_cli_capture_body(
&mock,
&[
"clickpipe",
"create",
"mongodb",
"svc-id",
"--name",
"test-pipe",
"--uri",
"mongodb://mongo.example.com:27017",
"--username",
"u",
"--password",
"p",
"--table-mapping",
"mydb.coll:t",
"--replication-mode",
"cdc",
"--org-id",
"11dfa1ec-767d-43cb-bfad-618ce2aaf959",
],
)
.await;
let dest = &body["destination"];
assert_eq!(dest["database"], "default");
for field in ["table", "columns", "managedTable", "tableDefinition"] {
assert!(
dest.get(field).is_none(),
"{field} leaked into Mongo destination body: {dest}",
);
}
}
#[tokio::test]
async fn s3_pipe_omits_iam_role_and_queue_url_when_not_passed() {
let mock = start_mock_clickpipes_api().await;
let body = invoke_cli_capture_body(
&mock,
&[
"clickpipe",
"create",
"object-storage",
"svc-id",
"--name",
"test-pipe",
"--source-url",
"https://bucket.s3.us-east-1.amazonaws.com/data/*.json",
"--format",
"JSONEachRow",
"--database",
"default",
"--table",
"events",
"--column",
"id:Int64",
"--access-key-id",
"AKIA000000000000FAKE",
"--secret-key",
"fake/secret/for/tests/0000000000000000",
"--org-id",
"11dfa1ec-767d-43cb-bfad-618ce2aaf959",
],
)
.await;
let s3 = &body["source"]["objectStorage"];
for field in [
"iamRole",
"queueUrl",
"connectionString",
"azureContainerName",
"path",
"serviceAccountKey",
"delimiter",
] {
assert!(
s3.get(field).is_none(),
"{field} leaked into S3 body when --{field} not passed: {s3}",
);
}
}
#[tokio::test]
async fn gcs_service_account_file_is_read_and_base64_encoded() {
use std::io::Write;
let mock = start_mock_clickpipes_api().await;
let dir = tempfile::tempdir().unwrap();
let sa_path = dir.path().join("service-account.json");
let sa_contents = br#"{"type":"service_account","project_id":"test"}"#;
let mut sa_file = std::fs::File::create(&sa_path).unwrap();
sa_file.write_all(sa_contents).unwrap();
let body = invoke_cli_capture_body(
&mock,
&[
"clickpipe",
"create",
"object-storage",
"svc-id",
"--name",
"gcs-pipe",
"--source-url",
"https://storage.googleapis.com/bucket/data/*.json",
"--format",
"JSONEachRow",
"--storage-type",
"gcs",
"--database",
"default",
"--table",
"events",
"--column",
"id:Int64",
"--service-account-file",
sa_path.to_str().unwrap(),
"--org-id",
"11dfa1ec-767d-43cb-bfad-618ce2aaf959",
],
)
.await;
let gcs = &body["source"]["objectStorage"];
assert_eq!(gcs["authentication"], "SERVICE_ACCOUNT");
let expected = base64::Engine::encode(&base64::engine::general_purpose::STANDARD, sa_contents);
assert_eq!(
gcs["serviceAccountKey"].as_str(),
Some(expected.as_str()),
"serviceAccountKey on the wire should be base64 of the file contents: {gcs}",
);
}
#[tokio::test]
async fn postgres_optional_fields_absent_when_flags_omitted() {
let mock = start_mock_clickpipes_api().await;
let body = invoke_cli_capture_body(
&mock,
&[
"clickpipe",
"create",
"postgres",
"svc-id",
"--name",
"t",
"--host",
"pg",
"--port",
"5432",
"--pg-database",
"test",
"--username",
"u",
"--password",
"p",
"--table-mapping",
"public.t:t",
"--replication-mode",
"cdc",
"--org-id",
"org",
],
)
.await;
let pg = &body["source"]["postgres"];
for field in ["iamRole", "tlsHost", "caCertificate"] {
assert!(
pg.get(field).is_none(),
"{field} leaked into postgres source body: {pg}",
);
}
}
#[tokio::test]
async fn mysql_optional_fields_absent_when_flags_omitted() {
let mock = start_mock_clickpipes_api().await;
let body = invoke_cli_capture_body(
&mock,
&[
"clickpipe",
"create",
"mysql",
"svc-id",
"--name",
"t",
"--host",
"mysql",
"--port",
"3306",
"--username",
"u",
"--password",
"p",
"--table-mapping",
"mydb.t:t",
"--replication-mode",
"cdc",
"--org-id",
"org",
],
)
.await;
let mysql = &body["source"]["mysql"];
for field in ["iamRole", "tlsHost", "caCertificate"] {
assert!(
mysql.get(field).is_none(),
"{field} leaked into mysql source body: {mysql}",
);
}
}
#[tokio::test]
async fn mongodb_tls_host_absent_when_not_passed() {
let mock = start_mock_clickpipes_api().await;
let body = invoke_cli_capture_body(
&mock,
&[
"clickpipe",
"create",
"mongodb",
"svc-id",
"--name",
"t",
"--uri",
"mongodb://m:27017",
"--username",
"u",
"--password",
"p",
"--table-mapping",
"db.c:t",
"--replication-mode",
"cdc",
"--org-id",
"org",
],
)
.await;
let mongo = &body["source"]["mongodb"];
assert!(
mongo.get("tlsHost").is_none(),
"tlsHost leaked into mongodb source body: {mongo}",
);
assert!(
mongo.get("caCertificate").is_none(),
"caCertificate leaked into mongodb source body: {mongo}",
);
}
fn kafka_args_minimal() -> Vec<&'static str> {
vec![
"clickpipe",
"create",
"kafka",
"svc-id",
"--name",
"t",
"--brokers",
"broker:9092",
"--topics",
"topic",
"--format",
"JSONEachRow",
"--database",
"default",
"--table",
"events",
"--column",
"id:Int64",
"--kafka-type",
"kafka",
"--auth",
"PLAIN",
"--username",
"u",
"--password",
"p",
"--org-id",
"org",
]
}
#[tokio::test]
async fn kafka_optional_fields_absent_when_flags_omitted() {
let mock = start_mock_clickpipes_api().await;
let body = invoke_cli_capture_body(&mock, &kafka_args_minimal()).await;
let kafka = &body["source"]["kafka"];
for field in [
"consumerGroup",
"iamRole",
"schemaRegistry",
"caCertificate",
] {
assert!(
kafka.get(field).is_none(),
"{field} leaked into kafka source body: {kafka}",
);
}
assert!(
kafka["offset"].get("timestamp").is_none(),
"offset.timestamp leaked when --offset-timestamp not passed: {kafka}",
);
}
#[tokio::test]
async fn kafka_plain_credentials_shape() {
let mock = start_mock_clickpipes_api().await;
let body = invoke_cli_capture_body(&mock, &kafka_args_minimal()).await;
let creds = &body["source"]["kafka"]["credentials"];
assert_eq!(creds["username"], "u");
assert_eq!(creds["password"], "p");
}
#[tokio::test]
async fn kafka_scram_sha_512_credentials_shape() {
let mock = start_mock_clickpipes_api().await;
let mut args = kafka_args_minimal();
let auth_idx = args.iter().position(|a| *a == "PLAIN").unwrap();
args[auth_idx] = "SCRAM-SHA-512";
let body = invoke_cli_capture_body(&mock, &args).await;
let creds = &body["source"]["kafka"]["credentials"];
assert_eq!(creds["username"], "u");
assert_eq!(creds["password"], "p");
}
#[tokio::test]
async fn kafka_iam_role_serializes_iam_role_field() {
let mock = start_mock_clickpipes_api().await;
let body = invoke_cli_capture_body(
&mock,
&[
"clickpipe",
"create",
"kafka",
"svc-id",
"--name",
"t",
"--brokers",
"broker:9092",
"--topics",
"topic",
"--format",
"JSONEachRow",
"--database",
"default",
"--table",
"events",
"--column",
"id:Int64",
"--kafka-type",
"msk",
"--auth",
"IAM_ROLE",
"--iam-role",
"arn:aws:iam::123:role/x",
"--org-id",
"org",
],
)
.await;
let kafka = &body["source"]["kafka"];
assert_eq!(kafka["iamRole"], "arn:aws:iam::123:role/x");
assert!(
kafka["credentials"].is_null(),
"IAM_ROLE credentials should be null, got: {}",
kafka["credentials"]
);
}
#[tokio::test]
async fn kafka_iam_user_credentials_shape() {
let mock = start_mock_clickpipes_api().await;
let body = invoke_cli_capture_body(
&mock,
&[
"clickpipe",
"create",
"kafka",
"svc-id",
"--name",
"t",
"--brokers",
"broker:9092",
"--topics",
"topic",
"--format",
"JSONEachRow",
"--database",
"default",
"--table",
"events",
"--column",
"id:Int64",
"--kafka-type",
"msk",
"--auth",
"IAM_USER",
"--access-key-id",
"AKIA000000000000FAKE",
"--secret-key",
"fake/secret/0000000000000000",
"--org-id",
"org",
],
)
.await;
let creds = &body["source"]["kafka"]["credentials"];
assert_eq!(creds["accessKeyId"], "AKIA000000000000FAKE");
assert_eq!(creds["secretKey"], "fake/secret/0000000000000000");
}
#[tokio::test]
async fn kafka_mutual_tls_credentials_use_cert_file_contents() {
use std::io::Write;
let mock = start_mock_clickpipes_api().await;
let dir = tempfile::tempdir().unwrap();
let cert_path = dir.path().join("client.crt");
let key_path = dir.path().join("client.key");
let mut cert_file = std::fs::File::create(&cert_path).unwrap();
let mut key_file = std::fs::File::create(&key_path).unwrap();
cert_file
.write_all(b"-----BEGIN CERTIFICATE-----\nCERT_PEM\n-----END CERTIFICATE-----\n")
.unwrap();
key_file
.write_all(b"-----BEGIN PRIVATE KEY-----\nKEY_PEM\n-----END PRIVATE KEY-----\n")
.unwrap();
let body = invoke_cli_capture_body(
&mock,
&[
"clickpipe",
"create",
"kafka",
"svc-id",
"--name",
"t",
"--brokers",
"broker:9092",
"--topics",
"topic",
"--format",
"JSONEachRow",
"--database",
"default",
"--table",
"events",
"--column",
"id:Int64",
"--kafka-type",
"kafka",
"--auth",
"MUTUAL_TLS",
"--client-certificate",
cert_path.to_str().unwrap(),
"--client-key",
key_path.to_str().unwrap(),
"--org-id",
"org",
],
)
.await;
let creds = &body["source"]["kafka"]["credentials"];
assert!(
creds["certificate"]
.as_str()
.map(|s| s.contains("CERT_PEM"))
.unwrap_or(false),
"MUTUAL_TLS certificate should contain file contents: {creds}",
);
assert!(
creds["privateKey"]
.as_str()
.map(|s| s.contains("KEY_PEM"))
.unwrap_or(false),
"MUTUAL_TLS privateKey should contain file contents: {creds}",
);
}
#[tokio::test]
async fn kinesis_iam_role_omits_access_key() {
let mock = start_mock_clickpipes_api().await;
let body = invoke_cli_capture_body(
&mock,
&[
"clickpipe",
"create",
"kinesis",
"svc-id",
"--name",
"t",
"--stream-name",
"s",
"--region",
"us-east-1",
"--format",
"JSONEachRow",
"--database",
"default",
"--table",
"events",
"--column",
"id:Int64",
"--auth",
"IAM_ROLE",
"--iam-role",
"arn:aws:iam::123:role/x",
"--iterator-type",
"TRIM_HORIZON",
"--org-id",
"org",
],
)
.await;
let kinesis = &body["source"]["kinesis"];
assert_eq!(kinesis["iamRole"], "arn:aws:iam::123:role/x");
assert!(
kinesis.get("accessKey").is_none(),
"accessKey leaked when --auth IAM_ROLE: {kinesis}",
);
}
#[tokio::test]
async fn kinesis_iam_user_omits_iam_role() {
let mock = start_mock_clickpipes_api().await;
let body = invoke_cli_capture_body(
&mock,
&[
"clickpipe",
"create",
"kinesis",
"svc-id",
"--name",
"t",
"--stream-name",
"s",
"--region",
"us-east-1",
"--format",
"JSONEachRow",
"--database",
"default",
"--table",
"events",
"--column",
"id:Int64",
"--auth",
"IAM_USER",
"--access-key-id",
"AKIA000000000000FAKE",
"--secret-key",
"fake/secret/0000000000000000",
"--iterator-type",
"TRIM_HORIZON",
"--org-id",
"org",
],
)
.await;
let kinesis = &body["source"]["kinesis"];
assert_eq!(kinesis["accessKey"]["accessKeyId"], "AKIA000000000000FAKE");
assert!(
kinesis.get("iamRole").is_none(),
"iamRole leaked when --auth IAM_USER: {kinesis}",
);
}
#[tokio::test]
async fn bigquery_destination_omits_table_columns_managed_table_definition() {
use std::io::Write;
let mock = start_mock_clickpipes_api().await;
let dir = tempfile::tempdir().unwrap();
let sa_path = dir.path().join("service-account.json");
let mut sa_file = std::fs::File::create(&sa_path).unwrap();
sa_file
.write_all(
br#"{
"type": "service_account",
"project_id": "test",
"private_key_id": "fake",
"private_key": "-----BEGIN PRIVATE KEY-----\nfake\n-----END PRIVATE KEY-----\n",
"client_email": "fake@test.iam.gserviceaccount.com",
"client_id": "0",
"auth_uri": "https://accounts.google.com/o/oauth2/auth",
"token_uri": "https://oauth2.googleapis.com/token"
}"#,
)
.unwrap();
let body = invoke_cli_capture_body(
&mock,
&[
"clickpipe",
"create",
"bigquery",
"svc-id",
"--name",
"t",
"--service-account-file",
sa_path.to_str().unwrap(),
"--staging-path",
"gs://bucket/staging",
"--table-mapping",
"dataset.t:t",
"--org-id",
"org",
],
)
.await;
let dest = &body["destination"];
assert_eq!(dest["database"], "default");
for field in ["table", "columns", "managedTable", "tableDefinition"] {
assert!(
dest.get(field).is_none(),
"{field} leaked into BigQuery destination body: {dest}",
);
}
}
fn postgres_args_minimal() -> Vec<String> {
[
"clickpipe",
"create",
"postgres",
"svc-id",
"--name",
"t",
"--host",
"pg",
"--port",
"5432",
"--pg-database",
"test",
"--username",
"u",
"--password",
"p",
"--table-mapping",
"public.t:t",
"--replication-mode",
"cdc",
"--org-id",
"org",
]
.iter()
.map(|s| s.to_string())
.collect()
}
#[tokio::test]
async fn postgres_publication_name_serializes_when_provided() {
let mock = start_mock_clickpipes_api().await;
let mut args = postgres_args_minimal();
args.push("--publication-name".into());
args.push("my_pub".into());
let arg_refs: Vec<&str> = args.iter().map(|s| s.as_str()).collect();
let body = invoke_cli_capture_body(&mock, &arg_refs).await;
assert_eq!(
body["source"]["postgres"]["settings"]["publicationName"], "my_pub",
"publicationName should round-trip the user-provided value"
);
}
#[tokio::test]
async fn postgres_replication_slot_name_serializes_when_provided() {
let mock = start_mock_clickpipes_api().await;
let mut args = postgres_args_minimal();
args.push("--replication-slot-name".into());
args.push("my_slot".into());
let arg_refs: Vec<&str> = args.iter().map(|s| s.as_str()).collect();
let body = invoke_cli_capture_body(&mock, &arg_refs).await;
assert_eq!(
body["source"]["postgres"]["settings"]["replicationSlotName"], "my_slot",
"replicationSlotName should round-trip the user-provided value"
);
}
#[tokio::test]
async fn postgres_tls_host_serializes_when_provided() {
let mock = start_mock_clickpipes_api().await;
let mut args = postgres_args_minimal();
args.push("--tls-host".into());
args.push("pg.example.com".into());
let arg_refs: Vec<&str> = args.iter().map(|s| s.as_str()).collect();
let body = invoke_cli_capture_body(&mock, &arg_refs).await;
assert_eq!(
body["source"]["postgres"]["tlsHost"], "pg.example.com",
"tlsHost should round-trip the user-provided value"
);
}
#[tokio::test]
async fn postgres_iam_role_serializes_when_provided() {
let mock = start_mock_clickpipes_api().await;
let mut args = postgres_args_minimal();
args.push("--iam-role".into());
args.push("arn:aws:iam::123:role/x".into());
let arg_refs: Vec<&str> = args.iter().map(|s| s.as_str()).collect();
let body = invoke_cli_capture_body(&mock, &arg_refs).await;
assert_eq!(
body["source"]["postgres"]["iamRole"], "arn:aws:iam::123:role/x",
"iamRole should round-trip the user-provided value"
);
}
#[tokio::test]
async fn postgres_ca_certificate_file_contents_flow_to_body() {
use std::io::Write;
let mock = start_mock_clickpipes_api().await;
let dir = tempfile::tempdir().unwrap();
let ca_path = dir.path().join("ca.pem");
let pem = "-----BEGIN CERTIFICATE-----\nCA_PEM_CONTENT\n-----END CERTIFICATE-----\n";
std::fs::File::create(&ca_path)
.unwrap()
.write_all(pem.as_bytes())
.unwrap();
let mut args = postgres_args_minimal();
args.push("--ca-certificate".into());
args.push(ca_path.to_str().unwrap().to_string());
let arg_refs: Vec<&str> = args.iter().map(|s| s.as_str()).collect();
let body = invoke_cli_capture_body(&mock, &arg_refs).await;
assert!(
body["source"]["postgres"]["caCertificate"]
.as_str()
.map(|s| s.contains("CA_PEM_CONTENT"))
.unwrap_or(false),
"caCertificate body should contain the file's PEM content, got {}",
body["source"]["postgres"]["caCertificate"]
);
}
#[tokio::test]
async fn postgres_replication_mode_snapshot_serializes() {
let mock = start_mock_clickpipes_api().await;
let mut args = postgres_args_minimal();
let idx = args.iter().position(|a| a == "cdc").unwrap();
args[idx] = "snapshot".into();
let arg_refs: Vec<&str> = args.iter().map(|s| s.as_str()).collect();
let body = invoke_cli_capture_body(&mock, &arg_refs).await;
assert_eq!(
body["source"]["postgres"]["settings"]["replicationMode"],
"snapshot",
);
}
#[tokio::test]
async fn postgres_replication_mode_cdc_only_serializes() {
let mock = start_mock_clickpipes_api().await;
let mut args = postgres_args_minimal();
let idx = args.iter().position(|a| a == "cdc").unwrap();
args[idx] = "cdc_only".into();
args.push("--publication-name".into());
args.push("p".into());
args.push("--replication-slot-name".into());
args.push("s".into());
let arg_refs: Vec<&str> = args.iter().map(|s| s.as_str()).collect();
let body = invoke_cli_capture_body(&mock, &arg_refs).await;
assert_eq!(
body["source"]["postgres"]["settings"]["replicationMode"],
"cdc_only",
);
}
#[tokio::test]
async fn postgres_multiple_table_mappings_serialize_as_array() {
let mock = start_mock_clickpipes_api().await;
let mut args = postgres_args_minimal();
args.push("--table-mapping".into());
args.push("public.t2:t2_dst".into());
args.push("--table-mapping".into());
args.push("other_schema.t3:t3_dst".into());
let arg_refs: Vec<&str> = args.iter().map(|s| s.as_str()).collect();
let body = invoke_cli_capture_body(&mock, &arg_refs).await;
let mappings = body["source"]["postgres"]["tableMappings"]
.as_array()
.unwrap_or_else(|| {
panic!(
"tableMappings should be an array, got: {}",
body["source"]["postgres"]["tableMappings"]
)
});
assert_eq!(
mappings.len(),
3,
"expected 3 table mappings (minimal default + 2 added), got {}: {:?}",
mappings.len(),
mappings
);
let target_tables: Vec<&str> = mappings
.iter()
.filter_map(|m| m["targetTable"].as_str())
.collect();
assert!(target_tables.contains(&"t"));
assert!(target_tables.contains(&"t2_dst"));
assert!(target_tables.contains(&"t3_dst"));
}
macro_rules! postgres_type_test {
($test_name:ident, $cli_value:literal, $wire_value:literal) => {
#[tokio::test]
async fn $test_name() {
let mock = start_mock_clickpipes_api().await;
let mut args = postgres_args_minimal();
args.push("--postgres-type".into());
args.push($cli_value.into());
let arg_refs: Vec<&str> = args.iter().map(|s| s.as_str()).collect();
let body = invoke_cli_capture_body(&mock, &arg_refs).await;
assert_eq!(
body["source"]["postgres"]["type"], $wire_value,
"--postgres-type {} should serialize to wire value {}",
$cli_value, $wire_value,
);
}
};
}
postgres_type_test!(postgres_type_postgres_serializes, "postgres", "postgres");
postgres_type_test!(postgres_type_supabase_serializes, "supabase", "supabase");
postgres_type_test!(postgres_type_neon_serializes, "neon", "neon");
postgres_type_test!(postgres_type_alloydb_serializes, "alloydb", "alloydb");
postgres_type_test!(
postgres_type_planetscale_serializes,
"planetscale",
"planetscale"
);
postgres_type_test!(
postgres_type_rdspostgres_serializes,
"rdspostgres",
"rdspostgres"
);
postgres_type_test!(
postgres_type_aurorapostgres_serializes,
"aurorapostgres",
"aurorapostgres"
);
postgres_type_test!(
postgres_type_cloudsqlpostgres_serializes,
"cloudsqlpostgres",
"cloudsqlpostgres"
);
postgres_type_test!(
postgres_type_azurepostgres_serializes,
"azurepostgres",
"azurepostgres"
);
postgres_type_test!(
postgres_type_crunchybridge_serializes,
"crunchybridge",
"crunchybridge"
);
postgres_type_test!(postgres_type_tigerdata_serializes, "tigerdata", "tigerdata");
#[tokio::test]
async fn dotenv_creds_produce_basic_auth_request() {
use std::io::Write;
let mock = MockServer::start().await;
let stub_orgs = serde_json::json!({
"result": [],
"status": 200,
"requestId": "stub-org-list",
});
Mock::given(method("GET"))
.and(path("/v1/organizations"))
.respond_with(ResponseTemplate::new(200).set_body_json(stub_orgs))
.mount(&mock)
.await;
let dir = tempfile::tempdir().unwrap();
let mut env_file = std::fs::File::create(dir.path().join(".env")).unwrap();
env_file
.write_all(
b"CLICKHOUSE_CLOUD_API_KEY=dotenv-key\nCLICKHOUSE_CLOUD_API_SECRET=dotenv-secret\n",
)
.unwrap();
drop(env_file);
let url = mock.uri();
let output = Command::new(clickhousectl_binary())
.env("DO_NOT_TRACK", "1")
.args(["cloud", "--url", &url, "--json", "org", "list"])
.current_dir(dir.path())
.env_remove("CLICKHOUSE_CLOUD_API_KEY")
.env_remove("CLICKHOUSE_CLOUD_API_SECRET")
.output()
.expect("failed to spawn clickhousectl");
assert_success(&output);
let requests = mock
.received_requests()
.await
.expect("mock requests log unavailable");
let auth = requests
.iter()
.find(|r| r.method == wiremock::http::Method::GET)
.and_then(|r| r.headers.get("Authorization"))
.expect("no Authorization header recorded");
let auth_str = auth.to_str().expect("non-utf8 auth header");
let expected = format!(
"Basic {}",
base64::Engine::encode(
&base64::engine::general_purpose::STANDARD,
"dotenv-key:dotenv-secret",
)
);
assert_eq!(
auth_str, expected,
"Authorization header should match the .env credentials exactly"
);
}
const QUERY_TEST_SERVICE_ID: &str = "11111111-2222-3333-4444-555555555555";
async fn start_mock_control_plane_with_service() -> MockServer {
let mock = MockServer::start().await;
let stub_service = serde_json::json!({
"result": { "id": QUERY_TEST_SERVICE_ID, "name": "demo" },
"status": 200,
"requestId": "stub-service-get",
});
Mock::given(method("GET"))
.and(path(format!(
"/v1/organizations/org-1/services/{QUERY_TEST_SERVICE_ID}"
)))
.respond_with(ResponseTemplate::new(200).set_body_json(stub_service))
.mount(&mock)
.await;
mock
}
async fn start_mock_query_host() -> MockServer {
let mock = MockServer::start().await;
Mock::given(method("POST"))
.and(path(format!("/service/{QUERY_TEST_SERVICE_ID}/run")))
.respond_with(ResponseTemplate::new(200).set_body_string("1\n"))
.mount(&mock)
.await;
mock
}
#[tokio::test]
async fn service_query_with_oauth_sends_bearer_and_never_provisions() {
let control = start_mock_control_plane_with_service().await;
let query_host = start_mock_query_host().await;
let dir = tempfile::tempdir().unwrap();
let home_dir = dir.path().join("home");
let ch_dir = home_dir.join(".clickhouse");
std::fs::create_dir_all(&ch_dir).unwrap();
let tokens = serde_json::json!({
"access_token": "test-bearer-token",
"refresh_token": "unused",
"expires_at": 4102444800u64, "api_url": format!("{}/v1", control.uri()),
});
std::fs::write(
ch_dir.join("tokens.json"),
serde_json::to_vec(&tokens).unwrap(),
)
.unwrap();
let url = control.uri();
let output = Command::new(clickhousectl_binary())
.env("DO_NOT_TRACK", "1")
.args([
"cloud",
"--url",
&url,
"service",
"query",
"--id",
QUERY_TEST_SERVICE_ID,
"--org-id",
"org-1",
"--query",
"SELECT 1",
])
.current_dir(dir.path())
.env("HOME", &home_dir)
.env_remove("CLICKHOUSE_CLOUD_API_KEY")
.env_remove("CLICKHOUSE_CLOUD_API_SECRET")
.env("CLICKHOUSE_CLOUD_QUERY_HOST", query_host.uri())
.output()
.expect("failed to spawn clickhousectl");
assert_success(&output);
let query_requests = query_host.received_requests().await.unwrap();
assert_eq!(query_requests.len(), 1);
let run = &query_requests[0];
let auth = run.headers.get("authorization").unwrap().to_str().unwrap();
assert_eq!(auth, "Bearer test-bearer-token");
assert!(
run.headers.get("auth-provider").is_none(),
"auth-provider header must not accompany a bearer token",
);
let control_requests = control.received_requests().await.unwrap();
assert!(
control_requests
.iter()
.all(|r| r.method == wiremock::http::Method::GET),
"OAuth service query made non-GET control-plane calls: {:?}",
control_requests
.iter()
.map(|r| format!("{} {}", r.method, r.url.path()))
.collect::<Vec<_>>(),
);
assert!(
!ch_dir.join("credentials.json").exists(),
"OAuth service query wrote .clickhouse/credentials.json",
);
assert_eq!(String::from_utf8_lossy(&output.stdout), "1\n");
}
#[tokio::test]
async fn service_query_with_stored_key_sends_basic_auth_with_that_key() {
let control = start_mock_control_plane_with_service().await;
let query_host = start_mock_query_host().await;
let dir = tempfile::tempdir().unwrap();
let ch_dir = dir.path().join(".clickhouse");
std::fs::create_dir_all(&ch_dir).unwrap();
let creds = serde_json::json!({
"service_query_keys": {
QUERY_TEST_SERVICE_ID: {
"key_id": "stored-key-id",
"key_secret": "stored-key-secret",
"endpoint_id": "ep-1",
"service_name": "demo",
"created_at": "2026-05-11T12:00:00Z",
}
}
});
std::fs::write(
ch_dir.join("credentials.json"),
serde_json::to_vec(&creds).unwrap(),
)
.unwrap();
let url = control.uri();
let output = Command::new(clickhousectl_binary())
.env("DO_NOT_TRACK", "1")
.args([
"cloud",
"--url",
&url,
"service",
"query",
"--id",
QUERY_TEST_SERVICE_ID,
"--org-id",
"org-1",
"--query",
"SELECT 1",
])
.current_dir(dir.path())
.env("CLICKHOUSE_CLOUD_API_KEY", "fake-key-for-tests")
.env("CLICKHOUSE_CLOUD_API_SECRET", "fake-secret-for-tests")
.env("CLICKHOUSE_CLOUD_QUERY_HOST", query_host.uri())
.output()
.expect("failed to spawn clickhousectl");
assert_success(&output);
let query_requests = query_host.received_requests().await.unwrap();
assert_eq!(query_requests.len(), 1);
let run = &query_requests[0];
let auth = run.headers.get("authorization").unwrap().to_str().unwrap();
let expected = format!(
"Basic {}",
base64::Engine::encode(
&base64::engine::general_purpose::STANDARD,
"stored-key-id:stored-key-secret",
)
);
assert_eq!(auth, expected);
assert_eq!(run.headers.get("auth-provider").unwrap(), "custom");
let control_requests = control.received_requests().await.unwrap();
assert!(
control_requests
.iter()
.all(|r| r.method == wiremock::http::Method::GET),
"stored-key service query made non-GET control-plane calls: {:?}",
control_requests
.iter()
.map(|r| format!("{} {}", r.method, r.url.path()))
.collect::<Vec<_>>(),
);
}
const QUERY_TEST_KEY_UUID: &str = "aaaaaaaa-bbbb-cccc-dddd-eeeeeeeeeeee";
async fn mount_key_create_and_delete(control: &MockServer, result: Value) {
Mock::given(method("POST"))
.and(path("/v1/organizations/org-1/keys"))
.respond_with(ResponseTemplate::new(200).set_body_json(serde_json::json!({
"result": result,
"status": 200,
"requestId": "stub-key-create",
})))
.mount(control)
.await;
Mock::given(method("DELETE"))
.and(path(format!(
"/v1/organizations/org-1/keys/{QUERY_TEST_KEY_UUID}"
)))
.respond_with(ResponseTemplate::new(200).set_body_json(serde_json::json!({
"status": 200,
"requestId": "stub-key-delete",
})))
.mount(control)
.await;
}
fn invoke_service_query_provisioning(control: &MockServer) -> (tempfile::TempDir, String) {
let dir = tempfile::tempdir().unwrap();
let url = control.uri();
let output = Command::new(clickhousectl_binary())
.env("DO_NOT_TRACK", "1")
.args([
"cloud",
"--url",
&url,
"service",
"query",
"--id",
QUERY_TEST_SERVICE_ID,
"--org-id",
"org-1",
"--query",
"SELECT 1",
])
.current_dir(dir.path())
.env("CLICKHOUSE_CLOUD_API_KEY", "fake-key-for-tests")
.env("CLICKHOUSE_CLOUD_API_SECRET", "fake-secret-for-tests")
.env("CLICKHOUSE_CLOUD_QUERY_HOST", "http://127.0.0.1:1")
.output()
.expect("failed to spawn clickhousectl");
assert!(
!output.status.success(),
"provisioning with an incomplete response must fail\nstdout:\n{}",
String::from_utf8_lossy(&output.stdout),
);
(dir, String::from_utf8_lossy(&output.stderr).to_string())
}
async fn recorded_key_deletes(control: &MockServer) -> Vec<String> {
control
.received_requests()
.await
.unwrap()
.iter()
.filter(|r| r.method == wiremock::http::Method::DELETE)
.map(|r| r.url.path().to_string())
.collect()
}
#[tokio::test]
async fn service_query_deletes_the_key_when_the_create_response_omits_the_secret() {
let control = start_mock_control_plane_with_service().await;
mount_key_create_and_delete(
&control,
serde_json::json!({
"key": { "id": QUERY_TEST_KEY_UUID },
"keyId": "provisioned-key-id",
}),
)
.await;
let (dir, stderr) = invoke_service_query_provisioning(&control);
assert!(
stderr.contains("keySecret"),
"stderr should name the missing field:\n{stderr}",
);
assert_eq!(
recorded_key_deletes(&control).await,
vec![format!(
"/v1/organizations/org-1/keys/{QUERY_TEST_KEY_UUID}"
)],
"the unusable key must be deleted exactly once",
);
let upserts = control
.received_requests()
.await
.unwrap()
.iter()
.filter(|r| r.url.path().ends_with("/serviceQueryEndpoint"))
.count();
assert_eq!(upserts, 0, "a keyless credential must not be bound");
assert!(!dir.path().join(".clickhouse/credentials.json").exists());
}
#[tokio::test]
async fn service_query_keeps_the_key_when_the_endpoint_response_omits_the_id() {
let control = start_mock_control_plane_with_service().await;
let query_host = start_mock_query_host().await;
mount_key_create_and_delete(
&control,
serde_json::json!({
"key": { "id": QUERY_TEST_KEY_UUID },
"keyId": "provisioned-key-id",
"keySecret": "provisioned-key-secret",
}),
)
.await;
let endpoint_path =
format!("/v1/organizations/org-1/services/{QUERY_TEST_SERVICE_ID}/serviceQueryEndpoint");
Mock::given(method("GET"))
.and(path(endpoint_path.clone()))
.respond_with(ResponseTemplate::new(404).set_body_json(serde_json::json!({
"error": "not found",
"status": 404,
"requestId": "stub-endpoint-get",
})))
.mount(&control)
.await;
Mock::given(method("POST"))
.and(path(endpoint_path))
.respond_with(ResponseTemplate::new(200).set_body_json(serde_json::json!({
"result": { "roles": ["sql_console_admin"] },
"status": 200,
"requestId": "stub-endpoint-upsert",
})))
.mount(&control)
.await;
let dir = tempfile::tempdir().unwrap();
let url = control.uri();
let output = Command::new(clickhousectl_binary())
.env("DO_NOT_TRACK", "1")
.args([
"cloud",
"--url",
&url,
"service",
"query",
"--id",
QUERY_TEST_SERVICE_ID,
"--org-id",
"org-1",
"--query",
"SELECT 1",
])
.current_dir(dir.path())
.env("CLICKHOUSE_CLOUD_API_KEY", "fake-key-for-tests")
.env("CLICKHOUSE_CLOUD_API_SECRET", "fake-secret-for-tests")
.env("CLICKHOUSE_CLOUD_QUERY_HOST", query_host.uri())
.output()
.expect("failed to spawn clickhousectl");
assert_success(&output);
assert!(
recorded_key_deletes(&control).await.is_empty(),
"a bound, usable key must not be discarded over an unused echoed id",
);
let stored: Value = serde_json::from_slice(
&std::fs::read(dir.path().join(".clickhouse/credentials.json")).unwrap(),
)
.unwrap();
let key = &stored["service_query_keys"][QUERY_TEST_SERVICE_ID];
assert_eq!(key["key_id"], "provisioned-key-id");
assert_eq!(key["key_secret"], "provisioned-key-secret");
assert!(
key.get("endpoint_id").is_none(),
"an absent endpoint id must not be stored: {stored}",
);
}
#[tokio::test]
async fn service_query_deletes_the_key_when_the_endpoint_get_omits_open_api_keys() {
let control = start_mock_control_plane_with_service().await;
mount_key_create_and_delete(
&control,
serde_json::json!({
"key": { "id": QUERY_TEST_KEY_UUID },
"keyId": "provisioned-key-id",
"keySecret": "provisioned-key-secret",
}),
)
.await;
let endpoint_path =
format!("/v1/organizations/org-1/services/{QUERY_TEST_SERVICE_ID}/serviceQueryEndpoint");
Mock::given(method("GET"))
.and(path(endpoint_path.clone()))
.respond_with(ResponseTemplate::new(200).set_body_json(serde_json::json!({
"result": { "id": "ep-1", "roles": ["sql_console_admin"] },
"status": 200,
"requestId": "stub-endpoint-get",
})))
.mount(&control)
.await;
Mock::given(method("POST"))
.and(path(endpoint_path.clone()))
.respond_with(ResponseTemplate::new(200).set_body_json(serde_json::json!({
"result": { "id": "ep-1" },
"status": 200,
"requestId": "stub-endpoint-upsert",
})))
.mount(&control)
.await;
let (dir, stderr) = invoke_service_query_provisioning(&control);
let upserts = control
.received_requests()
.await
.unwrap()
.iter()
.filter(|r| r.method == wiremock::http::Method::POST && r.url.path() == endpoint_path)
.count();
assert_eq!(
upserts, 0,
"the endpoint must not be rebound from an unknown key list",
);
assert_eq!(
recorded_key_deletes(&control).await,
vec![format!(
"/v1/organizations/org-1/keys/{QUERY_TEST_KEY_UUID}"
)],
"the unbindable key must be deleted exactly once",
);
assert!(!dir.path().join(".clickhouse/credentials.json").exists());
assert!(
stderr.contains("'openApiKeys'"),
"stderr should name the omitted field:\n{stderr}",
);
}
async fn provision_against_endpoint_with_keys(existing_keys: Value) -> (tempfile::TempDir, Value) {
let control = start_mock_control_plane_with_service().await;
let query_host = start_mock_query_host().await;
mount_key_create_and_delete(
&control,
serde_json::json!({
"key": { "id": QUERY_TEST_KEY_UUID },
"keyId": "provisioned-key-id",
"keySecret": "provisioned-key-secret",
}),
)
.await;
let endpoint_path =
format!("/v1/organizations/org-1/services/{QUERY_TEST_SERVICE_ID}/serviceQueryEndpoint");
Mock::given(method("GET"))
.and(path(endpoint_path.clone()))
.respond_with(ResponseTemplate::new(200).set_body_json(serde_json::json!({
"result": { "id": "ep-1", "openApiKeys": existing_keys },
"status": 200,
"requestId": "stub-endpoint-get",
})))
.mount(&control)
.await;
Mock::given(method("POST"))
.and(path(endpoint_path.clone()))
.respond_with(ResponseTemplate::new(200).set_body_json(serde_json::json!({
"result": { "id": "ep-1" },
"status": 200,
"requestId": "stub-endpoint-upsert",
})))
.mount(&control)
.await;
let dir = tempfile::tempdir().unwrap();
let url = control.uri();
let output = Command::new(clickhousectl_binary())
.env("DO_NOT_TRACK", "1")
.args([
"cloud",
"--url",
&url,
"service",
"query",
"--id",
QUERY_TEST_SERVICE_ID,
"--org-id",
"org-1",
"--query",
"SELECT 1",
])
.current_dir(dir.path())
.env("CLICKHOUSE_CLOUD_API_KEY", "fake-key-for-tests")
.env("CLICKHOUSE_CLOUD_API_SECRET", "fake-secret-for-tests")
.env("CLICKHOUSE_CLOUD_QUERY_HOST", query_host.uri())
.output()
.expect("failed to spawn clickhousectl");
assert_success(&output);
assert!(
recorded_key_deletes(&control).await.is_empty(),
"a successfully bound key must not be discarded",
);
let upsert = control
.received_requests()
.await
.unwrap()
.into_iter()
.find(|r| r.method == wiremock::http::Method::POST && r.url.path() == endpoint_path)
.expect("the endpoint upsert must be sent");
let body: Value = serde_json::from_slice(&upsert.body).unwrap();
(dir, body["openApiKeys"].clone())
}
#[tokio::test]
async fn service_query_binds_the_new_key_when_the_endpoint_reports_no_keys() {
let (dir, sent_keys) = provision_against_endpoint_with_keys(serde_json::json!([])).await;
assert_eq!(sent_keys, serde_json::json!([QUERY_TEST_KEY_UUID]));
let stored: Value = serde_json::from_slice(
&std::fs::read(dir.path().join(".clickhouse/credentials.json")).unwrap(),
)
.unwrap();
assert_eq!(
stored["service_query_keys"][QUERY_TEST_SERVICE_ID]["endpoint_id"], "ep-1",
"an echoed endpoint id is recorded",
);
}
#[tokio::test]
async fn service_query_merges_the_new_key_into_the_reported_keys() {
let existing = "99999999-8888-7777-6666-555555555555";
let (_dir, sent_keys) =
provision_against_endpoint_with_keys(serde_json::json!([existing])).await;
assert_eq!(
sent_keys,
serde_json::json!([existing, QUERY_TEST_KEY_UUID]),
"an existing binding must survive the upsert",
);
}
fn write_oauth_tokens(ch_dir: &std::path::Path, control_uri: &str) {
let tokens = serde_json::json!({
"access_token": "test-bearer-token",
"refresh_token": "unused",
"expires_at": 4102444800u64, "api_url": format!("{control_uri}/v1"),
});
std::fs::write(
ch_dir.join("tokens.json"),
serde_json::to_vec(&tokens).unwrap(),
)
.unwrap();
}
#[tokio::test]
async fn service_query_resends_with_wake_header_when_service_is_idle() {
let control = start_mock_control_plane_with_service().await;
let query_host = MockServer::start().await;
Mock::given(method("POST"))
.and(path(format!("/service/{QUERY_TEST_SERVICE_ID}/run")))
.and(header("wake-service", "true"))
.respond_with(ResponseTemplate::new(200).set_body_string("1\n"))
.with_priority(1)
.mount(&query_host)
.await;
Mock::given(method("POST"))
.and(path(format!("/service/{QUERY_TEST_SERVICE_ID}/run")))
.respond_with(
ResponseTemplate::new(206).set_body_string(r#"{"data":"Confirm wake service"}"#),
)
.with_priority(5)
.mount(&query_host)
.await;
let dir = tempfile::tempdir().unwrap();
let home_dir = dir.path().join("home");
let ch_dir = home_dir.join(".clickhouse");
std::fs::create_dir_all(&ch_dir).unwrap();
write_oauth_tokens(&ch_dir, &control.uri());
let url = control.uri();
let output = Command::new(clickhousectl_binary())
.env("DO_NOT_TRACK", "1")
.args([
"cloud",
"--url",
&url,
"service",
"query",
"--id",
QUERY_TEST_SERVICE_ID,
"--org-id",
"org-1",
"--query",
"SELECT 1",
])
.current_dir(dir.path())
.env("HOME", &home_dir)
.env_remove("CLICKHOUSE_CLOUD_API_KEY")
.env_remove("CLICKHOUSE_CLOUD_API_SECRET")
.env("CLICKHOUSE_CLOUD_QUERY_HOST", query_host.uri())
.output()
.expect("failed to spawn clickhousectl");
assert_success(&output);
let query_requests = query_host.received_requests().await.unwrap();
assert_eq!(query_requests.len(), 2);
assert!(
query_requests[0].headers.get("wake-service").is_none(),
"first attempt must not pre-emptively wake the service",
);
assert_eq!(
query_requests[1].headers.get("wake-service").unwrap(),
"true"
);
assert_eq!(String::from_utf8_lossy(&output.stdout), "1\n");
assert!(
String::from_utf8_lossy(&output.stderr).contains("idle"),
"stderr should mention the service is idle:\n{}",
String::from_utf8_lossy(&output.stderr),
);
}
#[tokio::test]
async fn service_query_fails_with_start_hint_when_service_is_stopped() {
let control = start_mock_control_plane_with_service().await;
let query_host = MockServer::start().await;
Mock::given(method("POST"))
.and(path(format!("/service/{QUERY_TEST_SERVICE_ID}/run")))
.respond_with(
ResponseTemplate::new(206).set_body_string(r#"{"data":"Service is stopped"}"#),
)
.mount(&query_host)
.await;
let dir = tempfile::tempdir().unwrap();
let home_dir = dir.path().join("home");
let ch_dir = home_dir.join(".clickhouse");
std::fs::create_dir_all(&ch_dir).unwrap();
write_oauth_tokens(&ch_dir, &control.uri());
let url = control.uri();
let output = Command::new(clickhousectl_binary())
.env("DO_NOT_TRACK", "1")
.args([
"cloud",
"--url",
&url,
"service",
"query",
"--id",
QUERY_TEST_SERVICE_ID,
"--org-id",
"org-1",
"--query",
"SELECT 1",
])
.current_dir(dir.path())
.env("HOME", &home_dir)
.env_remove("CLICKHOUSE_CLOUD_API_KEY")
.env_remove("CLICKHOUSE_CLOUD_API_SECRET")
.env("CLICKHOUSE_CLOUD_QUERY_HOST", query_host.uri())
.output()
.expect("failed to spawn clickhousectl");
assert!(
!output.status.success(),
"querying a stopped service must fail\nstdout:\n{}",
String::from_utf8_lossy(&output.stdout),
);
let stderr = String::from_utf8_lossy(&output.stderr);
assert!(
stderr.contains("stopped") && stderr.contains("service start"),
"stderr should say the service is stopped and hint at `service start`:\n{stderr}",
);
let query_requests = query_host.received_requests().await.unwrap();
assert_eq!(query_requests.len(), 1);
}
#[tokio::test]
async fn shell_env_overrides_dotenv_creds_in_request() {
use std::io::Write;
let mock = MockServer::start().await;
let stub_orgs = serde_json::json!({
"result": [],
"status": 200,
"requestId": "stub-org-list",
});
Mock::given(method("GET"))
.and(path("/v1/organizations"))
.respond_with(ResponseTemplate::new(200).set_body_json(stub_orgs))
.mount(&mock)
.await;
let dir = tempfile::tempdir().unwrap();
let mut env_file = std::fs::File::create(dir.path().join(".env")).unwrap();
env_file
.write_all(
b"CLICKHOUSE_CLOUD_API_KEY=dotenv-key\nCLICKHOUSE_CLOUD_API_SECRET=dotenv-secret\n",
)
.unwrap();
drop(env_file);
let url = mock.uri();
let output = Command::new(clickhousectl_binary())
.env("DO_NOT_TRACK", "1")
.args(["cloud", "--url", &url, "--json", "org", "list"])
.current_dir(dir.path())
.env("CLICKHOUSE_CLOUD_API_KEY", "shell-key")
.env("CLICKHOUSE_CLOUD_API_SECRET", "shell-secret")
.output()
.expect("failed to spawn clickhousectl");
assert_success(&output);
let requests = mock.received_requests().await.unwrap();
let auth = requests
.iter()
.find(|r| r.method == wiremock::http::Method::GET)
.and_then(|r| r.headers.get("Authorization"))
.expect("no Authorization header recorded");
let auth_str = auth.to_str().expect("non-utf8 auth header");
let expected = format!(
"Basic {}",
base64::Engine::encode(
&base64::engine::general_purpose::STANDARD,
"shell-key:shell-secret",
)
);
assert_eq!(
auth_str, expected,
"shell env vars must override .env values on the wire"
);
}
#[tokio::test]
async fn agent_session_and_trace_headers_are_forwarded() {
let mock = MockServer::start().await;
let stub_orgs = serde_json::json!({
"result": [],
"status": 200,
"requestId": "stub-org-list",
});
Mock::given(method("GET"))
.and(path("/v1/organizations"))
.respond_with(ResponseTemplate::new(200).set_body_json(stub_orgs))
.mount(&mock)
.await;
let url = mock.uri();
let traceparent = "00-0af7651916cd43dd8448eb211c80319c-b7ad6b7169203331-01";
let output = Command::new(clickhousectl_binary())
.env("DO_NOT_TRACK", "1")
.args(["cloud", "--url", &url, "--json", "org", "list"])
.env("CLICKHOUSE_CLOUD_API_KEY", "fake-key-for-tests")
.env("CLICKHOUSE_CLOUD_API_SECRET", "fake-secret-for-tests")
.env("AGENT", "claude-code")
.env("CLAUDE_CODE_SESSION_ID", "sess-test-267")
.env("TRACEPARENT", traceparent)
.output()
.expect("failed to spawn clickhousectl");
assert_success(&output);
let requests = mock
.received_requests()
.await
.expect("mock requests log unavailable");
let req = requests
.iter()
.find(|r| r.method == wiremock::http::Method::GET)
.expect("no GET request recorded");
assert_eq!(
req.headers
.get("agent-session-id")
.expect("agent-session-id header missing")
.to_str()
.unwrap(),
"sess-test-267",
);
assert_eq!(
req.headers
.get("traceparent")
.expect("traceparent header missing")
.to_str()
.unwrap(),
traceparent,
);
}
#[tokio::test]
async fn s3_skip_initial_load_serializes_when_passed() {
let mock = start_mock_clickpipes_api().await;
let body = invoke_cli_capture_body(
&mock,
&[
"clickpipe",
"create",
"object-storage",
"svc-id",
"--name",
"t",
"--source-url",
"https://bucket.s3.us-east-1.amazonaws.com/data/*.json",
"--format",
"JSONEachRow",
"--database",
"d",
"--table",
"t",
"--column",
"id:Int64",
"--continuous",
"--queue-url",
"https://sqs.us-east-1.amazonaws.com/123/q",
"--skip-initial-load",
"--org-id",
"org",
],
)
.await;
let s3 = &body["source"]["objectStorage"];
assert_eq!(s3["skipInitialLoad"], true);
assert_eq!(s3["queueUrl"], "https://sqs.us-east-1.amazonaws.com/123/q");
assert!(
s3.get("startAfter").is_none(),
"startAfter leaked when --start-after not passed: {s3}",
);
}
#[tokio::test]
async fn s3_start_after_serializes_when_passed() {
let mock = start_mock_clickpipes_api().await;
let body = invoke_cli_capture_body(
&mock,
&[
"clickpipe",
"create",
"object-storage",
"svc-id",
"--name",
"t",
"--source-url",
"https://bucket.s3.us-east-1.amazonaws.com/data/*.json",
"--format",
"JSONEachRow",
"--database",
"d",
"--table",
"t",
"--column",
"id:Int64",
"--continuous",
"--queue-url",
"https://sqs.us-east-1.amazonaws.com/123/q",
"--start-after",
"obj-key-001",
"--org-id",
"org",
],
)
.await;
let s3 = &body["source"]["objectStorage"];
assert_eq!(s3["startAfter"], "obj-key-001");
assert!(
s3.get("skipInitialLoad").is_none(),
"skipInitialLoad leaked when --skip-initial-load not passed: {s3}",
);
}
#[tokio::test]
async fn mysql_server_id_serializes_when_passed() {
let mock = start_mock_clickpipes_api().await;
let body = invoke_cli_capture_body(
&mock,
&[
"clickpipe",
"create",
"mysql",
"svc-id",
"--name",
"t",
"--host",
"mysql",
"--port",
"3306",
"--username",
"u",
"--password",
"p",
"--table-mapping",
"mydb.t:t",
"--replication-mode",
"cdc",
"--server-id",
"4242",
"--org-id",
"org",
],
)
.await;
let mysql = &body["source"]["mysql"];
assert_eq!(mysql["serverId"], 4242);
}
#[tokio::test]
async fn mysql_server_id_absent_when_not_passed() {
let mock = start_mock_clickpipes_api().await;
let body = invoke_cli_capture_body(
&mock,
&[
"clickpipe",
"create",
"mysql",
"svc-id",
"--name",
"t",
"--host",
"mysql",
"--port",
"3306",
"--username",
"u",
"--password",
"p",
"--table-mapping",
"mydb.t:t",
"--replication-mode",
"cdc",
"--org-id",
"org",
],
)
.await;
let mysql = &body["source"]["mysql"];
assert!(
mysql.get("serverId").is_none(),
"serverId leaked when --server-id not passed: {mysql}",
);
}
async fn start_mock_schema_discovery_api() -> MockServer {
let mock = MockServer::start().await;
let stub_response = serde_json::json!({
"result": {
"fields": [
{ "name": "id", "type": "Int64", "optional": false },
{ "name": "event", "type": "String", "optional": true },
],
},
"status": 200,
"requestId": "stub-schema-discovery",
});
Mock::given(method("POST"))
.and(path_regex(
r"^/v1/organizations/[^/]+/services/[^/]+/clickpipes/schemaDiscovery$",
))
.respond_with(ResponseTemplate::new(200).set_body_json(stub_response))
.mount(&mock)
.await;
mock
}
#[tokio::test]
async fn schema_discover_kafka_posts_source_body() {
let mock = start_mock_schema_discovery_api().await;
let body = invoke_cli_capture_body(
&mock,
&[
"clickpipe",
"schema-discover",
"svc-id",
"--org-id",
"org",
"kafka",
"--brokers",
"broker:9092",
"--topics",
"topic",
"--format",
"JSONEachRow",
"--auth",
"IAM_ROLE",
"--iam-role",
"arn:aws:iam::123:role/x",
],
)
.await;
let kafka = &body["source"]["kafka"];
assert_eq!(kafka["brokers"], "broker:9092");
assert_eq!(kafka["topics"], "topic");
assert_eq!(kafka["format"], "JSONEachRow");
assert!(
body["source"].get("kinesis").is_none(),
"kinesis leaked into kafka schema-discovery body: {}",
body["source"],
);
}
#[tokio::test]
async fn schema_discover_kinesis_posts_source_body() {
let mock = start_mock_schema_discovery_api().await;
let body = invoke_cli_capture_body(
&mock,
&[
"clickpipe",
"schema-discover",
"svc-id",
"--org-id",
"org",
"kinesis",
"--stream-name",
"mystream",
"--region",
"us-east-1",
"--format",
"JSONEachRow",
],
)
.await;
let kinesis = &body["source"]["kinesis"];
assert_eq!(kinesis["streamName"], "mystream");
assert_eq!(kinesis["region"], "us-east-1");
assert_eq!(kinesis["format"], "JSONEachRow");
assert!(
body["source"].get("kafka").is_none(),
"kafka leaked into kinesis schema-discovery body: {}",
body["source"],
);
}
async fn run_service_reset_password(
password_response: Value,
extra_args: &[&str],
) -> std::process::Output {
let mock = MockServer::start().await;
Mock::given(method("PATCH"))
.and(path_regex(
r"^/v1/organizations/[^/]+/services/[^/]+/password$",
))
.respond_with(ResponseTemplate::new(200).set_body_json(serde_json::json!({
"result": password_response,
"status": 200,
"requestId": "stub-reset-password",
})))
.mount(&mock)
.await;
let url = mock.uri();
let mut args: Vec<&str> = vec![
"cloud",
"--url",
&url,
"service",
"reset-password",
"svc-id",
"--org-id",
"org-1",
];
args.extend(extra_args);
Command::new(clickhousectl_binary())
.env("DO_NOT_TRACK", "1")
.args(&args)
.env("CLICKHOUSE_CLOUD_API_KEY", "fake-key-for-tests")
.env("CLICKHOUSE_CLOUD_API_SECRET", "fake-secret-for-tests")
.output()
.expect("failed to spawn clickhousectl")
}
#[tokio::test]
async fn service_reset_password_fails_when_the_generated_password_is_absent() {
for extra_args in [&[][..], &["--json"][..]] {
let output = run_service_reset_password(serde_json::json!({}), extra_args).await;
assert!(
!output.status.success(),
"a generation reset with no password must fail for args {extra_args:?}\nstdout:\n{}",
String::from_utf8_lossy(&output.stdout),
);
let stderr = String::from_utf8_lossy(&output.stderr);
assert!(
stderr.contains("omitted the generated password"),
"stderr should name the omitted password for args {extra_args:?}:\n{stderr}",
);
let stdout = String::from_utf8_lossy(&output.stdout);
assert!(
!stdout.contains("Password reset for service")
&& !stdout.contains("no plaintext password returned"),
"no success output may precede the failure for args {extra_args:?}:\n{stdout}",
);
}
}
#[tokio::test]
async fn service_reset_password_succeeds_for_a_hash_reset_without_a_password() {
let output =
run_service_reset_password(serde_json::json!({}), &["--new-password-hash", "e3b0c442"])
.await;
assert_success(&output);
let stdout = String::from_utf8_lossy(&output.stdout);
assert!(
!stdout.contains("New password"),
"a hash reset must not report a password:\n{stdout}",
);
}
#[tokio::test]
async fn service_reset_password_treats_a_double_sha1_only_reset_as_generation() {
let output = run_service_reset_password(
serde_json::json!({}),
&["--new-double-sha1-hash", "aabbccdd"],
)
.await;
assert!(
!output.status.success(),
"a double-SHA1-only reset with no password must fail\nstdout:\n{}",
String::from_utf8_lossy(&output.stdout),
);
let stderr = String::from_utf8_lossy(&output.stderr);
assert!(
stderr.contains("omitted the generated password"),
"stderr should name the omitted password:\n{stderr}",
);
}
#[tokio::test]
async fn service_reset_password_json_prints_the_generated_password() {
let output =
run_service_reset_password(serde_json::json!({ "password": "s3cret" }), &["--json"]).await;
assert_success(&output);
let body: Value =
serde_json::from_slice(&output.stdout).expect("--json output wasn't valid JSON");
assert_eq!(body["password"], "s3cret");
}
async fn run_key_create(result: Value, extra_args: &[&str]) -> std::process::Output {
let mock = MockServer::start().await;
Mock::given(method("POST"))
.and(path("/v1/organizations/org-1/keys"))
.respond_with(ResponseTemplate::new(200).set_body_json(serde_json::json!({
"result": result,
"status": 200,
"requestId": "stub-key-create",
})))
.mount(&mock)
.await;
let url = mock.uri();
let mut args: Vec<&str> = vec![
"cloud", "--url", &url, "key", "create", "--name", "ci", "--org-id", "org-1",
];
args.extend(extra_args);
Command::new(clickhousectl_binary())
.env("DO_NOT_TRACK", "1")
.args(&args)
.env("CLICKHOUSE_CLOUD_API_KEY", "fake-key-for-tests")
.env("CLICKHOUSE_CLOUD_API_SECRET", "fake-secret-for-tests")
.output()
.expect("failed to spawn clickhousectl")
}
#[tokio::test]
async fn key_create_fails_when_the_generated_material_is_absent() {
for extra_args in [&[][..], &["--json"][..]] {
let output = run_key_create(
serde_json::json!({
"key": { "id": "aaaaaaaa-bbbb-cccc-dddd-eeeeeeeeeeee", "name": "ci" },
"keyId": "generated-key-id",
}),
extra_args,
)
.await;
assert!(
!output.status.success(),
"a generated key create with no secret must fail for args {extra_args:?}\nstdout:\n{}",
String::from_utf8_lossy(&output.stdout),
);
let stderr = String::from_utf8_lossy(&output.stderr);
assert!(
stderr.contains("omitted the generated key material"),
"stderr should name the omitted material for args {extra_args:?}:\n{stderr}",
);
let stdout = String::from_utf8_lossy(&output.stdout);
assert!(
!stdout.contains("API key created!") && !stdout.contains("generated-key-id"),
"no success output may precede the failure for args {extra_args:?}:\n{stdout}",
);
}
}
#[tokio::test]
async fn key_create_json_prints_the_raw_response() {
let result = serde_json::json!({
"key": { "id": "aaaaaaaa-bbbb-cccc-dddd-eeeeeeeeeeee", "name": "ci" },
"keyId": "generated-key-id",
"keySecret": "generated-key-secret",
});
let output = run_key_create(result.clone(), &["--json"]).await;
assert_success(&output);
let body: Value =
serde_json::from_slice(&output.stdout).expect("--json output wasn't valid JSON");
assert_eq!(body, result);
}
#[tokio::test]
async fn key_create_succeeds_for_a_pre_hashed_key_without_generated_material() {
let output = run_key_create(
serde_json::json!({ "key": { "id": "aaaaaaaa-bbbb-cccc-dddd-eeeeeeeeeeee", "name": "ci" } }),
&[
"--hash-key-id",
"0f1e2d3c",
"--hash-key-id-suffix",
"3c",
"--hash-key-secret",
"4b5a6978",
],
)
.await;
assert_success(&output);
let stdout = String::from_utf8_lossy(&output.stdout);
assert!(
!stdout.contains("Key Secret"),
"a pre-hashed create must not report key material:\n{stdout}",
);
}