Fix JMS Inbound Endpoints for observation

The `JmsMessageDrivenEndpoint` delegates all the hard work to the
`ChannelPublishingJmsMessageListener`, but missed to propagate an `ObservationRegistry`
and other related options.
The `JmsInboundGateway` is worse: it delegated to the `JmsMessageDrivenEndpoint`

* Add `IntegrationObservation.HANDLER` observation to the `MessagingGatewaySupport.send()`
operation: used by the delegate in the `ChannelPublishingJmsMessageListener`
* Expose and propagate observation-related options from `JmsInboundGateway`
and `JmsMessageDrivenEndpoint`
* Expose `observationConvention()` option on the `MessagingGatewaySpec`
and `MessageProducerSpec`
* Remove unused imports
* Do not start a new `RECEIVER` observation if there is already `SERVER` one
* Fix `MessagingGatewaySupport` for `Observation.NOOP` check.
The parent process may still use `ObservationRegistry.NOOP` which sets
`Observation.NOOP` instance into the current context and thread local.

**Cherry-pick to `6.1.x` & `6.0.x`**
This commit is contained in:
Artem Bilan
2023-09-19 12:02:49 -04:00
committed by Christian Tzolov
parent 78b09ed610
commit 80990af0d1
8 changed files with 208 additions and 35 deletions

View File

@@ -18,6 +18,7 @@ package org.springframework.integration.jms;
import java.util.Map;
import io.micrometer.observation.ObservationRegistry;
import jakarta.jms.DeliveryMode;
import jakarta.jms.Destination;
import jakarta.jms.InvalidDestinationException;
@@ -38,6 +39,9 @@ import org.springframework.integration.gateway.MessagingGatewaySupport;
import org.springframework.integration.support.DefaultMessageBuilderFactory;
import org.springframework.integration.support.MessageBuilderFactory;
import org.springframework.integration.support.management.TrackableComponent;
import org.springframework.integration.support.management.metrics.MetricsCaptor;
import org.springframework.integration.support.management.observation.MessageReceiverObservationConvention;
import org.springframework.integration.support.management.observation.MessageRequestReplyReceiverObservationConvention;
import org.springframework.integration.support.utils.IntegrationUtils;
import org.springframework.jms.listener.SessionAwareMessageListener;
import org.springframework.jms.support.JmsUtils;
@@ -323,6 +327,26 @@ public class ChannelPublishingJmsMessageListener
this.extractReplyPayload = extractReplyPayload;
}
public void setMetricsCaptor(MetricsCaptor captor) {
this.gatewayDelegate.registerMetricsCaptor(captor);
}
public void setObservationRegistry(ObservationRegistry observationRegistry) {
this.gatewayDelegate.registerObservationRegistry(observationRegistry);
}
public void setRequestReplyObservationConvention(
@Nullable MessageRequestReplyReceiverObservationConvention observationConvention) {
this.gatewayDelegate.setObservationConvention(observationConvention);
}
public void setReceiverObservationConvention(
@Nullable MessageReceiverObservationConvention observationConvention) {
this.gatewayDelegate.setReceiverObservationConvention(observationConvention);
}
@Override
public void setBeanFactory(BeanFactory beanFactory) throws BeansException {
this.beanFactory = beanFactory;

View File

@@ -1,5 +1,5 @@
/*
* Copyright 2016-2022 the original author or authors.
* Copyright 2016-2023 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,10 +16,14 @@
package org.springframework.integration.jms;
import io.micrometer.observation.ObservationRegistry;
import org.springframework.beans.BeansException;
import org.springframework.context.ApplicationContext;
import org.springframework.integration.context.OrderlyShutdownCapable;
import org.springframework.integration.gateway.MessagingGatewaySupport;
import org.springframework.integration.support.management.metrics.MetricsCaptor;
import org.springframework.integration.support.management.observation.MessageRequestReplyReceiverObservationConvention;
import org.springframework.jms.listener.AbstractMessageListenerContainer;
import org.springframework.messaging.MessageChannel;
@@ -114,28 +118,38 @@ public class JmsInboundGateway extends MessagingGatewaySupport implements Orderl
this.endpoint.setShutdownContainerOnStop(shutdownContainerOnStop);
}
@Override
public void registerMetricsCaptor(MetricsCaptor metricsCaptorToRegister) {
super.registerMetricsCaptor(metricsCaptorToRegister);
this.endpoint.registerMetricsCaptor(metricsCaptorToRegister);
}
@Override
public void registerObservationRegistry(ObservationRegistry observationRegistry) {
super.registerObservationRegistry(observationRegistry);
this.endpoint.registerObservationRegistry(observationRegistry);
}
@Override
public void setObservationConvention(MessageRequestReplyReceiverObservationConvention observationConvention) {
super.setObservationConvention(observationConvention);
this.endpoint.getListener().setRequestReplyObservationConvention(observationConvention);
}
@Override
public String getComponentType() {
return this.endpoint.getComponentType();
}
@Override
public void setComponentName(String componentName) {
super.setComponentName(componentName);
this.endpoint.setComponentName(getComponentName());
}
@Override
public void setApplicationContext(ApplicationContext applicationContext) throws BeansException {
super.setApplicationContext(applicationContext);
this.endpoint.setApplicationContext(applicationContext);
this.endpoint.setBeanFactory(applicationContext);
this.endpoint.getListener().setBeanFactory(applicationContext);
}
@Override
protected void onInit() {
this.endpoint.setComponentName(getComponentName());
this.endpoint.afterPropertiesSet();
}

View File

@@ -1,5 +1,5 @@
/*
* Copyright 2002-2021 the original author or authors.
* Copyright 2002-2023 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,11 +16,15 @@
package org.springframework.integration.jms;
import io.micrometer.observation.ObservationRegistry;
import org.springframework.beans.BeansException;
import org.springframework.context.ApplicationContext;
import org.springframework.integration.context.OrderlyShutdownCapable;
import org.springframework.integration.endpoint.MessageProducerSupport;
import org.springframework.integration.jms.util.JmsAdapterUtils;
import org.springframework.integration.support.management.metrics.MetricsCaptor;
import org.springframework.integration.support.management.observation.MessageReceiverObservationConvention;
import org.springframework.jms.listener.AbstractMessageListenerContainer;
import org.springframework.jms.listener.DefaultMessageListenerContainer;
import org.springframework.messaging.MessageChannel;
@@ -91,7 +95,7 @@ public class JmsMessageDrivenEndpoint extends MessageProducerSupport implements
* container setting even if an external container is provided. Defaults to null
* (won't change container) if an external container is provided or `transacted` when
* the framework creates an implicit {@link DefaultMessageListenerContainer}.
* @param sessionAcknowledgeMode the acknowledge mode.
* @param sessionAcknowledgeMode the acknowledgement mode.
*/
public void setSessionAcknowledgeMode(String sessionAcknowledgeMode) {
this.sessionAcknowledgeMode = sessionAcknowledgeMode;
@@ -134,9 +138,9 @@ public class JmsMessageDrivenEndpoint extends MessageProducerSupport implements
}
/**
* Set to false to prevent listener container shutdown when the endpoint is stopped.
* Set to {@code false} to prevent listener container shutdown when the endpoint is stopped.
* Then, if so configured, any cached consumer(s) in the container will remain.
* Otherwise the shared connection and will be closed and the listener invokers shut
* Otherwise, the shared connection and will be closed and the listener invokers shut
* down; this behavior is new starting with version 5.1. Default: true.
* @param shutdownContainerOnStop false to not shutdown.
* @since 5.1
@@ -149,6 +153,24 @@ public class JmsMessageDrivenEndpoint extends MessageProducerSupport implements
return this.listener;
}
@Override
public void registerMetricsCaptor(MetricsCaptor captor) {
super.registerMetricsCaptor(captor);
this.listener.setMetricsCaptor(captor);
}
@Override
public void registerObservationRegistry(ObservationRegistry observationRegistry) {
super.registerObservationRegistry(observationRegistry);
this.listener.setObservationRegistry(observationRegistry);
}
@Override
public void setObservationConvention(MessageReceiverObservationConvention observationConvention) {
super.setObservationConvention(observationConvention);
this.listener.setReceiverObservationConvention(observationConvention);
}
@Override
public void setApplicationContext(ApplicationContext applicationContext) throws BeansException {
super.setApplicationContext(applicationContext);

View File

@@ -22,8 +22,11 @@ import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicBoolean;
import java.util.concurrent.atomic.AtomicInteger;
import io.micrometer.observation.tck.TestObservationRegistry;
import io.micrometer.observation.tck.TestObservationRegistryAssert;
import jakarta.jms.JMSException;
import jakarta.jms.TextMessage;
import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test;
import org.springframework.beans.factory.ListableBeanFactory;
@@ -40,6 +43,7 @@ import org.springframework.integration.channel.DirectChannel;
import org.springframework.integration.channel.FixedSubscriberChannel;
import org.springframework.integration.channel.QueueChannel;
import org.springframework.integration.config.EnableIntegration;
import org.springframework.integration.config.EnableIntegrationManagement;
import org.springframework.integration.config.GlobalChannelInterceptor;
import org.springframework.integration.core.MessageSource;
import org.springframework.integration.dsl.IntegrationFlow;
@@ -55,6 +59,7 @@ import org.springframework.integration.jms.JmsDestinationPollingSource;
import org.springframework.integration.jms.SubscribableJmsChannel;
import org.springframework.integration.scheduling.PollerMetadata;
import org.springframework.integration.support.MessageBuilder;
import org.springframework.integration.support.management.observation.IntegrationObservation;
import org.springframework.integration.test.util.TestUtils;
import org.springframework.jms.core.JmsTemplate;
import org.springframework.jms.listener.DefaultMessageListenerContainer;
@@ -144,6 +149,14 @@ public class JmsTests extends ActiveMQMultiContextTests {
@Autowired
JmsTemplate jmsTemplate;
@Autowired
TestObservationRegistry observationRegistry;
@BeforeEach
void setup() {
this.observationRegistry.clear();
}
@Test
public void testPollingFlow() {
this.controlBus.send("@'integerMessageSource.inboundChannelAdapter'.start()");
@@ -204,6 +217,14 @@ public class JmsTests extends ActiveMQMultiContextTests {
.isEqualTo("foo");
assertThat(this.jmsOutboundFlowTemplate).isNotNull();
TestObservationRegistryAssert.assertThat(this.observationRegistry)
.hasObservationWithNameEqualTo("spring.integration.handler")
.that()
.hasLowCardinalityKeyValue(IntegrationObservation.HandlerTags.COMPONENT_NAME.asString(),
"observedJmsMessageDrivenChannelAdapter")
.hasBeenStarted()
.hasBeenStopped();
}
@Test
@@ -239,6 +260,14 @@ public class JmsTests extends ActiveMQMultiContextTests {
.isNotNull()
.extracting(Message::getPayload)
.isEqualTo("error: junk is not convertible");
TestObservationRegistryAssert.assertThat(this.observationRegistry)
.hasObservationWithNameEqualTo("spring.integration.gateway")
.that()
.hasLowCardinalityKeyValue(IntegrationObservation.GatewayTags.COMPONENT_NAME.asString(),
"observedJmsInboundGateway")
.hasBeenStarted()
.hasBeenStopped();
}
@Test
@@ -294,8 +323,14 @@ public class JmsTests extends ActiveMQMultiContextTests {
@Configuration
@EnableIntegration
@IntegrationComponentScan
@EnableIntegrationManagement(observationPatterns = "observedJms*")
public static class ContextConfiguration {
@Bean
TestObservationRegistry observationRegistry() {
return TestObservationRegistry.create();
}
@Bean
public JmsTemplate jmsTemplate() {
JmsTemplate jmsTemplate = new JmsTemplate(connectionFactory);
@@ -407,9 +442,10 @@ public class JmsTests extends ActiveMQMultiContextTests {
public IntegrationFlow jmsMessageDrivenFlowWithContainer() {
return IntegrationFlow
.from(Jms.messageDrivenChannelAdapter(
Jms.container(amqFactory, "containerSpecDestination")
.pubSubDomain(false)
.taskExecutor(Executors.newCachedThreadPool())))
Jms.container(amqFactory, "containerSpecDestination")
.pubSubDomain(false)
.taskExecutor(Executors.newCachedThreadPool()))
.id("observedJmsMessageDrivenChannelAdapter"))
.transform(String::trim)
.channel(jmsOutboundInboundReplyChannel())
.get();
@@ -428,6 +464,7 @@ public class JmsTests extends ActiveMQMultiContextTests {
public IntegrationFlow jmsInboundGatewayFlow() {
return IntegrationFlow.from(
Jms.inboundGateway(amqFactory)
.id("observedJmsInboundGateway")
.requestChannel(jmsInboundGatewayInputChannel())
.replyTimeout(1)
.errorOnTimeout(true)