Skip to content

Commit fed6ff2

Browse files
committed
feat: implement performance-based peer selection for snap sync (#9722)
- Add peer performance metrics (latency, throughput) to EthPeer and EthPeerImmutableAttributes - Capture latency and throughput in AbstractPeerRequestTask - Add performance-based comparators (BY_LOW_LATENCY, BY_TRANSFER_SPEED) to EthPeers - Update AbstractRetryingSwitchingPeerTask to support custom comparators - Integrate performance-based selection into Snap Sync tasks and RequestDataStep - Add comprehensive unit tests for peer selection and metrics capturing Signed-off-by: Bruno Coste <bcoste@gmail.com>
1 parent ff1233a commit fed6ff2

8 files changed

Lines changed: 161 additions & 3 deletions

File tree

ethereum/eth/src/main/java/org/hyperledger/besu/ethereum/eth/manager/EthPeer.java

Lines changed: 30 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -94,6 +94,9 @@ protected boolean removeEldestEntry(final Map.Entry<Hash, Boolean> eldest) {
9494
private final AtomicInteger lastProtocolVersion = new AtomicInteger(0);
9595

9696
private volatile long lastRequestTimestamp = 0;
97+
private static final double ALPHA = 0.1;
98+
private volatile double averageLatencyMs = -1.0;
99+
private volatile double averageThroughputBytesPerSecond = -1.0;
97100

98101
private final Map<String, Map<Integer, RequestManager>> requestManagers;
99102

@@ -224,6 +227,33 @@ public void recordUsefulResponse() {
224227
reputation.recordUsefulResponse();
225228
}
226229

230+
public void recordLatency(final double latencyMs) {
231+
if (averageLatencyMs < 0) {
232+
averageLatencyMs = latencyMs;
233+
} else {
234+
averageLatencyMs = (ALPHA * latencyMs) + ((1.0 - ALPHA) * averageLatencyMs);
235+
}
236+
}
237+
238+
public void recordThroughput(final long bytes, final double durationMs) {
239+
if (durationMs <= 0) return;
240+
double throughput = (bytes * 1000.0) / durationMs;
241+
if (averageThroughputBytesPerSecond < 0) {
242+
averageThroughputBytesPerSecond = throughput;
243+
} else {
244+
averageThroughputBytesPerSecond =
245+
(ALPHA * throughput) + ((1.0 - ALPHA) * averageThroughputBytesPerSecond);
246+
}
247+
}
248+
249+
public double getAverageLatencyMs() {
250+
return averageLatencyMs;
251+
}
252+
253+
public double getAverageThroughputBytesPerSecond() {
254+
return averageThroughputBytesPerSecond;
255+
}
256+
227257
public void disconnect(final DisconnectReason reason) {
228258
connection.disconnect(reason);
229259
}

ethereum/eth/src/main/java/org/hyperledger/besu/ethereum/eth/manager/EthPeerImmutableAttributes.java

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -28,6 +28,8 @@ public record EthPeerImmutableAttributes(
2828
boolean isServingSnap,
2929
boolean hasAvailableRequestCapacity,
3030
boolean isInboundInitiated,
31+
double averageLatencyMs,
32+
double averageThroughputBytesPerSecond,
3133
EthPeer ethPeer) {
3234

3335
public static EthPeerImmutableAttributes from(final EthPeer peer) {
@@ -43,6 +45,8 @@ public static EthPeerImmutableAttributes from(final EthPeer peer) {
4345
peer.isServingSnap(),
4446
peer.hasAvailableRequestCapacity(),
4547
peer.getConnection().inboundInitiated(),
48+
peer.getAverageLatencyMs(),
49+
peer.getAverageThroughputBytesPerSecond(),
4650
peer);
4751
}
4852
}

ethereum/eth/src/main/java/org/hyperledger/besu/ethereum/eth/manager/EthPeers.java

Lines changed: 16 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -85,6 +85,16 @@ public class EthPeers implements PeerSelector {
8585
public static final Comparator<EthPeerImmutableAttributes> LEAST_TO_MOST_BUSY =
8686
Comparator.comparing(EthPeerImmutableAttributes::outstandingRequests)
8787
.thenComparing(EthPeerImmutableAttributes::lastRequestTimestamp);
88+
89+
public static final Comparator<EthPeerImmutableAttributes> BY_LOW_LATENCY =
90+
Comparator.comparing(
91+
(final EthPeerImmutableAttributes p) ->
92+
p.averageLatencyMs() < 0 ? Double.MAX_VALUE : p.averageLatencyMs())
93+
.reversed();
94+
95+
public static final Comparator<EthPeerImmutableAttributes> BY_HIGH_THROUGHPUT =
96+
Comparator.comparing(EthPeerImmutableAttributes::averageThroughputBytesPerSecond);
97+
8898
public static final int NODE_ID_LENGTH = 64;
8999
public static final int USEFULL_PEER_SCORE_THRESHOLD = 102;
90100

@@ -373,9 +383,14 @@ public Stream<EthPeerImmutableAttributes> streamAvailablePeers() {
373383
}
374384

375385
public Stream<EthPeerImmutableAttributes> streamBestPeers() {
386+
return streamBestPeers(getBestPeerComparator());
387+
}
388+
389+
public Stream<EthPeerImmutableAttributes> streamBestPeers(
390+
final Comparator<EthPeerImmutableAttributes> comparator) {
376391
return streamAvailablePeers()
377392
.filter(EthPeerImmutableAttributes::isFullyValidated)
378-
.sorted(getBestPeerComparator().reversed());
393+
.sorted(comparator.reversed());
379394
}
380395

381396
public Optional<EthPeer> bestPeer() {

ethereum/eth/src/main/java/org/hyperledger/besu/ethereum/eth/manager/task/AbstractPeerRequestTask.java

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -42,6 +42,7 @@ public abstract class AbstractPeerRequestTask<R> extends AbstractPeerTask<R> {
4242
private final String protocolName;
4343
private final int requestCode;
4444
private volatile PendingPeerRequest responseStream;
45+
private long startTimeNanos;
4546

4647
protected AbstractPeerRequestTask(
4748
final EthContext ethContext,
@@ -64,6 +65,7 @@ protected final void executeTask() {
6465
responseStream = sendRequest();
6566
responseStream.then(
6667
stream -> {
68+
startTimeNanos = System.nanoTime();
6769
// Start the timeout now that the request has actually been sent
6870
ethContext.getScheduler().failAfterTimeout(promise, timeout);
6971

@@ -108,6 +110,9 @@ private void handleMessage(
108110
final Optional<R> result = processResponse(streamClosed, message, peer);
109111
result.ifPresent(
110112
r -> {
113+
final double durationMs = (System.nanoTime() - startTimeNanos) / 1_000_000.0;
114+
peer.recordLatency(durationMs);
115+
peer.recordThroughput(message.getSize(), durationMs);
111116
promise.complete(r);
112117
peer.recordUsefulResponse();
113118
});

ethereum/eth/src/main/java/org/hyperledger/besu/ethereum/eth/manager/task/AbstractRetryingSwitchingPeerTask.java

Lines changed: 7 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -39,6 +39,7 @@ public abstract class AbstractRetryingSwitchingPeerTask<T> extends AbstractRetry
3939

4040
private final Set<EthPeer> triedPeers = new HashSet<>();
4141
private final Set<EthPeer> failedPeers = new HashSet<>();
42+
private Optional<Comparator<EthPeerImmutableAttributes>> peerComparator = Optional.empty();
4243

4344
protected AbstractRetryingSwitchingPeerTask(
4445
final EthContext ethContext,
@@ -100,6 +101,10 @@ protected CompletableFuture<T> executePeerTask(final Optional<EthPeer> assignedP
100101
});
101102
}
102103

104+
public void setPeerComparator(final Comparator<EthPeerImmutableAttributes> peerComparator) {
105+
this.peerComparator = Optional.of(peerComparator);
106+
}
107+
103108
@Override
104109
protected void handleTaskError(final Throwable error) {
105110
if (isPeerFailure(error)) {
@@ -129,7 +134,8 @@ private Optional<EthPeer> selectNextPeer() {
129134
protected Optional<EthPeer> nextPeerToTry() {
130135
return getEthContext()
131136
.getEthPeers()
132-
.streamBestPeers()
137+
.streamBestPeers(
138+
peerComparator.orElse(getEthContext().getEthPeers().getBestPeerComparator()))
133139
.filter((peer) -> isSuitablePeer(peer) && !triedPeers.contains(peer.ethPeer()))
134140
.map(EthPeerImmutableAttributes::ethPeer)
135141
.findFirst();

ethereum/eth/src/main/java/org/hyperledger/besu/ethereum/eth/sync/snapsync/RequestDataStep.java

Lines changed: 9 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -17,6 +17,7 @@
1717
import org.hyperledger.besu.datatypes.Hash;
1818
import org.hyperledger.besu.ethereum.core.BlockHeader;
1919
import org.hyperledger.besu.ethereum.eth.manager.EthContext;
20+
import org.hyperledger.besu.ethereum.eth.manager.EthPeers;
2021
import org.hyperledger.besu.ethereum.eth.manager.snap.RetryingGetAccountRangeFromPeerTask;
2122
import org.hyperledger.besu.ethereum.eth.manager.snap.RetryingGetBytecodeFromPeerTask;
2223
import org.hyperledger.besu.ethereum.eth.manager.snap.RetryingGetStorageRangeFromPeerTask;
@@ -92,9 +93,11 @@ public CompletableFuture<Task<SnapDataRequest>> requestAccount(
9293
ethContext,
9394
accountDataRequest.getStartKeyHash(),
9495
accountDataRequest.getEndKeyHash(),
95-
blockHeader,
96+
blockHeader,
9697
metricsSystem);
98+
((RetryingGetAccountRangeFromPeerTask) getAccountTask).setPeerComparator(EthPeers.BY_HIGH_THROUGHPUT);
9799
downloadState.addOutstandingTask(getAccountTask);
100+
98101
return getAccountTask
99102
.run()
100103
.orTimeout(10, TimeUnit.SECONDS)
@@ -142,6 +145,8 @@ public CompletableFuture<List<Task<SnapDataRequest>>> requestStorage(
142145
final EthTask<StorageRangeMessage.SlotRangeData> getStorageRangeTask =
143146
RetryingGetStorageRangeFromPeerTask.forStorageRange(
144147
ethContext, accountHashesAsBytes32, minRange, maxRange, blockHeader, metricsSystem);
148+
((RetryingGetStorageRangeFromPeerTask) getStorageRangeTask)
149+
.setPeerComparator(EthPeers.BY_HIGH_THROUGHPUT);
145150
downloadState.addOutstandingTask(getStorageRangeTask);
146151
return getStorageRangeTask
147152
.run()
@@ -203,6 +208,7 @@ public CompletableFuture<List<Task<SnapDataRequest>>> requestCode(
203208
final EthTask<Map<Bytes32, Bytes>> getByteCodeTask =
204209
RetryingGetBytecodeFromPeerTask.forByteCode(
205210
ethContext, codeHashes, blockHeader, metricsSystem);
211+
((RetryingGetBytecodeFromPeerTask) getByteCodeTask).setPeerComparator(EthPeers.BY_LOW_LATENCY);
206212
downloadState.addOutstandingTask(getByteCodeTask);
207213
return getByteCodeTask
208214
.run()
@@ -249,6 +255,8 @@ public CompletableFuture<List<Task<SnapDataRequest>>> requestTrieNodeByPath(
249255
final EthTask<Map<Bytes, Bytes>> getTrieNodeFromPeerTask =
250256
RetryingGetTrieNodeFromPeerTask.forTrieNodes(
251257
ethContext, message, blockHeader, metricsSystem);
258+
((RetryingGetTrieNodeFromPeerTask) getTrieNodeFromPeerTask)
259+
.setPeerComparator(EthPeers.BY_LOW_LATENCY);
252260
downloadState.addOutstandingTask(getTrieNodeFromPeerTask);
253261
return getTrieNodeFromPeerTask
254262
.run()

ethereum/eth/src/test/java/org/hyperledger/besu/ethereum/eth/manager/EthPeerTest.java

Lines changed: 33 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -392,6 +392,39 @@ public void recordUsefulResponse() {
392392
assertThat(peer.getReputation().compareTo(peer2.getReputation())).isGreaterThan(0);
393393
}
394394

395+
@Test
396+
public void shouldRecordLatencyAndThroughput() {
397+
final EthPeer peer = createPeer();
398+
assertThat(peer.getAverageLatencyMs()).isEqualTo(-1.0);
399+
assertThat(peer.getAverageThroughputBytesPerSecond()).isEqualTo(-1.0);
400+
401+
// First record sets the value
402+
peer.recordLatency(100.0);
403+
assertThat(peer.getAverageLatencyMs()).isEqualTo(100.0);
404+
405+
// Second record updates EMA: 0.1 * 200 + 0.9 * 100 = 110
406+
peer.recordLatency(200.0);
407+
assertThat(peer.getAverageLatencyMs()).isEqualTo(110.0);
408+
409+
// First record sets throughput
410+
// 1000 bytes in 100ms = 10,000 bytes/s
411+
peer.recordThroughput(1000, 100.0);
412+
assertThat(peer.getAverageThroughputBytesPerSecond()).isEqualTo(10000.0);
413+
414+
// Second record updates EMA
415+
// 2000 bytes in 100ms = 20,000 bytes/s
416+
// EMA: 0.1 * 20000 + 0.9 * 10000 = 11000
417+
peer.recordThroughput(2000, 100.0);
418+
assertThat(peer.getAverageThroughputBytesPerSecond()).isEqualTo(11000.0);
419+
}
420+
421+
@Test
422+
public void shouldHandleZeroDurationInThroughput() {
423+
final EthPeer peer = createPeer();
424+
peer.recordThroughput(1000, 0.0);
425+
assertThat(peer.getAverageThroughputBytesPerSecond()).isEqualTo(-1.0);
426+
}
427+
395428
private void messageStream(
396429
final ResponseStreamSupplier getStream,
397430
final MessageData targetMessage,

ethereum/eth/src/test/java/org/hyperledger/besu/ethereum/eth/manager/EthPeersTest.java

Lines changed: 57 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -125,6 +125,63 @@ public void comparesPeersWithTdAndNoHeight() {
125125
.isEmpty();
126126
}
127127

128+
@Test
129+
public void comparesPeersByLatency() {
130+
final EthPeer peerA =
131+
EthProtocolManagerTestUtil.createPeer(ethProtocolManager, 1000).getEthPeer();
132+
final EthPeer peerB =
133+
EthProtocolManagerTestUtil.createPeer(ethProtocolManager, 1000).getEthPeer();
134+
135+
peerA.recordLatency(100.0);
136+
peerB.recordLatency(200.0);
137+
138+
final EthPeerImmutableAttributes attrA = EthPeerImmutableAttributes.from(peerA);
139+
final EthPeerImmutableAttributes attrB = EthPeerImmutableAttributes.from(peerB);
140+
141+
// Smaller latency is better, so attrA > attrB in LOW_LATENCY comparator
142+
assertThat(EthPeers.BY_LOW_LATENCY.compare(attrA, attrB)).isGreaterThan(0);
143+
assertThat(EthPeers.BY_LOW_LATENCY.compare(attrB, attrA)).isLessThan(0);
144+
145+
// Peer with unknown latency should be worse than any known latency
146+
final EthPeer peerC =
147+
EthProtocolManagerTestUtil.createPeer(ethProtocolManager, 1000).getEthPeer();
148+
final EthPeerImmutableAttributes attrC = EthPeerImmutableAttributes.from(peerC);
149+
assertThat(EthPeers.BY_LOW_LATENCY.compare(attrA, attrC)).isGreaterThan(0);
150+
}
151+
152+
@Test
153+
public void comparesPeersByThroughput() {
154+
final EthPeer peerA =
155+
EthProtocolManagerTestUtil.createPeer(ethProtocolManager, 1000).getEthPeer();
156+
final EthPeer peerB =
157+
EthProtocolManagerTestUtil.createPeer(ethProtocolManager, 1000).getEthPeer();
158+
159+
peerA.recordThroughput(2000, 1000.0); // 2000 bytes/s
160+
peerB.recordThroughput(1000, 1000.0); // 1000 bytes/s
161+
162+
final EthPeerImmutableAttributes attrA = EthPeerImmutableAttributes.from(peerA);
163+
final EthPeerImmutableAttributes attrB = EthPeerImmutableAttributes.from(peerB);
164+
165+
// Larger throughput is better
166+
assertThat(EthPeers.BY_HIGH_THROUGHPUT.compare(attrA, attrB)).isGreaterThan(0);
167+
assertThat(EthPeers.BY_HIGH_THROUGHPUT.compare(attrB, attrA)).isLessThan(0);
168+
}
169+
170+
@Test
171+
public void shouldStreamBestPeersWithCustomComparator() {
172+
final EthPeer peerA =
173+
EthProtocolManagerTestUtil.createPeer(ethProtocolManager, 1000).getEthPeer();
174+
final EthPeer peerB =
175+
EthProtocolManagerTestUtil.createPeer(ethProtocolManager, 1000).getEthPeer();
176+
177+
peerA.recordLatency(200.0);
178+
peerB.recordLatency(100.0);
179+
180+
// BY_LOW_LATENCY reversed should put peerB first
181+
assertThat(ethPeers.streamBestPeers(EthPeers.BY_LOW_LATENCY).findFirst())
182+
.contains(EthPeerImmutableAttributes.from(peerB));
183+
}
184+
128185
@Test
129186
public void shouldExecutePeerRequestImmediatelyWhenPeerIsAvailable() throws Exception {
130187
final RespondingEthPeer peer = EthProtocolManagerTestUtil.createPeer(ethProtocolManager, 1000);

0 commit comments

Comments
 (0)