-
Notifications
You must be signed in to change notification settings - Fork 1.3k
Commit
This commit does not belong to any branch on this repository, and may belong to a fork outside of the repository.
Added general global bookmark provider for all KafkaMessageReceivedOv…
…erride Triggers (#3807) Co-authored-by: Yannick Laubscher <[email protected]>
- Loading branch information
1 parent
580fef3
commit 31575ac
Showing
2 changed files
with
53 additions
and
0 deletions.
There are no files selected for viewing
52 changes: 52 additions & 0 deletions
52
src/activities/Elsa.Activities.Kafka/Bookmarks/OverrideKafkaBookmarkProvider.cs
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,52 @@ | ||
using Elsa.Activities.Kafka.Activities.KafkaMessageReceived; | ||
using Elsa.Activities.Kafka.Configuration; | ||
using Elsa.Services; | ||
using System; | ||
using System.Collections.Generic; | ||
using System.Linq; | ||
using System.Text; | ||
using System.Threading; | ||
using System.Threading.Tasks; | ||
|
||
namespace Elsa.Activities.Kafka.Bookmarks | ||
{ | ||
/// <summary> | ||
/// General Bookmark provider for all KafkaMessageReceived override activities. | ||
/// </summary> | ||
internal class OverrideKafkaBookmarkProvider : BookmarkProvider<Elsa.Activities.Kafka.Bookmarks.MessageReceivedBookmark, KafkaMessageReceived> | ||
{ | ||
|
||
private readonly IKafkaCustomActivityProvider _customActivityProvider; | ||
public OverrideKafkaBookmarkProvider(IKafkaCustomActivityProvider customActivityProvider) | ||
{ | ||
_customActivityProvider = customActivityProvider; | ||
} | ||
|
||
public override bool SupportsActivity(BookmarkProviderContext<KafkaMessageReceived> context) | ||
{ | ||
if (_customActivityProvider != null && | ||
_customActivityProvider.KafkaOverrideTriggers != null | ||
&& _customActivityProvider.KafkaOverrideTriggers.Contains(context.ActivityExecutionContext.ActivityBlueprint.Type)) | ||
{ | ||
return true; | ||
} | ||
else | ||
{ | ||
return base.SupportsActivity(context); | ||
} | ||
} | ||
|
||
public override async ValueTask<IEnumerable<BookmarkResult>> GetBookmarksAsync(BookmarkProviderContext<KafkaMessageReceived> context, CancellationToken cancellationToken) => | ||
new[] | ||
{ | ||
Result( | ||
new Elsa.Activities.Kafka.Bookmarks.MessageReceivedBookmark( | ||
topic: (await context.ReadActivityPropertyAsync(x => x.Topic, cancellationToken))!, | ||
group: (await context.ReadActivityPropertyAsync(x => x.Group, cancellationToken))!, | ||
connectionString: (await context.ReadActivityPropertyAsync(x => x.ConnectionString, cancellationToken))!, | ||
headers: (await context.ReadActivityPropertyAsync(x => x.Headers, cancellationToken) ?? new Dictionary<string, string>())!, | ||
autoOffsetReset: Enum.Parse<Confluent.Kafka.AutoOffsetReset>(await context.ReadActivityPropertyAsync(x => x.AutoOffsetReset, cancellationToken) ?? ((int)Confluent.Kafka.AutoOffsetReset.Earliest).ToString())! | ||
)) | ||
}; | ||
} | ||
} |
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters