Add Workflow Heartbeat Middleware and update methods

Introduce `WorkflowHeartbeatMiddleware` to monitor workflow liveness and update timestamps during execution. Extend WorkflowInstanceStore implementations to support updating the `UpdatedAt` timestamp across various data stores. Enhance Dapper query builder and workflow execution pipeline with related functionality.
This commit is contained in:
Sipke Schoorstra 2025-02-22 16:05:36 +01:00
parent 0a21b5201a
commit 5ae23ce7e5
No known key found for this signature in database
GPG key ID: 5C10502B28A4268F
13 changed files with 236 additions and 20 deletions

View file

@ -73,6 +73,6 @@ public static class PropertyAccessorExtensions
: expression.Body is UnaryExpression unaryExpression
? unaryExpression.Operand is MemberExpression unaryMemberExpression
? unaryMemberExpression.Member as PropertyInfo
: default
: default;
: null
: null;
}

View file

@ -427,6 +427,21 @@ public static class ParameterizedQueryBuilderExtensions
.Select(x => x.Name)
.ToArray();
return Update(query, table, record, primaryKeyField, fields, getParameterName);
}
/// <summary>
/// Constructs an UPDATE query for the specified table and applies the provided record, primaryKeyField, and specified fields.
/// </summary>
/// <param name="query">The query being built.</param>
/// <param name="table">The name of the table to update.</param>
/// <param name="record">The object containing the values to be updated.</param>
/// <param name="primaryKeyField">The name of the primary key field to identify the record.</param>
/// <param name="fields">An array of field names to include in the update statement.</param>
/// <param name="getParameterName">An optional function to customize parameter names for the query.</param>
/// <returns>Returns the updated instance of <see cref="ParameterizedQuery"/>.</returns>
public static ParameterizedQuery Update(this ParameterizedQuery query, string table, object record, string primaryKeyField, string[] fields, Func<string, string>? getParameterName = default)
{
getParameterName ??= x => x;
query.Sql.AppendLine(query.Dialect.Update(table, primaryKeyField, fields, getParameterName));

View file

@ -160,6 +160,16 @@ internal class DapperWorkflowInstanceStore(Store<WorkflowInstanceRecord> store,
return await store.DeleteAsync(q => ApplyFilter(q, filter), cancellationToken);
}
public async Task UpdateUpdatedTimestampAsync(string workflowInstanceId, DateTimeOffset value, CancellationToken cancellationToken = default)
{
var record = new WorkflowInstanceRecord
{
Id = workflowInstanceId,
UpdatedAt = value
};
await store.UpdateAsync(record, [x => x.UpdatedAt], cancellationToken);
}
private void ApplyFilter(ParameterizedQuery query, WorkflowInstanceFilter filter)
{
query
@ -190,7 +200,7 @@ internal class DapperWorkflowInstanceStore(Store<WorkflowInstanceRecord> store,
[RequiresUnreferencedCode("Calls Elsa.Workflows.Contracts.IWorkflowStateSerializer.DeserializeAsync(String, CancellationToken)")]
private IEnumerable<WorkflowInstance> Map(IEnumerable<WorkflowInstanceRecord> source)
{
return source.Select( Map);
return source.Select(Map);
}
[RequiresUnreferencedCode("Calls Elsa.Workflows.Contracts.IWorkflowStateSerializer.DeserializeAsync(String, CancellationToken)")]

View file

@ -1,3 +1,4 @@
using System.Linq.Expressions;
using Dapper;
using Elsa.Common.Entities;
using Elsa.Common.Models;
@ -6,6 +7,7 @@ using Elsa.Dapper.Contracts;
using Elsa.Dapper.Extensions;
using Elsa.Dapper.Models;
using Elsa.Dapper.Records;
using Elsa.Extensions;
using JetBrains.Annotations;
namespace Elsa.Dapper.Services;
@ -410,6 +412,7 @@ public class Store<T>(IDbConnectionProvider dbConnectionProvider, ITenantAccesso
foreach (var record in recordsList)
{
var index = currentIndex;
SetTenantId(record);
query.Insert(TableName, record, field => $"{field}_{index}");
currentIndex++;
}
@ -426,7 +429,17 @@ public class Store<T>(IDbConnectionProvider dbConnectionProvider, ITenantAccesso
public async Task UpdateAsync(T record, CancellationToken cancellationToken = default)
{
using var connection = dbConnectionProvider.GetConnection();
var query = new ParameterizedQuery(dbConnectionProvider.Dialect).Insert(TableName, record);
SetTenantId(record);
var query = new ParameterizedQuery(dbConnectionProvider.Dialect).Update(TableName, record, PrimaryKey);
await query.ExecuteAsync(connection);
}
public async Task UpdateAsync(T record, Expression<Func<T, object>>[] props, CancellationToken cancellationToken = default)
{
using var connection = dbConnectionProvider.GetConnection();
SetTenantId(record);
var fields = props.Select(x => x.GetPropertyName()).ToArray();
var query = new ParameterizedQuery(dbConnectionProvider.Dialect).Update(TableName, record, PrimaryKey, fields);
await query.ExecuteAsync(connection);
}

View file

@ -35,16 +35,22 @@ public class ElasticWorkflowInstanceStore : IWorkflowInstanceStore
}
/// <inheritdoc />
public async ValueTask<Page<WorkflowInstance>> FindManyAsync(WorkflowInstanceFilter filter, PageArgs pageArgs, CancellationToken cancellationToken = default) =>
await _store.SearchAsync(d => Filter(d, filter), pageArgs, cancellationToken);
public async ValueTask<Page<WorkflowInstance>> FindManyAsync(WorkflowInstanceFilter filter, PageArgs pageArgs, CancellationToken cancellationToken = default)
{
return await _store.SearchAsync(d => Filter(d, filter), pageArgs, cancellationToken);
}
/// <inheritdoc />
public async ValueTask<Page<WorkflowInstance>> FindManyAsync<TOrderBy>(WorkflowInstanceFilter filter, PageArgs pageArgs, WorkflowInstanceOrder<TOrderBy> order, CancellationToken cancellationToken = default) =>
await _store.SearchAsync(d => Sort(Filter(d, filter), order), pageArgs, cancellationToken);
public async ValueTask<Page<WorkflowInstance>> FindManyAsync<TOrderBy>(WorkflowInstanceFilter filter, PageArgs pageArgs, WorkflowInstanceOrder<TOrderBy> order, CancellationToken cancellationToken = default)
{
return await _store.SearchAsync(d => Sort(Filter(d, filter), order), pageArgs, cancellationToken);
}
/// <inheritdoc />
public async ValueTask<IEnumerable<WorkflowInstance>> FindManyAsync(WorkflowInstanceFilter filter, CancellationToken cancellationToken = default) =>
await _store.SearchAsync(d => Filter(d, filter), cancellationToken);
public async ValueTask<IEnumerable<WorkflowInstance>> FindManyAsync(WorkflowInstanceFilter filter, CancellationToken cancellationToken = default)
{
return await _store.SearchAsync(d => Filter(d, filter), cancellationToken);
}
/// <inheritdoc />
public async ValueTask<IEnumerable<WorkflowInstance>> FindManyAsync<TOrderBy>(WorkflowInstanceFilter filter, WorkflowInstanceOrder<TOrderBy> order, CancellationToken cancellationToken = default)
@ -127,12 +133,21 @@ public class ElasticWorkflowInstanceStore : IWorkflowInstanceStore
}
/// <inheritdoc />
public async ValueTask SaveManyAsync(IEnumerable<WorkflowInstance> instances, CancellationToken cancellationToken = default) =>
public async ValueTask SaveManyAsync(IEnumerable<WorkflowInstance> instances, CancellationToken cancellationToken = default)
{
await _store.SaveManyAsync(instances, cancellationToken);
}
/// <inheritdoc />
public async ValueTask<long> DeleteAsync(WorkflowInstanceFilter filter, CancellationToken cancellationToken = default) =>
await _store.DeleteByQueryAsync(d => Filter(d, filter), cancellationToken);
public async ValueTask<long> DeleteAsync(WorkflowInstanceFilter filter, CancellationToken cancellationToken = default)
{
return await _store.DeleteByQueryAsync(d => Filter(d, filter), cancellationToken);
}
public Task UpdateUpdatedTimestampAsync(string workflowInstanceId, DateTimeOffset value, CancellationToken cancellationToken = default)
{
throw new NotImplementedException();
}
private static SearchRequestDescriptor<WorkflowInstance> Sort<TProp>(SearchRequestDescriptor<WorkflowInstance> descriptor, WorkflowInstanceOrder<TProp> order)
{
@ -145,9 +160,20 @@ public class ElasticWorkflowInstanceStore : IWorkflowInstanceStore
return descriptor;
}
private static SearchRequestDescriptor<WorkflowInstance> Filter(SearchRequestDescriptor<WorkflowInstance> descriptor, WorkflowInstanceFilter filter) => descriptor.Query(query => Filter(query, filter));
private static CountRequestDescriptor<WorkflowInstance> Filter(CountRequestDescriptor<WorkflowInstance> descriptor, WorkflowInstanceFilter filter) => descriptor.Query(query => Filter(query, filter));
private static DeleteByQueryRequestDescriptor<WorkflowInstance> Filter(DeleteByQueryRequestDescriptor<WorkflowInstance> descriptor, WorkflowInstanceFilter filter) => descriptor.Query(query => Filter(query, filter));
private static SearchRequestDescriptor<WorkflowInstance> Filter(SearchRequestDescriptor<WorkflowInstance> descriptor, WorkflowInstanceFilter filter)
{
return descriptor.Query(query => Filter(query, filter));
}
private static CountRequestDescriptor<WorkflowInstance> Filter(CountRequestDescriptor<WorkflowInstance> descriptor, WorkflowInstanceFilter filter)
{
return descriptor.Query(query => Filter(query, filter));
}
private static DeleteByQueryRequestDescriptor<WorkflowInstance> Filter(DeleteByQueryRequestDescriptor<WorkflowInstance> descriptor, WorkflowInstanceFilter filter)
{
return descriptor.Query(query => Filter(query, filter));
}
private static QueryDescriptor<WorkflowInstance> Filter(QueryDescriptor<WorkflowInstance> descriptor, WorkflowInstanceFilter filter)
{
@ -176,8 +202,9 @@ public class ElasticWorkflowInstanceStore : IWorkflowInstanceStore
.Query(filter.SearchTerm));
}
private static SearchRequestDescriptor<WorkflowInstance> Summarize(SearchRequestDescriptor<WorkflowInstance> descriptor) =>
descriptor.Fields(
private static SearchRequestDescriptor<WorkflowInstance> Summarize(SearchRequestDescriptor<WorkflowInstance> descriptor)
{
return descriptor.Fields(
field => field.Field(f => f.Id),
field => field.Field(f => f.DefinitionId),
field => field.Field(f => f.Status),
@ -191,7 +218,10 @@ public class ElasticWorkflowInstanceStore : IWorkflowInstanceStore
field => field.Field(f => f.DefinitionVersionId),
field => field.Field(f => f.UpdatedAt)
);
}
private static SearchRequestDescriptor<WorkflowInstance> SelectId(SearchRequestDescriptor<WorkflowInstance> descriptor) =>
descriptor.Fields(field => field.Field(f => f.Id));
private static SearchRequestDescriptor<WorkflowInstance> SelectId(SearchRequestDescriptor<WorkflowInstance> descriptor)
{
return descriptor.Fields(field => field.Field(f => f.Id));
}
}

View file

@ -225,6 +225,24 @@ public class Store<TDbContext, TEntity>(IDbContextFactory<TDbContext> dbContextF
await dbContext.SaveChangesAsync(cancellationToken);
}
/// <summary>
/// Updates specific properties of an entity in the database.
/// </summary>
/// <param name="entity">The entity to update.</param>
/// <param name="properties">An array of expressions indicating the properties to update.</param>
/// <param name="cancellationToken">The cancellation token.</param>
/// <returns>A task that represents the asynchronous operation.</returns>
public async Task UpdatePartialAsync(TEntity entity, Expression<Func<TEntity, object>>[] properties, CancellationToken cancellationToken = default)
{
await using var dbContext = await CreateDbContextAsync(cancellationToken);
dbContext.Attach(entity);
foreach (var property in properties)
dbContext.Entry(entity).Property(property).IsModified = true;
await dbContext.SaveChangesAsync(cancellationToken);
}
/// <summary>
/// Finds the entity matching the specified predicate.
/// </summary>

View file

@ -150,6 +150,16 @@ public class EFCoreWorkflowInstanceStore : IWorkflowInstanceStore
return await _store.DeleteWhereAsync(query => Filter(query, filter), cancellationToken);
}
public async Task UpdateUpdatedTimestampAsync(string workflowInstanceId, DateTimeOffset value, CancellationToken cancellationToken = default)
{
var entity = new WorkflowInstance
{
Id = workflowInstanceId,
UpdatedAt = value
};
await _store.UpdatePartialAsync(entity, [x => x.UpdatedAt], cancellationToken);
}
/// <inheritdoc />
[RequiresUnreferencedCode("Calls Elsa.Workflows.Contracts.IWorkflowStateSerializer.SerializeAsync(WorkflowState, CancellationToken)")]
public async ValueTask SaveAsync(WorkflowInstance instance, CancellationToken cancellationToken = default)

View file

@ -133,6 +133,29 @@ public class MongoDbStore<TDocument>(IMongoCollection<TDocument> collection, ITe
await collection.BulkWriteAsync(writes, cancellationToken: cancellationToken);
}
public async Task UpdatePartialAsync(
string id,
IDictionary<string, object> updatedFields,
string primaryKey = nameof(Entity.Id),
CancellationToken cancellationToken = default)
{
if (string.IsNullOrEmpty(id))
throw new ArgumentNullException(nameof(id));
if (updatedFields == null || updatedFields.Count == 0)
throw new ArgumentException("No fields to update were provided.", nameof(updatedFields));
var filter = Builders<TDocument>.Filter.Eq(primaryKey, id);
var updateDefinition = Builders<TDocument>.Update.Combine(
updatedFields.Select(field => Builders<TDocument>.Update.Set(field.Key, field.Value))
);
var updateResult = await collection.UpdateOneAsync(filter, updateDefinition, cancellationToken: cancellationToken);
if (updateResult.MatchedCount == 0)
throw new InvalidOperationException($"No document found with ID '{id}'.");
}
/// <summary>
/// Finds the document matching the specified predicate

View file

@ -123,6 +123,16 @@ public class MongoWorkflowInstanceStore(MongoDbStore<WorkflowInstance> mongoDbSt
return await mongoDbStore.DeleteWhereAsync<string>(query => Filter(query, filter), x => x.Id, cancellationToken);
}
public async Task UpdateUpdatedTimestampAsync(string workflowInstanceId, DateTimeOffset value, CancellationToken cancellationToken = default)
{
var props = new Dictionary<string, object>
{
[nameof(WorkflowInstance.UpdatedAt)] = value
};
await mongoDbStore.UpdatePartialAsync(workflowInstanceId, props, cancellationToken: cancellationToken);
}
/// <inheritdoc />
public async ValueTask SaveAsync(WorkflowInstance instance, CancellationToken cancellationToken = default)
{

View file

@ -170,4 +170,12 @@ public interface IWorkflowInstanceStore
/// <param name="cancellationToken">The cancellation token.</param>
/// <returns>The number of deleted workflow instances.</returns>
ValueTask<long> DeleteAsync(WorkflowInstanceFilter filter, CancellationToken cancellationToken = default);
/// <summary>
/// Updates the "LastUpdated" timestamp of a workflow instance.
/// </summary>
/// <param name="workflowInstanceId">The unique identifier of the workflow instance.</param>
/// <param name="value">The new timestamp value to set.</param>
/// <param name="cancellationToken">The cancellation token to observe during the operation.</param>
Task UpdateUpdatedTimestampAsync(string workflowInstanceId, DateTimeOffset value, CancellationToken cancellationToken = default);
}

View file

@ -154,6 +154,20 @@ public class MemoryWorkflowInstanceStore : IWorkflowInstanceStore
return ValueTask.FromResult(count);
}
/// <inheritdoc />
public async Task UpdateUpdatedTimestampAsync(string workflowInstanceId, DateTimeOffset value, CancellationToken cancellationToken = default)
{
var workflowInstance = await FindAsync(new()
{
Id = workflowInstanceId
}, cancellationToken);
if (workflowInstance == null)
throw new InvalidOperationException($"Workflow instance with ID '{workflowInstanceId}' does not exist.");
workflowInstance.UpdatedAt = value;
}
private static string GetId(WorkflowInstance workflowInstance) => workflowInstance.Id;
[RequiresUnreferencedCode("Calls Elsa.Workflows.Management.Filters.WorkflowInstanceFilter.Apply(IQueryable<WorkflowInstance>)")]

View file

@ -17,6 +17,7 @@ public static class WorkflowExecutionPipelineBuilderExtensions
public static IWorkflowExecutionPipelineBuilder UseDefaultPipeline(this IWorkflowExecutionPipelineBuilder pipelineBuilder) =>
pipelineBuilder
.Reset()
.UseWorkflowHeartbeat()
.UseEngineExceptionHandling()
.UsePersistentVariables()
.UseExceptionHandling()
@ -26,6 +27,8 @@ public static class WorkflowExecutionPipelineBuilderExtensions
/// Installs middleware that persists the workflow instance before and after workflow execution.
/// </summary>
public static IWorkflowExecutionPipelineBuilder UsePersistentVariables(this IWorkflowExecutionPipelineBuilder pipelineBuilder) => pipelineBuilder.UseMiddleware<PersistentVariablesMiddleware>();
public static IWorkflowExecutionPipelineBuilder UseWorkflowHeartbeat(this IWorkflowExecutionPipelineBuilder pipelineBuilder) => pipelineBuilder.UseMiddleware<WorkflowHeartbeatMiddleware>();
/// <summary>
/// Installs middleware that persists bookmarks after workflow execution.

View file

@ -0,0 +1,62 @@
using Elsa.Common;
using Elsa.Workflows.Management;
using Elsa.Workflows.Pipelines.WorkflowExecution;
using Elsa.Workflows.Runtime.Options;
using Microsoft.Extensions.Logging;
using Microsoft.Extensions.Options;
namespace Elsa.Workflows.Runtime.Middleware.Workflows;
public class WorkflowHeartbeatMiddleware(WorkflowMiddlewareDelegate next, IOptions<RuntimeOptions> options, ILogger<WorkflowHeartbeatMiddleware> logger, ILoggerFactory loggerFactory) : WorkflowExecutionMiddleware(next)
{
public override async ValueTask InvokeAsync(WorkflowExecutionContext context)
{
var livenessThreshold = options.Value.WorkflowLivenessThreshold;
var heartbeatInterval = TimeSpan.FromTicks((long)(livenessThreshold.Ticks * 0.6));
logger.LogDebug("Workflow heartbeat interval: {Interval}", heartbeatInterval);
using var heartbeat = new WorkflowHeartbeat(async () => await UpdateTimestampAsync(context), heartbeatInterval, loggerFactory);
await Next(context);
}
private async Task UpdateTimestampAsync(WorkflowExecutionContext context)
{
var workflowInstanceStore = context.GetRequiredService<IWorkflowInstanceStore>();
var clock = context.GetRequiredService<ISystemClock>();
var workflowInstanceId = context.Id;
await workflowInstanceStore.UpdateUpdatedTimestampAsync(workflowInstanceId, clock.UtcNow, context.CancellationToken);
}
private class WorkflowHeartbeat : IDisposable
{
private readonly Timer _timer;
private readonly Func<Task> _pulseAction;
private readonly TimeSpan _interval;
private readonly ILogger<WorkflowHeartbeat> _logger;
public WorkflowHeartbeat(Func<Task> pulseAction, TimeSpan interval, ILoggerFactory loggerFactory)
{
_pulseAction = pulseAction;
_interval = interval;
_logger = loggerFactory.CreateLogger<WorkflowHeartbeat>();
_timer = new(GeneratePulseAsync, null, _interval, Timeout.InfiniteTimeSpan);
}
private async void GeneratePulseAsync(object? state)
{
try
{
await _pulseAction();
_timer.Change(_interval, Timeout.InfiniteTimeSpan);
}
catch (Exception ex)
{
_logger.LogError(ex, "An error occurred while generating a workflow heartbeat.");
}
}
public void Dispose()
{
_timer.Dispose();
}
}
}