hanishi opened a new issue, #1620:
URL: https://github.com/apache/pekko-connectors/issues/1620

   # google-cloud-pub-sub-grpc: subscriber correctness follow-up to #1494
   
   ## What's going on
   
   PR #1494 added google-application-default credentials and StreamingPull 
auto-reconnect to the connector. After running the resulting subscriber under 
realistic load we hit three structural correctness gaps:
   
   1. **(introduced in #1494)** `autoExtendAckDeadlines` doesn't actually cover 
messages buffered behind a slow downstream stage, so Pub/Sub redelivers them.
   2. **(pre-existing since 2019, newly observable)** 
`GooglePubSub.subscribe(...)` propagates `maxOutstandingMessages`, 
`maxOutstandingBytes`, and `clientId` from the initial `StreamingPullRequest` 
into every keepalive request, which the server rejects with `INVALID_ARGUMENT`. 
The leak has been in `subscribe` since commit `3f1a4f7f28` (Jan 2019, 
akka.stream.alpakka era), but no caller had reason to set those fields until 
#1494's `RestartSettings` overload made long-lived streaming-pull subscribers 
viable. It's essentially a latent bug surfaced by the new use cases #1494 
enables.
   3. **(introduced in #1494)** `flowControlGate` is structurally identical to 
bug 1: it counts permits only when it pushes downstream, so its limit can't 
bound what the server delivers.
   
   All three are reproducible with stub `SubscriberClient`s, no live Pub/Sub 
needed. See the tests at the bottom.
   
   ---
   
   ## Bug 1: `autoExtendAckDeadlines` doesn't cover messages waiting in the 
buffer
   
   ### What we expected
   
   When you wrap a Pub/Sub source with `autoExtendAckDeadlines`, you'd expect 
every message that has reached your client to get its deadline extended in the 
background. That way Pub/Sub won't redeliver while you're still working on it.
   
   ### What actually happens
   
   Tracking only starts when a message reaches the inside of the operator. The 
current code uses a plain `.map`:
   
   ```scala
   Flow[ReceivedMessage]
     .via(killSwitch.flow)
     .map { msg =>
       tracked.put(msg.ackId, System.nanoTime())
       msg
     }
   ```
   
   A `.map` only runs its function when downstream asks for the next element. 
So in a normal pipeline like this:
   
   ```scala
   GooglePubSub.subscribe(request, 1.second)
     .via(GooglePubSub.autoExtendAckDeadlines(sub, 8.seconds, 30))
     .mapAsync(parallelism = 10)(processMessage)   // ~5 to 7 seconds per 
message
   ```
   
   if all 10 `mapAsync` slots are busy, nothing flows through `.map`. Meanwhile 
Pub/Sub has already pushed dozens more messages over the gRPC stream. Those 
messages sit in the gRPC adapter's buffer with their server-side deadline 
ticking, but our code never tracks them. The background ticker fires every 8 
seconds, looks at `tracked`, sees only the 10 messages currently held by 
`mapAsync`, and extends just those. The rest hit 60 seconds, expire, get 
redelivered. We end up processing duplicates.
   
   ### Why "just put a `.buffer` in front of it" doesn't help
   
   A `.buffer` upstream of `autoExtendAckDeadlines` lets messages queue up, but 
`.map` still won't fire on them until downstream pulls. A `.buffer` downstream 
pushes the bottleneck back by N elements, but once that buffer fills, the bug 
returns.
   
   ### Why eager-pull is the architecturally right answer
   
   #### What goes wrong
   
   Pub/Sub starts a per-message ack deadline timer at server dispatch. The 
client must declare possession of each message (track its ackId) before that 
deadline expires, or the server redelivers. The current `.map`-based tracker 
only declares possession when downstream pulls, and downstream pulls are paced 
by the slowest stage in the pipeline (typically `mapAsync(parallelism)`).
   
   When that pace is slower than the deadline, messages get redelivered. The 
threshold:
   
   ```
   (buffered_count × per_message_seconds / parallelism) > deadline_seconds
   ```
   
   Concretely: 60 messages buffered, 6s per message, `mapAsync(10)`, 60s 
deadline → messages at the back of the queue wait ~36s for tracking, some slip 
past 60s, get redelivered.
   
   Where that buffering lives matters for the fix:
   
   - **`mapConcat(_.receivedMessages.toVector)`** holds each 
`StreamingPullResponse`'s batch (typically 50 to 100 messages) and emits them 
one per downstream pull.
   - **gRPC-Java + Netty** hold messages below the Pekko stream layer entirely: 
HTTP/2 receive window (~1 MB default), `MessageDeframer`'s assembled-message 
queue. Pekko gRPC's own `GraphStage` is demand-driven (it asks gRPC for one 
message per `onPull`), so there's no Pekko-layer buffer to shrink. The buffer 
that lets the server outpace the consumer sits *underneath* Pekko Streams, 
where user code can't reach it.
   
   #### Why this is the deep cause
   
   The bug exists because two distinct timescales get conflated:
   
   - **Tracking** ("I have this message"): must complete in microseconds, 
because it has to beat the ack deadline (tens of seconds, fixed by the server).
   - **Processing** ("I am done with this message"): can take whatever it 
takes; the deadline is extendable via `ModifyAckDeadline` for as long as the 
user wants.
   
   A demand-driven `.map` puts both on the same clock: tracking only fires when 
downstream pulls, which is paced by processing. *That conflation is the bug.* 
The framework's default pull semantics happen to be incompatible with a 
protocol that has a server-controlled, deadlined unacked window. Same shape 
recurs in any at-least-once protocol with these properties (Kafka long-poll 
with manual commit, SQS visibility timeout, RabbitMQ prefetch). gRPC is 
incidental.
   
   #### Why eager-pull is uniquely the right shape
   
   Eager-pull splits the two clocks. A `GraphStage` that pulls from upstream 
into a bounded internal buffer whenever it has room runs the tracking work in 
`onPush` *on receipt*, then enqueues for downstream. Tracking now happens at 
adapter-delivery speed (microseconds); processing still flows through the 
buffer at its own pace.
   
   What this buys that nothing else does:
   
   - **Tracking decouples from processing.** Slow downstream no longer starves 
fast-required tracking.
   - **Backpressure reaches the right buffer.** When the eager-pull buffer 
fills, the stage stops pulling. That signal back-propagates through Pekko 
gRPC's GraphStage (stops calling `call.request(1)`), through gRPC-Java (no 
demand), into Netty (HTTP/2 receive window stops refreshing), and finally to 
the server (stops dispatching). Eager-pull is the only Pekko-layer construct 
that can tell the server "stop" without also breaking tracking. `.buffer(N)`, 
`.async`, larger `parallelism`, or a custom `Sink` for tracking either fail to 
apply backpressure all the way down or break tracking, downstream flow, or the 
abstraction.
   - **Bug 2 complements this at the protocol layer.** Server-side 
`maxOutstandingMessages` bounds the lower buffers via the protocol; eager-pull 
bounds the stream-layer buffer via backpressure. Together: the server has at 
most N messages in flight to the client, and the client tracks every one of 
them on receipt.
   
   Google's Java client uses the same shape 
(`MessageDispatcher.processReceivedMessages` registers in `pendingMessages` 
before the user callback) for the same reason. Independent corroboration of the 
constraint; the constraint itself comes from the protocol.
   
   ### Suggested fix
   
   Replace the `.map` with a small `GraphStage` that pulls from upstream 
eagerly into a bounded buffer, calling the tracking callback the moment it 
grabs each message. Google's official Java client (`google-cloud-pubsub`) does 
the same thing in `MessageDispatcher.processReceivedMessages`: it registers 
messages in `pendingMessages` before handing them to user code.
   
   Sketch:
   
   ```scala
   private[grpc] final class EagerPullTrackingStage[T](maxBuffer: Int, onTrack: 
T => Unit)
       extends GraphStage[FlowShape[T, T]] {
     // preStart: pull(in)
     // onPush:   onTrack(msg); buffer.offer(msg); pushIfPossible(); 
pullIfPossible()
     // onPull:   pushIfPossible(); pullIfPossible()
     // pullIfPossible: only pulls if buffer has room
   }
   ```
   
   The buffer needs an upper bound or it will OOM. With a sensible default 
(1000, matching Google's client) it's fine for almost everyone. New overloads 
can expose `maxBuffer` for tuning. The cleanest pairing is with server-side 
flow control via `maxOutstandingMessages`, which is what bug 2 is about.
   
   ---
   
   ## Bug 2: `subscribe()` makes `maxOutstandingMessages` unusable
   
   ### What we expected
   
   Pub/Sub's streaming pull protocol lets you cap delivery server-side via 
`maxOutstandingMessages` and `maxOutstandingBytes` on the initial 
`StreamingPullRequest`. We expected to set those on the request we pass to 
`GooglePubSub.subscribe()` and have everything just work.
   
   ### What actually happens
   
   `subscribe()` builds the keepalive request like this:
   
   ```scala
   val subsequentRequest = request
     .withSubscription("")
     .withStreamAckDeadlineSeconds(0)
   ```
   
   It clears two fields, but the Pub/Sub API has five fields that must only 
appear on the initial request:
   
   | Field                       | Cleared currently? |
   |-----------------------------|--------------------|
   | `subscription`              | yes                |
   | `streamAckDeadlineSeconds`  | yes                |
   | `clientId`                  | no                 |
   | `maxOutstandingMessages`    | no                 |
   | `maxOutstandingBytes`       | no                 |
   
   The keepalive tick re-sends our original request with three of those fields 
still set. About a second after the stream starts, the server replies:
   
   ```
   INVALID_ARGUMENT: You must only set max_outstanding_messages and
   max_outstanding_bytes on the initial request.
   ```
   
   The whole stream fails. There's no way to recover from inside the public 
API. The only workaround is to bypass `subscribe()` entirely and call 
`subscriberClient.streamingPull` directly.
   
   ### Suggested fix
   
   Build the subsequent request from the proto's default instance:
   
   ```scala
   val subsequentRequest = StreamingPullRequest.defaultInstance
   ```
   
   That clears every initial-only field at once. It's also safe against future 
proto revisions: if Google adds another initial-only field next year, this code 
won't accidentally leak it.
   
   The keepalive tick is purely a heartbeat in our implementation. It doesn't 
carry anything that needs to come from the original request.
   
   ---
   
   ## Bug 3: `flowControlGate` doesn't actually flow-control
   
   ### What we expected
   
   `flowControlGate(FlowControl(maxOutstandingMessages = 100))` looks like the 
protocol-level cap that Google's `FlowController` provides. We expected it to 
bound how many unacked messages the system holds at any given time, including 
ones still in flight from the server.
   
   ### What actually happens
   
   It only counts messages that have already been pushed past the gate to 
downstream. Looking at `impl/FlowControlGateStage.scala`:
   
   ```scala
   override def onPull(): Unit = {
     if (flowControl.outstanding.get() < flowControl.maxOutstandingMessages) {
       pull(in)
     } else {
       downstreamWaiting = true
     }
   }
   
   override def onPush(): Unit = {
     val msg = grab(in)
     flowControl.acquire()
     push(out, msg)
   }
   ```
   
   Same shape as the old `.map` tracker from bug 1: a 1-in-1-out passthrough 
whose `onPull` only pulls upstream when downstream asks. So when `mapAsync` 
saturates, the gate stops pulling, but messages keep arriving from the server 
into the gRPC adapter buffer with no permit cost. The `outstanding` counter 
stays low even when the system is holding hundreds of unacked messages.
   
   In practice this means three things go wrong:
   
   1. **It doesn't mitigate bug 1.** Wrapping the pipeline with 
`flowControlGate(FlowControl(100))` doesn't reduce the redelivery rate, because 
the gate isn't actually bounding what the server sends.
   2. **`outstandingCount` underreports.** Telemetry built on this counter 
shows you a tiny fraction of what's really in flight.
   3. **The name is misleading.** Users reach for "flow control" expecting 
Google's semantics and get downstream credit accounting instead.
   
   ### Suggested fix
   
   Rewrite `FlowControlGateStage` with the same eager-pull pattern as bug 1: 
acquire the permit on receipt (in `onPush`), buffer the message internally, 
push to downstream when demand arrives. The buffer is naturally bounded by the 
FlowControl limit, so no separate `maxBuffer` parameter is needed.
   
   ```scala
   override def onPush(): Unit = {
     val msg = grab(in)
     flowControl.acquire()       // count on receipt, not on push downstream
     buffer.offer(msg)
     pushIfPossible()
     pullIfPossible()
   }
   
   override def onPull(): Unit = {
     pushIfPossible()
     pullIfPossible()
   }
   
   private def pullIfPossible(): Unit =
     if (flowControl.outstanding.get() < flowControl.maxOutstandingMessages
         && !hasBeenPulled(in) && !isClosed(in)) pull(in)
     else if (buffer.isEmpty && isClosed(in)) completeStage()
   ```
   
   After this change, `outstandingCount` reflects messages that have been 
received from the server (whether or not they have flowed past the gate yet), 
the gate stops pulling once the limit is hit, and combined with gRPC 
backpressure the server eventually stops sending too. That's what most users 
reaching for "flow control" actually want.
   
   ### Note on the relationship with bug 2
   
   After bug 2 lands, `StreamingPullRequest.maxOutstandingMessages` gives you 
protocol-level flow control on the wire, with the server respecting the cap 
directly. That is strictly better than `flowControlGate` for the common case of 
"don't let the server send me more than N at a time." `flowControlGate` then 
survives as an application-level primitive for cases that genuinely need it, 
like per-tenant fairness across multiple subscribers sharing one Pub/Sub 
subscription, where server-side flow control can't tell callers apart.
   
   ---
   
   ## Why these three bugs go together
   
   Bug 1 and bug 3 are the same defect in two different operators: a 1-in-1-out 
passthrough whose `onPull` only pulls upstream when downstream demands. Both 
need the same eager-pull fix. Bug 2 is what unlocks the natural way to bound 
the eager-pull buffer for either: server-side 
`StreamingPullRequest.maxOutstandingMessages`, which makes the gRPC adapter 
stop pulling at the protocol level once the cap is reached.
   
   Without bug 2 fixed, you can't use the server-side cap, so the eager-pull 
buffers in bugs 1 and 3 have to be sized by hand. Without bug 1 or bug 3 fixed, 
even a correctly capped server-side flow doesn't help, because messages still 
pile up past the demand-bound stages the moment any downstream stage stalls.
   
   Fixing all three gives you: the server stops sending past the cap (bug 2), 
the eager-pull tracker covers every in-flight message for deadline extension 
(bug 1), the eager-pull gate enforces the in-flight limit on receipt (bug 3), 
and the keepalive tick stays valid forever (bug 2). Each fix on its own helps a 
little. All three together make the subscriber actually work under load.
   
   ---
   
   ## Reproduction
   
   All three tests use stub `SubscriberClient`s, no live Pub/Sub. They live in 
`src/test/scala/org/apache/pekko/stream/connectors/googlecloud/pubsub/grpc/AutoExtendAckDeadlinesSpec.scala`
 (despite the file name, it covers all three bugs since each is a small focused 
test that fits naturally beside the existing autoExtend tests).
   
   The stubs are defined at the bottom of that same file. Each subclasses a 
`TestSubscriberClientBase` trait that overrides every `SubscriberClient` method 
with `???`, so the concrete stubs only override the one or two methods the test 
actually exercises:
   
   - `SucceedingClient`: `modifyAckDeadline` returns 
`Future.successful(Empty())`. Used by bug 3.
   - `CapturingClient`: `modifyAckDeadline` records the incoming request into a 
`ConcurrentLinkedQueue` then returns success. Used by bug 1 to inspect which 
ackIds the ticker extended.
   - `CapturingStreamingPullClient`: overrides `streamingPull` to drain the 
input request `Source` into a queue and return a never-completing 
`Source.maybe[StreamingPullResponse]`, so the response stream stays open while 
the test inspects what was sent on the client → server side. Used by bug 2.
   
   ### Bug 1
   
   ```scala
   "track all in-flight messages eagerly even when downstream is backpressured" 
in {
     val captured = new ConcurrentLinkedQueue[ModifyAckDeadlineRequest]()
     val testSubscriber = new GrpcSubscriber(new CapturingClient(captured))
     val gate = Promise[ReceivedMessage]()  // mapAsync hangs forever
     val killSwitch = KillSwitches.shared("test")
   
     Source(1 to 10).map(i => makeMsg(i.toString))
       .via(GooglePubSub.autoExtendAckDeadlines(subscription, 200.millis, 30))
       .via(killSwitch.flow)
       .mapAsync(1)(_ => gate.future)
       .withAttributes(PubSubAttributes.subscriber(testSubscriber))
       .runWith(Sink.ignore)
   
     Thread.sleep(800)
     killSwitch.shutdown()
   
     val firstReq = captured.poll()
     // Today: firstReq.ackIds.size == 1     (only the message held by mapAsync)
     // Want:  firstReq.ackIds.size == 10    (everything we received)
     firstReq.ackIds.toSet shouldBe (1 to 10).map(_.toString).toSet
   }
   ```
   
   ### Bug 2
   
   ```scala
   "send only allowed fields on subsequent StreamingPullRequest messages" in {
     val captured = new ConcurrentLinkedQueue[StreamingPullRequest]()
     val testSubscriber = new GrpcSubscriber(new 
CapturingStreamingPullClient(captured)(system))
   
     val initial = StreamingPullRequest()
       .withSubscription(subscription)
       .withStreamAckDeadlineSeconds(60)
       .withClientId("test-client")
       .withMaxOutstandingMessages(100L)
       .withMaxOutstandingBytes(10485760L)
   
     val cancellable = GooglePubSub.subscribe(initial, 100.millis)
       .withAttributes(PubSubAttributes.subscriber(testSubscriber))
       .toMat(Sink.ignore)(Keep.left).run()
   
     Thread.sleep(500); cancellable.futureValue.cancel()
   
     val first = captured.poll()
     first.maxOutstandingMessages shouldBe 100L  // initial request kept 
verbatim
   
     val subsequent = 
Iterator.continually(Option(captured.poll())).takeWhile(_.isDefined).flatten.toList
     subsequent should not be empty
     // Today this fails: every subsequent request still carries clientId,
     // maxOutstandingMessages, maxOutstandingBytes from the initial request.
     subsequent.foreach { req =>
       req.subscription shouldBe ""
       req.streamAckDeadlineSeconds shouldBe 0
       req.clientId shouldBe ""
       req.maxOutstandingMessages shouldBe 0L
       req.maxOutstandingBytes shouldBe 0L
     }
   }
   ```
   
   ### Bug 3
   
   ```scala
   "acquire permits on receipt, not on push to downstream (eager-pull)" in {
     val testSubscriber = new GrpcSubscriber(new SucceedingClient())
     val limit = 5
     val flowControl = FlowControl(maxOutstandingMessages = limit.toLong)
     val gate = Promise[ReceivedMessage]() // mapAsync(1) hangs forever
     val killSwitch = KillSwitches.shared("flowControlGateTest")
   
     Source(1 to 100).map(i => makeMsg(i.toString))
       .via(GooglePubSub.flowControlGate(flowControl))
       .via(killSwitch.flow)
       .mapAsync(1)(_ => gate.future)
       .withAttributes(PubSubAttributes.subscriber(testSubscriber))
       .runWith(Sink.ignore)
   
     Thread.sleep(300)
   
     // Today: outstandingCount == 1 (only the message held by mapAsync got a 
permit acquired).
     // Want:  outstandingCount == 5 (the gate eagerly pulls and counts up to 
the limit).
     flowControl.outstandingCount shouldBe limit.toLong
   
     killSwitch.shutdown()
   }
   ```
   
   ---
   
   ## Recommended high-level API: `Subscriber`
   
   Once all three bugs are fixed at the operator level, the connector still has 
an ergonomic problem. To get a correct subscriber today the user has to compose 
four operators by hand:
   
   ```scala
   GooglePubSub.subscribe(request, 1.second, restartSettings)
     .via(GooglePubSub.autoExtendAckDeadlines(...))
     .via(GooglePubSub.flowControlGate(flowControl))
     .mapAsync(10)(processMessage)
     .map(...)
     .runWith(GooglePubSub.acknowledge(parallelism = 1, flowControl))
   ```
   
   There are several ways to compose these wrong: putting 
`autoExtendAckDeadlines` inside `RestartSource` instead of the 3-arg 
`subscribe` (loses tracking on each reconnect), forgetting to pair 
`acknowledge(_, flowControl)` with `flowControlGate(flowControl)` (permits 
never released), passing the wrong `subscription` string to 
`autoExtendAckDeadlines` and the request (extension goes to a different 
subscription), etc. The traps are silent.
   
   A `GooglePubSub.subscriber(...)` resource bundles the configuration once and 
exposes a clean `Source` and `Sink` that compose correctly by construction:
   
   ```scala
   val subscriber = GooglePubSub.subscriber(
     request          = StreamingPullRequest()
                          .withSubscription(subscriptionFqrs)
                          .withStreamAckDeadlineSeconds(60)
                          .withMaxOutstandingMessages(1000),
     pollInterval     = 1.second,
     ackDeadline      = AckDeadline.Fixed(extensionInterval = 8.seconds, 
deadlineSeconds = 30),
     restartSettings  = Some(RestartSettings(100.millis, 10.seconds, 0.2)),
     flowControl      = Some(FlowControl(maxOutstandingMessages = 1000)))
   
   subscriber.source
     .mapAsync(10)(processMessage)
     .map(msg => AcknowledgeRequest(subscriptionFqrs, Seq(msg.ackId)))
     .runWith(subscriber.acknowledge(parallelism = 1))
   ```
   
   Or even simpler, with the `run` convenience:
   
   ```scala
   subscriber.run(parallelism = 10)(processMessage)
   ```
   
   What `Subscriber` guarantees:
   
   - Restart wraps only the inner `subscribe`, not the deadline tracker, so 
tracking survives reconnect by construction. Users can't accidentally put 
`autoExtendAckDeadlines` inside their own `RestartSource` and lose state on 
every reconnect.
   - **In-flight messages are preserved across reconnect.** `RestartSource` 
re-materializes only `subscribe`; the eager-pull tracker, flow-control gate, 
and downstream operators stay alive. Messages buffered inside those operators 
at the moment a gRPC stream dies stay in their buffers, keep getting their 
deadlines extended by the long-lived ticker, and continue flowing downstream 
when the new stream comes up. This is actually better than Google's 
`MessageDispatcher`, which tracks ackIds for extension but doesn't buffer the 
message payload itself; here the payload survives in the stream and processing 
resumes on the original instance.
   - `subscriber.acknowledge(...)` releases flow-control permits if 
`flowControl` is set, and records completion latencies if `ackDeadline` is 
`Adaptive`. No risk of mis-pairing the operators.
   - The deadline tracker's ticker starts at construction and stops at 
`close()`. The source's `watchTermination` triggers `close()` automatically on 
stream completion.
   - Two strategies for ack deadlines via `AckDeadline.Fixed` or 
`AckDeadline.Adaptive`. Other config knobs stay optional.
   
   The low-level operators stay available for users who genuinely need 
different composition, like sharing one extender across multiple subscriptions 
or splitting the source for fan-out. For the 95% case, `Subscriber` is the 
entry point.
   
   `AckDeadlineExtender` survives as the underlying primitive that `Subscriber` 
is built on; users who want fine control over its lifecycle can keep using it 
directly.
   
   A Java DSL equivalent of `Subscriber` is straightforward to add but has been 
left for a follow-up issue. The Java DSL splits its `GrpcSubscriber` from the 
Scala one (each wraps a different generated client class), so a clean Java 
`Subscriber` either needs a `JavaSubscriberClient`-based extender or has to 
bridge to the Scala client. That's a small design call worth its own 
discussion. In the meantime, Java users get the bug 1, 2, and 3 fixes 
automatically through the existing Java DSL operators, and they can use the 
Java factories on `AckDeadlineExtender.create(...)` for restart-safe deadline 
tracking.
   
   # Visual summary: what was wrong and what this PR introduces
   
   ## The structural pattern shared by bugs 1 and 3
   
   ```
   ═══════════════════════════════════ BEFORE: demand-bound 1-in-1-out 
═════════════════════════════════════
   
     upstream                             this stage                          
downstream
     (gRPC adapter)                       (.map tracker / flowControlGate)    
(mapAsync, etc.)
     ────────────                         ──────────────                      
──────────────────
                                                                              
                                           onPull():                          
                                             if (room) pull(in)               
◄── pull
                                                                              
                                           onPush():                          
                                             grab(in)                         
                                             do work (track / acquire)        
                                             push(out, msg)                   
──► msg
                                                                              
                                           (work runs ONLY when downstream    
                                            pulls; nothing eager.)            
   
     When downstream stalls (e.g. mapAsync(parallelism=10) saturated):
       1. downstream stops pulling
       2. this stage's onPull() never fires        ───►  no upstream pull
       3. messages pile up in upstream buffers     ───►  but work never runs
       4. server-side deadline ticks anyway        ───►  REDELIVERY (bug 1)
                                                   ───►  outstanding 
undercounted (bug 3)
   
   
   ═══════════════════════════════════ AFTER: eager-pull with bounded buffer 
═══════════════════════════════
   
     upstream                             this stage                          
downstream
     (gRPC adapter)                       (EagerPullTrackingStage /           
(mapAsync, etc.)
                                           rewritten FlowControlGateStage)
     ────────────                         ──────────────                      
──────────────────
                                                                              
                                           preStart():                        
                                             pull(in)              ◄────── 
starts pulling on its own
                                                                              
     push msg ─────────►  drain  ─────►   onPush():                          
                                             grab(in)                         
                                             do work IMMEDIATELY              
◄── tracked / permit acquired
                                               on receipt                       
  on receipt
                                             buffer.offer(msg)                
                                             pushIfPossible()      ──────►   
──► msg (if downstream demand)
                                             pullIfPossible()                 
                                                                              
                                           onPull():                          
                                             pushIfPossible()                 
──► msg (drain buffer)
                                             pullIfPossible()                 
                                                                              
                                           pullIfPossible():                  
                                             if (buffer < cap                  
                                                 && upstream open) pull(in)   
   
     When downstream stalls:
       1. downstream stops pulling
       2. buffer fills up to its cap        ───►  backpressure flows upstream
       3. but every message we ever          ───►  tracked / permit acquired
          received has already had              correctly (no redelivery,
          its work done on receipt              outstanding reflects reality)
   
   
   ═══════════════════════════════════ Where the buffering actually lives 
═══════════════════════════════════
   
     The "upstream buffers" in the diagram above aren't a single thing, and 
most of
     them are not in Pekko Streams at all. Concretely:
   
         Server
           │ HTTP/2
           ▼
      ┌──────────────────────────────────────────────────────────┐
      │ Netty + gRPC-Java                                        │  ◄── BELOW 
Pekko.
      │   HTTP/2 receive window (~1 MB by default)               │      Holds 
bytes /
      │   MessageDeframer (assembled msgs, not yet delivered)    │      
assembled msgs.
      │   AsyncCallback queue (gRPC thread → GraphStage thread)  │      Not 
reachable
      └──────────────────────────────────────────────────────────┘      from 
user code.
           │   ▲
           │   │ call.request(1) per onPull
           ▼   │
      ┌──────────────────────────────────────────────────────────┐
      │ PekkoNettyGrpcClientGraphStage                           │  ◄── 
DEMAND-DRIVEN.
      │   onPull → call.request(1)                               │      At most 
1
      │   onMessage callback → emit(out, msg); requested -= 1    │      element 
held.
      └──────────────────────────────────────────────────────────┘
           │  one StreamingPullResponse (a BATCH of N msgs)
           ▼
      ┌──────────────────────────────────────────────────────────┐
      │ mapConcat(_.receivedMessages.toVector)                   │  ◄── DOMINANT
      │   holds the batch (typically 50 to 100 msgs)             │      
Pekko-visible
      │   emits one per downstream pull                          │      buffer.
      └──────────────────────────────────────────────────────────┘
           │  one ReceivedMessage per pull
           ▼
      ┌──────────────────────────────────────────────────────────┐
      │ .map { tracked.put(ackId, ...); msg }     ◄── BUG 1      │  ◄── 
TRACKING.
      │   demand-driven 1-in-1-out                               │      Fires 
only
      │   onPush only fires when downstream pulls                │      on 
downstream
      └──────────────────────────────────────────────────────────┘      demand.
           │
           ▼
      ┌──────────────────────────────────────────────────────────┐
      │ mapAsync(parallelism = 10)                               │  ◄── PACES 
THE
      │   slow processing (5 to 7s per message)                  │      ENTIRE 
STACK
      │   pulls only when a slot frees                           │      ABOVE.
      └──────────────────────────────────────────────────────────┘
   
     Why this matters for the fix:
   
     1. There is no Pekko-layer buffer to "make smaller" between Pekko gRPC and
        mapConcat. Tweaking .buffer(N) or async boundary defaults won't help.
   
     2. The buffer responsible for letting the server outpace the consumer lives
        BELOW Pekko entirely (Netty + gRPC-Java). The only way to bound it from
        stream code is by withholding upstream demand: when the eager-pull stage
        stops calling pull(in), Pekko gRPC stops calling call.request(1), 
gRPC-Java
        has no demand to satisfy, Netty's HTTP/2 window stops refreshing, and 
the
        server stops dispatching.
   
     3. Eager-pull is the only Pekko-layer construct that reaches all the way
        down to HTTP/2 without breaking tracking. That's why it's the right 
shape;
        other patterns (.buffer, .async, custom Sink) either fail to 
backpressure
        the lower buffers or break some other property.
   
     4. Bug 2's fix is the protocol-layer complement. Server-side
        maxOutstandingMessages caps the gRPC-Java + Netty buffers via the 
protocol
        itself, so backpressure has less work to do.
   
   
   ═══════════════════════════════════ Bug 2 (the keepalive request leak) 
══════════════════════════════════
   
     Initial request  ───►   server  ──► OK, stream opens
                                                                 
     Pre-fix subsequentRequest = 
request.withSubscription("").withStreamAckDeadlineSeconds(0)
                                             │
                                             │ still carries:
                                             ├─ clientId = "client-A"
                                             ├─ maxOutstandingMessages = 100   
◄── server forbids these
                                             └─ maxOutstandingBytes    = 10MB   
  on subsequent requests
                                             
     Tick #2 (1s later) ───►  server  ──► ❌ INVALID_ARGUMENT
                                             stream fails
   
   
     Post-fix subsequentRequest = StreamingPullRequest.defaultInstance
                                             │
                                             ├─ subscription          = ""
                                             ├─ streamAckDeadlineSecs = 0
                                             ├─ clientId              = ""
                                             ├─ maxOutstandingMessages= 0L
                                             └─ maxOutstandingBytes   = 0L
                                             (every initial-only field cleared,
                                              forward-safe against future 
protos)
                                             
     Every tick ─────────►   server  ──► ✓ accepted, stream stays open
   ```
   
   ---
   
   ## Composition before vs after
   
   ```
   ═══════════════════════════════════ BEFORE: hand-composed, easy to get wrong 
═════════════════════════════
   
     GooglePubSub.subscribe(req, 1.s, restartSettings)   ◄─── req carries 
maxOutstandingMessages
                                                               (bug 2 → 
INVALID_ARGUMENT)
       .via(GooglePubSub.autoExtendAckDeadlines(sub, ...)) ◄─── .map tracker 
(bug 1)
       .via(GooglePubSub.flowControlGate(fc))              ◄─── push-time 
permit (bug 3)
       .mapAsync(10)(processMessage)
       .map(_ => AcknowledgeRequest(...))
       .runWith(GooglePubSub.acknowledge(1, fc))         ◄─── must remember to 
pair fc here too
   
     Silent traps:
       - autoExtend inside RestartSource → tracking lost on every reconnect
       - acknowledge without fc → permits never released → gate locks forever
       - subscription mismatch between req and autoExtend → extends a different 
sub
   
   
   ═══════════════════════════════════ AFTER: one resource, correct by 
construction ═════════════════════════
   
     val subscriber = GooglePubSub.subscriber(
       request          = req w/ maxOutstandingMessages = N,    ◄── bug 2 fix 
automatic
       pollInterval     = 1.s,
       ackDeadline      = AckDeadline.Fixed(8.s, 30),           ◄── extender 
owns ticker
       restartSettings  = Some(...),                            ◄── wraps INNER 
subscribe only
       flowControl      = Some(FlowControl(N)))                 ◄── eager-pull 
gate
   
     subscriber.source                                          ◄── composes:
       .mapAsync(10)(processMessage)                                subscribe → 
autoExtend(extender)
       .map(_ => AcknowledgeRequest(...))                           → 
flowControlGate → watchTermination
       .runWith(subscriber.acknowledge(1))                       ◄── releases 
fc permits + records
                                                                      
completion latencies automatically
   
     By construction:
       + restart only re-materializes inner subscribe
       + extender + flowControl + buffers stay alive across reconnects
       + messages buffered at moment of disconnect KEEP BEING EXTENDED
       + no permit leaks possible (acknowledge knows about fc)
   
   
   ═══════════════════════════════════ Restart-safety detail 
═════════════════════════════════════════════════
   
     Inner gRPC stream fails:
                                 stays alive
                             ────────────────────►
     RestartSource ── re-mat ── EagerPullTracker ── EagerPullFlowGate ── 
mapAsync ── ack
          │                          ▲                    ▲
          │                          │                    │
          └─ subscribe restarts      buffer holds         outstanding still
             new gRPC stream         pre-disconnect       reflects in-flight
             comes up                messages             across the gap
   
     Extender (caller-owned):
          ────────────────────────────────────────────────►
          map<ackId,t>     ticker fires every extensionInterval
          persists         continues extending pre-disconnect ackIds
          across           during the backoff window
          reconnects
   
     This matches Google's MessageDispatcher lifecycle (built once per
     StreamingSubscriberConnection, reused across every gRPC stream restart).
     And it goes further: Pekko Streams' buffering means the message PAYLOAD
     also survives, not just the ackId, so processing resumes on the original
     message instance.
   ```
   
   


-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]


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


Reply via email to