Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
39 changes: 39 additions & 0 deletions Morpheus.Tests/LogsWriterServiceTests.cs
Original file line number Diff line number Diff line change
Expand Up @@ -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<DB> options = new DbContextOptionsBuilder<DB>()
.UseSqlite(connection)
.Options;
await using DB observer = new(options);
await observer.Database.EnsureCreatedAsync();
FailureState failureState = new();

ServiceCollection services = new();
services.AddScoped<DB>(_ => 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<IServiceScopeFactory>());

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;
Expand All @@ -54,6 +91,8 @@ private sealed class FailingDb(DbContextOptions<DB> options, FailureState state)
{
public override Task<int> SaveChangesAsync(CancellationToken cancellationToken = default)
{
cancellationToken.ThrowIfCancellationRequested();

if (state.FailNextSave)
{
state.FailNextSave = false;
Expand Down
16 changes: 15 additions & 1 deletion Services/LogsWriterService.cs
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand Down Expand Up @@ -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);
}
}
}
Expand Down