Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
3 changes: 3 additions & 0 deletions .cfnlintrc.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -108,6 +108,9 @@ ignore_templates:
- tests/translator/output/**/managed_policies_everything.json # intentionally contains wrong arns
- tests/translator/output/**/function_with_metrics_config.json
- tests/translator/output/**/function_with_self_managed_kafka_and_schema_registry.json # cfnlint is not updated to recognize the SchemaRegistryConfig property
- tests/translator/output/**/function_with_self_managed_kafka_oauth.json # cfnlint is not updated to recognize OAuth SourceAccessConfigurations types
- tests/translator/output/**/function_with_self_managed_kafka_iam_auth.json # cfnlint is not updated to recognize IAM_AUTH SourceAccessConfigurations type
- tests/translator/output/**/function_with_self_managed_kafka_iam_oauth.json # cfnlint is not updated to recognize IAM_OAUTHBEARER_AUTH SourceAccessConfigurations type
- tests/translator/output/**/function_with_msk_with_schema_registry_config.json # cfnlint is not updated to recognize the SchemaRegistryConfig property
- tests/translator/output/**/function_with_logging_config.json # cfnlint is not updated to recognize the LoggingConfig property
- tests/translator/output/aws-*/*capacity_provider*.json # Ignore Capacity Provider test format in non-aws partitions
Expand Down
2 changes: 1 addition & 1 deletion samtranslator/__init__.py
Original file line number Diff line number Diff line change
@@ -1 +1 @@
__version__ = "1.113.0"
__version__ = "1.114.0"
29 changes: 25 additions & 4 deletions samtranslator/model/eventsources/pull.py
Original file line number Diff line number Diff line change
Expand Up @@ -664,6 +664,19 @@ class SelfManagedKafka(PullEventSource):
"SASL_SCRAM_512_AUTH",
"BASIC_AUTH",
"CLIENT_CERTIFICATE_TLS_AUTH",
"OAUTHBEARER_AUTH",
"IAM_AUTH",
"IAM_OAUTHBEARER_AUTH",
]
NON_URI_TYPES = [
"IAM_AUTH",
"IAM_OAUTHBEARER_AUTH",
]
OAUTH_METADATA_TYPES = [
"OAUTHBEARER_SCOPE",
"OAUTHBEARER_AUDIENCE",
"OAUTHBEARER_LOGICAL_CLUSTER",
"OAUTHBEARER_IDENTITY_POOL",
]

def get_event_source_arn(self) -> PassThrough | None:
Expand Down Expand Up @@ -699,6 +712,8 @@ def get_policy_statements(
"No SourceAccessConfigurations for self managed kafka event provided.",
)
document = self.generate_policy_document(self.SourceAccessConfigurations, intrinsic_resolver)
if not document["PolicyDocument"]["Statement"]:
return None
return [document]

def generate_policy_document( # type: ignore[no-untyped-def]
Expand All @@ -711,7 +726,7 @@ def generate_policy_document( # type: ignore[no-untyped-def]
statements.append(secret_manager)

if authentication_uri_2:
secret_manager = self.get_secret_manager_secret(authentication_uri) # type: ignore[no-untyped-call]
secret_manager = self.get_secret_manager_secret(authentication_uri_2) # type: ignore[no-untyped-call]
statements.append(secret_manager)

if has_vpc_config:
Expand Down Expand Up @@ -741,6 +756,7 @@ def get_secret_key(self, source_access_configurations: list[Any]) -> tuple[str |
authentication_uri = None
has_vpc_subnet = False
has_vpc_security_group = False
has_auth_mechanism = False
authentication_uri_2 = None

if not isinstance(source_access_configurations, list):
Expand All @@ -759,18 +775,23 @@ def get_secret_key(self, source_access_configurations: list[Any]) -> tuple[str |
has_vpc_security_group = True

elif config.get("Type") in self.AUTH_MECHANISM:
if authentication_uri:
if has_auth_mechanism:
raise InvalidEventException(
self.relative_id,
"Multiple auth mechanism properties specified in SourceAccessConfigurations for self managed kafka event.",
)
self.validate_uri(config.get("URI"), "auth mechanism")
authentication_uri = config.get("URI")
has_auth_mechanism = True
if config.get("Type") not in self.NON_URI_TYPES:
self.validate_uri(config.get("URI"), "auth mechanism")
authentication_uri = config.get("URI")

elif config.get("Type") == "SERVER_ROOT_CA_CERTIFICATE":
self.validate_uri(config.get("URI"), "SERVER_ROOT_CA_CERTIFICATE")
authentication_uri_2 = config.get("URI")

elif config.get("Type") in self.OAUTH_METADATA_TYPES:
self.validate_uri(config.get("URI"), config.get("Type"))

else:
raise InvalidEventException(
self.relative_id,
Expand Down
Loading
Loading