Incremental work on bookmark processing

This commit is contained in:
Sipke Schoorstra 2022-01-18 18:24:52 +01:00
parent 6e969bfa5f
commit 59f44b4927
18 changed files with 162 additions and 62 deletions

View file

@ -8,9 +8,14 @@ using IWorkflowTriggerScheduler = Elsa.Activities.Scheduling.Contracts.IWorkflow
namespace Elsa.Activities.Scheduling.Handlers;
// Updates scheduled jobs based on the updated workflow triggers.
public class ScheduleWorkflowsHandler : INotificationHandler<TriggerIndexingFinished>
public class ScheduleWorkflowsHandler : INotificationHandler<TriggerIndexingFinished>, INotificationHandler<WorkflowExecuted>
{
private readonly IWorkflowTriggerScheduler _workflowTriggerScheduler;
public ScheduleWorkflowsHandler(IWorkflowTriggerScheduler workflowTriggerScheduler) => _workflowTriggerScheduler = workflowTriggerScheduler;
public async Task HandleAsync(TriggerIndexingFinished notification, CancellationToken cancellationToken) => await _workflowTriggerScheduler.ScheduleTriggersAsync(notification.Triggers, cancellationToken);
public Task HandleAsync(WorkflowExecuted notification, CancellationToken cancellationToken)
{
return Task.CompletedTask;
}
}

View file

@ -7,13 +7,14 @@ public static class UseActivitySchedulerMiddlewareExtensions
{
public static IWorkflowExecutionBuilder UseActivityScheduler(this IWorkflowExecutionBuilder builder) => builder.UseMiddleware<ActivitySchedulerMiddleware>();
}
public class ActivitySchedulerMiddleware : IWorkflowExecutionMiddleware
{
private readonly WorkflowMiddlewareDelegate _next;
public ActivitySchedulerMiddleware(WorkflowMiddlewareDelegate next) => _next = next;
public async ValueTask InvokeAsync(WorkflowExecutionContext context)
public class ActivitySchedulerMiddleware : WorkflowExecutionMiddleware
{
public ActivitySchedulerMiddleware(WorkflowMiddlewareDelegate next) : base(next)
{
}
public override async ValueTask InvokeAsync(WorkflowExecutionContext context)
{
var scheduler = context.Scheduler;
@ -28,6 +29,6 @@ public class ActivitySchedulerMiddleware : IWorkflowExecutionMiddleware
}
// Invoke next middleware.
await _next(context);
await Next(context);
}
}

View file

@ -0,0 +1,11 @@
using Elsa.Contracts;
using Elsa.Models;
namespace Elsa.Pipelines.WorkflowExecution.Components;
public abstract class WorkflowExecutionMiddleware : IWorkflowExecutionMiddleware
{
protected WorkflowExecutionMiddleware(WorkflowMiddlewareDelegate next) => Next = next;
protected WorkflowMiddlewareDelegate Next { get; }
public abstract ValueTask InvokeAsync(WorkflowExecutionContext context);
}

View file

@ -21,6 +21,7 @@ import FromJSONData = Model.FromJSONData;
})
export class FlowchartComponent implements ContainerActivityComponent {
private rootId: string = uuid();
private silent: boolean = false; // Whether to emit events or not.
@Prop({mutable: true}) public activityDescriptors: Array<ActivityDescriptor> = [];
@Prop({mutable: true}) public root?: Activity;
@ -87,20 +88,32 @@ export class FlowchartComponent implements ContainerActivityComponent {
);
}
public disableEvents = () => this.silent = true;
public enableEvents = async (emitWorkflowChanged: boolean): Promise<void> => {
this.silent = false;
if (emitWorkflowChanged === true) {
await this.onGraphChanged();
}
};
private createAndInitializeGraph = async () => {
const graph = this.graph = createGraph(this.container, {
nodeMovable: () => this.interactiveMode,
edgeMovable: () => this.interactiveMode,
arrowheadMovable: () => this.interactiveMode,
edgeLabelMovable: () => this.interactiveMode,
magnetConnectable: () => this.interactiveMode,
useEdgeTools: () => this.interactiveMode,
toolsAddable: () => this.interactiveMode,
stopDelegateOnDragging: () => this.interactiveMode,
vertexAddable: () => this.interactiveMode,
vertexDeletable: () => this.interactiveMode,
vertexMovable: () => this.interactiveMode,
});
nodeMovable: () => this.interactiveMode,
edgeMovable: () => this.interactiveMode,
arrowheadMovable: () => this.interactiveMode,
edgeLabelMovable: () => this.interactiveMode,
magnetConnectable: () => this.interactiveMode,
useEdgeTools: () => this.interactiveMode,
toolsAddable: () => this.interactiveMode,
stopDelegateOnDragging: () => this.interactiveMode,
vertexAddable: () => this.interactiveMode,
vertexDeletable: () => this.interactiveMode,
vertexMovable: () => this.interactiveMode,
},
this.disableEvents,
this.enableEvents);
graph.on('blank:click', this.onGraphClick);
graph.on('node:click', this.onNodeClick);
@ -316,7 +329,10 @@ export class FlowchartComponent implements ContainerActivityComponent {
} as Connection;
}
private onGraphChanged = async (e: any) => {
private onGraphChanged = async () => {
debugger;
if (this.silent)
return;
this.graphUpdated.emit({exportGraph: this.exportRootInternal});
}
}

View file

@ -1,7 +1,13 @@
import {CellView, Graph, Node, Shape} from '@antv/x6';
import {v4 as uuid} from 'uuid';
import './ports';
export function createGraph(container: HTMLElement, interacting: CellView.Interacting): Graph {
export function createGraph(
container: HTMLElement,
interacting: CellView.Interacting,
disableEvents: () => void,
enableEvents: (emitWorkflowChanged: boolean) => Promise<void>): Graph {
const graph = new Graph({
container: container,
interacting: interacting,
@ -15,7 +21,11 @@ export function createGraph(container: HTMLElement, interacting: CellView.Intera
},
height: 5000,
width: 5000,
async: true,
// Keep disabled for now until we find that performance degrades significantly when adding too many nodes.
// When we do enable async rendering, we need to take care of the selection rectangle after pasting nodes, which would be calculated too early (before rendering completed).
async: false,
autoResize: true,
keyboard: {
enabled: true,
@ -156,11 +166,20 @@ export function createGraph(container: HTMLElement, interacting: CellView.Intera
return false
});
graph.bindKey(['meta+v', 'ctrl+v'], () => {
graph.bindKey(['meta+v', 'ctrl+v'], async () => {
if (!graph.isClipboardEmpty()) {
const cells = graph.paste({offset: 32})
graph.cleanSelection()
graph.select(cells)
debugger;
disableEvents();
const cells = graph.paste({offset: 32});
debugger;
for (const cell of cells) {
cell.data.id = uuid();
}
debugger;
await enableEvents(true);
graph.cleanSelection();
graph.select(cells);
}
return false
});

View file

@ -1,6 +1,7 @@
using Elsa.Mediator.Contracts;
using Elsa.Models;
using Elsa.Persistence.Entities;
namespace Elsa.Persistence.Commands;
public record ReplaceWorkflowTriggers(ICollection<WorkflowTrigger> WorkflowTriggers) : ICommand;
public record ReplaceWorkflowTriggers(Workflow Workflow, ICollection<WorkflowTrigger> WorkflowTriggers) : ICommand;

View file

@ -4,34 +4,33 @@ using Elsa.Models;
using Elsa.Persistence.Commands;
using Elsa.Persistence.Entities;
using Elsa.Pipelines.WorkflowExecution;
using Elsa.Pipelines.WorkflowExecution.Components;
namespace Elsa.Persistence.Middleware.WorkflowExecution;
public static class PersistWorkflowExecutionLogMiddlewareExtensions
{
public static IWorkflowExecutionBuilder PersistWorkflowExecutionLog(this IWorkflowExecutionBuilder builder) => builder.UseMiddleware<PersistWorkflowExecutionLogMiddleware>();
public static IWorkflowExecutionBuilder UseWorkflowExecutionLogPersistence(this IWorkflowExecutionBuilder builder) => builder.UseMiddleware<PersistWorkflowExecutionLogMiddleware>();
}
/// <summary>
/// Takes care of persisting a workflow instance after workflow execution.
/// </summary>
public class PersistWorkflowExecutionLogMiddleware : IWorkflowExecutionMiddleware
public class PersistWorkflowExecutionLogMiddleware : WorkflowExecutionMiddleware
{
private readonly WorkflowMiddlewareDelegate _next;
private readonly ICommandSender _commandSender;
private readonly IIdentityGenerator _identityGenerator;
public PersistWorkflowExecutionLogMiddleware(WorkflowMiddlewareDelegate next, ICommandSender commandSender, IIdentityGenerator identityGenerator)
public PersistWorkflowExecutionLogMiddleware(WorkflowMiddlewareDelegate next, ICommandSender commandSender, IIdentityGenerator identityGenerator) : base(next)
{
_next = next;
_commandSender = commandSender;
_identityGenerator = identityGenerator;
}
public async ValueTask InvokeAsync(WorkflowExecutionContext context)
public override async ValueTask InvokeAsync(WorkflowExecutionContext context)
{
// Invoke next middleware.
await _next(context);
await Next(context);
// Persist workflow execution log entries.
var entries = context.ExecutionLog.Select(x => new WorkflowExecutionLogRecord

View file

@ -5,20 +5,20 @@ using Elsa.Persistence.Commands;
using Elsa.Persistence.Entities;
using Elsa.Persistence.Requests;
using Elsa.Pipelines.WorkflowExecution;
using Elsa.Pipelines.WorkflowExecution.Components;
namespace Elsa.Persistence.Middleware.WorkflowExecution;
public static class PersistWorkflowInstanceMiddlewareExtensions
{
public static IWorkflowExecutionBuilder PersistWorkflows(this IWorkflowExecutionBuilder builder) => builder.UseMiddleware<PersistWorkflowInstanceMiddleware>();
public static IWorkflowExecutionBuilder UsePersistence(this IWorkflowExecutionBuilder builder) => builder.UseMiddleware<PersistWorkflowInstanceMiddleware>();
}
/// <summary>
/// Takes care of persisting a workflow instance after workflow execution.
/// </summary>
public class PersistWorkflowInstanceMiddleware : IWorkflowExecutionMiddleware
public class PersistWorkflowInstanceMiddleware : WorkflowExecutionMiddleware
{
private readonly WorkflowMiddlewareDelegate _next;
private readonly IRequestSender _requestSender;
private readonly ICommandSender _commandSender;
private readonly IWorkflowStateSerializer _workflowStateSerializer;
@ -33,9 +33,8 @@ public class PersistWorkflowInstanceMiddleware : IWorkflowExecutionMiddleware
IWorkflowStateSerializer workflowStateSerializer,
IPayloadSerializer payloadSerializer,
IIdentityGenerator identityGenerator,
ISystemClock clock)
ISystemClock clock) : base(next)
{
_next = next;
_requestSender = requestSender;
_commandSender = commandSender;
_workflowStateSerializer = workflowStateSerializer;
@ -44,7 +43,7 @@ public class PersistWorkflowInstanceMiddleware : IWorkflowExecutionMiddleware
_clock = clock;
}
public async ValueTask InvokeAsync(WorkflowExecutionContext context)
public override async ValueTask InvokeAsync(WorkflowExecutionContext context)
{
var cancellationToken = context.CancellationToken;
var workflow = context.Workflow;
@ -85,7 +84,7 @@ public class PersistWorkflowInstanceMiddleware : IWorkflowExecutionMiddleware
await _commandSender.ExecuteAsync(new SaveWorkflowInstance(workflowInstance), cancellationToken);
// Invoke next middleware.
await _next(context);
await Next(context);
// Update workflow instance.
var workflowState = _workflowStateSerializer.ReadState(context);

View file

@ -12,8 +12,8 @@ public class ReplaceWorkflowTriggersHandler : ICommandHandler<ReplaceWorkflowTri
public async Task<Unit> HandleAsync(ReplaceWorkflowTriggers command, CancellationToken cancellationToken)
{
var workflowDefinitionIds = command.WorkflowTriggers.Select(x => x.WorkflowDefinitionId).Distinct().ToList();
await _store.DeleteWhereAsync(x => workflowDefinitionIds.Contains(x.WorkflowDefinitionId), cancellationToken);
var definitionId = command.Workflow.Identity.DefinitionId;
await _store.DeleteWhereAsync(x => x.WorkflowDefinitionId == definitionId, cancellationToken);
await _store.SaveManyAsync(command.WorkflowTriggers, cancellationToken);
return Unit.Instance;

View file

@ -1,4 +1,5 @@
syntax = "proto3";
import "google/protobuf/wrappers.proto";
package Elsa.Runtime.ProtoActor.Messages;
option csharp_namespace = "Elsa.Runtime.ProtoActor.Messages";
@ -34,9 +35,9 @@ message DispatchWorkflowResponse {
}
message HandleStimulusRequest {
string activityTypeName = 1;
optional string hash = 2;
map<string, Json> data = 3;
string activityTypeName = 1;
optional string hash = 2;
map<string, Json> data = 3;
}
message Bookmark {
@ -46,7 +47,7 @@ message Bookmark {
optional string payload = 4;
string activityId = 5;
string activityInstanceId = 6;
optional string callbackMethodName = 7;
optional google.protobuf.StringValue callbackMethodName = 7;
}
message Json {

View file

@ -40,10 +40,10 @@ public static class ServiceCollectionExtensions
// Workflow engine.
.AddSingleton<IWorkflowServer, WorkflowServer>()
// Domain event handlers
// Domain event handlers.
.AddNotificationHandlersFrom(typeof(ServiceCollectionExtensions))
// Hosted Sercices
// Hosted Services.
.AddHostedService<RegisterDescriptorsHostedService>()
.AddHostedService<RegisterExpressionSyntaxDescriptorsHostedService>();
;

View file

@ -0,0 +1,33 @@
using Elsa.Contracts;
using Elsa.Mediator.Contracts;
using Elsa.Models;
using Elsa.Pipelines.WorkflowExecution;
using Elsa.Pipelines.WorkflowExecution.Components;
using Elsa.Runtime.Notifications;
namespace Elsa.Runtime.Middleware;
public static class PersistWorkflowExecutionLogMiddlewareExtensions
{
public static IWorkflowExecutionBuilder UseWorkflowExecutionEvents(this IWorkflowExecutionBuilder builder) => builder.UseMiddleware<WorkflowExecutionEventsMiddleware>();
}
/// <summary>
/// Processes collected bookmarks.
/// </summary>
public class WorkflowExecutionEventsMiddleware : WorkflowExecutionMiddleware
{
private readonly IEventPublisher _eventPublisher;
public WorkflowExecutionEventsMiddleware(WorkflowMiddlewareDelegate next, IEventPublisher eventPublisher) : base(next)
{
_eventPublisher = eventPublisher;
}
public override async ValueTask InvokeAsync(WorkflowExecutionContext context)
{
await _eventPublisher.PublishAsync(new WorkflowExecuting(context));
await Next(context);
await _eventPublisher.PublishAsync(new WorkflowExecuted(context));
}
}

View file

@ -0,0 +1,6 @@
using Elsa.Mediator.Contracts;
using Elsa.Models;
namespace Elsa.Runtime.Notifications;
public record WorkflowExecuted(WorkflowExecutionContext WorkflowExecutionContext) : INotification;

View file

@ -0,0 +1,6 @@
using Elsa.Mediator.Contracts;
using Elsa.Models;
namespace Elsa.Runtime.Notifications;
public record WorkflowExecuting(WorkflowExecutionContext WorkflowExecutionContext) : INotification;

View file

@ -60,22 +60,22 @@ public class TriggerIndexer : ITriggerIndexer
// Only stream workflows from providers that are not "dynamic" (such as DatabaseWorkflowProvider).
var workflows = _workflowRegistry.StreamAllAsync(WorkflowRegistry.SkipDynamicProviders, cancellationToken);
var collectedTriggers = new List<WorkflowTrigger>();
//var collectedTriggers = new List<WorkflowTrigger>();
await foreach (var workflow in workflows.WithCancellation(cancellationToken))
{
var triggers = await GetTriggersAsync(workflow, cancellationToken).ToListAsync(cancellationToken);
collectedTriggers.AddRange(triggers);
//var triggers = await GetTriggersAsync(workflow, cancellationToken).ToListAsync(cancellationToken);
await IndexTriggersAsync(workflow, cancellationToken);
}
// Replace triggers for the specified workflow.
await _commandSender.ExecuteAsync(new ReplaceWorkflowTriggers(collectedTriggers), cancellationToken);
// // Replace triggers for the specified workflow.
// await _commandSender.ExecuteAsync(new ReplaceWorkflowTriggers(collectedTriggers), cancellationToken);
stopwatch.Stop();
_logger.LogInformation("Finished indexing workflow triggers in {ElapsedTime}", stopwatch.Elapsed);
// Publish event.
await _eventPublisher.PublishAsync(new TriggerIndexingFinished(collectedTriggers), cancellationToken);
// // Publish event.
// await _eventPublisher.PublishAsync(new TriggerIndexingFinished(collectedTriggers), cancellationToken);
}
public async Task<IEnumerable<WorkflowTrigger>> IndexTriggersAsync(Workflow workflow, CancellationToken cancellationToken = default)
@ -84,7 +84,7 @@ public class TriggerIndexer : ITriggerIndexer
var triggers = await GetTriggersAsync(workflow, cancellationToken).ToListAsync(cancellationToken);
// Replace triggers for the specified workflow.
await _commandSender.ExecuteAsync(new ReplaceWorkflowTriggers(triggers), cancellationToken);
await _commandSender.ExecuteAsync(new ReplaceWorkflowTriggers(workflow, triggers), cancellationToken);
// Publish event.
await _eventPublisher.PublishAsync(new TriggerIndexingFinished(triggers), cancellationToken);

View file

@ -16,6 +16,7 @@ using Elsa.Persistence.EntityFrameworkCore.Extensions;
using Elsa.Persistence.EntityFrameworkCore.Sqlite;
using Elsa.Persistence.Middleware.WorkflowExecution;
using Elsa.Pipelines.WorkflowExecution.Components;
using Elsa.Runtime.Middleware;
using Elsa.Runtime.ProtoActor.Extensions;
using Elsa.Samples.Web1.Activities;
using Elsa.Samples.Web1.Workflows;
@ -60,6 +61,7 @@ services
.AddActivity<If>()
.AddActivity<HttpTrigger>()
.AddActivity<Flowchart>()
.AddActivity<Delay>()
;
// Register available triggers.
@ -88,8 +90,9 @@ wellKnownTypeRegistry.RegisterType<string>("string");
// Configure workflow engine execution pipeline.
serviceProvider.ConfigureDefaultWorkflowExecutionPipeline(pipeline => pipeline
.PersistWorkflows()
.PersistWorkflowExecutionLog()
.UsePersistence()
.UseWorkflowExecutionLogPersistence()
.UseWorkflowExecutionEvents()
.UseActivityScheduler()
);

View file

@ -27,7 +27,7 @@ var app = builder.Build();
// Configure workflow engine execution pipeline.
app.Services.ConfigureDefaultWorkflowExecutionPipeline(pipeline => pipeline
.PersistWorkflows()
.UsePersistence()
.UseActivityScheduler()
);

View file

@ -31,7 +31,7 @@ class Program
// Configure workflow engine execution pipeline.
.ConfigureDefaultWorkflowExecutionPipeline(pipeline => pipeline
.PersistWorkflows()
.UsePersistence()
.UseActivityScheduler()
);