data "aws_caller_identity" "current" {}

data "aws_vpc" "vpc" {
  id = var.vpc_id
}

data "aws_ec2_managed_prefix_lists" "private_subnets" {
  tags = {
    environment  = var.environment
    subnet_group = "all"
    tier         = "private"
    vpc_id       = var.vpc_id
  }
}


data "aws_ec2_managed_prefix_list" "private_subnets" {
  id = one(data.aws_ec2_managed_prefix_lists.private_subnets.ids)
}

data "aws_ec2_managed_prefix_lists" "kafka_allowed_prefix_list" {
  filter {
    name   = "prefix-list-name"
    values = local.kafka_allowed_prefix_list_names
  }
}

data "aws_ec2_managed_prefix_lists" "kafka_allowed_custom_prefix_list" {
  filter {
    name   = "prefix-list-name"
    values = var.kafka_allowed_custom_prefix_list_names
  }
}

resource "aws_security_group" "kafka_sg" {
  name        = "${var.environment}-${var.cluster_name}-kafka-security-group"
  description = "Allow Kafka ports"
  vpc_id      = data.aws_vpc.vpc.id
  tags        = merge(local.combined_resource_tags, local.primsa_sg_exception_tags)

  egress {
    from_port = 0
    to_port   = 0
    protocol  = "-1"
    cidr_blocks = [
      "0.0.0.0/0"
    ]
  }
}

resource "aws_security_group_rule" "allow_all_self" {
  from_port         = 0
  to_port           = 0
  protocol          = "tcp"
  security_group_id = aws_security_group.kafka_sg.id
  self              = true
  type              = "ingress"
}

resource "aws_security_group_rule" "kafka_9094" {
  from_port         = 9094
  to_port           = 9094
  protocol          = "tcp"
  security_group_id = aws_security_group.kafka_sg.id
  cidr_blocks       = local.access_cidr_blocks[var.environment]
  prefix_list_ids = concat(
    data.aws_ec2_managed_prefix_lists.kafka_allowed_prefix_list.ids,
    data.aws_ec2_managed_prefix_lists.kafka_allowed_custom_prefix_list.ids,
  )
  type = "ingress"
}

resource "aws_security_group_rule" "kafka_2181" {
  from_port         = 2181
  to_port           = 2181
  protocol          = "tcp"
  security_group_id = aws_security_group.kafka_sg.id
  cidr_blocks       = local.access_cidr_blocks[var.environment]
  prefix_list_ids   = data.aws_ec2_managed_prefix_lists.kafka_allowed_prefix_list.ids
  type              = "ingress"
}

resource "aws_security_group_rule" "kafka_zookeeper_2182" {
  from_port         = 2182
  to_port           = 2182
  protocol          = "tcp"
  security_group_id = aws_security_group.kafka_sg.id
  cidr_blocks       = local.access_cidr_blocks[var.environment]
  prefix_list_ids   = data.aws_ec2_managed_prefix_lists.kafka_allowed_prefix_list.ids
  type              = "ingress"
}

resource "aws_security_group_rule" "kafka_ingress_security_group_rule" {
  for_each = { for idx, val in local.sg_port_rules : "${idx}-${val.source_sg_id}-${val.from_port}-${val.to_port}" => val }

  from_port                = each.value.from_port
  protocol                 = each.value.protocol
  security_group_id        = aws_security_group.kafka_sg.id
  source_security_group_id = each.value.source_sg_id
  to_port                  = each.value.to_port
  type                     = "ingress"
}

resource "aws_cloudwatch_log_group" "kafka_log_group" {
  name              = "${var.environment}-${var.cluster_name}"
  retention_in_days = var.log_retention_in_days
  tags              = local.combined_resource_tags
}

resource "aws_cloudwatch_log_subscription_filter" "datadog_lambda_function_log_filter" {
  count           = var.datadog_enabled ? 1 : 0
  name            = "${var.environment}-${var.cluster_name}-kafka-logs-subscription-filter"
  log_group_name  = aws_cloudwatch_log_group.kafka_log_group.name
  filter_pattern  = ""
  destination_arn = local.datadog_function_destination_arn
  distribution    = "ByLogStream"
}

resource "aws_kms_key" "kafka_kms_key" {
  description             = "${var.environment}-${var.cluster_name}-kafka-key"
  enable_key_rotation     = true
  deletion_window_in_days = var.kms_key_deletion_window
  tags                    = local.combined_resource_tags
}

resource "aws_msk_configuration" "managed_kafka_config" {
  kafka_versions    = [var.kafka_version]
  name              = "${var.environment}-${var.cluster_name}-${replace(var.kafka_version, ".", "")}-configuration"
  server_properties = file(var.kafka_properties)
  lifecycle {
    create_before_destroy = true
  }
}

resource "aws_msk_cluster" "managed_kafka" {
  cluster_name           = "${var.environment}-${var.cluster_name}"
  kafka_version          = var.kafka_version
  number_of_broker_nodes = var.kafka_num_nodes
  enhanced_monitoring    = var.cloudwatch_enhanced_monitoring

  broker_node_group_info {
    instance_type = var.kafka_instance_type

    storage_info {
      ebs_storage_info {
        volume_size = var.kafka_ebs_volume_size
      }
    }

    client_subnets  = var.subnet_ids
    security_groups = [aws_security_group.kafka_sg.id]
  }

  encryption_info {
    encryption_at_rest_kms_key_arn = aws_kms_key.kafka_kms_key.arn
    encryption_in_transit {
      in_cluster    = var.in_cluster_encryption
      client_broker = var.in_transit_encryption
    }
  }

  configuration_info {
    arn      = aws_msk_configuration.managed_kafka_config.arn
    revision = aws_msk_configuration.managed_kafka_config.latest_revision
  }

  logging_info {
    broker_logs {
      cloudwatch_logs {
        enabled   = true
        log_group = aws_cloudwatch_log_group.kafka_log_group.name
      }
    }
  }

  tags = local.combined_resource_tags

  # Allow storage scaling to do its thing without being overwritten.
  # This effectively relegates var.ebs_volume_size to only being utilized on cluster creation, but this is necessary to seamlessly scale.
  lifecycle {
    ignore_changes = [
      broker_node_group_info[0].storage_info[0].ebs_storage_info,
    ]
  }
}

# The min is not actually 1; the AWS APIs require it to be this value. It never scales down so this value is moot.
resource "aws_appautoscaling_target" "kafka_storage_scaling_target" {
  max_capacity       = var.kafka_storage_scaling_max_capacity
  min_capacity       = 1
  resource_id        = aws_msk_cluster.managed_kafka.arn
  scalable_dimension = "kafka:broker-storage:VolumeSize"
  service_namespace  = "kafka"
}

resource "aws_appautoscaling_policy" "kafka_storage_scaling_policy" {
  name               = "${var.environment}-${var.cluster_name}-storage-scaling"
  policy_type        = "TargetTrackingScaling"
  resource_id        = aws_msk_cluster.managed_kafka.arn
  scalable_dimension = aws_appautoscaling_target.kafka_storage_scaling_target.scalable_dimension
  service_namespace  = aws_appautoscaling_target.kafka_storage_scaling_target.service_namespace

  target_tracking_scaling_policy_configuration {
    predefined_metric_specification {
      predefined_metric_type = "KafkaBrokerStorageUtilization"
    }
    target_value = var.kafka_storage_scaling_target_value
  }
}
