allthingssecurity commented on code in PR #27302:
URL: https://github.com/apache/camel/pull/27302#discussion_r4172036930
##########
components/camel-kubernetes/src/main/java/org/apache/camel/component/kubernetes/KubernetesHelper.java:
##########
@@ -90,6 +93,45 @@ public static void close(Runnable runnable, Supplier<Watch>
watchGetter) {
}
}
+ /**
+ * The delay before a consumer watches again after the Kubernetes client
closed its watch with an error.
+ */
+ static final long WATCH_AGAIN_DELAY_MILLIS = 1000;
+
+ /**
+ * Watches again when the Kubernetes client closed the watch of a consumer
with an error. The client reconnects a
+ * watch by itself after transient errors, and only closes it with an
exception when it gives up: when the API
+ * server answers 410 Gone because the resource version of the watch is
too old (which happens to long-running
+ * watches), or when the reconnect limit is reached. The consumer would
then not receive any event anymore.
+ *
+ * @param consumer the consumer of the watch
+ * @param executor the executor of the consumer
+ * @param task the task that creates the watch of the consumer
+ */
+ public static void watchAgain(ServiceSupport consumer, ExecutorService
executor, Runnable task) {
+ if (consumer.isRunAllowed() && executor != null &&
!executor.isShutdown()) {
+ LOG.info("Watching again for {} in {} ms after its watch was
closed", consumer, WATCH_AGAIN_DELAY_MILLIS);
+ try {
+ 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();
+ }
+ });
Review 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_
##########
components/camel-kubernetes/src/main/java/org/apache/camel/component/kubernetes/KubernetesHelper.java:
##########
@@ -90,6 +93,45 @@ public static void close(Runnable runnable, Supplier<Watch>
watchGetter) {
}
}
+ /**
+ * The delay before a consumer watches again after the Kubernetes client
closed its watch with an error.
+ */
+ static final long WATCH_AGAIN_DELAY_MILLIS = 1000;
+
+ /**
+ * Watches again when the Kubernetes client closed the watch of a consumer
with an error. The client reconnects a
+ * watch by itself after transient errors, and only closes it with an
exception when it gives up: when the API
+ * server answers 410 Gone because the resource version of the watch is
too old (which happens to long-running
+ * watches), or when the reconnect limit is reached. The consumer would
then not receive any event anymore.
+ *
+ * @param consumer the consumer of the watch
+ * @param executor the executor of the consumer
+ * @param task the task that creates the watch of the consumer
+ */
+ public static void watchAgain(ServiceSupport consumer, ExecutorService
executor, Runnable task) {
+ if (consumer.isRunAllowed() && executor != null &&
!executor.isShutdown()) {
+ LOG.info("Watching again for {} in {} ms after its watch was
closed", consumer, WATCH_AGAIN_DELAY_MILLIS);
+ try {
+ 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();
+ }
Review 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_
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]