INT-3649: Support for Payload Application Events

JIRA: https://jira.spring.io/browse/INT-3649

Polishing according to the PR comments

Address PR comments
This commit is contained in:
Artem Bilan
2015-07-22 16:30:34 -04:00
committed by Gary Russell
parent 95a9e0d2dd
commit 4266bb53b4
11 changed files with 240 additions and 101 deletions

View File

@@ -1,5 +1,5 @@
/*
* Copyright 2002-2010 the original author or authors.
* Copyright 2002-2015 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.
@@ -22,17 +22,21 @@ import org.springframework.beans.factory.support.AbstractBeanDefinition;
import org.springframework.beans.factory.support.BeanDefinitionBuilder;
import org.springframework.beans.factory.xml.ParserContext;
import org.springframework.integration.config.xml.AbstractOutboundChannelAdapterParser;
import org.springframework.integration.config.xml.IntegrationNamespaceUtils;
import org.springframework.integration.event.outbound.ApplicationEventPublishingMessageHandler;
/**
* @author Oleg Zhurakousky
* @author Artem Bilan
* @since 2.0
*/
public class EventOutboundChannelAdapterParser extends AbstractOutboundChannelAdapterParser{
public class EventOutboundChannelAdapterParser extends AbstractOutboundChannelAdapterParser {
@Override
protected AbstractBeanDefinition parseConsumer(Element element, ParserContext parserContext) {
BeanDefinitionBuilder builder = BeanDefinitionBuilder.genericBeanDefinition(
"org.springframework.integration.event.outbound.ApplicationEventPublishingMessageHandler");
BeanDefinitionBuilder builder =
BeanDefinitionBuilder.genericBeanDefinition(ApplicationEventPublishingMessageHandler.class);
IntegrationNamespaceUtils.setValueIfAttributeDefined(builder, element, "publish-payload");
return builder.getBeanDefinition();
}

View File

@@ -16,18 +16,19 @@
package org.springframework.integration.event.inbound;
import java.util.Arrays;
import java.util.HashSet;
import java.util.Set;
import org.springframework.context.ApplicationEvent;
import org.springframework.context.ApplicationListener;
import org.springframework.context.PayloadApplicationEvent;
import org.springframework.context.event.ApplicationEventMulticaster;
import org.springframework.context.event.ContextClosedEvent;
import org.springframework.context.event.ContextStoppedEvent;
import org.springframework.context.event.SmartApplicationListener;
import org.springframework.context.event.GenericApplicationListener;
import org.springframework.context.support.AbstractApplicationContext;
import org.springframework.core.Ordered;
import org.springframework.core.ResolvableType;
import org.springframework.integration.endpoint.ExpressionMessageProducerSupport;
import org.springframework.messaging.Message;
import org.springframework.util.Assert;
@@ -41,14 +42,13 @@ import org.springframework.util.Assert;
* @author Mark Fisher
* @author Artem Bilan
* @author Gary Russell
*
* @see ApplicationEventMulticaster
* @see ExpressionMessageProducerSupport
*/
public class ApplicationEventListeningMessageProducer extends ExpressionMessageProducerSupport
implements SmartApplicationListener {
implements GenericApplicationListener {
private volatile Set<Class<? extends ApplicationEvent>> eventTypes;
private volatile Set<ResolvableType> eventTypes;
private ApplicationEventMulticaster applicationEventMulticaster;
@@ -63,15 +63,19 @@ public class ApplicationEventListeningMessageProducer extends ExpressionMessageP
* In addition, this method re-registers the current instance as a {@link ApplicationListener}
* with the {@link ApplicationEventMulticaster} which clears the listener cache. The cache will be
* refreshed on the next appropriate {@link ApplicationEvent}.
*
* @param eventTypes The event types.
* @see ApplicationEventMulticaster#addApplicationListener
* @see #supportsEventType
*/
@SafeVarargs
public final void setEventTypes(Class<? extends ApplicationEvent>... eventTypes) {
Set<Class<? extends ApplicationEvent>> eventSet = new HashSet<Class<? extends ApplicationEvent>>(
Arrays.asList(eventTypes));
eventSet.remove(null);
public final void setEventTypes(Class<?>... eventTypes) {
Assert.notNull(eventTypes, "'eventTypes' must not be null");
Set<ResolvableType> eventSet = new HashSet<ResolvableType>(eventTypes.length);
for (Class<?> eventType : eventTypes) {
if (eventType != null) {
eventSet.add(ResolvableType.forClass(eventType));
}
}
this.eventTypes = (eventSet.size() > 0 ? eventSet : null);
if (this.applicationEventMulticaster != null) {
@@ -98,13 +102,13 @@ public class ApplicationEventListeningMessageProducer extends ExpressionMessageP
@Override
public void onApplicationEvent(ApplicationEvent event) {
if (this.active || ((event instanceof ContextStoppedEvent || event instanceof ContextClosedEvent)
&& this.stoppedRecently())) {
&& this.stoppedRecently())) {
if (event.getSource() instanceof Message<?>) {
this.sendMessage((Message<?>) event.getSource());
}
else {
Message<?> message = null;
Object result = this.evaluatePayloadExpression(event);
Object result = extractObjectToSend(event);
if (result instanceof Message) {
message = (Message<?>) result;
}
@@ -116,20 +120,45 @@ public class ApplicationEventListeningMessageProducer extends ExpressionMessageP
}
}
private Object extractObjectToSend(Object root) {
if (root instanceof PayloadApplicationEvent) {
return ((PayloadApplicationEvent<?>) root).getPayload();
}
return evaluatePayloadExpression(root);
}
private boolean stoppedRecently() {
return this.stoppedAt > System.currentTimeMillis() - 5000;
}
@Override
public boolean supportsEventType(Class<? extends ApplicationEvent> eventType) {
public boolean supportsEventType(ResolvableType eventType) {
if (this.eventTypes == null) {
return true;
}
for (Class<? extends ApplicationEvent> type : this.eventTypes) {
for (ResolvableType type : this.eventTypes) {
if (type.isAssignableFrom(eventType)) {
return true;
}
}
if (eventType.getRawClass() != null
&& PayloadApplicationEvent.class.isAssignableFrom(eventType.getRawClass())) {
if (eventType.hasUnresolvableGenerics()) {
return true;
}
ResolvableType payloadType = eventType.as(PayloadApplicationEvent.class).getGeneric();
for (ResolvableType type : this.eventTypes) {
if (type.isAssignableFrom(payloadType)) {
return true;
}
}
}
return false;
}

View File

@@ -1,5 +1,5 @@
/*
* Copyright 2002-2010 the original author or authors.
* Copyright 2002-2015 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,17 +25,36 @@ import org.springframework.integration.handler.AbstractMessageHandler;
import org.springframework.util.Assert;
/**
* A {@link org.springframework.messaging.MessageHandler} that publishes each {@link Message} it receives as
* a {@link MessagingEvent}. The {@link MessagingEvent} is a subclass of
* A {@link org.springframework.messaging.MessageHandler} that publishes each {@link Message}
* it receives as a {@link MessagingEvent}. The {@link MessagingEvent} is a subclass of
* Spring's {@link ApplicationEvent} used by this adapter to simply wrap the
* {@link Message}.
*
* <p>
* If the {@link #publishPayload} flag is specified to {@code true}, the {@code payload}
* will be published as is without wrapping to any {@link ApplicationEvent}.
*
* @author Mark Fisher
* @author Artem Bilan
*/
public class ApplicationEventPublishingMessageHandler extends AbstractMessageHandler implements ApplicationEventPublisherAware {
public class ApplicationEventPublishingMessageHandler extends AbstractMessageHandler
implements ApplicationEventPublisherAware {
private ApplicationEventPublisher applicationEventPublisher;
private boolean publishPayload;
/**
* Specify if {@code payload} should be published as is
* or the whole {@code message} must be wrapped to the {@link MessagingEvent}.
* @param publishPayload the {@code boolean} flag to wrap the {@code message}
* to the {@link MessagingEvent} or publish {@code payload}
* as is. Defaults to {@code false}.
* @since 4.2
* @see ApplicationEventPublisher#publishEvent(Object)
*/
public void setPublishPayload(boolean publishPayload) {
this.publishPayload = publishPayload;
}
public void setApplicationEventPublisher(ApplicationEventPublisher applicationEventPublisher) {
this.applicationEventPublisher = applicationEventPublisher;
@@ -47,6 +66,9 @@ public class ApplicationEventPublishingMessageHandler extends AbstractMessageHan
if (message.getPayload() instanceof ApplicationEvent) {
this.applicationEventPublisher.publishEvent((ApplicationEvent) message.getPayload());
}
else if (this.publishPayload) {
this.applicationEventPublisher.publishEvent(message.getPayload());
}
else {
this.applicationEventPublisher.publishEvent(new MessagingEvent(message));
}

View File

@@ -43,8 +43,9 @@
<xsd:attribute name="event-types" type="xsd:string" use="optional">
<xsd:annotation>
<xsd:documentation>
Comma delimited list of event types (classes that extend ApplicationEvent) that this adapter
should send to the message channel. By default, all event types will be sent [OPTIONAL]
Comma delimited list of event types (classes that extend ApplicationEvent or any type
which can be treated as event 'payload') that this adapter should send to the message
channel. By default, all event types will be sent [OPTIONAL]
</xsd:documentation>
</xsd:annotation>
</xsd:attribute>
@@ -52,7 +53,8 @@
<xsd:annotation>
<xsd:documentation><![CDATA[
SpEL expression to be evaluated against the ApplicationEvent to create the payload instance.
If not provided, the ApplicationEvent itself will be the payload.
If not provided, the ApplicationEvent itself will be the 'payload'.
In case of 'PayloadApplicationEvent' the 'payload' is extracted from the event to send.
]]></xsd:documentation>
</xsd:annotation>
</xsd:attribute>
@@ -69,7 +71,8 @@
<xsd:complexType>
<xsd:choice minOccurs="0" maxOccurs="2">
<xsd:element ref="integration:poller" minOccurs="0" maxOccurs="1" />
<xsd:element name="request-handler-advice-chain" type="integration:handlerAdviceChainType" minOccurs="0" maxOccurs="1" />
<xsd:element name="request-handler-advice-chain" type="integration:handlerAdviceChainType"
minOccurs="0" maxOccurs="1" />
</xsd:choice>
<xsd:attributeGroup ref="integration:channelAdapterAttributes"/>
<xsd:attribute name="order">
@@ -80,6 +83,19 @@
</xsd:documentation>
</xsd:annotation>
</xsd:attribute>
<xsd:attribute name="publish-payload" default="false">
<xsd:annotation>
<xsd:documentation>
Specify if 'payload' should be published as is or the whole 'message'
must be wrapped to the 'MessagingEvent'.
See 'ApplicationEventPublisher#publishEvent(Object)' for more information.
Defaults to 'false'.
</xsd:documentation>
</xsd:annotation>
<xsd:simpleType>
<xsd:union memberTypes="xsd:boolean xsd:string" />
</xsd:simpleType>
</xsd:attribute>
</xsd:complexType>
</xsd:element>

View File

@@ -1,32 +1,34 @@
<?xml version="1.0" encoding="UTF-8"?>
<beans xmlns="http://www.springframework.org/schema/beans"
xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
xsi:schemaLocation="http://www.springframework.org/schema/beans http://www.springframework.org/schema/beans/spring-beans.xsd
xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
xsi:schemaLocation="http://www.springframework.org/schema/beans http://www.springframework.org/schema/beans/spring-beans.xsd
http://www.springframework.org/schema/context http://www.springframework.org/schema/context/spring-context.xsd
http://www.springframework.org/schema/integration http://www.springframework.org/schema/integration/spring-integration.xsd
http://www.springframework.org/schema/integration/event http://www.springframework.org/schema/integration/event/spring-integration-event.xsd"
xmlns:context="http://www.springframework.org/schema/context"
xmlns:int="http://www.springframework.org/schema/integration"
xmlns:int-event="http://www.springframework.org/schema/integration/event">
xmlns:context="http://www.springframework.org/schema/context"
xmlns:int="http://www.springframework.org/schema/integration"
xmlns:int-event="http://www.springframework.org/schema/integration/event">
<int:message-history/>
<int-event:inbound-channel-adapter id="eventAdapterSimple" channel="input"
error-channel="errorChannel"/>
error-channel="errorChannel"/>
<int:channel id="input">
<int:queue/>
</int:channel>
<int-event:inbound-channel-adapter id="eventAdapterFiltered" channel="inputFiltered" event-types="org.springframework.integration.event.config.EventInboundChannelAdapterParserTests$AnotherSampleEvent,
org.springframework.integration.event.config.EventInboundChannelAdapterParserTests$SampleEvent"/>
<int-event:inbound-channel-adapter id="eventAdapterFiltered" channel="inputFiltered"
event-types="org.springframework.integration.event.config.EventInboundChannelAdapterParserTests$AnotherSampleEvent,
org.springframework.integration.event.config.EventInboundChannelAdapterParserTests$SampleEvent,
java.util.Date"/>
<int:channel id="inputFiltered">
<int:queue/>
</int:channel>
<int-event:inbound-channel-adapter id="eventAdapterFilteredPlaceHolder" channel="inputFilteredPlaceHolder"
event-types="${event.types}"/>
<int-event:inbound-channel-adapter id="eventAdapterFilteredPlaceHolder" channel="inputFilteredPlaceHolder"
event-types="${event.types}"/>
<int:channel id="inputFilteredPlaceHolder">
<int:queue/>
@@ -42,6 +44,7 @@
<int:bridge input-channel="autoChannel" output-channel="nullChannel"/>
<context:property-placeholder location="classpath:org/springframework/integration/event/config/inbound-adapter.properties"/>
<context:property-placeholder
location="classpath:org/springframework/integration/event/config/inbound-adapter.properties"/>
</beans>

View File

@@ -1,5 +1,5 @@
/*
* Copyright 2002-2013 the original author or authors.
* Copyright 2002-2015 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.
@@ -22,11 +22,11 @@ import static org.junit.Assert.assertNull;
import static org.junit.Assert.assertSame;
import static org.junit.Assert.assertTrue;
import java.util.Date;
import java.util.Properties;
import java.util.Set;
import org.junit.Assert;
import org.junit.Test;
import org.junit.runner.RunWith;
@@ -36,13 +36,14 @@ import org.springframework.beans.factory.annotation.Qualifier;
import org.springframework.context.ApplicationContext;
import org.springframework.context.ApplicationEvent;
import org.springframework.context.event.ContextRefreshedEvent;
import org.springframework.core.ResolvableType;
import org.springframework.expression.Expression;
import org.springframework.messaging.Message;
import org.springframework.messaging.MessageChannel;
import org.springframework.messaging.PollableChannel;
import org.springframework.integration.event.inbound.ApplicationEventListeningMessageProducer;
import org.springframework.integration.history.MessageHistory;
import org.springframework.integration.test.util.TestUtils;
import org.springframework.messaging.Message;
import org.springframework.messaging.MessageChannel;
import org.springframework.messaging.PollableChannel;
import org.springframework.test.context.ContextConfiguration;
import org.springframework.test.context.junit4.SpringJUnit4ClassRunner;
@@ -51,6 +52,7 @@ import org.springframework.test.context.junit4.SpringJUnit4ClassRunner;
* @author Mark Fisher
* @author Gary Russell
* @author Gunnar Hillert
* @author Artem Bilan
* @since 2.0
*/
@RunWith(SpringJUnit4ClassRunner.class)
@@ -66,7 +68,8 @@ public class EventInboundChannelAdapterParserTests {
@Autowired
MessageChannel autoChannel;
@Autowired @Qualifier("autoChannel.adapter")
@Autowired
@Qualifier("autoChannel.adapter")
ApplicationEventListeningMessageProducer eventListener;
@Test
@@ -87,11 +90,12 @@ public class EventInboundChannelAdapterParserTests {
Assert.assertTrue(adapter instanceof ApplicationEventListeningMessageProducer);
DirectFieldAccessor adapterAccessor = new DirectFieldAccessor(adapter);
Assert.assertEquals(context.getBean("inputFiltered"), adapterAccessor.getPropertyValue("outputChannel"));
Set<Class<? extends ApplicationEvent>> eventTypes = (Set<Class<? extends ApplicationEvent>>) adapterAccessor.getPropertyValue("eventTypes");
Set<ResolvableType> eventTypes = (Set<ResolvableType>) adapterAccessor.getPropertyValue("eventTypes");
assertNotNull(eventTypes);
assertTrue(eventTypes.size() == 2);
assertTrue(eventTypes.contains(SampleEvent.class));
assertTrue(eventTypes.contains(AnotherSampleEvent.class));
assertTrue(eventTypes.size() == 3);
assertTrue(eventTypes.contains(ResolvableType.forClass(SampleEvent.class)));
assertTrue(eventTypes.contains(ResolvableType.forClass(AnotherSampleEvent.class)));
assertTrue(eventTypes.contains(ResolvableType.forClass(Date.class)));
assertNull(adapterAccessor.getPropertyValue("errorChannel"));
}
@@ -103,11 +107,11 @@ public class EventInboundChannelAdapterParserTests {
Assert.assertTrue(adapter instanceof ApplicationEventListeningMessageProducer);
DirectFieldAccessor adapterAccessor = new DirectFieldAccessor(adapter);
Assert.assertEquals(context.getBean("inputFilteredPlaceHolder"), adapterAccessor.getPropertyValue("outputChannel"));
Set<Class<? extends ApplicationEvent>> eventTypes = (Set<Class<? extends ApplicationEvent>>) adapterAccessor.getPropertyValue("eventTypes");
Set<ResolvableType> eventTypes = (Set<ResolvableType>) adapterAccessor.getPropertyValue("eventTypes");
assertNotNull(eventTypes);
assertTrue(eventTypes.size() == 2);
assertTrue(eventTypes.contains(SampleEvent.class));
assertTrue(eventTypes.contains(AnotherSampleEvent.class));
assertTrue(eventTypes.contains(ResolvableType.forClass(SampleEvent.class)));
assertTrue(eventTypes.contains(ResolvableType.forClass(AnotherSampleEvent.class)));
}
@Test
@@ -142,15 +146,20 @@ public class EventInboundChannelAdapterParserTests {
@SuppressWarnings("serial")
public static class SampleEvent extends ApplicationEvent {
public SampleEvent(Object source) {
super(source);
}
}
@SuppressWarnings("serial")
public static class AnotherSampleEvent extends ApplicationEvent {
public AnotherSampleEvent(Object source) {
super(source);
}
}
}

View File

@@ -1,21 +1,21 @@
<?xml version="1.0" encoding="UTF-8"?>
<beans xmlns="http://www.springframework.org/schema/beans"
xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
xmlns:int="http://www.springframework.org/schema/integration"
xmlns:int-event="http://www.springframework.org/schema/integration/event"
xsi:schemaLocation="http://www.springframework.org/schema/beans http://www.springframework.org/schema/beans/spring-beans.xsd
xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
xmlns:int="http://www.springframework.org/schema/integration"
xmlns:int-event="http://www.springframework.org/schema/integration/event"
xsi:schemaLocation="http://www.springframework.org/schema/beans http://www.springframework.org/schema/beans/spring-beans.xsd
http://www.springframework.org/schema/integration http://www.springframework.org/schema/integration/spring-integration.xsd
http://www.springframework.org/schema/integration/event http://www.springframework.org/schema/integration/event/spring-integration-event.xsd">
<int:channel id="input"/>
<int-event:outbound-channel-adapter id="eventAdapter" channel="input"/>
<int-event:outbound-channel-adapter id="eventAdapter" channel="input" publish-payload="true"/>
<int:channel id="inputAdvice"/>
<int-event:outbound-channel-adapter id="withAdvice" channel="inputAdvice">
<int-event:request-handler-advice-chain>
<bean class="org.springframework.integration.event.config.EventOutboundChannelAdapterParserTests$FooAdvice" />
<bean class="org.springframework.integration.event.config.EventOutboundChannelAdapterParserTests$FooAdvice"/>
</int-event:request-handler-advice-chain>
</int-event:outbound-channel-adapter>

View File

@@ -1,5 +1,5 @@
/*
* Copyright 2002-2013 the original author or authors.
* Copyright 2002-2015 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.
@@ -28,14 +28,16 @@ import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.context.ApplicationEvent;
import org.springframework.context.ApplicationListener;
import org.springframework.context.ConfigurableApplicationContext;
import org.springframework.context.PayloadApplicationEvent;
import org.springframework.context.support.ClassPathXmlApplicationContext;
import org.springframework.messaging.Message;
import org.springframework.integration.channel.DirectChannel;
import org.springframework.integration.channel.QueueChannel;
import org.springframework.messaging.MessageHandler;
import org.springframework.integration.endpoint.EventDrivenConsumer;
import org.springframework.integration.event.outbound.ApplicationEventPublishingMessageHandler;
import org.springframework.integration.handler.advice.AbstractRequestHandlerAdvice;
import org.springframework.integration.test.util.TestUtils;
import org.springframework.messaging.Message;
import org.springframework.messaging.MessageHandler;
import org.springframework.messaging.support.GenericMessage;
import org.springframework.test.context.ContextConfiguration;
import org.springframework.test.context.junit4.SpringJUnit4ClassRunner;
@@ -60,92 +62,101 @@ public class EventOutboundChannelAdapterParserTests {
@Test
public void validateEventParser() {
EventDrivenConsumer adapter = context.getBean("eventAdapter", EventDrivenConsumer.class);
EventDrivenConsumer adapter = this.context.getBean("eventAdapter", EventDrivenConsumer.class);
Assert.assertNotNull(adapter);
DirectFieldAccessor adapterAccessor = new DirectFieldAccessor(adapter);
MessageHandler handler = (MessageHandler) adapterAccessor.getPropertyValue("handler");
Assert.assertTrue(handler instanceof ApplicationEventPublishingMessageHandler);
Assert.assertEquals(context.getBean("input"), adapterAccessor.getPropertyValue("inputChannel"));
Assert.assertEquals(this.context.getBean("input"), adapterAccessor.getPropertyValue("inputChannel"));
Assert.assertTrue(TestUtils.getPropertyValue(handler, "publishPayload", Boolean.class));
}
@Test
public void validateUsage() {
ApplicationListener<?> listener = new ApplicationListener<ApplicationEvent>() {
@Override
public void onApplicationEvent(ApplicationEvent event) {
Object source = event.getSource();
if (source instanceof Message){
String payload = (String) ((Message<?>) source).getPayload();
if (event instanceof PayloadApplicationEvent) {
String payload = (String) ((PayloadApplicationEvent<?>) event).getPayload();
if (payload.equals("hello")) {
receivedEvent = true;
}
}
}
};
context.addApplicationListener(listener);
this.context.addApplicationListener(listener);
DirectChannel channel = context.getBean("input", DirectChannel.class);
channel.send(new GenericMessage<String>("hello"));
Assert.assertTrue(receivedEvent);
Assert.assertTrue(this.receivedEvent);
}
@Test
public void withAdvice() {
receivedEvent = false;
this.receivedEvent = false;
ApplicationListener<?> listener = new ApplicationListener<ApplicationEvent>() {
@Override
public void onApplicationEvent(ApplicationEvent event) {
Object source = event.getSource();
if (source instanceof Message){
if (source instanceof Message) {
String payload = (String) ((Message<?>) source).getPayload();
if (payload.equals("hello")) {
receivedEvent = true;
}
}
}
};
context.addApplicationListener(listener);
DirectChannel channel = context.getBean("inputAdvice", DirectChannel.class);
channel.send(new GenericMessage<String>("hello"));
Assert.assertTrue(receivedEvent);
Assert.assertTrue(this.receivedEvent);
Assert.assertEquals(1, adviceCalled);
}
@Test //INT-2275
public void testInsideChain() {
receivedEvent = false;
this.receivedEvent = false;
ApplicationListener<?> listener = new ApplicationListener<ApplicationEvent>() {
@Override
public void onApplicationEvent(ApplicationEvent event) {
Object source = event.getSource();
if (source instanceof Message){
if (source instanceof Message) {
String payload = (String) ((Message<?>) source).getPayload();
if (payload.equals("foobar")) {
receivedEvent = true;
}
}
}
};
context.addApplicationListener(listener);
this.context.addApplicationListener(listener);
DirectChannel channel = context.getBean("inputChain", DirectChannel.class);
channel.send(new GenericMessage<String>("foo"));
Assert.assertTrue(receivedEvent);
Assert.assertTrue(this.receivedEvent);
}
@Test(timeout=10000)
@Test(timeout = 10000)
public void validateUsageWithPollableChannel() throws Exception {
receivedEvent = false;
ConfigurableApplicationContext context = new ClassPathXmlApplicationContext("EventOutboundChannelAdapterParserTestsWithPollable-context.xml", EventOutboundChannelAdapterParserTests.class);
final CyclicBarrier barier = new CyclicBarrier(2);
this.receivedEvent = false;
ConfigurableApplicationContext context =
new ClassPathXmlApplicationContext("EventOutboundChannelAdapterParserTestsWithPollable-context.xml",
EventOutboundChannelAdapterParserTests.class);
final CyclicBarrier barrier = new CyclicBarrier(2);
ApplicationListener<?> listener = new ApplicationListener<ApplicationEvent>() {
@Override
public void onApplicationEvent(ApplicationEvent event) {
Object source = event.getSource();
if (source instanceof Message){
if (source instanceof Message) {
String payload = (String) ((Message<?>) source).getPayload();
if (payload.equals("hello")){
if (payload.equals("hello")) {
receivedEvent = true;
try {
barier.await();
barrier.await();
}
catch (InterruptedException e) {
Thread.currentThread().interrupt();
@@ -156,12 +167,14 @@ public class EventOutboundChannelAdapterParserTests {
}
}
}
};
context.addApplicationListener(listener);
QueueChannel channel = context.getBean("input", QueueChannel.class);
channel.send(new GenericMessage<String>("hello"));
barier.await();
Assert.assertTrue(receivedEvent);
barrier.await();
Assert.assertTrue(this.receivedEvent);
context.close();
}
public static class FooAdvice extends AbstractRequestHandlerAdvice {
@@ -173,4 +186,5 @@ public class EventOutboundChannelAdapterParserTests {
}
}
}

View File

@@ -43,6 +43,7 @@ import org.springframework.context.event.SimpleApplicationEventMulticaster;
import org.springframework.context.support.AbstractApplicationContext;
import org.springframework.context.support.ClassPathXmlApplicationContext;
import org.springframework.context.support.GenericApplicationContext;
import org.springframework.core.ResolvableType;
import org.springframework.messaging.Message;
import org.springframework.messaging.MessageHandlingException;
import org.springframework.integration.channel.DirectChannel;
@@ -68,9 +69,9 @@ public class ApplicationEventListeningMessageProducerTests {
adapter.start();
Message<?> message1 = channel.receive(0);
assertNull(message1);
assertTrue(adapter.supportsEventType(TestApplicationEvent1.class));
assertTrue(adapter.supportsEventType(ResolvableType.forClass(TestApplicationEvent1.class)));
adapter.onApplicationEvent(new TestApplicationEvent1());
assertTrue(adapter.supportsEventType(TestApplicationEvent2.class));
assertTrue(adapter.supportsEventType(ResolvableType.forClass(TestApplicationEvent2.class)));
adapter.onApplicationEvent(new TestApplicationEvent2());
Message<?> message2 = channel.receive(20);
assertNotNull(message2);
@@ -90,25 +91,25 @@ public class ApplicationEventListeningMessageProducerTests {
adapter.start();
Message<?> message1 = channel.receive(0);
assertNull(message1);
assertTrue(adapter.supportsEventType(TestApplicationEvent1.class));
assertTrue(adapter.supportsEventType(ResolvableType.forClass(TestApplicationEvent1.class)));
adapter.onApplicationEvent(new TestApplicationEvent1());
assertFalse(adapter.supportsEventType(TestApplicationEvent2.class));
assertFalse(adapter.supportsEventType(ResolvableType.forClass(TestApplicationEvent2.class)));
Message<?> message2 = channel.receive(20);
assertNotNull(message2);
assertEquals("event1", ((ApplicationEvent) message2.getPayload()).getSource());
assertNull(channel.receive(0));
adapter.setEventTypes((Class<? extends ApplicationEvent>) null);
assertTrue(adapter.supportsEventType(TestApplicationEvent1.class));
assertTrue(adapter.supportsEventType(TestApplicationEvent2.class));
assertTrue(adapter.supportsEventType(ResolvableType.forClass(TestApplicationEvent1.class)));
assertTrue(adapter.supportsEventType(ResolvableType.forClass(TestApplicationEvent2.class)));
adapter.setEventTypes(null, TestApplicationEvent2.class, null);
assertFalse(adapter.supportsEventType(TestApplicationEvent1.class));
assertTrue(adapter.supportsEventType(TestApplicationEvent2.class));
assertFalse(adapter.supportsEventType(ResolvableType.forClass(TestApplicationEvent1.class)));
assertTrue(adapter.supportsEventType(ResolvableType.forClass(TestApplicationEvent2.class)));
adapter.setEventTypes(null, null);
assertTrue(adapter.supportsEventType(TestApplicationEvent1.class));
assertTrue(adapter.supportsEventType(TestApplicationEvent2.class));
assertTrue(adapter.supportsEventType(ResolvableType.forClass(TestApplicationEvent1.class)));
assertTrue(adapter.supportsEventType(ResolvableType.forClass(TestApplicationEvent2.class)));
}
@Test
@@ -266,7 +267,7 @@ public class ApplicationEventListeningMessageProducerTests {
Set<?> listeners = TestUtils.getPropertyValue(entry.getValue(), "applicationListenerBeans", Set.class);
assertEquals(2, listeners.size());
for (Object listener : listeners) {
assertThat((String) listener,
assertThat((String) listener,
Matchers.is(Matchers.isOneOf("testListenerMessageProducer", "testListener")));
}
break;
@@ -282,6 +283,30 @@ public class ApplicationEventListeningMessageProducerTests {
assertNotNull(receive);
assertSame(event2, receive.getPayload());
assertNull(channel.receive(1));
ctx.close();
}
@Test
public void testPayloadEvents() {
GenericApplicationContext ctx = TestUtils.createTestApplicationContext();
ConfigurableListableBeanFactory beanFactory = ctx.getBeanFactory();
QueueChannel channel = new QueueChannel();
ApplicationEventListeningMessageProducer listenerMessageProducer =
new ApplicationEventListeningMessageProducer();
listenerMessageProducer.setOutputChannel(channel);
listenerMessageProducer.setEventTypes(String.class);
beanFactory.registerSingleton("testListenerMessageProducer", listenerMessageProducer);
ctx.refresh();
ctx.publishEvent("foo");
Message<?> receive = channel.receive(10000);
assertNotNull(receive);
assertEquals("foo", receive.getPayload());
ctx.close();
}