From a3e92aa6a8193ee2d7f90454f679532324989a8d Mon Sep 17 00:00:00 2001 From: Petr Heinz Date: Wed, 30 Sep 2026 17:19:16 +0200 Subject: [PATCH 1/2] Add failing tests: the queue grows without limit The tests use the new maxQueueSize option, so this commit does not compile yet. With only the property added, all 8 logs written while a batch is stuck are queued and delivered, nothing is reported, and maxQueueSize="0" is accepted. Co-Authored-By: Claude Opus 5.5 --- .../BetterStackLogsTargetTests.cs | 70 +++++++++++++++++++ 1 file changed, 70 insertions(+) diff --git a/tests/BetterStack.Logs.NLog.Tests/BetterStackLogsTargetTests.cs b/tests/BetterStack.Logs.NLog.Tests/BetterStackLogsTargetTests.cs index a7de84a..4d29616 100644 --- a/tests/BetterStack.Logs.NLog.Tests/BetterStackLogsTargetTests.cs +++ b/tests/BetterStack.Logs.NLog.Tests/BetterStackLogsTargetTests.cs @@ -2,6 +2,7 @@ using System.Linq; using Newtonsoft.Json.Linq; using NLog; +using NLog.Common; using NLog.Config; using NLog.Targets; using NLog.Targets.Wrappers; @@ -13,12 +14,16 @@ public class BetterStackLogsTargetTests : IDisposable { private readonly FakeIngestion ingestion = new FakeIngestion(); private readonly LogFactory logFactory = new LogFactory(); + private readonly System.IO.TextWriter originalInternalLogWriter = InternalLogger.LogWriter; + private readonly LogLevel originalInternalLogLevel = InternalLogger.LogLevel; public void Dispose() { logFactory.Shutdown(); ingestion.Dispose(); GlobalDiagnosticsContext.Clear(); + InternalLogger.LogWriter = originalInternalLogWriter; + InternalLogger.LogLevel = originalInternalLogLevel; } private BetterStackLogsTarget Target() => new BetterStackLogsTarget { @@ -331,6 +336,71 @@ public void RetriesTheBatchAfterServerError() Assert.Equal(failed.Body, retried.Body); } + [Fact] + public void DropsNewLogsOnceTheQueueIsFull() + { + var internalLog = new System.IO.StringWriter(); + InternalLogger.LogLevel = LogLevel.Error; + InternalLogger.LogWriter = internalLog; + ingestion.StatusCodes.Enqueue(500); + var target = Target(); + target.MaxQueueSize = 5; + var logger = LoggerFor(target); + + logger.Info("Stuck"); + ingestion.NextRequest(); + // Written while the failed batch waits a second for its retry + for (var i = 1; i <= 8; i++) logger.Info("Log " + i); + + Assert.Equal("Stuck", (string)Assert.Single(ingestion.NextRequest().Logs)["message"]); + Assert.Equal(new[] { "Log 1", "Log 2", "Log 3", "Log 4", "Log 5" }, ingestion.NextRequest().Logs.Select(log => (string)log["message"])); + var error = Assert.Single(internalLog.ToString().Split(new[] { Environment.NewLine }, StringSplitOptions.RemoveEmptyEntries)); + Assert.EndsWith("Error BetterStack.Logs: maximum number of logs in the queue reached (5). New logs will be dropped.", error); + } + + [Fact] + public void ReportsTheNextOverflowOnceTheQueueHasRoomAgain() + { + var internalLog = new System.IO.StringWriter(); + InternalLogger.LogLevel = LogLevel.Error; + InternalLogger.LogWriter = internalLog; + var target = Target(); + target.MaxQueueSize = 1; + var logger = LoggerFor(target); + + for (var round = 1; round <= 2; round++) { + ingestion.StatusCodes.Enqueue(500); + logger.Info("Stuck"); + ingestion.NextRequest(); + // Written while the failed batch waits a second for its retry + logger.Info("Queued"); + logger.Info("Dropped"); + ingestion.NextRequest(); + Assert.Equal("Queued", (string)Assert.Single(ingestion.NextRequest().Logs)["message"]); + } + + var errors = internalLog.ToString().Split(new[] { Environment.NewLine }, StringSplitOptions.RemoveEmptyEntries); + Assert.Equal(2, errors.Length); + Assert.All(errors, error => Assert.EndsWith("Error BetterStack.Logs: maximum number of logs in the queue reached (1). New logs will be dropped.", error)); + } + + [Fact] + public void KeepsAtMost100000LogsInTheQueueByDefault() + { + Assert.Equal(100000, new BetterStackLogsTarget().MaxQueueSize); + } + + [Fact] + public void ReportsMaxQueueSizeBelowOneAsConfigurationError() + { + logFactory.ThrowConfigExceptions = true; + var target = Target(); + target.MaxQueueSize = 0; + + var exception = Assert.Throws(() => LoggerFor(target)); + Assert.Equal("BetterStack.Logs: maxQueueSize is 0. Set it to 1 or more.", exception.Message); + } + [Fact] public void KeepsDeliveringAfterBatchRanOutOfRetries() { From ce8ce6da2d09fefff8351fed3a84fc0d0d45a043 Mon Sep 17 00:00:00 2001 From: Petr Heinz Date: Wed, 30 Sep 2026 17:21:28 +0200 Subject: [PATCH 2/2] Drop new logs once the queue holds maxQueueSize of them The drain queued without limit, so an application whose endpoint is down grew until it was killed. New target option maxQueueSize (100000 by default) bounds the queue like the Java client does: a full queue drops new logs and reports the first one in NLog's internal log, again for the next overflow once the drain has taken logs off the queue. The length is an Interlocked counter, since ConcurrentQueue.Count walks the segments. maxQueueSize below 1 is a configuration error. Drain keeps its constructor and gets an overload that takes the limit. Co-Authored-By: Claude Opus 5.5 --- .../BetterStackLogsTarget.cs | 13 ++++++- BetterStack.Logs/Drain.cs | 37 ++++++++++++++++++- example-project/README.md | 3 ++ 3 files changed, 51 insertions(+), 2 deletions(-) diff --git a/BetterStack.Logs.NLog/BetterStackLogsTarget.cs b/BetterStack.Logs.NLog/BetterStackLogsTarget.cs index 4071b2d..d87018b 100644 --- a/BetterStack.Logs.NLog/BetterStackLogsTarget.cs +++ b/BetterStack.Logs.NLog/BetterStackLogsTarget.cs @@ -56,6 +56,12 @@ public bool CaptureSourceLocation /// public bool IncludeGlobalDiagnosticContext { get; set; } = true; + /// + /// Maximum number of logs kept in memory while they wait to be sent. Once the queue holds that many, + /// new logs are dropped and an error is written to NLog's internal log. + /// + public int MaxQueueSize { get; set; } = 100000; + /// /// Control callsite capture of source-file and source-linenumber. /// @@ -95,6 +101,10 @@ protected override void InitializeTarget() { betterStackDrain?.Stop().Wait(); + if (MaxQueueSize < 1) { + throw new NLogConfigurationException($"BetterStack.Logs: maxQueueSize is {MaxQueueSize}. Set it to 1 or more."); + } + var sourceToken = RenderLogEvent(SourceToken, LogEventInfo.CreateNullEvent()); var endpoint = RenderLogEvent(Endpoint, LogEventInfo.CreateNullEvent()); @@ -107,7 +117,8 @@ protected override void InitializeTarget() betterStackDrain = new Drain( client, period: TimeSpan.FromMilliseconds(FlushPeriodMilliseconds), - maxBatchSize: MaxBatchSize + maxBatchSize: MaxBatchSize, + maxQueueSize: MaxQueueSize ); base.InitializeTarget(); diff --git a/BetterStack.Logs/Drain.cs b/BetterStack.Logs/Drain.cs index 3824444..becacb9 100644 --- a/BetterStack.Logs/Drain.cs +++ b/BetterStack.Logs/Drain.cs @@ -14,6 +14,7 @@ namespace BetterStack.Logs public sealed class Drain { private readonly int maxBatchSize; + private readonly int maxQueueSize; private readonly Client client; private readonly TimeSpan period; @@ -22,6 +23,10 @@ public sealed class Drain private ConcurrentQueue queue = new ConcurrentQueue(); private CancellationTokenSource cancellationTokenSource; + // Kept on every enqueue and dequeue: ConcurrentQueue.Count walks all segments of the queue + private int queueLength; + // 1 from the first log dropped on a full queue until the drain takes logs off the queue again + private int overflowReported; /// /// Initializes a Better Stack Logs drain and starts periodic logs delivery. @@ -31,11 +36,26 @@ public Drain( TimeSpan? period = null, int maxBatchSize = 1000, CancellationToken? cancellationToken = null + ) : this(client, period, maxBatchSize, 100000, cancellationToken) + { + } + + /// + /// Initializes a Better Stack Logs drain that holds at most maxQueueSize logs waiting to be delivered, + /// and starts periodic logs delivery. + /// + public Drain( + Client client, + TimeSpan? period, + int maxBatchSize, + int maxQueueSize, + CancellationToken? cancellationToken = null ) { this.client = client; this.period = period ?? TimeSpan.FromMilliseconds(250); this.maxBatchSize = maxBatchSize; + this.maxQueueSize = maxQueueSize; this.cancellationTokenSource = CancellationTokenSource.CreateLinkedTokenSource(cancellationToken ?? CancellationToken.None); runningTask = Task.Run(run); @@ -43,12 +63,22 @@ public Drain( /// /// Adds a single log event to a queue. The log event will be delivered later in a batch. + /// The log event is dropped when the queue already holds maxQueueSize of them. /// This method will throw an exception if the Drain is stopped. /// public void Enqueue(Log log) { if (cancellationTokenSource.IsCancellationRequested) throw new DrainIsClosedException(); + // Like in our Java client, a full queue drops the new log and keeps the ones it holds + if (Interlocked.Increment(ref queueLength) > maxQueueSize) { + Interlocked.Decrement(ref queueLength); + if (Interlocked.Exchange(ref overflowReported, 1) == 0) { + global::NLog.Common.InternalLogger.Error("BetterStack.Logs: maximum number of logs in the queue reached ({0}). New logs will be dropped.", maxQueueSize); + } + return; + } + queue.Enqueue(log); } @@ -86,10 +116,15 @@ private async Task flush() { var nextBatch = new List(expectedItemsCount); while (!queue.IsEmpty && nextBatch.Count < maxBatchSize) { - if (queue.TryDequeue(out var log)) nextBatch.Add(log); + if (queue.TryDequeue(out var log)) { + Interlocked.Decrement(ref queueLength); + nextBatch.Add(log); + } } if (nextBatch.Count > 0) { + // The queue has room again: its next overflow is reported again + Volatile.Write(ref overflowReported, 0); await client.Send(nextBatch); } } diff --git a/example-project/README.md b/example-project/README.md index a0c9a2a..354c778 100644 --- a/example-project/README.md +++ b/example-project/README.md @@ -220,6 +220,8 @@ This will create the following JSON output: The BetterStack.Logs target will send you logs periodically in batches to optimize network traffic with several retries in case of unexpected HTTP errors. You can adjust this behavior by setting the `maxBatchSize`, `flushPeriodMilliseconds`, and `retries` parameters to your custom values in your config. +While the logs wait to be sent, the target keeps at most `maxQueueSize` of them in memory, 100000 by default. When the endpoint cannot be reached for a while and the queue fills up, new logs are dropped and an error is written to NLog's internal log. + ```xml ```