diff --git a/src/core/Elsa.Abstractions/Models/WorkflowInstance.cs b/src/core/Elsa.Abstractions/Models/WorkflowInstance.cs index 7121e5dd5..e0a6ec88e 100644 --- a/src/core/Elsa.Abstractions/Models/WorkflowInstance.cs +++ b/src/core/Elsa.Abstractions/Models/WorkflowInstance.cs @@ -18,8 +18,10 @@ namespace Elsa.Models } public string DefinitionId { get; set; } = default!; + public string DefinitionVersionId { get; set; } = default!; public string? TenantId { get; set; } public int Version { get; set; } + public WorkflowStatus WorkflowStatus { get; set; } public string CorrelationId { get; set; } = default!; public string? ContextType { get; set; } diff --git a/src/core/Elsa.Abstractions/Services/Models/IWorkflowBlueprint.cs b/src/core/Elsa.Abstractions/Services/Models/IWorkflowBlueprint.cs index 6ebb9d170..ec3cc8ae0 100644 --- a/src/core/Elsa.Abstractions/Services/Models/IWorkflowBlueprint.cs +++ b/src/core/Elsa.Abstractions/Services/Models/IWorkflowBlueprint.cs @@ -4,6 +4,7 @@ namespace Elsa.Services.Models { public interface IWorkflowBlueprint : ICompositeActivityBlueprint { + string VersionId { get; } int Version { get; } string? TenantId { get; } bool IsSingleton { get; } diff --git a/src/core/Elsa.Abstractions/Services/Models/WorkflowBlueprint.cs b/src/core/Elsa.Abstractions/Services/Models/WorkflowBlueprint.cs index 04def0cec..aa7f6e0ac 100644 --- a/src/core/Elsa.Abstractions/Services/Models/WorkflowBlueprint.cs +++ b/src/core/Elsa.Abstractions/Services/Models/WorkflowBlueprint.cs @@ -16,6 +16,7 @@ namespace Elsa.Services.Models public WorkflowBlueprint( string id, int version, + string versionId, string? tenantId, bool isSingleton, string? name, @@ -33,21 +34,22 @@ namespace Elsa.Services.Models IEnumerable activities, IEnumerable connections, IActivityPropertyProviders activityPropertyValueProviders) : base( - id, - default, - name, - displayName, - description, - id, - true, - false, - false, + id, + default, + name, + displayName, + description, + id, + true, + false, + false, new Dictionary(), default) { Id = id; Parent = this; Version = version; + VersionId = versionId; TenantId = tenantId; IsSingleton = isSingleton; IsLatest = isLatest; @@ -66,6 +68,7 @@ namespace Elsa.Services.Models } public int Version { get; set; } + public string VersionId { get; set; } public string? TenantId { get; set; } public bool IsSingleton { get; set; } public bool IsPublished { get; set; } @@ -77,7 +80,7 @@ namespace Elsa.Services.Models /// The channel, or queue, to place workflow instances of this workflow blueprint in. Channels can be used by the workflow dispatcher to prioritize workflows. /// public string? Channel { get; } - + public Variables Variables { get; set; } public WorkflowContextOptions? ContextOptions { get; set; } public WorkflowPersistenceBehavior PersistenceBehavior { get; set; } diff --git a/src/core/Elsa.Core/Builders/WorkflowBuilder.cs b/src/core/Elsa.Core/Builders/WorkflowBuilder.cs index b89d7a76f..2f34e8e68 100644 --- a/src/core/Elsa.Core/Builders/WorkflowBuilder.cs +++ b/src/core/Elsa.Core/Builders/WorkflowBuilder.cs @@ -24,7 +24,7 @@ namespace Elsa.Builders PropertyValueProviders = new Dictionary(); PersistenceBehavior = WorkflowPersistenceBehavior.WorkflowBurst; } - + public int Version { get; private set; } public bool IsLatest { get; private set; } public bool IsPublished { get; private set; } @@ -45,25 +45,25 @@ namespace Elsa.Builders ActivityId = value!; return this; } - + public IWorkflowBuilder ForTenantId(string? value) { TenantId = value; return this; } - + public new IWorkflowBuilder WithDisplayName(string value) { DisplayName = value; return this; } - + public new IWorkflowBuilder WithDescription(string? value) { Description = value; return this; } - + public IWorkflowBuilder WithChannel(string? value) { Channel = value; @@ -93,7 +93,7 @@ namespace Elsa.Builders IsPublished = isPublished; return this; } - + public IWorkflowBuilder WithTag(string value) { Tag = value; @@ -148,7 +148,7 @@ namespace Elsa.Builders Name ??= workflowTypeName; DisplayName ??= workflowTypeName; - + WithId(workflowTypeName); workflow.Build(this); return BuildBlueprint(activityIdPrefix); @@ -165,10 +165,11 @@ namespace Elsa.Builders public IWorkflowBlueprint BuildBlueprint(string activityIdPrefix) { var compositeRoot = Build(activityIdPrefix); - + return new WorkflowBlueprint( ActivityId, Version, + $"{ActivityId}v{Version}", TenantId, IsSingleton, Name, diff --git a/src/core/Elsa.Core/Decorators/CachingWorkflowRegistry.cs b/src/core/Elsa.Core/Decorators/CachingWorkflowRegistry.cs index 6d0e76850..ce21c3566 100644 --- a/src/core/Elsa.Core/Decorators/CachingWorkflowRegistry.cs +++ b/src/core/Elsa.Core/Decorators/CachingWorkflowRegistry.cs @@ -26,9 +26,9 @@ namespace Elsa.Decorators private readonly IWorkflowInstanceStore _workflowInstanceStore; public CachingWorkflowRegistry( - IWorkflowRegistry workflowRegistry, - IMemoryCache memoryCache, - ICacheSignal cacheSignal, + IWorkflowRegistry workflowRegistry, + IMemoryCache memoryCache, + ICacheSignal cacheSignal, IWorkflowInstanceStore workflowInstanceStore) { _workflowRegistry = workflowRegistry; @@ -63,25 +63,28 @@ namespace Elsa.Decorators return await _workflowRegistry.ListAsync(cancellationToken).ToList(); }); } - + private async IAsyncEnumerable ListActiveInternalAsync([EnumeratorCancellation] CancellationToken cancellationToken) { var workflows = await ListInternalAsync(cancellationToken); - - foreach (var workflow in workflows) - { - // If a workflow is not published, only consider it for processing if it has at least one non-ended workflow instance. - if (!workflow.IsPublished && !await WorkflowHasUnfinishedWorkflowsAsync(workflow, cancellationToken)) - continue; + var publishedWorkflows = workflows.Where(x => x.IsPublished); - yield return workflow; - } - } - - private async Task WorkflowHasUnfinishedWorkflowsAsync(IWorkflowBlueprint workflowBlueprint, CancellationToken cancellationToken) - { - var count = await _workflowInstanceStore.CountAsync(new UnfinishedWorkflowSpecification().WithWorkflowDefinition(workflowBlueprint.Id), cancellationToken); - return count > 0; + foreach (var publishedWorkflow in publishedWorkflows) + yield return publishedWorkflow; + + // We also need to consider unpublished workflows for inclusion in case they still have associated active workflow instances. + var unpublishedWorkflows = workflows.Where(x => !x.IsPublished).ToDictionary(x => x.VersionId); + var unpublishedWorkflowIds = unpublishedWorkflows.Keys; + + if (!unpublishedWorkflowIds.Any()) + yield break; + + var activeWorkflowInstances = await _workflowInstanceStore.FindManyAsync(new UnfinishedWorkflowSpecification().WithWorkflowDefinitionVersionIds(unpublishedWorkflowIds), cancellationToken: cancellationToken).ToList(); + var activeUnpublishedWorkflowVersionIds = activeWorkflowInstances.Select(x => x.DefinitionVersionId).Distinct().ToList(); + var activeUnpublishedWorkflowVersions = unpublishedWorkflows.Where(x => activeUnpublishedWorkflowVersionIds.Contains(x.Key)).Select(x => x.Value); + + foreach (var unpublishedWorkflow in activeUnpublishedWorkflowVersions) + yield return unpublishedWorkflow; } Task INotificationHandler.Handle(WorkflowDefinitionSaved notification, CancellationToken cancellationToken) diff --git a/src/core/Elsa.Core/Persistence/Specifications/WorkflowInstances/WorkflowDefinitionVersionIdsSpecification.cs b/src/core/Elsa.Core/Persistence/Specifications/WorkflowInstances/WorkflowDefinitionVersionIdsSpecification.cs new file mode 100644 index 000000000..f67b65851 --- /dev/null +++ b/src/core/Elsa.Core/Persistence/Specifications/WorkflowInstances/WorkflowDefinitionVersionIdsSpecification.cs @@ -0,0 +1,14 @@ +using System; +using System.Collections.Generic; +using System.Linq; +using System.Linq.Expressions; +using Elsa.Models; + +namespace Elsa.Persistence.Specifications.WorkflowInstances; + +public class WorkflowDefinitionVersionIdsSpecification : Specification +{ + public ICollection WorkflowDefinitionVersionIds { get; set; } + public WorkflowDefinitionVersionIdsSpecification(IEnumerable workflowDefinitionVersionIds) => WorkflowDefinitionVersionIds = workflowDefinitionVersionIds.ToList(); + public override Expression> ToExpression() => x => WorkflowDefinitionVersionIds.Contains(x.DefinitionVersionId); +} \ No newline at end of file diff --git a/src/core/Elsa.Core/Persistence/Specifications/WorkflowInstances/WorkflowInstanceSpecificationExtensions.cs b/src/core/Elsa.Core/Persistence/Specifications/WorkflowInstances/WorkflowInstanceSpecificationExtensions.cs index d358372ca..d4e557667 100644 --- a/src/core/Elsa.Core/Persistence/Specifications/WorkflowInstances/WorkflowInstanceSpecificationExtensions.cs +++ b/src/core/Elsa.Core/Persistence/Specifications/WorkflowInstances/WorkflowInstanceSpecificationExtensions.cs @@ -1,21 +1,33 @@ using System; +using System.Collections.Generic; using Elsa.Models; namespace Elsa.Persistence.Specifications.WorkflowInstances { public static class WorkflowInstanceSpecificationExtensions { - public static ISpecification WithWorkflowDefinition(this ISpecification specification, string workflowDefinitionId) => specification.And(new WorkflowDefinitionIdSpecification(workflowDefinitionId)); + public static ISpecification WithWorkflowDefinition(this ISpecification specification, string workflowDefinitionId) => + specification.And(new WorkflowDefinitionIdSpecification(workflowDefinitionId)); + + public static ISpecification WithWorkflowDefinitionVersionIds(this ISpecification specification, IEnumerable workflowDefinitionVersionIds) => + specification.And(new WorkflowDefinitionVersionIdsSpecification(workflowDefinitionVersionIds)); + public static ISpecification WithWorkflowName(this ISpecification specification, string name) => specification.And(new WorkflowInstanceNameMatchSpecification(name)); public static ISpecification WithCorrelationId(this ISpecification specification, string correlationId) => specification.And(new CorrelationIdSpecification(correlationId)); - public static ISpecification WithContextId(this ISpecification specification, string contextType, string contextId) => specification.And(new WorkflowInstanceContextIdMatchSpecification(contextType, contextId)); - public static ISpecification WithContextId(this ISpecification specification, Type contextType, string contextId) => specification.And(new WorkflowInstanceContextIdMatchSpecification(contextType, contextId)); - public static ISpecification WithContextId(this ISpecification specification, string contextId) => specification.And(new WorkflowInstanceContextIdMatchSpecification(typeof(TContextType), contextId)); + + public static ISpecification WithContextId(this ISpecification specification, string contextType, string contextId) => + specification.And(new WorkflowInstanceContextIdMatchSpecification(contextType, contextId)); + + public static ISpecification WithContextId(this ISpecification specification, Type contextType, string contextId) => + specification.And(new WorkflowInstanceContextIdMatchSpecification(contextType, contextId)); + + public static ISpecification WithContextId(this ISpecification specification, string contextId) => + specification.And(new WorkflowInstanceContextIdMatchSpecification(typeof(TContextType), contextId)); + public static ISpecification WithContext(this ISpecification specification, string contextType) => specification.And(new WorkflowInstanceContextMatchSpecification(contextType)); public static ISpecification WithContext(this ISpecification specification, Type contextType) => specification.And(new WorkflowInstanceContextMatchSpecification(contextType)); public static ISpecification WithContextId(this ISpecification specification) => specification.And(new WorkflowInstanceContextMatchSpecification(typeof(TContextType))); public static ISpecification WithStatus(this ISpecification specification, WorkflowStatus status) => specification.And(new WorkflowStatusSpecification(status)); public static ISpecification WithSearchTerm(this ISpecification specification, string searchTerm) => specification.And(new WorkflowSearchTermSpecification(searchTerm)); - } -} +} \ No newline at end of file diff --git a/src/core/Elsa.Core/Services/Workflows/WorkflowBlueprintMaterializer.cs b/src/core/Elsa.Core/Services/Workflows/WorkflowBlueprintMaterializer.cs index 91215f340..7ff586c13 100644 --- a/src/core/Elsa.Core/Services/Workflows/WorkflowBlueprintMaterializer.cs +++ b/src/core/Elsa.Core/Services/Workflows/WorkflowBlueprintMaterializer.cs @@ -46,6 +46,7 @@ namespace Elsa.Services.Workflows var workflowBlueprint = new WorkflowBlueprint( workflowDefinition.DefinitionId, workflowDefinition.Version, + workflowDefinition.VersionId, workflowDefinition.TenantId, workflowDefinition.IsSingleton, workflowDefinition.Name, diff --git a/src/core/Elsa.Core/Services/Workflows/WorkflowFactory.cs b/src/core/Elsa.Core/Services/Workflows/WorkflowFactory.cs index 49dd0ac70..b18313c01 100644 --- a/src/core/Elsa.Core/Services/Workflows/WorkflowFactory.cs +++ b/src/core/Elsa.Core/Services/Workflows/WorkflowFactory.cs @@ -31,6 +31,7 @@ namespace Elsa.Services.Workflows DefinitionId = workflowBlueprint.Id, TenantId = tenantId ?? workflowBlueprint.TenantId, Version = workflowBlueprint.Version, + DefinitionVersionId = workflowBlueprint.VersionId, WorkflowStatus = WorkflowStatus.Idle, CorrelationId = correlationId ?? Guid.NewGuid().ToString("N"), ContextId = contextId, diff --git a/src/core/Elsa.Core/Services/Workflows/WorkflowRegistry.cs b/src/core/Elsa.Core/Services/Workflows/WorkflowRegistry.cs index f8be793bd..a4c7d2b94 100644 --- a/src/core/Elsa.Core/Services/Workflows/WorkflowRegistry.cs +++ b/src/core/Elsa.Core/Services/Workflows/WorkflowRegistry.cs @@ -46,16 +46,25 @@ namespace Elsa.Services.Workflows private async IAsyncEnumerable ListActiveInternalAsync([EnumeratorCancellation] CancellationToken cancellationToken) { - var workflows = await ListAsync(cancellationToken); + var workflows = await ListInternalAsync(cancellationToken).ToListAsync(cancellationToken); + var publishedWorkflows = workflows.Where(x => x.IsPublished); - foreach (var workflow in workflows) - { - // If a workflow is not published, only consider it for processing if it has at least one non-ended workflow instance. - if (!workflow.IsPublished && !await WorkflowHasUnfinishedWorkflowsAsync(workflow, cancellationToken)) - continue; + foreach (var publishedWorkflow in publishedWorkflows) + yield return publishedWorkflow; - yield return workflow; - } + // We also need to consider unpublished workflows for inclusion in case they still have associated active workflow instances. + var unpublishedWorkflows = workflows.Where(x => !x.IsPublished).ToDictionary(x => x.VersionId); + var unpublishedWorkflowIds = unpublishedWorkflows.Keys; + + if (!unpublishedWorkflowIds.Any()) + yield break; + + var activeWorkflowInstances = await _workflowInstanceStore.FindManyAsync(new UnfinishedWorkflowSpecification().WithWorkflowDefinitionVersionIds(unpublishedWorkflowIds), cancellationToken: cancellationToken).ToList(); + var activeUnpublishedWorkflowVersionIds = activeWorkflowInstances.Select(x => x.DefinitionVersionId).Distinct().ToList(); + var activeUnpublishedWorkflowVersions = unpublishedWorkflows.Where(x => activeUnpublishedWorkflowVersionIds.Contains(x.Key)).Select(x => x.Value); + + foreach (var unpublishedWorkflow in activeUnpublishedWorkflowVersions) + yield return unpublishedWorkflow; } private async IAsyncEnumerable ListInternalAsync([EnumeratorCancellation] CancellationToken cancellationToken) diff --git a/src/samples/server/Elsa.Samples.Server.Host/SampleActivity.cs b/src/samples/server/Elsa.Samples.Server.Host/SampleActivity.cs new file mode 100644 index 000000000..0cc68a9e0 --- /dev/null +++ b/src/samples/server/Elsa.Samples.Server.Host/SampleActivity.cs @@ -0,0 +1,11 @@ +using System.Collections.Generic; +using Elsa.Attributes; +using Elsa.Expressions; +using Elsa.Services; + +namespace Elsa.Samples.Server.Host; + +public class SampleActivity : Activity +{ + [ActivityInput(SupportedSyntaxes = new[]{SyntaxNames.JavaScript, SyntaxNames.Liquid})] public ICollection MyList { get; set; } +} \ No newline at end of file