diff --git a/src/modules/Elsa.Common/Extensions/PropertyAccessorExtensions.cs b/src/modules/Elsa.Common/Extensions/PropertyAccessorExtensions.cs index 912415a11..bea2a3bf6 100644 --- a/src/modules/Elsa.Common/Extensions/PropertyAccessorExtensions.cs +++ b/src/modules/Elsa.Common/Extensions/PropertyAccessorExtensions.cs @@ -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; } \ No newline at end of file diff --git a/src/modules/Elsa.Dapper/Extensions/ParameterizedQueryBuilderExtensions.cs b/src/modules/Elsa.Dapper/Extensions/ParameterizedQueryBuilderExtensions.cs index b2383ec54..6896bab9e 100644 --- a/src/modules/Elsa.Dapper/Extensions/ParameterizedQueryBuilderExtensions.cs +++ b/src/modules/Elsa.Dapper/Extensions/ParameterizedQueryBuilderExtensions.cs @@ -427,6 +427,21 @@ public static class ParameterizedQueryBuilderExtensions .Select(x => x.Name) .ToArray(); + return Update(query, table, record, primaryKeyField, fields, getParameterName); + } + + /// + /// Constructs an UPDATE query for the specified table and applies the provided record, primaryKeyField, and specified fields. + /// + /// The query being built. + /// The name of the table to update. + /// The object containing the values to be updated. + /// The name of the primary key field to identify the record. + /// An array of field names to include in the update statement. + /// An optional function to customize parameter names for the query. + /// Returns the updated instance of . + public static ParameterizedQuery Update(this ParameterizedQuery query, string table, object record, string primaryKeyField, string[] fields, Func? getParameterName = default) + { getParameterName ??= x => x; query.Sql.AppendLine(query.Dialect.Update(table, primaryKeyField, fields, getParameterName)); diff --git a/src/modules/Elsa.Dapper/Modules/Management/Stores/DapperWorkflowInstanceStore.cs b/src/modules/Elsa.Dapper/Modules/Management/Stores/DapperWorkflowInstanceStore.cs index e62794fa5..fd24b22cf 100644 --- a/src/modules/Elsa.Dapper/Modules/Management/Stores/DapperWorkflowInstanceStore.cs +++ b/src/modules/Elsa.Dapper/Modules/Management/Stores/DapperWorkflowInstanceStore.cs @@ -160,6 +160,16 @@ internal class DapperWorkflowInstanceStore(Store 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 store, [RequiresUnreferencedCode("Calls Elsa.Workflows.Contracts.IWorkflowStateSerializer.DeserializeAsync(String, CancellationToken)")] private IEnumerable Map(IEnumerable source) { - return source.Select( Map); + return source.Select(Map); } [RequiresUnreferencedCode("Calls Elsa.Workflows.Contracts.IWorkflowStateSerializer.DeserializeAsync(String, CancellationToken)")] diff --git a/src/modules/Elsa.Dapper/Services/Store.cs b/src/modules/Elsa.Dapper/Services/Store.cs index 6381d3d84..f51d47544 100644 --- a/src/modules/Elsa.Dapper/Services/Store.cs +++ b/src/modules/Elsa.Dapper/Services/Store.cs @@ -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(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(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>[] 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); } diff --git a/src/modules/Elsa.Elasticsearch/Modules/Management/WorkflowInstanceStore.cs b/src/modules/Elsa.Elasticsearch/Modules/Management/WorkflowInstanceStore.cs index 20197b4a3..e5f345edc 100644 --- a/src/modules/Elsa.Elasticsearch/Modules/Management/WorkflowInstanceStore.cs +++ b/src/modules/Elsa.Elasticsearch/Modules/Management/WorkflowInstanceStore.cs @@ -35,16 +35,22 @@ public class ElasticWorkflowInstanceStore : IWorkflowInstanceStore } /// - public async ValueTask> FindManyAsync(WorkflowInstanceFilter filter, PageArgs pageArgs, CancellationToken cancellationToken = default) => - await _store.SearchAsync(d => Filter(d, filter), pageArgs, cancellationToken); + public async ValueTask> FindManyAsync(WorkflowInstanceFilter filter, PageArgs pageArgs, CancellationToken cancellationToken = default) + { + return await _store.SearchAsync(d => Filter(d, filter), pageArgs, cancellationToken); + } /// - public async ValueTask> FindManyAsync(WorkflowInstanceFilter filter, PageArgs pageArgs, WorkflowInstanceOrder order, CancellationToken cancellationToken = default) => - await _store.SearchAsync(d => Sort(Filter(d, filter), order), pageArgs, cancellationToken); + public async ValueTask> FindManyAsync(WorkflowInstanceFilter filter, PageArgs pageArgs, WorkflowInstanceOrder order, CancellationToken cancellationToken = default) + { + return await _store.SearchAsync(d => Sort(Filter(d, filter), order), pageArgs, cancellationToken); + } /// - public async ValueTask> FindManyAsync(WorkflowInstanceFilter filter, CancellationToken cancellationToken = default) => - await _store.SearchAsync(d => Filter(d, filter), cancellationToken); + public async ValueTask> FindManyAsync(WorkflowInstanceFilter filter, CancellationToken cancellationToken = default) + { + return await _store.SearchAsync(d => Filter(d, filter), cancellationToken); + } /// public async ValueTask> FindManyAsync(WorkflowInstanceFilter filter, WorkflowInstanceOrder order, CancellationToken cancellationToken = default) @@ -127,12 +133,21 @@ public class ElasticWorkflowInstanceStore : IWorkflowInstanceStore } /// - public async ValueTask SaveManyAsync(IEnumerable instances, CancellationToken cancellationToken = default) => + public async ValueTask SaveManyAsync(IEnumerable instances, CancellationToken cancellationToken = default) + { await _store.SaveManyAsync(instances, cancellationToken); + } /// - public async ValueTask DeleteAsync(WorkflowInstanceFilter filter, CancellationToken cancellationToken = default) => - await _store.DeleteByQueryAsync(d => Filter(d, filter), cancellationToken); + public async ValueTask 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 Sort(SearchRequestDescriptor descriptor, WorkflowInstanceOrder order) { @@ -145,9 +160,20 @@ public class ElasticWorkflowInstanceStore : IWorkflowInstanceStore return descriptor; } - private static SearchRequestDescriptor Filter(SearchRequestDescriptor descriptor, WorkflowInstanceFilter filter) => descriptor.Query(query => Filter(query, filter)); - private static CountRequestDescriptor Filter(CountRequestDescriptor descriptor, WorkflowInstanceFilter filter) => descriptor.Query(query => Filter(query, filter)); - private static DeleteByQueryRequestDescriptor Filter(DeleteByQueryRequestDescriptor descriptor, WorkflowInstanceFilter filter) => descriptor.Query(query => Filter(query, filter)); + private static SearchRequestDescriptor Filter(SearchRequestDescriptor descriptor, WorkflowInstanceFilter filter) + { + return descriptor.Query(query => Filter(query, filter)); + } + + private static CountRequestDescriptor Filter(CountRequestDescriptor descriptor, WorkflowInstanceFilter filter) + { + return descriptor.Query(query => Filter(query, filter)); + } + + private static DeleteByQueryRequestDescriptor Filter(DeleteByQueryRequestDescriptor descriptor, WorkflowInstanceFilter filter) + { + return descriptor.Query(query => Filter(query, filter)); + } private static QueryDescriptor Filter(QueryDescriptor descriptor, WorkflowInstanceFilter filter) { @@ -176,8 +202,9 @@ public class ElasticWorkflowInstanceStore : IWorkflowInstanceStore .Query(filter.SearchTerm)); } - private static SearchRequestDescriptor Summarize(SearchRequestDescriptor descriptor) => - descriptor.Fields( + private static SearchRequestDescriptor Summarize(SearchRequestDescriptor 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 SelectId(SearchRequestDescriptor descriptor) => - descriptor.Fields(field => field.Field(f => f.Id)); + private static SearchRequestDescriptor SelectId(SearchRequestDescriptor descriptor) + { + return descriptor.Fields(field => field.Field(f => f.Id)); + } } \ No newline at end of file diff --git a/src/modules/Elsa.EntityFrameworkCore.Common/Store.cs b/src/modules/Elsa.EntityFrameworkCore.Common/Store.cs index 3387e912a..0a353a20e 100644 --- a/src/modules/Elsa.EntityFrameworkCore.Common/Store.cs +++ b/src/modules/Elsa.EntityFrameworkCore.Common/Store.cs @@ -225,6 +225,24 @@ public class Store(IDbContextFactory dbContextF await dbContext.SaveChangesAsync(cancellationToken); } + /// + /// Updates specific properties of an entity in the database. + /// + /// The entity to update. + /// An array of expressions indicating the properties to update. + /// The cancellation token. + /// A task that represents the asynchronous operation. + public async Task UpdatePartialAsync(TEntity entity, Expression>[] 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); + } + /// /// Finds the entity matching the specified predicate. /// diff --git a/src/modules/Elsa.EntityFrameworkCore/Modules/Management/WorkflowInstanceStore.cs b/src/modules/Elsa.EntityFrameworkCore/Modules/Management/WorkflowInstanceStore.cs index 80075637f..f8295372d 100644 --- a/src/modules/Elsa.EntityFrameworkCore/Modules/Management/WorkflowInstanceStore.cs +++ b/src/modules/Elsa.EntityFrameworkCore/Modules/Management/WorkflowInstanceStore.cs @@ -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); + } + /// [RequiresUnreferencedCode("Calls Elsa.Workflows.Contracts.IWorkflowStateSerializer.SerializeAsync(WorkflowState, CancellationToken)")] public async ValueTask SaveAsync(WorkflowInstance instance, CancellationToken cancellationToken = default) diff --git a/src/modules/Elsa.MongoDb/Common/MongoDbStore.cs b/src/modules/Elsa.MongoDb/Common/MongoDbStore.cs index 9c5c3b674..e7135b2c9 100644 --- a/src/modules/Elsa.MongoDb/Common/MongoDbStore.cs +++ b/src/modules/Elsa.MongoDb/Common/MongoDbStore.cs @@ -133,6 +133,29 @@ public class MongoDbStore(IMongoCollection collection, ITe await collection.BulkWriteAsync(writes, cancellationToken: cancellationToken); } + + public async Task UpdatePartialAsync( + string id, + IDictionary 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.Filter.Eq(primaryKey, id); + var updateDefinition = Builders.Update.Combine( + updatedFields.Select(field => Builders.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}'."); + } /// /// Finds the document matching the specified predicate diff --git a/src/modules/Elsa.MongoDb/Modules/Management/WorkflowInstanceStore.cs b/src/modules/Elsa.MongoDb/Modules/Management/WorkflowInstanceStore.cs index 9f140caf4..6913ff35f 100644 --- a/src/modules/Elsa.MongoDb/Modules/Management/WorkflowInstanceStore.cs +++ b/src/modules/Elsa.MongoDb/Modules/Management/WorkflowInstanceStore.cs @@ -123,6 +123,16 @@ public class MongoWorkflowInstanceStore(MongoDbStore mongoDbSt return await mongoDbStore.DeleteWhereAsync(query => Filter(query, filter), x => x.Id, cancellationToken); } + public async Task UpdateUpdatedTimestampAsync(string workflowInstanceId, DateTimeOffset value, CancellationToken cancellationToken = default) + { + var props = new Dictionary + { + [nameof(WorkflowInstance.UpdatedAt)] = value + }; + + await mongoDbStore.UpdatePartialAsync(workflowInstanceId, props, cancellationToken: cancellationToken); + } + /// public async ValueTask SaveAsync(WorkflowInstance instance, CancellationToken cancellationToken = default) { diff --git a/src/modules/Elsa.Workflows.Management/Contracts/IWorkflowInstanceStore.cs b/src/modules/Elsa.Workflows.Management/Contracts/IWorkflowInstanceStore.cs index 1c1e5b05b..7ba1c3647 100644 --- a/src/modules/Elsa.Workflows.Management/Contracts/IWorkflowInstanceStore.cs +++ b/src/modules/Elsa.Workflows.Management/Contracts/IWorkflowInstanceStore.cs @@ -170,4 +170,12 @@ public interface IWorkflowInstanceStore /// The cancellation token. /// The number of deleted workflow instances. ValueTask DeleteAsync(WorkflowInstanceFilter filter, CancellationToken cancellationToken = default); + + /// + /// Updates the "LastUpdated" timestamp of a workflow instance. + /// + /// The unique identifier of the workflow instance. + /// The new timestamp value to set. + /// The cancellation token to observe during the operation. + Task UpdateUpdatedTimestampAsync(string workflowInstanceId, DateTimeOffset value, CancellationToken cancellationToken = default); } \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Management/Stores/MemoryWorkflowInstanceStore.cs b/src/modules/Elsa.Workflows.Management/Stores/MemoryWorkflowInstanceStore.cs index 0514e28ba..09125daed 100644 --- a/src/modules/Elsa.Workflows.Management/Stores/MemoryWorkflowInstanceStore.cs +++ b/src/modules/Elsa.Workflows.Management/Stores/MemoryWorkflowInstanceStore.cs @@ -154,6 +154,20 @@ public class MemoryWorkflowInstanceStore : IWorkflowInstanceStore return ValueTask.FromResult(count); } + /// + 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)")] diff --git a/src/modules/Elsa.Workflows.Runtime/Extensions/WorkflowExecutionPipelineBuilderExtensions.cs b/src/modules/Elsa.Workflows.Runtime/Extensions/WorkflowExecutionPipelineBuilderExtensions.cs index 9affb2058..709dfb5c1 100644 --- a/src/modules/Elsa.Workflows.Runtime/Extensions/WorkflowExecutionPipelineBuilderExtensions.cs +++ b/src/modules/Elsa.Workflows.Runtime/Extensions/WorkflowExecutionPipelineBuilderExtensions.cs @@ -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. /// public static IWorkflowExecutionPipelineBuilder UsePersistentVariables(this IWorkflowExecutionPipelineBuilder pipelineBuilder) => pipelineBuilder.UseMiddleware(); + + public static IWorkflowExecutionPipelineBuilder UseWorkflowHeartbeat(this IWorkflowExecutionPipelineBuilder pipelineBuilder) => pipelineBuilder.UseMiddleware(); /// /// Installs middleware that persists bookmarks after workflow execution. diff --git a/src/modules/Elsa.Workflows.Runtime/Middleware/Workflows/HeartbeatMiddleware.cs b/src/modules/Elsa.Workflows.Runtime/Middleware/Workflows/HeartbeatMiddleware.cs new file mode 100644 index 000000000..1a071cda1 --- /dev/null +++ b/src/modules/Elsa.Workflows.Runtime/Middleware/Workflows/HeartbeatMiddleware.cs @@ -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 options, ILogger 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(); + var clock = context.GetRequiredService(); + var workflowInstanceId = context.Id; + await workflowInstanceStore.UpdateUpdatedTimestampAsync(workflowInstanceId, clock.UtcNow, context.CancellationToken); + } + + private class WorkflowHeartbeat : IDisposable + { + private readonly Timer _timer; + private readonly Func _pulseAction; + private readonly TimeSpan _interval; + private readonly ILogger _logger; + + public WorkflowHeartbeat(Func pulseAction, TimeSpan interval, ILoggerFactory loggerFactory) + { + _pulseAction = pulseAction; + _interval = interval; + _logger = loggerFactory.CreateLogger(); + _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(); + } + } +} \ No newline at end of file