davsclaus commented on code in PR #27302:
URL: https://github.com/apache/camel/pull/27302#discussion_r4171963074
##########
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:
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).
##########
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:
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.
--
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]