allthingssecurity opened a new pull request, #27636:
URL: https://github.com/apache/camel/pull/27636

   # Description
   
   [CAMEL-25503](https://issues.apache.org/jira/browse/CAMEL-25503)
   
   Follow-up to the point @davsclaus raised in the review of #27431 
(CAMEL-25351): `ReactiveStreamsConsumer` is not `Suspendable`, so a suspend 
acts as a stop, which since #27431 waits for the drain. A real suspend (stop 
requesting, resume by requesting again) avoids that wait for suspends from 
route policies and the graceful shutdown.
   
   What a suspend did before this change: `suspendRoute` (route controller, 
JMX) stopped the route (status `Stopped`, not `Suspended`) and waited until the 
items already taken from the stream were routed (up to 
`shutdownAwaitTermination`, 10 s by default); items the publisher still sent 
for the outstanding demand were discarded with a WARN; a `toStream` request to 
the route failed with `No consumers attached to the stream`. A route policy 
that suspends the consumer from one of its own exchanges 
(`RoutePolicySupport.suspendOrStopConsumer`, as `ThrottlingInflightRoutePolicy` 
does on the consumer's thread) stopped it from its own pool, so the queued 
items failed with a `RejectedExecutionException` (one WARN each) and were lost.
   
   This change:
   - the consumer keeps the items it receives in a FIFO queue 
(`ConcurrentLinkedQueue`); each item still gets one pool task, which routes the 
oldest queued item. Ordering with `concurrentConsumers=1` is unchanged;
   - `ReactiveStreamsConsumer implements Suspendable`. While the consumer is 
suspending or suspended, the pool tasks leave the items queued and 
`ReactiveStreamsCamelSubscriber.refill()` requests nothing. Items the publisher 
still sends for the demand requested before the suspend are queued too (at most 
`maxInflightExchanges`, as they count as inflight). The exchanges being routed 
complete normally; `doSuspend` waits for nothing;
   - `doResume` adds one task per queued item and calls `refill()`; a consumer 
suspended while not started (before its start, or after a stop) is started by 
the resume;
   - `doStop` adds one task per queued item before the drain from #27431, so 
the items held by a suspended consumer are routed by the stop. Otherwise the 
stop is unchanged (detach, `shutdownGraceful`, or `shutdown` from its own pool, 
then `super.doStop()`);
   - `doStart` also schedules items left queued by a stop that timed out 
(`shutdownNow` drops the pool tasks), so they are routed after the restart 
instead of staying inflight. This only applies after a forced stop.
   
   Visible effects: the route status after `suspendRoute` is `Suspended`; a 
graceful shutdown (route stop, context stop) now suspends the consumer first, 
waits for the exchanges in the route, then stops it and drains the queued 
items. The same items are routed as before, but the queued ones start only 
after the shutdown strategy's inflight wait (it polls every second), so such a 
stop can take up to about a second longer. With backpressure disabled 
(`maxInflightExchanges` not positive) the publisher cannot be paused, and the 
items it sends while the consumer is suspended are kept in memory until the 
resume or stop (component doc). Cost per item: one queue node and a volatile 
read of the consumer status.
   
   Docs: new "Suspending the consumer" section in the component page (catalog 
copy updated). Upgrade guide: new `=== camel-reactive-streams - suspending a 
consumer` right after the CAMEL-25351 entry.
   
   Tests: new `ConsumerSuspendTest` (latches, Awaitility, `MockEndpoint` assert 
periods, no sleeps). Without the main-code change:
   ```
   ConsumerSuspendTest.testSuspendDoesNotWaitForTheQueuedExchanges:88 » Timeout
   ConsumerSuspendTest.testSuspendAndResumeRoute:119 expected: <Suspended> but 
was: <Stopped>
   ConsumerSuspendTest.testRequestToASuspendedRouteWaitsForTheResume:145 » 
IllegalState No consumers attached to the stream controlled
   ConsumerSuspendTest.testRoutePolicySuspendsAndResumesTheConsumer:165 The 
consumer must be suspended ==> expected: <true> but was: <false>
   ```
   (on main the policy test also logs `Item ... not routed as the consumer is 
stopped` for the queued items). The first test also checks that the suspended 
consumer holds the queued items and requests nothing (with 
`exchangesRefillLowWatermark=1`, which would request one item per exchange 
done), and that a stop routes them. 
`testResumeStartsAConsumerThatWasNotStarted` (`autoStartup=false`, suspend then 
resume) passes on main too; it covers the not-started branch of `doResume`. 
`ConsumerStopTest`'s helper now also accepts a suspended consumer while the 
graceful stop waits for the gated exchange. With the change 
camel-reactive-streams (75 tests), camel-reactor (25) and camel-rxjava (25) 
pass.
   
   The branch merges cleanly with main and with our open PRs that touch the 
4.23 upgrade guide.
   
   # Target
   
   - [x] I checked that the commit is targeting the correct branch (Camel 4 
uses the `main` branch)
   
   # Tracking
   - [x] If this is a large change, bug fix, or code improvement, I checked 
there is a [JIRA issue](https://issues.apache.org/jira/browse/CAMEL) filed for 
the change (usually before you start working on it).
   
   # Apache Camel coding standards and style
   
   - [x] I checked that each commit in the pull request has a meaningful 
subject line and body.
   - [ ] I have run `mvn clean install -DskipTests` locally from root folder 
and I have committed all auto-generated changes.
     (I built and tested the affected modules, including the formatter and 
import-sort plugins. I did not run the full root build.)
   
   # AI-assisted contributions
   
   - [x] If this PR includes AI-generated code, commits have proper 
co-authorship attribution (e.g., `Co-authored-by` trailers) and the PR 
description identifies the AI tool used.
     This PR was prepared with Claude Code (Claude Opus 5.5). The commit 
carries a `Co-Authored-By` trailer.
   
   _Claude Code on behalf of allthingssecurity_
   
   🤖 Generated with [Claude Code](https://claude.com/claude-code)
   


-- 
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]

Reply via email to