From 9c072566ce3094219a5241adfa24d17f7ee4274e Mon Sep 17 00:00:00 2001 From: Raymond den Haan Date: Thu, 4 Jul 2024 15:47:34 +0200 Subject: [PATCH] Add workflow execution log sink Introduced an interface `IWorkflowExecutionLogSink` for storing workflow execution log records and implemented it as `StoreWorkflowExecutionLogSink`. This centralizes and streamlines the handling of persisting execution logs into the workflow's runtime, thus making the process more modular and maintainable. The `PersistWorkflowExecutionLogMiddleware` has also been refactored to use this sink, replacing the direct usage of the workflow execution log store. --- .../Activities/BulkDispatchWorkflows.cs | 6 +-- .../Activities/DispatchWorkflow.cs | 3 +- .../Bookmarks/EventBookmarkPayload.cs | 9 ++++ .../Contracts/IWorkflowExecutionLogSink.cs | 12 +++++ .../Features/WorkflowRuntimeFeature.cs | 7 +++ .../PersistWorkflowExecutionLogMiddleware.cs | 51 +------------------ .../Services/StoreWorkflowExecutionLogSink.cs | 43 ++++++++++++++++ 7 files changed, 76 insertions(+), 55 deletions(-) create mode 100644 src/modules/Elsa.Workflows.Runtime/Contracts/IWorkflowExecutionLogSink.cs create mode 100644 src/modules/Elsa.Workflows.Runtime/Services/StoreWorkflowExecutionLogSink.cs diff --git a/src/modules/Elsa.Workflows.Runtime/Activities/BulkDispatchWorkflows.cs b/src/modules/Elsa.Workflows.Runtime/Activities/BulkDispatchWorkflows.cs index 3db53171e..b04f1b296 100644 --- a/src/modules/Elsa.Workflows.Runtime/Activities/BulkDispatchWorkflows.cs +++ b/src/modules/Elsa.Workflows.Runtime/Activities/BulkDispatchWorkflows.cs @@ -1,3 +1,4 @@ +using System.Diagnostics.CodeAnalysis; using System.Runtime.CompilerServices; using Elsa.Common.Models; using Elsa.Expressions.Contracts; @@ -7,8 +8,6 @@ using Elsa.Extensions; using Elsa.Workflows.Activities.Flowchart.Attributes; using Elsa.Workflows.Attributes; using Elsa.Workflows.Contracts; -using Elsa.Workflows.Exceptions; -using Elsa.Workflows.UIHints; using Elsa.Workflows.Memory; using Elsa.Workflows.Models; using Elsa.Workflows.Options; @@ -18,8 +17,8 @@ using Elsa.Workflows.Runtime.Models; using Elsa.Workflows.Runtime.Requests; using Elsa.Workflows.Runtime.UIHints; using Elsa.Workflows.Services; +using Elsa.Workflows.UIHints; using JetBrains.Annotations; -using System.Diagnostics.CodeAnalysis; namespace Elsa.Workflows.Runtime.Activities; @@ -232,7 +231,6 @@ public class BulkDispatchWorkflows : Activity var input = context.WorkflowInput; var workflowInstanceId = input["WorkflowInstanceId"].ConvertTo()!; var workflowSubStatus = input["WorkflowSubStatus"].ConvertTo(); - var workflowOutput = input["WorkflowOutput"].ConvertTo>(); var finishedInstancesCount = context.GetProperty(CompletedInstancesCountKey) + 1; context.SetProperty(CompletedInstancesCountKey, finishedInstancesCount); diff --git a/src/modules/Elsa.Workflows.Runtime/Activities/DispatchWorkflow.cs b/src/modules/Elsa.Workflows.Runtime/Activities/DispatchWorkflow.cs index 5494f28da..7edd36f90 100644 --- a/src/modules/Elsa.Workflows.Runtime/Activities/DispatchWorkflow.cs +++ b/src/modules/Elsa.Workflows.Runtime/Activities/DispatchWorkflow.cs @@ -3,14 +3,13 @@ using Elsa.Common.Models; using Elsa.Extensions; using Elsa.Workflows.Attributes; using Elsa.Workflows.Contracts; -using Elsa.Workflows.Exceptions; -using Elsa.Workflows.UIHints; using Elsa.Workflows.Models; using Elsa.Workflows.Runtime.Bookmarks; using Elsa.Workflows.Runtime.Contracts; using Elsa.Workflows.Runtime.Models; using Elsa.Workflows.Runtime.Requests; using Elsa.Workflows.Runtime.UIHints; +using Elsa.Workflows.UIHints; using JetBrains.Annotations; namespace Elsa.Workflows.Runtime.Activities; diff --git a/src/modules/Elsa.Workflows.Runtime/Bookmarks/EventBookmarkPayload.cs b/src/modules/Elsa.Workflows.Runtime/Bookmarks/EventBookmarkPayload.cs index 1c032cd1a..d90314e8b 100644 --- a/src/modules/Elsa.Workflows.Runtime/Bookmarks/EventBookmarkPayload.cs +++ b/src/modules/Elsa.Workflows.Runtime/Bookmarks/EventBookmarkPayload.cs @@ -1,14 +1,23 @@ namespace Elsa.Workflows.Runtime.Bookmarks; +/// +/// The payload for an event bookmark. +/// public class EventBookmarkPayload { private readonly string _eventName = default!; + /// + /// The payload for an event bookmark. + /// public EventBookmarkPayload(string eventName) { EventName = eventName; } + /// + /// The name of the event. + /// public string EventName { get => _eventName; diff --git a/src/modules/Elsa.Workflows.Runtime/Contracts/IWorkflowExecutionLogSink.cs b/src/modules/Elsa.Workflows.Runtime/Contracts/IWorkflowExecutionLogSink.cs new file mode 100644 index 000000000..2c0601c79 --- /dev/null +++ b/src/modules/Elsa.Workflows.Runtime/Contracts/IWorkflowExecutionLogSink.cs @@ -0,0 +1,12 @@ +namespace Elsa.Workflows.Runtime.Contracts; + +/// +/// Represents a sink for storing workflow execution log records. +/// +public interface IWorkflowExecutionLogSink +{ + /// + /// Persists the execution logs of a workflow. + /// + Task PersistExecutionLogsAsync(WorkflowExecutionContext context, CancellationToken cancellationToken = default); +} \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Runtime/Features/WorkflowRuntimeFeature.cs b/src/modules/Elsa.Workflows.Runtime/Features/WorkflowRuntimeFeature.cs index d03639667..827c3920f 100644 --- a/src/modules/Elsa.Workflows.Runtime/Features/WorkflowRuntimeFeature.cs +++ b/src/modules/Elsa.Workflows.Runtime/Features/WorkflowRuntimeFeature.cs @@ -107,6 +107,11 @@ public class WorkflowRuntimeFeature : FeatureBase /// A factory that instantiates an . /// public Func BackgroundActivityScheduler { get; set; } = sp => ActivatorUtilities.CreateInstance(sp); + + /// + /// Represents a sink for workflow execution logs. + /// + public Func WorkflowExecutionLogSink { get; set; } = sp => sp.GetRequiredService(); /// /// A delegate to configure the . @@ -206,6 +211,7 @@ public class WorkflowRuntimeFeature : FeatureBase .AddScoped(WorkflowCancellationDispatcher) .AddScoped(WorkflowExecutionContextStore) .AddScoped(RunTaskDispatcher) + .AddScoped(WorkflowExecutionLogSink) .AddSingleton(BackgroundActivityScheduler) .AddSingleton() .AddScoped() @@ -218,6 +224,7 @@ public class WorkflowRuntimeFeature : FeatureBase .AddScoped() .AddScoped() .AddScoped() + .AddScoped() .AddScoped() .AddScoped() .AddScoped() diff --git a/src/modules/Elsa.Workflows.Runtime/Middleware/Workflows/PersistWorkflowExecutionLogMiddleware.cs b/src/modules/Elsa.Workflows.Runtime/Middleware/Workflows/PersistWorkflowExecutionLogMiddleware.cs index 642e4155b..36ba5f1e6 100644 --- a/src/modules/Elsa.Workflows.Runtime/Middleware/Workflows/PersistWorkflowExecutionLogMiddleware.cs +++ b/src/modules/Elsa.Workflows.Runtime/Middleware/Workflows/PersistWorkflowExecutionLogMiddleware.cs @@ -1,66 +1,19 @@ -using Elsa.Mediator.Contracts; -using Elsa.Workflows.Contracts; using Elsa.Workflows.Pipelines.WorkflowExecution; using Elsa.Workflows.Runtime.Contracts; -using Elsa.Workflows.Runtime.Entities; -using Elsa.Workflows.Runtime.Notifications; namespace Elsa.Workflows.Runtime.Middleware.Workflows; /// /// Takes care of persisting workflow execution log entries. /// -public class PersistWorkflowExecutionLogMiddleware : WorkflowExecutionMiddleware +public class PersistWorkflowExecutionLogMiddleware(WorkflowMiddlewareDelegate next, IWorkflowExecutionLogSink sink) : WorkflowExecutionMiddleware(next) { - private readonly IWorkflowExecutionLogStore _workflowExecutionLogStore; - private readonly INotificationSender _notificationSender; - private readonly IIdentityGenerator _identityGenerator; - - /// - public PersistWorkflowExecutionLogMiddleware( - WorkflowMiddlewareDelegate next, - IWorkflowExecutionLogStore workflowExecutionLogStore, - INotificationSender notificationSender, - IIdentityGenerator identityGenerator) : base(next) - { - _workflowExecutionLogStore = workflowExecutionLogStore; - _notificationSender = notificationSender; - _identityGenerator = identityGenerator; - } - /// public override async ValueTask InvokeAsync(WorkflowExecutionContext context) { // Invoke next middleware. await Next(context); - // Persist workflow execution log entries. - var entries = context.ExecutionLog.Select(x => new WorkflowExecutionLogRecord - { - Id = _identityGenerator.GenerateId(), - ActivityInstanceId = x.ActivityInstanceId, - ParentActivityInstanceId = x.ParentActivityInstanceId, - ActivityNodeId = x.NodeId, - ActivityId = x.ActivityId, - ActivityType = x.ActivityType, - ActivityTypeVersion = x.ActivityTypeVersion, - ActivityName = x.ActivityName, - Message = x.Message, - EventName = x.EventName, - WorkflowDefinitionId = context.Workflow.Identity.DefinitionId, - WorkflowDefinitionVersionId = context.Workflow.Identity.Id, - WorkflowInstanceId = context.Id, - WorkflowVersion = context.Workflow.Version, - Source = x.Source, - ActivityState = x.ActivityState, - Payload = x.Payload, - Timestamp = x.Timestamp, - Sequence = x.Sequence - }).ToList(); - - await _workflowExecutionLogStore.AddManyAsync(entries, context.CancellationTokens.SystemCancellationToken); - - // Publish notification. - await _notificationSender.SendAsync(new WorkflowExecutionLogUpdated(context), context.CancellationTokens.SystemCancellationToken); + await sink.PersistExecutionLogsAsync(context); } } \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Runtime/Services/StoreWorkflowExecutionLogSink.cs b/src/modules/Elsa.Workflows.Runtime/Services/StoreWorkflowExecutionLogSink.cs new file mode 100644 index 000000000..ecff347d0 --- /dev/null +++ b/src/modules/Elsa.Workflows.Runtime/Services/StoreWorkflowExecutionLogSink.cs @@ -0,0 +1,43 @@ +using Elsa.Mediator.Contracts; +using Elsa.Workflows.Contracts; +using Elsa.Workflows.Runtime.Contracts; +using Elsa.Workflows.Runtime.Entities; +using Elsa.Workflows.Runtime.Notifications; + +namespace Elsa.Workflows.Runtime.Services; + +/// +/// This implementation saves directly through the store. +/// +public class StoreWorkflowExecutionLogSink(IWorkflowExecutionLogStore store, IIdentityGenerator identityGenerator, INotificationSender notificationSender) : IWorkflowExecutionLogSink +{ + /// + public async Task PersistExecutionLogsAsync(WorkflowExecutionContext context, CancellationToken cancellationToken) + { + var records = context.ExecutionLog.Select(x => new WorkflowExecutionLogRecord + { + Id = identityGenerator.GenerateId(), + ActivityInstanceId = x.ActivityInstanceId, + ParentActivityInstanceId = x.ParentActivityInstanceId, + ActivityNodeId = x.NodeId, + ActivityId = x.ActivityId, + ActivityType = x.ActivityType, + ActivityTypeVersion = x.ActivityTypeVersion, + ActivityName = x.ActivityName, + Message = x.Message, + EventName = x.EventName, + WorkflowDefinitionId = context.Workflow.Identity.DefinitionId, + WorkflowDefinitionVersionId = context.Workflow.Identity.Id, + WorkflowInstanceId = context.Id, + WorkflowVersion = context.Workflow.Version, + Source = x.Source, + ActivityState = x.ActivityState, + Payload = x.Payload, + Timestamp = x.Timestamp, + Sequence = x.Sequence + }).ToList(); + + await store.AddManyAsync(records, context.CancellationTokens.SystemCancellationToken); + await notificationSender.SendAsync(new WorkflowExecutionLogUpdated(context), context.CancellationTokens.SystemCancellationToken); + } +} \ No newline at end of file