This is an automated email from the ASF dual-hosted git repository.
github-merge-queue[bot] pushed a commit to branch dev
in repository https://gitbox.apache.org/repos/asf/seatunnel.git
The following commit(s) were added to refs/heads/dev by this push:
new 70a6720fc3 [Improve][E2E] Reuse SeaTunnel container for selected test
classes (#11626)
70a6720fc3 is described below
commit 70a6720fc31867ef9fc00cc1a085dc30424bcf23
Author: Goutam Adwant <[email protected]>
AuthorDate: Tue Sep 15 13:02:42 2026 +0000
[Improve][E2E] Reuse SeaTunnel container for selected test classes (#11626)
Signed-off-by: goutamadwant <[email protected]>
---
.../seatunnel/e2e/connector/fake/FakeIT.java | 6 +
.../e2e/connector/fake/FakeSqlConfIT.java | 6 +
...ConfIT.java => SharedContainerTestSupport.java} | 20 +-
.../common/container/ReusableTestContainer.java | 39 +++
.../container/seatunnel/SeaTunnelContainer.java | 292 ++++++++++++++++++---
.../common/junit/ContainerTestingExtension.java | 89 ++++++-
.../e2e/common/junit/ReuseTestContainers.java | 42 +++
.../common/junit/SharedTestContainerResource.java | 100 +++++++
.../junit/SharedTestContainerResourceTest.java | 267 +++++++++++++++++++
.../junit/TestCaseInvocationContextProvider.java | 41 ++-
10 files changed, 850 insertions(+), 52 deletions(-)
diff --git
a/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-fake-e2e/src/test/java/org/apache/seatunnel/e2e/connector/fake/FakeIT.java
b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-fake-e2e/src/test/java/org/apache/seatunnel/e2e/connector/fake/FakeIT.java
index 6821b06a78..c1067161e6 100644
---
a/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-fake-e2e/src/test/java/org/apache/seatunnel/e2e/connector/fake/FakeIT.java
+++
b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-fake-e2e/src/test/java/org/apache/seatunnel/e2e/connector/fake/FakeIT.java
@@ -19,6 +19,8 @@ package org.apache.seatunnel.e2e.connector.fake;
import org.apache.seatunnel.e2e.common.TestSuiteBase;
import org.apache.seatunnel.e2e.common.container.TestContainer;
+import org.apache.seatunnel.e2e.common.container.TestContainerId;
+import org.apache.seatunnel.e2e.common.junit.ReuseTestContainers;
import org.junit.jupiter.api.Assertions;
import org.junit.jupiter.api.TestTemplate;
@@ -26,10 +28,14 @@ import org.testcontainers.containers.Container;
import java.io.IOException;
+@ReuseTestContainers(TestContainerId.SEATUNNEL)
public class FakeIT extends TestSuiteBase {
@TestTemplate
public void testFakeConnector(TestContainer container)
throws IOException, InterruptedException {
+ if (container.identifier() == TestContainerId.SEATUNNEL) {
+ SharedContainerTestSupport.assertSameContainer(container);
+ }
Container.ExecResult textWriteResult =
container.executeJob("/fake_to_assert.conf");
Assertions.assertEquals(0, textWriteResult.getExitCode());
Container.ExecResult fakeWithRange =
diff --git
a/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-fake-e2e/src/test/java/org/apache/seatunnel/e2e/connector/fake/FakeSqlConfIT.java
b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-fake-e2e/src/test/java/org/apache/seatunnel/e2e/connector/fake/FakeSqlConfIT.java
index cb603a9c4f..d080577d3f 100644
---
a/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-fake-e2e/src/test/java/org/apache/seatunnel/e2e/connector/fake/FakeSqlConfIT.java
+++
b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-fake-e2e/src/test/java/org/apache/seatunnel/e2e/connector/fake/FakeSqlConfIT.java
@@ -19,6 +19,8 @@ package org.apache.seatunnel.e2e.connector.fake;
import org.apache.seatunnel.e2e.common.TestSuiteBase;
import org.apache.seatunnel.e2e.common.container.TestContainer;
+import org.apache.seatunnel.e2e.common.container.TestContainerId;
+import org.apache.seatunnel.e2e.common.junit.ReuseTestContainers;
import org.junit.jupiter.api.Assertions;
import org.junit.jupiter.api.TestTemplate;
@@ -26,11 +28,15 @@ import org.testcontainers.containers.Container;
import java.io.IOException;
+@ReuseTestContainers(TestContainerId.SEATUNNEL)
public class FakeSqlConfIT extends TestSuiteBase {
@TestTemplate
public void testFakeConnector(TestContainer container)
throws IOException, InterruptedException {
+ if (container.identifier() == TestContainerId.SEATUNNEL) {
+ SharedContainerTestSupport.assertSameContainer(container);
+ }
Container.ExecResult textWriteResult =
container.executeJob("/fake_to_assert.sql");
Assertions.assertEquals(0, textWriteResult.getExitCode());
}
diff --git
a/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-fake-e2e/src/test/java/org/apache/seatunnel/e2e/connector/fake/FakeSqlConfIT.java
b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-fake-e2e/src/test/java/org/apache/seatunnel/e2e/connector/fake/SharedContainerTestSupport.java
similarity index 65%
copy from
seatunnel-e2e/seatunnel-connector-v2-e2e/connector-fake-e2e/src/test/java/org/apache/seatunnel/e2e/connector/fake/FakeSqlConfIT.java
copy to
seatunnel-e2e/seatunnel-connector-v2-e2e/connector-fake-e2e/src/test/java/org/apache/seatunnel/e2e/connector/fake/SharedContainerTestSupport.java
index cb603a9c4f..85a70fd0d5 100644
---
a/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-fake-e2e/src/test/java/org/apache/seatunnel/e2e/connector/fake/FakeSqlConfIT.java
+++
b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-fake-e2e/src/test/java/org/apache/seatunnel/e2e/connector/fake/SharedContainerTestSupport.java
@@ -17,21 +17,21 @@
package org.apache.seatunnel.e2e.connector.fake;
-import org.apache.seatunnel.e2e.common.TestSuiteBase;
import org.apache.seatunnel.e2e.common.container.TestContainer;
import org.junit.jupiter.api.Assertions;
-import org.junit.jupiter.api.TestTemplate;
-import org.testcontainers.containers.Container;
-import java.io.IOException;
+final class SharedContainerTestSupport {
-public class FakeSqlConfIT extends TestSuiteBase {
+ private static TestContainer firstContainer;
- @TestTemplate
- public void testFakeConnector(TestContainer container)
- throws IOException, InterruptedException {
- Container.ExecResult textWriteResult =
container.executeJob("/fake_to_assert.sql");
- Assertions.assertEquals(0, textWriteResult.getExitCode());
+ private SharedContainerTestSupport() {}
+
+ static synchronized void assertSameContainer(TestContainer container) {
+ if (firstContainer == null) {
+ firstContainer = container;
+ return;
+ }
+ Assertions.assertSame(firstContainer, container);
}
}
diff --git
a/seatunnel-e2e/seatunnel-e2e-common/src/test/java/org/apache/seatunnel/e2e/common/container/ReusableTestContainer.java
b/seatunnel-e2e/seatunnel-e2e-common/src/test/java/org/apache/seatunnel/e2e/common/container/ReusableTestContainer.java
new file mode 100644
index 0000000000..41aef13ac5
--- /dev/null
+++
b/seatunnel-e2e/seatunnel-e2e-common/src/test/java/org/apache/seatunnel/e2e/common/container/ReusableTestContainer.java
@@ -0,0 +1,39 @@
+/*
+ * 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
+ *
+ * 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.apache.seatunnel.e2e.common.container;
+
+/**
+ * A test container that can remove test-class state without stopping the
underlying service.
+ *
+ * <p>Implementations must be configuration-independent across test classes
because the first
+ * instance is shared with subsequent classes. Per-class setup must be
performed in {@link
+ * #prepareForTestClass()} or through {@link
TestContainer#executeExtraCommands} rather than during
+ * construction or startup. Implementations must detect filesystem or
classloader inputs that cannot
+ * be restored safely and fail cleanup so the shared resource is restarted. In
particular, a class
+ * must not opt into reuse when its setup replaces an existing same-path
connector artifact or
+ * mutates runtime libraries unless the implementation explicitly restores and
invalidates that
+ * state.
+ */
+public interface ReusableTestContainer extends TestContainer {
+
+ /** Records or verifies the clean baseline before a test class uses this
container. */
+ default void prepareForTestClass() throws Exception {}
+
+ /** Removes state owned by the completed test class and verifies the clean
baseline. */
+ void cleanUpAfterTestClass() throws Exception;
+}
diff --git
a/seatunnel-e2e/seatunnel-e2e-common/src/test/java/org/apache/seatunnel/e2e/common/container/seatunnel/SeaTunnelContainer.java
b/seatunnel-e2e/seatunnel-e2e-common/src/test/java/org/apache/seatunnel/e2e/common/container/seatunnel/SeaTunnelContainer.java
index f0adf75ceb..c43ea973d2 100644
---
a/seatunnel-e2e/seatunnel-e2e-common/src/test/java/org/apache/seatunnel/e2e/common/container/seatunnel/SeaTunnelContainer.java
+++
b/seatunnel-e2e/seatunnel-e2e-common/src/test/java/org/apache/seatunnel/e2e/common/container/seatunnel/SeaTunnelContainer.java
@@ -25,6 +25,7 @@ import org.apache.seatunnel.common.utils.FileUtils;
import org.apache.seatunnel.common.utils.JsonUtils;
import org.apache.seatunnel.e2e.common.container.AbstractTestContainer;
import org.apache.seatunnel.e2e.common.container.ContainerExtendedFactory;
+import org.apache.seatunnel.e2e.common.container.ReusableTestContainer;
import org.apache.seatunnel.e2e.common.container.TestContainer;
import org.apache.seatunnel.e2e.common.container.TestContainerId;
import org.apache.seatunnel.e2e.common.util.ContainerUtil;
@@ -32,6 +33,7 @@ import org.apache.seatunnel.e2e.common.util.MavenJarUtil;
import org.apache.commons.compress.utils.Lists;
import org.apache.http.HttpStatus;
+import org.apache.http.client.config.RequestConfig;
import org.apache.http.client.methods.CloseableHttpResponse;
import org.apache.http.client.methods.HttpGet;
import org.apache.http.client.methods.HttpPost;
@@ -60,12 +62,16 @@ import lombok.extern.slf4j.Slf4j;
import java.io.File;
import java.io.IOException;
import java.net.URL;
+import java.nio.file.Path;
import java.nio.file.Paths;
import java.util.ArrayList;
import java.util.Arrays;
import java.util.Collections;
+import java.util.Comparator;
import java.util.List;
import java.util.Map;
+import java.util.Set;
+import java.util.TreeSet;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicInteger;
import java.util.regex.Pattern;
@@ -78,21 +84,28 @@ import static
org.apache.seatunnel.e2e.common.util.ContainerUtil.copyAllConnecto
@NoArgsConstructor
@Slf4j
@AutoService(TestContainer.class)
-public class SeaTunnelContainer extends AbstractTestContainer {
+public class SeaTunnelContainer extends AbstractTestContainer implements
ReusableTestContainer {
public static final String SERVER_JVM_OPTION_PROPERTY =
"seatunnel.e2e.seatunnel.server.jvm.option";
public static final String CLIENT_JVM_OPTION_PROPERTY =
"seatunnel.e2e.seatunnel.client.jvm.option";
-
private static final ObjectMapper OBJECT_MAPPER = new ObjectMapper();
private static final String REST_STOP_JOB_PATH = "/stop-job";
private static final String REST_CHECKPOINT_OVERVIEW_PATH =
"/jobs/checkpoints";
protected static final String JDK_DOCKER_IMAGE =
"seatunnelhub/openjdk:8u342";
private static final String CLIENT_SHELL = "seatunnel.sh";
+ private static final String CONNECTOR_DIRECTORY =
"/tmp/seatunnel/connectors";
+ private static final String RUNTIME_LIBRARY_DIRECTORY =
"/tmp/seatunnel/lib";
+ private static final int REST_REQUEST_TIMEOUT_MILLIS = 5_000;
+ private static final int RUNNING_JOBS_TIMEOUT_SECONDS = 30;
protected static final String SERVER_SHELL = "seatunnel-cluster.sh";
protected static final String CONNECTOR_CHECK_SHELL =
"seatunnel-connector.sh";
protected GenericContainer<?> server;
private final AtomicInteger runningCount = new AtomicInteger();
+ private Set<String> connectorJarsBeforeTest;
+ private Set<String> connectorJarFingerprintsBeforeTest;
+ private Set<String> runtimeLibraryFingerprintsBeforeTest;
+ private Set<String> temporaryConfigsBeforeTest;
@Override
public void startUp() throws Exception {
@@ -258,6 +271,222 @@ public class SeaTunnelContainer extends
AbstractTestContainer {
FileUtils.deleteFile(HOST_VOLUME_MOUNT_PATH);
}
+ @Override
+ public void prepareForTestClass() throws Exception {
+ resetReuseBaselines();
+ try {
+ assertNoRunningJobs();
+ assertVolumeEmpty();
+ Set<String> connectorJars = listConnectorJars();
+ Set<String> connectorJarFingerprints =
listConnectorJarFingerprints();
+ Set<String> runtimeLibraryFingerprints =
listRuntimeLibraryFingerprints();
+ Set<String> temporaryConfigs = listTemporaryConfigs();
+ connectorJarsBeforeTest = connectorJars;
+ connectorJarFingerprintsBeforeTest = connectorJarFingerprints;
+ runtimeLibraryFingerprintsBeforeTest = runtimeLibraryFingerprints;
+ temporaryConfigsBeforeTest = temporaryConfigs;
+ } catch (Exception e) {
+ resetReuseBaselines();
+ throw e;
+ }
+ }
+
+ @Override
+ public void cleanUpAfterTestClass() throws Exception {
+ requireReuseBaselines();
+ if (runningCount.get() != 0) {
+ throw new IllegalStateException("Cannot clean a SeaTunnelContainer
with running jobs");
+ }
+ assertNoRunningJobs();
+ Container.ExecResult cleanupResult =
+ server.execInContainer(
+ "sh",
+ "-c",
+ "if [ ! -d /tmp/seatunnel_mnt ]; then "
+ + "echo '/tmp/seatunnel_mnt is missing' >&2;
exit 1; fi; "
+ + "find /tmp/seatunnel_mnt -mindepth 1
-delete");
+ if (cleanupResult.getExitCode() != 0) {
+ throw new IllegalStateException(
+ "Failed to clean shared SeaTunnelContainer: " +
cleanupResult.getStderr());
+ }
+ deleteAddedArtifacts(CONNECTOR_DIRECTORY, connectorJarsBeforeTest,
listConnectorJars());
+ deleteAddedArtifacts("/tmp", temporaryConfigsBeforeTest,
listTemporaryConfigs());
+ assertVolumeEmpty();
+ assertArtifactsRestored("connector JARs", connectorJarsBeforeTest,
listConnectorJars());
+ assertArtifactsRestored(
+ "connector JAR contents",
+ connectorJarFingerprintsBeforeTest,
+ listConnectorJarFingerprints());
+ assertArtifactsRestored(
+ "runtime libraries",
+ runtimeLibraryFingerprintsBeforeTest,
+ listRuntimeLibraryFingerprints());
+ assertArtifactsRestored(
+ "temporary configs", temporaryConfigsBeforeTest,
listTemporaryConfigs());
+ assertNoRunningJobs();
+ resetReuseBaselines();
+ }
+
+ private Set<String> listConnectorJars() throws IOException,
InterruptedException {
+ return listArtifacts(
+ "find " + CONNECTOR_DIRECTORY + " -type f -name '*.jar'
-printf '%P\\n'");
+ }
+
+ private Set<String> listConnectorJarFingerprints() throws IOException,
InterruptedException {
+ return listArtifacts(
+ "find " + CONNECTOR_DIRECTORY + " -type f -name '*.jar' -exec
sha256sum {} \\;");
+ }
+
+ private Set<String> listRuntimeLibraryFingerprints() throws IOException,
InterruptedException {
+ return listArtifacts(
+ "find "
+ + RUNTIME_LIBRARY_DIRECTORY
+ + " -type f -name '*.jar' -exec sha256sum {} \\;");
+ }
+
+ private Set<String> listTemporaryConfigs() throws IOException,
InterruptedException {
+ return listArtifacts(
+ "find /tmp -type f \\( -name '*.conf' -o -name '*.sql' \\) "
+ + "! -path '/tmp/seatunnel/*' ! -path
'/tmp/seatunnel_mnt/*' "
+ + "-printf '%P\\n'");
+ }
+
+ private Set<String> listArtifacts(String command) throws IOException,
InterruptedException {
+ Container.ExecResult result = server.execInContainer("sh", "-c",
command);
+ if (result.getExitCode() != 0) {
+ throw new IllegalStateException(
+ "Failed to inspect shared SeaTunnelContainer artifacts: "
+ result.getStderr());
+ }
+ return Arrays.stream(result.getStdout().split("\\R"))
+ .filter(name -> !name.isEmpty())
+ .collect(Collectors.toCollection(TreeSet::new));
+ }
+
+ private void deleteAddedArtifacts(
+ String directory, Set<String> baseline, Set<String>
currentArtifacts)
+ throws IOException, InterruptedException {
+ Set<String> artifactsToDelete = new TreeSet<>(currentArtifacts);
+ artifactsToDelete.removeAll(baseline);
+ Set<String> directoriesToDelete =
+ new TreeSet<>(
+ Comparator.comparingInt((String path) ->
Paths.get(path).getNameCount())
+ .reversed()
+ .thenComparing(Comparator.naturalOrder()));
+ for (String artifact : artifactsToDelete) {
+ Container.ExecResult result =
+ server.execInContainer("rm", "-f", directory + "/" +
artifact);
+ if (result.getExitCode() != 0) {
+ throw new IllegalStateException(
+ "Failed to remove shared SeaTunnelContainer artifact "
+ + artifact
+ + ": "
+ + result.getStderr());
+ }
+ Path parent = Paths.get(artifact).getParent();
+ while (parent != null) {
+ directoriesToDelete.add(parent.toString());
+ parent = parent.getParent();
+ }
+ }
+ for (String relativeDirectory : directoriesToDelete) {
+ Container.ExecResult directoryCleanup =
+ server.execInContainer(
+ "rmdir",
+ "--ignore-fail-on-non-empty",
+ directory + "/" + relativeDirectory);
+ if (directoryCleanup.getExitCode() != 0) {
+ throw new IllegalStateException(
+ "Failed to remove empty shared SeaTunnelContainer
directory "
+ + relativeDirectory
+ + ": "
+ + directoryCleanup.getStderr());
+ }
+ }
+ }
+
+ private void requireReuseBaselines() {
+ if (connectorJarsBeforeTest == null
+ || connectorJarFingerprintsBeforeTest == null
+ || runtimeLibraryFingerprintsBeforeTest == null
+ || temporaryConfigsBeforeTest == null) {
+ throw new IllegalStateException(
+ "Shared SeaTunnelContainer cleanup has no valid prepared
baseline");
+ }
+ }
+
+ private void resetReuseBaselines() {
+ connectorJarsBeforeTest = null;
+ connectorJarFingerprintsBeforeTest = null;
+ runtimeLibraryFingerprintsBeforeTest = null;
+ temporaryConfigsBeforeTest = null;
+ }
+
+ private void assertArtifactsRestored(
+ String artifactType, Set<String> expected, Set<String> actual) {
+ if (!actual.equals(expected)) {
+ throw new IllegalStateException(
+ "Shared SeaTunnelContainer did not restore "
+ + artifactType
+ + ", expected "
+ + expected
+ + " but found "
+ + actual);
+ }
+ }
+
+ private void assertVolumeEmpty() throws IOException, InterruptedException {
+ Container.ExecResult result =
+ server.execInContainer(
+ "sh",
+ "-c",
+ "if [ ! -d /tmp/seatunnel_mnt ]; then "
+ + "echo '/tmp/seatunnel_mnt is missing' >&2;
exit 1; fi; "
+ + "find /tmp/seatunnel_mnt -mindepth 1 -print
-quit");
+ if (result.getExitCode() != 0) {
+ throw new IllegalStateException(
+ "Failed to inspect shared SeaTunnelContainer volume: " +
result.getStderr());
+ }
+ if (!result.getStdout().trim().isEmpty()) {
+ throw new IllegalStateException(
+ "Shared SeaTunnelContainer volume is not empty: " +
result.getStdout());
+ }
+ }
+
+ private void assertNoRunningJobs() {
+ Awaitility.await()
+ .atMost(RUNNING_JOBS_TIMEOUT_SECONDS, TimeUnit.SECONDS)
+ .pollInterval(1, TimeUnit.SECONDS)
+ .ignoreExceptions()
+ .untilAsserted(this::assertNoRunningJobsOnce);
+ }
+
+ private void assertNoRunningJobsOnce() throws IOException {
+ HttpGet get =
+ new HttpGet(
+ String.format(
+
"http://%s:%d/hazelcast/rest/maps/running-jobs",
+ server.getHost(), server.getMappedPort(5801)));
+ RequestConfig requestConfig =
+ RequestConfig.custom()
+ .setConnectTimeout(REST_REQUEST_TIMEOUT_MILLIS)
+
.setConnectionRequestTimeout(REST_REQUEST_TIMEOUT_MILLIS)
+ .setSocketTimeout(REST_REQUEST_TIMEOUT_MILLIS)
+ .build();
+ try (CloseableHttpClient client =
+
HttpClients.custom().setDefaultRequestConfig(requestConfig).build();
+ CloseableHttpResponse response = client.execute(get)) {
+ String runningJobs =
+ response.getEntity() == null ? "" :
EntityUtils.toString(response.getEntity());
+ Assertions.assertEquals(
+ HttpStatus.SC_OK,
+ response.getStatusLine().getStatusCode(),
+ "Shared SeaTunnelContainer running-jobs request returned:
" + runningJobs);
+ Assertions.assertTrue(
+ OBJECT_MAPPER.readTree(runningJobs).isEmpty(),
+ "Shared SeaTunnelContainer still has running jobs: " +
runningJobs);
+ }
+ }
+
@Override
protected String getDockerImage() {
return JDK_DOCKER_IMAGE;
@@ -376,8 +605,14 @@ public class SeaTunnelContainer extends
AbstractTestContainer {
log.info("test in container: {}", identifier());
List<String> beforeThreads = ContainerUtil.getJVMThreadNames(server);
runningCount.incrementAndGet();
- Container.ExecResult result = executeJob(server, confFile, jobId,
variables);
- if (runningCount.decrementAndGet() > 0) {
+ Container.ExecResult result;
+ int remainingJobs;
+ try {
+ result = executeJob(server, confFile, jobId, variables);
+ } finally {
+ remainingJobs = runningCount.decrementAndGet();
+ }
+ if (remainingJobs > 0) {
// only check thread when job all finished.
return result;
}
@@ -687,14 +922,12 @@ public class SeaTunnelContainer extends
AbstractTestContainer {
public Container.ExecResult restoreJob(String confFile, String jobId,
String... variables)
throws IOException, InterruptedException {
runningCount.incrementAndGet();
- Container.ExecResult result =
- restoreJob(
- server,
- confFile,
- jobId,
- variables != null ? Arrays.asList(variables) : null);
- runningCount.decrementAndGet();
- return result;
+ try {
+ return restoreJob(
+ server, confFile, jobId, variables != null ?
Arrays.asList(variables) : null);
+ } finally {
+ runningCount.decrementAndGet();
+ }
}
@Override
@@ -702,15 +935,16 @@ public class SeaTunnelContainer extends
AbstractTestContainer {
String confFile, String jobId, String... variables)
throws IOException, InterruptedException {
runningCount.incrementAndGet();
- Container.ExecResult result =
- restoreJob(
- server,
- confFile,
- jobId,
- variables != null ? Arrays.asList(variables) : null,
- "--restore-with-checkpoint");
- runningCount.decrementAndGet();
- return result;
+ try {
+ return restoreJob(
+ server,
+ confFile,
+ jobId,
+ variables != null ? Arrays.asList(variables) : null,
+ "--restore-with-checkpoint");
+ } finally {
+ runningCount.decrementAndGet();
+ }
}
@Override
@@ -718,16 +952,12 @@ public class SeaTunnelContainer extends
AbstractTestContainer {
String confFile, String sourceJobId, String restoreJobId)
throws IOException, InterruptedException {
runningCount.incrementAndGet();
- Container.ExecResult result =
- restoreJob(
- server,
- confFile,
- sourceJobId,
- restoreJobId,
- null,
- "--restore-with-checkpoint");
- runningCount.decrementAndGet();
- return result;
+ try {
+ return restoreJob(
+ server, confFile, sourceJobId, restoreJobId, null,
"--restore-with-checkpoint");
+ } finally {
+ runningCount.decrementAndGet();
+ }
}
@Override
diff --git
a/seatunnel-e2e/seatunnel-e2e-common/src/test/java/org/apache/seatunnel/e2e/common/junit/ContainerTestingExtension.java
b/seatunnel-e2e/seatunnel-e2e-common/src/test/java/org/apache/seatunnel/e2e/common/junit/ContainerTestingExtension.java
index 164e76ff27..0dac40cdcc 100644
---
a/seatunnel-e2e/seatunnel-e2e-common/src/test/java/org/apache/seatunnel/e2e/common/junit/ContainerTestingExtension.java
+++
b/seatunnel-e2e/seatunnel-e2e-common/src/test/java/org/apache/seatunnel/e2e/common/junit/ContainerTestingExtension.java
@@ -18,7 +18,9 @@
package org.apache.seatunnel.e2e.common.junit;
import org.apache.seatunnel.e2e.common.container.ContainerExtendedFactory;
+import org.apache.seatunnel.e2e.common.container.ReusableTestContainer;
import org.apache.seatunnel.e2e.common.container.TestContainer;
+import org.apache.seatunnel.e2e.common.container.TestContainerId;
import org.apache.seatunnel.e2e.common.container.TestContainersFactory;
import org.junit.jupiter.api.extension.AfterAllCallback;
@@ -27,14 +29,22 @@ import org.junit.jupiter.api.extension.ExtensionContext;
import org.junit.platform.commons.support.AnnotationSupport;
import java.lang.annotation.Annotation;
+import java.util.ArrayList;
+import java.util.Arrays;
import java.util.Collection;
+import java.util.HashSet;
import java.util.List;
+import java.util.Set;
public class ContainerTestingExtension implements BeforeAllCallback,
AfterAllCallback {
public static final ExtensionContext.Namespace TEST_RESOURCE_NAMESPACE =
ExtensionContext.Namespace.create("testResourceNamespace");
+ private static final ExtensionContext.Namespace SHARED_CONTAINER_NAMESPACE
=
+ ExtensionContext.Namespace.create(ContainerTestingExtension.class);
public static final String TEST_CONTAINERS_STORE_KEY = "testContainers";
public static final String TEST_EXTENDED_FACTORY_STORE_KEY =
"testContainerExtendedFactory";
+ public static final String SHARED_CONTAINER_IDS_STORE_KEY =
"sharedContainerIds";
+ private static final String SHARED_CONTAINER_RESOURCES_STORE_KEY =
"sharedContainerResources";
@Override
public void beforeAll(ExtensionContext context) throws Exception {
@@ -63,12 +73,89 @@ public class ContainerTestingExtension implements
BeforeAllCallback, AfterAllCal
AnnotationUtil.filterDisabledContainers(
containersFactories.get(0).create(),
context.getRequiredTestInstance().getClass());
+ Set<TestContainerId> sharedContainerIds =
getSharedContainerIds(context);
+ List<SharedTestContainerResource> sharedResources = new ArrayList<>();
+ try {
+ // Keep the factory-defined order so classes sharing multiple
containers acquire their
+ // leases consistently and cannot create a lock-ordering cycle.
+ for (int i = 0; i < testContainers.size(); i++) {
+ TestContainer testContainer = testContainers.get(i);
+ if (!sharedContainerIds.contains(testContainer.identifier())) {
+ continue;
+ }
+ if (!(testContainer instanceof ReusableTestContainer)) {
+ throw new IllegalStateException(
+ String.format(
+ "TestContainer[%s] does not support reuse",
+ testContainer.identifier()));
+ }
+ SharedTestContainerResource sharedResource =
+ context.getRoot()
+ .getStore(SHARED_CONTAINER_NAMESPACE)
+ .getOrComputeIfAbsent(
+ testContainer.getClass().getName()
+ + ":"
+ + testContainer.identifier(),
+ ignored ->
+ new
SharedTestContainerResource(
+
(ReusableTestContainer) testContainer),
+ SharedTestContainerResource.class);
+ testContainers.set(i,
sharedResource.acquire(containerExtendedFactory));
+ sharedResources.add(sharedResource);
+ }
+ } catch (Exception acquireFailure) {
+ for (int i = sharedResources.size() - 1; i >= 0; i--) {
+ try {
+ sharedResources.get(i).release();
+ } catch (Exception releaseFailure) {
+ acquireFailure.addSuppressed(releaseFailure);
+ }
+ }
+ throw acquireFailure;
+ }
context.getStore(TEST_RESOURCE_NAMESPACE).put(TEST_CONTAINERS_STORE_KEY,
testContainers);
+ context.getStore(TEST_RESOURCE_NAMESPACE)
+ .put(SHARED_CONTAINER_IDS_STORE_KEY, sharedContainerIds);
+ context.getStore(TEST_RESOURCE_NAMESPACE)
+ .put(SHARED_CONTAINER_RESOURCES_STORE_KEY, sharedResources);
}
@Override
public void afterAll(ExtensionContext context) throws Exception {
-
context.getStore(TEST_RESOURCE_NAMESPACE).remove(TEST_CONTAINERS_STORE_KEY);
+ ExtensionContext.Store store =
context.getStore(TEST_RESOURCE_NAMESPACE);
+ @SuppressWarnings("unchecked")
+ List<SharedTestContainerResource> sharedResources =
+ (List<SharedTestContainerResource>)
+ store.remove(SHARED_CONTAINER_RESOURCES_STORE_KEY);
+ Exception cleanupFailure = null;
+ try {
+ if (sharedResources != null) {
+ for (SharedTestContainerResource sharedResource :
sharedResources) {
+ try {
+ sharedResource.release();
+ } catch (Exception e) {
+ if (cleanupFailure == null) {
+ cleanupFailure = e;
+ } else {
+ cleanupFailure.addSuppressed(e);
+ }
+ }
+ }
+ }
+ } finally {
+ store.remove(TEST_CONTAINERS_STORE_KEY);
+ store.remove(SHARED_CONTAINER_IDS_STORE_KEY);
+ }
+ if (cleanupFailure != null) {
+ throw cleanupFailure;
+ }
+ }
+
+ private Set<TestContainerId> getSharedContainerIds(ExtensionContext
context) {
+ return AnnotationSupport.findAnnotation(
+ context.getRequiredTestClass(),
ReuseTestContainers.class)
+ .map(annotation -> new
HashSet<>(Arrays.asList(annotation.value())))
+ .orElseGet(HashSet::new);
}
private void checkExactlyOneAnnotatedField(
diff --git
a/seatunnel-e2e/seatunnel-e2e-common/src/test/java/org/apache/seatunnel/e2e/common/junit/ReuseTestContainers.java
b/seatunnel-e2e/seatunnel-e2e-common/src/test/java/org/apache/seatunnel/e2e/common/junit/ReuseTestContainers.java
new file mode 100644
index 0000000000..30a46f64a2
--- /dev/null
+++
b/seatunnel-e2e/seatunnel-e2e-common/src/test/java/org/apache/seatunnel/e2e/common/junit/ReuseTestContainers.java
@@ -0,0 +1,42 @@
+/*
+ * 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
+ *
+ * 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.apache.seatunnel.e2e.common.junit;
+
+import org.apache.seatunnel.e2e.common.container.TestContainerId;
+
+import java.lang.annotation.ElementType;
+import java.lang.annotation.Inherited;
+import java.lang.annotation.Retention;
+import java.lang.annotation.RetentionPolicy;
+import java.lang.annotation.Target;
+
+/**
+ * Opts a test class into sharing selected containers with other opted-in
classes in the JVM.
+ *
+ * <p>Only use this for test classes whose container supports class-level
cleanup. Opted-in tests
+ * must not depend on another test class's files, connector artifacts, active
jobs, or finished-job
+ * history. Cleanup runs once per test class, not per test method. Test
methods in the same class
+ * share container state, so classes whose methods require a pristine
container must not opt in.
+ */
+@Target(ElementType.TYPE)
+@Retention(RetentionPolicy.RUNTIME)
+@Inherited
+public @interface ReuseTestContainers {
+
+ TestContainerId[] value();
+}
diff --git
a/seatunnel-e2e/seatunnel-e2e-common/src/test/java/org/apache/seatunnel/e2e/common/junit/SharedTestContainerResource.java
b/seatunnel-e2e/seatunnel-e2e-common/src/test/java/org/apache/seatunnel/e2e/common/junit/SharedTestContainerResource.java
new file mode 100644
index 0000000000..6b529a4980
--- /dev/null
+++
b/seatunnel-e2e/seatunnel-e2e-common/src/test/java/org/apache/seatunnel/e2e/common/junit/SharedTestContainerResource.java
@@ -0,0 +1,100 @@
+/*
+ * 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
+ *
+ * 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.apache.seatunnel.e2e.common.junit;
+
+import org.apache.seatunnel.e2e.common.container.ContainerExtendedFactory;
+import org.apache.seatunnel.e2e.common.container.ReusableTestContainer;
+
+import org.junit.jupiter.api.extension.ExtensionContext;
+
+import java.util.concurrent.Semaphore;
+
+/**
+ * Owns one reusable test container for a JUnit root store and serializes
test-class access to it. A
+ * failed preparation or cleanup invalidates the shared container so the next
acquisition starts
+ * from a fresh process.
+ */
+final class SharedTestContainerResource implements
ExtensionContext.Store.CloseableResource {
+
+ private final ReusableTestContainer container;
+ private final Semaphore classLease = new Semaphore(1, true);
+ private volatile boolean started;
+
+ SharedTestContainerResource(ReusableTestContainer container) {
+ this.container = container;
+ }
+
+ /** Acquires the class lease and prepares the shared container for one
test class. */
+ ReusableTestContainer acquire(ContainerExtendedFactory extendedFactory)
throws Exception {
+ classLease.acquire();
+ boolean acquired = false;
+ try {
+ if (!started) {
+ started = true;
+ container.startUp();
+ }
+ container.prepareForTestClass();
+ container.executeExtraCommands(extendedFactory);
+ acquired = true;
+ return container;
+ } catch (Exception acquireFailure) {
+ stopAfterFailure(acquireFailure);
+ throw acquireFailure;
+ } finally {
+ if (!acquired) {
+ classLease.release();
+ }
+ }
+ }
+
+ /** Cleans class-scoped state before releasing the lease to the next test
class. */
+ void release() throws Exception {
+ try {
+ container.cleanUpAfterTestClass();
+ } catch (Exception cleanupFailure) {
+ stopAfterFailure(cleanupFailure);
+ throw cleanupFailure;
+ } finally {
+ classLease.release();
+ }
+ }
+
+ @Override
+ public void close() throws Exception {
+ if (started) {
+ try {
+ container.tearDown();
+ } finally {
+ started = false;
+ }
+ }
+ }
+
+ private void stopAfterFailure(Exception failure) {
+ if (!started) {
+ return;
+ }
+ try {
+ container.tearDown();
+ } catch (Exception tearDownFailure) {
+ failure.addSuppressed(tearDownFailure);
+ } finally {
+ started = false;
+ }
+ }
+}
diff --git
a/seatunnel-e2e/seatunnel-e2e-common/src/test/java/org/apache/seatunnel/e2e/common/junit/SharedTestContainerResourceTest.java
b/seatunnel-e2e/seatunnel-e2e-common/src/test/java/org/apache/seatunnel/e2e/common/junit/SharedTestContainerResourceTest.java
new file mode 100644
index 0000000000..ea975e8b54
--- /dev/null
+++
b/seatunnel-e2e/seatunnel-e2e-common/src/test/java/org/apache/seatunnel/e2e/common/junit/SharedTestContainerResourceTest.java
@@ -0,0 +1,267 @@
+/*
+ * 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
+ *
+ * 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.apache.seatunnel.e2e.common.junit;
+
+import org.apache.seatunnel.e2e.common.container.ContainerExtendedFactory;
+import org.apache.seatunnel.e2e.common.container.ReusableTestContainer;
+import org.apache.seatunnel.e2e.common.container.TestContainerId;
+
+import org.junit.jupiter.api.Test;
+import org.testcontainers.containers.Container;
+
+import java.io.IOException;
+import java.util.List;
+import java.util.concurrent.CountDownLatch;
+import java.util.concurrent.ExecutorService;
+import java.util.concurrent.Executors;
+import java.util.concurrent.Future;
+import java.util.concurrent.TimeUnit;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertFalse;
+import static org.junit.jupiter.api.Assertions.assertSame;
+import static org.junit.jupiter.api.Assertions.assertThrows;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+
+class SharedTestContainerResourceTest {
+
+ private static final ContainerExtendedFactory NO_EXTENSION = container ->
{};
+
+ @Test
+ void shouldReuseContainerAndCleanEachClassLease() throws Exception {
+ CountingTestContainer container = new CountingTestContainer();
+ SharedTestContainerResource resource = new
SharedTestContainerResource(container);
+
+ assertSame(container, resource.acquire(NO_EXTENSION));
+ resource.release();
+ assertSame(container, resource.acquire(NO_EXTENSION));
+ resource.release();
+ resource.close();
+
+ assertEquals(1, container.startCount);
+ assertEquals(2, container.prepareCount);
+ assertEquals(2, container.extensionCount);
+ assertEquals(2, container.cleanupCount);
+ assertEquals(1, container.tearDownCount);
+ }
+
+ @Test
+ void shouldAllowStartupRetryAfterFailure() throws Exception {
+ CountingTestContainer container = new CountingTestContainer();
+ container.failFirstStartup = true;
+ SharedTestContainerResource resource = new
SharedTestContainerResource(container);
+
+ assertThrows(IOException.class, () -> resource.acquire(NO_EXTENSION));
+ assertSame(container, resource.acquire(NO_EXTENSION));
+ resource.release();
+ resource.close();
+
+ assertEquals(2, container.startCount);
+ assertEquals(1, container.cleanupCount);
+ assertEquals(2, container.tearDownCount);
+ }
+
+ @Test
+ void shouldRestartAfterPreparationFailure() throws Exception {
+ CountingTestContainer container = new CountingTestContainer();
+ container.failFirstPreparation = true;
+ SharedTestContainerResource resource = new
SharedTestContainerResource(container);
+
+ assertThrows(IOException.class, () -> resource.acquire(NO_EXTENSION));
+ assertSame(container, resource.acquire(NO_EXTENSION));
+ resource.release();
+ resource.close();
+
+ assertEquals(2, container.startCount);
+ assertEquals(2, container.prepareCount);
+ assertEquals(1, container.cleanupCount);
+ assertEquals(2, container.tearDownCount);
+ }
+
+ @Test
+ void shouldRestartAfterExtensionFailure() throws Exception {
+ CountingTestContainer container = new CountingTestContainer();
+ SharedTestContainerResource resource = new
SharedTestContainerResource(container);
+
+ assertThrows(
+ IOException.class,
+ () ->
+ resource.acquire(
+ ignored -> {
+ throw new IOException("extension failed");
+ }));
+ assertSame(container, resource.acquire(NO_EXTENSION));
+ resource.release();
+ resource.close();
+
+ assertEquals(2, container.startCount);
+ assertEquals(2, container.prepareCount);
+ assertEquals(2, container.extensionCount);
+ assertEquals(1, container.cleanupCount);
+ assertEquals(2, container.tearDownCount);
+ }
+
+ @Test
+ void shouldRestartAfterCleanupFailure() throws Exception {
+ CountingTestContainer container = new CountingTestContainer();
+ container.failFirstCleanup = true;
+ SharedTestContainerResource resource = new
SharedTestContainerResource(container);
+
+ assertSame(container, resource.acquire(NO_EXTENSION));
+ assertThrows(IOException.class, resource::release);
+ assertSame(container, resource.acquire(NO_EXTENSION));
+ resource.release();
+ resource.close();
+
+ assertEquals(2, container.startCount);
+ assertEquals(2, container.prepareCount);
+ assertEquals(2, container.cleanupCount);
+ assertEquals(2, container.tearDownCount);
+ }
+
+ @Test
+ void shouldSerializeClassLeases() throws Exception {
+ CountingTestContainer container = new CountingTestContainer();
+ SharedTestContainerResource resource = new
SharedTestContainerResource(container);
+ CountDownLatch firstLeaseAcquired = new CountDownLatch(1);
+ CountDownLatch releaseFirstLease = new CountDownLatch(1);
+ CountDownLatch secondAcquireStarted = new CountDownLatch(1);
+ CountDownLatch secondLeaseAcquired = new CountDownLatch(1);
+ ExecutorService executor = Executors.newFixedThreadPool(2);
+
+ try {
+ Future<?> firstLease =
+ executor.submit(
+ () -> {
+ resource.acquire(NO_EXTENSION);
+ firstLeaseAcquired.countDown();
+ releaseFirstLease.await();
+ resource.release();
+ return null;
+ });
+ assertTrue(firstLeaseAcquired.await(10, TimeUnit.SECONDS));
+
+ Future<?> secondLease =
+ executor.submit(
+ () -> {
+ secondAcquireStarted.countDown();
+ resource.acquire(NO_EXTENSION);
+ secondLeaseAcquired.countDown();
+ resource.release();
+ return null;
+ });
+
+ assertTrue(secondAcquireStarted.await(10, TimeUnit.SECONDS));
+ assertFalse(secondLeaseAcquired.await(200, TimeUnit.MILLISECONDS));
+ releaseFirstLease.countDown();
+ firstLease.get(10, TimeUnit.SECONDS);
+ secondLease.get(10, TimeUnit.SECONDS);
+ resource.close();
+
+ assertEquals(1, container.startCount);
+ assertEquals(2, container.prepareCount);
+ assertEquals(2, container.cleanupCount);
+ assertEquals(1, container.tearDownCount);
+ } finally {
+ releaseFirstLease.countDown();
+ executor.shutdownNow();
+ assertTrue(executor.awaitTermination(10, TimeUnit.SECONDS));
+ }
+ }
+
+ private static final class CountingTestContainer implements
ReusableTestContainer {
+
+ private int startCount;
+ private int prepareCount;
+ private int extensionCount;
+ private int cleanupCount;
+ private int tearDownCount;
+ private boolean failFirstStartup;
+ private boolean failFirstPreparation;
+ private boolean failFirstCleanup;
+
+ @Override
+ public void startUp() throws Exception {
+ startCount++;
+ if (failFirstStartup) {
+ failFirstStartup = false;
+ throw new IOException("startup failed");
+ }
+ }
+
+ @Override
+ public void tearDown() {
+ tearDownCount++;
+ }
+
+ @Override
+ public void prepareForTestClass() throws IOException {
+ prepareCount++;
+ if (failFirstPreparation) {
+ failFirstPreparation = false;
+ throw new IOException("preparation failed");
+ }
+ }
+
+ @Override
+ public void cleanUpAfterTestClass() throws IOException {
+ cleanupCount++;
+ if (failFirstCleanup) {
+ failFirstCleanup = false;
+ throw new IOException("cleanup failed");
+ }
+ }
+
+ @Override
+ public TestContainerId identifier() {
+ return TestContainerId.SEATUNNEL;
+ }
+
+ @Override
+ public void executeExtraCommands(ContainerExtendedFactory
extendedFactory)
+ throws IOException, InterruptedException {
+ extensionCount++;
+ extendedFactory.extend(null);
+ }
+
+ @Override
+ public Container.ExecResult executeJob(String confFile) {
+ throw new UnsupportedOperationException();
+ }
+
+ @Override
+ public Container.ExecResult executeJob(String confFile, List<String>
variables) {
+ throw new UnsupportedOperationException();
+ }
+
+ @Override
+ public String getServerLogs() {
+ throw new UnsupportedOperationException();
+ }
+
+ @Override
+ public void copyFileToContainer(String path, String targetPath) {
+ throw new UnsupportedOperationException();
+ }
+
+ @Override
+ public void copyAbsolutePathToContainer(String path, String
targetPath) {
+ throw new UnsupportedOperationException();
+ }
+ }
+}
diff --git
a/seatunnel-e2e/seatunnel-e2e-common/src/test/java/org/apache/seatunnel/e2e/common/junit/TestCaseInvocationContextProvider.java
b/seatunnel-e2e/seatunnel-e2e-common/src/test/java/org/apache/seatunnel/e2e/common/junit/TestCaseInvocationContextProvider.java
index abf8509889..552a20869c 100644
---
a/seatunnel-e2e/seatunnel-e2e-common/src/test/java/org/apache/seatunnel/e2e/common/junit/TestCaseInvocationContextProvider.java
+++
b/seatunnel-e2e/seatunnel-e2e-common/src/test/java/org/apache/seatunnel/e2e/common/junit/TestCaseInvocationContextProvider.java
@@ -19,6 +19,7 @@ package org.apache.seatunnel.e2e.common.junit;
import org.apache.seatunnel.e2e.common.container.ContainerExtendedFactory;
import org.apache.seatunnel.e2e.common.container.TestContainer;
+import org.apache.seatunnel.e2e.common.container.TestContainerId;
import org.junit.jupiter.api.extension.AfterTestExecutionCallback;
import org.junit.jupiter.api.extension.Extension;
@@ -34,8 +35,10 @@ import lombok.extern.slf4j.Slf4j;
import java.util.Arrays;
import java.util.List;
+import java.util.Set;
import java.util.stream.Stream;
+import static
org.apache.seatunnel.e2e.common.junit.ContainerTestingExtension.SHARED_CONTAINER_IDS_STORE_KEY;
import static
org.apache.seatunnel.e2e.common.junit.ContainerTestingExtension.TEST_CONTAINERS_STORE_KEY;
import static
org.apache.seatunnel.e2e.common.junit.ContainerTestingExtension.TEST_EXTENDED_FACTORY_STORE_KEY;
import static
org.apache.seatunnel.e2e.common.junit.ContainerTestingExtension.TEST_RESOURCE_NAMESPACE;
@@ -68,25 +71,35 @@ public class TestCaseInvocationContextProvider implements
TestTemplateInvocation
.get(TEST_EXTENDED_FACTORY_STORE_KEY);
int containerAmount = testContainers.size();
+ Set<TestContainerId> sharedContainerIds =
+ (Set<TestContainerId>)
+ context.getStore(TEST_RESOURCE_NAMESPACE)
+ .get(SHARED_CONTAINER_IDS_STORE_KEY);
return testContainers.stream()
.map(
testContainer ->
new TestResourceProvidingInvocationContext(
- testContainer,
containerExtendedFactory, containerAmount));
+ testContainer,
+ containerExtendedFactory,
+ containerAmount,
+
sharedContainerIds.contains(testContainer.identifier())));
}
static class TestResourceProvidingInvocationContext implements
TestTemplateInvocationContext {
private final TestContainer testContainer;
private final ContainerExtendedFactory containerExtendedFactory;
private final Integer containerAmount;
+ private final boolean shared;
public TestResourceProvidingInvocationContext(
TestContainer testContainer,
ContainerExtendedFactory containerExtendedFactory,
- int containerAmount) {
+ int containerAmount,
+ boolean shared) {
this.testContainer = testContainer;
this.containerExtendedFactory = containerExtendedFactory;
this.containerAmount = containerAmount;
+ this.shared = shared;
}
@Override
@@ -100,14 +113,16 @@ public class TestCaseInvocationContextProvider implements
TestTemplateInvocation
public List<Extension> getAdditionalExtensions() {
return Arrays.asList(
// Extension for injecting parameters
- new TestContainerResolver(testContainer,
containerExtendedFactory),
+ new TestContainerResolver(testContainer,
containerExtendedFactory, shared),
// Extension for closing test container
(AfterTestExecutionCallback)
ignore -> {
- testContainer.tearDown();
- log.info(
- "The TestContainer[{}] is closed.",
- testContainer.identifier());
+ if (!shared) {
+ testContainer.tearDown();
+ log.info(
+ "The TestContainer[{}] is closed.",
+ testContainer.identifier());
+ }
});
}
}
@@ -116,11 +131,15 @@ public class TestCaseInvocationContextProvider implements
TestTemplateInvocation
private final TestContainer testContainer;
private final ContainerExtendedFactory containerExtendedFactory;
+ private final boolean shared;
private TestContainerResolver(
- TestContainer testContainer, ContainerExtendedFactory
containerExtendedFactory) {
+ TestContainer testContainer,
+ ContainerExtendedFactory containerExtendedFactory,
+ boolean shared) {
this.testContainer = testContainer;
this.containerExtendedFactory = containerExtendedFactory;
+ this.shared = shared;
}
@Override
@@ -135,8 +154,10 @@ public class TestCaseInvocationContextProvider implements
TestTemplateInvocation
public Object resolveParameter(
ParameterContext parameterContext, ExtensionContext
extensionContext)
throws ParameterResolutionException {
- testContainer.startUp();
- testContainer.executeExtraCommands(containerExtendedFactory);
+ if (!shared) {
+ testContainer.startUp();
+ testContainer.executeExtraCommands(containerExtendedFactory);
+ }
log.info("The TestContainer[{}] is running.",
testContainer.identifier());
return this.testContainer;
}