// ==========================================================================
// 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