Refactor workflow runtime to use ObsoleteWorkflowRuntime delegation
Replaced direct method implementations in workflow runtimes with delegations to the new `ObsoleteWorkflowRuntime` class, simplifying the codebase. This consolidates logic and aligns the runtimes under a unified deprecated API layer.
This commit is contained in:
parent
e27101032b
commit
ee883dbd07
|
|
@ -1,237 +1,32 @@
|
|||
using Elsa.Common.Models;
|
||||
using Elsa.Workflows.Management.Filters;
|
||||
using Elsa.Workflows.Models;
|
||||
using Elsa.Workflows.Runtime.Deprecated;
|
||||
using Elsa.Workflows.Runtime.Entities;
|
||||
using Elsa.Workflows.Runtime.Filters;
|
||||
using Elsa.Workflows.Runtime.Matches;
|
||||
using Elsa.Workflows.Runtime.Messages;
|
||||
using Elsa.Workflows.Runtime.Options;
|
||||
using Elsa.Workflows.Runtime.Parameters;
|
||||
using Elsa.Workflows.Runtime.Params;
|
||||
using Elsa.Workflows.Runtime.Requests;
|
||||
using Elsa.Workflows.Runtime.Results;
|
||||
using Elsa.Workflows.State;
|
||||
using Open.Linq.AsyncExtensions;
|
||||
|
||||
namespace Elsa.Workflows.Runtime.Distributed;
|
||||
|
||||
public partial class DistributedWorkflowRuntime
|
||||
{
|
||||
public async Task<CanStartWorkflowResult> CanStartWorkflowAsync(string definitionId, StartWorkflowRuntimeParams? options = null)
|
||||
{
|
||||
var cancellationToken = options?.CancellationToken ?? CancellationToken.None;
|
||||
var workflowGraph = await workflowDefinitionService.FindWorkflowGraphAsync(definitionId, options?.VersionOptions ?? VersionOptions.Published, cancellationToken);
|
||||
var workflow = workflowGraph!.Workflow;
|
||||
private readonly ObsoleteWorkflowRuntime _obsoleteApi;
|
||||
|
||||
var canStart = await workflowActivationStrategyEvaluator.CanStartWorkflowAsync(new()
|
||||
{
|
||||
Workflow = workflow,
|
||||
CorrelationId = options?.CorrelationId,
|
||||
CancellationToken = cancellationToken
|
||||
});
|
||||
|
||||
return new(null, canStart);
|
||||
}
|
||||
|
||||
public async Task<WorkflowExecutionResult> StartWorkflowAsync(string definitionId, StartWorkflowRuntimeParams? options = null)
|
||||
{
|
||||
var cancellationToken = options?.CancellationToken ?? CancellationToken.None;
|
||||
var client = await CreateClientAsync(options?.InstanceId, cancellationToken);
|
||||
var createRequest = new CreateAndRunWorkflowInstanceRequest
|
||||
{
|
||||
Properties = options?.Properties,
|
||||
CorrelationId = options?.CorrelationId,
|
||||
Input = options?.Input,
|
||||
WorkflowDefinitionHandle = WorkflowDefinitionHandle.ByDefinitionId(definitionId, options?.VersionOptions ?? VersionOptions.Published),
|
||||
ParentId = options?.ParentWorkflowInstanceId,
|
||||
TriggerActivityId = options?.TriggerActivityId
|
||||
};
|
||||
var response = await client.CreateAndRunInstanceAsync(createRequest, cancellationToken);
|
||||
return new(response.WorkflowInstanceId, response.Status, response.SubStatus, response.Bookmarks, response.Incidents);
|
||||
}
|
||||
|
||||
public async Task<ICollection<WorkflowExecutionResult>> StartWorkflowsAsync(string activityTypeName, object bookmarkPayload, TriggerWorkflowsOptions? options = null)
|
||||
{
|
||||
var cancellationToken = options?.CancellationToken ?? CancellationToken.None;
|
||||
var metadata = new StimulusMetadata
|
||||
{
|
||||
CorrelationId = options?.CorrelationId,
|
||||
WorkflowInstanceId = options?.WorkflowInstanceId,
|
||||
Properties = options?.Properties,
|
||||
ActivityInstanceId = options?.ActivityInstanceId,
|
||||
Input = options?.Input
|
||||
};
|
||||
var result = await stimulusSender.SendAsync(activityTypeName, bookmarkPayload, metadata, cancellationToken);
|
||||
var results = result.WorkflowInstanceResponses.Select(x => new WorkflowExecutionResult(x.WorkflowInstanceId, x.Status, x.SubStatus, x.Bookmarks, x.Incidents)).ToList();
|
||||
return results;
|
||||
}
|
||||
|
||||
public async Task<WorkflowExecutionResult?> TryStartWorkflowAsync(string definitionId, StartWorkflowRuntimeParams? options = null)
|
||||
{
|
||||
return await StartWorkflowAsync(definitionId, options);
|
||||
}
|
||||
|
||||
public async Task<WorkflowExecutionResult?> ResumeWorkflowAsync(string workflowInstanceId, ResumeWorkflowRuntimeParams? options = null)
|
||||
{
|
||||
var cancellationToken = options?.CancellationToken ?? CancellationToken.None;
|
||||
var workflowClient = await CreateClientAsync(workflowInstanceId, cancellationToken);
|
||||
var exists = await workflowClient.InstanceExistsAsync(cancellationToken);
|
||||
|
||||
if (!exists)
|
||||
return null;
|
||||
|
||||
var runWorkflowRequest = new RunWorkflowInstanceRequest
|
||||
{
|
||||
Input = options?.Input,
|
||||
Properties = options?.Properties,
|
||||
ActivityHandle = options?.ActivityHandle,
|
||||
BookmarkId = options?.BookmarkId
|
||||
};
|
||||
|
||||
var response = await workflowClient.RunInstanceAsync(runWorkflowRequest, cancellationToken);
|
||||
|
||||
return new(response.WorkflowInstanceId, response.Status, response.SubStatus, response.Bookmarks, response.Incidents);
|
||||
}
|
||||
|
||||
public async Task<ICollection<WorkflowExecutionResult>> ResumeWorkflowsAsync(string activityTypeName, object bookmarkPayload, TriggerWorkflowsOptions? options = null)
|
||||
{
|
||||
var cancellationToken = options?.CancellationToken ?? CancellationToken.None;
|
||||
var metadata = new StimulusMetadata
|
||||
{
|
||||
CorrelationId = options?.CorrelationId,
|
||||
WorkflowInstanceId = options?.WorkflowInstanceId,
|
||||
Properties = options?.Properties,
|
||||
ActivityInstanceId = options?.ActivityInstanceId,
|
||||
Input = options?.Input
|
||||
};
|
||||
var result = await stimulusSender.SendAsync(activityTypeName, bookmarkPayload, metadata, cancellationToken);
|
||||
var results = result.WorkflowInstanceResponses.Select(x => new WorkflowExecutionResult(x.WorkflowInstanceId, x.Status, x.SubStatus, x.Bookmarks, x.Incidents)).ToList();
|
||||
return results;
|
||||
}
|
||||
|
||||
public async Task<TriggerWorkflowsResult> TriggerWorkflowsAsync(string activityTypeName, object bookmarkPayload, TriggerWorkflowsOptions? options = null)
|
||||
{
|
||||
var cancellationToken = options?.CancellationToken ?? CancellationToken.None;
|
||||
var metadata = new StimulusMetadata
|
||||
{
|
||||
CorrelationId = options?.CorrelationId,
|
||||
WorkflowInstanceId = options?.WorkflowInstanceId,
|
||||
Properties = options?.Properties,
|
||||
ActivityInstanceId = options?.ActivityInstanceId,
|
||||
Input = options?.Input
|
||||
};
|
||||
var result = await stimulusSender.SendAsync(activityTypeName, bookmarkPayload, metadata, cancellationToken);
|
||||
var results = result.WorkflowInstanceResponses.Select(x => new WorkflowExecutionResult(x.WorkflowInstanceId, x.Status, x.SubStatus, x.Bookmarks, x.Incidents)).ToList();
|
||||
return new(results);
|
||||
}
|
||||
|
||||
public async Task<WorkflowExecutionResult> ExecuteWorkflowAsync(WorkflowMatch match, ExecuteWorkflowParams? options = default)
|
||||
{
|
||||
var cancellationToken = options?.CancellationToken ?? CancellationToken.None;
|
||||
if (match is StartableWorkflowMatch collectedStartableWorkflow)
|
||||
{
|
||||
var startOptions = new StartWorkflowRuntimeParams
|
||||
{
|
||||
CorrelationId = collectedStartableWorkflow.CorrelationId,
|
||||
Input = options?.Input,
|
||||
Properties = options?.Properties,
|
||||
VersionOptions = VersionOptions.Published,
|
||||
TriggerActivityId = collectedStartableWorkflow.ActivityId,
|
||||
CancellationToken = cancellationToken
|
||||
};
|
||||
|
||||
var startResult = await StartWorkflowAsync(collectedStartableWorkflow.DefinitionId!, startOptions);
|
||||
return startResult with
|
||||
{
|
||||
TriggeredActivityId = collectedStartableWorkflow.ActivityId
|
||||
};
|
||||
}
|
||||
|
||||
var collectedResumableWorkflow = (match as ResumableWorkflowMatch)!;
|
||||
var runtimeOptions = new ResumeWorkflowRuntimeParams
|
||||
{
|
||||
CorrelationId = collectedResumableWorkflow.CorrelationId,
|
||||
BookmarkId = collectedResumableWorkflow.BookmarkId,
|
||||
Input = options?.Input,
|
||||
Properties = options?.Properties,
|
||||
CancellationToken = cancellationToken,
|
||||
};
|
||||
|
||||
return (await ResumeWorkflowAsync(collectedResumableWorkflow.WorkflowInstanceId, runtimeOptions))!;
|
||||
}
|
||||
|
||||
public async Task<CancellationResult> CancelWorkflowAsync(string workflowInstanceId, CancellationToken cancellationToken = default)
|
||||
{
|
||||
var client = await CreateClientAsync(workflowInstanceId, cancellationToken);
|
||||
await client.CancelAsync(cancellationToken);
|
||||
return new(true);
|
||||
}
|
||||
|
||||
public async Task<IEnumerable<WorkflowMatch>> FindWorkflowsAsync(WorkflowsFilter filter, CancellationToken cancellationToken = default)
|
||||
{
|
||||
var startableWorkflows = await FindStartableWorkflowsAsync(filter, cancellationToken);
|
||||
var resumableWorkflows = await FindResumableWorkflowsAsync(filter, cancellationToken);
|
||||
var results = startableWorkflows.Concat(resumableWorkflows).ToList();
|
||||
return results;
|
||||
}
|
||||
|
||||
public async Task<WorkflowState?> ExportWorkflowStateAsync(string workflowInstanceId, CancellationToken cancellationToken = default)
|
||||
{
|
||||
var client = await CreateClientAsync(workflowInstanceId, cancellationToken);
|
||||
return await client.ExportStateAsync(cancellationToken);
|
||||
}
|
||||
|
||||
public async Task ImportWorkflowStateAsync(WorkflowState workflowState, CancellationToken cancellationToken = default)
|
||||
{
|
||||
var client = await CreateClientAsync(workflowState.Id, cancellationToken);
|
||||
await client.ImportStateAsync(workflowState, cancellationToken);
|
||||
}
|
||||
|
||||
public async Task UpdateBookmarkAsync(StoredBookmark bookmark, CancellationToken cancellationToken = default)
|
||||
{
|
||||
await bookmarkStore.SaveAsync(bookmark, cancellationToken);
|
||||
}
|
||||
|
||||
public async Task<long> CountRunningWorkflowsAsync(CountRunningWorkflowsRequest request, CancellationToken cancellationToken = default)
|
||||
{
|
||||
var filter = new WorkflowInstanceFilter
|
||||
{
|
||||
DefinitionId = request.DefinitionId,
|
||||
Version = request.Version,
|
||||
CorrelationId = request.CorrelationId,
|
||||
WorkflowStatus = WorkflowStatus.Running
|
||||
};
|
||||
return await workflowInstanceStore.CountAsync(filter, cancellationToken);
|
||||
}
|
||||
|
||||
private async Task<IEnumerable<WorkflowMatch>> FindStartableWorkflowsAsync(WorkflowsFilter filter, CancellationToken cancellationToken = default)
|
||||
{
|
||||
var stimulusHash = stimulusHasher.Hash(filter.ActivityTypeName, filter.BookmarkPayload, filter.Options.ActivityInstanceId);
|
||||
var triggerBoundWorkflows = await triggerBoundWorkflowService.FindManyAsync(stimulusHash, cancellationToken).ToList();
|
||||
var correlationId = filter.Options.CorrelationId;
|
||||
|
||||
var query =
|
||||
from triggerBoundWorkflow in triggerBoundWorkflows
|
||||
from trigger in triggerBoundWorkflow.Triggers
|
||||
select new StartableWorkflowMatch(correlationId, trigger.ActivityId, triggerBoundWorkflow.WorkflowGraph.Workflow.Identity.DefinitionId, filter.BookmarkPayload);
|
||||
|
||||
return query.ToList();
|
||||
}
|
||||
|
||||
private async Task<IEnumerable<WorkflowMatch>> FindResumableWorkflowsAsync(WorkflowsFilter filter, CancellationToken cancellationToken)
|
||||
{
|
||||
var bookmarkOptions = new FindBookmarkOptions
|
||||
{
|
||||
CorrelationId = filter.Options.CorrelationId,
|
||||
WorkflowInstanceId = filter.Options.WorkflowInstanceId,
|
||||
ActivityInstanceId = filter.Options.ActivityInstanceId
|
||||
};
|
||||
var bookmarkBoundWorkflows = await bookmarkBoundWorkflowService.FindManyAsync(filter.ActivityTypeName, filter.BookmarkPayload, bookmarkOptions, cancellationToken).ToList();
|
||||
|
||||
return (
|
||||
from bookmarkBoundWorkflow in bookmarkBoundWorkflows
|
||||
from bookmark in bookmarkBoundWorkflow.Bookmarks
|
||||
select new ResumableWorkflowMatch(bookmarkBoundWorkflow.WorkflowInstanceId, bookmark.CorrelationId, bookmark.Id, bookmark.Payload))
|
||||
.ToList();
|
||||
}
|
||||
public Task<CanStartWorkflowResult> CanStartWorkflowAsync(string definitionId, StartWorkflowRuntimeParams? options = null) => _obsoleteApi.CanStartWorkflowAsync(definitionId, options);
|
||||
public Task<WorkflowExecutionResult> StartWorkflowAsync(string definitionId, StartWorkflowRuntimeParams? options = null) => _obsoleteApi.StartWorkflowAsync(definitionId, options);
|
||||
public Task<ICollection<WorkflowExecutionResult>> StartWorkflowsAsync(string activityTypeName, object bookmarkPayload, TriggerWorkflowsOptions? options = null) => _obsoleteApi.StartWorkflowsAsync(activityTypeName, bookmarkPayload, options);
|
||||
public Task<WorkflowExecutionResult?> TryStartWorkflowAsync(string definitionId, StartWorkflowRuntimeParams? options = null) => _obsoleteApi.TryStartWorkflowAsync(definitionId, options);
|
||||
public Task<WorkflowExecutionResult?> ResumeWorkflowAsync(string workflowInstanceId, ResumeWorkflowRuntimeParams? options = null) => _obsoleteApi.ResumeWorkflowAsync(workflowInstanceId, options);
|
||||
public Task<ICollection<WorkflowExecutionResult>> ResumeWorkflowsAsync(string activityTypeName, object bookmarkPayload, TriggerWorkflowsOptions? options = null) => _obsoleteApi.ResumeWorkflowsAsync(activityTypeName, bookmarkPayload, options);
|
||||
public Task<TriggerWorkflowsResult> TriggerWorkflowsAsync(string activityTypeName, object bookmarkPayload, TriggerWorkflowsOptions? options = null) => _obsoleteApi.TriggerWorkflowsAsync(activityTypeName, bookmarkPayload, options);
|
||||
public Task<WorkflowExecutionResult> ExecuteWorkflowAsync(WorkflowMatch match, ExecuteWorkflowParams? options = default) => _obsoleteApi.ExecuteWorkflowAsync(match, options);
|
||||
public Task<CancellationResult> CancelWorkflowAsync(string workflowInstanceId, CancellationToken cancellationToken = default) => _obsoleteApi.CancelWorkflowAsync(workflowInstanceId, cancellationToken);
|
||||
public Task<IEnumerable<WorkflowMatch>> FindWorkflowsAsync(WorkflowsFilter filter, CancellationToken cancellationToken = default) => _obsoleteApi.FindWorkflowsAsync(filter, cancellationToken);
|
||||
public Task<WorkflowState?> ExportWorkflowStateAsync(string workflowInstanceId, CancellationToken cancellationToken = default) => _obsoleteApi.ExportWorkflowStateAsync(workflowInstanceId, cancellationToken);
|
||||
public Task ImportWorkflowStateAsync(WorkflowState workflowState, CancellationToken cancellationToken = default) => _obsoleteApi.ImportWorkflowStateAsync(workflowState, cancellationToken);
|
||||
public Task UpdateBookmarkAsync(StoredBookmark bookmark, CancellationToken cancellationToken = default) => _obsoleteApi.UpdateBookmarkAsync(bookmark, cancellationToken);
|
||||
public Task<long> CountRunningWorkflowsAsync(CountRunningWorkflowsRequest request, CancellationToken cancellationToken = default) => _obsoleteApi.CountRunningWorkflowsAsync(request, cancellationToken);
|
||||
}
|
||||
|
|
@ -1,4 +1,5 @@
|
|||
using Elsa.Workflows.Management;
|
||||
using Elsa.Workflows.Runtime.Deprecated;
|
||||
using Microsoft.Extensions.DependencyInjection;
|
||||
|
||||
namespace Elsa.Workflows.Runtime.Distributed;
|
||||
|
|
@ -6,18 +7,21 @@ namespace Elsa.Workflows.Runtime.Distributed;
|
|||
/// <summary>
|
||||
/// Represents a distributed workflow runtime that can create <see cref="IWorkflowClient"/> instances connected to a workflow instance.
|
||||
/// </summary>
|
||||
public partial class DistributedWorkflowRuntime(
|
||||
IServiceProvider serviceProvider,
|
||||
IIdentityGenerator identityGenerator,
|
||||
IWorkflowDefinitionService workflowDefinitionService,
|
||||
IWorkflowActivationStrategyEvaluator workflowActivationStrategyEvaluator,
|
||||
IStimulusSender stimulusSender,
|
||||
IStimulusHasher stimulusHasher,
|
||||
IBookmarkStore bookmarkStore,
|
||||
IWorkflowInstanceStore workflowInstanceStore,
|
||||
ITriggerBoundWorkflowService triggerBoundWorkflowService,
|
||||
IBookmarkBoundWorkflowService bookmarkBoundWorkflowService) : IWorkflowRuntime
|
||||
public partial class DistributedWorkflowRuntime : IWorkflowRuntime
|
||||
{
|
||||
private readonly IServiceProvider _serviceProvider;
|
||||
private readonly IIdentityGenerator _identityGenerator;
|
||||
|
||||
/// <summary>
|
||||
/// Represents a distributed workflow runtime that can create <see cref="IWorkflowClient"/> instances connected to a workflow instance.
|
||||
/// </summary>
|
||||
public DistributedWorkflowRuntime(IServiceProvider serviceProvider, IIdentityGenerator identityGenerator)
|
||||
{
|
||||
_serviceProvider = serviceProvider;
|
||||
_identityGenerator = identityGenerator;
|
||||
_obsoleteApi = ActivatorUtilities.CreateInstance<ObsoleteWorkflowRuntime>(serviceProvider, (Func<string?, CancellationToken, ValueTask<IWorkflowClient>>)CreateClientAsync);
|
||||
}
|
||||
|
||||
/// <inheritdoc />
|
||||
public async ValueTask<IWorkflowClient> CreateClientAsync(CancellationToken cancellationToken = default)
|
||||
{
|
||||
|
|
@ -27,8 +31,8 @@ public partial class DistributedWorkflowRuntime(
|
|||
/// <inheritdoc />
|
||||
public ValueTask<IWorkflowClient> CreateClientAsync(string? workflowInstanceId, CancellationToken cancellationToken = default)
|
||||
{
|
||||
workflowInstanceId ??= identityGenerator.GenerateId();
|
||||
var client = (IWorkflowClient)ActivatorUtilities.CreateInstance(serviceProvider, typeof(DistributedWorkflowClient), workflowInstanceId);
|
||||
workflowInstanceId ??= _identityGenerator.GenerateId();
|
||||
var client = (IWorkflowClient)ActivatorUtilities.CreateInstance(_serviceProvider, typeof(DistributedWorkflowClient), workflowInstanceId);
|
||||
return new(client);
|
||||
}
|
||||
}
|
||||
|
|
@ -41,7 +41,7 @@ internal class WorkflowInstance(
|
|||
|
||||
public override Task OnStarted()
|
||||
{
|
||||
_linkedTokenSource = new CancellationTokenSource();
|
||||
_linkedTokenSource = new();
|
||||
_linkedCancellationToken = CancellationTokenSource.CreateLinkedTokenSource(Context.CancellationToken, _linkedTokenSource.Token).Token;
|
||||
return Task.CompletedTask;
|
||||
}
|
||||
|
|
@ -61,7 +61,7 @@ internal class WorkflowInstance(
|
|||
if (result.IsFaulted)
|
||||
onError(result.Exception.Message);
|
||||
else
|
||||
respond(new CreateWorkflowInstanceResponse());
|
||||
respond(new());
|
||||
});
|
||||
|
||||
return Task.CompletedTask;
|
||||
|
|
@ -170,7 +170,7 @@ internal class WorkflowInstance(
|
|||
{
|
||||
await EnsureStateAsync();
|
||||
var json = mappers.WorkflowStateJsonMapper.Map(WorkflowState);
|
||||
return new ExportWorkflowStateResponse
|
||||
return new()
|
||||
{
|
||||
SerializedWorkflowState = json
|
||||
};
|
||||
|
|
@ -186,12 +186,21 @@ internal class WorkflowInstance(
|
|||
await workflowInstanceManager.SaveAsync(WorkflowState, Context.CancellationToken);
|
||||
}
|
||||
|
||||
public override Task<InstanceExistsResponse> InstanceExists()
|
||||
{
|
||||
var exists = _workflowInstanceId != null;
|
||||
return Task.FromResult(new InstanceExistsResponse
|
||||
{
|
||||
Exists = exists
|
||||
});
|
||||
}
|
||||
|
||||
private async Task<RunWorkflowResult> RunAsync(RunWorkflowOptions runWorkflowOptions)
|
||||
{
|
||||
if (_isRunning)
|
||||
{
|
||||
_queuedRunWorkflowOptions.Enqueue(runWorkflowOptions);
|
||||
return new RunWorkflowResult(null!, null!, null);
|
||||
return new(null!, null!, null);
|
||||
}
|
||||
|
||||
_isRunning = true;
|
||||
|
|
|
|||
|
|
@ -1,6 +1,5 @@
|
|||
using System.Diagnostics.CodeAnalysis;
|
||||
using Elsa.Common.Models;
|
||||
using Elsa.Workflows.Management.Filters;
|
||||
using Elsa.Workflows.Runtime.Deprecated;
|
||||
using Elsa.Workflows.Runtime.Distributed;
|
||||
using Elsa.Workflows.Runtime.Entities;
|
||||
using Elsa.Workflows.Runtime.Filters;
|
||||
using Elsa.Workflows.Runtime.Matches;
|
||||
|
|
@ -10,255 +9,25 @@ using Elsa.Workflows.Runtime.Params;
|
|||
using Elsa.Workflows.Runtime.Requests;
|
||||
using Elsa.Workflows.Runtime.Results;
|
||||
using Elsa.Workflows.State;
|
||||
using Open.Linq.AsyncExtensions;
|
||||
|
||||
namespace Elsa.Workflows.Runtime.ProtoActor.Services;
|
||||
|
||||
public partial class ProtoActorWorkflowRuntime
|
||||
{
|
||||
/// <inheritdoc />
|
||||
public async Task<CanStartWorkflowResult> CanStartWorkflowAsync(string definitionId, StartWorkflowRuntimeParams? options = null)
|
||||
{
|
||||
var cancellationToken = options?.CancellationToken ?? CancellationToken.None;
|
||||
var workflowGraph = await workflowDefinitionService.FindWorkflowGraphAsync(definitionId, options?.VersionOptions ?? VersionOptions.Published, cancellationToken);
|
||||
var workflow = workflowGraph!.Workflow;
|
||||
private readonly ObsoleteWorkflowRuntime _obsoleteApi;
|
||||
|
||||
var canStart = await workflowActivationStrategyEvaluator.CanStartWorkflowAsync(new()
|
||||
{
|
||||
Workflow = workflow,
|
||||
CorrelationId = options?.CorrelationId,
|
||||
CancellationToken = cancellationToken
|
||||
});
|
||||
|
||||
return new(null, canStart);
|
||||
}
|
||||
|
||||
/// <inheritdoc />
|
||||
public async Task<WorkflowExecutionResult?> TryStartWorkflowAsync(string definitionId, StartWorkflowRuntimeParams? options = null)
|
||||
{
|
||||
var cancellationToken = options?.CancellationToken ?? CancellationToken.None;
|
||||
var client = (LocalWorkflowClient)await CreateClientAsync(options?.InstanceId, cancellationToken);
|
||||
var createRequest = new Messages.CreateAndRunWorkflowInstanceRequest
|
||||
{
|
||||
Properties = options?.Properties,
|
||||
CorrelationId = options?.CorrelationId,
|
||||
Input = options?.Input,
|
||||
WorkflowDefinitionHandle = Workflows.Models.WorkflowDefinitionHandle.ByDefinitionId(definitionId, options?.VersionOptions ?? VersionOptions.Published),
|
||||
ParentId = options?.ParentWorkflowInstanceId,
|
||||
TriggerActivityId = options?.TriggerActivityId
|
||||
};
|
||||
var response = await client.CreateAndRunInstanceAsync(createRequest, cancellationToken);
|
||||
return new(response.WorkflowInstanceId, response.Status, response.SubStatus, response.Bookmarks, response.Incidents);
|
||||
}
|
||||
|
||||
/// <inheritdoc />
|
||||
public async Task<WorkflowExecutionResult> StartWorkflowAsync(string definitionId, StartWorkflowRuntimeParams? options = null)
|
||||
{
|
||||
var cancellationToken = options?.CancellationToken ?? CancellationToken.None;
|
||||
var client = (LocalWorkflowClient)await CreateClientAsync(options?.InstanceId, cancellationToken);
|
||||
var createRequest = new Messages.CreateAndRunWorkflowInstanceRequest
|
||||
{
|
||||
Properties = options?.Properties,
|
||||
CorrelationId = options?.CorrelationId,
|
||||
Input = options?.Input,
|
||||
WorkflowDefinitionHandle = Workflows.Models.WorkflowDefinitionHandle.ByDefinitionId(definitionId, options?.VersionOptions ?? VersionOptions.Published),
|
||||
ParentId = options?.ParentWorkflowInstanceId,
|
||||
TriggerActivityId = options?.TriggerActivityId
|
||||
};
|
||||
var response = await client.CreateAndRunInstanceAsync(createRequest, cancellationToken);
|
||||
return new(response.WorkflowInstanceId, response.Status, response.SubStatus, response.Bookmarks, response.Incidents);
|
||||
}
|
||||
|
||||
/// <inheritdoc />
|
||||
public async Task<ICollection<WorkflowExecutionResult>> StartWorkflowsAsync(string activityTypeName, object bookmarkPayload, TriggerWorkflowsOptions? options = null)
|
||||
{
|
||||
var cancellationToken = options?.CancellationToken ?? CancellationToken.None;
|
||||
var metadata = new StimulusMetadata
|
||||
{
|
||||
CorrelationId = options?.CorrelationId,
|
||||
WorkflowInstanceId = options?.WorkflowInstanceId,
|
||||
Properties = options?.Properties,
|
||||
ActivityInstanceId = options?.ActivityInstanceId,
|
||||
Input = options?.Input
|
||||
};
|
||||
var result = await stimulusSender.SendAsync(activityTypeName, bookmarkPayload, metadata, cancellationToken);
|
||||
var results = result.WorkflowInstanceResponses.Select(x => new WorkflowExecutionResult(x.WorkflowInstanceId, x.Status, x.SubStatus, x.Bookmarks, x.Incidents)).ToList();
|
||||
return results;
|
||||
}
|
||||
|
||||
/// <inheritdoc />
|
||||
public async Task<WorkflowExecutionResult?> ResumeWorkflowAsync(string workflowInstanceId, ResumeWorkflowRuntimeParams? options = null)
|
||||
{
|
||||
var cancellationToken = options?.CancellationToken ?? CancellationToken.None;
|
||||
var workflowClient = await CreateClientAsync(workflowInstanceId, cancellationToken);
|
||||
var exists = await workflowClient.InstanceExistsAsync(cancellationToken);
|
||||
|
||||
if (!exists)
|
||||
return null;
|
||||
|
||||
var runWorkflowRequest = new Messages.RunWorkflowInstanceRequest
|
||||
{
|
||||
Input = options?.Input,
|
||||
Properties = options?.Properties,
|
||||
ActivityHandle = options?.ActivityHandle,
|
||||
BookmarkId = options?.BookmarkId
|
||||
};
|
||||
|
||||
var response = await workflowClient.RunInstanceAsync(runWorkflowRequest, cancellationToken);
|
||||
|
||||
return new(response.WorkflowInstanceId, response.Status, response.SubStatus, response.Bookmarks, response.Incidents);
|
||||
}
|
||||
|
||||
/// <inheritdoc />
|
||||
public async Task<ICollection<WorkflowExecutionResult>> ResumeWorkflowsAsync(string activityTypeName, object bookmarkPayload, TriggerWorkflowsOptions? options = null)
|
||||
{
|
||||
var cancellationToken = options?.CancellationToken ?? CancellationToken.None;
|
||||
var metadata = new StimulusMetadata
|
||||
{
|
||||
CorrelationId = options?.CorrelationId,
|
||||
WorkflowInstanceId = options?.WorkflowInstanceId,
|
||||
Properties = options?.Properties,
|
||||
ActivityInstanceId = options?.ActivityInstanceId,
|
||||
Input = options?.Input
|
||||
};
|
||||
var result = await stimulusSender.SendAsync(activityTypeName, bookmarkPayload, metadata, cancellationToken);
|
||||
var results = result.WorkflowInstanceResponses.Select(x => new WorkflowExecutionResult(x.WorkflowInstanceId, x.Status, x.SubStatus, x.Bookmarks, x.Incidents)).ToList();
|
||||
return results;
|
||||
}
|
||||
|
||||
/// <inheritdoc />
|
||||
public async Task<TriggerWorkflowsResult> TriggerWorkflowsAsync(string activityTypeName, object bookmarkPayload, TriggerWorkflowsOptions? options = null)
|
||||
{
|
||||
var cancellationToken = options?.CancellationToken ?? CancellationToken.None;
|
||||
var metadata = new StimulusMetadata
|
||||
{
|
||||
CorrelationId = options?.CorrelationId,
|
||||
WorkflowInstanceId = options?.WorkflowInstanceId,
|
||||
Properties = options?.Properties,
|
||||
ActivityInstanceId = options?.ActivityInstanceId,
|
||||
Input = options?.Input
|
||||
};
|
||||
var result = await stimulusSender.SendAsync(activityTypeName, bookmarkPayload, metadata, cancellationToken);
|
||||
var results = result.WorkflowInstanceResponses.Select(x => new WorkflowExecutionResult(x.WorkflowInstanceId, x.Status, x.SubStatus, x.Bookmarks, x.Incidents)).ToList();
|
||||
return new(results);
|
||||
}
|
||||
|
||||
/// <inheritdoc />
|
||||
public async Task<WorkflowExecutionResult> ExecuteWorkflowAsync(WorkflowMatch match, ExecuteWorkflowParams? options = null)
|
||||
{
|
||||
var cancellationToken = options?.CancellationToken ?? CancellationToken.None;
|
||||
if (match is StartableWorkflowMatch collectedStartableWorkflow)
|
||||
{
|
||||
var startOptions = new StartWorkflowRuntimeParams
|
||||
{
|
||||
CorrelationId = collectedStartableWorkflow.CorrelationId,
|
||||
Input = options?.Input,
|
||||
Properties = options?.Properties,
|
||||
VersionOptions = VersionOptions.Published,
|
||||
TriggerActivityId = collectedStartableWorkflow.ActivityId,
|
||||
CancellationToken = cancellationToken
|
||||
};
|
||||
|
||||
var startResult = await StartWorkflowAsync(collectedStartableWorkflow.DefinitionId!, startOptions);
|
||||
return startResult with
|
||||
{
|
||||
TriggeredActivityId = collectedStartableWorkflow.ActivityId
|
||||
};
|
||||
}
|
||||
|
||||
var collectedResumableWorkflow = (match as ResumableWorkflowMatch)!;
|
||||
var runtimeOptions = new ResumeWorkflowRuntimeParams
|
||||
{
|
||||
CorrelationId = collectedResumableWorkflow.CorrelationId,
|
||||
BookmarkId = collectedResumableWorkflow.BookmarkId,
|
||||
Input = options?.Input,
|
||||
Properties = options?.Properties,
|
||||
CancellationToken = cancellationToken,
|
||||
};
|
||||
|
||||
return (await ResumeWorkflowAsync(collectedResumableWorkflow.WorkflowInstanceId, runtimeOptions))!;
|
||||
}
|
||||
|
||||
/// <inheritdoc />
|
||||
public async Task<CancellationResult> CancelWorkflowAsync(string workflowInstanceId, CancellationToken cancellationToken)
|
||||
{
|
||||
var client = await CreateClientAsync(workflowInstanceId, cancellationToken);
|
||||
await client.CancelAsync(cancellationToken);
|
||||
return new(true);
|
||||
}
|
||||
|
||||
/// <inheritdoc />
|
||||
public async Task<IEnumerable<WorkflowMatch>> FindWorkflowsAsync(WorkflowsFilter filter, CancellationToken cancellationToken = default)
|
||||
{
|
||||
var startableWorkflows = await FindStartableWorkflowsAsync(filter, cancellationToken);
|
||||
var resumableWorkflows = await FindResumableWorkflowsAsync(filter, cancellationToken);
|
||||
var results = startableWorkflows.Concat(resumableWorkflows).ToList();
|
||||
return results;
|
||||
}
|
||||
|
||||
/// <inheritdoc />
|
||||
[RequiresUnreferencedCode("Calls Elsa.Workflows.Contracts.IWorkflowStateSerializer.DeserializeAsync(String, CancellationToken)")]
|
||||
public async Task<WorkflowState?> ExportWorkflowStateAsync(string workflowInstanceId, CancellationToken cancellationToken = default)
|
||||
{
|
||||
var client = await CreateClientAsync(workflowInstanceId, cancellationToken);
|
||||
return await client.ExportStateAsync(cancellationToken);
|
||||
}
|
||||
|
||||
/// <inheritdoc />
|
||||
[RequiresUnreferencedCode("Calls Elsa.Workflows.Contracts.IWorkflowStateSerializer.SerializeAsync(WorkflowState, CancellationToken)")]
|
||||
public async Task ImportWorkflowStateAsync(WorkflowState workflowState, CancellationToken cancellationToken = default)
|
||||
{
|
||||
var client = await CreateClientAsync(workflowState.Id, cancellationToken);
|
||||
await client.ImportStateAsync(workflowState, cancellationToken);
|
||||
}
|
||||
|
||||
/// <inheritdoc />
|
||||
public async Task UpdateBookmarkAsync(StoredBookmark bookmark, CancellationToken cancellationToken = default)
|
||||
{
|
||||
await bookmarkStore.SaveAsync(bookmark, cancellationToken);
|
||||
}
|
||||
|
||||
/// <inheritdoc />
|
||||
public async Task<long> CountRunningWorkflowsAsync(CountRunningWorkflowsRequest request, CancellationToken cancellationToken = default)
|
||||
{
|
||||
var filter = new WorkflowInstanceFilter
|
||||
{
|
||||
DefinitionId = request.DefinitionId,
|
||||
Version = request.Version,
|
||||
CorrelationId = request.CorrelationId,
|
||||
WorkflowStatus = WorkflowStatus.Running
|
||||
};
|
||||
return await workflowInstanceStore.CountAsync(filter, cancellationToken);
|
||||
}
|
||||
|
||||
private async Task<IEnumerable<WorkflowMatch>> FindStartableWorkflowsAsync(WorkflowsFilter filter, CancellationToken cancellationToken = default)
|
||||
{
|
||||
var stimulusHash = stimulusHasher.Hash(filter.ActivityTypeName, filter.BookmarkPayload, filter.Options.ActivityInstanceId);
|
||||
var triggerBoundWorkflows = await triggerBoundWorkflowService.FindManyAsync(stimulusHash, cancellationToken).ToList();
|
||||
var correlationId = filter.Options.CorrelationId;
|
||||
|
||||
var query =
|
||||
from triggerBoundWorkflow in triggerBoundWorkflows
|
||||
from trigger in triggerBoundWorkflow.Triggers
|
||||
select new StartableWorkflowMatch(correlationId, trigger.ActivityId, triggerBoundWorkflow.WorkflowGraph.Workflow.Identity.DefinitionId, filter.BookmarkPayload);
|
||||
|
||||
return query.ToList();
|
||||
}
|
||||
|
||||
private async Task<IEnumerable<WorkflowMatch>> FindResumableWorkflowsAsync(WorkflowsFilter filter, CancellationToken cancellationToken)
|
||||
{
|
||||
var bookmarkOptions = new FindBookmarkOptions
|
||||
{
|
||||
CorrelationId = filter.Options.CorrelationId,
|
||||
WorkflowInstanceId = filter.Options.WorkflowInstanceId,
|
||||
ActivityInstanceId = filter.Options.ActivityInstanceId
|
||||
};
|
||||
var bookmarkBoundWorkflows = await bookmarkBoundWorkflowService.FindManyAsync(filter.ActivityTypeName, filter.BookmarkPayload, bookmarkOptions, cancellationToken).ToList();
|
||||
|
||||
return (
|
||||
from bookmarkBoundWorkflow in bookmarkBoundWorkflows
|
||||
from bookmark in bookmarkBoundWorkflow.Bookmarks
|
||||
select new ResumableWorkflowMatch(bookmarkBoundWorkflow.WorkflowInstanceId, bookmark.CorrelationId, bookmark.Id, bookmark.Payload))
|
||||
.ToList();
|
||||
}
|
||||
public Task<CanStartWorkflowResult> CanStartWorkflowAsync(string definitionId, StartWorkflowRuntimeParams? options = null) => _obsoleteApi.CanStartWorkflowAsync(definitionId, options);
|
||||
public Task<WorkflowExecutionResult> StartWorkflowAsync(string definitionId, StartWorkflowRuntimeParams? options = null) => _obsoleteApi.StartWorkflowAsync(definitionId, options);
|
||||
public Task<ICollection<WorkflowExecutionResult>> StartWorkflowsAsync(string activityTypeName, object bookmarkPayload, TriggerWorkflowsOptions? options = null) => _obsoleteApi.StartWorkflowsAsync(activityTypeName, bookmarkPayload, options);
|
||||
public Task<WorkflowExecutionResult?> TryStartWorkflowAsync(string definitionId, StartWorkflowRuntimeParams? options = null) => _obsoleteApi.TryStartWorkflowAsync(definitionId, options);
|
||||
public Task<WorkflowExecutionResult?> ResumeWorkflowAsync(string workflowInstanceId, ResumeWorkflowRuntimeParams? options = null) => _obsoleteApi.ResumeWorkflowAsync(workflowInstanceId, options);
|
||||
public Task<ICollection<WorkflowExecutionResult>> ResumeWorkflowsAsync(string activityTypeName, object bookmarkPayload, TriggerWorkflowsOptions? options = null) => _obsoleteApi.ResumeWorkflowsAsync(activityTypeName, bookmarkPayload, options);
|
||||
public Task<TriggerWorkflowsResult> TriggerWorkflowsAsync(string activityTypeName, object bookmarkPayload, TriggerWorkflowsOptions? options = null) => _obsoleteApi.TriggerWorkflowsAsync(activityTypeName, bookmarkPayload, options);
|
||||
public Task<WorkflowExecutionResult> ExecuteWorkflowAsync(WorkflowMatch match, ExecuteWorkflowParams? options = default) => _obsoleteApi.ExecuteWorkflowAsync(match, options);
|
||||
public Task<CancellationResult> CancelWorkflowAsync(string workflowInstanceId, CancellationToken cancellationToken = default) => _obsoleteApi.CancelWorkflowAsync(workflowInstanceId, cancellationToken);
|
||||
public Task<IEnumerable<WorkflowMatch>> FindWorkflowsAsync(WorkflowsFilter filter, CancellationToken cancellationToken = default) => _obsoleteApi.FindWorkflowsAsync(filter, cancellationToken);
|
||||
public Task<WorkflowState?> ExportWorkflowStateAsync(string workflowInstanceId, CancellationToken cancellationToken = default) => _obsoleteApi.ExportWorkflowStateAsync(workflowInstanceId, cancellationToken);
|
||||
public Task ImportWorkflowStateAsync(WorkflowState workflowState, CancellationToken cancellationToken = default) => _obsoleteApi.ImportWorkflowStateAsync(workflowState, cancellationToken);
|
||||
public Task UpdateBookmarkAsync(StoredBookmark bookmark, CancellationToken cancellationToken = default) => _obsoleteApi.UpdateBookmarkAsync(bookmark, cancellationToken);
|
||||
public Task<long> CountRunningWorkflowsAsync(CountRunningWorkflowsRequest request, CancellationToken cancellationToken = default) => _obsoleteApi.CountRunningWorkflowsAsync(request, cancellationToken);
|
||||
}
|
||||
|
|
@ -1,4 +1,4 @@
|
|||
using Elsa.Workflows.Management;
|
||||
using Elsa.Workflows.Runtime.Deprecated;
|
||||
using Microsoft.Extensions.DependencyInjection;
|
||||
|
||||
namespace Elsa.Workflows.Runtime.ProtoActor.Services;
|
||||
|
|
@ -6,18 +6,22 @@ namespace Elsa.Workflows.Runtime.ProtoActor.Services;
|
|||
/// <summary>
|
||||
/// Represents a Proto.Actor implementation of the workflows runtime.
|
||||
/// </summary>
|
||||
public partial class ProtoActorWorkflowRuntime(
|
||||
IServiceProvider serviceProvider,
|
||||
IWorkflowDefinitionService workflowDefinitionService,
|
||||
IWorkflowActivationStrategyEvaluator workflowActivationStrategyEvaluator,
|
||||
IStimulusSender stimulusSender,
|
||||
IStimulusHasher stimulusHasher,
|
||||
ITriggerBoundWorkflowService triggerBoundWorkflowService,
|
||||
IBookmarkBoundWorkflowService bookmarkBoundWorkflowService,
|
||||
IBookmarkStore bookmarkStore,
|
||||
IWorkflowInstanceStore workflowInstanceStore,
|
||||
IIdentityGenerator identityGenerator) : IWorkflowRuntime
|
||||
public partial class ProtoActorWorkflowRuntime : IWorkflowRuntime
|
||||
{
|
||||
private readonly IServiceProvider _serviceProvider;
|
||||
private readonly IIdentityGenerator _identityGenerator;
|
||||
|
||||
/// <summary>
|
||||
/// Represents a Proto.Actor implementation of the workflows runtime.
|
||||
/// </summary>
|
||||
public ProtoActorWorkflowRuntime(IServiceProvider serviceProvider,
|
||||
IIdentityGenerator identityGenerator)
|
||||
{
|
||||
_serviceProvider = serviceProvider;
|
||||
_identityGenerator = identityGenerator;
|
||||
_obsoleteApi = ActivatorUtilities.CreateInstance<ObsoleteWorkflowRuntime>(serviceProvider, (Func<string?, CancellationToken, ValueTask<IWorkflowClient>>)CreateClientAsync);
|
||||
}
|
||||
|
||||
/// <inheritdoc />
|
||||
public async ValueTask<IWorkflowClient> CreateClientAsync(CancellationToken cancellationToken = default)
|
||||
{
|
||||
|
|
@ -27,8 +31,8 @@ public partial class ProtoActorWorkflowRuntime(
|
|||
/// <inheritdoc />
|
||||
public ValueTask<IWorkflowClient> CreateClientAsync(string? workflowInstanceId, CancellationToken cancellationToken = default)
|
||||
{
|
||||
workflowInstanceId ??= identityGenerator.GenerateId();
|
||||
var client = (IWorkflowClient)ActivatorUtilities.CreateInstance(serviceProvider, typeof(ProtoActorWorkflowClient), workflowInstanceId);
|
||||
workflowInstanceId ??= _identityGenerator.GenerateId();
|
||||
var client = (IWorkflowClient)ActivatorUtilities.CreateInstance(_serviceProvider, typeof(ProtoActorWorkflowClient), workflowInstanceId);
|
||||
return new(client);
|
||||
}
|
||||
}
|
||||
|
|
@ -14,10 +14,10 @@ using Elsa.Workflows.Runtime.Results;
|
|||
using Elsa.Workflows.State;
|
||||
using Open.Linq.AsyncExtensions;
|
||||
|
||||
namespace Elsa.Workflows.Runtime.Distributed;
|
||||
namespace Elsa.Workflows.Runtime.Deprecated;
|
||||
|
||||
public class ObsoleteWorkflowRuntime(
|
||||
Func<string?, CancellationToken, Task<IWorkflowClient>> createClientAsync,
|
||||
Func<string?, CancellationToken, ValueTask<IWorkflowClient>> createClientAsync,
|
||||
IWorkflowDefinitionService workflowDefinitionService,
|
||||
IWorkflowActivationStrategyEvaluator workflowActivationStrategyEvaluator,
|
||||
IStimulusSender stimulusSender,
|
||||
|
|
@ -1,237 +1,32 @@
|
|||
using Elsa.Common.Models;
|
||||
using Elsa.Workflows.Management.Filters;
|
||||
using Elsa.Workflows.Models;
|
||||
using Elsa.Workflows.Runtime.Deprecated;
|
||||
using Elsa.Workflows.Runtime.Entities;
|
||||
using Elsa.Workflows.Runtime.Filters;
|
||||
using Elsa.Workflows.Runtime.Matches;
|
||||
using Elsa.Workflows.Runtime.Messages;
|
||||
using Elsa.Workflows.Runtime.Options;
|
||||
using Elsa.Workflows.Runtime.Parameters;
|
||||
using Elsa.Workflows.Runtime.Params;
|
||||
using Elsa.Workflows.Runtime.Requests;
|
||||
using Elsa.Workflows.Runtime.Results;
|
||||
using Elsa.Workflows.State;
|
||||
using Open.Linq.AsyncExtensions;
|
||||
|
||||
namespace Elsa.Workflows.Runtime;
|
||||
|
||||
public partial class LocalWorkflowRuntime
|
||||
{
|
||||
public async Task<CanStartWorkflowResult> CanStartWorkflowAsync(string definitionId, StartWorkflowRuntimeParams? options = null)
|
||||
{
|
||||
var cancellationToken = options?.CancellationToken ?? CancellationToken.None;
|
||||
var workflowGraph = await workflowDefinitionService.FindWorkflowGraphAsync(definitionId, options?.VersionOptions ?? VersionOptions.Published, cancellationToken);
|
||||
var workflow = workflowGraph!.Workflow;
|
||||
private readonly ObsoleteWorkflowRuntime _obsoleteApi;
|
||||
|
||||
var canStart = await workflowActivationStrategyEvaluator.CanStartWorkflowAsync(new()
|
||||
{
|
||||
Workflow = workflow,
|
||||
CorrelationId = options?.CorrelationId,
|
||||
CancellationToken = cancellationToken
|
||||
});
|
||||
|
||||
return new(null, canStart);
|
||||
}
|
||||
|
||||
public async Task<WorkflowExecutionResult> StartWorkflowAsync(string definitionId, StartWorkflowRuntimeParams? options = null)
|
||||
{
|
||||
var cancellationToken = options?.CancellationToken ?? CancellationToken.None;
|
||||
var client = await CreateClientAsync(options?.InstanceId, cancellationToken);
|
||||
var createRequest = new CreateAndRunWorkflowInstanceRequest
|
||||
{
|
||||
Properties = options?.Properties,
|
||||
CorrelationId = options?.CorrelationId,
|
||||
Input = options?.Input,
|
||||
WorkflowDefinitionHandle = WorkflowDefinitionHandle.ByDefinitionId(definitionId, options?.VersionOptions ?? VersionOptions.Published),
|
||||
ParentId = options?.ParentWorkflowInstanceId,
|
||||
TriggerActivityId = options?.TriggerActivityId
|
||||
};
|
||||
var response = await client.CreateAndRunInstanceAsync(createRequest, cancellationToken);
|
||||
return new(response.WorkflowInstanceId, response.Status, response.SubStatus, response.Bookmarks, response.Incidents);
|
||||
}
|
||||
|
||||
public async Task<ICollection<WorkflowExecutionResult>> StartWorkflowsAsync(string activityTypeName, object bookmarkPayload, TriggerWorkflowsOptions? options = null)
|
||||
{
|
||||
var cancellationToken = options?.CancellationToken ?? CancellationToken.None;
|
||||
var metadata = new StimulusMetadata
|
||||
{
|
||||
CorrelationId = options?.CorrelationId,
|
||||
WorkflowInstanceId = options?.WorkflowInstanceId,
|
||||
Properties = options?.Properties,
|
||||
ActivityInstanceId = options?.ActivityInstanceId,
|
||||
Input = options?.Input
|
||||
};
|
||||
var result = await stimulusSender.SendAsync(activityTypeName, bookmarkPayload, metadata, cancellationToken);
|
||||
var results = result.WorkflowInstanceResponses.Select(x => new WorkflowExecutionResult(x.WorkflowInstanceId, x.Status, x.SubStatus, x.Bookmarks, x.Incidents)).ToList();
|
||||
return results;
|
||||
}
|
||||
|
||||
public async Task<WorkflowExecutionResult?> TryStartWorkflowAsync(string definitionId, StartWorkflowRuntimeParams? options = null)
|
||||
{
|
||||
return await StartWorkflowAsync(definitionId, options);
|
||||
}
|
||||
|
||||
public async Task<WorkflowExecutionResult?> ResumeWorkflowAsync(string workflowInstanceId, ResumeWorkflowRuntimeParams? options = null)
|
||||
{
|
||||
var cancellationToken = options?.CancellationToken ?? CancellationToken.None;
|
||||
var workflowClient = await CreateClientAsync(workflowInstanceId, cancellationToken);
|
||||
var exists = await workflowClient.InstanceExistsAsync(cancellationToken);
|
||||
|
||||
if (!exists)
|
||||
return null;
|
||||
|
||||
var runWorkflowRequest = new RunWorkflowInstanceRequest
|
||||
{
|
||||
Input = options?.Input,
|
||||
Properties = options?.Properties,
|
||||
ActivityHandle = options?.ActivityHandle,
|
||||
BookmarkId = options?.BookmarkId
|
||||
};
|
||||
|
||||
var response = await workflowClient.RunInstanceAsync(runWorkflowRequest, cancellationToken);
|
||||
|
||||
return new(response.WorkflowInstanceId, response.Status, response.SubStatus, response.Bookmarks, response.Incidents);
|
||||
}
|
||||
|
||||
public async Task<ICollection<WorkflowExecutionResult>> ResumeWorkflowsAsync(string activityTypeName, object bookmarkPayload, TriggerWorkflowsOptions? options = null)
|
||||
{
|
||||
var cancellationToken = options?.CancellationToken ?? CancellationToken.None;
|
||||
var metadata = new StimulusMetadata
|
||||
{
|
||||
CorrelationId = options?.CorrelationId,
|
||||
WorkflowInstanceId = options?.WorkflowInstanceId,
|
||||
Properties = options?.Properties,
|
||||
ActivityInstanceId = options?.ActivityInstanceId,
|
||||
Input = options?.Input
|
||||
};
|
||||
var result = await stimulusSender.SendAsync(activityTypeName, bookmarkPayload, metadata, cancellationToken);
|
||||
var results = result.WorkflowInstanceResponses.Select(x => new WorkflowExecutionResult(x.WorkflowInstanceId, x.Status, x.SubStatus, x.Bookmarks, x.Incidents)).ToList();
|
||||
return results;
|
||||
}
|
||||
|
||||
public async Task<TriggerWorkflowsResult> TriggerWorkflowsAsync(string activityTypeName, object bookmarkPayload, TriggerWorkflowsOptions? options = null)
|
||||
{
|
||||
var cancellationToken = options?.CancellationToken ?? CancellationToken.None;
|
||||
var metadata = new StimulusMetadata
|
||||
{
|
||||
CorrelationId = options?.CorrelationId,
|
||||
WorkflowInstanceId = options?.WorkflowInstanceId,
|
||||
Properties = options?.Properties,
|
||||
ActivityInstanceId = options?.ActivityInstanceId,
|
||||
Input = options?.Input
|
||||
};
|
||||
var result = await stimulusSender.SendAsync(activityTypeName, bookmarkPayload, metadata, cancellationToken);
|
||||
var results = result.WorkflowInstanceResponses.Select(x => new WorkflowExecutionResult(x.WorkflowInstanceId, x.Status, x.SubStatus, x.Bookmarks, x.Incidents)).ToList();
|
||||
return new(results);
|
||||
}
|
||||
|
||||
public async Task<WorkflowExecutionResult> ExecuteWorkflowAsync(WorkflowMatch match, ExecuteWorkflowParams? options = default)
|
||||
{
|
||||
var cancellationToken = options?.CancellationToken ?? CancellationToken.None;
|
||||
if (match is StartableWorkflowMatch collectedStartableWorkflow)
|
||||
{
|
||||
var startOptions = new StartWorkflowRuntimeParams
|
||||
{
|
||||
CorrelationId = collectedStartableWorkflow.CorrelationId,
|
||||
Input = options?.Input,
|
||||
Properties = options?.Properties,
|
||||
VersionOptions = VersionOptions.Published,
|
||||
TriggerActivityId = collectedStartableWorkflow.ActivityId,
|
||||
CancellationToken = cancellationToken
|
||||
};
|
||||
|
||||
var startResult = await StartWorkflowAsync(collectedStartableWorkflow.DefinitionId!, startOptions);
|
||||
return startResult with
|
||||
{
|
||||
TriggeredActivityId = collectedStartableWorkflow.ActivityId
|
||||
};
|
||||
}
|
||||
|
||||
var collectedResumableWorkflow = (match as ResumableWorkflowMatch)!;
|
||||
var runtimeOptions = new ResumeWorkflowRuntimeParams
|
||||
{
|
||||
CorrelationId = collectedResumableWorkflow.CorrelationId,
|
||||
BookmarkId = collectedResumableWorkflow.BookmarkId,
|
||||
Input = options?.Input,
|
||||
Properties = options?.Properties,
|
||||
CancellationToken = cancellationToken,
|
||||
};
|
||||
|
||||
return (await ResumeWorkflowAsync(collectedResumableWorkflow.WorkflowInstanceId, runtimeOptions))!;
|
||||
}
|
||||
|
||||
public async Task<CancellationResult> CancelWorkflowAsync(string workflowInstanceId, CancellationToken cancellationToken = default)
|
||||
{
|
||||
var client = await CreateClientAsync(workflowInstanceId, cancellationToken);
|
||||
await client.CancelAsync(cancellationToken);
|
||||
return new(true);
|
||||
}
|
||||
|
||||
public async Task<IEnumerable<WorkflowMatch>> FindWorkflowsAsync(WorkflowsFilter filter, CancellationToken cancellationToken = default)
|
||||
{
|
||||
var startableWorkflows = await FindStartableWorkflowsAsync(filter, cancellationToken);
|
||||
var resumableWorkflows = await FindResumableWorkflowsAsync(filter, cancellationToken);
|
||||
var results = startableWorkflows.Concat(resumableWorkflows).ToList();
|
||||
return results;
|
||||
}
|
||||
|
||||
public async Task<WorkflowState?> ExportWorkflowStateAsync(string workflowInstanceId, CancellationToken cancellationToken = default)
|
||||
{
|
||||
var client = await CreateClientAsync(workflowInstanceId, cancellationToken);
|
||||
return await client.ExportStateAsync(cancellationToken);
|
||||
}
|
||||
|
||||
public async Task ImportWorkflowStateAsync(WorkflowState workflowState, CancellationToken cancellationToken = default)
|
||||
{
|
||||
var client = await CreateClientAsync(workflowState.Id, cancellationToken);
|
||||
await client.ImportStateAsync(workflowState, cancellationToken);
|
||||
}
|
||||
|
||||
public async Task UpdateBookmarkAsync(StoredBookmark bookmark, CancellationToken cancellationToken = default)
|
||||
{
|
||||
await bookmarkStore.SaveAsync(bookmark, cancellationToken);
|
||||
}
|
||||
|
||||
public async Task<long> CountRunningWorkflowsAsync(CountRunningWorkflowsRequest request, CancellationToken cancellationToken = default)
|
||||
{
|
||||
var filter = new WorkflowInstanceFilter
|
||||
{
|
||||
DefinitionId = request.DefinitionId,
|
||||
Version = request.Version,
|
||||
CorrelationId = request.CorrelationId,
|
||||
WorkflowStatus = WorkflowStatus.Running
|
||||
};
|
||||
return await workflowInstanceStore.CountAsync(filter, cancellationToken);
|
||||
}
|
||||
|
||||
private async Task<IEnumerable<WorkflowMatch>> FindStartableWorkflowsAsync(WorkflowsFilter filter, CancellationToken cancellationToken = default)
|
||||
{
|
||||
var stimulusHash = stimulusHasher.Hash(filter.ActivityTypeName, filter.BookmarkPayload, filter.Options.ActivityInstanceId);
|
||||
var triggerBoundWorkflows = await triggerBoundWorkflowService.FindManyAsync(stimulusHash, cancellationToken).ToList();
|
||||
var correlationId = filter.Options.CorrelationId;
|
||||
|
||||
var query =
|
||||
from triggerBoundWorkflow in triggerBoundWorkflows
|
||||
from trigger in triggerBoundWorkflow.Triggers
|
||||
select new StartableWorkflowMatch(correlationId, trigger.ActivityId, triggerBoundWorkflow.WorkflowGraph.Workflow.Identity.DefinitionId, filter.BookmarkPayload);
|
||||
|
||||
return query.ToList();
|
||||
}
|
||||
|
||||
private async Task<IEnumerable<WorkflowMatch>> FindResumableWorkflowsAsync(WorkflowsFilter filter, CancellationToken cancellationToken)
|
||||
{
|
||||
var bookmarkOptions = new FindBookmarkOptions
|
||||
{
|
||||
CorrelationId = filter.Options.CorrelationId,
|
||||
WorkflowInstanceId = filter.Options.WorkflowInstanceId,
|
||||
ActivityInstanceId = filter.Options.ActivityInstanceId
|
||||
};
|
||||
var bookmarkBoundWorkflows = await bookmarkBoundWorkflowService.FindManyAsync(filter.ActivityTypeName, filter.BookmarkPayload, bookmarkOptions, cancellationToken).ToList();
|
||||
|
||||
return (
|
||||
from bookmarkBoundWorkflow in bookmarkBoundWorkflows
|
||||
from bookmark in bookmarkBoundWorkflow.Bookmarks
|
||||
select new ResumableWorkflowMatch(bookmarkBoundWorkflow.WorkflowInstanceId, bookmark.CorrelationId, bookmark.Id, bookmark.Payload))
|
||||
.ToList();
|
||||
}
|
||||
public Task<CanStartWorkflowResult> CanStartWorkflowAsync(string definitionId, StartWorkflowRuntimeParams? options = null) => _obsoleteApi.CanStartWorkflowAsync(definitionId, options);
|
||||
public Task<WorkflowExecutionResult> StartWorkflowAsync(string definitionId, StartWorkflowRuntimeParams? options = null) => _obsoleteApi.StartWorkflowAsync(definitionId, options);
|
||||
public Task<ICollection<WorkflowExecutionResult>> StartWorkflowsAsync(string activityTypeName, object bookmarkPayload, TriggerWorkflowsOptions? options = null) => _obsoleteApi.StartWorkflowsAsync(activityTypeName, bookmarkPayload, options);
|
||||
public Task<WorkflowExecutionResult?> TryStartWorkflowAsync(string definitionId, StartWorkflowRuntimeParams? options = null) => _obsoleteApi.TryStartWorkflowAsync(definitionId, options);
|
||||
public Task<WorkflowExecutionResult?> ResumeWorkflowAsync(string workflowInstanceId, ResumeWorkflowRuntimeParams? options = null) => _obsoleteApi.ResumeWorkflowAsync(workflowInstanceId, options);
|
||||
public Task<ICollection<WorkflowExecutionResult>> ResumeWorkflowsAsync(string activityTypeName, object bookmarkPayload, TriggerWorkflowsOptions? options = null) => _obsoleteApi.ResumeWorkflowsAsync(activityTypeName, bookmarkPayload, options);
|
||||
public Task<TriggerWorkflowsResult> TriggerWorkflowsAsync(string activityTypeName, object bookmarkPayload, TriggerWorkflowsOptions? options = null) => _obsoleteApi.TriggerWorkflowsAsync(activityTypeName, bookmarkPayload, options);
|
||||
public Task<WorkflowExecutionResult> ExecuteWorkflowAsync(WorkflowMatch match, ExecuteWorkflowParams? options = default) => _obsoleteApi.ExecuteWorkflowAsync(match, options);
|
||||
public Task<CancellationResult> CancelWorkflowAsync(string workflowInstanceId, CancellationToken cancellationToken = default) => _obsoleteApi.CancelWorkflowAsync(workflowInstanceId, cancellationToken);
|
||||
public Task<IEnumerable<WorkflowMatch>> FindWorkflowsAsync(WorkflowsFilter filter, CancellationToken cancellationToken = default) => _obsoleteApi.FindWorkflowsAsync(filter, cancellationToken);
|
||||
public Task<WorkflowState?> ExportWorkflowStateAsync(string workflowInstanceId, CancellationToken cancellationToken = default) => _obsoleteApi.ExportWorkflowStateAsync(workflowInstanceId, cancellationToken);
|
||||
public Task ImportWorkflowStateAsync(WorkflowState workflowState, CancellationToken cancellationToken = default) => _obsoleteApi.ImportWorkflowStateAsync(workflowState, cancellationToken);
|
||||
public Task UpdateBookmarkAsync(StoredBookmark bookmark, CancellationToken cancellationToken = default) => _obsoleteApi.UpdateBookmarkAsync(bookmark, cancellationToken);
|
||||
public Task<long> CountRunningWorkflowsAsync(CountRunningWorkflowsRequest request, CancellationToken cancellationToken = default) => _obsoleteApi.CountRunningWorkflowsAsync(request, cancellationToken);
|
||||
}
|
||||
|
|
@ -1,4 +1,4 @@
|
|||
using Elsa.Workflows.Management;
|
||||
using Elsa.Workflows.Runtime.Deprecated;
|
||||
using Microsoft.Extensions.DependencyInjection;
|
||||
|
||||
namespace Elsa.Workflows.Runtime;
|
||||
|
|
@ -8,18 +8,23 @@ namespace Elsa.Workflows.Runtime;
|
|||
/// It does not support clustering and is intended for single-node deployments only.
|
||||
/// For distributed deployments, use Proto.Actor or another distributed runtime.
|
||||
/// </summary>
|
||||
public partial class LocalWorkflowRuntime(
|
||||
IServiceProvider serviceProvider,
|
||||
IIdentityGenerator identityGenerator,
|
||||
IWorkflowDefinitionService workflowDefinitionService,
|
||||
IWorkflowActivationStrategyEvaluator workflowActivationStrategyEvaluator,
|
||||
IStimulusSender stimulusSender,
|
||||
IStimulusHasher stimulusHasher,
|
||||
IBookmarkStore bookmarkStore,
|
||||
IWorkflowInstanceStore workflowInstanceStore,
|
||||
ITriggerBoundWorkflowService triggerBoundWorkflowService,
|
||||
IBookmarkBoundWorkflowService bookmarkBoundWorkflowService) : IWorkflowRuntime
|
||||
public partial class LocalWorkflowRuntime : IWorkflowRuntime
|
||||
{
|
||||
private readonly IServiceProvider _serviceProvider;
|
||||
private readonly IIdentityGenerator _identityGenerator;
|
||||
|
||||
/// <summary>
|
||||
/// Represents a local implementation of the distributed runtime for running workflows.
|
||||
/// It does not support clustering and is intended for single-node deployments only.
|
||||
/// For distributed deployments, use Proto.Actor or another distributed runtime.
|
||||
/// </summary>
|
||||
public LocalWorkflowRuntime(IServiceProvider serviceProvider, IIdentityGenerator identityGenerator)
|
||||
{
|
||||
_serviceProvider = serviceProvider;
|
||||
_identityGenerator = identityGenerator;
|
||||
_obsoleteApi = ActivatorUtilities.CreateInstance<ObsoleteWorkflowRuntime>(serviceProvider, (Func<string?, CancellationToken, ValueTask<IWorkflowClient>>)CreateClientAsync);
|
||||
}
|
||||
|
||||
/// <inheritdoc />
|
||||
public async ValueTask<IWorkflowClient> CreateClientAsync(CancellationToken cancellationToken = default)
|
||||
{
|
||||
|
|
@ -29,8 +34,8 @@ public partial class LocalWorkflowRuntime(
|
|||
/// <inheritdoc />
|
||||
public ValueTask<IWorkflowClient> CreateClientAsync(string? workflowInstanceId, CancellationToken cancellationToken = default)
|
||||
{
|
||||
workflowInstanceId ??= identityGenerator.GenerateId();
|
||||
var client = (IWorkflowClient)ActivatorUtilities.CreateInstance(serviceProvider, typeof(LocalWorkflowClient), workflowInstanceId);
|
||||
workflowInstanceId ??= _identityGenerator.GenerateId();
|
||||
var client = (IWorkflowClient)ActivatorUtilities.CreateInstance(_serviceProvider, typeof(LocalWorkflowClient), workflowInstanceId);
|
||||
return new(client);
|
||||
}
|
||||
}
|
||||
Loading…
Reference in a new issue