diff --git a/docs/features/docker_compose.md b/docs/features/docker_compose.md index 6b874a348..1cb3a2bea 100644 --- a/docs/features/docker_compose.md +++ b/docs/features/docker_compose.md @@ -39,6 +39,43 @@ compose = DockerCompose( ) ``` +## Cleanup after process termination + +Normal context-manager exit stops the environment using Docker Compose. To also +clean up when the Python process is abruptly terminated, opt in to Ryuk: + +```python +with DockerCompose("path/to/compose/directory", ryuk=True) as compose: + # Run tests against the services. + pass +``` + +`ryuk=True` generates a unique project name for this `DockerCompose` object and +registers its `com.docker.compose.project` label with the process's shared Ryuk +container before running `compose up`. All commands use the generated project +name, including when the same object is stopped and restarted. Separate objects +use separate projects, even if they use the same Compose files. + +The generated name overrides `COMPOSE_PROJECT_NAME` and the Compose file's +top-level `name`. Hard-coded container names, published ports, and explicitly +named resources are not made unique; use a Compose file suitable for isolated, +disposable environments. External networks and volumes should be managed +separately and must not carry the generated project's cleanup label. + +With `keep_volumes=True`, leaving the context preserves the environment's volumes +for reuse by the same object. Those resources **remain eligible for Ryuk cleanup** +after the owning process loses its connection. This option does not guarantee +persistence beyond the process's lifetime when `ryuk=True`. + +`TESTCONTAINERS_RYUK_DISABLED=true` disables Ryuk startup and registration, even +with `ryuk=True`; the generated project name remains unchanged. When Ryuk is +enabled, startup or filter-registration errors prevent `compose up` from running. +Ryuk must have access to the same Docker daemon as the Compose command. + +Without `ryuk=True`, project naming and cleanup behavior are unchanged. Continue +using a context manager or calling `stop()` for normal cleanup; Ryuk is a fallback +for abrupt termination. Stopping one environment does not stop the shared Ryuk. + ## Accessing Services You can access service information and interact with containers: diff --git a/src/testcontainers/compose/compose.py b/src/testcontainers/compose/compose.py index 83d9d2e81..25d9f721c 100644 --- a/src/testcontainers/compose/compose.py +++ b/src/testcontainers/compose/compose.py @@ -11,9 +11,11 @@ from subprocess import run as subprocess_run from types import TracebackType from typing import Any, Callable, Literal, Optional, TypeVar, Union, cast +from uuid import uuid4 from typing_extensions import Self +from testcontainers.core.config import testcontainers_config from testcontainers.core.docker_client import DockerClient, get_docker_host_hostname, is_podman from testcontainers.core.exceptions import ContainerIsNotRunning, NoSuchPortExposed from testcontainers.core.inspect import ContainerInspectInfo, _ignore_properties @@ -225,6 +227,10 @@ class DockerCompose: Whether to suppress output when pulling images. quiet_build: Whether to suppress output when building images. + ryuk: + Run in an isolated project and register it for Ryuk cleanup before startup. + Defaults to False. TESTCONTAINERS_RYUK_DISABLED disables registration, + but the project remains isolated. Retained volumes are eligible for Ryuk cleanup. Example: @@ -260,10 +266,14 @@ class DockerCompose: profiles: Optional[list[str]] = None quiet_pull: bool = False quiet_build: bool = False + ryuk: bool = False + _project_name: Optional[str] = field(default=None, init=False, repr=False) _wait_strategies: Optional[dict[str, Any]] = field(default=None, init=False, repr=False) _docker_client: Optional[DockerClient] = field(default=None, init=False, repr=False) def __post_init__(self) -> None: + if self.ryuk: + self._project_name = f"testcontainers-{uuid4().hex}" if isinstance(self.compose_file_name, str): self.compose_file_name = [self.compose_file_name] if isinstance(self.env_file, str): @@ -295,6 +305,8 @@ def docker_compose_command(self) -> list[str]: def compose_command_property(self) -> list[str]: binary = self.docker_command_path or _default_compose_binary() docker_compose_cmd = [binary, "compose"] + if self._project_name is not None: + docker_compose_cmd += ["--project-name", self._project_name] if self.compose_file_name: for file in self.compose_file_name: docker_compose_cmd += ["-f", file] @@ -319,6 +331,12 @@ def start(self) -> None: """ Starts the docker compose environment. """ + if self._project_name is not None and not testcontainers_config.ryuk_disabled: + # container imports wait strategies, which in turn import compose. + from testcontainers.core.container import Reaper + + Reaper.get_instance().register_labels_filter({"com.docker.compose.project": self._project_name}) + base_cmd = self.compose_command_property or [] # pull means running a separate command before starting diff --git a/src/testcontainers/core/container.py b/src/testcontainers/core/container.py index 76fce9c14..5dc03b152 100644 --- a/src/testcontainers/core/container.py +++ b/src/testcontainers/core/container.py @@ -7,8 +7,11 @@ from dataclasses import dataclass from os import PathLike from socket import socket +from threading import Lock +from time import monotonic from types import TracebackType from typing import TYPE_CHECKING, Any, Optional, TypedDict, Union +from urllib.parse import urlencode import docker.errors from docker import version @@ -465,6 +468,47 @@ class Reaper: _instance: "Optional[Reaper]" = None _container: Optional[DockerContainer] = None _socket: Optional[socket] = None + _ACK_TIMEOUT = 10.0 + + def __init__(self) -> None: + self._filter_lock = Lock() + self._registration_failed = False + + def register_labels_filter(self, labels: dict[str, str]) -> None: + """Register an additional cleanup filter and wait for Ryuk to acknowledge it.""" + if not labels: + raise ValueError("A Ryuk cleanup filter must contain at least one label") + + with self._filter_lock: + rs = Reaper._socket + if rs is None or self._registration_failed: + raise ConnectionError("Ryuk connection is unavailable for filter registration") + + message = urlencode([("label", f"{key}={value}") for key, value in labels.items()]) + previous_timeout = rs.gettimeout() + deadline = monotonic() + self._ACK_TIMEOUT + try: + rs.settimeout(self._ACK_TIMEOUT) + rs.sendall((message + "\n").encode()) + response = b"" + while len(response) < len(b"ACK\n"): + remaining = deadline - monotonic() + if remaining <= 0: + raise TimeoutError("Timed out waiting for Ryuk to acknowledge a cleanup filter") + rs.settimeout(remaining) + chunk = rs.recv(len(b"ACK\n") - len(response)) + if not chunk: + raise ConnectionError("Ryuk disconnected before acknowledging a cleanup filter") + response += chunk + if response != b"ACK\n": + raise ConnectionError(f"Unexpected acknowledgement from Ryuk: {response!r}") + except OSError: + # A late ACK must not be mistaken for the next filter's ACK. Keep + # the socket open so existing environments are not reaped early. + self._registration_failed = True + raise + finally: + rs.settimeout(previous_timeout) @classmethod def get_instance(cls) -> "Reaper": @@ -536,11 +580,14 @@ def _create_instance(cls) -> "Reaper": 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()) + instance = Reaper() + try: + instance.register_labels_filter({LABEL_SESSION_ID: SESSION_ID}) + except OSError: + Reaper.delete_instance() + raise - Reaper._instance = Reaper() + Reaper._instance = instance atexit.register(Reaper.delete_instance) return Reaper._instance diff --git a/tests/core/compose_fixtures/ryuk/compose.yaml b/tests/core/compose_fixtures/ryuk/compose.yaml new file mode 100644 index 000000000..6208ee13b --- /dev/null +++ b/tests/core/compose_fixtures/ryuk/compose.yaml @@ -0,0 +1,21 @@ +name: testcontainers-ryuk-fixture +services: + service: + image: alpine:3.20 + init: true + command: ["sleep", "300"] + volumes: + - data:/data + - external_data:/external + networks: + - default + - external_network +volumes: + data: {} + external_data: + external: true + name: ${TC_RYUK_EXTERNAL_VOLUME} +networks: + external_network: + external: true + name: ${TC_RYUK_EXTERNAL_NETWORK} diff --git a/tests/core/compose_fixtures/ryuk/start_compose.py b/tests/core/compose_fixtures/ryuk/start_compose.py new file mode 100644 index 000000000..0a685e7e0 --- /dev/null +++ b/tests/core/compose_fixtures/ryuk/start_compose.py @@ -0,0 +1,53 @@ +"""Compose-only child process used by the abrupt-termination regression tests.""" + +import json +import sys +from pathlib import Path +from time import sleep + +from testcontainers.compose import DockerCompose +from testcontainers.core.container import Reaper +from testcontainers.core.labels import SESSION_ID +from testcontainers.core.waiting_utils import WaitStrategy + + +class FailedReadiness(WaitStrategy): + def wait_until_ready(self, container) -> None: + raise TimeoutError("Simulated failure after Compose created resources") + + +def main() -> None: + ready = Path(sys.argv[1]) + scenario = sys.argv[2] + compose = DockerCompose(Path(__file__).parent, ryuk=True, keep_volumes=True) + ready.with_suffix(".state.json").write_text( + json.dumps({"project": compose._project_name, "ryuk": f"testcontainers-ryuk-{SESSION_ID}"}) + ) + if scenario == "partial": + compose.waiting_for({"service": FailedReadiness()}) + try: + compose.start() + except TimeoutError: + pass + else: + raise AssertionError("Readiness should have failed") + else: + compose.start() + if scenario == "retained": + compose.exec_in_container(["sh", "-c", "echo retained > /data/value"]) + compose.stop(down=False) + with compose: + assert compose.exec_in_container(["cat", "/data/value"])[0].strip() == "retained" + + container = compose.get_container(include_all=True) + assert Reaper._container is not None + state = {"project": container.Project, "ryuk": Reaper._container.get_container_id()} + temporary = ready.with_suffix(".tmp") + temporary.write_text(json.dumps(state)) + temporary.replace(ready) + while True: + sleep(1) + + +if __name__ == "__main__": + main() diff --git a/tests/core/test_compose_ryuk.py b/tests/core/test_compose_ryuk.py new file mode 100644 index 000000000..954439d24 --- /dev/null +++ b/tests/core/test_compose_ryuk.py @@ -0,0 +1,282 @@ +import json +import os +import subprocess +import sys +from contextlib import ExitStack, suppress +from pathlib import Path +from time import monotonic, sleep +from uuid import uuid4 + +import docker +import pytest +from docker.errors import DockerException, NotFound +from pytest_mock import MockerFixture + +from testcontainers.compose import DockerCompose +from testcontainers.core.config import testcontainers_config +from testcontainers.core.container import DockerContainer, Reaper +from testcontainers.core.docker_client import DockerClient +from testcontainers.core.utils import is_mac + +FIXTURES = Path(__file__).parent / "compose_fixtures" + +_skip_if_mac_ryuk = pytest.mark.skipif( + is_mac(), + reason="Ryuk startup and cleanup are unreliable on Docker Desktop for macOS", +) + + +@pytest.mark.parametrize("enabled,disabled", [(False, False), (True, False), (True, True)]) +def test_compose_ryuk_registration_and_command_identity( + mocker: MockerFixture, monkeypatch: pytest.MonkeyPatch, enabled: bool, disabled: bool +): + monkeypatch.setattr(testcontainers_config, "ryuk_disabled", disabled) + get_reaper = mocker.patch.object(Reaper, "get_instance") + delete_reaper = mocker.patch.object(Reaper, "delete_instance") + compose = DockerCompose( + ".", + ryuk=enabled, + docker_command_path="docker", + pull=True, + compose_file_name=["compose.yaml", "override.yaml"], + profiles=["test"], + env_file="test.env", + services=["database"], + ) + events = [] + get_reaper.return_value.register_labels_filter.side_effect = lambda labels: events.append(("filter", labels)) + run = mocker.patch.object( + compose, "_run_command", side_effect=lambda **kwargs: events.append(("command", kwargs["cmd"])) + ) + base = compose.docker_compose_command()[:] + compose.start() + compose.stop() + compose.start() + + assert compose.docker_compose_command() == base + assert [call.kwargs["cmd"] for call in run.call_args_list] == [ + [*base, "pull"], + [*base, "up", "--wait", "database"], + [*base, "down", "--volumes", "database"], + [*base, "pull"], + [*base, "up", "--wait", "database"], + ] + if enabled: + assert base[:3] == ["docker", "compose", "--project-name"] + project = base[3] + assert DockerCompose(".", ryuk=True, docker_command_path="docker").docker_compose_command()[3] != project + else: + assert "--project-name" not in base + if enabled and not disabled: + assert events[0] == ("filter", {"com.docker.compose.project": project}) + assert get_reaper.call_count == 2 + else: + get_reaper.assert_not_called() + delete_reaper.assert_not_called() + + +@pytest.mark.parametrize("fail_start", [False, True]) +def test_ryuk_failure_prevents_compose_up(mocker: MockerFixture, monkeypatch: pytest.MonkeyPatch, fail_start: bool): + monkeypatch.setattr(testcontainers_config, "ryuk_disabled", False) + get_reaper = mocker.patch.object(Reaper, "get_instance") + failing_call = get_reaper if fail_start else get_reaper.return_value.register_labels_filter + failing_call.side_effect = ConnectionError("Ryuk unavailable") + compose = DockerCompose(".", ryuk=True, docker_command_path="docker") + run = mocker.patch.object(compose, "_run_command") + with pytest.raises(ConnectionError, match="Ryuk unavailable"): + compose.start() + run.assert_not_called() + + +@pytest.mark.parametrize("ryuk_removed", [False, True]) +def test_cleanup_timeout_reports_remaining_resources(mocker: MockerFixture, ryuk_removed: bool): + client = mocker.Mock(spec=docker.DockerClient) + container = mocker.Mock(id="container-id", status="running") + container.name = "service" + network = mocker.Mock(id="network-id") + network.name = "project-network" + volume = mocker.Mock() + volume.name = "project-volume" + client.containers.list.return_value = [container] + client.networks.list.return_value = [network] + client.volumes.list.return_value = [volume] + if ryuk_removed: + client.containers.get.side_effect = NotFound("Ryuk already removed") + else: + client.containers.get.return_value.logs.return_value = b"cleanup failed" + + with pytest.raises(TimeoutError) as exc: + _wait_for_project_removed(client, {"project": "test-project", "ryuk": "ryuk-id"}, timeout=0) + + message = str(exc.value) + for detail in ("test-project", "container-id", "service", "running", "network-id", "project-volume"): + assert detail in message + assert ("Unavailable:" if ryuk_removed else "cleanup failed") in message + + +@pytest.fixture +def compose_ryuk_resources(monkeypatch: pytest.MonkeyPatch): + with ExitStack() as cleanup: + client = DockerClient().client + cleanup.callback(client.close) + suffix = uuid4().hex + # Use the configured SDK client without adding Testcontainers ownership labels. + volume = client.volumes.create(name=f"tc-ryuk-external-{suffix}") + cleanup.callback(volume.remove) + network = client.networks.create(name=f"tc-ryuk-external-{suffix}") + cleanup.callback(network.remove) + assert network.name is not None + monkeypatch.setenv("TC_RYUK_EXTERNAL_VOLUME", volume.name) + monkeypatch.setenv("TC_RYUK_EXTERNAL_NETWORK", network.name) + monkeypatch.setenv("COMPOSE_PROJECT_NAME", f"tc-ryuk-unrelated-{suffix}") + yield client, volume, network + + +@pytest.fixture +def fresh_reaper(): + Reaper.delete_instance() + try: + yield + finally: + Reaper.delete_instance() + + +@_skip_if_mac_ryuk +def test_compose_ryuk_isolation_reuse_and_retention(compose_ryuk_resources, monkeypatch: pytest.MonkeyPatch): + monkeypatch.setattr(testcontainers_config, "ryuk_disabled", False) + client, external_volume, external_network = compose_ryuk_resources + first = DockerCompose(FIXTURES / "ryuk", ryuk=True, keep_volumes=True) + second = DockerCompose(FIXTURES / "ryuk", ryuk=True) + try: + with first: + reaper = Reaper.get_instance() + first_id = first.get_container().ID + first_project = first.get_container().Project + first.exec_in_container(["sh", "-c", "echo retained > /data/value"]) + with second: + assert Reaper.get_instance() is reaper + assert second.get_container().Project != first_project + assert second.get_container().ID != first_id + assert first.get_container().ID == first_id + assert Reaper.get_instance() is reaper + + # Reusing the same object retains its project and its data. + with first: + assert first.get_container().Project == first_project + assert first.exec_in_container(["cat", "/data/value"])[0].strip() == "retained" + first.stop() + assert not client.containers.list(all=True, filters={"label": f"com.docker.compose.project={first_project}"}) + assert not client.volumes.list(filters={"label": f"com.docker.compose.project={first_project}"}) + assert client.volumes.get(external_volume.name) + assert client.networks.get(external_network.id) + finally: + first.stop() + second.stop() + + +@_skip_if_mac_ryuk +def test_compose_reuses_reaper_started_by_docker_container( + fresh_reaper, compose_ryuk_resources, monkeypatch: pytest.MonkeyPatch +): + monkeypatch.setattr(testcontainers_config, "ryuk_disabled", False) + assert Reaper._instance is None + with DockerContainer("alpine:3.20", command="sleep 300") as ordinary: + reaper = Reaper._instance + assert reaper is not None + with DockerCompose(FIXTURES / "ryuk", ryuk=True): + assert Reaper.get_instance() is reaper + container = ordinary.get_wrapped_container() + container.reload() + assert container.status == "running" + assert Reaper.get_instance() is reaper + + +@_skip_if_mac_ryuk +@pytest.mark.parametrize("scenario", ["running", "retained", "partial"]) +def test_compose_ryuk_cleans_up_after_process_kill(compose_ryuk_resources, tmp_path: Path, scenario: str): + client, external_volume, external_network = compose_ryuk_resources + unrelated = DockerCompose(FIXTURES / "ryuk") + ready = tmp_path / "ready.json" + child = None + state = None + env = dict(os.environ, TESTCONTAINERS_RYUK_DISABLED="false", RYUK_RECONNECTION_TIMEOUT="1s") + try: + with unrelated: + unrelated_id = unrelated.get_container().ID + with (tmp_path / "child.log").open("w+") as log: + child = subprocess.Popen( + [sys.executable, str(FIXTURES / "ryuk" / "start_compose.py"), str(ready), scenario], + env=env, + stdout=log, + stderr=subprocess.STDOUT, + ) + deadline = monotonic() + 120 + while not ready.exists(): + if child.poll() is not None or monotonic() >= deadline: + log.seek(0) + pytest.fail(f"Compose child did not become ready:\n{log.read()}") + sleep(0.1) + state = json.loads(ready.read_text()) + filters = {"label": f"com.docker.compose.project={state['project']}"} + assert client.containers.list(all=True, filters=filters) + assert client.networks.list(filters=filters) + assert client.volumes.list(filters=filters) + assert client.containers.get(state["ryuk"]) + child.kill() + child.wait(timeout=10) + + _wait_for_project_removed(client, state) + + assert client.containers.get(unrelated_id).status == "running" + assert client.volumes.get(external_volume.name) + assert client.networks.get(external_network.id) + finally: + if child is not None and child.poll() is None: + child.kill() + child.wait(timeout=10) + if state is None and ready.exists(): + state = json.loads(ready.read_text()) + if state is None and ready.with_suffix(".state.json").exists(): + state = json.loads(ready.with_suffix(".state.json").read_text()) + if state is not None: + _remove_child_resources(client, state) + + +def _wait_for_project_removed(client: docker.DockerClient, state: dict[str, str], timeout: float = 60) -> None: + filters: dict[str, str | list[str] | bool] = {"label": f"com.docker.compose.project={state['project']}"} + deadline = monotonic() + timeout + while True: + containers = client.containers.list(all=True, filters=filters) + networks = client.networks.list(filters=filters) + volumes = client.volumes.list(filters=filters) + if not (containers or networks or volumes): + return + if monotonic() >= deadline: + try: + ryuk_logs = client.containers.get(state["ryuk"]).logs(tail=100).decode("utf-8", errors="replace") + except DockerException as exc: + ryuk_logs = f"Unavailable: {exc}" + raise TimeoutError( + f"Ryuk did not clean up project {state['project']} within {timeout}s.\n" + f"Containers (id, name, status): {[(c.id, c.name, c.status) for c in containers]}\n" + f"Networks (id, name): {[(n.id, n.name) for n in networks]}\n" + f"Volumes: {[v.name for v in volumes]}\n" + f"Ryuk logs:\n{ryuk_logs}" + ) + sleep(0.2) + + +def _remove_child_resources(client: docker.DockerClient, state: dict[str, str]) -> None: + # Only remove resources belonging to this child's generated project. + filters: dict[str, str | list[str] | bool] = {"label": f"com.docker.compose.project={state['project']}"} + for container in client.containers.list(all=True, filters=filters): + with suppress(NotFound): + container.remove(force=True, v=True) + for network in client.networks.list(filters=filters): + with suppress(NotFound): + network.remove() + for volume in client.volumes.list(filters=filters): + with suppress(NotFound): + volume.remove() + with suppress(NotFound): + client.containers.get(state["ryuk"]).remove(force=True) diff --git a/tests/core/test_reaper_filters.py b/tests/core/test_reaper_filters.py new file mode 100644 index 000000000..353f893a7 --- /dev/null +++ b/tests/core/test_reaper_filters.py @@ -0,0 +1,122 @@ +from socket import socket, socketpair +from urllib.parse import parse_qs + +import pytest +from pytest_mock import MockerFixture + +from testcontainers.core.container import Reaper +from testcontainers.core.labels import LABEL_SESSION_ID, SESSION_ID + + +def test_register_filters_handles_fragmented_acknowledgements(mocker: MockerFixture): + connection = mocker.Mock(spec=socket) + connection.gettimeout.return_value = 1.0 + connection.recv.side_effect = [b"A", b"C", b"K\n", b"ACK\n"] + mocker.patch.object(Reaper, "_socket", connection) + reaper = Reaper() + + reaper.register_labels_filter({LABEL_SESSION_ID: SESSION_ID}) + labels = {"com.docker.compose.project": "test-project", "custom": "spaces & equals= and +"} + reaper.register_labels_filter(labels) + + messages = [call.args[0].decode() for call in connection.sendall.call_args_list] + assert all(message.endswith("\n") for message in messages) + assert parse_qs(messages[0].strip()) == {"label": [f"{LABEL_SESSION_ID}={SESSION_ID}"]} + assert parse_qs(messages[1].strip()) == {"label": [f"{key}={value}" for key, value in labels.items()]} + assert connection.settimeout.call_args.args == (1.0,) + + +@pytest.mark.parametrize("response", [b"", b"NOPE", TimeoutError("no acknowledgement")]) +def test_failed_registration_does_not_reuse_a_late_ack_or_close_the_session(mocker: MockerFixture, response): + connection = mocker.Mock(spec=socket) + connection.gettimeout.return_value = 1.0 + connection.recv.side_effect = [response, b"ACK\n"] + mocker.patch.object(Reaper, "_socket", connection) + reaper = Reaper() + + with pytest.raises(OSError): + reaper.register_labels_filter({"project": "first"}) + with pytest.raises(ConnectionError, match="unavailable"): + reaper.register_labels_filter({"project": "second"}) + + assert connection.sendall.call_count == 1 + connection.close.assert_not_called() + assert connection.settimeout.call_args.args == (1.0,) + + +def test_registration_has_an_overall_deadline(mocker: MockerFixture): + connection = mocker.Mock(spec=socket) + connection.gettimeout.return_value = 1.0 + connection.recv.return_value = b"A" + mocker.patch.object(Reaper, "_socket", connection) + mocker.patch("testcontainers.core.container.monotonic", side_effect=[0, 1, 11]) + + with pytest.raises(TimeoutError, match="Timed out"): + Reaper().register_labels_filter({"project": "slow"}) + assert connection.recv.call_count == 1 + + +def test_empty_filter_is_rejected(mocker: MockerFixture): + connection = mocker.Mock(spec=socket) + mocker.patch.object(Reaper, "_socket", connection) + with pytest.raises(ValueError, match="at least one label"): + Reaper().register_labels_filter({}) + connection.sendall.assert_not_called() + + +def test_registration_requires_a_connection(mocker: MockerFixture): + mocker.patch.object(Reaper, "_socket", None) + with pytest.raises(ConnectionError, match="unavailable"): + Reaper().register_labels_filter({"project": "missing"}) + + +def test_unresponsive_peer_times_out_without_closing_the_connection(monkeypatch: pytest.MonkeyPatch): + connection, peer = socketpair() + with connection, peer: + connection.settimeout(1) + monkeypatch.setattr(Reaper, "_socket", connection) + monkeypatch.setattr(Reaper, "_ACK_TIMEOUT", 0.05) + reaper = Reaper() + with pytest.raises(TimeoutError): + reaper.register_labels_filter({"project": "unacknowledged"}) + assert connection.gettimeout() == 1 + assert connection.fileno() != -1 + peer.sendall(b"ACK\n") + with pytest.raises(ConnectionError, match="unavailable"): + reaper.register_labels_filter({"project": "another"}) + + +@pytest.mark.parametrize("failure", [False, True]) +def test_initial_session_registration_is_acknowledged(mocker: MockerFixture, failure: bool): + container = mocker.MagicMock() + container.get_container_host_ip.return_value = "localhost" + container.get_exposed_port.return_value = 12345 + constructor = mocker.patch("testcontainers.core.container.DockerContainer") + builder = constructor.return_value + for method in ("with_name", "with_exposed_ports", "with_volume_mapping", "with_kwargs", "with_env"): + getattr(builder, method).return_value = builder + builder.start.return_value = container + connection = mocker.patch("testcontainers.core.container.socket").return_value + connection.gettimeout.return_value = 1.0 + connection.recv.side_effect = [b""] if failure else [b"ACK\n"] + mocker.patch.object(Reaper, "_container", None) + mocker.patch.object(Reaper, "_instance", None) + mocker.patch.object(Reaper, "_socket", None) + register_exit = mocker.patch("testcontainers.core.container.atexit").register + + if failure: + with pytest.raises(ConnectionError): + Reaper.get_instance() + assert Reaper._instance is None + assert Reaper._socket is None + container.stop.assert_called_once() + register_exit.assert_not_called() + else: + instance = Reaper.get_instance() + assert instance is Reaper.get_instance() + container.stop.assert_not_called() + register_exit.assert_called_once_with(Reaper.delete_instance) + + message = connection.sendall.call_args.args[0].decode().strip() + assert parse_qs(message) == {"label": [f"{LABEL_SESSION_ID}={SESSION_ID}"]} + connection.recv.assert_called_once()