Skip to content
Merged
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
6 changes: 6 additions & 0 deletions docs/USAGE.md
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down
26 changes: 14 additions & 12 deletions src/expb/payloads/executor/executor.py
Original file line number Diff line number Diff line change
Expand Up @@ -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()
Expand Down Expand Up @@ -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(
Expand Down
11 changes: 11 additions & 0 deletions src/expb/payloads/executor/executor_config.py
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down
63 changes: 63 additions & 0 deletions tests/payloads/test_executor_cleanup.py
Original file line number Diff line number Diff line change
@@ -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")])
90 changes: 90 additions & 0 deletions tests/test_payload_skip_override.py
Original file line number Diff line number Diff line change
@@ -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)
Loading