Simplify activity definition activity

This commit is contained in:
Sipke Schoorstra 2022-08-11 20:17:50 +02:00
parent 98f195172b
commit c8a9a2ef8f
7 changed files with 9 additions and 96 deletions

View file

@ -49,6 +49,9 @@ public class EFCoreActivityDefinitionStore : IActivityDefinitionStore
return await query.PaginateAsync(x => ActivityDefinitionSummary.FromDefinition(x), pageArgs);
}
public async Task<ActivityDefinition?> FindByTypeAsync(string type, int version, CancellationToken cancellationToken = default) =>
await _store.FindAsync(x => x.Type == type && x.Version == version, cancellationToken);
public async Task<ActivityDefinition?> FindByDefinitionIdAsync(string definitionId, VersionOptions versionOptions, CancellationToken cancellationToken = default)
{
Expression<Func<ActivityDefinition, bool>> predicate = x => x.DefinitionId == definitionId;

View file

@ -10,16 +10,6 @@ namespace Elsa.ActivityDefinitions.Activities;
/// </summary>
public class ActivityDefinitionActivity : ActivityBase
{
/// <summary>
/// The activity definition ID to load & execute.
/// </summary>
public string DefinitionId { get; set; } = default!;
/// <summary>
/// The activity definition version number to load & execute.
/// </summary>
public int DefinitionVersion { get; set; }
protected override async ValueTask ExecuteAsync(ActivityExecutionContext context)
{
// Construct the root activity stored in the activity definitions.

View file

@ -1,73 +0,0 @@
using System.Text.Json;
using Elsa.ActivityDefinitions.Activities;
using Elsa.ActivityDefinitions.Extensions;
using Elsa.ActivityDefinitions.Services;
using Elsa.Mediator.Services;
using Elsa.Persistence.Common.Models;
using Elsa.Workflows.Core;
using Elsa.Workflows.Core.Serialization;
using Elsa.Workflows.Core.Services;
using Elsa.Workflows.Management.Notifications;
using Elsa.Workflows.Runtime.Services;
using Microsoft.Extensions.Logging;
using Open.Linq.AsyncExtensions;
namespace Elsa.ActivityDefinitions.Handlers;
/// <summary>
/// Updates the referenced definition version for all activity definitions referenced by the published workflow definition.
/// </summary>
public class UpdateReferencedVersionHandler : INotificationHandler<WorkflowDefinitionPublishing>
{
private readonly IWorkflowDefinitionService _workflowDefinitionService;
private readonly IActivityWalker _activityWalker;
private readonly IActivityDefinitionStore _activityDefinitionStore;
private readonly SerializerOptionsProvider _serializerOptionsProvider;
private readonly ILogger<UpdateReferencedVersionHandler> _logger;
public UpdateReferencedVersionHandler(
IWorkflowDefinitionService workflowDefinitionService,
IActivityWalker activityWalker,
IActivityDefinitionStore activityDefinitionStore,
SerializerOptionsProvider serializerOptionsProvider,
ILogger<UpdateReferencedVersionHandler> logger)
{
_workflowDefinitionService = workflowDefinitionService;
_activityWalker = activityWalker;
_activityDefinitionStore = activityDefinitionStore;
_serializerOptionsProvider = serializerOptionsProvider;
_logger = logger;
}
public async Task HandleAsync(WorkflowDefinitionPublishing notification, CancellationToken cancellationToken)
{
var workflowDefinition = notification.WorkflowDefinition;
var workflow = await _workflowDefinitionService.MaterializeWorkflowAsync(workflowDefinition, cancellationToken);
var root = workflow.Root;
var graph = await _activityWalker.WalkAsync(root, cancellationToken);
var nodes = graph.Flatten().Distinct().ToList();
var activityDefinitionLookup = await _activityDefinitionStore.ListAsync(VersionOptions.Published, cancellationToken).ToDictionary(x => x.DefinitionId);
var customActivities = nodes.Where(x => x.Activity is ActivityDefinitionActivity).Select(x => (ActivityDefinitionActivity)x.Activity).ToList();
foreach (var customActivity in customActivities)
{
if(!activityDefinitionLookup.TryGetValue(customActivity.DefinitionId, out var activityDefinition))
{
_logger.LogWarning(
"Workflow definition {WorkflowDefinitionId} version {WorkflowDefinitionVersion} references a custom activity with ID {ActivityDefinitionId}, but there is no (published) activity definition by that ID",
workflowDefinition.DefinitionId,
workflowDefinition.Version,
customActivity.DefinitionId);
continue;
}
// Update to latest published version.
customActivity.DefinitionVersion = activityDefinition.Version;
}
// Serialize the workflow and save changes.
var serializerOptions = _serializerOptionsProvider.CreateApiOptions();
workflowDefinition.StringData = JsonSerializer.Serialize(root, serializerOptions);
}
}

View file

@ -49,13 +49,6 @@ public class ActivityDefinitionActivityProvider : IActivityProvider
var activity = (ActivityDefinitionActivity)_activityFactory.Create(typeof(ActivityDefinitionActivity), context);
activity.Type = definition.Type;
activity.Version = definition.Version;
if (string.IsNullOrWhiteSpace(activity.DefinitionId))
activity.DefinitionId = definition.DefinitionId;
if (activity.DefinitionVersion == 0)
activity.DefinitionVersion = definition.Version;
return activity;
}
};

View file

@ -25,7 +25,7 @@ public class ActivityDefinitionMaterializer : IActivityDefinitionMaterializer
public async Task<IActivity> MaterializeAsync(ActivityDefinitionActivity activity, CancellationToken cancellationToken = default)
{
var definition = await _store.FindByDefinitionIdAsync(activity.DefinitionId, VersionOptions.SpecificVersion(activity.DefinitionVersion), cancellationToken);
var definition = await _store.FindByTypeAsync(activity.Type, activity.Version, cancellationToken);
if (definition == null)
return new Sequence();

View file

@ -32,15 +32,15 @@ public class MemoryActivityDefinitionStore : IActivityDefinitionStore
return Task.FromResult(page);
}
public Task<ActivityDefinition?> FindByDefinitionIdAsync(string definitionId, VersionOptions versionOptions, CancellationToken cancellationToken = default)
public Task<ActivityDefinition?> FindByTypeAsync(string type, int version, CancellationToken cancellationToken = default)
{
var definition = _store.Find(x => x.DefinitionId == definitionId && x.WithVersion(versionOptions));
var definition = _store.Find(x => x.Type == type && x.Version == version);
return Task.FromResult(definition);
}
public Task<ActivityDefinition?> FindByDefinitionVersionIdAsync(string definitionVersionId, CancellationToken cancellationToken = default)
public Task<ActivityDefinition?> FindByDefinitionIdAsync(string definitionId, VersionOptions versionOptions, CancellationToken cancellationToken = default)
{
var definition = _store.Find(x => x.Id == definitionVersionId);
var definition = _store.Find(x => x.DefinitionId == definitionId && x.WithVersion(versionOptions));
return Task.FromResult(definition);
}

View file

@ -16,8 +16,8 @@ public interface IActivityDefinitionStore
PageArgs? pageArgs = default,
CancellationToken cancellationToken = default);
Task<ActivityDefinition?> FindByTypeAsync(string type, int version, CancellationToken cancellationToken = default);
Task<ActivityDefinition?> FindByDefinitionIdAsync(string definitionId, VersionOptions versionOptions, CancellationToken cancellationToken = default);
Task<ActivityDefinition?> FindByDefinitionVersionIdAsync(string definitionVersionId, CancellationToken cancellationToken = default);
Task<IEnumerable<ActivityDefinition>> FindLatestAndPublishedByDefinitionIdAsync(string definitionId, CancellationToken cancellationToken = default);
Task SaveAsync(ActivityDefinition record, CancellationToken cancellationToken = default);
Task<int> DeleteByDefinitionIdAsync(string definitionId, CancellationToken cancellationToken = default);