[
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)