[ 
https://issues.apache.org/jira/browse/CAMEL-25053?page=com.atlassian.jira.plugin.system.issuetabpanels:all-tabpanel
 ]

Claus Ibsen updated CAMEL-25053:
--------------------------------
    Fix Version/s: 4.23.0

> camel-core - Stream Resequence EIP stops delivering after its route is 
> stopped and started while a message waits for its timeout, and callers 
> waiting for capacity are not released on stop
> -------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------
>
>                 Key: CAMEL-25053
>                 URL: https://issues.apache.org/jira/browse/CAMEL-25053
>             Project: Camel
>          Issue Type: Bug
>          Components: camel-core
>            Reporter: shashank
>            Priority: Minor
>             Fix For: 4.23.0
>
>
> The stream resequencer holds back a message that has a gap in front of it 
> until the missing message arrives or the timeout expires. The very first 
> message the resequencer receives (no message delivered yet) is held the same 
> way. The timeout is a {{TimerTask}} on the engine's {{java.util.Timer}}:
> * {{ResequencerEngine.insert}} 
> ({{core/camel-core-processor/.../resequencer/ResequencerEngine.java:241-281}})
>  calls {{element.schedule(new Timeout(timer, timeout))}} at {{:273}}.
> * {{deliverNext}} ({{:312-345}}) does not deliver the first element while it 
> still has a {{Timeout}} ({{:322}}), and nothing behind it is delivered either.
> When the route stops, {{StreamResequencer.doStop}} 
> ({{StreamResequencer.java:234-238}}) calls {{engine.stop()}}, which calls 
> {{timer.cancel()}} ({{ResequencerEngine.java:114-116}}). That discards the 
> pending {{TimerTask}}s. The elements stay in the engine, because the 
> processor, its engine and its sequence are kept across a route restart, and 
> each one still carries its {{Timeout}} ({{Element.scheduled()}} stays true). 
> On the next start, {{engine.start()}} ({{:106}}) creates a new timer but does 
> not schedule those timeouts again. So an element that was waiting at the stop 
> is never delivered, unless the one message it waits for still arrives (then 
> {{insert}} cancels its timeout). Everything inserted behind it waits as well. 
> In other words, the stall is permanent exactly when the gap is never filled, 
> which is the case the timeout exists for. Only removing and re-adding the 
> route recovers it.
> Once {{capacity}} elements are held, every caller blocks in 
> {{engine.waitUntil}} ({{StreamResequencer.java:254}}, 
> {{ResequencerEngine.java:139-152}}). That wait has no timeout, checks no 
> state, and is not released by {{stop()}}. The consumer threads stay blocked, 
> and so does a graceful stop of the route: it waits for the blocked exchanges 
> until the shutdown timeout, forces the stop, and the threads are still 
> blocked afterwards. This capacity wait was noted as a follow-up in the review 
> of CAMEL-24995 (PR #26851): "a caller blocked on a full stream resequencer 
> could hang during shutdown". The restart stall above was not mentioned there.
> What restarts the resequencer processor: {{ClusteredRoutePolicy}} (the route 
> is stopped when leadership is lost and started when it is regained), the 
> quartz {{ScheduledRoutePolicy}} ({{stopRoute}}/{{startRoute}} on a schedule), 
> and {{stopRoute}}/{{startRoute}} from the Java API, JMX or the controlbus. 
> Policies that only stop the consumer (the throttling route policies, 
> {{HazelcastRoutePolicy}}, the ZooKeeper {{MasterRoutePolicy}}) do not stop 
> the processor and do not trigger this. CAMEL-20435 (the batch resequencer 
> could not be restarted) came from a user restarting routes for active/passive 
> failover, which is the same kind of use.
> h3. Reproduction
> Route:
> {noformat}
> from("direct:in").routeId("r").resequence(header("seq")).stream().timeout(500).deliveryAttemptInterval(100).process(collect)
> {noformat}
> {noformat}
> no restart:     send 1, 3, then 4, 5, 6                                    -> 
> delivered [1, 3, 4, 5, 6]
> restart:        send 1, 3; stopRoute(r) + startRoute(r) before the 500 ms 
> timeout; send 4, 5, 6
>                                                                             
> -> 3 s later: []   8 s later: []
> capacity=3, same restart, then 4 more senders on 4 threads:
>                 6 s later: delivered [], producers returned 1/4
>                 stopRoute(r) returned after 3016 ms (forced, inflight=3); 
> afterwards still 1/4 returned
> capacity=2, 2 and 3 held for a missing 1, a third sender waits for capacity; 
> stopRoute, no restart:
>                 stopRoute returned after 3015 ms (forced); 2 s later the 
> third sender has still not returned
> {noformat}
> With the fix below, the restart case delivers {{[1, 3, 4, 5, 6]}}, all 
> producers of the capacity case return and everything is delivered, and the 
> waiting sender returns when the route stops. The unit tests in the PR cover 
> the same cases.
> A TLA+ model (insert, the timer, the delivery thread, the capacity wait, stop 
> and start) finds the same traces: {{NoDeadTimeoutWhileRunning}} is violated 
> in 4 states (Arrive(1), where 1 waits for its timeout, Stop, Start), 
> {{EventuallyDelivered}} is violated after a restart, and {{WaiterReleased}} 
> is violated after a stop. Without a stop, all properties hold. The fixed 
> model holds {{EventuallyDelivered}}, {{WaiterReleased}}, {{Ordered}} and 
> {{NoDup}}.
> h3. Proposed fix
> In {{ResequencerEngine}}:
> * {{start()}}: after creating the new timer, schedule a timeout again for 
> every element that still has one, under the engine lock. The element then 
> waits a full {{timeout}} from the restart, not the time that was left at the 
> stop.
> * {{stop()}}: mark the engine stopped and release all callers waiting in 
> {{waitUntil}}.
> * {{waitUntil}}: fail with {{RejectedExecutionException}} when the engine is 
> stopped, before and after the wait.
> In {{StreamResequencer}}: set that {{RejectedExecutionException}} on the 
> exchange (as the Delay and Throttle EIPs do when they are stopped). This is 
> visible to callers: a caller that was blocked on capacity when the route 
> stopped now gets this exception, where before it stayed blocked. Also end the 
> {{Delivery}} thread on stop with a flag (not an interrupt), so that quick 
> restarts do not leave old delivery threads running next to the new one.
> 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 (not 
> part of this ticket).
> Affected: long-standing. {{engine.stop()}} has cancelled the timer and 
> {{engine.start()}} has created a new one without rescheduling since at least 
> camel-3.21.0. The capacity wait was a {{Thread.sleep(timeout)}} loop without 
> a stop check up to 4.6 and has been a latch since "Fix busy-wait loops" 
> (12e7338a5c, 4.7).
> Duplicate check (2026-09-27): JIRA text "resequencer" returns 40 issues and 
> "StreamResequencer" returns 10. CAMEL-20435 is the batch resequencer that 
> could not restart at all. CAMEL-1034 and CAMEL-1037 (2008) are messages stuck 
> between JMS queues, a different cause and long fixed. GitHub PRs for 
> "resequencer" and "resequence": only #26851, which names the capacity wait as 
> a follow-up. No open PR touches {{ResequencerEngine}} or 
> {{StreamResequencer}}.
> _Filed with Claude Code on behalf of allthingssecurity._



--
This message was sent by Atlassian Jira
(v8.20.10#820010)

Reply via email to