Kill the connection if a peer sends multiple ping requests in parallel (#3172)
What changed, and why it matters
This update makes Eclair disconnect a Lightning peer if it sends more than one ping message before getting a reply. The goal is to stop a 'ping flood' attack, where a malicious peer could overwhelm the node with many ping requests and possibly cause it to run out of memory or become unresponsive. The change is defensive and includes a test showing the connection is closed when two pings arrive back-to-back.
Treat as a security hardening fix and include in release notes. Nodes should upgrade to avoid being a target of ping-flood resource exhaustion. No immediate incident response is required unless operators observe repeated disconnections from specific peers.
Security signals we found
Adds a counter to detect and terminate connections on parallel ping messages
Commit message describes the peer behavior as malicious and closes the connection
Includes a regression test named 'reject ping flooding'
Change is in the encrypted transport layer where all peer messages pass through
Evidence from the diff
TransportHandler.scala now tracks pendingPings, incremented on each received Ping and decremented on each sent Pong. If pendingPings exceeds 1 after decrypting a batch, the FSM stops and the TCP connection is dropped. The patch refactors outgoing message encoding into encodeAndSendToPeer so Pong dispatch can decrement the counter. A new unit test confirms that a second Ping before a Pong terminates the actor pipe. The commit message explicitly calls this a defense against a malicious peer.
Changed components
eclair-core/src/main/scala/fr/acinq/eclair/crypto/TransportHandler.scalaeclair-core/src/test/scala/fr/acinq/eclair/crypto/TransportHandlerSpec.scalaInspect captured patch +82 / −29
diff --git a/eclair-core/src/main/scala/fr/acinq/eclair/crypto/TransportHandler.scala b/eclair-core/src/main/scala/fr/acinq/eclair/crypto/TransportHandler.scala
index c05106d..74d9b2f 100644
--- a/eclair-core/src/main/scala/fr/acinq/eclair/crypto/TransportHandler.scala
+++ b/eclair-core/src/main/scala/fr/acinq/eclair/crypto/TransportHandler.scala
@@ -27,7 +27,7 @@ import fr.acinq.eclair.Logs.LogCategory
import fr.acinq.eclair.crypto.ChaCha20Poly1305.ChaCha20Poly1305Error
import fr.acinq.eclair.crypto.Noise._
import fr.acinq.eclair.remote.EclairInternalsSerializer.RemoteTypes
-import fr.acinq.eclair.wire.protocol.{AnnouncementSignatures, LightningMessage, RoutingMessage}
+import fr.acinq.eclair.wire.protocol._
import fr.acinq.eclair.{Diagnostics, FSMDiagnosticActorLogging, Logs, getSimpleClassName}
import scodec.bits.ByteVector
import scodec.{Attempt, Codec, DecodeResult}
@@ -92,6 +92,9 @@ class TransportHandler(keyPair: KeyPair, rs: Option[ByteVector], connection: Act
makeReader(keyPair)
}
+ /** We keep track of pending pings to defend against ping flooding. */
+ private var pendingPings = 0
+
private def decodeAndSendToListener(listener: ActorRef, plaintextMessages: Seq[ByteVector]): Map[LightningMessage, Int] = {
log.debug("decoding {} plaintext messages", plaintextMessages.size)
var m = Map.empty[LightningMessage, Int]
@@ -99,6 +102,13 @@ class TransportHandler(keyPair: KeyPair, rs: Option[ByteVector], connection: Act
case Attempt.Successful(DecodeResult(message, _)) =>
logMessage(message, "IN")
Monitoring.Metrics.MessageSize.withTag(Monitoring.Tags.MessageDirection, Monitoring.Tags.MessageDirections.IN).record(plaintext.size)
+ if (message.isInstanceOf[Ping]) {
+ pendingPings += 1
+ if (pendingPings > 1) {
+ // We will kill the connection anyway, no need to process remaining messages
+ return m
+ }
+ }
listener ! message
m += (message -> (m.getOrElse(message, 0) + 1))
case Attempt.Failure(err) =>
@@ -108,6 +118,18 @@ class TransportHandler(keyPair: KeyPair, rs: Option[ByteVector], connection: Act
m
}
+ private def encodeAndSendToPeer(encryptor: Encryptor, t: LightningMessage): Encryptor = {
+ if (t.isInstanceOf[Pong]) {
+ pendingPings -= 1
+ }
+ logMessage(t, "OUT")
+ val blob = codec.encode(t).require.toByteVector
+ Monitoring.Metrics.MessageSize.withTag(Monitoring.Tags.MessageDirection, Monitoring.Tags.MessageDirections.OUT).record(blob.size)
+ val (enc1, ciphertext) = encryptor.encrypt(blob)
+ connection ! Tcp.Write(buf(ciphertext), WriteAck)
+ enc1
+ }
+
startWith(Handshake, HandshakeData(reader))
when(Handshake) {
@@ -161,11 +183,16 @@ class TransportHandler(keyPair: KeyPair, rs: Option[ByteVector], connection: Act
context.watch(listener)
val (dec1, plaintextMessages) = dec.decrypt()
val unackedReceived1 = decodeAndSendToListener(listener, plaintextMessages)
- if (unackedReceived1.isEmpty) {
- log.debug("no decoded messages, resuming reading")
- connection ! Tcp.ResumeReading
+ if (pendingPings > 1) {
+ log.warning("ping flood detected (pendingPings={}): closing connection", pendingPings)
+ stop(FSM.Normal)
+ } else {
+ if (unackedReceived1.isEmpty) {
+ log.debug("no decoded messages, resuming reading")
+ connection ! Tcp.ResumeReading
+ }
+ goto(Normal) using NormalData(d.encryptor, dec1, listener, sendBuffer = SendBuffer(Queue.empty[LightningMessage], Queue.empty[LightningMessage]), unackedReceived = unackedReceived1, unackedSent = None)
}
- goto(Normal) using NormalData(d.encryptor, dec1, listener, sendBuffer = SendBuffer(Queue.empty[LightningMessage], Queue.empty[LightningMessage]), unackedReceived = unackedReceived1, unackedSent = None)
}
}
@@ -175,11 +202,16 @@ class TransportHandler(keyPair: KeyPair, rs: Option[ByteVector], connection: Act
log.debug("received chunk of size={}", data.size)
val (dec1, plaintextMessages) = d.decryptor.copy(buffer = d.decryptor.buffer ++ data).decrypt()
val unackedReceived1 = decodeAndSendToListener(d.listener, plaintextMessages)
- if (unackedReceived1.isEmpty) {
- log.debug("no decoded messages, resuming reading")
- connection ! Tcp.ResumeReading
+ if (pendingPings > 1) {
+ log.warning("ping flood detected (pendingPings={}): closing connection", pendingPings)
+ stop(FSM.Normal)
+ } else {
+ if (unackedReceived1.isEmpty) {
+ log.debug("no decoded messages, resuming reading")
+ connection ! Tcp.ResumeReading
+ }
+ stay() using d.copy(decryptor = dec1, unackedReceived = unackedReceived1)
}
- stay() using d.copy(decryptor = dec1, unackedReceived = unackedReceived1)
case Event(ReadAck(msg: LightningMessage), d: NormalData) =>
// how many occurrences of this message are still unacked?
@@ -209,32 +241,19 @@ class TransportHandler(keyPair: KeyPair, rs: Option[ByteVector], connection: Act
}
stay() using d.copy(sendBuffer = sendBuffer1)
} else {
- logMessage(t, "OUT")
- val blob = codec.encode(t).require.toByteVector
- Monitoring.Metrics.MessageSize.withTag(Monitoring.Tags.MessageDirection, Monitoring.Tags.MessageDirections.OUT).record(blob.size)
- val (enc1, ciphertext) = d.encryptor.encrypt(blob)
- connection ! Tcp.Write(buf(ciphertext), WriteAck)
+ val enc1 = encodeAndSendToPeer(d.encryptor, t)
stay() using d.copy(encryptor = enc1, unackedSent = Some(t))
}
case Event(WriteAck, d: NormalData) =>
- def send(t: LightningMessage) = {
- logMessage(t, "OUT")
- val blob = codec.encode(t).require.toByteVector
- Monitoring.Metrics.MessageSize.withTag(Monitoring.Tags.MessageDirection, Monitoring.Tags.MessageDirections.OUT).record(blob.size)
- val (enc1, ciphertext) = d.encryptor.encrypt(blob)
- connection ! Tcp.Write(buf(ciphertext), WriteAck)
- enc1
- }
-
d.sendBuffer.normalPriority.dequeueOption match {
case Some((t, normalPriority1)) =>
- val enc1 = send(t)
+ val enc1 = encodeAndSendToPeer(d.encryptor, t)
stay() using d.copy(encryptor = enc1, sendBuffer = d.sendBuffer.copy(normalPriority = normalPriority1), unackedSent = Some(t))
case None =>
d.sendBuffer.lowPriority.dequeueOption match {
case Some((t, lowPriority1)) =>
- val enc1 = send(t)
+ val enc1 = encodeAndSendToPeer(d.encryptor, t)
stay() using d.copy(encryptor = enc1, sendBuffer = d.sendBuffer.copy(lowPriority = lowPriority1), unackedSent = Some(t))
case None =>
stay() using d.copy(unackedSent = None)
diff --git a/eclair-core/src/test/scala/fr/acinq/eclair/crypto/TransportHandlerSpec.scala b/eclair-core/src/test/scala/fr/acinq/eclair/crypto/TransportHandlerSpec.scala
index 725233c..e0d8b72 100644
--- a/eclair-core/src/test/scala/fr/acinq/eclair/crypto/TransportHandlerSpec.scala
+++ b/eclair-core/src/test/scala/fr/acinq/eclair/crypto/TransportHandlerSpec.scala
@@ -19,11 +19,11 @@ package fr.acinq.eclair.crypto
import akka.actor.{Actor, ActorLogging, ActorRef, OneForOneStrategy, Props, Stash, SupervisorStrategy, Terminated}
import akka.io.Tcp
import akka.testkit.{TestActorRef, TestFSMRef, TestProbe}
-import fr.acinq.eclair.TestKitBaseClass
+import fr.acinq.eclair.{TestKitBaseClass, randomBytes32}
import fr.acinq.eclair.crypto.Noise.{Chacha20Poly1305CipherFunctions, CipherState}
import fr.acinq.eclair.crypto.TransportHandler.{Encryptor, ExtendedCipherState, Listener}
-import fr.acinq.eclair.wire.protocol.LightningMessageCodecs.{lightningMessageCodec, pingCodec}
-import fr.acinq.eclair.wire.protocol.{LightningMessage, Ping, Pong}
+import fr.acinq.eclair.wire.protocol.LightningMessageCodecs.{lightningMessageCodec, pingCodec, warningCodec}
+import fr.acinq.eclair.wire.protocol.{LightningMessage, Ping, Pong, Warning}
import org.scalatest.BeforeAndAfterAll
import org.scalatest.funsuite.AnyFunSuiteLike
import scodec.Codec
@@ -77,6 +77,7 @@ class TransportHandlerSpec extends TestKitBaseClass with AnyFunSuiteLike with Be
test("handle unknown messages") {
val incompleteCodec: Codec[LightningMessage] = discriminated[LightningMessage].by(uint16)
+ .typecase(1, warningCodec)
.typecase(18, pingCodec)
val pipe = system.actorOf(Props[MyPipePull]())
@@ -103,7 +104,7 @@ class TransportHandlerSpec extends TestKitBaseClass with AnyFunSuiteLike with Be
responder ! Pong(hex"deadbeef")
probe1.expectNoMessage(2 seconds) // unknown message
- val msg2 = Ping(42, hex"beefdead")
+ val msg2 = Warning(randomBytes32(), hex"beefdead")
responder ! msg2
probe1.expectMsg(msg2)
probe1.reply(TransportHandler.ReadAck(msg2))
@@ -199,6 +200,39 @@ class TransportHandlerSpec extends TestKitBaseClass with AnyFunSuiteLike with Be
assert(ciphertexts(1000) == hex"4a2f3cc3b5e78ddb83dcb426d9863d9d9a723b0337c89dd0b005d89f8d3c05c52b76b29b740f09")
assert(ciphertexts(1001) == hex"2ecd8c8a5629d0d02ab457a0fdd0f7b90a192cd46be5ecb6ca570bfc5e268338b1a16cf4ef2d36")
}
+
+ test("reject ping flooding") {
+ val pipe = system.actorOf(Props[MyPipe]())
+ val probe1 = TestProbe()
+ val probe2 = TestProbe()
+ val initiator = TestFSMRef(new TransportHandler(Initiator.s, Some(Responder.s.pub), pipe, lightningMessageCodec))
+ val responder = TestFSMRef(new TransportHandler(Responder.s, None, pipe, lightningMessageCodec))
+ pipe ! (initiator, responder)
+
+ awaitCond(initiator.stateName == TransportHandler.WaitingForListener)
+ awaitCond(responder.stateName == TransportHandler.WaitingForListener)
+
+ initiator ! Listener(probe1.ref)
+ responder ! Listener(probe2.ref)
+
+ awaitCond(initiator.stateName == TransportHandler.Normal)
+ awaitCond(responder.stateName == TransportHandler.Normal)
+
+ initiator.tell(Ping(1105, ByteVector("hello 1".getBytes)), probe1.ref)
+ probe2.expectMsg(Ping(1105, ByteVector("hello 1".getBytes)))
+
+ initiator.tell(Ping(1105, ByteVector("hello 2".getBytes)), probe1.ref)
+ probe2.expectNoMessage()
+
+ probe1.watch(initiator)
+ probe1.expectTerminated(initiator)
+
+ probe1.watch(responder)
+ probe1.expectTerminated(responder)
+
+ probe1.watch(pipe)
+ probe1.expectTerminated(pipe)
+ }
}
object TransportHandlerSpec {
Why this scored 64/100
Community notes
Notes can correct, qualify, or add evidence to the AI analysis. Every note shown here has been validated by a human moderator.
The AI analysis stands alone for now. Submit a note if you can add evidence or important context.