diff --git a/Morpheus.Tests/LogsWriterServiceTests.cs b/Morpheus.Tests/LogsWriterServiceTests.cs index 0b237a3..1092e22 100644 --- a/Morpheus.Tests/LogsWriterServiceTests.cs +++ b/Morpheus.Tests/LogsWriterServiceTests.cs @@ -45,6 +45,43 @@ public async Task ExecuteAsync_RetriesBatchAfterTransientPersistenceFailure() await writer.StopAsync(CancellationToken.None); } + [Fact] + public async Task ExecuteAsync_RetriesFinalBatchAfterTransientPersistenceFailure() + { + await using SqliteConnection connection = new("Data Source=:memory:"); + await connection.OpenAsync(); + + DbContextOptions options = new DbContextOptionsBuilder() + .UseSqlite(connection) + .Options; + await using DB observer = new(options); + await observer.Database.EnsureCreatedAsync(); + FailureState failureState = new(); + + ServiceCollection services = new(); + services.AddScoped(_ => new FailingDb(options, failureState)); + using ServiceProvider provider = services.BuildServiceProvider(); + + LogQueue queue = new(capacity: 10); + Assert.True(queue.TryEnqueue( + LogsService.CreateGeneralLog("shutdown retry", Discord.LogSeverity.Warning))); + LogsWriterService writer = new(queue, provider.GetRequiredService()); + + await writer.StartAsync(CancellationToken.None); + + DateTime deadline = DateTime.UtcNow.AddSeconds(2); + while (queue.Reader.Count > 0 && DateTime.UtcNow < deadline) + await Task.Delay(10); + Assert.Equal(0, queue.Reader.Count); + + await writer.StopAsync(CancellationToken.None); + + Assert.Equal(1, await observer.Logs.CountAsync()); + Assert.Equal( + LogsService.FormatGeneralLog("shutdown retry", Discord.LogSeverity.Warning), + (await observer.Logs.SingleAsync()).Message); + } + private sealed class FailureState { public bool FailNextSave { get; set; } = true; @@ -54,6 +91,8 @@ private sealed class FailingDb(DbContextOptions options, FailureState state) { public override Task SaveChangesAsync(CancellationToken cancellationToken = default) { + cancellationToken.ThrowIfCancellationRequested(); + if (state.FailNextSave) { state.FailNextSave = false; diff --git a/Services/LogsWriterService.cs b/Services/LogsWriterService.cs index 6e35bc9..29f3e3e 100644 --- a/Services/LogsWriterService.cs +++ b/Services/LogsWriterService.cs @@ -8,6 +8,7 @@ namespace Morpheus.Services; public sealed class LogsWriterService(LogQueue logQueue, IServiceScopeFactory scopeFactory) : BackgroundService { private const int BatchSize = 50; + private const int MaxFinalFlushAttempts = 3; private static readonly TimeSpan BatchDelay = TimeSpan.FromMilliseconds(250); protected override async Task ExecuteAsync(CancellationToken stoppingToken) @@ -52,11 +53,24 @@ protected override async Task ExecuteAsync(CancellationToken stoppingToken) } finally { + int failedFlushAttempts = 0; while (true) { DrainAvailable(batch); - if (batch.Count == 0 || !await FlushAsync(batch, CancellationToken.None)) + if (batch.Count == 0) break; + + if (await FlushAsync(batch, CancellationToken.None)) + { + failedFlushAttempts = 0; + continue; + } + + failedFlushAttempts++; + if (failedFlushAttempts >= MaxFinalFlushAttempts) + break; + + await Task.Delay(BatchDelay); } } }