shashank created CAMEL-25503:
--------------------------------

             Summary: camel-reactive-streams - Make the consumer Suspendable, 
so that a suspend does not stop it and wait for the drain
                 Key: CAMEL-25503
                 URL: https://issues.apache.org/jira/browse/CAMEL-25503
             Project: Camel
          Issue Type: Improvement
          Components: camel-reactive-streams
            Reporter: shashank


{{ReactiveStreamsConsumer}} (camel-reactive-streams, also used by camel-reactor 
and camel-rxjava) does not implement {{Suspendable}}. Camel therefore stops it 
whenever it would suspend it:

* {{suspendRoute}} (route controller, JMX) stops the route: its status is 
{{Stopped}}, not {{Suspended}}, and a {{toStream}} request to it fails with 
{{No consumers attached to the stream}};
* since CAMEL-25351 a stop waits until the items the consumer already took from 
the stream are routed (up to {{shutdownAwaitTermination}}, 10 s by default), so 
a suspend waits for them as well; items the publisher still sends for the 
outstanding demand are discarded with a WARN;
* a route policy that suspends the consumer from one of its own exchanges 
({{RoutePolicySupport.suspendOrStopConsumer}}, as 
{{ThrottlingInflightRoutePolicy}} does on the consumer's thread) stops it from 
its own thread pool: the queued items then fail with a 
{{RejectedExecutionException}} (one WARN each) and are lost.

In the review of the CAMEL-25351 pull request 
(https://github.com/apache/camel/pull/27431), Claus Ibsen suggested a real 
suspend as a follow-up (detach and stop requesting, then resume), as the sjms, 
kafka and cxf consumers do (their suspend was fixed in CAMEL-25312, CAMEL-25313 
and CAMEL-25314).

h3. Proposed change

* The consumer keeps the received items in a FIFO queue that its pool tasks 
take from in order (one task per item, as today), and implements 
{{Suspendable}}.
* While it is suspending or suspended, the tasks leave the items queued and the 
subscriber requests nothing; items sent for the demand requested before the 
suspend are queued too (at most {{maxInflightExchanges}}). The exchanges being 
routed complete normally and the suspend waits for nothing.
* {{doResume}} schedules the queued items and requests again; {{doStop}} 
schedules them before the drain added by CAMEL-25351, so a stop routes them; 
the stop is otherwise unchanged.

Visible effects (upgrade guide and component doc): the route status after 
{{suspendRoute}} is {{Suspended}}; a graceful shutdown suspends the consumer 
before stopping it; a {{toStream}} request to a suspended route waits for the 
resume. With backpressure disabled ({{maxInflightExchanges}} not positive) the 
publisher cannot be paused, so items it sends while the consumer is suspended 
are kept in memory until the resume or stop.

h3. Tests

New {{ConsumerSuspendTest}}: a suspend does not wait for the queued items, 
holds them and requests nothing, and a stop routes them; route controller 
suspend/resume; {{toStream}} to a suspended route; a route policy suspending 
from the consumer's thread and resuming. Four of them fail on main (suspend 
timeout, {{expected: <Suspended> but was: <Stopped>}}, {{No consumers attached 
to the stream}}, consumer not suspended). camel-reactive-streams (75 tests), 
camel-reactor (25) and camel-rxjava (25) pass with the change.

Duplicate check (2026-10-08): JIRA text "reactive-streams suspend" (5 issues: 
CAMEL-25351 about the stop, the others unrelated), "ReactiveStreamsConsumer 
Suspendable" (CAMEL-25351 only), component camel-reactive-streams (13 issues, 
none about suspend). GitHub pull requests "reactive-streams suspend" and 
"ReactiveStreamsConsumer": #27431 and #27434 (both CAMEL-25351, stop only); no 
open pull request touches camel-reactive-streams.

_Filed with Claude Code on behalf of allthingssecurity._



--
This message was sent by Atlassian Jira
(v8.20.10#820010)

Reply via email to