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:
- 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.
UseKafkaTransport() is registered with default options.
- 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.
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
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.AutoCreateTopicsdefaults totrue(src/NetEvolve.Pulse.Kafka/Outbox/KafkaTransportOptions.cs:24). Before every send,KafkaMessageTransport.EnsureTopicAsynccallsIAdminClient.CreateTopicsAsync(src/NetEvolve.Pulse.Kafka/Outbox/KafkaMessageTransport.cs:170). The only error it catches isTopicAlreadyExists(:173):A topic is added to
_ensuredTopicsonly on success or onTopicAlreadyExists. Every otherCreateTopicsExceptionescapes and nothing is cached, so the next message tries to create the topic again. BothSendAsync(:62) andSendBatchAsync(:85) callEnsureTopicAsyncbeforeProduceAsync/Produce, so every send fails before the producer runs.Failure scenario:
ordersis pre-provisioned. The producer's principal hasWRITEandDESCRIBEon it, but notCREATE. This is common for Confluent Cloud service accounts and other managed clusters.UseKafkaTransport()is registered with default options.CreateTopicswithTOPIC_AUTHORIZATION_FAILED(orCLUSTER_AUTHORIZATION_FAILED), notTOPIC_ALREADY_EXISTS, because it checks authorization before it checks whether the topic exists.EnsureTopicAsyncrethrows 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 mentionsAutoCreateTopics, theDefaultPartitionCount/DefaultReplicationFactordefaults of1(KafkaTransportOptions.cs:12,:18) or the need for aCREATEACL.No test covers
CreateTopicsExceptionwith an error code other thanTopicAlreadyExists.Specification
Apache Kafka broker,
CreateTopicshandling (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 lacksCREATEon the cluster, each topic name is checked forCREATEon the topic. Unauthorized names are answered with:These names are removed from the request before
controller.createTopics()runs, so the existence check never happens for them. An existing topic therefore still returnsTOPIC_AUTHORIZATION_FAILED.Kafka Connect hit the same failure mode, where only
TopicAlreadyExistswas tolerated: KAFKA-6250 "Kafka Connect requires permission to create internal topics even if they exist".Confluent.Kafka
IAdminClient.CreateTopicsAsyncthrowsCreateTopicsExceptionwith oneCreateTopicReportper topic, each carrying its ownError.Code(https://docs.confluent.io/platform/current/clients/confluent-kafka-dotnet/_site/api/Confluent.Kafka.Admin.CreateTopicsException.html).Requirements
EnsureTopicAsync, handleTopicAuthorizationFailedandClusterAuthorizationFailedwithout failing the send. Check whether the topic exists (for example viaIAdminClient.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.CreateTopicsfor every message after a non-transient failure.EnsureTopicAsyncso that bothSendAsyncandSendBatchAsyncare covered.PolicyViolation,InvalidReplicationFactorfor 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.AutoCreateTopics, theDefaultPartitionCount/DefaultReplicationFactordefaults of1and theCREATEACL requirement insrc/NetEvolve.Pulse.Kafka/README.md. Correct the admin client comment, which currently says the client is only used for health checks.Acceptance Criteria
CreateTopicsAsyncthrowsCreateTopicsExceptionwithTopicAuthorizationFailedfor an existing topic, andSendAsyncstill produces the message.SendBatchAsync.CreateTopicsAsyncis called at most once per topic, not once per message.TopicAlreadyExistsbehavior is unchanged.AutoCreateTopics, its partition and replication defaults and the required ACL, and shows how to disable auto-creation.