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

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


The following commit(s) were added to refs/heads/1.4.x by this push:
     new c23bad587 DeflateDecompressor: fix inflater memory leak (#1137) (#1140)
c23bad587 is described below

commit c23bad5870b3285bf756b22728edacf48b49ede0
Author: PJ Fanning <[email protected]>
AuthorDate: Wed Jul 8 01:44:44 2026 +0100

    DeflateDecompressor: fix inflater memory leak (#1137) (#1140)
    
    * fix: DeflateDecompressor inflater memory leak
    
    Motivation:
    DeflateDecompressor allocates an Inflater per stream but never explicitly
    releases its native memory. Under GCs with relaxed off-heap reclamation
    (e.g. ZGC), repeated request decompression can accumulate unreclaimed
    native buffers and eventually OOM.
    
    Modification:
    - Add createInflater(noWrap) factory method to DeflateDecompressor to
      allow test injection of a custom Inflater
    - Track the current inflater in a var within createLogic
    - Add idempotent cleanupInflater() that calls inflater.end()
    - Override postStop() to call cleanupInflater(), covering normal
      completion, failure, and cancellation
    - Add focused tests in DeflateSpec verifying the inflater is released
      on successful decode, early cancellation, and truncation
    
    Result:
    The Inflater's native memory is reclaimed promptly when the stage
    terminates, regardless of termination path.
    
    Tests:
    - Not run - sbt/builds skipped per user request to conserve credits
    
    References:
    Refs apache/pekko-http#1134
    
    * Update DeflateSpec.scala
    
    * close existing inflater
    
    * fix: await inflater.end() before asserting in cancellation test
    
    The Sink.ignore future completes when take(1) sends Complete downstream,
    but the Cancel signal that triggers postStop() (and thus end()) is a
    separate actor message dispatched afterwards. This means awaitResult can
    return before end() has been called, yielding endCalls = 0.
    
    Add a CountDownLatch to TrackingInflater that is counted down inside
    end(), and call inflater.awaitEnd() after awaitResult() in the
    cancellation test to eliminate the race.
    
    ---------
    
    Co-authored-by: copilot-swe-agent[bot] 
<[email protected]>
---
 .../pekko/http/scaladsl/coding/DeflateSpec.scala   | 60 ++++++++++++++++++++++
 .../http/scaladsl/coding/DeflateCompressor.scala   | 16 +++++-
 2 files changed, 75 insertions(+), 1 deletion(-)

diff --git 
a/http-tests/src/test/scala/org/apache/pekko/http/scaladsl/coding/DeflateSpec.scala
 
b/http-tests/src/test/scala/org/apache/pekko/http/scaladsl/coding/DeflateSpec.scala
index d7c4d3cf9..8d86bd40f 100644
--- 
a/http-tests/src/test/scala/org/apache/pekko/http/scaladsl/coding/DeflateSpec.scala
+++ 
b/http-tests/src/test/scala/org/apache/pekko/http/scaladsl/coding/DeflateSpec.scala
@@ -16,12 +16,16 @@ package org.apache.pekko.http.scaladsl.coding
 import org.apache.pekko
 import pekko.util.ByteString
 import java.io.{ InputStream, OutputStream }
+import java.util.concurrent.{ CountDownLatch, TimeUnit }
+import java.util.concurrent.atomic.AtomicInteger
 import java.util.zip._
 
 import pekko.http.scaladsl.model.HttpMethods.POST
 import pekko.http.scaladsl.model.{ HttpEntity, HttpRequest }
 import pekko.http.impl.util._
 import pekko.http.scaladsl.model.headers.{ `Content-Encoding`, HttpEncodings }
+import pekko.stream.SystemMaterializer
+import pekko.stream.scaladsl.{ Sink, Source }
 import pekko.testkit._
 import scala.annotation.nowarn
 
@@ -61,6 +65,62 @@ class DeflateSpec extends CoderSpec {
         3.seconds.dilated)
         .awaitResult(3.seconds.dilated) should equal(request)
     }
+    "release the inflater when decoding completes" in {
+      val inflater = new TrackingInflater
+      decodeWith(inflater, streamEncode(smallTextBytes)) should 
readAs(smallText)
+      inflater.endCalls.get() shouldEqual 1
+    }
+    "release the inflater when decoding is cancelled early" in {
+      val inflater = new TrackingInflater
+      val compressed = streamEncode(largeTextBytes)
+
+      Source.single(compressed)
+        .via(decoderWith(inflater).withMaxBytesPerChunk(1).decoderFlow)
+        .take(1)
+        .runWith(Sink.ignore)
+        .awaitResult(3.seconds.dilated)
+
+      // postStop() (which calls end()) is dispatched to the stage actor after 
the
+      // Sink.ignore future completes, so we must wait for end() itself rather 
than
+      // for the stream future to avoid a race.
+      inflater.awaitEnd(3.seconds.dilated)
+      inflater.endCalls.get() shouldEqual 1
+    }
+    "release the inflater when decoding is truncated" in {
+      val inflater = new TrackingInflater
+      // Truncated deflate streams complete without exception (completeStage 
on truncation)
+      decodeWith(inflater, streamEncode(smallTextBytes).dropRight(5))
+      inflater.endCalls.get() shouldEqual 1
+    }
+  }
+
+  private def decodeWith(inflater: TrackingInflater, bytes: ByteString): 
ByteString =
+    
decoderWith(inflater).decode(bytes)(SystemMaterializer(system).materializer).awaitResult(3.seconds.dilated)
+
+  @nowarn("msg=deprecated")
+  private def decoderWith(inflater: TrackingInflater): StreamDecoder =
+    new StreamDecoder {
+      override val encoding = HttpEncodings.deflate
+
+      override def newDecompressorStage(maxBytesPerChunk: Int) =
+        () =>
+          new DeflateDecompressor(maxBytesPerChunk) {
+            override protected[coding] def createInflater(noWrap: Boolean) = 
inflater
+          }
+    }
+
+  private class TrackingInflater extends java.util.zip.Inflater(false) {
+    val endCalls = new AtomicInteger
+    private val endLatch = new CountDownLatch(1)
+
+    def awaitEnd(atMost: FiniteDuration): Unit =
+      endLatch.await(atMost.toMillis, TimeUnit.MILLISECONDS)
+
+    override def end(): Unit = {
+      endCalls.incrementAndGet()
+      endLatch.countDown()
+      super.end()
+    }
   }
 
   private def encodeMessage(request: HttpRequest, compressionLevel: Int, 
noWrap: Boolean): HttpRequest = {
diff --git 
a/http/src/main/scala/org/apache/pekko/http/scaladsl/coding/DeflateCompressor.scala
 
b/http/src/main/scala/org/apache/pekko/http/scaladsl/coding/DeflateCompressor.scala
index d1abe53f8..a6f596744 100644
--- 
a/http/src/main/scala/org/apache/pekko/http/scaladsl/coding/DeflateCompressor.scala
+++ 
b/http/src/main/scala/org/apache/pekko/http/scaladsl/coding/DeflateCompressor.scala
@@ -104,7 +104,19 @@ private[coding] object DeflateCompressor {
 private[coding] class DeflateDecompressor(
     maxBytesPerChunk: Int = Decoder.MaxBytesPerChunkDefault) extends 
DeflateDecompressorBase(maxBytesPerChunk) {
 
+  protected[coding] def createInflater(noWrap: Boolean): Inflater = new 
Inflater(noWrap)
+
   override def createLogic(attr: Attributes) = new ParsingLogic {
+    private var currentInflater: Inflater = null
+    private var inflaterEnded = false
+
+    private def cleanupInflater(): Unit =
+      if (!inflaterEnded && currentInflater != null) {
+        inflaterEnded = true
+        currentInflater.end()
+      }
+
+    override def postStop(): Unit = cleanupInflater()
 
     /** Step that probes if the deflate stream contains a zlib wrapper */
     case object ProbeWrapping extends ParseStep[ByteString] {
@@ -112,6 +124,8 @@ private[coding] class DeflateDecompressor(
 
       override def parse(reader: ByteStringParser.ByteReader): 
ParseResult[ByteString] = {
         val inflater = examineAndBuildInflater(reader.remainingData)
+        cleanupInflater()
+        currentInflater = inflater
         ParseResult(None, new Inflate(inflater, noPostProcessing = true, 
ProbeWrapping))
       }
     }
@@ -132,7 +146,7 @@ private[coding] class DeflateDecompressor(
      */
     private def examineAndBuildInflater(bytes: ByteString): Inflater = {
       val wrapped = (bytes.head & 0x0F) == 0x08
-      new Inflater(!wrapped)
+      createInflater(!wrapped)
     }
 
     startWith(ProbeWrapping)


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

Reply via email to