From c413f210a9b203688bf8912254fe3cea028486ae Mon Sep 17 00:00:00 2001
From: raymonddenhaan <155616759+raymonddenhaan@users.noreply.github.com>
Date: Wed, 24 Jul 2024 22:30:14 +0200
Subject: [PATCH] Add activity execution log sink (#5833)
* Add activity execution log sink
Introduced an interface `IActivityExecutionLogSink` for storing workflow execution log records and implemented it as `StoreActivityExecutionLogSink`. This centralizes and streamlines the handling of persisting execution logs into the workflow's runtime, thus making the process more modular and maintainable. The `PersistActivityExecutionLogMiddleware` has also been refactored to use this sink, replacing the direct usage of the workflow execution log store.
* Refactor logging interfaces and implementations
Replaced specific logging interfaces and implementations with generic ones. The individual interfaces for workflow and activity execution logging have been replaced with a single ILogRecordExtractor and ILogRecordSink interface. Specific implementations have been adjusted to use these new interfaces. This allows for greater flexibility and reuse of logging code.
* Introduce ILogRecordStore interface
The `ILogRecordStore` interface has been introduced to provide a centralized place for handling log records. Both `IActivityExecutionStore` and `IWorkflowExecutionLogStore` have been updated to inherit from this new interface. As a result of this change, the `SaveManyAsync` methods in these two interfaces have been removed to avoid redundancy.
---
.../Elsa.Common/Contracts/ILogRecord.cs | 4 +++
.../Contracts/IActivityExecutionStore.cs | 12 +------
.../Contracts/ILogRecordExtractor.cs | 10 ++++++
.../Contracts/ILogRecordSink.cs | 10 ++++++
.../Contracts/ILogRecordStore.cs | 13 ++++++++
.../IWorkflowExecutionLogRecordExtractor.cs | 10 ------
.../Contracts/IWorkflowExecutionLogSink.cs | 12 -------
.../Contracts/IWorkflowExecutionLogStore.cs | 10 +-----
.../Entities/ActivityExecutionRecord.cs | 3 +-
.../Entities/WorkflowExecutionLogRecord.cs | 3 +-
.../Features/WorkflowRuntimeFeature.cs | 16 ++++++---
.../PersistActivityExecutionLogMiddleware.cs | 33 ++-----------------
.../PersistWorkflowExecutionLogMiddleware.cs | 4 +--
.../ActivityExecutionRecordExtractor.cs | 15 +++++++++
.../Services/StoreActivityExecutionLogSink.cs | 21 ++++++++++++
.../Services/StoreWorkflowExecutionLogSink.cs | 4 +--
...=> WorkflowExecutionLogRecordExtractor.cs} | 4 +--
17 files changed, 100 insertions(+), 84 deletions(-)
create mode 100644 src/modules/Elsa.Common/Contracts/ILogRecord.cs
create mode 100644 src/modules/Elsa.Workflows.Runtime/Contracts/ILogRecordExtractor.cs
create mode 100644 src/modules/Elsa.Workflows.Runtime/Contracts/ILogRecordSink.cs
create mode 100644 src/modules/Elsa.Workflows.Runtime/Contracts/ILogRecordStore.cs
delete mode 100644 src/modules/Elsa.Workflows.Runtime/Contracts/IWorkflowExecutionLogRecordExtractor.cs
delete mode 100644 src/modules/Elsa.Workflows.Runtime/Contracts/IWorkflowExecutionLogSink.cs
create mode 100644 src/modules/Elsa.Workflows.Runtime/Services/ActivityExecutionRecordExtractor.cs
create mode 100644 src/modules/Elsa.Workflows.Runtime/Services/StoreActivityExecutionLogSink.cs
rename src/modules/Elsa.Workflows.Runtime/Services/{WorkflowExecutionContext.cs => WorkflowExecutionLogRecordExtractor.cs} (87%)
diff --git a/src/modules/Elsa.Common/Contracts/ILogRecord.cs b/src/modules/Elsa.Common/Contracts/ILogRecord.cs
new file mode 100644
index 000000000..2af179591
--- /dev/null
+++ b/src/modules/Elsa.Common/Contracts/ILogRecord.cs
@@ -0,0 +1,4 @@
+namespace Elsa.Common.Contracts;
+
+/// Represents a log record.
+public interface ILogRecord;
\ No newline at end of file
diff --git a/src/modules/Elsa.Workflows.Runtime/Contracts/IActivityExecutionStore.cs b/src/modules/Elsa.Workflows.Runtime/Contracts/IActivityExecutionStore.cs
index 1cb89c04b..5b8bd01bf 100644
--- a/src/modules/Elsa.Workflows.Runtime/Contracts/IActivityExecutionStore.cs
+++ b/src/modules/Elsa.Workflows.Runtime/Contracts/IActivityExecutionStore.cs
@@ -7,7 +7,7 @@ namespace Elsa.Workflows.Runtime.Contracts;
///
/// Stores activity execution records.
///
-public interface IActivityExecutionStore
+public interface IActivityExecutionStore : ILogRecordStore
{
///
/// Adds or updates the specified in the persistence store.
@@ -19,16 +19,6 @@ public interface IActivityExecutionStore
///
Task SaveAsync(ActivityExecutionRecord record, CancellationToken cancellationToken = default);
- ///
- /// Adds or updates the specified set of objects in the persistence store.
- ///
- /// The activity execution records.
- /// An optional cancellation token.
- ///
- /// If a record does not already exist, it is added to the store; if it does exist, its existing entry is updated.
- ///
- Task SaveManyAsync(IEnumerable records, CancellationToken cancellationToken = default);
-
///
/// Finds an activity execution record matching the specified filter.
///
diff --git a/src/modules/Elsa.Workflows.Runtime/Contracts/ILogRecordExtractor.cs b/src/modules/Elsa.Workflows.Runtime/Contracts/ILogRecordExtractor.cs
new file mode 100644
index 000000000..43dd124ed
--- /dev/null
+++ b/src/modules/Elsa.Workflows.Runtime/Contracts/ILogRecordExtractor.cs
@@ -0,0 +1,10 @@
+using Elsa.Common.Contracts;
+
+namespace Elsa.Workflows.Runtime;
+
+/// Extracts execution log records.
+public interface ILogRecordExtractor where T: ILogRecord
+{
+ /// Extracts execution logs from a workflow execution context.
+ IEnumerable ExtractLogRecords(WorkflowExecutionContext context);
+}
\ No newline at end of file
diff --git a/src/modules/Elsa.Workflows.Runtime/Contracts/ILogRecordSink.cs b/src/modules/Elsa.Workflows.Runtime/Contracts/ILogRecordSink.cs
new file mode 100644
index 000000000..d49280a1d
--- /dev/null
+++ b/src/modules/Elsa.Workflows.Runtime/Contracts/ILogRecordSink.cs
@@ -0,0 +1,10 @@
+using Elsa.Common.Contracts;
+
+namespace Elsa.Workflows.Runtime;
+
+/// Represents a sink for storing log records.
+public interface ILogRecordSink where T : ILogRecord
+{
+ /// 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/Contracts/ILogRecordStore.cs b/src/modules/Elsa.Workflows.Runtime/Contracts/ILogRecordStore.cs
new file mode 100644
index 000000000..22f90bbc3
--- /dev/null
+++ b/src/modules/Elsa.Workflows.Runtime/Contracts/ILogRecordStore.cs
@@ -0,0 +1,13 @@
+using Elsa.Common.Contracts;
+
+namespace Elsa.Workflows.Runtime.Contracts;
+
+/// Represents a store of log records.
+public interface ILogRecordStore where T : ILogRecord
+{
+ /// Adds or updates the specified set oflog record objects in the persistence store.
+ ///
+ /// If a record does not already exist, it is added to the store; if it does exist, its existing entry is updated.
+ ///
+ Task SaveManyAsync(IEnumerable records, CancellationToken cancellationToken = default);
+}
\ No newline at end of file
diff --git a/src/modules/Elsa.Workflows.Runtime/Contracts/IWorkflowExecutionLogRecordExtractor.cs b/src/modules/Elsa.Workflows.Runtime/Contracts/IWorkflowExecutionLogRecordExtractor.cs
deleted file mode 100644
index 0e704bd55..000000000
--- a/src/modules/Elsa.Workflows.Runtime/Contracts/IWorkflowExecutionLogRecordExtractor.cs
+++ /dev/null
@@ -1,10 +0,0 @@
-using Elsa.Workflows.Runtime.Entities;
-
-namespace Elsa.Workflows.Runtime;
-
-/// Extracts workflow execution log records.
-public interface IWorkflowExecutionLogRecordExtractor
-{
- /// Extracts workflow execution logs from a workflow execution context.
- IEnumerable ExtractWorkflowExecutionLogs(WorkflowExecutionContext context);
-}
\ No newline at end of file
diff --git a/src/modules/Elsa.Workflows.Runtime/Contracts/IWorkflowExecutionLogSink.cs b/src/modules/Elsa.Workflows.Runtime/Contracts/IWorkflowExecutionLogSink.cs
deleted file mode 100644
index 2c0601c79..000000000
--- a/src/modules/Elsa.Workflows.Runtime/Contracts/IWorkflowExecutionLogSink.cs
+++ /dev/null
@@ -1,12 +0,0 @@
-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/Contracts/IWorkflowExecutionLogStore.cs b/src/modules/Elsa.Workflows.Runtime/Contracts/IWorkflowExecutionLogStore.cs
index 386a82b2a..eb5ae310a 100644
--- a/src/modules/Elsa.Workflows.Runtime/Contracts/IWorkflowExecutionLogStore.cs
+++ b/src/modules/Elsa.Workflows.Runtime/Contracts/IWorkflowExecutionLogStore.cs
@@ -8,7 +8,7 @@ namespace Elsa.Workflows.Runtime.Contracts;
///
/// Represents a store of .
///
-public interface IWorkflowExecutionLogStore
+public interface IWorkflowExecutionLogStore : ILogRecordStore
{
///
/// Adds the specified to te persistence store.
@@ -28,14 +28,6 @@ public interface IWorkflowExecutionLogStore
///
Task SaveAsync(WorkflowExecutionLogRecord record, CancellationToken cancellationToken = default);
- ///
- /// Adds or updates the specified set of objects in the persistence store.
- ///
- ///
- /// If a record does not already exist, it is added to the store; if it does exist, its existing entry is updated.
- ///
- Task SaveManyAsync(IEnumerable records, CancellationToken cancellationToken = default);
-
///
/// Returns the first workflow execution log record matching the specified filter.
///
diff --git a/src/modules/Elsa.Workflows.Runtime/Entities/ActivityExecutionRecord.cs b/src/modules/Elsa.Workflows.Runtime/Entities/ActivityExecutionRecord.cs
index 66e29fac2..ee46efef6 100644
--- a/src/modules/Elsa.Workflows.Runtime/Entities/ActivityExecutionRecord.cs
+++ b/src/modules/Elsa.Workflows.Runtime/Entities/ActivityExecutionRecord.cs
@@ -1,3 +1,4 @@
+using Elsa.Common.Contracts;
using Elsa.Common.Entities;
using Elsa.Workflows.State;
@@ -6,7 +7,7 @@ namespace Elsa.Workflows.Runtime.Entities;
///
/// Represents a single activity execution of an activity instance.
///
-public class ActivityExecutionRecord : Entity
+public class ActivityExecutionRecord : Entity, ILogRecord
{
///
/// Gets or sets the workflow instance ID.
diff --git a/src/modules/Elsa.Workflows.Runtime/Entities/WorkflowExecutionLogRecord.cs b/src/modules/Elsa.Workflows.Runtime/Entities/WorkflowExecutionLogRecord.cs
index c2eb0adfa..11194aeb4 100644
--- a/src/modules/Elsa.Workflows.Runtime/Entities/WorkflowExecutionLogRecord.cs
+++ b/src/modules/Elsa.Workflows.Runtime/Entities/WorkflowExecutionLogRecord.cs
@@ -1,3 +1,4 @@
+using Elsa.Common.Contracts;
using Elsa.Common.Entities;
namespace Elsa.Workflows.Runtime.Entities;
@@ -5,7 +6,7 @@ namespace Elsa.Workflows.Runtime.Entities;
///
/// Represents a workflow execution log entry.
///
-public class WorkflowExecutionLogRecord : Entity
+public class WorkflowExecutionLogRecord : Entity, ILogRecord
{
///
/// The ID of the workflow definition.
diff --git a/src/modules/Elsa.Workflows.Runtime/Features/WorkflowRuntimeFeature.cs b/src/modules/Elsa.Workflows.Runtime/Features/WorkflowRuntimeFeature.cs
index 6bdbb7d99..d84dd8f50 100644
--- a/src/modules/Elsa.Workflows.Runtime/Features/WorkflowRuntimeFeature.cs
+++ b/src/modules/Elsa.Workflows.Runtime/Features/WorkflowRuntimeFeature.cs
@@ -108,10 +108,15 @@ public class WorkflowRuntimeFeature : FeatureBase
public Func BackgroundActivityScheduler { get; set; } = sp => ActivatorUtilities.CreateInstance(sp);
///
- /// Represents a sink for workflow execution logs.
+ /// A factory that instantiates an for an .
///
- public Func WorkflowExecutionLogSink { get; set; } = sp => sp.GetRequiredService();
-
+ public Func> ActivityExecutionLogSink { get; set; } = sp => sp.GetRequiredService();
+
+ ///
+ /// A factory that instantiates an for an .
+ ///
+ public Func> WorkflowExecutionLogSink { get; set; } = sp => sp.GetRequiredService();
+
///
/// A delegate to configure the .
///
@@ -210,6 +215,7 @@ public class WorkflowRuntimeFeature : FeatureBase
.AddScoped(WorkflowCancellationDispatcher)
.AddScoped(WorkflowExecutionContextStore)
.AddScoped(RunTaskDispatcher)
+ .AddScoped(ActivityExecutionLogSink)
.AddScoped(WorkflowExecutionLogSink)
.AddSingleton(BackgroundActivityScheduler)
.AddSingleton()
@@ -225,13 +231,15 @@ public class WorkflowRuntimeFeature : FeatureBase
.AddScoped()
.AddScoped()
.AddScoped()
+ .AddScoped()
.AddScoped()
.AddScoped()
.AddScoped()
.AddScoped()
.AddScoped()
.AddScoped()
- .AddScoped()
+ .AddScoped, ActivityExecutionRecordExtractor>()
+ .AddScoped, WorkflowExecutionLogRecordExtractor>()
// Stores.
.AddScoped(BookmarkStore)
diff --git a/src/modules/Elsa.Workflows.Runtime/Middleware/Workflows/PersistActivityExecutionLogMiddleware.cs b/src/modules/Elsa.Workflows.Runtime/Middleware/Workflows/PersistActivityExecutionLogMiddleware.cs
index 0f989d19a..512fb616a 100644
--- a/src/modules/Elsa.Workflows.Runtime/Middleware/Workflows/PersistActivityExecutionLogMiddleware.cs
+++ b/src/modules/Elsa.Workflows.Runtime/Middleware/Workflows/PersistActivityExecutionLogMiddleware.cs
@@ -1,31 +1,11 @@
-using Elsa.Mediator.Contracts;
using Elsa.Workflows.Pipelines.WorkflowExecution;
-using Elsa.Workflows.Runtime.Contracts;
-using Elsa.Workflows.Runtime.Notifications;
+using Elsa.Workflows.Runtime.Entities;
namespace Elsa.Workflows.Runtime.Middleware.Workflows;
-///
/// Creates and updates activity execution records from activity execution contexts.
-///
-public class PersistActivityExecutionLogMiddleware : WorkflowExecutionMiddleware
+public class PersistActivityExecutionLogMiddleware(WorkflowMiddlewareDelegate next, ILogRecordSink sink) : WorkflowExecutionMiddleware(next)
{
- private readonly IActivityExecutionStore _activityExecutionStore;
- private readonly IActivityExecutionMapper _activityExecutionMapper;
- private readonly INotificationSender _notificationSender;
-
- ///
- public PersistActivityExecutionLogMiddleware(
- WorkflowMiddlewareDelegate next,
- IActivityExecutionStore activityExecutionStore,
- IActivityExecutionMapper activityExecutionMapper,
- INotificationSender notificationSender) : base(next)
- {
- _activityExecutionStore = activityExecutionStore;
- _activityExecutionMapper = activityExecutionMapper;
- _notificationSender = notificationSender;
- }
-
///
public override async ValueTask InvokeAsync(WorkflowExecutionContext context)
{
@@ -34,14 +14,7 @@ public class PersistActivityExecutionLogMiddleware : WorkflowExecutionMiddleware
// Get the managed cancellation token.
var cancellationToken = context.CancellationTokens.SystemCancellationToken;
-
- // Get all activity execution contexts.
- var activityExecutionContexts = context.ActivityExecutionContexts;
-
- // Persist activity execution entries.
- var entries = activityExecutionContexts.Select(_activityExecutionMapper.Map).ToList();
- await _activityExecutionStore.SaveManyAsync(entries, cancellationToken);
- await _notificationSender.SendAsync(new ActivityExecutionLogUpdated(context, entries), cancellationToken);
+ await sink.PersistExecutionLogsAsync(context, cancellationToken);
}
}
\ No newline at end of file
diff --git a/src/modules/Elsa.Workflows.Runtime/Middleware/Workflows/PersistWorkflowExecutionLogMiddleware.cs b/src/modules/Elsa.Workflows.Runtime/Middleware/Workflows/PersistWorkflowExecutionLogMiddleware.cs
index 36ba5f1e6..25d863054 100644
--- a/src/modules/Elsa.Workflows.Runtime/Middleware/Workflows/PersistWorkflowExecutionLogMiddleware.cs
+++ b/src/modules/Elsa.Workflows.Runtime/Middleware/Workflows/PersistWorkflowExecutionLogMiddleware.cs
@@ -1,12 +1,12 @@
using Elsa.Workflows.Pipelines.WorkflowExecution;
-using Elsa.Workflows.Runtime.Contracts;
+using Elsa.Workflows.Runtime.Entities;
namespace Elsa.Workflows.Runtime.Middleware.Workflows;
///
/// Takes care of persisting workflow execution log entries.
///
-public class PersistWorkflowExecutionLogMiddleware(WorkflowMiddlewareDelegate next, IWorkflowExecutionLogSink sink) : WorkflowExecutionMiddleware(next)
+public class PersistWorkflowExecutionLogMiddleware(WorkflowMiddlewareDelegate next, ILogRecordSink sink) : WorkflowExecutionMiddleware(next)
{
///
public override async ValueTask InvokeAsync(WorkflowExecutionContext context)
diff --git a/src/modules/Elsa.Workflows.Runtime/Services/ActivityExecutionRecordExtractor.cs b/src/modules/Elsa.Workflows.Runtime/Services/ActivityExecutionRecordExtractor.cs
new file mode 100644
index 000000000..85d02c69e
--- /dev/null
+++ b/src/modules/Elsa.Workflows.Runtime/Services/ActivityExecutionRecordExtractor.cs
@@ -0,0 +1,15 @@
+using Elsa.Workflows.Runtime.Contracts;
+using Elsa.Workflows.Runtime.Entities;
+
+namespace Elsa.Workflows.Runtime.Services;
+
+/// Extracts activity execution log records.
+public class ActivityExecutionRecordExtractor(IActivityExecutionMapper activityExecutionMapper) : ILogRecordExtractor
+{
+ ///
+ public IEnumerable ExtractLogRecords(WorkflowExecutionContext context)
+ {
+ var activityExecutionContexts = context.ActivityExecutionContexts;
+ return activityExecutionContexts.Select(activityExecutionMapper.Map).ToList();
+ }
+}
\ No newline at end of file
diff --git a/src/modules/Elsa.Workflows.Runtime/Services/StoreActivityExecutionLogSink.cs b/src/modules/Elsa.Workflows.Runtime/Services/StoreActivityExecutionLogSink.cs
new file mode 100644
index 000000000..fb32e08c5
--- /dev/null
+++ b/src/modules/Elsa.Workflows.Runtime/Services/StoreActivityExecutionLogSink.cs
@@ -0,0 +1,21 @@
+using Elsa.Mediator.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 StoreActivityExecutionLogSink(IActivityExecutionStore activityExecutionStore, ILogRecordExtractor extractor, INotificationSender notificationSender)
+ : ILogRecordSink
+{
+ ///
+ public async Task PersistExecutionLogsAsync(WorkflowExecutionContext context, CancellationToken cancellationToken = default)
+ {
+ var records = extractor.ExtractLogRecords(context).ToList();
+ await activityExecutionStore.SaveManyAsync(records, cancellationToken);
+ await notificationSender.SendAsync(new ActivityExecutionLogUpdated(context, records), cancellationToken);
+ }
+}
\ 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
index d32bb89db..b3d433e36 100644
--- a/src/modules/Elsa.Workflows.Runtime/Services/StoreWorkflowExecutionLogSink.cs
+++ b/src/modules/Elsa.Workflows.Runtime/Services/StoreWorkflowExecutionLogSink.cs
@@ -8,12 +8,12 @@ namespace Elsa.Workflows.Runtime.Services;
///
/// This implementation saves directly through the store.
///
-public class StoreWorkflowExecutionLogSink(IWorkflowExecutionLogStore store, IWorkflowExecutionLogRecordExtractor extractor, INotificationSender notificationSender) : IWorkflowExecutionLogSink
+public class StoreWorkflowExecutionLogSink(IWorkflowExecutionLogStore store, ILogRecordExtractor extractor, INotificationSender notificationSender) : ILogRecordSink
{
///
public async Task PersistExecutionLogsAsync(WorkflowExecutionContext context, CancellationToken cancellationToken)
{
- var records = extractor.ExtractWorkflowExecutionLogs(context).ToList();
+ var records = extractor.ExtractLogRecords(context).ToList();
await store.AddManyAsync(records, context.CancellationTokens.SystemCancellationToken);
await notificationSender.SendAsync(new WorkflowExecutionLogUpdated(context), context.CancellationTokens.SystemCancellationToken);
}
diff --git a/src/modules/Elsa.Workflows.Runtime/Services/WorkflowExecutionContext.cs b/src/modules/Elsa.Workflows.Runtime/Services/WorkflowExecutionLogRecordExtractor.cs
similarity index 87%
rename from src/modules/Elsa.Workflows.Runtime/Services/WorkflowExecutionContext.cs
rename to src/modules/Elsa.Workflows.Runtime/Services/WorkflowExecutionLogRecordExtractor.cs
index 577559bdb..d807fa3af 100644
--- a/src/modules/Elsa.Workflows.Runtime/Services/WorkflowExecutionContext.cs
+++ b/src/modules/Elsa.Workflows.Runtime/Services/WorkflowExecutionLogRecordExtractor.cs
@@ -4,10 +4,10 @@ using Elsa.Workflows.Runtime.Entities;
namespace Elsa.Workflows.Runtime.Services;
///
-public class WorkflowExecutionLogRecordExtractor(IIdentityGenerator identityGenerator) : IWorkflowExecutionLogRecordExtractor
+public class WorkflowExecutionLogRecordExtractor(IIdentityGenerator identityGenerator) : ILogRecordExtractor
{
///
- public IEnumerable ExtractWorkflowExecutionLogs(WorkflowExecutionContext context)
+ public IEnumerable ExtractLogRecords(WorkflowExecutionContext context)
{
return context.ExecutionLog.Select(x => new WorkflowExecutionLogRecord
{