From a0959ea2c77f59e9828f826918dc7b582bf19b43 Mon Sep 17 00:00:00 2001 From: EonaCat Date: Fri, 25 Sep 2026 15:13:02 +0200 Subject: [PATCH] Updated --- .../LogCentralClient.cs | 32 ++++++++- .../LogCentralEonaCatAdapter.cs | 1 + .../LogCentralOptions.cs | 12 ++++ EonaCat.LogStack/EonaCatLogger.cs | 1 + .../EonaCatLoggerCore/DeadLetterQueue.cs | 70 ++++++++++++++++++- .../Exceptions/AdvancedExceptionRenderer.cs | 35 +++------- .../Flows/DiagnosticsFlow.cs | 1 + .../EonaCatLoggerCore/Flows/DlqFlow.cs | 7 +- .../EonaCatLoggerCore/Flows/LokiFlow.cs | 29 +++++++- .../EonaCatLoggerCore/MessageTemplate.cs | 13 +++- EonaCat.LogStack/Telemetry/TelemetryServer.cs | 5 ++ 11 files changed, 171 insertions(+), 35 deletions(-) diff --git a/EonaCat.LogStack.LogClient/LogCentralClient.cs b/EonaCat.LogStack.LogClient/LogCentralClient.cs index e4d5ba6..03f5c85 100644 --- a/EonaCat.LogStack.LogClient/LogCentralClient.cs +++ b/EonaCat.LogStack.LogClient/LogCentralClient.cs @@ -86,6 +86,7 @@ namespace EonaCat.LogStack.LogClient entry.Properties = JsonHelper.ToJson(properties); _logQueue.Enqueue(entry); + TrimQueue(); if (_logQueue.Count >= _options.BatchSize) { @@ -93,6 +94,22 @@ namespace EonaCat.LogStack.LogClient } } + /// + /// Drops the oldest entries once the queue exceeds the configured maximum size, + /// so a server outage cannot grow the queue without bound. + /// + private void TrimQueue() + { + while (_logQueue.Count > _options.MaxQueueSize && _logQueue.TryDequeue(out _)) + { + } + } + + private bool IsExpired(LogEntry entry) + { + return DateTime.UtcNow - entry.Timestamp > _options.MaxRetention; + } + private async Task FlushAsync() { if (_logQueue.IsEmpty) @@ -106,6 +123,12 @@ namespace EonaCat.LogStack.LogClient var batch = new List(); while (batch.Count < _options.BatchSize && _logQueue.TryDequeue(out var entry)) { + // Drop stale entries instead of holding/sending them indefinitely. + if (IsExpired(entry)) + { + continue; + } + batch.Add(entry); } @@ -155,10 +178,17 @@ namespace EonaCat.LogStack.LogClient Console.WriteLine($"[LogCentral] Failed to send logs to EonaCat: {ex.Message}"); } + // Only hold onto entries that are still within the retention window; + // never retry/re-queue indefinitely, and never exceed the max queue size. foreach (var entry in batch) { - _logQueue.Enqueue(entry); + if (!IsExpired(entry)) + { + _logQueue.Enqueue(entry); + } } + + TrimQueue(); } } diff --git a/EonaCat.LogStack.LogClient/LogCentralEonaCatAdapter.cs b/EonaCat.LogStack.LogClient/LogCentralEonaCatAdapter.cs index 6801ef9..c85c945 100644 --- a/EonaCat.LogStack.LogClient/LogCentralEonaCatAdapter.cs +++ b/EonaCat.LogStack.LogClient/LogCentralEonaCatAdapter.cs @@ -17,6 +17,7 @@ namespace EonaCat.LogStack.LogClient public LogCentralEonaCatAdapter(LogBuilder logBuilder, LogCentralClient client) { _client = client; + _logBuilder = logBuilder ?? throw new ArgumentNullException(nameof(logBuilder)); _logBuilder.OnLog += LogSettings_OnLog; } diff --git a/EonaCat.LogStack.LogClient/LogCentralOptions.cs b/EonaCat.LogStack.LogClient/LogCentralOptions.cs index 46464d7..4510e8c 100644 --- a/EonaCat.LogStack.LogClient/LogCentralOptions.cs +++ b/EonaCat.LogStack.LogClient/LogCentralOptions.cs @@ -17,5 +17,17 @@ namespace EonaCat.LogStack.LogClient public int BatchSize { get; set; } = 1; public int FlushIntervalSeconds { get; set; } = 5; public bool EnableFallbackLogging { get; set; } = true; + + /// + /// Maximum number of entries kept in memory while the server is unreachable. + /// Oldest entries are dropped once this is exceeded to avoid unbounded memory growth. + /// + public int MaxQueueSize { get; set; } = 10000; + + /// + /// Maximum time an entry is kept in memory waiting to be (re)sent. Entries older + /// than this are dropped instead of retried/re-queued forever. + /// + public TimeSpan MaxRetention { get; set; } = TimeSpan.FromMinutes(10); } } diff --git a/EonaCat.LogStack/EonaCatLogger.cs b/EonaCat.LogStack/EonaCatLogger.cs index 1dea17a..eb88c52 100644 --- a/EonaCat.LogStack/EonaCatLogger.cs +++ b/EonaCat.LogStack/EonaCatLogger.cs @@ -864,6 +864,7 @@ namespace EonaCat.LogStack // Request cancellation of the async pipeline if it's still running _asyncCts?.Cancel(); _asyncCts?.Dispose(); + _deadLetterQueue?.Dispose(); GC.SuppressFinalize(this); } } diff --git a/EonaCat.LogStack/EonaCatLoggerCore/DeadLetterQueue.cs b/EonaCat.LogStack/EonaCatLoggerCore/DeadLetterQueue.cs index b19fd1f..c912e73 100644 --- a/EonaCat.LogStack/EonaCatLoggerCore/DeadLetterQueue.cs +++ b/EonaCat.LogStack/EonaCatLoggerCore/DeadLetterQueue.cs @@ -15,21 +15,53 @@ namespace EonaCat.LogStack.DeadLettering; /// Dead Letter Queue (DLQ) for capturing failed log attempts and enabling replay. /// Superior to data loss as it preserves logs that couldn't be delivered to any flow. /// -public class DeadLetterQueue +public class DeadLetterQueue : IDisposable { private readonly ConcurrentQueue _queue = new ConcurrentQueue(); private readonly int _maxCapacity; + private readonly TimeSpan _maxRetention; private int _currentCount; private long _totalDropped; private long _totalEnqueued; + private long _totalExpired; private readonly object _statsLock = new object(); + private bool _disposed; public event EventHandler EventEnqueued; public event EventHandler EventReplayed; - public DeadLetterQueue(int maxCapacity = 100000) + /// Maximum number of events held in memory. + /// + /// Maximum time an undeliverable event is kept in memory before being dropped. + /// Defaults to 10 minutes so a persistently unreachable destination cannot cause + /// unbounded memory growth. + /// + public DeadLetterQueue(int maxCapacity = 100000, TimeSpan? maxRetention = null) { _maxCapacity = Math.Max(1000, maxCapacity); + _maxRetention = maxRetention ?? TimeSpan.FromMinutes(10); + } + + /// + /// Removes events that have exceeded the retention window. The queue is FIFO, + /// so the oldest entries are always at the front. + /// + private void PurgeExpired() + { + var now = DateTime.UtcNow; + while (_queue.TryPeek(out var oldest) && now - oldest.EnqueuedAt > _maxRetention) + { + if (!_queue.TryDequeue(out _)) + { + break; + } + + _currentCount--; + lock (_statsLock) + { + _totalExpired++; + } + } } /// @@ -37,6 +69,7 @@ public class DeadLetterQueue /// public bool Enqueue(LogEvent logEvent, string failureReason, Exception failureException = null) { + PurgeExpired(); if (_currentCount >= _maxCapacity) { @@ -71,13 +104,18 @@ public class DeadLetterQueue /// /// Gets the count of events in the DLQ. /// - public int GetCount() => _currentCount; + public int GetCount() + { + PurgeExpired(); + return _currentCount; + } /// /// Peeks at the next event without removing it. /// public bool TryPeek(out DeadLetterEvent dlEvent) { + PurgeExpired(); return _queue.TryPeek(out dlEvent); } @@ -86,6 +124,7 @@ public class DeadLetterQueue /// public bool TryDequeue(out DeadLetterEvent dlEvent) { + PurgeExpired(); if (_queue.TryDequeue(out var evt)) { _currentCount--; @@ -102,6 +141,7 @@ public class DeadLetterQueue /// public List GetAll(int skip = 0, int take = 100) { + PurgeExpired(); return _queue.Skip(skip).Take(take).ToList(); } @@ -183,6 +223,7 @@ public class DeadLetterQueue /// public DeadLetterQueueStats GetStats() { + PurgeExpired(); lock (_statsLock) { return new DeadLetterQueueStats @@ -191,6 +232,7 @@ public class DeadLetterQueue MaxCapacity = _maxCapacity, TotalEnqueued = _totalEnqueued, TotalDropped = _totalDropped, + TotalExpired = _totalExpired, QueueUtilizationPercent = (_currentCount * 100.0) / _maxCapacity, OldestEventAge = _queue.TryPeek(out var oldest) ? DateTime.UtcNow - oldest.EnqueuedAt @@ -199,6 +241,23 @@ public class DeadLetterQueue } } + /// + /// Clears the queue and releases event subscribers. + /// + public void Dispose() + { + if (_disposed) + { + return; + } + + Clear(); + EventEnqueued = null; + EventReplayed = null; + _disposed = true; + GC.SuppressFinalize(this); + } + /// /// Gets events grouped by failure reason. /// @@ -300,6 +359,11 @@ public class DeadLetterQueueStats /// public long TotalDropped { get; set; } + /// + /// Total events dropped for exceeding the retention window. + /// + public long TotalExpired { get; set; } + /// /// Queue utilization as percentage. /// diff --git a/EonaCat.LogStack/EonaCatLoggerCore/Exceptions/AdvancedExceptionRenderer.cs b/EonaCat.LogStack/EonaCatLoggerCore/Exceptions/AdvancedExceptionRenderer.cs index f135e9a..f5e9e95 100644 --- a/EonaCat.LogStack/EonaCatLoggerCore/Exceptions/AdvancedExceptionRenderer.cs +++ b/EonaCat.LogStack/EonaCatLoggerCore/Exceptions/AdvancedExceptionRenderer.cs @@ -241,23 +241,19 @@ public sealed class ExceptionContextBuilder /// public static class ExceptionContextManager { - private static readonly object _lock = new object(); - private static readonly Dictionary _contextMap = new(); + // ConditionalWeakTable keys don't keep exceptions alive - entries are collected + // automatically once the exception itself is no longer referenced elsewhere, + // instead of accumulating forever in a regular static Dictionary. + // Not readonly: Clear() below swaps in a new instance (ConditionalWeakTable.Clear() + // isn't available on netstandard2.0/net4.8). + private static System.Runtime.CompilerServices.ConditionalWeakTable _contextMap = new(); /// /// Gets or creates context for an exception /// public static ExceptionContext GetOrCreateContext(Exception ex) { - lock (_lock) - { - if (!_contextMap.TryGetValue(ex, out var context)) - { - context = new ExceptionContext(); - _contextMap[ex] = context; - } - return context; - } + return _contextMap.GetValue(ex, static _ => new ExceptionContext()); } /// @@ -267,10 +263,8 @@ public static class ExceptionContextManager { var builder = ExceptionContext.Create(); configure(builder); - lock (_lock) - { - _contextMap[ex] = builder.Build(); - } + _contextMap.Remove(ex); + _contextMap.Add(ex, builder.Build()); return ex; } @@ -279,11 +273,7 @@ public static class ExceptionContextManager /// public static ExceptionContext? TryGetContext(Exception ex) { - lock (_lock) - { - _contextMap.TryGetValue(ex, out var context); - return context; - } + return _contextMap.TryGetValue(ex, out var context) ? context : null; } /// @@ -291,10 +281,7 @@ public static class ExceptionContextManager /// public static void Clear() { - lock (_lock) - { - _contextMap.Clear(); - } + _contextMap = new System.Runtime.CompilerServices.ConditionalWeakTable(); } } diff --git a/EonaCat.LogStack/EonaCatLoggerCore/Flows/DiagnosticsFlow.cs b/EonaCat.LogStack/EonaCatLoggerCore/Flows/DiagnosticsFlow.cs index 37319e6..5d17dd8 100644 --- a/EonaCat.LogStack/EonaCatLoggerCore/Flows/DiagnosticsFlow.cs +++ b/EonaCat.LogStack/EonaCatLoggerCore/Flows/DiagnosticsFlow.cs @@ -179,6 +179,7 @@ namespace EonaCat.LogStack.Flows _cts.Cancel(); _samplerThread.Join(TimeSpan.FromSeconds(3)); _cts.Dispose(); + _proc.Dispose(); await base.DisposeAsync().ConfigureAwait(false); } diff --git a/EonaCat.LogStack/EonaCatLoggerCore/Flows/DlqFlow.cs b/EonaCat.LogStack/EonaCatLoggerCore/Flows/DlqFlow.cs index f237040..32a461d 100644 --- a/EonaCat.LogStack/EonaCatLoggerCore/Flows/DlqFlow.cs +++ b/EonaCat.LogStack/EonaCatLoggerCore/Flows/DlqFlow.cs @@ -24,10 +24,11 @@ public class DlqFlow : FlowBase public DlqFlow( string name = "DLQ", LogLevel minimumLevel = LogLevel.Warning, - int maxCapacity = 100000) + int maxCapacity = 100000, + TimeSpan? maxRetention = null) : base(name, minimumLevel) { - _dlq = new DeadLetterQueue(maxCapacity); + _dlq = new DeadLetterQueue(maxCapacity, maxRetention); } /// @@ -214,6 +215,8 @@ public class DlqFlow : FlowBase catch { } } + _dlq.Dispose(); + await base.DisposeAsync().ConfigureAwait(false); } } diff --git a/EonaCat.LogStack/EonaCatLoggerCore/Flows/LokiFlow.cs b/EonaCat.LogStack/EonaCatLoggerCore/Flows/LokiFlow.cs index 5522f8c..2409220 100644 --- a/EonaCat.LogStack/EonaCatLoggerCore/Flows/LokiFlow.cs +++ b/EonaCat.LogStack/EonaCatLoggerCore/Flows/LokiFlow.cs @@ -49,7 +49,9 @@ public sealed class LokiFlow : FlowBase private readonly Dictionary _staticLabels; private readonly int _batchSize; private readonly TimeSpan _batchInterval; - private readonly ConcurrentQueue<(string Labels, string Line, long NanoTs)> _queue = new(); + private readonly int _maxQueueSize; + private readonly TimeSpan _maxRetention; + private readonly ConcurrentQueue<(string Labels, string Line, long NanoTs, DateTime EnqueuedAt)> _queue = new(); private readonly CancellationTokenSource _cts = new(); private readonly Task _flushTask; private readonly SemaphoreSlim _batchSignal = new(0, int.MaxValue); @@ -63,7 +65,9 @@ public sealed class LokiFlow : FlowBase string? basicUser = null, string? basicPassword = null, HttpClient? httpClient = null, - LogLevel minimumLevel = LogLevel.Trace) + LogLevel minimumLevel = LogLevel.Trace, + int maxQueueSize = 20000, + TimeSpan? maxRetention = null) : base("Loki", minimumLevel) { if (string.IsNullOrWhiteSpace(lokiUrl)) @@ -75,6 +79,10 @@ public sealed class LokiFlow : FlowBase _staticLabels = labels ?? new Dictionary(); _batchSize = batchSize > 0 ? batchSize : 50; _batchInterval = TimeSpan.FromMilliseconds(batchIntervalMs > 0 ? batchIntervalMs : 1000); + _maxQueueSize = maxQueueSize > 0 ? maxQueueSize : 20000; + // Never hold undeliverable data longer than this, to avoid unbounded memory growth + // when Loki is unreachable. + _maxRetention = maxRetention ?? TimeSpan.FromMinutes(10); if (httpClient != null) { @@ -112,9 +120,16 @@ public sealed class LokiFlow : FlowBase var line = BuildLine(logEvent); var nanoTs = logEvent.Timestamp * 100L; // ticks → nanoseconds - _queue.Enqueue((labels, line, nanoTs)); + _queue.Enqueue((labels, line, nanoTs, DateTime.UtcNow)); Interlocked.Increment(ref BlastedCount); + // Drop oldest entries once over capacity so a slow/unreachable Loki instance + // cannot grow this queue without bound. + while (_queue.Count > _maxQueueSize && _queue.TryDequeue(out _)) + { + Interlocked.Increment(ref DroppedCount); + } + if (_queue.Count >= _batchSize) { _batchSignal.Release(1); @@ -168,8 +183,16 @@ public sealed class LokiFlow : FlowBase // Group by label-set (Loki stream key) var groups = new Dictionary>(StringComparer.Ordinal); + var now = DateTime.UtcNow; while (_queue.TryDequeue(out var item)) { + // Drop stale entries instead of holding/sending them indefinitely. + if (now - item.EnqueuedAt > _maxRetention) + { + Interlocked.Increment(ref DroppedCount); + continue; + } + if (!groups.TryGetValue(item.Labels, out var list)) { list = new List<(string, long)>(16); diff --git a/EonaCat.LogStack/EonaCatLoggerCore/MessageTemplate.cs b/EonaCat.LogStack/EonaCatLoggerCore/MessageTemplate.cs index 6b4bdd5..8876f14 100644 --- a/EonaCat.LogStack/EonaCatLoggerCore/MessageTemplate.cs +++ b/EonaCat.LogStack/EonaCatLoggerCore/MessageTemplate.cs @@ -677,12 +677,21 @@ public sealed class MessageTemplate private static bool IsArgumentDriven(object?[] args) => args.Length > 0; // Template cache to avoid re-parsing the same strings + private const int MaxCacheSize = 5000; private static readonly System.Collections.Concurrent.ConcurrentDictionary _cache = new(StringComparer.Ordinal); /// Returns a cached parsed template (recommended for hot paths). - public static MessageTemplate FromCache(string template) => - _cache.GetOrAdd(template, static t => Parse(t)); + public static MessageTemplate FromCache(string template) + { + // Guard against unbounded growth if callers pass many unique/dynamic template strings. + if (_cache.Count >= MaxCacheSize) + { + _cache.Clear(); + } + + return _cache.GetOrAdd(template, static t => Parse(t)); + } /// Clears the template parse cache. public static void ClearCache() => _cache.Clear(); diff --git a/EonaCat.LogStack/Telemetry/TelemetryServer.cs b/EonaCat.LogStack/Telemetry/TelemetryServer.cs index 2d171ba..050c0a4 100644 --- a/EonaCat.LogStack/Telemetry/TelemetryServer.cs +++ b/EonaCat.LogStack/Telemetry/TelemetryServer.cs @@ -8,6 +8,7 @@ using System.Threading.Tasks; namespace EonaCat.LogStack.Telemetry; public sealed class TelemetryServer : IDisposable { + private const int MaxEvents = 5000; private readonly HttpListener _listener = new(); public ConcurrentQueue Events { get; } = new(); public int Count => Events.Count; @@ -30,6 +31,10 @@ public sealed class TelemetryServer : IDisposable if (item != null) { Events.Enqueue(item); + // Drop the oldest events once capacity is exceeded to avoid unbounded growth. + while (Events.Count > MaxEvents && Events.TryDequeue(out _)) + { + } } } else if (ctx.Request.Url?.AbsolutePath == "/")