diff --git a/BetterStack.Logs.NLog/BetterStackLogsTarget.cs b/BetterStack.Logs.NLog/BetterStackLogsTarget.cs
index d9fb0df..9b0ce2f 100644
--- a/BetterStack.Logs.NLog/BetterStackLogsTarget.cs
+++ b/BetterStack.Logs.NLog/BetterStackLogsTarget.cs
@@ -67,6 +67,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.
///
@@ -114,6 +120,10 @@ protected override void InitializeTarget()
{
stopDrain();
+ 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());
@@ -151,7 +161,8 @@ protected override void InitializeTarget()
betterStackDrain = new Drain(
client,
period: TimeSpan.FromMilliseconds(FlushPeriodMilliseconds),
- maxBatchSize: MaxBatchSize
+ maxBatchSize: MaxBatchSize,
+ maxQueueSize: MaxQueueSize
);
betterStackClient = client;
diff --git a/BetterStack.Logs/Drain.cs b/BetterStack.Logs/Drain.cs
index 5a2918f..697d7ce 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;
@@ -24,6 +25,10 @@ public sealed class Drain
private readonly ConcurrentQueue> flushRequests = new ConcurrentQueue>();
private readonly SemaphoreSlim flushSignal = new SemaphoreSlim(0);
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.
@@ -33,11 +38,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);
@@ -45,12 +65,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);
}
@@ -111,10 +141,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 fce1ed4..c64c01e 100644
--- a/example-project/README.md
+++ b/example-project/README.md
@@ -254,6 +254,8 @@ The exception, with its stack trace, is sent in the top-level `exception` field:
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. `retries` is how many times a failed request is retried after the first attempt, 10 by default.
+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
diff --git a/tests/BetterStack.Logs.NLog.Tests/BetterStackLogsTargetTests.cs b/tests/BetterStack.Logs.NLog.Tests/BetterStackLogsTargetTests.cs
index 3a08d0b..a4e796d 100644
--- a/tests/BetterStack.Logs.NLog.Tests/BetterStackLogsTargetTests.cs
+++ b/tests/BetterStack.Logs.NLog.Tests/BetterStackLogsTargetTests.cs
@@ -739,6 +739,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()
{