diff --git a/prowler/providers/aws/services/kafka/__init__.py b/prowler/providers/aws/services/kafka/__init__.py new file mode 100644 index 0000000000..e69de29bb2 diff --git a/prowler/providers/aws/services/kafka/kafka_client.py b/prowler/providers/aws/services/kafka/kafka_client.py new file mode 100644 index 0000000000..ceb8a7e7b4 --- /dev/null +++ b/prowler/providers/aws/services/kafka/kafka_client.py @@ -0,0 +1,4 @@ +from prowler.providers.aws.services.kafka.kafka_service import Kafka +from prowler.providers.common.common import get_global_provider + +kafka_client = Kafka(get_global_provider()) diff --git a/prowler/providers/aws/services/kafka/kafka_cluster_uses_latest_version/__init__.py b/prowler/providers/aws/services/kafka/kafka_cluster_uses_latest_version/__init__.py new file mode 100644 index 0000000000..e69de29bb2 diff --git a/prowler/providers/aws/services/kafka/kafka_cluster_uses_latest_version/kafka_cluster_uses_latest_version.metadata.json b/prowler/providers/aws/services/kafka/kafka_cluster_uses_latest_version/kafka_cluster_uses_latest_version.metadata.json new file mode 100644 index 0000000000..64cf77e06f --- /dev/null +++ b/prowler/providers/aws/services/kafka/kafka_cluster_uses_latest_version/kafka_cluster_uses_latest_version.metadata.json @@ -0,0 +1,32 @@ +{ + "Provider": "aws", + "CheckID": "kafka_cluster_uses_latest_version", + "CheckTitle": "MSK cluster should use the latest version.", + "CheckType": [ + "Infrastructure Security" + ], + "ServiceName": "kafka", + "SubServiceName": "cluster", + "ResourceIdTemplate": "arn:partition:kafka:region:account-id:cluster", + "Severity": "medium", + "ResourceType": "", + "Description": "Ensure that your Amazon Managed Streaming for Apache Kafka (MSK) cluster is using the latest version to benefit from the latest security features, bug fixes, and performance improvements.", + "Risk": "Running an outdated version of Amazon MSK may expose your cluster to security vulnerabilities, bugs, and performance issues.", + "RelatedUrl": "https://docs.aws.amazon.com/lightsail/latest/userguide/amazon-lightsail-databases.html", + "Remediation": { + "Code": { + "CLI": "aws kafka update-cluster-configuration --cluster-arn --current-version --target-version ", + "NativeIaC": "", + "Other": "https://www.trendmicro.com/cloudoneconformity/knowledge-base/aws/MSK/enable-apache-kafka-latest-security-features.html", + "Terraform": "" + }, + "Recommendation": { + "Text": "To upgrade your Amazon MSK cluster to the latest version, use the AWS Management Console, AWS CLI, or SDKs to update the cluster configuration. For more information, refer to the official Amazon MSK documentation.", + "Url": "https://docs.aws.amazon.com/msk/latest/developerguide/version-support.html#version-upgrades" + } + }, + "Categories": [], + "DependsOn": [], + "RelatedTo": [], + "Notes": "" +} diff --git a/prowler/providers/aws/services/kafka/kafka_cluster_uses_latest_version/kafka_cluster_uses_latest_version.py b/prowler/providers/aws/services/kafka/kafka_cluster_uses_latest_version/kafka_cluster_uses_latest_version.py new file mode 100644 index 0000000000..10de7bf597 --- /dev/null +++ b/prowler/providers/aws/services/kafka/kafka_cluster_uses_latest_version/kafka_cluster_uses_latest_version.py @@ -0,0 +1,28 @@ +from prowler.lib.check.models import Check, Check_Report_AWS +from prowler.providers.aws.services.kafka.kafka_client import kafka_client + + +class kafka_cluster_uses_latest_version(Check): + def execute(self): + findings = [] + + for arn_cluster, cluster in kafka_client.clusters.items(): + report = Check_Report_AWS(self.metadata()) + report.region = cluster.region + report.resource_id = cluster.id + report.resource_arn = arn_cluster + report.resource_tags = cluster.tags + report.status = "PASS" + report.status_extended = ( + f"Kafka cluster '{cluster.name}' is using the latest version." + ) + + if cluster.kafka_version != kafka_client.kafka_versions[-1].version: + report.status = "FAIL" + report.status_extended = ( + f"Kafka cluster '{cluster.name}' is not using the latest version." + ) + + 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 new file mode 100644 index 0000000000..65592fcf5e --- /dev/null +++ b/prowler/providers/aws/services/kafka/kafka_service.py @@ -0,0 +1,115 @@ +from pydantic import BaseModel + +from prowler.lib.logger import logger +from prowler.lib.scan_filters.scan_filters import is_resource_filtered +from prowler.providers.aws.lib.service.service import AWSService + + +class Kafka(AWSService): + def __init__(self, provider): + super().__init__(__class__.__name__, provider) + self.account_arn_template = f"arn:{self.audited_partition}:kafka:{self.region}:{self.audited_account}:cluster" + self.clusters = {} + self.__threading_call__(self.__list_clusters__) + self.kafka_versions = [] + self.__threading_call__(self.__list_kafka_versions__) + + def __list_clusters__(self, regional_client): + try: + cluster_paginator = regional_client.get_paginator("list_clusters") + + for page in cluster_paginator.paginate(): + for cluster in page["ClusterInfoList"]: + arn = cluster.get( + "ClusterArn", + f"{self.account_arn_template}/{cluster.get('ClusterName', '')}", + ) + + if not self.audit_resources or is_resource_filtered( + arn, self.audit_resources + ): + self.clusters[cluster.get("ClusterArn", "")] = Cluster( + id=arn.split(":")[-1].split("/")[-1], + name=cluster.get("ClusterName", ""), + region=regional_client.region, + tags=list(cluster.get("Tags", {})), + state=cluster.get("State", ""), + kafka_version=cluster.get( + "CurrentBrokerSoftwareInfo", {} + ).get("KafkaVersion", ""), + data_volume_kms_key_id=cluster.get("EncryptionInfo", {}) + .get("EncryptionAtRest", {}) + .get("DataVolumeKMSKeyId", ""), + encryption_in_transit=EncryptionInTransit( + client_broker=cluster.get("EncryptionInfo", {}) + .get("EncryptionInTransit", {}) + .get("ClientBroker", "PLAINTEXT"), + in_cluster=cluster.get("EncryptionInfo", {}) + .get("EncryptionInTransit", {}) + .get("InCluster", False), + ), + tls_authentication=cluster.get("ClientAuthentication", {}) + .get("Tls", {}) + .get("Enabled", False), + public_access=cluster.get("BrokerNodeGroupInfo", {}) + .get("ConenctivityInfo", {}) + .get("PublicAccess", {}) + .get("Type", "SERVICE_PROVIDED_EIPS") + != "DISABLED", + unauthentication_access=cluster.get( + "ClientAuthentication", {} + ) + .get("Unauthenticated", {}) + .get("Enabled", False), + enhanced_monitoring=cluster.get( + "EnhancedMonitoring", "DEFAULT" + ), + ) + except Exception as error: + logger.error( + f"{regional_client.region} -- {error.__class__.__name__}[{error.__traceback__.tb_lineno}]: {error}" + ) + + def __list_kafka_versions__(self, regional_client): + try: + kafka_versions_paginator = regional_client.get_paginator( + "list_kafka_versions" + ) + + for page in kafka_versions_paginator.paginate(): + for version in page["KafkaVersions"]: + self.kafka_versions.append( + KafkaVersion( + version=version.get("Version", "UNKNOWN"), + status=version.get("Status", "UNKNOWN"), + ) + ) + except Exception as error: + logger.error( + f"{regional_client.region} -- {error.__class__.__name__}[{error.__traceback__.tb_lineno}]: {error}" + ) + + +class EncryptionInTransit(BaseModel): + client_broker: str + in_cluster: bool + + +class Cluster(BaseModel): + id: str + name: str + region: str + tags: list + kafka_version: str + state: str + data_volume_kms_key_id: str + encryption_in_transit: EncryptionInTransit + tls_authentication: bool + public_access: bool + unauthentication_access: bool + enhanced_monitoring: str + + +class KafkaVersion(BaseModel): + version: str + status: str diff --git a/tests/providers/aws/services/kafka/kafka_cluster_uses_latest_version/kafka_cluster_uses_latest_version_test.py b/tests/providers/aws/services/kafka/kafka_cluster_uses_latest_version/kafka_cluster_uses_latest_version_test.py new file mode 100644 index 0000000000..9ef950045d --- /dev/null +++ b/tests/providers/aws/services/kafka/kafka_cluster_uses_latest_version/kafka_cluster_uses_latest_version_test.py @@ -0,0 +1,140 @@ +from unittest.mock import MagicMock, patch + +from prowler.providers.aws.services.kafka.kafka_service import ( + Cluster, + EncryptionInTransit, + KafkaVersion, +) +from tests.providers.aws.utils import AWS_REGION_US_EAST_1, set_mocked_aws_provider + + +class Test_kafka_cluster_latest_version: + def test_kafka_no_clusters(self): + kafka_client = MagicMock + kafka_client.clusters = {} + + with patch( + "prowler.providers.common.common.get_global_provider", + return_value=set_mocked_aws_provider([AWS_REGION_US_EAST_1]), + ), patch( + "prowler.providers.aws.services.kafka.kafka_service.Kafka", + new=kafka_client, + ): + from prowler.providers.aws.services.kafka.kafka_cluster_uses_latest_version.kafka_cluster_uses_latest_version import ( + kafka_cluster_uses_latest_version, + ) + + check = kafka_cluster_uses_latest_version() + result = check.execute() + + assert len(result) == 0 + + def test_kafka_cluster_not_using_latest_version(self): + kafka_client = MagicMock + kafka_client.clusters = { + "arn:aws:kafka:us-east-1:123456789012:cluster/demo-cluster-1/6357e0b2-0e6a-4b86-a0b4-70df934c2e31-5": Cluster( + id="6357e0b2-0e6a-4b86-a0b4-70df934c2e31-5", + name="demo-cluster-1", + region=AWS_REGION_US_EAST_1, + tags=[], + state="ACTIVE", + kafka_version="2.2.1", + data_volume_kms_key_id=f"arn:aws:kms:{AWS_REGION_US_EAST_1}:123456789012:key/a7ca56d5-0768-4b64-a670-339a9fbef81c", + encryption_in_transit=EncryptionInTransit( + client_broker="TLS_PLAINTEXT", + in_cluster=True, + ), + tls_authentication=True, + public_access=True, + unauthentication_access=False, + enhanced_monitoring="DEFAULT", + ) + } + + kafka_client.kafka_versions = [ + KafkaVersion(version="1.0.0", status="DEPRECATED"), + KafkaVersion(version="2.8.0", status="ACTIVE"), + ] + + with patch( + "prowler.providers.common.common.get_global_provider", + return_value=set_mocked_aws_provider([AWS_REGION_US_EAST_1]), + ), patch( + "prowler.providers.aws.services.kafka.kafka_service.Kafka", + new=kafka_client, + ): + from prowler.providers.aws.services.kafka.kafka_cluster_uses_latest_version.kafka_cluster_uses_latest_version import ( + kafka_cluster_uses_latest_version, + ) + + check = kafka_cluster_uses_latest_version() + result = check.execute() + + assert len(result) == 1 + assert result[0].status == "FAIL" + assert ( + result[0].status_extended + == "Kafka cluster 'demo-cluster-1' is not using the latest version." + ) + assert result[0].resource_id == "6357e0b2-0e6a-4b86-a0b4-70df934c2e31-5" + assert ( + result[0].resource_arn + == "arn:aws:kafka:us-east-1:123456789012:cluster/demo-cluster-1/6357e0b2-0e6a-4b86-a0b4-70df934c2e31-5" + ) + assert result[0].resource_tags == [] + assert result[0].region == AWS_REGION_US_EAST_1 + + def test_kafka_cluster_using_latest_version_pass(self): + kafka_client = MagicMock + kafka_client.clusters = { + "arn:aws:kafka:us-east-1:123456789012:cluster/demo-cluster-1/6357e0b2-0e6a-4b86-a0b4-70df934c2e31-5": Cluster( + id="6357e0b2-0e6a-4b86-a0b4-70df934c2e31-5", + name="demo-cluster-1", + region=AWS_REGION_US_EAST_1, + tags=[], + state="ACTIVE", + kafka_version="2.8.0", + data_volume_kms_key_id=f"arn:aws:kms:{AWS_REGION_US_EAST_1}:123456789012:key/a7ca56d5-0768-4b64-a670-339a9fbef81c", + encryption_in_transit=EncryptionInTransit( + client_broker="TLS_PLAINTEXT", + in_cluster=True, + ), + tls_authentication=True, + public_access=True, + unauthentication_access=False, + enhanced_monitoring="DEFAULT", + ) + } + + kafka_client.kafka_versions = [ + KafkaVersion(version="1.0.0", status="DEPRECATED"), + KafkaVersion(version="2.8.0", status="ACTIVE"), + ] + + with patch( + "prowler.providers.common.common.get_global_provider", + return_value=set_mocked_aws_provider([AWS_REGION_US_EAST_1]), + ), patch( + "prowler.providers.aws.services.kafka.kafka_service.Kafka", + new=kafka_client, + ): + from prowler.providers.aws.services.kafka.kafka_cluster_uses_latest_version.kafka_cluster_uses_latest_version import ( + kafka_cluster_uses_latest_version, + ) + + check = kafka_cluster_uses_latest_version() + result = check.execute() + + assert len(result) == 1 + assert result[0].status == "PASS" + assert ( + result[0].status_extended + == "Kafka cluster 'demo-cluster-1' is using the latest version." + ) + assert result[0].resource_id == "6357e0b2-0e6a-4b86-a0b4-70df934c2e31-5" + assert ( + result[0].resource_arn + == "arn:aws:kafka:us-east-1:123456789012:cluster/demo-cluster-1/6357e0b2-0e6a-4b86-a0b4-70df934c2e31-5" + ) + assert result[0].resource_tags == [] + assert result[0].region == AWS_REGION_US_EAST_1 diff --git a/tests/providers/aws/services/kafka/kafka_service_test.py b/tests/providers/aws/services/kafka/kafka_service_test.py new file mode 100644 index 0000000000..0f8c216b24 --- /dev/null +++ b/tests/providers/aws/services/kafka/kafka_service_test.py @@ -0,0 +1,109 @@ +from unittest.mock import patch + +import botocore + +from prowler.providers.aws.services.kafka.kafka_service import Kafka +from tests.providers.aws.utils import ( + AWS_ACCOUNT_NUMBER, + AWS_REGION_US_EAST_1, + set_mocked_aws_provider, +) + +make_api_call = botocore.client.BaseClient._make_api_call + + +def mock_make_api_call(self, operation_name, kwarg): + if operation_name == "ListClusters": + return { + "ClusterInfoList": [ + { + "BrokerNodeGroupInfo": { + "BrokerAZDistribution": "DEFAULT", + "ClientSubnets": ["subnet-cbfff283", "subnet-6746046b"], + "InstanceType": "kafka.m5.large", + "SecurityGroups": ["sg-f839b688"], + "StorageInfo": {"EbsStorageInfo": {"VolumeSize": 100}}, + }, + "ClusterArn": f"arn:aws:kafka:{AWS_REGION_US_EAST_1}:123456789012:cluster/demo-cluster-1/6357e0b2-0e6a-4b86-a0b4-70df934c2e31-5", + "ClusterName": "demo-cluster-1", + "CreationTime": "2020-07-09T02:31:36.223000+00:00", + "CurrentBrokerSoftwareInfo": {"KafkaVersion": "2.2.1"}, + "CurrentVersion": "K3AEGXETSR30VB", + "EncryptionInfo": { + "EncryptionAtRest": { + "DataVolumeKMSKeyId": f"arn:aws:kms:{AWS_REGION_US_EAST_1}:123456789012:key/a7ca56d5-0768-4b64-a670-339a9fbef81c" + }, + "EncryptionInTransit": { + "ClientBroker": "TLS_PLAINTEXT", + "InCluster": True, + }, + }, + "ClientAuthentication": { + "Tls": {"CertificateAuthorityArnList": [], "Enabled": True}, + "Unauthenticated": {"Enabled": False}, + }, + "EnhancedMonitoring": "DEFAULT", + "OpenMonitoring": { + "Prometheus": { + "JmxExporter": {"EnabledInBroker": False}, + "NodeExporter": {"EnabledInBroker": False}, + } + }, + "NumberOfBrokerNodes": 2, + "State": "ACTIVE", + "Tags": {}, + "ZookeeperConnectString": f"z-2.demo-cluster-1.xuy0sb.c5.kafka.{AWS_REGION_US_EAST_1}.amazonaws.com:2181,z-1.demo-cluster-1.xuy0sb.c5.kafka.{AWS_REGION_US_EAST_1}.amazonaws.com:2181,z-3.demo-cluster-1.xuy0sb.c5.kafka.{AWS_REGION_US_EAST_1}.amazonaws.com:2181", + } + ] + } + elif operation_name == "ListKafkaVersions": + return { + "KafkaVersions": [ + {"Version": "1.0.0", "Status": "DEPRECATED"}, + {"Version": "2.8.0", "Status": "ACTIVE"}, + ] + } + return make_api_call(self, operation_name, kwarg) + + +class TestKafkaService: + @patch("botocore.client.BaseClient._make_api_call", new=mock_make_api_call) + def test_service(self): + kafka = Kafka(set_mocked_aws_provider([AWS_REGION_US_EAST_1])) + + # General assertions + assert kafka.service == "kafka" + assert kafka.__class__.__name__ == "Kafka" + assert kafka.session.__class__.__name__ == "Session" + assert kafka.audited_account == AWS_ACCOUNT_NUMBER + # Clusters assertions + assert len(kafka.clusters) == 1 + cluster_arn = f"arn:aws:kafka:{AWS_REGION_US_EAST_1}:123456789012:cluster/demo-cluster-1/6357e0b2-0e6a-4b86-a0b4-70df934c2e31-5" + assert cluster_arn in kafka.clusters + assert ( + kafka.clusters[cluster_arn].id == "6357e0b2-0e6a-4b86-a0b4-70df934c2e31-5" + ) + assert kafka.clusters[cluster_arn].name == "demo-cluster-1" + assert kafka.clusters[cluster_arn].region == AWS_REGION_US_EAST_1 + assert kafka.clusters[cluster_arn].tags == [] + assert kafka.clusters[cluster_arn].state == "ACTIVE" + assert kafka.clusters[cluster_arn].kafka_version == "2.2.1" + assert ( + kafka.clusters[cluster_arn].data_volume_kms_key_id + == f"arn:aws:kms:{AWS_REGION_US_EAST_1}:123456789012:key/a7ca56d5-0768-4b64-a670-339a9fbef81c" + ) + assert ( + kafka.clusters[cluster_arn].encryption_in_transit.client_broker + == "TLS_PLAINTEXT" + ) + assert kafka.clusters[cluster_arn].encryption_in_transit.in_cluster + assert kafka.clusters[cluster_arn].enhanced_monitoring == "DEFAULT" + assert kafka.clusters[cluster_arn].tls_authentication + assert kafka.clusters[cluster_arn].public_access + assert not kafka.clusters[cluster_arn].unauthentication_access + # Kafka versions assertions + assert len(kafka.kafka_versions) == 2 + assert kafka.kafka_versions[0].version == "1.0.0" + assert kafka.kafka_versions[0].status == "DEPRECATED" + assert kafka.kafka_versions[1].version == "2.8.0" + assert kafka.kafka_versions[1].status == "ACTIVE"