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-grpc.git


The following commit(s) were added to refs/heads/main by this push:
     new 85ec44b0 Update pekkoVersion to 2.0.0-M4 (#856)
85ec44b0 is described below

commit 85ec44b046f79b46c57a25296c48b33f8375df51
Author: PJ Fanning <[email protected]>
AuthorDate: Sat Aug 22 12:01:39 2026 +0100

    Update pekkoVersion to 2.0.0-M4 (#856)
    
    * Update pekkoVersion to 2.0.0-M4
    
    * avoid deprecated method
    
    * more versions to update
---
 plugin-tester-java/build.gradle                       |  2 +-
 plugin-tester-java/pom.xml                            |  2 +-
 plugin-tester-scala/build.gradle                      |  2 +-
 plugin-tester-scala/pom.xml                           |  2 +-
 project/PekkoCoreDependency.scala                     |  2 +-
 .../pekko/grpc/internal/PekkoHttpClientUtils.scala    | 19 +++++++++++++------
 6 files changed, 18 insertions(+), 11 deletions(-)

diff --git a/plugin-tester-java/build.gradle b/plugin-tester-java/build.gradle
index 1af6e07c..bef38aaa 100644
--- a/plugin-tester-java/build.gradle
+++ b/plugin-tester-java/build.gradle
@@ -27,7 +27,7 @@ repositories {
 def scalaFullVersion = "2.13.18"
 def scalaVersion = org.gradle.util.VersionNumber.parse(scalaFullVersion)
 def scalaBinaryVersion = "${scalaVersion.major}.${scalaVersion.minor}"
-def pekkoVersion = "2.0.0-M3"
+def pekkoVersion = "2.0.0-M4"
 def pekkoHttpVersion = "2.0.0-M1"
 
 dependencies {
diff --git a/plugin-tester-java/pom.xml b/plugin-tester-java/pom.xml
index ec5cf938..e4de1afc 100644
--- a/plugin-tester-java/pom.xml
+++ b/plugin-tester-java/pom.xml
@@ -24,7 +24,7 @@
     <maven.compiler.target>17</maven.compiler.target>
     <maven-dependency-plugin.version>3.11.0</maven-dependency-plugin.version>
     <maven-exec-plugin.version>3.6.3</maven-exec-plugin.version>
-    <pekko.version>2.0.0-M3</pekko.version>
+    <pekko.version>2.0.0-M4</pekko.version>
     <pekko.http.version>2.0.0-M1</pekko.http.version>
     <grpc.version>1.83.1</grpc.version> <!-- checked synced by 
VersionSyncCheckPlugin -->
     <project.encoding>UTF-8</project.encoding>
diff --git a/plugin-tester-scala/build.gradle b/plugin-tester-scala/build.gradle
index 07e576f6..be8e310b 100644
--- a/plugin-tester-scala/build.gradle
+++ b/plugin-tester-scala/build.gradle
@@ -22,7 +22,7 @@ repositories {
 def scalaFullVersion = "2.13.18"
 def scalaVersion = org.gradle.util.VersionNumber.parse(scalaFullVersion)
 def scalaBinaryVersion = "${scalaVersion.major}.${scalaVersion.minor}"
-def pekkoVersion = "2.0.0-M3"
+def pekkoVersion = "2.0.0-M4"
 def pekkoHttpVersion = "2.0.0-M1"
 
 dependencies {
diff --git a/plugin-tester-scala/pom.xml b/plugin-tester-scala/pom.xml
index 85ff145c..d30793c4 100644
--- a/plugin-tester-scala/pom.xml
+++ b/plugin-tester-scala/pom.xml
@@ -22,7 +22,7 @@
   <properties>
     <maven.compiler.source>17</maven.compiler.source>
     <maven.compiler.target>17</maven.compiler.target>
-    <pekko.version>2.0.0-M3</pekko.version>
+    <pekko.version>2.0.0-M4</pekko.version>
     <pekko.http.version>2.0.0-M1</pekko.http.version>
     <grpc.version>1.83.1</grpc.version> <!-- checked synced by 
VersionSyncCheckPlugin -->
     <project.encoding>UTF-8</project.encoding>
diff --git a/project/PekkoCoreDependency.scala 
b/project/PekkoCoreDependency.scala
index f0a948fb..ed180d0a 100644
--- a/project/PekkoCoreDependency.scala
+++ b/project/PekkoCoreDependency.scala
@@ -22,5 +22,5 @@ import com.github.pjfanning.pekkobuild.PekkoDependency
 object PekkoCoreDependency extends PekkoDependency {
   override val checkProject: String = "pekko-cluster-sharding-typed"
   override val module: Option[String] = None
-  override val currentVersion: String = "2.0.0-M3"
+  override val currentVersion: String = "2.0.0-M4"
 }
diff --git 
a/runtime/src/main/scala/org/apache/pekko/grpc/internal/PekkoHttpClientUtils.scala
 
b/runtime/src/main/scala/org/apache/pekko/grpc/internal/PekkoHttpClientUtils.scala
index 4cda5ef2..7b1cf9a0 100644
--- 
a/runtime/src/main/scala/org/apache/pekko/grpc/internal/PekkoHttpClientUtils.scala
+++ 
b/runtime/src/main/scala/org/apache/pekko/grpc/internal/PekkoHttpClientUtils.scala
@@ -29,7 +29,7 @@ import pekko.http.scaladsl.{ ClientTransport, 
ConnectionContext, Http }
 import pekko.http.scaladsl.model._
 import pekko.http.scaladsl.model.headers.RawHeader
 import pekko.http.scaladsl.settings.ClientConnectionSettings
-import pekko.stream.{ Materializer, OverflowStrategy, QueueOfferResult }
+import pekko.stream.{ Materializer, QueueOfferResult }
 import pekko.stream.scaladsl.{ Keep, Sink, Source }
 import pekko.util.ByteString
 import io.grpc.{ CallOptions, MethodDescriptor, Status, StatusRuntimeException 
}
@@ -133,7 +133,7 @@ object PekkoHttpClientUtils {
 
     val (queue, doneFuture) =
       Source
-        .queue[HttpRequest](4242, OverflowStrategy.fail)
+        .queue[HttpRequest](4242)
         .via(http2client)
         .toMat(Sink.foreach { res =>
           res.attribute(ResponsePromise.Key).get.promise.trySuccess(res)
@@ -142,11 +142,18 @@ object PekkoHttpClientUtils {
 
     def singleRequest(request: HttpRequest): Future[HttpResponse] = {
       val p = Promise[HttpResponse]()
-      queue.offer(request.addAttribute(ResponsePromise.Key, 
ResponsePromise(p))).foreach {
-        case QueueOfferResult.Enqueued => // promise will be completed by the 
response sink
-        case _                         => p.tryFailure(new 
IllegalStateException("Request queue closed"))
+      queue.offer(request.addAttribute(ResponsePromise.Key, 
ResponsePromise(p))) match {
+        case QueueOfferResult.Enqueued =>
+          p.future
+        case QueueOfferResult.Dropped =>
+          Future.failed(
+            new StatusRuntimeException(
+              Status.RESOURCE_EXHAUSTED.withDescription("Too many concurrent 
requests, buffer is full")))
+        case QueueOfferResult.QueueClosed =>
+          Future.failed(new 
StatusRuntimeException(Status.UNAVAILABLE.withDescription("Request queue 
closed")))
+        case QueueOfferResult.Failure(cause) =>
+          Future.failed(new 
StatusRuntimeException(Status.UNAVAILABLE.withCause(cause)))
       }
-      p.future
     }
 
     implicit def serializerFromMethodDescriptor[I, O](descriptor: 
MethodDescriptor[I, O]): ProtobufSerializer[I] =


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

Reply via email to