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

Reply via email to