From d136dc1e01154b481886ba069a02de830d3ff1a9 Mon Sep 17 00:00:00 2001 From: Alexey Osipov Date: Mon, 24 Aug 2026 09:00:31 +0300 Subject: [PATCH 1/8] Fix pubsub TTL cache eviction --- .../TtlCacheTests.cs | 45 +++++++ .../Libp2p.Protocols.Pubsub/TtlCache.cs | 110 +++++++++++++----- 2 files changed, 125 insertions(+), 30 deletions(-) create mode 100644 src/libp2p/Libp2p.Protocols.Pubsub.Tests/TtlCacheTests.cs diff --git a/src/libp2p/Libp2p.Protocols.Pubsub.Tests/TtlCacheTests.cs b/src/libp2p/Libp2p.Protocols.Pubsub.Tests/TtlCacheTests.cs new file mode 100644 index 00000000..e6ee9aee --- /dev/null +++ b/src/libp2p/Libp2p.Protocols.Pubsub.Tests/TtlCacheTests.cs @@ -0,0 +1,45 @@ +// SPDX-FileCopyrightText: 2026 Demerzel Solutions Limited +// SPDX-License-Identifier: MIT + +namespace Nethermind.Libp2p.Protocols.Pubsub.Tests; + +[TestFixture] +public class TtlCacheTests +{ + [Test] + public void RemoveExpired_RemovesEntriesRegardlessOfKeyOrder() + { + using TtlCache cache = new(50); + MessageId expiredHigh = new([0xFF]); + MessageId liveLow = new([0x01]); + + cache.Add(expiredHigh); + Thread.Sleep(80); + cache.Add(liveLow); + + cache.RemoveExpired(DateTimeOffset.UtcNow); + + Assert.Multiple(() => + { + Assert.That(cache.Contains(expiredHigh), Is.False); + Assert.That(cache.Contains(liveLow), Is.True); + }); + } + + [Test] + public void ExpiredEntries_AreNotReturnedBeforeTheSweeperRuns() + { + using TtlCache cache = new(25); + MessageId id = new([0x01]); + cache.Add(id, "value"); + + Thread.Sleep(60); + + Assert.Multiple(() => + { + Assert.That(cache.Contains(id), Is.False); + Assert.That(cache.TryGet(id, out _), Is.False); + Assert.That(cache.ToList(), Is.Empty); + }); + } +} diff --git a/src/libp2p/Libp2p.Protocols.Pubsub/TtlCache.cs b/src/libp2p/Libp2p.Protocols.Pubsub/TtlCache.cs index 1c777b74..f93d1fcf 100644 --- a/src/libp2p/Libp2p.Protocols.Pubsub/TtlCache.cs +++ b/src/libp2p/Libp2p.Protocols.Pubsub/TtlCache.cs @@ -6,68 +6,118 @@ namespace Nethermind.Libp2p.Protocols.Pubsub; internal class TtlCache : IDisposable where TKey : notnull { private readonly int ttl; + private readonly object sync = new(); + private readonly Dictionary items = []; + private readonly CancellationTokenSource sweeperCancellation = new(); + private int disposed; - private struct CachedItem - { - public TItem Item { get; set; } - public DateTimeOffset ValidTill { get; set; } - } - - private readonly SortedDictionary items = []; - private bool isDisposed; + private readonly record struct CachedItem(TItem Item, DateTimeOffset ValidTill); public TtlCache(int ttl) { + ArgumentOutOfRangeException.ThrowIfNegativeOrZero(ttl); this.ttl = ttl; - Task.Run(async () => + _ = Task.Run(async () => { - while (!isDisposed) + try { - await Task.Delay(5_000); - DateTimeOffset now = DateTimeOffset.UtcNow; - lock (items) + while (true) { - TKey[] keys = items.TakeWhile(i => i.Value.ValidTill < now).Select(i => i.Key).ToArray(); - foreach (TKey keyToRemove in keys) - { - items.Remove(keyToRemove); - } + await Task.Delay(Math.Min(5_000, ttl), sweeperCancellation.Token); + RemoveExpired(DateTimeOffset.UtcNow); } } + catch (OperationCanceledException) when (sweeperCancellation.IsCancellationRequested) + { + } }); } - public bool Contains(TKey key) => items.ContainsKey(key); + public bool Contains(TKey key) => TryGet(key, out _); - public TItem Get(TKey key) => items.GetValueOrDefault(key).Item; + public TItem Get(TKey key) => TryGet(key, out TItem item) ? item : default!; - public void Add(TKey key, TItem item) + public bool TryGet(TKey key, out TItem item) { - lock (items) + lock (sync) { - items.TryAdd(key, new CachedItem + if (items.TryGetValue(key, out CachedItem cachedItem) && cachedItem.ValidTill > DateTimeOffset.UtcNow) { - Item = item, - ValidTill = DateTimeOffset.UtcNow.AddMilliseconds(ttl), - }); + item = cachedItem.Item; + return true; + } + } + + item = default!; + return false; + } + + internal void RemoveExpired(DateTimeOffset now) + { + lock (sync) + { + if (items.Count == 0) + { + return; + } + + List? expired = null; + foreach ((TKey key, CachedItem item) in items) + { + if (item.ValidTill <= now) + { + (expired ??= []).Add(key); + } + } + + if (expired is not null) + { + foreach (TKey key in expired) + { + items.Remove(key); + } + } + } + } + + public void Add(TKey key, TItem item) + { + lock (sync) + { + items.TryAdd(key, new CachedItem(item, DateTimeOffset.UtcNow.AddMilliseconds(ttl))); } } public void Dispose() { - isDisposed = true; + if (Interlocked.Exchange(ref disposed, 1) != 0) + { + return; + } + + sweeperCancellation.Cancel(); + sweeperCancellation.Dispose(); + + lock (sync) + { + items.Clear(); + } } internal IList ToList() { - lock (items) + DateTimeOffset now = DateTimeOffset.UtcNow; + lock (sync) { - return items.Values.Select(i => i.Item).ToList(); + return items.Values + .Where(item => item.ValidTill > now) + .Select(item => item.Item) + .ToList(); } } } internal class TtlCache(int ttl) : TtlCache(ttl) where TKey : notnull { - public void Add(TKey key) => Add(key, false); + public void Add(TKey key) => Add(key, true); } From e6f0db95306a1c5f0a56bc6b303474dc4ffaa36c Mon Sep 17 00:00:00 2001 From: Alexey Osipov Date: Mon, 24 Aug 2026 10:30:45 +0300 Subject: [PATCH 2/8] Repair expired TTL cache entries --- .../Libp2p.Protocols.Pubsub.Tests/TtlCacheTests.cs | 13 +++++++++++++ src/libp2p/Libp2p.Protocols.Pubsub/TtlCache.cs | 9 ++++++++- 2 files changed, 21 insertions(+), 1 deletion(-) diff --git a/src/libp2p/Libp2p.Protocols.Pubsub.Tests/TtlCacheTests.cs b/src/libp2p/Libp2p.Protocols.Pubsub.Tests/TtlCacheTests.cs index e6ee9aee..8295e63b 100644 --- a/src/libp2p/Libp2p.Protocols.Pubsub.Tests/TtlCacheTests.cs +++ b/src/libp2p/Libp2p.Protocols.Pubsub.Tests/TtlCacheTests.cs @@ -42,4 +42,17 @@ public void ExpiredEntries_AreNotReturnedBeforeTheSweeperRuns() Assert.That(cache.ToList(), Is.Empty); }); } + + [Test] + public void Add_ReplacesAnExpiredEntry() + { + using TtlCache cache = new(25); + MessageId id = new([0x01]); + cache.Add(id, "expired"); + + Thread.Sleep(60); + cache.Add(id, "replacement"); + + Assert.That(cache.Get(id), Is.EqualTo("replacement")); + } } diff --git a/src/libp2p/Libp2p.Protocols.Pubsub/TtlCache.cs b/src/libp2p/Libp2p.Protocols.Pubsub/TtlCache.cs index f93d1fcf..dab96e98 100644 --- a/src/libp2p/Libp2p.Protocols.Pubsub/TtlCache.cs +++ b/src/libp2p/Libp2p.Protocols.Pubsub/TtlCache.cs @@ -9,6 +9,7 @@ internal class TtlCache : IDisposable where TKey : notnull private readonly object sync = new(); private readonly Dictionary items = []; private readonly CancellationTokenSource sweeperCancellation = new(); + private readonly Task sweeperTask; private int disposed; private readonly record struct CachedItem(TItem Item, DateTimeOffset ValidTill); @@ -17,7 +18,7 @@ public TtlCache(int ttl) { ArgumentOutOfRangeException.ThrowIfNegativeOrZero(ttl); this.ttl = ttl; - _ = Task.Run(async () => + sweeperTask = Task.Run(async () => { try { @@ -84,6 +85,11 @@ public void Add(TKey key, TItem item) { lock (sync) { + if (items.TryGetValue(key, out CachedItem cachedItem) && cachedItem.ValidTill <= DateTimeOffset.UtcNow) + { + items.Remove(key); + } + items.TryAdd(key, new CachedItem(item, DateTimeOffset.UtcNow.AddMilliseconds(ttl))); } } @@ -96,6 +102,7 @@ public void Dispose() } sweeperCancellation.Cancel(); + sweeperTask.GetAwaiter().GetResult(); sweeperCancellation.Dispose(); lock (sync) From ff883a863820a9d40251beb99a0d797624814f0d Mon Sep 17 00:00:00 2001 From: Alexey Osipov Date: Mon, 24 Aug 2026 10:55:09 +0300 Subject: [PATCH 3/8] Keep the TTL cache sweep cadence stable --- src/libp2p/Libp2p.Protocols.Pubsub.Tests/TtlCacheTests.cs | 2 +- src/libp2p/Libp2p.Protocols.Pubsub/TtlCache.cs | 2 +- 2 files changed, 2 insertions(+), 2 deletions(-) diff --git a/src/libp2p/Libp2p.Protocols.Pubsub.Tests/TtlCacheTests.cs b/src/libp2p/Libp2p.Protocols.Pubsub.Tests/TtlCacheTests.cs index 8295e63b..33ca2625 100644 --- a/src/libp2p/Libp2p.Protocols.Pubsub.Tests/TtlCacheTests.cs +++ b/src/libp2p/Libp2p.Protocols.Pubsub.Tests/TtlCacheTests.cs @@ -27,7 +27,7 @@ public void RemoveExpired_RemovesEntriesRegardlessOfKeyOrder() } [Test] - public void ExpiredEntries_AreNotReturnedBeforeTheSweeperRuns() + public void ExpiredEntries_AreNotReturned() { using TtlCache cache = new(25); MessageId id = new([0x01]); diff --git a/src/libp2p/Libp2p.Protocols.Pubsub/TtlCache.cs b/src/libp2p/Libp2p.Protocols.Pubsub/TtlCache.cs index dab96e98..92e49a82 100644 --- a/src/libp2p/Libp2p.Protocols.Pubsub/TtlCache.cs +++ b/src/libp2p/Libp2p.Protocols.Pubsub/TtlCache.cs @@ -24,7 +24,7 @@ public TtlCache(int ttl) { while (true) { - await Task.Delay(Math.Min(5_000, ttl), sweeperCancellation.Token); + await Task.Delay(5_000, sweeperCancellation.Token); RemoveExpired(DateTimeOffset.UtcNow); } } From 2026c7b8e14be5257329e4e7c215a1519755cd2b Mon Sep 17 00:00:00 2001 From: Alexey Osipov Date: Mon, 24 Aug 2026 11:00:04 +0300 Subject: [PATCH 4/8] Stabilize TTL cache timing tests --- .../Libp2p.Protocols.Pubsub.Tests/TtlCacheTests.cs | 12 ++++++------ 1 file changed, 6 insertions(+), 6 deletions(-) diff --git a/src/libp2p/Libp2p.Protocols.Pubsub.Tests/TtlCacheTests.cs b/src/libp2p/Libp2p.Protocols.Pubsub.Tests/TtlCacheTests.cs index 33ca2625..4efe490a 100644 --- a/src/libp2p/Libp2p.Protocols.Pubsub.Tests/TtlCacheTests.cs +++ b/src/libp2p/Libp2p.Protocols.Pubsub.Tests/TtlCacheTests.cs @@ -9,12 +9,12 @@ public class TtlCacheTests [Test] public void RemoveExpired_RemovesEntriesRegardlessOfKeyOrder() { - using TtlCache cache = new(50); + using TtlCache cache = new(500); MessageId expiredHigh = new([0xFF]); MessageId liveLow = new([0x01]); cache.Add(expiredHigh); - Thread.Sleep(80); + Thread.Sleep(750); cache.Add(liveLow); cache.RemoveExpired(DateTimeOffset.UtcNow); @@ -29,11 +29,11 @@ public void RemoveExpired_RemovesEntriesRegardlessOfKeyOrder() [Test] public void ExpiredEntries_AreNotReturned() { - using TtlCache cache = new(25); + using TtlCache cache = new(500); MessageId id = new([0x01]); cache.Add(id, "value"); - Thread.Sleep(60); + Thread.Sleep(750); Assert.Multiple(() => { @@ -46,11 +46,11 @@ public void ExpiredEntries_AreNotReturned() [Test] public void Add_ReplacesAnExpiredEntry() { - using TtlCache cache = new(25); + using TtlCache cache = new(500); MessageId id = new([0x01]); cache.Add(id, "expired"); - Thread.Sleep(60); + Thread.Sleep(750); cache.Add(id, "replacement"); Assert.That(cache.Get(id), Is.EqualTo("replacement")); From 7fda16b6f5c3a75cbdb69f7e853264c7d103cb4e Mon Sep 17 00:00:00 2001 From: Alexey Osipov Date: Mon, 24 Aug 2026 13:14:46 +0300 Subject: [PATCH 5/8] Bound TTL cache entries --- .../TtlCacheTests.cs | 20 +++++ .../Libp2p.Protocols.Pubsub/TtlCache.cs | 84 ++++++++++++++----- 2 files changed, 81 insertions(+), 23 deletions(-) diff --git a/src/libp2p/Libp2p.Protocols.Pubsub.Tests/TtlCacheTests.cs b/src/libp2p/Libp2p.Protocols.Pubsub.Tests/TtlCacheTests.cs index 4efe490a..b7079efd 100644 --- a/src/libp2p/Libp2p.Protocols.Pubsub.Tests/TtlCacheTests.cs +++ b/src/libp2p/Libp2p.Protocols.Pubsub.Tests/TtlCacheTests.cs @@ -55,4 +55,24 @@ public void Add_ReplacesAnExpiredEntry() Assert.That(cache.Get(id), Is.EqualTo("replacement")); } + + [Test] + public void Add_EvictsTheOldestLiveEntryAtCapacity() + { + using TtlCache cache = new(ttl: 1_000, maxEntries: 2); + MessageId first = new([0x01]); + MessageId second = new([0x02]); + MessageId third = new([0x03]); + + cache.Add(first, "first"); + cache.Add(second, "second"); + cache.Add(third, "third"); + + Assert.Multiple(() => + { + Assert.That(cache.Contains(first), Is.False); + Assert.That(cache.Get(second), Is.EqualTo("second")); + Assert.That(cache.Get(third), Is.EqualTo("third")); + }); + } } diff --git a/src/libp2p/Libp2p.Protocols.Pubsub/TtlCache.cs b/src/libp2p/Libp2p.Protocols.Pubsub/TtlCache.cs index 92e49a82..b7b8613e 100644 --- a/src/libp2p/Libp2p.Protocols.Pubsub/TtlCache.cs +++ b/src/libp2p/Libp2p.Protocols.Pubsub/TtlCache.cs @@ -6,18 +6,23 @@ namespace Nethermind.Libp2p.Protocols.Pubsub; internal class TtlCache : IDisposable where TKey : notnull { private readonly int ttl; + private readonly int maxEntries; private readonly object sync = new(); private readonly Dictionary items = []; + private readonly Queue<(TKey Key, long Sequence)> insertionOrder = []; private readonly CancellationTokenSource sweeperCancellation = new(); private readonly Task sweeperTask; private int disposed; + private long sequence; - private readonly record struct CachedItem(TItem Item, DateTimeOffset ValidTill); + private readonly record struct CachedItem(TItem Item, DateTimeOffset ValidTill, long Sequence); - public TtlCache(int ttl) + public TtlCache(int ttl, int maxEntries = int.MaxValue) { ArgumentOutOfRangeException.ThrowIfNegativeOrZero(ttl); + ArgumentOutOfRangeException.ThrowIfNegativeOrZero(maxEntries); this.ttl = ttl; + this.maxEntries = maxEntries; sweeperTask = Task.Run(async () => { try @@ -42,10 +47,15 @@ public bool TryGet(TKey key, out TItem item) { lock (sync) { - if (items.TryGetValue(key, out CachedItem cachedItem) && cachedItem.ValidTill > DateTimeOffset.UtcNow) + if (items.TryGetValue(key, out CachedItem cachedItem)) { - item = cachedItem.Item; - return true; + if (cachedItem.ValidTill > DateTimeOffset.UtcNow) + { + item = cachedItem.Item; + return true; + } + + items.Remove(key); } } @@ -57,41 +67,68 @@ internal void RemoveExpired(DateTimeOffset now) { lock (sync) { - if (items.Count == 0) - { - return; - } + RemoveExpiredLocked(now); + } + } - List? expired = null; - foreach ((TKey key, CachedItem item) in items) + public void Add(TKey key, TItem item) + { + lock (sync) + { + DateTimeOffset now = DateTimeOffset.UtcNow; + if (items.TryGetValue(key, out CachedItem cachedItem)) { - if (item.ValidTill <= now) + if (cachedItem.ValidTill > now) { - (expired ??= []).Add(key); + return; } + + items.Remove(key); } - if (expired is not null) + while (items.Count >= maxEntries) { - foreach (TKey key in expired) - { - items.Remove(key); - } + EvictOldest(); } + + long itemSequence = ++sequence; + items.Add(key, new CachedItem(item, now.AddMilliseconds(ttl), itemSequence)); + insertionOrder.Enqueue((key, itemSequence)); } } - public void Add(TKey key, TItem item) + private void RemoveExpiredLocked(DateTimeOffset now) { - lock (sync) + List? expired = null; + foreach ((TKey key, CachedItem item) in items) + { + if (item.ValidTill <= now) + { + (expired ??= []).Add(key); + } + } + + if (expired is not null) { - if (items.TryGetValue(key, out CachedItem cachedItem) && cachedItem.ValidTill <= DateTimeOffset.UtcNow) + foreach (TKey key in expired) { items.Remove(key); } + } + } - items.TryAdd(key, new CachedItem(item, DateTimeOffset.UtcNow.AddMilliseconds(ttl))); + private void EvictOldest() + { + while (insertionOrder.TryDequeue(out (TKey Key, long Sequence) oldest)) + { + if (items.TryGetValue(oldest.Key, out CachedItem item) && item.Sequence == oldest.Sequence) + { + items.Remove(oldest.Key); + return; + } } + + throw new InvalidOperationException("TTL cache insertion order was unexpectedly empty."); } public void Dispose() @@ -108,6 +145,7 @@ public void Dispose() lock (sync) { items.Clear(); + insertionOrder.Clear(); } } @@ -124,7 +162,7 @@ internal IList ToList() } } -internal class TtlCache(int ttl) : TtlCache(ttl) where TKey : notnull +internal class TtlCache(int ttl, int maxEntries = int.MaxValue) : TtlCache(ttl, maxEntries) where TKey : notnull { public void Add(TKey key) => Add(key, true); } From 187c9cde0c5af99f2e4e4f867cdab2869d407311 Mon Sep 17 00:00:00 2001 From: Alexey Osipov Date: Mon, 24 Aug 2026 13:40:28 +0300 Subject: [PATCH 6/8] Stabilize TTL cache regression coverage --- .../Libp2p.Protocols.Pubsub.Tests/TtlCacheTests.cs | 8 +++++--- src/libp2p/Libp2p.Protocols.Pubsub/TtlCache.cs | 11 +++++++++++ 2 files changed, 16 insertions(+), 3 deletions(-) diff --git a/src/libp2p/Libp2p.Protocols.Pubsub.Tests/TtlCacheTests.cs b/src/libp2p/Libp2p.Protocols.Pubsub.Tests/TtlCacheTests.cs index b7079efd..7e2f69fd 100644 --- a/src/libp2p/Libp2p.Protocols.Pubsub.Tests/TtlCacheTests.cs +++ b/src/libp2p/Libp2p.Protocols.Pubsub.Tests/TtlCacheTests.cs @@ -14,13 +14,15 @@ public void RemoveExpired_RemovesEntriesRegardlessOfKeyOrder() MessageId liveLow = new([0x01]); cache.Add(expiredHigh); - Thread.Sleep(750); + DateTimeOffset expiredAfter = DateTimeOffset.UtcNow.AddMilliseconds(500); + Assert.That(() => DateTimeOffset.UtcNow >= expiredAfter, Is.True.After(2_000, 25)); cache.Add(liveLow); cache.RemoveExpired(DateTimeOffset.UtcNow); Assert.Multiple(() => { + Assert.That(cache.Count, Is.EqualTo(1)); Assert.That(cache.Contains(expiredHigh), Is.False); Assert.That(cache.Contains(liveLow), Is.True); }); @@ -33,7 +35,7 @@ public void ExpiredEntries_AreNotReturned() MessageId id = new([0x01]); cache.Add(id, "value"); - Thread.Sleep(750); + Assert.That(() => cache.Contains(id), Is.False.After(2_000, 25)); Assert.Multiple(() => { @@ -50,7 +52,7 @@ public void Add_ReplacesAnExpiredEntry() MessageId id = new([0x01]); cache.Add(id, "expired"); - Thread.Sleep(750); + Assert.That(() => cache.Contains(id), Is.False.After(2_000, 25)); cache.Add(id, "replacement"); Assert.That(cache.Get(id), Is.EqualTo("replacement")); diff --git a/src/libp2p/Libp2p.Protocols.Pubsub/TtlCache.cs b/src/libp2p/Libp2p.Protocols.Pubsub/TtlCache.cs index b7b8613e..99c56c52 100644 --- a/src/libp2p/Libp2p.Protocols.Pubsub/TtlCache.cs +++ b/src/libp2p/Libp2p.Protocols.Pubsub/TtlCache.cs @@ -43,6 +43,17 @@ public TtlCache(int ttl, int maxEntries = int.MaxValue) public TItem Get(TKey key) => TryGet(key, out TItem item) ? item : default!; + internal int Count + { + get + { + lock (sync) + { + return items.Count; + } + } + } + public bool TryGet(TKey key, out TItem item) { lock (sync) From 73b1d01a6b0bd5e42fad7284c9601a7c87c8638b Mon Sep 17 00:00:00 2001 From: Alexey Osipov Date: Mon, 24 Aug 2026 14:01:23 +0300 Subject: [PATCH 7/8] Remove expired TTL cache entries from eviction order --- .../TtlCacheTests.cs | 1 + .../Libp2p.Protocols.Pubsub/TtlCache.cs | 49 ++++++++++++------- 2 files changed, 33 insertions(+), 17 deletions(-) diff --git a/src/libp2p/Libp2p.Protocols.Pubsub.Tests/TtlCacheTests.cs b/src/libp2p/Libp2p.Protocols.Pubsub.Tests/TtlCacheTests.cs index 7e2f69fd..8328d26a 100644 --- a/src/libp2p/Libp2p.Protocols.Pubsub.Tests/TtlCacheTests.cs +++ b/src/libp2p/Libp2p.Protocols.Pubsub.Tests/TtlCacheTests.cs @@ -23,6 +23,7 @@ public void RemoveExpired_RemovesEntriesRegardlessOfKeyOrder() Assert.Multiple(() => { Assert.That(cache.Count, Is.EqualTo(1)); + Assert.That(cache.EntryOrderCount, Is.EqualTo(1)); Assert.That(cache.Contains(expiredHigh), Is.False); Assert.That(cache.Contains(liveLow), Is.True); }); diff --git a/src/libp2p/Libp2p.Protocols.Pubsub/TtlCache.cs b/src/libp2p/Libp2p.Protocols.Pubsub/TtlCache.cs index 99c56c52..00acdd00 100644 --- a/src/libp2p/Libp2p.Protocols.Pubsub/TtlCache.cs +++ b/src/libp2p/Libp2p.Protocols.Pubsub/TtlCache.cs @@ -9,13 +9,11 @@ internal class TtlCache : IDisposable where TKey : notnull private readonly int maxEntries; private readonly object sync = new(); private readonly Dictionary items = []; - private readonly Queue<(TKey Key, long Sequence)> insertionOrder = []; + private readonly LinkedList insertionOrder = []; private readonly CancellationTokenSource sweeperCancellation = new(); private readonly Task sweeperTask; private int disposed; - private long sequence; - - private readonly record struct CachedItem(TItem Item, DateTimeOffset ValidTill, long Sequence); + private readonly record struct CachedItem(TItem Item, DateTimeOffset ValidTill, LinkedListNode Node); public TtlCache(int ttl, int maxEntries = int.MaxValue) { @@ -54,6 +52,17 @@ internal int Count } } + internal int EntryOrderCount + { + get + { + lock (sync) + { + return insertionOrder.Count; + } + } + } + public bool TryGet(TKey key, out TItem item) { lock (sync) @@ -66,7 +75,7 @@ public bool TryGet(TKey key, out TItem item) return true; } - items.Remove(key); + Remove(key, cachedItem); } } @@ -94,7 +103,7 @@ public void Add(TKey key, TItem item) return; } - items.Remove(key); + Remove(key, cachedItem); } while (items.Count >= maxEntries) @@ -102,9 +111,8 @@ public void Add(TKey key, TItem item) EvictOldest(); } - long itemSequence = ++sequence; - items.Add(key, new CachedItem(item, now.AddMilliseconds(ttl), itemSequence)); - insertionOrder.Enqueue((key, itemSequence)); + LinkedListNode node = insertionOrder.AddLast(key); + items.Add(key, new CachedItem(item, now.AddMilliseconds(ttl), node)); } } @@ -123,23 +131,30 @@ private void RemoveExpiredLocked(DateTimeOffset now) { foreach (TKey key in expired) { - items.Remove(key); + Remove(key, items[key]); } } } private void EvictOldest() { - while (insertionOrder.TryDequeue(out (TKey Key, long Sequence) oldest)) + LinkedListNode? oldest = insertionOrder.First; + if (oldest is null) { - if (items.TryGetValue(oldest.Key, out CachedItem item) && item.Sequence == oldest.Sequence) - { - items.Remove(oldest.Key); - return; - } + throw new InvalidOperationException("TTL cache insertion order was unexpectedly empty."); + } + + insertionOrder.RemoveFirst(); + if (!items.Remove(oldest.Value)) + { + throw new InvalidOperationException("TTL cache insertion order was out of sync with its entries."); } + } - throw new InvalidOperationException("TTL cache insertion order was unexpectedly empty."); + private void Remove(TKey key, CachedItem item) + { + items.Remove(key); + insertionOrder.Remove(item.Node); } public void Dispose() From 64f118c236f0276874125cc9d26b7c6539f14b0d Mon Sep 17 00:00:00 2001 From: Alexey Osipov Date: Tue, 22 Sep 2026 18:24:55 +0300 Subject: [PATCH 8/8] Address independent review feedback --- .../TtlCacheTests.cs | 115 +++++++++++++----- .../Libp2p.Protocols.Pubsub/TtlCache.cs | 70 +++-------- 2 files changed, 101 insertions(+), 84 deletions(-) diff --git a/src/libp2p/Libp2p.Protocols.Pubsub.Tests/TtlCacheTests.cs b/src/libp2p/Libp2p.Protocols.Pubsub.Tests/TtlCacheTests.cs index 8328d26a..f65014bd 100644 --- a/src/libp2p/Libp2p.Protocols.Pubsub.Tests/TtlCacheTests.cs +++ b/src/libp2p/Libp2p.Protocols.Pubsub.Tests/TtlCacheTests.cs @@ -1,6 +1,8 @@ // SPDX-FileCopyrightText: 2026 Demerzel Solutions Limited // SPDX-License-Identifier: MIT +using NSubstitute; + namespace Nethermind.Libp2p.Protocols.Pubsub.Tests; [TestFixture] @@ -9,73 +11,128 @@ public class TtlCacheTests [Test] public void RemoveExpired_RemovesEntriesRegardlessOfKeyOrder() { - using TtlCache cache = new(500); + TestTimeProvider clock = new(); + using TtlCache cache = new(500, clock); MessageId expiredHigh = new([0xFF]); MessageId liveLow = new([0x01]); cache.Add(expiredHigh); - DateTimeOffset expiredAfter = DateTimeOffset.UtcNow.AddMilliseconds(500); - Assert.That(() => DateTimeOffset.UtcNow >= expiredAfter, Is.True.After(2_000, 25)); + clock.UtcNow = clock.UtcNow.AddMilliseconds(500); cache.Add(liveLow); - cache.RemoveExpired(DateTimeOffset.UtcNow); + cache.RemoveExpired(clock.UtcNow); Assert.Multiple(() => { Assert.That(cache.Count, Is.EqualTo(1)); - Assert.That(cache.EntryOrderCount, Is.EqualTo(1)); Assert.That(cache.Contains(expiredHigh), Is.False); Assert.That(cache.Contains(liveLow), Is.True); }); } - [Test] - public void ExpiredEntries_AreNotReturned() + [TestCase("Contains")] + [TestCase("TryGet")] + [TestCase("Get")] + [TestCase("ToList")] + public void ExpiredEntries_AreNotReturnedBeforeSweeping(string operation) { - using TtlCache cache = new(500); + TestTimeProvider clock = new(); + using TtlCache cache = new(500, clock); MessageId id = new([0x01]); cache.Add(id, "value"); + clock.UtcNow = clock.UtcNow.AddMilliseconds(500); - Assert.That(() => cache.Contains(id), Is.False.After(2_000, 25)); - - Assert.Multiple(() => + // No cache read or sweep may remove the expired entry before the operation under test. + Assert.That(cache.Count, Is.EqualTo(1)); + switch (operation) { - Assert.That(cache.Contains(id), Is.False); - Assert.That(cache.TryGet(id, out _), Is.False); - Assert.That(cache.ToList(), Is.Empty); - }); + case "Contains": + Assert.That(cache.Contains(id), Is.False); + break; + case "TryGet": + Assert.That(cache.TryGet(id, out string value), Is.False); + Assert.That(value, Is.Null); + break; + case "Get": + Assert.That(cache.Get(id), Is.Null); + break; + case "ToList": + Assert.That(cache.ToList(), Is.Empty); + break; + } + } + + [Test] + public void ToList_ReturnsOnlyLiveEntries() + { + TestTimeProvider clock = new(); + using TtlCache cache = new(500, clock); + cache.Add(new([0x01]), "expired"); + clock.UtcNow = clock.UtcNow.AddMilliseconds(500); + cache.Add(new([0x02]), "live"); + + Assert.That(cache.Count, Is.EqualTo(2)); + Assert.That(cache.ToList(), Is.EqualTo(new[] { "live" })); } [Test] public void Add_ReplacesAnExpiredEntry() { - using TtlCache cache = new(500); + TestTimeProvider clock = new(); + using TtlCache cache = new(500, clock); MessageId id = new([0x01]); cache.Add(id, "expired"); + clock.UtcNow = clock.UtcNow.AddMilliseconds(500); - Assert.That(() => cache.Contains(id), Is.False.After(2_000, 25)); + Assert.That(cache.Count, Is.EqualTo(1)); cache.Add(id, "replacement"); + Assert.That(cache.Count, Is.EqualTo(1)); Assert.That(cache.Get(id), Is.EqualTo("replacement")); } [Test] - public void Add_EvictsTheOldestLiveEntryAtCapacity() + public void Add_DoesNotReplaceOrRefreshALiveEntry() { - using TtlCache cache = new(ttl: 1_000, maxEntries: 2); + TestTimeProvider clock = new(); + using TtlCache cache = new(500, clock); + MessageId id = new([0x01]); + cache.Add(id, "original"); + clock.UtcNow = clock.UtcNow.AddMilliseconds(250); + cache.Add(id, "replacement"); + + Assert.That(cache.Get(id), Is.EqualTo("original")); + clock.UtcNow = clock.UtcNow.AddMilliseconds(250); + Assert.That(cache.Contains(id), Is.False); + } + + [Test] + public void RemoveExpired_RemovesLaterEntriesAfterClockMovesBackward() + { + TestTimeProvider clock = new(); + using TtlCache cache = new(500, clock); MessageId first = new([0x01]); MessageId second = new([0x02]); - MessageId third = new([0x03]); + cache.Add(first); + clock.UtcNow = clock.UtcNow.AddMilliseconds(-250); + cache.Add(second); + clock.UtcNow = clock.UtcNow.AddMilliseconds(500); - cache.Add(first, "first"); - cache.Add(second, "second"); - cache.Add(third, "third"); + cache.RemoveExpired(clock.UtcNow); - Assert.Multiple(() => - { - Assert.That(cache.Contains(first), Is.False); - Assert.That(cache.Get(second), Is.EqualTo("second")); - Assert.That(cache.Get(third), Is.EqualTo("third")); - }); + Assert.That(cache.Count, Is.EqualTo(1)); + Assert.That(cache.Contains(first), Is.True); + Assert.That(cache.Contains(second), Is.False); + } + + private sealed class TestTimeProvider : TimeProvider + { + public DateTimeOffset UtcNow { get; set; } = new(2026, 1, 1, 0, 0, 0, TimeSpan.Zero); + + public override DateTimeOffset GetUtcNow() => UtcNow; + + // Keep the background sweeper dormant so each test controls expiry explicitly. + public override ITimer CreateTimer(TimerCallback callback, object? state, TimeSpan dueTime, TimeSpan period) + => Substitute.For(); } } diff --git a/src/libp2p/Libp2p.Protocols.Pubsub/TtlCache.cs b/src/libp2p/Libp2p.Protocols.Pubsub/TtlCache.cs index 00acdd00..8cbfc585 100644 --- a/src/libp2p/Libp2p.Protocols.Pubsub/TtlCache.cs +++ b/src/libp2p/Libp2p.Protocols.Pubsub/TtlCache.cs @@ -6,29 +6,27 @@ namespace Nethermind.Libp2p.Protocols.Pubsub; internal class TtlCache : IDisposable where TKey : notnull { private readonly int ttl; - private readonly int maxEntries; + private readonly TimeProvider timeProvider; private readonly object sync = new(); private readonly Dictionary items = []; - private readonly LinkedList insertionOrder = []; private readonly CancellationTokenSource sweeperCancellation = new(); private readonly Task sweeperTask; private int disposed; - private readonly record struct CachedItem(TItem Item, DateTimeOffset ValidTill, LinkedListNode Node); + private readonly record struct CachedItem(TItem Item, DateTimeOffset ValidTill); - public TtlCache(int ttl, int maxEntries = int.MaxValue) + public TtlCache(int ttl, TimeProvider? timeProvider = null) { ArgumentOutOfRangeException.ThrowIfNegativeOrZero(ttl); - ArgumentOutOfRangeException.ThrowIfNegativeOrZero(maxEntries); this.ttl = ttl; - this.maxEntries = maxEntries; + this.timeProvider = timeProvider ?? TimeProvider.System; sweeperTask = Task.Run(async () => { try { while (true) { - await Task.Delay(5_000, sweeperCancellation.Token); - RemoveExpired(DateTimeOffset.UtcNow); + await Task.Delay(TimeSpan.FromSeconds(5), this.timeProvider, sweeperCancellation.Token); + RemoveExpired(this.timeProvider.GetUtcNow()); } } catch (OperationCanceledException) when (sweeperCancellation.IsCancellationRequested) @@ -52,30 +50,19 @@ internal int Count } } - internal int EntryOrderCount - { - get - { - lock (sync) - { - return insertionOrder.Count; - } - } - } - public bool TryGet(TKey key, out TItem item) { lock (sync) { if (items.TryGetValue(key, out CachedItem cachedItem)) { - if (cachedItem.ValidTill > DateTimeOffset.UtcNow) + if (cachedItem.ValidTill > timeProvider.GetUtcNow()) { item = cachedItem.Item; return true; } - Remove(key, cachedItem); + items.Remove(key); } } @@ -95,7 +82,7 @@ public void Add(TKey key, TItem item) { lock (sync) { - DateTimeOffset now = DateTimeOffset.UtcNow; + DateTimeOffset now = timeProvider.GetUtcNow(); if (items.TryGetValue(key, out CachedItem cachedItem)) { if (cachedItem.ValidTill > now) @@ -103,21 +90,16 @@ public void Add(TKey key, TItem item) return; } - Remove(key, cachedItem); - } - - while (items.Count >= maxEntries) - { - EvictOldest(); + items.Remove(key); } - LinkedListNode node = insertionOrder.AddLast(key); - items.Add(key, new CachedItem(item, now.AddMilliseconds(ttl), node)); + items.Add(key, new CachedItem(item, now.AddMilliseconds(ttl))); } } private void RemoveExpiredLocked(DateTimeOffset now) { + // Wall-clock adjustments can make expiration order differ from insertion order. List? expired = null; foreach ((TKey key, CachedItem item) in items) { @@ -131,32 +113,11 @@ private void RemoveExpiredLocked(DateTimeOffset now) { foreach (TKey key in expired) { - Remove(key, items[key]); + items.Remove(key); } } } - private void EvictOldest() - { - LinkedListNode? oldest = insertionOrder.First; - if (oldest is null) - { - throw new InvalidOperationException("TTL cache insertion order was unexpectedly empty."); - } - - insertionOrder.RemoveFirst(); - if (!items.Remove(oldest.Value)) - { - throw new InvalidOperationException("TTL cache insertion order was out of sync with its entries."); - } - } - - private void Remove(TKey key, CachedItem item) - { - items.Remove(key); - insertionOrder.Remove(item.Node); - } - public void Dispose() { if (Interlocked.Exchange(ref disposed, 1) != 0) @@ -171,15 +132,14 @@ public void Dispose() lock (sync) { items.Clear(); - insertionOrder.Clear(); } } internal IList ToList() { - DateTimeOffset now = DateTimeOffset.UtcNow; lock (sync) { + DateTimeOffset now = timeProvider.GetUtcNow(); return items.Values .Where(item => item.ValidTill > now) .Select(item => item.Item) @@ -188,7 +148,7 @@ internal IList ToList() } } -internal class TtlCache(int ttl, int maxEntries = int.MaxValue) : TtlCache(ttl, maxEntries) where TKey : notnull +internal class TtlCache(int ttl, TimeProvider? timeProvider = null) : TtlCache(ttl, timeProvider) where TKey : notnull { public void Add(TKey key) => Add(key, true); }