// ========================================================================== // Squidex Headless CMS // ========================================================================== // Copyright (c) Squidex UG (haftungsbeschraenkt) // All rights reserved. Licensed under the MIT license. // ========================================================================== #if INCLUDE_KAFKA using System.ComponentModel.DataAnnotations; using Squidex.Domain.Apps.Core.HandleRules; using Squidex.Domain.Apps.Core.Rules.Deprecated; using Squidex.Flows; using Squidex.Infrastructure.Reflection; using Squidex.Infrastructure.Validation; namespace Squidex.Extensions.Actions.Kafka; [FlowStep( Title = "Kafka", IconImage = "", IconColor = "#404244", Display = "Push to kafka", Description = "Connect to Kafka stream and push data to that stream.", ReadMore = "https://kafka.apache.org/quickstart")] #pragma warning disable CS0618 // Type or member is obsolete public sealed record KafkaFlowStep : FlowStep, IConvertibleToAction #pragma warning restore CS0618 // Type or member is obsolete { [LocalizedRequired] [Display(Name = "Topic Name", Description = "The name of the topic.")] [Editor(FlowStepEditor.Text)] [Expression] public string TopicName { get; set; } [Display(Name = "Payload (Optional)", Description = "Leave it empty to use the full event as body.")] [Editor(FlowStepEditor.TextArea)] [Expression(ExpressionFallback.Envelope)] public string? Payload { get; set; } [Display(Name = "Key", Description = "The message key, commonly used for partitioning.")] [Editor(FlowStepEditor.Text)] [Expression] public string? Key { get; set; } [Display(Name = "Partition Key", Description = "The partition key, only used when we don't want to define partiontionig with key.")] [Editor(FlowStepEditor.Text)] [Expression] public string? PartitionKey { get; set; } [Display(Name = "Partition Count", Description = "Define the number of partitions for specific topic.")] [Editor(FlowStepEditor.Text)] public int PartitionCount { get; set; } [Display(Name = "Headers (Optional)", Description = "The message headers in the format '[Key]=[Value]', one entry per line.")] [Editor(FlowStepEditor.TextArea)] [Expression] public string? Headers { get; set; } [Display(Name = "Schema (Optional)", Description = "Define a specific AVRO schema in JSON format.")] [Editor(FlowStepEditor.TextArea)] public string? Schema { get; set; } public override async ValueTask ExecuteAsync(FlowExecutionContext executionContext, CancellationToken ct) { if (executionContext.IsSimulation) { executionContext.LogSkipSimulation(); return Next(); } var @event = ((FlowEventContext)executionContext.Context).Event; var key = Key; if (string.IsNullOrWhiteSpace(key)) { key = @event.Name; } try { var request = new KafkaMessageRequest { Headers = ParseHeaders(Headers), MessageKey = key, MessageValue = Payload, PartitionCount = PartitionCount, PartitionKey = PartitionKey, Schema = Schema, TopicName = TopicName, }; await executionContext.Resolve() .SendAsync(request, ct); executionContext.Log($"Event pushed to {TopicName} kafka topic with '{key}' message key."); return Next(); } catch (Exception ex) { executionContext.Log("Push to kafka failed", ex.Message); throw; } } private static Dictionary? ParseHeaders(string? headers) { if (string.IsNullOrWhiteSpace(headers)) { return null; } var headersDictionary = new Dictionary(); foreach (var line in headers.Split('\n')) { var indexEqual = line.IndexOf('=', StringComparison.Ordinal); if (indexEqual > 0 && indexEqual < line.Length - 1) { var headerKey = line[..indexEqual]; var headerValue = line[(indexEqual + 1)..]; headersDictionary[headerKey] = headerValue!; } } return headersDictionary; } #pragma warning disable CS0618 // Type or member is obsolete public RuleAction ToAction() { return SimpleMapper.Map(this, new KafkaAction()); } #pragma warning restore CS0618 // Type or member is obsolete } #endif