shashank created CAMEL-25351:
--------------------------------

             Summary: 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


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