From dd4c33c8cb71fb12dfb1599d930097a7d42bc876 Mon Sep 17 00:00:00 2001 From: Raymond den Haan Date: Fri, 9 Feb 2024 14:07:04 +0100 Subject: [PATCH] Added hosted services for monitoring instance lifetimes --- .../Elsa.Hosting.Management.csproj | 33 ++++++++ .../InstanceHeartbeatMonitorService.cs | 81 +++++++++++++++++++ .../InstanceHeartbeatService.cs | 66 +++++++++++++++ .../Notifications/InstanceDeactivated.cs | 5 ++ .../Options/HeartbeatSettings.cs | 7 ++ 5 files changed, 192 insertions(+) create mode 100644 src/modules/Elsa.Hosting.Management/Elsa.Hosting.Management.csproj create mode 100644 src/modules/Elsa.Hosting.Management/HostedServices/InstanceHeartbeatMonitorService.cs create mode 100644 src/modules/Elsa.Hosting.Management/HostedServices/InstanceHeartbeatService.cs create mode 100644 src/modules/Elsa.Hosting.Management/Notifications/InstanceDeactivated.cs create mode 100644 src/modules/Elsa.Hosting.Management/Options/HeartbeatSettings.cs diff --git a/src/modules/Elsa.Hosting.Management/Elsa.Hosting.Management.csproj b/src/modules/Elsa.Hosting.Management/Elsa.Hosting.Management.csproj new file mode 100644 index 000000000..8d004ae72 --- /dev/null +++ b/src/modules/Elsa.Hosting.Management/Elsa.Hosting.Management.csproj @@ -0,0 +1,33 @@ + + + + enable + enable + net6.0;net7.0;net8.0 + latest + 2024 + https://github.com/elsa-workflows/elsa-core + icon.png + https://github.com/elsa-workflows/elsa-core + + + + + Provides hosting management functionality. + + elsa module hosting + + + + + + + + + + + + + + + diff --git a/src/modules/Elsa.Hosting.Management/HostedServices/InstanceHeartbeatMonitorService.cs b/src/modules/Elsa.Hosting.Management/HostedServices/InstanceHeartbeatMonitorService.cs new file mode 100644 index 000000000..fca6b0df4 --- /dev/null +++ b/src/modules/Elsa.Hosting.Management/HostedServices/InstanceHeartbeatMonitorService.cs @@ -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; + +/// +/// Service to check the heartbeats of all running instances and determine whether instances have stopped working. +/// +public class InstanceHeartbeatMonitorService : IHostedService, IDisposable +{ + private readonly IServiceProvider _serviceProvider; + private readonly HeartbeatSettings _heartbeatSettings; + private Timer? _timer; + + /// + /// Creates a new instance of the + /// + public InstanceHeartbeatMonitorService(IServiceProvider serviceProvider, IOptions 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(); + var store = scope.ServiceProvider.GetRequiredService(); + var notificationSender = scope.ServiceProvider.GetRequiredService(); + + 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)); + } + } +} \ No newline at end of file diff --git a/src/modules/Elsa.Hosting.Management/HostedServices/InstanceHeartbeatService.cs b/src/modules/Elsa.Hosting.Management/HostedServices/InstanceHeartbeatService.cs new file mode 100644 index 000000000..9492eb269 --- /dev/null +++ b/src/modules/Elsa.Hosting.Management/HostedServices/InstanceHeartbeatService.cs @@ -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; + +/// +/// Service to write heartbeat messages per running instance. +/// +public class InstanceHeartbeatService : IHostedService, IDisposable +{ + private readonly IServiceProvider _serviceProvider; + private readonly HeartbeatSettings _heartbeatSettings; + private Timer? _timer; + + internal static string HeartbeatKeyPrefix = "Heartbeat_"; + /// + /// Creates a new instance of the + /// + public InstanceHeartbeatService(IServiceProvider serviceProvider, IOptions 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(); + var store = scope.ServiceProvider.GetRequiredService(); + + await store.SaveAsync(new SerializedKeyValuePair + { + Key = $"{HeartbeatKeyPrefix}{instanceNameRetriever.GetName()}", + SerializedValue = DateTime.UtcNow.ToString("o") + }, + default); + } +} \ No newline at end of file diff --git a/src/modules/Elsa.Hosting.Management/Notifications/InstanceDeactivated.cs b/src/modules/Elsa.Hosting.Management/Notifications/InstanceDeactivated.cs new file mode 100644 index 000000000..3eb9f1e8a --- /dev/null +++ b/src/modules/Elsa.Hosting.Management/Notifications/InstanceDeactivated.cs @@ -0,0 +1,5 @@ +using Elsa.Mediator.Contracts; + +namespace Elsa.Hosting.Management.Notifications; + +public record InstanceDeactivated(string InstanceName) : INotification; diff --git a/src/modules/Elsa.Hosting.Management/Options/HeartbeatSettings.cs b/src/modules/Elsa.Hosting.Management/Options/HeartbeatSettings.cs new file mode 100644 index 000000000..0951173c1 --- /dev/null +++ b/src/modules/Elsa.Hosting.Management/Options/HeartbeatSettings.cs @@ -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); +} \ No newline at end of file