This is an automated email from the ASF dual-hosted git repository.

He-Pin pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/pekko.git


The following commit(s) were added to refs/heads/main by this push:
     new 05bdc601b8 fix(remote): buffer messages during passive handoff (#3361)
05bdc601b8 is described below

commit 05bdc601b81bf5ff38b3c2c622247616b51bce4f
Author: He-Pin(kerr) <[email protected]>
AuthorDate: Mon Jul 20 10:55:41 2026 +0800

    fix(remote): buffer messages during passive handoff (#3361)
    
    Motivation:
    Classic remoting discards unreliable inbound messages while the original 
reader is suspended for a passive connection handoff. If the replacement fails, 
those messages are permanently lost.
    
    Modification:
    Buffer a bounded number of unreliable payloads on the suspended reader and 
replay them only when the replacement has not dispatched anything. If it has, 
reconnect the writer while preserving reliable-delivery state to avoid 
reordering. Add configuration and deterministic regression coverage.
    
    Result:
    Failed passive handoffs recover buffered unreliable payloads when safe 
without violating actor message ordering or losing buffered system messages.
    
    Tests:
    - sbt "remote / Test / testOnly 
org.apache.pekko.remote.classic.RemotingSpec" (22 passed)
    - sbt "remote / Test / testOnly 
org.apache.pekko.remote.ReliableDeliverySupervisorSpec" (5 passed)
    - sbt "remote / Test / testOnly org.apache.pekko.remote.RemoteConfigSpec" 
(6 passed)
    - sbt "++3.3.8" focused Remoting and ReliableDeliverySupervisor tests 
(passed)
    - sbt +mimaReportBinaryIssues (passed for Scala 2.13 and 3.3)
    - sbt headerCreateAll +headerCheckAll checkCodeStyle (passed)
    - sbt docs/paradox (passed)
    - scalafmt --list --mode diff-ref=origin/main (passed)
    - sbt sortImports (failed: existing Scalafix/scala-meta NoSuchMethodError; 
unrelated generated changes reverted)
    - validatePullRequest not run (per user request)
    
    References:
    Fixes #3208
---
 remote/src/main/resources/reference.conf           |   6 +
 .../scala/org/apache/pekko/remote/Endpoint.scala   | 124 +++++++++++++++----
 .../org/apache/pekko/remote/RemoteSettings.scala   |   5 +
 .../scala/org/apache/pekko/remote/Remoting.scala   |  51 +++++---
 .../remote/transport/PekkoProtocolTransport.scala  |   3 +
 .../remote/ReliableDeliverySupervisorSpec.scala    |  42 +++++++
 .../org/apache/pekko/remote/RemoteConfigSpec.scala |   1 +
 .../apache/pekko/remote/classic/RemotingSpec.scala | 131 +++++++++++++++++----
 8 files changed, 302 insertions(+), 61 deletions(-)

diff --git a/remote/src/main/resources/reference.conf 
b/remote/src/main/resources/reference.conf
index 2a24300e81..308bb3064a 100644
--- a/remote/src/main/resources/reference.conf
+++ b/remote/src/main/resources/reference.conf
@@ -355,6 +355,12 @@ pekko {
       # Reuse inbound connections for outbound messages
       use-passive-connections = on
 
+      # Maximum number of unreliable inbound messages to buffer on the original
+      # outbound connection while its reader is suspended during a passive
+      # connection handoff. The messages are replayed if the passive reader 
fails
+      # before delivering a message; additional messages are discarded.
+      passive-connection-buffer-size = 128
+
       # Controls the backoff interval after a refused write is reattempted.
       # (Transports may refuse writes if their internal buffer is full)
       backoff-interval = 5 ms
diff --git a/remote/src/main/scala/org/apache/pekko/remote/Endpoint.scala 
b/remote/src/main/scala/org/apache/pekko/remote/Endpoint.scala
index 8e0208bade..866607d52c 100644
--- a/remote/src/main/scala/org/apache/pekko/remote/Endpoint.scala
+++ b/remote/src/main/scala/org/apache/pekko/remote/Endpoint.scala
@@ -15,6 +15,7 @@ package org.apache.pekko.remote
 
 import java.io.NotSerializableException
 import java.util.concurrent.{ ConcurrentHashMap, TimeoutException }
+import java.util.concurrent.atomic.AtomicBoolean
 import java.util.concurrent.locks.LockSupport
 
 import scala.annotation.nowarn
@@ -224,6 +225,7 @@ private[remote] class OversizedPayloadException(msg: 
String) extends EndpointExc
  */
 private[remote] object ReliableDeliverySupervisor {
   case object Ungate
+  case object Reconnect
   case object AttemptSysMsgRedelivery
   final case class GotUid(uid: Int, remoteAddres: Address)
   @nowarn("msg=deprecated")
@@ -308,6 +310,11 @@ private[remote] class ReliableDeliverySupervisor(
     bailoutAt = None
   }
 
+  private def prepareForReconnect(): Unit = {
+    currentHandle = None
+    uidConfirmed = false
+  }
+
   reset()
 
   def nextSeq(): SeqNo = {
@@ -387,6 +394,11 @@ private[remote] class ReliableDeliverySupervisor(
     case s: EndpointWriter.StopReading =>
       writer.forward(s)
 
+    case Reconnect =>
+      prepareForReconnect()
+      context.stop(writer)
+      context.become(gated(writerTerminated = false, earlyUngateRequested = 
true))
+
     case EndpointWriter.ResumeReading =>
       writer ! EndpointWriter.ResumeReading
 
@@ -436,6 +448,13 @@ private[remote] class ReliableDeliverySupervisor(
     case EndpointWriter.StopReading(w, replyTo) =>
       replyTo ! EndpointWriter.StoppedReading(w)
       sender() ! EndpointWriter.StoppedReading(w)
+    case Reconnect =>
+      prepareForReconnect()
+      if (writerTerminated) self ! Ungate
+      else {
+        context.stop(writer)
+        context.become(gated(writerTerminated = false, earlyUngateRequested = 
true))
+      }
   }
 
   def idle: Receive = {
@@ -465,6 +484,8 @@ private[remote] class ReliableDeliverySupervisor(
       goToActive()
     case EndpointWriter.StopReading(w, replyTo) =>
       replyTo ! EndpointWriter.StoppedReading(w)
+    case Reconnect =>
+      prepareForReconnect()
     case Ungate => // ok, not gated
   }
 
@@ -1069,7 +1090,8 @@ private[remote] class EndpointWriter(
                 inbound,
                 handle.handshakeInfo.uid,
                 reliableDeliverySupervisor,
-                receiveBuffers))
+                receiveBuffers,
+                handle.inboundMessageDispatched))
             .withDeploy(Deploy.local),
           "endpointReader-" + AddressUrlEncoder(remoteAddress) + "-" + 
readerId.next()))
     handle.readHandlerPromise.success(ActorHandleEventListener(newReader))
@@ -1104,6 +1126,31 @@ private[remote] object EndpointReader {
       uid: Int,
       reliableDeliverySupervisor: Option[ActorRef],
       receiveBuffers: ConcurrentHashMap[Link, ResendState]): Props =
+    props(
+      localAddress,
+      remoteAddress,
+      transport,
+      settings,
+      codec,
+      msgDispatch,
+      inbound,
+      uid,
+      reliableDeliverySupervisor,
+      receiveBuffers,
+      new AtomicBoolean(false))
+
+  def props(
+      localAddress: Address,
+      remoteAddress: Address,
+      transport: Transport,
+      settings: RemoteSettings,
+      codec: PekkoPduCodec,
+      msgDispatch: InboundMessageDispatcher,
+      inbound: Boolean,
+      uid: Int,
+      reliableDeliverySupervisor: Option[ActorRef],
+      receiveBuffers: ConcurrentHashMap[Link, ResendState],
+      inboundMessageDispatched: AtomicBoolean): Props =
     Props(
       classOf[EndpointReader],
       localAddress,
@@ -1115,7 +1162,8 @@ private[remote] object EndpointReader {
       inbound,
       uid,
       reliableDeliverySupervisor,
-      receiveBuffers)
+      receiveBuffers,
+      inboundMessageDispatched)
 
 }
 
@@ -1133,13 +1181,39 @@ private[remote] class EndpointReader(
     val inbound: Boolean,
     val uid: Int,
     val reliableDeliverySupervisor: Option[ActorRef],
-    val receiveBuffers: ConcurrentHashMap[Link, ResendState])
+    val receiveBuffers: ConcurrentHashMap[Link, ResendState],
+    val inboundMessageDispatched: AtomicBoolean)
     extends EndpointActor(localAddress, remoteAddress, transport, settings, 
codec) {
 
+  def this(
+      localAddress: Address,
+      remoteAddress: Address,
+      transport: Transport,
+      settings: RemoteSettings,
+      codec: PekkoPduCodec,
+      msgDispatch: InboundMessageDispatcher,
+      inbound: Boolean,
+      uid: Int,
+      reliableDeliverySupervisor: Option[ActorRef],
+      receiveBuffers: ConcurrentHashMap[Link, ResendState]) =
+    this(
+      localAddress,
+      remoteAddress,
+      transport,
+      settings,
+      codec,
+      msgDispatch,
+      inbound,
+      uid,
+      reliableDeliverySupervisor,
+      receiveBuffers,
+      new AtomicBoolean(false))
+
   import EndpointWriter.{ OutboundAck, ResumeReading, StopReading, 
StoppedReading }
 
   val provider = RARP(context.system).provider
   var ackedReceiveBuffer = new AckedReceiveBuffer[Message]
+  private var suspendedMessages = Vector.empty[Message]
 
   override def preStart(): Unit = {
     receiveBuffers.get(Link(localAddress, remoteAddress)) match {
@@ -1189,10 +1263,7 @@ private[remote] class EndpointReader(
             ackedReceiveBuffer = ackedReceiveBuffer.receive(msg)
             deliverAndAck()
           } else
-            try msgDispatch.dispatch(msg.recipient, msg.recipientAddress, 
msg.serializedMessage, msg.senderOption)
-            catch {
-              case NonFatal(e) => logTransientSerializationError(msg, e)
-            }
+            dispatchMessage(msg)
 
         case None =>
       }
@@ -1230,6 +1301,8 @@ private[remote] class EndpointReader(
       replyTo ! StoppedReading(writer)
 
     case ResumeReading =>
+      suspendedMessages.foreach(dispatchMessage)
+      suspendedMessages = Vector.empty
       context.become(receive)
 
     case InboundPayload(p) if p.length <= transport.maximumPayloadBytes =>
@@ -1237,13 +1310,21 @@ private[remote] class EndpointReader(
       for (ack <- ackOption; reliableDelivery <- reliableDeliverySupervisor)
         reliableDelivery ! ReliableDeliverySupervisor.AckFromReader(uid, ack)
 
-      if (log.isWarningEnabled)
-        log.warning(
-          "Discarding inbound message to [{}] in read-only association to 
[{}]. " +
-          "If this happens often you may consider using 
pekko.remote.classic.use-passive-connections=off " +
-          "or use Artery TCP.",
-          msgOption.map(_.recipient).getOrElse("unknown"),
-          remoteAddress)
+      msgOption.foreach { msg =>
+        if (!msg.reliableDeliveryEnabled) {
+          if (suspendedMessages.size < settings.PassiveConnectionBufferSize)
+            suspendedMessages :+= msg
+          else if (log.isWarningEnabled)
+            log.warning(
+              "Discarding inbound message to [{}] in read-only association to 
[{}] because the passive connection " +
+              "buffer of [{}] messages is full. If this happens often you may 
consider increasing " +
+              "pekko.remote.classic.passive-connection-buffer-size, setting " +
+              "pekko.remote.classic.use-passive-connections=off, or using 
Artery TCP.",
+              msg.recipient,
+              remoteAddress,
+              settings.PassiveConnectionBufferSize)
+        }
+      }
 
     case InboundPayload(oversized) =>
       log.error(
@@ -1279,14 +1360,17 @@ private[remote] class EndpointReader(
 
     // Notify writer that some messages can be acked
     context.parent ! OutboundAck(ack)
-    deliver.foreach { m =>
-      try msgDispatch.dispatch(m.recipient, m.recipientAddress, 
m.serializedMessage, m.senderOption)
-      catch {
-        case NonFatal(e) => logTransientSerializationError(m, e)
-      }
-    }
+    deliver.foreach(dispatchMessage)
   }
 
+  private def dispatchMessage(msg: Message): Unit =
+    try {
+      inboundMessageDispatched.set(true)
+      msgDispatch.dispatch(msg.recipient, msg.recipientAddress, 
msg.serializedMessage, msg.senderOption)
+    } catch {
+      case NonFatal(e) => logTransientSerializationError(msg, e)
+    }
+
   private def tryDecodeMessageAndAck(pdu: ByteString): (Option[Ack], 
Option[Message]) =
     try {
       codec.decodeMessage(pdu, provider, localAddress)
diff --git a/remote/src/main/scala/org/apache/pekko/remote/RemoteSettings.scala 
b/remote/src/main/scala/org/apache/pekko/remote/RemoteSettings.scala
index 1c20795bb6..ab6651ffb0 100644
--- a/remote/src/main/scala/org/apache/pekko/remote/RemoteSettings.scala
+++ b/remote/src/main/scala/org/apache/pekko/remote/RemoteSettings.scala
@@ -110,6 +110,11 @@ final class RemoteSettings(val config: Config) {
   @deprecated("Classic remoting is deprecated, use Artery", "Akka 2.6.0")
   val UsePassiveConnections: Boolean = 
getBoolean("pekko.remote.classic.use-passive-connections")
 
+  @deprecated("Classic remoting is deprecated, use Artery", "Akka 2.6.0")
+  val PassiveConnectionBufferSize: Int = {
+    getInt("pekko.remote.classic.passive-connection-buffer-size")
+  }.requiring(_ >= 0, "passive-connection-buffer-size must be >= 0")
+
   @deprecated("Classic remoting is deprecated, use Artery", "Akka 2.6.0")
   val BackoffPeriod: FiniteDuration = {
     config.getMillisDuration("pekko.remote.classic.backoff-interval")
diff --git a/remote/src/main/scala/org/apache/pekko/remote/Remoting.scala 
b/remote/src/main/scala/org/apache/pekko/remote/Remoting.scala
index 3aed2fcc75..d314c1bc6a 100644
--- a/remote/src/main/scala/org/apache/pekko/remote/Remoting.scala
+++ b/remote/src/main/scala/org/apache/pekko/remote/Remoting.scala
@@ -538,6 +538,7 @@ private[remote] class EndpointManager(conf: Config, log: 
LoggingAdapter)
 
   var pendingReadHandoffs = Map[ActorRef, PekkoProtocolHandle]()
   private var readOnlyReaderResumptions = Map[ActorRef, ResumableReader]()
+  private var readOnlyEndpointHandles = Map[ActorRef, PekkoProtocolHandle]()
   var stashedInbound = Map[ActorRef, Vector[InboundAssociation]]()
 
   def handleStashedInbound(endpoint: ActorRef, writerIsIdle: Boolean): Unit = {
@@ -850,6 +851,10 @@ private[remote] class EndpointManager(conf: Config, log: 
LoggingAdapter)
           pendingReadHandoffs
             .get(endpoint)
             .foreach(_.disassociate("the existing readOnly association was 
replaced by a new incoming one", log))
+          readOnlyEndpointHandles.get(endpoint).foreach { previousHandle =>
+            handle.inboundMessageDispatched = 
previousHandle.inboundMessageDispatched
+            readOnlyEndpointHandles += endpoint -> handle
+          }
           pendingReadHandoffs += endpoint -> handle
           endpoint ! EndpointWriter.TakeOver(handle, self)
           endpoints.writableEndpointWithPolicyFor(handle.remoteAddress) match {
@@ -971,6 +976,11 @@ private[remote] class EndpointManager(conf: Config, log: 
LoggingAdapter)
       pendingReadHandoffs -= takingOverFrom
       eventPublisher.notifyListeners(AssociatedEvent(handle.localAddress, 
handle.remoteAddress, inbound = true))
 
+      val inheritedResumableReader = 
readOnlyReaderResumptions.get(takingOverFrom).filter(
+        resumableReaderMatchesHandle(_, readOnlyEndpoint = takingOverFrom, 
handle = handle))
+      val resumableReader = resumeReadingFrom
+        .map(writer => ResumableReader(writer, handle.remoteAddress, 
handle.handshakeInfo.uid))
+        .orElse(inheritedResumableReader)
       val endpoint = createEndpoint(
         handle.remoteAddress,
         handle.localAddress,
@@ -979,15 +989,12 @@ private[remote] class EndpointManager(conf: Config, log: 
LoggingAdapter)
         Some(handle),
         writing = false)
       endpoints.registerReadOnlyEndpoint(handle.remoteAddress, endpoint, 
handle.handshakeInfo.uid)
-      val inheritedResumableReader = 
readOnlyReaderResumptions.get(takingOverFrom).filter(
-        resumableReaderMatchesHandle(_, readOnlyEndpoint = takingOverFrom, 
handle = handle))
-      val resumableReader = resumeReadingFrom
-        .map(writer => ResumableReader(writer, handle.remoteAddress, 
handle.handshakeInfo.uid))
-        .orElse(inheritedResumableReader)
       readOnlyReaderResumptions -= takingOverFrom
+      readOnlyEndpointHandles -= takingOverFrom
       resumableReader.foreach {
-        case ResumableReader(writer, remoteAddress, uid) =>
-          readOnlyReaderResumptions += endpoint -> ResumableReader(writer, 
remoteAddress, uid)
+        case resumableReader @ ResumableReader(writer, remoteAddress, uid) =>
+          readOnlyReaderResumptions += endpoint -> resumableReader
+          readOnlyEndpointHandles += endpoint -> handle
           log.debug(
             "Registered passive read handoff fallback for [{}] with UID [{}] 
from writable endpoint [{}] to read-only endpoint [{}]",
             remoteAddress,
@@ -1004,13 +1011,24 @@ private[remote] class EndpointManager(conf: Config, 
log: LoggingAdapter)
       case ResumableReader(writer, remoteAddress, uid) =>
         endpoints.writableEndpointWithPolicyFor(remoteAddress) match {
           case Some(Pass(`writer`, Some(`uid`))) =>
-            log.info(
-              "Resuming outbound reader [{}] for [{}] with UID [{}] after 
passive read-only endpoint [{}] stopped",
-              writer,
-              remoteAddress,
-              uid,
-              readOnlyEndpoint)
-            writer ! EndpointWriter.ResumeReading
+            if 
(readOnlyEndpointHandles.get(readOnlyEndpoint).forall(_.inboundMessageDispatched.get()))
 {
+              log.warning(
+                "Reconnecting outbound endpoint [{}] for [{}] with UID [{}] 
after passive read-only endpoint [{}] " +
+                "dispatched messages and then stopped, because resuming the 
original reader could violate message ordering",
+                writer,
+                remoteAddress,
+                uid,
+                readOnlyEndpoint)
+              writer ! ReliableDeliverySupervisor.Reconnect
+            } else {
+              log.info(
+                "Resuming outbound reader [{}] for [{}] with UID [{}] after 
passive read-only endpoint [{}] stopped",
+                writer,
+                remoteAddress,
+                uid,
+                readOnlyEndpoint)
+              writer ! EndpointWriter.ResumeReading
+            }
           case currentPolicy =>
             log.debug(
               "Not resuming outbound reader [{}] for [{}] with UID [{}] after 
passive read-only endpoint [{}] stopped",
@@ -1025,6 +1043,7 @@ private[remote] class EndpointManager(conf: Config, log: 
LoggingAdapter)
         }
     }
     readOnlyReaderResumptions -= readOnlyEndpoint
+    readOnlyEndpointHandles -= readOnlyEndpoint
   }
 
   private def resumableReaderMatchesHandle(
@@ -1047,8 +1066,10 @@ private[remote] class EndpointManager(conf: Config, log: 
LoggingAdapter)
       endpoints.registerReadOnlyEndpoint(withHandle.remoteAddress, 
takingOverFrom, withHandle.handshakeInfo.uid)
       if (!readOnlyReaderResumptions
           .get(takingOverFrom)
-          .forall(resumableReaderMatchesHandle(_, readOnlyEndpoint = 
takingOverFrom, handle = withHandle)))
+          .forall(resumableReaderMatchesHandle(_, readOnlyEndpoint = 
takingOverFrom, handle = withHandle))) {
         readOnlyReaderResumptions -= takingOverFrom
+        readOnlyEndpointHandles -= takingOverFrom
+      }
     }
   }
 
diff --git 
a/remote/src/main/scala/org/apache/pekko/remote/transport/PekkoProtocolTransport.scala
 
b/remote/src/main/scala/org/apache/pekko/remote/transport/PekkoProtocolTransport.scala
index cba06c4ccf..6ef39f4a2c 100644
--- 
a/remote/src/main/scala/org/apache/pekko/remote/transport/PekkoProtocolTransport.scala
+++ 
b/remote/src/main/scala/org/apache/pekko/remote/transport/PekkoProtocolTransport.scala
@@ -14,6 +14,7 @@
 package org.apache.pekko.remote.transport
 
 import java.util.concurrent.TimeoutException
+import java.util.concurrent.atomic.AtomicBoolean
 
 import scala.annotation.nowarn
 import scala.collection.immutable
@@ -237,6 +238,8 @@ private[remote] class PekkoProtocolHandle(
     override val addedSchemeIdentifier: String)
     extends AbstractTransportAdapterHandle(_localAddress, _remoteAddress, 
_wrappedHandle, addedSchemeIdentifier) {
 
+  private[remote] var inboundMessageDispatched = new AtomicBoolean(false)
+
   override def write(payload: ByteString): Boolean = 
wrappedHandle.write(codec.constructPayload(payload))
 
   override def disassociate(reason: String, log: LoggingAdapter): Unit = {
diff --git 
a/remote/src/test/scala/org/apache/pekko/remote/ReliableDeliverySupervisorSpec.scala
 
b/remote/src/test/scala/org/apache/pekko/remote/ReliableDeliverySupervisorSpec.scala
index 4529eb21fe..684c85e24a 100644
--- 
a/remote/src/test/scala/org/apache/pekko/remote/ReliableDeliverySupervisorSpec.scala
+++ 
b/remote/src/test/scala/org/apache/pekko/remote/ReliableDeliverySupervisorSpec.scala
@@ -25,6 +25,7 @@ import scala.concurrent.duration._
 
 import org.apache.pekko
 import pekko.actor.{ ActorRef, Address, Nobody, RootActorPath, Terminated }
+import pekko.dispatch.sysmsg.Watch
 import pekko.remote.EndpointManager.{ Link, ResendState, Send }
 import pekko.remote.ReliableDeliverySupervisor.AckFromReader
 import pekko.remote.transport._
@@ -146,6 +147,47 @@ class ReliableDeliverySupervisorSpec extends 
PekkoSpec(ReliableDeliverySuperviso
 
       writerProbe.expectMsg(EndpointWriter.ResumeReading)
     }
+
+    "reconnect without discarding buffered system messages" in {
+      val parentProbe = TestProbe()
+      val supervisor = newSupervisor(initialUid = oldUid, parentProbe.ref)
+      val underlying = supervisor.underlyingActor
+      val oldWriter = underlying.writer
+      val writerProbe = TestProbe()
+      val supervisorProbe = TestProbe()
+
+      underlying.resendBuffer = bufferWith(0L, 1L)
+      underlying.seqCounter = 2L
+      writerProbe.watch(oldWriter)
+      supervisorProbe.watch(supervisor)
+      supervisor.unwatch(oldWriter)
+
+      supervisor.receive(ReliableDeliverySupervisor.Reconnect)
+      writerProbe.expectTerminated(oldWriter)
+      underlying.currentHandle shouldBe None
+      underlying.uidConfirmed shouldBe false
+      underlying.resendBuffer.nonAcked.map(_.seq) should ===(Vector(SeqNo(0), 
SeqNo(1)))
+      underlying.seqCounter should ===(2L)
+
+      underlying.writer = writerProbe.ref
+      underlying.context.become(underlying.receive)
+      val systemSend = Send(Watch(remoteRef, remoteRef), OptionVal.None, 
remoteRef, seqOpt = None)
+      supervisor.receive(systemSend)
+
+      writerProbe.expectNoMessage(100.millis)
+      underlying.resendBuffer.nonAcked.map(_.seq) should ===(Vector(SeqNo(0), 
SeqNo(1), SeqNo(2)))
+      underlying.seqCounter should ===(3L)
+
+      supervisor.receive(ReliableDeliverySupervisor.GotUid(oldUid, 
remoteAddress))
+      parentProbe.expectMsg(ReliableDeliverySupervisor.GotUid(oldUid, 
remoteAddress))
+      writerProbe.expectMsg(send(0L))
+      writerProbe.expectMsg(send(1L))
+      writerProbe.expectMsg(systemSend.copy(seqOpt = Some(SeqNo(2L))))
+      underlying.uidConfirmed shouldBe true
+
+      underlying.resendBuffer = new AckedSendBuffer[Send](0)
+      supervisorProbe.expectNoMessage(100.millis)
+    }
   }
 
   private def newSupervisor(initialUid: Int, parent: ActorRef): 
TestActorRef[ReliableDeliverySupervisor] = {
diff --git 
a/remote/src/test/scala/org/apache/pekko/remote/RemoteConfigSpec.scala 
b/remote/src/test/scala/org/apache/pekko/remote/RemoteConfigSpec.scala
index 85e2594287..9a992d8078 100644
--- a/remote/src/test/scala/org/apache/pekko/remote/RemoteConfigSpec.scala
+++ b/remote/src/test/scala/org/apache/pekko/remote/RemoteConfigSpec.scala
@@ -47,6 +47,7 @@ class RemoteConfigSpec extends PekkoSpec("""
       RetryGateClosedFor should ===(5.seconds)
       Dispatcher should ===("pekko.remote.default-remote-dispatcher")
       UsePassiveConnections should ===(true)
+      PassiveConnectionBufferSize should ===(128)
       BackoffPeriod should ===(5.millis)
       LogBufferSizeExceeding should ===(50000)
       SysMsgAckTimeout should ===(0.3.seconds)
diff --git 
a/remote/src/test/scala/org/apache/pekko/remote/classic/RemotingSpec.scala 
b/remote/src/test/scala/org/apache/pekko/remote/classic/RemotingSpec.scala
index ae79f265c6..8f451ece03 100644
--- a/remote/src/test/scala/org/apache/pekko/remote/classic/RemotingSpec.scala
+++ b/remote/src/test/scala/org/apache/pekko/remote/classic/RemotingSpec.scala
@@ -31,7 +31,7 @@ import pekko.remote.transport.AssociationHandle.{ 
HandleEvent, HandleEventListen
 import pekko.remote.transport.Transport.InvalidAssociationException
 import pekko.testkit._
 import pekko.testkit.SocketUtil.temporaryServerAddress
-import pekko.util.ByteString
+import pekko.util.{ ByteString, OptionVal }
 
 import com.typesafe.config._
 
@@ -688,7 +688,7 @@ class RemotingSpec extends PekkoSpec(RemotingSpec.cfg) with 
ImplicitSender with
 
     }
 
-    "should resume the outbound reader when passive read handoff fails" in {
+    "should recover buffered unreliable messages without reordering across a 
failed passive read handoff" in {
       val localAddress = Address("pekko.test", "resume-system1", "localhost", 
101)
       val rawLocalAddress = localAddress.copy(protocol = "test")
       val remoteAddress = Address("pekko.test", "resume-system2", "localhost", 
102)
@@ -700,6 +700,7 @@ class RemotingSpec extends PekkoSpec(RemotingSpec.cfg) with 
ImplicitSender with
         pekko.remote.classic.enabled-transports = ["pekko.remote.classic.test"]
         pekko.remote.classic.retry-gate-closed-for = 5s
         pekko.remote.classic.log-remote-lifecycle-events = on
+        pekko.remote.classic.passive-connection-buffer-size = 1
 
         pekko.remote.classic.test {
           registry-key = TwHFcESm
@@ -713,7 +714,7 @@ class RemotingSpec extends PekkoSpec(RemotingSpec.cfg) with 
ImplicitSender with
         val receiverProbe = TestProbe()(thisSystem)
         thisSystem.actorOf(Props(new Actor {
             def receive = {
-              case message => receiverProbe.ref ! message
+              case message => receiverProbe.ref ! (message -> sender().path)
             }
           }).withDeploy(Deploy.local), receiverName)
         val associationEventProbe = TestProbe()(thisSystem)
@@ -751,7 +752,13 @@ class RemotingSpec extends PekkoSpec(RemotingSpec.cfg) 
with ImplicitSender with
         }
         outboundRemoteHandle.association.write(originalHandshakePacket)
 
-        def payload(message: String): ByteString = {
+        val remoteSender =
+          RARP(thisSystem).provider.resolveActorRef(remoteAddress.toString + 
"/user/handoff-sender")
+
+        def payload(
+            message: String,
+            seqOption: Option[SeqNo] = None,
+            senderOption: OptionVal[ActorRef] = OptionVal.None): ByteString = {
           val serializedMessage =
             
MessageSerializer.serialize(thisSystem.asInstanceOf[ExtendedActorSystem], 
message)
           val recipient = new EmptyLocalActorRef(
@@ -764,41 +771,113 @@ class RemotingSpec extends PekkoSpec(RemotingSpec.cfg) 
with ImplicitSender with
               localAddress,
               recipient,
               serializedMessage,
-              pekko.util.OptionVal.None))
+              senderOption,
+              seqOption = seqOption))
         }
 
-        outboundRemoteHandle.association.write(payload("before-handoff"))
-        receiverProbe.expectMsg("before-handoff")
+        def establishPassiveAssociation(): (TestAssociationHandle, TestProbe) 
= {
+          val inboundHandleProbe = TestProbe()
+          val inboundHandle =
+            Await.result(remoteTransport.associate(rawLocalAddress), 
3.seconds).asInstanceOf[TestAssociationHandle]
+          inboundHandle.readHandlerPromise.success(new 
AssociationHandle.HandleEventListener {
+            override def notify(ev: HandleEvent): Unit = 
inboundHandleProbe.ref ! ev
+          })
 
-        val inboundHandleProbe = TestProbe()
-        val inboundHandle = 
Await.result(remoteTransport.associate(rawLocalAddress), 3.seconds)
-        inboundHandle.readHandlerPromise.success(new 
AssociationHandle.HandleEventListener {
-          override def notify(ev: HandleEvent): Unit = inboundHandleProbe.ref 
! ev
-        })
-
-        awaitAssert {
-          
registry.getRemoteReadHandlerFor(inboundHandle.asInstanceOf[TestAssociationHandle]).get
-        }
+          awaitAssert {
+            registry.getRemoteReadHandlerFor(inboundHandle).get
+          }
 
-        inboundHandle.write(originalHandshakePacket)
-        associationEventProbe.fishForMessage(3.seconds, "passive association") 
{
-          case AssociatedEvent(`localAddress`, `remoteAddress`, true) => true
-          case _                                                      => false
+          inboundHandle.write(originalHandshakePacket)
+          associationEventProbe.fishForMessage(3.seconds, "passive 
association") {
+            case AssociatedEvent(`localAddress`, `remoteAddress`, true) => true
+            case _                                                      => 
false
+          }
+          inboundHandle -> inboundHandleProbe
         }
 
         val brokenPacket = 
PekkoPduProtobufCodec.constructPayload(ByteString(0, 1, 2, 3, 4, 5, 6))
-        inboundHandle.write(brokenPacket)
-        inboundHandleProbe.fishForMessage(3.seconds, "read-only endpoint 
disassociation") {
-          case _: AssociationHandle.Disassociated => true
-          case _                                  => false
+
+        def breakAssociation(handle: TestAssociationHandle, probe: TestProbe): 
Unit = {
+          handle.write(brokenPacket)
+          probe.fishForMessage(3.seconds, "read-only endpoint disassociation") 
{
+            case _: AssociationHandle.Disassociated => true
+            case _                                  => false
+          }
         }
 
+        outboundRemoteHandle.association.write(
+          payload("before-handoff", senderOption = 
OptionVal.Some(remoteSender)))
+        receiverProbe.expectMsg("before-handoff" -> remoteSender.path)
+
+        val (failedHandle, failedHandleProbe) = establishPassiveAssociation()
+
+        EventFilter
+          .warning(pattern = ".*passive connection buffer.*full.*", 
occurrences = 1)
+          .intercept {
+            
outboundRemoteHandle.association.write(payload("reliable-during-handoff", 
Some(SeqNo(0))))
+            outboundRemoteHandle.association.write(
+              payload("during-handoff", senderOption = 
OptionVal.Some(remoteSender)))
+            outboundRemoteHandle.association.write(
+              payload("buffer-overflow", senderOption = 
OptionVal.Some(remoteSender)))
+          }(thisSystem)
+
+        breakAssociation(failedHandle, failedHandleProbe)
+        receiverProbe.expectMsg("during-handoff" -> remoteSender.path)
+        receiverProbe.expectNoMessage(200.millis)
+
         // The read-only handle disassociation and resuming the outbound 
reader are handled by different actors.
         // Retry until the outbound reader has processed ResumeReading.
         awaitAssert {
-          
outboundRemoteHandle.association.write(payload("after-handoff-failure"))
-          receiverProbe.expectMsg(500.millis, "after-handoff-failure")
+          outboundRemoteHandle.association.write(
+            payload("after-handoff-failure", senderOption = 
OptionVal.Some(remoteSender)))
+          receiverProbe.expectMsg(500.millis, "after-handoff-failure" -> 
remoteSender.path)
+        }
+
+        val (orderingHandle, orderingHandleProbe) = 
establishPassiveAssociation()
+
+        EventFilter
+          .warning(pattern = ".*passive connection buffer.*full.*", 
occurrences = 1)
+          .intercept {
+            outboundRemoteHandle.association.write(
+              payload("older-message", senderOption = 
OptionVal.Some(remoteSender)))
+            outboundRemoteHandle.association.write(
+              payload("ordering-buffer-overflow", senderOption = 
OptionVal.Some(remoteSender)))
+          }(thisSystem)
+
+        orderingHandle.write(payload("newer-message", senderOption = 
OptionVal.Some(remoteSender)))
+        receiverProbe.expectMsg("newer-message" -> remoteSender.path)
+
+        EventFilter
+          .warning(pattern = "Reconnecting outbound endpoint.*message 
ordering", occurrences = 1)
+          .intercept {
+            breakAssociation(orderingHandle, orderingHandleProbe)
+          }(thisSystem)
+        outboundHandleProbe.fishForMessage(3.seconds, "original endpoint 
disassociation") {
+          case _: AssociationHandle.Disassociated => true
+          case _                                  => false
+        }
+        receiverProbe.expectNoMessage(200.millis)
+
+        val reconnectedRemoteHandle = awaitAssert {
+          dummySelection.tell("reconnect", system.deadLetters)
+          
remoteTransportProbe.expectMsgType[Transport.InboundAssociation](500.millis)
+        }
+        val reconnectedHandleProbe = TestProbe()
+        reconnectedRemoteHandle.association.readHandlerPromise.success(new 
HandleEventListener {
+          override def notify(ev: HandleEvent): Unit = 
reconnectedHandleProbe.ref ! ev
+        })
+        reconnectedHandleProbe.fishForMessage(3.seconds, "reconnected outbound 
handshake") {
+          case AssociationHandle.InboundPayload(_) => true
+          case _                                   => false
+        }
+        reconnectedRemoteHandle.association.write(originalHandshakePacket)
+        associationEventProbe.fishForMessage(3.seconds, "reconnected outbound 
association") {
+          case AssociatedEvent(`localAddress`, `remoteAddress`, false) => true
+          case _                                                       => false
         }
+        reconnectedRemoteHandle.association.write(
+          payload("after-ordering-reconnect", senderOption = 
OptionVal.Some(remoteSender)))
+        receiverProbe.expectMsg("after-ordering-reconnect" -> 
remoteSender.path)
       } finally shutdown(thisSystem)
 
     }


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to