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-persistence-r2dbc.git


The following commit(s) were added to refs/heads/main by this push:
     new 0378532  more timeout checks (#460)
0378532 is described below

commit 037853274842e2d1bafcf142de3251669c9dfa10
Author: PJ Fanning <[email protected]>
AuthorDate: Mon Jul 27 15:41:56 2026 +0100

    more timeout checks (#460)
---
 .../pekko/persistence/r2dbc/internal/R2dbcExecutor.scala      | 11 +++++++++--
 1 file changed, 9 insertions(+), 2 deletions(-)

diff --git 
a/core/src/main/scala/org/apache/pekko/persistence/r2dbc/internal/R2dbcExecutor.scala
 
b/core/src/main/scala/org/apache/pekko/persistence/r2dbc/internal/R2dbcExecutor.scala
index a2808a0..66d286a 100644
--- 
a/core/src/main/scala/org/apache/pekko/persistence/r2dbc/internal/R2dbcExecutor.scala
+++ 
b/core/src/main/scala/org/apache/pekko/persistence/r2dbc/internal/R2dbcExecutor.scala
@@ -250,6 +250,10 @@ class R2dbcExecutor(
       logPrefix: String)(statement: Connection => Statement, mapRow: Row => 
A): Future[immutable.IndexedSeq[A]] = {
     getConnection(logPrefix).flatMap { connection =>
       val startTime = nanoTime()
+      val timeoutTask = closeCallsExceeding.map { timeout =>
+        system.scheduler.scheduleOnce(timeout, () => 
closeAfterTimeout(connection))
+      }
+
       val mappedRows =
         try {
           val boundStmt = statement(connection)
@@ -262,17 +266,20 @@ class R2dbcExecutor(
 
       mappedRows.failed.foreach { exc =>
         log.debug("{} - Select failed: {}", logPrefix: Any, exc: Any)
-        connection.close().asFutureDone()
+        val done = connection.close().asFutureDone()
+        timeoutTask.foreach { task => done.onComplete(_ => task.cancel()) }
       }
 
       mappedRows.flatMap { r =>
-        connection.close().asFutureDone().map { _ =>
+        val done = connection.close().asFutureDone().map { _ =>
           val durationMicros = durationInMicros(startTime)
           if (durationMicros >= logDbCallsExceedingMicros)
             log.info("{} - Selected [{}] rows in [{}] µs", logPrefix, r.size: 
java.lang.Integer,
               durationMicros: java.lang.Long)
           r
         }
+        timeoutTask.foreach { task => done.onComplete(_ => task.cancel()) }
+        done
       }
 
     }


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

Reply via email to