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
