Added hosted services for monitoring instance lifetimes
This commit is contained in:
parent
9397fa1192
commit
dd4c33c8cb
|
|
@ -0,0 +1,33 @@
|
|||
<Project Sdk="Microsoft.NET.Sdk">
|
||||
|
||||
<PropertyGroup>
|
||||
<ImplicitUsings>enable</ImplicitUsings>
|
||||
<Nullable>enable</Nullable>
|
||||
<TargetFrameworks>net6.0;net7.0;net8.0</TargetFrameworks>
|
||||
<LangVersion>latest</LangVersion>
|
||||
<Copyright>2024</Copyright>
|
||||
<PackageProjectUrl>https://github.com/elsa-workflows/elsa-core</PackageProjectUrl>
|
||||
<PackageIcon>icon.png</PackageIcon>
|
||||
<RepositoryUrl>https://github.com/elsa-workflows/elsa-core</RepositoryUrl>
|
||||
</PropertyGroup>
|
||||
|
||||
<PropertyGroup>
|
||||
<Description>
|
||||
Provides hosting management functionality.
|
||||
</Description>
|
||||
<PackageTags>elsa module hosting</PackageTags>
|
||||
</PropertyGroup>
|
||||
|
||||
<ItemGroup>
|
||||
<PackageReference Include="Microsoft.Extensions.Hosting.Abstractions" Version="8.0.0" />
|
||||
</ItemGroup>
|
||||
|
||||
<ItemGroup>
|
||||
<ProjectReference Include="..\Elsa.Workflows.Runtime\Elsa.Workflows.Runtime.csproj" />
|
||||
</ItemGroup>
|
||||
|
||||
<ItemGroup>
|
||||
<ProjectReference Include="..\Elsa.Workflows.Runtime\Elsa.Workflows.Runtime.csproj" />
|
||||
</ItemGroup>
|
||||
|
||||
</Project>
|
||||
|
|
@ -0,0 +1,81 @@
|
|||
using Elsa.Hosting.Management.Notifications;
|
||||
using Elsa.Hosting.Management.Options;
|
||||
using Elsa.Mediator.Contracts;
|
||||
using Elsa.Workflows.Runtime.Contracts;
|
||||
using Elsa.Workflows.Runtime.Models;
|
||||
using Medallion.Threading;
|
||||
using Microsoft.Extensions.DependencyInjection;
|
||||
using Microsoft.Extensions.Hosting;
|
||||
using Microsoft.Extensions.Options;
|
||||
|
||||
namespace Elsa.Hosting.Management.HostedServices;
|
||||
|
||||
/// <summary>
|
||||
/// Service to check the heartbeats of all running instances and determine whether instances have stopped working.
|
||||
/// </summary>
|
||||
public class InstanceHeartbeatMonitorService : IHostedService, IDisposable
|
||||
{
|
||||
private readonly IServiceProvider _serviceProvider;
|
||||
private readonly HeartbeatSettings _heartbeatSettings;
|
||||
private Timer? _timer;
|
||||
|
||||
/// <summary>
|
||||
/// Creates a new instance of the <see cref="InstanceHeartbeatService"/>
|
||||
/// </summary>
|
||||
public InstanceHeartbeatMonitorService(IServiceProvider serviceProvider, IOptions<HeartbeatSettings> heartbeatSettings)
|
||||
{
|
||||
_serviceProvider = serviceProvider;
|
||||
_heartbeatSettings = heartbeatSettings.Value;
|
||||
}
|
||||
|
||||
public Task StartAsync(CancellationToken cancellationToken)
|
||||
{
|
||||
_timer = new Timer(MonitorHeartbeats, null, TimeSpan.Zero, _heartbeatSettings.InstanceHeartbeatRythm);
|
||||
return Task.CompletedTask;
|
||||
}
|
||||
|
||||
public Task StopAsync(CancellationToken cancellationToken)
|
||||
{
|
||||
_timer?.Change(Timeout.Infinite, Timeout.Infinite);
|
||||
return Task.CompletedTask;
|
||||
}
|
||||
|
||||
public void Dispose()
|
||||
{
|
||||
_timer?.Dispose();
|
||||
}
|
||||
|
||||
private void MonitorHeartbeats(object? state)
|
||||
{
|
||||
_ = Task.Run(async () => await MonitorHeartbeatsAsync());
|
||||
}
|
||||
|
||||
private async Task MonitorHeartbeatsAsync()
|
||||
{
|
||||
using var scope = _serviceProvider.CreateScope();
|
||||
|
||||
var lockProvider = scope.ServiceProvider.GetRequiredService<IDistributedLockProvider>();
|
||||
var store = scope.ServiceProvider.GetRequiredService<IKeyValueStore>();
|
||||
var notificationSender = scope.ServiceProvider.GetRequiredService<INotificationSender>();
|
||||
|
||||
var lockKey = "InstanceHeartbeatMonitorService";
|
||||
await using var monitorLock = await lockProvider.TryAcquireLockAsync(lockKey, TimeSpan.Zero);
|
||||
if (monitorLock == null)
|
||||
return;
|
||||
|
||||
var filter = new KeyValueFilter { StartsWith = true, Key = "Heartbeat_" };
|
||||
var heartbeats = await store.FindManyAsync(filter, default);
|
||||
|
||||
foreach (var heartbeat in heartbeats)
|
||||
{
|
||||
var lastHeartbeat = DateTime.Parse(heartbeat.SerializedValue);
|
||||
if (DateTime.UtcNow - lastHeartbeat <= _heartbeatSettings.InstanceDeactivatedPeriod)
|
||||
{
|
||||
continue;
|
||||
}
|
||||
|
||||
var instanceName = heartbeat.Key.Substring(InstanceHeartbeatService.HeartbeatKeyPrefix.Length);
|
||||
await notificationSender.SendAsync(new InstanceDeactivated(instanceName));
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
@ -0,0 +1,66 @@
|
|||
using Elsa.Hosting.Management.Options;
|
||||
using Elsa.Workflows.Contracts;
|
||||
using Elsa.Workflows.Runtime.Contracts;
|
||||
using Elsa.Workflows.Runtime.Entities;
|
||||
using Microsoft.Extensions.DependencyInjection;
|
||||
using Microsoft.Extensions.Hosting;
|
||||
using Microsoft.Extensions.Options;
|
||||
|
||||
namespace Elsa.Hosting.Management.HostedServices;
|
||||
|
||||
/// <summary>
|
||||
/// Service to write heartbeat messages per running instance.
|
||||
/// </summary>
|
||||
public class InstanceHeartbeatService : IHostedService, IDisposable
|
||||
{
|
||||
private readonly IServiceProvider _serviceProvider;
|
||||
private readonly HeartbeatSettings _heartbeatSettings;
|
||||
private Timer? _timer;
|
||||
|
||||
internal static string HeartbeatKeyPrefix = "Heartbeat_";
|
||||
/// <summary>
|
||||
/// Creates a new instance of the <see cref="InstanceHeartbeatService"/>
|
||||
/// </summary>
|
||||
public InstanceHeartbeatService(IServiceProvider serviceProvider, IOptions<HeartbeatSettings> heartbeatSettings)
|
||||
{
|
||||
_serviceProvider = serviceProvider;
|
||||
_heartbeatSettings = heartbeatSettings.Value;
|
||||
}
|
||||
|
||||
public Task StartAsync(CancellationToken cancellationToken)
|
||||
{
|
||||
_timer = new Timer(WriteHeartbeat, null, TimeSpan.Zero, _heartbeatSettings.InstanceHeartbeatRythm);
|
||||
return Task.CompletedTask;
|
||||
}
|
||||
|
||||
public Task StopAsync(CancellationToken cancellationToken)
|
||||
{
|
||||
_timer?.Change(Timeout.Infinite, Timeout.Infinite);
|
||||
return Task.CompletedTask;
|
||||
}
|
||||
|
||||
public void Dispose()
|
||||
{
|
||||
_timer?.Dispose();
|
||||
}
|
||||
|
||||
private void WriteHeartbeat(object? state)
|
||||
{
|
||||
_ = Task.Run(async () => await WriteHeartbeatAsync());
|
||||
}
|
||||
|
||||
private async Task WriteHeartbeatAsync()
|
||||
{
|
||||
using var scope = _serviceProvider.CreateScope();
|
||||
|
||||
var instanceNameRetriever = scope.ServiceProvider.GetRequiredService<IInstanceNameRetriever>();
|
||||
var store = scope.ServiceProvider.GetRequiredService<IKeyValueStore>();
|
||||
|
||||
await store.SaveAsync(new SerializedKeyValuePair
|
||||
{
|
||||
Key = $"{HeartbeatKeyPrefix}{instanceNameRetriever.GetName()}",
|
||||
SerializedValue = DateTime.UtcNow.ToString("o")
|
||||
},
|
||||
default);
|
||||
}
|
||||
}
|
||||
|
|
@ -0,0 +1,5 @@
|
|||
using Elsa.Mediator.Contracts;
|
||||
|
||||
namespace Elsa.Hosting.Management.Notifications;
|
||||
|
||||
public record InstanceDeactivated(string InstanceName) : INotification;
|
||||
|
|
@ -0,0 +1,7 @@
|
|||
namespace Elsa.Hosting.Management.Options;
|
||||
|
||||
public class HeartbeatSettings
|
||||
{
|
||||
public TimeSpan InstanceHeartbeatRythm { get; set; } = TimeSpan.FromMinutes(1);
|
||||
public TimeSpan InstanceDeactivatedPeriod { get; set; } = TimeSpan.FromHours(1);
|
||||
}
|
||||
Loading…
Reference in a new issue