This commit is contained in:
EonaCat authored and Jeroen Saey committed 2026-09-25 15:13:02 +02:00
1 parent 17d7426b44
commit a0959ea2c7
11 files changed
+171 -35

No files matched your search

+31 -1
View File
@@ -86,6 +86,7 @@ namespace EonaCat.LogStack.LogClient
entry.Properties = JsonHelper.ToJson(properties); entry.Properties = JsonHelper.ToJson(properties);
_logQueue.Enqueue(entry); _logQueue.Enqueue(entry);
TrimQueue();
if (_logQueue.Count >= _options.BatchSize) if (_logQueue.Count >= _options.BatchSize)
{ {
@@ -93,6 +94,22 @@ namespace EonaCat.LogStack.LogClient
} }
} }
/// <summary>
/// Drops the oldest entries once the queue exceeds the configured maximum size,
/// so a server outage cannot grow the queue without bound.
/// </summary>
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() private async Task FlushAsync()
{ {
if (_logQueue.IsEmpty) if (_logQueue.IsEmpty)
@@ -106,6 +123,12 @@ namespace EonaCat.LogStack.LogClient
var batch = new List<LogEntry>(); var batch = new List<LogEntry>();
while (batch.Count < _options.BatchSize && _logQueue.TryDequeue(out var entry)) 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); batch.Add(entry);
} }
@@ -155,10 +178,17 @@ namespace EonaCat.LogStack.LogClient
Console.WriteLine($"[LogCentral] Failed to send logs to EonaCat: {ex.Message}"); 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) foreach (var entry in batch)
{ {
_logQueue.Enqueue(entry); if (!IsExpired(entry))
{
_logQueue.Enqueue(entry);
}
} }
TrimQueue();
} }
} }
@@ -17,6 +17,7 @@ namespace EonaCat.LogStack.LogClient
public LogCentralEonaCatAdapter(LogBuilder logBuilder, LogCentralClient client) public LogCentralEonaCatAdapter(LogBuilder logBuilder, LogCentralClient client)
{ {
_client = client; _client = client;
_logBuilder = logBuilder ?? throw new ArgumentNullException(nameof(logBuilder));
_logBuilder.OnLog += LogSettings_OnLog; _logBuilder.OnLog += LogSettings_OnLog;
} }
@@ -17,5 +17,17 @@ namespace EonaCat.LogStack.LogClient
public int BatchSize { get; set; } = 1; public int BatchSize { get; set; } = 1;
public int FlushIntervalSeconds { get; set; } = 5; public int FlushIntervalSeconds { get; set; } = 5;
public bool EnableFallbackLogging { get; set; } = true; public bool EnableFallbackLogging { get; set; } = true;
/// <summary>
/// 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.
/// </summary>
public int MaxQueueSize { get; set; } = 10000;
/// <summary>
/// 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.
/// </summary>
public TimeSpan MaxRetention { get; set; } = TimeSpan.FromMinutes(10);
} }
} }
+1
View File
@@ -864,6 +864,7 @@ namespace EonaCat.LogStack
// Request cancellation of the async pipeline if it's still running // Request cancellation of the async pipeline if it's still running
_asyncCts?.Cancel(); _asyncCts?.Cancel();
_asyncCts?.Dispose(); _asyncCts?.Dispose();
_deadLetterQueue?.Dispose();
GC.SuppressFinalize(this); GC.SuppressFinalize(this);
} }
} }
@@ -15,21 +15,53 @@ namespace EonaCat.LogStack.DeadLettering;
/// Dead Letter Queue (DLQ) for capturing failed log attempts and enabling replay. /// 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. /// Superior to data loss as it preserves logs that couldn't be delivered to any flow.
/// </summary> /// </summary>
public class DeadLetterQueue public class DeadLetterQueue : IDisposable
{ {
private readonly ConcurrentQueue<DeadLetterEvent> _queue = new ConcurrentQueue<DeadLetterEvent>(); private readonly ConcurrentQueue<DeadLetterEvent> _queue = new ConcurrentQueue<DeadLetterEvent>();
private readonly int _maxCapacity; private readonly int _maxCapacity;
private readonly TimeSpan _maxRetention;
private int _currentCount; private int _currentCount;
private long _totalDropped; private long _totalDropped;
private long _totalEnqueued; private long _totalEnqueued;
private long _totalExpired;
private readonly object _statsLock = new object(); private readonly object _statsLock = new object();
private bool _disposed;
public event EventHandler<DeadLetterEvent> EventEnqueued; public event EventHandler<DeadLetterEvent> EventEnqueued;
public event EventHandler<DeadLetterEventReplayed> EventReplayed; public event EventHandler<DeadLetterEventReplayed> EventReplayed;
public DeadLetterQueue(int maxCapacity = 100000) /// <param name="maxCapacity">Maximum number of events held in memory.</param>
/// <param name="maxRetention">
/// 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.
/// </param>
public DeadLetterQueue(int maxCapacity = 100000, TimeSpan? maxRetention = null)
{ {
_maxCapacity = Math.Max(1000, maxCapacity); _maxCapacity = Math.Max(1000, maxCapacity);
_maxRetention = maxRetention ?? TimeSpan.FromMinutes(10);
}
/// <summary>
/// Removes events that have exceeded the retention window. The queue is FIFO,
/// so the oldest entries are always at the front.
/// </summary>
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++;
}
}
} }
/// <summary> /// <summary>
@@ -37,6 +69,7 @@ public class DeadLetterQueue
/// </summary> /// </summary>
public bool Enqueue(LogEvent logEvent, string failureReason, Exception failureException = null) public bool Enqueue(LogEvent logEvent, string failureReason, Exception failureException = null)
{ {
PurgeExpired();
if (_currentCount >= _maxCapacity) if (_currentCount >= _maxCapacity)
{ {
@@ -71,13 +104,18 @@ public class DeadLetterQueue
/// <summary> /// <summary>
/// Gets the count of events in the DLQ. /// Gets the count of events in the DLQ.
/// </summary> /// </summary>
public int GetCount() => _currentCount; public int GetCount()
{
PurgeExpired();
return _currentCount;
}
/// <summary> /// <summary>
/// Peeks at the next event without removing it. /// Peeks at the next event without removing it.
/// </summary> /// </summary>
public bool TryPeek(out DeadLetterEvent dlEvent) public bool TryPeek(out DeadLetterEvent dlEvent)
{ {
PurgeExpired();
return _queue.TryPeek(out dlEvent); return _queue.TryPeek(out dlEvent);
} }
@@ -86,6 +124,7 @@ public class DeadLetterQueue
/// </summary> /// </summary>
public bool TryDequeue(out DeadLetterEvent dlEvent) public bool TryDequeue(out DeadLetterEvent dlEvent)
{ {
PurgeExpired();
if (_queue.TryDequeue(out var evt)) if (_queue.TryDequeue(out var evt))
{ {
_currentCount--; _currentCount--;
@@ -102,6 +141,7 @@ public class DeadLetterQueue
/// </summary> /// </summary>
public List<DeadLetterEvent> GetAll(int skip = 0, int take = 100) public List<DeadLetterEvent> GetAll(int skip = 0, int take = 100)
{ {
PurgeExpired();
return _queue.Skip(skip).Take(take).ToList(); return _queue.Skip(skip).Take(take).ToList();
} }
@@ -183,6 +223,7 @@ public class DeadLetterQueue
/// </summary> /// </summary>
public DeadLetterQueueStats GetStats() public DeadLetterQueueStats GetStats()
{ {
PurgeExpired();
lock (_statsLock) lock (_statsLock)
{ {
return new DeadLetterQueueStats return new DeadLetterQueueStats
@@ -191,6 +232,7 @@ public class DeadLetterQueue
MaxCapacity = _maxCapacity, MaxCapacity = _maxCapacity,
TotalEnqueued = _totalEnqueued, TotalEnqueued = _totalEnqueued,
TotalDropped = _totalDropped, TotalDropped = _totalDropped,
TotalExpired = _totalExpired,
QueueUtilizationPercent = (_currentCount * 100.0) / _maxCapacity, QueueUtilizationPercent = (_currentCount * 100.0) / _maxCapacity,
OldestEventAge = _queue.TryPeek(out var oldest) OldestEventAge = _queue.TryPeek(out var oldest)
? DateTime.UtcNow - oldest.EnqueuedAt ? DateTime.UtcNow - oldest.EnqueuedAt
@@ -199,6 +241,23 @@ public class DeadLetterQueue
} }
} }
/// <summary>
/// Clears the queue and releases event subscribers.
/// </summary>
public void Dispose()
{
if (_disposed)
{
return;
}
Clear();
EventEnqueued = null;
EventReplayed = null;
_disposed = true;
GC.SuppressFinalize(this);
}
/// <summary> /// <summary>
/// Gets events grouped by failure reason. /// Gets events grouped by failure reason.
/// </summary> /// </summary>
@@ -300,6 +359,11 @@ public class DeadLetterQueueStats
/// </summary> /// </summary>
public long TotalDropped { get; set; } public long TotalDropped { get; set; }
/// <summary>
/// Total events dropped for exceeding the retention window.
/// </summary>
public long TotalExpired { get; set; }
/// <summary> /// <summary>
/// Queue utilization as percentage. /// Queue utilization as percentage.
/// </summary> /// </summary>
@@ -241,23 +241,19 @@ public sealed class ExceptionContextBuilder
/// </summary> /// </summary>
public static class ExceptionContextManager public static class ExceptionContextManager
{ {
private static readonly object _lock = new object(); // ConditionalWeakTable keys don't keep exceptions alive - entries are collected
private static readonly Dictionary<Exception, ExceptionContext> _contextMap = new(); // 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<Exception, ExceptionContext> _contextMap = new();
/// <summary> /// <summary>
/// Gets or creates context for an exception /// Gets or creates context for an exception
/// </summary> /// </summary>
public static ExceptionContext GetOrCreateContext(Exception ex) public static ExceptionContext GetOrCreateContext(Exception ex)
{ {
lock (_lock) return _contextMap.GetValue(ex, static _ => new ExceptionContext());
{
if (!_contextMap.TryGetValue(ex, out var context))
{
context = new ExceptionContext();
_contextMap[ex] = context;
}
return context;
}
} }
/// <summary> /// <summary>
@@ -267,10 +263,8 @@ public static class ExceptionContextManager
{ {
var builder = ExceptionContext.Create(); var builder = ExceptionContext.Create();
configure(builder); configure(builder);
lock (_lock) _contextMap.Remove(ex);
{ _contextMap.Add(ex, builder.Build());
_contextMap[ex] = builder.Build();
}
return ex; return ex;
} }
@@ -279,11 +273,7 @@ public static class ExceptionContextManager
/// </summary> /// </summary>
public static ExceptionContext? TryGetContext(Exception ex) public static ExceptionContext? TryGetContext(Exception ex)
{ {
lock (_lock) return _contextMap.TryGetValue(ex, out var context) ? context : null;
{
_contextMap.TryGetValue(ex, out var context);
return context;
}
} }
/// <summary> /// <summary>
@@ -291,10 +281,7 @@ public static class ExceptionContextManager
/// </summary> /// </summary>
public static void Clear() public static void Clear()
{ {
lock (_lock) _contextMap = new System.Runtime.CompilerServices.ConditionalWeakTable<Exception, ExceptionContext>();
{
_contextMap.Clear();
}
} }
} }
@@ -179,6 +179,7 @@ namespace EonaCat.LogStack.Flows
_cts.Cancel(); _cts.Cancel();
_samplerThread.Join(TimeSpan.FromSeconds(3)); _samplerThread.Join(TimeSpan.FromSeconds(3));
_cts.Dispose(); _cts.Dispose();
_proc.Dispose();
await base.DisposeAsync().ConfigureAwait(false); await base.DisposeAsync().ConfigureAwait(false);
} }
@@ -24,10 +24,11 @@ public class DlqFlow : FlowBase
public DlqFlow( public DlqFlow(
string name = "DLQ", string name = "DLQ",
LogLevel minimumLevel = LogLevel.Warning, LogLevel minimumLevel = LogLevel.Warning,
int maxCapacity = 100000) int maxCapacity = 100000,
TimeSpan? maxRetention = null)
: base(name, minimumLevel) : base(name, minimumLevel)
{ {
_dlq = new DeadLetterQueue(maxCapacity); _dlq = new DeadLetterQueue(maxCapacity, maxRetention);
} }
/// <summary> /// <summary>
@@ -214,6 +215,8 @@ public class DlqFlow : FlowBase
catch { } catch { }
} }
_dlq.Dispose();
await base.DisposeAsync().ConfigureAwait(false); await base.DisposeAsync().ConfigureAwait(false);
} }
} }
@@ -49,7 +49,9 @@ public sealed class LokiFlow : FlowBase
private readonly Dictionary<string, string> _staticLabels; private readonly Dictionary<string, string> _staticLabels;
private readonly int _batchSize; private readonly int _batchSize;
private readonly TimeSpan _batchInterval; 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 CancellationTokenSource _cts = new();
private readonly Task _flushTask; private readonly Task _flushTask;
private readonly SemaphoreSlim _batchSignal = new(0, int.MaxValue); private readonly SemaphoreSlim _batchSignal = new(0, int.MaxValue);
@@ -63,7 +65,9 @@ public sealed class LokiFlow : FlowBase
string? basicUser = null, string? basicUser = null,
string? basicPassword = null, string? basicPassword = null,
HttpClient? httpClient = null, HttpClient? httpClient = null,
LogLevel minimumLevel = LogLevel.Trace) LogLevel minimumLevel = LogLevel.Trace,
int maxQueueSize = 20000,
TimeSpan? maxRetention = null)
: base("Loki", minimumLevel) : base("Loki", minimumLevel)
{ {
if (string.IsNullOrWhiteSpace(lokiUrl)) if (string.IsNullOrWhiteSpace(lokiUrl))
@@ -75,6 +79,10 @@ public sealed class LokiFlow : FlowBase
_staticLabels = labels ?? new Dictionary<string, string>(); _staticLabels = labels ?? new Dictionary<string, string>();
_batchSize = batchSize > 0 ? batchSize : 50; _batchSize = batchSize > 0 ? batchSize : 50;
_batchInterval = TimeSpan.FromMilliseconds(batchIntervalMs > 0 ? batchIntervalMs : 1000); _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) if (httpClient != null)
{ {
@@ -112,9 +120,16 @@ public sealed class LokiFlow : FlowBase
var line = BuildLine(logEvent); var line = BuildLine(logEvent);
var nanoTs = logEvent.Timestamp * 100L; // ticks → nanoseconds var nanoTs = logEvent.Timestamp * 100L; // ticks → nanoseconds
_queue.Enqueue((labels, line, nanoTs)); _queue.Enqueue((labels, line, nanoTs, DateTime.UtcNow));
Interlocked.Increment(ref BlastedCount); 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) if (_queue.Count >= _batchSize)
{ {
_batchSignal.Release(1); _batchSignal.Release(1);
@@ -168,8 +183,16 @@ public sealed class LokiFlow : FlowBase
// Group by label-set (Loki stream key) // Group by label-set (Loki stream key)
var groups = new Dictionary<string, List<(string Line, long NanoTs)>>(StringComparer.Ordinal); var groups = new Dictionary<string, List<(string Line, long NanoTs)>>(StringComparer.Ordinal);
var now = DateTime.UtcNow;
while (_queue.TryDequeue(out var item)) 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)) if (!groups.TryGetValue(item.Labels, out var list))
{ {
list = new List<(string, long)>(16); list = new List<(string, long)>(16);
@@ -677,12 +677,21 @@ public sealed class MessageTemplate
private static bool IsArgumentDriven(object?[] args) => args.Length > 0; private static bool IsArgumentDriven(object?[] args) => args.Length > 0;
// Template cache to avoid re-parsing the same strings // Template cache to avoid re-parsing the same strings
private const int MaxCacheSize = 5000;
private static readonly System.Collections.Concurrent.ConcurrentDictionary<string, MessageTemplate> _cache private static readonly System.Collections.Concurrent.ConcurrentDictionary<string, MessageTemplate> _cache
= new(StringComparer.Ordinal); = new(StringComparer.Ordinal);
/// <summary>Returns a cached parsed template (recommended for hot paths).</summary> /// <summary>Returns a cached parsed template (recommended for hot paths).</summary>
public static MessageTemplate FromCache(string template) => public static MessageTemplate FromCache(string template)
_cache.GetOrAdd(template, static t => Parse(t)); {
// 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));
}
/// <summary>Clears the template parse cache.</summary> /// <summary>Clears the template parse cache.</summary>
public static void ClearCache() => _cache.Clear(); public static void ClearCache() => _cache.Clear();
@@ -8,6 +8,7 @@ using System.Threading.Tasks;
namespace EonaCat.LogStack.Telemetry; namespace EonaCat.LogStack.Telemetry;
public sealed class TelemetryServer : IDisposable public sealed class TelemetryServer : IDisposable
{ {
private const int MaxEvents = 5000;
private readonly HttpListener _listener = new(); private readonly HttpListener _listener = new();
public ConcurrentQueue<TelemetryEvent> Events { get; } = new(); public ConcurrentQueue<TelemetryEvent> Events { get; } = new();
public int Count => Events.Count; public int Count => Events.Count;
@@ -30,6 +31,10 @@ public sealed class TelemetryServer : IDisposable
if (item != null) if (item != null)
{ {
Events.Enqueue(item); 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 == "/") else if (ctx.Request.Url?.AbsolutePath == "/")