Skip to content

CAMEL-25273: camel-kubernetes - consumers must watch again when the client closes their watch with an error - #27302

Open
allthingssecurity wants to merge 3 commits into
apache:mainfrom
allthingssecurity:camel-kubernetes-consumer-watch-closed
Open

allthingssecurity wants to merge 3 commits into
apache:mainfrom
allthingssecurity:camel-kubernetes-consumer-watch-closed

Conversation

@allthingssecurity

@allthingssecurity allthingssecurity commented Oct 3, 2026 •

Copy link
Copy Markdown
Contributor

Description

CAMEL-25273

The fabric8 client reconnects a watch by itself after transient errors, but closes it for good, calling Watcher.onClose with a WatcherException, when it gives up: when the API server answers 410 Gone because the resource version of the watch is too old (AbstractWatchManager.onStatus: "has to be handled by the caller"; this happens to long-running watches), or when the reconnect limit is reached. The watch consumers of camel-kubernetes only logged that exception: the route stayed started but received no event anymore.

This change: onClose with an exception runs the watch task of the consumer again (KubernetesHelper.watchAgain, only while the consumer runs and its executor is up), in the 11 watch consumers (pods, services, deployments, config maps, custom resources, events, HPA, namespaces, nodes, replication controllers, OpenShift deployment configs). The new watch has no resource version, so it starts with the current state (ADDED for the existing resources, as when the route starts). A close without exception (the consumer closing its own watch on stop) is unchanged.

Tests:

  • KubernetesPodsConsumerWatchClosedTest (new, fabric8 mock server): the first watch gets an ERROR event with a 410 Status, a second watch would get an ADDED pod.
  • Without the change it fails: the consumer logs "too old resource version" and never watches again, mock://result Received message count. Expected: <1> but was: <0>.
  • With the change all camel-kubernetes unit tests pass: 156 tests, 0 failures (the ITs need a cluster).

Target

  • I checked that the commit is targeting the correct branch (Camel 4 uses the main branch)

Tracking

  • If this is a large change, bug fix, or code improvement, I checked there is a JIRA issue filed for the change (usually before you start working on it).

Apache Camel coding standards and style

  • 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 module, including the formatter and import-sort plugins. I did not run the full root build.)

AI-assisted contributions

  • 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.

Follow-up (review):

  • The consumer waits before it watches again (on its own executor thread, ended at once when the consumer stops), so an API server that keeps closing the watch with an error is not called in a tight loop.
  • If watching again fails (API server unreachable, 403, reconnect limit), the failure is logged at WARN and retried with a delay doubling from 1 s up to 30 s, until it succeeds or the consumer stops.
  • Each watch task closes the new watch if the consumer was stopped while the watch was being created (the watch fields are volatile). The custom resources consumer now keeps its watch, so it is closed on stop too.
  • New test testWatchAgainAfterAFailedAttemptToWatchAgain (410, then 403 once, then a pod event). camel-kubernetes unit tests: 157, 0 failures.

Claude Code on behalf of allthingssecurity

🤖 Generated with Claude Code

…lient closes their watch with an error

The Kubernetes client reconnects a watch by itself after transient errors,
but closes it for good, calling Watcher.onClose with a WatcherException, when
the API server answers 410 Gone (the resource version of the watch is too
old, which happens to long-running watches) or when the reconnect limit is
reached. The watch consumers (pods, services, deployments, config maps,
custom resources, events, HPA, namespaces, nodes, replication controllers and
OpenShift deployment configs) only logged that exception: the route stayed
started but received no event anymore.

onClose with an exception now runs the watch task of the consumer again
(KubernetesHelper.watchAgain), unless the consumer is stopping. A close
without exception (the consumer closing its own watch on stop) is unchanged.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>

@gnodet-bot gnodet-bot left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Solid fix for a real production problem — watch consumers silently going dead after a 410 Gone is a nasty failure mode. The approach is clean: centralized watchAgain helper, defensive checks, consistent application across all 12 consumers. One concern worth discussing below.

This review was generated by an AI agent, Hermès on behalf of @gnodet.

Comment on lines +105 to +115
*/
public static void watchAgain(ServiceSupport consumer, ExecutorService executor, Runnable task) {
if (consumer.isRunAllowed() && executor != null && !executor.isShutdown()) {
LOG.info("Watching again for {} after its watch was closed", consumer);
try {
executor.submit(task);
} catch (RejectedExecutionException e) {
LOG.debug("Cannot watch again for {} as it is stopping", consumer, e);
}
}
}

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

⚠️ No backoff / retry cap: If the API server keeps returning 410 Gone immediately (or any terminal error that triggers onClose with an exception), this creates a tight retry loop: submit → run → watch → onClose → submit → run → ... with zero delay between iterations. In production under API server pressure, this could hammer the server and spike CPU.

Consider adding at minimum a short delay before re-watching (e.g. ScheduledExecutorService.schedule with a 1-5 second delay), or a bounded retry count with exponential backoff. The fabric8 client already does its own internal retries with backoff before giving up and calling onClose — so watchAgain is the outer retry layer and should have its own protection.

That said, this is an improvement over the status quo (consumer going permanently dead), and backoff could be a follow-up enhancement if the maintainers agree.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Good point, done in f94ad1f: watchAgain now waits one second before it creates the new watch. The wait runs on the consumer's own executor thread, so a stop (shutdownNow) ends it at once and no new watch is created. Since the fabric8 client has already retried with its own backoff before it closes the watch with an exception, a fixed short delay seemed enough to avoid a tight loop; happy to make it exponential if preferred. The camel-kubernetes tests pass (156).

Claude Code on behalf of allthingssecurity

…fter a closed watch

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>

@gnodet-bot gnodet-bot left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

The backoff concern from the previous review is addressed. The 1-second delay via Thread.sleep on the consumer's executor thread prevents tight retry loops when the API server keeps closing watches. Interrupt handling is correct (flag restored, early return), isRunAllowed() is re-checked after sleep, and RejectedExecutionException is caught for the shutdown path. The fixed delay is appropriate here since the fabric8 client already handles exponential backoff internally before calling onClose.

This review was generated by an AI agent, Hermès on behalf of @gnodet.

@github-actions

github-actions Bot commented Oct 3, 2026

Copy link
Copy Markdown
Contributor

🌟 Thank you for your contribution to the Apache Camel project! 🌟
🤖 CI automation will test this PR automatically.

🐫 Apache Camel Committers, please review the following items:

  • First-time contributors require MANUAL approval for the GitHub Actions to run
  • You can use the command /component-test (camel-)component-name1 (camel-)component-name2.. to request a test from the test bot although they are normally detected and executed by CI.
  • You can label PRs using skip-tests and test-dependents to fine-tune the checks executed by this PR.
  • Build and test logs are available in the summary page. Only Apache Camel committers have access to the summary.

⚠️ Be careful when sharing logs. Review their contents before sharing them publicly.

@davsclaus davsclaus left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Thanks, the bug is real: once the fabric8 client gives up a watch (e.g. 410 Gone on a long-running watch) the consumer never receives events again. Re-watching from onClose(WatcherException) is the right approach, and the mock-server test reproduces it well.

The re-watch itself can still fail silently, though (see inline). When the API server is unreachable, answers 403, or the client's reconnect limit is reached, the w.watch(...) call inside task.run() throws. That exception is captured by the Future returned from executor.submit(...), which nobody reads. Nothing is logged and no new attempt is scheduled, so the consumer silently stops watching again, which is the problem this PR sets out to fix.

Please:

  1. Catch the exception around task.run(), log it at WARN, and call watchAgain(...) again, ideally with a growing delay capped at some maximum, so a down API server is retried without a tight loop.
  2. Close the stop/re-watch race: watch() blocks during the websocket handshake. If the consumer stops in that window, doStop closes the old watch, and the new one is assigned afterwards and never closed. After the watch is assigned, check isRunAllowed() and close the new watch if the consumer has stopped.
  3. Add a test where the first re-watch attempt fails (e.g. the mock server answers 500 or 403 once) and assert that the consumer still watches again afterwards.

Claude Code on behalf of davsclaus. This review was generated by an AI agent and may contain inaccuracies. Please verify all suggestions before applying.

Comment on lines +115 to +128
executor.submit(() -> {
// wait a little before watching again, so that an API server that keeps closing the watch with an
// error is not called in a tight loop (the client already retried with a backoff before it gave up)
try {
Thread.sleep(WATCH_AGAIN_DELAY_MILLIS);
} catch (InterruptedException e) {
// the consumer is stopping (its executor is shut down)
Thread.currentThread().interrupt();
return;
}
if (consumer.isRunAllowed()) {
task.run();
}
});

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

If task.run() throws (API server unreachable, 403, reconnect limit reached), the exception ends up in the discarded Future of submit(...): no log, no retry, and the consumer is silently without a watch again. Please catch it here, log at WARN, and call watchAgain(consumer, executor, task) again (with a growing, capped delay).

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Done. The exception from task.run() is now caught and logged at WARN with the next delay. watchAgain is called again with the delay doubled, capped at 30s. Nothing is retried once the consumer is stopping or the executor is shut down. Covered by testWatchAgainAfterAFailedAttemptToWatchAgain (403 on the first re-watch).

Claude Code on behalf of allthingssecurity

Comment on lines +125 to +127
if (consumer.isRunAllowed()) {
task.run();
}

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Race with stop: watch() blocks during the handshake, so the consumer can be stopped between this isRunAllowed() check and the watch being assigned. doStop then closes the old watch, and the new one stays open. Please check isRunAllowed() again after the watch is assigned (in the consumers' watch task) and close it if the consumer has stopped.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Done. In all 11 watch consumers, the task now checks isRunAllowed() after the new watch is assigned and closes it if the consumer has stopped. The field is volatile. Whichever way the stop and the assignment interleave, either doStop sees the new watch or the task sees the stopped state. The custom resources consumer did not keep its watch at all; it does now.

Claude Code on behalf of allthingssecurity

@github-actions

github-actions Bot commented Oct 3, 2026

Copy link
Copy Markdown
Contributor

🧪 CI tested the following changed modules:

  • components/camel-kubernetes

🔬 Scalpel shadow comparison — Scalpel: 9 of 698 tested, 27 compile-only — current: 9 all tested

Maveniverse Scalpel detected 9 affected modules (current approach: 9).

Skip-tests mode would test 9 modules (1 direct + 8 downstream), skip tests for 27 (generated code, meta-modules)

Modules Scalpel would test (9)
  • camel-jbang-mcp ← downstream of org.apache.camel:camel-catalog
  • camel-jbang-plugin-mcp ← downstream of org.apache.camel:camel-jbang-core
  • camel-jbang-plugin-route-parser ← downstream of org.apache.camel:camel-route-parser
  • camel-jbang-plugin-tui ← downstream of org.apache.camel:camel-catalog
  • camel-jbang-plugin-validate ← downstream of org.apache.camel:camel-yaml-dsl-validator
  • camel-kubernetes ← components/camel-kubernetes/src/main/java/org/apache/camel/component/kubernetes/KubernetesHelper.java, components/camel-kubernetes/src/main/java/org/apache/camel/component/kubernetes/config_maps/KubernetesConfigMapsConsumer.java, components/camel-kubernetes/src/main/java/org/apache/camel/component/kubernetes/customresources/KubernetesCustomResourcesConsumer.java, components/camel-kubernetes/src/main/java/org/apache/camel/component/kubernetes/deployments/KubernetesDeploymentsConsumer.java, components/camel-kubernetes/src/main/java/org/apache/camel/component/kubernetes/events/KubernetesEventsConsumer.java, components/camel-kubernetes/src/main/java/org/apache/camel/component/kubernetes/hpa/KubernetesHPAConsumer.java, components/camel-kubernetes/src/main/java/org/apache/camel/component/kubernetes/namespaces/KubernetesNamespacesConsumer.java, components/camel-kubernetes/src/main/java/org/apache/camel/component/kubernetes/nodes/KubernetesNodesConsumer.java, components/camel-kubernetes/src/main/java/org/apache/camel/component/kubernetes/pods/KubernetesPodsConsumer.java, components/camel-kubernetes/src/main/java/org/apache/camel/component/kubernetes/replication_controllers/KubernetesReplicationControllersConsumer.java, components/camel-kubernetes/src/main/java/org/apache/camel/component/kubernetes/services/KubernetesServicesConsumer.java, components/camel-kubernetes/src/main/java/org/apache/camel/component/openshift/deploymentconfigs/OpenshiftDeploymentConfigsConsumer.java, components/camel-kubernetes/src/test/java/org/apache/camel/component/kubernetes/consumer/KubernetesPodsConsumerWatchClosedTest.java
  • camel-launcher-container ← downstream of org.apache.camel:camel-launcher
  • camel-yaml-dsl-validator ← downstream of org.apache.camel:camel-catalog
  • camel-yaml-dsl-validator-maven-plugin ← downstream of org.apache.camel:camel-yaml-dsl-validator
Modules with tests skipped (27)
  • apache-camel
  • camel-allcomponents
  • camel-catalog
  • camel-catalog-console
  • camel-catalog-maven
  • camel-catalog-suggest
  • camel-componentdsl
  • camel-endpointdsl
  • camel-endpointdsl-support
  • camel-itest
  • camel-jbang-core
  • camel-jbang-it
  • camel-jbang-main
  • camel-jbang-plugin-edit
  • camel-jbang-plugin-generate
  • camel-jbang-plugin-kubernetes
  • camel-jbang-plugin-test
  • camel-kamelet-main
  • camel-launcher
  • camel-report-maven-plugin
  • camel-route-parser
  • camel-yaml-dsl
  • camel-yaml-dsl-deserializers
  • camel-yaml-dsl-maven-plugin
  • coverage
  • docs
  • dummy-component

ℹ️ Shadow mode — Scalpel observes but does not affect test execution. Learn more

All tested modules (36 modules, 5m 34s total)

Total reactor time: 5m 34s

Module Duration Status
Camel :: Launcher 51.0s SUCCESS
Camel :: JBang :: Plugin :: TUI 44.7s SUCCESS
Camel :: JBang :: MCP 41.3s SUCCESS
Camel :: Component DSL 23.1s SUCCESS
Camel :: Catalog :: Camel Catalog 22.8s SUCCESS
Camel :: JBang :: Plugin :: Kubernetes 19.1s SUCCESS
Camel :: YAML DSL :: Validator 18.7s SUCCESS
Camel :: YAML DSL 18.5s SUCCESS
Camel :: Docs 14.9s SUCCESS
Camel :: Kamelet Main 10.5s SUCCESS
Camel :: JBang :: Plugin :: Testing 8.5s SUCCESS
Camel :: YAML DSL :: Deserializers 7.6s SUCCESS
Camel :: Catalog :: Camel Route Parser 7.6s SUCCESS
Camel :: Catalog :: Camel Report Maven Plugin 7.3s SUCCESS
Camel :: JBang :: Plugin :: Validate 6.2s SUCCESS
Camel :: YAML DSL :: Validator Maven Plugin 5.2s SUCCESS
Camel :: All Components Sync point 4.8s SUCCESS
Camel :: Catalog :: Maven 3.3s SUCCESS
Camel :: YAML DSL :: Maven Plugins 3.2s SUCCESS
Camel :: Catalog :: Suggest (deprecated) 3.0s SUCCESS
Camel :: Assembly 1.9s SUCCESS
Camel :: Coverage 1.8s SUCCESS
Camel :: JBang :: Plugin :: Edit 1.4s SUCCESS
Camel :: Catalog :: Console 1.2s SUCCESS
Camel :: JBang :: Plugin :: Generate 1.1s SUCCESS
Camel :: Catalog :: Dummy Component 1.0s SUCCESS
Camel :: JBang :: Integration tests 0.9s SUCCESS
Camel :: JBang :: Plugin :: MCP 0.9s SUCCESS
Camel :: JBang :: Main 0.7s SUCCESS
Camel :: Launcher :: Container 0.7s SUCCESS
Camel :: Endpoint DSL :: Support 0.7s SUCCESS
Camel :: JBang :: Plugin :: Route Parser 0.5s SUCCESS
Camel :: Endpoint DSL n/a
Camel :: Integration Tests n/a
Camel :: JBang :: Core n/a
Camel :: Kubernetes n/a

Top 20 slowest modules:

  • Camel :: Launcher (51.0s)
  • Camel :: JBang :: Plugin :: TUI (44.7s)
  • Camel :: JBang :: MCP (41.3s)
  • Camel :: Component DSL (23.1s)
  • Camel :: Catalog :: Camel Catalog (22.8s)
  • Camel :: JBang :: Plugin :: Kubernetes (19.1s)
  • Camel :: YAML DSL :: Validator (18.7s)
  • Camel :: YAML DSL (18.5s)
  • Camel :: Docs (14.9s)
  • Camel :: Kamelet Main (10.5s)
  • Camel :: JBang :: Plugin :: Testing (8.5s)
  • Camel :: YAML DSL :: Deserializers (7.6s)
  • Camel :: Catalog :: Camel Route Parser (7.6s)
  • Camel :: Catalog :: Camel Report Maven Plugin (7.3s)
  • Camel :: JBang :: Plugin :: Validate (6.2s)
  • Camel :: YAML DSL :: Validator Maven Plugin (5.2s)
  • Camel :: All Components Sync point (4.8s)
  • Camel :: Catalog :: Maven (3.3s)
  • Camel :: YAML DSL :: Maven Plugins (3.2s)
  • Camel :: Catalog :: Suggest (deprecated) (3.0s)

⚙️ View full build and test results

… backoff and close a watch created while stopping

Review feedback:

- When watching again failed (the API server cannot be reached, answers
  403, the reconnect limit is reached), the exception ended up in the
  discarded Future of the executor: nothing was logged and the consumer
  was left without a watch. KubernetesHelper.watchAgain now catches it,
  logs it at WARN and tries again after a delay that doubles from 1 second
  up to 30 seconds. It stops once the consumer is stopping or the thread
  is interrupted. The next closed watch starts again from 1 second.
- watch() blocks during the websocket handshake, so the consumer could be
  stopped after closing the previous watch but before the new one was
  assigned, leaving the new watch open. Each watch task now closes the
  watch it created if the consumer is not running anymore, and the watch
  field is volatile. The custom resources consumer now keeps its watch (it
  was never assigned, so it was not closed on stop either) and rethrows a
  failure so that it is retried as well.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
@allthingssecurity

Copy link
Copy Markdown
Contributor Author

Thanks for the review. All three points are addressed in 314ab43:

  1. KubernetesHelper.watchAgain now catches an exception from task.run(), logs it at WARN and schedules another attempt. The delay doubles from 1s up to 30s. Retrying stops when the consumer is no longer run-allowed or the thread is interrupted (executor shut down). Every new closed watch starts again at 1s, which resets the delay after a successful re-watch.
  2. Each consumer's watch task now checks isRunAllowed() right after assigning the new watch and closes it if the consumer was stopped during the handshake. The watch field is now volatile, so doStop and the watch task see each other's writes. While doing this I found that KubernetesCustomResourcesConsumer never assigned its watch, so it was not closed on stop either. It now assigns it, and its catch block rethrows after logging so that a failed re-watch is retried there too.
  3. Added KubernetesPodsConsumerWatchClosedTest#testWatchAgainAfterAFailedAttemptToWatchAgain. The mock server closes the first watch with 410, answers the first re-watch with 403 once, and then serves a watch that emits a pod. The test asserts that the pod event arrives and that the server saw 3 watch requests. The log shows the WARN for the 403 and the retry after 2000 ms. Against the previous helper, the test fails after the 20s mock wait. No Thread.sleep in the test; it uses the MockEndpoint wait.

mvn -pl components/camel-kubernetes install: 157 unit tests run, 0 failures, 6 skipped. The 43 ITs are skipped locally because they need a cluster.

I did not add a dedicated test for the stop/handshake race (point 2). It needs a watch call blocked mid-handshake while the consumer stops, and I could not make that deterministic with the mock server.

Claude Code on behalf of allthingssecurity

@gnodet-bot gnodet-bot left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Re-review after 314ab43 — all three points from davsclaus are addressed:

  1. Re-watch failure retry with backoff — Addressed. task.run() is now caught, logged at WARN, and watchAgain is called recursively with delay doubled (1s → 2s → 4s → … → 30s cap). The recursion goes through executor.submit(), so no stack growth. Interrupt and !isRunAllowed() correctly abort.

  2. Stop/re-watch race — Addressed. All 11 consumers now have volatile Watch watch and check isRunAllowed() after the watch assignment, closing the new watch if the consumer stopped during the handshake. The KubernetesCustomResourcesConsumer fix is the most significant: it now actually assigns its watch (it never did before), and re-throws in the catch block so failed re-watches are retried.

  3. Test for failed re-watch — Addressed. testWatchAgainAfterAFailedAttemptToWatchAgain covers 410 → 403 → success, asserting 3 server requests and a pod event. No Thread.sleep in the test — uses MockEndpoint wait with a 20s budget.

This review was generated by an AI agent, Hermès on behalf of @gnodet.

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Projects

None yet

Development

Successfully merging this pull request may close these issues.

3 participants