[
https://issues.apache.org/jira/browse/CAMEL-25315?page=com.atlassian.jira.plugin.system.issuetabpanels:all-tabpanel
]
Claus Ibsen resolved CAMEL-25315.
---------------------------------
Resolution: Fixed
Merged in https://github.com/apache/camel/pull/27391 (commit e670025a49a4) for
4.23.0.
_Claude Code on behalf of davsclaus_
> camel-zookeeper-master - a master whose consumer failed to start never
> consumes, and a consumer started while the leadership is lost or the route
> stops keeps running
> ---------------------------------------------------------------------------------------------------------------------------------------------------------------------
>
> Key: CAMEL-25315
> URL: https://issues.apache.org/jira/browse/CAMEL-25315
> Project: Camel
> Issue Type: Bug
> Components: camel-zookeeper-master
> Reporter: shashank
> Assignee: shashank
> Priority: Major
> Fix For: 4.23.0
>
>
> {{MasterConsumer}} starts the delegate consumer from the {{CHANGED}} group
> event (the group's own thread, {{onLockOwned}}) and stops it from
> {{DISCONNECTED}} (called directly on the Curator connection-state thread, and
> by {{ZooKeeperGroup.close()}} on the closing thread) and from {{doStop}}.
> Nothing synchronises them, and {{onLockOwned}} only acts {{if (delegate ==
> null)}}:
> {code:java}
> if (delegate == null) {
> try {
> ServiceHelper.startService(endpoint.getConsumerEndpoint());
> delegate = endpoint.getConsumerEndpoint().createConsumer(processor);
> ...
> groupListener.updateState(thisNodeState);
> ServiceHelper.startService(delegate);
> } catch (Exception e) {
> LOG.error("Failed to start master consumer for: {}", endpoint, e);
> }
> }
> {code}
> * *A failed start is never retried.* When the delegate cannot start at the
> moment the node is elected (JMS broker, FTP or database server not reachable
> yet), the failed consumer stays in {{delegate}}, so every later leadership
> event skips the start. The node keeps the ZooKeeper leadership, the other
> nodes stay standby, and nobody consumes until this node loses its connection
> or the route is restarted. The only trace is one ERROR log line.
> * *A consumer started after a disconnect keeps running on a node that is not
> the master.* If the connection is lost ({{SUSPENDED}}/{{LOST}} ->
> {{DISCONNECTED}} -> {{stopConsumer()}}) between the null check and the
> assignment of the new delegate (endpoint start, {{createConsumer}}), the stop
> finds nothing, and the delegate is then assigned and started. Another node
> becomes the master: two nodes consume. When this node reconnects as a standby
> nothing stops it (standby events do nothing).
> * *The same after a route stop:* {{doStop}} stops nothing, and {{close()}}
> only calls {{DISCONNECTED}} again after waiting up to 5 seconds for the group
> thread, so a start that takes longer leaves a consumer running after the
> route (or the CamelContext) stopped.
> camel-master's {{MasterConsumer}} handles these cases (leadership lock, start
> attempts with a backoff, a start that finishes after the leadership was lost
> is stopped).
> h3. Reproduction
> {{MasterConsumerLeadershipTest}} (camel-zookeeper-master, no ZooKeeper: a
> {{ManagedGroupFactory}} in the registry returns a fake group whose events the
> test fires; the delegate endpoint can fail a start or hold {{createConsumer}}
> on a latch), on main:
> {noformat}
> testFailedStartIsRetried: the first start fails
> The master must start its consumer after a failed start ==> expected: <1>
> but was: <0> within 20 seconds.
> testDisconnectedWhileCreatingConsumer: DISCONNECTED while the group thread is
> in createConsumer
> A node that lost the leadership must not consume ==> expected: <0> but was:
> <1>
> testStoppedWhileCreatingConsumer: stopRoute while the group thread is in
> createConsumer
> A stopped route must not consume ==> expected: <0> but was: <1>
> {noformat}
> No sleeps: latches, Awaitility.
> The defect was found with a TLA+ model of the consumer (group thread, Curator
> connection thread, stopping thread; the null check, the creation, the start
> and the stop as separate steps; {{close()}} waiting at most 5 s): "a master
> that received a leadership event consumes" is violated in 3 steps with one
> failed start, and "no delegate runs on a stopped route or a node that is not
> the master" in 5 steps (event, connection lost, stopConsumer, create, start)
> and in 7 steps for the route stop. A first version of the model let the start
> read a stale delegate; the Java reads the field again
> ({{startService(delegate)}}), which the corrected model follows. With the fix
> all properties hold, with up to 3 events, 2 failed starts, a disconnect and a
> stop, and TLC's deadlock check on; mutations of the fix without the
> generation check (a disconnect during the start leaves a zombie) or without
> the retry violate them. The model does not cover reconnecting as master after
> a disconnect.
> h3. Proposed fix
> A {{leadershipLock}} and a generation counter (incremented by
> {{onDisconnected}} and {{doStop}}): a start is marked under the lock after
> checking again that the group is connected and this node is the master
> ({{ZooKeeperGroup}} clears {{connected}} before it calls {{DISCONNECTED}});
> the consumer is created and started without holding the lock (it can take
> long); it is published under the lock only if the generation did not change,
> otherwise it is stopped. A failed start is retried after 5 seconds by a
> scheduled task of the same generation (cancelled by a disconnect or stop),
> and the started state is published to the group once per leadership term, so
> a retry does not trigger another group event (a new state with a random
> container id would). This follows camel-master's {{MasterConsumer}}
> (CAMEL-24583, CAMEL-24626) in its locking and lost-leadership checks, with
> two differences: camel-zookeeper-master has no
> {{backOffDelay}}/{{backOffMaxAttempts}} options, so the delay is fixed at
> camel-master's default of 5 seconds and the retries go on while this node is
> the master (camel-master gives up after 10 attempts and stays leader without
> consuming); and each attempt creates a new delegate consumer (a failed
> {{start()}} already stopped it), where camel-master reuses one. Each failed
> attempt is logged at WARN with its stack trace.
> With the fix the new tests and the camel-zookeeper-master unit tests (14)
> pass; the integration tests need Docker and were not run.
> Affected: 4.14.x, 4.18.x and main (same code, GitHub contents API);
> {{MasterConsumer}} unchanged since 2021.
> Duplicate check (2026-10-04): JIRA component camel-zookeeper-master and text
> "zookeeper-master"/"zookeepermaster" (53 issues; CAMEL-24457 and CAMEL-24545
> are about the camel-zookeeper cluster service, CAMEL-17226 namespaces; none
> about the master consumer's start or races). No open PR touches the component.
> _Filed with Claude Code on behalf of allthingssecurity._
--
This message was sent by Atlassian Jira
(v8.20.10#820010)