diff --git a/crates/clickhouse-cloud-api/README.md b/crates/clickhouse-cloud-api/README.md index 22c14f28..c391c89c 100644 --- a/crates/clickhouse-cloud-api/README.md +++ b/crates/clickhouse-cloud-api/README.md @@ -52,6 +52,8 @@ Postgres slow-query aggregate durations (`*DurationUs`) and execution `durationU Kinesis source format enums now include `Protobuf`. Set `ClickPipePostKinesisSource.protobuf_schema` to the base64-encoded `.proto` source or serialized `FileDescriptorSet` for that format; omit it for other formats. Organization Prometheus discovery has graduated from beta and is no longer listed in `BETA_OPERATIONS`. +`Organization.capabilities.snapshots` reports snapshot eligibility as an optional boolean. Kinesis create requests accept `ClickPipeKinesisSchemaRegistry` for AWS Glue; supply `type`, `glue_region`, and `glue_registry_name`, and optionally `glue_role_arn` to assume a different role from the Kinesis source. Kinesis responses use `ClickPipeKinesisSchemaRegistryResponse`, whose fields tolerate absence and null; convert it with `TryFrom` before writing it back. Kafka create requests accept `tombstone_mode: Some(Delete)` to delete matching destination rows for tombstone records. This requires exactly-once delivery and is set only at creation. Struct-literal callers of `Organization`, `ClickPipeKinesisSource`, `ClickPipeKafkaSource`, `ClickPipePostKinesisSource`, and `ClickPipePostKafkaSource` need to supply the new optional fields as `None` or use `..Default::default()`. + The beta Query API endpoint management methods are `query_api_endpoint_create`, `query_api_endpoint_get`, `query_api_endpoint_list`, `query_api_endpoint_update`, and `query_api_endpoint_delete`. Create and update take `PublicQueryApiEndpointRequest`; list accepts an optional cursor and limit (1–100) and returns `items` with `pagination.next_cursor`. User-owned endpoints can be listed and read, but cannot be updated or deleted through this API. ### ClickHouse settings models diff --git a/crates/clickhouse-cloud-api/clickhouse_cloud_openapi.json b/crates/clickhouse-cloud-api/clickhouse_cloud_openapi.json index 6ba7c530..c5d9bbfb 100644 --- a/crates/clickhouse-cloud-api/clickhouse_cloud_openapi.json +++ b/crates/clickhouse-cloud-api/clickhouse_cloud_openapi.json @@ -5,7 +5,7 @@ "version": "1.0", "contact": { "name": "ClickHouse Support", - "url": "https://clickhouse.com/docs/en/cloud/manage/openapi?referrer=openapi-1152576", + "url": "https://clickhouse.com/docs/en/cloud/manage/openapi?referrer=openapi-1156400", "email": "support@clickhouse.com" } }, @@ -8887,7 +8887,7 @@ }, "patch": { "summary": "Update ClickPipe", - "description": "Update the specified ClickPipe. Source fields not present in the per-source update schemas are immutable after creation. For Kafka sources, values submitted for immutable fields (type, format, brokers, topics, consumerGroup, offset, schemaRegistry, exactlyOnce) are not applied, except schema registry credentials, which are rejected.", + "description": "Update the specified ClickPipe. Source fields not present in the per-source update schemas are immutable after creation. For Kafka sources, values submitted for immutable fields (type, format, brokers, topics, consumerGroup, offset, schemaRegistry, exactlyOnce, tombstoneMode) are not applied, except schema registry credentials, which are rejected.", "operationId": "clickPipeUpdate", "parameters": [ { @@ -23024,6 +23024,14 @@ } } }, + "OrganizationCapabilities": { + "properties": { + "snapshots": { + "description": "Whether the organization is eligible to use service snapshots: true only when the organization has the snapshots feature enabled, is on a PPv2 tier, and has the backups entitlement — the same conditions enforced when a snapshotConfiguration is saved. Check this before configuring snapshots on a service. Snapshots apply to primary services only, so a secondary/replica service is rejected regardless of organization eligibility.", + "type": "boolean" + } + } + }, "Organization": { "properties": { "id": { @@ -23057,6 +23065,9 @@ "enableCoreDumps": { "description": "Whether crash reports (core dumps) collection is enabled for services in the organization. When disabled at the organization level, individual services cannot enable crash reports.", "type": "boolean" + }, + "capabilities": { + "$ref": "#/components/schemas/OrganizationCapabilities" } } }, @@ -23914,8 +23925,9 @@ "type": "string" }, "topics": { - "description": "Topics of the Kafka source.", - "type": "string" + "description": "One or more Kafka topics as a comma-separated string. All topics must have the same schema and are ingested into the same destination table by a single ClickPipe.", + "type": "string", + "example": "topic1,topic2" }, "consumerGroup": { "description": "Consumer group of the Kafka source. If not provided \"clickpipes-<>\" will be used.", @@ -23986,6 +23998,17 @@ "boolean", "null" ] + }, + "tombstoneMode": { + "description": "How Kafka tombstone records are handled. Set to \"delete\" to delete the matching destination row. Requires exactly-once delivery and can only be set at pipe creation.", + "type": [ + "string", + "null" + ], + "enum": [ + "delete" + ], + "example": "delete" } } }, @@ -24020,8 +24043,9 @@ "type": "string" }, "topics": { - "description": "Topics of the Kafka source.", - "type": "string" + "description": "One or more Kafka topics as a comma-separated string. All topics must have the same schema and are ingested into the same destination table by a single ClickPipe.", + "type": "string", + "example": "topic1,topic2" }, "consumerGroup": { "description": "Consumer group of the Kafka source. If not provided \"clickpipes-<>\" will be used.", @@ -24093,6 +24117,17 @@ "null" ] }, + "tombstoneMode": { + "description": "How Kafka tombstone records are handled. Set to \"delete\" to delete the matching destination row. Requires exactly-once delivery and can only be set at pipe creation.", + "type": [ + "string", + "null" + ], + "enum": [ + "delete" + ], + "example": "delete" + }, "credentials": { "description": "Credentials for Kafka source. Choose one that is supported by the authentication method.", "oneOf": [ @@ -24179,6 +24214,40 @@ } } }, + "ClickPipeKinesisSchemaRegistry": { + "properties": { + "type": { + "description": "Type of the schema registry. Kinesis ClickPipes support the AWS Glue Schema Registry, which authenticates with IAM instead of credentials.", + "type": "string", + "enum": [ + "glue" + ] + }, + "glueRegion": { + "description": "AWS region of the Glue Schema Registry.", + "type": "string", + "example": "us-east-1" + }, + "glueRegistryName": { + "description": "Name of the Glue Schema Registry.", + "type": "string", + "example": "my-registry" + }, + "glueRoleArn": { + "description": "IAM role to assume for Glue Schema Registry access. Defaults to the IAM identity of the Kinesis source.", + "type": [ + "string", + "null" + ], + "example": "arn:aws:iam::123456789012:role/MyGlueRegistryRole" + } + }, + "required": [ + "type", + "glueRegion", + "glueRegistryName" + ] + }, "ClickPipeKinesisSource": { "properties": { "format": { @@ -24240,6 +24309,16 @@ "null" ], "example": "arn:aws:iam::123456789012:role/MyRole" + }, + "schemaRegistry": { + "oneOf": [ + { + "$ref": "#/components/schemas/ClickPipeKinesisSchemaRegistry" + }, + { + "type": "null" + } + ] } } }, @@ -24305,6 +24384,16 @@ ], "example": "arn:aws:iam::123456789012:role/MyRole" }, + "schemaRegistry": { + "oneOf": [ + { + "$ref": "#/components/schemas/ClickPipeKinesisSchemaRegistry" + }, + { + "type": "null" + } + ] + }, "accessKey": { "oneOf": [ { @@ -24316,7 +24405,7 @@ ] }, "protobufSchema": { - "description": "Base64-encoded .proto source or serialized FileDescriptorSet. Required with Protobuf format and not supported with other formats.", + "description": "Base64-encoded .proto source or serialized FileDescriptorSet. Required with Protobuf format unless a schema registry is configured, and not supported with other formats.", "type": "string", "example": "c3ludGF4ID0gInByb3RvMyI7IG1lc3NhZ2UgRXZlbnQge30=", "maxLength": 1048576, @@ -33796,7 +33885,7 @@ }, "groups": { "type": "array", - "description": "A list of groups to which the user belongs. Role may be derived from group display or value.", + "description": "A list of groups to which the user belongs. Read-only; ignored on write. Membership is managed via the Groups endpoints.", "items": { "$ref": "#/components/schemas/ScimUserGroup" } @@ -33810,7 +33899,7 @@ }, "roles": { "type": "array", - "description": "A list of roles for the user.", + "description": "A list of roles for the user. Read-only; ignored on write. Membership is managed via the Groups endpoints.", "items": { "$ref": "#/components/schemas/ScimUserRole" } @@ -34018,7 +34107,7 @@ }, "groups": { "type": "array", - "description": "A list of groups to which the user belongs. Role may be derived from group display or value.", + "description": "A list of groups to which the user belongs. Read-only; ignored on write. Membership is managed via the Groups endpoints.", "items": { "$ref": "#/components/schemas/ScimUserGroup" } @@ -34032,7 +34121,7 @@ }, "roles": { "type": "array", - "description": "A list of roles for the user.", + "description": "A list of roles for the user. Read-only; ignored on write. Membership is managed via the Groups endpoints.", "items": { "$ref": "#/components/schemas/ScimUserRole" } @@ -38902,4 +38991,4 @@ ] } ] -} +} \ No newline at end of file diff --git a/crates/clickhouse-cloud-api/src/convert.rs b/crates/clickhouse-cloud-api/src/convert.rs index 9da84411..2f2a0059 100644 --- a/crates/clickhouse-cloud-api/src/convert.rs +++ b/crates/clickhouse-cloud-api/src/convert.rs @@ -19,6 +19,7 @@ use std::fmt; +mod clickpipes; mod clickstack; mod postgres; mod service; diff --git a/crates/clickhouse-cloud-api/src/convert/clickpipes.rs b/crates/clickhouse-cloud-api/src/convert/clickpipes.rs new file mode 100644 index 00000000..c47a11c9 --- /dev/null +++ b/crates/clickhouse-cloud-api/src/convert/clickpipes.rs @@ -0,0 +1,29 @@ +use super::MissingRequiredFields; +use crate::models::{ClickPipeKinesisSchemaRegistry, ClickPipeKinesisSchemaRegistryResponse}; + +impl TryFrom for ClickPipeKinesisSchemaRegistry { + type Error = MissingRequiredFields; + + fn try_from(value: ClickPipeKinesisSchemaRegistryResponse) -> Result { + let mut missing = Vec::new(); + if value.r#type.is_none() { + missing.push("type"); + } + if value.glue_region.is_none() { + missing.push("glueRegion"); + } + if value.glue_registry_name.is_none() { + missing.push("glueRegistryName"); + } + if !missing.is_empty() { + return Err(MissingRequiredFields::new(missing)); + } + + Ok(Self { + r#type: value.r#type.expect("checked above"), + glue_region: value.glue_region.expect("checked above"), + glue_registry_name: value.glue_registry_name.expect("checked above"), + glue_role_arn: value.glue_role_arn, + }) + } +} diff --git a/crates/clickhouse-cloud-api/src/models.rs b/crates/clickhouse-cloud-api/src/models.rs index b5b1fce7..a80d5750 100644 --- a/crates/clickhouse-cloud-api/src/models.rs +++ b/crates/clickhouse-cloud-api/src/models.rs @@ -185,7 +185,8 @@ pub use organization_private_endpoints::{ }; pub use organizations::{ ActiveBalance, ActiveBalances, CreditBalance, CreditBalanceType, CreditBalances, Organization, - OrganizationPatchRequest, PrometheusDiscoveryLabels, PrometheusDiscoveryTargetGroup, + OrganizationCapabilities, OrganizationPatchRequest, PrometheusDiscoveryLabels, + PrometheusDiscoveryTargetGroup, }; pub use postgres::{ BasePostgresService, PgBouncerConfig, PgBouncerConfigResponse, PgConfig, diff --git a/crates/clickhouse-cloud-api/src/models/clickpipes.rs b/crates/clickhouse-cloud-api/src/models/clickpipes.rs index fb035130..3b33e3c4 100644 --- a/crates/clickhouse-cloud-api/src/models/clickpipes.rs +++ b/crates/clickhouse-cloud-api/src/models/clickpipes.rs @@ -2275,6 +2275,8 @@ pub struct ClickPipeKafkaSource { pub schema_registry: Option, #[serde(skip_serializing_if = "Option::is_none")] pub topics: Option, + #[serde(rename = "tombstoneMode", skip_serializing_if = "Option::is_none")] + pub tombstone_mode: Option, #[serde(skip_serializing_if = "Option::is_none")] pub r#type: Option, } @@ -2292,6 +2294,8 @@ pub struct ClickPipeKinesisSource { pub iterator_type: Option, #[serde(skip_serializing_if = "Option::is_none")] pub region: Option, + #[serde(rename = "schemaRegistry", skip_serializing_if = "Option::is_none")] + pub schema_registry: Option, #[serde(rename = "streamName", skip_serializing_if = "Option::is_none")] pub stream_name: Option, #[serde(skip_serializing_if = "Option::is_none")] @@ -2300,6 +2304,52 @@ pub struct ClickPipeKinesisSource { pub use_enhanced_fan_out: Option, } +/// Values of `ClickPipeKinesisSchemaRegistry.type` in the Cloud API. +#[derive(Debug, Clone, PartialEq, Default, Serialize, Deserialize)] +pub enum ClickPipeKinesisSchemaRegistryType { + #[serde(rename = "glue")] + #[default] + Glue, + #[serde(untagged)] + Unknown(String), +} + +impl std::fmt::Display for ClickPipeKinesisSchemaRegistryType { + fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + match self { + Self::Glue => write!(f, "glue"), + Self::Unknown(value) => write!(f, "{value}"), + } + } +} + +/// `ClickPipeKinesisSchemaRegistry` for Kinesis create requests. +#[derive(Debug, Clone, PartialEq, Default, Serialize, Deserialize)] +pub struct ClickPipeKinesisSchemaRegistry { + #[serde(rename = "type")] + pub r#type: ClickPipeKinesisSchemaRegistryType, + #[serde(rename = "glueRegion")] + pub glue_region: String, + #[serde(rename = "glueRegistryName")] + pub glue_registry_name: String, + /// Omit to use the Kinesis source's IAM identity for Glue access. + #[serde(rename = "glueRoleArn", skip_serializing_if = "Option::is_none")] + pub glue_role_arn: Option, +} + +/// `ClickPipeKinesisSchemaRegistry` in a Kinesis source response. +#[derive(Debug, Clone, PartialEq, Default, Serialize, Deserialize)] +pub struct ClickPipeKinesisSchemaRegistryResponse { + #[serde(rename = "type", skip_serializing_if = "Option::is_none")] + pub r#type: Option, + #[serde(rename = "glueRegion", skip_serializing_if = "Option::is_none")] + pub glue_region: Option, + #[serde(rename = "glueRegistryName", skip_serializing_if = "Option::is_none")] + pub glue_registry_name: Option, + #[serde(rename = "glueRoleArn", skip_serializing_if = "Option::is_none")] + pub glue_role_arn: Option, +} + /// `ClickPipeMongoDBPipeSettings` from the ClickHouse Cloud API. #[derive(Debug, Clone, PartialEq, Default, Serialize, Deserialize)] pub struct ClickPipeMongoDBPipeSettings { @@ -3124,9 +3174,50 @@ pub struct ClickPipePostKafkaSource { #[serde(rename = "schemaRegistry", skip_serializing_if = "Option::is_none")] pub schema_registry: Option, pub topics: String, + /// Deleting Kafka tombstones requires exactly-once delivery and can only be set on creation. + #[serde(rename = "tombstoneMode", skip_serializing_if = "Option::is_none")] + pub tombstone_mode: Option, pub r#type: ClickPipePostKafkaSourceType, } +/// Inline enum for `ClickPipePostKafkaSource.tombstoneMode`. +#[derive(Debug, Clone, PartialEq, Default, Serialize, Deserialize)] +pub enum ClickPipePostKafkaSourceTombstonemode { + #[serde(rename = "delete")] + #[default] + Delete, + #[serde(untagged)] + Unknown(String), +} + +impl std::fmt::Display for ClickPipePostKafkaSourceTombstonemode { + fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + match self { + Self::Delete => write!(f, "delete"), + Self::Unknown(value) => write!(f, "{value}"), + } + } +} + +/// Inline enum for `ClickPipeKafkaSource.tombstoneMode`. +#[derive(Debug, Clone, PartialEq, Default, Serialize, Deserialize)] +pub enum ClickPipeKafkaSourceTombstonemode { + #[serde(rename = "delete")] + #[default] + Delete, + #[serde(untagged)] + Unknown(String), +} + +impl std::fmt::Display for ClickPipeKafkaSourceTombstonemode { + fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + match self { + Self::Delete => write!(f, "delete"), + Self::Unknown(value) => write!(f, "{value}"), + } + } +} + /// `ClickPipePostKinesisSource` from the ClickHouse Cloud API. #[derive(Debug, Clone, PartialEq, Default, Serialize, Deserialize)] pub struct ClickPipePostKinesisSource { @@ -3139,6 +3230,8 @@ pub struct ClickPipePostKinesisSource { #[serde(rename = "iteratorType")] pub iterator_type: ClickPipePostKinesisSourceIteratortype, pub region: String, + #[serde(rename = "schemaRegistry", skip_serializing_if = "Option::is_none")] + pub schema_registry: Option, #[serde(rename = "streamName")] pub stream_name: String, #[serde(skip_serializing_if = "Option::is_none")] diff --git a/crates/clickhouse-cloud-api/src/models/organizations.rs b/crates/clickhouse-cloud-api/src/models/organizations.rs index 020ac814..0f692852 100644 --- a/crates/clickhouse-cloud-api/src/models/organizations.rs +++ b/crates/clickhouse-cloud-api/src/models/organizations.rs @@ -73,6 +73,8 @@ pub struct PrometheusDiscoveryTargetGroup { pub struct Organization { #[serde(rename = "byocConfig", skip_serializing_if = "Option::is_none")] pub byoc_config: Option>, + #[serde(skip_serializing_if = "Option::is_none")] + pub capabilities: Option, #[serde(rename = "createdAt", skip_serializing_if = "Option::is_none")] pub created_at: Option>, #[serde(rename = "enableCoreDumps", skip_serializing_if = "Option::is_none")] @@ -85,6 +87,13 @@ pub struct Organization { pub private_endpoints: Option>, } +/// `OrganizationCapabilities` from the ClickHouse Cloud API. +#[derive(Debug, Clone, PartialEq, Default, Serialize, Deserialize)] +pub struct OrganizationCapabilities { + #[serde(skip_serializing_if = "Option::is_none")] + pub snapshots: Option, +} + /// `OrganizationPatchRequest` from the ClickHouse Cloud API. #[derive(Debug, Clone, PartialEq, Default, Serialize, Deserialize)] pub struct OrganizationPatchRequest { diff --git a/crates/clickhouse-cloud-api/tests/client_test.rs b/crates/clickhouse-cloud-api/tests/client_test.rs index 868ee086..782c78a8 100644 --- a/crates/clickhouse-cloud-api/tests/client_test.rs +++ b/crates/clickhouse-cloud-api/tests/client_test.rs @@ -94,6 +94,25 @@ async fn get_organization() { assert_eq!(org.name.as_deref(), Some("My Org")); } +#[tokio::test] +async fn get_organization_preserves_snapshot_capability() { + let (server, client) = setup().await; + Mock::given(method("GET")) + .and(path("/v1/organizations/org-1")) + .respond_with(ok_json(serde_json::json!({ + "capabilities": {"snapshots": false} + }))) + .mount(&server) + .await; + let org = client + .organization_get("org-1") + .await + .unwrap() + .result + .unwrap(); + assert_eq!(org.capabilities.unwrap().snapshots, Some(false)); +} + #[tokio::test] async fn get_active_balances_with_pagination() { let (s, c) = setup().await; @@ -1910,6 +1929,83 @@ async fn create_click_pipe() { assert_eq!(pipe.name.as_deref(), Some("new-pipe")); } +#[tokio::test] +async fn create_kinesis_click_pipe_sends_glue_registry() { + let (server, client) = setup().await; + Mock::given(method("POST")) + .and(path("/v1/organizations/org-1/services/svc-1/clickpipes")) + .and(body_partial_json(serde_json::json!({ + "source": {"kinesis": {"schemaRegistry": { + "type": "glue", "glueRegion": "us-east-1", "glueRegistryName": "events" + }}} + }))) + .respond_with(ok_json(serde_json::json!({"name": "kinesis-pipe"}))) + .expect(1) + .mount(&server) + .await; + let request = ClickPipePostRequest { + name: "kinesis-pipe".into(), + source: ClickPipePostSource { + kinesis: Some(ClickPipePostKinesisSource { + schema_registry: Some(ClickPipeKinesisSchemaRegistry { + r#type: ClickPipeKinesisSchemaRegistryType::Glue, + glue_region: "us-east-1".into(), + glue_registry_name: "events".into(), + glue_role_arn: None, + }), + ..Default::default() + }), + ..Default::default() + }, + ..Default::default() + }; + let response = client + .click_pipe_create("org-1", "svc-1", &request) + .await + .unwrap(); + assert_eq!( + response.result.unwrap().name.as_deref(), + Some("kinesis-pipe") + ); +} + +#[tokio::test] +async fn get_click_pipe_preserves_nullable_kinesis_registry() { + let (server, client) = setup().await; + Mock::given(method("GET")) + .and(path( + "/v1/organizations/org-1/services/svc-1/clickpipes/pipe-1", + )) + .respond_with(ok_json(serde_json::json!({ + "source": {"kinesis": {"schemaRegistry": { + "type": "glue", "glueRegion": "us-east-1", "glueRegistryName": null + }}} + }))) + .mount(&server) + .await; + let pipe = client + .click_pipe_get("org-1", "svc-1", "pipe-1") + .await + .unwrap() + .result + .unwrap(); + let registry = pipe + .source + .unwrap() + .kinesis + .unwrap() + .schema_registry + .unwrap(); + assert_eq!(registry.glue_region.as_deref(), Some("us-east-1")); + assert_eq!(registry.glue_registry_name, None); + assert_eq!( + ClickPipeKinesisSchemaRegistry::try_from(registry) + .unwrap_err() + .fields(), + &["glueRegistryName"] + ); +} + #[tokio::test] async fn create_click_pipe_sends_start_paused_and_table_ttl() { let (server, client) = setup().await; diff --git a/crates/clickhouse-cloud-api/tests/model_facade_test.rs b/crates/clickhouse-cloud-api/tests/model_facade_test.rs index c8040c4b..9c85af99 100644 --- a/crates/clickhouse-cloud-api/tests/model_facade_test.rs +++ b/crates/clickhouse-cloud-api/tests/model_facade_test.rs @@ -69,6 +69,18 @@ fn extracted_models_keep_root_and_models_paths() { api::ActiveBalances::default(), api::models::ActiveBalances::default(), ); + assert_same_type( + api::OrganizationCapabilities::default(), + api::models::OrganizationCapabilities::default(), + ); + assert_same_type( + api::ClickPipeKinesisSchemaRegistry::default(), + api::models::ClickPipeKinesisSchemaRegistry::default(), + ); + assert_same_type( + api::ClickPipeKinesisSchemaRegistryResponse::default(), + api::models::ClickPipeKinesisSchemaRegistryResponse::default(), + ); assert_same_type(api::Activity::default(), api::models::Activity::default()); assert_same_type(api::ApiKey::default(), api::models::ApiKey::default()); assert_same_type( diff --git a/crates/clickhouse-cloud-api/tests/models_test.rs b/crates/clickhouse-cloud-api/tests/models_test.rs index 4580dda0..62cb4dcd 100644 --- a/crates/clickhouse-cloud-api/tests/models_test.rs +++ b/crates/clickhouse-cloud-api/tests/models_test.rs @@ -7096,6 +7096,157 @@ fn kinesis_protobuf_schema_round_trips_and_other_formats_omit_it() { ); } +#[test] +fn organization_capabilities_tolerate_absence_null_and_false() { + for wire in [ + serde_json::json!({}), + serde_json::json!({"capabilities": null}), + ] { + let org: Organization = serde_json::from_value(wire).unwrap(); + assert_eq!(org.capabilities, None); + assert_eq!(serde_json::to_value(org).unwrap(), serde_json::json!({})); + } + for wire in [ + serde_json::json!({}), + serde_json::json!({"snapshots": null}), + ] { + let capabilities: OrganizationCapabilities = serde_json::from_value(wire).unwrap(); + assert_eq!(capabilities.snapshots, None); + assert_eq!( + serde_json::to_value(capabilities).unwrap(), + serde_json::json!({}) + ); + } + let org: Organization = serde_json::from_value(serde_json::json!({ + "capabilities": {"snapshots": false} + })) + .unwrap(); + assert_eq!(org.capabilities.unwrap().snapshots, Some(false)); +} + +#[test] +fn kinesis_registry_request_is_strict_and_response_is_tolerant() { + let full = serde_json::json!({ + "type": "glue", "glueRegion": "us-east-1", "glueRegistryName": "events" + }); + let request: ClickPipeKinesisSchemaRegistry = serde_json::from_value(full.clone()).unwrap(); + assert_eq!(serde_json::to_value(request).unwrap(), full); + for field in ["type", "glueRegion", "glueRegistryName"] { + let mut missing = full.clone(); + missing.as_object_mut().unwrap().remove(field); + assert!(serde_json::from_value::(missing).is_err()); + let mut null = full.clone(); + null[field] = serde_json::Value::Null; + assert!(serde_json::from_value::(null).is_err()); + } + for wire in [ + serde_json::json!({}), + serde_json::json!({ + "type": null, "glueRegion": null, "glueRegistryName": null, "glueRoleArn": null + }), + ] { + let response: ClickPipeKinesisSchemaRegistryResponse = + serde_json::from_value(wire).unwrap(); + assert_eq!(response, ClickPipeKinesisSchemaRegistryResponse::default()); + assert_eq!( + serde_json::to_value(response).unwrap(), + serde_json::json!({}) + ); + } +} + +#[test] +fn kinesis_registry_conversion_reports_nested_missing_fields_and_preserves_unknown_type() { + let response: ClickPipeKinesisSchemaRegistryResponse = + serde_json::from_value(serde_json::json!({ + "type": "future", "glueRegion": "us-east-1", "glueRoleArn": "arn:aws:iam::123:role/Glue" + })) + .unwrap(); + let error = ClickPipeKinesisSchemaRegistry::try_from(response.clone()).unwrap_err(); + assert_eq!(error.fields(), &["glueRegistryName"]); + let mut complete = response; + complete.glue_registry_name = Some("events".into()); + let request = ClickPipeKinesisSchemaRegistry::try_from(complete).unwrap(); + assert_eq!( + serde_json::to_value(request).unwrap(), + serde_json::json!({ + "type": "future", "glueRegion": "us-east-1", "glueRegistryName": "events", + "glueRoleArn": "arn:aws:iam::123:role/Glue" + }) + ); + assert_eq!( + ClickPipeKinesisSchemaRegistryType::Unknown("future".into()).to_string(), + "future" + ); +} + +#[test] +fn kafka_tombstone_mode_is_lossless_and_optional() { + let request = ClickPipePostKafkaSource::default(); + assert!( + serde_json::to_value(request) + .unwrap() + .get("tombstoneMode") + .is_none() + ); + let response: ClickPipeKafkaSource = serde_json::from_value(serde_json::json!({ + "tombstoneMode": "future" + })) + .unwrap(); + assert!( + matches!(response.tombstone_mode, Some(ClickPipeKafkaSourceTombstonemode::Unknown(ref value)) if value == "future") + ); + let request = ClickPipePostKafkaSource { + tombstone_mode: Some(ClickPipePostKafkaSourceTombstonemode::Delete), + ..Default::default() + }; + assert_eq!( + serde_json::to_value(request).unwrap()["tombstoneMode"], + "delete" + ); + assert_eq!( + ClickPipePostKafkaSourceTombstonemode::Delete.to_string(), + "delete" + ); + for wire in [ + serde_json::json!({}), + serde_json::json!({"tombstoneMode": null}), + ] { + let response: ClickPipeKafkaSource = serde_json::from_value(wire).unwrap(); + assert_eq!(response.tombstone_mode, None); + assert!( + serde_json::to_value(response) + .unwrap() + .get("tombstoneMode") + .is_none() + ); + } +} + +#[test] +fn kinesis_sources_omit_missing_or_null_registry() { + let request = ClickPipePostKinesisSource::default(); + assert!( + serde_json::to_value(request) + .unwrap() + .get("schemaRegistry") + .is_none() + ); + for wire in [ + serde_json::json!({}), + serde_json::json!({"schemaRegistry": null}), + ] { + let response: ClickPipeKinesisSource = serde_json::from_value(wire).unwrap(); + assert_eq!(response.schema_registry, None); + assert!( + serde_json::to_value(response) + .unwrap() + .get("schemaRegistry") + .is_none() + ); + } +} + #[test] fn query_api_endpoint_request_is_strict_and_omits_optional_fields() { let required = serde_json::json!({ diff --git a/crates/clickhousectl/src/cloud/clickpipes.rs b/crates/clickhousectl/src/cloud/clickpipes.rs index cf6e7a86..4922dcc5 100644 --- a/crates/clickhousectl/src/cloud/clickpipes.rs +++ b/crates/clickhousectl/src/cloud/clickpipes.rs @@ -2847,6 +2847,7 @@ fn build_kafka_source_with_exactly_once( }), schema_registry, protobuf_schema, + tombstone_mode: None, ca_certificate, reverse_private_endpoint_ids: args.reverse_private_endpoint_ids.clone(), }) @@ -2930,6 +2931,7 @@ fn build_kinesis_source( authentication: parse_enum(auth)?, iam_role: args.iam_role.clone(), access_key, + schema_registry: None, use_enhanced_fan_out: if args.enhanced_fan_out { Some(true) } else { diff --git a/scripts/classify-cloud-integration.py b/scripts/classify-cloud-integration.py index 5e5a3f57..6ebb6c44 100644 --- a/scripts/classify-cloud-integration.py +++ b/scripts/classify-cloud-integration.py @@ -42,6 +42,9 @@ "crates/clickhouse-cloud-api/src/client/udfs.rs": NO_SUITES, "crates/clickhouse-cloud-api/src/convert.rs": ALL_SUITES, "crates/clickhouse-cloud-api/src/convert/clickstack.rs": NO_SUITES, + "crates/clickhouse-cloud-api/src/convert/clickpipes.rs": frozenset( + {"clickpipes"} + ), "crates/clickhouse-cloud-api/src/convert/postgres.rs": frozenset( {"postgres", "clickpipes"} ), diff --git a/scripts/classify-install-integration.py b/scripts/classify-install-integration.py index 24f3f59f..631221d7 100644 --- a/scripts/classify-install-integration.py +++ b/scripts/classify-install-integration.py @@ -46,6 +46,7 @@ "crates/clickhouse-openapi-analyzer/tests/fixtures/operation_contracts/before.json", "crates/clickhouse-openapi-analyzer/tests/fixtures/operation_contracts/after.json", "crates/clickhouse-cloud-api/src/client/query_api_endpoints.rs", + "crates/clickhouse-cloud-api/src/convert/clickpipes.rs", "crates/clickhouse-cloud-api/src/models/query_api_endpoints.rs", "crates/clickhousectl/src/cloud/clickstack.rs", "crates/clickhousectl/src/cloud/config.rs", diff --git a/scripts/tests/test_classify_cloud_integration.py b/scripts/tests/test_classify_cloud_integration.py index 752a4ec8..705909b3 100644 --- a/scripts/tests/test_classify_cloud_integration.py +++ b/scripts/tests/test_classify_cloud_integration.py @@ -78,6 +78,7 @@ def test_current_source_mapping_values(self): }, frozenset({"clickpipes"}): { "crates/clickhouse-cloud-api/src/client/clickpipes.rs", + "crates/clickhouse-cloud-api/src/convert/clickpipes.rs", "crates/clickhouse-cloud-api/src/models/clickpipes.rs", }, frozenset({"service", "organization"}): { diff --git a/scripts/tests/test_classify_install_integration.py b/scripts/tests/test_classify_install_integration.py index 7c688924..3409a9f6 100644 --- a/scripts/tests/test_classify_install_integration.py +++ b/scripts/tests/test_classify_install_integration.py @@ -120,6 +120,7 @@ def test_exact_path_mappings(self): "crates/clickhouse-openapi-analyzer/tests/fixtures/operation_contracts/before.json", "crates/clickhouse-openapi-analyzer/tests/fixtures/operation_contracts/after.json", "crates/clickhouse-cloud-api/src/client/query_api_endpoints.rs", + "crates/clickhouse-cloud-api/src/convert/clickpipes.rs", "crates/clickhouse-cloud-api/src/models/query_api_endpoints.rs", "crates/clickhousectl/src/cloud/clickstack.rs", "crates/clickhousectl/src/cloud/config.rs",