diff --git a/src/activities/Elsa.Activities.Entity/Bookmarks/EntityChangedBookmark.cs b/src/activities/Elsa.Activities.Entity/Bookmarks/EntityChangedBookmark.cs index 3aa7556b1..2bb953592 100644 --- a/src/activities/Elsa.Activities.Entity/Bookmarks/EntityChangedBookmark.cs +++ b/src/activities/Elsa.Activities.Entity/Bookmarks/EntityChangedBookmark.cs @@ -2,34 +2,34 @@ using System.Collections.Generic; using System.Threading; using System.Threading.Tasks; using Elsa.Services; -using Elsa.Services.Bookmarks; namespace Elsa.Activities.Entity.Bookmarks { public class EntityChangedBookmark : IBookmark { - public EntityChangedBookmark(string? entityName, EntityChangedAction? action, string? contextId) + public EntityChangedBookmark(string? entityName, EntityChangedAction? action) { EntityName = entityName; Action = action; - ContextId = contextId; } public string? EntityName { get; } public EntityChangedAction? Action { get; } - public string? ContextId { get; } } public class EntityChangedWorkflowTriggerProvider : BookmarkProvider { - public override async ValueTask> GetBookmarksAsync(BookmarkProviderContext context, CancellationToken cancellationToken) => - new[] - { - Result(new EntityChangedBookmark( - await context.ReadActivityPropertyAsync(x => x.EntityName, cancellationToken), - await context.ReadActivityPropertyAsync(x => x.Action, cancellationToken), - context.ActivityExecutionContext.WorkflowExecutionContext.WorkflowInstance.ContextId - )) - }; + public override async ValueTask> GetBookmarksAsync(BookmarkProviderContext context, CancellationToken cancellationToken) + { + var entityName = await context.ReadActivityPropertyAsync(x => x.EntityName, cancellationToken); + var action = await context.ReadActivityPropertyAsync(x => x.Action, cancellationToken); + + var bookmark = new EntityChangedBookmark( + entityName, + action + ); + + return new[] { Result(bookmark) }; + } } } \ No newline at end of file diff --git a/src/activities/Elsa.Activities.Entity/Extensions/WorkflowRunnerExtensions.cs b/src/activities/Elsa.Activities.Entity/Extensions/WorkflowRunnerExtensions.cs index 79eb382b0..edd15fba6 100644 --- a/src/activities/Elsa.Activities.Entity/Extensions/WorkflowRunnerExtensions.cs +++ b/src/activities/Elsa.Activities.Entity/Extensions/WorkflowRunnerExtensions.cs @@ -4,16 +4,14 @@ using Elsa.Activities.Entity.Bookmarks; using Elsa.Activities.Entity.Models; using Elsa.Models; using Elsa.Services; +using Elsa.Services.Models; namespace Elsa.Activities.Entity.Extensions { public static class WorkflowRunnerExtensions { - // TODO: Design multi-tenancy. - private const string? TenantId = default; - public static async Task TriggerEntityChangedWorkflowsAsync( - this IWorkflowDispatcher workflowDispatcher, + this IWorkflowLaunchpad workflowLaunchpad, string entityId, string entityName, EntityChangedAction changedAction, @@ -21,16 +19,39 @@ namespace Elsa.Activities.Entity.Extensions string? contextId = default, CancellationToken cancellationToken = default) { - const string activityType = nameof(EntityChanged); var input = new EntityChangedContext(entityId, entityName, changedAction); + var query = QueryEntityChangedWorkflowsAsync(entityName, changedAction, correlationId, contextId); + await workflowLaunchpad.CollectAndExecuteWorkflowsAsync(query, new WorkflowInput(input), cancellationToken); + } + + public static async Task DispatchEntityChangedWorkflowsAsync( + this IWorkflowLaunchpad workflowLaunchpad, + string entityId, + string entityName, + EntityChangedAction changedAction, + string? correlationId = default, + string? contextId = default, + CancellationToken cancellationToken = default) + { + var input = new EntityChangedContext(entityId, entityName, changedAction); + var query = QueryEntityChangedWorkflowsAsync(entityName, changedAction, correlationId, contextId); + await workflowLaunchpad.CollectAndExecuteWorkflowsAsync(query, new WorkflowInput(input), cancellationToken); + } + + private static WorkflowsQuery QueryEntityChangedWorkflowsAsync( + string entityName, + EntityChangedAction changedAction, + string? correlationId = default, + string? contextId = default) + { + const string activityType = nameof(EntityChanged); var bookmark = new EntityChangedBookmark( entityName, - changedAction, - contextId + changedAction ); - await workflowDispatcher.DispatchAsync(new TriggerWorkflowsRequest(activityType, bookmark, new WorkflowInput(input), correlationId, default, contextId, TenantId), cancellationToken); + return new WorkflowsQuery(activityType, bookmark, correlationId, ContextId: contextId); } } } \ No newline at end of file diff --git a/src/core/Elsa.Core/StartupTasks/IndexTriggers.cs b/src/core/Elsa.Core/StartupTasks/IndexTriggers.cs index 46c0fdb5b..b452d6c55 100644 --- a/src/core/Elsa.Core/StartupTasks/IndexTriggers.cs +++ b/src/core/Elsa.Core/StartupTasks/IndexTriggers.cs @@ -1,7 +1,6 @@ using System.Threading; using System.Threading.Tasks; using Elsa.Services; -using Elsa.Services.Triggers; namespace Elsa.StartupTasks { diff --git a/src/samples/console/Elsa.Samples.EntityChanged/SomeRepository.cs b/src/samples/console/Elsa.Samples.EntityChanged/SomeRepository.cs index 231f5ca64..64a5c88bd 100644 --- a/src/samples/console/Elsa.Samples.EntityChanged/SomeRepository.cs +++ b/src/samples/console/Elsa.Samples.EntityChanged/SomeRepository.cs @@ -9,10 +9,10 @@ namespace Elsa.Samples.EntityChanged { public class SomeRepository { - private readonly IWorkflowDispatcher _workflowRunner; + private readonly IWorkflowLaunchpad _workflowLaunchpad; private readonly ICollection _collection = new List(); - public SomeRepository(IWorkflowDispatcher workflowRunner) => _workflowRunner = workflowRunner; + public SomeRepository(IWorkflowLaunchpad workflowLaunchpad) => _workflowLaunchpad = workflowLaunchpad; public Task AddAsync(Entity entity) { _collection.Add(entity); @@ -28,7 +28,7 @@ namespace Elsa.Samples.EntityChanged public Task GetAsync(string id) => Task.FromResult(_collection.FirstOrDefault(x => x.Id == id)); private async Task TriggerWorkflowsAsync(Entity entity, EntityChangedAction changedAction) => - await _workflowRunner.TriggerEntityChangedWorkflowsAsync( + await _workflowLaunchpad.TriggerEntityChangedWorkflowsAsync( entity.Id, entity.GetType().GetEntityName(), changedAction,