From 07e51805e387611bb145c19621da2914b24613af Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Kamil=20Chodo=C5=82a?= Date: Sun, 13 Sep 2026 11:58:57 +0200 Subject: [PATCH 1/2] Close RPC consumers before client shutdown --- src/expb/payloads/executor/executor.py | 26 +++++----- tests/payloads/test_executor_cleanup.py | 63 +++++++++++++++++++++++++ 2 files changed, 77 insertions(+), 12 deletions(-) create mode 100644 tests/payloads/test_executor_cleanup.py diff --git a/src/expb/payloads/executor/executor.py b/src/expb/payloads/executor/executor.py index 54d60b2..2a30525 100644 --- a/src/expb/payloads/executor/executor.py +++ b/src/expb/payloads/executor/executor.py @@ -1665,6 +1665,20 @@ def _collect_k6_metric(decoded_line: str) -> None: line_callback=_collect_k6_metric, ) + # Stop RPC-consuming infrastructure before the client so its connections do not + # delay the client's graceful shutdown (the payload server also owns the SSE stream). + self._teardown_container( + self.config.get_payload_server_container_name(), + log_file=self.config.outputs_dir / "payload-server.log", + stop_timeout=3, + print_console=print_logs_to_console, + ) + + self._teardown_container( + self.config.get_alloy_container_name(), + stop_timeout=3, + ) + # Stop the collector while the client runtime is still alive: the stop request makes # the runtime emit its method rundown, without which the .nettrace stacks do not resolve. self._stop_dotnet_trace_collector() @@ -1703,18 +1717,6 @@ def _collect_k6_metric(decoded_line: str) -> None: if print_logs_to_console and print_per_payload_metrics_table: self._print_per_payload_metrics_table(per_payload_metrics_rows) - self._teardown_container( - self.config.get_payload_server_container_name(), - log_file=self.config.outputs_dir / "payload-server.log", - stop_timeout=3, - print_console=print_logs_to_console, - ) - - self._teardown_container( - self.config.get_alloy_container_name(), - stop_timeout=3, - ) - # Clean docker network try: containers_network = self.config.docker_client.networks.get( diff --git a/tests/payloads/test_executor_cleanup.py b/tests/payloads/test_executor_cleanup.py new file mode 100644 index 0000000..7a5ca64 --- /dev/null +++ b/tests/payloads/test_executor_cleanup.py @@ -0,0 +1,63 @@ +from pathlib import Path +from unittest.mock import Mock, patch + +from expb.payloads.executor.executor import Executor + + +def test_cleanup_stops_rpc_consumers_before_client_and_preserves_metrics(tmp_path: Path): + config = Mock( + executor_name="scenario", + outputs_dir=tmp_path, + docker_client=Mock(), + ) + config.get_k6_container_name.return_value = "k6" + config.get_payload_server_container_name.return_value = "payload-server" + config.get_alloy_container_name.return_value = "alloy" + config.get_execution_client_container_name.return_value = "nethermind" + config.get_execution_client_name.return_value = "nethermind" + config.docker_client.networks.get.return_value = Mock() + + executor = Executor(config=config, logger=Mock()) + lifecycle: list[str] = [] + + def teardown(name: str, **kwargs): + lifecycle.append(name) + if name == "k6": + kwargs["line_callback"]( + 'EXPB_PER_PAYLOAD_METRIC idx=7 gas_used=123 processing_ms=45.6\n' + ) + return [] + + def stop_trace(): + lifecycle.append("eventpipe") + + def finalize_perf(_container): + lifecycle.append("perf") + + teardown_mock = Mock(side_effect=teardown) + with ( + patch.object(executor, "stop_extra_commands"), + patch.object(executor, "_teardown_container", teardown_mock), + patch.object(executor, "_finalize_perf", side_effect=finalize_perf), + patch.object(executor, "_stop_dotnet_trace_collector", side_effect=stop_trace), + patch.object(executor, "_print_per_payload_metrics_table") as print_table, + patch.object(executor, "remove_directories"), + ): + executor.cleanup_scenario( + print_logs_to_console=True, + print_per_payload_metrics_table=True, + ) + + assert lifecycle == [ + "perf", + "k6", + "payload-server", + "alloy", + "eventpipe", + "nethermind", + ] + teardown_calls = teardown_mock.call_args_list + assert teardown_calls[0].kwargs["log_file"] == tmp_path / "k6.log" + assert teardown_calls[0].kwargs["line_callback"] is not None + assert teardown_calls[1].kwargs["log_file"] == tmp_path / "payload-server.log" + print_table.assert_called_once_with([(7, "123", "45.6")]) From 4a7ef676493fedfc4973c2d3420442d71dedd255 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Kamil=20Chodo=C5=82a?= Date: Sun, 13 Sep 2026 12:06:26 +0200 Subject: [PATCH 2/2] Allow explicit payload skip override for snapshot alignment --- docs/USAGE.md | 6 ++ src/expb/payloads/executor/executor_config.py | 11 +++ tests/test_payload_skip_override.py | 90 +++++++++++++++++++ 3 files changed, 107 insertions(+) create mode 100644 tests/test_payload_skip_override.py diff --git a/docs/USAGE.md b/docs/USAGE.md index 9e29973..8d150b1 100644 --- a/docs/USAGE.md +++ b/docs/USAGE.md @@ -105,6 +105,12 @@ expb execute-scenarios [OPTIONS] * `--print-logs / --no-print-logs`: Print K6 and Execution Client logs to console. [default: no-print-logs] * `--help`: Show this message and exit. +`EXPB_SKIP_OVERRIDE` sets the number of payloads skipped before the configured +warmup and measured payloads. It must be a nonnegative integer; for example, +`EXPB_SKIP_OVERRIDE=10 EXPB_WARMUP_OVERRIDE=0` starts measurement at the +eleventh payload without changing the scenario file. `EXPB_WARMUP_OVERRIDE` +continues to override the number of unmeasured warmup payloads. + ## `expb compress-payloads` Compress execution payloads txs for a given block range into bigger blocks. diff --git a/src/expb/payloads/executor/executor_config.py b/src/expb/payloads/executor/executor_config.py index 43b6e30..2f2bc4a 100644 --- a/src/expb/payloads/executor/executor_config.py +++ b/src/expb/payloads/executor/executor_config.py @@ -101,6 +101,17 @@ def __init__( self.k6_warmup_wait: int = scenario.warmup_wait self.k6_payloads_skip: int | None = scenario.payloads_skip self.k6_payloads_warmup: int | None = scenario.payloads_warmup + # EXPB_SKIP_OVERRIDE changes the number of payloads skipped before the + # configured warmup and measured payloads, without editing the scenario. + _skip_override = os.environ.get("EXPB_SKIP_OVERRIDE") + if _skip_override is not None: + if not _skip_override.isdecimal(): + raise ValueError( + "EXPB_SKIP_OVERRIDE must be a nonnegative integer, " + f"got '{_skip_override}'" + ) + self.k6_payloads_skip = int(_skip_override) + # EXPB_WARMUP_OVERRIDE overrides the number of unmeasured warmup payloads # without editing the scenario config. A larger warmup lets the OS page # cache and client caches reach a warm steady state before measurement. diff --git a/tests/test_payload_skip_override.py b/tests/test_payload_skip_override.py new file mode 100644 index 0000000..faa2363 --- /dev/null +++ b/tests/test_payload_skip_override.py @@ -0,0 +1,90 @@ +from pathlib import Path +from unittest.mock import Mock + +import pytest + +from expb.configs.scenarios import Scenario, ScenariosPaths +from expb.payloads.executor.executor_config import ExecutorConfig + + +@pytest.fixture +def scenario_inputs(tmp_path: Path) -> tuple[Scenario, ScenariosPaths]: + payloads_file = tmp_path / "payloads.jsonl" + fcus_file = tmp_path / "fcus.jsonl" + payloads_file.touch() + fcus_file.touch() + scenario = Scenario( + name="skip-override", + client="nethermind", + payloads=payloads_file, + fcus=fcus_file, + snapshot_source="snapshot", + skip=7, + amount=100, + warmup=3, + ) + return scenario, ScenariosPaths(work=tmp_path / "work", outputs=tmp_path / "outputs") + + +def create_config( + monkeypatch: pytest.MonkeyPatch, + scenario_inputs: tuple[Scenario, ScenariosPaths], +) -> ExecutorConfig: + scenario, paths = scenario_inputs + monkeypatch.setattr("expb.payloads.executor.executor_config.docker.from_env", Mock()) + monkeypatch.setattr( + "expb.payloads.executor.executor_config.os.getuid", lambda: 1000, raising=False + ) + monkeypatch.setattr( + "expb.payloads.executor.executor_config.os.getgid", lambda: 1000, raising=False + ) + return ExecutorConfig( + scenario=scenario, + snapshot_service=Mock(), + paths=paths, + ) + + +def test_skip_override_preserves_configured_skip_and_downstream_totals( + monkeypatch: pytest.MonkeyPatch, + scenario_inputs: tuple[Scenario, ScenariosPaths], +): + monkeypatch.delenv("EXPB_SKIP_OVERRIDE", raising=False) + + config = create_config(monkeypatch, scenario_inputs) + + assert config.k6_payloads_skip == 7 + environment = config.get_payload_server_environment() + assert environment["EXPB_SKIP"] == "7" + assert environment["EXPB_TOTAL"] == "103" + + +@pytest.mark.parametrize("override", ["0", "11"]) +def test_skip_override_is_used_by_payload_server( + monkeypatch: pytest.MonkeyPatch, + scenario_inputs: tuple[Scenario, ScenariosPaths], + override: str, +): + monkeypatch.setenv("EXPB_SKIP_OVERRIDE", override) + + config = create_config(monkeypatch, scenario_inputs) + + assert config.k6_payloads_skip == int(override) + environment = config.get_payload_server_environment() + assert environment["EXPB_SKIP"] == override + assert environment["EXPB_TOTAL"] == "103" + + +@pytest.mark.parametrize("override", ["", "-1", "1.5", "many"]) +def test_skip_override_rejects_invalid_values( + monkeypatch: pytest.MonkeyPatch, + scenario_inputs: tuple[Scenario, ScenariosPaths], + override: str, +): + monkeypatch.setenv("EXPB_SKIP_OVERRIDE", override) + + with pytest.raises( + ValueError, + match="EXPB_SKIP_OVERRIDE must be a nonnegative integer", + ): + create_config(monkeypatch, scenario_inputs)