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

pjfanning pushed a commit to branch 1.7.x
in repository https://gitbox.apache.org/repos/asf/pekko.git


The following commit(s) were added to refs/heads/1.7.x by this push:
     new bf54f9bd22 fix: use a Seq rather than a Set to hold TCP magic 
ByteStrings (#3484) (#3498)
bf54f9bd22 is described below

commit bf54f9bd22496921e439ea6f591b2a99393a9dd4
Author: PJ Fanning <[email protected]>
AuthorDate: Tue Sep 1 08:13:49 2026 +0100

    fix: use a Seq rather than a Set to hold TCP magic ByteStrings (#3484) 
(#3498)
    
    Motivation:
    default to HashSet, which hashes its members, and hashing a ByteString walks
    every one of its bytes. Comparing this handful of 4-byte values is better 
done
    with a couple of equality checks over a small ordered collection.
    
    Modification:
    Change ArterySettings.Advanced.TcpMagicValues from Set[ByteString] to
    immutable.Seq[ByteString], built with distinct.toList rather than toSet so 
the
    de-duplication the Set provided incidentally is kept, and change the
    TcpFraming acceptedMagic parameter to match. These were the only two
    Set[ByteString] in the tree; both are private[pekko] and both arrived in 
#3425.
    
    Add ArterySettingsSpec, which did not exist: the tcp-magic parsing added by
    
    Result:
    TcpFraming's magic check on each inbound connection is a small number of
    ByteString equality comparisons instead of a hashCode over the received 
bytes.
    Acceptance order is now the configured order, where the Set left it
    unspecified; this is consistent with TcpMagic, the outbound magic, already
    being the first configured value. No other behaviour changes.
    
    Tests:
    - sbt "remote/testOnly org.apache.pekko.remote.artery.ArterySettingsSpec" - 
7
      tests succeeded, 0 failed
    - sbt "remote/testOnly org.apache.pekko.remote.artery.tcp.TcpFramingSpec" - 
11
      tests succeeded, 0 failed, covering the default, custom and legacy AKKA 
magic
      paths through the changed contains call
    - sbt "remote/mimaReportBinaryIssues" - no issues, no filters needed
    - ArterySettingsSpec pins the configuration order, the de-duplication 
including
      after 4-byte truncation, and both require rejections, so the collection 
type
      cannot be changed back without a failure.
    
    References:
    Fixes #3483, Refs #3425
    
    Update TcpFramingSpec.scala
    
    Update ArterySettingsSpec.scala
---
 .../pekko/remote/artery/ArterySettings.scala       |  6 +-
 .../pekko/remote/artery/tcp/TcpFraming.scala       |  4 +-
 .../pekko/remote/artery/ArterySettingsSpec.scala   | 81 ++++++++++++++++++++++
 .../pekko/remote/artery/tcp/TcpFramingSpec.scala   |  4 +-
 4 files changed, 90 insertions(+), 5 deletions(-)

diff --git 
a/remote/src/main/scala/org/apache/pekko/remote/artery/ArterySettings.scala 
b/remote/src/main/scala/org/apache/pekko/remote/artery/ArterySettings.scala
index fd847b3fe2..4e051d5732 100644
--- a/remote/src/main/scala/org/apache/pekko/remote/artery/ArterySettings.scala
+++ b/remote/src/main/scala/org/apache/pekko/remote/artery/ArterySettings.scala
@@ -131,13 +131,15 @@ private[pekko] final class ArterySettings private 
(config: Config) {
         s"tcp-magic value [$first] must produce at least 4 UTF-8 bytes, but 
produced [${bytes.length}] bytes")
       bytes.take(4)
     }
-    val TcpMagicValues: Set[ByteString] = {
+    // A Seq rather than a Set: a Set hashes its members, and hashing a 
ByteString walks all its
+    // bytes, where comparing this handful of 4-byte values is a couple of 
cheap equality checks.
+    val TcpMagicValues: immutable.Seq[ByteString] = {
       tcpMagicList.map { s =>
         val bytes = ByteString(s.getBytes(StandardCharsets.UTF_8))
         require(bytes.length >= 4,
           s"tcp-magic value [$s] must produce at least 4 UTF-8 bytes, but 
produced [${bytes.length}] bytes")
         bytes.take(4)
-      }.toSet
+      }.distinct.toList
     }
     val Dispatcher: String = getString("use-dispatcher")
     val ControlStreamDispatcher: String = 
getString("use-control-stream-dispatcher")
diff --git 
a/remote/src/main/scala/org/apache/pekko/remote/artery/tcp/TcpFraming.scala 
b/remote/src/main/scala/org/apache/pekko/remote/artery/tcp/TcpFraming.scala
index 39da7b75c6..d752ecd015 100644
--- a/remote/src/main/scala/org/apache/pekko/remote/artery/tcp/TcpFraming.scala
+++ b/remote/src/main/scala/org/apache/pekko/remote/artery/tcp/TcpFraming.scala
@@ -17,6 +17,8 @@ package tcp
 import java.nio.ByteBuffer
 import java.nio.ByteOrder
 
+import scala.collection.immutable
+
 import org.apache.pekko
 import pekko.annotation.InternalApi
 import pekko.stream.Attributes
@@ -76,7 +78,7 @@ import pekko.util.ByteString
  * INTERNAL API
  */
 @InternalApi private[pekko] class TcpFraming(
-    acceptedMagic: Set[ByteString] = Set(TcpFraming.DefaultMagic),
+    acceptedMagic: immutable.Seq[ByteString] = List(TcpFraming.DefaultMagic),
     flightRecorder: RemotingFlightRecorder = NoOpRemotingFlightRecorder)
     extends ByteStringParser[EnvelopeBuffer] {
 
diff --git 
a/remote/src/test/scala/org/apache/pekko/remote/artery/ArterySettingsSpec.scala 
b/remote/src/test/scala/org/apache/pekko/remote/artery/ArterySettingsSpec.scala
new file mode 100644
index 0000000000..23f2344c91
--- /dev/null
+++ 
b/remote/src/test/scala/org/apache/pekko/remote/artery/ArterySettingsSpec.scala
@@ -0,0 +1,81 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ *    http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.pekko.remote.artery
+
+import java.nio.charset.StandardCharsets
+
+import org.apache.pekko
+import pekko.util.ByteString
+
+import com.typesafe.config.ConfigFactory
+import org.scalatest.matchers.should.Matchers
+import org.scalatest.wordspec.AnyWordSpec
+
+class ArterySettingsSpec extends AnyWordSpec with Matchers {
+
+  private def settings(tcpMagic: String): ArterySettings =
+    ArterySettings(
+      ConfigFactory
+        .parseString(s"pekko.remote.artery.advanced.tcp-magic = $tcpMagic")
+        .withFallback(ConfigFactory.load())
+        .resolve()
+        .getConfig("pekko.remote.artery"))
+
+  private def magic(s: String): ByteString = 
ByteString(s.getBytes(StandardCharsets.UTF_8))
+
+  "ArterySettings.TcpMagicValues" must {
+
+    // Held as a Seq rather than a Set: a Set hashes its members, and hashing 
a ByteString walks
+    // all of its bytes. These assertions pin the ordering and de-duplication 
that the Set
+    // previously provided incidentally, so the collection type cannot be 
changed back silently.
+    "default to AKKA then AKKA, in configuration order" in {
+      val defaults = 
ArterySettings(ConfigFactory.load().getConfig("pekko.remote.artery"))
+      defaults.Advanced.TcpMagicValues should ===(List(magic("AKKA"), 
magic("PEKK")))
+      defaults.Advanced.TcpMagic should ===(magic("AKKA"))
+    }
+
+    "preserve configuration order" in {
+      settings("""["AKKA", "PEKK"]""").Advanced.TcpMagicValues should 
===(List(magic("AKKA"), magic("PEKK")))
+    }
+
+    "use the first configured value as the outbound magic" in {
+      settings("""["AKKA", "PEKK"]""").Advanced.TcpMagic should 
===(magic("AKKA"))
+    }
+
+    "drop duplicates, keeping the first occurrence" in {
+      settings("""["PEKK", "AKKA", "PEKK"]""").Advanced.TcpMagicValues should 
===(
+        List(magic("PEKK"), magic("AKKA")))
+    }
+
+    "truncate each value to 4 bytes, and de-duplicate after truncating" in {
+      settings("""["PEKKO"]""").Advanced.TcpMagicValues should 
===(List(magic("PEKK")))
+      settings("""["PEKKO", "PEKK"]""").Advanced.TcpMagicValues should 
===(List(magic("PEKK")))
+    }
+
+    // Advanced is an object, so it is initialised lazily: the requires do not 
run until one of
+    // its members is touched, which is why these force TcpMagicValues rather 
than just building
+    // the settings.
+    "reject a value shorter than 4 UTF-8 bytes" in {
+      an[IllegalArgumentException] should be thrownBy 
settings("""["PEK"]""").Advanced.TcpMagicValues
+    }
+
+    "reject an empty list" in {
+      an[IllegalArgumentException] should be thrownBy 
settings("[]").Advanced.TcpMagicValues
+    }
+  }
+}
diff --git 
a/remote/src/test/scala/org/apache/pekko/remote/artery/tcp/TcpFramingSpec.scala 
b/remote/src/test/scala/org/apache/pekko/remote/artery/tcp/TcpFramingSpec.scala
index 39f5e2ee18..f27d0547d8 100644
--- 
a/remote/src/test/scala/org/apache/pekko/remote/artery/tcp/TcpFramingSpec.scala
+++ 
b/remote/src/test/scala/org/apache/pekko/remote/artery/tcp/TcpFramingSpec.scala
@@ -31,7 +31,7 @@ class TcpFramingSpec extends PekkoSpec("""
   import TcpFraming.encodeFrameHeader
 
   private val magic = TcpFraming.DefaultMagic
-  private val acceptedMagic = Set(magic, TcpFraming.PekkoMagic)
+  private val acceptedMagic = List(magic, TcpFraming.PekkoMagic)
   private val framingFlow = Flow[ByteString].via(new TcpFraming(acceptedMagic))
 
   private val payload5 = ByteString((1 to 5).map(_.toByte).toArray)
@@ -117,7 +117,7 @@ class TcpFramingSpec extends PekkoSpec("""
 
     "accept custom magic" in {
       val customMagic = ByteString('T'.toByte, 'E'.toByte, 'S'.toByte, 
'T'.toByte)
-      val customFramingFlow = Flow[ByteString].via(new 
TcpFraming(Set(customMagic)))
+      val customFramingFlow = Flow[ByteString].via(new 
TcpFraming(List(customMagic)))
       val bytes = TcpFraming.encodeConnectionHeader(customMagic, 2) ++ 
frameBytes(1)
       val frames = 
Source(List(bytes)).via(customFramingFlow).runWith(Sink.seq).futureValue
       frames.head.streamId should ===(2)


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

Reply via email to