diff --git a/modules/kafka/src/main/java/org/testcontainers/containers/KafkaContainer.java b/modules/kafka/src/main/java/org/testcontainers/containers/KafkaContainer.java index 7eb836ade18..aa04d877186 100644 --- a/modules/kafka/src/main/java/org/testcontainers/containers/KafkaContainer.java +++ b/modules/kafka/src/main/java/org/testcontainers/containers/KafkaContainer.java @@ -6,6 +6,7 @@ import org.testcontainers.utility.ComparableVersion; import org.testcontainers.utility.DockerImageName; +import java.io.IOException; import java.util.ArrayList; import java.util.Arrays; import java.util.HashSet; @@ -45,6 +46,8 @@ public class KafkaContainer extends GenericContainer { private static final String STARTER_SCRIPT = "/tmp/testcontainers_start.sh"; + private static final String STARTER_SCRIPT_SENTINEL = STARTER_SCRIPT + ".ready"; + // https://docs.confluent.io/platform/7.0.0/release-notes/index.html#ak-raft-kraft private static final String MIN_KRAFT_TAG = "7.0.0"; @@ -200,6 +203,11 @@ protected void containerIsStarting(InspectContainerResponse containerInfo) { // Run the original command command += "/etc/confluent/docker/run \n"; copyFileToContainer(Transferable.of(command, 0777), STARTER_SCRIPT); + try { + execInContainer("touch", STARTER_SCRIPT_SENTINEL); + } catch (IOException | InterruptedException e) { + throw new RuntimeException(e); + } } protected String commandKraft() { @@ -272,7 +280,7 @@ private static class KafkaContainerDef extends ContainerDef { addExposedTcpPort(KAFKA_PORT); setEntrypoint("sh"); - setCommand("-c", "while [ ! -f " + STARTER_SCRIPT + " ]; do sleep 0.1; done; " + STARTER_SCRIPT); + setCommand("-c", "while [ ! -f " + STARTER_SCRIPT_SENTINEL + " ]; do sleep 0.1; done; " + STARTER_SCRIPT); setWaitStrategy(Wait.forLogMessage(".*\\[KafkaServer id=\\d+\\] started.*", 1)); } diff --git a/modules/kafka/src/main/java/org/testcontainers/kafka/ConfluentKafkaContainer.java b/modules/kafka/src/main/java/org/testcontainers/kafka/ConfluentKafkaContainer.java index 381ba836715..7697c4934cd 100644 --- a/modules/kafka/src/main/java/org/testcontainers/kafka/ConfluentKafkaContainer.java +++ b/modules/kafka/src/main/java/org/testcontainers/kafka/ConfluentKafkaContainer.java @@ -5,6 +5,7 @@ import org.testcontainers.images.builder.Transferable; import org.testcontainers.utility.DockerImageName; +import java.io.IOException; import java.util.ArrayList; import java.util.LinkedHashSet; import java.util.List; @@ -66,6 +67,11 @@ protected void containerIsStarting(InspectContainerResponse containerInfo) { command += "/etc/confluent/docker/run \n"; copyFileToContainer(Transferable.of(command, 0777), KafkaHelper.STARTER_SCRIPT); + try { + execInContainer("touch", KafkaHelper.STARTER_SCRIPT_SENTINEL); + } catch (IOException | InterruptedException e) { + throw new RuntimeException(e); + } } /** diff --git a/modules/kafka/src/main/java/org/testcontainers/kafka/KafkaContainer.java b/modules/kafka/src/main/java/org/testcontainers/kafka/KafkaContainer.java index 375fd132f6c..00ed8cfe688 100644 --- a/modules/kafka/src/main/java/org/testcontainers/kafka/KafkaContainer.java +++ b/modules/kafka/src/main/java/org/testcontainers/kafka/KafkaContainer.java @@ -5,6 +5,7 @@ import org.testcontainers.images.builder.Transferable; import org.testcontainers.utility.DockerImageName; +import java.io.IOException; import java.util.ArrayList; import java.util.LinkedHashSet; import java.util.List; @@ -72,6 +73,11 @@ protected void containerIsStarting(InspectContainerResponse containerInfo) { command += "/etc/kafka/docker/run \n"; copyFileToContainer(Transferable.of(command, 0777), STARTER_SCRIPT); + try { + execInContainer("touch", KafkaHelper.STARTER_SCRIPT_SENTINEL); + } catch (IOException | InterruptedException e) { + throw new RuntimeException(e); + } } /** diff --git a/modules/kafka/src/main/java/org/testcontainers/kafka/KafkaHelper.java b/modules/kafka/src/main/java/org/testcontainers/kafka/KafkaHelper.java index 61e790d474f..12a8a6eb47c 100644 --- a/modules/kafka/src/main/java/org/testcontainers/kafka/KafkaHelper.java +++ b/modules/kafka/src/main/java/org/testcontainers/kafka/KafkaHelper.java @@ -25,10 +25,12 @@ class KafkaHelper { static final String STARTER_SCRIPT = "/tmp/testcontainers_start.sh"; + static final String STARTER_SCRIPT_SENTINEL = STARTER_SCRIPT + ".ready"; + static final String[] COMMAND = { "sh", "-c", - "while [ ! -f " + STARTER_SCRIPT + " ]; do sleep 0.1; done; " + STARTER_SCRIPT, + "while [ ! -f " + STARTER_SCRIPT_SENTINEL + " ]; do sleep 0.1; done; " + STARTER_SCRIPT, }; static final WaitStrategy WAIT_STRATEGY = Wait.forLogMessage(".*Transitioning from RECOVERY to RUNNING.*", 1);