Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
56 commits
Select commit Hold shift + click to select a range
d136dc1
Fix pubsub TTL cache eviction
flcl42 Aug 24, 2026
e6f0db9
Repair expired TTL cache entries
flcl42 Aug 24, 2026
ff883a8
Keep the TTL cache sweep cadence stable
flcl42 Aug 24, 2026
2026c7b
Stabilize TTL cache timing tests
flcl42 Aug 24, 2026
7fda16b
Bound TTL cache entries
flcl42 Aug 24, 2026
187c9cd
Stabilize TTL cache regression coverage
flcl42 Aug 24, 2026
73b1d01
Remove expired TTL cache entries from eviction order
flcl42 Aug 24, 2026
e7f6b7d
Bound incoming pubsub RPC frames
flcl42 Aug 24, 2026
a4dad34
Validate pubsub RPC frame limits
flcl42 Aug 24, 2026
23455df
Reject invalid pubsub RPC frame settings
flcl42 Aug 24, 2026
3baa4fc
Preserve the varint reader overload
flcl42 Aug 24, 2026
b2712d8
Tighten pubsub frame limit validation
flcl42 Aug 24, 2026
3a51110
Repair pubsub topic lifecycle
flcl42 Aug 24, 2026
41f5a08
Complete pubsub topic unsubscribe lifecycle
flcl42 Aug 24, 2026
cfb12a2
Use a clear topic lifecycle test name
flcl42 Aug 24, 2026
8d4c668
Gossip publish-only topics
flcl42 Aug 24, 2026
418d501
Ignore gossip for inactive topics
flcl42 Aug 24, 2026
ab72e3e
Remove stale pubsub fanout state
flcl42 Aug 24, 2026
344edad
Serialize pubsub topic lifecycle transitions
flcl42 Aug 24, 2026
367e146
Serialize pubsub publish routing
flcl42 Aug 24, 2026
063c3fe
Add Gossipsub v1.3 extensions control support
flcl42 Aug 24, 2026
57c3ee0
Harden strict pubsub no-sign validation
flcl42 Aug 24, 2026
5e8c219
Use opaque extension topic IDs
flcl42 Aug 24, 2026
0cb8a79
Require message IDs for unsigned pubsub
flcl42 Aug 24, 2026
7f6400b
Reject extensions outside Gossipsub v1.3
flcl42 Aug 24, 2026
34d12d5
Support Gossipsub partial messages
flcl42 Aug 24, 2026
a5835b5
Harden partial message negotiation
flcl42 Aug 24, 2026
1b034fd
Honor partial message subscription flags
flcl42 Aug 24, 2026
167fe48
Complete partial message routing
flcl42 Aug 24, 2026
ead2e85
Serialize partial subscription announcements
flcl42 Aug 24, 2026
dd5a697
Keep partial-message topics opt-in
flcl42 Aug 24, 2026
4524361
Add Gossipsub direct peer support
flcl42 Aug 24, 2026
ea94233
Back off direct peer grafts
flcl42 Aug 24, 2026
d4f4e5f
Test direct peer graft backoff
flcl42 Aug 24, 2026
4ef6232
Allow publishing without a local subscription
flcl42 Aug 24, 2026
afad1db
Name invalid direct peer configuration
flcl42 Aug 24, 2026
19c8581
Address independent review feedback
flcl42 Sep 22, 2026
265a0b3
Address independent review feedback
flcl42 Sep 22, 2026
64f118c
Address independent review feedback
flcl42 Sep 22, 2026
4da538e
Merge branch 'pubsub-ttl-cache' into pubsub-rpc-limit
flcl42 Sep 22, 2026
10a37dd
Address independent review feedback
flcl42 Sep 22, 2026
035a1df
Merge branch 'pubsub-topic-lifecycle' into pubsub-strict-no-sign
flcl42 Sep 22, 2026
cc08c01
Address independent review feedback
flcl42 Sep 22, 2026
4899641
Merge branch 'pubsub-strict-no-sign' into gossipsub-v13-extensions
flcl42 Sep 22, 2026
a8e70fe
Address independent review feedback
flcl42 Sep 22, 2026
38b0731
Merge branch 'gossipsub-v13-partial-messages' into gossipsub-direct-p…
flcl42 Sep 22, 2026
5e3d446
Address independent review feedback
flcl42 Sep 22, 2026
52630fc
Merge corrected RPC limits base
flcl42 Sep 22, 2026
c5fa2ed
Merge branch 'pubsub-topic-lifecycle' into pubsub-strict-no-sign
flcl42 Sep 22, 2026
332948c
Merge branch 'pubsub-strict-no-sign' into gossipsub-v13-extensions
flcl42 Sep 22, 2026
980e229
Merge corrected Gossipsub extensions base
flcl42 Sep 22, 2026
9c035a2
Merge corrected partial messages base
flcl42 Sep 22, 2026
93a518e
Use spelling-friendly partial message test name
flcl42 Sep 22, 2026
7442e46
Merge branch 'gossipsub-v13-partial-messages' into gossipsub-direct-p…
flcl42 Sep 22, 2026
679c538
Merge restacked partial messages base
flcl42 Sep 23, 2026
1a92986
Dial direct peers from a single deduplicated path
flcl42 Sep 23, 2026
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
322 changes: 322 additions & 0 deletions src/libp2p/Libp2p.Protocols.Pubsub.Tests/DirectPeersTests.cs
Original file line number Diff line number Diff line change
@@ -0,0 +1,322 @@
// SPDX-FileCopyrightText: 2026 Demerzel Solutions Limited
// SPDX-License-Identifier: MIT

using Multiformats.Address;
using Nethermind.Libp2p.Core.Discovery;
using Nethermind.Libp2p.Protocols;
using Nethermind.Libp2p.Protocols.Pubsub.Dto;
using System.Collections.ObjectModel;

namespace Nethermind.Libp2p.Protocols.Pubsub.Tests;

[TestFixture]
public class DirectPeersTests
{
[Test]
public void DirectPeers_ForwardValidMessagesDespitePeerScores()
{
const string topic = "topic";
Multiaddress senderAddress = TestPeers.Multiaddr(1);
Multiaddress receiverAddress = TestPeers.Multiaddr(2);
PeerId senderPeerId = senderAddress.GetPeerId()!;
PeerId receiverPeerId = receiverAddress.GetPeerId()!;
PubsubRouter router = new(
new PeerStore(),
new PubsubSettings { DirectPeers = [senderAddress, receiverAddress] });
ITopic localTopic = router.GetTopic(topic);
List<Rpc> senderRpcs = [];
List<Rpc> receiverRpcs = [];
TaskCompletionSource senderConnection = new();
TaskCompletionSource receiverConnection = new();
router.OutboundConnection(senderAddress, PubsubRouter.GossipsubProtocolVersionV11, senderConnection.Task, senderRpcs.Add);
router.OutboundConnection(receiverAddress, PubsubRouter.GossipsubProtocolVersionV11, receiverConnection.Task, receiverRpcs.Add);
router.OnRpc(senderPeerId, new Rpc().WithTopics([topic], []));
router.OnRpc(receiverPeerId, new Rpc().WithTopics([topic], []));

router.SetAppSpecificScore(senderPeerId, -20);
router.SetAppSpecificScore(receiverPeerId, -20);
senderRpcs.Clear();
receiverRpcs.Clear();

PeerId? receivedFrom = null;
localTopic.OnMessage += (peerId, _) => receivedFrom = peerId;
Identity author = TestPeers.Identity(1);
router.OnRpc(senderPeerId, new Rpc().WithMessages(topic, 1, author.PeerId.Bytes, [1, 2, 3], author));

Assert.Multiple(() =>
{
Assert.That(receivedFrom, Is.EqualTo(senderPeerId));
Assert.That(receiverRpcs.Single().Publish.Single().Data.ToByteArray(), Is.EqualTo(new byte[] { 1, 2, 3 }));
Assert.That(senderRpcs, Is.Empty);
});

senderConnection.SetResult();
receiverConnection.SetResult();
}

[Test]
public void DirectPeers_AreNeverAddedToTheMesh()
{
const string topic = "topic";
Multiaddress directAddress = TestPeers.Multiaddr(1);
Multiaddress firstMeshAddress = TestPeers.Multiaddr(2);
Multiaddress secondMeshAddress = TestPeers.Multiaddr(3);
PubsubRouter router = new(
new PeerStore(),
new PubsubSettings { DirectPeers = [directAddress] });
IRoutingStateContainer state = router;
_ = router.GetTopic(topic);
TaskCompletionSource connection = new();

foreach (Multiaddress address in new[] { directAddress, firstMeshAddress, secondMeshAddress })
{
router.OutboundConnection(address, PubsubRouter.GossipsubProtocolVersionV11, connection.Task, _ => { });
router.OnRpc(address.GetPeerId()!, new Rpc().WithTopics([topic], []));
}

router.Heartbeat().GetAwaiter().GetResult();

Assert.Multiple(() =>
{
Assert.That(state.GossipsubPeers[topic], Has.Member(directAddress.GetPeerId()));
Assert.That(state.Mesh[topic], Does.Not.Contain(directAddress.GetPeerId()));
Assert.That(state.Mesh[topic], Has.Count.EqualTo(2));
});

connection.SetResult();
}

[Test]
public void DirectPeerGrafts_AreRejectedWithPrune()
{
const string topic = "topic";
Multiaddress directAddress = TestPeers.Multiaddr(1);
PeerId directPeerId = directAddress.GetPeerId()!;
PubsubRouter router = new(
new PeerStore(),
new PubsubSettings { DirectPeers = [directAddress], PruneBackoff = 2_000 });
_ = router.GetTopic(topic);
List<Rpc> sentRpcs = [];
TaskCompletionSource connection = new();
router.OutboundConnection(directAddress, PubsubRouter.GossipsubProtocolVersionV11, connection.Task, sentRpcs.Add);
sentRpcs.Clear();

Rpc graft = new() { Control = new ControlMessage() };
graft.Control.Graft.Add(new ControlGraft { TopicID = topic });
router.OnRpc(directPeerId, graft);

ControlPrune prune = sentRpcs.Single().Control.Prune.Single();
Assert.Multiple(() =>
{
Assert.That(prune.TopicID, Is.EqualTo(topic));
Assert.That(prune.Backoff, Is.EqualTo(2));
});
connection.SetResult();
}

[TestCase(false)]
[TestCase(true)]
public async Task Router_ConnectsConfiguredDirectPeersAtStartup(bool addressesAlreadyKnown)
{
PeerStore peerStore = new();
Multiaddress directAddress = TestPeers.Multiaddr(1);
PeerId directPeerId = directAddress.GetPeerId()!;
peerStore.GetPeerInfo(directPeerId).SupportedProtocols = [PubsubRouter.GossipsubProtocolVersionV12];
if (addressesAlreadyKnown)
{
peerStore.Discover([directAddress]);
}
PubsubRouter router = new(peerStore, new PubsubSettings { DirectPeers = [directAddress] });

TaskCompletionSource protocolDialed = new(TaskCreationOptions.RunContinuationsAsynchronously);
ISession session = Substitute.For<ISession>();
session.RemoteAddress.Returns(directAddress);
session.DialAsync<GossipsubProtocolV12>(Arg.Any<CancellationToken>()).Returns(_ =>
{
protocolDialed.TrySetResult();
return Task.CompletedTask;
});

ILocalPeer localPeer = Substitute.For<ILocalPeer>();
localPeer.Identity.Returns(TestPeers.Identity(2));
localPeer.ListenAddresses.Returns(new ObservableCollection<Multiaddress>());
localPeer.DialAsync(Arg.Any<Multiaddress[]>(), Arg.Any<CancellationToken>()).Returns(session);

using CancellationTokenSource cancellation = new();
await router.StartAsync(localPeer, cancellation.Token);

await protocolDialed.Task.WaitAsync(TimeSpan.FromSeconds(2));
_ = localPeer.Received(1).DialAsync(Arg.Any<Multiaddress[]>(), Arg.Any<CancellationToken>());
_ = session.Received(1).DialAsync<GossipsubProtocolV12>(Arg.Any<CancellationToken>());

cancellation.Cancel();
}

[Test]
public async Task PublishPartial_ExcludesDirectPeersFromFanoutAndMesh()
{
const string topic = "topic";
Multiaddress directAddress = TestPeers.Multiaddr(1);
Multiaddress otherAddress = TestPeers.Multiaddr(2);
using PubsubRouter router = new(new PeerStore(), new PubsubSettings
{
DirectPeers = [directAddress],
EnablePartialMessages = true,
});
IRoutingStateContainer state = router;
router.GetPartialMessagesTopic(topic,
new PartialMessagesTopicOptions { SupportsSendingPartialMessages = true }, subscribe: false);
ILocalPeer localPeer = Substitute.For<ILocalPeer>();
localPeer.Identity.Returns(TestPeers.Identity(3));
localPeer.ListenAddresses.Returns(new ObservableCollection<Multiaddress>());
using CancellationTokenSource cancellation = new();
TaskCompletionSource connection = new();
try
{
await router.StartAsync(localPeer, cancellation.Token);
foreach (Multiaddress address in new[] { directAddress, otherAddress })
{
router.OutboundConnection(address, PubsubRouter.GossipsubProtocolVersionV13, connection.Task, _ => { });
router.OnRpc(address.GetPeerId()!, new Rpc().WithTopics([topic], []));
}

router.PublishPartial(topic, [1], partialMessage: [2]);
Assert.That(state.Fanout[topic], Is.EquivalentTo(new[] { otherAddress.GetPeerId() }));

router.Subscribe(topic);
Assert.That(state.Mesh[topic], Is.EquivalentTo(new[] { otherAddress.GetPeerId() }));
}
finally
{
cancellation.Cancel();
connection.TrySetResult();
}
}

[Test]
public void Subscribe_ExcludesDirectPeersWhenMovingFanoutToMesh()
{
const string topic = "topic";
Multiaddress directAddress = TestPeers.Multiaddr(1);
Multiaddress otherAddress = TestPeers.Multiaddr(2);
PeerId otherPeerId = otherAddress.GetPeerId()!;
using PubsubRouter router = new(new PeerStore(), new PubsubSettings { DirectPeers = [directAddress] });
IRoutingStateContainer state = router;
TaskCompletionSource connection = new();
try
{
router.OutboundConnection(otherAddress, PubsubRouter.GossipsubProtocolVersionV12, connection.Task, _ => { });
router.OnRpc(otherPeerId, new Rpc().WithTopics([topic], []));
state.Fanout[topic] = [directAddress.GetPeerId()!, otherPeerId];

router.Subscribe(topic);

Assert.That(state.Mesh[topic], Is.EquivalentTo(new[] { otherPeerId }));
Assert.That(state.Fanout.ContainsKey(topic), Is.False);
}
finally
{
connection.SetResult();
}
}

[Test]
public async Task Router_ReconnectsDisconnectedDirectPeersIndependentlyOfOtherIntervals()
{
Multiaddress directAddress = TestPeers.Multiaddr(1);
PeerStore peerStore = new();
peerStore.GetPeerInfo(directAddress.GetPeerId()!).SupportedProtocols = [PubsubRouter.GossipsubProtocolVersionV12];
using PubsubRouter router = new(peerStore, new PubsubSettings
{
DirectPeers = [directAddress],
DirectConnectPeriod = 50,
ReconnectionPeriod = 60_000,
HeartbeatInterval = 60_000,
});
TaskCompletionSource firstConnection = new();
TaskCompletionSource secondConnection = new();
TaskCompletionSource redialed = new(TaskCreationOptions.RunContinuationsAsynchronously);
int protocolDials = 0;
ISession session = Substitute.For<ISession>();
session.RemoteAddress.Returns(directAddress);
session.DialAsync<GossipsubProtocolV12>(Arg.Any<CancellationToken>()).Returns(_ =>
{
int attempt = Interlocked.Increment(ref protocolDials);
router.OutboundConnection(directAddress, PubsubRouter.GossipsubProtocolVersionV12,
attempt == 1 ? firstConnection.Task : secondConnection.Task, _ => { });
if (attempt > 1)
{
redialed.TrySetResult();
}
return Task.CompletedTask;
});
ILocalPeer localPeer = Substitute.For<ILocalPeer>();
localPeer.Identity.Returns(TestPeers.Identity(2));
localPeer.ListenAddresses.Returns(new ObservableCollection<Multiaddress>());
localPeer.DialAsync(Arg.Any<Multiaddress[]>(), Arg.Any<CancellationToken>()).Returns(session);
using CancellationTokenSource cancellation = new();
try
{
await router.StartAsync(localPeer, cancellation.Token);
await Task.Delay(200);
_ = localPeer.Received(1).DialAsync(Arg.Any<Multiaddress[]>(), Arg.Any<CancellationToken>());

firstConnection.SetResult();
await redialed.Task.WaitAsync(TimeSpan.FromSeconds(3));

await Task.Delay(200);
_ = localPeer.Received(2).DialAsync(Arg.Any<Multiaddress[]>(), Arg.Any<CancellationToken>());
_ = session.Received(2).DialAsync<GossipsubProtocolV12>(Arg.Any<CancellationToken>());
}
finally
{
cancellation.Cancel();
firstConnection.TrySetResult();
secondConnection.TrySetResult();
}
}

[Test]
public async Task Router_DoesNotRepeatInFlightDirectPeerDials()
{
Multiaddress directAddress = TestPeers.Multiaddr(1);
PeerStore peerStore = new();
peerStore.Discover([directAddress]);
await using PubsubRouter router = new(peerStore, new PubsubSettings
{
DirectPeers = [directAddress],
DirectConnectPeriod = 20,
ReconnectionPeriod = 60_000,
HeartbeatInterval = 60_000,
});
TaskCompletionSource<ISession> slowDial = new(TaskCreationOptions.RunContinuationsAsynchronously);
ILocalPeer localPeer = Substitute.For<ILocalPeer>();
localPeer.Identity.Returns(TestPeers.Identity(2));
localPeer.ListenAddresses.Returns(new ObservableCollection<Multiaddress>());
localPeer.DialAsync(Arg.Any<Multiaddress[]>(), Arg.Any<CancellationToken>()).Returns(slowDial.Task);

try
{
await router.StartAsync(localPeer);
await Task.Delay(200);

_ = localPeer.Received(1).DialAsync(Arg.Any<Multiaddress[]>(), Arg.Any<CancellationToken>());
}
finally
{
slowDial.TrySetCanceled();
}
}

[Test]
public void DirectPeers_RequirePeerIdsInTheirAddresses()
{
PubsubSettings settings = new()
{
DirectPeers = new[] { Multiaddress.Decode("/ip4/127.0.0.1/tcp/4001") },
};

ArgumentException exception = Assert.Throws<ArgumentException>(() => new PubsubRouter(new PeerStore(), settings))!;
Assert.That(exception.ParamName, Is.EqualTo(nameof(PubsubSettings.DirectPeers)));
}
}
12 changes: 12 additions & 0 deletions src/libp2p/Libp2p.Protocols.Pubsub.Tests/PubsubProtocolTests.cs
Original file line number Diff line number Diff line change
Expand Up @@ -39,6 +39,18 @@ public void Publish_WithNullMessage_ThrowsArgumentNullException()
Assert.Throws<ArgumentNullException>(() => router.Publish("test-topic", null!));
}

[Test]
public async Task Publish_WithoutSubscription_DoesNotThrow()
{
PubsubRouter router = new(new PeerStore());
using CancellationTokenSource cancellation = new();
await router.StartAsync(new LocalPeerStub(), cancellation.Token);

Assert.DoesNotThrow(() => router.Publish("test-topic", [1, 2, 3]));

cancellation.Cancel();
}

[Test]
public void Topic_OnMessage_IncludesReceivedFromPeerId()
{
Expand Down
13 changes: 13 additions & 0 deletions src/libp2p/Libp2p.Protocols.Pubsub/PubSubSettings.cs
Original file line number Diff line number Diff line change
Expand Up @@ -3,6 +3,7 @@

using Nethermind.Libp2p.Core;
using Nethermind.Libp2p.Protocols.Pubsub.Dto;
using Multiformats.Address;

namespace Nethermind.Libp2p.Protocols.Pubsub;

Expand All @@ -22,6 +23,18 @@ public class PubsubSettings

public int MaxConnections { get; set; }

/// <summary>
/// Peers with reciprocal explicit peering agreements. Each address must
/// contain a peer ID and is configured before the router starts.
/// </summary>
public Multiaddress[] DirectPeers { get; set; } = [];

/// <summary>
/// Interval in milliseconds for reconnecting disconnected direct peers. Gossipsub recommends
/// five minutes.
/// </summary>
public int DirectConnectPeriod { get; set; } = 5 * 60 * 1000;

public int HeartbeatInterval { get; set; } = 1_000; // Time between heartbeats 1 second
public int FanoutTtl { get; set; } = 60 * 1000; // Time-to-live for each topic's fanout state 60 seconds
public int mcache_len { get; set; } = 5; // Number of history windows in message cache 5
Expand Down
Loading
Loading