Fix IntegrationMBeanExporter logic for endpoints

Related to https://stackoverflow.com/questions/72851234/error-in-startup-application-when-using-serviceactivator-in-spring-cloud-stream

The `MessagingAnnotationPostProcessor` register an `endpoint` bean for Messaging Annotation on POJO methods.
The `IntegrationMBeanExporter` post-process such a bean and registers respective MBean.
It does this not in optimal way scanning all the `IntegrationConsumer` beans for requested `MessageHandler`
which may cause a `BeanCurrentlyInCreationException`.

* Rework the logic of the `IntegrationMBeanExporter.postProcessAbstractEndpoint()` to propagate provided endpoint
for the monitor registration to bypass application context scanning for matched name.
* Swap `equals()` for `monitor` since an `extractTarget()` may return `null`
* Some other code clean in the `IntegrationMBeanExporter`
* Change one of the `@ServiceActivator` in the configuration for the `ScatterGatherHandlerIntegrationTests`
to POJO method.
However, this didn't fail for me with original code unlike in the sample application provided in the mentioned SO thread.

**Cherry-pick to `main`**
This commit is contained in:
Artem Bilan
2022-07-07 15:30:43 -04:00
committed by Gary Russell
parent 8228269fde
commit 3bfaa3868b
2 changed files with 50 additions and 60 deletions

View File

@@ -1,5 +1,5 @@
/*
* Copyright 2002-2021 the original author or authors.
* Copyright 2002-2022 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.
@@ -264,8 +264,7 @@ public class IntegrationMBeanExporter extends MBeanExporter
}
private void populateMessageSources() {
this.applicationContext.getBeansOfType(
IntegrationInboundManagement.class)
this.applicationContext.getBeansOfType(IntegrationInboundManagement.class)
.values()
.stream()
// If the source is proxied, we have to extract the target to expose as an MBean.
@@ -337,7 +336,7 @@ public class IntegrationMBeanExporter extends MBeanExporter
MessageHandler handler = integrationConsumer.getHandler();
MessageHandler monitor = (MessageHandler) extractTarget(handler);
if (monitor instanceof IntegrationManagement) {
registerHandler((IntegrationManagement) monitor);
registerHandler((IntegrationManagement) monitor, integrationConsumer);
this.handlers.put(((IntegrationManagement) monitor).getComponentName(),
(IntegrationManagement) monitor);
this.runtimeBeans.add(monitor);
@@ -691,8 +690,12 @@ public class IntegrationMBeanExporter extends MBeanExporter
}
private void registerHandler(IntegrationManagement monitor2) {
IntegrationManagement monitor = enhanceHandlerMonitor(monitor2);
private void registerHandler(IntegrationManagement monitor) {
registerHandler(monitor, null);
}
private void registerHandler(IntegrationManagement monitor2, @Nullable IntegrationConsumer consumer) {
IntegrationManagement monitor = enhanceHandlerMonitor(monitor2, consumer);
String name = monitor.getComponentName();
if (!this.objectNames.containsKey(monitor2) && matches(this.componentNamePatterns, name)) {
String beanKey = getHandlerBeanKey(monitor);
@@ -778,7 +781,7 @@ public class IntegrationMBeanExporter extends MBeanExporter
*/
private boolean matches(String[] patterns, String name) {
Boolean match = PatternMatchUtils.smartMatch(name, patterns);
return match == null ? false : match;
return match != null && match;
}
private Object extractTarget(Object bean) {
@@ -846,51 +849,53 @@ public class IntegrationMBeanExporter extends MBeanExporter
.collect(Collectors.joining(","));
}
@SuppressWarnings("unlikely-arg-type")
private IntegrationManagement enhanceHandlerMonitor(IntegrationManagement monitor2) {
private IntegrationManagement enhanceHandlerMonitor(IntegrationManagement monitor,
@Nullable IntegrationConsumer consumer) {
if (monitor2.getManagedName() != null && monitor2.getManagedType() != null) {
return monitor2;
if (monitor.getManagedName() != null && monitor.getManagedType() != null) {
return monitor;
}
// Assignment algorithm and bean id, with bean id pulled reflectively out of enclosing endpoint if possible
String[] names = this.applicationContext.getBeanNamesForType(IntegrationConsumer.class);
String name = null;
String endpointName = null;
String source = "endpoint";
IntegrationConsumer endpoint = null;
IntegrationConsumer endpoint = consumer;
for (String beanName : names) {
endpoint = this.applicationContext.getBean(beanName, IntegrationConsumer.class);
try {
MessageHandler handler = endpoint.getHandler();
if (handler.equals(monitor2) ||
extractTarget(handlerInAnonymousWrapper(handler)).equals(monitor2)) {
name = beanName;
endpointName = beanName;
break;
if (endpoint == null) {
// Assignment algorithm and bean id, with bean id pulled reflectively out of enclosing endpoint if possible
String[] names = this.applicationContext.getBeanNamesForType(IntegrationConsumer.class);
for (String beanName : names) {
endpoint = this.applicationContext.getBean(beanName, IntegrationConsumer.class);
try {
MessageHandler handler = endpoint.getHandler();
if (handler.equals(monitor) || monitor.equals(extractTarget(handlerInAnonymousWrapper(handler)))) {
endpointName = beanName;
break;
}
}
catch (Exception ex) {
logger.trace("Could not get handler from bean = " + beanName, ex);
endpoint = null;
}
}
catch (Exception e) {
logger.trace("Could not get handler from bean = " + beanName, e);
endpoint = null;
}
}
else {
endpointName = endpoint.getBeanName();
}
IntegrationManagement messageHandlerMetrics =
buildMessageHandlerMetrics(monitor2, name, source, endpoint);
buildMessageHandlerMetrics(monitor, endpointName, source, endpoint);
if (endpointName != null) {
this.endpointsByMonitor.put(messageHandlerMetrics, endpointName);
}
return messageHandlerMetrics;
}
private IntegrationManagement buildMessageHandlerMetrics(
IntegrationManagement monitor2,
String name, String source, IntegrationConsumer endpoint) {
IntegrationManagement result = monitor2;
String managedType = source;
String managedName = name;
@@ -907,16 +912,16 @@ public class IntegrationMBeanExporter extends MBeanExporter
}
if (managedName == null) {
managedName = ((NamedComponent) monitor2).getComponentName();
managedName = monitor2.getComponentName();
if (managedName == null) {
managedName = monitor2.toString();
}
managedType = "handler";
}
result.setManagedType(managedType);
result.setManagedName(managedName);
return result;
monitor2.setManagedType(managedType);
monitor2.setManagedName(managedName);
return monitor2;
}
private String buildAnonymousManagedName(Map<Object, AtomicLong> anonymousCache, MessageChannel messageChannel) {
@@ -940,7 +945,6 @@ public class IntegrationMBeanExporter extends MBeanExporter
}
private IntegrationInboundManagement enhanceSourceMonitor(IntegrationInboundManagement source2) {
if (source2.getManagedName() != null) {
return source2;
}
@@ -968,7 +972,6 @@ public class IntegrationMBeanExporter extends MBeanExporter
@SuppressWarnings("unlikely-arg-type")
private AbstractEndpoint getEndpointForMonitor(IntegrationInboundManagement source2) {
for (AbstractEndpoint endpoint : this.applicationContext.getBeansOfType(AbstractEndpoint.class).values()) {
Object target = null;
if (source2 instanceof MessagingGatewaySupport && endpoint.equals(source2)) {
@@ -988,7 +991,6 @@ public class IntegrationMBeanExporter extends MBeanExporter
IntegrationInboundManagement source2, String name,
String source, Object endpoint) {
IntegrationInboundManagement result = source2;
String managedType = source;
String managedName = name;
@@ -999,8 +1001,8 @@ public class IntegrationMBeanExporter extends MBeanExporter
try {
target = targetSource.getTarget();
}
catch (Exception e) {
logger.error("Could not get handler from bean = " + managedName);
catch (Exception ex) {
logger.error("Could not get handler from bean = " + managedName, ex);
}
}
@@ -1019,13 +1021,13 @@ public class IntegrationMBeanExporter extends MBeanExporter
}
if (managedName == null) {
managedName = result.toString();
managedName = source2.toString();
managedType = "source";
}
result.setManagedType(managedType);
result.setManagedName(managedName);
return result;
source2.setManagedType(managedType);
source2.setManagedName(managedName);
return source2;
}
}

View File

@@ -1,5 +1,5 @@
/*
* Copyright 2014-2019 the original author or authors.
* Copyright 2014-2022 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.
@@ -336,21 +336,9 @@ public class ScatterGatherHandlerIntegrationTests {
return new DirectChannel();
}
@Bean
@ServiceActivator(inputChannel = "serviceChannel2")
public MessageHandler service2() {
return new AbstractReplyProducingMessageHandler() {
{
setOutputChannel(gatherChannel());
}
@Override
protected Object handleRequestMessage(Message<?> requestMessage) {
return Math.random();
}
};
@ServiceActivator(inputChannel = "serviceChannel2", outputChannel = "gatherChannel")
public double service2() {
return Math.random();
}
}