"""Test Kafka Connect heplers.""" from unittest.mock import Mock from unittest.mock import patch import pytest from dbdeploy.util.kafka import connect @pytest.fixture def mock_requests(): """Mock requests.""" requests_path = 'dbdeploy.util.kafka.connect.requests' with patch(requests_path) as requests: yield requests class TestKafkaConnect: """Test Kafka connect helpers.""" def test_create_connector(self, mock_requests): """Test kafka connect connector created.""" response = Mock() response.status_code = 201 mock_requests.put.return_value = response uri = 'http://test_domain.io' connector_name = 'test_connector' connector_config = {'...': '...'} test_result = connect.create_connector( uri=uri, connector_name=connector_name, connector_config=connector_config) assert test_result == (True, None) def test_create_connector_fails(self, mock_requests): """Test kafka connect connector creation fails.""" response = Mock() response.status_code = 403 mock_requests.put.return_value = response uri = 'http://test_domain.io' connector_name = 'test_connector' connector_config = {'...': '...'} test_result = connect.create_connector( uri=uri, connector_name=connector_name, connector_config=connector_config) expected_msg = ('Cannot create connector using provided configuration.' f' Status code: {response.status_code}.') assert test_result == (False, expected_msg) def test_delete_connector(self, mock_requests): """Test kafka connect connector deleted.""" response = Mock() response.status_code = 200 mock_requests.delete.return_value = response uri = 'http://test_domain.io' connector_name = 'test_connector' test_result = connect.delete_connector( uri=uri, connector_name=connector_name) assert test_result == (True, None) def test_delete_connector_failed(self, mock_requests): """Test kafka connect connector deletion fails.""" response = Mock() response.status_code = 403 mock_requests.delete.return_value = response uri = 'http://test_domain.io' connector_name = 'test_connector' test_result = connect.delete_connector( uri=uri, connector_name=connector_name) expected_msg = ('Cannot delete connector. ' f'Status code: {response.status_code}.') assert test_result == (False, expected_msg) def test_get_neo4j_sink_connector_config(self): """Generate connector configuration.""" test_result = connect.get_neo4j_sink_connector_config( cypher_query='MERGE (n:Node)', dlq_topic_name='dlq.topic', environment='test', kafka_bootstrap_brokers='b-1:9094,b-2:9094,b-3:9094', neo4j_server_uri='neo4j+s://test_domain.io', topic='test_topic', batch_size=2500, tasks_max=1, aws_region='us-east-1', kafka_security_protocol='SSL', service_name='kafka-connect-db-deploy', secret_key='NEO4J_CREDENTIALS' ) expected_config_keys = { 'tasks.max', 'topics', 'neo4j.cypher.topic.test_topic', 'errors.deadletterqueue.topic.name', 'kafka.bootstrap.servers', 'kafka.security.protocol', 'neo4j.uri', 'config.providers', 'config.providers.aws.class', 'config.providers.aws.param.aws.region', 'config.providers.aws.param.aws.secret.key', 'config.providers.aws.param.aws.auth.method', 'config.providers.aws.param.aws.access.key', 'consumer.override.max.poll.interval.ms', 'neo4j.authentication.basic.username', 'neo4j.authentication.basic.password', 'connector.class', 'neo4j.batch-size', 'neo4j.batch-timeout', 'neo4j.max-retry-time', 'key.converter', 'value.converter', 'key.converter.schemas.enable', 'value.converter.schemas.enable', 'errors.retry.timeout', 'neo4j.encryption.enabled', 'errors.retry.delay.max.ms', 'errors.tolerance', 'errors.deadletterqueue.context.headers.enable' } assert isinstance(test_result, dict) assert set(test_result.keys()) == expected_config_keys def test_get_jdbc_sink_connector_config_default_insert_mode(self): """Test JDBC connector configuration with default insert mode.""" test_result = connect.get_jdbc_sink_connector_config( environment='test', topic='test_topic', dlq_topic_name='dlq.topic', table_name='test_table', primary_keys='id,name', batch_size=2500, tasks_max=1, aws_region='us-east-1', service_name='kafka-connect-db-deploy', secret_key='MYSQL_CREDENTIALS' ) expected_config_keys = { 'tasks.max', 'config.providers', 'config.providers.aws.class', 'config.providers.aws.param.aws.auth.method', 'config.providers.aws.param.aws.access.key', 'config.providers.aws.param.aws.secret.key', 'config.providers.aws.param.aws.region', 'connection.user', 'connection.password', 'connection.url', 'topics', 'errors.deadletterqueue.topic.name', 'connector.class', 'key.converter', 'value.converter', 'key.converter.schemas.enable', 'value.converter.schemas.enable', 'transforms', 'transforms.ValueToKey.type', 'transforms.ValueToKey.fields', 'transforms.RegexRouter.type', 'transforms.RegexRouter.regex', 'transforms.RegexRouter.replacement', 'pk.mode', 'pk.fields', 'insert.mode', 'auto.create', 'auto.evolve', 'dialect.name', 'table.name.format', 'delete.enabled', 'batch.size', 'errors.tolerance' } assert isinstance(test_result, dict) assert set(test_result.keys()) == expected_config_keys assert test_result['insert.mode'] == 'update' # Default value def test_get_jdbc_sink_connector_config_with_upsert_mode(self): """Test JDBC connector configuration with upsert insert mode.""" test_result = connect.get_jdbc_sink_connector_config( environment='test', topic='test_topic', dlq_topic_name='dlq.topic', table_name='test_table', primary_keys='id,name', insert_mode='upsert', batch_size=2500, tasks_max=1, aws_region='us-east-1', service_name='kafka-connect-db-deploy', secret_key='MYSQL_CREDENTIALS' ) assert isinstance(test_result, dict) assert test_result['insert.mode'] == 'upsert' def test_get_jdbc_sink_connector_config_with_insert_mode(self): """Test JDBC connector configuration with insert mode.""" test_result = connect.get_jdbc_sink_connector_config( environment='test', topic='test_topic', dlq_topic_name='dlq.topic', table_name='test_table', primary_keys='id,name', insert_mode='insert', batch_size=2500, tasks_max=1, aws_region='us-east-1', service_name='kafka-connect-db-deploy', secret_key='MYSQL_CREDENTIALS' ) assert isinstance(test_result, dict) assert test_result['insert.mode'] == 'insert' def test_get_jdbc_sink_connector_config_with_delete_mode(self): """Test JDBC connector configuration with delete mode.""" test_result = connect.get_jdbc_sink_connector_config( environment='test', topic='test_topic', dlq_topic_name='dlq.topic', table_name='test_table', primary_keys='id,name', insert_mode='delete', batch_size=2500, tasks_max=1, aws_region='us-east-1', service_name='kafka-connect-db-deploy', secret_key='MYSQL_CREDENTIALS' ) assert isinstance(test_result, dict) assert test_result['insert.mode'] == 'delete' def test_get_jdbc_sink_connector_config_invalid_insert_mode_fails(self): """Test JDBC connector configuration fails with invalid insert mode.""" with pytest.raises(ValueError, match="Invalid insert_mode 'invalid_mode'"): connect.get_jdbc_sink_connector_config( environment='test', topic='test_topic', dlq_topic_name='dlq.topic', table_name='test_table', primary_keys='id,name', insert_mode='invalid_mode', batch_size=2500, tasks_max=1, aws_region='us-east-1', service_name='kafka-connect-db-deploy', secret_key='MYSQL_CREDENTIALS' ) def test_get_jdbc_sink_connector_config_all_valid_modes(self): """Test JDBC connector configuration with all valid insert modes.""" from dbdeploy.base import constants for insert_mode in constants.JDBC_INSERT_MODES: test_result = connect.get_jdbc_sink_connector_config( environment='test', topic='test_topic', dlq_topic_name='dlq.topic', table_name='test_table', primary_keys='id,name', insert_mode=insert_mode, batch_size=2500, tasks_max=1, aws_region='us-east-1', service_name='kafka-connect-db-deploy', secret_key='MYSQL_CREDENTIALS' ) assert isinstance(test_result, dict) assert test_result['insert.mode'] == insert_mode