mirror of
https://github.com/prowler-cloud/prowler.git
synced 2026-07-23 12:31:54 +00:00
feat(kafka): add new check kafka_connector_in_transit_encryption_enabled (#5577)
Co-authored-by: Sergio Garcia <38561120+sergargar@users.noreply.github.com>
This commit is contained in:
committed by
GitHub
parent
bc89f4383e
commit
694cee1afb
+34
@@ -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": ""
|
||||
}
|
||||
+23
@@ -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
|
||||
@@ -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
|
||||
|
||||
@@ -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())
|
||||
+130
@@ -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 == []
|
||||
@@ -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"
|
||||
|
||||
Reference in New Issue
Block a user