Add generic BackgroundWorker for running tasks asynchronously in the background
This commit is contained in:
parent
705979f94e
commit
a0c9f16aae
13
src/core/Elsa.Abstractions/Services/IBackgroundWorker.cs
Normal file
13
src/core/Elsa.Abstractions/Services/IBackgroundWorker.cs
Normal file
|
|
@ -0,0 +1,13 @@
|
|||
using System;
|
||||
using System.Threading;
|
||||
using System.Threading.Tasks;
|
||||
|
||||
namespace Elsa.Services
|
||||
{
|
||||
public interface IBackgroundWorker
|
||||
{
|
||||
Task ScheduleTask(Func<ValueTask> task, CancellationToken cancellationToken = default);
|
||||
Task ScheduleTask(Action task, CancellationToken cancellationToken = default);
|
||||
Task RunAsync(CancellationToken cancellationToken = default);
|
||||
}
|
||||
}
|
||||
|
|
@ -37,6 +37,7 @@
|
|||
<PackageReference Include="Scrutor" Version="3.3.0" />
|
||||
<PackageReference Include="Storage.Net" Version="9.2.7" />
|
||||
<PackageReference Include="System.Linq.Async" Version="5.0.0" />
|
||||
<PackageReference Include="System.Threading.Channels" Version="4.7.1" />
|
||||
<PackageReference Include="YamlDotNet" Version="9.1.0" />
|
||||
<PackageReference Include="Rebus.ServiceProvider" Version="6.0.0" />
|
||||
</ItemGroup>
|
||||
|
|
|
|||
|
|
@ -36,6 +36,7 @@ namespace Microsoft.Extensions.DependencyInjection
|
|||
this IServiceCollection services,
|
||||
Action<ElsaOptions>? configure = default)
|
||||
{
|
||||
services.AddHostedService<StartBackgroundWorker>();
|
||||
services.AddStartupRunner();
|
||||
|
||||
var options = new ElsaOptions(services);
|
||||
|
|
@ -95,6 +96,7 @@ namespace Microsoft.Extensions.DependencyInjection
|
|||
.AddSingleton<IActivityFactory, ActivityFactory>()
|
||||
.AddSingleton<IWorkflowBlueprintMaterializer, WorkflowBlueprintMaterializer>()
|
||||
.AddSingleton<IWorkflowBlueprintReflector, WorkflowBlueprintReflector>()
|
||||
.AddSingleton<IBackgroundWorker, BackgroundWorker>()
|
||||
.AddScoped<IWorkflowSelector, WorkflowSelector>()
|
||||
.AddScoped<IWorkflowPublisher, WorkflowPublisher>()
|
||||
.AddScoped<IWorkflowContextManager, WorkflowContextManager>()
|
||||
|
|
|
|||
14
src/core/Elsa.Core/HostedServices/StartBackgroundWorker.cs
Normal file
14
src/core/Elsa.Core/HostedServices/StartBackgroundWorker.cs
Normal file
|
|
@ -0,0 +1,14 @@
|
|||
using System.Threading;
|
||||
using System.Threading.Tasks;
|
||||
using Elsa.Services;
|
||||
using Microsoft.Extensions.Hosting;
|
||||
|
||||
namespace Elsa.HostedServices
|
||||
{
|
||||
public class StartBackgroundWorker : BackgroundService
|
||||
{
|
||||
private readonly IBackgroundWorker _backgroundWorker;
|
||||
public StartBackgroundWorker(IBackgroundWorker backgroundWorker) => _backgroundWorker = backgroundWorker;
|
||||
protected override async Task ExecuteAsync(CancellationToken stoppingToken) => await _backgroundWorker.RunAsync(stoppingToken);
|
||||
}
|
||||
}
|
||||
29
src/core/Elsa.Core/Services/BackgroundWorker.cs
Normal file
29
src/core/Elsa.Core/Services/BackgroundWorker.cs
Normal file
|
|
@ -0,0 +1,29 @@
|
|||
using System;
|
||||
using System.Threading;
|
||||
using System.Threading.Channels;
|
||||
using System.Threading.Tasks;
|
||||
|
||||
namespace Elsa.Services
|
||||
{
|
||||
public class BackgroundWorker : IBackgroundWorker
|
||||
{
|
||||
private readonly Channel<Func<ValueTask>> _channel;
|
||||
public BackgroundWorker() => _channel = Channel.CreateBounded<Func<ValueTask>>(10);
|
||||
|
||||
public async Task ScheduleTask(Func<ValueTask> task, CancellationToken cancellationToken) => await _channel.Writer.WriteAsync(task, cancellationToken);
|
||||
|
||||
public async Task ScheduleTask(Action task, CancellationToken cancellationToken) =>
|
||||
await _channel.Writer.WriteAsync(() =>
|
||||
{
|
||||
task();
|
||||
return new ValueTask();
|
||||
}, cancellationToken);
|
||||
|
||||
public async Task RunAsync(CancellationToken cancellationToken)
|
||||
{
|
||||
while (await _channel.Reader.WaitToReadAsync(cancellationToken))
|
||||
while (_channel.Reader.TryRead(out var task))
|
||||
await task();
|
||||
}
|
||||
}
|
||||
}
|
||||
Loading…
Reference in a new issue