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.
This commit is contained in:
raymonddenhaan 2024-07-24 22:30:14 +02:00 committed by GitHub
parent 53cb8e75c3
commit c413f210a9
No known key found for this signature in database
GPG key ID: B5690EEEBB952194
17 changed files with 100 additions and 84 deletions

View file

@ -0,0 +1,4 @@
namespace Elsa.Common.Contracts;
/// Represents a log record.
public interface ILogRecord;

View file

@ -7,7 +7,7 @@ namespace Elsa.Workflows.Runtime.Contracts;
/// <summary>
/// Stores activity execution records.
/// </summary>
public interface IActivityExecutionStore
public interface IActivityExecutionStore : ILogRecordStore<ActivityExecutionRecord>
{
/// <summary>
/// Adds or updates the specified <see cref="ActivityExecutionRecord"/> in the persistence store.
@ -19,16 +19,6 @@ public interface IActivityExecutionStore
/// </remarks>
Task SaveAsync(ActivityExecutionRecord record, CancellationToken cancellationToken = default);
/// <summary>
/// Adds or updates the specified set of <see cref="ActivityExecutionRecord"/> objects in the persistence store.
/// </summary>
/// <param name="records">The activity execution records.</param>
/// <param name="cancellationToken">An optional cancellation token.</param>
/// <remarks>
/// If a record does not already exist, it is added to the store; if it does exist, its existing entry is updated.
/// </remarks>
Task SaveManyAsync(IEnumerable<ActivityExecutionRecord> records, CancellationToken cancellationToken = default);
/// <summary>
/// Finds an activity execution record matching the specified filter.
/// </summary>

View file

@ -0,0 +1,10 @@
using Elsa.Common.Contracts;
namespace Elsa.Workflows.Runtime;
/// Extracts execution log records.
public interface ILogRecordExtractor<T> where T: ILogRecord
{
/// Extracts execution logs from a workflow execution context.
IEnumerable<T> ExtractLogRecords(WorkflowExecutionContext context);
}

View file

@ -0,0 +1,10 @@
using Elsa.Common.Contracts;
namespace Elsa.Workflows.Runtime;
/// Represents a sink for storing log records.
public interface ILogRecordSink<in T> where T : ILogRecord
{
/// Persists the execution logs of a workflow.
Task PersistExecutionLogsAsync(WorkflowExecutionContext context, CancellationToken cancellationToken = default);
}

View file

@ -0,0 +1,13 @@
using Elsa.Common.Contracts;
namespace Elsa.Workflows.Runtime.Contracts;
/// Represents a store of log records.
public interface ILogRecordStore<in T> where T : ILogRecord
{
/// Adds or updates the specified set oflog record objects in the persistence store.
/// <remarks>
/// If a record does not already exist, it is added to the store; if it does exist, its existing entry is updated.
/// </remarks>
Task SaveManyAsync(IEnumerable<T> records, CancellationToken cancellationToken = default);
}

View file

@ -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<WorkflowExecutionLogRecord> ExtractWorkflowExecutionLogs(WorkflowExecutionContext context);
}

View file

@ -1,12 +0,0 @@
namespace Elsa.Workflows.Runtime.Contracts;
/// <summary>
/// Represents a sink for storing workflow execution log records.
/// </summary>
public interface IWorkflowExecutionLogSink
{
/// <summary>
/// Persists the execution logs of a workflow.
/// </summary>
Task PersistExecutionLogsAsync(WorkflowExecutionContext context, CancellationToken cancellationToken = default);
}

View file

@ -8,7 +8,7 @@ namespace Elsa.Workflows.Runtime.Contracts;
/// <summary>
/// Represents a store of <see cref="WorkflowExecutionLogRecord"/>.
/// </summary>
public interface IWorkflowExecutionLogStore
public interface IWorkflowExecutionLogStore : ILogRecordStore<WorkflowExecutionLogRecord>
{
/// <summary>
/// Adds the specified <see cref="WorkflowExecutionLogRecord"/> to te persistence store.
@ -28,14 +28,6 @@ public interface IWorkflowExecutionLogStore
/// </remarks>
Task SaveAsync(WorkflowExecutionLogRecord record, CancellationToken cancellationToken = default);
/// <summary>
/// Adds or updates the specified set of <see cref="WorkflowExecutionLogRecord"/> objects in the persistence store.
/// </summary>
/// <remarks>
/// If a record does not already exist, it is added to the store; if it does exist, its existing entry is updated.
/// </remarks>
Task SaveManyAsync(IEnumerable<WorkflowExecutionLogRecord> records, CancellationToken cancellationToken = default);
/// <summary>
/// Returns the first workflow execution log record matching the specified filter.
/// </summary>

View file

@ -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;
/// <summary>
/// Represents a single activity execution of an activity instance.
/// </summary>
public class ActivityExecutionRecord : Entity
public class ActivityExecutionRecord : Entity, ILogRecord
{
/// <summary>
/// Gets or sets the workflow instance ID.

View file

@ -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;
/// <summary>
/// Represents a workflow execution log entry.
/// </summary>
public class WorkflowExecutionLogRecord : Entity
public class WorkflowExecutionLogRecord : Entity, ILogRecord
{
/// <summary>
/// The ID of the workflow definition.

View file

@ -108,10 +108,15 @@ public class WorkflowRuntimeFeature : FeatureBase
public Func<IServiceProvider, IBackgroundActivityScheduler> BackgroundActivityScheduler { get; set; } = sp => ActivatorUtilities.CreateInstance<LocalBackgroundActivityScheduler>(sp);
/// <summary>
/// Represents a sink for workflow execution logs.
/// A factory that instantiates an <see cref="ILogRecordSink"/> for an <see cref="ActivityExecutionRecord"/>.
/// </summary>
public Func<IServiceProvider, IWorkflowExecutionLogSink> WorkflowExecutionLogSink { get; set; } = sp => sp.GetRequiredService<StoreWorkflowExecutionLogSink>();
public Func<IServiceProvider, ILogRecordSink<ActivityExecutionRecord>> ActivityExecutionLogSink { get; set; } = sp => sp.GetRequiredService<StoreActivityExecutionLogSink>();
/// <summary>
/// A factory that instantiates an <see cref="ILogRecordSink"/> for an <see cref="WorkflowExecutionLogRecord"/>.
/// </summary>
public Func<IServiceProvider, ILogRecordSink<WorkflowExecutionLogRecord>> WorkflowExecutionLogSink { get; set; } = sp => sp.GetRequiredService<StoreWorkflowExecutionLogSink>();
/// <summary>
/// A delegate to configure the <see cref="DistributedLockingOptions"/>.
/// </summary>
@ -210,6 +215,7 @@ public class WorkflowRuntimeFeature : FeatureBase
.AddScoped(WorkflowCancellationDispatcher)
.AddScoped(WorkflowExecutionContextStore)
.AddScoped(RunTaskDispatcher)
.AddScoped(ActivityExecutionLogSink)
.AddScoped(WorkflowExecutionLogSink)
.AddSingleton(BackgroundActivityScheduler)
.AddSingleton<RandomLongIdentityGenerator>()
@ -225,13 +231,15 @@ public class WorkflowRuntimeFeature : FeatureBase
.AddScoped<ITaskReporter, TaskReporter>()
.AddScoped<SynchronousTaskDispatcher>()
.AddScoped<BackgroundTaskDispatcher>()
.AddScoped<StoreActivityExecutionLogSink>()
.AddScoped<StoreWorkflowExecutionLogSink>()
.AddScoped<IEventPublisher, EventPublisher>()
.AddScoped<IWorkflowInbox, DefaultWorkflowInbox>()
.AddScoped<IBookmarkUpdater, BookmarkUpdater>()
.AddScoped<IBookmarksPersister, BookmarksPersister>()
.AddScoped<IWorkflowCancellationService, WorkflowCancellationService>()
.AddScoped<IWorkflowExecutionLogRecordExtractor, WorkflowExecutionLogRecordExtractor>()
.AddScoped<ILogRecordExtractor<ActivityExecutionRecord>, ActivityExecutionRecordExtractor>()
.AddScoped<ILogRecordExtractor<WorkflowExecutionLogRecord>, WorkflowExecutionLogRecordExtractor>()
// Stores.
.AddScoped(BookmarkStore)

View file

@ -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;
/// <summary>
/// Creates and updates activity execution records from activity execution contexts.
/// </summary>
public class PersistActivityExecutionLogMiddleware : WorkflowExecutionMiddleware
public class PersistActivityExecutionLogMiddleware(WorkflowMiddlewareDelegate next, ILogRecordSink<ActivityExecutionRecord> sink) : WorkflowExecutionMiddleware(next)
{
private readonly IActivityExecutionStore _activityExecutionStore;
private readonly IActivityExecutionMapper _activityExecutionMapper;
private readonly INotificationSender _notificationSender;
/// <inheritdoc />
public PersistActivityExecutionLogMiddleware(
WorkflowMiddlewareDelegate next,
IActivityExecutionStore activityExecutionStore,
IActivityExecutionMapper activityExecutionMapper,
INotificationSender notificationSender) : base(next)
{
_activityExecutionStore = activityExecutionStore;
_activityExecutionMapper = activityExecutionMapper;
_notificationSender = notificationSender;
}
/// <inheritdoc />
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);
}
}

View file

@ -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;
/// <summary>
/// Takes care of persisting workflow execution log entries.
/// </summary>
public class PersistWorkflowExecutionLogMiddleware(WorkflowMiddlewareDelegate next, IWorkflowExecutionLogSink sink) : WorkflowExecutionMiddleware(next)
public class PersistWorkflowExecutionLogMiddleware(WorkflowMiddlewareDelegate next, ILogRecordSink<WorkflowExecutionLogRecord> sink) : WorkflowExecutionMiddleware(next)
{
/// <inheritdoc />
public override async ValueTask InvokeAsync(WorkflowExecutionContext context)

View file

@ -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<ActivityExecutionRecord>
{
/// <inheritdoc />
public IEnumerable<ActivityExecutionRecord> ExtractLogRecords(WorkflowExecutionContext context)
{
var activityExecutionContexts = context.ActivityExecutionContexts;
return activityExecutionContexts.Select(activityExecutionMapper.Map).ToList();
}
}

View file

@ -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;
/// <summary>
/// This implementation saves <see cref="ActivityExecutionRecord"/> directly through the store.
/// </summary>
public class StoreActivityExecutionLogSink(IActivityExecutionStore activityExecutionStore, ILogRecordExtractor<ActivityExecutionRecord> extractor, INotificationSender notificationSender)
: ILogRecordSink<ActivityExecutionRecord>
{
/// <inheritdoc />
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);
}
}

View file

@ -8,12 +8,12 @@ namespace Elsa.Workflows.Runtime.Services;
/// <summary>
/// This implementation saves <see cref="WorkflowExecutionLogRecord"/> directly through the store.
/// </summary>
public class StoreWorkflowExecutionLogSink(IWorkflowExecutionLogStore store, IWorkflowExecutionLogRecordExtractor extractor, INotificationSender notificationSender) : IWorkflowExecutionLogSink
public class StoreWorkflowExecutionLogSink(IWorkflowExecutionLogStore store, ILogRecordExtractor<WorkflowExecutionLogRecord> extractor, INotificationSender notificationSender) : ILogRecordSink<WorkflowExecutionLogRecord>
{
/// <inheritdoc />
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);
}

View file

@ -4,10 +4,10 @@ using Elsa.Workflows.Runtime.Entities;
namespace Elsa.Workflows.Runtime.Services;
/// <inheritdoc />
public class WorkflowExecutionLogRecordExtractor(IIdentityGenerator identityGenerator) : IWorkflowExecutionLogRecordExtractor
public class WorkflowExecutionLogRecordExtractor(IIdentityGenerator identityGenerator) : ILogRecordExtractor<WorkflowExecutionLogRecord>
{
/// <inheritdoc />
public IEnumerable<WorkflowExecutionLogRecord> ExtractWorkflowExecutionLogs(WorkflowExecutionContext context)
public IEnumerable<WorkflowExecutionLogRecord> ExtractLogRecords(WorkflowExecutionContext context)
{
return context.ExecutionLog.Select(x => new WorkflowExecutionLogRecord
{