CAMEL-25273: camel-kubernetes - consumers must watch again when the client closes their watch with an error - #27302
Conversation
…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
left a comment
There was a problem hiding this comment.
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.
| */ | ||
| 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); | ||
| } | ||
| } | ||
| } |
There was a problem hiding this comment.
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.
There was a problem hiding this comment.
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
left a comment
There was a problem hiding this comment.
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.
|
🌟 Thank you for your contribution to the Apache Camel project! 🌟 🐫 Apache Camel Committers, please review the following items:
|
davsclaus
left a comment
There was a problem hiding this comment.
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:
- Catch the exception around
task.run(), log it at WARN, and callwatchAgain(...)again, ideally with a growing delay capped at some maximum, so a down API server is retried without a tight loop. - Close the stop/re-watch race:
watch()blocks during the websocket handshake. If the consumer stops in that window,doStopcloses the old watch, and the new one is assigned afterwards and never closed. After the watch is assigned, checkisRunAllowed()and close the new watch if the consumer has stopped. - 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.
| 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(); | ||
| } | ||
| }); |
There was a problem hiding this comment.
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).
There was a problem hiding this comment.
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
| if (consumer.isRunAllowed()) { | ||
| task.run(); | ||
| } |
There was a problem hiding this comment.
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.
There was a problem hiding this comment.
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
|
🧪 CI tested the following changed modules:
🔬 Scalpel shadow comparison — Scalpel: 9 of 698 tested, 27 compile-only — current: 9 all testedMaveniverse 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)
Modules with tests skipped (27)
All tested modules (36 modules, 5m 34s total)Total reactor time: 5m 34s
Top 20 slowest modules:
|
… 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>
|
Thanks for the review. All three points are addressed in 314ab43:
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
left a comment
There was a problem hiding this comment.
Re-review after 314ab43 — all three points from davsclaus are addressed:
-
Re-watch failure retry with backoff — Addressed.
task.run()is now caught, logged at WARN, andwatchAgainis called recursively with delay doubled (1s → 2s → 4s → … → 30s cap). The recursion goes throughexecutor.submit(), so no stack growth. Interrupt and!isRunAllowed()correctly abort. -
Stop/re-watch race — Addressed. All 11 consumers now have
volatile Watch watchand checkisRunAllowed()after the watch assignment, closing the new watch if the consumer stopped during the handshake. TheKubernetesCustomResourcesConsumerfix 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. -
Test for failed re-watch — Addressed.
testWatchAgainAfterAFailedAttemptToWatchAgaincovers 410 → 403 → success, asserting 3 server requests and a pod event. NoThread.sleepin the test — uses MockEndpoint wait with a 20s budget.
This review was generated by an AI agent, Hermès on behalf of @gnodet.
Description
CAMEL-25273
The fabric8 client reconnects a watch by itself after transient errors, but closes it for good, calling
Watcher.onClosewith aWatcherException, 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:
onClosewith 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 (ADDEDfor 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 anERRORevent with a 410Status, a second watch would get anADDEDpod.mock://result Received message count. Expected: <1> but was: <0>.Target
mainbranch)Tracking
Apache Camel coding standards and style
mvn clean install -DskipTestslocally 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
Co-authored-bytrailers) 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-Bytrailer.Follow-up (review):
watchfields are volatile). The custom resources consumer now keeps its watch, so it is closed on stop too.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