diff --git a/serialization-jackson/src/main/resources/reference.conf b/serialization-jackson/src/main/resources/reference.conf index 5ad52a7110..593a4ad0aa 100644 --- a/serialization-jackson/src/main/resources/reference.conf +++ b/serialization-jackson/src/main/resources/reference.conf @@ -210,7 +210,11 @@ pekko.serialization.jackson { # payload that decompresses to more than this is rejected rather than # allocated, guarding against a small message that inflates without bound. # This applies on deserialization regardless of the `algorithm` setting above. - max-decompressed-size = 256 MiB + # The default of `unlimited` applies no limit, preserving the behaviour of + # earlier releases; a negative number such as -1 also means unlimited. Set a + # size such as `256 MiB` to bound decompression, choosing a value larger than + # any payload the system legitimately exchanges. + max-decompressed-size = unlimited } # Whether the type should be written to the manifest. diff --git a/serialization-jackson/src/main/scala/org/apache/pekko/serialization/jackson/JacksonSerializer.scala b/serialization-jackson/src/main/scala/org/apache/pekko/serialization/jackson/JacksonSerializer.scala index 82fb9cd3f8..22145470a1 100644 --- a/serialization-jackson/src/main/scala/org/apache/pekko/serialization/jackson/JacksonSerializer.scala +++ b/serialization-jackson/src/main/scala/org/apache/pekko/serialization/jackson/JacksonSerializer.scala @@ -207,7 +207,16 @@ import pekko.util.OptionVal """"off" or "gzip"""") } } - private val maxDecompressedSize: Long = conf.getBytes("compression.max-decompressed-size") + // "unlimited" or a negative number means no limit; getBytes refuses both, so read them first + private val maxDecompressedSize: Long = { + val raw = conf.getString("compression.max-decompressed-size") + if (raw == "unlimited") -1L + else + raw.toLongOption match { + case Some(n) if n < 0 => n + case _ => conf.getBytes("compression.max-decompressed-size") + } + } private val migrations: Map[String, JacksonMigration] = { import scala.jdk.CollectionConverters._ conf.getConfig("migrations").root.unwrapped.asScala.toMap.map { @@ -544,7 +553,10 @@ import pekko.util.OptionVal // meta.length is the decompressed size declared on the wire; a small // message can declare a huge (or negative) size and drive a large // allocation, so bound it before decompressing. - if (meta.length < 0 || meta.length > maxDecompressedSize) + if (meta.length < 0) + throw new IllegalArgumentException( + s"Compressed message declares a negative decompressed size [${meta.length}] bytes") + if (maxDecompressedSize >= 0 && meta.length > maxDecompressedSize) throw new IllegalArgumentException( s"Compressed message declares decompressed size [${meta.length}] bytes, which exceeds the maximum " + s"of [$maxDecompressedSize] bytes (pekko.serialization.jackson.compression.max-decompressed-size)") @@ -556,7 +568,7 @@ import pekko.util.OptionVal } // gunzip with a bound on the decompressed size, so a small gzip payload cannot - // inflate without limit (a "zip bomb"). + // inflate without limit (a "zip bomb"). A negative maximum applies no bound. private def gunzip(in: GZIPInputStream): Array[Byte] = { val out = new ByteArrayOutputStream() val buffer = new Array[Byte](BufferSize) @@ -564,7 +576,7 @@ import pekko.util.OptionVal var n = in.read(buffer) while (n != -1) { total += n - if (total > maxDecompressedSize) + if (maxDecompressedSize >= 0 && total > maxDecompressedSize) throw new IllegalArgumentException( s"Decompressed message exceeds the maximum of [$maxDecompressedSize] bytes " + "(pekko.serialization.jackson.compression.max-decompressed-size)") diff --git a/serialization-jackson/src/test/scala/org/apache/pekko/serialization/jackson/JacksonSerializerSpec.scala b/serialization-jackson/src/test/scala/org/apache/pekko/serialization/jackson/JacksonSerializerSpec.scala index b2b6c2eb2c..87dcef6b6d 100644 --- a/serialization-jackson/src/test/scala/org/apache/pekko/serialization/jackson/JacksonSerializerSpec.scala +++ b/serialization-jackson/src/test/scala/org/apache/pekko/serialization/jackson/JacksonSerializerSpec.scala @@ -687,6 +687,42 @@ class JacksonJsonSerializerSpec extends JacksonSerializerSpec("jackson-json") { } ex.getMessage should include("max-decompressed-size") } + + "apply no gzip decompression limit when max-decompressed-size is -1" in withSystem(""" + pekko.serialization.jackson.jackson-json.compression { + algorithm = gzip + compress-larger-than = 0 KiB + max-decompressed-size = -1 + } + """) { sys => + val msg = SimpleCommand("0" * (8 * 1024)) + JacksonSerializer.isGZipped(serializeToBinary(msg, sys)) should ===(true) + checkSerialization(msg, sys) + } + + "apply no lz4 decompression limit when max-decompressed-size is -1" in withSystem(""" + pekko.serialization.jackson.jackson-json.compression { + algorithm = lz4 + compress-larger-than = 0 KiB + max-decompressed-size = -1 + } + """) { sys => + val msg = SimpleCommand("0" * (8 * 1024)) + JacksonSerializer.isLZ4(serializeToBinary(msg, sys)) should ===(true) + checkSerialization(msg, sys) + } + + "apply no decompression limit when max-decompressed-size is unlimited" in withSystem(""" + pekko.serialization.jackson.jackson-json.compression { + algorithm = gzip + compress-larger-than = 0 KiB + max-decompressed-size = unlimited + } + """) { sys => + val msg = SimpleCommand("0" * (8 * 1024)) + JacksonSerializer.isGZipped(serializeToBinary(msg, sys)) should ===(true) + checkSerialization(msg, sys) + } } "JacksonJsonSerializer without type in manifest" should { diff --git a/serialization-jackson3/src/main/resources/reference.conf b/serialization-jackson3/src/main/resources/reference.conf index 0a94af43f2..18c37b5e9e 100644 --- a/serialization-jackson3/src/main/resources/reference.conf +++ b/serialization-jackson3/src/main/resources/reference.conf @@ -191,7 +191,11 @@ pekko.serialization.jackson3 { # payload that decompresses to more than this is rejected rather than # allocated, guarding against a small message that inflates without bound. # This applies on deserialization regardless of the `algorithm` setting above. - max-decompressed-size = 256 MiB + # The default of `unlimited` applies no limit, preserving the behaviour of + # earlier releases; a negative number such as -1 also means unlimited. Set a + # size such as `256 MiB` to bound decompression, choosing a value larger than + # any payload the system legitimately exchanges. + max-decompressed-size = unlimited } # Whether the type should be written to the manifest. diff --git a/serialization-jackson3/src/main/scala/org/apache/pekko/serialization/jackson3/JacksonSerializer.scala b/serialization-jackson3/src/main/scala/org/apache/pekko/serialization/jackson3/JacksonSerializer.scala index 6a9db1da2c..3b2adc16ce 100644 --- a/serialization-jackson3/src/main/scala/org/apache/pekko/serialization/jackson3/JacksonSerializer.scala +++ b/serialization-jackson3/src/main/scala/org/apache/pekko/serialization/jackson3/JacksonSerializer.scala @@ -208,7 +208,16 @@ import pekko.util.OptionVal """"off" or "gzip"""") } } - private val maxDecompressedSize: Long = conf.getBytes("compression.max-decompressed-size") + // "unlimited" or a negative number means no limit; getBytes refuses both, so read them first + private val maxDecompressedSize: Long = { + val raw = conf.getString("compression.max-decompressed-size") + if (raw == "unlimited") -1L + else + raw.toLongOption match { + case Some(n) if n < 0 => n + case _ => conf.getBytes("compression.max-decompressed-size") + } + } private val migrations: Map[String, JacksonMigration] = { import scala.jdk.CollectionConverters._ conf.getConfig("migrations").root.unwrapped.asScala.toMap.map { @@ -545,7 +554,10 @@ import pekko.util.OptionVal // meta.length is the decompressed size declared on the wire; a small // message can declare a huge (or negative) size and drive a large // allocation, so bound it before decompressing. - if (meta.length < 0 || meta.length > maxDecompressedSize) + if (meta.length < 0) + throw new IllegalArgumentException( + s"Compressed message declares a negative decompressed size [${meta.length}] bytes") + if (maxDecompressedSize >= 0 && meta.length > maxDecompressedSize) throw new IllegalArgumentException( s"Compressed message declares decompressed size [${meta.length}] bytes, which exceeds the maximum " + s"of [$maxDecompressedSize] bytes (pekko.serialization.jackson3.compression.max-decompressed-size)") @@ -557,7 +569,7 @@ import pekko.util.OptionVal } // gunzip with a bound on the decompressed size, so a small gzip payload cannot - // inflate without limit (a "zip bomb"). + // inflate without limit (a "zip bomb"). A negative maximum applies no bound. private def gunzip(in: GZIPInputStream): Array[Byte] = { val out = new ByteArrayOutputStream() val buffer = new Array[Byte](BufferSize) @@ -565,7 +577,7 @@ import pekko.util.OptionVal var n = in.read(buffer) while (n != -1) { total += n - if (total > maxDecompressedSize) + if (maxDecompressedSize >= 0 && total > maxDecompressedSize) throw new IllegalArgumentException( s"Decompressed message exceeds the maximum of [$maxDecompressedSize] bytes " + "(pekko.serialization.jackson3.compression.max-decompressed-size)") diff --git a/serialization-jackson3/src/test/scala/org/apache/pekko/serialization/jackson3/JacksonSerializerSpec.scala b/serialization-jackson3/src/test/scala/org/apache/pekko/serialization/jackson3/JacksonSerializerSpec.scala index 8338f13c40..d094246f2f 100644 --- a/serialization-jackson3/src/test/scala/org/apache/pekko/serialization/jackson3/JacksonSerializerSpec.scala +++ b/serialization-jackson3/src/test/scala/org/apache/pekko/serialization/jackson3/JacksonSerializerSpec.scala @@ -632,6 +632,42 @@ class JacksonJsonSerializerSpec extends JacksonSerializerSpec("jackson-json") { } ex.getMessage should include("max-decompressed-size") } + + "apply no gzip decompression limit when max-decompressed-size is -1" in withSystem(""" + pekko.serialization.jackson3.jackson-json.compression { + algorithm = gzip + compress-larger-than = 0 KiB + max-decompressed-size = -1 + } + """) { sys => + val msg = SimpleCommand("0" * (8 * 1024)) + JacksonSerializer.isGZipped(serializeToBinary(msg, sys)) should ===(true) + checkSerialization(msg, sys) + } + + "apply no lz4 decompression limit when max-decompressed-size is -1" in withSystem(""" + pekko.serialization.jackson3.jackson-json.compression { + algorithm = lz4 + compress-larger-than = 0 KiB + max-decompressed-size = -1 + } + """) { sys => + val msg = SimpleCommand("0" * (8 * 1024)) + JacksonSerializer.isLZ4(serializeToBinary(msg, sys)) should ===(true) + checkSerialization(msg, sys) + } + + "apply no decompression limit when max-decompressed-size is unlimited" in withSystem(""" + pekko.serialization.jackson3.jackson-json.compression { + algorithm = gzip + compress-larger-than = 0 KiB + max-decompressed-size = unlimited + } + """) { sys => + val msg = SimpleCommand("0" * (8 * 1024)) + JacksonSerializer.isGZipped(serializeToBinary(msg, sys)) should ===(true) + checkSerialization(msg, sys) + } } "JacksonJsonSerializer without type in manifest" should {