allthingssecurity opened a new pull request, #27470:
URL: https://github.com/apache/camel/pull/27470

   # Description
   
   [CAMEL-25350](https://issues.apache.org/jira/browse/CAMEL-25350)
   
   `CamelSubscription.flush()` moves up to the requested number of exchanges 
from the buffer to a local queue, then sends them one by one and stops when the 
subscription was cancelled. The exchanges still in the local queue were 
dropped: they are no longer in the buffer, so `cancel()` does not discard them, 
and their dispatch callback is never called. The producer of 
`to("reactive-streams:x")` completes its exchange only from that callback, so a 
synchronous caller blocks forever and the route keeps the exchanges inflight (a 
graceful shutdown waits for them until its timeout).
   
   It happens whenever a subscriber that requested several exchanges cancels 
while a batch of more than one is being sent, in its `onNext` (Reactor 
`next()`, `takeWhile`, a subscriber that has seen enough) or from another 
thread. With `Flux.from(camel.fromStream("numbers", Integer.class)).next()` and 
8 exchanges sent at once with `asyncSend`, 194 of 200 such subscriptions left 
exchanges that never completed (1037 in total) before this change, none after 
(a measurement, not a committed test).
   
   This change takes the exchanges from the local queue one at a time and, when 
the loop stops because of the cancel, hands the ones not sent to 
`discardBuffer`, as `cancel()` does with the buffer: they fail with the same 
`IllegalStateException` instead of never completing. A request of zero or less 
(rule 3.9 of the specification, the subscription ends with an error) left the 
buffered exchanges behind in the same way; they are now discarded too.
   
   No upgrade-guide note: the exchanges concerned never completed before, they 
now fail with the exception that `cancel()` already uses for the buffered ones.
   
   The defect was found with a TLA+ model of `CamelSubscription` (publish, 
request, checkAndFlush, the flush task with one step per locked section and per 
`onNext`, cancel in `onNext` or from another thread): "a published exchange 
that was 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, cancel in the first `onNext`, the loop breaks). 
It holds without a cancel, and with the fix.
   
   Tests: new `CancelSubscriptionTest` (Awaitility on the buffer size of the 
subscription, no sleeps). Without the main-code change:
   ```
   CancelSubscriptionTest.testCancelInOnNextCompletesTheOtherExchanges:56 » 
ConditionTimeout ... Exchange 2 never completed ==> expected: <true> but was: 
<false> within 5 seconds.
   CancelSubscriptionTest.testNonPositiveRequestCompletesTheBufferedExchanges » 
ConditionTimeout ... A buffered exchange never completed ==> expected: <true> 
but was: <false> within 5 seconds.
   ```
   With the change the camel-reactive-streams tests pass (67).
   
   Not changed: a subscriber whose `onNext` throws (forbidden by rule 2.13) 
still ends the flush task, leaving the rest of its batch and the subscription's 
`sending` flag behind (the `TODO` in `flush()`); that needs a decision on how 
to signal it and is left for a separate issue.
   
   # Target
   
   - [x] I checked that the commit is targeting the correct branch (Camel 4 
uses the `main` branch)
   
   # Tracking
   - [x] If this is a large change, bug fix, or code improvement, I checked 
there is a [JIRA issue](https://issues.apache.org/jira/browse/CAMEL) filed for 
the change (usually before you start working on it).
   
   # Apache Camel coding standards and style
   
   - [x] I checked that each commit in the pull request has a meaningful 
subject line and body.
   - [ ] I have run `mvn clean install -DskipTests` locally from root folder 
and I have committed all auto-generated changes.
     (I built and tested the affected module, including the formatter and 
import-sort plugins. I did not run the full root build.)
   
   # AI-assisted contributions
   
   - [x] If this PR includes AI-generated code, commits have proper 
co-authorship attribution (e.g., `Co-authored-by` trailers) and the PR 
description identifies the AI tool used.
     This PR was prepared with Claude Code (Claude Opus 5.5). The commit 
carries a `Co-Authored-By` trailer.
   
   _Claude Code on behalf of allthingssecurity_
   
   🤖 Generated with [Claude Code](https://claude.com/claude-code)
   


-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]

Reply via email to