Add caching layer to WorkflowRegistry
This commit is contained in:
parent
386879781d
commit
420e2e65c5
|
|
@ -40,7 +40,7 @@ namespace Elsa.Activities.AzureServiceBus.StartupTasks
|
|||
private async IAsyncEnumerable<string> GetQueueNamesAsync([EnumeratorCancellation] CancellationToken cancellationToken)
|
||||
{
|
||||
var workflowRegistry = _serviceProvider.GetRequiredService<IWorkflowRegistry>();
|
||||
var workflows = await workflowRegistry.GetWorkflowsAsync(cancellationToken).ToListAsync(cancellationToken);
|
||||
var workflows = await workflowRegistry.ListAsync(cancellationToken);
|
||||
|
||||
var query =
|
||||
from workflow in workflows
|
||||
|
|
|
|||
|
|
@ -134,7 +134,7 @@ namespace Elsa.Activities.Http.Middleware
|
|||
CancellationToken cancellationToken)
|
||||
{
|
||||
// Find workflows starting with HttpRequestReceived.
|
||||
var workflows = await workflowRegistry.GetWorkflowsAsync(cancellationToken).ToListAsync(cancellationToken);
|
||||
var workflows = await workflowRegistry.ListAsync(cancellationToken);
|
||||
|
||||
var httpWorkflows =
|
||||
from workflow in workflows
|
||||
|
|
|
|||
|
|
@ -25,7 +25,7 @@ namespace Elsa.Activities.Timers.Hangfire.Jobs
|
|||
|
||||
public async Task ExecuteAsync(RunHangfireWorkflowJobModel data)
|
||||
{
|
||||
var workflowBlueprint = (await _workflowRegistry.GetWorkflowAsync(data.WorkflowDefinitionId, data.TenantId, VersionOptions.Published));
|
||||
var workflowBlueprint = (await _workflowRegistry.GetAsync(data.WorkflowDefinitionId, data.TenantId, VersionOptions.Published));
|
||||
|
||||
if(workflowBlueprint == null)
|
||||
return;
|
||||
|
|
|
|||
|
|
@ -58,7 +58,7 @@ namespace Elsa.Activities.Timers.Quartz.Jobs
|
|||
{
|
||||
if (workflowInstanceId == null)
|
||||
{
|
||||
var workflowBlueprint = (await _workflowRegistry.GetWorkflowAsync(workflowDefinitionId, tenantId, VersionOptions.Published, cancellationToken))!;
|
||||
var workflowBlueprint = (await _workflowRegistry.GetAsync(workflowDefinitionId, tenantId, VersionOptions.Published, cancellationToken))!;
|
||||
|
||||
if (!workflowBlueprint.IsSingleton || await GetWorkflowIsAlreadyExecutingAsync(tenantId, workflowDefinitionId) == false)
|
||||
await _workflowQueue.EnqueueWorkflowDefinition(workflowDefinitionId, tenantId, activityId, null, null, null, cancellationToken);
|
||||
|
|
|
|||
|
|
@ -1,5 +1,4 @@
|
|||
using System;
|
||||
using System.Collections.Generic;
|
||||
using System.Collections.Generic;
|
||||
using System.Linq;
|
||||
using System.Threading;
|
||||
using System.Threading.Tasks;
|
||||
|
|
|
|||
|
|
@ -2,6 +2,9 @@
|
|||
|
||||
namespace Elsa.Caching
|
||||
{
|
||||
/// <summary>
|
||||
/// Provides change tokens for memory caches, allowing code to evict cache entries by triggering a signal.
|
||||
/// </summary>
|
||||
public interface ISignal
|
||||
{
|
||||
IChangeToken GetToken(string key);
|
||||
|
|
|
|||
|
|
@ -2,7 +2,6 @@
|
|||
using System.Threading;
|
||||
using System.Threading.Tasks;
|
||||
using Elsa.Services.Models;
|
||||
using Microsoft.Extensions.DependencyInjection;
|
||||
|
||||
namespace Elsa.Services
|
||||
{
|
||||
|
|
|
|||
|
|
@ -1,3 +1,4 @@
|
|||
using System;
|
||||
using System.Collections.Generic;
|
||||
using System.Threading;
|
||||
using System.Threading.Tasks;
|
||||
|
|
@ -8,7 +9,9 @@ namespace Elsa.Services
|
|||
{
|
||||
public interface IWorkflowRegistry
|
||||
{
|
||||
IAsyncEnumerable<IWorkflowBlueprint> GetWorkflowsAsync(CancellationToken cancellationToken = default);
|
||||
Task<IWorkflowBlueprint?> GetWorkflowAsync(string id, string? tenantId, VersionOptions version, CancellationToken cancellationToken = default);
|
||||
Task<IEnumerable<IWorkflowBlueprint>> ListAsync(CancellationToken cancellationToken = default);
|
||||
Task<IWorkflowBlueprint?> GetAsync(string id, string? tenantId, VersionOptions version, CancellationToken cancellationToken = default);
|
||||
Task<IEnumerable<IWorkflowBlueprint>> FindManyAsync(Func<IWorkflowBlueprint, bool> predicate, CancellationToken cancellationToken = default);
|
||||
Task<IWorkflowBlueprint?> FindAsync(Func<IWorkflowBlueprint, bool> predicate, CancellationToken cancellationToken = default);
|
||||
}
|
||||
}
|
||||
|
|
@ -1,5 +1,4 @@
|
|||
using System;
|
||||
using System.Collections.Generic;
|
||||
using System.Linq;
|
||||
using System.Threading;
|
||||
using System.Threading.Tasks;
|
||||
|
|
@ -69,7 +68,7 @@ namespace Elsa.Activities.Workflows
|
|||
|
||||
private async Task<IWorkflowBlueprint?> FindWorkflowBlueprintAsync(CancellationToken cancellationToken)
|
||||
{
|
||||
var query = (IEnumerable<IWorkflowBlueprint>)(await _workflowRegistry.GetWorkflowsAsync(cancellationToken).ToListAsync(cancellationToken));
|
||||
var query = await _workflowRegistry.ListAsync(cancellationToken);
|
||||
|
||||
query = query.Where(x => x.WithVersion(VersionOptions.Published));
|
||||
|
||||
|
|
|
|||
|
|
@ -10,7 +10,6 @@ using Elsa.Persistence.Specifications.Bookmarks;
|
|||
using Elsa.Serialization;
|
||||
using Elsa.Services;
|
||||
using Elsa.Services.Models;
|
||||
using Microsoft.Extensions.DependencyInjection;
|
||||
using Microsoft.Extensions.Logging;
|
||||
using Open.Linq.AsyncExtensions;
|
||||
using Rebus.Extensions;
|
||||
|
|
@ -61,7 +60,7 @@ namespace Elsa.Bookmarks
|
|||
var workflowInstanceIds = workflowInstanceList.Select(x => x.Id).ToList();
|
||||
await DeleteBookmarksAsync(workflowInstanceIds, cancellationToken);
|
||||
|
||||
var workflowBlueprints = await _workflowRegistry.GetWorkflowsAsync(cancellationToken).ToDictionaryAsync(x => (x.Id, x.Version), cancellationToken);
|
||||
var workflowBlueprints = (await _workflowRegistry.ListAsync(cancellationToken)).ToDictionary(x => (x.Id, x.Version));
|
||||
var entities = new List<Bookmark>();
|
||||
|
||||
foreach (var workflowInstance in workflowInstanceList.Where(x => x.WorkflowStatus == WorkflowStatus.Suspended))
|
||||
|
|
|
|||
|
|
@ -32,7 +32,7 @@ namespace Elsa.Consumers
|
|||
{
|
||||
var workflowDefinitionId = message.WorkflowDefinitionId;
|
||||
var tenantId = message.TenantId;
|
||||
var workflowBlueprint = await _workflowRegistry.GetWorkflowAsync(workflowDefinitionId, tenantId, VersionOptions.Published);
|
||||
var workflowBlueprint = await _workflowRegistry.GetAsync(workflowDefinitionId, tenantId, VersionOptions.Published);
|
||||
|
||||
if (!ValidatePreconditions(workflowDefinitionId, workflowBlueprint))
|
||||
return;
|
||||
|
|
|
|||
55
src/core/Elsa.Core/Decorators/CachingWorkflowRegistry.cs
Normal file
55
src/core/Elsa.Core/Decorators/CachingWorkflowRegistry.cs
Normal file
|
|
@ -0,0 +1,55 @@
|
|||
using System;
|
||||
using System.Collections.Generic;
|
||||
using System.Linq;
|
||||
using System.Threading;
|
||||
using System.Threading.Tasks;
|
||||
using Elsa.Caching;
|
||||
using Elsa.Models;
|
||||
using Elsa.Services;
|
||||
using Elsa.Services.Models;
|
||||
using Microsoft.Extensions.Caching.Memory;
|
||||
using Open.Linq.AsyncExtensions;
|
||||
|
||||
namespace Elsa.Decorators
|
||||
{
|
||||
public class CachingWorkflowRegistry : IWorkflowRegistry
|
||||
{
|
||||
private const string CacheKey = "WorkflowRegistry";
|
||||
private readonly IWorkflowRegistry _workflowRegistry;
|
||||
private readonly IMemoryCache _memoryCache;
|
||||
private readonly ISignal _signal;
|
||||
|
||||
public CachingWorkflowRegistry(IWorkflowRegistry workflowRegistry, IMemoryCache memoryCache, ISignal signal)
|
||||
{
|
||||
_workflowRegistry = workflowRegistry;
|
||||
_memoryCache = memoryCache;
|
||||
_signal = signal;
|
||||
}
|
||||
|
||||
public async Task<IEnumerable<IWorkflowBlueprint>> ListAsync(CancellationToken cancellationToken) => await GetWorkflowBlueprints(cancellationToken);
|
||||
|
||||
public async Task<IWorkflowBlueprint?> GetAsync(string id, string? tenantId, VersionOptions version, CancellationToken cancellationToken) =>
|
||||
await FindAsync(x => x.Id == id && x.TenantId == tenantId && x.WithVersion(version), cancellationToken);
|
||||
|
||||
public async Task<IEnumerable<IWorkflowBlueprint>> FindManyAsync(Func<IWorkflowBlueprint, bool> predicate, CancellationToken cancellationToken)
|
||||
{
|
||||
var workflows = await GetWorkflowBlueprints(cancellationToken);
|
||||
return workflows.Where(predicate);
|
||||
}
|
||||
|
||||
public async Task<IWorkflowBlueprint?> FindAsync(Func<IWorkflowBlueprint, bool> predicate, CancellationToken cancellationToken)
|
||||
{
|
||||
var workflows = await GetWorkflowBlueprints(cancellationToken);
|
||||
return workflows.FirstOrDefault(predicate);
|
||||
}
|
||||
|
||||
private async Task<ICollection<IWorkflowBlueprint>> GetWorkflowBlueprints(CancellationToken cancellationToken)
|
||||
{
|
||||
return await _memoryCache.GetOrCreateAsync(CacheKey, async entry =>
|
||||
{
|
||||
entry.Monitor(_signal.GetToken(CacheKey));
|
||||
return await _workflowRegistry.ListAsync(cancellationToken).ToList();
|
||||
});
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
@ -10,7 +10,9 @@ using Elsa.ActivityTypeProviders;
|
|||
using Elsa.Bookmarks;
|
||||
using Elsa.Builders;
|
||||
using Elsa.Consumers;
|
||||
using Elsa.Decorators;
|
||||
using Elsa.Expressions;
|
||||
using Elsa.Handlers;
|
||||
using Elsa.HostedServices;
|
||||
using Elsa.Mapping;
|
||||
using Elsa.Messages;
|
||||
|
|
@ -30,7 +32,6 @@ using Microsoft.Extensions.DependencyInjection.Extensions;
|
|||
using Newtonsoft.Json;
|
||||
using NodaTime;
|
||||
using Rebus.Handlers;
|
||||
using Rebus.ServiceProvider;
|
||||
|
||||
// ReSharper disable once CheckNamespace
|
||||
namespace Microsoft.Extensions.DependencyInjection
|
||||
|
|
@ -60,7 +61,6 @@ namespace Microsoft.Extensions.DependencyInjection
|
|||
.AddWorkflowsCore()
|
||||
.AddCoreActivities();
|
||||
|
||||
options.AddMediatR();
|
||||
options.AddAutoMapper();
|
||||
options.AddConsumer<RunWorkflowDefinitionConsumer, RunWorkflowDefinition>();
|
||||
options.AddConsumer<RunWorkflowInstanceConsumer, RunWorkflowInstance>();
|
||||
|
|
@ -71,12 +71,12 @@ namespace Microsoft.Extensions.DependencyInjection
|
|||
|
||||
return services;
|
||||
}
|
||||
|
||||
|
||||
/// <summary>
|
||||
/// Starts the specified workflow upon application startup.
|
||||
/// </summary>
|
||||
public static IServiceCollection StartWorkflow<T>(this IServiceCollection services) where T : class, IWorkflow => services.AddHostedService<StartWorkflow<T>>();
|
||||
|
||||
|
||||
public static ElsaOptions AddConsumer<TConsumer, TMessage>(this ElsaOptions elsaOptions) where TConsumer : class, IHandleMessages<TMessage>
|
||||
{
|
||||
elsaOptions.Services.AddTransient<IHandleMessages<TMessage>, TConsumer>();
|
||||
|
|
@ -84,70 +84,101 @@ namespace Microsoft.Extensions.DependencyInjection
|
|||
return elsaOptions;
|
||||
}
|
||||
|
||||
private static ElsaOptions AddMediatR(this ElsaOptions options)
|
||||
{
|
||||
options.Services.AddMediatR(mediatr => mediatr.AsScoped(), typeof(IActivity));
|
||||
return options;
|
||||
}
|
||||
|
||||
private static ElsaOptions AddWorkflowsCore(this ElsaOptions options)
|
||||
{
|
||||
var services = options.Services;
|
||||
services.TryAddSingleton<IClock>(SystemClock.Instance);
|
||||
|
||||
services
|
||||
.TryAddSingleton<IClock>(SystemClock.Instance);
|
||||
|
||||
services
|
||||
.AddLogging()
|
||||
.AddLocalization()
|
||||
.AddMemoryCache()
|
||||
.AddSingleton<IIdGenerator, IdGenerator>()
|
||||
.AddTransient<Func<JsonSerializer>>(sp => sp.GetRequiredService<JsonSerializer>)
|
||||
.AddTransient(sp => sp.GetRequiredService<ElsaOptions>().CreateJsonSerializer(sp))
|
||||
.AddSingleton<IContentSerializer, DefaultContentSerializer>()
|
||||
.AddSingleton<TypeJsonConverter>()
|
||||
.TryAddProvider<IExpressionHandler, LiteralHandler>(ServiceLifetime.Singleton)
|
||||
.TryAddProvider<IExpressionHandler, VariableHandler>(ServiceLifetime.Singleton)
|
||||
.AddScoped<IExpressionEvaluator, ExpressionEvaluator>()
|
||||
.AddScoped<IWorkflowRegistry, WorkflowRegistry>()
|
||||
.AddSingleton<IActivityActivator, ActivityActivator>()
|
||||
.AddScoped<IWorkflowRunner, WorkflowRunner>()
|
||||
.AddScoped<IWorkflowTriggerInterruptor, WorkflowTriggerInterruptor>()
|
||||
.AddScoped<IWorkflowReviver, WorkflowReviver>()
|
||||
.AddSingleton<IActivityDescriber, ActivityDescriber>()
|
||||
.AddSingleton<IWorkflowFactory, WorkflowFactory>()
|
||||
.AddSingleton<IWorkflowBlueprintMaterializer, WorkflowBlueprintMaterializer>()
|
||||
.AddSingleton<IWorkflowBlueprintReflector, WorkflowBlueprintReflector>()
|
||||
.AddSingleton<IBackgroundWorker, BackgroundWorker>()
|
||||
.AddScoped<IWorkflowPublisher, WorkflowPublisher>()
|
||||
.AddScoped<IWorkflowContextManager, WorkflowContextManager>()
|
||||
.AddSingleton<IActivityTypeService, ActivityTypeService>()
|
||||
.AddSingleton<IActivityTypeProvider, TypeBasedActivityProvider>()
|
||||
.AddScoped<IWorkflowExecutionLog, WorkflowExecutionLog>()
|
||||
;
|
||||
|
||||
// Serialization.
|
||||
services
|
||||
.AddTransient<Func<JsonSerializer>>(sp => sp.GetRequiredService<JsonSerializer>)
|
||||
.AddTransient(sp => sp.GetRequiredService<ElsaOptions>().CreateJsonSerializer(sp))
|
||||
.AddSingleton<IContentSerializer, DefaultContentSerializer>()
|
||||
.AddSingleton<TypeJsonConverter>();
|
||||
|
||||
// Expressions.
|
||||
services
|
||||
.TryAddProvider<IExpressionHandler, LiteralHandler>(ServiceLifetime.Singleton)
|
||||
.TryAddProvider<IExpressionHandler, VariableHandler>(ServiceLifetime.Singleton)
|
||||
.AddScoped<IExpressionEvaluator, ExpressionEvaluator>();
|
||||
|
||||
// Workflow providers.
|
||||
services
|
||||
.AddWorkflowProvider<ProgrammaticWorkflowProvider>()
|
||||
.AddWorkflowProvider<StorageWorkflowProvider>()
|
||||
.AddWorkflowProvider<DatabaseWorkflowProvider>();
|
||||
|
||||
// Metadata.
|
||||
services
|
||||
.AddSingleton<IActivityDescriber, ActivityDescriber>()
|
||||
.AddMetadataHandlers();
|
||||
|
||||
// Bookmarks.
|
||||
services
|
||||
.AddSingleton<IBookmarkHasher, BookmarkHasher>()
|
||||
.AddScoped<IBookmarkIndexer, BookmarkIndexer>()
|
||||
.AddScoped<IBookmarkFinder, BookmarkFinder>()
|
||||
.AddScoped<ITriggerIndexer, TriggerIndexer>()
|
||||
.AddSingleton<ITriggerStore, TriggerStore>()
|
||||
.AddScoped<ITriggerFinder, TriggerFinder>()
|
||||
.AddSingleton<IBackgroundWorker, BackgroundWorker>()
|
||||
.AddBookmarkProvider<SignalReceivedBookmarkProvider>()
|
||||
.AddBookmarkProvider<RunWorkflowBookmarkProvider>();
|
||||
|
||||
// Mediator.
|
||||
services
|
||||
.AddMediatR(mediatr => mediatr.AsScoped(), typeof(IActivity), typeof(LogWorkflowExecution));
|
||||
|
||||
// Service Bus.
|
||||
services
|
||||
.AddScoped<IWorkflowQueue, WorkflowQueue>()
|
||||
.AddScoped<IWorkflowPublisher, WorkflowPublisher>()
|
||||
.AddScoped<IWorkflowContextManager, WorkflowContextManager>()
|
||||
.AddSingleton<IActivityTypeService, ActivityTypeService>()
|
||||
.AddSingleton<IActivityTypeProvider, TypeBasedActivityProvider>()
|
||||
.AddWorkflowProvider<ProgrammaticWorkflowProvider>()
|
||||
.AddWorkflowProvider<StorageWorkflowProvider>()
|
||||
.AddWorkflowProvider<DatabaseWorkflowProvider>()
|
||||
.AddTransient<IWorkflowBuilder, WorkflowBuilder>()
|
||||
.AddTransient<ICompositeActivityBuilder, CompositeActivityBuilder>()
|
||||
.AddTransient<Func<IWorkflowBuilder>>(sp => sp.GetRequiredService<IWorkflowBuilder>)
|
||||
.AddAutoMapperProfile<NodaTimeProfile>()
|
||||
.AddAutoMapperProfile<CloningProfile>()
|
||||
.AddSingleton<ICloner, AutoMapperCloner>()
|
||||
.AddNotificationHandlers(typeof(ElsaServiceCollectionExtensions))
|
||||
.AddSingleton<ServiceBusFactory>()
|
||||
.AddSingleton<IServiceBusFactory, ServiceBusFactory>()
|
||||
.AddSingleton<ICommandSender, CommandSender>()
|
||||
.AddSingleton<IEventPublisher, EventPublisher>()
|
||||
.AddSingleton<IEventPublisher, EventPublisher>();
|
||||
|
||||
options
|
||||
.AddConsumer<RunWorkflowDefinitionConsumer, RunWorkflowDefinition>()
|
||||
.AddConsumer<RunWorkflowInstanceConsumer, RunWorkflowInstance>();
|
||||
|
||||
// AutoMapper.
|
||||
services
|
||||
.AddAutoMapperProfile<NodaTimeProfile>()
|
||||
.AddAutoMapperProfile<CloningProfile>()
|
||||
.AddSingleton<ICloner, AutoMapperCloner>();
|
||||
|
||||
// Caching.
|
||||
services
|
||||
.AddMemoryCache()
|
||||
.AddScoped<ISignaler, Signaler>()
|
||||
.AddScoped<IWorkflowExecutionLog, WorkflowExecutionLog>()
|
||||
.AutoRegisterHandlersFromAssemblyOf<RunWorkflowInstanceConsumer>()
|
||||
.AddBookmarkProvider<SignalReceivedBookmarkProvider>()
|
||||
.AddBookmarkProvider<RunWorkflowBookmarkProvider>()
|
||||
.AddMetadataHandlers();
|
||||
.Decorate<IWorkflowRegistry, CachingWorkflowRegistry>();
|
||||
|
||||
// Builder API.
|
||||
services
|
||||
.AddTransient<IWorkflowBuilder, WorkflowBuilder>()
|
||||
.AddTransient<ICompositeActivityBuilder, CompositeActivityBuilder>()
|
||||
.AddTransient<Func<IWorkflowBuilder>>(sp => sp.GetRequiredService<IWorkflowBuilder>);
|
||||
|
||||
return options;
|
||||
}
|
||||
|
|
|
|||
|
|
@ -12,7 +12,7 @@ namespace Elsa
|
|||
this IWorkflowRegistry workflowRegistry,
|
||||
string? tenantId,
|
||||
CancellationToken cancellationToken = default) =>
|
||||
workflowRegistry.GetWorkflowAsync(typeof(T).Name, tenantId, VersionOptions.Latest, cancellationToken);
|
||||
workflowRegistry.GetAsync(typeof(T).Name, tenantId, VersionOptions.Latest, cancellationToken);
|
||||
|
||||
public static Task<IWorkflowBlueprint?> GetWorkflowAsync<T>(
|
||||
this IWorkflowRegistry workflowRegistry,
|
||||
|
|
@ -24,7 +24,7 @@ namespace Elsa
|
|||
string id,
|
||||
VersionOptions versionOptions,
|
||||
CancellationToken cancellationToken = default) =>
|
||||
workflowRegistry.GetWorkflowAsync(id, default, versionOptions, cancellationToken);
|
||||
workflowRegistry.GetAsync(id, default, versionOptions, cancellationToken);
|
||||
|
||||
// public static async Task<IEnumerable<(IWorkflowBlueprint Workflow, IActivityBlueprint Activity)>>
|
||||
// GetWorkflowsByStartActivityAsync<T>(
|
||||
|
|
|
|||
|
|
@ -2,7 +2,6 @@
|
|||
using System.Threading;
|
||||
using System.Threading.Tasks;
|
||||
using Elsa.Services.Models;
|
||||
using Microsoft.Extensions.DependencyInjection;
|
||||
|
||||
namespace Elsa.Services
|
||||
{
|
||||
|
|
|
|||
|
|
@ -6,6 +6,7 @@ using System.Threading;
|
|||
using System.Threading.Tasks;
|
||||
using Elsa.Models;
|
||||
using Elsa.Services.Models;
|
||||
using Open.Linq.AsyncExtensions;
|
||||
|
||||
namespace Elsa.Services
|
||||
{
|
||||
|
|
@ -18,22 +19,11 @@ namespace Elsa.Services
|
|||
_workflowProviders = workflowProviders;
|
||||
}
|
||||
|
||||
public async IAsyncEnumerable<IWorkflowBlueprint> GetWorkflowsAsync([EnumeratorCancellation] CancellationToken cancellationToken)
|
||||
{
|
||||
var providers = _workflowProviders;
|
||||
public async Task<IEnumerable<IWorkflowBlueprint>> ListAsync(CancellationToken cancellationToken) => await GetWorkflowsInternalAsync(cancellationToken).ToListAsync(cancellationToken);
|
||||
|
||||
foreach (var provider in providers)
|
||||
await foreach (var workflow in provider.GetWorkflowsAsync(cancellationToken).WithCancellation(cancellationToken))
|
||||
yield return workflow;
|
||||
}
|
||||
|
||||
public async Task<IWorkflowBlueprint?> GetWorkflowAsync(
|
||||
string id,
|
||||
string? tenantId,
|
||||
VersionOptions version,
|
||||
CancellationToken cancellationToken)
|
||||
public async Task<IWorkflowBlueprint?> GetAsync(string id, string? tenantId, VersionOptions version, CancellationToken cancellationToken)
|
||||
{
|
||||
var workflows = await GetWorkflowsAsync(cancellationToken).ToListAsync(cancellationToken);
|
||||
var workflows = await ListAsync(cancellationToken).ToList();
|
||||
var query = workflows.Where(workflow => workflow.Id == id && workflow.WithVersion(version));
|
||||
|
||||
if (tenantId != null)
|
||||
|
|
@ -44,10 +34,19 @@ namespace Elsa.Services
|
|||
.FirstOrDefault();
|
||||
}
|
||||
|
||||
public async Task<IEnumerable<IWorkflowBlueprint>> FindWorkflowsAsync(Func<IWorkflowBlueprint, bool> predicate, CancellationToken cancellationToken) =>
|
||||
await GetWorkflowsAsync(cancellationToken).Where(predicate).OrderByDescending(x => x.Version).ToListAsync(cancellationToken);
|
||||
public async Task<IEnumerable<IWorkflowBlueprint>> FindManyAsync(Func<IWorkflowBlueprint, bool> predicate, CancellationToken cancellationToken) =>
|
||||
(await ListAsync(cancellationToken).Where(predicate).OrderByDescending(x => x.Version)).ToList();
|
||||
|
||||
public async Task<IWorkflowBlueprint?> FindWorkflowAsync(Func<IWorkflowBlueprint, bool> predicate, CancellationToken cancellationToken) =>
|
||||
await GetWorkflowsAsync(cancellationToken).Where(predicate).OrderByDescending(x => x.Version).FirstOrDefaultAsync(cancellationToken);
|
||||
public async Task<IWorkflowBlueprint?> FindAsync(Func<IWorkflowBlueprint, bool> predicate, CancellationToken cancellationToken) =>
|
||||
(await ListAsync(cancellationToken).Where(predicate).OrderByDescending(x => x.Version)).FirstOrDefault();
|
||||
|
||||
private async IAsyncEnumerable<IWorkflowBlueprint> GetWorkflowsInternalAsync([EnumeratorCancellation] CancellationToken cancellationToken)
|
||||
{
|
||||
var providers = _workflowProviders;
|
||||
|
||||
foreach (var provider in providers)
|
||||
await foreach (var workflow in provider.GetWorkflowsAsync(cancellationToken).WithCancellation(cancellationToken))
|
||||
yield return workflow;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
@ -148,7 +148,7 @@ namespace Elsa.Services
|
|||
object? input = default,
|
||||
CancellationToken cancellationToken = default)
|
||||
{
|
||||
var workflowBlueprint = await _workflowRegistry.GetWorkflowAsync(
|
||||
var workflowBlueprint = await _workflowRegistry.GetAsync(
|
||||
workflowInstance.DefinitionId,
|
||||
workflowInstance.TenantId,
|
||||
VersionOptions.SpecificVersion(workflowInstance.Version),
|
||||
|
|
|
|||
|
|
@ -64,7 +64,7 @@ namespace Elsa.Services
|
|||
await InterruptActivityTypeInternalAsync(activityType, input, cancellationToken).ToListAsync(cancellationToken);
|
||||
|
||||
private async Task<IWorkflowBlueprint?> GetWorkflowBlueprintAsync(WorkflowInstance workflowInstance, CancellationToken cancellationToken) =>
|
||||
await _workflowRegistry.GetWorkflowAsync(workflowInstance.DefinitionId, workflowInstance.TenantId, VersionOptions.SpecificVersion(workflowInstance.Version), cancellationToken);
|
||||
await _workflowRegistry.GetAsync(workflowInstance.DefinitionId, workflowInstance.TenantId, VersionOptions.SpecificVersion(workflowInstance.Version), cancellationToken);
|
||||
|
||||
private async IAsyncEnumerable<WorkflowInstance> InterruptActivityTypeInternalAsync(string activityType, object? input, [EnumeratorCancellation] CancellationToken cancellationToken)
|
||||
{
|
||||
|
|
|
|||
|
|
@ -7,7 +7,6 @@ using System.Threading.Tasks;
|
|||
using Elsa.Bookmarks;
|
||||
using Elsa.Services;
|
||||
using Elsa.Services.Models;
|
||||
using Microsoft.Extensions.DependencyInjection;
|
||||
using Microsoft.Extensions.Logging;
|
||||
|
||||
namespace Elsa.Triggers
|
||||
|
|
@ -43,7 +42,7 @@ namespace Elsa.Triggers
|
|||
|
||||
public async Task IndexTriggersAsync(CancellationToken cancellationToken = default)
|
||||
{
|
||||
var workflowBlueprints = await _workflowRegistry.GetWorkflowsAsync(cancellationToken).ToListAsync(cancellationToken);
|
||||
var workflowBlueprints = await _workflowRegistry.ListAsync(cancellationToken);
|
||||
await IndexTriggersAsync(workflowBlueprints, cancellationToken);
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -3,8 +3,6 @@ using System;
|
|||
using Elsa.Persistence.EntityFramework.Core;
|
||||
using Microsoft.EntityFrameworkCore;
|
||||
using Microsoft.EntityFrameworkCore.Infrastructure;
|
||||
using Microsoft.EntityFrameworkCore.Metadata;
|
||||
using Microsoft.EntityFrameworkCore.Storage.ValueConversion;
|
||||
|
||||
namespace Elsa.Persistence.EntityFramework.SqlServer.Migrations
|
||||
{
|
||||
|
|
|
|||
|
|
@ -3,7 +3,6 @@ using System;
|
|||
using Elsa.Persistence.EntityFramework.Core;
|
||||
using Microsoft.EntityFrameworkCore;
|
||||
using Microsoft.EntityFrameworkCore.Infrastructure;
|
||||
using Microsoft.EntityFrameworkCore.Storage.ValueConversion;
|
||||
|
||||
namespace Elsa.Persistence.EntityFramework.Sqlite.Migrations
|
||||
{
|
||||
|
|
|
|||
|
|
@ -6,7 +6,6 @@ using Elsa.Samples.Timers.Workflows;
|
|||
using Microsoft.EntityFrameworkCore;
|
||||
using Microsoft.Extensions.DependencyInjection;
|
||||
using Microsoft.Extensions.Hosting;
|
||||
using NodaTime;
|
||||
|
||||
namespace Elsa.Samples.Timers
|
||||
{
|
||||
|
|
|
|||
|
|
@ -42,7 +42,7 @@ namespace Elsa.Server.Api.Endpoints.WorkflowRegistry
|
|||
public async Task<ActionResult<WorkflowBlueprintModel>> Handle(string id, VersionOptions? versionOptions = default, CancellationToken cancellationToken = default)
|
||||
{
|
||||
versionOptions ??= VersionOptions.Latest;
|
||||
var workflowBlueprint = await _workflowRegistry.GetWorkflowAsync(id, null, versionOptions.Value, cancellationToken);
|
||||
var workflowBlueprint = await _workflowRegistry.GetAsync(id, null, versionOptions.Value, cancellationToken);
|
||||
|
||||
if (workflowBlueprint == null)
|
||||
return NotFound();
|
||||
|
|
|
|||
|
|
@ -10,6 +10,7 @@ using Elsa.Services;
|
|||
using Microsoft.AspNetCore.Http;
|
||||
using Microsoft.AspNetCore.Mvc;
|
||||
using Microsoft.Extensions.DependencyInjection;
|
||||
using Open.Linq.AsyncExtensions;
|
||||
using Swashbuckle.AspNetCore.Annotations;
|
||||
|
||||
namespace Elsa.Server.Api.Endpoints.WorkflowRegistry
|
||||
|
|
@ -42,7 +43,7 @@ namespace Elsa.Server.Api.Endpoints.WorkflowRegistry
|
|||
public async Task<ActionResult<PagedList<WorkflowBlueprintSummaryModel>>> Handle(int? page = default, int? pageSize = default, VersionOptions? version = default, CancellationToken cancellationToken = default)
|
||||
{
|
||||
version ??= VersionOptions.Latest;
|
||||
var workflowBlueprints = await _workflowRegistry.GetWorkflowsAsync(cancellationToken).Where(x => x.WithVersion(version.Value)).ToListAsync(cancellationToken);
|
||||
var workflowBlueprints = await _workflowRegistry.FindManyAsync(x => x.WithVersion(version.Value), cancellationToken).ToList();
|
||||
var totalCount = workflowBlueprints.Count;
|
||||
var skip = page * pageSize;
|
||||
var items = workflowBlueprints.AsEnumerable();
|
||||
|
|
|
|||
|
|
@ -1,5 +1,4 @@
|
|||
using System;
|
||||
using System.Linq;
|
||||
using System.Linq;
|
||||
using System.Threading;
|
||||
using System.Threading.Tasks;
|
||||
using AutoMapper;
|
||||
|
|
|
|||
Loading…
Reference in a new issue