using Confluent.Kafka; using Confluent.Kafka.Admin; using Microsoft.Extensions.Options; public class KafkaTopicProvisioner : IHostedService { private readonly KafkaOptions _options; private readonly ILogger _logger; public KafkaTopicProvisioner(IOptions options, ILogger logger) { _options = options.Value; _logger = logger; } public async Task StartAsync(CancellationToken cancellationToken) { using var admin = new AdminClientBuilder(new AdminClientConfig { BootstrapServers = _options.BootstrapServers }).Build(); var topicNames = new[] { _options.Topics.ObservationRecorded, _options.Topics.AlertGenerated, _options.Topics.AlertAcknowledged, _options.Topics.EncounterStatusChanged, _options.Topics.SepsisBundleCreated, _options.Topics.SepsisBundleUpdated, _options.Topics.GcsScored }; var specs = topicNames.Select(name => new TopicSpecification { Name = name, NumPartitions = _options.NumPartitions, ReplicationFactor = 1 }).ToList(); try { await admin.CreateTopicsAsync(specs); _logger.LogInformation("Kafka topics provisioned: {Topics}", string.Join(", ", topicNames)); } catch (CreateTopicsException ex) { var errors = ex.Results .Where(r => r.Error.Code is not (ErrorCode.NoError or ErrorCode.TopicAlreadyExists)) .ToList(); if (errors.Count > 0) throw new InvalidOperationException( $"Failed to create Kafka topics: {string.Join(", ", errors.Select(e => e.Error.Reason))}"); _logger.LogInformation("Kafka topics already exist — skipping creation"); } } public Task StopAsync(CancellationToken cancellationToken) => Task.CompletedTask; }