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

   # Description
   
   [CAMEL-25053](https://issues.apache.org/jira/browse/CAMEL-25053)
   
   The stream resequencer holds back a message that has a gap in front of it 
until the missing message arrives or its timeout expires. The very first 
message it receives is held the same way. The timeout is a `TimerTask` on the 
engine's `java.util.Timer`, and `deliverNext()` does not deliver past an 
element that still has one.
   
   `StreamResequencer.doStop()` calls `ResequencerEngine.stop()`, which cancels 
the timer and with it every pending timeout. The processor, its engine and the 
held elements survive a route restart, and the elements still carry their dead 
`Timeout`. `start()` creates a new timer but does not schedule those timeouts 
again. So after a route stop and start while a message waits for its timeout, 
that message is never delivered unless the gap is filled, and nothing behind it 
is delivered either. Once `capacity` elements are held, every caller blocks in 
`ResequencerEngine.waitUntil()`, which has no timeout and is not released by 
`stop()`: the graceful stop of the route waits for them until the shutdown 
timeout, forces the stop, and the callers are still blocked afterwards.
   
   The capacity wait was noted as a follow-up in the review of CAMEL-24995 
(#26851). The restart stall was not.
   
   What restarts the processor: `ClusteredRoutePolicy` (leadership lost and 
regained), the quartz `ScheduledRoutePolicy`, and `stopRoute`/`startRoute` from 
the API, JMX or the controlbus. Route policies that only stop the consumer do 
not trigger it.
   
   This change:
   - `ResequencerEngine.start()` schedules a new timeout, on the new timer, for 
every element that still had one when the engine was stopped. The element then 
waits a full `timeout` from the restart, not the time that was left at the 
stop. This runs under the engine lock, like `insert()`.
   - `ResequencerEngine.stop()` marks the engine stopped and releases every 
caller waiting in `waitUntil()`. `waitUntil()` throws 
`RejectedExecutionException` when the engine is stopped before the wait, and 
when the wait was ended by a stop, also if the engine has been started again 
before the caller wakes up (a stop counter taken with the wait, not the flag, 
decides this).
   - `ResequencerEngine.insert()` throws the same exception once the engine is 
stopped, before the element is added. So a caller that passed the capacity wait 
just before the stop fails without leaving its element queued. Before, such an 
element was added, then its timeout could not be scheduled on the cancelled 
timer, so the exchange failed but the element was delivered after the restart 
(a duplicate if the caller retried).
   - `StreamResequencer.process()` sets that `RejectedExecutionException` on 
the exchange, as the Delay and Throttle EIPs do when they are stopped, also 
when `ignoreInvalidExchanges` is set (the exchange is not queued, so it must 
not look successful). This is visible to callers: a caller that was blocked on 
capacity when the route stopped now fails with this exception instead of 
staying blocked.
   - `stop()` does not wait for the ready elements to be processed downstream. 
It sets the stopped flag before it takes the engine lock, and `deliverNext()` 
returns `false` once the flag is set, so a `deliver()` in progress returns 
after the element it is sending, and `stop()` waits only for that element. The 
other elements stay held and are delivered after a restart. The timer is 
cancelled and the waiters are released under the engine lock, so a caller 
cannot miss the release, and an `insert()` that has not seen the flag still 
schedules its timeout, which `start()` schedules again. On main `stop()` did 
not take the lock at all, but a delivery in progress then kept sending the 
ready elements to the processor that was being stopped.
   - `stop()` no longer throws a `NullPointerException` when the engine was 
never started (`BaseService` calls `stop()` when a start fails).
   - `StreamResequencer.doStop()` ends the `Delivery` thread with a flag (not 
an interrupt, which could hit an exchange being delivered). Before, each quick 
stop and start left the old delivery thread running next to the new one until a 
stop longer than `deliveryAttemptInterval`. The engine lock already prevented 
duplicate delivery, so this only removes the extra threads.
   
   Not changed: messages held when the route stops stay in memory and are 
delivered after the restart. Delivering or waiting for them on stop is a 
separate question and not part of this PR. A downstream processor that blocks 
forever on the element being delivered at the stop still blocks `stop()`.
   
   Upgrade guide: a short 4.23 entry says that a caller blocked on capacity now 
fails with `RejectedExecutionException` when the route stops, instead of 
staying blocked.
   
   Tests:
   - `ResequencerEngineTest.testTimeoutAfterRestart`: an element waiting for 
its timeout (2 s, so that it cannot expire before the stop even if the test 
pauses), an engine stop and start, then the element must still be delivered, 
followed by the one inserted after the restart.
   - `ResequencerEngineTest.testWaitUntilReleasedOnStop`: a thread waiting in 
`waitUntil()` must get a `RejectedExecutionException` when the engine stops.
   - `ResequencerEngineTest.testStopDoesNotWaitForTheReadyBacklog`: five ready 
elements, and a sender that blocks on each of them (the first until the test 
releases it). `stop()` is called while the first element is being sent, then 
that element is released. `stop()` must return, `deliver()` must return after 
the first element, and the other four must stay queued.
   - `ResequencerEngineTest.testInsertRejectedWhenStopped`: `insert()` on a 
stopped engine throws `RejectedExecutionException` and does not queue the 
element.
   - `ResequencerEngineTest.testStopWithoutStart`: `stop()` on an engine that 
was never started does not throw.
   - New `StreamResequencerRestartTest` (ContextTestSupport):
     - `testTimeoutExpiresAfterRouteRestart`: msg1 is delivered, msg3 waits for 
the missing msg2, the route is stopped and started, then msg4 and msg5 are 
sent. The mock must receive msg3, msg4 and msg5.
     - `testCallerWaitingForCapacityIsReleasedOnStop`: capacity 1, msg2 is 
held, a second sender waits for capacity, and the route is stopped. The sender 
must return with a `RejectedExecutionException`.
     - `testDeliveryThreadEndsOnStop`: a route with a 1 s 
`deliveryAttemptInterval` is stopped and started three times. Only one live 
`Resequencer Delivery` thread of the context may be left, and none after a 
final stop.
   
     No sleeps: `stopRoute` is synchronous, the thread states, thread counts 
and the inflight count are awaited with Awaitility, and the deliveries with the 
mock.
   
   Without the main-code change all eight tests fail:
   ```
   ResequencerEngineTest.testTimeoutAfterRestart                              
ConditionTimeoutException: not fulfilled within 10 seconds
   ResequencerEngineTest.testWaitUntilReleasedOnStop                          
waitUntil is still blocked after the resequencer was stopped
   ResequencerEngineTest.testStopDoesNotWaitForTheReadyBacklog                
java.util.concurrent.TimeoutException (deliver() keeps sending after the stop)
   ResequencerEngineTest.testInsertRejectedWhenStopped                        
Expected java.util.concurrent.RejectedExecutionException to be thrown, but 
nothing was thrown.
   ResequencerEngineTest.testStopWithoutStart                                 
NullPointerException: Cannot invoke "java.util.Timer.cancel()" because 
"this.timer" is null
   StreamResequencerRestartTest.testTimeoutExpiresAfterRouteRestart           
mock://result Received message count. Expected: <3> but was: <0>
   StreamResequencerRestartTest.testCallerWaitingForCapacityIsReleasedOnStop  
java.util.concurrent.TimeoutException
   StreamResequencerRestartTest.testDeliveryThreadEndsOnStop                  
ConditionTimeoutException: not fulfilled within 10 seconds
   ```
   With the change, `*Resequenc*` in camel-core passes: 49 tests, 0 failures (2 
skipped: the existing load tests of `ResequencerEngineTest`, disabled by 
default).
   
   Found with a TLA+ model of the stream resequencer (insert, the timer, the 
delivery thread, the capacity wait, stop and start). For the current code, 
`NoDeadTimeoutWhileRunning`, `EventuallyDelivered` and `WaiterReleased` are 
violated after a stop and start. The fixed model holds them, as well as 
`Ordered` and `NoDup`. I then reproduced the stall and the blocked callers 
against the real classes.
   
   # 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 modules, 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