diff --git a/example-expb.yaml b/example-expb.yaml index d813b24..3d100db 100644 --- a/example-expb.yaml +++ b/example-expb.yaml @@ -86,6 +86,8 @@ scenarios: network: mainnet # Optional: Snapshot backend to use. Available snapshot backends: overlay, zfs, copy. Defaults to overlay snapshot_backend: overlay + # Optional: Snapshot mount target. Use /execution-data/geth when the source contains the contents of Geth's datadir/geth. + # snapshot_mount_path: /execution-data/geth # Optional: Snapshot path for copy backend. If defined, this path will be used instead of work_dir/snapshot. Only used when snapshot_backend is copy snapshot_path: ./work_dir/snapshot # Optional: Number of times to repeat the scenario. Defaults to 1. diff --git a/src/expb/configs/scenarios.py b/src/expb/configs/scenarios.py index 83f1fb1..9e55866 100644 --- a/src/expb/configs/scenarios.py +++ b/src/expb/configs/scenarios.py @@ -1,4 +1,5 @@ from pathlib import Path +from typing import Literal from pydantic import ( BaseModel, @@ -139,6 +140,10 @@ class Scenario(BaseModel): description="Snapshot backend to use.", default=SnapshotBackend.OVERLAY, ) + snapshot_mount_path: Literal["/execution-data", "/execution-data/geth"] = Field( + description="Path where the snapshot is mounted in the execution client container.", + default="/execution-data", + ) snapshot_path: Path | None = Field( description="Path to the snapshot directory for copy backend (overrides work_dir).", default=None, diff --git a/src/expb/payloads/executor/executor_config.py b/src/expb/payloads/executor/executor_config.py index 43b6e30..72bd3bd 100644 --- a/src/expb/payloads/executor/executor_config.py +++ b/src/expb/payloads/executor/executor_config.py @@ -14,7 +14,6 @@ CLIENT_METRICS_PORT, CLIENT_RPC_PORT, CLIENT_RPC_WS_PORT, - CLIENTS_DATA_DIR, CLIENTS_JWT_SECRET_DIR, Client, ) @@ -124,6 +123,7 @@ def __init__( ## Snapshot config self.snapshot_source: str = scenario.snapshot_source + self.snapshot_mount_path: str = scenario.snapshot_mount_path self.snapshot_service: SnapshotService = snapshot_service ## Outputs directory @@ -315,7 +315,7 @@ def get_execution_client_volumes(self) -> list[dict[str, Any]]: ) execution_container_volumes.append( { - "bind": CLIENTS_DATA_DIR, + "bind": self.snapshot_mount_path, "config": { "name": f"{container_name}-overlay-merged", "driver": "local", diff --git a/src/expb/payloads/executor/services/k6.py b/src/expb/payloads/executor/services/k6.py index 1036de8..ed3ffe0 100644 --- a/src/expb/payloads/executor/services/k6.py +++ b/src/expb/payloads/executor/services/k6.py @@ -52,6 +52,8 @@ def build_k6_script_config( "thresholds": { "http_req_failed{kind:newPayload}": ["rate < 0.01"], "http_req_failed{kind:forkchoiceUpdated}": ["rate < 0.01"], + "checks{kind:newPayload}": ["rate == 1"], + "checks{kind:forkchoiceUpdated}": ["rate == 1"], }, "systemTags": [ "scenario", diff --git a/src/expb/payloads/executor/services/templates/k6-script.js.j2 b/src/expb/payloads/executor/services/templates/k6-script.js.j2 index c4c13fc..37d9e2d 100644 --- a/src/expb/payloads/executor/services/templates/k6-script.js.j2 +++ b/src/expb/payloads/executor/services/templates/k6-script.js.j2 @@ -15,10 +15,10 @@ const abortOnEOF = (__ENV.EXPB_ABORT_ON_EOF || '0') === '1'; const addCorrelationHeader = (__ENV.EXPB_ADD_CID || '1') === '1'; const engineEndpoint = __ENV.EXPB_ENGINE_ENDPOINT; const perPayloadMetrics = (__ENV.EXPB_PER_PAYLOAD_METRICS || '0') === '1'; -const discardResponses = (__ENV.EXPB_DISCARD_RESPONSES || '1') === '1'; const enableLogging = (__ENV.EXPB_ENABLE_LOGGING || '0') === '1'; const perPayloadMetricsLogs = (__ENV.EXPB_PER_PAYLOAD_METRICS_LOGS || '0') === '1'; const verbosePostLogs = (__ENV.EXPB_VERBOSE_POST_LOGS || '0') === '1'; +const maxDiagnosticBodyLength = 512; // Load k6 options JSON const config = JSON.parse(open(__ENV.EXPB_CONFIG_FILE_PATH)); @@ -76,6 +76,62 @@ function emitPerPayloadMetricsRow(idx, gasUsed, processingMs) { ); } +function boundedResponseBody(response) { + const body = typeof response.body === 'string' ? response.body : ''; + const compactBody = body.replace(/\s+/g, ' '); + if (compactBody.length <= maxDiagnosticBodyLength) return compactBody; + return compactBody.slice(0, maxDiagnosticBodyLength) + '...'; +} + +function checkEngineResponse(response, kind, idx, tags) { + let payloadStatus = null; + let valid = response.status === 200; + let diagnostic = response.status === 200 ? '' : `HTTP status ${response.status}`; + + if (valid) { + let body; + try { + body = JSON.parse(response.body || ''); + } catch (error) { + valid = false; + diagnostic = 'response body is not valid JSON'; + } + + if (valid && (!body || body.error !== undefined)) { + valid = false; + diagnostic = 'JSON-RPC error'; + } + + if (valid && (!body || body.result === undefined)) { + valid = false; + diagnostic = 'response has no JSON-RPC result'; + } + + if (valid) { + // engine_newPayload returns the payloadStatus object directly, while + // engine_forkchoiceUpdated wraps it in result.payloadStatus. + payloadStatus = kind === 'newPayload' + ? body.result && body.result.status + : body.result && body.result.payloadStatus && body.result.payloadStatus.status; + if (payloadStatus !== 'VALID') { + valid = false; + diagnostic = `payload status ${payloadStatus || 'missing'}`; + } + } + } + + if (!valid) { + console.error( + `[engine] ${kind} idx=${idx} ${diagnostic}; response=${boundedResponseBody(response)}` + ); + } + + check(response, { + 'status_200': (x) => x.status === 200, + 'payload_status_valid': () => valid, + }, tags); +} + // --- Fetch next pair from payload server --- // Line format: {metadata_json}\t{raw_NP}\t{raw_FCU} // Metadata is pre-extracted by the executor, no regex needed here. @@ -170,9 +226,9 @@ function processNextPayload(base_tags = {}, warmup = false) { const r = http.post(engineEndpoint, pair.rawPayload, { headers: headers, tags: tags, - responseType: discardResponses ? 'none' : 'text', + responseType: 'text', }); - check(r, { 'status_200': (x) => x.status === 200 }, tags); + checkEngineResponse(r, 'newPayload', idx, tags); if (!warmup) { emitPerPayloadMetricsRow(idx, pair.gasUsed, r.timings.waiting); @@ -213,9 +269,9 @@ function processNextPayload(base_tags = {}, warmup = false) { const r = http.post(engineEndpoint, pair.rawFcu, { headers: headers, tags: tags, - responseType: discardResponses ? 'none' : 'text', + responseType: 'text', }); - check(r, { 'status_200': (x) => x.status === 200 }, tags); + checkEngineResponse(r, 'forkchoiceUpdated', idx, tags); }); // Emit client-side processing metric for the PREVIOUS block. diff --git a/tests/test_executor_config.py b/tests/test_executor_config.py new file mode 100644 index 0000000..b84639e --- /dev/null +++ b/tests/test_executor_config.py @@ -0,0 +1,89 @@ +from pathlib import Path + +import pytest + +from expb.clients import Client +from expb.configs.scenarios import Scenario, ScenariosPaths +from expb.payloads.executor import executor_config as executor_config_module +from expb.payloads.executor.executor_config import ExecutorConfig +from expb.payloads.executor.services.snapshots import SnapshotService + + +class StubSnapshotService(SnapshotService): + def __init__(self, snapshot_path: Path): + self.snapshot_path = snapshot_path + + def get_snapshot(self, name: str, source: str) -> Path: + return self.snapshot_path + + +def make_scenario(tmp_path: Path, **overrides) -> Scenario: + payloads = tmp_path / "payloads.jsonl" + fcus = tmp_path / "fcus.jsonl" + payloads.touch() + fcus.touch() + values = { + "client": Client.GETH, + "payloads": payloads, + "fcus": fcus, + "snapshot_source": "snapshot", + } + values.update(overrides) + return Scenario.model_validate(values) + + +@pytest.mark.parametrize( + "mount_path", + ["/execution-data", "/execution-data/geth"], +) +def test_snapshot_mount_path_allows_supported_client_layouts( + tmp_path: Path, mount_path: str +) -> None: + scenario = make_scenario(tmp_path, snapshot_mount_path=mount_path) + + assert scenario.snapshot_mount_path == mount_path + + +def test_snapshot_mount_path_defaults_to_client_data_directory(tmp_path: Path) -> None: + scenario = make_scenario(tmp_path) + + assert scenario.snapshot_mount_path == "/execution-data" + + +def test_snapshot_mount_path_rejects_arbitrary_container_paths(tmp_path: Path) -> None: + with pytest.raises(ValueError): + make_scenario(tmp_path, snapshot_mount_path="/tmp/snapshot") + + +def test_executor_mounts_snapshot_at_configured_path( + tmp_path: Path, monkeypatch: pytest.MonkeyPatch +) -> None: + monkeypatch.setattr(executor_config_module.docker, "from_env", lambda: object()) + monkeypatch.setattr(executor_config_module.os, "getuid", lambda: 0, raising=False) + monkeypatch.setattr(executor_config_module.os, "getgid", lambda: 0, raising=False) + + snapshot_path = tmp_path / "snapshot" + snapshot_path.mkdir() + scenario = make_scenario(tmp_path, name="geth-fusaka", snapshot_mount_path="/execution-data/geth") + config = ExecutorConfig( + scenario=scenario, + snapshot_service=StubSnapshotService(snapshot_path), + paths=ScenariosPaths( + work=tmp_path / "work", + outputs=tmp_path / "outputs", + ), + ) + + volumes = config.get_execution_client_volumes() + snapshot_volume = next( + volume + for volume in volumes + if volume["config"]["name"].endswith("-overlay-merged") + ) + + assert snapshot_volume["bind"] == "/execution-data/geth" + assert snapshot_volume["config"]["driver_opts"]["o"].startswith("bind,rw") + assert snapshot_volume["config"]["driver_opts"]["device"] == str( + snapshot_path.resolve() + ) + assert all(volume["bind"] != "/execution-data" for volume in volumes) diff --git a/tests/test_k6.py b/tests/test_k6.py new file mode 100644 index 0000000..36e3e7d --- /dev/null +++ b/tests/test_k6.py @@ -0,0 +1,95 @@ +import json +import re +import shutil +import subprocess + +import pytest + +from expb.clients import Client +from expb.payloads.executor.services.k6 import ( + build_k6_script_config, + get_k6_script_content, +) + + +def test_k6_thresholds_fail_non_valid_engine_responses() -> None: + config = build_k6_script_config( + test_id="test", + scenario_name="scenario", + client=Client.GETH, + iterations=1, + ) + + thresholds = config["options"]["thresholds"] + assert thresholds["checks{kind:newPayload}"] == ["rate == 1"] + assert thresholds["checks{kind:forkchoiceUpdated}"] == ["rate == 1"] + + +def test_k6_engine_response_validation_covers_engine_statuses() -> None: + node = shutil.which("node") + if node is None: + pytest.skip("node is required to execute the K6 response validation") + + script = get_k6_script_content() + script = re.sub(r"^import .*;\r?\n", "", script, flags=re.MULTILINE) + script = script.replace("export const options", "const options") + script = script.replace("export function setup", "function setup") + script = script.replace("export default function ()", "function defaultFunction()") + script += "\nglobalThis.checkEngineResponse = checkEngineResponse;\n" + + driver = f""" +const vm = require('node:vm'); +const checks = []; +const errors = []; +const context = {{ + __ENV: {{ + EXPB_CONFIG_FILE_PATH: '/config.json', + EXPB_JWTSECRET_FILE_PATH: '/jwt.hex', + }}, + open: (path) => path === '/config.json' ? '{{"options":{{}}}}' : '00'.repeat(32), + Gauge: function Gauge() {{ this.add = () => {{}}; }}, + encoding: {{ b64encode: () => '' }}, + check: (response, checkFunctions, tags) => {{ + const result = {{}}; + for (const [name, checkFunction] of Object.entries(checkFunctions)) {{ + result[name] = Boolean(checkFunction(response)); + }} + checks.push({{ result, tags }}); + return Object.values(result).every(Boolean); + }}, + console: {{ error: (message) => errors.push(message) }}, +}}; +vm.runInNewContext({json.dumps(script)}, context); +const cases = [ + ['newPayload', {{ status: 200, body: '{{"result":{{"status":"VALID"}}}}' }}, true], + ['forkchoiceUpdated', {{ status: 200, body: '{{"result":{{"payloadStatus":{{"status":"VALID"}}}}}}' }}, true], + ['newPayload', {{ status: 200, body: '{{"result":{{"status":"INVALID"}}}}' }}, false], + ['newPayload', {{ status: 200, body: '{{"result":{{"status":"SYNCING"}}}}' }}, false], + ['newPayload', {{ status: 200, body: '{{"result":{{"status":"ACCEPTED"}}}}' }}, false], + ['forkchoiceUpdated', {{ status: 200, body: '{{"result":{{"payloadStatus":{{"status":"INVALID"}}}}}}' }}, false], + ['forkchoiceUpdated', {{ status: 200, body: '{{"result":{{"payloadStatus":{{"status":"SYNCING"}}}}}}' }}, false], + ['forkchoiceUpdated', {{ status: 200, body: '{{"result":{{"payloadStatus":{{"status":"ACCEPTED"}}}}}}' }}, false], + ['newPayload', {{ status: 200, body: '{{"error":{{"code":-32000,"message":"failed"}}}}' }}, false], + ['forkchoiceUpdated', {{ status: 200, body: '{{' }}, false], + ['newPayload', {{ status: 503, body: 'temporarily unavailable' }}, false], + ['forkchoiceUpdated', {{ status: 200, body: 'x'.repeat(2000) }}, false], +]; +for (const [kind, response, expected] of cases) {{ + context.checkEngineResponse(response, kind, 42, {{ kind }}); + if (checks.at(-1).result.payload_status_valid !== expected) {{ + throw new Error(`unexpected result for ${{kind}}: ${{JSON.stringify(response)}}`); + }} +}} +if (errors.length !== cases.length - 2) throw new Error(`unexpected diagnostic count: ${{errors.length}}`); +if (errors.some((message) => message.length > 700)) throw new Error('diagnostic was not bounded'); +if (!errors.some((message) => message.includes('...'))) throw new Error('long diagnostic was not truncated'); +process.stdout.write(JSON.stringify({{ checks, errors }})); +""" + result = subprocess.run( + [node, "-e", driver], + capture_output=True, + text=True, + ) + + assert result.returncode == 0, result.stderr + assert len(json.loads(result.stdout)["checks"]) == 12