From 9360dfeeb8ee36cc93a967df39c24739ba36bd14 Mon Sep 17 00:00:00 2001 From: mmiscevic <47240710+mmiscevic@users.noreply.github.com> Date: Wed, 26 Feb 2020 18:49:31 +0100 Subject: [PATCH] Added Tag property to BlockingActivity.cs and BlockingActivityEntity.cs; (#254) Implements #241 --- .../Models/BlockingActivity.cs | 3 +- .../Persistence/IWorkflowInstanceStore.cs | 3 +- .../Memory/MemoryWorkflowInstanceStore.cs | 21 +++++++++++++- .../Services/CosmosDbWorkflowInstanceStore.cs | 27 +++++++++++++++-- .../Entities/BlockingActivityEntity.cs | 1 + ...ntityFrameworkCoreWorkflowInstanceStore.cs | 29 +++++++++++++++++-- .../Services/MongoWorkflowInstanceStore.cs | 22 +++++++++++++- .../Indexes/WorkflowInstanceIndex.cs | 2 ++ .../Services/YesSqlWorkflowInstanceStore.cs | 21 ++++++++++++++ .../StartupTasks/InitializeStoreTask.cs | 1 + 10 files changed, 121 insertions(+), 9 deletions(-) diff --git a/src/core/Elsa.Abstractions/Models/BlockingActivity.cs b/src/core/Elsa.Abstractions/Models/BlockingActivity.cs index 17d55880e..6d0fdf997 100644 --- a/src/core/Elsa.Abstractions/Models/BlockingActivity.cs +++ b/src/core/Elsa.Abstractions/Models/BlockingActivity.cs @@ -11,8 +11,9 @@ namespace Elsa.Models ActivityId = activityId; ActivityType = activityType; } - + public string? ActivityId { get; set; } public string? ActivityType { get; set; } + public string? Tag { get; set; } } } \ No newline at end of file diff --git a/src/core/Elsa.Abstractions/Persistence/IWorkflowInstanceStore.cs b/src/core/Elsa.Abstractions/Persistence/IWorkflowInstanceStore.cs index 33104ffe6..72a9b6b27 100644 --- a/src/core/Elsa.Abstractions/Persistence/IWorkflowInstanceStore.cs +++ b/src/core/Elsa.Abstractions/Persistence/IWorkflowInstanceStore.cs @@ -6,12 +6,13 @@ using Elsa.Models; namespace Elsa.Persistence { public interface IWorkflowInstanceStore - { + { Task SaveAsync(WorkflowInstance instance, CancellationToken cancellationToken = default); Task GetByIdAsync(string id, CancellationToken cancellationToken = default); Task GetByCorrelationIdAsync(string correlationId, CancellationToken cancellationToken = default); Task> ListByDefinitionAsync(string definitionId, CancellationToken cancellationToken = default); Task> ListAllAsync(CancellationToken cancellationToken = default); + Task> ListByBlockingActivityTagAsync(string activityType, string tag, string? correlationId = default, CancellationToken cancellationToken = default); Task> ListByBlockingActivityAsync(string activityType, string? correlationId = default, CancellationToken cancellationToken = default); Task> ListByStatusAsync(string definitionId, WorkflowStatus status, CancellationToken cancellationToken = default); Task> ListByStatusAsync(WorkflowStatus status, CancellationToken cancellationToken = default); diff --git a/src/core/Elsa.Core/Persistence/Memory/MemoryWorkflowInstanceStore.cs b/src/core/Elsa.Core/Persistence/Memory/MemoryWorkflowInstanceStore.cs index c418db817..c4f96d189 100644 --- a/src/core/Elsa.Core/Persistence/Memory/MemoryWorkflowInstanceStore.cs +++ b/src/core/Elsa.Core/Persistence/Memory/MemoryWorkflowInstanceStore.cs @@ -45,10 +45,29 @@ namespace Elsa.Persistence.Memory var workflows = workflowInstances.Values.AsEnumerable(); return Task.FromResult(workflows); } + public Task> ListByBlockingActivityTagAsync( + string activityType, + string tag, + string? correlationId = null, + CancellationToken cancellationToken = default) + { + var query = workflowInstances.Values.AsQueryable(); + + query = query.Where(x => x.Status == WorkflowStatus.Suspended); + + if (!string.IsNullOrWhiteSpace(correlationId)) + query = query.Where(x => x.CorrelationId == correlationId); + + query = query.Where( + x => x.BlockingActivities.Any(y => y.ActivityType == activityType && y.Tag == tag) + ); + + return Task.FromResult(query.AsEnumerable().GetBlockingActivities()); + } public Task> ListByBlockingActivityAsync( string activityType, - string? correlationId = default, + string? correlationId = default, CancellationToken cancellationToken = default) { var query = workflowInstances.Values.AsQueryable(); diff --git a/src/providers/Elsa.Persistence.DocumentDb/Services/CosmosDbWorkflowInstanceStore.cs b/src/providers/Elsa.Persistence.DocumentDb/Services/CosmosDbWorkflowInstanceStore.cs index fe43d68af..4aa641b45 100644 --- a/src/providers/Elsa.Persistence.DocumentDb/Services/CosmosDbWorkflowInstanceStore.cs +++ b/src/providers/Elsa.Persistence.DocumentDb/Services/CosmosDbWorkflowInstanceStore.cs @@ -65,6 +65,27 @@ namespace Elsa.Persistence.DocumentDb.Services .OrderByDescending(x => x.CreatedAt); return mapper.Map>(query); } + public async Task> ListByBlockingActivityTagAsync( + string activityType, + string tag, + string? correlationId = null, + CancellationToken cancellationToken = default) + { + var client = storage.Client; + var collectionUrl = await GetCollectionUriAsync(cancellationToken); + var query = client + .CreateDocumentQuery(collectionUrl) + .Where(x => x.Status == WorkflowStatus.Suspended); + + if (!string.IsNullOrWhiteSpace(correlationId)) + query = query.Where(x => x.CorrelationId == correlationId); + + query = query.Where(x => x.BlockingActivities.Any(y => y.ActivityType == activityType && y.Tag == tag)); + query = query.OrderByDescending(x => x.CreatedAt); + + var instances = Map(query.ToList()); + return instances.GetBlockingActivities(activityType); + } public async Task> ListByBlockingActivityAsync( string activityType, @@ -77,7 +98,7 @@ namespace Elsa.Persistence.DocumentDb.Services .CreateDocumentQuery(collectionUrl) .Where(x => x.Status == WorkflowStatus.Suspended); - if (!string.IsNullOrWhiteSpace(correlationId)) + if (!string.IsNullOrWhiteSpace(correlationId)) query = query.Where(x => x.CorrelationId == correlationId); query = query.Where(x => x.BlockingActivities.Any(y => y.ActivityType == activityType)); @@ -135,10 +156,10 @@ namespace Elsa.Persistence.DocumentDb.Services document = (dynamic)response.Resource; return Map(document); } - + private async Task GetCollectionUriAsync(CancellationToken cancellationToken) { - if (collectionUrl == null) + if (collectionUrl == null) collectionUrl = await storage.GetCollectionAsync("WorkflowInstances", cancellationToken); return collectionUrl; diff --git a/src/providers/Elsa.Persistence.EntityFrameworkCore/Entities/BlockingActivityEntity.cs b/src/providers/Elsa.Persistence.EntityFrameworkCore/Entities/BlockingActivityEntity.cs index 557be9282..3a8960d91 100644 --- a/src/providers/Elsa.Persistence.EntityFrameworkCore/Entities/BlockingActivityEntity.cs +++ b/src/providers/Elsa.Persistence.EntityFrameworkCore/Entities/BlockingActivityEntity.cs @@ -6,5 +6,6 @@ namespace Elsa.Persistence.EntityFrameworkCore.Entities public WorkflowInstanceEntity WorkflowInstance { get; set; } public string ActivityId { get; set; } public string ActivityType { get; set; } + public string? Tag { get; set; } } } \ No newline at end of file diff --git a/src/providers/Elsa.Persistence.EntityFrameworkCore/Services/EntityFrameworkCoreWorkflowInstanceStore.cs b/src/providers/Elsa.Persistence.EntityFrameworkCore/Services/EntityFrameworkCoreWorkflowInstanceStore.cs index 1c5984f76..afddf638b 100644 --- a/src/providers/Elsa.Persistence.EntityFrameworkCore/Services/EntityFrameworkCoreWorkflowInstanceStore.cs +++ b/src/providers/Elsa.Persistence.EntityFrameworkCore/Services/EntityFrameworkCoreWorkflowInstanceStore.cs @@ -107,6 +107,31 @@ namespace Elsa.Persistence.EntityFrameworkCore.Services .ToListAsync(cancellationToken); return Map(documents); } + public async Task> ListByBlockingActivityTagAsync( + string activityType, + string tag, + string correlationId = default, + CancellationToken cancellationToken = default) + { + var query = dbContext + .WorkflowInstances + .Include(x => x.Activities) + .Include(x => x.BlockingActivities) + .AsQueryable(); + + query = query.Where(x => x.Status == WorkflowStatus.Suspended); + + if (!string.IsNullOrWhiteSpace(correlationId)) + query = query.Where(x => x.CorrelationId == correlationId); + + query = query.Where(x => x.BlockingActivities.Any(y => y.ActivityType == activityType && y.Tag == tag)); + query = query.OrderByDescending(x => x.CreatedAt); + + var documents = await query.ToListAsync(cancellationToken); + var instances = Map(documents); + + return instances.GetBlockingActivities(activityType); + } public async Task> ListByBlockingActivityAsync( string activityType, @@ -176,11 +201,11 @@ namespace Elsa.Persistence.EntityFrameworkCore.Services var blockingActivityRecords = await dbContext.BlockingActivities .Where(x => x.WorkflowInstance.InstanceId == id) .ToListAsync(cancellationToken); - + dbContext.ActivityInstances.RemoveRange(activityInstanceRecords); dbContext.BlockingActivities.RemoveRange(blockingActivityRecords); dbContext.WorkflowInstances.Remove(record); - + await dbContext.SaveChangesAsync(cancellationToken); } diff --git a/src/providers/Elsa.Persistence.MongoDb/Services/MongoWorkflowInstanceStore.cs b/src/providers/Elsa.Persistence.MongoDb/Services/MongoWorkflowInstanceStore.cs index c438cb3c9..4435ef74b 100644 --- a/src/providers/Elsa.Persistence.MongoDb/Services/MongoWorkflowInstanceStore.cs +++ b/src/providers/Elsa.Persistence.MongoDb/Services/MongoWorkflowInstanceStore.cs @@ -63,6 +63,26 @@ namespace Elsa.Persistence.MongoDb.Services .ToListAsync(cancellationToken); } + public async Task> ListByBlockingActivityTagAsync( + string activityType, + string tag, + string correlationId = null, + CancellationToken cancellationToken = default) + { + var query = collection.AsQueryable(); + + query = query.Where(x => x.Status == WorkflowStatus.Suspended); + + if (!string.IsNullOrWhiteSpace(correlationId)) + query = query.Where(x => x.CorrelationId == correlationId); + + query = query.Where(x => x.BlockingActivities.Any(y => y.ActivityType == activityType && y.Tag == tag)); + query = query.OrderByDescending(x => x.CreatedAt); + + var instances = await query.ToListAsync(cancellationToken); + + return instances.GetBlockingActivities(activityType); + } public async Task> ListByBlockingActivityAsync( string activityType, string correlationId = default, @@ -77,7 +97,7 @@ namespace Elsa.Persistence.MongoDb.Services query = query.Where(x => x.BlockingActivities.Any(y => y.ActivityType == activityType)); query = query.OrderByDescending(x => x.CreatedAt); - + var instances = await query.ToListAsync(cancellationToken); return instances.GetBlockingActivities(activityType); diff --git a/src/providers/Elsa.Persistence.YesSql/Indexes/WorkflowInstanceIndex.cs b/src/providers/Elsa.Persistence.YesSql/Indexes/WorkflowInstanceIndex.cs index 4cc541af7..5930537ac 100644 --- a/src/providers/Elsa.Persistence.YesSql/Indexes/WorkflowInstanceIndex.cs +++ b/src/providers/Elsa.Persistence.YesSql/Indexes/WorkflowInstanceIndex.cs @@ -19,6 +19,7 @@ namespace Elsa.Persistence.YesSql.Indexes { public string ActivityId { get; set; } public string ActivityType { get; set; } + public string Tag { get; set; } public string CorrelationId { get; set; } public WorkflowStatus ProcessStatus { get; set; } public DateTime CreatedAt { get; set; } @@ -47,6 +48,7 @@ namespace Elsa.Persistence.YesSql.Indexes { ActivityId = activity.ActivityId, ActivityType = activity.ActivityType, + Tag = activity.Tag, CorrelationId = workflowInstance.CorrelationId, ProcessStatus = workflowInstance.Status, CreatedAt = workflowInstance.CreatedAt diff --git a/src/providers/Elsa.Persistence.YesSql/Services/YesSqlWorkflowInstanceStore.cs b/src/providers/Elsa.Persistence.YesSql/Services/YesSqlWorkflowInstanceStore.cs index c9a692a29..f7aa6ac1c 100644 --- a/src/providers/Elsa.Persistence.YesSql/Services/YesSqlWorkflowInstanceStore.cs +++ b/src/providers/Elsa.Persistence.YesSql/Services/YesSqlWorkflowInstanceStore.cs @@ -78,6 +78,27 @@ namespace Elsa.Persistence.YesSql.Services .ListAsync(); return mapper.Map>(documents); } + public async Task> ListByBlockingActivityTagAsync( + string activityType, + string tag, + string correlationId = null, + CancellationToken cancellationToken = default) + { + var query = session.Query(); + + query = query.Where(x => x.ProcessStatus == WorkflowStatus.Suspended); + + if (!string.IsNullOrWhiteSpace(correlationId)) + query = query.Where(x => x.CorrelationId == correlationId); + + query = query.Where(x => x.ActivityType == activityType && x.Tag == tag); + query = query.OrderByDescending(x => x.CreatedAt); + + var documents = await query.ListAsync(); + var instances = mapper.Map>(documents); + + return instances.GetBlockingActivities(activityType); + } public async Task> ListByBlockingActivityAsync( string activityType, diff --git a/src/providers/Elsa.Persistence.YesSql/StartupTasks/InitializeStoreTask.cs b/src/providers/Elsa.Persistence.YesSql/StartupTasks/InitializeStoreTask.cs index 32cd32a43..46fcbc8e0 100644 --- a/src/providers/Elsa.Persistence.YesSql/StartupTasks/InitializeStoreTask.cs +++ b/src/providers/Elsa.Persistence.YesSql/StartupTasks/InitializeStoreTask.cs @@ -67,6 +67,7 @@ namespace Elsa.Persistence.YesSql.StartupTasks .CreateMapIndexTable(nameof(WorkflowInstanceBlockingActivitiesIndex), table => table .Column("ActivityId") .Column("ActivityType") + .Column("Tag") .Column("CorrelationId") .Column("WorkflowStatus") .Column("CreatedAt")