Fix Azure Service Bus queue worker
This fixes an issue with message / workflow correlation.
This commit is contained in:
parent
254ec93263
commit
f57b834d36
|
|
@ -1,12 +1,18 @@
|
|||
using System;
|
||||
using System.Linq;
|
||||
using System.Threading;
|
||||
using System.Threading.Tasks;
|
||||
using Elsa.Activities.AzureServiceBus.Triggers;
|
||||
using Elsa.Models;
|
||||
using Elsa.Persistence;
|
||||
using Elsa.Persistence.Specifications;
|
||||
using Elsa.Services;
|
||||
using Elsa.Triggers;
|
||||
using Microsoft.Azure.ServiceBus;
|
||||
using Microsoft.Azure.ServiceBus.Core;
|
||||
using Microsoft.Extensions.DependencyInjection;
|
||||
using Microsoft.Extensions.Logging;
|
||||
using Open.Linq.AsyncExtensions;
|
||||
|
||||
namespace Elsa.Activities.AzureServiceBus.Services
|
||||
{
|
||||
|
|
@ -40,16 +46,38 @@ namespace Elsa.Activities.AzureServiceBus.Services
|
|||
{
|
||||
using var scope = _serviceProvider.CreateScope();
|
||||
var workflowRunner = scope.ServiceProvider.GetRequiredService<IWorkflowRunner>();
|
||||
var queueName = _messageReceiver.Path;
|
||||
|
||||
async Task TriggerNewWorkflowAsync()
|
||||
{
|
||||
await workflowRunner!.TriggerWorkflowsAsync<MessageReceivedTrigger>(
|
||||
x => x.QueueName == queueName && x.CorrelationId == null,
|
||||
message,
|
||||
message.CorrelationId,
|
||||
cancellationToken: cancellationToken);
|
||||
}
|
||||
|
||||
Func<MessageReceivedTrigger, bool> predicate = string.IsNullOrWhiteSpace(message.CorrelationId)
|
||||
? x => x.QueueName == _messageReceiver.Path && x.CorrelationId == null
|
||||
: x => x.QueueName == _messageReceiver.Path && x.CorrelationId == message.CorrelationId;
|
||||
|
||||
await workflowRunner.TriggerWorkflowsAsync(
|
||||
predicate,
|
||||
message,
|
||||
message.CorrelationId,
|
||||
cancellationToken: cancellationToken);
|
||||
if (string.IsNullOrWhiteSpace(message.CorrelationId))
|
||||
{
|
||||
await TriggerNewWorkflowAsync();
|
||||
return;
|
||||
}
|
||||
|
||||
var workflowSelector = scope.ServiceProvider.GetRequiredService<IWorkflowSelector>();
|
||||
var workflowInstanceStore = scope.ServiceProvider.GetRequiredService<IWorkflowInstanceStore>();
|
||||
var correlatedWorkflowInstanceCount = await workflowInstanceStore.CountAsync(new CorrelationIdSpecification<WorkflowInstance>(message.CorrelationId), cancellationToken);
|
||||
|
||||
if (correlatedWorkflowInstanceCount > 0)
|
||||
{
|
||||
// Trigger existing workflows (if blocked on this message).
|
||||
var existingWorkflows = await workflowSelector.SelectWorkflowsAsync<MessageReceivedTrigger>(x => x.QueueName == queueName && x.CorrelationId == message.CorrelationId, cancellationToken).ToList();
|
||||
await workflowRunner.TriggerWorkflowsAsync(existingWorkflows, message, message.CorrelationId, cancellationToken: cancellationToken);
|
||||
}
|
||||
else
|
||||
{
|
||||
// Trigger new workflow.
|
||||
await TriggerNewWorkflowAsync();
|
||||
}
|
||||
}
|
||||
|
||||
private Task ExceptionReceivedHandler(ExceptionReceivedEventArgs e)
|
||||
|
|
|
|||
|
|
@ -40,7 +40,7 @@ namespace Elsa.Activities.Timers
|
|||
|
||||
if (ExecuteAt <= now)
|
||||
{
|
||||
_logger.LogDebug("Scheduled trigger time lies in the past ('{Delta}'). Skipping scheduling.", now - ExecuteAt);
|
||||
_logger.LogDebug("Scheduled trigger time lies in the past ('{Delta}'). Skipping scheduling", now - ExecuteAt);
|
||||
return Done();
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -1,4 +1,5 @@
|
|||
using System;
|
||||
using System.Collections.Generic;
|
||||
using System.Threading;
|
||||
using System.Threading.Tasks;
|
||||
using Elsa.Builders;
|
||||
|
|
@ -18,6 +19,13 @@ namespace Elsa.Services
|
|||
CancellationToken cancellationToken = default)
|
||||
where TTrigger : ITrigger;
|
||||
|
||||
Task TriggerWorkflowsAsync(
|
||||
IEnumerable<WorkflowSelectorResult> results,
|
||||
object? input = default,
|
||||
string? correlationId = default,
|
||||
string? contextId = default,
|
||||
CancellationToken cancellationToken = default);
|
||||
|
||||
ValueTask<WorkflowInstance> RunWorkflowAsync(
|
||||
WorkflowInstance workflowInstance,
|
||||
string? activityId = default,
|
||||
|
|
|
|||
|
|
@ -68,7 +68,16 @@ namespace Elsa.Services
|
|||
where TTrigger : ITrigger
|
||||
{
|
||||
var results = await _workflowSelector.SelectWorkflowsAsync(predicate, cancellationToken).ToList();
|
||||
|
||||
await TriggerWorkflowsAsync(results, input, correlationId, contextId, cancellationToken);
|
||||
}
|
||||
|
||||
public async Task TriggerWorkflowsAsync(
|
||||
IEnumerable<WorkflowSelectorResult> results,
|
||||
object? input = default,
|
||||
string? correlationId = default,
|
||||
string? contextId = default,
|
||||
CancellationToken cancellationToken = default)
|
||||
{
|
||||
foreach (var result in results)
|
||||
{
|
||||
if (result.WorkflowInstanceId != null)
|
||||
|
|
|
|||
Loading…
Reference in a new issue