From fdeb04db02a713e852bf19923a0f4ca35c6a3572 Mon Sep 17 00:00:00 2001 From: PJ Fanning Date: Mon, 21 Sep 2026 15:58:44 +0100 Subject: [PATCH] test: make RemotingSpec bind retry match the classic Netty BindException Motivation: RemotingSpec "allow other system to connect even if it's not there at first" still fails intermittently in CI with `java.net.BindException: Address already in use` when the port from temporaryServerAddress() is taken before the new ActorSystem binds. The retry in selectionAndBind only fires when the exception message contains "Failed to bind", but classic Netty remoting rethrows the raw BindException from ActorSystem startup, so the guard never matched and the retry added in #3015 was dead code on this path. Modification: Add an isBindFailure helper that walks the cause chain and matches on java.net.BindException (keeping the "Failed to bind" message check for other transports), and use it as the retry guard in selectionAndBind. Add a directional test that starts an ActorSystem on a port already held by a ServerSocket and asserts the thrown exception is recognized by isBindFailure. Result: The lazy-connect tests retry on a fresh port when the temporary port is already in use instead of failing the suite. Tests: - scalafmt remote/src/test/scala/org/apache/pekko/remote/classic/RemotingSpec.scala - sbt "remote/testOnly org.apache.pekko.remote.classic.RemotingSpec" (23 tests, all passing) - git diff --check References: Refs #1679, Refs #3015 - port race in RemotingSpec lazy-connect tests --- .../pekko/remote/classic/RemotingSpec.scala | 32 +++++++++++++++++-- 1 file changed, 30 insertions(+), 2 deletions(-) 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 8f451ece030..770c4e9507b 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 @@ -14,9 +14,10 @@ package org.apache.pekko.remote.classic import java.io.NotSerializableException +import java.net.{ BindException, InetAddress, ServerSocket } import java.util.concurrent.ThreadLocalRandom -import scala.annotation.nowarn +import scala.annotation.{ nowarn, tailrec } import scala.concurrent.{ Await, Future } import scala.concurrent.duration._ import scala.util.control.NonFatal @@ -966,6 +967,17 @@ class RemotingSpec extends PekkoSpec(RemotingSpec.cfg) with ImplicitSender with } + // classic Netty rethrows the raw java.net.BindException from ActorSystem startup, so look for it in + // the cause chain rather than relying on the message text + @tailrec + def isBindFailure(t: Throwable): Boolean = + t match { + case null => false + case _: BindException => true + case _ if t.getMessage != null && t.getMessage.contains("Failed to bind") => true + case _ => isBindFailure(t.getCause) + } + // retry a few times as the temporaryServerAddress can be taken by the time the new actor system // binds def selectionAndBind( @@ -984,13 +996,29 @@ class RemotingSpec extends PekkoSpec(RemotingSpec.cfg) with ImplicitSender with try { (ActorSystem("other-system", otherConfig), otherSelection) } catch { - case NonFatal(ex) if ex.getMessage.contains("Failed to bind") && retries > 0 => + case NonFatal(ex) if isBindFailure(ex) && retries > 0 => selectionAndBind(config, thisSystem, probe, retries = retries - 1) case other => throw other } } + "recognize a bind failure when the configured port is already taken" in { + val taken = new ServerSocket(0, 1, InetAddress.getByName("localhost")) + try { + val config = ConfigFactory.parseString(s""" + pekko.remote.classic.enabled-transports = ["pekko.remote.classic.netty.tcp"] + pekko.remote.classic.netty.tcp.port = ${taken.getLocalPort} + """).withFallback(remoteSystem.settings.config) + val ex = intercept[Exception] { + val sys = ActorSystem("taken-port-system", config) + shutdown(sys) + } + isBindFailure(ex) shouldBe true + isBindFailure(new IllegalStateException("something else")) shouldBe false + } finally taken.close() + } + "be able to connect to system even if it's not there at first" in { val config = ConfigFactory.parseString(s""" pekko.remote.classic.enabled-transports = ["pekko.remote.classic.netty.tcp"]