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