Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
5 changes: 5 additions & 0 deletions boms/geode-all-bom/src/test/resources/expected-pom.xml
Original file line number Diff line number Diff line change
Expand Up @@ -197,6 +197,11 @@
<artifactId>micrometer-core</artifactId>
<version>1.16.7</version>
</dependency>
<dependency>
<groupId>io.micrometer</groupId>
<artifactId>micrometer-observation</artifactId>
<version>1.16.7</version>
</dependency>
<dependency>
<groupId>io.swagger.core.v3</groupId>
<artifactId>swagger-annotations</artifactId>
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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')
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
1 change: 1 addition & 0 deletions geode-core/build.gradle
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand All @@ -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;

/**
Expand All @@ -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 <code>DistributedSystem</code> with the given configuration properties.
*/
Expand Down Expand Up @@ -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";
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -49,5 +50,7 @@ <K, V> Region<K, V> basicCreateRegion(String name, RegionAttributes<K, V> attrs)

MeterRegistry getMeterRegistry();

ObservationRegistry getObservationRegistry();

ClientMetadataService getClientMetadataService();
}
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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;
Expand Down Expand Up @@ -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.
Expand Down Expand Up @@ -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);
}

/**
Expand All @@ -221,12 +235,22 @@ public static InternalDistributedSystem connectInternal(
SecurityConfig securityConfig,
MetricsService.Builder metricsSessionBuilder,
final MembershipLocator<InternalDistributedMember> 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<InternalDistributedMember> 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();
Expand Down Expand Up @@ -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;
}
Expand Down Expand Up @@ -654,6 +679,7 @@ public MemoryAllocator getOffHeapStore() {
@VisibleForTesting
void initialize(SecurityManager securityManager, PostProcessor postProcessor,
MetricsService.Builder metricsServiceBuilder,
ObservationService.Builder observationServiceBuilder,
final MembershipLocator<InternalDistributedMember> membershipLocatorArg,
ClusterDistributionManagerConstructor clusterDistributionManagerConstructor) {

Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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.
*/
Expand Down Expand Up @@ -1618,6 +1650,7 @@ protected void disconnect(boolean preparingForReconnect, String reason, boolean

functionStatsManager.close();
metricsService.stop();
observationService.stop();

InternalFunctionService.unregisterAllFunctions();

Expand Down Expand Up @@ -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()) {
Expand Down Expand Up @@ -3001,15 +3034,23 @@ public static class Builder {

private SecurityConfig securityConfig;
private final MetricsService.Builder metricsServiceBuilder;
private final ObservationService.Builder observationServiceBuilder;

private MembershipLocator<InternalDistributedMember> locator;

private ClusterDistributionManagerConstructor clusterDistributionManagerConstructor =
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) {
Expand Down Expand Up @@ -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;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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));
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -565,6 +566,8 @@ void invokeRegionEntrySynchronizationListenersAfterSynchronization(

MeterRegistry getMeterRegistry();

ObservationRegistry getObservationRegistry();

/**
* Generate XML for the cache before shutting down due to forced disconnect.
*/
Expand Down
Loading
Loading