[
https://issues.apache.org/jira/browse/CAMEL-25350?page=com.atlassian.jira.plugin.system.issuetabpanels:all-tabpanel
]
shashank reassigned CAMEL-25350:
--------------------------------
Assignee: shashank
> 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
> Assignee: shashank
> Priority: Major
>
> {{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)