Files
Rendezvous/src/FinalFactory.Rendezvous.Server/Deployment/GracefulDrainService.cs
T
KyuubiYoru cc5793f935
quality-gate / quality (push) Failing after 1m50s
quality-gate / container (push) Has been skipped
feat(release): add reproducible signed artifacts (#19)
2026-07-16 17:48:21 +02:00

103 lines
3.8 KiB
C#

using System.Diagnostics;
using FinalFactory.Rendezvous.Server.State;
using Microsoft.Extensions.Options;
namespace FinalFactory.Rendezvous.Server.Deployment;
internal sealed class GracefulDrainService : IHostedService, IDisposable
{
private static readonly TimeSpan PollInterval = TimeSpan.FromMilliseconds(50);
private readonly InMemoryEphemeralRendezvousStore _store;
private readonly IHostApplicationLifetime _lifetime;
private readonly DeploymentOptions _options;
private readonly ILogger<GracefulDrainService> _logger;
private readonly object _gate = new();
private CancellationTokenRegistration _stoppingRegistration;
private Task? _drainTask;
public GracefulDrainService(
InMemoryEphemeralRendezvousStore store,
IHostApplicationLifetime lifetime,
IOptions<DeploymentOptions> options,
ILogger<GracefulDrainService> logger)
{
_store = store;
_lifetime = lifetime;
_options = options.Value;
_logger = logger;
}
public Task StartAsync(CancellationToken cancellationToken)
{
cancellationToken.ThrowIfCancellationRequested();
_stoppingRegistration = _lifetime.ApplicationStopping.Register(
() => EnsureDrainAsync().GetAwaiter().GetResult());
return Task.CompletedTask;
}
public Task StopAsync(CancellationToken cancellationToken)
{
// ApplicationStopping callbacks run before hosted services and listeners
// stop. StopAsync is the idempotent fallback for directly driven hosts.
_ = cancellationToken;
return EnsureDrainAsync();
}
public void Dispose() => _stoppingRegistration.Dispose();
private Task EnsureDrainAsync()
{
lock (_gate)
{
return _drainTask ??= DrainAsync();
}
}
private async Task DrainAsync()
{
_store.BeginDrain(CancellationToken.None);
TimeSpan deadline = TimeSpan.FromSeconds(_options.DrainDeadlineSeconds);
TimeSpan minimum = TimeSpan.FromSeconds(_options.MinimumDrainSeconds);
long startedAt = Stopwatch.GetTimestamp();
LogDrainStarted(_logger, _options.DrainDeadlineSeconds);
try
{
while (Stopwatch.GetElapsedTime(startedAt) < deadline)
{
TimeSpan elapsed = Stopwatch.GetElapsedTime(startedAt);
if (elapsed >= minimum && _store.GetActiveJoinAttemptCountForDrain() == 0)
{
break;
}
TimeSpan remaining = deadline - elapsed;
await Task.Delay(
remaining < PollInterval ? remaining : PollInterval,
CancellationToken.None).ConfigureAwait(false);
}
}
finally
{
_store.MarkUnavailable();
double elapsedMilliseconds = Stopwatch.GetElapsedTime(startedAt).TotalMilliseconds;
LogDrainFinished(_logger, elapsedMilliseconds);
}
}
private static readonly Action<ILogger, int, Exception?> DrainStarted = LoggerMessage.Define<int>(
LogLevel.Information,
new EventId(1, nameof(LogDrainStarted)),
"Graceful drain started with a {DrainDeadlineSeconds}-second deadline");
private static readonly Action<ILogger, double, Exception?> DrainFinished = LoggerMessage.Define<double>(
LogLevel.Information,
new EventId(2, nameof(LogDrainFinished)),
"Graceful drain finished after {ElapsedMilliseconds:F0} ms; ephemeral state was cleared");
private static void LogDrainStarted(ILogger logger, int drainDeadlineSeconds) =>
DrainStarted(logger, drainDeadlineSeconds, null);
private static void LogDrainFinished(ILogger logger, double elapsedMilliseconds) =>
DrainFinished(logger, elapsedMilliseconds, null);
}