What changed, and why it matters
This commit changes how Eclair (a Bitcoin Lightning node) loads stored peer data when the node restarts. Previously, every peer's stored data was read from the database immediately at startup. Now it is loaded only when the peer first reconnects. The stated goal is to avoid a 'herd effect'—a large burst of database reads all at once after a restart. There is no direct security vulnerability in the diff; it is a performance and reliability improvement. The change does not alter who can request data or what data is stored.
No security action required. Treat as a normal performance/reliability refactor. Reviewers may want to verify that the lazy load does not skip sending PeerStorageRetrieval in any edge case where it was previously sent, and that PendingWrite data is still flushed on peer stop.
Security signals we found
No authentication or authorization changes
No new network inputs or message handlers
No change to cryptographic operations
Refactoring only: state-machine change with equivalent observable behavior
Potential minor reliability improvement: reduced startup DB load
No explicit security fix or vulnerability disclosure in commit message
Evidence from the diff
The patch refactors PeerStorage from a simple case class (Option[ByteVector], written: Boolean) into a sealed trait with states Uninitialized, Empty, PendingWrite, and WrittenToDb. On Peer initialization, storage is now Uninitialized instead of being eagerly fetched from nodeParams.db.peers.getStorage(remoteNodeId). The DB read is deferred until the first ConnectionReady event, at which point the stored blob is sent to the peer via PeerStorageRetrieval and the state becomes WrittenToDb. Write scheduling is consolidated behind maybeScheduleWritePeerStorageTick and writePeerStorage helpers. Tests confirm lazy loading and that storage is not re-read on reconnection.
Changed components
eclair-core/src/main/scala/fr/acinq/eclair/io/Peer.scalaPeerStorage state machinePeer reconnection / ConnectionReady flowPeer storage write schedulingInspect captured patch +91 / −30
diff --git a/eclair-core/src/main/scala/fr/acinq/eclair/io/Peer.scala b/eclair-core/src/main/scala/fr/acinq/eclair/io/Peer.scala
index 7529e07..ccf59b1 100644
--- a/eclair-core/src/main/scala/fr/acinq/eclair/io/Peer.scala
+++ b/eclair-core/src/main/scala/fr/acinq/eclair/io/Peer.scala
@@ -89,14 +89,9 @@ class Peer(val nodeParams: NodeParams,
FinalChannelId(state.channelId) -> channel
}.toMap
context.system.eventStream.publish(PeerCreated(self, remoteNodeId))
- val peerStorageData = if (nodeParams.features.hasFeature(Features.ProvideStorage)) {
- nodeParams.db.peers.getStorage(remoteNodeId)
- } else {
- None
- }
// When we restart, we will attempt to reconnect right away, but then we'll wait.
// We don't fetch our peer's features from the DB: if the connection succeeds, we will get them from their init message, which saves a DB call.
- goto(DISCONNECTED) using DisconnectedData(channels, activeChannels = Set.empty, PeerStorage(peerStorageData, written = true), remoteFeatures_opt = None)
+ goto(DISCONNECTED) using DisconnectedData(channels, activeChannels = Set.empty, peerStorage = PeerStorage.Uninitialized, remoteFeatures_opt = None)
}
when(DISCONNECTED) {
@@ -150,13 +145,14 @@ class Peer(val nodeParams: NodeParams,
case Event(_: LightningMessage, _) => stay() // we probably just got disconnected and that's the last messages we received
- case Event(WritePeerStorage, d: DisconnectedData) =>
- d.peerStorage.data.foreach(nodeParams.db.peers.updateStorage(remoteNodeId, _))
- stay() using d.copy(peerStorage = d.peerStorage.copy(written = true))
+ case Event(TickWritePeerStorage, d: DisconnectedData) =>
+ val peerStorage1 = writePeerStorage(d.peerStorage)
+ stay() using d.copy(peerStorage = peerStorage1)
case Event(e: ChannelReadyForPayments, d: DisconnectedData) =>
- if (!d.peerStorage.written && !isTimerActive(WritePeerStorageTimerKey)) {
- startSingleTimer(WritePeerStorageTimerKey, WritePeerStorage, nodeParams.peerStorageConfig.getWriteDelay(remoteNodeId, d.remoteFeatures_opt.map(_.features)))
+ d.peerStorage match {
+ case _: PeerStorage.PendingWrite => maybeScheduleWritePeerStorageTick(d.remoteFeatures_opt.map(_.features))
+ case _ => ()
}
val remoteFeatures_opt = d.remoteFeatures_opt match {
case Some(remoteFeatures) if !remoteFeatures.written =>
@@ -503,8 +499,9 @@ class Peer(val nodeParams: NodeParams,
}
}
}
- if (!d.peerStorage.written && !isTimerActive(WritePeerStorageTimerKey)) {
- startSingleTimer(WritePeerStorageTimerKey, WritePeerStorage, nodeParams.peerStorageConfig.getWriteDelay(remoteNodeId, Some(d.remoteFeatures)))
+ d.peerStorage match {
+ case _: PeerStorage.PendingWrite => maybeScheduleWritePeerStorageTick(Some(d.remoteFeatures))
+ case _ => ()
}
if (!d.remoteFeaturesWritten) {
// We have a channel, so we can write to the DB without any DoS risk.
@@ -624,18 +621,18 @@ class Peer(val nodeParams: NodeParams,
// writing to the DB and may never store our peer's backup.
if (d.activeChannels.isEmpty) {
log.debug("received peer storage from peer with no active channel")
- } else if (!isTimerActive(WritePeerStorageTimerKey)) {
- startSingleTimer(WritePeerStorageTimerKey, WritePeerStorage, nodeParams.peerStorageConfig.getWriteDelay(remoteNodeId, Some(d.remoteFeatures)))
+ } else {
+ maybeScheduleWritePeerStorageTick(Some(d.remoteFeatures))
}
- stay() using d.copy(peerStorage = PeerStorage(Some(store.blob), written = false))
+ stay() using d.copy(peerStorage = PeerStorage.PendingWrite(store.blob))
} else {
log.debug("ignoring peer storage, feature disabled")
stay()
}
- case Event(WritePeerStorage, d: ConnectedData) =>
- d.peerStorage.data.foreach(nodeParams.db.peers.updateStorage(remoteNodeId, _))
- stay() using d.copy(peerStorage = d.peerStorage.copy(written = true))
+ case Event(TickWritePeerStorage, d: ConnectedData) =>
+ val peerStorage1 = writePeerStorage(d.peerStorage)
+ stay() using d.copy(peerStorage = peerStorage1)
case Event(unhandledMsg: LightningMessage, _) =>
log.warning("ignoring message {}", unhandledMsg)
@@ -895,7 +892,23 @@ class Peer(val nodeParams: NodeParams,
}
// If we have some data stored from our peer, we send it to them before doing anything else.
- peerStorage.data.foreach(connectionReady.peerConnection ! PeerStorageRetrieval(_))
+ val peerStorage1 = peerStorage match {
+ case PeerStorage.Uninitialized =>
+ val peerStorageData_opt = if (nodeParams.features.hasFeature(Features.ProvideStorage)) {
+ nodeParams.db.peers.getStorage(remoteNodeId)
+ } else {
+ None
+ }
+ peerStorageData_opt.map(PeerStorage.WrittenToDb(_)).getOrElse(PeerStorage.Empty)
+ case other => other
+ }
+ val peerStorageData_opt = peerStorage1 match {
+ case PeerStorage.Uninitialized => None // impossible!
+ case PeerStorage.Empty => None
+ case PeerStorage.PendingWrite(data) => Some(data)
+ case PeerStorage.WrittenToDb(data) => Some(data)
+ }
+ peerStorageData_opt.foreach(connectionReady.peerConnection ! PeerStorageRetrieval(_))
// let's bring existing/requested channels online
channels.values.toSet[ActorRef].foreach(_ ! INPUT_RECONNECTED(connectionReady.peerConnection, connectionReady.localInit, connectionReady.remoteInit)) // we deduplicate with toSet because there might be two entries per channel (tmp id and final id)
@@ -914,7 +927,7 @@ class Peer(val nodeParams: NodeParams,
connectionReady.peerConnection ! CurrentFeeCredit(nodeParams.chainHash, feeCredit.getOrElse(0 msat))
}
- goto(CONNECTED) using ConnectedData(connectionReady.address, connectionReady.peerConnection, connectionReady.localInit, connectionReady.remoteInit, channels, activeChannels, feerates, None, peerStorage, remoteFeaturesWritten = connectionReady.outgoing)
+ goto(CONNECTED) using ConnectedData(connectionReady.address, connectionReady.peerConnection, connectionReady.localInit, connectionReady.remoteInit, channels, activeChannels, feerates, None, peerStorage1, remoteFeaturesWritten = connectionReady.outgoing)
}
/**
@@ -988,8 +1001,9 @@ class Peer(val nodeParams: NodeParams,
private val openChannelInterceptor = context.spawnAnonymous(Behaviors.supervise(OpenChannelInterceptor(context.self.toTyped, nodeParams, remoteNodeId, wallet, pendingChannelsRateLimiter)).onFailure(typed.SupervisorStrategy.resume))
private def stopPeer(peerStorage: PeerStorage): State = {
- if (!peerStorage.written) {
- peerStorage.data.foreach(nodeParams.db.peers.updateStorage(remoteNodeId, _))
+ peerStorage match {
+ case PeerStorage.PendingWrite(data) => nodeParams.db.peers.updateStorage(remoteNodeId, data)
+ case _ => ()
}
log.info("removing peer from db")
cancelUnsignedOnTheFlyFunding()
@@ -1025,13 +1039,29 @@ class Peer(val nodeParams: NodeParams,
Logs.mdc(LogCategory(currentMessage), Some(remoteNodeId), Logs.channelId(currentMessage), nodeAlias_opt = Some(nodeParams.alias))
}
- private val WritePeerStorageTimerKey = "peer-storage-write"
+ private def maybeScheduleWritePeerStorageTick(remoteFeatures_opt: Option[Features[InitFeature]]): Unit = {
+ if (!isTimerActive(WRITE_PEER_STORAGE_TIMER_KEY)) {
+ startSingleTimer(WRITE_PEER_STORAGE_TIMER_KEY, TickWritePeerStorage, nodeParams.peerStorageConfig.getWriteDelay(remoteNodeId, remoteFeatures_opt))
+ }
+ }
+
+ private def writePeerStorage(peerStorage: PeerStorage): PeerStorage = {
+ peerStorage match {
+ case PeerStorage.PendingWrite(data) =>
+ nodeParams.db.peers.updateStorage(remoteNodeId, data)
+ PeerStorage.WrittenToDb(data)
+ case other => other
+ }
+ }
+
}
object Peer {
val CHANNELID_ZERO: ByteVector32 = ByteVector32.Zeroes
+ private val WRITE_PEER_STORAGE_TIMER_KEY = "tick-write-peer-storage"
+
trait ChannelFactory {
def spawn(context: ActorContext, remoteNodeId: PublicKey, channelKeys: ChannelKeys): ActorRef
}
@@ -1053,7 +1083,13 @@ object Peer {
case class TemporaryChannelId(id: ByteVector32) extends ChannelId
case class FinalChannelId(id: ByteVector32) extends ChannelId
- case class PeerStorage(data: Option[ByteVector], written: Boolean)
+ sealed trait PeerStorage
+ object PeerStorage {
+ case object Uninitialized extends PeerStorage
+ case object Empty extends PeerStorage
+ case class PendingWrite(data: ByteVector) extends PeerStorage
+ case class WrittenToDb(data: ByteVector) extends PeerStorage
+ }
case class LastRemoteFeatures(features: Features[InitFeature], written: Boolean)
@@ -1065,7 +1101,7 @@ object Peer {
case object Nothing extends Data {
override def channels: Map[_ <: ChannelId, ActorRef] = Map.empty
override def activeChannels: Set[ByteVector32] = Set.empty
- override def peerStorage: PeerStorage = PeerStorage(None, written = true)
+ override def peerStorage: PeerStorage = PeerStorage.Uninitialized
}
case class DisconnectedData(channels: Map[FinalChannelId, ActorRef], activeChannels: Set[ByteVector32], peerStorage: PeerStorage, remoteFeatures_opt: Option[LastRemoteFeatures]) extends Data
case class ConnectedData(address: NodeAddress, peerConnection: ActorRef, localInit: protocol.Init, remoteInit: protocol.Init, channels: Map[ChannelId, ActorRef], activeChannels: Set[ByteVector32], currentFeerates: RecommendedFeerates, previousFeerates_opt: Option[RecommendedFeerates], peerStorage: PeerStorage, remoteFeaturesWritten: Boolean) extends Data {
@@ -1185,6 +1221,6 @@ object Peer {
case class RelayUnknownMessage(unknownMessage: UnknownMessage)
- case object WritePeerStorage
+ private case object TickWritePeerStorage
// @formatter:on
}
diff --git a/eclair-core/src/test/scala/fr/acinq/eclair/io/PeerSpec.scala b/eclair-core/src/test/scala/fr/acinq/eclair/io/PeerSpec.scala
index af9a163..164767d 100644
--- a/eclair-core/src/test/scala/fr/acinq/eclair/io/PeerSpec.scala
+++ b/eclair-core/src/test/scala/fr/acinq/eclair/io/PeerSpec.scala
@@ -777,6 +777,31 @@ class PeerSpec extends FixtureSpec {
assert(nodeParams.db.peers.getStorage(remoteNodeId).contains(hex"1111"))
}
+ test("lazily load peer storage on first connection") { f =>
+ import f._
+
+ // Peer storage is not loaded after initialization.
+ switchboard.send(peer, Peer.Init(Set.empty, Map.empty))
+ assert(peer.stateData.peerStorage == PeerStorage.Uninitialized)
+
+ // Store some data in the DB for this peer.
+ nodeParams.db.peers.updateStorage(remoteNodeId, hex"abcdef")
+
+ // On first connection, peer storage is lazily loaded from DB.
+ val localInit = protocol.Init(peer.underlyingActor.nodeParams.features.initFeatures())
+ switchboard.send(peer, PeerConnection.ConnectionReady(peerConnection.ref, remoteNodeId, fakeIPAddress, outgoing = true, localInit, protocol.Init(Bob.nodeParams.features.initFeatures())))
+ peerConnection.expectMsg(PeerStorageRetrieval(hex"abcdef"))
+ peerConnection.expectMsgType[RecommendedFeerates]
+ assert(peer.stateData.peerStorage == PeerStorage.WrittenToDb(hex"abcdef"))
+
+ // On reconnection, peer storage stays loaded (no DB re-read).
+ val peerConnection2 = TestProbe()
+ switchboard.send(peer, PeerConnection.ConnectionReady(peerConnection2.ref, remoteNodeId, fakeIPAddress, outgoing = true, localInit, protocol.Init(Bob.nodeParams.features.initFeatures())))
+ peerConnection2.expectMsg(PeerStorageRetrieval(hex"abcdef"))
+ peerConnection2.expectMsgType[RecommendedFeerates]
+ assert(peer.stateData.peerStorage == PeerStorage.WrittenToDb(hex"abcdef"))
+ }
+
test("store remote features when channel confirms") { f =>
import f._
diff --git a/eclair-core/src/test/scala/fr/acinq/eclair/io/ReconnectionTaskSpec.scala b/eclair-core/src/test/scala/fr/acinq/eclair/io/ReconnectionTaskSpec.scala
index 9f2611f..dcdf1d8 100644
--- a/eclair-core/src/test/scala/fr/acinq/eclair/io/ReconnectionTaskSpec.scala
+++ b/eclair-core/src/test/scala/fr/acinq/eclair/io/ReconnectionTaskSpec.scala
@@ -38,8 +38,8 @@ class ReconnectionTaskSpec extends TestKitBaseClass with FixtureAnyFunSuiteLike
private val recommendedFeerates = RecommendedFeerates(Block.RegtestGenesisBlock.hash, TestConstants.feeratePerKw, TestConstants.anchorOutputsFeeratePerKw)
private val PeerNothingData = Peer.Nothing
- private val PeerDisconnectedData = Peer.DisconnectedData(channels, activeChannels = Set.empty, PeerStorage(None, written = true), remoteFeatures_opt = None)
- private val PeerConnectedData = Peer.ConnectedData(fakeIPAddress, system.deadLetters, null, null, channels.map { case (k: ChannelId, v) => (k, v) }, activeChannels = Set.empty, recommendedFeerates, None, PeerStorage(None, written = true), remoteFeaturesWritten = true)
+ private val PeerDisconnectedData = Peer.DisconnectedData(channels, activeChannels = Set.empty, PeerStorage.Empty, remoteFeatures_opt = None)
+ private val PeerConnectedData = Peer.ConnectedData(fakeIPAddress, system.deadLetters, null, null, channels.map { case (k: ChannelId, v) => (k, v) }, activeChannels = Set.empty, recommendedFeerates, None, PeerStorage.Empty, remoteFeaturesWritten = true)
case class FixtureParam(nodeParams: NodeParams, remoteNodeId: PublicKey, reconnectionTask: TestFSMRef[ReconnectionTask.State, ReconnectionTask.Data, ReconnectionTask], monitor: TestProbe)
@@ -82,7 +82,7 @@ class ReconnectionTaskSpec extends TestKitBaseClass with FixtureAnyFunSuiteLike
import f._
val peer = TestProbe()
- peer.send(reconnectionTask, Peer.Transition(PeerNothingData, Peer.DisconnectedData(Map.empty, activeChannels = Set.empty, PeerStorage(None, written = true), None)))
+ peer.send(reconnectionTask, Peer.Transition(PeerNothingData, Peer.DisconnectedData(Map.empty, activeChannels = Set.empty, PeerStorage.Empty, None)))
monitor.expectNoMessage()
}
Why this scored 18/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.