using System.Diagnostics; using Elsa.Diagnostics.StructuredLogs.Contracts; using Elsa.Diagnostics.StructuredLogs.Models; using Elsa.Diagnostics.StructuredLogs.Persistence.Relational.Contracts; using Elsa.Diagnostics.StructuredLogs.Persistence.Relational.Options; using Elsa.Diagnostics.StructuredLogs.Persistence.Relational.Stores; using Microsoft.Extensions.Hosting; using Microsoft.Extensions.Options; namespace Elsa.Diagnostics.StructuredLogs.Persistence.Relational.Services; public class StructuredLogWriteBuffer( RelationalStructuredLogStore store, IOptions options) : IStructuredLogStore, IStructuredLogWriteBuffer, IHostedService, IAsyncDisposable { private readonly object _lifecycleLock = new(); private readonly Queue _queue = new(); private readonly SemaphoreSlim _signal = new(0); private CancellationTokenSource _stopTokenSource = new(); private Task? _backgroundTask; private long _droppedWriteCount; private int _activeStartCount; private int _disposed; public long DroppedWriteCount => Interlocked.Read(ref _droppedWriteCount); public ValueTask WriteAsync(StructuredLogEvent logEvent, CancellationToken cancellationToken = default) { var shouldSignal = false; lock (_queue) { if (_queue.Count >= Math.Max(1, options.Value.WriteQueue.Capacity)) { Interlocked.Increment(ref _droppedWriteCount); return ValueTask.CompletedTask; } shouldSignal = _queue.Count == 0; _queue.Enqueue(logEvent); } if (shouldSignal) _signal.Release(); return ValueTask.CompletedTask; } public async ValueTask WriteManyAsync(IReadOnlyCollection logEvents, CancellationToken cancellationToken = default) { foreach (var logEvent in logEvents) await WriteAsync(logEvent, cancellationToken); } public ValueTask QueryAsync(StructuredLogFilter filter, CancellationToken cancellationToken = default) { return store.QueryAsync(filter, cancellationToken); } public ValueTask> ListSourcesAsync(CancellationToken cancellationToken = default) { return store.ListSourcesAsync(cancellationToken); } public Task StartAsync(CancellationToken cancellationToken) { lock (_lifecycleLock) { _activeStartCount++; if (_backgroundTask is { IsCompleted: false }) return Task.CompletedTask; if (_stopTokenSource.IsCancellationRequested) { _stopTokenSource.Dispose(); _stopTokenSource = new(); } var stopToken = _stopTokenSource.Token; _backgroundTask = Task.Run(() => ProcessQueueAsync(stopToken), CancellationToken.None); } return Task.CompletedTask; } public async Task StopAsync(CancellationToken cancellationToken) { Task? backgroundTask; CancellationTokenSource stopTokenSource; lock (_lifecycleLock) { if (_activeStartCount == 0) return; _activeStartCount--; if (_activeStartCount > 0) return; backgroundTask = _backgroundTask; stopTokenSource = _stopTokenSource; } await stopTokenSource.CancelAsync(); if (backgroundTask != null) { try { await backgroundTask.WaitAsync(cancellationToken); } catch (OperationCanceledException) when (cancellationToken.IsCancellationRequested || stopTokenSource.IsCancellationRequested) { // Expected during shutdown; remaining queued writes are flushed below. } } using var timeoutTokenSource = CancellationTokenSource.CreateLinkedTokenSource(cancellationToken); timeoutTokenSource.CancelAfter(options.Value.WriteQueue.ShutdownFlushTimeout); try { await FlushAsync(timeoutTokenSource.Token); } catch (OperationCanceledException) when (timeoutTokenSource.IsCancellationRequested) { CountPendingWritesAsDropped(); } } public async ValueTask FlushAsync(CancellationToken cancellationToken = default) { while (true) { var batch = DequeueBatch(); if (batch.Count == 0) return; try { await store.WriteManyAsync(batch, cancellationToken); } catch { Interlocked.Add(ref _droppedWriteCount, batch.Count); throw; } } } public async ValueTask DisposeAsync() { if (Interlocked.Exchange(ref _disposed, 1) == 1) return; Task? backgroundTask; CancellationTokenSource stopTokenSource; lock (_lifecycleLock) { backgroundTask = _backgroundTask; stopTokenSource = _stopTokenSource; } await stopTokenSource.CancelAsync(); using var timeoutTokenSource = new CancellationTokenSource(options.Value.WriteQueue.ShutdownFlushTimeout); if (backgroundTask != null) { try { await backgroundTask.WaitAsync(timeoutTokenSource.Token); } catch (OperationCanceledException) when (timeoutTokenSource.IsCancellationRequested || stopTokenSource.IsCancellationRequested) { CountPendingWritesAsDropped(); } } try { await FlushAsync(timeoutTokenSource.Token); } catch (OperationCanceledException) when (timeoutTokenSource.IsCancellationRequested) { CountPendingWritesAsDropped(); } _signal.Dispose(); stopTokenSource.Dispose(); } private async Task ProcessQueueAsync(CancellationToken cancellationToken) { while (!cancellationToken.IsCancellationRequested) { try { await _signal.WaitAsync(options.Value.WriteQueue.FlushInterval, cancellationToken); if (cancellationToken.IsCancellationRequested) return; await FlushAsync(); } catch (OperationCanceledException) when (cancellationToken.IsCancellationRequested) { return; } catch (ObjectDisposedException) when (cancellationToken.IsCancellationRequested) { return; } catch (Exception e) { Trace.TraceError("Failed to flush structured log writes: {0}", e); if (cancellationToken.IsCancellationRequested) return; } } } private IReadOnlyCollection DequeueBatch() { var batchSize = Math.Max(1, options.Value.WriteQueue.BatchSize); var batch = new List(batchSize); lock (_queue) { while (_queue.Count > 0 && batch.Count < batchSize) batch.Add(_queue.Dequeue()); } return batch; } private void CountPendingWritesAsDropped() { lock (_queue) { if (_queue.Count == 0) return; Interlocked.Add(ref _droppedWriteCount, _queue.Count); _queue.Clear(); } } }