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
{