From fcead29f8646c7bb0ed4cd8b3269e5ad1fdb2446 Mon Sep 17 00:00:00 2001 From: PJ Fanning Date: Mon, 31 Aug 2026 17:52:37 +0100 Subject: [PATCH 1/2] fix: bound Jackson payload decompression size (#3491) Motivation: JacksonSerializer.decompress inflated gzip payloads with an unbounded transferTo, and passed the lz4 decompressed length declared on the wire straight to the decompressor as the allocation size. A small, well-formed message could therefore declare (or expand to) an arbitrarily large size and drive an OutOfMemoryError on deserialization. Modification: Add a `compression.max-decompressed-size` setting (default 256 MiB) to the jackson and jackson3 modules. On deserialization the gzip path copies through a bounded loop and the lz4 path rejects a declared length that is negative or over the cap, before allocating. Applies regardless of the `algorithm` setting, since decompression is chosen by the payload's magic bytes. Result: A payload that would decompress beyond the cap is rejected with an IllegalArgumentException instead of exhausting the heap. Tests: - sbt "serialization-jackson/testOnly *JacksonJsonSerializerSpec" "serialization-jackson3/testOnly *JacksonJsonSerializerSpec" - 71 passed each, incl. new gzip/lz4 cap tests - sbt "serialization-jackson/mimaReportBinaryIssues" - no issues (changed symbols are @InternalApi/private) - sbt scalafmt for changed main and test sources References: None - found while reviewing the draft threat model in #3478 --- .../src/main/resources/reference.conf | 6 +++ .../jackson/JacksonSerializer.scala | 40 +++++++++++++------ .../jackson/JacksonSerializerSpec.scala | 34 ++++++++++++++++ 3 files changed, 68 insertions(+), 12 deletions(-) diff --git a/serialization-jackson/src/main/resources/reference.conf b/serialization-jackson/src/main/resources/reference.conf index 5db0ac82b3e..5ad52a71100 100644 --- a/serialization-jackson/src/main/resources/reference.conf +++ b/serialization-jackson/src/main/resources/reference.conf @@ -205,6 +205,12 @@ pekko.serialization.jackson { # If compression is enabled with the `algorithm` setting the payload is compressed # when it's larger than this value. compress-larger-than = 0 KiB + + # Maximum size of a payload after decompression. A compressed (gzip or lz4) + # 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 } # 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 ebadef6a7ca..72ce32564d4 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,6 +207,7 @@ import pekko.util.OptionVal """"off" or "gzip"""") } } + private val maxDecompressedSize: Long = conf.getBytes("compression.max-decompressed-size") private val migrations: Map[String, JacksonMigration] = { import pekko.util.ccompat.JavaConverters._ conf.getConfig("migrations").root.unwrapped.asScala.toMap.map { @@ -535,22 +536,18 @@ import pekko.util.OptionVal def decompress(bytes: Array[Byte]): Array[Byte] = { if (isGZipped(bytes)) { val in = new GZIPInputStream(new UnsynchronizedByteArrayInputStream(bytes)) - val out = new ByteArrayOutputStream() - val buffer = new Array[Byte](BufferSize) - - @tailrec def readChunk(): Unit = in.read(buffer) match { - case -1 => () - case n => - out.write(buffer, 0, n) - readChunk() - } - - try readChunk() + try gunzip(in) finally in.close() - out.toByteArray } else { LZ4Meta.get(bytes) match { case OptionVal.Some(meta) => + // 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) + 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)") val srcLen = bytes.length - meta.offset lz4Decompressor.decompress(bytes, meta.offset, srcLen, meta.length) case _ => bytes @@ -558,4 +555,23 @@ import pekko.util.OptionVal } } + // gunzip with a bound on the decompressed size, so a small gzip payload cannot + // inflate without limit (a "zip bomb"). + private def gunzip(in: GZIPInputStream): Array[Byte] = { + val out = new ByteArrayOutputStream() + val buffer = new Array[Byte](BufferSize) + var total = 0L + var n = in.read(buffer) + while (n != -1) { + total += n + if (total > maxDecompressedSize) + throw new IllegalArgumentException( + s"Decompressed message exceeds the maximum of [$maxDecompressedSize] bytes " + + "(pekko.serialization.jackson.compression.max-decompressed-size)") + out.write(buffer, 0, n) + n = in.read(buffer) + } + out.toByteArray + } + } 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 60fab76c931..53d4a2735f0 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 @@ -649,6 +649,40 @@ class JacksonJsonSerializerSpec extends JacksonSerializerSpec("jackson-json") { check(SimpleCommand("Bob"), false) check(new SimpleCommandNotCaseClass("Bob"), false) } + + "reject a gzip payload that decompresses beyond max-decompressed-size" in withSystem(""" + pekko.serialization.jackson.jackson-json.compression { + algorithm = gzip + compress-larger-than = 0 KiB + max-decompressed-size = 1 KiB + } + """) { sys => + val msg = SimpleCommand("0" * (8 * 1024)) + val serializer = serializerFor(msg, sys) + val blob = serializeToBinary(msg, sys) + JacksonSerializer.isGZipped(blob) should ===(true) + val ex = intercept[IllegalArgumentException] { + deserializeFromBinary(blob, serializer.identifier, serializer.manifest(msg), sys) + } + ex.getMessage should include("max-decompressed-size") + } + + "reject an lz4 payload that declares a size beyond max-decompressed-size" in withSystem(""" + pekko.serialization.jackson.jackson-json.compression { + algorithm = lz4 + compress-larger-than = 0 KiB + max-decompressed-size = 1 KiB + } + """) { sys => + val msg = SimpleCommand("0" * (8 * 1024)) + val serializer = serializerFor(msg, sys) + val blob = serializeToBinary(msg, sys) + JacksonSerializer.isLZ4(blob) should ===(true) + val ex = intercept[IllegalArgumentException] { + deserializeFromBinary(blob, serializer.identifier, serializer.manifest(msg), sys) + } + ex.getMessage should include("max-decompressed-size") + } } "JacksonJsonSerializer without type in manifest" should { From 843ef0aaee47a33b7cbc39e8cdae2106fa97692d Mon Sep 17 00:00:00 2001 From: PJ Fanning Date: Thu, 3 Sep 2026 14:31:58 +0100 Subject: [PATCH 2/2] fix: default the Jackson max-decompressed-size to unlimited (#3515) * fix: default the Jackson max-decompressed-size to unlimited Motivation: A bounded default could reject a payload an existing system legitimately exchanges, so a patch release carrying the 256 MiB default from #3491 could break running systems on upgrade. The bound should be opt-in, matching the change made to pekko.serialization.max-decompressed-size in #3502. Modification: Default pekko.serialization.jackson.compression.max-decompressed-size (and the jackson3 equivalent) to -1, meaning no limit and matching the behaviour of releases before #3491. A negative maximum skips the gzip size check and the LZ4 declared-size check; a negative declared LZ4 size is still rejected, since it is malformed regardless of the limit. Config's getBytes refuses negative numbers, so the setting is read as a plain long first and as a memory size only when that is not a negative number. Result: Jackson payload decompression is unbounded by default; configuring a size such as 256 MiB bounds it. Tests: - sbt "serialization-jackson/testOnly org.apache.pekko.serialization.jackson.*" - 122 passed - sbt "serialization-jackson3/testOnly org.apache.pekko.serialization.jackson3.*" - 120 passed - sbt "serialization-jackson/scalafmtCheckAll" "serialization-jackson3/scalafmtCheckAll" - clean - sbt "serialization-jackson/mimaReportBinaryIssues" - no issues References: Refs #3491, Refs #3502 * also accept "unlimited" for max-decompressed-size Motivation: Review on #3515 noted that Pekko is inconsistent about unlimited spellings and an explicit keyword is clearer than a magic number, while the neighbouring read.max-document-length and read.max-token-count settings use -1. Accept both. Modification: The setting reads "unlimited" or any negative number as no limit; the reference.conf default is written as `unlimited`. Applied to both serialization-jackson and serialization-jackson3, with a test each for the keyword. Result: `max-decompressed-size = unlimited` and `= -1` both disable the bound. Tests: - sbt "serialization-jackson/testOnly org.apache.pekko.serialization.jackson.*" - 123 passed - sbt "serialization-jackson3/testOnly org.apache.pekko.serialization.jackson3.*" - 121 passed - sbt "serialization-jackson/scalafmtCheckAll" "serialization-jackson3/scalafmtCheckAll" - clean References: Refs #3515 --- .../src/main/resources/reference.conf | 6 +++- .../jackson/JacksonSerializer.scala | 23 +++++++++--- .../jackson/JacksonSerializerSpec.scala | 36 +++++++++++++++++++ 3 files changed, 60 insertions(+), 5 deletions(-) diff --git a/serialization-jackson/src/main/resources/reference.conf b/serialization-jackson/src/main/resources/reference.conf index 5ad52a71100..593a4ad0aa8 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 72ce32564d4..9530aeb74b1 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,19 @@ 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. Not toLongOption, which Scala 2.12 does not have. + private val maxDecompressedSize: Long = { + val raw = conf.getString("compression.max-decompressed-size") + if (raw == "unlimited") -1L + else + try { + val n = raw.toLong + if (n < 0) n else conf.getBytes("compression.max-decompressed-size") + } catch { + case _: NumberFormatException => conf.getBytes("compression.max-decompressed-size") + } + } private val migrations: Map[String, JacksonMigration] = { import pekko.util.ccompat.JavaConverters._ conf.getConfig("migrations").root.unwrapped.asScala.toMap.map { @@ -544,7 +556,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 +571,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 +579,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 53d4a2735f0..30ed2818f33 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 @@ -683,6 +683,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 {