[
https://issues.apache.org/jira/browse/CAMEL-25312?page=com.atlassian.jira.plugin.system.issuetabpanels:all-tabpanel
]
Claus Ibsen reassigned CAMEL-25312:
-----------------------------------
Assignee: shashank
> camel-sjms - a suspended consumer keeps consuming (graceful shutdown,
> suspendRoute, route policies), and a consumer started again receives nothing
> --------------------------------------------------------------------------------------------------------------------------------------------------
>
> Key: CAMEL-25312
> URL: https://issues.apache.org/jira/browse/CAMEL-25312
> Project: Camel
> Issue Type: Bug
> Components: camel-sjms, camel-sjms2
> Reporter: shashank
> Assignee: shashank
> Priority: Major
> Fix For: 4.23.0
>
>
> {{SjmsConsumer}} implements {{Suspendable}} but has no
> {{doSuspend}}/{{doResume}}, so {{suspend()}} only changes its status: the
> listener container keeps receiving and routing messages while the route is
> reported as suspended. Everything that suspends a consumer is affected:
> * the graceful shutdown (default): {{DefaultShutdownStrategy}} suspends a
> {{Suspendable}} consumer instead of stopping it, defers its shutdown and
> waits until no exchange is inflight. The sjms consumer keeps taking new
> messages meanwhile, so with a steady flow the route (or the CamelContext)
> only stops at the shutdown timeout, and the exchanges inflight then are
> forced. This is what CAMEL-9567 fixed in 2016 for the old consumer.
> * {{suspendRoute}} (route controller, JMX) and the route policies that
> suspend the consumer ({{ThrottlingInflightRoutePolicy}},
> {{ThrottlingExceptionRoutePolicy}} with an open circuit, ...): nothing is
> throttled, an open circuit keeps routing (failing) messages.
> Separately, {{SimpleMessageListenerContainer.stopConsumers}} closes the
> consumers and sessions but keeps them in the {{consumers}}/{{sessions}}
> fields, and {{initConsumers}} only creates consumers when {{consumers ==
> null}}. A container that is started again therefore keeps its closed
> consumers and never receives a message: this happens when the same consumer
> instance is stopped and started (the stop/start operations of the managed
> consumer in JMX; a route restart creates a new consumer and is not affected).
> h3. Reproduction
> {{SjmsConsumerRestartTest}} (camel-sjms, embedded Artemis), on main:
> {noformat}
> testSuspendResumeRoute: suspendRoute("suspend"), send a message
> mock://suspend Received message count. Expected: <0> but was: <1>
> testSuspendResumeConsumer: ServiceHelper.suspendService(consumer) as the
> route policies do
> mock://policy Received message count. Expected: <0> but was: <1>
> testGracefulStopTakesNoNewMessage: message A blocks in the route, stopRoute
> (graceful) begins and suspends
> the consumer, the test waits for the connection to stop, then sends message
> B and releases A
> ConditionTimeout: the connection is never stopped within 20 seconds
> (before the waits for the stop thread were added: A message sent after the
> route began to stop must not be
> consumed ==> expected: <1> but was: <2>)
> testStopStartConsumer: the consumer instance stopped and started twice, two
> messages sent
> mock://restart Received message count. Expected: <2> but was: <0>
> {noformat}
> {{SjmsConsumerRoutePolicyTest}} (the route policy suspends the consumer when
> an exchange is done), on main:
> {noformat}
> testCircuitBreakerSuspendsConsumer: a failure opens the circuit
> (ThrottlingExceptionRoutePolicy) on the listener thread
> mock://result Received message count. Expected: <0> but was: <1>
> testSuspendWhileThePolicySuspends: suspendService(consumer) while the
> exchange is in flight, then the policy opens the circuit
> mock://blocking Received message count. Expected: <0> but was: <1>
> testCircuitBreakerOnAnotherThread: the same with the route continued by
> threads(1)
> mock://threads Received message count. Expected: <0> but was: <1>
> {noformat}
> (These three tests and the waits of {{SjmsConsumerRestartTest}} use two
> package-private hooks of the container ({{isSuspendResumePending}},
> {{isConnectionStarted}}) and a getter on the consumer
> ({{getListenerContainer}}), so on main they ran with those three methods
> added.) The control {{testStopStartRoute}}, a route restart, passes;
> {{testGracefulStopAsyncReply}} guards the fix: an {{asyncConsumer}} InOut
> exchange in flight during the graceful stop must still send its reply with
> the consumer's session. No sleeps: latches, {{MockEndpoint}} assert period,
> Awaitility.
> The defect was found with a TLA+ model of the consumer and its listener
> container (route start/stop/suspend/resume of one consumer instance, the
> connection recovery task, each Java statement one step): "a suspended route
> receives nothing" is violated in 5 steps (start, suspend), and "a started
> route receives messages" in 13 steps for start, stop, start (the closed
> consumers are kept). With both parts of the fix all properties hold, also
> with connection failures and recovery during the operations; the model also
> showed that {{doShutdown}} must forget the consumers too (the recovery may
> create consumers after {{stopConsumers}}). That model makes the suspend one
> atomic step on a management thread; a second model of the threads (message
> listener, route policy, management thread, the container's stop/start thread,
> a {{threads()}} thread; the JMS rule that the listener cannot stop its
> connection and that {{Connection.stop()}} waits for it; the consumer lock)
> shows the ignored stop on the listener thread and the two waiting cycles of a
> suspend on the calling thread, and that a single ordered stop/start thread
> has none (TLC deadlock check on; a mutation that starts the connection on the
> calling thread while a stop is queued is caught).
> h3. Proposed fix
> * {{SjmsConsumer.doSuspend}}/{{doResume}} suspend/resume the listener
> container, whose {{doSuspend}} stops the connection ({{Connection.stop()}}:
> no message is delivered any more) and {{doResume}} starts it again. The
> consumers and sessions stay open, as exchanges in flight may still use them
> (an {{asyncConsumer}} sends the reply with the consumer's session when the
> exchange is done); closing them on suspend breaks that, which the guard test
> shows.
> * The connection is stopped and started by a single thread of the container
> ({{SjmsSuspendResume}}), in the order of the suspend and resume calls, and
> suspend/resume do not wait for it. {{Connection.stop()}} waits for the
> message listeners in progress and must not be called by a message listener of
> the connection (JMS 2.0; Artemis throws {{IllegalStateException}}), while a
> route policy suspends the consumer when an exchange is done: for a
> synchronous consumer on the listener thread itself (a suspend on the calling
> thread did nothing there: the exception was ignored and the consumer kept
> receiving), or on a thread the listener waits for (route continued with
> {{threads()}}), and a graceful shutdown or JMX suspend waiting in
> {{Connection.stop()}} would hold the lock of the consumer that the listener's
> route policy needs. Both cycles waited until Artemis gave up after its
> {{onMessageCloseTimeout}} (10 s, "AMQ212002: Timed out after waiting 10000ms
> for handler to complete processing"). {{doStop}} lets a pending stop/start
> finish first.
> * {{SimpleMessageListenerContainer.stopConsumers}} sets {{consumers}} and
> {{sessions}} to null after closing them; {{doShutdown}} does the same after
> closing the connection (the connection recovery may have created consumers
> after {{stopConsumers}}; the model found this).
> Limits: messages the JMS client has prefetched stay with the stopped
> connection until it is started again or closed. A connection factory that
> hands out a shared connection and does not stop it
> ({{JmsPoolConnectionFactory}} of pooled-jms: {{stop()}} is a no-op; Spring's
> {{SingleConnectionFactory}} only stops the shared connection when the last
> user stops it) keeps delivering to a suspended consumer, as today. Stopping
> the connection of one consumer never suspends another route: each listener
> container creates its own connection, and the pooling factories above do not
> stop a shared one.
> Upgrade guide note (next to the other camel-sjms note): a suspended sjms
> route now stops consuming, and the limits above.
> With the fix the new tests and the camel-sjms (132 tests, 5 skipped) and
> camel-sjms2 (23 tests) suites pass; {{MllpTcpServerConsumerTransactionTest}}
> (camel-mllp, uses sjms) is disabled ({{@Disabled}}, 2 skipped).
> Affected: 4.14.x, 4.18.x and main (same code, GitHub contents API), since
> 3.8.0.
> Duplicate check (2026-10-04): JIRA text "sjms" with "suspend" (CAMEL-9567
> only, the old consumer), "sjms" with restart/stopRoute/startRoute (none),
> "SimpleMessageListenerContainer" (CAMEL-22026, thread leak on stop, fixed by
> closing the connection on shutdown; not this), sjms issues since 2024 (10,
> none). GitHub pull requests "sjms suspend" (none); open PR #27173 (batch
> consumer) changes {{SimpleMessageListenerContainer}} in other places.
> _Filed with Claude Code on behalf of allthingssecurity._
--
This message was sent by Atlassian Jira
(v8.20.10#820010)