elsa-core/src/modules/Elsa.Workflows.Runtime/Services/LocalBackgroundActivityScheduler.cs
Sipke Schoorstra ac48e54ab8 Introduce job unscheduling functionality
Add a new method `UnscheduledAsync` to unschedule jobs across various components, including the job queue, background activity scheduler, and Hangfire integration. This enhancement provides a more robust way to remove jobs from scheduling, complementing the existing cancellation functionality.

This fixes an issue where a background activity execution job that resumed a workflow, which in turn would remove any associated bookmarks, which in turn would cancel the background job while that job is still executing and has to e.g. persist changes made to the DB.
2024-12-05 19:25:16 +01:00

49 lines
2 KiB
C#

using Elsa.Mediator.Contracts;
using Microsoft.Extensions.DependencyInjection;
namespace Elsa.Workflows.Runtime;
/// <summary>
/// Invokes activities from a background worker within the context of its workflow instance using a local background worker.
/// </summary>
public class LocalBackgroundActivityScheduler(IJobQueue jobQueue, IServiceScopeFactory scopeFactory) : IBackgroundActivityScheduler
{
public Task<string> CreateAsync(ScheduledBackgroundActivity scheduledBackgroundActivity, CancellationToken cancellationToken = default)
{
var jobId = jobQueue.Create(async ct => await InvokeBackgroundActivity(scheduledBackgroundActivity, ct));
return Task.FromResult(jobId);
}
public Task ScheduleAsync(string jobId, CancellationToken cancellationToken = default)
{
jobQueue.Enqueue(jobId);
return Task.CompletedTask;
}
/// <inheritdoc />
public Task<string> ScheduleAsync(ScheduledBackgroundActivity scheduledBackgroundActivity, CancellationToken cancellationToken = default)
{
var jobId = jobQueue.Enqueue(async ct => await InvokeBackgroundActivity(scheduledBackgroundActivity, ct));
return Task.FromResult(jobId);
}
public Task UnscheduledAsync(string jobId, CancellationToken cancellationToken = default)
{
jobQueue.Dequeue(jobId);
return Task.CompletedTask;
}
/// <inheritdoc />
public Task CancelAsync(string jobId, CancellationToken cancellationToken = default)
{
jobQueue.Cancel(jobId);
return Task.CompletedTask;
}
private async Task InvokeBackgroundActivity(ScheduledBackgroundActivity scheduledBackgroundActivity, CancellationToken cancellationToken)
{
using var scope = scopeFactory.CreateScope();
var backgroundActivityInvoker = scope.ServiceProvider.GetRequiredService<IBackgroundActivityInvoker>();
await backgroundActivityInvoker.ExecuteAsync(scheduledBackgroundActivity, cancellationToken);
}
}