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]
