Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
109 changes: 65 additions & 44 deletions src/testcontainers/core/container.py
Original file line number Diff line number Diff line change
Expand Up @@ -7,6 +7,7 @@
from dataclasses import dataclass
from os import PathLike
from socket import socket
from time import monotonic, sleep
from types import TracebackType
from typing import TYPE_CHECKING, Any, Optional, TypedDict, Union

Expand Down Expand Up @@ -479,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:
Expand All @@ -491,56 +493,75 @@ 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)
.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)
.start()
# 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))
)
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))

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}"
)

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())
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)

return Reaper._instance

@staticmethod
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
deadline = monotonic() + timeout
while last_exception is None or monotonic() < deadline:
s = socket()
try:
s.settimeout(reply_timeout)
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
99 changes: 99 additions & 0 deletions tests/core/test_ryuk.py
Original file line number Diff line number Diff line change
@@ -1,3 +1,5 @@
import socket
import threading
from time import perf_counter, sleep

import pytest
Expand All @@ -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

Expand Down Expand Up @@ -95,3 +98,99 @@ 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:
# 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()
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, 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
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, 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")