From 5499ea9ce79487578a4f34e80841ea12bfd851f9 Mon Sep 17 00:00:00 2001 From: Makarand Milind Hinge Date: Thu, 24 Sep 2026 16:57:49 +0530 Subject: [PATCH 1/3] GEODE-10440: add failing observation tests --- ...ernalDistributedSystemIntegrationTest.java | 52 +++++++ .../internal/cache/GemFireCacheImplTest.java | 10 ++ ...CacheBuilderAllowsMultipleSystemsTest.java | 2 +- .../cache/InternalCacheBuilderTest.java | 50 ++++++- .../util/InternalCacheBuilderTestUtil.java | 4 +- .../internal/GeodeObservationSupportTest.java | 133 +++++++++++++++++ ...edSystemObservationServiceBuilderTest.java | 116 +++++++++++++++ ...stributedSystemObservationServiceTest.java | 140 ++++++++++++++++++ .../ObservationRegistrySupplierTest.java | 82 ++++++++++ 9 files changed, 584 insertions(+), 5 deletions(-) create mode 100644 geode-core/src/test/java/org/apache/geode/metrics/internal/GeodeObservationSupportTest.java create mode 100644 geode-core/src/test/java/org/apache/geode/metrics/internal/InternalDistributedSystemObservationServiceBuilderTest.java create mode 100644 geode-core/src/test/java/org/apache/geode/metrics/internal/InternalDistributedSystemObservationServiceTest.java create mode 100644 geode-core/src/test/java/org/apache/geode/metrics/internal/ObservationRegistrySupplierTest.java diff --git a/geode-core/src/integrationTest/java/org/apache/geode/distributed/internal/InternalDistributedSystemIntegrationTest.java b/geode-core/src/integrationTest/java/org/apache/geode/distributed/internal/InternalDistributedSystemIntegrationTest.java index 96ba25a520fe..21734df0907e 100644 --- a/geode-core/src/integrationTest/java/org/apache/geode/distributed/internal/InternalDistributedSystemIntegrationTest.java +++ b/geode-core/src/integrationTest/java/org/apache/geode/distributed/internal/InternalDistributedSystemIntegrationTest.java @@ -56,6 +56,7 @@ import java.util.logging.Level; import io.micrometer.core.instrument.MeterRegistry; +import io.micrometer.observation.ObservationRegistry; import org.junit.jupiter.api.AfterEach; import org.junit.jupiter.api.Tag; import org.junit.jupiter.api.Test; @@ -70,6 +71,7 @@ import org.apache.geode.internal.ConfigSource; import org.apache.geode.internal.logging.InternalLogWriter; import org.apache.geode.metrics.internal.MetricsService; +import org.apache.geode.metrics.internal.ObservationService; import org.apache.geode.util.internal.GeodeGlossary; /** @@ -93,6 +95,15 @@ private InternalDistributedSystem createSystem(Properties props, return system; } + private InternalDistributedSystem createSystem(Properties props, + MetricsService.Builder metricsSessionBuilder, + ObservationService.Builder observationSessionBuilder) { + system = new InternalDistributedSystem.Builder(props, metricsSessionBuilder, + observationSessionBuilder) + .build(); + return system; + } + /** * Creates a DistributedSystem with the given configuration properties. */ @@ -740,6 +751,47 @@ public void getMeterRegistry_returnsMetricsSessionMeterRegistry() { assertThat(system.getMeterRegistry()).isSameAs(sessionMeterRegistry); } + @Test + public void usesSessionBuilderToCreateObservationSession() { + MetricsService.Builder metricsSessionBuilder = mock(MetricsService.Builder.class); + when(metricsSessionBuilder.build(any())).thenReturn(mock(MetricsService.class)); + ObservationService observationSession = mock(ObservationService.class); + ObservationService.Builder observationSessionBuilder = mock(ObservationService.Builder.class); + when(observationSessionBuilder.build(any())).thenReturn(observationSession); + + createSystem(getCommonProperties(), metricsSessionBuilder, observationSessionBuilder); + + verify(observationSessionBuilder).build(system); + } + + @Test + public void startsObservationSession() { + MetricsService.Builder metricsSessionBuilder = mock(MetricsService.Builder.class); + when(metricsSessionBuilder.build(any())).thenReturn(mock(MetricsService.class)); + ObservationService observationSession = mock(ObservationService.class); + ObservationService.Builder observationSessionBuilder = mock(ObservationService.Builder.class); + when(observationSessionBuilder.build(any())).thenReturn(observationSession); + + createSystem(getCommonProperties(), metricsSessionBuilder, observationSessionBuilder); + + verify(observationSession).start(); + } + + @Test + public void getObservationRegistry_returnsObservationSessionObservationRegistry() { + ObservationRegistry observationRegistry = mock(ObservationRegistry.class); + MetricsService.Builder metricsSessionBuilder = mock(MetricsService.Builder.class); + when(metricsSessionBuilder.build(any())).thenReturn(mock(MetricsService.class)); + ObservationService observationSession = mock(ObservationService.class); + when(observationSession.getObservationRegistry()).thenReturn(observationRegistry); + ObservationService.Builder observationSessionBuilder = mock(ObservationService.Builder.class); + when(observationSessionBuilder.build(any())).thenReturn(observationSession); + + createSystem(getCommonProperties(), metricsSessionBuilder, observationSessionBuilder); + + assertThat(system.getObservationRegistry()).isSameAs(observationRegistry); + } + @Test public void connect() { String theName = "theName"; diff --git a/geode-core/src/test/java/org/apache/geode/internal/cache/GemFireCacheImplTest.java b/geode-core/src/test/java/org/apache/geode/internal/cache/GemFireCacheImplTest.java index 901e37112071..d12eaea1360c 100644 --- a/geode-core/src/test/java/org/apache/geode/internal/cache/GemFireCacheImplTest.java +++ b/geode-core/src/test/java/org/apache/geode/internal/cache/GemFireCacheImplTest.java @@ -43,6 +43,7 @@ import java.util.stream.IntStream; import io.micrometer.core.instrument.MeterRegistry; +import io.micrometer.observation.ObservationRegistry; import org.junit.After; import org.junit.Before; import org.junit.Rule; @@ -453,6 +454,15 @@ public void getMeterRegistry_returnsTheSystemMeterRegistry() { .isSameAs(systemMeterRegistry); } + @Test + public void getObservationRegistry_returnsTheSystemObservationRegistry() { + ObservationRegistry systemObservationRegistry = mock(ObservationRegistry.class); + when(internalDistributedSystem.getObservationRegistry()).thenReturn(systemObservationRegistry); + + assertThat(gemFireCacheImpl.getObservationRegistry()) + .isSameAs(systemObservationRegistry); + } + @Test public void addGatewayReceiverServer_requiresPreviouslyAddedGatewayReceiver() { Throwable thrown = catchThrowable( diff --git a/geode-core/src/test/java/org/apache/geode/internal/cache/InternalCacheBuilderAllowsMultipleSystemsTest.java b/geode-core/src/test/java/org/apache/geode/internal/cache/InternalCacheBuilderAllowsMultipleSystemsTest.java index a9c0edb44b1e..604a4afc7f35 100644 --- a/geode-core/src/test/java/org/apache/geode/internal/cache/InternalCacheBuilderAllowsMultipleSystemsTest.java +++ b/geode-core/src/test/java/org/apache/geode/internal/cache/InternalCacheBuilderAllowsMultipleSystemsTest.java @@ -109,7 +109,7 @@ configProperties, new CacheConfig(), metricsSessionBuilder, THROWING_SYSTEM_SUPP internalCacheBuilder.create(); - verify(systemConstructor).construct(same(configProperties), any(), any()); + verify(systemConstructor).construct(same(configProperties), any(), any(), any()); } @Test diff --git a/geode-core/src/test/java/org/apache/geode/internal/cache/InternalCacheBuilderTest.java b/geode-core/src/test/java/org/apache/geode/internal/cache/InternalCacheBuilderTest.java index 76d9331c93b2..07526da25f3b 100644 --- a/geode-core/src/test/java/org/apache/geode/internal/cache/InternalCacheBuilderTest.java +++ b/geode-core/src/test/java/org/apache/geode/internal/cache/InternalCacheBuilderTest.java @@ -56,6 +56,7 @@ import org.apache.geode.internal.cache.InternalCacheBuilder.InternalCacheConstructor; import org.apache.geode.internal.cache.InternalCacheBuilder.InternalDistributedSystemConstructor; import org.apache.geode.metrics.internal.MetricsService; +import org.apache.geode.metrics.internal.ObservationService; /** * Unit tests for {@link InternalCacheBuilder}. @@ -71,6 +72,9 @@ public class InternalCacheBuilderTest { @Mock private MetricsService.Builder metricsServiceBuilder; + @Mock + private ObservationService.Builder observationServiceBuilder; + @Before public void setUp() { initMocks(this); @@ -90,6 +94,19 @@ nullSingletonSystemSupplier, constructorOf(constructedSystem()), nullSingletonCa verify(theMetricsServiceBuilder).setIsClient(false); } + @Test + public void setsObservationServiceBuilderIsClientFalseByDefault() { + ObservationService.Builder theObservationServiceBuilder = + mock(ObservationService.Builder.class); + + new InternalCacheBuilder(new Properties(), new CacheConfig(), metricsServiceBuilder, + theObservationServiceBuilder, nullSingletonSystemSupplier, + constructorOf(constructedSystem()), + nullSingletonCacheSupplier, constructorOf(constructedCache())); + + verify(theObservationServiceBuilder).setIsClient(false); + } + @Test public void addMeterSubregistry_addsGivenRegistryToMetricsServiceBuilder() { InternalCacheBuilder internalCacheBuilder = new InternalCacheBuilder( @@ -129,6 +146,21 @@ public void setIsClient_setsIsClientInMetricsServiceBuilder() { verify(theMetricsServiceBuilder).setIsClient(true); } + @Test + public void setIsClient_setsIsClientInObservationServiceBuilder() { + ObservationService.Builder theObservationServiceBuilder = + mock(ObservationService.Builder.class); + + InternalCacheBuilder internalCacheBuilder = new InternalCacheBuilder( + new Properties(), new CacheConfig(), metricsServiceBuilder, theObservationServiceBuilder, + nullSingletonSystemSupplier, constructorOf(constructedSystem()), nullSingletonCacheSupplier, + constructorOf(constructedCache())); + + internalCacheBuilder.setIsClient(true); + + verify(theObservationServiceBuilder).setIsClient(true); + } + @Test public void create_throwsNullPointerException_ifConfigPropertiesIsNull() { InternalCacheBuilder internalCacheBuilder = new InternalCacheBuilder( @@ -163,7 +195,7 @@ configProperties, new CacheConfig(), metricsServiceBuilder, nullSingletonSystemS internalCacheBuilder .create(); - verify(systemConstructor).construct(same(configProperties), any(), any()); + verify(systemConstructor).construct(same(configProperties), any(), any(), any()); } @Test @@ -192,7 +224,21 @@ public void create_passesMetricsServiceBuilderToSystemConstructor_ifNoSystemExis internalCacheBuilder.create(); - verify(systemConstructor).construct(any(), any(), same(theMetricsServiceBuilder)); + verify(systemConstructor).construct(any(), any(), same(theMetricsServiceBuilder), any()); + } + + @Test + public void create_passesObservationServiceBuilderToSystemConstructor_ifNoSystemExists() { + InternalDistributedSystemConstructor systemConstructor = constructorOf(constructedSystem()); + + InternalCacheBuilder internalCacheBuilder = new InternalCacheBuilder( + new Properties(), new CacheConfig(), metricsServiceBuilder, observationServiceBuilder, + nullSingletonSystemSupplier, systemConstructor, nullSingletonCacheSupplier, + constructorOf(constructedCache())); + + internalCacheBuilder.create(); + + verify(systemConstructor).construct(any(), any(), any(), same(observationServiceBuilder)); } @Test diff --git a/geode-core/src/test/java/org/apache/geode/internal/util/InternalCacheBuilderTestUtil.java b/geode-core/src/test/java/org/apache/geode/internal/util/InternalCacheBuilderTestUtil.java index eed9a21eb90d..aac104249e37 100644 --- a/geode-core/src/test/java/org/apache/geode/internal/util/InternalCacheBuilderTestUtil.java +++ b/geode-core/src/test/java/org/apache/geode/internal/util/InternalCacheBuilderTestUtil.java @@ -104,7 +104,7 @@ public static InternalDistributedSystemConstructor constructorOf( InternalDistributedSystem constructedSystem) { InternalDistributedSystemConstructor constructor = mock(InternalDistributedSystemConstructor.class, "internal distributed system constructor"); - when(constructor.construct(any(), any(), any())).thenReturn(constructedSystem); + when(constructor.construct(any(), any(), any(), any())).thenReturn(constructedSystem); return constructor; } @@ -137,7 +137,7 @@ public static CacheConfig throwingCacheConfig(Throwable throwable) { }; public static final InternalDistributedSystemConstructor THROWING_SYSTEM_CONSTRUCTOR = - (configProperties, securityConfig, userMeterRegistries) -> { + (configProperties, securityConfig, userMeterRegistries, observationServiceBuilder) -> { throw new AssertionError("throwing system constructor"); }; diff --git a/geode-core/src/test/java/org/apache/geode/metrics/internal/GeodeObservationSupportTest.java b/geode-core/src/test/java/org/apache/geode/metrics/internal/GeodeObservationSupportTest.java new file mode 100644 index 000000000000..0a3c390dd4ce --- /dev/null +++ b/geode-core/src/test/java/org/apache/geode/metrics/internal/GeodeObservationSupportTest.java @@ -0,0 +1,133 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more contributor license + * agreements. See the NOTICE file distributed with this work for additional information regarding + * copyright ownership. The ASF licenses this file to You under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance with the License. You may obtain a + * copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software distributed under the License + * is distributed on an "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express + * or implied. See the License for the specific language governing permissions and limitations under + * the License. + */ +package org.apache.geode.metrics.internal; + +import static org.assertj.core.api.Assertions.assertThat; +import static org.assertj.core.api.Assertions.assertThatThrownBy; + +import java.util.ArrayList; +import java.util.List; + +import io.micrometer.observation.Observation; +import io.micrometer.observation.ObservationHandler; +import io.micrometer.observation.ObservationRegistry; +import org.junit.Before; +import org.junit.Test; + +public class GeodeObservationSupportTest { + private final ObservationRegistry observationRegistry = ObservationRegistry.create(); + private final RecordingObservationHandler observationHandler = new RecordingObservationHandler(); + + @Before + public void setUp() { + observationRegistry.observationConfig().observationHandler(observationHandler); + } + + @Test + public void observeRunnable_startsAndStopsObservation() { + GeodeObservationSupport.observe(observationRegistry, "geode.test.operation", + observation -> observation.lowCardinalityKeyValue("operation", "test"), + () -> { + }); + + assertThat(observationHandler.events) + .containsExactly("start:geode.test.operation", "stop:geode.test.operation"); + assertThat(observationHandler.lowCardinalityKeyValue("operation")) + .isEqualTo("test"); + } + + @Test + public void observeRunnable_recordsErrorAndStopsObservation() { + RuntimeException failure = new RuntimeException("expected failure"); + + assertThatThrownBy(() -> GeodeObservationSupport.observe(observationRegistry, + "geode.test.operation", observation -> { + }, () -> { + throw failure; + })).isSameAs(failure); + + assertThat(observationHandler.events) + .containsExactly("start:geode.test.operation", "error:geode.test.operation", + "stop:geode.test.operation"); + assertThat(observationHandler.error) + .isSameAs(failure); + } + + @Test + public void observeSupplier_returnsSupplierValue() { + String value = GeodeObservationSupport.observe(observationRegistry, "geode.test.operation", + observation -> { + }, () -> "result"); + + assertThat(value) + .isEqualTo("result"); + assertThat(observationHandler.events) + .containsExactly("start:geode.test.operation", "stop:geode.test.operation"); + } + + @Test + public void observeSupplier_recordsErrorAndStopsObservation() { + RuntimeException failure = new RuntimeException("expected failure"); + + assertThatThrownBy(() -> GeodeObservationSupport.observe(observationRegistry, + "geode.test.operation", observation -> { + }, () -> { + throw failure; + })).isSameAs(failure); + + assertThat(observationHandler.events) + .containsExactly("start:geode.test.operation", "error:geode.test.operation", + "stop:geode.test.operation"); + assertThat(observationHandler.error) + .isSameAs(failure); + } + + private static class RecordingObservationHandler + implements ObservationHandler { + private final List events = new ArrayList<>(); + private Observation.Context context; + private Throwable error; + + @Override + public void onStart(Observation.Context context) { + this.context = context; + events.add("start:" + context.getName()); + } + + @Override + public void onError(Observation.Context context) { + error = context.getError(); + events.add("error:" + context.getName()); + } + + @Override + public void onStop(Observation.Context context) { + events.add("stop:" + context.getName()); + } + + @Override + public boolean supportsContext(Observation.Context context) { + return true; + } + + private String lowCardinalityKeyValue(String key) { + return context.getLowCardinalityKeyValues().stream() + .filter(keyValue -> keyValue.getKey().equals(key)) + .findFirst() + .map(keyValue -> keyValue.getValue()) + .orElse(null); + } + } +} diff --git a/geode-core/src/test/java/org/apache/geode/metrics/internal/InternalDistributedSystemObservationServiceBuilderTest.java b/geode-core/src/test/java/org/apache/geode/metrics/internal/InternalDistributedSystemObservationServiceBuilderTest.java new file mode 100644 index 000000000000..178d88a084f0 --- /dev/null +++ b/geode-core/src/test/java/org/apache/geode/metrics/internal/InternalDistributedSystemObservationServiceBuilderTest.java @@ -0,0 +1,116 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more contributor license + * agreements. See the NOTICE file distributed with this work for additional information regarding + * copyright ownership. The ASF licenses this file to You under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance with the License. You may obtain a + * copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software distributed under the License + * is distributed on an "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express + * or implied. See the License for the specific language governing permissions and limitations under + * the License. + */ +package org.apache.geode.metrics.internal; + +import static org.assertj.core.api.Assertions.assertThat; +import static org.mockito.ArgumentMatchers.any; +import static org.mockito.ArgumentMatchers.same; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.verify; +import static org.mockito.Mockito.when; + +import io.micrometer.observation.ObservationRegistry; +import org.apache.logging.log4j.Logger; +import org.junit.Rule; +import org.junit.Test; +import org.mockito.Mock; +import org.mockito.junit.MockitoJUnit; +import org.mockito.junit.MockitoRule; + +import org.apache.geode.distributed.internal.InternalDistributedSystem; +import org.apache.geode.internal.util.CollectingServiceLoader; +import org.apache.geode.metrics.ObservationPublishingService; + +public class InternalDistributedSystemObservationServiceBuilderTest { + @Rule + public MockitoRule mockitoRule = MockitoJUnit.rule(); + + @Mock + private InternalDistributedSystem system; + + @Mock + private InternalDistributedSystemObservationService.Factory observationServiceFactory; + + private final InternalDistributedSystemObservationService.Builder serviceBuilder = + new InternalDistributedSystemObservationService.Builder(); + + @Test + public void createsInternalDistributedSystemObservationService() { + ObservationService observationService = serviceBuilder.build(system); + + assertThat(observationService) + .isInstanceOf(InternalDistributedSystemObservationService.class); + } + + @Test + public void usesFactoryToCreateSession_ifFactorySet() { + ObservationService observationServiceCreatedByFactory = mock(ObservationService.class); + when(observationServiceFactory.create(any(), any(), any(), any())) + .thenReturn(observationServiceCreatedByFactory); + + ObservationService observationService = serviceBuilder + .setObservationServiceFactory(observationServiceFactory) + .build(system); + + assertThat(observationService) + .isSameAs(observationServiceCreatedByFactory); + } + + @Test + public void passesItselfToFactory() { + serviceBuilder.setObservationServiceFactory(observationServiceFactory) + .build(system); + + verify(observationServiceFactory) + .create(same(serviceBuilder), any(), any(), any()); + } + + @Test + public void passesGivenServiceLoaderToFactory() { + CollectingServiceLoader serviceLoader = + mock(CollectingServiceLoader.class); + + serviceBuilder.setObservationServiceFactory(observationServiceFactory) + .setServiceLoader(serviceLoader) + .build(system); + + verify(observationServiceFactory) + .create(any(), any(), same(serviceLoader), any()); + } + + @Test + public void passesGivenObservationRegistryToFactory() { + ObservationRegistry observationRegistry = ObservationRegistry.create(); + + serviceBuilder.setObservationServiceFactory(observationServiceFactory) + .setObservationRegistry(observationRegistry) + .build(system); + + verify(observationServiceFactory) + .create(any(), any(), any(), same(observationRegistry)); + } + + @Test + public void passesGivenLoggerToFactory() { + Logger logger = mock(Logger.class); + + serviceBuilder.setObservationServiceFactory(observationServiceFactory) + .setLogger(logger) + .build(system); + + verify(observationServiceFactory) + .create(any(), same(logger), any(), any()); + } +} diff --git a/geode-core/src/test/java/org/apache/geode/metrics/internal/InternalDistributedSystemObservationServiceTest.java b/geode-core/src/test/java/org/apache/geode/metrics/internal/InternalDistributedSystemObservationServiceTest.java new file mode 100644 index 000000000000..099c3a727f0f --- /dev/null +++ b/geode-core/src/test/java/org/apache/geode/metrics/internal/InternalDistributedSystemObservationServiceTest.java @@ -0,0 +1,140 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more contributor license + * agreements. See the NOTICE file distributed with this work for additional information regarding + * copyright ownership. The ASF licenses this file to You under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance with the License. You may obtain a + * copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software distributed under the License + * is distributed on an "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express + * or implied. See the License for the specific language governing permissions and limitations under + * the License. + */ +package org.apache.geode.metrics.internal; + +import static java.util.Collections.singletonList; +import static org.assertj.core.api.Assertions.assertThat; +import static org.mockito.ArgumentMatchers.contains; +import static org.mockito.ArgumentMatchers.same; +import static org.mockito.Mockito.doThrow; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.verify; +import static org.mockito.Mockito.when; + +import io.micrometer.observation.ObservationRegistry; +import org.apache.logging.log4j.Logger; +import org.junit.Rule; +import org.junit.Test; +import org.mockito.Mock; +import org.mockito.junit.MockitoJUnit; +import org.mockito.junit.MockitoRule; + +import org.apache.geode.internal.util.CollectingServiceLoader; +import org.apache.geode.metrics.ObservationPublishingService; + +public class InternalDistributedSystemObservationServiceTest { + @Rule + public MockitoRule mockitoRule = MockitoJUnit.rule(); + + @Mock + private Logger logger; + + @Mock + private CollectingServiceLoader publishingServiceLoader; + + @Mock + private ObservationService.Builder observationServiceBuilder; + + private final ObservationRegistry observationRegistry = ObservationRegistry.create(); + + @Test + public void remembersObservationRegistry() { + ObservationService observationService = + new InternalDistributedSystemObservationService(observationServiceBuilder, logger, + publishingServiceLoader, observationRegistry); + + assertThat(observationService.getObservationRegistry()) + .isSameAs(observationRegistry); + } + + @Test + public void remembersObservationServiceBuilder() { + ObservationService.Builder builder = mock(ObservationService.Builder.class); + + ObservationService observationService = + new InternalDistributedSystemObservationService(builder, logger, publishingServiceLoader, + observationRegistry); + + assertThat(observationService.getRebuilder()) + .isSameAs(builder); + } + + @Test + public void start_loadsAndStartsObservationPublishingServices() { + ObservationPublishingService publishingService = mock(ObservationPublishingService.class); + when(publishingServiceLoader.loadServices(ObservationPublishingService.class)) + .thenReturn(singletonList(publishingService)); + + ObservationService observationService = + new InternalDistributedSystemObservationService(observationServiceBuilder, logger, + publishingServiceLoader, observationRegistry); + + observationService.start(); + + verify(publishingService).start(same(observationService)); + } + + @Test + public void start_logsAndContinues_ifPublishingServiceThrows() { + ObservationPublishingService publishingService = mock(ObservationPublishingService.class); + RuntimeException failure = new RuntimeException("start failure"); + when(publishingServiceLoader.loadServices(ObservationPublishingService.class)) + .thenReturn(singletonList(publishingService)); + ObservationService observationService = + new InternalDistributedSystemObservationService(observationServiceBuilder, logger, + publishingServiceLoader, observationRegistry); + + doThrow(failure).when(publishingService).start(same(observationService)); + observationService.start(); + + verify(logger).error(contains("Exception while starting observation publishing service"), + same(failure)); + } + + @Test + public void stop_stopsObservationPublishingServices() { + ObservationPublishingService publishingService = mock(ObservationPublishingService.class); + when(publishingServiceLoader.loadServices(ObservationPublishingService.class)) + .thenReturn(singletonList(publishingService)); + + ObservationService observationService = + new InternalDistributedSystemObservationService(observationServiceBuilder, logger, + publishingServiceLoader, observationRegistry); + + observationService.start(); + observationService.stop(); + + verify(publishingService).stop(same(observationService)); + } + + @Test + public void stop_logsAndContinues_ifPublishingServiceThrows() { + ObservationPublishingService publishingService = mock(ObservationPublishingService.class); + RuntimeException failure = new RuntimeException("stop failure"); + when(publishingServiceLoader.loadServices(ObservationPublishingService.class)) + .thenReturn(singletonList(publishingService)); + + ObservationService observationService = + new InternalDistributedSystemObservationService(observationServiceBuilder, logger, + publishingServiceLoader, observationRegistry); + + observationService.start(); + doThrow(failure).when(publishingService).stop(same(observationService)); + observationService.stop(); + + verify(logger).error(contains("Exception while stopping observation publishing service"), + same(failure)); + } +} diff --git a/geode-core/src/test/java/org/apache/geode/metrics/internal/ObservationRegistrySupplierTest.java b/geode-core/src/test/java/org/apache/geode/metrics/internal/ObservationRegistrySupplierTest.java new file mode 100644 index 000000000000..15c8c7e53cd3 --- /dev/null +++ b/geode-core/src/test/java/org/apache/geode/metrics/internal/ObservationRegistrySupplierTest.java @@ -0,0 +1,82 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more contributor license + * agreements. See the NOTICE file distributed with this work for additional information regarding + * copyright ownership. The ASF licenses this file to You under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance with the License. You may obtain a + * copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software distributed under the License + * is distributed on an "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express + * or implied. See the License for the specific language governing permissions and limitations under + * the License. + */ +package org.apache.geode.metrics.internal; + +import static org.assertj.core.api.Assertions.assertThat; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.when; + +import io.micrometer.observation.ObservationRegistry; +import org.junit.Test; + +import org.apache.geode.distributed.internal.InternalDistributedSystem; +import org.apache.geode.internal.cache.InternalCache; + +public class ObservationRegistrySupplierTest { + @Test + public void get_internalDistributedSystemIsNull_expectNull() { + ObservationRegistrySupplier observationRegistrySupplier = new ObservationRegistrySupplier( + () -> null); + + ObservationRegistry value = observationRegistrySupplier.get(); + + assertThat(value) + .isNull(); + } + + @Test + public void get_internalCacheIsNull_expectNull() { + InternalDistributedSystem internalDistributedSystem = mock(InternalDistributedSystem.class); + when(internalDistributedSystem.getCache()).thenReturn(null); + ObservationRegistrySupplier observationRegistrySupplier = + new ObservationRegistrySupplier(() -> internalDistributedSystem); + + ObservationRegistry value = observationRegistrySupplier.get(); + + assertThat(value) + .isNull(); + } + + @Test + public void get_observationRegistryIsNull_expectNull() { + InternalDistributedSystem internalDistributedSystem = mock(InternalDistributedSystem.class); + InternalCache internalCache = mock(InternalCache.class); + when(internalDistributedSystem.getCache()).thenReturn(internalCache); + when(internalCache.getObservationRegistry()).thenReturn(null); + ObservationRegistrySupplier observationRegistrySupplier = + new ObservationRegistrySupplier(() -> internalDistributedSystem); + + ObservationRegistry value = observationRegistrySupplier.get(); + + assertThat(value) + .isNull(); + } + + @Test + public void get_observationRegistryExists_expectActualObservationRegistry() { + InternalDistributedSystem internalDistributedSystem = mock(InternalDistributedSystem.class); + InternalCache internalCache = mock(InternalCache.class); + ObservationRegistry observationRegistry = mock(ObservationRegistry.class); + when(internalDistributedSystem.getCache()).thenReturn(internalCache); + when(internalCache.getObservationRegistry()).thenReturn(observationRegistry); + ObservationRegistrySupplier observationRegistrySupplier = + new ObservationRegistrySupplier(() -> internalDistributedSystem); + + ObservationRegistry value = observationRegistrySupplier.get(); + + assertThat(value) + .isSameAs(observationRegistry); + } +} From cb87c2c72429990f9b911b941143f9498b7b2d2b Mon Sep 17 00:00:00 2001 From: Makarand Milind Hinge Date: Thu, 24 Sep 2026 17:33:57 +0530 Subject: [PATCH 2/3] GEODE-10440: add Micrometer Observation support --- .../src/test/resources/expected-pom.xml | 5 + .../plugins/DependencyConstraints.groovy | 1 + geode-core/build.gradle | 1 + .../client/internal/InternalClientCache.java | 3 + .../internal/InternalDistributedSystem.java | 58 ++++++- .../internal/cache/GemFireCacheImpl.java | 6 + .../geode/internal/cache/InternalCache.java | 3 + .../internal/cache/InternalCacheBuilder.java | 27 +++- .../cache/InternalCacheForClientAccess.java | 6 + .../geode/internal/cache/LocalRegion.java | 38 +++-- .../cache/tier/sockets/ServerConnection.java | 7 +- .../cache/xmlcache/CacheCreation.java | 6 + .../metrics/ObservationPublishingService.java | 53 +++++++ .../geode/metrics/ObservationSession.java | 36 +++++ .../internal/GeodeObservationSupport.java | 73 +++++++++ ...alDistributedSystemObservationService.java | 143 ++++++++++++++++++ .../internal/ObservationRegistrySupplier.java | 47 ++++++ .../metrics/internal/ObservationService.java | 65 ++++++++ .../src/test/resources/expected-pom.xml | 15 ++ 19 files changed, 568 insertions(+), 25 deletions(-) create mode 100644 geode-core/src/main/java/org/apache/geode/metrics/ObservationPublishingService.java create mode 100644 geode-core/src/main/java/org/apache/geode/metrics/ObservationSession.java create mode 100644 geode-core/src/main/java/org/apache/geode/metrics/internal/GeodeObservationSupport.java create mode 100644 geode-core/src/main/java/org/apache/geode/metrics/internal/InternalDistributedSystemObservationService.java create mode 100644 geode-core/src/main/java/org/apache/geode/metrics/internal/ObservationRegistrySupplier.java create mode 100644 geode-core/src/main/java/org/apache/geode/metrics/internal/ObservationService.java diff --git a/boms/geode-all-bom/src/test/resources/expected-pom.xml b/boms/geode-all-bom/src/test/resources/expected-pom.xml index 84d481c32760..22511894e3d9 100644 --- a/boms/geode-all-bom/src/test/resources/expected-pom.xml +++ b/boms/geode-all-bom/src/test/resources/expected-pom.xml @@ -197,6 +197,11 @@ micrometer-core 1.16.7 + + io.micrometer + micrometer-observation + 1.16.7 + io.swagger.core.v3 swagger-annotations diff --git a/build-tools/geode-dependency-management/src/main/groovy/org/apache/geode/gradle/plugins/DependencyConstraints.groovy b/build-tools/geode-dependency-management/src/main/groovy/org/apache/geode/gradle/plugins/DependencyConstraints.groovy index 1a6e6f032e63..6e9427d3e6b5 100644 --- a/build-tools/geode-dependency-management/src/main/groovy/org/apache/geode/gradle/plugins/DependencyConstraints.groovy +++ b/build-tools/geode-dependency-management/src/main/groovy/org/apache/geode/gradle/plugins/DependencyConstraints.groovy @@ -152,6 +152,7 @@ class DependencyConstraints { api(group: 'io.github.resilience4j', name: 'resilience4j-retry', version: '1.7.1') api(group: 'io.lettuce', name: 'lettuce-core', version: '6.1.8.RELEASE') api(group: 'io.micrometer', name: 'micrometer-core', version: get('micrometer.version')) + api(group: 'io.micrometer', name: 'micrometer-observation', version: get('micrometer.version')) // Pin Reactor Core (pulled in via spring-shell-core) to 3.8.7 api(group: 'io.projectreactor', name: 'reactor-core', version: get('reactor-core.version')) api(group: 'io.swagger.core.v3', name: 'swagger-annotations', version: '2.2.22') diff --git a/geode-core/build.gradle b/geode-core/build.gradle index 453c32139306..dabe25765c51 100755 --- a/geode-core/build.gradle +++ b/geode-core/build.gradle @@ -264,6 +264,7 @@ dependencies { //micrometer is used for micrometer based metrics from geode geode api('io.micrometer:micrometer-core') + api('io.micrometer:micrometer-observation') //FastUtil contains optimized collections that are used in multiple places in core diff --git a/geode-core/src/main/java/org/apache/geode/cache/client/internal/InternalClientCache.java b/geode-core/src/main/java/org/apache/geode/cache/client/internal/InternalClientCache.java index 7e10090ccc20..1b1523778c78 100644 --- a/geode-core/src/main/java/org/apache/geode/cache/client/internal/InternalClientCache.java +++ b/geode-core/src/main/java/org/apache/geode/cache/client/internal/InternalClientCache.java @@ -15,6 +15,7 @@ package org.apache.geode.cache.client.internal; import io.micrometer.core.instrument.MeterRegistry; +import io.micrometer.observation.ObservationRegistry; import org.apache.geode.cache.Region; import org.apache.geode.cache.RegionAttributes; @@ -49,5 +50,7 @@ Region basicCreateRegion(String name, RegionAttributes attrs) MeterRegistry getMeterRegistry(); + ObservationRegistry getObservationRegistry(); + ClientMetadataService getClientMetadataService(); } diff --git a/geode-core/src/main/java/org/apache/geode/distributed/internal/InternalDistributedSystem.java b/geode-core/src/main/java/org/apache/geode/distributed/internal/InternalDistributedSystem.java index cf9043979163..3e3b54ed0ad8 100644 --- a/geode-core/src/main/java/org/apache/geode/distributed/internal/InternalDistributedSystem.java +++ b/geode-core/src/main/java/org/apache/geode/distributed/internal/InternalDistributedSystem.java @@ -42,6 +42,7 @@ import java.util.concurrent.atomic.AtomicReference; import io.micrometer.core.instrument.MeterRegistry; +import io.micrometer.observation.ObservationRegistry; import org.apache.logging.log4j.Logger; import org.apache.geode.CancelCriterion; @@ -117,8 +118,10 @@ import org.apache.geode.logging.internal.spi.LogConfigSupplier; import org.apache.geode.logging.internal.spi.LogFile; import org.apache.geode.management.ManagementException; +import org.apache.geode.metrics.internal.InternalDistributedSystemObservationService; import org.apache.geode.metrics.internal.MeterRegistrySupplier; import org.apache.geode.metrics.internal.MetricsService; +import org.apache.geode.metrics.internal.ObservationService; import org.apache.geode.pdx.internal.TypeRegistry; import org.apache.geode.security.GemFireSecurityException; import org.apache.geode.security.PostProcessor; @@ -170,6 +173,7 @@ public class InternalDistributedSystem extends DistributedSystem private final StatisticsManager statisticsManager; private MetricsService metricsService; + private ObservationService observationService; private final FunctionStatsManager functionStatsManager; /** * True if the user is allowed lock when memory resources appear to be overcommitted. @@ -206,7 +210,17 @@ public static InternalDistributedSystem connectInternal( Properties config, SecurityConfig securityConfig, MetricsService.Builder metricsSessionBuilder) { - return connectInternal(config, securityConfig, metricsSessionBuilder, null); + return connectInternal(config, securityConfig, metricsSessionBuilder, + new InternalDistributedSystemObservationService.Builder(), null); + } + + public static InternalDistributedSystem connectInternal( + Properties config, + SecurityConfig securityConfig, + MetricsService.Builder metricsSessionBuilder, + ObservationService.Builder observationSessionBuilder) { + return connectInternal(config, securityConfig, metricsSessionBuilder, observationSessionBuilder, + null); } /** @@ -221,12 +235,22 @@ public static InternalDistributedSystem connectInternal( SecurityConfig securityConfig, MetricsService.Builder metricsSessionBuilder, final MembershipLocator locator) { + return connectInternal(config, securityConfig, metricsSessionBuilder, + new InternalDistributedSystemObservationService.Builder(), locator); + } + + public static InternalDistributedSystem connectInternal( + Properties config, + SecurityConfig securityConfig, + MetricsService.Builder metricsSessionBuilder, + ObservationService.Builder observationSessionBuilder, + final MembershipLocator locator) { if (config == null) { config = new Properties(); } if (Boolean.getBoolean(ALLOW_MULTIPLE_SYSTEMS_PROPERTY)) { - return new Builder(config, metricsSessionBuilder) + return new Builder(config, metricsSessionBuilder, observationSessionBuilder) .setSecurityConfig(securityConfig) .setLocator(locator) .build(); @@ -277,10 +301,11 @@ public static InternalDistributedSystem connectInternal( } // Make a new connection to the distributed system - InternalDistributedSystem newSystem = new Builder(config, metricsSessionBuilder) - .setSecurityConfig(securityConfig) - .setLocator(locator) - .build(); + InternalDistributedSystem newSystem = + new Builder(config, metricsSessionBuilder, observationSessionBuilder) + .setSecurityConfig(securityConfig) + .setLocator(locator) + .build(); addSystem(newSystem); return newSystem; } @@ -654,6 +679,7 @@ public MemoryAllocator getOffHeapStore() { @VisibleForTesting void initialize(SecurityManager securityManager, PostProcessor postProcessor, MetricsService.Builder metricsServiceBuilder, + ObservationService.Builder observationServiceBuilder, final MembershipLocator membershipLocatorArg, ClusterDistributionManagerConstructor clusterDistributionManagerConstructor) { @@ -805,6 +831,8 @@ void initialize(SecurityManager securityManager, PostProcessor postProcessor, metricsService = metricsServiceBuilder.build(this); metricsService.start(); + observationService = observationServiceBuilder.build(this); + observationService.start(); // Log any instantiators that were registered before the log writer // was created @@ -1161,6 +1189,10 @@ public MeterRegistry getMeterRegistry() { return metricsService.getMeterRegistry(); } + public ObservationRegistry getObservationRegistry() { + return observationService.getObservationRegistry(); + } + /** * This class defers to the DM. If we don't have a DM, we're dead. */ @@ -1618,6 +1650,7 @@ protected void disconnect(boolean preparingForReconnect, String reason, boolean functionStatsManager.close(); metricsService.stop(); + observationService.stop(); InternalFunctionService.unregisterAllFunctions(); @@ -2593,7 +2626,7 @@ private void reconnect(boolean forcedDisconnect, String reason) { try { newDS = connectInternal(configProps, null, metricsService.getRebuilder(), - membershipLocator); + observationService.getRebuilder(), membershipLocator); } catch (CancelException e) { if (isReconnectCancelled()) { @@ -3001,6 +3034,7 @@ public static class Builder { private SecurityConfig securityConfig; private final MetricsService.Builder metricsServiceBuilder; + private final ObservationService.Builder observationServiceBuilder; private MembershipLocator locator; @@ -3008,8 +3042,15 @@ public static class Builder { new DefaultClusterDistributionManagerConstructor(); public Builder(Properties configProperties, MetricsService.Builder metricsServiceBuilder) { + this(configProperties, metricsServiceBuilder, + new InternalDistributedSystemObservationService.Builder()); + } + + public Builder(Properties configProperties, MetricsService.Builder metricsServiceBuilder, + ObservationService.Builder observationServiceBuilder) { this.configProperties = configProperties; this.metricsServiceBuilder = metricsServiceBuilder; + this.observationServiceBuilder = observationServiceBuilder; } public Builder setSecurityConfig(SecurityConfig securityConfig) { @@ -3047,7 +3088,8 @@ public InternalDistributedSystem build() { FunctionStatsManager::new); newSystem .initialize(securityConfig.getSecurityManager(), securityConfig.getPostProcessor(), - metricsServiceBuilder, locator, clusterDistributionManagerConstructor); + metricsServiceBuilder, observationServiceBuilder, locator, + clusterDistributionManagerConstructor); notifyConnectListeners(newSystem); stopThreads = false; return newSystem; diff --git a/geode-core/src/main/java/org/apache/geode/internal/cache/GemFireCacheImpl.java b/geode-core/src/main/java/org/apache/geode/internal/cache/GemFireCacheImpl.java index 6e6c6c568df4..3918e0626917 100755 --- a/geode-core/src/main/java/org/apache/geode/internal/cache/GemFireCacheImpl.java +++ b/geode-core/src/main/java/org/apache/geode/internal/cache/GemFireCacheImpl.java @@ -111,6 +111,7 @@ import com.sun.jna.Native; import com.sun.jna.Platform; import io.micrometer.core.instrument.MeterRegistry; +import io.micrometer.observation.ObservationRegistry; import jakarta.transaction.TransactionManager; import org.apache.commons.lang3.StringUtils; import org.apache.logging.log4j.Logger; @@ -1174,6 +1175,11 @@ public MeterRegistry getMeterRegistry() { return system.getMeterRegistry(); } + @Override + public ObservationRegistry getObservationRegistry() { + return system.getObservationRegistry(); + } + @Override public void saveCacheXmlForReconnect() { prepareForReconnect((pw) -> CacheXmlGenerator.generate((Cache) this, pw, false)); diff --git a/geode-core/src/main/java/org/apache/geode/internal/cache/InternalCache.java b/geode-core/src/main/java/org/apache/geode/internal/cache/InternalCache.java index 2607b603bb82..7f9b6bb6f900 100644 --- a/geode-core/src/main/java/org/apache/geode/internal/cache/InternalCache.java +++ b/geode-core/src/main/java/org/apache/geode/internal/cache/InternalCache.java @@ -26,6 +26,7 @@ import java.util.concurrent.Executor; import io.micrometer.core.instrument.MeterRegistry; +import io.micrometer.observation.ObservationRegistry; import jakarta.transaction.TransactionManager; import org.apache.geode.cache.Cache; @@ -565,6 +566,8 @@ void invokeRegionEntrySynchronizationListenersAfterSynchronization( MeterRegistry getMeterRegistry(); + ObservationRegistry getObservationRegistry(); + /** * Generate XML for the cache before shutting down due to forced disconnect. */ diff --git a/geode-core/src/main/java/org/apache/geode/internal/cache/InternalCacheBuilder.java b/geode-core/src/main/java/org/apache/geode/internal/cache/InternalCacheBuilder.java index 14bf28072731..0e600def5ae9 100644 --- a/geode-core/src/main/java/org/apache/geode/internal/cache/InternalCacheBuilder.java +++ b/geode-core/src/main/java/org/apache/geode/internal/cache/InternalCacheBuilder.java @@ -41,7 +41,9 @@ import org.apache.geode.distributed.internal.SecurityConfig; import org.apache.geode.logging.internal.log4j.api.LogService; import org.apache.geode.metrics.internal.InternalDistributedSystemMetricsService; +import org.apache.geode.metrics.internal.InternalDistributedSystemObservationService; import org.apache.geode.metrics.internal.MetricsService; +import org.apache.geode.metrics.internal.ObservationService; import org.apache.geode.pdx.PdxSerializer; import org.apache.geode.pdx.internal.TypeRegistry; import org.apache.geode.security.AuthenticationFailedException; @@ -66,6 +68,7 @@ public class InternalCacheBuilder { private final InternalDistributedSystemConstructor internalDistributedSystemConstructor; private final InternalCacheConstructor internalCacheConstructor; private final MetricsService.Builder metricsSessionBuilder; + private final ObservationService.Builder observationSessionBuilder; private boolean isExistingOk = IS_EXISTING_OK_DEFAULT; private boolean isClient = IS_CLIENT_DEFAULT; @@ -109,6 +112,7 @@ private InternalCacheBuilder(Properties configProperties, CacheConfig cacheConfi this(configProperties, cacheConfig, new InternalDistributedSystemMetricsService.Builder(), + new InternalDistributedSystemObservationService.Builder(), InternalDistributedSystem::getConnectedInstance, InternalDistributedSystem::connectInternal, GemFireCacheImpl::getInstance, @@ -119,6 +123,7 @@ private InternalCacheBuilder(Properties configProperties, CacheConfig cacheConfi InternalCacheBuilder(Properties configProperties, CacheConfig cacheConfig, MetricsService.Builder metricsSessionBuilder, + ObservationService.Builder observationSessionBuilder, Supplier singletonSystemSupplier, InternalDistributedSystemConstructor internalDistributedSystemConstructor, Supplier singletonCacheSupplier, @@ -130,7 +135,22 @@ private InternalCacheBuilder(Properties configProperties, CacheConfig cacheConfi this.internalCacheConstructor = internalCacheConstructor; this.singletonCacheSupplier = singletonCacheSupplier; this.metricsSessionBuilder = metricsSessionBuilder; + this.observationSessionBuilder = observationSessionBuilder; this.metricsSessionBuilder.setIsClient(isClient); + this.observationSessionBuilder.setIsClient(isClient); + } + + @VisibleForTesting + InternalCacheBuilder(Properties configProperties, + CacheConfig cacheConfig, + MetricsService.Builder metricsSessionBuilder, + Supplier singletonSystemSupplier, + InternalDistributedSystemConstructor internalDistributedSystemConstructor, + Supplier singletonCacheSupplier, + InternalCacheConstructor internalCacheConstructor) { + this(configProperties, cacheConfig, metricsSessionBuilder, + new InternalDistributedSystemObservationService.Builder(), singletonSystemSupplier, + internalDistributedSystemConstructor, singletonCacheSupplier, internalCacheConstructor); } /** @@ -290,6 +310,7 @@ public InternalCacheBuilder setIsExistingOk(boolean isExistingOk) { public InternalCacheBuilder setIsClient(boolean isClient) { this.isClient = isClient; metricsSessionBuilder.setIsClient(isClient); + observationSessionBuilder.setIsClient(isClient); return this; } @@ -343,7 +364,8 @@ private InternalDistributedSystem createInternalDistributedSystem() { cacheConfig.getPostProcessor()); return internalDistributedSystemConstructor - .construct(configProperties, securityConfig, metricsSessionBuilder); + .construct(configProperties, securityConfig, metricsSessionBuilder, + observationSessionBuilder); } private InternalCache existingCache(Supplier systemCacheSupplier, @@ -415,6 +437,7 @@ InternalCache construct(boolean isClient, PoolFactory poolFactory, @VisibleForTesting public interface InternalDistributedSystemConstructor { InternalDistributedSystem construct(Properties configProperties, SecurityConfig securityConfig, - MetricsService.Builder metricsSessionBuilder); + MetricsService.Builder metricsSessionBuilder, + ObservationService.Builder observationSessionBuilder); } } diff --git a/geode-core/src/main/java/org/apache/geode/internal/cache/InternalCacheForClientAccess.java b/geode-core/src/main/java/org/apache/geode/internal/cache/InternalCacheForClientAccess.java index a42be04ec9ab..13f3d95c0700 100644 --- a/geode-core/src/main/java/org/apache/geode/internal/cache/InternalCacheForClientAccess.java +++ b/geode-core/src/main/java/org/apache/geode/internal/cache/InternalCacheForClientAccess.java @@ -34,6 +34,7 @@ import javax.naming.Context; import io.micrometer.core.instrument.MeterRegistry; +import io.micrometer.observation.ObservationRegistry; import jakarta.transaction.TransactionManager; import org.apache.geode.CancelCriterion; @@ -1283,6 +1284,11 @@ public MeterRegistry getMeterRegistry() { return delegate.getMeterRegistry(); } + @Override + public ObservationRegistry getObservationRegistry() { + return delegate.getObservationRegistry(); + } + @Override public void saveCacheXmlForReconnect() { delegate.saveCacheXmlForReconnect(); diff --git a/geode-core/src/main/java/org/apache/geode/internal/cache/LocalRegion.java b/geode-core/src/main/java/org/apache/geode/internal/cache/LocalRegion.java index 3e912efcd759..c46c528afbeb 100644 --- a/geode-core/src/main/java/org/apache/geode/internal/cache/LocalRegion.java +++ b/geode-core/src/main/java/org/apache/geode/internal/cache/LocalRegion.java @@ -220,6 +220,7 @@ import org.apache.geode.internal.util.concurrent.FutureResult; import org.apache.geode.internal.util.concurrent.StoppableCountDownLatch; import org.apache.geode.logging.internal.log4j.api.LogService; +import org.apache.geode.metrics.internal.GeodeObservationSupport; import org.apache.geode.pdx.JSONFormatter; import org.apache.geode.pdx.PdxInstance; import org.apache.geode.util.internal.GeodeGlossary; @@ -1304,12 +1305,17 @@ Object getDeserialized(RegionEntry regionEntry, boolean updateStats, boolean dis @Override public Object get(Object key, Object aCallbackArgument, boolean generateCallbacks, EntryEventImpl clientEvent) throws TimeoutException, CacheLoaderException { - Object result = - get(key, aCallbackArgument, generateCallbacks, false, false, null, clientEvent, false); - if (Token.isInvalid(result)) { - result = null; - } - return result; + return GeodeObservationSupport.observe(cache.getObservationRegistry(), "geode.cache.region.get", + observation -> observation.lowCardinalityKeyValue("operation", "get"), + () -> { + Object result = + get(key, aCallbackArgument, generateCallbacks, false, false, null, clientEvent, + false); + if (Token.isInvalid(result)) { + result = null; + } + return result; + }); } /** @@ -1630,14 +1636,18 @@ public void invalidateRegion(Object aCallbackArgument) throws TimeoutException { @Override public Object put(Object key, Object value, Object aCallbackArgument) throws TimeoutException, CacheWriterException { - long startPut = getStatisticsClock().getTime(); - @Released - EntryEventImpl event = newUpdateEntryEvent(key, value, aCallbackArgument); - try { - return validatedPut(event, startPut); - } finally { - event.release(); - } + return GeodeObservationSupport.observe(cache.getObservationRegistry(), "geode.cache.region.put", + observation -> observation.lowCardinalityKeyValue("operation", "put"), + () -> { + long startPut = getStatisticsClock().getTime(); + @Released + EntryEventImpl event = newUpdateEntryEvent(key, value, aCallbackArgument); + try { + return validatedPut(event, startPut); + } finally { + event.release(); + } + }); } Object validatedPut(EntryEventImpl event, long startPut) diff --git a/geode-core/src/main/java/org/apache/geode/internal/cache/tier/sockets/ServerConnection.java b/geode-core/src/main/java/org/apache/geode/internal/cache/tier/sockets/ServerConnection.java index debcbb7646c9..80ea9d37a4a8 100644 --- a/geode-core/src/main/java/org/apache/geode/internal/cache/tier/sockets/ServerConnection.java +++ b/geode-core/src/main/java/org/apache/geode/internal/cache/tier/sockets/ServerConnection.java @@ -80,6 +80,7 @@ import org.apache.geode.internal.serialization.KnownVersion; import org.apache.geode.internal.util.Breadcrumbs; import org.apache.geode.logging.internal.log4j.api.LogService; +import org.apache.geode.metrics.internal.GeodeObservationSupport; import org.apache.geode.security.AuthenticationExpiredException; import org.apache.geode.security.AuthenticationFailedException; import org.apache.geode.security.AuthenticationRequiredException; @@ -877,7 +878,11 @@ void doNormalMessage() { // if a subject exists for this uniqueId, binds the subject to this thread so that we can do // authorization later threadState = bindSubject(command); - command.execute(message, this, securityService); + Command commandToExecute = command; + GeodeObservationSupport.observe(getCache().getObservationRegistry(), "geode.server.command", + observation -> observation.lowCardinalityKeyValue("message.type", + message.getMessageType().name()), + () -> commandToExecute.execute(message, this, securityService)); } } finally { suspendThreadMonitoring(); diff --git a/geode-core/src/main/java/org/apache/geode/internal/cache/xmlcache/CacheCreation.java b/geode-core/src/main/java/org/apache/geode/internal/cache/xmlcache/CacheCreation.java index 434882e86829..820d2a0b4ffe 100755 --- a/geode-core/src/main/java/org/apache/geode/internal/cache/xmlcache/CacheCreation.java +++ b/geode-core/src/main/java/org/apache/geode/internal/cache/xmlcache/CacheCreation.java @@ -45,6 +45,7 @@ import javax.naming.Context; import io.micrometer.core.instrument.MeterRegistry; +import io.micrometer.observation.ObservationRegistry; import jakarta.transaction.TransactionManager; import org.apache.logging.log4j.Logger; @@ -2486,6 +2487,11 @@ public MeterRegistry getMeterRegistry() { throw new UnsupportedOperationException("Should not be invoked"); } + @Override + public ObservationRegistry getObservationRegistry() { + throw new UnsupportedOperationException("Should not be invoked"); + } + @Override public void saveCacheXmlForReconnect() { throw new UnsupportedOperationException("Should not be invoked"); diff --git a/geode-core/src/main/java/org/apache/geode/metrics/ObservationPublishingService.java b/geode-core/src/main/java/org/apache/geode/metrics/ObservationPublishingService.java new file mode 100644 index 000000000000..907e80db75b2 --- /dev/null +++ b/geode-core/src/main/java/org/apache/geode/metrics/ObservationPublishingService.java @@ -0,0 +1,53 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more contributor license + * agreements. See the NOTICE file distributed with this work for additional information regarding + * copyright ownership. The ASF licenses this file to You under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance with the License. You may obtain a + * copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software distributed under the License + * is distributed on an "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express + * or implied. See the License for the specific language governing permissions and limitations under + * the License. + */ +package org.apache.geode.metrics; + +import java.util.ServiceLoader; + +import io.micrometer.observation.ObservationHandler; + +import org.apache.geode.annotations.Experimental; + +/** + * Configures observation publishing when an {@link ObservationSession} starts. + * + *

+ * Geode discovers {@code ObservationPublishingService}s during system creation, using the standard + * Java {@link ServiceLoader} mechanism. + * + *

+ * A typical implementation registers {@link ObservationHandler}s, predicates, filters, or + * conventions with {@link ObservationSession#getObservationRegistry()}. + * + *

+ * Experimental: Micrometer Observation is a new addition to Geode and the API may change. + */ +@Experimental("Micrometer Observation is a new addition to Geode and the API may change") +public interface ObservationPublishingService { + + /** + * Invoked when an observation session starts. + * + * @param session the observation session to configure + */ + void start(ObservationSession session); + + /** + * Invoked when an observation session stops. + * + * @param session the observation session this publishing service configured + */ + void stop(ObservationSession session); +} diff --git a/geode-core/src/main/java/org/apache/geode/metrics/ObservationSession.java b/geode-core/src/main/java/org/apache/geode/metrics/ObservationSession.java new file mode 100644 index 000000000000..3f345e018158 --- /dev/null +++ b/geode-core/src/main/java/org/apache/geode/metrics/ObservationSession.java @@ -0,0 +1,36 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more contributor license + * agreements. See the NOTICE file distributed with this work for additional information regarding + * copyright ownership. The ASF licenses this file to You under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance with the License. You may obtain a + * copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software distributed under the License + * is distributed on an "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express + * or implied. See the License for the specific language governing permissions and limitations under + * the License. + */ +package org.apache.geode.metrics; + +import io.micrometer.observation.ObservationRegistry; + +import org.apache.geode.annotations.Experimental; + +/** + * A session that manages Micrometer Observation instrumentation for Geode. + * + *

+ * Experimental: Micrometer Observation is a new addition to Geode and the API may change. + */ +@Experimental("Micrometer Observation is a new addition to Geode and the API may change") +public interface ObservationSession { + + /** + * Returns the registry used by this session to create observations. + * + * @return the observation registry + */ + ObservationRegistry getObservationRegistry(); +} diff --git a/geode-core/src/main/java/org/apache/geode/metrics/internal/GeodeObservationSupport.java b/geode-core/src/main/java/org/apache/geode/metrics/internal/GeodeObservationSupport.java new file mode 100644 index 000000000000..9f30a7e43d60 --- /dev/null +++ b/geode-core/src/main/java/org/apache/geode/metrics/internal/GeodeObservationSupport.java @@ -0,0 +1,73 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more contributor license + * agreements. See the NOTICE file distributed with this work for additional information regarding + * copyright ownership. The ASF licenses this file to You under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance with the License. You may obtain a + * copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software distributed under the License + * is distributed on an "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express + * or implied. See the License for the specific language governing permissions and limitations under + * the License. + */ +package org.apache.geode.metrics.internal; + +import java.util.function.Consumer; + +import io.micrometer.observation.Observation; +import io.micrometer.observation.ObservationRegistry; + +/** + * Utility methods for creating Geode observations with consistent error and stop handling. + */ +public final class GeodeObservationSupport { + + private GeodeObservationSupport() {} + + public static void observe(ObservationRegistry observationRegistry, String name, + Consumer observationConfigurer, ObservationRunnable runnable) { + observe(observationRegistry, name, observationConfigurer, () -> { + runnable.run(); + return null; + }); + } + + public static T observe(ObservationRegistry observationRegistry, String name, + Consumer observationConfigurer, ObservationSupplier supplier) { + Observation observation = startObservation(observationRegistry, name, observationConfigurer); + try (Observation.Scope ignored = observation.openScope()) { + return supplier.get(); + } catch (Throwable thrown) { + observation.error(thrown); + return sneakyThrow(thrown); + } finally { + observation.stop(); + } + } + + public static Observation startObservation(ObservationRegistry observationRegistry, String name, + Consumer observationConfigurer) { + Observation observation = observationRegistry == null + ? Observation.NOOP + : Observation.createNotStarted(name, observationRegistry); + observationConfigurer.accept(observation); + return observation.start(); + } + + @FunctionalInterface + public interface ObservationRunnable { + void run() throws Exception; + } + + @FunctionalInterface + public interface ObservationSupplier { + T get() throws Exception; + } + + @SuppressWarnings("unchecked") + private static T sneakyThrow(Throwable thrown) throws E { + throw (E) thrown; + } +} diff --git a/geode-core/src/main/java/org/apache/geode/metrics/internal/InternalDistributedSystemObservationService.java b/geode-core/src/main/java/org/apache/geode/metrics/internal/InternalDistributedSystemObservationService.java new file mode 100644 index 000000000000..33a1bb644874 --- /dev/null +++ b/geode-core/src/main/java/org/apache/geode/metrics/internal/InternalDistributedSystemObservationService.java @@ -0,0 +1,143 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more contributor license + * agreements. See the NOTICE file distributed with this work for additional information regarding + * copyright ownership. The ASF licenses this file to You under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance with the License. You may obtain a + * copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software distributed under the License + * is distributed on an "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express + * or implied. See the License for the specific language governing permissions and limitations under + * the License. + */ +package org.apache.geode.metrics.internal; + +import java.util.ArrayList; +import java.util.Collection; +import java.util.function.Supplier; + +import io.micrometer.observation.ObservationRegistry; +import org.apache.logging.log4j.Logger; + +import org.apache.geode.annotations.VisibleForTesting; +import org.apache.geode.distributed.internal.InternalDistributedSystem; +import org.apache.geode.internal.util.CollectingServiceLoader; +import org.apache.geode.internal.util.ListCollectingServiceLoader; +import org.apache.geode.logging.internal.log4j.api.LogService; +import org.apache.geode.metrics.ObservationPublishingService; + +/** + * Manages Micrometer Observation on behalf of an {@code InternalDistributedSystem}. + */ +public class InternalDistributedSystemObservationService implements ObservationService { + private final ObservationService.Builder builder; + private final Logger logger; + private final CollectingServiceLoader publishingServiceLoader; + private final ObservationRegistry observationRegistry; + private final Collection publishingServices = new ArrayList<>(); + + @FunctionalInterface + @VisibleForTesting + interface Factory { + ObservationService create(ObservationService.Builder builder, Logger logger, + CollectingServiceLoader publishingServiceLoader, + ObservationRegistry observationRegistry); + } + + @VisibleForTesting + InternalDistributedSystemObservationService(ObservationService.Builder builder, Logger logger, + CollectingServiceLoader publishingServiceLoader, + ObservationRegistry observationRegistry) { + this.builder = builder; + this.logger = logger; + this.publishingServiceLoader = publishingServiceLoader; + this.observationRegistry = observationRegistry; + } + + @Override + public void start() { + publishingServices.addAll( + publishingServiceLoader.loadServices(ObservationPublishingService.class)); + publishingServices.forEach(this::startObservationPublishingService); + } + + @Override + public void stop() { + publishingServices.forEach(this::stopObservationPublishingService); + publishingServices.clear(); + } + + @Override + public ObservationRegistry getObservationRegistry() { + return observationRegistry; + } + + @Override + public ObservationService.Builder getRebuilder() { + return builder; + } + + private void startObservationPublishingService(ObservationPublishingService service) { + try { + service.start(this); + } catch (Exception thrown) { + logger.error("Exception while starting observation publishing service " + + service.getClass().getName(), thrown); + } + } + + private void stopObservationPublishingService(ObservationPublishingService service) { + try { + service.stop(this); + } catch (Exception thrown) { + logger.error("Exception while stopping observation publishing service " + + service.getClass().getName(), thrown); + } + } + + public static class Builder implements ObservationService.Builder { + private Supplier loggerSupplier = LogService::getLogger; + private Factory observationServiceFactory = InternalDistributedSystemObservationService::new; + private Supplier observationRegistrySupplier = ObservationRegistry::create; + private Supplier> serviceLoaderSupplier = + ListCollectingServiceLoader::new; + + @Override + public ObservationService build(InternalDistributedSystem system) { + return observationServiceFactory.create(this, loggerSupplier.get(), + serviceLoaderSupplier.get(), + observationRegistrySupplier.get()); + } + + @Override + public ObservationService.Builder setIsClient(boolean isClient) { + return this; + } + + @VisibleForTesting + Builder setLogger(Logger logger) { + loggerSupplier = () -> logger; + return this; + } + + @VisibleForTesting + Builder setObservationRegistry(ObservationRegistry observationRegistry) { + observationRegistrySupplier = () -> observationRegistry; + return this; + } + + @VisibleForTesting + Builder setObservationServiceFactory(Factory factory) { + observationServiceFactory = factory; + return this; + } + + @VisibleForTesting + Builder setServiceLoader(CollectingServiceLoader serviceLoader) { + serviceLoaderSupplier = () -> serviceLoader; + return this; + } + } +} diff --git a/geode-core/src/main/java/org/apache/geode/metrics/internal/ObservationRegistrySupplier.java b/geode-core/src/main/java/org/apache/geode/metrics/internal/ObservationRegistrySupplier.java new file mode 100644 index 000000000000..5d19b6ed6254 --- /dev/null +++ b/geode-core/src/main/java/org/apache/geode/metrics/internal/ObservationRegistrySupplier.java @@ -0,0 +1,47 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more contributor license + * agreements. See the NOTICE file distributed with this work for additional information regarding + * copyright ownership. The ASF licenses this file to You under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance with the License. You may obtain a + * copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software distributed under the License + * is distributed on an "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express + * or implied. See the License for the specific language governing permissions and limitations under + * the License. + */ +package org.apache.geode.metrics.internal; + +import java.util.function.Supplier; + +import io.micrometer.observation.ObservationRegistry; + +import org.apache.geode.distributed.internal.InternalDistributedSystem; +import org.apache.geode.internal.cache.InternalCache; + +public class ObservationRegistrySupplier implements Supplier { + + private final Supplier internalDistributedSystemSupplier; + + public ObservationRegistrySupplier( + Supplier internalDistributedSystemSupplier) { + this.internalDistributedSystemSupplier = internalDistributedSystemSupplier; + } + + @Override + public ObservationRegistry get() { + InternalDistributedSystem system = internalDistributedSystemSupplier.get(); + if (system == null) { + return null; + } + + InternalCache internalCache = system.getCache(); + if (internalCache == null) { + return null; + } + + return internalCache.getObservationRegistry(); + } +} diff --git a/geode-core/src/main/java/org/apache/geode/metrics/internal/ObservationService.java b/geode-core/src/main/java/org/apache/geode/metrics/internal/ObservationService.java new file mode 100644 index 000000000000..18420bcb62ed --- /dev/null +++ b/geode-core/src/main/java/org/apache/geode/metrics/internal/ObservationService.java @@ -0,0 +1,65 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more contributor license + * agreements. See the NOTICE file distributed with this work for additional information regarding + * copyright ownership. The ASF licenses this file to You under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance with the License. You may obtain a + * copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software distributed under the License + * is distributed on an "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express + * or implied. See the License for the specific language governing permissions and limitations under + * the License. + */ +package org.apache.geode.metrics.internal; + +import io.micrometer.observation.ObservationRegistry; + +import org.apache.geode.distributed.internal.InternalDistributedSystem; +import org.apache.geode.metrics.ObservationSession; + +/** + * An observation session that can be started and stopped, and that manages an observation registry. + */ +public interface ObservationService extends ObservationSession { + + /** + * Starts this observation service and loads configured publishing services. + */ + void start(); + + /** + * Stops this observation service, freeing any resources. + */ + void stop(); + + /** + * Returns this service's observation registry. + * + * @return the observation registry + */ + ObservationRegistry getObservationRegistry(); + + /** + * Returns the builder that built this observation service. The builder can be used during + * reconnect to create an observation service configured similarly to this one. + * + * @return the builder that built this observation service + */ + Builder getRebuilder(); + + interface Builder { + /** + * Informs this builder whether it is building an observation service on behalf of a client. + * + * @return this builder + */ + Builder setIsClient(boolean isClient); + + /** + * Builds an observation service associated with the given system. + */ + ObservationService build(InternalDistributedSystem system); + } +} diff --git a/geode-core/src/test/resources/expected-pom.xml b/geode-core/src/test/resources/expected-pom.xml index 23f621c02619..686d53a8d67d 100644 --- a/geode-core/src/test/resources/expected-pom.xml +++ b/geode-core/src/test/resources/expected-pom.xml @@ -76,6 +76,21 @@ + + io.micrometer + micrometer-observation + compile + + + log4j-to-slf4j + org.apache.logging.log4j + + + * + ch.qos.logback + + + jakarta.resource jakarta.resource-api From 757afe9a1ca2a607289120f4eefef3bf16d0a05e Mon Sep 17 00:00:00 2001 From: Makarand Milind Hinge Date: Sun, 27 Sep 2026 10:11:40 +0530 Subject: [PATCH 3/3] GEODE-10440: update assembly content for observation API --- .../src/integrationTest/resources/assembly_content.txt | 2 ++ 1 file changed, 2 insertions(+) diff --git a/geode-assembly/src/integrationTest/resources/assembly_content.txt b/geode-assembly/src/integrationTest/resources/assembly_content.txt index bb903288049b..c8d5711404ab 100644 --- a/geode-assembly/src/integrationTest/resources/assembly_content.txt +++ b/geode-assembly/src/integrationTest/resources/assembly_content.txt @@ -766,6 +766,8 @@ javadoc/org/apache/geode/memcached/package-summary.html javadoc/org/apache/geode/memcached/package-tree.html javadoc/org/apache/geode/metrics/MetricsPublishingService.html javadoc/org/apache/geode/metrics/MetricsSession.html +javadoc/org/apache/geode/metrics/ObservationPublishingService.html +javadoc/org/apache/geode/metrics/ObservationSession.html javadoc/org/apache/geode/metrics/package-summary.html javadoc/org/apache/geode/metrics/package-tree.html javadoc/org/apache/geode/modules/gatewaydelta/AbstractGatewayDeltaEvent.html