From 694cee1afbf1d3549563d7cf2ce1e1783356cfe4 Mon Sep 17 00:00:00 2001 From: Mario Rodriguez Lopez <101330800+MarioRgzLpz@users.noreply.github.com> Date: Mon, 4 Nov 2024 18:46:32 +0100 Subject: [PATCH] feat(kafka): add new check `kafka_connector_in_transit_encryption_enabled` (#5577) Co-authored-by: Sergio Garcia <38561120+sergargar@users.noreply.github.com> --- .../__init__.py | 0 ...n_transit_encryption_enabled.metadata.json | 34 +++++ ...connector_in_transit_encryption_enabled.py | 23 ++++ .../aws/services/kafka/kafka_service.py | 38 +++++ .../aws/services/kafka/kafkaconnect_client.py | 4 + ...ctor_in_transit_encryption_enabled_test.py | 130 ++++++++++++++++++ .../aws/services/kafka/kafka_service_test.py | 24 +++- 7 files changed, 252 insertions(+), 1 deletion(-) create mode 100644 prowler/providers/aws/services/kafka/kafka_connector_in_transit_encryption_enabled/__init__.py create mode 100644 prowler/providers/aws/services/kafka/kafka_connector_in_transit_encryption_enabled/kafka_connector_in_transit_encryption_enabled.metadata.json create mode 100644 prowler/providers/aws/services/kafka/kafka_connector_in_transit_encryption_enabled/kafka_connector_in_transit_encryption_enabled.py create mode 100644 prowler/providers/aws/services/kafka/kafkaconnect_client.py create mode 100644 tests/providers/aws/services/kafka/kafka_connector_in_transit_encryption_enabled/kafka_connector_in_transit_encryption_enabled_test.py diff --git a/prowler/providers/aws/services/kafka/kafka_connector_in_transit_encryption_enabled/__init__.py b/prowler/providers/aws/services/kafka/kafka_connector_in_transit_encryption_enabled/__init__.py new file mode 100644 index 0000000000..e69de29bb2 diff --git a/prowler/providers/aws/services/kafka/kafka_connector_in_transit_encryption_enabled/kafka_connector_in_transit_encryption_enabled.metadata.json b/prowler/providers/aws/services/kafka/kafka_connector_in_transit_encryption_enabled/kafka_connector_in_transit_encryption_enabled.metadata.json new file mode 100644 index 0000000000..36f40c4c18 --- /dev/null +++ b/prowler/providers/aws/services/kafka/kafka_connector_in_transit_encryption_enabled/kafka_connector_in_transit_encryption_enabled.metadata.json @@ -0,0 +1,34 @@ +{ + "Provider": "aws", + "CheckID": "kafka_connector_in_transit_encryption_enabled", + "CheckTitle": "MSK Connect connectors should be encrypted in transit", + "CheckType": [ + "Software and Configuration Checks/AWS Security Best Practices" + ], + "ServiceName": "kafka", + "SubServiceName": "", + "ResourceIdTemplate": "arn:aws:kafkaconnect:{region}:{account-id}:connector/{connector-name}/{connector-id}", + "Severity": "medium", + "ResourceType": "Other", + "Description": "This control checks whether an Amazon MSK Connect connector is encrypted in transit. This control fails if the connector isn't encrypted in transit.", + "Risk": "Data in transit can be intercepted or eavesdropped on by unauthorized users. Ensuring encryption in transit helps to protect sensitive data as it moves between nodes in a network or from your MSK cluster to connected applications.", + "RelatedUrl": "https://docs.aws.amazon.com/msk/latest/developerguide/msk-connect.html", + "Remediation": { + "Code": { + "CLI": "aws kafkaconnect create-connector --encryption-in-transit-config 'EncryptionInTransitType=TLS'", + "NativeIaC": "", + "Other": "https://docs.aws.amazon.com/securityhub/latest/userguide/msk-controls.html#msk-3", + "Terraform": "" + }, + "Recommendation": { + "Text": "Enable encryption in transit for MSK Connect connectors to secure data as it moves across networks.", + "Url": "https://docs.aws.amazon.com/msk/latest/developerguide/mkc-create-connector-intro.html" + } + }, + "Categories": [ + "encryption" + ], + "DependsOn": [], + "RelatedTo": [], + "Notes": "" +} diff --git a/prowler/providers/aws/services/kafka/kafka_connector_in_transit_encryption_enabled/kafka_connector_in_transit_encryption_enabled.py b/prowler/providers/aws/services/kafka/kafka_connector_in_transit_encryption_enabled/kafka_connector_in_transit_encryption_enabled.py new file mode 100644 index 0000000000..fd058bffce --- /dev/null +++ b/prowler/providers/aws/services/kafka/kafka_connector_in_transit_encryption_enabled/kafka_connector_in_transit_encryption_enabled.py @@ -0,0 +1,23 @@ +from prowler.lib.check.models import Check, Check_Report_AWS +from prowler.providers.aws.services.kafka.kafkaconnect_client import kafkaconnect_client + + +class kafka_connector_in_transit_encryption_enabled(Check): + def execute(self): + findings = [] + + for arn_connector, connector in kafkaconnect_client.connectors.items(): + report = Check_Report_AWS(self.metadata()) + report.region = connector.region + report.resource_id = connector.name + report.resource_arn = arn_connector + report.status = "FAIL" + report.status_extended = f"Kafka connector {connector.name} does not have encryption in transit enabled." + + if connector.encryption_in_transit == "TLS": + report.status = "PASS" + report.status_extended = f"Kafka connector {connector.name} has encryption in transit enabled." + + findings.append(report) + + return findings diff --git a/prowler/providers/aws/services/kafka/kafka_service.py b/prowler/providers/aws/services/kafka/kafka_service.py index fc15990990..a30e8ccafc 100644 --- a/prowler/providers/aws/services/kafka/kafka_service.py +++ b/prowler/providers/aws/services/kafka/kafka_service.py @@ -113,3 +113,41 @@ class Cluster(BaseModel): class KafkaVersion(BaseModel): version: str status: str + + +class KafkaConnect(AWSService): + def __init__(self, provider): + super().__init__(__class__.__name__, provider) + self.connectors = {} + self.__threading_call__(self._list_connectors) + + def _list_connectors(self, regional_client): + try: + connector_paginator = regional_client.get_paginator("list_connectors") + + for page in connector_paginator.paginate(): + for connector in page["connectors"]: + connector_arn = connector["connectorArn"] + + if not self.audit_resources or is_resource_filtered( + connector_arn, self.audit_resources + ): + self.connectors[connector_arn] = Connector( + arn=connector_arn, + name=connector.get("connectorName", ""), + region=regional_client.region, + encryption_in_transit=connector.get( + "kafkaClusterEncryptionInTransit", {} + ).get("encryptionType", "PLAINTEXT"), + ) + except Exception as error: + logger.error( + f"{regional_client.region} -- {error.__class__.__name__}[{error.__traceback__.tb_lineno}]: {error}" + ) + + +class Connector(BaseModel): + name: str + arn: str + region: str + encryption_in_transit: str diff --git a/prowler/providers/aws/services/kafka/kafkaconnect_client.py b/prowler/providers/aws/services/kafka/kafkaconnect_client.py new file mode 100644 index 0000000000..45a456191a --- /dev/null +++ b/prowler/providers/aws/services/kafka/kafkaconnect_client.py @@ -0,0 +1,4 @@ +from prowler.providers.aws.services.kafka.kafka_service import KafkaConnect +from prowler.providers.common.provider import Provider + +kafkaconnect_client = KafkaConnect(Provider.get_global_provider()) diff --git a/tests/providers/aws/services/kafka/kafka_connector_in_transit_encryption_enabled/kafka_connector_in_transit_encryption_enabled_test.py b/tests/providers/aws/services/kafka/kafka_connector_in_transit_encryption_enabled/kafka_connector_in_transit_encryption_enabled_test.py new file mode 100644 index 0000000000..9efc27b9e3 --- /dev/null +++ b/tests/providers/aws/services/kafka/kafka_connector_in_transit_encryption_enabled/kafka_connector_in_transit_encryption_enabled_test.py @@ -0,0 +1,130 @@ +from unittest.mock import patch + +import botocore +from boto3 import client + +from prowler.providers.aws.services.kafka.kafka_service import KafkaConnect +from tests.providers.aws.utils import ( + AWS_ACCOUNT_NUMBER, + AWS_REGION_US_EAST_1, + set_mocked_aws_provider, +) + +orig = botocore.client.BaseClient._make_api_call + + +def mock_make_api_call(self, operation_name, kwarg): + if operation_name == "ListConnectors": + return { + "connectors": [ + { + "connectorName": "connector-plaintext", + "connectorArn": f"arn:aws:kafkaconnect:{AWS_REGION_US_EAST_1}:{AWS_ACCOUNT_NUMBER}:connector/connector-plaintext/058406e6-a8f7-4135-8860-d4786220a395-3", + "kafkaClusterEncryptionInTransit": {"encryptionType": "PLAINTEXT"}, + }, + ], + } + return orig(self, operation_name, kwarg) + + +def mock_make_api_call_v2(self, operation_name, kwarg): + if operation_name == "ListConnectors": + return { + "connectors": [ + { + "connectorName": "connector-tls", + "connectorArn": f"arn:aws:kafkaconnect:{AWS_REGION_US_EAST_1}:{AWS_ACCOUNT_NUMBER}:connector/connector-tls/058406e6-a8f7-4135-8860-d4786220a395-3", + "kafkaClusterEncryptionInTransit": {"encryptionType": "TLS"}, + }, + ], + } + return orig(self, operation_name, kwarg) + + +class Test_kafka_connector_in_transit_encryption_enabled: + def test_kafka_no_connector(self): + + mocked_aws_provider = set_mocked_aws_provider([AWS_REGION_US_EAST_1]) + + with patch( + "prowler.providers.common.provider.Provider.get_global_provider", + return_value=mocked_aws_provider, + ), patch( + "prowler.providers.aws.services.kafka.kafka_connector_in_transit_encryption_enabled.kafka_connector_in_transit_encryption_enabled.kafkaconnect_client", + new=KafkaConnect(mocked_aws_provider), + ): + from prowler.providers.aws.services.kafka.kafka_connector_in_transit_encryption_enabled.kafka_connector_in_transit_encryption_enabled import ( + kafka_connector_in_transit_encryption_enabled, + ) + + check = kafka_connector_in_transit_encryption_enabled() + result = check.execute() + + assert len(result) == 0 + + @patch("botocore.client.BaseClient._make_api_call", new=mock_make_api_call) + def test_kafka_cluster_not_using_in_transit_encryption(self): + client("kafkaconnect", region_name=AWS_REGION_US_EAST_1) + + mocked_aws_provider = set_mocked_aws_provider([AWS_REGION_US_EAST_1]) + + with patch( + "prowler.providers.common.provider.Provider.get_global_provider", + return_value=mocked_aws_provider, + ), patch( + "prowler.providers.aws.services.kafka.kafka_connector_in_transit_encryption_enabled.kafka_connector_in_transit_encryption_enabled.kafkaconnect_client", + new=KafkaConnect(mocked_aws_provider), + ): + from prowler.providers.aws.services.kafka.kafka_connector_in_transit_encryption_enabled.kafka_connector_in_transit_encryption_enabled import ( + kafka_connector_in_transit_encryption_enabled, + ) + + check = kafka_connector_in_transit_encryption_enabled() + result = check.execute() + + assert len(result) == 1 + assert result[0].status == "FAIL" + assert ( + result[0].status_extended + == "Kafka connector connector-plaintext does not have encryption in transit enabled." + ) + assert result[0].resource_id == "connector-plaintext" + assert ( + result[0].resource_arn + == f"arn:aws:kafkaconnect:{AWS_REGION_US_EAST_1}:{AWS_ACCOUNT_NUMBER}:connector/connector-plaintext/058406e6-a8f7-4135-8860-d4786220a395-3" + ) + assert result[0].region == AWS_REGION_US_EAST_1 + assert result[0].resource_tags == [] + + @patch("botocore.client.BaseClient._make_api_call", new=mock_make_api_call_v2) + def test_kafka_cluster_using_in_transit_encryption(self): + + mocked_aws_provider = set_mocked_aws_provider([AWS_REGION_US_EAST_1]) + + with patch( + "prowler.providers.common.provider.Provider.get_global_provider", + return_value=mocked_aws_provider, + ), patch( + "prowler.providers.aws.services.kafka.kafka_connector_in_transit_encryption_enabled.kafka_connector_in_transit_encryption_enabled.kafkaconnect_client", + new=KafkaConnect(mocked_aws_provider), + ): + from prowler.providers.aws.services.kafka.kafka_connector_in_transit_encryption_enabled.kafka_connector_in_transit_encryption_enabled import ( + kafka_connector_in_transit_encryption_enabled, + ) + + check = kafka_connector_in_transit_encryption_enabled() + result = check.execute() + + assert len(result) == 1 + assert result[0].status == "PASS" + assert ( + result[0].status_extended + == "Kafka connector connector-tls has encryption in transit enabled." + ) + assert result[0].resource_id == "connector-tls" + assert ( + result[0].resource_arn + == f"arn:aws:kafkaconnect:{AWS_REGION_US_EAST_1}:{AWS_ACCOUNT_NUMBER}:connector/connector-tls/058406e6-a8f7-4135-8860-d4786220a395-3" + ) + assert result[0].region == AWS_REGION_US_EAST_1 + assert result[0].resource_tags == [] diff --git a/tests/providers/aws/services/kafka/kafka_service_test.py b/tests/providers/aws/services/kafka/kafka_service_test.py index 0f8c216b24..83e79d6d28 100644 --- a/tests/providers/aws/services/kafka/kafka_service_test.py +++ b/tests/providers/aws/services/kafka/kafka_service_test.py @@ -2,7 +2,7 @@ from unittest.mock import patch import botocore -from prowler.providers.aws.services.kafka.kafka_service import Kafka +from prowler.providers.aws.services.kafka.kafka_service import Kafka, KafkaConnect from tests.providers.aws.utils import ( AWS_ACCOUNT_NUMBER, AWS_REGION_US_EAST_1, @@ -63,6 +63,16 @@ def mock_make_api_call(self, operation_name, kwarg): {"Version": "2.8.0", "Status": "ACTIVE"}, ] } + elif operation_name == "ListConnectors": + return { + "connectors": [ + { + "connectorName": "demo-connector", + "connectorArn": f"arn:aws:kafkaconnect:{AWS_REGION_US_EAST_1}:{AWS_ACCOUNT_NUMBER}:connector/demo-connector/058406e6-a8f7-4135-8860-d4786220a395-3", + "kafkaClusterEncryptionInTransit": {"encryptionType": "PLAINTEXT"}, + }, + ], + } return make_api_call(self, operation_name, kwarg) @@ -107,3 +117,15 @@ class TestKafkaService: assert kafka.kafka_versions[0].status == "DEPRECATED" assert kafka.kafka_versions[1].version == "2.8.0" assert kafka.kafka_versions[1].status == "ACTIVE" + + @patch("botocore.client.BaseClient._make_api_call", new=mock_make_api_call) + def test_list_connectors(self): + kafka = KafkaConnect(set_mocked_aws_provider([AWS_REGION_US_EAST_1])) + + assert len(kafka.connectors) == 1 + connector_arn = f"arn:aws:kafkaconnect:{AWS_REGION_US_EAST_1}:{AWS_ACCOUNT_NUMBER}:connector/demo-connector/058406e6-a8f7-4135-8860-d4786220a395-3" + assert connector_arn in kafka.connectors + assert kafka.connectors[connector_arn].name == "demo-connector" + assert kafka.connectors[connector_arn].arn == connector_arn + assert kafka.connectors[connector_arn].region == AWS_REGION_US_EAST_1 + assert kafka.connectors[connector_arn].encryption_in_transit == "PLAINTEXT"