diff --git a/build.gradle b/build.gradle index 0c202aa538..c71190a468 100644 --- a/build.gradle +++ b/build.gradle @@ -121,7 +121,7 @@ subprojects { subproject -> kryoShadedVersion = '3.0.3' lettuceVersion = '5.1.0.RELEASE' log4jVersion = '2.11.1' - micrometerVersion = '1.0.6' + micrometerVersion = '1.1.0-SNAPSHOT' mockitoVersion = '2.22.0' mysqlVersion = '8.0.11' pahoMqttClientVersion = '1.2.0' diff --git a/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/channel/PollableAmqpChannel.java b/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/channel/PollableAmqpChannel.java index 3f274857b8..22332181d4 100644 --- a/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/channel/PollableAmqpChannel.java +++ b/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/channel/PollableAmqpChannel.java @@ -352,4 +352,12 @@ public class PollableAmqpChannel extends AbstractAmqpChannel return this.executorInterceptorsSize > 0; } + @Override + public void destroy() throws Exception { + super.destroy(); + if (this.receiveCounter != null) { + this.receiveCounter.remove(); + } + } + } diff --git a/spring-integration-amqp/src/test/java/org/springframework/integration/amqp/outbound/AmqpOutboundEndpointTests.java b/spring-integration-amqp/src/test/java/org/springframework/integration/amqp/outbound/AmqpOutboundEndpointTests.java index 68449635ce..bf6b09933e 100644 --- a/spring-integration-amqp/src/test/java/org/springframework/integration/amqp/outbound/AmqpOutboundEndpointTests.java +++ b/spring-integration-amqp/src/test/java/org/springframework/integration/amqp/outbound/AmqpOutboundEndpointTests.java @@ -177,7 +177,7 @@ public class AmqpOutboundEndpointTests { .setHeader(AmqpHeaders.CONTENT_TYPE, "application/json") .build(); this.ctRequestChannel.send(message); - org.springframework.amqp.core.Message m = template.receive(); + org.springframework.amqp.core.Message m = receive(template); assertNotNull(m); assertEquals("\"hello\"", new String(m.getBody(), "UTF-8")); assertEquals("application/json", m.getMessageProperties().getContentType()); @@ -186,7 +186,7 @@ public class AmqpOutboundEndpointTests { message = MessageBuilder.withPayload("hello") .build(); this.ctRequestChannel.send(message); - m = template.receive(); + m = receive(template); assertNotNull(m); assertEquals("hello", new String(m.getBody(), "UTF-8")); assertEquals("text/plain", m.getMessageProperties().getContentType()); @@ -195,4 +195,15 @@ public class AmqpOutboundEndpointTests { } } + private org.springframework.amqp.core.Message receive(RabbitTemplate template) throws Exception { + int n = 0; + org.springframework.amqp.core.Message message = template.receive(); + while (message == null && n++ < 100) { + Thread.sleep(100); + message = template.receive(); + } + assertNotNull(message); + return message; + } + } diff --git a/spring-integration-core/src/main/java/org/springframework/integration/channel/AbstractMessageChannel.java b/spring-integration-core/src/main/java/org/springframework/integration/channel/AbstractMessageChannel.java index 52352f1399..68c77f836c 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/channel/AbstractMessageChannel.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/channel/AbstractMessageChannel.java @@ -22,6 +22,8 @@ import java.util.Comparator; import java.util.Deque; import java.util.Iterator; import java.util.List; +import java.util.Set; +import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.CopyOnWriteArrayList; import org.apache.commons.logging.Log; @@ -39,6 +41,7 @@ import org.springframework.integration.support.management.MessageChannelMetrics; import org.springframework.integration.support.management.MetricsContext; import org.springframework.integration.support.management.Statistics; import org.springframework.integration.support.management.TrackableComponent; +import org.springframework.integration.support.management.metrics.MeterFacade; import org.springframework.integration.support.management.metrics.MetricsCaptor; import org.springframework.integration.support.management.metrics.SampleFacade; import org.springframework.integration.support.management.metrics.TimerFacade; @@ -73,6 +76,8 @@ public abstract class AbstractMessageChannel extends IntegrationObjectSupport private final ManagementOverrides managementOverrides = new ManagementOverrides(); + protected final Set meters = ConcurrentHashMap.newKeySet(); // NOSONAR + private volatile boolean shouldTrack = false; private volatile Class[] datatypes = new Class[0]; @@ -493,13 +498,15 @@ public abstract class AbstractMessageChannel extends IntegrationObjectSupport } private TimerFacade buildSendTimer(boolean success, String exception) { - return this.metricsCaptor.timerBuilder(SEND_TIMER_NAME) + TimerFacade timer = this.metricsCaptor.timerBuilder(SEND_TIMER_NAME) .tag("type", "channel") .tag("name", getComponentName() == null ? "unknown" : getComponentName()) .tag("result", success ? "success" : "failure") .tag("exception", exception) .description("Send processing time") .build(); + this.meters.add(timer); + return timer; } private Message convertPayloadIfNecessary(Message message) { @@ -544,6 +551,10 @@ public abstract class AbstractMessageChannel extends IntegrationObjectSupport */ protected abstract boolean doSend(Message message, long timeout); + @Override + public void destroy() throws Exception { + this.meters.forEach(t -> t.remove()); + } /** * A convenience wrapper class for the list of ChannelInterceptors. diff --git a/spring-integration-core/src/main/java/org/springframework/integration/channel/AbstractPollableChannel.java b/spring-integration-core/src/main/java/org/springframework/integration/channel/AbstractPollableChannel.java index f6a5a9a406..220cf8a25f 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/channel/AbstractPollableChannel.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/channel/AbstractPollableChannel.java @@ -139,14 +139,15 @@ public abstract class AbstractPollableChannel extends AbstractMessageChannel catch (RuntimeException e) { if (countsEnabled && !counted) { if (getMetricsCaptor() != null) { - getMetricsCaptor().counterBuilder(RECEIVE_COUNTER_NAME) + CounterFacade counter = getMetricsCaptor().counterBuilder(RECEIVE_COUNTER_NAME) .tag("name", getComponentName() == null ? "unknown" : getComponentName()) .tag("type", "channel") .tag("result", "failure") .tag("exception", e.getClass().getSimpleName()) .description("Messages received") - .build() - .increment(); + .build(); + this.meters.add(counter); + counter.increment(); } getMetrics().afterError(); } @@ -231,4 +232,12 @@ public abstract class AbstractPollableChannel extends AbstractMessageChannel @Nullable protected abstract Message doReceive(long timeout); + @Override + public void destroy() throws Exception { + super.destroy(); + if (this.receiveCounter != null) { + this.receiveCounter.remove(); + } + } + } diff --git a/spring-integration-core/src/main/java/org/springframework/integration/channel/NullChannel.java b/spring-integration-core/src/main/java/org/springframework/integration/channel/NullChannel.java index 416d1a985b..e66c9170a2 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/channel/NullChannel.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/channel/NullChannel.java @@ -273,4 +273,11 @@ public class NullChannel implements PollableChannel, MessageChannelMetrics, return (this.beanName != null) ? this.beanName : super.toString(); } + @Override + public void destroy() throws Exception { + if (this.successTimer != null) { + this.successTimer.remove(); + } + } + } diff --git a/spring-integration-core/src/main/java/org/springframework/integration/config/annotation/AbstractMethodAnnotationPostProcessor.java b/spring-integration-core/src/main/java/org/springframework/integration/config/annotation/AbstractMethodAnnotationPostProcessor.java index 0dcfdacd3b..6a8f8aff61 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/config/annotation/AbstractMethodAnnotationPostProcessor.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/config/annotation/AbstractMethodAnnotationPostProcessor.java @@ -34,6 +34,7 @@ import org.springframework.aop.framework.ProxyFactory; import org.springframework.aop.support.DefaultBeanFactoryPointcutAdvisor; import org.springframework.aop.support.NameMatchMethodPointcut; import org.springframework.aop.support.NameMatchMethodPointcutAdvisor; +import org.springframework.beans.factory.DisposableBean; import org.springframework.beans.factory.NoSuchBeanDefinitionException; import org.springframework.beans.factory.config.ConfigurableListableBeanFactory; import org.springframework.beans.factory.support.BeanDefinitionValidationException; @@ -108,6 +109,8 @@ public abstract class AbstractMethodAnnotationPostProcessor annotationType; + protected final Disposables disposables; // NOSONAR + @SuppressWarnings("unchecked") public AbstractMethodAnnotationPostProcessor(ConfigurableListableBeanFactory beanFactory) { Assert.notNull(beanFactory, "'beanFactory' must not be null"); @@ -123,6 +126,14 @@ public abstract class AbstractMethodAnnotationPostProcessor) GenericTypeResolver.resolveTypeArgument(this.getClass(), MethodAnnotationPostProcessor.class); + Disposables disposables = null; + try { + disposables = beanFactory.getBean(Disposables.class); + } + catch (Exception e) { + // NOSONAR - only for test cases + } + this.disposables = disposables; } @@ -177,6 +188,9 @@ public abstract class AbstractMethodAnnotationPostProcessor disposables = new ArrayList<>(); + + public void add(DisposableBean... disposables) { + this.disposables.addAll(Arrays.asList(disposables)); + } + + @Override + public void destroy() throws Exception { + this.disposables.forEach(d -> { + try { + d.destroy(); + } + catch (Exception e) { + // NOSONAR + } + }); + } +} diff --git a/spring-integration-core/src/main/java/org/springframework/integration/config/annotation/InboundChannelAdapterAnnotationPostProcessor.java b/spring-integration-core/src/main/java/org/springframework/integration/config/annotation/InboundChannelAdapterAnnotationPostProcessor.java index 11debb52cb..b983191328 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/config/annotation/InboundChannelAdapterAnnotationPostProcessor.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/config/annotation/InboundChannelAdapterAnnotationPostProcessor.java @@ -131,6 +131,9 @@ public class InboundChannelAdapterAnnotationPostProcessor extends this.beanFactory.registerSingleton(messageSourceBeanName, methodInvokingMessageSource); messageSource = (MessageSource) this.beanFactory .initializeBean(methodInvokingMessageSource, messageSourceBeanName); + if (this.disposables != null) { + this.disposables.add(methodInvokingMessageSource); + } } return messageSource; } diff --git a/spring-integration-core/src/main/java/org/springframework/integration/config/annotation/MessagingAnnotationPostProcessor.java b/spring-integration-core/src/main/java/org/springframework/integration/config/annotation/MessagingAnnotationPostProcessor.java index 558408bd41..8bc7c46823 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/config/annotation/MessagingAnnotationPostProcessor.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/config/annotation/MessagingAnnotationPostProcessor.java @@ -38,6 +38,8 @@ import org.springframework.beans.factory.BeanFactoryAware; import org.springframework.beans.factory.InitializingBean; import org.springframework.beans.factory.config.BeanPostProcessor; import org.springframework.beans.factory.config.ConfigurableListableBeanFactory; +import org.springframework.beans.factory.support.BeanDefinitionBuilder; +import org.springframework.beans.factory.support.BeanDefinitionRegistry; import org.springframework.core.annotation.AnnotatedElementUtils; import org.springframework.core.annotation.AnnotationUtils; import org.springframework.integration.annotation.Aggregator; @@ -50,6 +52,7 @@ import org.springframework.integration.annotation.Router; import org.springframework.integration.annotation.ServiceActivator; import org.springframework.integration.annotation.Splitter; import org.springframework.integration.annotation.Transformer; +import org.springframework.integration.context.IntegrationContextUtils; import org.springframework.integration.endpoint.AbstractEndpoint; import org.springframework.integration.util.MessagingAnnotationUtils; import org.springframework.util.Assert; @@ -94,6 +97,10 @@ public class MessagingAnnotationPostProcessor implements BeanPostProcessor, Bean @Override public void afterPropertiesSet() { Assert.notNull(this.beanFactory, "BeanFactory must not be null"); + ((BeanDefinitionRegistry) this.beanFactory).registerBeanDefinition( + IntegrationContextUtils.DISPOSABLES_BEAN_NAME, + BeanDefinitionBuilder.genericBeanDefinition(Disposables.class, () -> new Disposables()) + .getRawBeanDefinition()); this.postProcessors.put(Filter.class, new FilterAnnotationPostProcessor(this.beanFactory)); this.postProcessors.put(Router.class, new RouterAnnotationPostProcessor(this.beanFactory)); this.postProcessors.put(Transformer.class, new TransformerAnnotationPostProcessor(this.beanFactory)); diff --git a/spring-integration-core/src/main/java/org/springframework/integration/context/IntegrationContextUtils.java b/spring-integration-core/src/main/java/org/springframework/integration/context/IntegrationContextUtils.java index 1c8d25f137..8d7ba5e458 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/context/IntegrationContextUtils.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/context/IntegrationContextUtils.java @@ -99,6 +99,8 @@ public abstract class IntegrationContextUtils { public static final String LIST_ARGUMENT_RESOLVERS_BEAN_NAME = "integrationListArgumentResolvers"; + public static final String DISPOSABLES_BEAN_NAME = "integrationDisposableAutoCreatedBeans"; + /** * @param beanFactory BeanFactory for lookup, must not be null. * @return The {@link MetadataStore} bean whose name is "metadataStore". diff --git a/spring-integration-core/src/main/java/org/springframework/integration/endpoint/AbstractMessageSource.java b/spring-integration-core/src/main/java/org/springframework/integration/endpoint/AbstractMessageSource.java index 20e57a72a7..7087092047 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/endpoint/AbstractMessageSource.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/endpoint/AbstractMessageSource.java @@ -233,4 +233,11 @@ public abstract class AbstractMessageSource extends AbstractExpressionEvaluat */ protected abstract Object doReceive(); + @Override + public void destroy() throws Exception { + if (this.receiveCounter != null) { + this.receiveCounter.remove(); + } + } + } diff --git a/spring-integration-core/src/main/java/org/springframework/integration/handler/AbstractMessageHandler.java b/spring-integration-core/src/main/java/org/springframework/integration/handler/AbstractMessageHandler.java index f6c068996a..b7a12fe6fe 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/handler/AbstractMessageHandler.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/handler/AbstractMessageHandler.java @@ -16,6 +16,9 @@ package org.springframework.integration.handler; +import java.util.Set; +import java.util.concurrent.ConcurrentHashMap; + import org.reactivestreams.Subscription; import org.springframework.core.Ordered; @@ -59,6 +62,8 @@ public abstract class AbstractMessageHandler extends IntegrationObjectSupport private final ManagementOverrides managementOverrides = new ManagementOverrides(); + private final Set timers = ConcurrentHashMap.newKeySet(); + private volatile boolean shouldTrack = false; private volatile int order = Ordered.LOWEST_PRECEDENCE; @@ -90,7 +95,6 @@ public abstract class AbstractMessageHandler extends IntegrationObjectSupport this.managementOverrides.loggingConfigured = true; } - @SuppressWarnings("unchecked") @Override public void registerMetricsCaptor(MetricsCaptor metricsCaptor) { this.metricsCaptor = metricsCaptor; @@ -185,13 +189,15 @@ public abstract class AbstractMessageHandler extends IntegrationObjectSupport } private TimerFacade buildSendTimer(boolean success, String exception) { - return this.metricsCaptor.timerBuilder(SEND_TIMER_NAME) + TimerFacade timer = this.metricsCaptor.timerBuilder(SEND_TIMER_NAME) .tag("type", "handler") .tag("name", getComponentName() == null ? "unknown" : getComponentName()) .tag("result", success ? "success" : "failure") .tag("exception", exception) .description("Send processing time") .build(); + this.timers.add(timer); + return timer; } @Override @@ -330,4 +336,9 @@ public abstract class AbstractMessageHandler extends IntegrationObjectSupport return this.managedType; } + @Override + public void destroy() throws Exception { + this.timers.forEach(t -> t.remove()); + } + } diff --git a/spring-integration-core/src/main/java/org/springframework/integration/support/management/IntegrationManagement.java b/spring-integration-core/src/main/java/org/springframework/integration/support/management/IntegrationManagement.java index bb1ade6d58..6ad332d29d 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/support/management/IntegrationManagement.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/support/management/IntegrationManagement.java @@ -1,5 +1,5 @@ /* - * Copyright 2015-2016 the original author or authors. + * Copyright 2015-2018 the original author or authors. * * Licensed under the Apache License, Version 2.0 (the "License"); * you may not use this file except in compliance with the License. @@ -16,6 +16,7 @@ package org.springframework.integration.support.management; +import org.springframework.beans.factory.DisposableBean; import org.springframework.integration.support.management.metrics.MetricsCaptor; import org.springframework.jmx.export.annotation.ManagedAttribute; import org.springframework.jmx.export.annotation.ManagedOperation; @@ -27,7 +28,7 @@ import org.springframework.jmx.export.annotation.ManagedOperation; * @since 4.2 * */ -public interface IntegrationManagement { +public interface IntegrationManagement extends DisposableBean { String METER_PREFIX = "spring.integration."; diff --git a/spring-integration-core/src/main/java/org/springframework/integration/support/management/LifecycleMessageHandlerMetrics.java b/spring-integration-core/src/main/java/org/springframework/integration/support/management/LifecycleMessageHandlerMetrics.java index 317d0117b9..353a064144 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/support/management/LifecycleMessageHandlerMetrics.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/support/management/LifecycleMessageHandlerMetrics.java @@ -187,4 +187,9 @@ public class LifecycleMessageHandlerMetrics implements MessageHandlerMetrics, Li return this.delegate.getOverrides(); } + @Override + public void destroy() throws Exception { + this.delegate.destroy(); + } + } diff --git a/spring-integration-core/src/main/java/org/springframework/integration/support/management/LifecycleMessageSourceMetrics.java b/spring-integration-core/src/main/java/org/springframework/integration/support/management/LifecycleMessageSourceMetrics.java index 506146eb9a..09abea24f8 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/support/management/LifecycleMessageSourceMetrics.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/support/management/LifecycleMessageSourceMetrics.java @@ -126,4 +126,9 @@ public class LifecycleMessageSourceMetrics implements MessageSourceMetrics, Life return this.delegate.getOverrides(); } + @Override + public void destroy() throws Exception { + this.delegate.destroy(); + } + } diff --git a/spring-integration-core/src/main/java/org/springframework/integration/support/management/MessageSourceMetrics.java b/spring-integration-core/src/main/java/org/springframework/integration/support/management/MessageSourceMetrics.java index bf8ca7fa55..0e44480661 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/support/management/MessageSourceMetrics.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/support/management/MessageSourceMetrics.java @@ -16,7 +16,6 @@ package org.springframework.integration.support.management; -import org.springframework.integration.support.management.metrics.CounterFacade; import org.springframework.jmx.export.annotation.ManagedMetric; import org.springframework.jmx.support.MetricType; @@ -50,16 +49,4 @@ public interface MessageSourceMetrics extends IntegrationManagement { String getManagedType(); - /** - * Set a micrometer counter to count messages produced. - * @param counter the counter. - * @since 5.0.2 - * @deprecated in favor of built-in counter registration via {@code MeterRegistry} callbacks. - * Will be remove in the next release. - */ - @Deprecated - default void setCounter(CounterFacade counter) { - // no op - } - } diff --git a/spring-integration-core/src/main/java/org/springframework/integration/support/management/metrics/CounterFacade.java b/spring-integration-core/src/main/java/org/springframework/integration/support/management/metrics/CounterFacade.java index b8263e2261..76ab06ad7a 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/support/management/metrics/CounterFacade.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/support/management/metrics/CounterFacade.java @@ -21,7 +21,7 @@ package org.springframework.integration.support.management.metrics; * @since 5.0.4 * */ -public interface CounterFacade { +public interface CounterFacade extends MeterFacade { void increment(); diff --git a/spring-integration-core/src/main/java/org/springframework/integration/support/management/metrics/GaugeFacade.java b/spring-integration-core/src/main/java/org/springframework/integration/support/management/metrics/GaugeFacade.java index 69dada42e1..ff330957da 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/support/management/metrics/GaugeFacade.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/support/management/metrics/GaugeFacade.java @@ -21,6 +21,6 @@ package org.springframework.integration.support.management.metrics; * @since 5.0.4 * */ -public interface GaugeFacade { +public interface GaugeFacade extends MeterFacade { } diff --git a/spring-integration-core/src/main/java/org/springframework/integration/support/management/metrics/MeterFacade.java b/spring-integration-core/src/main/java/org/springframework/integration/support/management/metrics/MeterFacade.java new file mode 100644 index 0000000000..f9057461e1 --- /dev/null +++ b/spring-integration-core/src/main/java/org/springframework/integration/support/management/metrics/MeterFacade.java @@ -0,0 +1,40 @@ +/* + * Copyright 2018 the original author or authors. + * + * Licensed 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.springframework.integration.support.management.metrics; + +import org.springframework.lang.Nullable; + +/** + * Facade for Meters. + * + * @author Gary Russell + * @since 5.1 + * + */ +public interface MeterFacade { + + /** + * Remove this meter facade. + * @param the type of meter removed. + * @return the facade that was removed, or null. + */ + @Nullable + default T remove() { + return null; + } + +} diff --git a/spring-integration-core/src/main/java/org/springframework/integration/support/management/metrics/MetricsCaptor.java b/spring-integration-core/src/main/java/org/springframework/integration/support/management/metrics/MetricsCaptor.java index ae9dde6e1f..4a62ebda6a 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/support/management/metrics/MetricsCaptor.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/support/management/metrics/MetricsCaptor.java @@ -58,6 +58,17 @@ public interface MetricsCaptor { */ SampleFacade start(); + /** + * Remove a meter facade. + * @param facade the facade to remove. + * @return the removed facade, or null. + * @since 5.1 + */ + @Nullable + default MeterFacade removeMeter(MeterFacade facade) { + return null; + } + /** * A builder for a timer. */ diff --git a/spring-integration-core/src/main/java/org/springframework/integration/support/management/metrics/TimerFacade.java b/spring-integration-core/src/main/java/org/springframework/integration/support/management/metrics/TimerFacade.java index d0ce7c4720..7f3fa2c860 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/support/management/metrics/TimerFacade.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/support/management/metrics/TimerFacade.java @@ -23,7 +23,7 @@ import java.util.concurrent.TimeUnit; * @since 5.0.4 * */ -public interface TimerFacade { +public interface TimerFacade extends MeterFacade { void record(long time, TimeUnit unit); diff --git a/spring-integration-core/src/main/java/org/springframework/integration/support/management/micrometer/MicrometerMetricsCaptor.java b/spring-integration-core/src/main/java/org/springframework/integration/support/management/micrometer/MicrometerMetricsCaptor.java index 6a9f9befd6..fdfd57a5f1 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/support/management/micrometer/MicrometerMetricsCaptor.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/support/management/micrometer/MicrometerMetricsCaptor.java @@ -24,6 +24,7 @@ import org.springframework.context.ApplicationContext; import org.springframework.context.support.GenericApplicationContext; import org.springframework.integration.support.management.metrics.CounterFacade; import org.springframework.integration.support.management.metrics.GaugeFacade; +import org.springframework.integration.support.management.metrics.MeterFacade; import org.springframework.integration.support.management.metrics.MetricsCaptor; import org.springframework.integration.support.management.metrics.SampleFacade; import org.springframework.integration.support.management.metrics.TimerFacade; @@ -31,6 +32,7 @@ import org.springframework.util.Assert; import io.micrometer.core.instrument.Counter; import io.micrometer.core.instrument.Gauge; +import io.micrometer.core.instrument.Meter; import io.micrometer.core.instrument.MeterRegistry; import io.micrometer.core.instrument.Timer; @@ -73,6 +75,11 @@ public class MicrometerMetricsCaptor implements MetricsCaptor { return new MicroSample(Timer.start(this.meterRegistry)); } + @Override + public MeterFacade removeMeter(MeterFacade facade) { + return facade.remove(); + } + /** * Add a MicrometerMetricsCaptor to the context if there's a MeterRegistry. * @param applicationContext the application context. @@ -135,19 +142,51 @@ public class MicrometerMetricsCaptor implements MetricsCaptor { @Override public MicroTimer build() { - return new MicroTimer(this.builder.register(this.meterRegistry)); + return new MicroTimer(this.builder.register(this.meterRegistry), this.meterRegistry); } } - private static class MicroTimer implements TimerFacade { + private static abstract class AbstractMeter implements MeterFacade { + + protected final MeterRegistry meterRegistry; // NOSONAR + + protected AbstractMeter(MeterRegistry meterRegistry) { + this.meterRegistry = meterRegistry; + } + + /** + * Get the meter. + * @return the meter. + */ + protected abstract Meter getMeter(); + + @SuppressWarnings("unchecked") + @Override + public T remove() { + if (this.meterRegistry.remove(getMeter()) != null) { + return (T) this; + } + else { + return null; + } + } + + } + private static class MicroTimer extends AbstractMeter implements TimerFacade { private final Timer timer; - MicroTimer(Timer timer) { + MicroTimer(Timer timer, MeterRegistry meterRegistry) { + super(meterRegistry); this.timer = timer; } + @Override + protected Meter getMeter() { + return this.timer; + } + @Override public void record(long time, TimeUnit unit) { this.timer.record(time, unit); @@ -180,19 +219,25 @@ public class MicrometerMetricsCaptor implements MetricsCaptor { @Override public CounterFacade build() { - return new MicroCounter(this.builder.register(this.meterRegistry)); + return new MicroCounter(this.builder.register(this.meterRegistry), this.meterRegistry); } } - private static class MicroCounter implements CounterFacade { + private static class MicroCounter extends AbstractMeter implements CounterFacade { private final Counter counter; - MicroCounter(Counter counter) { + MicroCounter(Counter counter, MeterRegistry meterRegistry) { + super(meterRegistry); this.counter = counter; } + @Override + protected Meter getMeter() { + return this.counter; + } + @Override public void increment() { this.counter.increment(); @@ -225,20 +270,25 @@ public class MicrometerMetricsCaptor implements MetricsCaptor { @Override public GaugeFacade build() { - return new MicroGauge(this.builder.register(this.meterRegistry)); + return new MicroGauge(this.builder.register(this.meterRegistry), this.meterRegistry); } } - private static class MicroGauge implements GaugeFacade { + private static class MicroGauge extends AbstractMeter implements GaugeFacade { - @SuppressWarnings("unused") private final Gauge gauge; - MicroGauge(Gauge gauge) { + MicroGauge(Gauge gauge, MeterRegistry meterRegistry) { + super(meterRegistry); this.gauge = gauge; } + @Override + protected Meter getMeter() { + return this.gauge; + } + } } diff --git a/spring-integration-core/src/test/java/org/springframework/integration/config/annotation/MessagingAnnotationPostProcessorChannelCreationTests.java b/spring-integration-core/src/test/java/org/springframework/integration/config/annotation/MessagingAnnotationPostProcessorChannelCreationTests.java index 95045c7213..d3e71dbb04 100644 --- a/spring-integration-core/src/test/java/org/springframework/integration/config/annotation/MessagingAnnotationPostProcessorChannelCreationTests.java +++ b/spring-integration-core/src/test/java/org/springframework/integration/config/annotation/MessagingAnnotationPostProcessorChannelCreationTests.java @@ -1,5 +1,5 @@ /* - * Copyright 2017 the original author or authors. + * Copyright 2017-2018 the original author or authors. * * Licensed under the Apache License, Version 2.0 (the "License"); * you may not use this file except in compliance with the License. @@ -25,12 +25,14 @@ import static org.mockito.BDDMockito.given; import static org.mockito.BDDMockito.willAnswer; import static org.mockito.Mockito.mock; import static org.mockito.Mockito.verify; +import static org.mockito.Mockito.withSettings; import org.junit.Test; import org.springframework.beans.factory.BeanCreationException; import org.springframework.beans.factory.NoSuchBeanDefinitionException; import org.springframework.beans.factory.config.ConfigurableListableBeanFactory; +import org.springframework.beans.factory.support.BeanDefinitionRegistry; import org.springframework.integration.annotation.ServiceActivator; import org.springframework.integration.channel.DirectChannel; import org.springframework.messaging.MessageChannel; @@ -46,7 +48,8 @@ public class MessagingAnnotationPostProcessorChannelCreationTests { @Test public void testAutoCreateChannel() { - ConfigurableListableBeanFactory beanFactory = mock(ConfigurableListableBeanFactory.class); + ConfigurableListableBeanFactory beanFactory = mock(ConfigurableListableBeanFactory.class, + withSettings().extraInterfaces(BeanDefinitionRegistry.class)); given(beanFactory.getBean("channel", MessageChannel.class)).willThrow(NoSuchBeanDefinitionException.class); willAnswer(invocation -> invocation.getArgument(0)) .given(beanFactory).initializeBean(any(DirectChannel.class), eq("channel")); @@ -61,7 +64,8 @@ public class MessagingAnnotationPostProcessorChannelCreationTests { @Test public void testDontCreateChannelWhenChannelHasBadDefinition() { - ConfigurableListableBeanFactory beanFactory = mock(ConfigurableListableBeanFactory.class); + ConfigurableListableBeanFactory beanFactory = mock(ConfigurableListableBeanFactory.class, + withSettings().extraInterfaces(BeanDefinitionRegistry.class)); given(beanFactory.getBean("channel", MessageChannel.class)).willThrow(BeanCreationException.class); willAnswer(invocation -> invocation.getArgument(0)) .given(beanFactory).initializeBean(any(DirectChannel.class), eq("channel")); diff --git a/spring-integration-core/src/test/java/org/springframework/integration/support/management/micrometer/MicrometerMetricsTests.java b/spring-integration-core/src/test/java/org/springframework/integration/support/management/micrometer/MicrometerMetricsTests.java index 89b51450bb..3749c30e83 100644 --- a/spring-integration-core/src/test/java/org/springframework/integration/support/management/micrometer/MicrometerMetricsTests.java +++ b/spring-integration-core/src/test/java/org/springframework/integration/support/management/micrometer/MicrometerMetricsTests.java @@ -45,10 +45,10 @@ import org.springframework.messaging.MessageHandler; import org.springframework.messaging.MessagingException; import org.springframework.messaging.PollableChannel; import org.springframework.messaging.support.GenericMessage; -import org.springframework.test.annotation.DirtiesContext; import org.springframework.test.context.junit4.SpringRunner; import io.micrometer.core.instrument.MeterRegistry; +import io.micrometer.core.instrument.search.MeterNotFoundException; import io.micrometer.core.instrument.simple.SimpleMeterRegistry; /** @@ -58,7 +58,6 @@ import io.micrometer.core.instrument.simple.SimpleMeterRegistry; * */ @RunWith(SpringRunner.class) -@DirtiesContext public class MicrometerMetricsTests { @Autowired @@ -86,7 +85,7 @@ public class MicrometerMetricsTests { private NullChannel nullChannel; @Test - public void testSend() { + public void testSend() throws Exception { GenericMessage message = new GenericMessage<>("foo"); this.channel.send(message); try { @@ -170,6 +169,39 @@ public class MicrometerMetricsTests { .tag("name", "newChannel") .tag("result", "success") .timer().count()).isEqualTo(1); + + // Test meter removal + registry.get("spring.integration.send") + .tag("name", "newChannel") + .tag("result", "success") + .timer(); + newChannel.destroy(); + try { + registry.get("spring.integration.send") + .tag("name", "newChannel") + .tag("result", "success") + .timer(); + fail("Expected MeterNotFoundException"); + } + catch (MeterNotFoundException e) { + assertThat(e).hasMessageContaining("A meter with name 'spring.integration.send' was found"); + assertThat(e).hasMessageContaining("No meters have a tag 'name' with value 'newChannel'"); + } + this.context.close(); + try { + registry.get("spring.integration.send").timers(); + fail("Expected MeterNotFoundException"); + } + catch (MeterNotFoundException e) { + assertThat(e).hasMessageContaining("No meter with name 'spring.integration.send' was found"); + } + try { + registry.get("spring.integration.receive").counters(); + fail("Expected MeterNotFoundException"); + } + catch (MeterNotFoundException e) { + assertThat(e).hasMessageContaining("No meter with name 'spring.integration.receive' was found"); + } } @Configuration