INT-3926: Add <poller> for all outbound tags
JIRA: https://jira.spring.io/browse/INT-3926 * Add `<poller>` for all `<outbound-adapter(gateway)>` * Add/modify tests for `<queue>` as an input channel and particular `<poller>` * Move XSD from the `xml` package to the `config` level to achieve the consistency throughout the project.
This commit is contained in:
committed by
Gary Russell
parent
afc286a66a
commit
74cbc58698
@@ -721,6 +721,7 @@
|
||||
<xsd:complexType>
|
||||
<xsd:sequence>
|
||||
<xsd:sequence>
|
||||
<xsd:element ref="integration:poller" minOccurs="0" maxOccurs="1" />
|
||||
<xsd:element name="sql-parameter-definition" minOccurs="0"
|
||||
maxOccurs="unbounded" type="sqlParameterDefinitionType">
|
||||
<xsd:annotation>
|
||||
|
||||
@@ -18,7 +18,13 @@
|
||||
|
||||
<int:poller id="defaultPoller" default="true" fixed-rate="5000"/>
|
||||
|
||||
<int:channel id="startChannel"/>
|
||||
<int:channel id="startChannel">
|
||||
<int:queue/>
|
||||
</int:channel>
|
||||
|
||||
<int:channel id="startErrorsChannel">
|
||||
<int:queue/>
|
||||
</int:channel>
|
||||
|
||||
<int-jdbc:stored-proc-outbound-gateway request-channel="startChannel"
|
||||
data-source="dataSource"
|
||||
@@ -29,6 +35,7 @@
|
||||
is-function="false"
|
||||
expect-single-result="true"
|
||||
reply-channel="outputChannel">
|
||||
<int:poller fixed-delay="100" error-channel="startErrorsChannel"/>
|
||||
<int-jdbc:parameter name="username" expression="payload.username"/>
|
||||
<int-jdbc:parameter name="password" expression="payload.password"/>
|
||||
<int-jdbc:parameter name="email" expression="payload.email"/>
|
||||
|
||||
@@ -1,5 +1,5 @@
|
||||
/*
|
||||
* Copyright 2002-2014 the original author or authors.
|
||||
* Copyright 2002-2016 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,8 +16,10 @@
|
||||
|
||||
package org.springframework.integration.jdbc;
|
||||
|
||||
import static org.hamcrest.Matchers.instanceOf;
|
||||
import static org.junit.Assert.assertEquals;
|
||||
import static org.junit.Assert.assertNotNull;
|
||||
import static org.junit.Assert.assertThat;
|
||||
import static org.junit.Assert.assertTrue;
|
||||
|
||||
import java.sql.CallableStatement;
|
||||
@@ -48,8 +50,10 @@ import org.springframework.integration.support.json.JsonOutboundMessageMapper;
|
||||
import org.springframework.jdbc.core.JdbcTemplate;
|
||||
import org.springframework.jdbc.core.SqlReturnType;
|
||||
import org.springframework.messaging.Message;
|
||||
import org.springframework.messaging.MessageChannel;
|
||||
import org.springframework.messaging.MessageHandlingException;
|
||||
import org.springframework.messaging.PollableChannel;
|
||||
import org.springframework.messaging.support.ErrorMessage;
|
||||
import org.springframework.messaging.support.GenericMessage;
|
||||
import org.springframework.test.annotation.DirtiesContext;
|
||||
import org.springframework.test.context.ContextConfiguration;
|
||||
@@ -73,10 +77,14 @@ public class StoredProcOutboundGatewayWithSpelIntegrationTests {
|
||||
|
||||
@Autowired
|
||||
@Qualifier("startChannel")
|
||||
DirectChannel channel;
|
||||
MessageChannel channel;
|
||||
|
||||
@Autowired
|
||||
DirectChannel getMessageChannel;
|
||||
@Qualifier("startErrorsChannel")
|
||||
PollableChannel startErrorsChannel;
|
||||
|
||||
@Autowired
|
||||
MessageChannel getMessageChannel;
|
||||
|
||||
@Autowired
|
||||
PollableChannel output2Channel;
|
||||
@@ -95,8 +103,8 @@ public class StoredProcOutboundGatewayWithSpelIntegrationTests {
|
||||
User user2 = new User("Second User", "my second password", "email2");
|
||||
|
||||
Message<User> user1Message = MessageBuilder.withPayload(user1)
|
||||
.setHeader("my_stored_procedure", "CREATE_USER")
|
||||
.build();
|
||||
.setHeader("my_stored_procedure", "CREATE_USER")
|
||||
.build();
|
||||
Message<User> user2Message = MessageBuilder.withPayload(user2)
|
||||
.setHeader("my_stored_procedure", "CREATE_USER_RETURN_ALL")
|
||||
.build();
|
||||
@@ -132,19 +140,18 @@ public class StoredProcOutboundGatewayWithSpelIntegrationTests {
|
||||
|
||||
Message<User> user1Message = MessageBuilder.withPayload(user1).build();
|
||||
|
||||
try {
|
||||
channel.send(user1Message);
|
||||
} catch (MessageHandlingException e) {
|
||||
this.channel.send(user1Message);
|
||||
|
||||
String expectedMessage = "Unable to resolve Stored Procedure/Function name " +
|
||||
"for the provided Expression 'headers['my_stored_procedure']'.";
|
||||
String actualMessage = e.getCause().getMessage();
|
||||
Assert.assertEquals(expectedMessage, actualMessage);
|
||||
return;
|
||||
}
|
||||
Message<?> receive = this.startErrorsChannel.receive(1000);
|
||||
assertNotNull(receive);
|
||||
assertThat(receive, instanceOf(ErrorMessage.class));
|
||||
|
||||
Assert.fail("Expected a MessageHandlingException to be thrown.");
|
||||
MessageHandlingException exception = (MessageHandlingException) receive.getPayload();
|
||||
|
||||
String expectedMessage = "Unable to resolve Stored Procedure/Function name " +
|
||||
"for the provided Expression 'headers['my_stored_procedure']'.";
|
||||
String actualMessage = exception.getCause().getMessage();
|
||||
Assert.assertEquals(expectedMessage, actualMessage);
|
||||
}
|
||||
|
||||
@Test
|
||||
@@ -161,7 +168,8 @@ public class StoredProcOutboundGatewayWithSpelIntegrationTests {
|
||||
assertNotNull(resultMessage);
|
||||
Object resultPayload = resultMessage.getPayload();
|
||||
assertTrue(resultPayload instanceof String);
|
||||
Message<?> message = new JsonInboundMessageMapper(String.class, new Jackson2JsonMessageParser()).toMessage((String) resultPayload);
|
||||
Message<?> message = new JsonInboundMessageMapper(String.class, new Jackson2JsonMessageParser())
|
||||
.toMessage((String) resultPayload);
|
||||
assertEquals(testMessage.getPayload(), message.getPayload());
|
||||
assertEquals(testMessage.getHeaders().get("FOO"), message.getHeaders().get("FOO"));
|
||||
Mockito.verify(clobSqlReturnType).getTypeValue(Mockito.any(CallableStatement.class),
|
||||
@@ -173,11 +181,11 @@ public class StoredProcOutboundGatewayWithSpelIntegrationTests {
|
||||
private final AtomicInteger count = new AtomicInteger();
|
||||
|
||||
public Integer next() throws InterruptedException {
|
||||
if (count.get()>2){
|
||||
if (count.get() > 2) {
|
||||
//prevent message overload
|
||||
return null;
|
||||
}
|
||||
return Integer.valueOf(count.incrementAndGet());
|
||||
return count.incrementAndGet();
|
||||
}
|
||||
}
|
||||
|
||||
@@ -200,4 +208,5 @@ public class StoredProcOutboundGatewayWithSpelIntegrationTests {
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user