shashank created CAMEL-25350:
--------------------------------
Summary: camel-reactive-streams - exchanges of a batch never
complete when the subscriber cancels while the batch is sent
Key: CAMEL-25350
URL: https://issues.apache.org/jira/browse/CAMEL-25350
Project: Camel
Issue Type: Bug
Components: camel-reactive-streams
Reporter: shashank
{{CamelSubscription.flush()}} (default reactive streams engine) moves up to the
requested number of exchanges from the buffer of a subscription to a local
queue, then calls {{onNext}} for each and stops when the subscription was
cancelled. The exchanges still in the local queue are then dropped: they are
not in the buffer anymore, so {{cancel()}} does not discard them, and their
dispatch callback is never called.
The producer of {{to("reactive-streams:x")}} (also used by
{{CamelReactiveStreamsService.from(uri)}}) completes its exchange only from
that callback, so:
* a synchronous caller (a consumer thread, {{ProducerTemplate.sendBody}})
blocks forever;
* the exchanges stay inflight in the route, and a graceful shutdown waits for
them until its timeout.
It happens whenever a subscriber that has requested several exchanges cancels
while a batch with more than one exchange is being sent: in its {{onNext}} (for
example Reactor's {{next()}}, {{takeWhile}}, a subscriber that has seen enough)
or from another thread (dispose, timeout). A batch holds several exchanges as
soon as the producers publish faster than one {{onNext}} at a time, or
exchanges were buffered while the subscriber had no demand.
h3. Reproduction
{{CancelSubscriptionTest}}: three exchanges are sent to {{direct:numbers ->
reactive-streams:numbers}} while the subscriber has no demand; it then requests
3 and cancels in its first {{onNext}}. On main:
{noformat}
CancelSubscriptionTest.testCancelInOnNextCompletesTheOtherExchanges:56
ConditionTimeout
Exchange 2 never completed ==> expected: <true> but was: <false> within 5
seconds.
{noformat}
The control {{testDeliveredWithoutCancel}} passes. No sleeps (Awaitility on the
buffer size of the subscription).
How often it happens in practice (a measurement, not a committed test): with
{{Flux.from(camel.fromStream("numbers", Integer.class)).next().subscribe()}}
and 8 exchanges sent at once with {{ProducerTemplate.asyncSend}} to
{{direct:numbers -> reactive-streams:numbers}}, 194 of 200 such subscriptions
left exchanges that never completed on main (1037 exchanges in total); with the
fix none.
The defect was found with a TLA+ model of {{CamelSubscription}} (producers'
{{publish}}, {{request}}, {{checkAndFlush}}, the flush task with one step per
locked section and per {{onNext}}, {{cancel}} in {{onNext}} or from the
subscriber's thread): "an exchange that was published and not called back is
held in the buffer, the sending queue, or being delivered or discarded" is
violated in 14 steps (two exchanges buffered, request, the flush takes both,
the subscriber cancels in the first {{onNext}}, the loop breaks). Without a
cancel it holds; with the fix it holds for both cancel modes, with 3 and 4
exchanges.
h3. Proposed fix
{{flush()}} takes the exchanges from its local queue one by one and, when it
stops because the subscription was cancelled, passes the ones not sent to
{{discardBuffer}}, as {{cancel()}} does with the buffer: their exchanges fail
with the same {{IllegalStateException}} ("subscription cancelled") instead of
never completing. No exchange is called back twice (the local queue and the
buffer are disjoint).
With the fix the camel-reactive-streams tests pass (67).
Affected: 4.14.x, 4.18.x and main (same code, GitHub contents API); the loop
dates from the component's first version.
Duplicate check (2026-10-05, again by the review): JIRA component
camel-reactive-streams (11 issues since 2017: CAMEL-23918 flaky tests,
CAMEL-18212, CAMEL-15294, CAMEL-14219, CAMEL-11615, ...), text
"CamelSubscription" (CAMEL-11140 uuid, CAMEL-11615, CAMEL-11124): none about
this. GitHub pull requests "reactive-streams", "CamelSubscription": 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)