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)