This is an automated email from the ASF dual-hosted git repository.

github-merge-queue[bot] pushed a commit to branch 
gh-readonly-queue/dev/pr-12481-7ba68325cb7ea1bf1a28f56c4f9c20695d221df3
in repository https://gitbox.apache.org/repos/asf/seatunnel.git

commit 0e4cf52896ee3604a8e7899d7736ffd914ee6ce7
Author: Goutam Adwant <[email protected]>
AuthorDate: Fri Oct 2 07:54:59 2026 +0000

    [Fix][E2E] Always stop E2E engine containers in teardown (#12481)
---
 .../common/container/AbstractTestContainer.java    | 137 +++++++++++
 .../container/AbstractTestContainerTest.java       | 254 +++++++++++++++++++++
 .../flink/AbstractTestFlinkContainer.java          |  16 +-
 .../flink/AbstractTestFlinkContainerTest.java      | 158 +++++++++++++
 .../container/seatunnel/SeaTunnelContainer.java    |   7 +-
 .../spark/AbstractTestSparkContainer.java          |   7 +-
 6 files changed, 556 insertions(+), 23 deletions(-)

diff --git 
a/seatunnel-e2e/seatunnel-e2e-common/src/test/java/org/apache/seatunnel/e2e/common/container/AbstractTestContainer.java
 
b/seatunnel-e2e/seatunnel-e2e-common/src/test/java/org/apache/seatunnel/e2e/common/container/AbstractTestContainer.java
index 8bc2dca3f4..1aea31348f 100644
--- 
a/seatunnel-e2e/seatunnel-e2e-common/src/test/java/org/apache/seatunnel/e2e/common/container/AbstractTestContainer.java
+++ 
b/seatunnel-e2e/seatunnel-e2e-common/src/test/java/org/apache/seatunnel/e2e/common/container/AbstractTestContainer.java
@@ -19,6 +19,7 @@ package org.apache.seatunnel.e2e.common.container;
 
 import org.apache.seatunnel.shade.org.apache.commons.lang3.StringUtils;
 
+import org.apache.seatunnel.common.utils.FileUtils;
 import org.apache.seatunnel.e2e.common.util.ContainerUtil;
 import org.apache.seatunnel.e2e.common.util.MavenJarUtil;
 
@@ -141,6 +142,142 @@ public abstract class AbstractTestContainer implements 
TestContainer {
                 container, this.startModuleFullPath, SEATUNNEL_HOME);
     }
 
+    /**
+     * Stops the given engine containers and then deletes {@link 
#HOST_VOLUME_MOUNT_PATH}. Every
+     * step runs even if an earlier one fails: a container left running keeps 
its network alias on
+     * the shared network, and the next test case can then talk to it instead 
of its own container.
+     * The first failure is rethrown after all steps have run; later ones are 
added to it as
+     * suppressed. If an interrupt is only recorded as suppressed, the 
thread's interrupt flag is
+     * set again before rethrowing.
+     *
+     * @param containers containers to stop in order; {@code null} entries are 
skipped
+     * @throws Exception the first failure to stop a container or to delete 
the host path
+     */
+    protected void stopContainersAndDeleteVolume(GenericContainer<?>... 
containers)
+            throws Exception {
+        Exception failure = null;
+        for (GenericContainer<?> container : containers) {
+            try {
+                stopContainer(container);
+            } catch (Exception e) {
+                failure = addFailure(failure, e);
+            }
+        }
+        try {
+            deleteHostVolumeMountPath();
+        } catch (Exception e) {
+            failure = addFailure(failure, e);
+        }
+        if (failure != null) {
+            restoreInterruptIfSuppressed(failure);
+            throw failure;
+        }
+    }
+
+    /**
+     * Removes {@link #CONTAINER_VOLUME_MOUNT_PATH} inside the container, then 
stops it. The removal
+     * is best effort: if it fails or is skipped, a warning is logged and the 
container is stopped
+     * anyway.
+     */
+    static void stopContainer(GenericContainer<?> container) throws Exception {
+        if (container == null) {
+            return;
+        }
+        Exception failure = null;
+        try {
+            removeContainerVolumeMountPath(container);
+        } catch (InterruptedException e) {
+            failure = e;
+        }
+        try {
+            container.stop();
+        } catch (Exception e) {
+            failure = addFailure(failure, e);
+        }
+        if (failure != null) {
+            throw failure;
+        }
+    }
+
+    /**
+     * Files the engine writes to the bind-mounted volume can be owned by 
root, so they are removed
+     * from inside the container before it stops. The exit code is not 
checked: {@code rm} always
+     * fails on the mount point itself ("Device or resource busy") after 
removing its contents.
+     * Anything left behind is reported by {@link #deleteHostPath(String)}.
+     *
+     * @return {@code true} if the removal ran, {@code false} if it was 
skipped or failed; both
+     *     cases are logged
+     */
+    static boolean removeContainerVolumeMountPath(GenericContainer<?> 
container)
+            throws InterruptedException {
+        try {
+            if (!container.isRunning()) {
+                LOG.warn(
+                        "Container{} {} is not running, skipping the removal 
of {} inside it",
+                        container.getNetworkAliases(),
+                        container.getContainerId(),
+                        CONTAINER_VOLUME_MOUNT_PATH);
+                return false;
+            }
+            container.execInContainer("rm", "-rf", 
CONTAINER_VOLUME_MOUNT_PATH);
+            return true;
+        } catch (InterruptedException e) {
+            throw e;
+        } catch (Exception e) {
+            // The container can stop between isRunning() and the exec.
+            LOG.warn(
+                    "Failed to remove {} inside container{} {}",
+                    CONTAINER_VOLUME_MOUNT_PATH,
+                    container.getNetworkAliases(),
+                    container.getContainerId(),
+                    e);
+        }
+        return false;
+    }
+
+    /** Deletes {@link #HOST_VOLUME_MOUNT_PATH} and logs a warning if it could 
not be deleted. */
+    protected void deleteHostVolumeMountPath() {
+        deleteHostPath(HOST_VOLUME_MOUNT_PATH);
+    }
+
+    /** @return {@code true} if the path no longer exists */
+    static boolean deleteHostPath(String path) {
+        FileUtils.deleteFile(path);
+        if (!new File(path).exists()) {
+            return true;
+        }
+        LOG.warn(
+                "Could not delete {} on the host, the next test case that 
mounts it will see the"
+                        + " files left there",
+                path);
+        return false;
+    }
+
+    /**
+     * Sets the interrupt flag again if an {@link InterruptedException} was 
only recorded as
+     * suppressed. Called after every container was stopped, because a set 
flag can make the
+     * remaining {@code stop()} calls fail.
+     */
+    private static void restoreInterruptIfSuppressed(Exception failure) {
+        if (failure instanceof InterruptedException) {
+            return;
+        }
+        for (Throwable suppressed : failure.getSuppressed()) {
+            if (suppressed instanceof InterruptedException) {
+                Thread.currentThread().interrupt();
+                return;
+            }
+        }
+    }
+
+    private static Exception addFailure(Exception first, Exception next) {
+        if (first == null) {
+            return next;
+        }
+        first.addSuppressed(next);
+        return first;
+    }
+
     protected Container.ExecResult executeJob(GenericContainer<?> container, 
String confFile)
             throws IOException, InterruptedException {
         return executeJob(container, confFile, null, null);
diff --git 
a/seatunnel-e2e/seatunnel-e2e-common/src/test/java/org/apache/seatunnel/e2e/common/container/AbstractTestContainerTest.java
 
b/seatunnel-e2e/seatunnel-e2e-common/src/test/java/org/apache/seatunnel/e2e/common/container/AbstractTestContainerTest.java
new file mode 100644
index 0000000000..a4dd4cf26f
--- /dev/null
+++ 
b/seatunnel-e2e/seatunnel-e2e-common/src/test/java/org/apache/seatunnel/e2e/common/container/AbstractTestContainerTest.java
@@ -0,0 +1,254 @@
+/*
+ * 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;
+
+import org.apache.seatunnel.e2e.common.container.seatunnel.SeaTunnelContainer;
+import org.apache.seatunnel.e2e.common.container.spark.Spark3Container;
+
+import org.junit.jupiter.api.Assertions;
+import org.junit.jupiter.api.Assumptions;
+import org.junit.jupiter.api.Test;
+import org.junit.jupiter.api.io.TempDir;
+import org.mockito.Mockito;
+import org.testcontainers.containers.Container;
+import org.testcontainers.containers.GenericContainer;
+
+import java.io.File;
+import java.nio.file.FileSystems;
+import java.nio.file.Files;
+import java.nio.file.Path;
+
+class AbstractTestContainerTest {
+
+    private static final String VOLUME = 
AbstractTestContainer.CONTAINER_VOLUME_MOUNT_PATH;
+
+    @Test
+    void shouldSkipVolumeCleanupWhenContainerIsNotRunning() throws Exception {
+        GenericContainer<?> container = stoppedContainer();
+
+        
Assertions.assertFalse(AbstractTestContainer.removeContainerVolumeMountPath(container));
+
+        Mockito.verify(container, Mockito.never()).execInContainer("rm", 
"-rf", VOLUME);
+    }
+
+    /** The container can stop between {@code isRunning()} and the exec. */
+    @Test
+    void shouldNotThrowWhenVolumeCleanupFails() throws Exception {
+        GenericContainer<?> container = runningContainer(0);
+        Mockito.when(container.execInContainer("rm", "-rf", VOLUME))
+                .thenThrow(new IllegalStateException("container is not 
running"));
+
+        
Assertions.assertFalse(AbstractTestContainer.removeContainerVolumeMountPath(container));
+    }
+
+    /**
+     * {@code rm -rf} exits with 1 on the bind-mount point itself even after 
removing its contents,
+     * so the exit code does not mean the cleanup failed.
+     */
+    @Test
+    void shouldRunVolumeCleanupOnRunningContainer() throws Exception {
+        GenericContainer<?> container = runningContainer(1);
+
+        
Assertions.assertTrue(AbstractTestContainer.removeContainerVolumeMountPath(container));
+
+        Mockito.verify(container).execInContainer("rm", "-rf", VOLUME);
+    }
+
+    @Test
+    void shouldReportHostPathThatCannotBeDeleted(@TempDir Path tempDir) throws 
Exception {
+        // On Windows a read-only directory does not stop its entries from 
being deleted.
+        Assumptions.assumeTrue(
+                
FileSystems.getDefault().supportedFileAttributeViews().contains("posix"),
+                "needs POSIX directory permissions");
+        File volume = tempDir.resolve("volume").toFile();
+        File readOnlyDir = new File(volume, "written-by-container");
+        Assertions.assertTrue(readOnlyDir.mkdirs());
+        Files.write(new File(readOnlyDir, "part-0").toPath(), new byte[] {1});
+        Assertions.assertTrue(readOnlyDir.setWritable(false));
+        try {
+            Assumptions.assumeFalse(readOnlyDir.canWrite(), "running as a user 
that ignores mode");
+
+            
Assertions.assertFalse(AbstractTestContainer.deleteHostPath(volume.getPath()));
+            Assertions.assertTrue(volume.exists());
+        } finally {
+            readOnlyDir.setWritable(true);
+        }
+        
Assertions.assertTrue(AbstractTestContainer.deleteHostPath(volume.getPath()));
+        Assertions.assertFalse(volume.exists());
+    }
+
+    @Test
+    void shouldStopSparkMasterWhenItIsNotRunning() throws Exception {
+        GenericContainer<?> master = stoppedContainer();
+        TestSparkContainer container = new TestSparkContainer(master);
+
+        Assertions.assertDoesNotThrow(container::tearDown);
+
+        Mockito.verify(master, Mockito.never()).execInContainer("rm", "-rf", 
VOLUME);
+        Mockito.verify(master).stop();
+        Assertions.assertTrue(container.hostVolumeDeleted);
+    }
+
+    @Test
+    void shouldStopSparkMasterWhenItStopsBeforeVolumeCleanup() throws 
Exception {
+        GenericContainer<?> master = runningContainer(0);
+        Mockito.when(master.execInContainer("rm", "-rf", VOLUME))
+                .thenThrow(new IllegalStateException("container is not 
running"));
+        TestSparkContainer container = new TestSparkContainer(master);
+
+        Assertions.assertDoesNotThrow(container::tearDown);
+
+        Mockito.verify(master).stop();
+        Assertions.assertTrue(container.hostVolumeDeleted);
+    }
+
+    @Test
+    void shouldCleanVolumeAndStopRunningSparkMaster() throws Exception {
+        GenericContainer<?> master = runningContainer(0);
+        TestSparkContainer container = new TestSparkContainer(master);
+
+        Assertions.assertDoesNotThrow(container::tearDown);
+
+        Mockito.verify(master).execInContainer("rm", "-rf", VOLUME);
+        Mockito.verify(master).stop();
+        Assertions.assertTrue(container.hostVolumeDeleted);
+    }
+
+    @Test
+    void shouldStopSeaTunnelServerWhenItIsNotRunning() throws Exception {
+        GenericContainer<?> server = stoppedContainer();
+        TestSeaTunnelContainer container = new TestSeaTunnelContainer(server);
+
+        Assertions.assertDoesNotThrow(container::tearDown);
+
+        Mockito.verify(server, Mockito.never()).execInContainer("rm", "-rf", 
VOLUME);
+        Mockito.verify(server).stop();
+        Assertions.assertTrue(container.hostVolumeDeleted);
+    }
+
+    @Test
+    void shouldStopSeaTunnelServerWhenItStopsBeforeVolumeCleanup() throws 
Exception {
+        GenericContainer<?> server = runningContainer(0);
+        Mockito.when(server.execInContainer("rm", "-rf", VOLUME))
+                .thenThrow(new IllegalStateException("container is not 
running"));
+        TestSeaTunnelContainer container = new TestSeaTunnelContainer(server);
+
+        Assertions.assertDoesNotThrow(container::tearDown);
+
+        Mockito.verify(server).stop();
+        Assertions.assertTrue(container.hostVolumeDeleted);
+    }
+
+    @Test
+    void shouldCleanVolumeAndStopRunningSeaTunnelServer() throws Exception {
+        GenericContainer<?> server = runningContainer(0);
+        TestSeaTunnelContainer container = new TestSeaTunnelContainer(server);
+
+        Assertions.assertDoesNotThrow(container::tearDown);
+
+        Mockito.verify(server).execInContainer("rm", "-rf", VOLUME);
+        Mockito.verify(server).stop();
+        Assertions.assertTrue(container.hostVolumeDeleted);
+    }
+
+    @Test
+    void shouldRethrowSeaTunnelServerStopFailure() {
+        GenericContainer<?> server = stoppedContainer();
+        IllegalStateException stopFailure = new IllegalStateException("stop 
failed");
+        Mockito.doThrow(stopFailure).when(server).stop();
+        TestSeaTunnelContainer container = new TestSeaTunnelContainer(server);
+
+        IllegalStateException thrown =
+                Assertions.assertThrows(IllegalStateException.class, 
container::tearDown);
+
+        Assertions.assertSame(stopFailure, thrown);
+        Assertions.assertTrue(container.hostVolumeDeleted);
+    }
+
+    /** An interrupt that is only recorded as suppressed must still reach the 
caller's thread. */
+    @Test
+    void shouldStopAllContainersAndKeepInterruptWhenItIsNotTheFirstFailure() 
throws Exception {
+        GenericContainer<?> first = stoppedContainer();
+        IllegalStateException stopFailure = new IllegalStateException("stop 
failed");
+        Mockito.doThrow(stopFailure).when(first).stop();
+        GenericContainer<?> second = runningContainer(0);
+        InterruptedException interrupt = new 
InterruptedException("interrupted");
+        Mockito.when(second.execInContainer("rm", "-rf", 
VOLUME)).thenThrow(interrupt);
+        TestSeaTunnelContainer container = new TestSeaTunnelContainer(null);
+
+        try {
+            IllegalStateException thrown =
+                    Assertions.assertThrows(
+                            IllegalStateException.class,
+                            () -> 
container.stopContainersAndDeleteVolume(first, second));
+
+            Assertions.assertSame(stopFailure, thrown);
+            Assertions.assertSame(interrupt, thrown.getSuppressed()[0]);
+            Assertions.assertTrue(Thread.currentThread().isInterrupted());
+        } finally {
+            Thread.interrupted();
+        }
+        Mockito.verify(second).stop();
+        Assertions.assertTrue(container.hostVolumeDeleted);
+    }
+
+    private static GenericContainer<?> stoppedContainer() {
+        GenericContainer<?> container = Mockito.mock(GenericContainer.class);
+        Mockito.when(container.isRunning()).thenReturn(false);
+        return container;
+    }
+
+    private static GenericContainer<?> runningContainer(int cleanupExitCode) 
throws Exception {
+        GenericContainer<?> container = Mockito.mock(GenericContainer.class);
+        Mockito.when(container.isRunning()).thenReturn(true);
+        Container.ExecResult result = Mockito.mock(Container.ExecResult.class);
+        Mockito.when(result.getExitCode()).thenReturn(cleanupExitCode);
+        Mockito.when(container.execInContainer("rm", "-rf", 
VOLUME)).thenReturn(result);
+        return container;
+    }
+
+    /** Uses a mocked master and records the host cleanup instead of deleting 
a real path. */
+    private static class TestSparkContainer extends Spark3Container {
+
+        boolean hostVolumeDeleted;
+
+        TestSparkContainer(GenericContainer<?> master) {
+            this.master = master;
+        }
+
+        @Override
+        protected void deleteHostVolumeMountPath() {
+            hostVolumeDeleted = true;
+        }
+    }
+
+    /** Uses a mocked server and records the host cleanup instead of deleting 
a real path. */
+    private static class TestSeaTunnelContainer extends SeaTunnelContainer {
+
+        boolean hostVolumeDeleted;
+
+        TestSeaTunnelContainer(GenericContainer<?> server) {
+            this.server = server;
+        }
+
+        @Override
+        protected void deleteHostVolumeMountPath() {
+            hostVolumeDeleted = true;
+        }
+    }
+}
diff --git 
a/seatunnel-e2e/seatunnel-e2e-common/src/test/java/org/apache/seatunnel/e2e/common/container/flink/AbstractTestFlinkContainer.java
 
b/seatunnel-e2e/seatunnel-e2e-common/src/test/java/org/apache/seatunnel/e2e/common/container/flink/AbstractTestFlinkContainer.java
index 3ae21d8722..88e20e6dd3 100644
--- 
a/seatunnel-e2e/seatunnel-e2e-common/src/test/java/org/apache/seatunnel/e2e/common/container/flink/AbstractTestFlinkContainer.java
+++ 
b/seatunnel-e2e/seatunnel-e2e-common/src/test/java/org/apache/seatunnel/e2e/common/container/flink/AbstractTestFlinkContainer.java
@@ -158,17 +158,11 @@ public abstract class AbstractTestFlinkContainer extends 
AbstractTestContainer {
 
     @Override
     public void tearDown() throws Exception {
-        if (taskManager != null) {
-            // delete the volume
-            taskManager.execInContainer("rm", "-rf", 
CONTAINER_VOLUME_MOUNT_PATH);
-            taskManager.stop();
-        }
-        if (jobManager != null) {
-            // delete the volume
-            jobManager.execInContainer("rm", "-rf", 
CONTAINER_VOLUME_MOUNT_PATH);
-            jobManager.stop();
-        }
-        FileUtils.deleteFile(HOST_VOLUME_MOUNT_PATH);
+        // Stop both containers even if one of them never started or the 
volume cleanup fails. A
+        // JobManager left running keeps the "jobmanager" alias on the shared 
network, and the
+        // TaskManager of the next test case can then register with it instead 
of its own
+        // JobManager, which leaves that case's job waiting for slots forever.
+        stopContainersAndDeleteVolume(taskManager, jobManager);
     }
 
     @Override
diff --git 
a/seatunnel-e2e/seatunnel-e2e-common/src/test/java/org/apache/seatunnel/e2e/common/container/flink/AbstractTestFlinkContainerTest.java
 
b/seatunnel-e2e/seatunnel-e2e-common/src/test/java/org/apache/seatunnel/e2e/common/container/flink/AbstractTestFlinkContainerTest.java
new file mode 100644
index 0000000000..77ca4e9c02
--- /dev/null
+++ 
b/seatunnel-e2e/seatunnel-e2e-common/src/test/java/org/apache/seatunnel/e2e/common/container/flink/AbstractTestFlinkContainerTest.java
@@ -0,0 +1,158 @@
+/*
+ * 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.flink;
+
+import org.junit.jupiter.api.Assertions;
+import org.junit.jupiter.api.Test;
+import org.mockito.Mockito;
+import org.testcontainers.containers.Container;
+import org.testcontainers.containers.GenericContainer;
+
+class AbstractTestFlinkContainerTest {
+
+    /**
+     * A TaskManager that failed to start must not stop the JobManager from 
being stopped. A leaked
+     * JobManager keeps the {@code jobmanager} alias on the shared network, so 
the TaskManager of
+     * the next test case can register with it and leave that case's job 
waiting for slots forever.
+     */
+    @Test
+    void shouldStopJobManagerWhenTaskManagerIsNotRunning() throws Exception {
+        GenericContainer<?> jobManager = runningContainer();
+        GenericContainer<?> taskManager = Mockito.mock(GenericContainer.class);
+        Mockito.when(taskManager.isRunning()).thenReturn(false);
+        TestFlinkContainer container = new TestFlinkContainer(jobManager, 
taskManager);
+
+        Assertions.assertDoesNotThrow(container::tearDown);
+
+        Mockito.verify(taskManager, Mockito.never())
+                .execInContainer("rm", "-rf", TestFlinkContainer.VOLUME);
+        Mockito.verify(taskManager).stop();
+        Mockito.verify(jobManager).execInContainer("rm", "-rf", 
TestFlinkContainer.VOLUME);
+        Mockito.verify(jobManager).stop();
+        Assertions.assertTrue(container.hostVolumeDeleted);
+    }
+
+    /**
+     * The TaskManager can die between {@code isRunning()} and the exec, for 
example when it runs
+     * out of metaspace. Both containers are still stopped, so the test case 
must not fail on the
+     * best-effort volume cleanup.
+     */
+    @Test
+    void shouldNotFailWhenTaskManagerStopsBeforeVolumeCleanup() throws 
Exception {
+        GenericContainer<?> jobManager = runningContainer();
+        GenericContainer<?> taskManager = runningContainer();
+        Mockito.when(taskManager.execInContainer("rm", "-rf", 
TestFlinkContainer.VOLUME))
+                .thenThrow(new IllegalStateException("container is not 
running"));
+        TestFlinkContainer container = new TestFlinkContainer(jobManager, 
taskManager);
+
+        Assertions.assertDoesNotThrow(container::tearDown);
+
+        Mockito.verify(taskManager).stop();
+        Mockito.verify(jobManager).execInContainer("rm", "-rf", 
TestFlinkContainer.VOLUME);
+        Mockito.verify(jobManager).stop();
+        Assertions.assertTrue(container.hostVolumeDeleted);
+    }
+
+    /** {@code rm -rf} exits with 1 on the bind-mount point itself, even after 
emptying it. */
+    @Test
+    void shouldNotFailWhenVolumeCleanupExitsWithError() throws Exception {
+        GenericContainer<?> jobManager = runningContainer();
+        GenericContainer<?> taskManager = runningContainer();
+        Container.ExecResult failedResult = execResult(1);
+        Mockito.when(taskManager.execInContainer("rm", "-rf", 
TestFlinkContainer.VOLUME))
+                .thenReturn(failedResult);
+        TestFlinkContainer container = new TestFlinkContainer(jobManager, 
taskManager);
+
+        Assertions.assertDoesNotThrow(container::tearDown);
+
+        Mockito.verify(taskManager).stop();
+        Mockito.verify(jobManager).stop();
+        Assertions.assertTrue(container.hostVolumeDeleted);
+    }
+
+    @Test
+    void shouldStopJobManagerAndRethrowWhenVolumeCleanupIsInterrupted() throws 
Exception {
+        GenericContainer<?> jobManager = runningContainer();
+        GenericContainer<?> taskManager = runningContainer();
+        InterruptedException interrupted = new 
InterruptedException("interrupted");
+        Mockito.when(taskManager.execInContainer("rm", "-rf", 
TestFlinkContainer.VOLUME))
+                .thenThrow(interrupted);
+        TestFlinkContainer container = new TestFlinkContainer(jobManager, 
taskManager);
+
+        InterruptedException thrown =
+                Assertions.assertThrows(InterruptedException.class, 
container::tearDown);
+
+        Assertions.assertSame(interrupted, thrown);
+        Mockito.verify(taskManager).stop();
+        Mockito.verify(jobManager).stop();
+        Assertions.assertTrue(container.hostVolumeDeleted);
+    }
+
+    @Test
+    void shouldKeepFirstFailureAndSuppressLaterOnes() throws Exception {
+        GenericContainer<?> jobManager = runningContainer();
+        GenericContainer<?> taskManager = runningContainer();
+        IllegalStateException taskManagerStopFailure = new 
IllegalStateException("stop failed");
+        IllegalStateException jobManagerStopFailure = new 
IllegalStateException("stop failed");
+        Mockito.doThrow(taskManagerStopFailure).when(taskManager).stop();
+        Mockito.doThrow(jobManagerStopFailure).when(jobManager).stop();
+        TestFlinkContainer container = new TestFlinkContainer(jobManager, 
taskManager);
+
+        IllegalStateException thrown =
+                Assertions.assertThrows(IllegalStateException.class, 
container::tearDown);
+
+        Assertions.assertSame(taskManagerStopFailure, thrown);
+        Assertions.assertArrayEquals(
+                new Throwable[] {jobManagerStopFailure}, 
thrown.getSuppressed());
+        Mockito.verify(jobManager).stop();
+        Assertions.assertTrue(container.hostVolumeDeleted);
+    }
+
+    private static GenericContainer<?> runningContainer() throws Exception {
+        GenericContainer<?> container = Mockito.mock(GenericContainer.class);
+        Mockito.when(container.isRunning()).thenReturn(true);
+        Container.ExecResult result = execResult(0);
+        Mockito.when(container.execInContainer("rm", "-rf", 
TestFlinkContainer.VOLUME))
+                .thenReturn(result);
+        return container;
+    }
+
+    private static Container.ExecResult execResult(int exitCode) {
+        Container.ExecResult result = Mockito.mock(Container.ExecResult.class);
+        Mockito.when(result.getExitCode()).thenReturn(exitCode);
+        return result;
+    }
+
+    /** Uses mocked containers and records the host cleanup instead of 
deleting a real path. */
+    private static class TestFlinkContainer extends Flink18Container {
+
+        static final String VOLUME = CONTAINER_VOLUME_MOUNT_PATH;
+
+        boolean hostVolumeDeleted;
+
+        TestFlinkContainer(GenericContainer<?> jobManager, GenericContainer<?> 
taskManager) {
+            this.jobManager = jobManager;
+            this.taskManager = taskManager;
+        }
+
+        @Override
+        protected void deleteHostVolumeMountPath() {
+            hostVolumeDeleted = true;
+        }
+    }
+}
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 4b6ac82570..b617d44c34 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
@@ -263,12 +263,7 @@ public class SeaTunnelContainer extends 
AbstractTestContainer implements Reusabl
 
     @Override
     public void tearDown() throws Exception {
-        if (server != null) {
-            // delete the volume
-            server.execInContainer("rm", "-rf", CONTAINER_VOLUME_MOUNT_PATH);
-            server.close();
-        }
-        FileUtils.deleteFile(HOST_VOLUME_MOUNT_PATH);
+        stopContainersAndDeleteVolume(server);
     }
 
     @Override
diff --git 
a/seatunnel-e2e/seatunnel-e2e-common/src/test/java/org/apache/seatunnel/e2e/common/container/spark/AbstractTestSparkContainer.java
 
b/seatunnel-e2e/seatunnel-e2e-common/src/test/java/org/apache/seatunnel/e2e/common/container/spark/AbstractTestSparkContainer.java
index ac201f2ae0..dbed35a8b9 100644
--- 
a/seatunnel-e2e/seatunnel-e2e-common/src/test/java/org/apache/seatunnel/e2e/common/container/spark/AbstractTestSparkContainer.java
+++ 
b/seatunnel-e2e/seatunnel-e2e-common/src/test/java/org/apache/seatunnel/e2e/common/container/spark/AbstractTestSparkContainer.java
@@ -86,12 +86,7 @@ public abstract class AbstractTestSparkContainer extends 
AbstractTestContainer {
 
     @Override
     public void tearDown() throws Exception {
-        if (master != null) {
-            // delete the volume
-            master.execInContainer("rm", "-rf", CONTAINER_VOLUME_MOUNT_PATH);
-            master.stop();
-        }
-        FileUtils.deleteFile(HOST_VOLUME_MOUNT_PATH);
+        stopContainersAndDeleteVolume(master);
     }
 
     @Override

Reply via email to