2
0
mirror of https://github.com/esiur/esiur-dotnet.git synced 2026-06-13 14:38:43 +00:00

Prevent unbounded queue history retention

This commit is contained in:
2026-06-04 18:23:10 +03:00
parent 3cd611970a
commit 2963c0505b
5 changed files with 146 additions and 47 deletions
+40 -8
View File
@@ -50,7 +50,8 @@ public class AsyncQueue<T> : AsyncReply<T>
int currentId = 0;
int currentFlushId;
public List<AsyncQueueItem<T>> Processed = new();
bool captureProcessedItems;
List<AsyncQueueItem<T>> processed = new();
List<AsyncQueueItem<T>> list = new List<AsyncQueueItem<T>>();
//Action<T> callback;
@@ -63,6 +64,34 @@ public class AsyncQueue<T> : AsyncReply<T>
//return this;
//}
/// <summary>
/// Enables or disables retaining delivered queue items for diagnostics.
/// Capture is disabled by default, and disabling it discards retained items.
/// </summary>
public void SetProcessedCapture(bool enabled)
{
lock (queueLock)
{
captureProcessedItems = enabled;
if (!enabled)
processed = new();
}
}
/// <summary>
/// Atomically returns and removes the delivered queue items retained since the previous drain.
/// </summary>
public List<AsyncQueueItem<T>> DrainProcessed()
{
lock (queueLock)
{
var result = processed;
processed = new();
return result;
}
}
public void Add(AsyncReply<T> reply)
{
lock (queueLock)
@@ -121,13 +150,16 @@ public class AsyncQueue<T> : AsyncReply<T>
Trigger(list[i].Reply.Result);
resultReady = false;
var p = list[i];
p.Delivered = DateTime.Now;
p.Ready = p.Reply.ReadyTime;
p.BatchSize = batchSize;
p.FlushId = flushId;
//p.HasResource = p.Reply. (p.Ready - p.Arrival).TotalMilliseconds > 5;
Processed.Add(p);
if (captureProcessedItems)
{
var p = list[i];
p.Delivered = DateTime.Now;
p.Ready = p.Reply.ReadyTime;
p.BatchSize = batchSize;
p.FlushId = flushId;
//p.HasResource = p.Reply. (p.Ready - p.Arrival).TotalMilliseconds > 5;
processed.Add(p);
}
list.RemoveAt(i);
+13 -3
View File
@@ -378,11 +378,21 @@ public partial class EpConnection : NetworkConnection, IStore
}
/// <summary>
/// Enables or disables retaining delivered resource queue items for diagnostics.
/// Capture is disabled by default.
/// </summary>
public void SetFinishedQueueCapture(bool enabled)
{
_queue.SetProcessedCapture(enabled);
}
/// <summary>
/// Atomically returns and removes the retained delivered resource queue items.
/// </summary>
public List<AsyncQueueItem<EpResourceQueueItem>> GetFinishedQueue()
{
var l = _queue.Processed.ToArray().ToList();
_queue.Processed.Clear();
return l;
return _queue.DrainProcessed();
}
void init()