feat(aws): Add new kafka service (#4001)

Co-authored-by: Sergio Garcia <38561120+sergargar@users.noreply.github.com>
This commit is contained in:
Rubén De la Torre Vico
2024-05-16 14:29:05 +02:00
committed by GitHub
parent 416e406394
commit 4aedba71fd
8 changed files with 428 additions and 0 deletions
@@ -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())
@@ -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 <arn_cluster> --current-version <current_version> --target-version <latest_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": ""
}
@@ -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
@@ -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
@@ -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
@@ -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"