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]
