[ 
https://issues.apache.org/jira/browse/CAMEL-25351?page=com.atlassian.jira.plugin.system.issuetabpanels:comment-tabpanel&focusedCommentId=18123712#comment-18123712
 ] 

Guillaume Nodet commented on CAMEL-25351:
-----------------------------------------

This issue is being investigated by a coding agent (on behalf of gnodet).

Preliminary analysis confirms the bug in ReactiveStreamsConsumer.doStop(): 
shutdownNow drops queued exchanges that were already received from the stream, 
and the dropped items callbacks never fire, causing inflightCount to 
permanently inflate. After a restart, refill() never requests more items 
because the inflated inflight count keeps newRequest below minRequests.

_Note: This comment was generated by an AI coding agent and requires manual 
verification._

> camel-reactive-streams - stopping a consumer drops the items it already 
> received, and after a restart it may never request again
> --------------------------------------------------------------------------------------------------------------------------------
>
>                 Key: CAMEL-25351
>                 URL: https://issues.apache.org/jira/browse/CAMEL-25351
>             Project: Camel
>          Issue Type: Bug
>          Components: camel-reactive-streams
>            Reporter: shashank
>            Assignee: Guillaume Nodet
>            Priority: Major
>
> {{ReactiveStreamsCamelSubscriber}} hands each item it receives from the 
> stream to {{ReactiveStreamsConsumer}}, which queues it on its own thread pool 
> ({{concurrentConsumers}} threads, 1 by default); the subscriber counts the 
> item as inflight until the exchange is done and requests more only while 
> fewer than {{maxInflightExchanges}} (128 by default) are requested or 
> inflight, in batches above the refill watermark 
> ({{exchangesRefillLowWatermark}}, 0.25: at least 96 items).
> {{ReactiveStreamsConsumer.doStop()}} calls {{super.doStop()}} (which stops 
> the route processor), detaches the consumer and calls {{shutdownNow}} on the 
> pool. The consumer is not {{Suspendable}}, so a graceful shutdown (route 
> stop, CamelContext stop, route policies, JMX) stops it right away:
> * the items queued in the pool are dropped without a log; they were already 
> delivered by the publisher and cannot be requested again (lost messages), and 
> the item being routed is interrupted;
> * the dropped items stay counted as inflight in the subscriber for good. When 
> the route is started again the subscriber requests fewer items, and nothing 
> at all once more than 32 items were dropped (128 - dropped < 96): the route 
> is started and never consumes again.
> This applies to the default engine and to camel-reactor and camel-rxjava, 
> which use the same consumer and subscriber.
> h3. Reproduction
> {{ConsumerStopTest}} ({{Flux.range}} into 
> {{reactive-streams:queued?maxInflightExchanges=10}}, the first exchange waits 
> on a latch so 9 are queued; the route is stopped and the latch released once 
> the consumer is stopping). On main:
> {noformat}
> ConsumerStopTest.testStopProcessesTheExchangesTakenFromTheStream:53
>   The exchanges taken from the stream before the stop must be processed ==> 
> expected: <[1, 2, 3, 4, 5, 6, 7, 8, 9, 10]> but was: <[]>
> ConsumerStopTest.testConsumesAgainAfterRestart:64 ConditionTimeout
>   The route must consume again after its restart ==> expected: <true> but 
> was: <false> within 10 seconds.
> ConsumerStopTest.testStopFromTheRouteDoesNotWaitForItself:92 ConditionTimeout
>   (a route policy stops the consumer when the first exchange is done: the 9 
> queued exchanges stay inflight)
> {noformat}
> The control {{testConsumesAgainAfterRestartWithoutQueuedExchanges}} (nothing 
> queued when the route stops) passes on main. No sleeps (latches, Awaitility).
> The defect was found with a TLA+ model of the subscriber (onNext, refill and 
> the callback as locked sections), the consumer's pool and the route 
> stop/start: "every item taken from the stream is routed" is violated in 8 
> steps, and "a restarted, idle consumer whose publisher has items left has 
> requested some" in 19 steps (two items queued, stop drops them, restart: the 
> refill computes 3 - 0 - 2 = 1 < 2). A first fix that only replaced 
> {{shutdownNow}} with {{shutdownGraceful}} still loses the items, as the model 
> and the test showed: {{DefaultConsumer.doStop()}} stops the route processor 
> first, so {{RedeliveryErrorHandler}} rejects the drained exchanges with a 
> {{RejectedExecutionException}}.
> h3. Proposed fix
> {{doStop()}} detaches the consumer first (items arriving later are discarded 
> with a WARN, as before), then shuts the pool down with {{shutdownGraceful}} 
> (the queued exchanges are routed, up to the {{shutdownAwaitTermination}} of 
> the {{ExecutorServiceManager}}, 10 s by default, then the rest is dropped as 
> before), then calls {{super.doStop()}}.
> When the stop is called on a thread of the consumer's own pool, for example 
> by a route policy that stops the consumer when an exchange is done 
> ({{RoutePolicySupport.stopConsumer}}, used by 
> {{ThrottlingInflightRoutePolicy}} and {{ThrottlingExceptionRoutePolicy}}), 
> waiting for the pool would wait for the calling thread itself until the 
> timeout and then drop the queued items. In that case the pool is only shut 
> down: the queued exchanges run after the current one and complete (failed 
> with a {{RejectedExecutionException}} that the consumer logs, or routed if 
> the consumer is started again first), so they are not left inflight.
> Stop time: a stop now waits for the queued exchanges and no longer interrupts 
> the running exchange at once, so a stop with a long-running exchange takes up 
> to {{shutdownAwaitTermination}} before it is interrupted (measured: 10 s with 
> the fix, 1 s before). A route that synchronously stops itself from one of its 
> exchanges waits for itself until that timeout and is then stopped forcibly 
> (before: forcibly at once); as for {{onCompletion().parallelProcessing()}}, 
> it should stop the route asynchronously. Upgrade guide note (new {{=== 
> camel-reactive-streams}} section).
> With the fix the camel-reactive-streams (69), camel-reactor (25) and 
> camel-rxjava (25) tests pass; the model holds for one and two stop/start 
> cycles and with maxInflightExchanges 3 and 4, and, in a review extension 
> where a route policy stops the consumer on its own thread, the stop does not 
> wait for itself and nothing is dropped (the first version of the fix 
> deadlocks there). The model does not cover the timeout path or several 
> {{concurrentConsumers}}.
> Affected: 4.14.x, 4.18.x and main (same code, GitHub contents API).
> Duplicate check (2026-10-05, again by the review): JIRA component 
> camel-reactive-streams (11 issues since 2017, none about stop or refill), 
> text "ReactiveStreamsConsumer" (only CAMEL-10806, the rxjava2 component), 
> "ReactiveStreamsCamelSubscriber", "maxInflightExchanges": none. GitHub pull 
> requests "reactive-streams", "ReactiveStreamsConsumer": none (#24448 and 
> #27012 only fixed flaky tests).
> _Filed with Claude Code on behalf of allthingssecurity._



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

Reply via email to