From e4d98d965b7ef2a9402cbd6467aa85c903b9086e Mon Sep 17 00:00:00 2001 From: Cosimo Damiano Prete Date: Fri, 4 Sep 2026 11:56:24 +0200 Subject: [PATCH] fix: Make sure that StatusUpdateTrigger and InfoUpdateTrigger keep running on rolling updates --- .../config/AdminServerAutoConfiguration.java | 36 +++++++++++++++-- .../EventsourcingInstanceRepository.java | 5 ++- .../SnapshottingInstanceRepository.java | 8 +++- .../eventstore/ConcurrentMapEventStore.java | 4 +- .../eventstore/HazelcastEventStore.java | 11 +++++- .../server/services/InfoUpdateTrigger.java | 17 ++++++++ .../server/services/StatusUpdateTrigger.java | 17 ++++++++ .../SnapshottingInstanceRepositoryTest.java | 21 ++++++++++ .../eventstore/HazelcastEventStoreTest.java | 39 +++++++++++++++++++ .../services/InfoUpdateTriggerTest.java | 18 +++++++++ .../services/StatusUpdateTriggerTest.java | 18 +++++++++ 11 files changed, 186 insertions(+), 8 deletions(-) diff --git a/spring-boot-admin-server/src/main/java/de/codecentric/boot/admin/server/config/AdminServerAutoConfiguration.java b/spring-boot-admin-server/src/main/java/de/codecentric/boot/admin/server/config/AdminServerAutoConfiguration.java index e58a4d12d48..cae496aa82d 100644 --- a/spring-boot-admin-server/src/main/java/de/codecentric/boot/admin/server/config/AdminServerAutoConfiguration.java +++ b/spring-boot-admin-server/src/main/java/de/codecentric/boot/admin/server/config/AdminServerAutoConfiguration.java @@ -35,10 +35,13 @@ import org.springframework.context.annotation.Conditional; import org.springframework.context.annotation.Configuration; import org.springframework.context.annotation.Lazy; +import reactor.core.publisher.Flux; +import de.codecentric.boot.admin.server.domain.entities.Instance; import de.codecentric.boot.admin.server.domain.entities.InstanceRepository; import de.codecentric.boot.admin.server.domain.entities.SnapshottingInstanceRepository; import de.codecentric.boot.admin.server.domain.events.InstanceEvent; +import de.codecentric.boot.admin.server.domain.values.InstanceId; import de.codecentric.boot.admin.server.eventstore.InMemoryEventStore; import de.codecentric.boot.admin.server.eventstore.InstanceEventPublisher; import de.codecentric.boot.admin.server.eventstore.InstanceEventStore; @@ -162,7 +165,8 @@ public StatusUpdater statusUpdater(InstanceRepository instanceRepository, @Bean(initMethod = "start", destroyMethod = "stop") @ConditionalOnMissingBean - public StatusUpdateTrigger statusUpdateTrigger(StatusUpdater statusUpdater, Publisher events) { + public StatusUpdateTrigger statusUpdateTrigger(StatusUpdater statusUpdater, Publisher events, + InstanceRegistry instanceRegistry) { AdminServerProperties.MonitorProperties monitorProperties = this.adminServerProperties.getMonitor(); Duration defaultTimeout = monitorProperties.getDefaultTimeout(); @@ -175,7 +179,7 @@ public StatusUpdateTrigger statusUpdateTrigger(StatusUpdater statusUpdater, Publ } return new StatusUpdateTrigger(statusUpdater, events, statusInterval, monitorProperties.getStatusLifetime(), - monitorProperties.getStatusMaxBackoff()); + monitorProperties.getStatusMaxBackoff(), getExistingInstanceIds(instanceRegistry)); } @Bean @@ -212,10 +216,11 @@ public InfoUpdater infoUpdater(InstanceRepository instanceRepository, @Bean(initMethod = "start", destroyMethod = "stop") @ConditionalOnMissingBean - public InfoUpdateTrigger infoUpdateTrigger(InfoUpdater infoUpdater, Publisher events) { + public InfoUpdateTrigger infoUpdateTrigger(InfoUpdater infoUpdater, Publisher events, + InstanceRegistry instanceRegistry) { return new InfoUpdateTrigger(infoUpdater, events, this.adminServerProperties.getMonitor().getInfoInterval(), this.adminServerProperties.getMonitor().getInfoLifetime(), - this.adminServerProperties.getMonitor().getInfoMaxBackoff()); + this.adminServerProperties.getMonitor().getInfoMaxBackoff(), getExistingInstanceIds(instanceRegistry)); } @Bean @@ -230,4 +235,27 @@ public SnapshottingInstanceRepository instanceRepository(InstanceEventStore even return new SnapshottingInstanceRepository(eventStore); } + /* + * Fetches the existing registered instance IDs from the instance registry to use them + * as initial data set for the StatusUpdateTrigger and InfoUpdaterTrigger. This + * ensures that the triggers will update the status and info for all existing + * instances on startup and correctly start polling for the updates. This is necessary + * because the IntervalCheck used in the triggers only updates the status and info for + * instances that have been updated since the last check by checking the local + * "lastChecked" map. On rolling updates with Hazelcast, the details about the + * instances will be migrated from an instance to another, but the "lastChecked" map + * will be empty for the new instance, so the triggers will not update the status and + * info for the existing instances. As such, the existing instance IDs are fetched and + * passed to the triggers to ensure that the "lastChecked" map is aware of them + * accordingly. + * + * @param instanceRegistry the registry to fetch the existing registered instance IDs + * from + * + * @return a Flux of existing registered instance IDs + */ + private static Flux getExistingInstanceIds(InstanceRegistry instanceRegistry) { + return instanceRegistry.getInstances().filter(Instance::isRegistered).map(Instance::getId).distinct(); + } + } diff --git a/spring-boot-admin-server/src/main/java/de/codecentric/boot/admin/server/domain/entities/EventsourcingInstanceRepository.java b/spring-boot-admin-server/src/main/java/de/codecentric/boot/admin/server/domain/entities/EventsourcingInstanceRepository.java index 6afb541a77a..f2090afc30b 100644 --- a/spring-boot-admin-server/src/main/java/de/codecentric/boot/admin/server/domain/entities/EventsourcingInstanceRepository.java +++ b/spring-boot-admin-server/src/main/java/de/codecentric/boot/admin/server/domain/entities/EventsourcingInstanceRepository.java @@ -16,6 +16,7 @@ package de.codecentric.boot.admin.server.domain.entities; +import java.util.Comparator; import java.util.function.BiFunction; import org.slf4j.Logger; @@ -38,6 +39,8 @@ public class EventsourcingInstanceRepository implements InstanceRepository { private static final Logger log = LoggerFactory.getLogger(EventsourcingInstanceRepository.class); + private static final Comparator byVersion = Comparator.comparingLong(InstanceEvent::getVersion); + private final InstanceEventStore eventStore; private final Retry retryOptimisticLockException = Retry.max(10) @@ -57,7 +60,7 @@ public Mono save(Instance instance) { public Flux findAll() { return this.eventStore.findAll() .groupBy(InstanceEvent::getInstance) - .flatMap((f) -> f.reduce(Instance.create(f.key()), Instance::apply)); + .flatMap((f) -> f.sort(byVersion).reduce(Instance.create(f.key()), Instance::apply)); } @Override diff --git a/spring-boot-admin-server/src/main/java/de/codecentric/boot/admin/server/domain/entities/SnapshottingInstanceRepository.java b/spring-boot-admin-server/src/main/java/de/codecentric/boot/admin/server/domain/entities/SnapshottingInstanceRepository.java index bb40b90bf1b..d0d1b419858 100644 --- a/spring-boot-admin-server/src/main/java/de/codecentric/boot/admin/server/domain/entities/SnapshottingInstanceRepository.java +++ b/spring-boot-admin-server/src/main/java/de/codecentric/boot/admin/server/domain/entities/SnapshottingInstanceRepository.java @@ -16,6 +16,7 @@ package de.codecentric.boot.admin.server.domain.entities; +import java.util.Comparator; import java.util.Set; import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ConcurrentMap; @@ -42,6 +43,8 @@ public class SnapshottingInstanceRepository extends EventsourcingInstanceReposit private static final Logger log = LoggerFactory.getLogger(SnapshottingInstanceRepository.class); + private static final Comparator byVersion = Comparator.comparingLong(InstanceEvent::getVersion); + private final ConcurrentMap snapshots = new ConcurrentHashMap<>(); private final Set outdatedSnapshots = ConcurrentHashMap.newKeySet(); @@ -79,7 +82,10 @@ public Mono save(Instance instance) { } public void start() { - this.subscription = this.eventStore.findAll().concatWith(this.eventStore).subscribe(this::updateSnapshot); + Flux initialEvents = this.eventStore.findAll() + .groupBy(InstanceEvent::getInstance) + .flatMap((events) -> events.sort(byVersion)); + this.subscription = initialEvents.concatWith(this.eventStore).subscribe(this::updateSnapshot); } public void stop() { diff --git a/spring-boot-admin-server/src/main/java/de/codecentric/boot/admin/server/eventstore/ConcurrentMapEventStore.java b/spring-boot-admin-server/src/main/java/de/codecentric/boot/admin/server/eventstore/ConcurrentMapEventStore.java index be059760ecd..ff5fdddb134 100644 --- a/spring-boot-admin-server/src/main/java/de/codecentric/boot/admin/server/eventstore/ConcurrentMapEventStore.java +++ b/spring-boot-admin-server/src/main/java/de/codecentric/boot/admin/server/eventstore/ConcurrentMapEventStore.java @@ -45,6 +45,8 @@ public abstract class ConcurrentMapEventStore extends InstanceEventPublisher imp private static final Logger log = LoggerFactory.getLogger(ConcurrentMapEventStore.class); + protected static final long NO_LATEST_VERSION = -1; + private static final Comparator byTimestampAndIdAndVersion = comparing(InstanceEvent::getTimestamp) .thenComparing(InstanceEvent::getInstance) .thenComparing(InstanceEvent::getVersion); @@ -130,7 +132,7 @@ private OptimisticLockingException createOptimisticLockException(InstanceEvent e } protected static long getLastVersion(List events) { - return events.isEmpty() ? -1 : events.get(events.size() - 1).getVersion(); + return events.isEmpty() ? NO_LATEST_VERSION : events.get(events.size() - 1).getVersion(); } private static DistinctEventType getDistinctEventTypeFor(InstanceEvent event) { diff --git a/spring-boot-admin-server/src/main/java/de/codecentric/boot/admin/server/eventstore/HazelcastEventStore.java b/spring-boot-admin-server/src/main/java/de/codecentric/boot/admin/server/eventstore/HazelcastEventStore.java index 1fc6b754372..3d0beedde52 100644 --- a/spring-boot-admin-server/src/main/java/de/codecentric/boot/admin/server/eventstore/HazelcastEventStore.java +++ b/spring-boot-admin-server/src/main/java/de/codecentric/boot/admin/server/eventstore/HazelcastEventStore.java @@ -44,10 +44,19 @@ public HazelcastEventStore(int maxLogSizePerAggregate, IMap>() { + @Override + public void entryAdded(EntryEvent> event) { + log.debug("Added {}", event); + publishNewEvents(event, NO_LATEST_VERSION); + } + @Override public void entryUpdated(EntryEvent> event) { log.debug("Updated {}", event); - long lastKnownVersion = getLastVersion(event.getOldValue()); + publishNewEvents(event, getLastVersion(event.getOldValue())); + } + + private void publishNewEvents(EntryEvent> event, long lastKnownVersion) { List newEvents = event.getValue() .stream() .filter((e) -> e.getVersion() > lastKnownVersion) diff --git a/spring-boot-admin-server/src/main/java/de/codecentric/boot/admin/server/services/InfoUpdateTrigger.java b/spring-boot-admin-server/src/main/java/de/codecentric/boot/admin/server/services/InfoUpdateTrigger.java index fa7ce27b166..4e37de29e39 100644 --- a/spring-boot-admin-server/src/main/java/de/codecentric/boot/admin/server/services/InfoUpdateTrigger.java +++ b/spring-boot-admin-server/src/main/java/de/codecentric/boot/admin/server/services/InfoUpdateTrigger.java @@ -18,9 +18,11 @@ import java.time.Duration; +import org.jspecify.annotations.Nullable; import org.reactivestreams.Publisher; import org.slf4j.Logger; import org.slf4j.LoggerFactory; +import reactor.core.Disposable; import reactor.core.publisher.Flux; import reactor.core.publisher.Mono; @@ -38,11 +40,21 @@ public class InfoUpdateTrigger extends AbstractEventHandler { private final IntervalCheck intervalCheck; + private final Publisher existingInstanceIds; + + @Nullable private Disposable startupSubscription; + public InfoUpdateTrigger(InfoUpdater infoUpdater, Publisher publisher, Duration updateInterval, Duration infoLifetime, Duration maxBackoff) { + this(infoUpdater, publisher, updateInterval, infoLifetime, maxBackoff, Flux.empty()); + } + + public InfoUpdateTrigger(InfoUpdater infoUpdater, Publisher publisher, Duration updateInterval, + Duration infoLifetime, Duration maxBackoff, Publisher existingInstanceIds) { super(publisher, InstanceEvent.class); this.infoUpdater = infoUpdater; this.intervalCheck = new IntervalCheck("info", this::updateInfo, updateInterval, infoLifetime, maxBackoff); + this.existingInstanceIds = existingInstanceIds; } @Override @@ -64,10 +76,15 @@ protected Mono updateInfo(InstanceId instanceId) { public void start() { super.start(); this.intervalCheck.start(); + this.startupSubscription = Flux.from(this.existingInstanceIds).flatMap(this::updateInfo).subscribe(); } @Override public void stop() { + if (this.startupSubscription != null) { + this.startupSubscription.dispose(); + this.startupSubscription = null; + } super.stop(); this.intervalCheck.stop(); } diff --git a/spring-boot-admin-server/src/main/java/de/codecentric/boot/admin/server/services/StatusUpdateTrigger.java b/spring-boot-admin-server/src/main/java/de/codecentric/boot/admin/server/services/StatusUpdateTrigger.java index 80470f96957..88d361b6570 100644 --- a/spring-boot-admin-server/src/main/java/de/codecentric/boot/admin/server/services/StatusUpdateTrigger.java +++ b/spring-boot-admin-server/src/main/java/de/codecentric/boot/admin/server/services/StatusUpdateTrigger.java @@ -18,9 +18,11 @@ import java.time.Duration; +import org.jspecify.annotations.Nullable; import org.reactivestreams.Publisher; import org.slf4j.Logger; import org.slf4j.LoggerFactory; +import reactor.core.Disposable; import reactor.core.publisher.Flux; import reactor.core.publisher.Mono; @@ -37,12 +39,22 @@ public class StatusUpdateTrigger extends AbstractEventHandler { private final IntervalCheck intervalCheck; + private final Publisher existingInstanceIds; + + @Nullable private Disposable startupSubscription; + public StatusUpdateTrigger(StatusUpdater statusUpdater, Publisher publisher, Duration updateInterval, Duration statusLifetime, Duration maxBackoff) { + this(statusUpdater, publisher, updateInterval, statusLifetime, maxBackoff, Flux.empty()); + } + + public StatusUpdateTrigger(StatusUpdater statusUpdater, Publisher publisher, Duration updateInterval, + Duration statusLifetime, Duration maxBackoff, Publisher existingInstanceIds) { super(publisher, InstanceEvent.class); this.statusUpdater = statusUpdater; this.intervalCheck = new IntervalCheck("status", this::updateStatus, updateInterval, statusLifetime, maxBackoff); + this.existingInstanceIds = existingInstanceIds; } @Override @@ -67,10 +79,15 @@ protected Mono updateStatus(InstanceId instanceId) { public void start() { super.start(); this.intervalCheck.start(); + this.startupSubscription = Flux.from(this.existingInstanceIds).flatMap(this::updateStatus).subscribe(); } @Override public void stop() { + if (this.startupSubscription != null) { + this.startupSubscription.dispose(); + this.startupSubscription = null; + } super.stop(); this.intervalCheck.stop(); } diff --git a/spring-boot-admin-server/src/test/java/de/codecentric/boot/admin/server/domain/entities/SnapshottingInstanceRepositoryTest.java b/spring-boot-admin-server/src/test/java/de/codecentric/boot/admin/server/domain/entities/SnapshottingInstanceRepositoryTest.java index 4fd7bc83e60..ed5383fc6b3 100644 --- a/spring-boot-admin-server/src/test/java/de/codecentric/boot/admin/server/domain/entities/SnapshottingInstanceRepositoryTest.java +++ b/spring-boot-admin-server/src/test/java/de/codecentric/boot/admin/server/domain/entities/SnapshottingInstanceRepositoryTest.java @@ -16,6 +16,8 @@ package de.codecentric.boot.admin.server.domain.entities; +import java.time.Instant; + import org.junit.jupiter.api.AfterEach; import org.junit.jupiter.api.BeforeEach; import org.junit.jupiter.api.Test; @@ -24,6 +26,7 @@ import reactor.test.StepVerifier; import de.codecentric.boot.admin.server.domain.events.InstanceRegisteredEvent; +import de.codecentric.boot.admin.server.domain.events.InstanceStatusChangedEvent; import de.codecentric.boot.admin.server.domain.values.InstanceId; import de.codecentric.boot.admin.server.domain.values.Registration; import de.codecentric.boot.admin.server.domain.values.StatusInfo; @@ -100,6 +103,24 @@ void should_update_cache_after_error() { .verifyComplete(); } + @Test + void should_replay_initial_events_in_version_order() { + this.repository.stop(); + InstanceId id = InstanceId.of("clock-skewed"); + Instant now = Instant.now(); + Registration registration = Registration.create("app", "https://health").build(); + when(this.eventStore.findAll()) + .thenReturn(Flux.just(new InstanceStatusChangedEvent(id, 1L, now.minusSeconds(30), StatusInfo.ofDown()), + new InstanceRegisteredEvent(id, 0L, now, registration))); + + this.repository.start(); + + StepVerifier.create(this.repository.find(id)).assertNext((instance) -> { + assertThat(instance.isRegistered()).isTrue(); + assertThat(instance.getStatusInfo().getStatus()).isEqualTo(StatusInfo.STATUS_DOWN); + }).verifyComplete(); + } + @Test void should_return_outdated_instance_not_present_in_cache() { this.repository.stop(); diff --git a/spring-boot-admin-server/src/test/java/de/codecentric/boot/admin/server/eventstore/HazelcastEventStoreTest.java b/spring-boot-admin-server/src/test/java/de/codecentric/boot/admin/server/eventstore/HazelcastEventStoreTest.java index 3d831f28727..37ce7f04564 100644 --- a/spring-boot-admin-server/src/test/java/de/codecentric/boot/admin/server/eventstore/HazelcastEventStoreTest.java +++ b/spring-boot-admin-server/src/test/java/de/codecentric/boot/admin/server/eventstore/HazelcastEventStoreTest.java @@ -16,9 +16,22 @@ package de.codecentric.boot.admin.server.eventstore; +import java.time.Duration; +import java.util.List; + import com.hazelcast.config.Config; import com.hazelcast.core.Hazelcast; import com.hazelcast.core.HazelcastInstance; +import com.hazelcast.map.IMap; +import org.junit.jupiter.api.Test; +import reactor.test.StepVerifier; + +import de.codecentric.boot.admin.server.domain.events.InstanceEvent; +import de.codecentric.boot.admin.server.domain.events.InstanceRegisteredEvent; +import de.codecentric.boot.admin.server.domain.values.InstanceId; +import de.codecentric.boot.admin.server.domain.values.Registration; + +import static java.util.Collections.singletonList; public class HazelcastEventStoreTest extends AbstractEventStoreTest { @@ -41,4 +54,30 @@ protected void shutdownStore() { } } + @Test + public void should_publish_events_for_added_entries() { + Config config = new Config(); + config.getNetworkConfig().getJoin().getMulticastConfig().setEnabled(false); + config.getNetworkConfig().getJoin().getAutoDetectionConfig().setEnabled(false); + HazelcastInstance hazelcastInstance = Hazelcast.newHazelcastInstance(config); + try { + InstanceId id = InstanceId.of("id"); + Registration registration = Registration.create("foo", "https://health").build(); + InstanceEvent event = new InstanceRegisteredEvent(id, 0L, registration); + IMap> eventLog = hazelcastInstance + .getMap("testList" + System.currentTimeMillis()); + InstanceEventStore store = new HazelcastEventStore(100, eventLog); + + StepVerifier.create(store) + .expectSubscription() + .then(() -> eventLog.put(id, singletonList(event))) + .expectNext(event) + .thenCancel() + .verify(Duration.ofSeconds(10)); + } + finally { + hazelcastInstance.shutdown(); + } + } + } diff --git a/spring-boot-admin-server/src/test/java/de/codecentric/boot/admin/server/services/InfoUpdateTriggerTest.java b/spring-boot-admin-server/src/test/java/de/codecentric/boot/admin/server/services/InfoUpdateTriggerTest.java index ed9f55c5eaf..db067126aac 100644 --- a/spring-boot-admin-server/src/test/java/de/codecentric/boot/admin/server/services/InfoUpdateTriggerTest.java +++ b/spring-boot-admin-server/src/test/java/de/codecentric/boot/admin/server/services/InfoUpdateTriggerTest.java @@ -20,6 +20,7 @@ import org.junit.jupiter.api.BeforeEach; import org.junit.jupiter.api.Test; +import reactor.core.publisher.Flux; import reactor.core.publisher.Mono; import reactor.test.publisher.TestPublisher; @@ -170,4 +171,21 @@ void should_continue_update_after_error() { verify(this.updater, times(2)).updateInfo(this.instance.getId()); } + @Test + void should_update_existing_instances_on_start() { + this.trigger.stop(); + clearInvocations(this.updater); + + InfoUpdateTrigger trigger = new InfoUpdateTrigger(this.updater, Flux.empty(), Duration.ofDays(1), + Duration.ofDays(1), Duration.ofDays(1), Flux.just(this.instance.getId())); + try { + trigger.start(); + await().atMost(Duration.ofSeconds(1)) + .untilAsserted(() -> verify(this.updater, times(1)).updateInfo(this.instance.getId())); + } + finally { + trigger.stop(); + } + } + } diff --git a/spring-boot-admin-server/src/test/java/de/codecentric/boot/admin/server/services/StatusUpdateTriggerTest.java b/spring-boot-admin-server/src/test/java/de/codecentric/boot/admin/server/services/StatusUpdateTriggerTest.java index e1899409f20..cfc6c1a9e70 100644 --- a/spring-boot-admin-server/src/test/java/de/codecentric/boot/admin/server/services/StatusUpdateTriggerTest.java +++ b/spring-boot-admin-server/src/test/java/de/codecentric/boot/admin/server/services/StatusUpdateTriggerTest.java @@ -20,6 +20,7 @@ import org.junit.jupiter.api.BeforeEach; import org.junit.jupiter.api.Test; +import reactor.core.publisher.Flux; import reactor.core.publisher.Mono; import reactor.test.publisher.TestPublisher; @@ -147,4 +148,21 @@ void should_continue_update_after_error() { verify(this.updater, times(2)).updateStatus(this.instance.getId()); } + @Test + void should_update_existing_instances_on_start() { + this.trigger.stop(); + clearInvocations(this.updater); + + StatusUpdateTrigger trigger = new StatusUpdateTrigger(this.updater, Flux.empty(), Duration.ofDays(1), + Duration.ofDays(1), Duration.ofDays(1), Flux.just(this.instance.getId())); + try { + trigger.start(); + await().atMost(Duration.ofSeconds(1)) + .untilAsserted(() -> verify(this.updater, times(1)).updateStatus(this.instance.getId())); + } + finally { + trigger.stop(); + } + } + }