diff --git a/README.adoc b/README.adoc
index 65b02fe..b43dd7d 100644
--- a/README.adoc
+++ b/README.adoc
@@ -2,10 +2,6 @@
This repository contains extended testing for https://github.com/spring-cloud/stream-applications[stream applications].
-## Integration Tests
-
-Standalone integration tests for testing the various applications.
-
## Acceptance tests on Kubernetes
We test a handful of stream pipelines standalone to ensure that they can be run on Kubernetes.
diff --git a/stream-applications-integration-tests/.mvn/wrapper/MavenWrapperDownloader.java b/stream-applications-integration-tests/.mvn/wrapper/MavenWrapperDownloader.java
deleted file mode 100644
index 0b1fd00..0000000
--- a/stream-applications-integration-tests/.mvn/wrapper/MavenWrapperDownloader.java
+++ /dev/null
@@ -1,118 +0,0 @@
-/*
- * Copyright 2007-present 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.
- */
-
-import java.net.*;
-import java.io.*;
-import java.nio.channels.*;
-import java.util.Properties;
-
-public class MavenWrapperDownloader {
-
- private static final String WRAPPER_VERSION = "0.5.6";
- /**
- * Default URL to download the maven-wrapper.jar from, if no 'downloadUrl' is provided.
- */
- private static final String DEFAULT_DOWNLOAD_URL = "https://repo.maven.apache.org/maven2/io/takari/maven-wrapper/"
- + WRAPPER_VERSION + "/maven-wrapper-" + WRAPPER_VERSION + ".jar";
-
- /**
- * Path to the maven-wrapper.properties file, which might contain a downloadUrl property to
- * use instead of the default one.
- */
- private static final String MAVEN_WRAPPER_PROPERTIES_PATH =
- ".mvn/wrapper/maven-wrapper.properties";
-
- /**
- * Path where the maven-wrapper.jar will be saved to.
- */
- private static final String MAVEN_WRAPPER_JAR_PATH =
- ".mvn/wrapper/maven-wrapper.jar";
-
- /**
- * Name of the property which should be used to override the default download url for the wrapper.
- */
- private static final String PROPERTY_NAME_WRAPPER_URL = "wrapperUrl";
-
- public static void main(String args[]) {
- System.out.println("- Downloader started");
- File baseDirectory = new File(args[0]);
- System.out.println("- Using base directory: " + baseDirectory.getAbsolutePath());
-
- // If the maven-wrapper.properties exists, read it and check if it contains a custom
- // wrapperUrl parameter.
- File mavenWrapperPropertyFile = new File(baseDirectory, MAVEN_WRAPPER_PROPERTIES_PATH);
- String url = DEFAULT_DOWNLOAD_URL;
- if (mavenWrapperPropertyFile.exists()) {
- FileInputStream mavenWrapperPropertyFileInputStream = null;
- try {
- mavenWrapperPropertyFileInputStream = new FileInputStream(mavenWrapperPropertyFile);
- Properties mavenWrapperProperties = new Properties();
- mavenWrapperProperties.load(mavenWrapperPropertyFileInputStream);
- url = mavenWrapperProperties.getProperty(PROPERTY_NAME_WRAPPER_URL, url);
- } catch (IOException e) {
- System.out.println("- ERROR loading '" + MAVEN_WRAPPER_PROPERTIES_PATH + "'");
- } finally {
- try {
- if (mavenWrapperPropertyFileInputStream != null) {
- mavenWrapperPropertyFileInputStream.close();
- }
- } catch (IOException e) {
- // Ignore ...
- }
- }
- }
- System.out.println("- Downloading from: " + url);
-
- File outputFile = new File(baseDirectory.getAbsolutePath(), MAVEN_WRAPPER_JAR_PATH);
- if (!outputFile.getParentFile().exists()) {
- if (!outputFile.getParentFile().mkdirs()) {
- System.out.println(
- "- ERROR creating output directory '" + outputFile.getParentFile().getAbsolutePath() + "'");
- }
- }
- System.out.println("- Downloading to: " + outputFile.getAbsolutePath());
- try {
- downloadFileFromURL(url, outputFile);
- System.out.println("Done");
- System.exit(0);
- } catch (Throwable e) {
- System.out.println("- Error downloading");
- e.printStackTrace();
- System.exit(1);
- }
- }
-
- private static void downloadFileFromURL(String urlString, File destination) throws Exception {
- if (System.getenv("MVNW_USERNAME") != null && System.getenv("MVNW_PASSWORD") != null) {
- String username = System.getenv("MVNW_USERNAME");
- char[] password = System.getenv("MVNW_PASSWORD").toCharArray();
- Authenticator.setDefault(new Authenticator() {
- @Override
- protected PasswordAuthentication getPasswordAuthentication() {
- return new PasswordAuthentication(username, password);
- }
- });
- }
- URL website = new URL(urlString);
- ReadableByteChannel rbc;
- rbc = Channels.newChannel(website.openStream());
- FileOutputStream fos = new FileOutputStream(destination);
- fos.getChannel().transferFrom(rbc, 0, Long.MAX_VALUE);
- fos.close();
- rbc.close();
- }
-
-}
diff --git a/stream-applications-integration-tests/.mvn/wrapper/maven-wrapper.jar b/stream-applications-integration-tests/.mvn/wrapper/maven-wrapper.jar
deleted file mode 100644
index 2cc7d4a..0000000
Binary files a/stream-applications-integration-tests/.mvn/wrapper/maven-wrapper.jar and /dev/null differ
diff --git a/stream-applications-integration-tests/.mvn/wrapper/maven-wrapper.properties b/stream-applications-integration-tests/.mvn/wrapper/maven-wrapper.properties
deleted file mode 100644
index 642d572..0000000
--- a/stream-applications-integration-tests/.mvn/wrapper/maven-wrapper.properties
+++ /dev/null
@@ -1,2 +0,0 @@
-distributionUrl=https://repo.maven.apache.org/maven2/org/apache/maven/apache-maven/3.6.3/apache-maven-3.6.3-bin.zip
-wrapperUrl=https://repo.maven.apache.org/maven2/io/takari/maven-wrapper/0.5.6/maven-wrapper-0.5.6.jar
diff --git a/stream-applications-integration-tests/README.adoc b/stream-applications-integration-tests/README.adoc
deleted file mode 100644
index 361577f..0000000
--- a/stream-applications-integration-tests/README.adoc
+++ /dev/null
@@ -1,42 +0,0 @@
-= Stream Applications Integration Tests
-
-This contains integration tests for pre-packaged stream-applications for Docker using https://www.testcontainers.org/[TestContainers].
-These are end-to-end integration tests running apps and required resources, using docker-compose.
-The goal is to have an end-to-end integration test for each pre-packaged application.
-We don't aim to test all different configuration options, as this is the responsibility of the stream application and function components.
-One of the major benefits is to verify the built Docker images run correctly, especially when we introduce global changes,
-such as upgrading the base JDK image, the maven plugins, or other pervasive changes.
-
-== Test Strategy
-
-See https://github.com/spring-cloud/stream-applications/tree/master/applications/stream-applications-core/common/stream-applications-test-support[] for a full description.
-
-
-The tests use following patterns:
-
-== Source
-To test a source, we may require some application specific setup or event to trigger the source.
-For example, the jdbc source needs some data in the database to which it is listening.
-Then use an `OutputMatcher` to verify the output.
-
-== Sink
-To test a sink, we need to publish a message to its input. Simply use the provided TestTopicSender.
-Then we need to verify the result by checking the sink's external resource.
-
-== Processor
-To test a processor we publish a message and use an `OutputMatcher` to verify the output.
-
-== Configuration
-See link:src/test/java/org/springframework/cloud/stream/apps/integration/test/common/Configuration.java[Configuration] for configuration
-options. These tests use Spring but not boot currently.
-The most important setting is the image versions to test.
-To override it, set the System property, e.g.,
-
-./mvnw clean test -Dspring.cloud.stream.applications.version=3.0.0-M3
-
-
-
-
-
-
-
diff --git a/stream-applications-integration-tests/mvnw b/stream-applications-integration-tests/mvnw
deleted file mode 100755
index a16b543..0000000
--- a/stream-applications-integration-tests/mvnw
+++ /dev/null
@@ -1,310 +0,0 @@
-#!/bin/sh
-# ----------------------------------------------------------------------------
-# Licensed to the Apache Software Foundation (ASF) under one
-# or more contributor license agreements. See the NOTICE file
-# distributed with this work for additional information
-# regarding copyright ownership. The ASF licenses this file
-# to you 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.
-# ----------------------------------------------------------------------------
-
-# ----------------------------------------------------------------------------
-# Maven Start Up Batch script
-#
-# Required ENV vars:
-# ------------------
-# JAVA_HOME - location of a JDK home dir
-#
-# Optional ENV vars
-# -----------------
-# M2_HOME - location of maven2's installed home dir
-# MAVEN_OPTS - parameters passed to the Java VM when running Maven
-# e.g. to debug Maven itself, use
-# set MAVEN_OPTS=-Xdebug -Xrunjdwp:transport=dt_socket,server=y,suspend=y,address=8000
-# MAVEN_SKIP_RC - flag to disable loading of mavenrc files
-# ----------------------------------------------------------------------------
-
-if [ -z "$MAVEN_SKIP_RC" ] ; then
-
- if [ -f /etc/mavenrc ] ; then
- . /etc/mavenrc
- fi
-
- if [ -f "$HOME/.mavenrc" ] ; then
- . "$HOME/.mavenrc"
- fi
-
-fi
-
-# OS specific support. $var _must_ be set to either true or false.
-cygwin=false;
-darwin=false;
-mingw=false
-case "`uname`" in
- CYGWIN*) cygwin=true ;;
- MINGW*) mingw=true;;
- Darwin*) darwin=true
- # Use /usr/libexec/java_home if available, otherwise fall back to /Library/Java/Home
- # See https://developer.apple.com/library/mac/qa/qa1170/_index.html
- if [ -z "$JAVA_HOME" ]; then
- if [ -x "/usr/libexec/java_home" ]; then
- export JAVA_HOME="`/usr/libexec/java_home`"
- else
- export JAVA_HOME="/Library/Java/Home"
- fi
- fi
- ;;
-esac
-
-if [ -z "$JAVA_HOME" ] ; then
- if [ -r /etc/gentoo-release ] ; then
- JAVA_HOME=`java-config --jre-home`
- fi
-fi
-
-if [ -z "$M2_HOME" ] ; then
- ## resolve links - $0 may be a link to maven's home
- PRG="$0"
-
- # need this for relative symlinks
- while [ -h "$PRG" ] ; do
- ls=`ls -ld "$PRG"`
- link=`expr "$ls" : '.*-> \(.*\)$'`
- if expr "$link" : '/.*' > /dev/null; then
- PRG="$link"
- else
- PRG="`dirname "$PRG"`/$link"
- fi
- done
-
- saveddir=`pwd`
-
- M2_HOME=`dirname "$PRG"`/..
-
- # make it fully qualified
- M2_HOME=`cd "$M2_HOME" && pwd`
-
- cd "$saveddir"
- # echo Using m2 at $M2_HOME
-fi
-
-# For Cygwin, ensure paths are in UNIX format before anything is touched
-if $cygwin ; then
- [ -n "$M2_HOME" ] &&
- M2_HOME=`cygpath --unix "$M2_HOME"`
- [ -n "$JAVA_HOME" ] &&
- JAVA_HOME=`cygpath --unix "$JAVA_HOME"`
- [ -n "$CLASSPATH" ] &&
- CLASSPATH=`cygpath --path --unix "$CLASSPATH"`
-fi
-
-# For Mingw, ensure paths are in UNIX format before anything is touched
-if $mingw ; then
- [ -n "$M2_HOME" ] &&
- M2_HOME="`(cd "$M2_HOME"; pwd)`"
- [ -n "$JAVA_HOME" ] &&
- JAVA_HOME="`(cd "$JAVA_HOME"; pwd)`"
-fi
-
-if [ -z "$JAVA_HOME" ]; then
- javaExecutable="`which javac`"
- if [ -n "$javaExecutable" ] && ! [ "`expr \"$javaExecutable\" : '\([^ ]*\)'`" = "no" ]; then
- # readlink(1) is not available as standard on Solaris 10.
- readLink=`which readlink`
- if [ ! `expr "$readLink" : '\([^ ]*\)'` = "no" ]; then
- if $darwin ; then
- javaHome="`dirname \"$javaExecutable\"`"
- javaExecutable="`cd \"$javaHome\" && pwd -P`/javac"
- else
- javaExecutable="`readlink -f \"$javaExecutable\"`"
- fi
- javaHome="`dirname \"$javaExecutable\"`"
- javaHome=`expr "$javaHome" : '\(.*\)/bin'`
- JAVA_HOME="$javaHome"
- export JAVA_HOME
- fi
- fi
-fi
-
-if [ -z "$JAVACMD" ] ; then
- if [ -n "$JAVA_HOME" ] ; then
- if [ -x "$JAVA_HOME/jre/sh/java" ] ; then
- # IBM's JDK on AIX uses strange locations for the executables
- JAVACMD="$JAVA_HOME/jre/sh/java"
- else
- JAVACMD="$JAVA_HOME/bin/java"
- fi
- else
- JAVACMD="`which java`"
- fi
-fi
-
-if [ ! -x "$JAVACMD" ] ; then
- echo "Error: JAVA_HOME is not defined correctly." >&2
- echo " We cannot execute $JAVACMD" >&2
- exit 1
-fi
-
-if [ -z "$JAVA_HOME" ] ; then
- echo "Warning: JAVA_HOME environment variable is not set."
-fi
-
-CLASSWORLDS_LAUNCHER=org.codehaus.plexus.classworlds.launcher.Launcher
-
-# traverses directory structure from process work directory to filesystem root
-# first directory with .mvn subdirectory is considered project base directory
-find_maven_basedir() {
-
- if [ -z "$1" ]
- then
- echo "Path not specified to find_maven_basedir"
- return 1
- fi
-
- basedir="$1"
- wdir="$1"
- while [ "$wdir" != '/' ] ; do
- if [ -d "$wdir"/.mvn ] ; then
- basedir=$wdir
- break
- fi
- # workaround for JBEAP-8937 (on Solaris 10/Sparc)
- if [ -d "${wdir}" ]; then
- wdir=`cd "$wdir/.."; pwd`
- fi
- # end of workaround
- done
- echo "${basedir}"
-}
-
-# concatenates all lines of a file
-concat_lines() {
- if [ -f "$1" ]; then
- echo "$(tr -s '\n' ' ' < "$1")"
- fi
-}
-
-BASE_DIR=`find_maven_basedir "$(pwd)"`
-if [ -z "$BASE_DIR" ]; then
- exit 1;
-fi
-
-##########################################################################################
-# Extension to allow automatically downloading the maven-wrapper.jar from Maven-central
-# This allows using the maven wrapper in projects that prohibit checking in binary data.
-##########################################################################################
-if [ -r "$BASE_DIR/.mvn/wrapper/maven-wrapper.jar" ]; then
- if [ "$MVNW_VERBOSE" = true ]; then
- echo "Found .mvn/wrapper/maven-wrapper.jar"
- fi
-else
- if [ "$MVNW_VERBOSE" = true ]; then
- echo "Couldn't find .mvn/wrapper/maven-wrapper.jar, downloading it ..."
- fi
- if [ -n "$MVNW_REPOURL" ]; then
- jarUrl="$MVNW_REPOURL/io/takari/maven-wrapper/0.5.6/maven-wrapper-0.5.6.jar"
- else
- jarUrl="https://repo.maven.apache.org/maven2/io/takari/maven-wrapper/0.5.6/maven-wrapper-0.5.6.jar"
- fi
- while IFS="=" read key value; do
- case "$key" in (wrapperUrl) jarUrl="$value"; break ;;
- esac
- done < "$BASE_DIR/.mvn/wrapper/maven-wrapper.properties"
- if [ "$MVNW_VERBOSE" = true ]; then
- echo "Downloading from: $jarUrl"
- fi
- wrapperJarPath="$BASE_DIR/.mvn/wrapper/maven-wrapper.jar"
- if $cygwin; then
- wrapperJarPath=`cygpath --path --windows "$wrapperJarPath"`
- fi
-
- if command -v wget > /dev/null; then
- if [ "$MVNW_VERBOSE" = true ]; then
- echo "Found wget ... using wget"
- fi
- if [ -z "$MVNW_USERNAME" ] || [ -z "$MVNW_PASSWORD" ]; then
- wget "$jarUrl" -O "$wrapperJarPath"
- else
- wget --http-user=$MVNW_USERNAME --http-password=$MVNW_PASSWORD "$jarUrl" -O "$wrapperJarPath"
- fi
- elif command -v curl > /dev/null; then
- if [ "$MVNW_VERBOSE" = true ]; then
- echo "Found curl ... using curl"
- fi
- if [ -z "$MVNW_USERNAME" ] || [ -z "$MVNW_PASSWORD" ]; then
- curl -o "$wrapperJarPath" "$jarUrl" -f
- else
- curl --user $MVNW_USERNAME:$MVNW_PASSWORD -o "$wrapperJarPath" "$jarUrl" -f
- fi
-
- else
- if [ "$MVNW_VERBOSE" = true ]; then
- echo "Falling back to using Java to download"
- fi
- javaClass="$BASE_DIR/.mvn/wrapper/MavenWrapperDownloader.java"
- # For Cygwin, switch paths to Windows format before running javac
- if $cygwin; then
- javaClass=`cygpath --path --windows "$javaClass"`
- fi
- if [ -e "$javaClass" ]; then
- if [ ! -e "$BASE_DIR/.mvn/wrapper/MavenWrapperDownloader.class" ]; then
- if [ "$MVNW_VERBOSE" = true ]; then
- echo " - Compiling MavenWrapperDownloader.java ..."
- fi
- # Compiling the Java class
- ("$JAVA_HOME/bin/javac" "$javaClass")
- fi
- if [ -e "$BASE_DIR/.mvn/wrapper/MavenWrapperDownloader.class" ]; then
- # Running the downloader
- if [ "$MVNW_VERBOSE" = true ]; then
- echo " - Running MavenWrapperDownloader.java ..."
- fi
- ("$JAVA_HOME/bin/java" -cp .mvn/wrapper MavenWrapperDownloader "$MAVEN_PROJECTBASEDIR")
- fi
- fi
- fi
-fi
-##########################################################################################
-# End of extension
-##########################################################################################
-
-export MAVEN_PROJECTBASEDIR=${MAVEN_BASEDIR:-"$BASE_DIR"}
-if [ "$MVNW_VERBOSE" = true ]; then
- echo $MAVEN_PROJECTBASEDIR
-fi
-MAVEN_OPTS="$(concat_lines "$MAVEN_PROJECTBASEDIR/.mvn/jvm.config") $MAVEN_OPTS"
-
-# For Cygwin, switch paths to Windows format before running java
-if $cygwin; then
- [ -n "$M2_HOME" ] &&
- M2_HOME=`cygpath --path --windows "$M2_HOME"`
- [ -n "$JAVA_HOME" ] &&
- JAVA_HOME=`cygpath --path --windows "$JAVA_HOME"`
- [ -n "$CLASSPATH" ] &&
- CLASSPATH=`cygpath --path --windows "$CLASSPATH"`
- [ -n "$MAVEN_PROJECTBASEDIR" ] &&
- MAVEN_PROJECTBASEDIR=`cygpath --path --windows "$MAVEN_PROJECTBASEDIR"`
-fi
-
-# Provide a "standardized" way to retrieve the CLI args that will
-# work with both Windows and non-Windows executions.
-MAVEN_CMD_LINE_ARGS="$MAVEN_CONFIG $@"
-export MAVEN_CMD_LINE_ARGS
-
-WRAPPER_LAUNCHER=org.apache.maven.wrapper.MavenWrapperMain
-
-exec "$JAVACMD" \
- $MAVEN_OPTS \
- -classpath "$MAVEN_PROJECTBASEDIR/.mvn/wrapper/maven-wrapper.jar" \
- "-Dmaven.home=${M2_HOME}" "-Dmaven.multiModuleProjectDirectory=${MAVEN_PROJECTBASEDIR}" \
- ${WRAPPER_LAUNCHER} $MAVEN_CONFIG "$@"
diff --git a/stream-applications-integration-tests/mvnw.cmd b/stream-applications-integration-tests/mvnw.cmd
deleted file mode 100644
index c8d4337..0000000
--- a/stream-applications-integration-tests/mvnw.cmd
+++ /dev/null
@@ -1,182 +0,0 @@
-@REM ----------------------------------------------------------------------------
-@REM Licensed to the Apache Software Foundation (ASF) under one
-@REM or more contributor license agreements. See the NOTICE file
-@REM distributed with this work for additional information
-@REM regarding copyright ownership. The ASF licenses this file
-@REM to you under the Apache License, Version 2.0 (the
-@REM "License"); you may not use this file except in compliance
-@REM with the License. You may obtain a copy of the License at
-@REM
-@REM https://www.apache.org/licenses/LICENSE-2.0
-@REM
-@REM Unless required by applicable law or agreed to in writing,
-@REM software distributed under the License is distributed on an
-@REM "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
-@REM KIND, either express or implied. See the License for the
-@REM specific language governing permissions and limitations
-@REM under the License.
-@REM ----------------------------------------------------------------------------
-
-@REM ----------------------------------------------------------------------------
-@REM Maven Start Up Batch script
-@REM
-@REM Required ENV vars:
-@REM JAVA_HOME - location of a JDK home dir
-@REM
-@REM Optional ENV vars
-@REM M2_HOME - location of maven2's installed home dir
-@REM MAVEN_BATCH_ECHO - set to 'on' to enable the echoing of the batch commands
-@REM MAVEN_BATCH_PAUSE - set to 'on' to wait for a keystroke before ending
-@REM MAVEN_OPTS - parameters passed to the Java VM when running Maven
-@REM e.g. to debug Maven itself, use
-@REM set MAVEN_OPTS=-Xdebug -Xrunjdwp:transport=dt_socket,server=y,suspend=y,address=8000
-@REM MAVEN_SKIP_RC - flag to disable loading of mavenrc files
-@REM ----------------------------------------------------------------------------
-
-@REM Begin all REM lines with '@' in case MAVEN_BATCH_ECHO is 'on'
-@echo off
-@REM set title of command window
-title %0
-@REM enable echoing by setting MAVEN_BATCH_ECHO to 'on'
-@if "%MAVEN_BATCH_ECHO%" == "on" echo %MAVEN_BATCH_ECHO%
-
-@REM set %HOME% to equivalent of $HOME
-if "%HOME%" == "" (set "HOME=%HOMEDRIVE%%HOMEPATH%")
-
-@REM Execute a user defined script before this one
-if not "%MAVEN_SKIP_RC%" == "" goto skipRcPre
-@REM check for pre script, once with legacy .bat ending and once with .cmd ending
-if exist "%HOME%\mavenrc_pre.bat" call "%HOME%\mavenrc_pre.bat"
-if exist "%HOME%\mavenrc_pre.cmd" call "%HOME%\mavenrc_pre.cmd"
-:skipRcPre
-
-@setlocal
-
-set ERROR_CODE=0
-
-@REM To isolate internal variables from possible post scripts, we use another setlocal
-@setlocal
-
-@REM ==== START VALIDATION ====
-if not "%JAVA_HOME%" == "" goto OkJHome
-
-echo.
-echo Error: JAVA_HOME not found in your environment. >&2
-echo Please set the JAVA_HOME variable in your environment to match the >&2
-echo location of your Java installation. >&2
-echo.
-goto error
-
-:OkJHome
-if exist "%JAVA_HOME%\bin\java.exe" goto init
-
-echo.
-echo Error: JAVA_HOME is set to an invalid directory. >&2
-echo JAVA_HOME = "%JAVA_HOME%" >&2
-echo Please set the JAVA_HOME variable in your environment to match the >&2
-echo location of your Java installation. >&2
-echo.
-goto error
-
-@REM ==== END VALIDATION ====
-
-:init
-
-@REM Find the project base dir, i.e. the directory that contains the folder ".mvn".
-@REM Fallback to current working directory if not found.
-
-set MAVEN_PROJECTBASEDIR=%MAVEN_BASEDIR%
-IF NOT "%MAVEN_PROJECTBASEDIR%"=="" goto endDetectBaseDir
-
-set EXEC_DIR=%CD%
-set WDIR=%EXEC_DIR%
-:findBaseDir
-IF EXIST "%WDIR%"\.mvn goto baseDirFound
-cd ..
-IF "%WDIR%"=="%CD%" goto baseDirNotFound
-set WDIR=%CD%
-goto findBaseDir
-
-:baseDirFound
-set MAVEN_PROJECTBASEDIR=%WDIR%
-cd "%EXEC_DIR%"
-goto endDetectBaseDir
-
-:baseDirNotFound
-set MAVEN_PROJECTBASEDIR=%EXEC_DIR%
-cd "%EXEC_DIR%"
-
-:endDetectBaseDir
-
-IF NOT EXIST "%MAVEN_PROJECTBASEDIR%\.mvn\jvm.config" goto endReadAdditionalConfig
-
-@setlocal EnableExtensions EnableDelayedExpansion
-for /F "usebackq delims=" %%a in ("%MAVEN_PROJECTBASEDIR%\.mvn\jvm.config") do set JVM_CONFIG_MAVEN_PROPS=!JVM_CONFIG_MAVEN_PROPS! %%a
-@endlocal & set JVM_CONFIG_MAVEN_PROPS=%JVM_CONFIG_MAVEN_PROPS%
-
-:endReadAdditionalConfig
-
-SET MAVEN_JAVA_EXE="%JAVA_HOME%\bin\java.exe"
-set WRAPPER_JAR="%MAVEN_PROJECTBASEDIR%\.mvn\wrapper\maven-wrapper.jar"
-set WRAPPER_LAUNCHER=org.apache.maven.wrapper.MavenWrapperMain
-
-set DOWNLOAD_URL="https://repo.maven.apache.org/maven2/io/takari/maven-wrapper/0.5.6/maven-wrapper-0.5.6.jar"
-
-FOR /F "tokens=1,2 delims==" %%A IN ("%MAVEN_PROJECTBASEDIR%\.mvn\wrapper\maven-wrapper.properties") DO (
- IF "%%A"=="wrapperUrl" SET DOWNLOAD_URL=%%B
-)
-
-@REM Extension to allow automatically downloading the maven-wrapper.jar from Maven-central
-@REM This allows using the maven wrapper in projects that prohibit checking in binary data.
-if exist %WRAPPER_JAR% (
- if "%MVNW_VERBOSE%" == "true" (
- echo Found %WRAPPER_JAR%
- )
-) else (
- if not "%MVNW_REPOURL%" == "" (
- SET DOWNLOAD_URL="%MVNW_REPOURL%/io/takari/maven-wrapper/0.5.6/maven-wrapper-0.5.6.jar"
- )
- if "%MVNW_VERBOSE%" == "true" (
- echo Couldn't find %WRAPPER_JAR%, downloading it ...
- echo Downloading from: %DOWNLOAD_URL%
- )
-
- powershell -Command "&{"^
- "$webclient = new-object System.Net.WebClient;"^
- "if (-not ([string]::IsNullOrEmpty('%MVNW_USERNAME%') -and [string]::IsNullOrEmpty('%MVNW_PASSWORD%'))) {"^
- "$webclient.Credentials = new-object System.Net.NetworkCredential('%MVNW_USERNAME%', '%MVNW_PASSWORD%');"^
- "}"^
- "[Net.ServicePointManager]::SecurityProtocol = [Net.SecurityProtocolType]::Tls12; $webclient.DownloadFile('%DOWNLOAD_URL%', '%WRAPPER_JAR%')"^
- "}"
- if "%MVNW_VERBOSE%" == "true" (
- echo Finished downloading %WRAPPER_JAR%
- )
-)
-@REM End of extension
-
-@REM Provide a "standardized" way to retrieve the CLI args that will
-@REM work with both Windows and non-Windows executions.
-set MAVEN_CMD_LINE_ARGS=%*
-
-%MAVEN_JAVA_EXE% %JVM_CONFIG_MAVEN_PROPS% %MAVEN_OPTS% %MAVEN_DEBUG_OPTS% -classpath %WRAPPER_JAR% "-Dmaven.multiModuleProjectDirectory=%MAVEN_PROJECTBASEDIR%" %WRAPPER_LAUNCHER% %MAVEN_CONFIG% %*
-if ERRORLEVEL 1 goto error
-goto end
-
-:error
-set ERROR_CODE=1
-
-:end
-@endlocal & set ERROR_CODE=%ERROR_CODE%
-
-if not "%MAVEN_SKIP_RC%" == "" goto skipRcPost
-@REM check for post script, once with legacy .bat ending and once with .cmd ending
-if exist "%HOME%\mavenrc_post.bat" call "%HOME%\mavenrc_post.bat"
-if exist "%HOME%\mavenrc_post.cmd" call "%HOME%\mavenrc_post.cmd"
-:skipRcPost
-
-@REM pause the script if MAVEN_BATCH_PAUSE is set to 'on'
-if "%MAVEN_BATCH_PAUSE%" == "on" pause
-
-if "%MAVEN_TERMINATE_CMD%" == "on" exit %ERROR_CODE%
-
-exit /B %ERROR_CODE%
diff --git a/stream-applications-integration-tests/pom.xml b/stream-applications-integration-tests/pom.xml
deleted file mode 100644
index fd4b158..0000000
--- a/stream-applications-integration-tests/pom.xml
+++ /dev/null
@@ -1,287 +0,0 @@
-
-
- 4.0.0
-
- org.springframework.boot
- spring-boot-starter-parent
- 2.3.4.RELEASE
-
-
- org.springframework.cloud.stream.apps
- stream-applications-integration-tests
- 1.0.0-SNAPSHOT
- stream-applications-integration-tests
- Integration Tests for stream applications
-
-
- 1.8
- 1.0.3-SNAPSHOT
- 3.0.2-SNAPSHOT
- 2.6.2
- 1.15.1
- 3.1.0
- false
- true
- true
- true
- 8.29
- https://raw.githubusercontent.com/spring-cloud/stream-applications/master/etc/checkstyle
-
-
- ${checkstyle.location}/checkstyle-suppressions.xml
-
-
- ${checkstyle.location}/nohttp-checkstyle.xml
-
- ${checkstyle.location}/checkstyle-suppressions.xml
-
- 0.0.2.RELEASE
- true
- 0.0.7
- 2.27.1
- 1.11.415
- 8.0.16
- 5.6.2
-
-
-
-
- org.springframework.boot
- spring-boot-starter-webflux
- test
-
-
- org.springframework.boot
- spring-boot-starter-test
- test
-
-
- org.junit.vintage
- junit-vintage-engine
-
-
-
-
- org.testcontainers
- testcontainers
- test
-
-
-
- org.testcontainers
- junit-jupiter
- test
-
-
-
- org.junit.jupiter
- junit-jupiter-api
- ${junit-jupiter-api.version}
- test
-
-
- org.testcontainers
- kafka
- test
-
-
- org.testcontainers
- rabbitmq
-
-
- org.testcontainers
- mongodb
- test
-
-
- org.testcontainers
- mysql
- test
-
-
- mysql
- mysql-connector-java
- ${mysql-connector-java.version}
- test
-
-
- org.springframework.cloud.stream.app
- stream-applications-test-support
- ${stream-applications.version}
- test
-
-
- org.springframework.cloud.fn
- function-test-support
- ${java-functions.version}
- test
-
-
- org.springframework.data
- spring-data-geode
-
-
- org.apache.logging.log4j
- log4j
-
-
- test
-
-
- com.amazonaws
- aws-java-sdk-s3
- ${aws.version}
- test
-
-
- com.squareup.okhttp3
- mockwebserver
- test
-
-
- org.springframework.boot
- spring-boot-starter-data-mongodb
- test
-
-
- org.mariadb.jdbc
- mariadb-java-client
- ${mariadb-client.version}
- test
-
-
- org.springframework.boot
- spring-boot-starter-jdbc
- test
-
-
-
-
-
- org.apache.maven.plugins
- maven-surefire-plugin
-
- 1C
- all
-
-
-
- org.apache.maven.plugins
- maven-checkstyle-plugin
- ${maven-checkstyle-plugin.version}
-
-
- com.puppycrawl.tools
- checkstyle
- ${puppycrawl-tools-checkstyle.version}
-
-
- io.spring.javaformat
- spring-javaformat-checkstyle
- ${spring-javaformat-checkstyle.version}
-
-
- io.spring.nohttp
- nohttp-checkstyle
- ${nohttp-checkstyle.version}
-
-
-
-
- checkstyle-validation
- validate
- true
-
- ${disable.checks}
- ${checkstyle.location}/checkstyle.xml
- ${checkstyle.location}/checkstyle-header.txt
-
- checkstyle.build.directory=${project.build.directory}
- checkstyle.suppressions.file=${checkstyle.suppressions.file}
- checkstyle.additional.suppressions.file=${checkstyle.additional.suppressions.file}
-
- true
-
- ${maven-checkstyle-plugin.includeTestSourceDirectory}
-
- ${maven-checkstyle-plugin.failsOnError}
-
-
- ${maven-checkstyle-plugin.failOnViolation}
-
-
-
- check
-
-
-
- no-http-checkstyle-validation
- validate
- true
-
- ${disable.nohttp.checks}
- ${checkstyle.nohttp.file}
- **/*
- **/.idea/**/*,**/.git/**/*,**/target/**/*,**/*.log
- ./
-
-
- check
-
-
-
-
-
-
-
-
-
-
- org.testcontainers
- testcontainers-bom
- ${test-containers.version}
- pom
- import
-
-
-
-
-
-
-
- true
-
- spring-snapshots
- Spring Snapshots
- https://repo.spring.io/snapshot
-
-
-
- false
-
- spring-milestone-release
- Spring Milestone Release
- https://repo.spring.io/milestone
-
-
-
-
-
- true
-
- spring-snapshots
- Spring Snapshots
- https://repo.spring.io/snapshot
-
-
-
- false
-
- spring-milestones
- Spring Milestones
- https://repo.spring.io/milestone
-
-
-
-
diff --git a/stream-applications-integration-tests/src/test/java/org/springframework/cloud/stream/apps/integration/test/common/Configuration.java b/stream-applications-integration-tests/src/test/java/org/springframework/cloud/stream/apps/integration/test/common/Configuration.java
deleted file mode 100644
index 16b4d44..0000000
--- a/stream-applications-integration-tests/src/test/java/org/springframework/cloud/stream/apps/integration/test/common/Configuration.java
+++ /dev/null
@@ -1,42 +0,0 @@
-/*
- * Copyright 2020-2020 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.apps.integration.test.common;
-
-import java.time.Duration;
-import java.util.function.Supplier;
-
-public abstract class Configuration {
-
- public static String VERSION;
-
- public static final Duration DEFAULT_DURATION = Duration.ofMinutes(1);
-
- private static final String SPRING_CLOUD_STREAM_APPLICATIONS_VERSION = "spring.cloud.stream.applications.version";
-
- static {
- VERSION = System.getProperty(SPRING_CLOUD_STREAM_APPLICATIONS_VERSION, "latest");
- }
-
- public static class VersionSupplier implements Supplier {
-
- @Override
- public String get() {
- return VERSION;
- }
- }
-
-}
diff --git a/stream-applications-integration-tests/src/test/java/org/springframework/cloud/stream/apps/integration/test/processor/httprequest/HttpRequestProcessorTests.java b/stream-applications-integration-tests/src/test/java/org/springframework/cloud/stream/apps/integration/test/processor/httprequest/HttpRequestProcessorTests.java
deleted file mode 100644
index cca1a4a..0000000
--- a/stream-applications-integration-tests/src/test/java/org/springframework/cloud/stream/apps/integration/test/processor/httprequest/HttpRequestProcessorTests.java
+++ /dev/null
@@ -1,96 +0,0 @@
-/*
- * Copyright 2020-2020 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.apps.integration.test.processor.httprequest;
-
-import java.io.IOException;
-import java.net.InetAddress;
-
-import okhttp3.mockwebserver.Dispatcher;
-import okhttp3.mockwebserver.MockResponse;
-import okhttp3.mockwebserver.MockWebServer;
-import okhttp3.mockwebserver.RecordedRequest;
-import org.junit.jupiter.api.AfterAll;
-import org.junit.jupiter.api.BeforeAll;
-import org.junit.jupiter.api.Test;
-
-import org.springframework.beans.factory.annotation.Autowired;
-import org.springframework.cloud.stream.app.test.integration.OutputMatcher;
-import org.springframework.cloud.stream.app.test.integration.StreamAppContainer;
-import org.springframework.cloud.stream.app.test.integration.StreamAppContainerTestUtils;
-import org.springframework.cloud.stream.app.test.integration.TestTopicSender;
-import org.springframework.http.HttpHeaders;
-import org.springframework.http.HttpStatus;
-import org.springframework.http.MediaType;
-import org.springframework.messaging.MessageHeaders;
-
-import static org.awaitility.Awaitility.await;
-import static org.springframework.cloud.stream.app.test.integration.AppLog.appLog;
-import static org.springframework.cloud.stream.apps.integration.test.common.Configuration.DEFAULT_DURATION;
-
-abstract class HttpRequestProcessorTests {
-
- private static MockWebServer server;
-
- private static int serverPort;
-
- @Autowired
- private TestTopicSender testTopicSender;
-
- @Autowired
- private OutputMatcher outputMatcher;
-
- private static StreamAppContainer processor;
-
- protected static StreamAppContainer configureProcessor(StreamAppContainer baseContainer) {
- serverPort = StreamAppContainerTestUtils.findAvailablePort();
- processor = baseContainer.withLogConsumer(appLog("http-request-processor"))
- .withEnv("HTTP_REQUEST_URL_EXPRESSION",
- "'http://" + StreamAppContainerTestUtils.localHostAddress() + ":" + serverPort + "'")
- .withEnv("HTTP_REQUEST_HTTP_METHOD_EXPRESSION", "'POST'");
- return processor;
- }
-
- @BeforeAll
- static void startServer() throws Exception {
- server = new MockWebServer();
- server.start(InetAddress.getLocalHost(), serverPort);
- }
-
- @Test
- void get() {
- server.setDispatcher(new Dispatcher() {
- @Override
- public MockResponse dispatch(RecordedRequest recordedRequest) {
- return new MockResponse().setHeader(HttpHeaders.CONTENT_TYPE,
- MediaType.APPLICATION_JSON_VALUE)
- .setBody("{\"response\":\"" + recordedRequest.getBody().readUtf8() + "\"}")
- .setResponseCode(HttpStatus.OK.value());
- }
- });
- testTopicSender.send(processor.getInputDestination(), "ping");
- await().atMost(DEFAULT_DURATION)
- .until(outputMatcher.messageMatches(message -> message.getPayload().equals("{\"response\":\"ping\"}")
- && message.getHeaders().get(MessageHeaders.CONTENT_TYPE)
- .equals(MediaType.APPLICATION_JSON_VALUE)));
- }
-
- @AfterAll
- static void cleanUp() throws IOException {
- server.shutdown();
- }
-
-}
diff --git a/stream-applications-integration-tests/src/test/java/org/springframework/cloud/stream/apps/integration/test/processor/httprequest/KafkaHttpRequestProcessorTests.java b/stream-applications-integration-tests/src/test/java/org/springframework/cloud/stream/apps/integration/test/processor/httprequest/KafkaHttpRequestProcessorTests.java
deleted file mode 100644
index 3ad5e1e..0000000
--- a/stream-applications-integration-tests/src/test/java/org/springframework/cloud/stream/apps/integration/test/processor/httprequest/KafkaHttpRequestProcessorTests.java
+++ /dev/null
@@ -1,36 +0,0 @@
-/*
- * Copyright 2020-2020 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.apps.integration.test.processor.httprequest;
-
-import org.testcontainers.junit.jupiter.Container;
-
-import org.springframework.cloud.stream.app.test.integration.StreamAppContainer;
-import org.springframework.cloud.stream.app.test.integration.StreamAppContainerTestUtils;
-import org.springframework.cloud.stream.app.test.integration.junit.jupiter.KafkaStreamAppTest;
-import org.springframework.cloud.stream.app.test.integration.kafka.KafkaStreamAppContainer;
-
-
-import static org.springframework.cloud.stream.apps.integration.test.common.Configuration.VERSION;
-
-@KafkaStreamAppTest
-class KafkaHttpRequestProcessorTests extends HttpRequestProcessorTests {
- @Container
- private static StreamAppContainer container = configureProcessor(
- new KafkaStreamAppContainer(StreamAppContainerTestUtils.imageName(
- "http-request-processor-kafka", VERSION)));
-
-}
diff --git a/stream-applications-integration-tests/src/test/java/org/springframework/cloud/stream/apps/integration/test/processor/httprequest/RabbitMQHttpRequestProcessorTests.java b/stream-applications-integration-tests/src/test/java/org/springframework/cloud/stream/apps/integration/test/processor/httprequest/RabbitMQHttpRequestProcessorTests.java
deleted file mode 100644
index 6781ad0..0000000
--- a/stream-applications-integration-tests/src/test/java/org/springframework/cloud/stream/apps/integration/test/processor/httprequest/RabbitMQHttpRequestProcessorTests.java
+++ /dev/null
@@ -1,35 +0,0 @@
-/*
- * Copyright 2020-2020 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.apps.integration.test.processor.httprequest;
-
-import org.testcontainers.junit.jupiter.Container;
-
-import org.springframework.cloud.stream.app.test.integration.StreamAppContainer;
-import org.springframework.cloud.stream.app.test.integration.StreamAppContainerTestUtils;
-import org.springframework.cloud.stream.app.test.integration.junit.jupiter.RabbitMQStreamAppTest;
-import org.springframework.cloud.stream.app.test.integration.rabbitmq.RabbitMQStreamAppContainer;
-
-import static org.springframework.cloud.stream.apps.integration.test.common.Configuration.VERSION;
-
-@RabbitMQStreamAppTest
-class RabbitMQHttpRequestProcessorTests extends HttpRequestProcessorTests {
-
- @Container
- private static StreamAppContainer container = configureProcessor(
- new RabbitMQStreamAppContainer(StreamAppContainerTestUtils.imageName(
- "http-request-processor-rabbit", VERSION)));
-}
diff --git a/stream-applications-integration-tests/src/test/java/org/springframework/cloud/stream/apps/integration/test/sink/jdbc/JdbcSinkTests.java b/stream-applications-integration-tests/src/test/java/org/springframework/cloud/stream/apps/integration/test/sink/jdbc/JdbcSinkTests.java
deleted file mode 100644
index 8615b34..0000000
--- a/stream-applications-integration-tests/src/test/java/org/springframework/cloud/stream/apps/integration/test/sink/jdbc/JdbcSinkTests.java
+++ /dev/null
@@ -1,113 +0,0 @@
-/*
- * Copyright 2020-2020 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.apps.integration.test.sink.jdbc;
-
-import com.zaxxer.hikari.HikariDataSource;
-import org.junit.jupiter.api.AfterAll;
-import org.junit.jupiter.api.BeforeAll;
-import org.junit.jupiter.api.Test;
-import org.junit.jupiter.api.extension.ExtendWith;
-import org.testcontainers.containers.BindMode;
-import org.testcontainers.containers.MySQLContainer;
-import org.testcontainers.containers.wait.strategy.Wait;
-import org.testcontainers.junit.jupiter.Container;
-import org.testcontainers.utility.DockerImageName;
-
-import org.springframework.beans.factory.annotation.Autowired;
-import org.springframework.cloud.stream.app.test.integration.StreamAppContainer;
-import org.springframework.cloud.stream.app.test.integration.TestTopicSender;
-import org.springframework.cloud.stream.app.test.integration.junit.jupiter.BaseContainerExtension;
-import org.springframework.cloud.stream.app.test.integration.kafka.KafkaConfig;
-import org.springframework.jdbc.core.JdbcTemplate;
-
-import static org.assertj.core.api.Assertions.assertThat;
-import static org.awaitility.Awaitility.await;
-import static org.springframework.cloud.stream.app.test.integration.AppLog.appLog;
-import static org.springframework.cloud.stream.apps.integration.test.common.Configuration.DEFAULT_DURATION;
-
-@ExtendWith(BaseContainerExtension.class)
-public abstract class JdbcSinkTests {
-
- private static JdbcTemplate jdbcTemplate;
-
- private static StreamAppContainer sink;
-
- @Autowired
- private TestTopicSender testTopicSender;
-
- @Container
- private static MySQLContainer mySQL = new MySQLContainer<>(DockerImageName.parse("mysql:5.7"))
- .withUsername("test")
- .withPassword("secret")
- .withExposedPorts(3306)
- .withNetwork(KafkaConfig.kafka.getNetwork())
- .withNetworkAliases("mysql-for-sink")
- .withClasspathResourceMapping("init.sql", "/init.sql", BindMode.READ_ONLY)
- .withLogConsumer(appLog("mysql-for-sink"))
- .withCommand("--init-file", "/init.sql");
-
- @BeforeAll
- static void init() {
- sink = BaseContainerExtension.containerInstance()
- .dependsOn(mySQL)
- .withEnv("JDBC_CONSUMER_COLUMNS", "name,city:address.city,street:address.street")
- .withEnv("JDBC_CONSUMER_TABLE_NAME", "People")
- .withEnv("SPRING_DATASOURCE_USERNAME", "test")
- .withEnv("SPRING_DATASOURCE_PASSWORD", "secret")
- .withEnv("SPRING_DATASOURCE_DRIVER_CLASS_NAME", "org.mariadb.jdbc.Driver")
- .withEnv("SPRING_DATASOURCE_URL",
- "jdbc:mariadb://mysql-for-sink:3306/test")
- .waitingFor(Wait.forLogMessage(".*Started JdbcSink.*", 1));
- startSink();
- }
-
- static void startSink() {
-
- HikariDataSource dataSource = new HikariDataSource();
- dataSource.setDriverClassName("org.mariadb.jdbc.Driver");
- dataSource.setUsername(mySQL.getUsername());
- dataSource.setPassword(mySQL.getPassword());
- dataSource.setJdbcUrl("jdbc:mysql://localhost:" + mySQL.getMappedPort(3306) + "/test");
- jdbcTemplate = new JdbcTemplate(dataSource);
- jdbcTemplate.execute("DELETE FROM People");
- await().atMost(DEFAULT_DURATION)
- .until(() -> jdbcTemplate.queryForObject("SELECT COUNT(*) from People", Integer.class)
- .intValue() == 0);
- sink.start();
- }
-
- @Test
- void test() {
-
- String json = "{\"name\":\"My Name\",\"address\":{ \"city\": \"Big City\",\"street\":\"Narrow Alley\"}}";
- testTopicSender.send(sink.getInputDestination(), json);
-
- await().atMost(DEFAULT_DURATION)
- .untilAsserted(
- () -> assertThat(
- jdbcTemplate.queryForObject("SELECT COUNT(*) from People", Integer.class).intValue())
- .isOne());
- assertThat(jdbcTemplate.queryForObject("SELECT name from People",
- String.class)).isEqualTo("My Name");
- }
-
- @AfterAll
- static void cleanUp() {
- sink.stop();
- }
-
-}
diff --git a/stream-applications-integration-tests/src/test/java/org/springframework/cloud/stream/apps/integration/test/sink/jdbc/KafkaJdbcSinkTests.java b/stream-applications-integration-tests/src/test/java/org/springframework/cloud/stream/apps/integration/test/sink/jdbc/KafkaJdbcSinkTests.java
deleted file mode 100644
index a8d1174..0000000
--- a/stream-applications-integration-tests/src/test/java/org/springframework/cloud/stream/apps/integration/test/sink/jdbc/KafkaJdbcSinkTests.java
+++ /dev/null
@@ -1,26 +0,0 @@
-/*
- * Copyright 2020-2020 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.apps.integration.test.sink.jdbc;
-
-import org.springframework.cloud.stream.app.test.integration.junit.jupiter.KafkaBaseContainer;
-import org.springframework.cloud.stream.app.test.integration.junit.jupiter.KafkaStreamAppTest;
-import org.springframework.cloud.stream.apps.integration.test.common.Configuration;
-
-@KafkaStreamAppTest
-@KafkaBaseContainer(name = "jdbc-sink-kafka", versionSupplier = Configuration.VersionSupplier.class)
-public class KafkaJdbcSinkTests extends JdbcSinkTests {
-}
diff --git a/stream-applications-integration-tests/src/test/java/org/springframework/cloud/stream/apps/integration/test/sink/jdbc/RabbitMQJdbcSinkTests.java b/stream-applications-integration-tests/src/test/java/org/springframework/cloud/stream/apps/integration/test/sink/jdbc/RabbitMQJdbcSinkTests.java
deleted file mode 100644
index aca720b..0000000
--- a/stream-applications-integration-tests/src/test/java/org/springframework/cloud/stream/apps/integration/test/sink/jdbc/RabbitMQJdbcSinkTests.java
+++ /dev/null
@@ -1,26 +0,0 @@
-/*
- * Copyright 2020-2020 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.apps.integration.test.sink.jdbc;
-
-import org.springframework.cloud.stream.app.test.integration.junit.jupiter.RabbitMQBaseContainer;
-import org.springframework.cloud.stream.app.test.integration.junit.jupiter.RabbitMQStreamAppTest;
-import org.springframework.cloud.stream.apps.integration.test.common.Configuration;
-
-@RabbitMQStreamAppTest
-@RabbitMQBaseContainer(name = "jdbc-sink-rabbit", versionSupplier = Configuration.VersionSupplier.class)
-public class RabbitMQJdbcSinkTests extends JdbcSinkTests {
-}
diff --git a/stream-applications-integration-tests/src/test/java/org/springframework/cloud/stream/apps/integration/test/sink/mongodb/KafkaMongoDBSinkTests.java b/stream-applications-integration-tests/src/test/java/org/springframework/cloud/stream/apps/integration/test/sink/mongodb/KafkaMongoDBSinkTests.java
deleted file mode 100644
index 89b3f73..0000000
--- a/stream-applications-integration-tests/src/test/java/org/springframework/cloud/stream/apps/integration/test/sink/mongodb/KafkaMongoDBSinkTests.java
+++ /dev/null
@@ -1,26 +0,0 @@
-/*
- * Copyright 2020-2020 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.apps.integration.test.sink.mongodb;
-
-import org.springframework.cloud.stream.app.test.integration.junit.jupiter.KafkaBaseContainer;
-import org.springframework.cloud.stream.app.test.integration.junit.jupiter.KafkaStreamAppTest;
-import org.springframework.cloud.stream.apps.integration.test.common.Configuration;
-
-@KafkaStreamAppTest
-@KafkaBaseContainer(name = "mongodb-sink-kafka", versionSupplier = Configuration.VersionSupplier.class)
-public class KafkaMongoDBSinkTests extends MongoDBSinkTests {
-}
diff --git a/stream-applications-integration-tests/src/test/java/org/springframework/cloud/stream/apps/integration/test/sink/mongodb/MongoDBSinkTests.java b/stream-applications-integration-tests/src/test/java/org/springframework/cloud/stream/apps/integration/test/sink/mongodb/MongoDBSinkTests.java
deleted file mode 100644
index c9c2661..0000000
--- a/stream-applications-integration-tests/src/test/java/org/springframework/cloud/stream/apps/integration/test/sink/mongodb/MongoDBSinkTests.java
+++ /dev/null
@@ -1,98 +0,0 @@
-/*
- * Copyright 2020-2020 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.apps.integration.test.sink.mongodb;
-
-import java.time.Duration;
-import java.util.List;
-
-import org.bson.Document;
-import org.junit.jupiter.api.AfterAll;
-import org.junit.jupiter.api.BeforeAll;
-import org.junit.jupiter.api.Test;
-import org.junit.jupiter.api.extension.ExtendWith;
-import org.testcontainers.containers.MongoDBContainer;
-import org.testcontainers.containers.wait.strategy.Wait;
-
-import org.springframework.beans.factory.annotation.Autowired;
-import org.springframework.cloud.stream.app.test.integration.StreamAppContainer;
-import org.springframework.cloud.stream.app.test.integration.StreamAppContainerTestUtils;
-import org.springframework.cloud.stream.app.test.integration.TestTopicSender;
-import org.springframework.cloud.stream.app.test.integration.junit.jupiter.BaseContainerExtension;
-import org.springframework.data.mongodb.MongoDatabaseFactory;
-import org.springframework.data.mongodb.core.MongoTemplate;
-import org.springframework.data.mongodb.core.SimpleMongoClientDatabaseFactory;
-
-import static org.assertj.core.api.Assertions.assertThat;
-import static org.awaitility.Awaitility.await;
-import static org.springframework.cloud.stream.apps.integration.test.common.Configuration.DEFAULT_DURATION;
-
-@ExtendWith(BaseContainerExtension.class)
-abstract class MongoDBSinkTests {
-
- private static MongoTemplate mongoTemplate;
-
- @Autowired
- private TestTopicSender testTopicSender;
-
- private static final MongoDBContainer mongoDBContainer = new MongoDBContainer("mongo:4.0.10")
- .withExposedPorts(27017)
- .withStartupTimeout(Duration.ofMinutes(2));
-
- private static String mongoConnectionString() {
- return String.format("mongodb://%s:%s/%s", StreamAppContainerTestUtils.localHostAddress(),
- mongoDBContainer.getMappedPort(27017), "test");
- }
-
- private static StreamAppContainer sink;
-
- @BeforeAll
- protected static void configureSink() {
- mongoDBContainer.start();
- sink = BaseContainerExtension.containerInstance()
- .withEnv("MONGODB_CONSUMER_COLLECTION", "test")
- .withEnv("SPRING_DATA_MONGODB_URL", mongoConnectionString())
- .waitingFor(Wait.forLogMessage(".*Started MongodbSink.*", 1));
-
- sink.start();
- buildMongoTemplate();
- }
-
- static void buildMongoTemplate() {
- mongoDBContainer.start();
- MongoDatabaseFactory mongoDatabaseFactory = new SimpleMongoClientDatabaseFactory(
- mongoConnectionString());
- mongoTemplate = new MongoTemplate(mongoDatabaseFactory);
- }
-
- @Test
- void postData() {
- String json = "{\"name\":\"My Name\",\"address\":{ \"city\": \"Big City\", \"street\":\"Narrow Alley\"}}";
- testTopicSender.send(sink.getInputDestination(), json);
-
- await().atMost(DEFAULT_DURATION).untilAsserted(() -> {
- List docs = mongoTemplate.findAll(Document.class, "test");
- assertThat(docs).allMatch(document -> document.get("name", String.class).equals("My Name"));
- });
- }
-
- @AfterAll
- static void cleanUp() {
- mongoDBContainer.close();
- sink.stop();
- }
-
-}
diff --git a/stream-applications-integration-tests/src/test/java/org/springframework/cloud/stream/apps/integration/test/sink/mongodb/RabbitMQMongoDBSinkTests.java b/stream-applications-integration-tests/src/test/java/org/springframework/cloud/stream/apps/integration/test/sink/mongodb/RabbitMQMongoDBSinkTests.java
deleted file mode 100644
index f6fab69..0000000
--- a/stream-applications-integration-tests/src/test/java/org/springframework/cloud/stream/apps/integration/test/sink/mongodb/RabbitMQMongoDBSinkTests.java
+++ /dev/null
@@ -1,26 +0,0 @@
-/*
- * Copyright 2020-2020 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.apps.integration.test.sink.mongodb;
-
-import org.springframework.cloud.stream.app.test.integration.junit.jupiter.RabbitMQBaseContainer;
-import org.springframework.cloud.stream.app.test.integration.junit.jupiter.RabbitMQStreamAppTest;
-import org.springframework.cloud.stream.apps.integration.test.common.Configuration;
-
-@RabbitMQStreamAppTest
-@RabbitMQBaseContainer(name = "mongodb-sink-rabbit", versionSupplier = Configuration.VersionSupplier.class)
-public class RabbitMQMongoDBSinkTests extends MongoDBSinkTests {
-}
diff --git a/stream-applications-integration-tests/src/test/java/org/springframework/cloud/stream/apps/integration/test/sink/tcp/KafkaTcpSinkTests.java b/stream-applications-integration-tests/src/test/java/org/springframework/cloud/stream/apps/integration/test/sink/tcp/KafkaTcpSinkTests.java
deleted file mode 100644
index c5cac73..0000000
--- a/stream-applications-integration-tests/src/test/java/org/springframework/cloud/stream/apps/integration/test/sink/tcp/KafkaTcpSinkTests.java
+++ /dev/null
@@ -1,26 +0,0 @@
-/*
- * Copyright 2020-2020 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.apps.integration.test.sink.tcp;
-
-import org.springframework.cloud.stream.app.test.integration.junit.jupiter.KafkaBaseContainer;
-import org.springframework.cloud.stream.app.test.integration.junit.jupiter.KafkaStreamAppTest;
-import org.springframework.cloud.stream.apps.integration.test.common.Configuration;
-
-@KafkaStreamAppTest
-@KafkaBaseContainer(name = "tcp-sink-kafka", versionSupplier = Configuration.VersionSupplier.class)
-public class KafkaTcpSinkTests extends TcpSinkTests {
-}
diff --git a/stream-applications-integration-tests/src/test/java/org/springframework/cloud/stream/apps/integration/test/sink/tcp/RabbitMQTcpSinkTests.java b/stream-applications-integration-tests/src/test/java/org/springframework/cloud/stream/apps/integration/test/sink/tcp/RabbitMQTcpSinkTests.java
deleted file mode 100644
index e1b731c..0000000
--- a/stream-applications-integration-tests/src/test/java/org/springframework/cloud/stream/apps/integration/test/sink/tcp/RabbitMQTcpSinkTests.java
+++ /dev/null
@@ -1,26 +0,0 @@
-/*
- * Copyright 2020-2020 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.apps.integration.test.sink.tcp;
-
-import org.springframework.cloud.stream.app.test.integration.junit.jupiter.RabbitMQBaseContainer;
-import org.springframework.cloud.stream.app.test.integration.junit.jupiter.RabbitMQStreamAppTest;
-import org.springframework.cloud.stream.apps.integration.test.common.Configuration;
-
-@RabbitMQStreamAppTest
-@RabbitMQBaseContainer(name = "tcp-sink-rabbit", versionSupplier = Configuration.VersionSupplier.class)
-public class RabbitMQTcpSinkTests extends TcpSinkTests {
-}
diff --git a/stream-applications-integration-tests/src/test/java/org/springframework/cloud/stream/apps/integration/test/sink/tcp/TcpSinkTests.java b/stream-applications-integration-tests/src/test/java/org/springframework/cloud/stream/apps/integration/test/sink/tcp/TcpSinkTests.java
deleted file mode 100644
index 84fa884..0000000
--- a/stream-applications-integration-tests/src/test/java/org/springframework/cloud/stream/apps/integration/test/sink/tcp/TcpSinkTests.java
+++ /dev/null
@@ -1,98 +0,0 @@
-/*
- * Copyright 2020-2020 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.apps.integration.test.sink.tcp;
-
-import java.io.BufferedReader;
-import java.io.IOException;
-import java.io.InputStreamReader;
-import java.net.InetAddress;
-import java.net.ServerSocket;
-import java.net.Socket;
-import java.time.Duration;
-import java.util.concurrent.atomic.AtomicBoolean;
-
-import org.junit.jupiter.api.AfterAll;
-import org.junit.jupiter.api.BeforeAll;
-import org.junit.jupiter.api.Test;
-import org.junit.jupiter.api.extension.ExtendWith;
-import org.testcontainers.containers.wait.strategy.Wait;
-
-import org.springframework.beans.factory.annotation.Autowired;
-import org.springframework.cloud.stream.app.test.integration.StreamAppContainer;
-import org.springframework.cloud.stream.app.test.integration.StreamAppContainerTestUtils;
-import org.springframework.cloud.stream.app.test.integration.TestTopicSender;
-import org.springframework.cloud.stream.app.test.integration.junit.jupiter.BaseContainerExtension;
-
-import static org.awaitility.Awaitility.await;
-import static org.springframework.cloud.stream.apps.integration.test.common.Configuration.DEFAULT_DURATION;
-
-@ExtendWith(BaseContainerExtension.class)
-abstract class TcpSinkTests {
-
- private static int tcpPort;
-
- private static Socket socket;
-
- private static final AtomicBoolean socketReady = new AtomicBoolean();
-
- private static StreamAppContainer sink;
-
- @Autowired
- private TestTopicSender testTopicSender;
-
- @BeforeAll
- static void configureSink() {
- tcpPort = StreamAppContainerTestUtils.findAvailablePort();
- startTcpServer();
- sink = BaseContainerExtension.containerInstance()
- .withEnv("TCP_CONSUMER_HOST", StreamAppContainerTestUtils.localHostAddress())
- .withEnv("TCP_PORT", String.valueOf(tcpPort))
- .withEnv("TCP_CONSUMER_ENCODER", "CRLF")
- .waitingFor(Wait.forLogMessage(".*Started TcpSink.*", 1));
- sink.start();
- }
-
- static void startTcpServer() {
- socketReady.set(false);
- new Thread(() -> {
- try {
- socket = new ServerSocket(tcpPort, 50, InetAddress.getLocalHost()).accept();
- socketReady.set(true);
- }
- catch (IOException e) {
- throw new RuntimeException("failed to bind to port " + tcpPort + ": " + e.getMessage(), e);
- }
- }).start();
- }
-
- @Test
- void postData() throws IOException {
- // Sink will not connect until it receives a message.
- String text = "Hello, world!";
- testTopicSender.send(sink.getInputDestination(), text);
-
- await().atMost(DEFAULT_DURATION).untilTrue(socketReady);
- BufferedReader reader = new BufferedReader(new InputStreamReader(socket.getInputStream()));
- await().atMost(Duration.ofSeconds(10)).until(() -> reader.readLine().equals(text));
- }
-
- @AfterAll
- static void cleanUp() throws IOException {
- sink.stop();
- socket.close();
- }
-}
diff --git a/stream-applications-integration-tests/src/test/java/org/springframework/cloud/stream/apps/integration/test/source/geode/GeodeSourceTests.java b/stream-applications-integration-tests/src/test/java/org/springframework/cloud/stream/apps/integration/test/source/geode/GeodeSourceTests.java
deleted file mode 100644
index 7cdd7c8..0000000
--- a/stream-applications-integration-tests/src/test/java/org/springframework/cloud/stream/apps/integration/test/source/geode/GeodeSourceTests.java
+++ /dev/null
@@ -1,118 +0,0 @@
-/*
- * Copyright 2020-2020 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.apps.integration.test.source.geode;
-
-import java.time.Duration;
-import java.util.UUID;
-import java.util.function.Consumer;
-
-import com.github.dockerjava.api.command.CreateContainerCmd;
-import com.github.dockerjava.api.model.ExposedPort;
-import com.github.dockerjava.api.model.HostConfig;
-import com.github.dockerjava.api.model.PortBinding;
-import com.github.dockerjava.api.model.Ports;
-import org.apache.geode.cache.Region;
-import org.apache.geode.cache.client.ClientCache;
-import org.apache.geode.cache.client.ClientCacheFactory;
-import org.apache.geode.cache.client.ClientRegionShortcut;
-import org.junit.jupiter.api.AfterAll;
-import org.junit.jupiter.api.BeforeAll;
-import org.junit.jupiter.api.Test;
-import org.junit.jupiter.api.extension.ExtendWith;
-import org.testcontainers.containers.wait.strategy.Wait;
-import org.testcontainers.images.builder.ImageFromDockerfile;
-import org.testcontainers.junit.jupiter.Container;
-
-import org.springframework.beans.factory.annotation.Autowired;
-import org.springframework.cloud.fn.test.support.geode.GeodeContainer;
-import org.springframework.cloud.stream.app.test.integration.OutputMatcher;
-import org.springframework.cloud.stream.app.test.integration.StreamAppContainer;
-import org.springframework.cloud.stream.app.test.integration.StreamAppContainerTestUtils;
-import org.springframework.cloud.stream.app.test.integration.junit.jupiter.BaseContainerExtension;
-
-import static org.awaitility.Awaitility.await;
-import static org.springframework.cloud.stream.apps.integration.test.common.Configuration.DEFAULT_DURATION;
-
-@ExtendWith(BaseContainerExtension.class)
-abstract class GeodeSourceTests {
- private static int locatorPort = StreamAppContainerTestUtils.findAvailablePort();
-
- private static int cacheServerPort = StreamAppContainerTestUtils.findAvailablePort();
-
- private static Region