diff --git a/README.adoc b/README.adoc
index 089a43b1..ae5b2045 100644
--- a/README.adoc
+++ b/README.adoc
@@ -165,21 +165,17 @@ The following are the various components of this repository.
You can build everything from the root of the repository.
-NOTE: The build depends on the global property `${revision}` to configure the current project version.
-This value is resolved with the `maven-flatten-plugin`.
-Because of this, you need to install some core poms in the local maven repository before running a build command.
-To perform this necessary step, run `setup-build.sh`.
-
-`./setup-build.sh`
-
`./mvnw clean install`
TIP: By default the build skips long running tests, annotated with `@Tag("integration")`,including tests requiring Testcontainers. To run the full test suite, activate the `integration` Maven profile:
`./mvnw clean install -Pintegration`. This applies to building functions or individual applications as well.
However, this may not be what you are interested in since you are probably interested in a single application or a few of them.
+To build the functions and applications that you are interested in, you need to first build the base module with the following.
-To build the functions and applications that you are interested in, you need to build them selectively, as shown below.
+`./mvnw clean install -f stream-applications-build`
+
+You can then build the desired functions/apps selectively, as shown below.
==== Building Functions
diff --git a/applications/sink/pgcopy-sink/pom.xml b/applications/sink/pgcopy-sink/pom.xml
index d97cf9ba..4ed38b82 100644
--- a/applications/sink/pgcopy-sink/pom.xml
+++ b/applications/sink/pgcopy-sink/pom.xml
@@ -30,9 +30,10 @@
org.springframework.cloud
- spring-cloud-stream-test-support-internal
- 3.0.9.RELEASE
+ spring-cloud-stream
+ test-jar
test
+ test-binder
diff --git a/applications/sink/pgcopy-sink/src/main/java/org/springframework/cloud/stream/app/pgcopy/sink/PgcopySinkConfiguration.java b/applications/sink/pgcopy-sink/src/main/java/org/springframework/cloud/stream/app/pgcopy/sink/PgcopySinkConfiguration.java
index 7c74a87a..80ae6e7b 100644
--- a/applications/sink/pgcopy-sink/src/main/java/org/springframework/cloud/stream/app/pgcopy/sink/PgcopySinkConfiguration.java
+++ b/applications/sink/pgcopy-sink/src/main/java/org/springframework/cloud/stream/app/pgcopy/sink/PgcopySinkConfiguration.java
@@ -20,6 +20,7 @@ import java.sql.Connection;
import java.sql.SQLException;
import java.util.Collection;
import java.util.Collections;
+import java.util.function.Consumer;
import javax.sql.DataSource;
@@ -70,10 +71,10 @@ import org.springframework.util.StringUtils;
*
* @author Thomas Risberg
* @author Janne Valkealahti
+ * @author Chris Bono
*/
@Configuration
@EnableScheduling
-//@EnableBinding(Sink.class)
@EnableConfigurationProperties(PgcopySinkProperties.class)
public class PgcopySinkConfiguration {
@@ -83,13 +84,12 @@ public class PgcopySinkConfiguration {
private PgcopySinkProperties properties;
@Bean
- public MessageChannel toSink() {
- return new DirectChannel();
+ public Consumer> pgcopyConsumer(MessageHandler aggregatingMessageHandler) {
+ return aggregatingMessageHandler::handleMessage;
}
@Bean
@Primary
- @ServiceActivator(inputChannel = "input")
FactoryBean aggregatorFactoryBean(MessageChannel toSink, MessageGroupStore messageGroupStore) {
AggregatorFactoryBean aggregatorFactoryBean = new AggregatorFactoryBean();
aggregatorFactoryBean.setCorrelationStrategy(
@@ -103,6 +103,11 @@ public class PgcopySinkConfiguration {
return aggregatorFactoryBean;
}
+ @Bean
+ public MessageChannel toSink() {
+ return new DirectChannel();
+ }
+
@Bean
@ServiceActivator(inputChannel = "toSink")
public MessageHandler datasetSinkMessageHandler(final JdbcTemplate jdbcTemplate,
diff --git a/applications/sink/pgcopy-sink/src/test/java/org/springframework/cloud/stream/app/pgcopy/sink/PgcopyBadErrorTableIntegrationTests.java b/applications/sink/pgcopy-sink/src/test/java/org/springframework/cloud/stream/app/pgcopy/sink/PgcopyBadErrorTableIntegrationTests.java
index 187951e8..188ddff2 100644
--- a/applications/sink/pgcopy-sink/src/test/java/org/springframework/cloud/stream/app/pgcopy/sink/PgcopyBadErrorTableIntegrationTests.java
+++ b/applications/sink/pgcopy-sink/src/test/java/org/springframework/cloud/stream/app/pgcopy/sink/PgcopyBadErrorTableIntegrationTests.java
@@ -21,16 +21,15 @@ import java.util.ArrayList;
import java.util.List;
import java.util.Properties;
-import org.junit.Assert;
-import org.junit.Before;
-import org.junit.ClassRule;
-import org.junit.Test;
+import org.junit.jupiter.api.BeforeEach;
+import org.junit.jupiter.api.Test;
+import org.junit.jupiter.api.extension.ExtendWith;
import org.postgresql.util.PSQLException;
import org.springframework.boot.SpringApplication;
import org.springframework.boot.autoconfigure.SpringBootApplication;
import org.springframework.boot.test.util.TestPropertyValues;
-import org.springframework.cloud.stream.app.pgcopy.test.PostgresTestSupport;
+import org.springframework.cloud.stream.app.pgcopy.test.PostgresAvailableExtension;
import org.springframework.context.annotation.AnnotationConfigApplicationContext;
import org.springframework.core.io.support.PropertiesLoaderUtils;
import org.springframework.dao.DataAccessException;
@@ -38,7 +37,7 @@ import org.springframework.jdbc.core.JdbcOperations;
import org.springframework.jdbc.core.JdbcTemplate;
import org.springframework.jdbc.datasource.DriverManagerDataSource;
-import static org.hamcrest.Matchers.is;
+import static org.assertj.core.api.Assertions.assertThat;
/**
* Integration Tests testing bad error table specified for PgcopySink. Only runs if PostgreSQL database is available.
@@ -46,18 +45,16 @@ import static org.hamcrest.Matchers.is;
* @author Thomas Risberg
* @author Artem Bilan
*/
+@ExtendWith(PostgresAvailableExtension.class)
public class PgcopyBadErrorTableIntegrationTests {
- @ClassRule
- public static PostgresTestSupport postgresAvailable = new PostgresTestSupport();
-
private String[] env = { "pgcopy.tableName=names", "pgcopy.columns=id,name,age", "pgcopy.format=CSV" };
private String[] jdbc = { };
private Properties appProperties = new Properties();
- @Before
+ @BeforeEach
public void setup() {
try {
appProperties = PropertiesLoaderUtils.loadAllProperties("application.properties");
@@ -95,12 +92,11 @@ public class PgcopyBadErrorTableIntegrationTests {
dae = cause;
}
}
- Assert.assertThat(cause.getClass().getName(), is(PSQLException.class.getName()));
- Assert.assertNotNull(ise);
- Assert.assertTrue(ise.getMessage().contains("Invalid error table specified"));
- Assert.assertNotNull(dae);
- Assert.assertTrue(cause.getMessage().contains("relation"));
- Assert.assertTrue(cause.getMessage().contains("does not exist"));
+ assertThat(cause).isInstanceOf(PSQLException.class)
+ .hasMessageContaining("relation")
+ .hasMessageContaining("does not exist");
+ assertThat(ise).hasMessageContaining("Invalid error table specified");
+ assertThat(dae).isNotNull();
}
context.close();
}
@@ -148,12 +144,11 @@ public class PgcopyBadErrorTableIntegrationTests {
dae = cause;
}
}
- Assert.assertThat(cause.getClass().getName(), is(PSQLException.class.getName()));
- Assert.assertNotNull(ise);
- Assert.assertTrue(ise.getMessage().contains("Invalid error table specified"));
- Assert.assertNotNull(dae);
- Assert.assertTrue(cause.getMessage().contains("column"));
- Assert.assertTrue(cause.getMessage().contains("does not exist"));
+ assertThat(cause).isInstanceOf(PSQLException.class)
+ .hasMessageContaining("column")
+ .hasMessageContaining("does not exist");
+ assertThat(ise).hasMessageContaining("Invalid error table specified");
+ assertThat(dae).isNotNull();
}
context.close();
}
diff --git a/applications/sink/pgcopy-sink/src/test/java/org/springframework/cloud/stream/app/pgcopy/sink/PgcopyErrorTableIntegrationTests.java b/applications/sink/pgcopy-sink/src/test/java/org/springframework/cloud/stream/app/pgcopy/sink/PgcopyErrorTableIntegrationTests.java
index d4eeef51..e06b8205 100644
--- a/applications/sink/pgcopy-sink/src/test/java/org/springframework/cloud/stream/app/pgcopy/sink/PgcopyErrorTableIntegrationTests.java
+++ b/applications/sink/pgcopy-sink/src/test/java/org/springframework/cloud/stream/app/pgcopy/sink/PgcopyErrorTableIntegrationTests.java
@@ -16,60 +16,69 @@
package org.springframework.cloud.stream.app.pgcopy.sink;
-import org.junit.ClassRule;
-import org.junit.Test;
-import org.junit.runner.RunWith;
+import java.util.function.Consumer;
+
+import org.junit.jupiter.api.Test;
+import org.junit.jupiter.api.extension.ExtendWith;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.boot.SpringApplication;
-import org.springframework.boot.autoconfigure.SpringBootApplication;
+import org.springframework.boot.SpringBootConfiguration;
+import org.springframework.boot.autoconfigure.EnableAutoConfiguration;
import org.springframework.boot.test.context.SpringBootTest;
-import org.springframework.cloud.stream.app.pgcopy.test.PostgresTestSupport;
+import org.springframework.cloud.stream.app.pgcopy.test.PostgresAvailableExtension;
+import org.springframework.cloud.stream.binder.test.TestChannelBinderConfiguration;
+import org.springframework.context.annotation.Import;
import org.springframework.jdbc.core.JdbcOperations;
+import org.springframework.messaging.Message;
+import org.springframework.messaging.support.MessageBuilder;
import org.springframework.test.annotation.DirtiesContext;
-import org.springframework.test.context.TestPropertySource;
-import org.springframework.test.context.junit4.SpringRunner;
+
+import static org.assertj.core.api.Assertions.assertThat;
/**
* Integration Tests for PgcopySink with error table. Only runs if PostgreSQL database is available.
*
* @author Thomas Risberg
* @author Janne Valkealahti
+ * @author Chris Bono
*/
-@RunWith(SpringRunner.class)
-@SpringBootTest(webEnvironment = SpringBootTest.WebEnvironment.NONE,
- classes = PgcopyErrorTableIntegrationTests.PgcopySinkApplication.class)
-@TestPropertySource(properties = { "pgcopy.tableName=names", "pgcopy.batch-size=3", "pgcopy.initialize=true",
- "pgcopy.columns=id,name,age", "pgcopy.format=CSV", "pgcopy.error-table=test_errors",
- "spring.datasource.initialization-mode=always", "spring.datasource.schema=classpath:error-table-ddl.sql",
- "spring.datasource.continue-on-error=true" })
-@DirtiesContext(classMode = DirtiesContext.ClassMode.BEFORE_EACH_TEST_METHOD)
+@SpringBootTest(
+ webEnvironment = SpringBootTest.WebEnvironment.NONE,
+ classes = PgcopyErrorTableIntegrationTests.PgcopySinkApplication.class,
+ properties = {
+ "spring.cloud.function.definition=pgcopyConsumer",
+ "pgcopy.tableName=names", "pgcopy.batch-size=3", "pgcopy.initialize=true",
+ "pgcopy.columns=id,name,age", "pgcopy.format=CSV", "pgcopy.error-table=test_errors",
+ "spring.sql.init.mode=always", "spring.sql.init.schema-locations=classpath:error-table-ddl.sql",
+ "spring.sql.init.continue-on-error=true"
+ })
+@ExtendWith(PostgresAvailableExtension.class)
+@DirtiesContext
public class PgcopyErrorTableIntegrationTests {
- @ClassRule
- public static PostgresTestSupport postgresAvailable = new PostgresTestSupport();
-
-// @Autowired
-// protected Sink channels;
+ @Autowired
+ private Consumer> pgcopyConsumer;
@Autowired
- protected JdbcOperations jdbcOperations;
+ private JdbcOperations jdbcOperations;
@Test
public void testCopyCSV() {
- /*
- TODO sink fix
- channels.input().send(MessageBuilder.withPayload("123,Nisse,25").build());
- channels.input().send(MessageBuilder.withPayload("GARBAGE").build());
- channels.input().send(MessageBuilder.withPayload("125,Bubba,22").build());
- int result = jdbcOperations.queryForObject("select count(*) from names", Integer.class);
- int errors = jdbcOperations.queryForObject("select count(*) from test_errors", Integer.class);
- Assert.assertThat(result, is(2));
- Assert.assertThat(errors, is(1));
- */
+ this.pgcopyConsumer.accept(MessageBuilder.withPayload("123,Nisse,25").build());
+ this.pgcopyConsumer.accept(MessageBuilder.withPayload("GARBAGE").build());
+ this.pgcopyConsumer.accept(MessageBuilder.withPayload("125,Bubba,22").build());
+
+ int result = this.jdbcOperations.queryForObject("select count(*) from names", Integer.class);
+ int errors = this.jdbcOperations.queryForObject("select count(*) from test_errors", Integer.class);
+
+ assertThat(result).isEqualTo(2);
+ assertThat(errors).isEqualTo(1);
}
- @SpringBootApplication
+ @SpringBootConfiguration
+ @EnableAutoConfiguration
+ @Import({ PgcopySinkConfiguration.class, TestChannelBinderConfiguration.class })
public static class PgcopySinkApplication {
public static void main(String[] args) {
SpringApplication.run(PgcopySinkApplication.class, args);
diff --git a/applications/sink/pgcopy-sink/src/test/java/org/springframework/cloud/stream/app/pgcopy/sink/PgcopySinkIntegrationTests.java b/applications/sink/pgcopy-sink/src/test/java/org/springframework/cloud/stream/app/pgcopy/sink/PgcopySinkIntegrationTests.java
index f1160447..a5320f67 100644
--- a/applications/sink/pgcopy-sink/src/test/java/org/springframework/cloud/stream/app/pgcopy/sink/PgcopySinkIntegrationTests.java
+++ b/applications/sink/pgcopy-sink/src/test/java/org/springframework/cloud/stream/app/pgcopy/sink/PgcopySinkIntegrationTests.java
@@ -16,39 +16,45 @@
package org.springframework.cloud.stream.app.pgcopy.sink;
-import org.junit.Assert;
-import org.junit.ClassRule;
-import org.junit.Test;
-import org.junit.runner.RunWith;
+import java.util.function.Consumer;
+
+import org.junit.jupiter.api.Test;
+import org.junit.jupiter.api.extension.ExtendWith;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.boot.SpringApplication;
-import org.springframework.boot.autoconfigure.SpringBootApplication;
+import org.springframework.boot.SpringBootConfiguration;
+import org.springframework.boot.autoconfigure.EnableAutoConfiguration;
import org.springframework.boot.test.context.SpringBootTest;
-import org.springframework.cloud.stream.app.pgcopy.test.PostgresTestSupport;
+import org.springframework.cloud.stream.app.pgcopy.test.PostgresAvailableExtension;
+import org.springframework.cloud.stream.binder.test.TestChannelBinderConfiguration;
+import org.springframework.context.annotation.Import;
import org.springframework.jdbc.core.JdbcOperations;
+import org.springframework.messaging.Message;
+import org.springframework.messaging.support.MessageBuilder;
import org.springframework.test.annotation.DirtiesContext;
import org.springframework.test.context.TestPropertySource;
-import org.springframework.test.context.junit4.SpringRunner;
-import static org.hamcrest.Matchers.is;
+import static org.assertj.core.api.Assertions.assertThat;
/**
* Integration Tests for PgcopySink. Only runs if PostgreSQL database is available.
*
* @author Thomas Risberg
+ * @author Chris Bono
*/
-@RunWith(SpringRunner.class)
-@SpringBootTest(webEnvironment = SpringBootTest.WebEnvironment.NONE,
- classes = PgcopySinkIntegrationTests.PgcopySinkApplication.class)
+@SpringBootTest(
+ webEnvironment = SpringBootTest.WebEnvironment.NONE,
+ classes = PgcopyErrorTableIntegrationTests.PgcopySinkApplication.class,
+ properties = {
+ "spring.cloud.function.definition=pgcopyConsumer"
+ })
+@ExtendWith(PostgresAvailableExtension.class)
@DirtiesContext
public abstract class PgcopySinkIntegrationTests {
- @ClassRule
- public static PostgresTestSupport postgresAvailable = new PostgresTestSupport();
-
-// @Autowired
-// protected Sink channels;
+ @Autowired
+ protected Consumer> pgcopyConsumer;
@Autowired
protected JdbcOperations jdbcOperations;
@@ -59,9 +65,9 @@ public abstract class PgcopySinkIntegrationTests {
@Test
public void testBasicCopy() {
String sent = "hello42";
-// channels.input().send(MessageBuilder.withPayload(sent).build());
+ this.pgcopyConsumer.accept(MessageBuilder.withPayload(sent).build());
String result = jdbcOperations.queryForObject("select payload from test", String.class);
- Assert.assertThat(result, is("hello42"));
+ assertThat(result).isEqualTo("hello42");
}
}
@@ -71,12 +77,12 @@ public abstract class PgcopySinkIntegrationTests {
@Test
public void testCopyText() {
-// channels.input().send(MessageBuilder.withPayload("123\tNisse\t25").build());
-// channels.input().send(MessageBuilder.withPayload("124\tAnna\t21").build());
-// channels.input().send(MessageBuilder.withPayload("125\tBubba\t22").build());
-// channels.input().send(MessageBuilder.withPayload("126\tPelle\t32").build());
+ this.pgcopyConsumer.accept(MessageBuilder.withPayload("123\tNisse\t25").build());
+ this.pgcopyConsumer.accept(MessageBuilder.withPayload("124\tAnna\t21").build());
+ this.pgcopyConsumer.accept(MessageBuilder.withPayload("125\tBubba\t22").build());
+ this.pgcopyConsumer.accept(MessageBuilder.withPayload("126\tPelle\t32").build());
int result = jdbcOperations.queryForObject("select count(*) from names", Integer.class);
- Assert.assertThat(result, is(4));
+ assertThat(result).isEqualTo(4);
}
}
@@ -86,11 +92,11 @@ public abstract class PgcopySinkIntegrationTests {
@Test
public void testCopyCSV() {
-// channels.input().send(MessageBuilder.withPayload("123,\"Nisse\",25").build());
-// channels.input().send(MessageBuilder.withPayload("124,\"Anna\",21").build());
-// channels.input().send(MessageBuilder.withPayload("125,\"Bubba\",22").build());
+ this.pgcopyConsumer.accept(MessageBuilder.withPayload("123,\"Nisse\",25").build());
+ this.pgcopyConsumer.accept(MessageBuilder.withPayload("124,\"Anna\",21").build());
+ this.pgcopyConsumer.accept(MessageBuilder.withPayload("125,\"Bubba\",22").build());
int result = jdbcOperations.queryForObject("select count(*) from names", Integer.class);
- Assert.assertThat(result, is(3));
+ assertThat(result).isEqualTo(3);
}
}
@@ -100,13 +106,13 @@ public abstract class PgcopySinkIntegrationTests {
@Test
public void testCopyCSV() {
-// channels.input().send(MessageBuilder.withPayload("123,\"Nisse\",25").build());
-// channels.input().send(MessageBuilder.withPayload("124,,21").build());
-// channels.input().send(MessageBuilder.withPayload("125,\"Bubba\",22").build());
+ this.pgcopyConsumer.accept(MessageBuilder.withPayload("123,\"Nisse\",25").build());
+ this.pgcopyConsumer.accept(MessageBuilder.withPayload("124,,21").build());
+ this.pgcopyConsumer.accept(MessageBuilder.withPayload("125,\"Bubba\",22").build());
int result = jdbcOperations.queryForObject("select count(*) from names", Integer.class);
int nulls = jdbcOperations.queryForObject("select count(*) from names where name is null", Integer.class);
- Assert.assertThat(result, is(3));
- Assert.assertThat(nulls, is(1));
+ assertThat(result).isEqualTo(3);
+ assertThat(nulls).isEqualTo(1);
}
}
@@ -116,13 +122,13 @@ public abstract class PgcopySinkIntegrationTests {
@Test
public void testCopyCSV() {
-// channels.input().send(MessageBuilder.withPayload("123,\"Nisse\",25").build());
-// channels.input().send(MessageBuilder.withPayload("124,null,21").build());
-// channels.input().send(MessageBuilder.withPayload("125,\"Bubba\",22").build());
+ this.pgcopyConsumer.accept(MessageBuilder.withPayload("123,\"Nisse\",25").build());
+ this.pgcopyConsumer.accept(MessageBuilder.withPayload("124,null,21").build());
+ this.pgcopyConsumer.accept(MessageBuilder.withPayload("125,\"Bubba\",22").build());
int result = jdbcOperations.queryForObject("select count(*) from names", Integer.class);
int nulls = jdbcOperations.queryForObject("select count(*) from names where name is null", Integer.class);
- Assert.assertThat(result, is(3));
- Assert.assertThat(nulls, is(1));
+ assertThat(result).isEqualTo(3);
+ assertThat(nulls).isEqualTo(1);
}
}
@@ -132,11 +138,11 @@ public abstract class PgcopySinkIntegrationTests {
@Test
public void testCopyCSV() {
-// channels.input().send(MessageBuilder.withPayload("123|\"Nisse\"|25").build());
-// channels.input().send(MessageBuilder.withPayload("124|\"Anna\"|21").build());
-// channels.input().send(MessageBuilder.withPayload("125|\"Bubba\"|22").build());
+ this.pgcopyConsumer.accept(MessageBuilder.withPayload("123|\"Nisse\"|25").build());
+ this.pgcopyConsumer.accept(MessageBuilder.withPayload("124|\"Anna\"|21").build());
+ this.pgcopyConsumer.accept(MessageBuilder.withPayload("125|\"Bubba\"|22").build());
int result = jdbcOperations.queryForObject("select count(*) from names", Integer.class);
- Assert.assertThat(result, is(3));
+ assertThat(result).isEqualTo(3);
}
}
@@ -146,11 +152,11 @@ public abstract class PgcopySinkIntegrationTests {
@Test
public void testCopyCSV() {
-// channels.input().send(MessageBuilder.withPayload("123\t\"Nisse\"\t25").build());
-// channels.input().send(MessageBuilder.withPayload("124\t\"Anna\"\t21").build());
-// channels.input().send(MessageBuilder.withPayload("125\t\"Bubba\"\t22").build());
+ this.pgcopyConsumer.accept(MessageBuilder.withPayload("123\t\"Nisse\"\t25").build());
+ this.pgcopyConsumer.accept(MessageBuilder.withPayload("124\t\"Anna\"\t21").build());
+ this.pgcopyConsumer.accept(MessageBuilder.withPayload("125\t\"Bubba\"\t22").build());
int result = jdbcOperations.queryForObject("select count(*) from names", Integer.class);
- Assert.assertThat(result, is(3));
+ assertThat(result).isEqualTo(3);
}
}
@@ -160,13 +166,13 @@ public abstract class PgcopySinkIntegrationTests {
@Test
public void testCopyCSV() {
-// channels.input().send(MessageBuilder.withPayload("123,Nisse,25").build());
-// channels.input().send(MessageBuilder.withPayload("124,'Anna',21").build());
-// channels.input().send(MessageBuilder.withPayload("125,Bubba,22").build());
+ this.pgcopyConsumer.accept(MessageBuilder.withPayload("123,Nisse,25").build());
+ this.pgcopyConsumer.accept(MessageBuilder.withPayload("124,'Anna',21").build());
+ this.pgcopyConsumer.accept(MessageBuilder.withPayload("125,Bubba,22").build());
int result = jdbcOperations.queryForObject("select count(*) from names", Integer.class);
int quoted = jdbcOperations.queryForObject("select count(*) from names where name = 'Anna'", Integer.class);
- Assert.assertThat(result, is(3));
- Assert.assertThat(quoted, is(1));
+ assertThat(result).isEqualTo(3);
+ assertThat(quoted).isEqualTo(1);
}
}
@@ -176,17 +182,19 @@ public abstract class PgcopySinkIntegrationTests {
@Test
public void testCopyCSV() {
-// channels.input().send(MessageBuilder.withPayload("123,Nisse,25").build());
-// channels.input().send(MessageBuilder.withPayload("124,\"Anna\\\"\",21").build());
-// channels.input().send(MessageBuilder.withPayload("125,Bubba,22").build());
+ this.pgcopyConsumer.accept(MessageBuilder.withPayload("123,Nisse,25").build());
+ this.pgcopyConsumer.accept(MessageBuilder.withPayload("124,\"Anna\\\"\",21").build());
+ this.pgcopyConsumer.accept(MessageBuilder.withPayload("125,Bubba,22").build());
int result = jdbcOperations.queryForObject("select count(*) from names", Integer.class);
int quoted = jdbcOperations.queryForObject("select count(*) from names where name = 'Anna\"'", Integer.class);
- Assert.assertThat(result, is(3));
- Assert.assertThat(quoted, is(1));
+ assertThat(result).isEqualTo(3);
+ assertThat(quoted).isEqualTo(1);
}
}
- @SpringBootApplication
+ @SpringBootConfiguration
+ @EnableAutoConfiguration
+ @Import({ PgcopySinkConfiguration.class, TestChannelBinderConfiguration.class })
public static class PgcopySinkApplication {
public static void main(String[] args) {
SpringApplication.run(PgcopySinkApplication.class, args);
diff --git a/applications/sink/pgcopy-sink/src/test/java/org/springframework/cloud/stream/app/pgcopy/sink/PgcopySinkPropertiesTests.java b/applications/sink/pgcopy-sink/src/test/java/org/springframework/cloud/stream/app/pgcopy/sink/PgcopySinkPropertiesTests.java
index 8bed14e2..aa1da398 100644
--- a/applications/sink/pgcopy-sink/src/test/java/org/springframework/cloud/stream/app/pgcopy/sink/PgcopySinkPropertiesTests.java
+++ b/applications/sink/pgcopy-sink/src/test/java/org/springframework/cloud/stream/app/pgcopy/sink/PgcopySinkPropertiesTests.java
@@ -16,11 +16,9 @@
package org.springframework.cloud.stream.app.pgcopy.sink;
-import org.junit.After;
-import org.junit.Before;
-import org.junit.Rule;
-import org.junit.Test;
-import org.junit.rules.ExpectedException;
+import org.junit.jupiter.api.AfterEach;
+import org.junit.jupiter.api.BeforeEach;
+import org.junit.jupiter.api.Test;
import org.springframework.beans.factory.BeanCreationException;
import org.springframework.boot.context.properties.EnableConfigurationProperties;
@@ -28,6 +26,7 @@ import org.springframework.boot.test.util.TestPropertyValues;
import org.springframework.context.annotation.AnnotationConfigApplicationContext;
import org.springframework.context.annotation.Configuration;
+import static org.assertj.core.api.Assertions.assertThatThrownBy;
import static org.hamcrest.CoreMatchers.equalTo;
import static org.hamcrest.MatcherAssert.assertThat;
@@ -39,25 +38,23 @@ public class PgcopySinkPropertiesTests {
private AnnotationConfigApplicationContext context;
- @Rule
- public ExpectedException thrown = ExpectedException.none();
-
- @Before
+ @BeforeEach
public void setUp() {
this.context = new AnnotationConfigApplicationContext();
}
- @After
+ @AfterEach
public void tearDown() {
this.context.close();
}
@Test
public void tableNameIsRequired() {
- this.thrown.expect(BeanCreationException.class);
- this.thrown.expectMessage("Failed to bind properties under 'pgcopy' to org.springframework.cloud.stream.app.pgcopy.sink.PgcopySinkProperties");
this.context.register(Conf.class);
- this.context.refresh();
+ assertThatThrownBy(() -> this.context.refresh())
+ .isInstanceOf(BeanCreationException.class)
+ .cause()
+ .hasMessageContaining("Failed to bind properties under 'pgcopy' to org.springframework.cloud.stream.app.pgcopy.sink.PgcopySinkProperties");
}
@Test
diff --git a/applications/sink/pgcopy-sink/src/test/java/org/springframework/cloud/stream/app/pgcopy/test/PostgresAvailableExtension.java b/applications/sink/pgcopy-sink/src/test/java/org/springframework/cloud/stream/app/pgcopy/test/PostgresAvailableExtension.java
new file mode 100644
index 00000000..91d8505a
--- /dev/null
+++ b/applications/sink/pgcopy-sink/src/test/java/org/springframework/cloud/stream/app/pgcopy/test/PostgresAvailableExtension.java
@@ -0,0 +1,64 @@
+/*
+ * Copyright 2022-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.
+ * You may obtain a copy of the License at
+ *
+ * https://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.cloud.stream.app.pgcopy.test;
+
+import java.sql.Connection;
+
+import javax.sql.DataSource;
+
+import org.junit.jupiter.api.extension.ConditionEvaluationResult;
+import org.junit.jupiter.api.extension.ExecutionCondition;
+import org.junit.jupiter.api.extension.ExtensionContext;
+
+import org.springframework.boot.WebApplicationType;
+import org.springframework.boot.autoconfigure.EnableAutoConfiguration;
+import org.springframework.boot.builder.SpringApplicationBuilder;
+import org.springframework.context.ConfigurableApplicationContext;
+import org.springframework.context.annotation.Configuration;
+import org.springframework.core.log.LogAccessor;
+import org.springframework.dao.DataAccessException;
+import org.springframework.jdbc.datasource.DataSourceUtils;
+
+public class PostgresAvailableExtension implements ExecutionCondition {
+
+ private final LogAccessor logger = new LogAccessor(this.getClass());
+
+ @Override
+ public ConditionEvaluationResult evaluateExecutionCondition(ExtensionContext context) {
+
+ try (ConfigurableApplicationContext applicationContext = new SpringApplicationBuilder(PostgresAvailableExtension.Config.class)
+ .properties("spring.integration.jdbc.initialize-schema=never",
+ "spring.integration.jdbc.platform=postgres")
+ .web(WebApplicationType.NONE)
+ .run()) {
+ DataSource dataSource = applicationContext.getBean(DataSource.class);
+ Connection con = DataSourceUtils.getConnection(dataSource);
+ DataSourceUtils.releaseConnection(con, dataSource);
+ }
+ catch (DataAccessException ex) {
+ logger.warn(ex, () -> "Postgres not available - " + ex.getMessage());
+ return ConditionEvaluationResult.disabled("Postgres not available");
+ }
+ return ConditionEvaluationResult.enabled("Postgres available");
+ }
+
+ @Configuration
+ @EnableAutoConfiguration
+ public static class Config {
+
+ }
+}
diff --git a/applications/sink/pgcopy-sink/src/test/java/org/springframework/cloud/stream/app/pgcopy/test/PostgresTestSupport.java b/applications/sink/pgcopy-sink/src/test/java/org/springframework/cloud/stream/app/pgcopy/test/PostgresTestSupport.java
deleted file mode 100644
index f3eab03e..00000000
--- a/applications/sink/pgcopy-sink/src/test/java/org/springframework/cloud/stream/app/pgcopy/test/PostgresTestSupport.java
+++ /dev/null
@@ -1,66 +0,0 @@
-/*
- * 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.
- * You may obtain a copy of the License at
- *
- * https://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.cloud.stream.app.pgcopy.test;
-
-import java.sql.Connection;
-
-import javax.sql.DataSource;
-
-import org.springframework.boot.WebApplicationType;
-import org.springframework.boot.autoconfigure.EnableAutoConfiguration;
-import org.springframework.boot.builder.SpringApplicationBuilder;
-import org.springframework.cloud.stream.test.junit.AbstractExternalResourceTestSupport;
-import org.springframework.context.ConfigurableApplicationContext;
-import org.springframework.context.annotation.Configuration;
-import org.springframework.jdbc.datasource.DataSourceUtils;
-
-/**
- * JUnit {@link org.junit.Rule} that detects the fact that a PostgreSQL server is running on localhost.
- *
- * @author Thomas Risberg
- * @author Artem Bilan
- */
-public class PostgresTestSupport extends AbstractExternalResourceTestSupport {
-
- private ConfigurableApplicationContext context;
-
- public PostgresTestSupport() {
- super("POSTGRES");
- }
-
- @Override
- protected void cleanupResource() {
- context.close();
- }
-
- @Override
- protected void obtainResource() {
- context = new SpringApplicationBuilder(Config.class)
- .web(WebApplicationType.NONE)
- .run();
- DataSource dataSource = context.getBean(DataSource.class);
- Connection con = DataSourceUtils.getConnection(dataSource);
- DataSourceUtils.releaseConnection(con, dataSource);
- }
-
- @Configuration
- @EnableAutoConfiguration
- public static class Config {
-
- }
-
-}
diff --git a/applications/sink/rabbit-sink/src/test/java/org/springframework/cloud/stream/app/sink/rabbit/RabbitSinkInvalidConfigTests.java b/applications/sink/rabbit-sink/src/test/java/org/springframework/cloud/stream/app/sink/rabbit/RabbitSinkInvalidConfigTests.java
index d4ad34c4..cacf29b3 100644
--- a/applications/sink/rabbit-sink/src/test/java/org/springframework/cloud/stream/app/sink/rabbit/RabbitSinkInvalidConfigTests.java
+++ b/applications/sink/rabbit-sink/src/test/java/org/springframework/cloud/stream/app/sink/rabbit/RabbitSinkInvalidConfigTests.java
@@ -30,7 +30,7 @@ import org.springframework.validation.FieldError;
import static org.hamcrest.MatcherAssert.assertThat;
import static org.hamcrest.Matchers.containsString;
import static org.hamcrest.Matchers.instanceOf;
-import static org.junit.Assert.fail;
+import static org.junit.jupiter.api.Assertions.fail;
/**
* Tests for RabbitSource with invalid config.
diff --git a/applications/source/rabbit-source/src/test/java/org/springframework/cloud/stream/app/source/rabbit/RabbitSourceInvalidConfigTests.java b/applications/source/rabbit-source/src/test/java/org/springframework/cloud/stream/app/source/rabbit/RabbitSourceInvalidConfigTests.java
index 23e7035c..e36569e7 100644
--- a/applications/source/rabbit-source/src/test/java/org/springframework/cloud/stream/app/source/rabbit/RabbitSourceInvalidConfigTests.java
+++ b/applications/source/rabbit-source/src/test/java/org/springframework/cloud/stream/app/source/rabbit/RabbitSourceInvalidConfigTests.java
@@ -16,7 +16,7 @@
package org.springframework.cloud.stream.app.source.rabbit;
-import org.junit.Test;
+import org.junit.jupiter.api.Test;
import org.springframework.beans.factory.BeanCreationException;
import org.springframework.boot.context.properties.EnableConfigurationProperties;
@@ -31,7 +31,7 @@ import org.springframework.validation.FieldError;
import static org.hamcrest.MatcherAssert.assertThat;
import static org.hamcrest.Matchers.containsString;
import static org.hamcrest.Matchers.instanceOf;
-import static org.junit.Assert.fail;
+import static org.junit.jupiter.api.Assertions.fail;
/**
* Tests for RabbitSource with invalid config.
diff --git a/functions/consumer/twitter-consumer/src/test/java/org/springframework/cloud/fn/consumer/twitter/status/update/TwitterUpdateSinkFunctionConfigurationTests.java b/functions/consumer/twitter-consumer/src/test/java/org/springframework/cloud/fn/consumer/twitter/status/update/TwitterUpdateSinkFunctionConfigurationTests.java
index a3bdcf3e..b24363c4 100644
--- a/functions/consumer/twitter-consumer/src/test/java/org/springframework/cloud/fn/consumer/twitter/status/update/TwitterUpdateSinkFunctionConfigurationTests.java
+++ b/functions/consumer/twitter-consumer/src/test/java/org/springframework/cloud/fn/consumer/twitter/status/update/TwitterUpdateSinkFunctionConfigurationTests.java
@@ -19,11 +19,12 @@ package org.springframework.cloud.fn.consumer.twitter.status.update;
import java.util.function.Consumer;
import java.util.function.Function;
-import org.junit.Test;
import twitter4j.StatusUpdate;
import twitter4j.Twitter;
import twitter4j.TwitterException;
+import org.junit.jupiter.api.Test;
+
import org.springframework.expression.Expression;
import org.springframework.expression.ExpressionParser;
import org.springframework.expression.spel.standard.SpelExpressionParser;
diff --git a/functions/supplier/twitter-supplier/src/test/java/org/springframework/cloud/fn/supplier/twitter/status/search/SearchPaginationTests.java b/functions/supplier/twitter-supplier/src/test/java/org/springframework/cloud/fn/supplier/twitter/status/search/SearchPaginationTests.java
index 67e48ce9..0c984910 100644
--- a/functions/supplier/twitter-supplier/src/test/java/org/springframework/cloud/fn/supplier/twitter/status/search/SearchPaginationTests.java
+++ b/functions/supplier/twitter-supplier/src/test/java/org/springframework/cloud/fn/supplier/twitter/status/search/SearchPaginationTests.java
@@ -20,7 +20,6 @@ import java.util.ArrayList;
import java.util.Date;
import java.util.List;
-import org.junit.Test;
import twitter4j.GeoLocation;
import twitter4j.HashtagEntity;
import twitter4j.MediaEntity;
@@ -33,6 +32,8 @@ import twitter4j.URLEntity;
import twitter4j.User;
import twitter4j.UserMentionEntity;
+import org.junit.jupiter.api.Test;
+
import static org.assertj.core.api.Assertions.assertThat;
import static org.springframework.cloud.fn.supplier.twitter.status.search.SearchPagination.UNBOUNDED;
diff --git a/functions/supplier/twitter-supplier/src/test/java/org/springframework/cloud/fn/supplier/twitter/status/stream/TwitterStreamSupplierTests.java b/functions/supplier/twitter-supplier/src/test/java/org/springframework/cloud/fn/supplier/twitter/status/stream/TwitterStreamSupplierTests.java
index ff75ac54..f3cf79f1 100644
--- a/functions/supplier/twitter-supplier/src/test/java/org/springframework/cloud/fn/supplier/twitter/status/stream/TwitterStreamSupplierTests.java
+++ b/functions/supplier/twitter-supplier/src/test/java/org/springframework/cloud/fn/supplier/twitter/status/stream/TwitterStreamSupplierTests.java
@@ -20,10 +20,9 @@ import java.util.concurrent.TimeUnit;
import java.util.function.Function;
import java.util.function.Supplier;
-import org.junit.AfterClass;
-import org.junit.BeforeClass;
-import org.junit.Test;
-import org.junit.runner.RunWith;
+import org.junit.jupiter.api.AfterAll;
+import org.junit.jupiter.api.BeforeAll;
+import org.junit.jupiter.api.Test;
import org.mockserver.client.MockServerClient;
import org.mockserver.integration.ClientAndServer;
import org.mockserver.model.Header;
@@ -56,7 +55,6 @@ import static org.mockserver.verify.VerificationTimes.once;
/**
* @author Christian Tzolov
*/
-@RunWith(SpringRunner.class)
@SpringBootTest(
webEnvironment = SpringBootTest.WebEnvironment.NONE,
properties = {
@@ -82,7 +80,7 @@ public abstract class TwitterStreamSupplierTests {
@Autowired
protected Supplier>> twitterStreamSupplier;
- @BeforeClass
+ @BeforeAll
public static void startServer() {
mockServer = ClientAndServer.startClientAndServer(MOCK_SERVER_PORT);
@@ -109,7 +107,7 @@ public abstract class TwitterStreamSupplierTests {
.withBody(new StringBody("count=0&stall_warnings=true")));
}
- @AfterClass
+ @AfterAll
public static void stopServer() {
mockServer.stop();
}
diff --git a/stream-applications-build/pom.xml b/stream-applications-build/pom.xml
index 23161176..7d4ffba9 100644
--- a/stream-applications-build/pom.xml
+++ b/stream-applications-build/pom.xml
@@ -14,7 +14,7 @@
17
3.3.1
3.2.1
- 2.22.2
+ 3.0.0-M7
UTF-8
UTF-8
${java.version}
@@ -44,7 +44,6 @@
4.0.0-SNAPSHOT
4.0.0-SNAPSHOT
1.17.5
- 2.22.2
5.13.2
4.0.5
@@ -150,15 +149,10 @@
maven-surefire-plugin
${maven-surefire-plugin.version}
-
- **/*Tests.java
- **/*Test.java
-
-
- **/Abstract*.java
-
+ true
+
org.apache.maven.plugins
maven-checkstyle-plugin