# In-cluster Kafka implementation (Bitnami/Strimzi etc.) via the Mongey/kafka # provider (Kafka Admin API). The topic/user flattening locals mirror # modules/kafka-topics-yc 1:1 so the same infrastructure.yaml declaration # produces the same logical result; only the backing resources differ: # yandex_mdb_kafka_topic -> kafka_topic # yandex_mdb_kafka_user -> kafka_user_scram_credential # (YC role permissions) -> kafka_acl (roles mapped to ACL operations) # # NOTE (design open question): Mongey/kafka SCRAM + ACL support must be verified # against the target Kafka version before production use. locals { raw_topics = var.topics topics = { for topic in local.raw_topics : topic.name => { name = topic.name owner = topic.owner cluster_ref = topic.clusterRef cluster = var.kafka_cluster_refs[topic.clusterRef] partitions = try( tonumber(topic.partitions[var.environment]), try(tonumber(topic.partitions._default), try(tonumber(topic.partitions), var.kafka_cluster_refs[topic.clusterRef].default_partitions)) ) replication_factor = try( tonumber(topic.replicationFactor[var.environment]), try(tonumber(topic.replicationFactor._default), try(tonumber(topic.replicationFactor), var.kafka_cluster_refs[topic.clusterRef].default_replication_factor)) ) config = { for key, value in merge( try(topic.inheritDefaultConfig, true) ? try(var.kafka_cluster_refs[topic.clusterRef].default_topic_config, {}) : {}, try(topic.config, {}) ) : key => try(tostring(value[var.environment]), try(tostring(value._default), tostring(value))) } deletion_policy = try(topic.deletionPolicy, "orphan") # Roles granted to the topic owner on its own topic; falls back to the # platform default when ownerRoles is not declared. owner_roles = try(topic.ownerRoles, null) == null ? var.default_user_roles : topic.ownerRoles } } # Access grants declared per user under environments..kafka.users. user_grants = flatten([ for user in var.users : [ for permission in try(user.permissions, []) : [ for role in permission.roles : { cluster_ref = user.clusterRef user = user.name topic_name = permission.topic role = role } ] ] ]) # Every (cluster, user, topic, role) tuple: owner roles on own topics plus # user-side grants. permission_tuples = concat( flatten([ for topic in values(local.topics) : [ for role in topic.owner_roles : { cluster_ref = topic.cluster_ref user = topic.owner topic_name = topic.name role = role } ] ]), local.user_grants ) user_clusters = { for item in distinct(concat( [for tuple in local.permission_tuples : "${tuple.cluster_ref}:${tuple.user}"], [for user in var.users : "${user.clusterRef}:${user.name}"] )) : item => { cluster_ref = split(":", item)[0] user = split(":", item)[1] cluster = var.kafka_cluster_refs[split(":", item)[0]] } } user_permissions = { for user_key, user in local.user_clusters : user_key => distinct([ for tuple in local.permission_tuples : { topic_name = tuple.topic_name role = tuple.role } if "${tuple.cluster_ref}:${tuple.user}" == user_key ]) } # Abstract role -> concrete Kafka ACL operations. Kept deliberately small; # extend as the platform role vocabulary grows. role_operations = { ACCESS_ROLE_PRODUCER = [ { resource_type = "Topic", operation = "Write" }, { resource_type = "Topic", operation = "Describe" }, ] ACCESS_ROLE_CONSUMER = [ { resource_type = "Topic", operation = "Read" }, { resource_type = "Topic", operation = "Describe" }, { resource_type = "Group", operation = "Read" }, ] ACCESS_ROLE_TOPIC_ADMIN = [ { resource_type = "Topic", operation = "All" }, ] } # Flatten (user, topic, role) permissions into individual ACL entries. acl_entries = merge([ for user_key, perms in local.user_permissions : { for e in flatten([ for p in perms : [ for op in try(local.role_operations[p.role], []) : { key = "${user_key}:${p.topic_name}:${op.resource_type}:${op.operation}" principal = "User:${local.user_clusters[user_key].user}" resource_type = op.resource_type resource_name = op.resource_type == "Group" ? "*" : p.topic_name operation = op.operation } ] ]) : e.key => e } ]...) } resource "random_password" "kafka_user" { for_each = var.create_users ? local.user_clusters : {} length = var.user_password_length special = false upper = true lower = true numeric = true lifecycle { ignore_changes = all } } resource "kafka_topic" "this" { for_each = local.topics name = each.value.name replication_factor = each.value.replication_factor partitions = each.value.partitions config = each.value.config lifecycle { prevent_destroy = true } } resource "kafka_user_scram_credential" "this" { for_each = var.create_users ? local.user_clusters : {} username = each.value.user scram_mechanism = try(each.value.cluster.sasl_mechanism, "SCRAM-SHA-512") scram_iterations = 4096 password = random_password.kafka_user[each.key].result lifecycle { prevent_destroy = true ignore_changes = [password] } depends_on = [kafka_topic.this] } resource "kafka_acl" "this" { for_each = var.create_users ? local.acl_entries : {} resource_name = each.value.resource_name resource_type = each.value.resource_type resource_pattern_type_filter = "Literal" acl_principal = each.value.principal acl_host = "*" acl_operation = each.value.operation acl_permission_type = "Allow" depends_on = [kafka_user_scram_credential.this] }