< Summary

Information
Class: Elsa.Diagnostics.StructuredLogs.Persistence.Relational.Services.StructuredLogWriteBuffer
Assembly: Elsa.Diagnostics.StructuredLogs.Persistence.Relational
File(s): /home/runner/work/elsa-core/elsa-core/src/modules/Elsa.Diagnostics.StructuredLogs.Persistence.Relational/Services/StructuredLogWriteBuffer.cs
Line coverage
25%
Covered lines: 32
Uncovered lines: 93
Coverable lines: 125
Total lines: 253
Line coverage: 25.6%
Branch coverage
21%
Covered branches: 8
Total branches: 38
Branch coverage: 21%
Method coverage

Feature is only available for sponsors

Upgrade to PRO version

Metrics

MethodBranch coverage Crap Score Cyclomatic complexity Line coverage
.ctor(...)100%11100%
get_DroppedWriteCount()100%210%
WriteAsync(...)0%2040%
WriteManyAsync()0%620%
QueryAsync(...)100%210%
ListSourcesAsync(...)100%210%
StartAsync(...)0%4260%
StopAsync()0%4260%
FlushAsync()50%3236.36%
DisposeAsync()83.33%8663.63%
ProcessQueueAsync()0%4260%
DequeueBatch()50%4485.71%
CountPendingWritesAsDropped()0%620%

File(s)

/home/runner/work/elsa-core/elsa-core/src/modules/Elsa.Diagnostics.StructuredLogs.Persistence.Relational/Services/StructuredLogWriteBuffer.cs

#LineLine coverage
 1using System.Diagnostics;
 2using Elsa.Diagnostics.StructuredLogs.Contracts;
 3using Elsa.Diagnostics.StructuredLogs.Models;
 4using Elsa.Diagnostics.StructuredLogs.Persistence.Relational.Contracts;
 5using Elsa.Diagnostics.StructuredLogs.Persistence.Relational.Options;
 6using Elsa.Diagnostics.StructuredLogs.Persistence.Relational.Stores;
 7using Microsoft.Extensions.Hosting;
 8using Microsoft.Extensions.Options;
 9
 10namespace Elsa.Diagnostics.StructuredLogs.Persistence.Relational.Services;
 11
 112public class StructuredLogWriteBuffer(
 113    RelationalStructuredLogStore store,
 114    IStructuredLogSourceRegistry sourceRegistry,
 115    IOptions<RelationalStructuredLogOptions> options) : IStructuredLogStore, IStructuredLogWriteBuffer, IHostedService, 
 16{
 117    private readonly object _lifecycleLock = new();
 118    private readonly Queue<StructuredLogEvent> _queue = new();
 119    private readonly SemaphoreSlim _signal = new(0);
 120    private CancellationTokenSource _stopTokenSource = new();
 21    private Task? _backgroundTask;
 22    private long _droppedWriteCount;
 23    private int _activeStartCount;
 24    private int _disposed;
 25
 026    public long DroppedWriteCount => Interlocked.Read(ref _droppedWriteCount);
 27
 28    public ValueTask WriteAsync(StructuredLogEvent logEvent, CancellationToken cancellationToken = default)
 29    {
 030        var shouldSignal = false;
 31
 032        lock (_queue)
 33        {
 034            if (_queue.Count >= Math.Max(1, options.Value.WriteQueue.Capacity))
 35            {
 036                Interlocked.Increment(ref _droppedWriteCount);
 037                return ValueTask.CompletedTask;
 38            }
 39
 040            shouldSignal = _queue.Count == 0;
 041            _queue.Enqueue(logEvent);
 042        }
 43
 044        sourceRegistry.MarkSeen(logEvent.SourceId, logEvent.ReceivedAt);
 45
 046        if (shouldSignal)
 047            _signal.Release();
 48
 049        return ValueTask.CompletedTask;
 050    }
 51
 52    public async ValueTask WriteManyAsync(IReadOnlyCollection<StructuredLogEvent> logEvents, CancellationToken cancellat
 53    {
 054        foreach (var logEvent in logEvents)
 055            await WriteAsync(logEvent, cancellationToken);
 056    }
 57
 58    public ValueTask<RecentStructuredLogsResult> QueryAsync(StructuredLogFilter filter, CancellationToken cancellationTo
 59    {
 060        return store.QueryAsync(filter, cancellationToken);
 61    }
 62
 63    public ValueTask<IReadOnlyCollection<StructuredLogSource>> ListSourcesAsync(CancellationToken cancellationToken = de
 64    {
 065        return store.ListSourcesAsync(cancellationToken);
 66    }
 67
 68    public Task StartAsync(CancellationToken cancellationToken)
 69    {
 070        lock (_lifecycleLock)
 71        {
 072            _activeStartCount++;
 73
 074            if (_backgroundTask is { IsCompleted: false })
 075                return Task.CompletedTask;
 76
 077            if (_stopTokenSource.IsCancellationRequested)
 78            {
 079                _stopTokenSource.Dispose();
 080                _stopTokenSource = new();
 81            }
 82
 083            var stopToken = _stopTokenSource.Token;
 084            _backgroundTask = Task.Run(() => ProcessQueueAsync(stopToken), CancellationToken.None);
 085        }
 86
 087        return Task.CompletedTask;
 088    }
 89
 90    public async Task StopAsync(CancellationToken cancellationToken)
 91    {
 92        Task? backgroundTask;
 93        CancellationTokenSource stopTokenSource;
 94
 095        lock (_lifecycleLock)
 96        {
 097            if (_activeStartCount == 0)
 098                return;
 99
 0100            _activeStartCount--;
 101
 0102            if (_activeStartCount > 0)
 0103                return;
 104
 0105            backgroundTask = _backgroundTask;
 0106            stopTokenSource = _stopTokenSource;
 0107        }
 108
 0109        await stopTokenSource.CancelAsync();
 110
 0111        if (backgroundTask != null)
 112        {
 113            try
 114            {
 0115                await backgroundTask.WaitAsync(cancellationToken);
 0116            }
 0117            catch (OperationCanceledException) when (cancellationToken.IsCancellationRequested || stopTokenSource.IsCanc
 118            {
 119                // Expected during shutdown; remaining queued writes are flushed below.
 0120            }
 121        }
 122
 0123        using var timeoutTokenSource = CancellationTokenSource.CreateLinkedTokenSource(cancellationToken);
 0124        timeoutTokenSource.CancelAfter(options.Value.WriteQueue.ShutdownFlushTimeout);
 125        try
 126        {
 0127            await FlushAsync(timeoutTokenSource.Token);
 0128        }
 0129        catch (OperationCanceledException) when (timeoutTokenSource.IsCancellationRequested)
 130        {
 0131            CountPendingWritesAsDropped();
 0132        }
 0133    }
 134
 135    public async ValueTask FlushAsync(CancellationToken cancellationToken = default)
 136    {
 0137        while (true)
 138        {
 1139            var batch = DequeueBatch();
 1140            if (batch.Count == 0)
 1141                return;
 142
 143            try
 144            {
 0145                await store.WriteManyAsync(batch, cancellationToken);
 0146            }
 0147            catch
 148            {
 0149                Interlocked.Add(ref _droppedWriteCount, batch.Count);
 0150                throw;
 151            }
 0152        }
 1153    }
 154
 155    public async ValueTask DisposeAsync()
 156    {
 2157        if (Interlocked.Exchange(ref _disposed, 1) == 1)
 1158            return;
 159
 160        Task? backgroundTask;
 161        CancellationTokenSource stopTokenSource;
 162
 1163        lock (_lifecycleLock)
 164        {
 1165            backgroundTask = _backgroundTask;
 1166            stopTokenSource = _stopTokenSource;
 1167        }
 168
 1169        await stopTokenSource.CancelAsync();
 1170        using var timeoutTokenSource = new CancellationTokenSource(options.Value.WriteQueue.ShutdownFlushTimeout);
 171
 1172        if (backgroundTask != null)
 173        {
 174            try
 175            {
 0176                await backgroundTask.WaitAsync(timeoutTokenSource.Token);
 0177            }
 0178            catch (OperationCanceledException) when (timeoutTokenSource.IsCancellationRequested || stopTokenSource.IsCan
 179            {
 0180                CountPendingWritesAsDropped();
 0181            }
 182        }
 183
 184        try
 185        {
 1186            await FlushAsync(timeoutTokenSource.Token);
 1187        }
 0188        catch (OperationCanceledException) when (timeoutTokenSource.IsCancellationRequested)
 189        {
 0190            CountPendingWritesAsDropped();
 0191        }
 192
 1193        _signal.Dispose();
 1194        stopTokenSource.Dispose();
 2195    }
 196
 197    private async Task ProcessQueueAsync(CancellationToken cancellationToken)
 198    {
 0199        while (!cancellationToken.IsCancellationRequested)
 200        {
 201            try
 202            {
 0203                await _signal.WaitAsync(options.Value.WriteQueue.FlushInterval, cancellationToken);
 204
 0205                if (cancellationToken.IsCancellationRequested)
 0206                    return;
 207
 0208                await FlushAsync();
 0209            }
 0210            catch (OperationCanceledException) when (cancellationToken.IsCancellationRequested)
 211            {
 0212                return;
 213            }
 0214            catch (ObjectDisposedException) when (cancellationToken.IsCancellationRequested)
 215            {
 0216                return;
 217            }
 0218            catch (Exception e)
 219            {
 0220                Trace.TraceError("Failed to flush structured log writes: {0}", e);
 221
 0222                if (cancellationToken.IsCancellationRequested)
 0223                    return;
 0224            }
 225        }
 0226    }
 227
 228    private IReadOnlyCollection<StructuredLogEvent> DequeueBatch()
 229    {
 1230        var batchSize = Math.Max(1, options.Value.WriteQueue.BatchSize);
 1231        var batch = new List<StructuredLogEvent>(batchSize);
 232
 1233        lock (_queue)
 234        {
 1235            while (_queue.Count > 0 && batch.Count < batchSize)
 0236                batch.Add(_queue.Dequeue());
 1237        }
 238
 1239        return batch;
 240    }
 241
 242    private void CountPendingWritesAsDropped()
 243    {
 0244        lock (_queue)
 245        {
 0246            if (_queue.Count == 0)
 0247                return;
 248
 0249            Interlocked.Add(ref _droppedWriteCount, _queue.Count);
 0250            _queue.Clear();
 0251        }
 0252    }
 253}