diff --git a/spring-integration-etcd/README.md b/spring-integration-etcd/README.md new file mode 100644 index 0000000..b323cd0 --- /dev/null +++ b/spring-integration-etcd/README.md @@ -0,0 +1,23 @@ +SPRING INTEGRATION ETCD SUPPORT +==================================== + +## ETCD LEADER ELECTION + +If you need to elect a leader (e.g. for highly available message consumer where only one node should receive messages) you just need to create a `LeaderInitiator`. +Example: + +```java +@Bean +public EtcdClient etcdClient() { + return EtcdClient.forEndpoint("localhost",2379).withPlainText().build(); +} + +@Bean +public LeaderInitiator initiator() { + LeaderInitiator initiator = new LeaderInitiator(etcdClient()); + return initiator; +} +``` + +Then when a node is elected leader it will send `OnGrantedEvent` to all application listeners. +See the [Spring Integration User Guide](https://docs.spring.io/spring-integration/reference/html/messaging-endpoints-chapter.html#endpoint-roles) for more information on how to use those events to control messaging endpoints. diff --git a/spring-integration-etcd/build.gradle b/spring-integration-etcd/build.gradle new file mode 100644 index 0000000..81bf982 --- /dev/null +++ b/spring-integration-etcd/build.gradle @@ -0,0 +1,169 @@ +plugins { + id 'java' + id 'eclipse' + id 'idea' + id 'jacoco' + id 'org.sonarqube' version '2.6.2' +} + +apply from: "${rootProject.projectDir}/publish-maven.gradle" + +description = 'Spring Integration Etcd Support' +group = 'org.springframework.integration' + +repositories { + if (version.endsWith('BUILD-SNAPSHOT')) { + maven { url 'http://repo.spring.io/libs-snapshot' } + } + maven { url 'http://repo.spring.io/libs-milestone' } +} + +sourceCompatibility = targetCompatibility = 1.8 + +ext { + etcdVersion = '0.0.7' + log4jVersion = '2.11.1' + springIntegrationVersion = '5.0.8.RELEASE' + + + idPrefix = 'etcd' + + linkHomepage = 'https://github.com/spring-projects/spring-integration-extensions' + linkCi = 'https://build.spring.io/browse/INTEXT' + linkIssue = 'https://jira.spring.io/browse/INTEXT' + linkScmUrl = 'https://github.com/spring-projects/spring-integration-extensions' + linkScmConnection = 'https://github.com/spring-projects/spring-integration-extensions.git' + linkScmDevConnection = 'git@github.com:spring-projects/spring-integration-extensions.git' +} + +eclipse.project.natures += 'org.springframework.ide.eclipse.core.springnature' + +sourceSets { + test { + resources { + srcDirs = ['src/test/resources', 'src/test/java'] + } + } +} + +jacoco { + toolVersion = "0.8.2" +} + +dependencies { + compile "org.springframework.integration:spring-integration-core:$springIntegrationVersion" + compile "com.ibm.etcd:etcd-java:$etcdVersion" + + testCompile "org.springframework.integration:spring-integration-test:$springIntegrationVersion" + + testRuntime "org.apache.logging.log4j:log4j-slf4j-impl:$log4jVersion" +} + +// enable all compiler warnings; individual projects may customize further +[compileJava, compileTestJava]*.options*.compilerArgs = ['-Xlint:all,-options'] + +test { + // suppress all console output during testing unless running `gradle -i` + logging.captureStandardOutput(LogLevel.INFO) + + maxHeapSize = "1024m" + jacoco { + append = false + destinationFile = file("$buildDir/jacoco.exec") + } +} + +jacocoTestReport { + reports { + xml.enabled false + csv.enabled false + html.destination file("${buildDir}/reports/jacoco/html") + } +} + +task sourcesJar(type: Jar) { + classifier = 'sources' + from sourceSets.main.allJava +} + +task javadocJar(type: Jar) { + classifier = 'javadoc' + from javadoc +} + +artifacts { + archives sourcesJar + archives javadocJar +} + +sonarqube { + properties { + property "sonar.jacoco.reportPath", "${buildDir.name}/jacoco.exec" + property "sonar.links.homepage", linkHomepage + property "sonar.links.ci", linkCi + property "sonar.links.issue", linkIssue + property "sonar.links.scm", linkScmUrl + property "sonar.links.scm_dev", linkScmDevConnection + property "sonar.java.coveragePlugin", "jacoco" + } +} + +task distZip(type: Zip) { + group = 'Distribution' + classifier = 'dist' + description = "Builds -${classifier} archive, containing all jars and docs, " + + "suitable for community download page." + + ext.baseDir = "${project.name}-${project.version}"; + + from('src/dist') { + include 'readme.txt' + include 'license.txt' + include 'notice.txt' + into "${baseDir}" + } + + into("${baseDir}/libs") { + from project.jar + from project.sourcesJar + from project.javadocJar + } +} + +// Create an optional "with dependencies" distribution. +// Not published by default; only for use when building from source. +task depsZip(type: Zip, dependsOn: distZip) { zipTask -> + group = 'Distribution' + classifier = 'dist-with-deps' + description = "Builds -${classifier} archive, containing everything " + + "in the -${distZip.classifier} archive plus all dependencies." + + from zipTree(distZip.archivePath) + + gradle.taskGraph.whenReady { taskGraph -> + if (taskGraph.hasTask(":${zipTask.name}")) { + def projectName = rootProject.name + def artifacts = new HashSet() + + rootProject.configurations.runtime.resolvedConfiguration.resolvedArtifacts.each { artifact -> + def dependency = artifact.moduleVersion.id + if (!projectName.equals(dependency.name)) { + artifacts << artifact.file + } + } + + zipTask.from(artifacts) { + into "${distZip.baseDir}/deps" + } + } + } +} + +artifacts { + archives distZip +} + +task dist(dependsOn: assemble) { + group = 'Distribution' + description = 'Builds -dist distribution archives.' +} diff --git a/spring-integration-etcd/docker-compose.yml b/spring-integration-etcd/docker-compose.yml new file mode 100644 index 0000000..7bbe351 --- /dev/null +++ b/spring-integration-etcd/docker-compose.yml @@ -0,0 +1,8 @@ +version: '3' +services: + etcd: + image: quay.io/coreos/etcd:v3.3 + ports: + - "2379:2379" + - "2380:2380" + command: /usr/local/bin/etcd --name spring-integration-etcd-test --advertise-client-urls http://0.0.0.0:2379 --listen-client-urls http://0.0.0.0:2379 --debug \ No newline at end of file diff --git a/spring-integration-etcd/gradle.properties b/spring-integration-etcd/gradle.properties new file mode 100644 index 0000000..d45ca0e --- /dev/null +++ b/spring-integration-etcd/gradle.properties @@ -0,0 +1 @@ +version=1.0.0-BUILD-SNAPSHOT diff --git a/spring-integration-etcd/gradle/wrapper/gradle-wrapper.jar b/spring-integration-etcd/gradle/wrapper/gradle-wrapper.jar new file mode 100644 index 0000000..29953ea Binary files /dev/null and b/spring-integration-etcd/gradle/wrapper/gradle-wrapper.jar differ diff --git a/spring-integration-etcd/gradle/wrapper/gradle-wrapper.properties b/spring-integration-etcd/gradle/wrapper/gradle-wrapper.properties new file mode 100644 index 0000000..e0b3fb8 --- /dev/null +++ b/spring-integration-etcd/gradle/wrapper/gradle-wrapper.properties @@ -0,0 +1,5 @@ +distributionBase=GRADLE_USER_HOME +distributionPath=wrapper/dists +distributionUrl=https\://services.gradle.org/distributions/gradle-4.10.2-bin.zip +zipStoreBase=GRADLE_USER_HOME +zipStorePath=wrapper/dists diff --git a/spring-integration-etcd/gradlew b/spring-integration-etcd/gradlew new file mode 100755 index 0000000..cccdd3d --- /dev/null +++ b/spring-integration-etcd/gradlew @@ -0,0 +1,172 @@ +#!/usr/bin/env sh + +############################################################################## +## +## Gradle start up script for UN*X +## +############################################################################## + +# Attempt to set APP_HOME +# Resolve links: $0 may be a link +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 +SAVED="`pwd`" +cd "`dirname \"$PRG\"`/" >/dev/null +APP_HOME="`pwd -P`" +cd "$SAVED" >/dev/null + +APP_NAME="Gradle" +APP_BASE_NAME=`basename "$0"` + +# Add default JVM options here. You can also use JAVA_OPTS and GRADLE_OPTS to pass JVM options to this script. +DEFAULT_JVM_OPTS="" + +# Use the maximum available, or set MAX_FD != -1 to use that value. +MAX_FD="maximum" + +warn () { + echo "$*" +} + +die () { + echo + echo "$*" + echo + exit 1 +} + +# OS specific support (must be 'true' or 'false'). +cygwin=false +msys=false +darwin=false +nonstop=false +case "`uname`" in + CYGWIN* ) + cygwin=true + ;; + Darwin* ) + darwin=true + ;; + MINGW* ) + msys=true + ;; + NONSTOP* ) + nonstop=true + ;; +esac + +CLASSPATH=$APP_HOME/gradle/wrapper/gradle-wrapper.jar + +# Determine the Java command to use to start the JVM. +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 + if [ ! -x "$JAVACMD" ] ; then + die "ERROR: JAVA_HOME is set to an invalid directory: $JAVA_HOME + +Please set the JAVA_HOME variable in your environment to match the +location of your Java installation." + fi +else + JAVACMD="java" + which java >/dev/null 2>&1 || die "ERROR: JAVA_HOME is not set and no 'java' command could be found in your PATH. + +Please set the JAVA_HOME variable in your environment to match the +location of your Java installation." +fi + +# Increase the maximum file descriptors if we can. +if [ "$cygwin" = "false" -a "$darwin" = "false" -a "$nonstop" = "false" ] ; then + MAX_FD_LIMIT=`ulimit -H -n` + if [ $? -eq 0 ] ; then + if [ "$MAX_FD" = "maximum" -o "$MAX_FD" = "max" ] ; then + MAX_FD="$MAX_FD_LIMIT" + fi + ulimit -n $MAX_FD + if [ $? -ne 0 ] ; then + warn "Could not set maximum file descriptor limit: $MAX_FD" + fi + else + warn "Could not query maximum file descriptor limit: $MAX_FD_LIMIT" + fi +fi + +# For Darwin, add options to specify how the application appears in the dock +if $darwin; then + GRADLE_OPTS="$GRADLE_OPTS \"-Xdock:name=$APP_NAME\" \"-Xdock:icon=$APP_HOME/media/gradle.icns\"" +fi + +# For Cygwin, switch paths to Windows format before running java +if $cygwin ; then + APP_HOME=`cygpath --path --mixed "$APP_HOME"` + CLASSPATH=`cygpath --path --mixed "$CLASSPATH"` + JAVACMD=`cygpath --unix "$JAVACMD"` + + # We build the pattern for arguments to be converted via cygpath + ROOTDIRSRAW=`find -L / -maxdepth 1 -mindepth 1 -type d 2>/dev/null` + SEP="" + for dir in $ROOTDIRSRAW ; do + ROOTDIRS="$ROOTDIRS$SEP$dir" + SEP="|" + done + OURCYGPATTERN="(^($ROOTDIRS))" + # Add a user-defined pattern to the cygpath arguments + if [ "$GRADLE_CYGPATTERN" != "" ] ; then + OURCYGPATTERN="$OURCYGPATTERN|($GRADLE_CYGPATTERN)" + fi + # Now convert the arguments - kludge to limit ourselves to /bin/sh + i=0 + for arg in "$@" ; do + CHECK=`echo "$arg"|egrep -c "$OURCYGPATTERN" -` + CHECK2=`echo "$arg"|egrep -c "^-"` ### Determine if an option + + if [ $CHECK -ne 0 ] && [ $CHECK2 -eq 0 ] ; then ### Added a condition + eval `echo args$i`=`cygpath --path --ignore --mixed "$arg"` + else + eval `echo args$i`="\"$arg\"" + fi + i=$((i+1)) + done + case $i in + (0) set -- ;; + (1) set -- "$args0" ;; + (2) set -- "$args0" "$args1" ;; + (3) set -- "$args0" "$args1" "$args2" ;; + (4) set -- "$args0" "$args1" "$args2" "$args3" ;; + (5) set -- "$args0" "$args1" "$args2" "$args3" "$args4" ;; + (6) set -- "$args0" "$args1" "$args2" "$args3" "$args4" "$args5" ;; + (7) set -- "$args0" "$args1" "$args2" "$args3" "$args4" "$args5" "$args6" ;; + (8) set -- "$args0" "$args1" "$args2" "$args3" "$args4" "$args5" "$args6" "$args7" ;; + (9) set -- "$args0" "$args1" "$args2" "$args3" "$args4" "$args5" "$args6" "$args7" "$args8" ;; + esac +fi + +# Escape application args +save () { + for i do printf %s\\n "$i" | sed "s/'/'\\\\''/g;1s/^/'/;\$s/\$/' \\\\/" ; done + echo " " +} +APP_ARGS=$(save "$@") + +# Collect all arguments for the java command, following the shell quoting and substitution rules +eval set -- $DEFAULT_JVM_OPTS $JAVA_OPTS $GRADLE_OPTS "\"-Dorg.gradle.appname=$APP_BASE_NAME\"" -classpath "\"$CLASSPATH\"" org.gradle.wrapper.GradleWrapperMain "$APP_ARGS" + +# by default we should be in the correct project dir, but when run from Finder on Mac, the cwd is wrong +if [ "$(uname)" = "Darwin" ] && [ "$HOME" = "$PWD" ]; then + cd "$(dirname "$0")" +fi + +exec "$JAVACMD" "$@" diff --git a/spring-integration-etcd/gradlew.bat b/spring-integration-etcd/gradlew.bat new file mode 100644 index 0000000..e95643d --- /dev/null +++ b/spring-integration-etcd/gradlew.bat @@ -0,0 +1,84 @@ +@if "%DEBUG%" == "" @echo off +@rem ########################################################################## +@rem +@rem Gradle startup script for Windows +@rem +@rem ########################################################################## + +@rem Set local scope for the variables with windows NT shell +if "%OS%"=="Windows_NT" setlocal + +set DIRNAME=%~dp0 +if "%DIRNAME%" == "" set DIRNAME=. +set APP_BASE_NAME=%~n0 +set APP_HOME=%DIRNAME% + +@rem Add default JVM options here. You can also use JAVA_OPTS and GRADLE_OPTS to pass JVM options to this script. +set DEFAULT_JVM_OPTS= + +@rem Find java.exe +if defined JAVA_HOME goto findJavaFromJavaHome + +set JAVA_EXE=java.exe +%JAVA_EXE% -version >NUL 2>&1 +if "%ERRORLEVEL%" == "0" goto init + +echo. +echo ERROR: JAVA_HOME is not set and no 'java' command could be found in your PATH. +echo. +echo Please set the JAVA_HOME variable in your environment to match the +echo location of your Java installation. + +goto fail + +:findJavaFromJavaHome +set JAVA_HOME=%JAVA_HOME:"=% +set JAVA_EXE=%JAVA_HOME%/bin/java.exe + +if exist "%JAVA_EXE%" goto init + +echo. +echo ERROR: JAVA_HOME is set to an invalid directory: %JAVA_HOME% +echo. +echo Please set the JAVA_HOME variable in your environment to match the +echo location of your Java installation. + +goto fail + +:init +@rem Get command-line arguments, handling Windows variants + +if not "%OS%" == "Windows_NT" goto win9xME_args + +:win9xME_args +@rem Slurp the command line arguments. +set CMD_LINE_ARGS= +set _SKIP=2 + +:win9xME_args_slurp +if "x%~1" == "x" goto execute + +set CMD_LINE_ARGS=%* + +:execute +@rem Setup the command line + +set CLASSPATH=%APP_HOME%\gradle\wrapper\gradle-wrapper.jar + +@rem Execute Gradle +"%JAVA_EXE%" %DEFAULT_JVM_OPTS% %JAVA_OPTS% %GRADLE_OPTS% "-Dorg.gradle.appname=%APP_BASE_NAME%" -classpath "%CLASSPATH%" org.gradle.wrapper.GradleWrapperMain %CMD_LINE_ARGS% + +:end +@rem End local scope for the variables with windows NT shell +if "%ERRORLEVEL%"=="0" goto mainEnd + +:fail +rem Set variable GRADLE_EXIT_CONSOLE if you need the _script_ return code instead of +rem the _cmd.exe /c_ return code! +if not "" == "%GRADLE_EXIT_CONSOLE%" exit 1 +exit /b 1 + +:mainEnd +if "%OS%"=="Windows_NT" endlocal + +:omega diff --git a/spring-integration-etcd/publish-maven.gradle b/spring-integration-etcd/publish-maven.gradle new file mode 100644 index 0000000..9fa8bfb --- /dev/null +++ b/spring-integration-etcd/publish-maven.gradle @@ -0,0 +1,62 @@ +apply plugin: 'maven' + +ext.optionalDeps = [] +ext.providedDeps = [] + +ext.optional = { optionalDeps << it } +ext.provided = { providedDeps << it } + +install { + repositories.mavenInstaller { + customizePom(pom, project) + } +} + +def customizePom(pom, gradleProject) { + pom.whenConfigured { generatedPom -> + // respect 'optional' and 'provided' dependencies + gradleProject.optionalDeps.each { dep -> + generatedPom.dependencies.find { it.artifactId == dep.name }?.optional = true + } + gradleProject.providedDeps.each { dep -> + generatedPom.dependencies.find { it.artifactId == dep.name }?.scope = 'provided' + } + + // eliminate test-scoped dependencies (no need in maven central poms) + generatedPom.dependencies.removeAll { dep -> + dep.scope == 'test' + } + + // add all items necessary for maven central publication + generatedPom.project { + name = gradleProject.description + description = gradleProject.description + url = linkHomepage + organization { + name = 'SpringIO' + url = 'http://spring.io' + } + licenses { + license { + name 'The Apache Software License, Version 2.0' + url 'http://www.apache.org/licenses/LICENSE-2.0.txt' + distribution 'repo' + } + } + + scm { + url = linkScmUrl + connection = 'scm:git:' + linkScmConnection + developerConnection = 'scm:git:' + linkScmDevConnection + } + + developers { + developer { + id = 'venilnoronha' + name = 'Venil Noronha' + email = 'venil.noronha@gmail.com' + } + } + } + } +} diff --git a/spring-integration-etcd/settings.gradle b/spring-integration-etcd/settings.gradle new file mode 100644 index 0000000..054b5c5 --- /dev/null +++ b/spring-integration-etcd/settings.gradle @@ -0,0 +1 @@ +rootProject.name = 'spring-integration-etcd' diff --git a/spring-integration-etcd/src/dist/license.txt b/spring-integration-etcd/src/dist/license.txt new file mode 100644 index 0000000..261eeb9 --- /dev/null +++ b/spring-integration-etcd/src/dist/license.txt @@ -0,0 +1,201 @@ + Apache License + Version 2.0, January 2004 + http://www.apache.org/licenses/ + + TERMS AND CONDITIONS FOR USE, REPRODUCTION, AND DISTRIBUTION + + 1. Definitions. + + "License" shall mean the terms and conditions for use, reproduction, + and distribution as defined by Sections 1 through 9 of this document. + + "Licensor" shall mean the copyright owner or entity authorized by + the copyright owner that is granting the License. + + "Legal Entity" shall mean the union of the acting entity and all + other entities that control, are controlled by, or are under common + control with that entity. For the purposes of this definition, + "control" means (i) the power, direct or indirect, to cause the + direction or management of such entity, whether by contract or + otherwise, or (ii) ownership of fifty percent (50%) or more of the + outstanding shares, or (iii) beneficial ownership of such entity. + + "You" (or "Your") shall mean an individual or Legal Entity + exercising permissions granted by this License. + + "Source" form shall mean the preferred form for making modifications, + including but not limited to software source code, documentation + source, and configuration files. + + "Object" form shall mean any form resulting from mechanical + transformation or translation of a Source form, including but + not limited to compiled object code, generated documentation, + and conversions to other media types. + + "Work" shall mean the work of authorship, whether in Source or + Object form, made available under the License, as indicated by a + copyright notice that is included in or attached to the work + (an example is provided in the Appendix below). + + "Derivative Works" shall mean any work, whether in Source or Object + form, that is based on (or derived from) the Work and for which the + editorial revisions, annotations, elaborations, or other modifications + represent, as a whole, an original work of authorship. For the purposes + of this License, Derivative Works shall not include works that remain + separable from, or merely link (or bind by name) to the interfaces of, + the Work and Derivative Works thereof. + + "Contribution" shall mean any work of authorship, including + the original version of the Work and any modifications or additions + to that Work or Derivative Works thereof, that is intentionally + submitted to Licensor for inclusion in the Work by the copyright owner + or by an individual or Legal Entity authorized to submit on behalf of + the copyright owner. For the purposes of this definition, "submitted" + means any form of electronic, verbal, or written communication sent + to the Licensor or its representatives, including but not limited to + communication on electronic mailing lists, source code control systems, + and issue tracking systems that are managed by, or on behalf of, the + Licensor for the purpose of discussing and improving the Work, but + excluding communication that is conspicuously marked or otherwise + designated in writing by the copyright owner as "Not a Contribution." + + "Contributor" shall mean Licensor and any individual or Legal Entity + on behalf of whom a Contribution has been received by Licensor and + subsequently incorporated within the Work. + + 2. Grant of Copyright License. Subject to the terms and conditions of + this License, each Contributor hereby grants to You a perpetual, + worldwide, non-exclusive, no-charge, royalty-free, irrevocable + copyright license to reproduce, prepare Derivative Works of, + publicly display, publicly perform, sublicense, and distribute the + Work and such Derivative Works in Source or Object form. + + 3. Grant of Patent License. Subject to the terms and conditions of + this License, each Contributor hereby grants to You a perpetual, + worldwide, non-exclusive, no-charge, royalty-free, irrevocable + (except as stated in this section) patent license to make, have made, + use, offer to sell, sell, import, and otherwise transfer the Work, + where such license applies only to those patent claims licensable + by such Contributor that are necessarily infringed by their + Contribution(s) alone or by combination of their Contribution(s) + with the Work to which such Contribution(s) was submitted. If You + institute patent litigation against any entity (including a + cross-claim or counterclaim in a lawsuit) alleging that the Work + or a Contribution incorporated within the Work constitutes direct + or contributory patent infringement, then any patent licenses + granted to You under this License for that Work shall terminate + as of the date such litigation is filed. + + 4. Redistribution. You may reproduce and distribute copies of the + Work or Derivative Works thereof in any medium, with or without + modifications, and in Source or Object form, provided that You + meet the following conditions: + + (a) You must give any other recipients of the Work or + Derivative Works a copy of this License; and + + (b) You must cause any modified files to carry prominent notices + stating that You changed the files; and + + (c) You must retain, in the Source form of any Derivative Works + that You distribute, all copyright, patent, trademark, and + attribution notices from the Source form of the Work, + excluding those notices that do not pertain to any part of + the Derivative Works; and + + (d) If the Work includes a "NOTICE" text file as part of its + distribution, then any Derivative Works that You distribute must + include a readable copy of the attribution notices contained + within such NOTICE file, excluding those notices that do not + pertain to any part of the Derivative Works, in at least one + of the following places: within a NOTICE text file distributed + as part of the Derivative Works; within the Source form or + documentation, if provided along with the Derivative Works; or, + within a display generated by the Derivative Works, if and + wherever such third-party notices normally appear. The contents + of the NOTICE file are for informational purposes only and + do not modify the License. You may add Your own attribution + notices within Derivative Works that You distribute, alongside + or as an addendum to the NOTICE text from the Work, provided + that such additional attribution notices cannot be construed + as modifying the License. + + You may add Your own copyright statement to Your modifications and + may provide additional or different license terms and conditions + for use, reproduction, or distribution of Your modifications, or + for any such Derivative Works as a whole, provided Your use, + reproduction, and distribution of the Work otherwise complies with + the conditions stated in this License. + + 5. Submission of Contributions. Unless You explicitly state otherwise, + any Contribution intentionally submitted for inclusion in the Work + by You to the Licensor shall be under the terms and conditions of + this License, without any additional terms or conditions. + Notwithstanding the above, nothing herein shall supersede or modify + the terms of any separate license agreement you may have executed + with Licensor regarding such Contributions. + + 6. Trademarks. This License does not grant permission to use the trade + names, trademarks, service marks, or product names of the Licensor, + except as required for reasonable and customary use in describing the + origin of the Work and reproducing the content of the NOTICE file. + + 7. Disclaimer of Warranty. Unless required by applicable law or + agreed to in writing, Licensor provides the Work (and each + Contributor provides its Contributions) on an "AS IS" BASIS, + WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or + implied, including, without limitation, any warranties or conditions + of TITLE, NON-INFRINGEMENT, MERCHANTABILITY, or FITNESS FOR A + PARTICULAR PURPOSE. You are solely responsible for determining the + appropriateness of using or redistributing the Work and assume any + risks associated with Your exercise of permissions under this License. + + 8. Limitation of Liability. In no event and under no legal theory, + whether in tort (including negligence), contract, or otherwise, + unless required by applicable law (such as deliberate and grossly + negligent acts) or agreed to in writing, shall any Contributor be + liable to You for damages, including any direct, indirect, special, + incidental, or consequential damages of any character arising as a + result of this License or out of the use or inability to use the + Work (including but not limited to damages for loss of goodwill, + work stoppage, computer failure or malfunction, or any and all + other commercial damages or losses), even if such Contributor + has been advised of the possibility of such damages. + + 9. Accepting Warranty or Additional Liability. While redistributing + the Work or Derivative Works thereof, You may choose to offer, + and charge a fee for, acceptance of support, warranty, indemnity, + or other liability obligations and/or rights consistent with this + License. However, in accepting such obligations, You may act only + on Your own behalf and on Your sole responsibility, not on behalf + of any other Contributor, and only if You agree to indemnify, + defend, and hold each Contributor harmless for any liability + incurred by, or claims asserted against, such Contributor by reason + of your accepting any such warranty or additional liability. + + END OF TERMS AND CONDITIONS + + APPENDIX: How to apply the Apache License to your work. + + To apply the Apache License to your work, attach the following + boilerplate notice, with the fields enclosed by brackets "[]" + replaced with your own identifying information. (Don't include + the brackets!) The text should be enclosed in the appropriate + comment syntax for the file format. We also recommend that a + file or class name and description of purpose be included on the + same "printed page" as the copyright notice for easier + identification within third-party archives. + + Copyright [yyyy] [name of copyright owner] + + 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 + + http://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. diff --git a/spring-integration-etcd/src/dist/notice.txt b/spring-integration-etcd/src/dist/notice.txt new file mode 100644 index 0000000..882d618 --- /dev/null +++ b/spring-integration-etcd/src/dist/notice.txt @@ -0,0 +1,21 @@ + ======================================================================== + == NOTICE file corresponding to section 4 d of the Apache License, == + == Version 2.0, in this case for the Spring Integration distribution. == + ======================================================================== + + This product includes software developed by + the Apache Software Foundation (http://www.apache.org). + + The end-user documentation included with a redistribution, if any, + must include the following acknowledgement: + + "This product includes software developed by the Spring Framework + Project (http://www.spring.io)." + + Alternatively, this acknowledgement may appear in the software itself, + if and wherever such third-party acknowledgements normally appear. + + The names "Spring", "Spring Framework", and "Spring Integration" must + not be used to endorse or promote products derived from this software + without prior written permission. For written permission, please contact + enquiries@springsource.com. diff --git a/spring-integration-etcd/src/main/java/org/springframework/integration/etcd/leader/LeaderInitiator.java b/spring-integration-etcd/src/main/java/org/springframework/integration/etcd/leader/LeaderInitiator.java new file mode 100644 index 0000000..9a783ba --- /dev/null +++ b/spring-integration-etcd/src/main/java/org/springframework/integration/etcd/leader/LeaderInitiator.java @@ -0,0 +1,453 @@ +/* + * Copyright 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 + * + * http://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.integration.etcd.leader; + +import static com.ibm.etcd.client.KeyUtils.bs; + +import java.util.concurrent.Callable; +import java.util.concurrent.CancellationException; +import java.util.concurrent.ExecutionException; +import java.util.concurrent.ExecutorService; +import java.util.concurrent.Executors; +import java.util.concurrent.Future; +import java.util.concurrent.TimeUnit; + +import org.apache.commons.logging.Log; +import org.apache.commons.logging.LogFactory; + +import org.springframework.beans.factory.DisposableBean; +import org.springframework.context.ApplicationEventPublisher; +import org.springframework.context.ApplicationEventPublisherAware; +import org.springframework.context.Lifecycle; +import org.springframework.integration.leader.Candidate; +import org.springframework.integration.leader.Context; +import org.springframework.integration.leader.DefaultCandidate; +import org.springframework.integration.leader.event.DefaultLeaderEventPublisher; +import org.springframework.integration.leader.event.LeaderEventPublisher; +import org.springframework.util.Assert; + +import com.ibm.etcd.api.DeleteRangeRequest; +import com.ibm.etcd.api.PutRequest; +import com.ibm.etcd.api.TxnResponse; +import com.ibm.etcd.client.EtcdClient; +import io.grpc.StatusRuntimeException; + +/** + * Bootstrap leadership {@link Candidate candidates} with etcd. Upon construction, {@link #start} + * must be invoked to register the candidate for leadership election. + * + * @author Venil Noronha + * @author Patrick Peralta + * @author Lewis Headden + */ +public class LeaderInitiator implements Lifecycle, DisposableBean, ApplicationEventPublisherAware { + + private final Log logger = LogFactory.getLog(getClass()); + + /** + * TTL for etcd entry in seconds. + */ + private static final int TTL = 10; + + /** + * Number of seconds to sleep between issuing heartbeats. + */ + private static final int HEART_BEAT_SLEEP = TTL / 2; + + /** + * Default namespace for etcd entry. + */ + private static final String DEFAULT_NAMESPACE = "spring-integration"; + + /** + * {@link EtcdClient} instance. + */ + private final EtcdClient client; + + /** + * Candidate for leader election. + */ + private final Candidate candidate; + + /** + * Executor service for running leadership daemon. + */ + private final ExecutorService leaderExecutorService = + Executors.newSingleThreadExecutor( + r -> { + Thread thread = new Thread(r, "Etcd-Leadership"); + thread.setDaemon(true); + return thread; + }); + + /** + * Executor service for running leadership worker daemon. + */ + private final ExecutorService workerExecutorService = + Executors.newSingleThreadExecutor( + r -> { + Thread thread = new Thread(r, "Etcd-Leadership-Worker"); + thread.setDaemon(true); + return thread; + }); + + /** + * Flag that indicates whether the current candidate is the leader. + */ + private volatile boolean isLeader = false; + + /** + * Flag that indicates whether the current candidate's leadership should be relinquished. + */ + private volatile boolean relinquishLeadership = false; + + /** + * Future returned by submitting a {@link Initiator} to {@link #leaderExecutorService}. + * This is used to cancel leadership. + */ + private volatile Future initiatorFuture; + + /** + * Future returned by submitting a {@link Worker} to {@link #workerExecutorService}. + * This is used to notify leadership revocation. + */ + private volatile Future workerFuture; + + /** + * Flag that indicates whether the leadership election for this {@link #candidate} is running. + */ + private volatile boolean running; + + /** + * Leader event publisher. + */ + private volatile LeaderEventPublisher leaderEventPublisher = new DefaultLeaderEventPublisher(); + + private boolean customPublisher = false; + + /** + * The {@link EtcdContext} instance. + */ + private final EtcdContext context; + + /** + * The base etcd path where candidate id is to be stored. + */ + private final String baseEtcdPath; + + /** + * Construct a {@link LeaderInitiator}. + * @param client {@link EtcdClient} instance + * @param namespace Etcd namespace + */ + public LeaderInitiator(EtcdClient client, String namespace) { + this(client, new DefaultCandidate(), namespace); + } + + /** + * Construct a {@link LeaderInitiator}. + * @param client {@link EtcdClient} instance + * @param candidate leadership election candidate + * @param namespace Etcd namespace + */ + public LeaderInitiator(EtcdClient client, Candidate candidate, String namespace) { + Assert.notNull(client, "'etcdClient' must not be null"); + Assert.notNull(candidate, "'candidate' must not be null"); + this.client = client; + this.candidate = candidate; + this.context = new EtcdContext(); + this.baseEtcdPath = (namespace == null ? DEFAULT_NAMESPACE : namespace) + "/" + candidate.getRole(); + } + + @Override + public void setApplicationEventPublisher(ApplicationEventPublisher applicationEventPublisher) { + if (!this.customPublisher) { + this.leaderEventPublisher = new DefaultLeaderEventPublisher(applicationEventPublisher); + } + } + + /** + * Start the registration of the {@link #candidate} for leader election. + */ + @Override + public synchronized void start() { + if (!this.running) { + this.initiatorFuture = this.leaderExecutorService.submit(new Initiator()); + this.running = true; + } + } + + /** + * Stop the registration of the {@link #candidate} for leader election. + * If the candidate is currently leader, its leadership will be revoked. + */ + @Override + public synchronized void stop() { + if (this.running) { + this.running = false; + this.initiatorFuture.cancel(true); + } + } + + /** + * @return true if leadership election for this {@link #candidate} is running + */ + @Override + public boolean isRunning() { + return this.running; + } + + @Override + public void destroy() throws Exception { + stop(); + this.workerExecutorService.shutdown(); + this.leaderExecutorService.shutdown(); + this.workerExecutorService.awaitTermination(30, TimeUnit.SECONDS); + this.leaderExecutorService.awaitTermination(30, TimeUnit.SECONDS); + } + + /** + * Sets the {@link LeaderEventPublisher}. + * @param leaderEventPublisher the event publisher + */ + public void setLeaderEventPublisher(LeaderEventPublisher leaderEventPublisher) { + Assert.notNull(leaderEventPublisher, "leaderEventPublisher cannot be null"); + this.leaderEventPublisher = leaderEventPublisher; + } + + /** + * Notifies that the candidate has acquired leadership. + */ + private void notifyGranted() { + this.isLeader = true; + this.leaderEventPublisher.publishOnGranted( + LeaderInitiator.this, this.context, this.candidate.getRole()); + this.workerFuture = this.workerExecutorService.submit(new Worker()); + } + + /** + * Notifies that the candidate's leadership was revoked. + * @throws InterruptedException if the current thread was interrupted while waiting for the worker + * thread to finish. + */ + private void notifyRevoked() throws InterruptedException { + this.isLeader = false; + this.leaderEventPublisher.publishOnRevoked( + LeaderInitiator.this, this.context, this.candidate.getRole()); + this.workerFuture.cancel(true); + try { + this.workerFuture.get(); + } + catch (InterruptedException e) { + throw e; + } + catch (CancellationException e) { + // Consume + } + catch (ExecutionException e) { + logger.error("Exception thrown by candidate", e.getCause()); + } + } + + /** + * Tries to delete the candidate's entry from etcd. + */ + private void tryDeleteCandidateEntry() { + try { + final TxnResponse response = + this.client + .getKvClient() + .txnIf() + .cmpEqual(bs(this.baseEtcdPath)) + .value(bs(this.candidate.getId())) + .then() + .delete(DeleteRangeRequest.newBuilder().setKey(bs(this.baseEtcdPath)).build()) + .sync(); + + if (!response.getSucceeded()) { + logger.warn("Couldn't delete candidate's entry from etcd because candidate was not leader"); + } + } + catch (StatusRuntimeException e) { + logger.warn("Failed deleting candidate entry from etcd", e); + } + } + + /** + * Callable that invokes {@link Candidate#onGranted(Context)} when the candidate is granted + * leadership. + */ + class Worker implements Callable { + + @Override + public Void call() throws InterruptedException { + try { + LeaderInitiator.this.candidate.onGranted(LeaderInitiator.this.context); + Thread.sleep(Long.MAX_VALUE); + } + finally { + LeaderInitiator.this.relinquishLeadership = true; + LeaderInitiator.this.candidate.onRevoked(LeaderInitiator.this.context); + } + return null; + } + + } + + /** + * Callable that manages the etcd heart beats for leadership election. + */ + class Initiator implements Callable { + + @Override + public Void call() throws InterruptedException { + try { + while (LeaderInitiator.this.running) { + if (LeaderInitiator.this.relinquishLeadership) { + relinquishLeadership(); + LeaderInitiator.this.relinquishLeadership = false; + } + else if (LeaderInitiator.this.isLeader) { + sendHeartBeat(); + } + else { + tryAcquire(); + } + TimeUnit.SECONDS.sleep(HEART_BEAT_SLEEP); + } + } + finally { + if (LeaderInitiator.this.isLeader) { + relinquishLeadership(); + } + } + return null; + } + + /** + * Relinquishes leadership of current candidate by deleting candidate's entry from etcd and then + * notifies that the current candidate is no longer leader. + * @throws InterruptedException if the current thread was interrupted while notifying + * revocation. + */ + private void relinquishLeadership() throws InterruptedException { + tryDeleteCandidateEntry(); + notifyRevoked(); + } + + /** + * Sends a heart beat to maintain leadership by refreshing the ttl of the etcd key. If the key + * has a different value during the call, it is assumed that the current candidate's leadership + * is revoked. If access to etcd fails, then the the current candidate's leadership is + * relinquished. + * @throws InterruptedException if the current thread was interrupted while notifying + * revocation. + */ + private void sendHeartBeat() throws InterruptedException { + try { + final TxnResponse response = + LeaderInitiator.this.client + .getKvClient() + .txnIf() + .cmpEqual(bs(LeaderInitiator.this.baseEtcdPath)) + .value(bs(LeaderInitiator.this.candidate.getId())) + .then() + .put( + PutRequest.newBuilder() + .setKey(bs(LeaderInitiator.this.baseEtcdPath)) + .setValue(bs(LeaderInitiator.this.candidate.getId())) + .build()) + .sync(); + + if (!response.getSucceeded()) { + notifyRevoked(); + } + } + catch (StatusRuntimeException e) { + logger.warn("Failed to send leadership heartbeat", e); + } + } + + /** + * Tries to acquire leadership by posting the candidate's id to etcd. If the etcd call is + * successful, it is assumed that the current candidate is now leader. + */ + private void tryAcquire() { + try { + final long leaseId = + LeaderInitiator.this.client + .getLeaseClient() + .create(TTL) + .get() + .getID(); + + final TxnResponse response = + LeaderInitiator.this.client + .getKvClient() + .txnIf() + .notExists(bs(LeaderInitiator.this.baseEtcdPath)) + .then() + .put( + PutRequest.newBuilder() + .setKey(bs(LeaderInitiator.this.baseEtcdPath)) + .setValue(bs(LeaderInitiator.this.candidate.getId())) + .setLease(leaseId) + .build()) + .sync(); + if (response.getSucceeded()) { + notifyGranted(); + } + else { + logger.info("Tried to acquire leadership but another candidate is leader"); + } + } + catch (InterruptedException | ExecutionException e) { + logger.warn("Failed trying to acquire leadership", e.getCause()); + } + } + + } + + /** + * Implementation of leadership context backed by Etcd. + */ + class EtcdContext implements Context { + + @Override + public boolean isLeader() { + return LeaderInitiator.this.isLeader; + } + + @Override + public void yield() { + if (LeaderInitiator.this.isLeader) { + LeaderInitiator.this.relinquishLeadership = true; + } + } + + @Override + public String toString() { + return String.format( + "EtcdContext{role=%s, id=%s, isLeader=%s}", + LeaderInitiator.this.candidate.getRole(), + LeaderInitiator.this.candidate.getId(), + isLeader()); + } + + } + +} diff --git a/spring-integration-etcd/src/test/java/org/springframework/integration/etcd/leader/EtcdTests.java b/spring-integration-etcd/src/test/java/org/springframework/integration/etcd/leader/EtcdTests.java new file mode 100644 index 0000000..299a06e --- /dev/null +++ b/spring-integration-etcd/src/test/java/org/springframework/integration/etcd/leader/EtcdTests.java @@ -0,0 +1,346 @@ +/* + * Copyright 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 + * + * http://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.integration.etcd.leader; + +import static org.hamcrest.CoreMatchers.is; +import static org.junit.Assert.assertThat; + +import java.util.ArrayList; +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.TimeUnit; + +import org.apache.commons.logging.Log; +import org.apache.commons.logging.LogFactory; +import org.junit.Test; + +import org.springframework.context.ApplicationListener; +import org.springframework.context.annotation.AnnotationConfigApplicationContext; +import org.springframework.context.annotation.Bean; +import org.springframework.context.annotation.Configuration; +import org.springframework.integration.leader.AbstractCandidate; +import org.springframework.integration.leader.Context; +import org.springframework.integration.leader.DefaultCandidate; +import org.springframework.integration.leader.event.AbstractLeaderEvent; +import org.springframework.integration.leader.event.DefaultLeaderEventPublisher; +import org.springframework.integration.leader.event.LeaderEventPublisher; + +import com.ibm.etcd.client.EtcdClient; + +/** + * Tests for etcd leader election. + * + * @author Venil Noronha + * @author Lewis Headden + */ +public class EtcdTests { + + @Test + public void testSimpleLeader() throws InterruptedException { + AnnotationConfigApplicationContext ctx = new AnnotationConfigApplicationContext(SimpleTestConfig.class); + TestCandidate candidate = ctx.getBean(TestCandidate.class); + TestEventListener listener = ctx.getBean(TestEventListener.class); + assertThat(candidate.onGrantedLatch.await(5, TimeUnit.SECONDS), is(true)); + assertThat(listener.onEventLatch.await(5, TimeUnit.SECONDS), is(true)); + assertThat(listener.events.size(), is(1)); + ctx.close(); + } + + @Test + public void testLeaderYield() throws InterruptedException { + AnnotationConfigApplicationContext ctx = new AnnotationConfigApplicationContext(YieldTestConfig.class); + YieldTestCandidate candidate = ctx.getBean(YieldTestCandidate.class); + YieldTestEventListener listener = ctx.getBean(YieldTestEventListener.class); + assertThat(candidate.onGrantedLatch.await(5, TimeUnit.SECONDS), is(true)); + assertThat(candidate.onRevokedLatch.await(10, TimeUnit.SECONDS), is(true)); + assertThat(listener.onEventsLatch.await(1, TimeUnit.MILLISECONDS), is(true)); + assertThat(listener.events.size(), is(2)); + ctx.close(); + } + + @Test + public void testBlockingThreadLeader() throws InterruptedException { + AnnotationConfigApplicationContext ctx = + new AnnotationConfigApplicationContext(BlockingThreadTestConfig.class); + BlockingThreadTestCandidate candidate = ctx.getBean(BlockingThreadTestCandidate.class); + YieldTestEventListener listener = ctx.getBean(YieldTestEventListener.class); + assertThat(candidate.onGrantedLatch.await(5, TimeUnit.SECONDS), is(true)); + Thread.sleep(2000); // Let the grant-notification thread run for a while + candidate.ctx.yield(); // Internally interrupts the grant notification thread + assertThat(candidate.onRevokedLatch.await(10, TimeUnit.SECONDS), is(true)); + assertThat(listener.onEventsLatch.await(1, TimeUnit.MILLISECONDS), is(true)); + assertThat(listener.events.size(), is(2)); + ctx.close(); + } + + @Test + public void testFailingCandidateGrantCallback() throws InterruptedException { + AnnotationConfigApplicationContext ctx = + new AnnotationConfigApplicationContext(FailingCandidateTestConfig.class); + FailingTestCandidate candidate = ctx.getBean(FailingTestCandidate.class); + YieldTestEventListener listener = ctx.getBean(YieldTestEventListener.class); + assertThat(candidate.onGrantedLatch.await(5, TimeUnit.SECONDS), is(true)); + assertThat(candidate.onRevokedLatch.await(10, TimeUnit.SECONDS), is(true)); + assertThat(listener.onEventsLatch.await(10, TimeUnit.SECONDS), is(true)); + assertThat(listener.events.size(), is(2)); + ctx.close(); + } + + @Configuration + static class SimpleTestConfig { + + @Bean + public TestCandidate candidate() { + return new TestCandidate(); + } + + @Bean + public EtcdClient etcdInstance() { + return EtcdClient.forEndpoint("localhost", 2379).withPlainText().build(); + } + + @Bean + public LeaderInitiator initiator() { + LeaderInitiator initiator = new LeaderInitiator(etcdInstance(), candidate(), "etcd-test"); + initiator.start(); + return initiator; + } + + @Bean + public TestEventListener testEventListener() { + return new TestEventListener(); + } + + } + + static class TestCandidate extends DefaultCandidate { + + CountDownLatch onGrantedLatch = new CountDownLatch(1); + + @Override + public void onGranted(Context ctx) { + this.onGrantedLatch.countDown(); + super.onGranted(ctx); + } + + } + + static class TestEventListener implements ApplicationListener { + + CountDownLatch onEventLatch = new CountDownLatch(1); + + ArrayList events = new ArrayList<>(); + + @Override + public void onApplicationEvent(AbstractLeaderEvent event) { + this.events.add(event); + this.onEventLatch.countDown(); + } + + } + + @Configuration + static class YieldTestConfig { + + @Bean + public YieldTestCandidate candidate() { + return new YieldTestCandidate(); + } + + @Bean + public EtcdClient etcdInstance() { + return EtcdClient.forEndpoint("localhost", 2379).withPlainText().build(); + } + + @Bean + public LeaderInitiator initiator() { + LeaderInitiator initiator = new LeaderInitiator(etcdInstance(), candidate(), "etcd-yield-test"); + initiator.setLeaderEventPublisher(leaderEventPublisher()); + initiator.start(); + return initiator; + } + + @Bean + public LeaderEventPublisher leaderEventPublisher() { + return new DefaultLeaderEventPublisher(); + } + + @Bean + public YieldTestEventListener testEventListener() { + return new YieldTestEventListener(); + } + + } + + static class YieldTestCandidate extends DefaultCandidate { + + CountDownLatch onGrantedLatch = new CountDownLatch(1); + + CountDownLatch onRevokedLatch = new CountDownLatch(1); + + @Override + public void onGranted(Context ctx) { + super.onGranted(ctx); + this.onGrantedLatch.countDown(); + ctx.yield(); + } + + @Override + public void onRevoked(Context ctx) { + super.onRevoked(ctx); + this.onRevokedLatch.countDown(); + } + + } + + static class YieldTestEventListener implements ApplicationListener { + + CountDownLatch onEventsLatch = new CountDownLatch(2); + + ArrayList events = new ArrayList<>(); + + @Override + public void onApplicationEvent(AbstractLeaderEvent event) { + this.events.add(event); + this.onEventsLatch.countDown(); + } + + } + + @Configuration + static class BlockingThreadTestConfig { + + @Bean + public BlockingThreadTestCandidate candidate() { + return new BlockingThreadTestCandidate(); + } + + @Bean + public EtcdClient etcdInstance() { + return EtcdClient.forEndpoint("localhost", 2379).withPlainText().build(); + } + + @Bean + public LeaderInitiator initiator() { + LeaderInitiator initiator = new LeaderInitiator(etcdInstance(), candidate(), "etcd-blocking-thread-test"); + initiator.setLeaderEventPublisher(leaderEventPublisher()); + initiator.start(); + return initiator; + } + + @Bean + public LeaderEventPublisher leaderEventPublisher() { + return new DefaultLeaderEventPublisher(); + } + + @Bean + public YieldTestEventListener testEventListener() { + return new YieldTestEventListener(); + } + + } + + static class BlockingThreadTestCandidate extends AbstractCandidate { + + private final Log logger = LogFactory.getLog(getClass()); + + CountDownLatch onGrantedLatch = new CountDownLatch(1); + + CountDownLatch onRevokedLatch = new CountDownLatch(1); + + Context ctx = null; + + @Override + public void onGranted(Context ctx) throws InterruptedException { + this.ctx = ctx; + logger.info(this + " has been granted leadership; context: " + ctx); + this.onGrantedLatch.countDown(); + while (true) { + logger.info(this + " is doing some heavy lifting"); + try { + Thread.sleep(1000); // Mock heavy lifting + } + catch (InterruptedException e) { + logger.info(this + " was interrupted, rethrowing the exception"); + throw e; + } + } + } + + @Override + public void onRevoked(Context ctx) { + logger.info(this + " leadership has been revoked"); + this.onRevokedLatch.countDown(); + } + + } + + @Configuration + static class FailingCandidateTestConfig { + + @Bean + public FailingTestCandidate candidate() { + return new FailingTestCandidate(); + } + + @Bean + public EtcdClient etcdInstance() { + return EtcdClient.forEndpoint("localhost", 2379).withPlainText().build(); + } + + @Bean + public LeaderInitiator initiator() { + LeaderInitiator initiator = + new LeaderInitiator(etcdInstance(), candidate(), "etcd-failing-candidate-test"); + initiator.setLeaderEventPublisher(leaderEventPublisher()); + initiator.start(); + return initiator; + } + + @Bean + public LeaderEventPublisher leaderEventPublisher() { + return new DefaultLeaderEventPublisher(); + } + + @Bean + public YieldTestEventListener testEventListener() { + return new YieldTestEventListener(); + } + + } + + static class FailingTestCandidate extends DefaultCandidate { + + CountDownLatch onGrantedLatch = new CountDownLatch(1); + + CountDownLatch onRevokedLatch = new CountDownLatch(1); + + @Override + public void onGranted(Context ctx) { + super.onGranted(ctx); + this.onGrantedLatch.countDown(); + throw new RuntimeException("Candidate grant callback failure"); + } + + @Override + public void onRevoked(Context ctx) { + super.onRevoked(ctx); + this.onRevokedLatch.countDown(); + } + + } + +} diff --git a/spring-integration-etcd/src/test/resources/log4j2-test.xml b/spring-integration-etcd/src/test/resources/log4j2-test.xml new file mode 100644 index 0000000..b04559d --- /dev/null +++ b/spring-integration-etcd/src/test/resources/log4j2-test.xml @@ -0,0 +1,17 @@ + + + + + + + + + + + + + + + + + \ No newline at end of file