# Generic event source mapping for Kafka. Used with MSK clusters in a different AWS account or self-managed Kafka.
resource "aws_lambda_event_source_mapping" "kafka_event_source_mapping" {
  count = var.kafka_event_enabled ? length(var.kafka_topics) : 0

  batch_size                         = var.event_source_mapping_batch_size
  function_name                      = local.lambda_function_arn
  topics                             = [var.kafka_topics[count.index]]
  starting_position                  = var.event_source_mapping_starting_position
  maximum_batching_window_in_seconds = var.event_source_mapping_batch_window

  self_managed_event_source {
    endpoints = {
      KAFKA_BOOTSTRAP_SERVERS = join(",", var.kafka_bootstrap_servers)
    }
  }

  dynamic "destination_config" {
    for_each = var.event_source_mapping_failure_arn != null ? [""] : []
    content {
      on_failure {
        destination_arn = var.event_source_mapping_failure_arn
      }
    }
  }

  dynamic "self_managed_kafka_event_source_config" {
    for_each = var.kafka_use_lambda_name_as_consumer_id ? [var.lambda_name] : []
    content {
      consumer_group_id = "${var.lambda_name}-${var.kafka_topics[count.index]}"
    }
  }

  dynamic "filter_criteria" {
    for_each = var.event_source_mapping_filter_criteria_pattern != "" ? [""] : []
    content {
      filter {
        pattern = var.event_source_mapping_filter_criteria_pattern
      }
    }
  }

  dynamic "source_access_configuration" {
    for_each = var.vpc_enabled ? var.vpc_subnet_ids : []
    content {
      type = "VPC_SUBNET"
      uri  = "subnet:${source_access_configuration.value}"
    }
  }

  dynamic "source_access_configuration" {
    for_each = var.vpc_enabled ? (var.vpc_create_security_group ? aws_security_group.lambda_security_group[*].id : data.aws_security_group.default[*].id) : []
    content {
      type = "VPC_SECURITY_GROUP"
      uri  = "security_group:${source_access_configuration.value}"
    }
  }
}
