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

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


The following commit(s) were added to refs/heads/main by this push:
     new 817f614e3 S3: release disk-buffered chunks when the upload is done 
with them (#1930)
817f614e3 is described below

commit 817f614e39f0d9e2beb753001bbc0c9a6bb9cc34
Author: PJ Fanning <[email protected]>
AuthorDate: Wed Sep 30 14:12:30 2026 +0100

    S3: release disk-buffered chunks when the upload is done with them (#1930)
    
    * S3: release disk-buffered chunks when the upload is done with them
    
    DiskBuffer deleted its temp file only once the chunk had been materialized 
as
    many times as the retry budget allows, so a chunk that uploaded first time -
    two materializations out of eight - was never deleted, and survived on disk
    until the JVM exited via deleteOnExit.
    
    Give Chunk a dispose() that DiskChunk implements as a file delete, and call 
it
    from the RetryFlow decision branch that declines to retry, which is the 
point
    the chunk is definitively finished with. Also delete the temp file in 
postStop
    when the stage never emitted a chunk. The existing countdown stays as the 
net
    for retries being exhausted, which is precisely when it fires.
    
    Co-Authored-By: Claude Opus 5 (1M context) <[email protected]>
    
    * S3: delete disk-buffered temp files in code, not at JVM exit
    
    deleteOnExit only fires on an orderly shutdown, never on SIGKILL, and its
    JVM-global registration is never pruned. Replace it with a registry, held by
    the DiskBuffer for the upload it belongs to, that tracks the temp file of 
every
    chunk it emits.
    
    A chunk releases its own file as soon as the upload is finished with it. Any
    file still registered when the upload stream terminates - a chunk emitted 
and
    then abandoned by a cancelled or failed upload - is deleted by a cleanUp() 
hung
    off watchTermination, so nothing waits for the JVM to exit.
    
    Co-Authored-By: Claude Opus 5 (1M context) <[email protected]>
    
    ---------
    
    Co-authored-by: Claude Opus 5 (1M context) <[email protected]>
---
 .../pekko/stream/connectors/s3/impl/Chunk.scala    | 28 +++++++++++-
 .../stream/connectors/s3/impl/DiskBuffer.scala     | 50 +++++++++++++++++++---
 .../stream/connectors/s3/impl/MemoryBuffer.scala   |  2 +-
 .../pekko/stream/connectors/s3/impl/S3Stream.scala | 18 ++++++--
 .../stream/connectors/s3/impl/DiskBufferSpec.scala | 48 +++++++++++++++++++++
 5 files changed, 135 insertions(+), 11 deletions(-)

diff --git 
a/s3/src/main/scala/org/apache/pekko/stream/connectors/s3/impl/Chunk.scala 
b/s3/src/main/scala/org/apache/pekko/stream/connectors/s3/impl/Chunk.scala
index 237a432b6..5648447df 100644
--- a/s3/src/main/scala/org/apache/pekko/stream/connectors/s3/impl/Chunk.scala
+++ b/s3/src/main/scala/org/apache/pekko/stream/connectors/s3/impl/Chunk.scala
@@ -18,18 +18,44 @@ import pekko.stream.scaladsl.Source
 import pekko.NotUsed
 import pekko.annotation.InternalApi
 import pekko.http.scaladsl.model.{ ContentTypes, HttpEntity, RequestEntity }
+import pekko.stream.FlowShape
+import pekko.stream.stage.GraphStage
 import pekko.util.ByteString
 
+import java.io.File
+
 /**
  * Internal Api
  */
 @InternalApi private[impl] sealed trait Chunk {
   def asEntity(): RequestEntity
   def size: Int
+
+  /**
+   * Releases whatever backs this chunk, once it is known that `data` will not 
be materialized again.
+   * Must not be called while an upload attempt for the chunk may still be 
retried.
+   */
+  def dispose(): Unit = ()
 }
 
-@InternalApi private[impl] final case class DiskChunk(data: Source[ByteString, 
NotUsed], size: Int) extends Chunk {
+@InternalApi private[impl] final class DiskChunk(val data: Source[ByteString, 
NotUsed],
+    val size: Int,
+    file: File,
+    registry: TempFileRegistry) extends Chunk {
   def asEntity(): RequestEntity = 
HttpEntity(ContentTypes.`application/octet-stream`, size, data)
+  override def dispose(): Unit = registry.release(file)
+}
+
+/**
+ * A stage that buffers a chunk, and can release whatever any chunk it emitted 
still holds.
+ */
+@InternalApi private[impl] trait ChunkBuffer extends 
GraphStage[FlowShape[ByteString, Chunk]] {
+
+  /**
+   * Releases the resources of every chunk this buffer emitted that was not 
disposed of individually.
+   * Called once the stream the chunks were emitted into has terminated, 
however it terminated.
+   */
+  def cleanUp(): Unit = ()
 }
 
 @InternalApi private[impl] final case class MemoryChunk(data: ByteString) 
extends Chunk {
diff --git 
a/s3/src/main/scala/org/apache/pekko/stream/connectors/s3/impl/DiskBuffer.scala 
b/s3/src/main/scala/org/apache/pekko/stream/connectors/s3/impl/DiskBuffer.scala
index 75fc824a6..14841584c 100644
--- 
a/s3/src/main/scala/org/apache/pekko/stream/connectors/s3/impl/DiskBuffer.scala
+++ 
b/s3/src/main/scala/org/apache/pekko/stream/connectors/s3/impl/DiskBuffer.scala
@@ -15,8 +15,10 @@ package org.apache.pekko.stream.connectors.s3.impl
 
 import java.io.{ File, FileOutputStream }
 import java.nio.BufferOverflowException
+import java.io.File
 import java.nio.file.Files
 import java.nio.file.Path
+import java.util.concurrent.ConcurrentHashMap
 import java.util.concurrent.atomic.AtomicInteger
 
 import org.apache.pekko
@@ -28,7 +30,6 @@ import pekko.stream.FlowShape
 import pekko.stream.Inlet
 import pekko.stream.Outlet
 import pekko.stream.scaladsl.FileIO
-import pekko.stream.stage.GraphStage
 import pekko.stream.stage.GraphStageLogic
 import pekko.stream.stage.InHandler
 import pekko.stream.stage.OutHandler
@@ -36,6 +37,29 @@ import pekko.util.ByteString
 
 import scala.concurrent.ExecutionContext
 
+/**
+ * Internal Api
+ *
+ * Tracks the temp files of chunks that have been emitted but not yet 
released, so that they can be deleted
+ * as soon as they are no longer needed, and at the latest when the stream 
they were emitted into ends.
+ */
+@InternalApi private[impl] final class TempFileRegistry {
+  private val files = ConcurrentHashMap.newKeySet[File]()
+
+  def register(file: File): Unit = { files.add(file): Unit }
+
+  /** Deletes the file and stops tracking it. Deleting a file that is already 
gone is a no-op. */
+  def release(file: File): Unit = {
+    files.remove(file)
+    file.delete(): Unit
+  }
+
+  def releaseAll(): Unit = {
+    files.forEach { file => file.delete(): Unit }
+    files.clear()
+  }
+}
+
 /**
  * Internal Api
  *
@@ -44,14 +68,20 @@ import scala.concurrent.ExecutionContext
  * The stage waits for the incoming stream to complete. After that, it emits a 
single Chunk item on its output. The Chunk
  * contains a bytestream source that can be materialized multiple times, and 
the total size of the file.
  *
- * @param maxMaterializations Number of expected materializations for the 
completed chunk. After this, the temp file is deleted.
+ * @param maxMaterializations Maximum number of materializations the completed 
chunk may see, which is reached only
+ *                            when every upload retry is used. After this, the 
temp file is deleted. In the ordinary
+ *                            case the chunk is disposed of as soon as the 
upload it belongs to is finished with it.
  * @param maxSize Maximum size on disk to buffer
  */
 @InternalApi private[impl] final class DiskBuffer(maxMaterializations: Int, 
maxSize: Int, tempPath: Option[Path])
-    extends GraphStage[FlowShape[ByteString, Chunk]] {
+    extends ChunkBuffer {
   require(maxMaterializations > 0, "maxMaterializations should be at least 1")
   require(maxSize > 0, "maximumSize should be at least 1")
 
+  private val registry = new TempFileRegistry
+
+  override def cleanUp(): Unit = registry.releaseAll()
+
   val in = Inlet[ByteString]("DiskBuffer.in")
   val out = Outlet[Chunk]("DiskBuffer.out")
   override val shape = FlowShape.of(in, out)
@@ -65,7 +95,6 @@ import scala.concurrent.ExecutionContext
         .map(dir => Files.createTempFile(dir, "s3-buffer-", ".bin"))
         .getOrElse(Files.createTempFile("s3-buffer-", ".bin"))
         .toFile
-      path.deleteOnExit()
       var length = 0
       val pathOut = new FileOutputStream(path)
 
@@ -82,6 +111,8 @@ import scala.concurrent.ExecutionContext
         pull(in)
       }
 
+      private var emitted = false
+
       override def onUpstreamFinish(): Unit = {
         if (isAvailable(out)) emit()
         completeStage()
@@ -92,20 +123,27 @@ import scala.concurrent.ExecutionContext
         try {
           pathOut.close()
         } catch { case x: Throwable => () }
+        finally {
+          // nothing downstream can reference the file if the chunk was never 
emitted
+          if (!emitted) path.delete(): Unit
+        }
 
       private def emit(): Unit = {
         pathOut.close()
+        emitted = true
+
+        registry.register(path)
 
         val deleteCounter = new AtomicInteger(maxMaterializations)
         val src = FileIO.fromPath(path.toPath, 65536).mapMaterializedValue { f 
=>
           if (deleteCounter.decrementAndGet() <= 0)
             f.onComplete { _ =>
-              path.delete()
+              registry.release(path)
 
             }(ExecutionContext.parasitic)
           NotUsed
         }
-        emit(out, DiskChunk(src, length), () => completeStage())
+        emit(out, new DiskChunk(src, length, path, registry), () => 
completeStage())
       }
       setHandlers(in, out, this)
     }
diff --git 
a/s3/src/main/scala/org/apache/pekko/stream/connectors/s3/impl/MemoryBuffer.scala
 
b/s3/src/main/scala/org/apache/pekko/stream/connectors/s3/impl/MemoryBuffer.scala
index 682116f68..3646762f6 100644
--- 
a/s3/src/main/scala/org/apache/pekko/stream/connectors/s3/impl/MemoryBuffer.scala
+++ 
b/s3/src/main/scala/org/apache/pekko/stream/connectors/s3/impl/MemoryBuffer.scala
@@ -29,7 +29,7 @@ import pekko.util.ByteString
  *
  * @param maxSize Maximum size to buffer
  */
-@InternalApi private[impl] final class MemoryBuffer(maxSize: Int) extends 
GraphStage[FlowShape[ByteString, Chunk]] {
+@InternalApi private[impl] final class MemoryBuffer(maxSize: Int) extends 
ChunkBuffer {
   val in = Inlet[ByteString]("MemoryBuffer.in")
   val out = Outlet[Chunk]("MemoryBuffer.out")
   override val shape = FlowShape.of(in, out)
diff --git 
a/s3/src/main/scala/org/apache/pekko/stream/connectors/s3/impl/S3Stream.scala 
b/s3/src/main/scala/org/apache/pekko/stream/connectors/s3/impl/S3Stream.scala
index ce974da44..9c79c7385 100644
--- 
a/s3/src/main/scala/org/apache/pekko/stream/connectors/s3/impl/S3Stream.scala
+++ 
b/s3/src/main/scala/org/apache/pekko/stream/connectors/s3/impl/S3Stream.scala
@@ -1126,7 +1126,8 @@ import scala.util.{ Failure, Success, Try }
       initialUploadState: Option[(String, Int)] = None)(
       parallelism: Int): Flow[ByteString, UploadPartResponse, NotUsed] = {
 
-    def getChunkBuffer(chunkSize: Int, bufferSize: Int, maxRetriesPerChunk: 
Int)(implicit settings: S3Settings) =
+    def getChunkBuffer(chunkSize: Int, bufferSize: Int, maxRetriesPerChunk: 
Int)(
+        implicit settings: S3Settings): ChunkBuffer =
       settings.bufferType match {
         case MemoryBufferType =>
           new MemoryBuffer(bufferSize)
@@ -1185,8 +1186,10 @@ import scala.util.{ Failure, Success, Try }
 
         import conf.multipartUploadSettings.retrySettings._
 
+        val chunkBuffer = getChunkBuffer(chunkSize, chunkBufferSize, 
maxRetries) // creates the chunks
+
         SplitAfterSize(chunkSize, chunkBufferSize)(atLeastOneByteString)
-          .via(getChunkBuffer(chunkSize, chunkBufferSize, maxRetries)) // 
creates the chunks
+          .via(chunkBuffer)
           .mergeSubstreamsWithParallelism(parallelism)
           .filter(_.size > 0)
           .via(atLeastOne)
@@ -1198,8 +1201,11 @@ import scala.util.{ Failure, Success, Try }
               if (isTransientError(r.status)) {
                 r.entity.discardBytes()
                 Some(chunkAndUploadInfo)
-              } else
+              } else {
+                // the request has been sent and will not be retried, so the 
buffered chunk is no longer needed
+                chunkAndUploadInfo._1.dispose()
                 None
+              }
             case (chunkAndUploadInfo, (Failure(_), _)) =>
               // Treat any exception as transient.
               Some(chunkAndUploadInfo)
@@ -1209,6 +1215,12 @@ import scala.util.{ Failure, Success, Try }
               handleChunkResponse(response, upload, index, 
conf.multipartUploadSettings.retrySettings)
           }
           .mergeSubstreamsWithParallelism(parallelism)
+          .watchTermination((_, done) => {
+            // a chunk that was emitted and then abandoned - a cancelled or 
failed upload - is never disposed
+            // of individually, so release whatever is still held once this 
upload has finished either way
+            done.onComplete(_ => 
chunkBuffer.cleanUp())(ExecutionContext.parasitic)
+            NotUsed
+          })
       }
       .mapMaterializedValue(_ => NotUsed)
   }
diff --git 
a/s3/src/test/scala/org/apache/pekko/stream/connectors/s3/impl/DiskBufferSpec.scala
 
b/s3/src/test/scala/org/apache/pekko/stream/connectors/s3/impl/DiskBufferSpec.scala
index 16651bfdb..5f8358708 100644
--- 
a/s3/src/test/scala/org/apache/pekko/stream/connectors/s3/impl/DiskBufferSpec.scala
+++ 
b/s3/src/test/scala/org/apache/pekko/stream/connectors/s3/impl/DiskBufferSpec.scala
@@ -95,4 +95,52 @@ class DiskBufferSpec(_system: ActorSystem)
     }
 
   }
+
+  it should "delete its temp file when the chunk is disposed" in {
+    val tmpDir = Files.createTempDirectory("DiskBufferSpec").toFile()
+    val before = tmpDir.list().size
+    val chunk = Source(Vector(ByteString(1, 2, 3)))
+      .via(new DiskBuffer(8, 200, Some(tmpDir.toPath)))
+      .runWith(Sink.seq)
+      .futureValue
+      .head
+
+    // a successful upload materializes the chunk far fewer times than the 
retry budget allows,
+    // so disposal, not the materialization count, is what releases the file
+    chunk.asInstanceOf[DiskChunk].data.runWith(Sink.ignore).futureValue
+    tmpDir.list().size should be(before + 1)
+
+    chunk.dispose()
+    tmpDir.list().size should be(before)
+  }
+
+  it should "delete the temp files of chunks that were never disposed of when 
it is cleaned up" in {
+    val tmpDir = Files.createTempDirectory("DiskBufferSpec").toFile()
+    val before = tmpDir.list().size
+    val buffer = new DiskBuffer(8, 200, Some(tmpDir.toPath))
+
+    Source(Vector(ByteString(1, 2, 
3))).via(buffer).runWith(Sink.seq).futureValue
+
+    // the chunk was emitted and then abandoned, as it would be by a cancelled 
upload
+    tmpDir.list().size should be(before + 1)
+
+    buffer.cleanUp()
+    tmpDir.list().size should be(before)
+  }
+
+  it should "delete its temp file if it fails before emitting a chunk" in {
+    val tmpDir = Files.createTempDirectory("DiskBufferSpec").toFile()
+    val before = tmpDir.list().size
+
+    Source
+      .failed(new RuntimeException("boom"))
+      .via(new DiskBuffer(8, 200, Some(tmpDir.toPath)))
+      .runWith(Sink.seq)
+      .failed
+      .futureValue shouldBe a[RuntimeException]
+
+    eventually {
+      tmpDir.list().size should be(before)
+    }
+  }
 }


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

Reply via email to