Skip to content

fix: Kafka topic auto-creation fails every send when the principal lacks CREATE in NetEvolve.Pulse.Kafka #845

Description

@samtrion

User Story

As a developer running the Kafka outbox transport against a locked-down or managed cluster with pre-provisioned topics, I want a missing topic-creation permission to not block delivery, so that messages reach topics that already exist and the producer is allowed to write to.


Problem

KafkaTransportOptions.AutoCreateTopics defaults to true (src/NetEvolve.Pulse.Kafka/Outbox/KafkaTransportOptions.cs:24). Before every send, KafkaMessageTransport.EnsureTopicAsync calls IAdminClient.CreateTopicsAsync (src/NetEvolve.Pulse.Kafka/Outbox/KafkaMessageTransport.cs:170). The only error it catches is TopicAlreadyExists (:173):

catch (CreateTopicsException ex) when (ex.Results.All(static r => r.Error.Code == ErrorCode.TopicAlreadyExists))

A topic is added to _ensuredTopics only on success or on TopicAlreadyExists. Every other CreateTopicsException escapes and nothing is cached, so the next message tries to create the topic again. Both SendAsync (:62) and SendBatchAsync (:85) call EnsureTopicAsync before ProduceAsync/Produce, so every send fails before the producer runs.

Failure scenario:

  1. A topic orders is pre-provisioned. The producer's principal has WRITE and DESCRIBE on it, but not CREATE. This is common for Confluent Cloud service accounts and other managed clusters.
  2. UseKafkaTransport() is registered with default options.
  3. The broker answers CreateTopics with TOPIC_AUTHORIZATION_FAILED (or CLUSTER_AUTHORIZATION_FAILED), not TOPIC_ALREADY_EXISTS, because it checks authorization before it checks whether the topic exists.
  4. EnsureTopicAsync rethrows on every call and no message is ever produced. Every outbox message is retried until it ends up dead-lettered.

The workaround is UseKafkaTransport(o => o.AutoCreateTopics = false), but it is not documented. The README (src/NetEvolve.Pulse.Kafka/README.md:18) says the admin client is "used for health checks". It never mentions AutoCreateTopics, the DefaultPartitionCount/DefaultReplicationFactor defaults of 1 (KafkaTransportOptions.cs:12, :18) or the need for a CREATE ACL.

No test covers CreateTopicsException with an error code other than TopicAlreadyExists.


Specification

Apache Kafka broker, CreateTopics handling (core/src/main/scala/kafka/server/ControllerApis.scala, createTopics, https://github.com/apache/kafka/blob/trunk/core/src/main/scala/kafka/server/ControllerApis.scala). If the principal lacks CREATE on the cluster, each topic name is checked for CREATE on the topic. Unauthorized names are answered with:

setErrorCode(TOPIC_AUTHORIZATION_FAILED.code).setErrorMessage("Authorization failed.")

These names are removed from the request before controller.createTopics() runs, so the existence check never happens for them. An existing topic therefore still returns TOPIC_AUTHORIZATION_FAILED.

Kafka Connect hit the same failure mode, where only TopicAlreadyExists was tolerated: KAFKA-6250 "Kafka Connect requires permission to create internal topics even if they exist".

Confluent.Kafka IAdminClient.CreateTopicsAsync throws CreateTopicsException with one CreateTopicReport per topic, each carrying its own Error.Code (https://docs.confluent.io/platform/current/clients/confluent-kafka-dotnet/_site/api/Confluent.Kafka.Admin.CreateTopicsException.html).


Requirements

  • In EnsureTopicAsync, handle TopicAuthorizationFailed and ClusterAuthorizationFailed without failing the send. Check whether the topic exists (for example via IAdminClient.GetMetadata(topic, timeout)). If it exists, treat it as ensured and cache it. If existence cannot be confirmed, log a warning, cache the attempt and let the produce call decide.
  • Do not retry CreateTopics for every message after a non-transient failure.
  • Fix this in EnsureTopicAsync so that both SendAsync and SendBatchAsync are covered.
  • Other error codes (for example PolicyViolation, InvalidReplicationFactor for a topic that does not exist yet) may keep failing, but the log or exception must name the topic and the option (AutoCreateTopics) to change.
  • Document AutoCreateTopics, the DefaultPartitionCount/DefaultReplicationFactor defaults of 1 and the CREATE ACL requirement in src/NetEvolve.Pulse.Kafka/README.md. Correct the admin client comment, which currently says the client is only used for health checks.

Acceptance Criteria

  • A failing unit test first: CreateTopicsAsync throws CreateTopicsException with TopicAuthorizationFailed for an existing topic, and SendAsync still produces the message.
  • The same scenario is covered for SendBatchAsync.
  • After an authorization failure, CreateTopicsAsync is called at most once per topic, not once per message.
  • TopicAlreadyExists behavior is unchanged.
  • The README documents AutoCreateTopics, its partition and replication defaults and the required ACL, and shows how to disable auto-creation.

Activity

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Metadata

Metadata

Assignees

No one assigned

    Labels

    type:bugIndicates an issue or flaw that needs to be fixed.

    Type

    No type

    Projects

    No projects

      Milestone

      No milestone

      Relationships

      None yet

      Development

      No branches or pull requests

      Issue actions