[ 
https://issues.apache.org/jira/browse/CAMEL-25503?page=com.atlassian.jira.plugin.system.issuetabpanels:all-tabpanel
 ]

shashank reassigned CAMEL-25503:
--------------------------------

    Assignee: shashank

> 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
>            Assignee: shashank
>            Priority: Minor
>
> {{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