Added Tag property to BlockingActivity.cs and BlockingActivityEntity.cs; (#254)

Implements #241
This commit is contained in:
mmiscevic 2020-02-26 18:49:31 +01:00 committed by GitHub
parent 299d465570
commit 9360dfeeb8
No known key found for this signature in database
GPG key ID: 4AEE18F83AFDEB23
10 changed files with 121 additions and 9 deletions

View file

@ -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; }
}
}

View file

@ -6,12 +6,13 @@ using Elsa.Models;
namespace Elsa.Persistence
{
public interface IWorkflowInstanceStore
{
{
Task<WorkflowInstance> SaveAsync(WorkflowInstance instance, CancellationToken cancellationToken = default);
Task<WorkflowInstance> GetByIdAsync(string id, CancellationToken cancellationToken = default);
Task<WorkflowInstance> GetByCorrelationIdAsync(string correlationId, CancellationToken cancellationToken = default);
Task<IEnumerable<WorkflowInstance>> ListByDefinitionAsync(string definitionId, CancellationToken cancellationToken = default);
Task<IEnumerable<WorkflowInstance>> ListAllAsync(CancellationToken cancellationToken = default);
Task<IEnumerable<(WorkflowInstance WorkflowInstance, BlockingActivity BlockingActivity)>> ListByBlockingActivityTagAsync(string activityType, string tag, string? correlationId = default, CancellationToken cancellationToken = default);
Task<IEnumerable<(WorkflowInstance WorkflowInstance, BlockingActivity BlockingActivity)>> ListByBlockingActivityAsync(string activityType, string? correlationId = default, CancellationToken cancellationToken = default);
Task<IEnumerable<WorkflowInstance>> ListByStatusAsync(string definitionId, WorkflowStatus status, CancellationToken cancellationToken = default);
Task<IEnumerable<WorkflowInstance>> ListByStatusAsync(WorkflowStatus status, CancellationToken cancellationToken = default);

View file

@ -45,10 +45,29 @@ namespace Elsa.Persistence.Memory
var workflows = workflowInstances.Values.AsEnumerable();
return Task.FromResult(workflows);
}
public Task<IEnumerable<(WorkflowInstance, BlockingActivity)>> 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<IEnumerable<(WorkflowInstance, BlockingActivity)>> ListByBlockingActivityAsync(
string activityType,
string? correlationId = default,
string? correlationId = default,
CancellationToken cancellationToken = default)
{
var query = workflowInstances.Values.AsQueryable();

View file

@ -65,6 +65,27 @@ namespace Elsa.Persistence.DocumentDb.Services
.OrderByDescending(x => x.CreatedAt);
return mapper.Map<IEnumerable<WorkflowInstance>>(query);
}
public async Task<IEnumerable<(WorkflowInstance, BlockingActivity)>> ListByBlockingActivityTagAsync(
string activityType,
string tag,
string? correlationId = null,
CancellationToken cancellationToken = default)
{
var client = storage.Client;
var collectionUrl = await GetCollectionUriAsync(cancellationToken);
var query = client
.CreateDocumentQuery<WorkflowInstanceDocument>(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<IEnumerable<(WorkflowInstance, BlockingActivity)>> ListByBlockingActivityAsync(
string activityType,
@ -77,7 +98,7 @@ namespace Elsa.Persistence.DocumentDb.Services
.CreateDocumentQuery<WorkflowInstanceDocument>(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<Uri> GetCollectionUriAsync(CancellationToken cancellationToken)
{
if (collectionUrl == null)
if (collectionUrl == null)
collectionUrl = await storage.GetCollectionAsync("WorkflowInstances", cancellationToken);
return collectionUrl;

View file

@ -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; }
}
}

View file

@ -107,6 +107,31 @@ namespace Elsa.Persistence.EntityFrameworkCore.Services
.ToListAsync(cancellationToken);
return Map(documents);
}
public async Task<IEnumerable<(WorkflowInstance, BlockingActivity)>> 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<IEnumerable<(WorkflowInstance, BlockingActivity)>> 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);
}

View file

@ -63,6 +63,26 @@ namespace Elsa.Persistence.MongoDb.Services
.ToListAsync(cancellationToken);
}
public async Task<IEnumerable<(WorkflowInstance, BlockingActivity)>> 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<IEnumerable<(WorkflowInstance, BlockingActivity)>> 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);

View file

@ -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

View file

@ -78,6 +78,27 @@ namespace Elsa.Persistence.YesSql.Services
.ListAsync();
return mapper.Map<IEnumerable<WorkflowInstance>>(documents);
}
public async Task<IEnumerable<(WorkflowInstance, BlockingActivity)>> ListByBlockingActivityTagAsync(
string activityType,
string tag,
string correlationId = null,
CancellationToken cancellationToken = default)
{
var query = session.Query<WorkflowInstanceDocument, WorkflowInstanceBlockingActivitiesIndex>();
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<IEnumerable<WorkflowInstance>>(documents);
return instances.GetBlockingActivities(activityType);
}
public async Task<IEnumerable<(WorkflowInstance, BlockingActivity)>> ListByBlockingActivityAsync(
string activityType,

View file

@ -67,6 +67,7 @@ namespace Elsa.Persistence.YesSql.StartupTasks
.CreateMapIndexTable(nameof(WorkflowInstanceBlockingActivitiesIndex), table => table
.Column<string>("ActivityId")
.Column<string>("ActivityType")
.Column<string>("Tag")
.Column<string>("CorrelationId")
.Column<string>("WorkflowStatus")
.Column<DateTime>("CreatedAt")