From 75104cd281f24c2324e4f4b11d63753076311c4b Mon Sep 17 00:00:00 2001 From: Riti Grover Date: Thu, 24 Sep 2026 17:18:17 +0530 Subject: [PATCH 1/2] fix(core): make sure Ryuk acknowledges the session filter --- src/testcontainers/core/container.py | 63 ++++++++++++++++------------ tests/core/test_ryuk.py | 50 ++++++++++++++++++++++ 2 files changed, 87 insertions(+), 26 deletions(-) diff --git a/src/testcontainers/core/container.py b/src/testcontainers/core/container.py index 76fce9c14..893809b9f 100644 --- a/src/testcontainers/core/container.py +++ b/src/testcontainers/core/container.py @@ -7,6 +7,7 @@ from dataclasses import dataclass from os import PathLike from socket import socket +from time import sleep from types import TracebackType from typing import TYPE_CHECKING, Any, Optional, TypedDict, Union @@ -498,11 +499,13 @@ def _create_instance(cls) -> "Reaper": .with_volume_mapping(c.ryuk_docker_socket, "/var/run/docker.sock", "rw") .with_kwargs(privileged=c.ryuk_privileged, auto_remove=True) .with_env("RYUK_RECONNECTION_TIMEOUT", c.ryuk_reconnection_timeout) + # Set before start(): a wait strategy only applies to a container that has not started yet. + # Ryuk 0.8 logs "Started!", newer versions log "msg=Started". + .waiting_for(LogMessageWaitStrategy(r"\bStarted\b").with_startup_timeout(20)) .start() ) rc = Reaper._container assert rc is not None - rc.waiting_for(LogMessageWaitStrategy(r".* Started!").with_startup_timeout(20)) container_host = rc.get_container_host_ip() container_port = int(rc.get_exposed_port(8080)) @@ -514,33 +517,41 @@ def _create_instance(cls) -> "Reaper": f"Could not obtain network details for {rcc.id}. Host: {container_host} Port: {container_port}" ) - last_connection_exception: Optional[Exception] = None - for _ in range(50): - try: - s = socket() - Reaper._socket = s - s.settimeout(1) - s.connect((container_host, container_port)) - last_connection_exception = None - break - except (ConnectionRefusedError, OSError) as e: - if Reaper._socket is not None: - with contextlib.suppress(Exception): - Reaper._socket.close() - Reaper._socket = None - last_connection_exception = e - - from time import sleep - - sleep(0.5) - if last_connection_exception: - raise last_connection_exception - - rs = Reaper._socket - assert rs is not None - rs.send(f"label={LABEL_SESSION_ID}={SESSION_ID}\r\n".encode()) + Reaper._socket = Reaper._connect_and_register(container_host, container_port) Reaper._instance = Reaper() atexit.register(Reaper.delete_instance) return Reaper._instance + + @staticmethod + def _connect_and_register(host: str, port: int, attempts: int = 50, retry_delay: float = 0.5) -> socket: + """Connect to Ryuk and register the session filter, retrying until Ryuk acknowledges it. + + A published port can accept a connection before Ryuk listens behind it (docker-proxy on Linux), + so a successful connect is not enough: only Ryuk's ACK shows that the filter was registered. + """ + filter_line = f"label={LABEL_SESSION_ID}={SESSION_ID}\r\n".encode() + last_exception: Optional[Exception] = None + for _ in range(attempts): + s = socket() + try: + s.settimeout(1) + s.connect((host, port)) + s.sendall(filter_line) + reply = b"" + while not reply.endswith(b"\n"): + chunk = s.recv(64) + if not chunk: + raise ConnectionResetError("Ryuk closed the connection before acknowledging the filter") + reply += chunk + if reply.strip() != b"ACK": + raise ConnectionError(f"Unexpected reply from Ryuk: {reply!r}") + return s + except OSError as e: + with contextlib.suppress(Exception): + s.close() + last_exception = e + sleep(retry_delay) + assert last_exception is not None + raise last_exception diff --git a/tests/core/test_ryuk.py b/tests/core/test_ryuk.py index 81bb736ca..33ef9732a 100644 --- a/tests/core/test_ryuk.py +++ b/tests/core/test_ryuk.py @@ -1,3 +1,5 @@ +import socket +import threading from time import perf_counter, sleep import pytest @@ -6,6 +8,7 @@ from testcontainers.core.config import testcontainers_config from testcontainers.core.container import DockerContainer, Reaper +from testcontainers.core.labels import LABEL_SESSION_ID, SESSION_ID from testcontainers.core.utils import is_mac from testcontainers.core.waiting_utils import wait_for_logs @@ -95,3 +98,50 @@ def test_ryuk_is_reused_in_same_process(): with DockerContainer("hello-world") as container: wait_for_logs(container, "Hello from Docker!") assert reaper_instance is Reaper._instance + + +def _serve(connections: list) -> tuple[str, int, list]: + """Serve one scripted behaviour per incoming connection on a local port. + + Each entry is "reset" (accept and close without reading, like docker-proxy before Ryuk is up), + "silent" (read the line but never answer) or "ack" (answer like Ryuk). + """ + server = socket.socket() + server.bind(("127.0.0.1", 0)) + server.listen() + received: list = [] + + def run() -> None: + for behaviour in connections: + conn, _ = server.accept() + with conn: + if behaviour == "reset": + continue + received.append(conn.recv(1024)) + if behaviour == "ack": + conn.sendall(b"ACK\n") + conn.recv(1024) # keep the connection open until the client closes it + else: + sleep(1.5) + server.close() + + threading.Thread(target=run, daemon=True).start() + host, port = server.getsockname() + return host, port, received + + +def test_reaper_retries_until_ryuk_acknowledges_the_filter(): + # https://github.com/testcontainers/testcontainers-python/issues/1114 + host, port, received = _serve(["reset", "silent", "ack"]) + s = Reaper._connect_and_register(host, port, attempts=5, retry_delay=0.01) + try: + assert received[-1] == f"label={LABEL_SESSION_ID}={SESSION_ID}\r\n".encode() + assert len(received) == 2 + finally: + s.close() + + +def test_reaper_raises_when_ryuk_never_acknowledges(): + host, port, _ = _serve(["reset", "reset"]) + with pytest.raises(OSError): + Reaper._connect_and_register(host, port, attempts=2, retry_delay=0.01) From bff1678e474f7ae9bb8c9767d01a41b77463d7dc Mon Sep 17 00:00:00 2001 From: Riti Grover Date: Sun, 27 Sep 2026 16:26:53 +0530 Subject: [PATCH 2/2] fix(core): remove Ryuk when setup fails and bound the registration time If start() or the registration failed, the Ryuk container kept running under its fixed name and the next start hit a name conflict. It is now removed before the error is re-raised. The registration retries now stop after a 30 s deadline instead of 50 attempts, which could take about two minutes against a silent peer. Ryuk itself exits 60 s after start without a client. Tests: the silent fake peer now waits for the client to close, which removes a timing race. New tests cover the deadline, the cleanup on both failure paths, and the wait strategy set before start(). --- src/testcontainers/core/container.py | 56 ++++++++++++++++------------ tests/core/test_ryuk.py | 55 +++++++++++++++++++++++++-- 2 files changed, 85 insertions(+), 26 deletions(-) diff --git a/src/testcontainers/core/container.py b/src/testcontainers/core/container.py index 893809b9f..9db4fe4a6 100644 --- a/src/testcontainers/core/container.py +++ b/src/testcontainers/core/container.py @@ -7,7 +7,7 @@ from dataclasses import dataclass from os import PathLike from socket import socket -from time import sleep +from time import monotonic, sleep from types import TracebackType from typing import TYPE_CHECKING, Any, Optional, TypedDict, Union @@ -480,9 +480,10 @@ def delete_instance(cls) -> None: Reaper._socket.close() Reaper._socket = None - if Reaper._container is not None and Reaper._container._container is not None: - with contextlib.suppress(docker.errors.NotFound): - Reaper._container.stop() + if Reaper._container is not None: + if Reaper._container._container is not None: + with contextlib.suppress(docker.errors.NotFound): + Reaper._container.stop() Reaper._container = None if Reaper._instance is not None: @@ -492,7 +493,9 @@ def delete_instance(cls) -> None: def _create_instance(cls) -> "Reaper": logger.debug(f"Creating new Reaper for session: {SESSION_ID}") - Reaper._container = ( + # Keep the reference before start(), so a failed start or registration can remove the container. + # Otherwise it keeps its fixed name and the next start fails with a name conflict. + rc = Reaper._container = ( DockerContainer(c.ryuk_image) .with_name(f"testcontainers-ryuk-{SESSION_ID}") .with_exposed_ports(8080) @@ -502,22 +505,25 @@ def _create_instance(cls) -> "Reaper": # Set before start(): a wait strategy only applies to a container that has not started yet. # Ryuk 0.8 logs "Started!", newer versions log "msg=Started". .waiting_for(LogMessageWaitStrategy(r"\bStarted\b").with_startup_timeout(20)) - .start() ) - rc = Reaper._container - assert rc is not None - - container_host = rc.get_container_host_ip() - container_port = int(rc.get_exposed_port(8080)) - - if not container_host or not container_port: - rcc = rc._container - assert rcc - raise ContainerConnectException( - f"Could not obtain network details for {rcc.id}. Host: {container_host} Port: {container_port}" - ) - - Reaper._socket = Reaper._connect_and_register(container_host, container_port) + try: + rc.start() + + container_host = rc.get_container_host_ip() + container_port = int(rc.get_exposed_port(8080)) + + if not container_host or not container_port: + rcc = rc._container + assert rcc + raise ContainerConnectException( + f"Could not obtain network details for {rcc.id}. Host: {container_host} Port: {container_port}" + ) + + Reaper._socket = Reaper._connect_and_register(container_host, container_port) + except BaseException: + with contextlib.suppress(Exception): + Reaper.delete_instance() + raise Reaper._instance = Reaper() atexit.register(Reaper.delete_instance) @@ -525,18 +531,22 @@ def _create_instance(cls) -> "Reaper": return Reaper._instance @staticmethod - def _connect_and_register(host: str, port: int, attempts: int = 50, retry_delay: float = 0.5) -> socket: + def _connect_and_register( + host: str, port: int, timeout: float = 30, retry_delay: float = 0.5, reply_timeout: float = 1 + ) -> socket: """Connect to Ryuk and register the session filter, retrying until Ryuk acknowledges it. A published port can accept a connection before Ryuk listens behind it (docker-proxy on Linux), so a successful connect is not enough: only Ryuk's ACK shows that the filter was registered. + The retries stop after ``timeout`` seconds, well within the 60 s that Ryuk waits for a first client. """ filter_line = f"label={LABEL_SESSION_ID}={SESSION_ID}\r\n".encode() last_exception: Optional[Exception] = None - for _ in range(attempts): + deadline = monotonic() + timeout + while last_exception is None or monotonic() < deadline: s = socket() try: - s.settimeout(1) + s.settimeout(reply_timeout) s.connect((host, port)) s.sendall(filter_line) reply = b"" diff --git a/tests/core/test_ryuk.py b/tests/core/test_ryuk.py index 33ef9732a..01500b5eb 100644 --- a/tests/core/test_ryuk.py +++ b/tests/core/test_ryuk.py @@ -122,7 +122,9 @@ def run() -> None: conn.sendall(b"ACK\n") conn.recv(1024) # keep the connection open until the client closes it else: - sleep(1.5) + # Wait for the client to time out and close, so the next connection is accepted right away. + while conn.recv(1024): + pass server.close() threading.Thread(target=run, daemon=True).start() @@ -133,7 +135,7 @@ def run() -> None: def test_reaper_retries_until_ryuk_acknowledges_the_filter(): # https://github.com/testcontainers/testcontainers-python/issues/1114 host, port, received = _serve(["reset", "silent", "ack"]) - s = Reaper._connect_and_register(host, port, attempts=5, retry_delay=0.01) + s = Reaper._connect_and_register(host, port, retry_delay=0.01, reply_timeout=0.2) try: assert received[-1] == f"label={LABEL_SESSION_ID}={SESSION_ID}\r\n".encode() assert len(received) == 2 @@ -144,4 +146,51 @@ def test_reaper_retries_until_ryuk_acknowledges_the_filter(): def test_reaper_raises_when_ryuk_never_acknowledges(): host, port, _ = _serve(["reset", "reset"]) with pytest.raises(OSError): - Reaper._connect_and_register(host, port, attempts=2, retry_delay=0.01) + Reaper._connect_and_register(host, port, timeout=0.5, retry_delay=0.01) + + +def test_reaper_gives_up_after_the_timeout(): + # A peer that accepts and never answers must not block for longer than the timeout and one attempt. + host, port, _ = _serve(["silent"] * 100) + start = perf_counter() + with pytest.raises(OSError): + Reaper._connect_and_register(host, port, timeout=0.5, retry_delay=0.01, reply_timeout=0.2) + assert perf_counter() - start < 2 + + +@pytest.mark.parametrize("fail_in", ["start", "register"]) +def test_reaper_removes_its_container_when_setup_fails(monkeypatch: pytest.MonkeyPatch, fail_in: str): + Reaper.delete_instance() + seen: dict = {} + + def fake_start(self: DockerContainer) -> DockerContainer: + seen["wait"] = self._wait_strategy + self._container = object() # set once create() has run, before the wait + if fail_in == "start": + raise TimeoutError("Ryuk did not log Started") + return self + + def fake_stop(self: DockerContainer, **kwargs: object) -> None: + seen["stopped"] = True + + def fail_register(host: str, port: int) -> socket.socket: + raise ConnectionResetError("no ACK") + + monkeypatch.setattr("testcontainers.core.container.DockerClient", lambda **kwargs: None) + monkeypatch.setattr(DockerContainer, "start", fake_start) + monkeypatch.setattr(DockerContainer, "stop", fake_stop) + monkeypatch.setattr(DockerContainer, "get_container_host_ip", lambda self: "127.0.0.1") + monkeypatch.setattr(DockerContainer, "get_exposed_port", lambda self, port: "8080") + monkeypatch.setattr(Reaper, "_connect_and_register", staticmethod(fail_register)) + + with pytest.raises((TimeoutError, ConnectionResetError)): + Reaper.get_instance() + + assert seen.get("stopped") + assert Reaper._container is None + assert Reaper._instance is None + # The wait is set before start(), and it matches the log line of old and new Ryuk versions. + pattern = seen["wait"]._message + assert pattern.search("Started!") + assert pattern.search("level=INFO msg=Started address=[::]:8080") + assert not pattern.search("Starting")