diff --git a/rust_snuba/benches/processors.rs b/rust_snuba/benches/processors.rs index 8deebb5755..4f72d07702 100644 --- a/rust_snuba/benches/processors.rs +++ b/rust_snuba/benches/processors.rs @@ -57,6 +57,7 @@ fn create_factory( user: "test".into(), password: "test".into(), database: "test".into(), + verify: None, }, message_processor: MessageProcessorConfig { python_class_name: python_class_name.into(), diff --git a/rust_snuba/src/config.rs b/rust_snuba/src/config.rs index acd19c5559..0e820dff37 100644 --- a/rust_snuba/src/config.rs +++ b/rust_snuba/src/config.rs @@ -101,6 +101,11 @@ pub struct ClickhouseConfig { pub user: String, pub password: String, pub database: String, + /// Mirrors the Python cluster `verify` setting: `None` (unset) and + /// `Some(true)` keep reqwest's default certificate verification, only + /// `Some(false)` disables it. + #[serde(default)] + pub verify: Option, } #[derive(Deserialize, Clone, Debug)] @@ -135,4 +140,18 @@ mod tests { "10000" ); } + + #[test] + fn clickhouse_config_verify_defaults_to_none() { + let raw = r#"{"host": "h", "port": 9000, "secure": true, "user": "u", "password": "p", "database": "d"}"#; + let config: ClickhouseConfig = serde_json::from_str(raw).unwrap(); + assert_eq!(config.verify, None); + } + + #[test] + fn clickhouse_config_deserializes_verify() { + let raw = r#"{"host": "h", "port": 9000, "secure": true, "user": "u", "password": "p", "database": "d", "verify": false}"#; + let config: ClickhouseConfig = serde_json::from_str(raw).unwrap(); + assert_eq!(config.verify, Some(false)); + } } diff --git a/rust_snuba/src/strategies/clickhouse/writer_v2.rs b/rust_snuba/src/strategies/clickhouse/writer_v2.rs index 118c7196a3..312030a983 100644 --- a/rust_snuba/src/strategies/clickhouse/writer_v2.rs +++ b/rust_snuba/src/strategies/clickhouse/writer_v2.rs @@ -247,6 +247,13 @@ pub struct ClickhouseClient { query: String, } +/// Matches the Python clients, which skip certificate and hostname +/// verification when the cluster sets `verify=False` (self-signed / +/// private-CA HTTPS endpoints). +fn tls_verification_disabled(config: &ClickhouseConfig) -> bool { + config.secure && config.verify == Some(false) +} + impl ClickhouseClient { pub fn new( config: &ClickhouseConfig, @@ -291,12 +298,23 @@ impl ClickhouseClient { ); let timeouts = get_clickhouse_write_client_timeouts(&storage_name); - let client = Client::builder() + let mut builder = Client::builder() .connect_timeout(timeouts.connect) .pool_idle_timeout(timeouts.pool_idle) .tcp_keepalive(timeouts.tcp_keepalive) .tcp_keepalive_interval(timeouts.tcp_keepalive_interval) - .tcp_keepalive_retries(timeouts.tcp_keepalive_retries) + .tcp_keepalive_retries(timeouts.tcp_keepalive_retries); + + if tls_verification_disabled(config) { + tracing::warn!( + "ClickHouse TLS certificate and hostname verification disabled (verify=false)" + ); + builder = builder + .tls_danger_accept_invalid_certs(true) + .tls_danger_accept_invalid_hostnames(true); + } + + let client = builder .build() .expect("failed to build ClickHouse HTTP client"); @@ -503,9 +521,29 @@ mod tests { user: std::env::var("CLICKHOUSE_USER").unwrap_or("default".to_string()), password: std::env::var("CLICKHOUSE_PASSWORD").unwrap_or("".to_string()), database: std::env::var("CLICKHOUSE_DATABASE").unwrap_or("default".to_string()), + verify: None, } } + #[test] + fn test_tls_verification_disabled() { + let mut config = make_test_config(); + + config.secure = true; + config.verify = Some(false); + assert!(tls_verification_disabled(&config)); + + config.verify = Some(true); + assert!(!tls_verification_disabled(&config)); + + config.verify = None; + assert!(!tls_verification_disabled(&config)); + + config.secure = false; + config.verify = Some(false); + assert!(!tls_verification_disabled(&config)); + } + #[tokio::test] async fn test_compressed_insert_against_live_clickhouse() { crate::testutils::initialize_python(); @@ -729,6 +767,7 @@ mod tests { user: "default".to_string(), password: "".to_string(), database: "default".to_string(), + verify: None, }; let client = ClickhouseClient::new( @@ -781,6 +820,7 @@ mod tests { user: "default".to_string(), password: "".to_string(), database: "default".to_string(), + verify: None, }; let client = ClickhouseClient::new( &config, diff --git a/snuba/clusters/cluster.py b/snuba/clusters/cluster.py index 168b01918c..f668144a22 100644 --- a/snuba/clusters/cluster.py +++ b/snuba/clusters/cluster.py @@ -305,7 +305,7 @@ def __init__( database: str, secure: bool, ca_certs: str | None, - verify: bool | None, + verify: bool | str | None, storage_sets: set[str], single_node: bool, # The cluster name and distributed cluster name only apply if single_node is set to False @@ -371,7 +371,7 @@ def get_node_connection( self.__database, self.__secure, self.__ca_certs, - self.__verify, + self.get_verify(), ) def get_deleter(self) -> Reader: @@ -418,7 +418,7 @@ def get_batch_writer( password=self.__password, secure=self.__secure, ca_certs=self.__ca_certs, - verify=self.__verify, + verify=self.get_verify(), metrics=metrics, statement=insert_statement.with_database(self.__database), encoding=encoding, @@ -504,7 +504,12 @@ def get_ca_certs(self) -> str | None: return self.__ca_certs def get_verify(self) -> bool | None: - return self.__verify + # CLICKHOUSE_VERIFY arrives as a raw env string; coerce once here so + # every client sees the same value. Unset (None) stays None. + verify = self.__verify + if isinstance(verify, str): + return verify.strip().lower() not in ("false", "0") + return verify CLUSTERS = [ @@ -517,7 +522,7 @@ def get_verify(self) -> bool | None: database=cluster.get("database", "default"), secure=cluster.get("secure", False), ca_certs=cluster.get("ca_certs", None), - verify=cluster.get("verify", False), + verify=cluster.get("verify"), storage_sets=cluster["storage_sets"], single_node=cluster["single_node"], cluster_name=cluster.get("cluster_name", None), @@ -558,7 +563,7 @@ def _build_sliced_cluster(cluster: Mapping[str, Any]) -> ClickhouseCluster: database=cluster.get("database", "default"), secure=cluster.get("secure", False), ca_certs=cluster.get("ca_certs", None), - verify=cluster.get("verify", False), + verify=cluster.get("verify"), storage_sets={storage_tuple[0] for storage_tuple in cluster["storage_set_slices"]}, single_node=cluster["single_node"], cluster_name=cluster.get("cluster_name", None), diff --git a/snuba/consumers/consumer_config.py b/snuba/consumers/consumer_config.py index 4d52086eee..4b0c80940a 100644 --- a/snuba/consumers/consumer_config.py +++ b/snuba/consumers/consumer_config.py @@ -20,6 +20,7 @@ class ClickhouseClusterConfig: password: str database: str secure: bool + verify: bool | None @dataclass(frozen=True) @@ -283,6 +284,7 @@ def resolve_storage_config(storage_name: str, storage: WritableTableStorage) -> password=password, secure=cluster.get_secure(), database=cluster.get_database(), + verify=cluster.get_verify(), ) processor = storage.get_table_writer().get_stream_loader().get_processor() diff --git a/tests/clusters/test_verify.py b/tests/clusters/test_verify.py new file mode 100644 index 0000000000..f3241911e5 --- /dev/null +++ b/tests/clusters/test_verify.py @@ -0,0 +1,37 @@ +import pytest + +from snuba.clusters.cluster import ClickhouseCluster + + +@pytest.mark.parametrize( + "raw,expected", + [ + (None, None), + (True, True), + (False, False), + ("true", True), + ("1", True), + ("false", False), + ("FALSE", False), + ("0", False), + (" false ", False), + ("", True), + ("yes", True), + ("garbage", True), + ], +) +def test_get_verify_coercion(raw: bool | str | None, expected: bool | None) -> None: + cluster = ClickhouseCluster( + "127.0.0.1", + 8001, + "default", + "", + "default", + True, + None, + raw, + {"events"}, + True, + ) + + assert cluster.get_verify() == expected diff --git a/tests/consumers/test_consumer_config.py b/tests/consumers/test_consumer_config.py index 222a44ff6c..cbc3803b45 100644 --- a/tests/consumers/test_consumer_config.py +++ b/tests/consumers/test_consumer_config.py @@ -1,6 +1,10 @@ +import dataclasses + import pytest -from snuba.consumers.consumer_config import resolve_consumer_config +from snuba.consumers.consumer_config import resolve_consumer_config, resolve_storage_config +from snuba.datasets.storages.factory import get_writable_storage +from snuba.datasets.storages.storage_key import StorageKey def test_consumer_config() -> None: @@ -19,6 +23,7 @@ def test_consumer_config() -> None: assert len(resolved.storages) == 1 assert resolved.storages[0].clickhouse_table_name in ("errors_local", "errors_dist") + assert resolved.storages[0].clickhouse_cluster.verify is None assert resolved.raw_topic.broker_config["bootstrap.servers"] == "some_server:9092" assert resolved.raw_topic.physical_topic_name == "new-events" assert resolved.raw_topic.logical_topic_name == "events" @@ -48,6 +53,18 @@ def test_consumer_config() -> None: ) +def test_resolve_storage_config_propagates_verify_false( + monkeypatch: pytest.MonkeyPatch, +) -> None: + storage = get_writable_storage(StorageKey.ERRORS) + monkeypatch.setattr(storage.get_cluster(), "get_verify", lambda: False) + + resolved = resolve_storage_config("errors", storage) + + assert resolved.clickhouse_cluster.verify is False + assert dataclasses.asdict(resolved)["clickhouse_cluster"]["verify"] is False + + def test_group_instance_id_in_broker_config() -> None: """Static membership: --group-instance-id lands in librdkafka broker config.""" resolved = resolve_consumer_config(